//! The one wrapper around atrium. Everything else in this crate talks to //! ATProto through the `Atproto` struct below; only this module and its //! submodules import `atrium_*` or `jose_jwk` — the crates are 0.x and churn, //! and this is the wall the churn stops at (plan/complete/sign-in.md). mod client; mod store; mod verify; use std::sync::Arc; use atrium_api::agent::{Agent, SessionManager}; use atrium_api::com::atproto::repo::{ create_record, delete_record, get_record, put_record, upload_blob, }; use atrium_api::types::TryIntoUnknown; use atrium_api::types::string::{Datetime, Did, Handle}; use atrium_common::resolver::Resolver; use atrium_oauth::{AuthorizeOptions, CallbackParams}; use atrium_xrpc::{InputDataOrBytes, OutputDataOrBytes, XrpcClient, XrpcRequest}; use crate::config::Config; use crate::db::Db; /// The camo collection. Named here because the OAuth scope, the record's /// `$type` and every repo call have to agree. pub const CAMO_NSID: &str = "blue.lance.camo"; /// A flare's record type: userinput.app's discussion, in the player's repo. pub const DISCUSSION_NSID: &str = "app.userinput.discussion"; /// The space collection a flare's target lives in, and the shape FLARE_SPACE /// must point at. pub const SPACE_NSID: &str = "app.userinput.space"; /// The scope that covers writing a discussion — userinput.app's own /// published permission set. /// /// Asked for at sign-in and never checked against afterwards: the /// authorization server expands an `include:` into the permissions the set /// holds before minting the token, so this string appears in no granted /// scope. See `repo_did`. pub const FLARE_SCOPE: &str = "include:app.userinput.authBasic"; pub enum LoginError { /// The input is not a handle. InvalidHandle, /// The account's own infrastructure could not be reached or refused. Upstream, } /// The exchange failed after consent. Denial never reaches the wrapper: the /// authorization server signals it with `error=` and no code, which the /// route answers by itself. pub struct CallbackFailed; /// A repository write did not happen. #[derive(Debug)] pub enum WriteError { /// There is no usable session to write with: the stored one could not be /// restored, or its refresh token is dead. Only signing in again fixes /// it, and no server was reached to say anything more. NeedsSignIn, /// The account's own server refused the write, in its own words. Refused(Refusal), /// No such camo in that repository. NotFound, /// The record cannot be represented in the atproto data model. In /// practice that means a fractional number: the model has integers and /// no float, and a record carrying one is refused on the way out rather /// than by the PDS. Unstorable, /// The record changed since it was read. The caller read a stale copy and /// has to look again rather than overwrite whatever happened in between. Conflict, /// The PDS refused the write or could not be reached. Upstream, } /// A refusal from the account's own server, in the server's own words. /// /// The PDS is the authority on what a token may write, so what it says is /// relayed rather than guessed at from here. A missing permission comes back /// as 403 `ScopeMissingError` naming the exact scope — `repo:` or /// `blob:` — which is the one thing a paraphrase loses, and which a /// permission set makes impossible to work out on this side anyway: the /// authorization server expands `include:` into member permissions /// before minting the token, so what was asked for and what was granted are /// never the same string. #[derive(Debug)] pub struct Refusal { /// What it answered with. 401 and 403 are the ones fresh consent can fix. pub status: u16, /// The server's `error` name, e.g. `ScopeMissingError`. Sanitised. pub error: Option, /// The server's `message`. Sanitised. pub message: Option, } impl Refusal { /// Whether signing in again could plausibly change the answer. pub fn fixable_by_signing_in(&self) -> bool { matches!(self.status, 401 | 403) } /// The server's words as one line, or None if it gave none worth showing. pub fn detail(&self) -> Option { match (&self.error, &self.message) { // The message repeats the error name often enough that showing // both reads as a stutter; the message is the specific one. (_, Some(message)) => Some(message.clone()), (Some(error), None) => Some(error.clone()), (None, None) => None, } } } /// Another server's text on its way to our users: control characters out, and /// a cap, because nothing downstream should have to size for an essay. fn sanitize_detail(raw: Option<&str>) -> Option { const MAX_CHARS: usize = 200; let text: String = raw? .chars() .filter(|c| !c.is_control()) .take(MAX_CHARS) .collect(); let text = text.trim(); (!text.is_empty()).then(|| text.to_owned()) } /// Where a written camo ended up. pub struct Written { pub uri: String, pub cid: String, } /// A camo as its owner published it. pub struct Published { /// The blob's content address. This is what the record points at, and /// what the launch manifest records the camo as having been published /// under. pub cid: String, /// The bytes the PDS served, unexamined. Untrusted: see `published_camo`. pub png: Vec, } /// One followed account, from the public AppView's follow list. pub struct FollowedProfile { pub did: String, pub handle: String, } /// A public read did not happen. Nothing here distinguishes whose fault it /// was to the player, because the answer is the same either way: the match /// starts without a camo. #[derive(Debug)] pub enum ReadError { /// No such record, or one with no image in it. NotFound, /// The account's server refused or could not be reached. Upstream, } /// A future boxed so it can sit behind a trait object. `Resolver::resolve` /// returns `impl Future`, which has no fixed size and so cannot appear in a /// `dyn Resolver` vtable; boxing the future itself is what makes a resolver /// nameable as one field's type regardless of which concrete resolver it is. type BoxFuture<'a, T> = std::pin::Pin + Send + 'a>>; /// A handle resolver, type-erased. `Atproto` holds one of these rather than /// a concrete `client::HandleResolver` so that `resolve_player`'s two /// production callers (`routes::plan_players`'s `player` control, /// `matches::lobby::Lobbies::invite_handle`) can be tested against a /// scripted answer instead of a real DNS-over-HTTPS round trip — the same /// tension `verify_handle` resolves by taking its resolvers generically /// (`atproto/verify.rs`). A generic *function* argument cannot help here: /// both callers reach resolution through an `Atproto` that is already built /// by the time they see it (an `AppState` field, a `&Atproto` parameter), so /// what needs to be swappable is the field itself, and `Resolver`'s RPITIT /// method rules out `Box>` directly. /// /// Takes the handle by value rather than `&Handle`, unlike `Resolver` /// itself: `Handle` is cheap to clone (one `String`), and taking it owned /// means the boxed future below only has to outlive `&self`, not a second, /// independent borrow of the argument — the two-lifetime shape rustc /// otherwise cannot reconcile with a single boxed trait object ("hidden type /// captures lifetime that does not appear in bounds"). trait DynHandleResolver: Send + Sync { fn resolve(&self, handle: Handle) -> BoxFuture<'_, Result>; } impl DynHandleResolver for client::HandleResolver { fn resolve(&self, handle: Handle) -> BoxFuture<'_, Result> { Box::pin(async move { Resolver::resolve(self, &handle).await.map_err(|_| ()) }) } } /// Bridges the boxed form back into `Resolver`, so `verify_handle` — generic /// over the trait for the identical reason — can still be called with /// `&self.handle_resolver` without knowing the concrete resolver underneath /// was erased. impl<'x> Resolver for Box { type Input = Handle; type Output = Did; type Error = (); async fn resolve(&self, input: &Handle) -> Result { self.as_ref().resolve(input.clone()).await } } /// An account, as far as an address is allowed to state it. /// /// `handle` is None where none resolves back to the DID, which is not an /// error: the DID is the account, and it is always true. pub struct PlayerActor { pub did: String, pub handle: Option, } pub struct Atproto { client: Arc, did_resolver: client::DidResolver, handle_resolver: Box, /// Unauthenticated reads against whichever PDS owns the record. Separate /// from the OAuth client's own: nothing here carries a credential. http: reqwest::Client, confidential: bool, } impl Atproto { pub fn new(config: &Config, db: Db) -> Result { let http = client::HttpClient::default(); Ok(Atproto { client: Arc::new(client::build(config, db, &http)?), // A second pair: the first is owned by the OAuthClient. did_resolver: client::did_resolver(&http), handle_resolver: Box::new(client::handle_resolver(&http)), http: reqwest::Client::new(), confidential: config.public_url.is_some(), }) } /// Starts the flow for a handle: resolves it, pushes the authorization /// request, returns the URL to send the browser to — the consent screen /// on the account's own server. pub async fn begin_login(&self, handle: &str) -> Result { let handle = handle.trim(); if Handle::new(handle.to_owned()).is_err() { return Err(LoginError::InvalidHandle); } // Handles are public identifiers; logging them is the audit trail of // who tries to sign in. Tokens and codes stay out of logs as ever. tracing::info!(handle, "login: attempt"); // Not AuthorizeOptions::default(): its scopes are bare `atproto`, and // the pushed authorization request asks for these, not the ones in // the client metadata. let options = AuthorizeOptions { scopes: client::scopes(), ..AuthorizeOptions::default() }; self.client.authorize(handle, options).await.map_err(|_| { // atrium's error values can embed request URIs and server // responses; log the fact, not the value. tracing::warn!(handle, "login: authorization could not be started"); LoginError::Upstream }) } /// Completes the flow from the callback query: code for tokens, tokens /// into the store, handle verified both directions. Returns the DID and /// the verified handle. pub async fn complete_login( &self, code: String, state: Option, iss: Option, ) -> Result<(String, Option), CallbackFailed> { let client = Arc::clone(&self.client); let params = CallbackParams { code, state, iss }; // atrium 0.1.7 hits a todo!() when the token exchange itself fails — // a panic, not an Err. Run the callback in its own task so that // surfaces as a JoinError and the user gets the failure redirect // instead of a dead connection. let session = match tokio::spawn(async move { client.callback(params).await }).await { Ok(Ok((session, _app_state))) => session, Ok(Err(_)) => { tracing::warn!("login: callback failed"); return Err(CallbackFailed); } Err(_) => { tracing::warn!("login: token exchange failed"); return Err(CallbackFailed); } }; let Some(did) = session.did().await else { tracing::warn!("login: session has no DID"); return Err(CallbackFailed); }; let handle = verify::verify_handle(&self.did_resolver, &self.handle_resolver, &did).await; tracing::info!( did = did.as_str(), handle = handle.as_deref().unwrap_or("(unverified)"), "login: completed" ); Ok((did.as_str().to_owned(), handle)) } /// Verifies the handle for an already signed-in DID, the same way /// sign-in did. Two network round trips, so callers run it off the /// request path (routes.rs). pub async fn reverify_handle(&self, did: &str) -> Option { let did = Did::new(did.to_owned()).ok()?; verify::verify_handle(&self.did_resolver, &self.handle_resolver, &did).await } /// Resolves a typed handle to the account it names, for putting another /// player in a match. Returns the DID and the canonical handle. /// /// Forward resolution only — DNS TXT or well-known, the mechanism the /// domain owner controls — which is the authoritative direction when the /// question starts from the name. The bidirectional check exists for the /// opposite question (login hands us a DID and the handle is a claim); /// here a claim in a DID document is never consulted. pub async fn resolve_player(&self, handle: &str) -> Option<(String, String)> { let trimmed = handle.trim().trim_start_matches('@'); let handle = Handle::new(trimmed.to_ascii_lowercase()).ok()?; match self.handle_resolver.resolve(&handle).await { Ok(did) => Some((did.as_str().to_owned(), handle.as_str().to_owned())), Err(_) => { // Handles are public identifiers, and which name failed is // the one thing a challenge that will not start needs logged. tracing::info!( handle = handle.as_str(), "challenge: handle did not resolve" ); None } } } /// Who a `/profile/<...>` address names, resolved from either form. /// /// The web app does exactly this in `web/src/identity.ts`, and it has to /// be done here as well rather than shared: the page is drawn in the /// browser and the unfurl is answered before any browser is involved. /// Both follow the same rule, which is that the DID document is the /// authority in both directions - a resolver's answer is confirmed /// against the document it points at, and a handle a document claims is /// resolved back to it. A handle that cannot be confirmed is dropped /// rather than shown: a handle is rented, and the account that used to /// hold one is not the account that holds it now. /// /// None means nobody. A DID with no document is nobody as much as a /// handle nothing resolves is. pub async fn resolve_actor(&self, actor: &str) -> Option { let trimmed = actor.trim().trim_start_matches('@'); if trimmed.is_empty() { return None; } if trimmed.starts_with("did:") { let did = Did::new(trimmed.to_owned()).ok()?; // The document is what says this DID exists at all, and // verify_handle reads it anyway - but it swallows the difference // between "no document" and "no handle", and those are a 404 and // a page. self.did_resolver.resolve(&did).await.ok()?; let handle = verify::verify_handle(&self.did_resolver, &self.handle_resolver, &did).await; return Some(PlayerActor { did: did.as_str().to_owned(), handle, }); } let (did, handle) = self.resolve_player(trimmed).await?; let confirmed = { let did = Did::new(did.clone()).ok()?; verify::verify_handle(&self.did_resolver, &self.handle_resolver, &did).await }; // The handle the document claims, not the one that was asked for: // they are the same string in the ordinary case, and where they are // not, the document is the one to believe. Some(PlayerActor { did, handle: confirmed.filter(|claimed| *claimed == handle), }) } /// How many camo an account has published, or None where its repository /// could not be read. /// /// The same public listing the camo picker does from the browser, done /// here because a card is drawn where no browser is. One page: a figure /// on a card does not need to be exact past a hundred, and a repository /// that answers slowly must not hold up an unfurl. pub async fn camo_count(&self, did: &str) -> Option { let did = Did::new(did.to_owned()).ok()?; let document = self.did_resolver.resolve(&did).await.ok()?; let pds = document.get_pds_endpoint()?; let listing: serde_json::Value = self .fetch_json(&format!( "{pds}/xrpc/com.atproto.repo.listRecords\ ?repo={did}&collection={CAMO_NSID}&limit=100", did = did.as_str() )) .await .ok()?; Some(listing.get("records")?.as_array()?.len()) } /// The account's own picture, as it uploaded it. /// /// From the PDS rather than from the appview's CDN: the appview only /// knows accounts it has indexed, and this has to answer for an account /// that has never posted. The bytes are whatever was uploaded, so the /// card decodes both formats it can be (see card.rs). /// /// None where there is no picture, where the repository could not be /// read, or where the blob is larger than a portrait is worth: an avatar /// is a small square by the time anybody sees it, and a card must not /// wait on somebody's ten-megabyte upload. pub async fn avatar(&self, did: &str) -> Option> { const MOST: usize = 4 * 1024 * 1024; let did = Did::new(did.to_owned()).ok()?; let document = self.did_resolver.resolve(&did).await.ok()?; let pds = document.get_pds_endpoint()?; let record: serde_json::Value = self .fetch_json(&format!( "{pds}/xrpc/com.atproto.repo.getRecord\ ?repo={did}&collection=app.bsky.actor.profile&rkey=self", did = did.as_str() )) .await .ok()?; // Both spellings, for the reason published_camo reads both: `ref.$link` // is the current blob shape and a bare `cid` is what the format looked // like before it settled, and both are in live repositories. let cid = record .pointer("/value/avatar/ref/$link") .or_else(|| record.pointer("/value/avatar/cid")) .and_then(serde_json::Value::as_str)?; if !is_blob_cid(cid) { tracing::warn!(did = did.as_str(), "avatar: reference is not a CID"); return None; } let response = self .get(&format!( "{pds}/xrpc/com.atproto.sync.getBlob?did={did}&cid={cid}", did = did.as_str() )) .await .ok()?; if !response.status().is_success() { return None; } // Refused on the declared length where there is one, and on the body // otherwise: a length header is a claim, and the check that matters // is on what actually arrived. if response.content_length().is_some_and(|n| n as usize > MOST) { tracing::info!(did = did.as_str(), "avatar: larger than a portrait needs"); return None; } let bytes = response.bytes().await.ok()?; if bytes.len() > MOST { return None; } Some(bytes.to_vec()) } /// Best-effort server-side revocation. The stored session row is deleted /// by atrium on success; the caller deletes it unconditionally as /// belt-and-braces either way. pub async fn logout(&self, did: &str) { let Ok(did) = Did::new(did.to_owned()) else { return; }; tracing::info!(did = did.as_str(), "logout"); if self.client.revoke(&did).await.is_err() { tracing::info!("logout: server-side revocation failed; session row deleted anyway"); } } /// Writes a new camo to the player's own repository. /// /// The blob and the record are one operation on purpose: a blob no record /// refers to is not downloadable and the PDS collects it within hours, so /// the record write is the commit point and a failure before it leaves /// nothing behind that lasts. pub async fn create_camo( &self, did: &str, name: &str, png: Vec, source: Option, ) -> Result { let did = self.repo_did(did)?; if source.as_ref().is_some_and(has_fraction) { tracing::warn!("camo: settings carry a fractional number and cannot be stored"); return Err(WriteError::Unstorable); } let session = self.session(&did).await?; let blob = upload_png(&session, png).await?; let record = camo_record(name, blob_json(blob)?, Datetime::now().as_str(), source); let input = create_record::InputData { collection: camo_nsid(), record: to_unknown(record)?, repo: did.clone().into(), // Unset: the PDS mints the TID the lexicon's key type asks for, // and it is the only party that can order one correctly. rkey: None, swap_commit: None, validate: None, }; let agent = Agent::new(session); match agent.api.com.atproto.repo.create_record(input.into()).await { Ok(output) => { tracing::info!(did = did.as_str(), "camo: record written"); Ok(Written { uri: output.data.uri.clone(), cid: output.data.cid.as_ref().to_string(), }) } Err(e) => Err(write_failed("create record", &e)), } } /// Changes a camo that already exists. /// /// The record put back is the record that was read, with the named fields /// replaced. Rebuilding it from the fields this service knows about would /// delete the editor's settings: they are not in the lexicon, this code /// has no reason to understand them, and reading, mutating and writing /// the whole value is what keeps them — including fields written by a /// version of the site newer than this one. /// /// `createdAt` is never touched. It says when the camo was made, not when /// it was last edited. pub async fn update_camo( &self, did: &str, rkey: &str, name: Option<&str>, png: Option>, source: Option>, swap_record: Option<&str>, ) -> Result { let did = self.repo_did(did)?; if source .as_ref() .and_then(|s| s.as_ref()) .is_some_and(has_fraction) { tracing::warn!("camo: settings carry a fractional number and cannot be stored"); return Err(WriteError::Unstorable); } let session = self.session(&did).await?; // Before the agent, which consumes the session. An upload for a write // that then fails leaves an unreferenced blob, which the PDS collects // on its own within hours. let uploaded = match png { Some(png) => Some(upload_png(&session, png).await?), None => None, }; let rkey = record_key(rkey)?; let agent = Agent::new(session); let current = agent .api .com .atproto .repo .get_record( get_record::ParametersData { cid: None, collection: camo_nsid(), repo: did.clone().into(), rkey: rkey.clone(), } .into(), ) .await .map_err(|e| read_failed("get record", e))?; let Ok(serde_json::Value::Object(mut record)) = serde_json::to_value(¤t.data.value) else { tracing::warn!("camo: stored record is not an object"); return Err(WriteError::Upstream); }; if let Some(name) = name { record.insert("name".into(), name.into()); } if let Some(blob) = uploaded { record.insert("image".into(), blob_json(blob)?); } match source { Some(Some(source)) => { record.insert("source".into(), source); } Some(None) => { record.remove("source"); } None => {} } let input = put_record::InputData { collection: camo_nsid(), record: to_unknown(record)?, repo: did.clone().into(), rkey, swap_commit: None, // The caller read a copy. Refuse if it is no longer the current // one, rather than overwriting whatever another tab did. swap_record: match swap_record { Some(cid) => Some(cid.parse().map_err(|_| WriteError::Conflict)?), None => None, }, validate: None, }; match agent.api.com.atproto.repo.put_record(input.into()).await { Ok(output) => { tracing::info!(did = did.as_str(), "camo: record changed"); Ok(Written { uri: output.data.uri.clone(), cid: output.data.cid.as_ref().to_string(), }) } Err(e) => Err(match &e { atrium_xrpc::Error::XrpcResponse(r) if r.status.as_u16() == 400 => { // InvalidSwap is the only 400 a well-formed put can get // here, and it means the record moved under the caller. tracing::info!("camo: put refused, the record had changed"); WriteError::Conflict } _ => write_failed("put record", &e), }), } } /// Removes a camo. pub async fn delete_camo(&self, did: &str, rkey: &str) -> Result<(), WriteError> { let did = self.repo_did(did)?; let session = self.session(&did).await?; let agent = Agent::new(session); let input = delete_record::InputData { collection: camo_nsid(), repo: did.clone().into(), rkey: record_key(rkey)?, swap_commit: None, swap_record: None, }; match agent.api.com.atproto.repo.delete_record(input.into()).await { Ok(_) => { tracing::info!(did = did.as_str(), "camo: record deleted"); Ok(()) } Err(e) => Err(write_failed("delete record", &e)), } } /// A published camo, read from the account that owns it. /// /// No session and no token. A record and its blob are public, which is /// the same reason the browser lists a collection without asking us; here /// it also means a match can be painted long after the OAuth session /// behind it has expired, and that a camo belonging to someone who is not /// the player launching could be fetched the same way when matches have /// two humans in them. /// /// The bytes come back exactly as the PDS served them and are untrusted. /// A player owns their repository and can put anything in this collection /// with any client, so re-encoding on the way *in* was never a boundary — /// the boundary is where these bytes are handed to other players, which /// is the caller of this. pub async fn published_camo(&self, did: &str, rkey: &str) -> Result { let Ok(did) = Did::new(did.to_owned()) else { tracing::error!("camo: session DID is not a DID"); return Err(ReadError::Upstream); }; record_key(rkey).map_err(|_| ReadError::NotFound)?; let document = self.did_resolver.resolve(&did).await.map_err(|_| { tracing::warn!(did = did.as_str(), "camo: DID document could not be read"); ReadError::Upstream })?; let Some(pds) = document.get_pds_endpoint() else { tracing::warn!(did = did.as_str(), "camo: DID document names no PDS"); return Err(ReadError::Upstream); }; let record: serde_json::Value = self .fetch_json(&format!( "{pds}/xrpc/com.atproto.repo.getRecord\ ?repo={did}&collection={CAMO_NSID}&rkey={rkey}", did = did.as_str() )) .await?; // The blob reference, not the record's own CID: the record's says // which version of the description this is, and the image is a // separate object with an address of its own. // // Both spellings, because both are in live repositories. `ref.$link` // is the current blob shape; a bare `cid` is what the format looked // like before it settled, and `web/src/camo/repo.ts` reads both for // the same reason. A camo this could not find would still list and // still be selectable in the picker, and then silently never paint. let Some(cid) = record .pointer("/value/image/ref/$link") .or_else(|| record.pointer("/value/image/cid")) .and_then(serde_json::Value::as_str) else { tracing::warn!("camo: record carries no image blob reference"); return Err(ReadError::NotFound); }; // The one value in this whole path that a hostile server picks freely. // It goes into a query string and then into the launch manifest, so it // is checked for shape before either. if !is_blob_cid(cid) { tracing::warn!("camo: record's image reference is not a CID"); return Err(ReadError::NotFound); } // getBlob takes the DID as well as the CID. A blob reference on its // own does not say which repository it is in, so a record copied // between accounts would point at bytes only the original can serve. let response = self .get(&format!( "{pds}/xrpc/com.atproto.sync.getBlob?did={did}&cid={cid}", did = did.as_str() )) .await .map_err(|_| upstream("blob fetch"))?; if !response.status().is_success() { return Err(status_error("blob fetch", response.status())); } let png = read_capped(response, BLOB_BYTE_LIMIT, "blob body").await?; tracing::info!( did = did.as_str(), rkey, bytes = png.len(), "camo: read from the owner's repository" ); Ok(Published { cid: cid.to_owned(), png, }) } /// Accounts `did` follows, from the public AppView — a follow graph is /// the AppView's own index, not a fact any one repository holds, so /// unlike published_camo there is no PDS to resolve first. /// /// Paged and capped at two pages of 100: a follow list has no upper /// bound, and this only ever feeds a best-effort suggestion, not a /// completeness guarantee. No OAuth scope involved — the same class of /// public, unauthenticated read as published_camo's blob fetch. pub async fn followed(&self, did: &str) -> Result, ReadError> { const PAGE_LIMIT: u32 = 100; const MAX_PAGES: usize = 2; let mut out = Vec::new(); let mut cursor: Option = None; for _ in 0..MAX_PAGES { let mut url = reqwest::Url::parse("https://public.api.bsky.app/xrpc/app.bsky.graph.getFollows") .expect("constant URL parses"); url.query_pairs_mut() .append_pair("actor", did) .append_pair("limit", &PAGE_LIMIT.to_string()); if let Some(cursor) = &cursor { url.query_pairs_mut().append_pair("cursor", cursor); } let page = self.fetch_json(url.as_str()).await?; let Some(follows) = page.get("follows").and_then(|v| v.as_array()) else { break; }; for entry in follows { let (Some(did), Some(handle)) = ( entry.get("did").and_then(|v| v.as_str()), entry.get("handle").and_then(|v| v.as_str()), ) else { continue; }; out.push(FollowedProfile { did: did.to_owned(), handle: handle.to_owned(), }); } cursor = page .get("cursor") .and_then(|v| v.as_str()) .map(str::to_owned); if cursor.is_none() { break; } } Ok(out) } /// Files a flare: an `app.userinput.discussion` in the player's own /// repository, pointing at the board's space. /// /// The space's strongRef is fetched fresh from its owner's PDS on every /// flare — the ref must carry the space record's current CID, and caching /// one would break every flare the moment the board is renamed. pub async fn create_flare( &self, did: &str, title: &str, body: &str, tags: &[String], png: Option>, space_uri: &str, ) -> Result { let did = self.repo_did(did)?; let space = self.space_ref(space_uri).await?; let session = self.session(&did).await?; // Before the agent consumes the session; an unreferenced blob from a // write that then fails is collected by the PDS within hours. let image = match png { Some(png) => Some(upload_png(&session, png).await?), None => None, }; let mut record = serde_json::Map::new(); record.insert("$type".into(), DISCUSSION_NSID.into()); record.insert("space".into(), space); record.insert("title".into(), title.into()); if !body.is_empty() { record.insert("body".into(), body.into()); } // What the player ticked. The lexicon says these are drawn from the // space's tag list, and both forms offer that list; a value from // somewhere else lands in the sender's own repo and the board simply // does not filter on it. if !tags.is_empty() { record.insert("tags".into(), tags.into()); } if let Some(image) = image { record.insert( "images".into(), serde_json::json!([{ "image": blob_json(image)?, "alt": "" }]), ); } record.insert("createdAt".into(), Datetime::now().as_str().into()); let input = create_record::InputData { collection: discussion_nsid(), record: to_unknown(record)?, repo: did.clone().into(), rkey: None, swap_commit: None, validate: None, }; let agent = Agent::new(session); match agent.api.com.atproto.repo.create_record(input.into()).await { Ok(output) => { tracing::info!(did = did.as_str(), "flare: discussion written"); Ok(Written { uri: output.data.uri.clone(), cid: output.data.cid.as_ref().to_string(), }) } Err(e) => Err(write_failed("create discussion", &e)), } } /// The board space's strongRef `{uri, cid}`, read from its owner's PDS. async fn space_ref(&self, space_uri: &str) -> Result { let (space_did, rkey) = split_space_uri(space_uri).ok_or_else(|| { // FLARE_SPACE was validated at startup; reaching this is a bug. tracing::error!("flare: configured space URI does not parse"); WriteError::Upstream })?; let Ok(space_did) = Did::new(space_did.to_owned()) else { tracing::error!("flare: configured space DID is not a DID"); return Err(WriteError::Upstream); }; let document = self.did_resolver.resolve(&space_did).await.map_err(|_| { tracing::warn!("flare: space owner's DID document could not be read"); WriteError::Upstream })?; let Some(pds) = document.get_pds_endpoint() else { tracing::warn!("flare: space owner's DID document names no PDS"); return Err(WriteError::Upstream); }; let record = self .fetch_json(&format!( "{pds}/xrpc/com.atproto.repo.getRecord\ ?repo={did}&collection={SPACE_NSID}&rkey={rkey}", did = space_did.as_str() )) .await .map_err(|_| { tracing::warn!("flare: space record could not be read"); WriteError::Upstream })?; let Some(cid) = record.get("cid").and_then(serde_json::Value::as_str) else { tracing::warn!("flare: space record answered without a CID"); return Err(WriteError::Upstream); }; if !is_blob_cid(cid) { tracing::warn!("flare: space record's CID is not shaped like one"); return Err(WriteError::Upstream); } Ok(serde_json::json!({ "uri": space_uri, "cid": cid })) } /// One GET at the other end of the world, with a deadline on it. /// /// The host comes out of a DID document the player controls, so it may be /// anyone's server and may never answer. Without a timeout the request /// holding it — `POST /api/matches` — waits as long as that server likes. /// `Matches::advance` puts a deadline on its own probe for the same /// reason. async fn get(&self, url: &str) -> Result { self.http .get(url) .timeout(std::time::Duration::from_secs(10)) .send() .await } async fn fetch_json(&self, url: &str) -> Result { let response = self.get(url).await.map_err(|_| upstream("record fetch"))?; if !response.status().is_success() { return Err(status_error("record fetch", response.status())); } // from_slice rather than reqwest's own json(): reqwest is built here // without that feature, and this is the whole of what it would add. let body = read_capped(response, RECORD_BYTE_LIMIT, "record body").await?; serde_json::from_slice(&body).map_err(|_| upstream("record decode")) } /// The DID as a DID, for a write about to be attempted. /// /// Nothing here asks whether the grant covers the collection. It used to, /// against the scope stored with the token, and that check could only ever /// be wrong in one of two directions: the authorization server expands an /// `include:` permission set into its member permissions before /// minting the token, so the string asked for is not the string granted /// and a literal comparison refuses a write the token would have been /// allowed to make — which is what happened to every flare. The PDS knows /// the answer for certain, says which permission is missing when it says /// no, and had to be asked anyway. So the write is attempted and its /// refusal is what a player is told. fn repo_did(&self, did: &str) -> Result { Did::new(did.to_owned()).map_err(|_| { // The DID came out of a cookie we signed ourselves. tracing::error!("repo: session DID is not a DID"); WriteError::Upstream }) } async fn session(&self, did: &Did) -> Result { self.client.restore(did).await.map_err(|_| { // atrium's errors can embed request URIs and server responses; // log the fact, not the value. tracing::warn!(did = did.as_str(), "repo: session could not be restored"); WriteError::NeedsSignIn }) } pub fn is_confidential(&self) -> bool { self.confidential } /// Served at /oauth/client-metadata.json in confidential mode. pub fn client_metadata(&self) -> serde_json::Value { let mut metadata = serde_json::to_value(&self.client.client_metadata).expect("metadata serializes"); // atrium's metadata type has no client_name, and without one the // consent screen shows the raw client_id URL. The authorization // server reads the document we serve, so adding it here is enough. metadata["client_name"] = "lance.blue".into(); metadata } /// Served at /.well-known/jwks.json in confidential mode. atrium strips /// the private members before handing this out. pub fn jwks(&self) -> serde_json::Value { serde_json::to_value(self.client.jwks()).expect("jwks serializes") } } /// A camo is a couple of kilobytes; this is the same ceiling `/api/camo` /// puts on an upload, and generous by three orders of magnitude. const BLOB_BYTE_LIMIT: usize = 1024 * 1024; /// The record is a name, a blob reference and the editor's settings. atproto /// caps a record at 64KiB anyway; this is room for a PDS that does not. const RECORD_BYTE_LIMIT: usize = 256 * 1024; /// Read a response body, giving up past `limit`. /// /// Not `Response::bytes()`, which reads whatever arrives. The server at the /// other end is the player's own, chosen by them, and reqwest is built here /// with gzip — so a few kilobytes on the wire can be gigabytes in memory, and /// nothing but this says otherwise. Chunk by chunk because the decompressed /// size is not knowable from a header, and a `Content-Length` a hostile server /// wrote is not evidence of anything. async fn read_capped( mut response: reqwest::Response, limit: usize, step: &'static str, ) -> Result, ReadError> { let mut body = Vec::new(); while let Some(chunk) = response.chunk().await.map_err(|_| upstream(step))? { if body.len() + chunk.len() > limit { tracing::warn!(step, limit, "camo: public read exceeded its size limit"); return Err(ReadError::Upstream); } body.extend_from_slice(&chunk); } Ok(body) } /// Whether a string is shaped like the CID of an atproto blob. /// /// Not a parse, and deliberately not a claim about the bytes: it says the /// value can be put in a query string and written into the launch manifest /// without carrying anything else in with it. The value is read out of a /// record served by a host the player controls, so it is the one input here /// that is neither ours nor checked by something else downstream. /// /// CIDv1 in base32 is `b` and then lowercase RFC 4648 without padding. The /// length bound is a sanity check rather than a spec: a raw sha-256 blob CID /// is 59 characters. fn is_blob_cid(cid: &str) -> bool { cid.len() >= 8 && cid.len() <= 128 && cid.starts_with('b') && cid[1..] .bytes() .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit()) } /// The step is named because these are three different requests to two /// different hosts and "camo read failed" says nothing. The URL is not logged: /// it carries a DID and a record key. fn upstream(step: &'static str) -> ReadError { tracing::warn!(step, "camo: public read failed"); ReadError::Upstream } fn status_error(step: &'static str, status: reqwest::StatusCode) -> ReadError { tracing::warn!(step, status = status.as_u16(), "camo: public read refused"); if status == reqwest::StatusCode::NOT_FOUND || status == reqwest::StatusCode::BAD_REQUEST { // A PDS answers a missing record with 400 InvalidRequest rather than // 404, so both mean the same thing here. ReadError::NotFound } else { ReadError::Upstream } } /// The collection, parsed. The constant is a valid NSID or this crate does /// not work at all, which a test asserts. fn camo_nsid() -> atrium_api::types::string::Nsid { CAMO_NSID.parse().expect("CAMO_NSID is a valid NSID") } fn discussion_nsid() -> atrium_api::types::string::Nsid { DISCUSSION_NSID .parse() .expect("DISCUSSION_NSID is a valid NSID") } /// `at://did/app.userinput.space/rkey` into its DID and record key. Public /// so config.rs can refuse a FLARE_SPACE that is not shaped like this at /// startup instead of on the first flare. pub fn split_space_uri(uri: &str) -> Option<(&str, &str)> { let rest = uri.strip_prefix("at://")?; let mut parts = rest.splitn(3, '/'); let did = parts.next()?; let collection = parts.next()?; let rkey = parts.next()?; (did.starts_with("did:") && collection == SPACE_NSID && !rkey.is_empty() && !rkey.contains('/')) .then_some((did, rkey)) } /// A record key from the path. Attacker-controlled, so it is parsed rather /// than pasted into a request. fn record_key(rkey: &str) -> Result { rkey.parse().map_err(|_| WriteError::NotFound) } /// A record as the data model, from the JSON this service assembled. /// /// The record travels to the PDS as JSON, so a blob's `{"$link": …}` is /// already the shape that goes on the wire; this conversion is what checks it /// is a shape atproto allows at all — a float anywhere in it fails here /// rather than at the PDS. fn to_unknown( record: serde_json::Map, ) -> Result { serde_json::Value::Object(record) .try_into_unknown() .map_err(|_| { tracing::error!("camo: record does not fit the atproto data model"); WriteError::Upstream }) } /// A new `blue.lance.camo`, as JSON. /// /// `source` goes in beside the declared fields and is not looked at. The /// lexicon does not declare it, validation ignores it, and what it means is /// the editor's business — see the `lexicons` repo's TODO/02-camo.md for why /// it is deliberately not in the schema. fn camo_record( name: &str, image: serde_json::Value, created_at: &str, source: Option, ) -> serde_json::Map { let mut record = serde_json::Map::new(); record.insert("$type".into(), CAMO_NSID.into()); record.insert("name".into(), name.into()); record.insert("image".into(), image); record.insert("createdAt".into(), created_at.into()); if let Some(source) = source { record.insert("source".into(), source); } record } /// Whether a value can be written at all. /// /// The atproto data model has no float. atrium enforces that when it /// serialises the request, which is far enough downstream that the failure /// arrives as an opaque transport error — that is exactly how a camo whose /// settings held slider fractions came back to players as "your account's /// server would not accept the change" about a server that was never asked. /// Checking here turns it into a sentence about the actual problem. fn has_fraction(value: &serde_json::Value) -> bool { match value { serde_json::Value::Number(n) => n.as_i64().is_none() && n.as_u64().is_none(), serde_json::Value::Array(items) => items.iter().any(has_fraction), serde_json::Value::Object(fields) => fields.values().any(has_fraction), _ => false, } } /// A blob reference as JSON, for splicing into a record being assembled. fn blob_json(blob: atrium_api::types::BlobRef) -> Result { serde_json::to_value(blob).map_err(|_| { tracing::error!("camo: blob reference does not serialize"); WriteError::Upstream }) } /// atrium's own `upload_blob` sends `Content-Type: */*`. A PDS that cannot /// sniff the type records that string as the blob's mimeType, and the lexicon /// accepts `image/png` and nothing else — so the request is built here in /// order to declare it. async fn upload_png( session: &client::Session, png: Vec, ) -> Result { let request = XrpcRequest::<(), Vec> { method: atrium_xrpc::http::Method::POST, nsid: upload_blob::NSID.into(), parameters: None, input: Some(InputDataOrBytes::Bytes(png)), encoding: Some("image/png".to_owned()), }; let response = session .send_xrpc::<(), Vec, upload_blob::Output, upload_blob::Error>(&request) .await .map_err(|e| write_failed("upload blob", &e))?; match response { OutputDataOrBytes::Data(output) => Ok(output.data.blob.clone()), OutputDataOrBytes::Bytes(_) => { tracing::warn!("camo: uploadBlob answered with bytes, not a blob reference"); Err(WriteError::Upstream) } } } /// A refusal from the account's own server, carried out rather than reduced /// to a guess: the status, the error name and the message all travel, because /// the PDS is the only party that knows why it said no. `ScopeMissingError` /// naming the exact permission is the case this exists for. fn write_failed( step: &'static str, error: &atrium_xrpc::Error, ) -> WriteError { match error { // 401 with a `WWW-Authenticate` challenge. The header is where OAuth // puts `error="invalid_token"` and its description (RFC 6750), so it // is read for those rather than shown raw. atrium_xrpc::Error::Authentication(challenge) => { let refusal = Refusal { status: 401, error: sanitize_detail(challenge_param(challenge, "error").as_deref()), message: sanitize_detail( challenge_param(challenge, "error_description").as_deref(), ), }; tracing::warn!( step, error = refusal.error.as_deref().unwrap_or("(none)"), "repo: PDS rejected our authentication" ); WriteError::Refused(refusal) } atrium_xrpc::Error::XrpcResponse(response) => { let status = response.status.as_u16(); let (error, message) = match &response.error { Some(atrium_xrpc::error::XrpcErrorKind::Undefined(body)) => ( sanitize_detail(body.error.as_deref()), sanitize_detail(body.message.as_deref()), ), // A lexicon-defined error: its Display is the error name and // whatever the schema gave it, which is all there is to say. Some(atrium_xrpc::error::XrpcErrorKind::Custom(custom)) => { (sanitize_detail(Some(&custom.to_string())), None) } None => (None, None), }; tracing::warn!( step, status, error = error.as_deref().unwrap_or("(none)"), "repo: repository call refused" ); WriteError::Refused(Refusal { status, error, message, }) } // Everything that is not an answer from the server: the request could // not be built or serialised, the response could not be parsed, the // connection failed. The variant is the only thing that separates // them, and it carries no token or URI, so it is logged. other => { let kind = match other { atrium_xrpc::Error::HttpRequest(_) => "http request", atrium_xrpc::Error::HttpClient(_) => "http client", atrium_xrpc::Error::SerdeJson(_) => "serde json", atrium_xrpc::Error::SerdeHtmlForm(_) => "serde html form", atrium_xrpc::Error::UnexpectedResponseType => "unexpected response type", _ => "other", }; tracing::warn!(step, kind, "camo: repository call failed"); WriteError::Upstream } } } /// One quoted parameter out of a `WWW-Authenticate` challenge, e.g. the /// `invalid_token` in `DPoP error="invalid_token", error_description="…"`. /// /// Deliberately small: it finds `name="` and reads to the closing quote, /// which is the whole of what RFC 6750 challenges use. An unquoted or /// escaped value returns None rather than something half-parsed, and the /// refusal then travels with its status alone. fn challenge_param(challenge: &atrium_xrpc::http::HeaderValue, name: &str) -> Option { let text = challenge.to_str().ok()?; let start = text.find(&format!("{name}=\""))? + name.len() + 2; let rest = &text[start..]; let end = rest.find('"')?; // A backslash before the closing quote means the value was escaped, and // this reader would cut it in the wrong place. if rest[..end].ends_with('\\') { return None; } Some(rest[..end].to_owned()) } /// The same, for a read: a 400 from getRecord is RecordNotFound, which is a /// missing camo rather than a broken one. fn read_failed( step: &'static str, error: atrium_xrpc::Error, ) -> WriteError { match &error { atrium_xrpc::Error::XrpcResponse(response) if response.status.as_u16() == 400 => { WriteError::NotFound } _ => write_failed(step, &error), } } /// A `DynHandleResolver` scripted rather than networked: `resolve` answers /// by exact match against `(handle, did)` pairs given at construction, and /// refuses everything else — the same "unresolvable" answer a genuine DNS /// miss gives `resolve_player`. Lives outside `mod tests` because /// `Atproto::for_test`, which uses it, is `pub(crate)`: `routes.rs`'s and /// `matches::lobby`'s own test modules need to build an `Atproto` around one /// without reaching into this module's private types (`Handle`, `Did`) — /// see this crate's own note in this module's doc comment about atrium_* /// staying behind this module's wall. #[cfg(test)] struct FakeHandleResolver(std::collections::HashMap); #[cfg(test)] impl DynHandleResolver for FakeHandleResolver { fn resolve(&self, handle: Handle) -> BoxFuture<'_, Result> { let result = self .0 .get(handle.as_str()) .and_then(|did| Did::new(did.clone()).ok()) .ok_or(()); Box::pin(async move { result }) } } #[cfg(test)] impl Atproto { /// An `Atproto` whose handle resolution is scripted, not networked — /// the test seam `resolve_player`'s two production callers /// (`routes::plan_players`'s `player` control, /// `matches::lobby::Lobbies::invite_handle`) need to be exercised /// without live DNS-over-HTTPS egress, which `cargo test` must not /// depend on (plan/lobby.md). `resolves` is `(handle, did)` pairs, /// matched the same way `resolve_player` normalizes its input — lower /// the handle and drop any leading `@` before passing it in here; a /// handle with no entry refuses, like one DNS never heard of. pub(crate) fn for_test(db: Db, resolves: &[(&str, &str)]) -> Atproto { let config = Config::from_lookup(|_| None).expect("loopback config always builds"); let http = client::HttpClient::default(); let resolves = resolves .iter() .map(|(handle, did)| ((*handle).to_owned(), (*did).to_owned())) .collect(); Atproto { client: Arc::new(client::build(&config, db, &http).expect("loopback client builds")), did_resolver: client::did_resolver(&http), handle_resolver: Box::new(FakeHandleResolver(resolves)), http: reqwest::Client::new(), confidential: false, } } } #[cfg(test)] mod tests { use super::*; use std::collections::HashMap; fn temp_db() -> (tempfile::TempDir, Db) { let dir = tempfile::tempdir().unwrap(); let db = Db::open(&dir.path().join("test.sqlite")).unwrap(); (dir, db) } fn config(vars: &[(&str, &str)]) -> Config { let map: HashMap = vars .iter() .map(|(k, v)| (k.to_string(), v.to_string())) .collect(); Config::from_lookup(|key| map.get(key).cloned()).unwrap() } // The RFC 7515 A.3 example key with a kid added: a published test // vector, useful precisely because it is not a secret. Nothing may ever // deploy it. const TEST_JWK: &str = r#"{ "kty": "EC", "crv": "P-256", "kid": "test-1", "x": "f83OJ3D2xF1Bg8vub9tLe1gHMzV76e8Tus9uPHvRVEU", "y": "x_FEzRu9m36HLN_tue659LNpXW6pCyStikYjKIWI5a0", "d": "jpsQnnGQmL-YBIffH1136cspYG6-0iY7X1fCE9-E9LI" }"#; #[tokio::test] async fn loopback_shape() { let (_dir, db) = temp_db(); let atproto = Atproto::new(&config(&[]), db).unwrap(); assert!(!atproto.is_confidential()); let metadata = atproto.client_metadata(); let client_id = metadata["client_id"].as_str().unwrap(); assert!( client_id.starts_with("http://localhost?"), "got {client_id}" ); assert!( client_id.contains("127.0.0.1"), "redirect must be loopback IP, got {client_id}" ); assert!( client_id.contains("repo%3Ablue.lance.camo"), "loopback client_id carries its own scope, got {client_id}" ); } /// The scope has to be identical in the client metadata and in the /// authorization request; atrium takes them from different places and /// they only agree because both call `client::scopes`. #[test] fn login_asks_for_the_scopes_the_metadata_declares() { let scopes = client::scopes(); let declared: Vec<&str> = scopes.iter().map(AsRef::as_ref).collect(); assert_eq!( declared, [ "atproto", "repo:blue.lance.camo", "blob:image/png", "include:app.userinput.authBasic" ] ); } #[test] fn camo_nsid_is_the_lexicon_id() { assert_eq!(CAMO_NSID, "blue.lance.camo"); assert!(CAMO_NSID.parse::().is_ok()); assert!( DISCUSSION_NSID .parse::() .is_ok() ); } #[test] fn space_uris_split_or_refuse() { assert_eq!( split_space_uri("at://did:plc:abc/app.userinput.space/3kxyz"), Some(("did:plc:abc", "3kxyz")) ); assert_eq!( split_space_uri("at://did:plc:abc/app.userinput.space/"), None ); assert_eq!( split_space_uri("at://did:plc:abc/blue.lance.camo/3kxyz"), None, "wrong collection" ); assert_eq!(split_space_uri("https://userinput.app/s/x/y"), None); assert_eq!( split_space_uri("at://alice.example/app.userinput.space/3kxyz"), None, "a handle is not a DID" ); } /// The bytes that go on the wire, pinned. /// /// Two things are being held still. The blob has to serialise as a link — /// `{"$link": …}` — because a nested object there is a different record /// that no client will read as an image. And `source` has to survive the /// trip through the data model untouched, since the whole point of it is /// that this service does not understand it. #[test] fn camo_record_matches_the_lexicon() { const CID: &str = "bafkreibme22gw2h7y2h7tg2fhqotaqjucnbc24deqo72b6mkl2egezxhvy"; let image = serde_json::json!({ "$type": "blob", "ref": { "$link": CID }, "mimeType": "image/png", "size": 1234, }); let source = serde_json::json!({ "v": 1, "generator": "dpm", "params": { "seed": 7 }, }); let record = camo_record( "Highland", image, "2026-08-06T12:00:00.000Z", Some(source.clone()), ); let unknown = to_unknown(record).expect("record fits the data model"); assert_eq!( serde_json::to_value(&unknown).unwrap(), serde_json::json!({ "$type": "blue.lance.camo", "name": "Highland", "createdAt": "2026-08-06T12:00:00.000Z", "image": { "$type": "blob", "ref": { "$link": CID }, "mimeType": "image/png", "size": 1234, }, "source": source, }) ); } /// The atproto data model has no float, and atrium enforces that when it /// serialises the request — far downstream of anything that could name /// the problem. This is what production hit: a camo whose settings held /// slider fractions never reached the PDS at all. #[test] fn a_fractional_parameter_is_rejected_by_the_data_model() { let record = camo_record( "Highland", serde_json::json!({"$type": "blob", "ref": {"$link": "bafkreibme22gw2h7y2h7tg2fhqotaqjucnbc24deqo72b6mkl2egezxhvy"}, "mimeType": "image/png", "size": 12}), "2026-08-06T12:00:00.000Z", Some(serde_json::json!({"v": 1, "params": {"scale": 0.51}})), ); let unknown = to_unknown(record).expect("conversion itself succeeds"); assert!( serde_json::to_vec(&unknown).is_err(), "a float has to fail before the request is built" ); } /// So the write refuses it first, with something a person can act on. #[test] fn a_fraction_anywhere_in_the_settings_is_caught() { for (why, source) in [ ("a slider", serde_json::json!({"params": {"scale": 0.51}})), ( "nested in an array", serde_json::json!({"palette": [[1, 2, 0.5]]}), ), ("at the top", serde_json::json!({"v": 1.5})), ] { assert!(has_fraction(&source), "{why}"); } // What the editor sends now: whole numbers throughout. let good = serde_json::json!({ "v": 2, "generator": "dpm", "params": {"seed": 7, "scale": 510, "contrast": 590, "detail": 360, "angle": 0}, "palette": {"name": "Woodland", "colors": [[74, 83, 58]]}, }); assert!(!has_fraction(&good)); // And it survives the trip the fractional one could not make. let record = camo_record( "Highland", serde_json::json!({"$type": "blob", "ref": {"$link": "bafkreibme22gw2h7y2h7tg2fhqotaqjucnbc24deqo72b6mkl2egezxhvy"}, "mimeType": "image/png", "size": 12}), "2026-08-06T12:00:00.000Z", Some(good), ); let unknown = to_unknown(record).expect("record fits the data model"); serde_json::to_vec(&unknown).expect("and serialises onto the wire"); } /// A camo with no settings — one imported from a file — has no `source` /// key at all, rather than a null one. #[test] fn an_imported_camo_writes_no_source() { let record = camo_record( "From a file", serde_json::json!({"$type": "blob", "ref": {"$link": "bafkreibme22gw2h7y2h7tg2fhqotaqjucnbc24deqo72b6mkl2egezxhvy"}, "mimeType": "image/png", "size": 12}), "2026-08-06T12:00:00.000Z", None, ); assert!(!record.contains_key("source")); } #[tokio::test] async fn confidential_shape() { let (_dir, db) = temp_db(); let atproto = Atproto::new( &config(&[ ("PUBLIC_URL", "https://api.lance.blue"), ("WEB_ORIGIN", "https://lance.blue"), ("SESSION_SECRET", "0123456789abcdef0123456789abcdef"), ("PRIVATE_KEY_JWK", TEST_JWK), ]), db, ) .unwrap(); assert!(atproto.is_confidential()); let metadata = atproto.client_metadata(); assert_eq!( metadata["client_id"], "https://api.lance.blue/oauth/client-metadata.json" ); assert_eq!(metadata["token_endpoint_auth_method"], "private_key_jwt"); assert_eq!(metadata["client_name"], "lance.blue"); assert_eq!( metadata["jwks_uri"], "https://api.lance.blue/.well-known/jwks.json" ); assert_eq!( metadata["redirect_uris"][0], "https://api.lance.blue/oauth/callback" ); // Every collection is named in full: atproto has no partial // wildcard, so this string grows a term per collection we write. assert_eq!( metadata["scope"], "atproto repo:blue.lance.camo blob:image/png include:app.userinput.authBasic" ); let jwks = atproto.jwks(); let key = &jwks["keys"][0]; assert_eq!(key["kid"], "test-1"); assert!(key.get("d").is_none(), "private member served: {jwks}"); } #[tokio::test] async fn confidential_requires_kid() { let (_dir, db) = temp_db(); let no_kid = TEST_JWK.replace(r#""kid": "test-1","#, ""); let result = Atproto::new( &config(&[ ("PUBLIC_URL", "https://api.lance.blue"), ("SESSION_SECRET", "0123456789abcdef0123456789abcdef"), ("PRIVATE_KEY_JWK", &no_kid), ]), db, ); assert!(result.is_err(), "a key without a kid must be rejected"); } /// `PRIVATE_KEY_JWK` also accepts a JSON array — the shape a rotation /// publishes so a new key can start signing before the old one is /// withdrawn (plan/complete/sign-in.md). Both keys must be accepted and both must be /// served, or a rotation silently drops one signer. #[tokio::test] async fn confidential_array_of_keys_are_all_accepted() { let (_dir, db) = temp_db(); let second = TEST_JWK.replace("\"test-1\"", "\"test-2\""); let both = format!("[{TEST_JWK}, {second}]"); let atproto = Atproto::new( &config(&[ ("PUBLIC_URL", "https://api.lance.blue"), ("WEB_ORIGIN", "https://lance.blue"), ("SESSION_SECRET", "0123456789abcdef0123456789abcdef"), ("PRIVATE_KEY_JWK", &both), ]), db, ) .unwrap(); let jwks = atproto.jwks(); let mut kids: Vec<&str> = jwks["keys"] .as_array() .unwrap() .iter() .map(|key| key["kid"].as_str().unwrap()) .collect(); kids.sort_unstable(); assert_eq!(kids, ["test-1", "test-2"]); } /// One bad key in the array must fail construction the same way a bad /// single object does — not silently drop the offending key and carry /// on with only the good one. #[tokio::test] async fn confidential_array_with_a_missing_kid_is_rejected() { let (_dir, db) = temp_db(); let no_kid = TEST_JWK.replace(r#""kid": "test-1","#, ""); let mixed = format!("[{TEST_JWK}, {no_kid}]"); let result = Atproto::new( &config(&[ ("PUBLIC_URL", "https://api.lance.blue"), ("SESSION_SECRET", "0123456789abcdef0123456789abcdef"), ("PRIVATE_KEY_JWK", &mixed), ]), db, ); assert!( result.is_err(), "an array containing a key without a kid must be rejected" ); } #[tokio::test] async fn begin_login_rejects_non_handles() { let (_dir, db) = temp_db(); let atproto = Atproto::new(&config(&[]), db).unwrap(); assert!(matches!( atproto.begin_login("not a handle").await, Err(LoginError::InvalidHandle) )); assert!(matches!( atproto.begin_login("").await, Err(LoginError::InvalidHandle) )); } /// The blob reference is read out of a record served by a host the player /// controls, and it is put in a query string and then in the launch /// manifest. Everything else crossing that boundary is parsed — the DID, /// the record key, the image itself — so this one is checked too. #[test] fn only_something_shaped_like_a_cid_reaches_the_query_string() { // The blob CID of a real camo in @permadeath.com's repo. assert!(is_blob_cid( "bafkreicps3ralra6g43dzfcqaf7afwdbi3vlm6maoei3slesnrd54zdte4" )); // A second parameter smuggled in behind the first. assert!(!is_blob_cid("bafkrei…&did=did:plc:someone-else")); assert!(!is_blob_cid("bafkreicps3ralra6g43&cid=x")); // Header injection, path traversal, and an empty or absurd value. assert!(!is_blob_cid("bafkrei\r\nX-Evil: 1")); assert!(!is_blob_cid("../../etc/passwd")); assert!(!is_blob_cid("")); assert!(!is_blob_cid(&"b".repeat(200))); // Base32 here is lowercase; uppercase is a different encoding and not // one a PDS emits. assert!(!is_blob_cid("bAFKREICPS3RALRA6G43")); } /// The refusal that started all this: a token minted from an `include:` /// permission set carries the set's members, never the `include:` string, /// so a scope check on this side refused every flare and told the player /// to sign in again — which could not change the stored scope, because /// nothing was wrong with it. The PDS says what is missing; that sentence /// is what has to survive the trip. #[test] fn a_missing_permission_arrives_in_the_servers_own_words() { let error: atrium_xrpc::Error = atrium_xrpc::Error::XrpcResponse(atrium_xrpc::error::XrpcError { status: atrium_xrpc::http::StatusCode::FORBIDDEN, error: Some(atrium_xrpc::error::XrpcErrorKind::Undefined( atrium_xrpc::error::ErrorResponseBody { error: Some("ScopeMissingError".into()), message: Some( r#"Missing required scope "repo:app.userinput.discussion?action=create""# .into(), ), }, )), }); let WriteError::Refused(refusal) = write_failed("create discussion", &error) else { panic!("a refusal from the PDS must arrive as one"); }; assert_eq!(refusal.status, 403); assert!(refusal.fixable_by_signing_in()); assert_eq!( refusal.detail().unwrap(), r#"Missing required scope "repo:app.userinput.discussion?action=create""# ); } /// Another server's text, on its way to our users. #[test] fn relayed_text_is_stripped_and_capped() { assert_eq!(sanitize_detail(Some(" refused ")).unwrap(), "refused"); assert_eq!(sanitize_detail(Some("one\ntwo")).unwrap(), "onetwo"); assert_eq!(sanitize_detail(Some(" ")), None); assert_eq!(sanitize_detail(None), None); assert_eq!(sanitize_detail(Some(&"a".repeat(500))).unwrap().len(), 200); } /// 401s carry their reason in the challenge header, which is where OAuth /// puts it, and the quoted value is the only part worth showing. #[test] fn a_challenge_header_gives_up_its_reason() { let header = atrium_xrpc::http::HeaderValue::from_static( r#"DPoP error="invalid_token", error_description="Token has expired""#, ); assert_eq!(challenge_param(&header, "error").unwrap(), "invalid_token"); assert_eq!( challenge_param(&header, "error_description").unwrap(), "Token has expired" ); assert_eq!(challenge_param(&header, "realm"), None); // Unquoted and escaped values are left alone rather than cut wrong. let bare = atrium_xrpc::http::HeaderValue::from_static("DPoP error=invalid_token"); assert_eq!(challenge_param(&bare, "error"), None); let escaped = atrium_xrpc::http::HeaderValue::from_static(r#"DPoP error_description="a \" b""#); assert_eq!(challenge_param(&escaped, "error_description"), None); } /// A missing record comes back as InvalidRequest, not 404, and the caller /// has to tell "no such camo" from "the server is down" to know whether /// retrying is worth anything. #[test] fn a_pds_saying_invalid_request_means_no_such_record() { assert!(matches!( status_error("record fetch", reqwest::StatusCode::BAD_REQUEST), ReadError::NotFound )); assert!(matches!( status_error("record fetch", reqwest::StatusCode::NOT_FOUND), ReadError::NotFound )); assert!(matches!( status_error("record fetch", reqwest::StatusCode::BAD_GATEWAY), ReadError::Upstream )); } }