//! Transport helpers shared by the generated lexicon builder functions. //! //! The generated functions in [`crate::lexicon`] return plain //! [`reqwest::RequestBuilder`]s with the URL, HTTP method, and body already //! set. Callers attach whatever headers/auth they need (e.g. //! `.bearer_auth(token)`) and then either drive the request themselves or use //! [`send`] / [`send_unit`] / [`send_bytes`] for XRPC-aware status handling //! and typed deserialization. use serde::de::DeserializeOwned; use serde::Deserialize; #[derive(Debug, thiserror::Error)] pub enum XrpcError { #[error(transparent)] Transport(#[from] reqwest::Error), #[error("{status}{}{}", error.as_deref().map(|e| format!(" {e}")).unwrap_or_default(), message.as_deref().map(|m| format!(": {m}")).unwrap_or_default())] Status { status: reqwest::StatusCode, /// XRPC error code (the `error` field of the error body), if any. error: Option, /// XRPC error message, if any. message: Option, /// Raw response body. body: String, }, #[error("failed to decode response from {url}: {source}; body: {snippet}")] Decode { url: String, snippet: String, #[source] source: serde_json::Error, }, } /// Builds `/xrpc/`, defaulting to `https://` when `base` has no /// scheme. pub fn xrpc_url(base: &str, nsid: &str) -> String { let base = base.trim_end_matches('/'); if base.starts_with("http://") || base.starts_with("https://") { format!("{base}/xrpc/{nsid}") } else { format!("https://{base}/xrpc/{nsid}") } } /// Appends `params` to the request's query string. /// /// `params` is serialized through `serde_json`; `null`s are skipped and /// arrays repeat the key (`?tag=a&tag=b`), matching XRPC conventions. pub fn push_query( mut rb: reqwest::RequestBuilder, params: &T, ) -> reqwest::RequestBuilder { let value = match serde_json::to_value(params) { Ok(serde_json::Value::Object(map)) => map, // A params struct always serializes to an object; anything else is a // programming error surfaced at request time by reqwest. _ => return rb, }; for (key, value) in value { match value { serde_json::Value::Null => {} serde_json::Value::Array(items) => { for item in items { rb = rb.query(&[(key.as_str(), scalar_to_string(&item))]); } } other => { rb = rb.query(&[(key.as_str(), scalar_to_string(&other))]); } } } rb } fn scalar_to_string(v: &serde_json::Value) -> String { match v { serde_json::Value::String(s) => s.clone(), other => other.to_string(), } } /// Sends the request and deserializes a JSON response body into `T`. pub async fn send(rb: reqwest::RequestBuilder) -> Result { parse(rb.send().await?).await } /// Sends the request, checking status but discarding any response body. pub async fn send_unit(rb: reqwest::RequestBuilder) -> Result<(), XrpcError> { parse_unit(rb.send().await?).await } /// Sends the request and returns the raw response bytes (for binary outputs /// such as `com.atproto.sync.getBlob`). pub async fn send_bytes(rb: reqwest::RequestBuilder) -> Result, XrpcError> { parse_bytes(rb.send().await?).await } /// Checks the status of an already-received response and deserializes its /// JSON body into `T`. Useful when the caller needs to inspect response /// headers (e.g. DPoP nonces) before handing off parsing. pub async fn parse(res: reqwest::Response) -> Result { let (url, body) = success_body(res).await?; serde_json::from_slice(&body).map_err(|source| XrpcError::Decode { url, snippet: snippet(&body), source, }) } /// Checks the status of an already-received response, discarding any body. pub async fn parse_unit(res: reqwest::Response) -> Result<(), XrpcError> { success_body(res).await.map(|_| ()) } /// Checks the status of an already-received response and returns the raw /// body bytes. pub async fn parse_bytes(res: reqwest::Response) -> Result, XrpcError> { success_body(res).await.map(|(_, body)| body.to_vec()) } async fn success_body(res: reqwest::Response) -> Result<(String, bytes::Bytes), XrpcError> { let status = res.status(); let url = res.url().to_string(); let body = res.bytes().await?; if status.is_success() { return Ok((url, body)); } #[derive(Deserialize)] struct ErrorBody { error: Option, message: Option, } let parsed: Option = serde_json::from_slice(&body).ok(); Err(XrpcError::Status { status, error: parsed.as_ref().and_then(|e| e.error.clone()), message: parsed.as_ref().and_then(|e| e.message.clone()), body: String::from_utf8_lossy(&body).into_owned(), }) } fn snippet(body: &[u8]) -> String { let s = String::from_utf8_lossy(body); s.chars().take(300).collect() }