From 73d5ec0176ca6420c69336db5911234f0eef6b9e Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Thu, 3 Sep 2026 18:37:26 -0400 Subject: [PATCH] fix(pds)!: announce the account's handle on the identity frame `#identity` now carries `claimed_handle`'s handle, the one `alsoKnownAs` and `/.well-known/atproto-did` are both computed from, so an indexer can run atproto's bidirectional check on what the stream gave it. `relay_view.rs` drives the server from outside the process and is what caught the omission. Change-Id: Id50703bd504f91407a0f5cdb20383188cb64e901 --- crates/didbot-pds/src/provision.rs | 11 +- crates/didbot-serve/src/routes.rs | 12 +- crates/didbot/tests/relay_view.rs | 448 +++++++++++++++++++++++++++++ docs/conformance.md | 22 +- plan/pds-xrpc.md | 20 ++ 5 files changed, 498 insertions(+), 15 deletions(-) create mode 100644 crates/didbot/tests/relay_view.rs diff --git a/crates/didbot-pds/src/provision.rs b/crates/didbot-pds/src/provision.rs index 40b1e5e1..fbe7c450 100644 --- a/crates/didbot-pds/src/provision.rs +++ b/crates/didbot-pds/src/provision.rs @@ -2285,6 +2285,15 @@ where /// Only ever called once the document actually serves. A `#identity` for /// a DID that does not resolve yet is a lie a consumer acts on: it /// re-resolves, gets nothing, and caches the nothing. + /// + /// The handle is [`claimed_handle`]'s — the same function `alsoKnownAs` + /// and `/.well-known/atproto-did` are computed from, so what the stream + /// announces is what the bidirectional handle rule resolves either way. + /// [`Self::issued_handle`] answers a narrower question, and answers it + /// for DNS: it withholds a handle that is already the DID's own hostname, + /// because publishing and withdrawing that separately would be publishing + /// one name twice. The lexicon asks for "the current handle for the + /// account", which an account storing no handle of its own still has. fn announce_repo_identity(&self, account: &AgentAccount) { let Some(repos) = &self.repos else { return; @@ -2294,7 +2303,7 @@ where } repos.identity(&RepoIdentity { did: account.did.as_str().to_owned(), - handle: self.issued_handle(account).map(ToOwned::to_owned), + handle: claimed_handle(&account.did, account.handle.as_deref()).map(ToOwned::to_owned), }); } diff --git a/crates/didbot-serve/src/routes.rs b/crates/didbot-serve/src/routes.rs index c3626493..a35e4487 100644 --- a/crates/didbot-serve/src/routes.rs +++ b/crates/didbot-serve/src/routes.rs @@ -2889,12 +2889,12 @@ async fn delete_session(State(state): State, headers: HeaderMap) -> Re /// `GET /xrpc/com.atproto.server.getSession` /// -/// The one route today that actually requires -/// [`Credential::LegacySession`](crate::auth::Credential) and checks it: the -/// account admin surface and every repository write route stay `Public` for -/// now (see `auth::ROUTE_CREDENTIALS`'s own note), so this is where a -/// session answering for the account it names, and no other, is exercised -/// end to end rather than only at the extractor level. +/// Where a session answering for the account it names, and no other, is +/// exercised end to end rather than only at the extractor level: this route +/// reads [`Credential::LegacySession`](crate::auth::Credential) and answers +/// with the account that credential stands for. The repository write routes +/// take an agent token instead — see `auth::ROUTE_CREDENTIALS` for which +/// credential each method takes. async fn get_session(State(state): State, headers: HeaderMap) -> Response { let did = match auth::require_access_token(&state.sessions, &headers) { Ok(did) => did, diff --git a/crates/didbot/tests/relay_view.rs b/crates/didbot/tests/relay_view.rs new file mode 100644 index 00000000..6108250c --- /dev/null +++ b/crates/didbot/tests/relay_view.rs @@ -0,0 +1,448 @@ +//! What this server looks like to somebody who only has its address. +//! +//! Every other test of `com.atproto.sync.subscribeRepos` reaches into the +//! process to make the events it then reads off the wire — +//! `tests/subscribe_repos.rs` calls `Registry::provision` and +//! `Registry::put_record` directly. That checks the producer. It cannot check +//! that the wire surface a stranger *drives* and the wire surface a stranger +//! *reads* describe the same repository, because only one of the two is on +//! the wire. +//! +//! This binary closes that loop. A relay's half is a WebSocket opened on +//! `/xrpc/com.atproto.sync.subscribeRepos` with no credential; a client's +//! half is `bot.did.provisionAgent`, `com.atproto.repo.createRecord`, +//! `com.atproto.repo.getRecord`, `com.atproto.repo.listRecords`, +//! `com.atproto.server.describeServer`, `com.atproto.sync.getLatestCommit`, +//! `/.well-known/atproto-did` and `/.well-known/did.json`, over HTTP with +//! nothing but a bearer token the server itself handed back. The assertions +//! are the joins between the two: the `cid` `createRecord` answered is the +//! `cid` the `#commit`'s op carries and the `cid` `getRecord` returns, the +//! `commit` link on the frame is the one `getLatestCommit` names, and the +//! handle the `#identity` frame announced is the handle +//! `/.well-known/atproto-did` maps back to the DID the frame named. +//! +//! The WebSocket client and the server harness are [`support`]'s, shared with +//! `tests/subscribe_repos.rs` and `tests/subscribe_backpressure.rs`. The +//! account-making and record-writing helpers on that harness are deliberately +//! unused here: using them would put the thing under test back inside the +//! process. + +mod support; + +use std::time::Duration; + +use didbot::data::{dag_cbor, Value}; +use support::stack::{Stack, ZONE}; +use support::ws; + +/// A collection this deployment holds no schema for, used here as an ordinary +/// record that is not one of its own. +const THING: &str = "com.example.thing"; + +/// How long a step may take before it counts as a hang. +const PATIENCE: Duration = Duration::from_secs(10); + +/// One decoded frame: the header and the body a frame's two concatenated +/// DAG-CBOR values are. +struct Frame { + header: Value, + body: Value, +} + +impl Frame { + /// Splits and decodes a frame, the way a consumer with no framing library + /// has to: try every split point until both halves decode. + fn parse(bytes: &[u8]) -> Self { + for split in 1..bytes.len() { + let (Ok(header), Ok(body)) = ( + dag_cbor::decode(&bytes[..split]), + dag_cbor::decode(&bytes[split..]), + ) else { + continue; + }; + return Self { header, body }; + } + panic!("a frame was not two concatenated DAG-CBOR values"); + } + + /// The `t` this frame's header names. + fn tag(&self) -> String { + match field(&self.header, "t") { + Some(Value::String(tag)) => tag.clone(), + other => panic!("a frame header carries no string t: {other:?}"), + } + } + + /// The `seq` its body carries. + fn seq(&self) -> i64 { + match field(&self.body, "seq") { + Some(Value::Integer(seq)) => *seq, + other => panic!("a frame carries no integer seq: {other:?}"), + } + } + + /// A string field of the body. + fn text(&self, key: &str) -> String { + match field(&self.body, key) { + Some(Value::String(value)) => value.clone(), + other => panic!("{key} is not a string on this frame: {other:?}"), + } + } + + /// A CID-link field of the body, as the base32 string every JSON route + /// on this server spells the same CID with. + fn link(&self, key: &str) -> String { + match field(&self.body, key) { + Some(Value::Link(cid)) => cid.to_string(), + other => panic!("{key} is not a CID link on this frame: {other:?}"), + } + } +} + +fn field<'a>(value: &'a Value, key: &str) -> Option<&'a Value> { + match value { + Value::Map(fields) => fields.get(key), + _ => None, + } +} + +/// The single `ops` entry of a `#commit`, as `(action, path, cid)`. +/// +/// A `cid` of `None` is the null the lexicon requires on a removal; every +/// other action carries the record's CID. +fn only_op(frame: &Frame) -> (String, String, Option) { + let Some(Value::List(ops)) = field(&frame.body, "ops") else { + panic!("a #commit carries no ops list: {:?}", frame.body); + }; + assert_eq!(ops.len(), 1, "expected one op, got {ops:?}"); + let op = &ops[0]; + let action = match field(op, "action") { + Some(Value::String(action)) => action.clone(), + other => panic!("an op carries no string action: {other:?}"), + }; + let path = match field(op, "path") { + Some(Value::String(path)) => path.clone(), + other => panic!("an op carries no string path: {other:?}"), + }; + let cid = match field(op, "cid") { + Some(Value::Link(cid)) => Some(cid.to_string()), + Some(Value::Null) | None => None, + other => panic!("an op's cid is neither a link nor null: {other:?}"), + }; + (action, path, cid) +} + +/// The next frame, or a failure rather than a hang. +async fn next(client: &mut ws::Client) -> Frame { + let bytes = tokio::time::timeout(PATIENCE, client.next()) + .await + .expect("a frame arrives before the deadline") + .expect("the server did not close the stream"); + Frame::parse(&bytes) +} + +/// A stranger's HTTP client and the one address it was given. +struct Outsider { + http: reqwest::Client, + base: String, +} + +impl Outsider { + fn new(stack: &Stack) -> Self { + Self { + http: reqwest::Client::builder() + .timeout(PATIENCE) + .build() + .expect("a default client builds"), + base: format!("http://{}", stack.addr), + } + } + + /// `GET`, with an optional `Host` header for the routes that read one. + async fn get(&self, path: &str, host: Option<&str>) -> (reqwest::StatusCode, String) { + let mut request = self.http.get(format!("{}{path}", self.base)); + if let Some(host) = host { + request = request.header(reqwest::header::HOST, host); + } + let response = request.send().await.expect("the server answers"); + let status = response.status(); + (status, response.text().await.expect("a body reads")) + } + + /// `GET`, expecting JSON and a 200. + async fn get_json(&self, path: &str, host: Option<&str>) -> serde_json::Value { + let (status, body) = self.get(path, host).await; + assert_eq!(status, 200, "GET {path} answered {status}: {body}"); + serde_json::from_str(&body).unwrap_or_else(|err| panic!("GET {path} is not JSON: {err}")) + } + + /// `POST` of a JSON body, with an optional bearer token, expecting a 200. + async fn post_json( + &self, + path: &str, + token: Option<&str>, + body: serde_json::Value, + ) -> serde_json::Value { + let mut request = self.http.post(format!("{}{path}", self.base)).json(&body); + if let Some(token) = token { + request = request.bearer_auth(token); + } + let response = request.send().await.expect("the server answers"); + let status = response.status(); + let text = response.text().await.expect("a body reads"); + assert_eq!(status, 200, "POST {path} answered {status}: {text}"); + serde_json::from_str(&text).unwrap_or_else(|err| panic!("POST {path} is not JSON: {err}")) + } +} + +/// A `str` field of a JSON object, or a failure naming what was there. +fn json_str(value: &serde_json::Value, key: &str) -> String { + value[key] + .as_str() + .unwrap_or_else(|| panic!("{key} is not a string in {value}")) + .to_owned() +} + +/// A consumer that opened the stream before anything existed sees an account +/// appear and write, in the order and with the identifiers the same server's +/// own XRPC routes report. +/// +/// This is the shape of an alpha's first federated moment: a relay is +/// subscribed, an agent is minted, the agent writes, and what the relay reads +/// has to be the same repository the agent's client just wrote to. Every +/// identifier asserted here is compared against a *different* route's answer +/// rather than against a constant, so the two halves cannot agree by both +/// being wrong in the same way. +#[tokio::test(flavor = "multi_thread")] +async fn a_stranger_on_the_stream_watches_an_account_appear_and_write() { + let stack = Stack::start(64).await; + let outsider = Outsider::new(&stack); + + // Subscribed with no cursor and no credential, before the account it is + // about to be told about exists. + let mut relay = stack.follow(None).await; + + let provisioned = outsider + .post_json( + "/xrpc/bot.did.provisionAgent", + None, + serde_json::json!({ "agentId": "kestrel" }), + ) + .await; + let did = json_str(&provisioned, "did"); + let handle = json_str(&provisioned, "handle"); + let token = json_str(&provisioned, "agentToken"); + + // Identity, then account, then whatever commits provisioning itself + // writes. The order is the one a consumer needs: nothing may name a + // repository before the frame that says the repository exists. + let identity = next(&mut relay).await; + assert_eq!(identity.tag(), "#identity"); + assert_eq!(identity.text("did"), did); + assert_eq!( + identity.text("handle"), + handle, + "the handle on the stream is not the handle provisioning answered with" + ); + + let account = next(&mut relay).await; + assert_eq!(account.tag(), "#account"); + assert_eq!(account.text("did"), did); + assert_eq!(field(&account.body, "active"), Some(&Value::Bool(true))); + + let mut seqs = vec![identity.seq(), account.seq()]; + + // The write a client makes. Provisioning writes records of its own first, + // so the frame this is looking for is found by the record it names rather + // than by counting: a consumer cannot know how many records an account is + // born holding, only which one it asked for. + let written = outsider + .post_json( + "/xrpc/com.atproto.repo.createRecord", + Some(&token), + serde_json::json!({ + "repo": did, + "collection": THING, + "record": { + "$type": THING, + "text": "reading the cart", + "emoji": "\u{1f50e}", + "createdAt": "2026-01-01T00:00:00Z", + }, + }), + ) + .await; + let uri = json_str(&written, "uri"); + let record_cid = json_str(&written, "cid"); + let rkey = uri + .rsplit_once('/') + .expect("an at:// URI ends in a record key") + .1 + .to_owned(); + let wanted = format!("{THING}/{rkey}"); + + let commit = loop { + let frame = next(&mut relay).await; + assert_eq!( + frame.tag(), + "#commit", + "a frame that is neither identity nor account nor commit reached a consumer" + ); + assert_eq!(frame.text("repo"), did); + seqs.push(frame.seq()); + if only_op(&frame).1 == wanted { + break frame; + } + }; + + // Sequence numbers: one space, strictly increasing, no gap across + // everything one account's first minute produced. + assert!( + seqs.len() >= 3, + "an account was born and wrote in fewer than three frames: {seqs:?}" + ); + for pair in seqs.windows(2) { + assert_eq!( + pair[1], + pair[0] + 1, + "a consumer resuming inside this run would be told its cursor is \ + unreachable over a number nothing spent: {seqs:?}" + ); + } + + // The op names the record the client just wrote, by the CID the route + // that wrote it answered with. + let (action, path, op_cid) = only_op(&commit); + assert_eq!(action, "create"); + assert_eq!(path, format!("{THING}/{rkey}")); + assert_eq!( + op_cid.as_deref(), + Some(record_cid.as_str()), + "the cid on the stream is not the cid createRecord answered with" + ); + + // And a stranger reading the record back over XRPC gets the same CID + // again -- so the stream, the write and the read all name one record. + let read = outsider + .get_json( + &format!("/xrpc/com.atproto.repo.getRecord?repo={did}&collection={THING}&rkey={rkey}"), + None, + ) + .await; + assert_eq!(json_str(&read, "uri"), uri); + assert_eq!(json_str(&read, "cid"), record_cid); + assert_eq!(read["value"]["text"], "reading the cart"); + + // The commit the frame points at is the repository's current head, by the + // route a relay uses to check exactly that after a gap. + let head = outsider + .get_json( + &format!("/xrpc/com.atproto.sync.getLatestCommit?did={did}"), + None, + ) + .await; + assert_eq!( + json_str(&head, "cid"), + commit.link("commit"), + "getLatestCommit names a different commit than the frame that announced it" + ); + assert_eq!(json_str(&head, "rev"), commit.text("rev")); +} + +/// The identifiers the stream announced resolve, over the routes an +/// off-the-shelf client uses before it has ever heard of this project. +/// +/// A relay hands an indexer a DID and a handle. Everything downstream of that +/// -- an app showing a profile, a client logging in -- starts by resolving +/// them, and by asking the server what it is. This drives that path with the +/// values taken off the wire rather than constructed here, so a handle the +/// stream announces that nothing can resolve fails. +#[tokio::test(flavor = "multi_thread")] +async fn the_identifiers_the_stream_announced_resolve_from_outside() { + let stack = Stack::start(64).await; + let outsider = Outsider::new(&stack); + let mut relay = stack.follow(None).await; + + outsider + .post_json( + "/xrpc/bot.did.provisionAgent", + None, + serde_json::json!({ "agentId": "quernstone" }), + ) + .await; + + let identity = next(&mut relay).await; + assert_eq!(identity.tag(), "#identity"); + let did = identity.text("did"); + let handle = identity.text("handle"); + + // Handle to DID, the mapping the bidirectional-verification rule needs. + let (status, resolved) = outsider + .get("/.well-known/atproto-did", Some(&handle)) + .await; + assert_eq!( + status, 200, + "the announced handle does not resolve: {resolved}" + ); + assert_eq!(resolved.trim(), did); + + // DID to document, and the service entry that says where this repository + // is served from. A client with the document and nothing else has to be + // able to find the server. + let document = outsider + .get_json("/.well-known/did.json", Some(&handle)) + .await; + assert_eq!(json_str(&document, "id"), did); + assert_eq!( + document["alsoKnownAs"][0], + serde_json::Value::String(format!("at://{handle}")), + "the document does not point back at the handle the stream announced" + ); + let service = document["service"] + .as_array() + .expect("the document carries a service list") + .iter() + .find(|entry| entry["type"] == "AtprotoPersonalDataServer") + .expect("the document carries an #atproto_pds service entry"); + assert_eq!(json_str(service, "id"), format!("{did}#atproto_pds")); + assert!( + !json_str(service, "serviceEndpoint").is_empty(), + "the service entry names no endpoint, so nothing can find this repository" + ); + + // What the server says it is. `availableUserDomains` is what a client + // uses to decide a handle belongs here at all. + let described = outsider + .get_json("/xrpc/com.atproto.server.describeServer", None) + .await; + assert!( + json_str(&described, "did").starts_with("did:"), + "describeServer answers no service DID: {described}" + ); + let domains: Vec = described["availableUserDomains"] + .as_array() + .expect("availableUserDomains is a list") + .iter() + .map(|value| value.as_str().unwrap_or_default().to_owned()) + .collect(); + assert!( + domains + .iter() + .any(|domain| domain.trim_start_matches('.') == ZONE), + "describeServer's availableUserDomains {domains:?} does not include the zone \ + the handle it just minted is under" + ); + + // Finally, the handle itself as an at-identifier on a repository route: + // an off-the-shelf client browses by handle, not by DID. + let listed = outsider + .get_json( + &format!("/xrpc/com.atproto.repo.listRecords?repo={handle}&collection={THING}"), + None, + ) + .await; + assert!( + listed["records"].is_array(), + "listRecords by handle answered no records list: {listed}" + ); +} diff --git a/docs/conformance.md b/docs/conformance.md index f40c113f..e64ff136 100644 --- a/docs/conformance.md +++ b/docs/conformance.md @@ -422,12 +422,11 @@ field for the same reason. `com.atproto.server.{create,refresh,delete,get}Session` are new surface, added by [auth-types](../plan/auth-types.md). `getSession` is not in that -epic's original list — it was added because nothing else could check the -credential end to end: every `com.atproto.repo.*` write route stays public -for now (see that epic's own write-surface note), so `getSession` is the one -route that actually requires `Authorization` and answers with the account it -names, which is what proves a session works rather than merely that it -parses. They join `wire.rs`'s suite the same way every other route does: +epic's original list — it was added because it is the one route that answers +with the account a credential names, which is what proves a session works +rather than merely that it parses. The repository write routes take the agent +token `provisionAgent` hands back instead; see `ATPROTO_CREDENTIALS` in +`crates/didbot-serve/src/routes.rs` for which credential each method takes. They join `wire.rs`'s suite the same way every other route does: `procedure` validates the request against `createSession`'s vendored `input` schema before it is sent, and `.conforms` validates the response against each method's vendored `output` schema — `createSession` and @@ -678,8 +677,15 @@ being complete are two different facts, and `#identity` says the second: Sending it any earlier — while the account exists but its initial records are still landing — would be a lie a consumer acts on: it fetches the repository, finds it empty or partial, and caches that. See -`AccountState::Provisioning` and `Provisioner::announceable`. It carries the -handle this deployment issued, when it issued one. +`AccountState::Provisioning` and `Provisioner::announceable`. + +It carries the account's handle: the one `claimed_handle` computes, which is +what the DID document's `alsoKnownAs` claims and what +`/.well-known/atproto-did` maps back to this DID. One function answers all +three, so an indexer that takes the handle off the frame and runs atproto's +bidirectional check against this server gets an agreement. An account that +asked for no handle of its own is announced under its DID's own hostname, +which is the name it answers to. `#account` goes out twice in an account's life: `active: true` beside the `#identity` above, and `active: false, status: "deleted"` once the hostname is diff --git a/plan/pds-xrpc.md b/plan/pds-xrpc.md index a00a84de..e2b0acaa 100644 --- a/plan/pds-xrpc.md +++ b/plan/pds-xrpc.md @@ -36,6 +36,26 @@ checked today. ## Done +- [x] **The stream and the routes agree, checked from outside the process.** + `crates/didbot/tests/relay_view.rs` drives this server the way an + outside implementation would: a WebSocket on `subscribeRepos` with no + credential, and `bot.did.provisionAgent`, + `com.atproto.repo.{createRecord,getRecord,listRecords}`, + `com.atproto.sync.getLatestCommit`, + `com.atproto.server.describeServer`, `/.well-known/atproto-did` and + `/.well-known/did.json` over HTTP, carrying nothing but the token the + server itself handed back. Every identifier is joined to a second + route's answer rather than to a constant: the `cid` on the commit's op + is the one `createRecord` returned and `getRecord` returns, the frame's + `commit` link is `getLatestCommit`'s head, the sequence numbers across + an account's first minute are contiguous, and the handle on + `#identity` is the handle `/.well-known/atproto-did` maps back to the + DID the same frame named. That last join is what the test was written + for and what it caught: `#identity` now carries `claimed_handle`'s + handle — the one `alsoKnownAs` and the well-known are both computed + from — so an indexer can run atproto's bidirectional check on what the + stream gave it. + - [x] **Announce a deletion.** A `deleteRecord` now goes out on both streams, the same way a write does. On `com.atproto.sync.subscribeRepos` it is a `#commit` whose op is a `delete` — null `cid`, and the removed record's -- 2.51.2