Something went wrong. Try again.
Rust CLI for tangled
Something went wrong. Try again.
5.2 kB · 148 lines
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149//! 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<String>, /// XRPC error message, if any. message: Option<String>, /// 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 `<base>/xrpc/<nsid>`, 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<T: serde::Serialize>( 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<T: DeserializeOwned>(rb: reqwest::RequestBuilder) -> Result<T, XrpcError> { 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<Vec<u8>, 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<T: DeserializeOwned>(res: reqwest::Response) -> Result<T, XrpcError> { 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<Vec<u8>, 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<String>, message: Option<String>, } let parsed: Option<ErrorBody> = 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()}