//! Writing `bot.did.operator` into the operator's own repository. //! //! This is the one step in this command that touches the operator's own //! account with write authority, and it is deliberately the last thing this //! command does -- see [`crate::operate::orchestrate`] for the ordering and why //! nothing calls this module until every earlier check has passed. The //! request itself is `com.atproto.repo.putRecord`, not `createRecord`: the //! record key is fixed at the server's own authority (see [`crate::operate::record`]), //! so //! `putRecord` is what makes a re-run an *update* of the same claim rather //! than a second, conflicting one -- the idempotence //! `plan/handshake.md` asks for falls out of using the right verb rather //! than needing its own check. //! //! # Why this hand-rolls an `XrpcRequest` //! //! `jacquard-oauth`'s `OAuthSession` calls typed `jacquard_common::xrpc::XrpcRequest` //! values -- that is what gets a request DPoP-signed and retried on a stale //! nonce or an expired token without this module reimplementing either. No //! generated binding for `com.atproto.repo.putRecord` is a dependency here //! (`jacquard-api`'s lexicon codegen is not part of this workspace), so //! `PutRecordInput` and `PutRecordResponse` below implement the trait by //! hand, for this one method, rather than pulling in a full generated API //! surface for a single call. use didbot_claim_check::OperatorClaim; use jacquard_common::xrpc::{XrpcClient, XrpcMethod, XrpcRequest, XrpcResp}; use serde::{Deserialize, Serialize}; use smol_str::SmolStr; /// `com.atproto.repo.putRecord`'s input. #[derive(Debug, Serialize)] struct PutRecordInput { repo: String, collection: String, rkey: String, record: serde_json::Value, } /// The one field this command reads back out of a successful write. #[derive(Debug, Deserialize)] pub struct PutRecordOutput { /// The `at://` URI the record now lives at. pub uri: String, } /// The error shape `com.atproto.repo.putRecord` answers with when it is not /// a bare auth failure (those are [`jacquard_common::error::AuthError`], /// handled by `jacquard-oauth` itself) -- `{"error": ..., "message": ...}`, /// read as one opaque string rather than matched by name, since this /// command's own refusal message (see [`WriteError::Refused`]) already /// tells the operator what to check regardless of which specific atproto /// error code came back. #[derive(Debug, Deserialize, Serialize, thiserror::Error)] #[error("{0}")] struct PutRecordErr(SmolStr); struct PutRecordResponse; impl XrpcResp for PutRecordResponse { const NSID: &'static str = "com.atproto.repo.putRecord"; const ENCODING: &'static str = "application/json"; type Output = PutRecordOutput; type Err = PutRecordErr; } impl XrpcRequest for PutRecordInput { const NSID: &'static str = "com.atproto.repo.putRecord"; const METHOD: XrpcMethod = XrpcMethod::Procedure("application/json"); type Response = PutRecordResponse; } /// `com.atproto.repo.deleteRecord`'s input. #[derive(Debug, Serialize)] struct DeleteRecordInput { repo: String, collection: String, rkey: String, } /// `deleteRecord`'s answer, none of which this command reads. #[derive(Debug, Deserialize)] pub struct DeleteRecordOutput {} struct DeleteRecordResponse; impl XrpcResp for DeleteRecordResponse { const NSID: &'static str = "com.atproto.repo.deleteRecord"; const ENCODING: &'static str = "application/json"; type Output = DeleteRecordOutput; type Err = PutRecordErr; } impl XrpcRequest for DeleteRecordInput { const NSID: &'static str = "com.atproto.repo.deleteRecord"; const METHOD: XrpcMethod = XrpcMethod::Procedure("application/json"); type Response = DeleteRecordResponse; } /// `com.atproto.repo.getRecord`'s parameters. #[derive(Debug, Serialize)] struct GetRecordParams { repo: String, collection: String, rkey: String, } /// `getRecord`'s answer, reduced to the record. #[derive(Debug, Deserialize)] pub struct GetRecordOutput { /// The record body. pub value: serde_json::Value, } struct GetRecordResponse; impl XrpcResp for GetRecordResponse { const NSID: &'static str = "com.atproto.repo.getRecord"; const ENCODING: &'static str = "application/json"; type Output = GetRecordOutput; type Err = PutRecordErr; } impl XrpcRequest for GetRecordParams { const NSID: &'static str = "com.atproto.repo.getRecord"; const METHOD: XrpcMethod = XrpcMethod::Query; type Response = GetRecordResponse; } /// Why the write did not happen. #[derive(Debug, thiserror::Error)] pub enum WriteError { /// The record's own fields could not be serialized to RFC 3339. Should /// not happen for any timestamp `operate` constructs; see /// [`crate::operate::record::to_record_value`]. #[error("could not format the record to write: {0}")] Malformed(#[from] time::error::Format), /// The operator's PDS refused the write. The likeliest cause named /// explicitly, because this is a one-shot command run once by a person /// at the moment nothing else works yet -- a bare transport error here /// is expensive to leave generic. #[error( "the operator's PDS refused to write bot.did.operator ({detail}) -- check that the \ OAuth consent you approved actually granted repo:bot.did.operator write access, and \ that this account's PDS supports atproto's granular OAuth scope grammar rather than \ only the legacy transition:generic scope." )] Refused { /// What the server said. detail: String, }, /// The record that may already stand at the key could not be read /// back, so nothing was written over it. #[error("could not read bot.did.operator/{rkey} back from the operator's PDS ({detail})")] Read { /// The record key. rkey: String, /// What the PDS said. detail: String, }, } /// Reads a `bot.did.operator` record back before it is written over, so /// an admission the server refuses can put back what stood there rather /// than deleting it. #[allow(async_fn_in_trait)] pub trait RecordReader { /// The claim at `rkey`, or `None` when the repository holds none. async fn read_operator_record( &self, operator_did: &str, rkey: &str, ) -> Result, WriteError>; } impl RecordReader for jacquard_oauth::client::OAuthSession where T: jacquard_oauth::resolver::OAuthResolver + jacquard_oauth::dpop::DpopExt + jacquard_common::xrpc::XrpcExt + Send + Sync + 'static, S: jacquard_oauth::authstore::ClientAuthStore + Send + Sync + 'static, W: Send + Sync, { async fn read_operator_record( &self, operator_did: &str, rkey: &str, ) -> Result, WriteError> { let trouble = |detail: String| WriteError::Read { rkey: rkey.to_owned(), detail, }; let response = self .send(GetRecordParams { repo: operator_did.to_owned(), collection: crate::operate::record::COLLECTION.to_owned(), rkey: rkey.to_owned(), }) .await .map_err(|error| trouble(error.to_string()))?; if response.status().is_success() { let output: GetRecordOutput = serde_json::from_slice(response.buffer()) .map_err(|error| trouble(error.to_string()))?; return Ok(OperatorClaim::from_json(&output.value)); } // No record is a 400 naming `RecordNotFound`, or a 404 from a // server that prefers it; anything else is a repository that was // not read. let named = serde_json::from_slice::(response.buffer()) .ok() .and_then(|body| body.get("error")?.as_str().map(str::to_owned)); if response.status().as_u16() == 404 || named.as_deref() == Some("RecordNotFound") { return Ok(None); } Err(trouble(format!( "{} {}", response.status(), String::from_utf8_lossy(response.buffer()).trim() ))) } } /// Writes `claim` into the operator's own repository, at the record key /// [`crate::operate::record::COLLECTION`]/`rkey`. Implemented for anything that /// can make an authenticated XRPC call -- in production, a /// `jacquard_oauth::client::OAuthSession` -- so tests can substitute a fake /// and this module never has to know how the session was authenticated. // See `crate::operate::tls::TlsProbe`'s own note on why `async fn` in this trait is // fine here: always used generically, never as a trait object. #[allow(async_fn_in_trait)] pub trait RecordWriter { /// Writes the claim at record key `rkey`, returning the `at://` URI it /// now lives at. /// /// `rkey` is the *decoded* authority of the server's own DID /// ([`didbot_identity::did::AccountDid::authority`]), which is what the /// server's operator poll reads at -- not the hostname the command was /// given. /// The two differ for a development server on a port, and the encoded /// form is not a legal atproto record key. async fn put_operator_record( &self, operator_did: &str, rkey: &str, claim: &OperatorClaim, ) -> Result; } /// Takes a `bot.did.operator` record back. An admission writes its record /// before the server creates the account and deletes it if the server /// refuses, so no record names an account that does not exist. Its own /// trait rather than a method of [`RecordWriter`] because a claim never /// deletes, and a writer that only claims need not be able to. #[allow(async_fn_in_trait)] pub trait RecordDeleter { /// Removes the record at `rkey`. async fn delete_operator_record( &self, operator_did: &str, rkey: &str, ) -> Result<(), WriteError>; } impl RecordWriter for jacquard_oauth::client::OAuthSession where T: jacquard_oauth::resolver::OAuthResolver + jacquard_oauth::dpop::DpopExt + jacquard_common::xrpc::XrpcExt + Send + Sync + 'static, S: jacquard_oauth::authstore::ClientAuthStore + Send + Sync + 'static, W: Send + Sync, { async fn put_operator_record( &self, operator_did: &str, rkey: &str, claim: &OperatorClaim, ) -> Result { let record = crate::operate::record::to_record_value(claim)?; let input = PutRecordInput { repo: operator_did.to_owned(), collection: crate::operate::record::COLLECTION.to_owned(), rkey: rkey.to_owned(), record, }; let response = self .send(input) .await .map_err(|error| WriteError::Refused { detail: error.to_string(), })?; let output = response .parse::() .map_err(|error| WriteError::Refused { detail: error.to_string(), })?; Ok(output.uri) } } impl RecordDeleter for jacquard_oauth::client::OAuthSession where T: jacquard_oauth::resolver::OAuthResolver + jacquard_oauth::dpop::DpopExt + jacquard_common::xrpc::XrpcExt + Send + Sync + 'static, S: jacquard_oauth::authstore::ClientAuthStore + Send + Sync + 'static, W: Send + Sync, { async fn delete_operator_record( &self, operator_did: &str, rkey: &str, ) -> Result<(), WriteError> { let input = DeleteRecordInput { repo: operator_did.to_owned(), collection: crate::operate::record::COLLECTION.to_owned(), rkey: rkey.to_owned(), }; let refused = |error: &dyn std::fmt::Display| WriteError::Refused { detail: error.to_string(), }; let response = self.send(input).await.map_err(|error| refused(&error))?; response .parse::() .map_err(|error| refused(&error))?; Ok(()) } }