diff --git a/knot2/crates/knot-atproto/Cargo.toml b/knot2/crates/knot-atproto/Cargo.toml index efe199c5c..9f967b783 100644 --- a/knot2/crates/knot-atproto/Cargo.toml +++ b/knot2/crates/knot-atproto/Cargo.toml @@ -9,6 +9,7 @@ license.workspace = true knot-types = { workspace = true } knot-runtime = { workspace = true } knot-cache = { workspace = true } +cid = { workspace = true } futures = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } @@ -21,9 +22,9 @@ url = { workspace = true } base64 = { workspace = true } bytes = { workspace = true } http = { workspace = true } +k256 = { workspace = true } [dev-dependencies] knot-lexicons = { workspace = true } tokio = { workspace = true } -k256 = { workspace = true } proptest = { workspace = true } diff --git a/knot2/crates/knot-atproto/src/identity.rs b/knot2/crates/knot-atproto/src/identity.rs index 30ec6c122..81e6094bc 100644 --- a/knot2/crates/knot-atproto/src/identity.rs +++ b/knot2/crates/knot-atproto/src/identity.rs @@ -8,9 +8,10 @@ use serde::Serialize; use sha2::{Digest, Sha256}; const PLC_OP_TYPE: &str = "plc_operation"; -const ATPROTO_METHOD: &str = "atproto"; -const KNOT_HOME_SERVICE: &str = "tangled_knot"; -const KNOT_HOME_TYPE: &str = "TangledKnot"; +pub(crate) const ATPROTO_METHOD: &str = "atproto"; +pub(crate) const KNOT_HOME_SERVICE: &str = "tangled_knot"; +pub(crate) const KNOT_HOME_TYPE: &str = "TangledKnot"; +const PDS_SERVICE_ID: &str = crate::plc::ATPROTO_PDS_SERVICE; const DID_PLC_PREFIX: &str = "did:plc:"; const DID_SUFFIX_LEN: usize = 24; const MINT_NONCE_LEN: usize = 16; @@ -19,6 +20,8 @@ const MINT_NONCE_LEN: usize = 16; pub enum IdentityError { #[error("plc operation couldn't be encoded: {0}")] Encode(String), + #[error("plc operation shape unsupported: {0}")] + Op(String), #[error("derived did:plc isn't valid DID: {0}")] Did(#[from] knot_types::ParseError), } @@ -47,16 +50,16 @@ struct PlcOperation { pub struct PreparedRepoDid { pub did: RepoDid, - operation_json: Vec, + operation: serde_json::Value, } impl PreparedRepoDid { - pub fn operation_json(&self) -> &[u8] { - &self.operation_json + pub fn operation(&self) -> &serde_json::Value { + &self.operation } } -fn did_key(public: &PublicKeyBytes) -> String { +pub(crate) fn did_key_for(public: &PublicKeyBytes) -> String { format!("did:key:{}", multikey_secp256k1(public)) } @@ -104,7 +107,7 @@ pub fn prepare_repo_did( let base = knot_service_url.as_str(); let mut operation = PlcOperation { op_type: PLC_OP_TYPE, - rotation_keys: vec![did_key(&signer.public_key())], + rotation_keys: vec![did_key_for(&signer.public_key())], verification_methods: BTreeMap::new(), also_known_as: Vec::new(), services: BTreeMap::from([( @@ -123,12 +126,9 @@ pub fn prepare_repo_did( let signed = encode_cbor(&operation)?; let did = derive_did_plc(&signed)?; - let operation_json = - serde_json::to_vec(&operation).map_err(|error| IdentityError::Encode(error.to_string()))?; - Ok(PreparedRepoDid { - did, - operation_json, - }) + let operation = serde_json::to_value(operation) + .map_err(|error| IdentityError::Encode(error.to_string()))?; + Ok(PreparedRepoDid { did, operation }) } fn derive_did_plc(signed_cbor: &[u8]) -> Result { @@ -169,6 +169,10 @@ fn did_web_document( "id": format!("#{KNOT_HOME_SERVICE}"), "type": KNOT_HOME_TYPE, "serviceEndpoint": service_url + }, { + "id": format!("#{PDS_SERVICE_ID}"), + "type": crate::plc::PDS_SERVICE_TYPE, + "serviceEndpoint": service_url }] }) } @@ -195,7 +199,7 @@ mod tests { ) .unwrap(); assert_eq!(first.did, again.did); - assert_eq!(first.operation_json(), again.operation_json()); + assert_eq!(first.operation(), again.operation()); let slashed = prepare_repo_did( &key, @@ -252,8 +256,7 @@ mod tests { &repo_nonce(106), ) .unwrap(); - let operation: serde_json::Value = - serde_json::from_slice(prepared.operation_json()).unwrap(); + let operation = prepared.operation(); let signature_b64 = operation["sig"].as_str().unwrap(); let signature = @@ -278,8 +281,7 @@ mod tests { &repo_nonce(107), ) .unwrap(); - let operation: serde_json::Value = - serde_json::from_slice(prepared.operation_json()).unwrap(); + let operation = prepared.operation(); assert_eq!(operation["type"], "plc_operation"); assert_eq!(operation["prev"], serde_json::Value::Null); @@ -300,7 +302,7 @@ mod tests { serde_json::json!({}), "repos have no atproto signing key" ); - let expected_key = did_key(&key.public_key()); + let expected_key = did_key_for(&key.public_key()); assert_eq!(operation["rotationKeys"][0], expected_key); assert_eq!(operation["rotationKeys"].as_array().unwrap().len(), 1); } @@ -333,5 +335,113 @@ mod tests { }), "knot self-declares its tangled_knot service at the knot root" ); + assert_eq!( + document["service"][1], + serde_json::json!({ + "id": "#atproto_pds", + "type": crate::plc::PDS_SERVICE_TYPE, + "serviceEndpoint": "https://knot.oyster.cafe" + }), + "knot publishes itself as its own PDS beside the service entry" + ); + } + + #[test] + fn the_update_operation_reaches_the_desired_shape_and_chains_prev() { + let key = runtime_signer(8); + let target = crate::plc::RepoTarget::new( + &key.public_key(), + &KnotServiceUrl::new("https://knot.oyster.cafe").unwrap(), + ); + let genesis = prepare_repo_did( + &key, + &KnotServiceUrl::new("https://knot.oyster.cafe").unwrap(), + &repo_nonce(108), + ) + .unwrap(); + let last = genesis.operation(); + + assert!( + !target.satisfied_by(last).unwrap(), + "a fresh genesis op isn't yet in the desired account shape" + ); + let update = crate::plc::update_operation(&key, last, &target).unwrap(); + assert_eq!( + update["rotationKeys"], + serde_json::json!([did_key_for(&key.public_key())]), + "the update pins rotation to the knot key" + ); + assert_eq!( + update["verificationMethods"]["atproto"], + did_key_for(&key.public_key()) + ); + let pds = &update["services"]["atproto_pds"]; + assert_eq!(pds["type"], crate::plc::PDS_SERVICE_TYPE); + assert_eq!(pds["endpoint"], "https://knot.oyster.cafe"); + assert!( + update["services"][KNOT_HOME_SERVICE]["endpoint"] + .as_str() + .unwrap() + .starts_with("https://knot.oyster.cafe/repo/"), + "the tangled_knot service survives the snapshot with its repo path" + ); + assert!( + update["prev"] + .as_str() + .expect("a chained update carries prev") + .starts_with("bafyre"), + "prev must be a CIDv1 dag-cbor string the way the directory serves them" + ); + assert_ne!(update["sig"], serde_json::Value::Null); + assert!( + target.satisfied_by(&update).unwrap(), + "the signed update lands in the desired shape" + ); + + let mut tampered = update.clone(); + tampered["verificationMethods"]["atproto"] = + serde_json::Value::String(did_key_for(&runtime_signer(9).public_key())); + assert!( + !target.satisfied_by(&tampered).unwrap(), + "a substituted signing key no longer satisfies" + ); + } + + #[test] + fn legacy_knot1_document_is_rewritten_to_desired_shape() { + let per_repo = runtime_signer(10); + let knot = runtime_signer(11); + let legacy = serde_json::json!({ + "type": "plc_operation", + "rotationKeys": [did_key_for(&per_repo.public_key())], + "verificationMethods": {"atproto": did_key_for(&per_repo.public_key())}, + "alsoKnownAs": [], + "services": {"atproto_pds": { + "type": crate::plc::PDS_SERVICE_TYPE, + "endpoint": "https://old.knot.example" + }}, + "prev": null, + "sig": "AAAA" + }); + let target = crate::plc::RepoTarget::new( + &knot.public_key(), + &KnotServiceUrl::new("https://new.knot.example").unwrap(), + ); + assert!(!target.satisfied_by(&legacy).unwrap()); + let replacement = crate::plc::update_operation(&per_repo, &legacy, &target).unwrap(); + assert_eq!( + replacement["rotationKeys"], + serde_json::json!([did_key_for(&knot.public_key())]), + "the replacement moves rotation to the knot key in the same operation" + ); + assert_eq!( + replacement["services"]["tangled_knot"], + serde_json::json!({ + "type": "TangledKnot", + "endpoint": "https://new.knot.example", + }), + "a knot1 document without a home service gains one pointing here" + ); + assert!(target.satisfied_by(&replacement).unwrap()); } } diff --git a/knot2/crates/knot-atproto/src/jwt.rs b/knot2/crates/knot-atproto/src/jwt.rs index 15696e2ca..99f7e8ca2 100644 --- a/knot2/crates/knot-atproto/src/jwt.rs +++ b/knot2/crates/knot-atproto/src/jwt.rs @@ -298,8 +298,16 @@ mod tests { let signing = signer(1); let repo = RepoDid::new("did:plc:3fwecdnvtcscjnrx2p4n7alz").unwrap(); let checked = |aud: &str| { - let token = mint_claims(&signing, &claims(SQUID, aud, UnixSeconds::new(1_000), METHOD)); - check_claims(&parse(&token).unwrap(), &repo, &member_method(), UnixSeconds::new(900)) + let token = mint_claims( + &signing, + &claims(SQUID, aud, UnixSeconds::new(1_000), METHOD), + ); + check_claims( + &parse(&token).unwrap(), + &repo, + &member_method(), + UnixSeconds::new(900), + ) }; checked(repo.as_str()).unwrap(); diff --git a/knot2/crates/knot-atproto/src/lib.rs b/knot2/crates/knot-atproto/src/lib.rs index 2fa2b0275..85b655a77 100644 --- a/knot2/crates/knot-atproto/src/lib.rs +++ b/knot2/crates/knot-atproto/src/lib.rs @@ -1,6 +1,7 @@ mod auth; mod identity; mod jwt; +mod plc; mod pointer; mod pubkeys; mod resolve; @@ -12,6 +13,7 @@ pub use identity::{ IdentityError, MintNonce, PreparedRepoDid, knot_did_document, prepare_repo_did, }; pub use jwt::{JwtError, JwtNonce, ServiceJwt}; +pub use plc::{PDS_SERVICE_TYPE, RepoTarget, update_operation}; pub use pointer::PointerReceipt; pub use pubkeys::{KeyParseError, parse_authorized_key}; pub use resolve::{Identity, PdsEndpoint, PlcDirectory, ResolveError}; @@ -27,10 +29,13 @@ pub mod fuzz { let did = knot_types::AccountDid::new("did:plc:nel").expect("constant test did is valid"); let _ = crate::resolve::identity_from_document(&did, data); let repo_did = knot_types::RepoDid::new("did:plc:nel").expect("constant test did is valid"); - let _ = crate::resolve::document_publishes_key( + let pds = knot_types::KnotServiceUrl::new("https://knot.oyster.cafe") + .expect("constant test url is valid"); + let _ = crate::resolve::document_publishes_account_shape( &repo_did, data, &knot_runtime::PublicKeyBytes::from_bytes(data.to_vec()), + &pds, ); } } @@ -109,6 +114,8 @@ pub enum AtprotoError { Identity(#[from] IdentityError), #[error("plc submission for {did} returned HTTP {status}")] PlcSubmit { did: RepoDid, status: HttpStatus }, + #[error("plc log fetch for {did} returned HTTP {status}")] + PlcFetch { did: RepoDid, status: HttpStatus }, #[error("putRecord for {subject} returned HTTP {status}")] PutRecord { subject: AccountDid, @@ -134,6 +141,7 @@ impl AtprotoError { AtprotoError::Resolve(error) => error.is_transient(), AtprotoError::ListRecords { status, .. } | AtprotoError::PlcSubmit { status, .. } + | AtprotoError::PlcFetch { status, .. } | AtprotoError::PutRecord { status, .. } | AtprotoError::GetRecord { status, .. } => status.is_transient(), _ => false, @@ -153,6 +161,7 @@ impl AtprotoError { | AtprotoError::ReplayShareExhausted { .. } | AtprotoError::Identity(_) | AtprotoError::PlcSubmit { .. } + | AtprotoError::PlcFetch { .. } | AtprotoError::PutRecord { .. } | AtprotoError::GetRecord { .. } | AtprotoError::PointerEncode(_) @@ -544,7 +553,9 @@ impl Atproto { rkey: &str, ) -> Result { let cached = self.identity(owner).await?; - let read = self.read_record(&cached.value, owner, collection, rkey).await?; + let read = self + .read_record(&cached.value, owner, collection, rkey) + .await?; match read { RecordPresence::Absent if !cached.fresh => { let current = self.fetch_identity(owner).await?; @@ -639,26 +650,48 @@ impl Atproto { pub async fn submit_plc_operation( &self, - prepared: &PreparedRepoDid, + did: &RepoDid, + operation: &serde_json::Value, ) -> Result<(), AtprotoError> { - let account = AccountDid::from(prepared.did.clone()); - let url = resolve::document_url(&account, &self.plc_directory)?; + let url = self.document_url(&AccountDid::from(did.clone()))?; + let body = serde_json::to_vec(operation) + .map_err(|error| AtprotoError::PointerEncode(error.to_string()))?; resolve::guard_fetch_url(&url)?; - let request = json_post( - url, - bytes::Bytes::copy_from_slice(prepared.operation_json()), - ); + let request = json_post(url, bytes::Bytes::from(body)); let response = self.http.execute(request).await?; if response.status.is_success() { Ok(()) } else { Err(AtprotoError::PlcSubmit { - did: prepared.did.clone(), + did: did.clone(), status: HttpStatus::from(response.status), }) } } + pub async fn last_plc_operation( + &self, + did: &RepoDid, + ) -> Result, AtprotoError> { + let mut url = self.document_url(&AccountDid::from(did.clone()))?; + url.set_path(&format!("{}/log/last", url.path().trim_end_matches('/'))); + resolve::guard_fetch_url(&url)?; + let response = self.http.execute(HttpRequest::get(url)).await?; + let status = HttpStatus::from(response.status); + if status.get() == 404 { + return Ok(None); + } + if !response.status.is_success() { + return Err(AtprotoError::PlcFetch { + did: did.clone(), + status, + }); + } + serde_json::from_slice(&response.body) + .map(Some) + .map_err(|error| AtprotoError::MalformedRecords(error.to_string())) + } + pub async fn publish_pointer( &self, authorizer: &dyn PointerAuthorizer, @@ -693,10 +726,11 @@ impl Atproto { pointer::receipt_from_response(&response.body) } - pub async fn verify_did_web_publishes_key( + pub async fn verify_did_web_account_document( &self, did: &RepoDid, expected: &PublicKeyBytes, + expected_pds: &knot_types::KnotServiceUrl, ) -> Result<(), AtprotoError> { let url = resolve::web_document_url_for(did)?; resolve::guard_fetch_url(&url)?; @@ -711,7 +745,8 @@ impl Atproto { } .into()); } - resolve::document_publishes_key(did, &response.body, expected).map_err(Into::into) + resolve::document_publishes_account_shape(did, &response.body, expected, expected_pds) + .map_err(Into::into) } fn record_jti( @@ -2123,6 +2158,7 @@ mod tests { &repo_nonce(111), ) .unwrap(); + let operation = prepared.operation(); let expected_path = format!("/{}", prepared.did.as_str()); let http = FakeHttp::new(move |request| { assert_eq!(request.method, http::Method::POST); @@ -2131,7 +2167,7 @@ mod tests { Ok(ok(Bytes::new())) }); Atproto::new(http, clock(), knot_did(KNOT), plc()) - .submit_plc_operation(&prepared) + .submit_plc_operation(&prepared.did.clone(), operation) .await .unwrap(); @@ -2141,17 +2177,19 @@ mod tests { &repo_nonce(112), ) .unwrap(); + let rejected_op = rejected.operation(); let http = FakeHttp::new(move |_| Ok(status(StatusCode::BAD_REQUEST, Bytes::new()))); let atproto = Atproto::new(http, clock(), knot_did(KNOT), plc()); assert!(matches!( - atproto.submit_plc_operation(&rejected).await, + atproto.submit_plc_operation(&rejected.did.clone(), rejected_op).await, Err(AtprotoError::PlcSubmit { status, .. }) if status.get() == 400 )); } #[tokio::test] - async fn a_did_web_document_verification_covers_present_absent_and_missing() { + async fn a_did_web_document_verification_demands_atproto_and_pds_entries() { let signing = signer(9); + let expected_pds = knot_types::KnotServiceUrl::new("https://knot.oyster.cafe").unwrap(); let doc = signing.clone(); let http = FakeHttp::new(move |request| { assert_eq!( @@ -2162,15 +2200,16 @@ mod tests { id: "did:web:limpet.olaren.dev", signing: &doc, handle: "nel.pet", - pds: "https://pds.oyster.cafe", + pds: "https://knot.oyster.cafe", method: MethodKind::Multikey, }))) }); let atproto = Atproto::new(http, clock(), knot_did(KNOT), plc()); atproto - .verify_did_web_publishes_key( + .verify_did_web_account_document( &repo_did("did:web:limpet.olaren.dev"), &PublicKeyBytes::from_bytes(sec1(&signing)), + &expected_pds, ) .await .unwrap(); @@ -2187,9 +2226,10 @@ mod tests { }); let atproto = Atproto::new(http, clock(), knot_did(KNOT), plc()); let error = atproto - .verify_did_web_publishes_key( + .verify_did_web_account_document( &repo_did("did:web:limpet.olaren.dev"), &PublicKeyBytes::from_bytes(sec1(&signer(3))), + &expected_pds, ) .await .unwrap_err(); @@ -2198,12 +2238,38 @@ mod tests { AtprotoError::Resolve(ResolveError::ExpectedKeyAbsent { .. }) )); + let stray_pds = knot_types::KnotServiceUrl::new("https://elsewhere.example").unwrap(); + let doc = signing.clone(); + let http = FakeHttp::new(move |_| { + Ok(ok(did_doc(DocSpec { + id: "did:web:limpet.olaren.dev", + signing: &doc, + handle: "nel.pet", + pds: "https://pds.oyster.cafe", + method: MethodKind::Multikey, + }))) + }); + let atproto = Atproto::new(http, clock(), knot_did(KNOT), plc()); + let error = atproto + .verify_did_web_account_document( + &repo_did("did:web:limpet.olaren.dev"), + &PublicKeyBytes::from_bytes(sec1(&signing)), + &stray_pds, + ) + .await + .unwrap_err(); + assert!(matches!( + error, + AtprotoError::Resolve(ResolveError::BadPds { .. }) + )); + let http = FakeHttp::new(|_| Ok(status(StatusCode::NOT_FOUND, Bytes::new()))); let atproto = Atproto::new(http, clock(), knot_did(KNOT), plc()); let error = atproto - .verify_did_web_publishes_key( + .verify_did_web_account_document( &repo_did("did:web:limpet.olaren.dev"), &PublicKeyBytes::from_bytes(vec![1, 2, 3]), + &expected_pds, ) .await .unwrap_err(); diff --git a/knot2/crates/knot-atproto/src/plc.rs b/knot2/crates/knot-atproto/src/plc.rs new file mode 100644 index 000000000..7bf667ca5 --- /dev/null +++ b/knot2/crates/knot-atproto/src/plc.rs @@ -0,0 +1,308 @@ +use base64::Engine; +use base64::engine::general_purpose::URL_SAFE_NO_PAD; +use serde::{Deserialize, Serialize}; +use serde_json::{Map, Value, json}; +use sha2::Digest; + +use knot_runtime::PublicKeyBytes; +use knot_runtime::Signer; +use knot_types::KnotServiceUrl; + +use crate::identity::IdentityError; +use crate::identity::KNOT_HOME_SERVICE; +use crate::identity::KNOT_HOME_TYPE; +use crate::identity::did_key_for; + +pub const ATPROTO_PDS_SERVICE: &str = "atproto_pds"; +pub const PDS_SERVICE_TYPE: &str = "AtprotoPersonalDataServer"; + +const PLC_OP_TYPE: &str = "plc_operation"; +const CID_DAG_CBOR_CODEC: u64 = 0x71; +const MULTIHASH_SHA2_256_CODE: u64 = 0x12; +const PLC_MAX_OPERATION_BYTES: usize = 7500; + +pub struct RepoTarget { + knot_did_key: String, + pds_endpoint: String, + home_origin: String, +} + +impl RepoTarget { + pub fn new(knot_public: &PublicKeyBytes, service_url: &KnotServiceUrl) -> Self { + let pds_endpoint = service_url.as_str().trim_end_matches('/').to_owned(); + let home_origin = url::Url::parse(service_url.as_str()) + .map(|url| url.origin().ascii_serialization()) + .unwrap_or_else(|_| pds_endpoint.clone()); + Self { + knot_did_key: did_key_for(knot_public), + pds_endpoint, + home_origin, + } + } + + pub fn knot_did_key(&self) -> &str { + &self.knot_did_key + } + + pub fn satisfied_by(&self, last_op: &Value) -> Result { + let current = Operation::parse(last_op).ok_or_else(unsupported_shape)?; + let desired = desired_operation(self, ¤t); + Ok(current.rotation_keys == desired.rotation_keys + && current.verification_methods == desired.verification_methods + && current.services == desired.services) + } +} + +fn unsupported_shape() -> IdentityError { + IdentityError::Op("previous operation isn't a supported plc_operation shape".into()) +} + +#[derive(Serialize, Deserialize)] +struct Operation { + #[serde(rename = "type")] + op_type: String, + #[serde(rename = "rotationKeys")] + rotation_keys: Vec, + #[serde(rename = "verificationMethods")] + verification_methods: Map, + #[serde(rename = "alsoKnownAs")] + also_known_as: Vec, + services: Map, +} + +impl Operation { + fn parse(op: &Value) -> Option { + let parsed: Operation = serde_json::from_value(op.clone()).ok()?; + (parsed.op_type == PLC_OP_TYPE).then_some(parsed) + } +} + +#[derive(Deserialize)] +struct ServiceEndpoint { + endpoint: String, +} + +fn points_at_origin(endpoint: &str, origin: &str) -> bool { + url::Url::parse(endpoint).is_ok_and(|parsed| parsed.origin().ascii_serialization() == origin) +} + +fn desired_operation(target: &RepoTarget, current: &Operation) -> Operation { + let home = current + .services + .get(KNOT_HOME_SERVICE) + .and_then(|entry| serde_json::from_value::(entry.clone()).ok()) + .map(|home| { + if points_at_origin(&home.endpoint, &target.home_origin) { + json!({ + "type": KNOT_HOME_TYPE, + "endpoint": home.endpoint, + }) + } else { + fresh_home(target) + } + }) + .unwrap_or_else(|| fresh_home(target)); + let mut desired = Operation { + op_type: PLC_OP_TYPE.to_owned(), + rotation_keys: vec![target.knot_did_key.clone()], + verification_methods: current.verification_methods.clone(), + also_known_as: current.also_known_as.clone(), + services: current.services.clone(), + }; + desired.verification_methods.insert( + crate::identity::ATPROTO_METHOD.to_owned(), + json!(target.knot_did_key), + ); + desired.services.insert( + ATPROTO_PDS_SERVICE.to_owned(), + json!({ "type": PDS_SERVICE_TYPE, "endpoint": target.pds_endpoint }), + ); + desired.services.insert(KNOT_HOME_SERVICE.to_owned(), home); + desired +} + +fn fresh_home(target: &RepoTarget) -> Value { + json!({ "type": KNOT_HOME_TYPE, "endpoint": target.pds_endpoint }) +} + +fn encode_unsigned_cbor(op: &Value) -> Result, IdentityError> { + let mut unsigned = op.clone(); + if let Some(object) = unsigned.as_object_mut() { + object.remove("sig"); + } + serde_ipld_dagcbor::to_vec(&unsigned).map_err(|error| IdentityError::Encode(error.to_string())) +} + +fn operation_cid(op: &Value) -> Result { + let bytes = + serde_ipld_dagcbor::to_vec(op).map_err(|error| IdentityError::Encode(error.to_string()))?; + let digest = sha2::Sha256::digest(&bytes); + let multihash = cid::multihash::Multihash::wrap(MULTIHASH_SHA2_256_CODE, &digest) + .map_err(|error| IdentityError::Encode(error.to_string()))?; + Ok(cid::Cid::new_v1(CID_DAG_CBOR_CODEC, multihash).to_string()) +} + +fn signed(mut op: Value, signer: &dyn Signer) -> Result { + let unsigned = encode_unsigned_cbor(&op)?; + let sig = URL_SAFE_NO_PAD.encode(signer.sign(&unsigned).as_bytes()); + match op.as_object_mut() { + Some(object) => { + object.insert("sig".into(), Value::String(sig)); + } + None => return Err(IdentityError::Op("operation isn't a JSON object".into())), + } + let bytes = serde_ipld_dagcbor::to_vec(&op) + .map_err(|error| IdentityError::Encode(error.to_string()))?; + if bytes.len() > PLC_MAX_OPERATION_BYTES { + return Err(IdentityError::Op(format!( + "encoded operation is {} bytes over the {PLC_MAX_OPERATION_BYTES}-byte limit", + bytes.len() + ))); + } + Ok(op) +} + +pub fn update_operation( + signer: &dyn Signer, + last_op: &Value, + target: &RepoTarget, +) -> Result { + let current = Operation::parse(last_op).ok_or_else(unsupported_shape)?; + let desired = desired_operation(target, ¤t); + let mut op = + serde_json::to_value(desired).map_err(|error| IdentityError::Encode(error.to_string()))?; + op["prev"] = json!(operation_cid(last_op)?); + signed(op, signer) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + fn test_signer(seed: u8) -> knot_runtime::K256Signer { + knot_runtime::K256Signer::from_slice(&[seed; 32]).expect("a fixed scalar signs") + } + + fn cafe_url() -> KnotServiceUrl { + KnotServiceUrl::new("https://knot.oyster.cafe").expect("constant test url is valid") + } + + #[test] + fn a_second_rotation_key_defeats_satisfaction_until_collapsed() { + let knot = test_signer(3); + let stranger = test_signer(4); + let target = RepoTarget::new(&knot.public_key(), &cafe_url()); + let last = json!({ + "type": PLC_OP_TYPE, + "rotationKeys": [did_key_for(&knot.public_key()), did_key_for(&stranger.public_key())], + "verificationMethods": {"atproto": did_key_for(&knot.public_key())}, + "alsoKnownAs": [], + "services": {}, + }); + assert!( + !target.satisfied_by(&last).unwrap(), + "a document keeping a second rotation authority isn't converged" + ); + let collapsed = update_operation(&stranger, &last, &target).unwrap(); + assert_eq!( + collapsed["rotationKeys"], + json!([did_key_for(&knot.public_key())]), + "the update collapses rotation onto the knot key alone" + ); + assert!(target.satisfied_by(&collapsed).unwrap()); + } + + #[test] + fn an_oversized_operation_refuses_to_encode() { + let signer = test_signer(6); + let target = RepoTarget::new(&signer.public_key(), &cafe_url()); + let filler = "p".repeat(9000); + let last = json!({ + "type": PLC_OP_TYPE, + "rotationKeys": [did_key_for(&signer.public_key())], + "verificationMethods": {"atproto": did_key_for(&signer.public_key())}, + "alsoKnownAs": [format!("at://{filler}")], + "services": {"atproto_pds": { + "type": PDS_SERVICE_TYPE, + "endpoint": format!("https://pad.example/{filler}"), + }}, + }); + let error = update_operation(&signer, &last, &target).unwrap_err(); + assert!(matches!(&error, IdentityError::Op(message) if message.contains("byte limit"))); + } + + #[test] + fn operation_cids_reproduce_the_directory_over_a_captured_operation() { + let captured = json!({ + "sig": "ZznbxHinpBI3NgEYWzXUXLA65s2U1ezJooreZlscHYQRfd4EMlijqhpikGMabs-81Tsy4Mt9Iscpmk7Uz13aAg", + "prev": null, + "type": PLC_OP_TYPE, + "services": {}, + "alsoKnownAs": ["at://op0"], + "rotationKeys": [ + "did:key:zQ3shmWf4f6ZwzNyjUDYw4oFQjgHWZoYZDjJYtz75YfYYrphB", + "did:key:zQ3shr2PwdoF6kbuYYymyZq6YbGWShiTXAkdMioCPdSdPL9NV" + ], + "verificationMethods": {}, + }); + let stored_cid = "bafyreihqa4o4gsyai22ydms6j4cua6yr6tuc7trq2temvnycadk2daniru"; + assert_eq!( + operation_cid(&captured).unwrap(), + stored_cid, + "the directory hashes signed operations including sig" + ); + } + + #[test] + fn a_foreign_tangled_knot_origin_is_replaced_instead_of_kept() { + let knot = test_signer(7); + let target = RepoTarget::new(&knot.public_key(), &cafe_url()); + let last = json!({ + "type": PLC_OP_TYPE, + "rotationKeys": [did_key_for(&knot.public_key())], + "verificationMethods": {"atproto": did_key_for(&knot.public_key())}, + "alsoKnownAs": [], + "services": {"tangled_knot": { + "type": KNOT_HOME_TYPE, + "endpoint": "https://knot.oyster.cafe.evil.example/repo/lure", + }}, + }); + assert!( + !target.satisfied_by(&last).unwrap(), + "a home service pointing at a stranger's origin isn't converged" + ); + let replaced = update_operation(&knot, &last, &target).unwrap(); + assert_eq!( + replaced["services"]["tangled_knot"]["endpoint"], + "https://knot.oyster.cafe" + ); + assert!(target.satisfied_by(&replaced).unwrap()); + } + + #[test] + fn a_same_origin_home_keeps_its_path_and_gains_the_canonical_type() { + let knot = test_signer(8); + let target = RepoTarget::new(&knot.public_key(), &cafe_url()); + let last = json!({ + "type": PLC_OP_TYPE, + "rotationKeys": [did_key_for(&knot.public_key())], + "verificationMethods": {"atproto": did_key_for(&knot.public_key())}, + "alsoKnownAs": [], + "services": {"tangled_knot": { + "type": "TangledKnotV1", + "endpoint": "https://knot.oyster.cafe/repo/tagged", + }}, + }); + let update = update_operation(&knot, &last, &target).unwrap(); + assert_eq!( + update["services"]["tangled_knot"], + json!({ + "type": "TangledKnot", + "endpoint": "https://knot.oyster.cafe/repo/tagged", + }), + "same-origin paths survive; only the type moves to the canonical one" + ); + assert!(target.satisfied_by(&update).unwrap()); + } +} diff --git a/knot2/crates/knot-atproto/src/resolve.rs b/knot2/crates/knot-atproto/src/resolve.rs index fca998a17..0b22f87b1 100644 --- a/knot2/crates/knot-atproto/src/resolve.rs +++ b/knot2/crates/knot-atproto/src/resolve.rs @@ -250,10 +250,11 @@ pub(crate) fn web_document_url_for(did: &RepoDid) -> Result { }) } -pub(crate) fn document_publishes_key( +pub(crate) fn document_publishes_account_shape( did: &RepoDid, body: &[u8], expected: &PublicKeyBytes, + expected_pds: &knot_types::KnotServiceUrl, ) -> Result<(), ResolveError> { let document: DidDocument = serde_json::from_slice(body).map_err(|error| ResolveError::Malformed(error.to_string()))?; @@ -267,16 +268,25 @@ pub(crate) fn document_publishes_key( .verification_method .as_ref() .ok_or(ResolveError::MissingSigningKey)?; - let published = methods - .iter() - .filter_map(|method| method_key(method).ok()) - .any(|key| { - matches!(key.codec, KeyCodec::Secp256k1) && key.bytes.as_ref() == expected.as_bytes() - }); - if published { - Ok(()) - } else { - Err(ResolveError::ExpectedKeyAbsent { did: did.clone() }) + let published = methods.iter().any(|method| { + method.id.as_str().ends_with("#atproto") + && method_key(method).is_ok_and(|key| { + matches!(key.codec, KeyCodec::Secp256k1) + && key.bytes.as_ref() == expected.as_bytes() + }) + }); + if !published { + return Err(ResolveError::ExpectedKeyAbsent { did: did.clone() }); + } + let endpoint = document + .pds_endpoint() + .ok_or(ResolveError::MissingPds)? + .as_ref() + .trim_end_matches('/') + .to_owned(); + match endpoint == expected_pds.as_str().trim_end_matches('/') { + true => Ok(()), + false => Err(ResolveError::BadPds { value: endpoint }), } } @@ -461,37 +471,24 @@ mod tests { } #[test] - fn document_publishes_key_matches_only_on_codec_and_bytes() { + fn an_account_shape_rejects_the_same_bytes_under_a_foreign_codec() { let published = did_doc(DocSpec { - id: SQUID, - signing: &signer(5), - handle: "nel.pet", - pds: "https://pds.oyster.cafe", - method: MethodKind::LegacyK256, - }); - document_publishes_key( - &RepoDid::new(SQUID).unwrap(), - &published, - &PublicKeyBytes::from_bytes(sec1(&signer(5))), - ) - .unwrap(); - - let foreign = did_doc(DocSpec { id: SQUID, signing: &signer(5), handle: "nel.pet", pds: "https://pds.oyster.cafe", method: MethodKind::LegacyP256, }); - let error = document_publishes_key( + let error = crate::resolve::document_publishes_account_shape( &RepoDid::new(SQUID).unwrap(), - &foreign, + &published, &PublicKeyBytes::from_bytes(sec1(&signer(5))), + &knot_types::KnotServiceUrl::new("https://pds.oyster.cafe").unwrap(), ) .unwrap_err(); assert!( matches!(error, ResolveError::ExpectedKeyAbsent { .. }), - "the same bytes under a foreign codec mustn't satisfy publishes_key, got {error:?}" + "the same bytes under a P-256 codec mustn't satisfy the account shape, got {error:?}" ); } diff --git a/knot2/crates/knot-cob/src/backend.rs b/knot2/crates/knot-cob/src/backend.rs index 809758da9..1a437ff43 100644 --- a/knot2/crates/knot-cob/src/backend.rs +++ b/knot2/crates/knot-cob/src/backend.rs @@ -17,6 +17,7 @@ const TYPE_HEADER: &str = "cob-type"; const SIG_HEADER: &str = "cob-sig"; const AUTHOR_HEADER: &str = "cob-author"; const PAYLOAD_BLOB: &str = "payload"; +const RECORD_BLOB: &str = "record"; pub(crate) const MAX_GRAPH_CHANGES: usize = 100_000; const REBUILD_CHANGE_BYTES: u64 = 2560; const REBUILD_FOLD_DIVISOR: u64 = 4; @@ -87,6 +88,7 @@ pub(crate) fn list_objects(repo: &Repo, type_name: &TypeName) -> Result, signer: &dyn Signer, timestamp: UnixSeconds, +) -> Result { + write_change_with_record( + home, repo, type_name, payload, None, parents, object, signer, timestamp, + ) +} + +#[allow(clippy::too_many_arguments)] +pub(crate) fn write_change_with_record( + home: &CobHome, + repo: &Repo, + type_name: &TypeName, + payload: &[u8], + record: Option<&[u8]>, + parents: &[ChangeId], + object: Option, + signer: &dyn Signer, + timestamp: UnixSeconds, ) -> Result { let git = repo.git(); let payload_oid = git .write_blob(payload) .map_err(|error| CobError::Write(error.to_string()))? .detach(); + let record_oid = record + .map(|bytes| { + git.write_blob(bytes) + .map(|oid| oid.detach()) + .map_err(|error| CobError::Write(error.to_string())) + }) + .transpose()?; let revision = git - .write_object(build_tree(payload_oid)) + .write_object(build_tree(payload_oid, record_oid)) .map_err(|error| CobError::Write(error.to_string()))? .detach(); let author = ActorId::from_secp256k1(signer.public_key().as_bytes()); @@ -315,7 +341,7 @@ pub(crate) fn read_change(repo: &Repo, id: ChangeId) -> Result .and_then(|raw| { knot_types::decode_hex(raw).ok_or_else(|| malformed("cob-sig isn't valid hex".into())) })?; - let payload = read_payload(repo, revision)?; + let (payload, record) = read_tree_blobs(repo, revision)?; Ok(Change { id, revision, @@ -324,11 +350,12 @@ pub(crate) fn read_change(repo: &Repo, id: ChangeId) -> Result author, signature: Signature::from_bytes(signature), payload: Payload::new(payload), + record: record.map(Payload::new), timestamp, }) } -fn read_payload(repo: &Repo, revision: Oid) -> Result, CobError> { +fn read_tree_blobs(repo: &Repo, revision: Oid) -> Result<(Vec, Option>), CobError> { let malformed = |reason: String| CobError::MalformedChange { oid: revision, reason, @@ -343,15 +370,19 @@ fn read_payload(repo: &Repo, revision: Oid) -> Result, CobError> { .data; let tree = gix::objs::TreeRef::from_bytes(&data, repo.git().object_hash()) .map_err(|error| malformed(error.to_string()))?; - let payload_oid = - entry_oid(&tree, PAYLOAD_BLOB).ok_or_else(|| malformed("missing payload blob".into()))?; - let payload = repo - .git() - .find_object(payload_oid) - .map_err(|error| malformed(error.to_string()))? - .detach() - .data; - Ok(payload) + let read = |name: &str| { + entry_oid(&tree, name).map(|oid| { + repo.git() + .find_object(oid) + .map(|object| object.detach().data) + .map_err(|error| malformed(error.to_string())) + }) + }; + let payload = read(PAYLOAD_BLOB) + .transpose()? + .ok_or_else(|| malformed("missing payload blob".into()))?; + let record = read(RECORD_BLOB).transpose()?; + Ok((payload, record)) } fn entry_oid(tree: &gix::objs::TreeRef<'_>, name: &str) -> Option { @@ -361,9 +392,15 @@ fn entry_oid(tree: &gix::objs::TreeRef<'_>, name: &str) -> Option .map(|entry| entry.oid.to_owned()) } -fn build_tree(payload_oid: gix::ObjectId) -> gix::objs::Tree { +fn build_tree(payload_oid: gix::ObjectId, record_oid: Option) -> gix::objs::Tree { gix::objs::Tree { - entries: vec![blob_entry(PAYLOAD_BLOB, payload_oid)], + entries: [ + Some(blob_entry(PAYLOAD_BLOB, payload_oid)), + record_oid.map(|oid| blob_entry(RECORD_BLOB, oid)), + ] + .into_iter() + .flatten() + .collect(), } } diff --git a/knot2/crates/knot-cob/src/change.rs b/knot2/crates/knot-cob/src/change.rs index 1e6a27b68..85c5a0b20 100644 --- a/knot2/crates/knot-cob/src/change.rs +++ b/knot2/crates/knot-cob/src/change.rs @@ -160,6 +160,7 @@ pub struct Change { pub author: ActorId, pub signature: Signature, pub payload: Payload, + pub record: Option, pub timestamp: UnixSeconds, } @@ -226,6 +227,7 @@ mod tests { author, signature: signer.sign(&bytes), payload: Payload::new(Vec::new()), + record: None, timestamp, } } @@ -349,6 +351,7 @@ mod tests { author, signature: signer.sign(&bytes), payload: Payload::new(Vec::new()), + record: None, timestamp, }; assert!(bound.verify(&cob_home(), &bound.author, Some(home))); diff --git a/knot2/crates/knot-cob/src/lib.rs b/knot2/crates/knot-cob/src/lib.rs index 4ab638753..578df86aa 100644 --- a/knot2/crates/knot-cob/src/lib.rs +++ b/knot2/crates/knot-cob/src/lib.rs @@ -37,6 +37,21 @@ pub struct Created { pub tip: ChangeId, } +#[derive(Debug, Clone)] +pub struct Materialized

{ + pub change: P, + pub record: Option>, +} + +impl

Materialized

{ + pub fn bare(change: P) -> Self { + Self { + change, + record: None, + } + } +} + #[derive(Debug)] pub struct Delta { pub changes: Vec, @@ -94,6 +109,10 @@ impl<'r> CobStore<'r> { Self { repo } } + pub fn repo(&self) -> &'r Repo { + self.repo + } + pub fn create( &self, home: &CobHome, @@ -101,20 +120,50 @@ impl<'r> CobStore<'r> { signer: &dyn Signer, timestamp: UnixSeconds, ) -> Result { - let type_name = P::type_name(); let bytes = payload.encode()?; - let tip = backend::write_change( + self.create_from_bytes(home, &P::type_name(), &bytes, None, signer, timestamp) + } + + pub fn create_materialized( + &self, + home: &CobHome, + payload: &Materialized

, + signer: &dyn Signer, + timestamp: UnixSeconds, + ) -> Result { + let bytes = payload.change.encode()?; + self.create_from_bytes( home, - self.repo, - &type_name, + &P::type_name(), &bytes, + payload.record.as_deref(), + signer, + timestamp, + ) + } + + fn create_from_bytes( + &self, + home: &CobHome, + type_name: &TypeName, + bytes: &[u8], + record: Option<&[u8]>, + signer: &dyn Signer, + timestamp: UnixSeconds, + ) -> Result { + let tip = backend::write_change_with_record( + home, + self.repo, + type_name, + bytes, + record, &[], None, signer, timestamp, )?; let object = CobId::new(tip.oid()); - let name = backend::cob_ref_name(&type_name, object)?; + let name = backend::cob_ref_name(type_name, object)?; self.repo.update_ref(&RefUpdate::Create { name, new: tip.oid(), @@ -150,7 +199,16 @@ impl<'r> CobStore<'r> { let expected = backend::resolve_tip(self.repo, &type_name, object)? .map(ChangeId::new) .ok_or(CobError::NoSuchObject(object))?; - self.chain(home, object, &type_name, expected, changes, signer) + self.chain( + home, + object, + &type_name, + expected, + changes + .into_iter() + .map(|(payload, timestamp)| (payload, None, timestamp)), + signer, + ) } #[allow(clippy::too_many_arguments)] @@ -169,7 +227,29 @@ impl<'r> CobStore<'r> { object, type_name, expected, - std::iter::once((payload, timestamp)), + std::iter::once((payload, None, timestamp)), + signer, + ) + .map(|tip| tip.expect("one change chains onto one tip")) + } + + #[allow(clippy::too_many_arguments)] + fn append_materialized( + &self, + home: &CobHome, + object: CobId, + type_name: &TypeName, + expected: ChangeId, + payload: &Materialized

, + signer: &dyn Signer, + timestamp: UnixSeconds, + ) -> Result { + self.chain( + home, + object, + type_name, + expected, + std::iter::once((&payload.change, payload.record.as_deref(), timestamp)), signer, ) .map(|tip| tip.expect("one change chains onto one tip")) @@ -181,20 +261,21 @@ impl<'r> CobStore<'r> { object: CobId, type_name: &TypeName, expected: ChangeId, - changes: impl IntoIterator, + changes: impl IntoIterator, UnixSeconds)>, signer: &dyn Signer, ) -> Result, CobError> { let written = match changes.into_iter().try_fold( Vec::::new(), - |mut written, (payload, timestamp)| { + |mut written, (payload, record, timestamp)| { let parent = written.last().copied().unwrap_or(expected); let write = || -> Result { let bytes = payload.encode()?; - backend::write_change( + backend::write_change_with_record( home, self.repo, type_name, &bytes, + record, &[parent], Some(object), signer, @@ -417,6 +498,24 @@ impl<'r> CobStore<'r> { timestamp: UnixSeconds, decide: impl Fn(&E::State) -> Result, D>, ) -> Result, D> + where + E: Checkpoint, + E::State: Serialize + DeserializeOwned, + D: From, + { + self.update_materialized_checkpointed::(home, object, signer, timestamp, |state| { + decide(state).map(|change| change.map(Materialized::bare)) + }) + } + + pub fn update_materialized_checkpointed( + &self, + home: &CobHome, + object: CobId, + signer: &dyn Signer, + timestamp: UnixSeconds, + decide: impl Fn(&E::State) -> Result>, D>, + ) -> Result, D> where E: Checkpoint, E::State: Serialize + DeserializeOwned, @@ -425,31 +524,36 @@ impl<'r> CobStore<'r> { let type_name = E::Change::type_name(); let attempt = || -> Result>, D> { let (state, expected, suffix) = self.checkpointed_state::(object)?; - match decide(&state)? { - None => Ok(Some(None)), - Some(change) => match self.append( - home, object, &type_name, expected, &change, signer, timestamp, - ) { - Ok(tip) => { - let stride = - checkpoint_stride(E::SNAPSHOT_STRIDE, E::checkpoint_size(&state)); - if suffix.saturating_add(1) >= stride { - let author = ActorId::from_secp256k1(signer.public_key().as_bytes()); - let folded = E::apply(state, change, &author); - if let Err(error) = self.write_checkpoint::(object, tip, &folded) { - tracing::warn!( - cob = type_name.as_str(), - object = %object.oid().to_hex(), - %error, - "checkpoint write failed" - ); - } + let Some(materialized) = decide(&state)? else { + return Ok(Some(None)); + }; + match self.append_materialized( + home, + object, + &type_name, + expected, + &materialized, + signer, + timestamp, + ) { + Ok(tip) => { + let stride = checkpoint_stride(E::SNAPSHOT_STRIDE, E::checkpoint_size(&state)); + if suffix.saturating_add(1) >= stride { + let author = ActorId::from_secp256k1(signer.public_key().as_bytes()); + let folded = E::apply(state, materialized.change, &author); + if let Err(error) = self.write_checkpoint::(object, tip, &folded) { + tracing::warn!( + cob = type_name.as_str(), + object = %object.oid().to_hex(), + %error, + "checkpoint write failed" + ); } - Ok(Some(Some(tip))) } - Err(CobError::StaleTip { .. }) => Ok(None), - Err(other) => Err(D::from(other)), - }, + Ok(Some(Some(tip))) + } + Err(CobError::StaleTip { .. }) => Ok(None), + Err(other) => Err(D::from(other)), } }; (0..MAX_CAS_RETRIES) @@ -1391,6 +1495,7 @@ mod tests { author: actor.clone(), signature: Signature::from_bytes(Vec::new()), payload: Payload::new(Vec::new()), + record: None, timestamp: UnixSeconds::new(index as i64), }; (id, change) diff --git a/knot2/crates/knot-cobs/src/grant.rs b/knot2/crates/knot-cobs/src/grant.rs index 672c715cb..9adb19c0a 100644 --- a/knot2/crates/knot-cobs/src/grant.rs +++ b/knot2/crates/knot-cobs/src/grant.rs @@ -171,22 +171,6 @@ pub enum Effect<'a> { Revoke(&'a AccountDid), } -impl Effect<'_> { - pub fn announcement(&self) -> Option { - match self { - Self::Admit(Offer::Granted, _) | Self::Accept(_) => Some(Announcement::Effective), - Self::Admit(Offer::Invited, _) => None, - Self::Revoke(_) => Some(Announcement::Cleared), - } - } -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum Announcement { - Effective, - Cleared, -} - pub trait GrantChange { fn effect(&self) -> Effect<'_>; @@ -441,7 +425,10 @@ mod tests { .for_each(|(subject, standing, since)| { let entry = roster.get(&did(subject)).unwrap(); assert_eq!( - (entry.standing(), entry.effective_since().map(EffectiveSince::seconds)), + ( + entry.standing(), + entry.effective_since().map(EffectiveSince::seconds) + ), (standing, since.map(UnixSeconds::new)), "effective_since comes from the offer, so accepting mustn't shift {subject} in \ a page a reader is already walking" @@ -496,26 +483,6 @@ mod tests { ); } - #[test] - fn an_invite_emits_no_announcement_and_the_other_three_set_a_direction() { - [ - (b'g', Some(Announcement::Effective)), - (b'a', Some(Announcement::Effective)), - (b'r', Some(Announcement::Cleared)), - (b'i', None), - ] - .into_iter() - .for_each(|(op, announcement)| { - let signed = change(op, "nel", "olaren", 1); - assert_eq!( - GrantChange::effect(&signed).announcement(), - announcement, - "spindle opens a repository off an announcement, so {signed:?} would let in an \ - account that never signed" - ); - }); - } - proptest! { #![proptest_config(ProptestConfig { cases: 256, ..ProptestConfig::default() })] diff --git a/knot2/crates/knot-cobs/src/lib.rs b/knot2/crates/knot-cobs/src/lib.rs index 3abf90593..97f013a57 100644 --- a/knot2/crates/knot-cobs/src/lib.rs +++ b/knot2/crates/knot-cobs/src/lib.rs @@ -8,8 +8,8 @@ mod registry; pub use blocklist::{Blocklist, BlocklistChange, BlocklistCob}; pub use collaborators::{Collaborators, CollaboratorsChange, CollaboratorsCob}; pub use grant::{ - Accept, Announcement, Effect, EffectiveSince, Entry, Grant, GrantChange, Invite, Offer, - Removal, Roster, Standing, + Accept, Effect, EffectiveSince, Entry, Grant, GrantChange, Invite, Offer, Removal, Roster, + Standing, }; pub use import::{ImportError, verify_cob_ref}; pub use members::{Members, MembersChange, MembersCob}; diff --git a/knot2/crates/knot-config/src/lib.rs b/knot2/crates/knot-config/src/lib.rs index 1a972a371..db993002b 100644 --- a/knot2/crates/knot-config/src/lib.rs +++ b/knot2/crates/knot-config/src/lib.rs @@ -237,6 +237,9 @@ pub struct HttpConfig { pub struct AtprotoConfig { #[config(env = "KNOT_PLC_DIRECTORY")] pub plc_directory: Url, + + #[config(env = "KNOT_REPO_SIGNING_KEYS")] + pub repo_signing_keys: Option, } #[derive(Debug, Config)] @@ -264,9 +267,6 @@ pub struct XrpcConfig { #[config(env = "KNOT_XRPC_LANGUAGES_BUDGET_MS", default = 1_000)] pub languages_budget_ms: u64, - #[config(env = "KNOT_XRPC_LANGUAGES_PUSH_BUDGET_MS", default = 2_000)] - pub languages_push_budget_ms: u64, - /// Body limit for the merge and mergeCheck procedures, whose patch payloads /// routinely exceed the general XRPC body limit. #[config(env = "KNOT_XRPC_MAX_PATCH_BYTES", default = 16_777_216)] @@ -676,6 +676,13 @@ impl KnotConfig { self.atproto.plc_directory.host().is_some(), "atproto.plc_directory must have host", ), + check( + self.atproto + .repo_signing_keys + .as_ref() + .is_none_or(|path| !path.as_os_str().is_empty()), + "atproto.repo_signing_keys must name the migration key archive", + ), check( self.xrpc.max_body_bytes > 0, "xrpc.max_body_bytes must be greater than zero", @@ -700,10 +707,6 @@ impl KnotConfig { self.xrpc.languages_budget_ms > 0, "xrpc.languages_budget_ms must be greater than zero", ), - check( - self.xrpc.languages_push_budget_ms > 0, - "xrpc.languages_push_budget_ms must be greater than zero", - ), check( self.xrpc.max_patch_bytes > 0, "xrpc.max_patch_bytes must be greater than zero", @@ -1270,6 +1273,7 @@ mod tests { }, atproto: AtprotoConfig { plc_directory: Url::parse("https://plc.nel.pet/").unwrap(), + repo_signing_keys: None, }, xrpc: XrpcConfig { max_body_bytes: 65_536, @@ -1278,7 +1282,6 @@ mod tests { tree_last_commit_budget_ms: 300, blob_last_commit_budget_ms: 2_000, languages_budget_ms: 1_000, - languages_push_budget_ms: 2_000, max_patch_bytes: 16_777_216, max_patch_decompressed_bytes: 134_217_728, preauth_burst: 20, @@ -1594,11 +1597,6 @@ mod tests { |config| config.xrpc.languages_budget_ms = 0, "languages_budget_ms", ), - ( - "zero_languages_push_budget", - |config| config.xrpc.languages_push_budget_ms = 0, - "languages_push_budget_ms", - ), ( "zero_events_replay_buffer", |config| config.xrpc.events_replay_buffer = 0, diff --git a/knot2/crates/knot-consent/src/lib.rs b/knot2/crates/knot-consent/src/lib.rs index 41f188552..381722856 100644 --- a/knot2/crates/knot-consent/src/lib.rs +++ b/knot2/crates/knot-consent/src/lib.rs @@ -80,7 +80,11 @@ mod tests { async fn reading_for(status: StatusCode, scope: Acceptance<'_>) -> Read { let atproto = Atproto::new( FakeHttp::new(move |_| { - Ok(HttpResponse { status, headers: http::HeaderMap::new(), body: Bytes::new() }) + Ok(HttpResponse { + status, + headers: http::HeaderMap::new(), + body: Bytes::new(), + }) }), ManualClock::new(UnixMicros::new(NOW)), knot(), diff --git a/knot2/crates/knot-index/src/consent.rs b/knot2/crates/knot-index/src/consent.rs index d1b95d368..aa95ac926 100644 --- a/knot2/crates/knot-index/src/consent.rs +++ b/knot2/crates/knot-index/src/consent.rs @@ -337,7 +337,9 @@ impl Consents<'_> { match scope { Acceptance::OfMembership => index.members.slot(&index.interner, subject, read), Acceptance::OfCollaboration(repo) => { - index.collaborators.slot(&index.interner, repo, subject, read) + index + .collaborators + .slot(&index.interner, repo, subject, read) } } } diff --git a/knot2/crates/knot-index/src/lib.rs b/knot2/crates/knot-index/src/lib.rs index c7c9af12b..c9cf122b3 100644 --- a/knot2/crates/knot-index/src/lib.rs +++ b/knot2/crates/knot-index/src/lib.rs @@ -667,7 +667,11 @@ impl Index { ) { (Resolved::Ready(owner), Resolved::Ready(granted), Resolved::Ready(invited)) => { Folded::Seated(Seats { - grants: owner.map(AccountDid::from).into_iter().chain(granted).collect(), + grants: owner + .map(AccountDid::from) + .into_iter() + .chain(granted) + .collect(), invited, }) } @@ -794,7 +798,12 @@ impl KeySet<'_> { .into_iter() .chain(index.due_among(&vouched, now, against)) .collect(); - let kept = pushers.iter().chain(vouched.iter()).cloned().chain(outside).collect(); + let kept = pushers + .iter() + .chain(vouched.iter()) + .cloned() + .chain(outside) + .collect(); MemberWork { unread: UnreadMembers(unread), due: StaleMembers(due), diff --git a/knot2/crates/knot-runtime/src/lib.rs b/knot2/crates/knot-runtime/src/lib.rs index 5b8aa94e4..8ce3d5d1c 100644 --- a/knot2/crates/knot-runtime/src/lib.rs +++ b/knot2/crates/knot-runtime/src/lib.rs @@ -48,6 +48,12 @@ mod contract { fn signer_contract(signer: &dyn Signer) { let signature = signer.sign(b"the message"); + let parsed = k256::ecdsa::Signature::from_slice(signature.as_bytes()) + .expect("signatures are compact"); + assert!( + parsed.normalize_s().is_none(), + "signatures are canonical low-S" + ); assert!(verify(&signer.public_key(), b"the message", &signature)); assert!(!verify( &signer.public_key(), diff --git a/knot2/crates/knot-secrets/src/lib.rs b/knot2/crates/knot-secrets/src/lib.rs index 17356674d..40f8dd20c 100644 --- a/knot2/crates/knot-secrets/src/lib.rs +++ b/knot2/crates/knot-secrets/src/lib.rs @@ -38,6 +38,8 @@ pub enum SecretsError { Malformed(String), #[error("no signing key sealed for {0}")] Missing(String), + #[error("{0} doesn't seal a valid secp256k1 scalar")] + InvalidScalar(String), #[error("signing key is already sealed for {0}")] Occupied(String), #[error("{len}-byte master key is shorter than the {MIN_MASTER_KEY_LEN}-byte minimum")] @@ -237,6 +239,29 @@ impl SealedStore { self.commit_locked(&guard, staged) } + pub fn ingest(&self, provided: I) -> Result + where + I: IntoIterator, + { + let guard = self.persist_guard(); + let mut staged = self.staged(); + let imported = provided + .into_iter() + .try_fold(0usize, |count, (id, bytes)| { + if staged.contains_key(&id) { + return Ok::(count); + } + K256Signer::from_slice(&bytes) + .map_err(|_| SecretsError::InvalidScalar(id.0.clone()))?; + staged.insert(id, SecretScalar(bytes)); + Ok(count + 1) + })?; + match imported { + 0 => Ok(0), + _ => self.commit_locked(&guard, staged).map(|()| imported), + } + } + pub fn ensure(&self, id: impl Into) -> Result { let id = id.into(); let guard = self.persist_guard(); diff --git a/knot2/crates/knot-server/src/main.rs b/knot2/crates/knot-server/src/main.rs index 2f620b38d..221bd1fbc 100644 --- a/knot2/crates/knot-server/src/main.rs +++ b/knot2/crates/knot-server/src/main.rs @@ -35,6 +35,7 @@ use tower_http::services::ServeFile; const MAINTENANCE_SHUTDOWN_DRAIN: Duration = Duration::from_secs(30); const KEYFILL_SHUTDOWN_DRAIN: Duration = Duration::from_secs(10); const EDGE_SHUTDOWN_DRAIN: Duration = Duration::from_secs(40); +const SEQ_PERSIST_INTERVAL_SECS: u64 = 10; const DEFAULT_HOMEPAGE: &str = include_str!("homepage.html"); @@ -290,6 +291,13 @@ async fn main() -> anyhow::Result<()> { ) .context("open sealed key store")?, ); + if let Some(archive) = &config.atproto.repo_signing_keys { + knot_xrpc::plc::ingest_legacy_keys(&secrets, archive) + .map_err(|error| anyhow::anyhow!(error.to_string()).context("sealing the migration per-repo signing keys")) + .map(|count| { + tracing::info!(count, path = %archive.display(), "sealed migrated per-repo signing keys") + })?; + } let knot_signing_key = secrets .ensure(&knot_did) .context("seal knot's own signing key")?; @@ -437,9 +445,6 @@ async fn main() -> anyhow::Result<()> { languages: knot_xrpc::LanguagesReadBudget::new(knot_xrpc::ReadBudget::Within( Duration::from_millis(config.xrpc.languages_budget_ms), )), - languages_push: knot_xrpc::LanguagesPushBudget::new(Duration::from_millis( - config.xrpc.languages_push_budget_ms, - )), }; let committer = knot_xrpc::Committer { name: AuthorName::new(config.git.user_name.clone()), @@ -456,7 +461,9 @@ async fn main() -> anyhow::Result<()> { knot_events::ReplayBytes::new(config.xrpc.events_replay_bytes as usize) .context("xrpc.events_replay_bytes must be greater than zero")?, ); - let events = Arc::new(knot_events::EventLog::new(SystemClock, replay_bounds)); + let firehose_floor = knot_git::Repo::open(&meta_path) + .map(|meta| knot_xrpc::firehose_floor(&meta, &layout, &index.hosted_repos())) + .context("open the meta repo to seed the firehose seq floor")?; let subscriber_gate = Arc::new(knot_events::SubscriberGate::new( knot_events::GlobalSubscriberLimit::new(config.xrpc.events_max_subscribers as usize), knot_events::PerPeerSubscriberLimit::new(config.xrpc.events_max_per_peer as usize), @@ -556,14 +563,12 @@ async fn main() -> anyhow::Result<()> { index: Arc::clone(&index), atproto: Arc::clone(&atproto), knot_actor, - events: Arc::clone(&events), hostname: hostname.clone(), appview: appview_endpoint.clone(), admins: admins.clone(), admission, max_pack_bytes: byte_limits.pack, archive_limit: byte_limits.archive, - languages_push_budget: budgets.languages_push, ci_logs: ci_logs.clone(), }) .with_maintenance(maintenance_handle.clone()) @@ -571,11 +576,6 @@ async fn main() -> anyhow::Result<()> { .with_slots(slots.clone()) .with_key_ttl(key_pace.ttl) .with_catalog(Arc::clone(&catalog)); - let ssh_state = Arc::new(match &lfs_handle { - Some(handle) => ssh_base.with_lfs(handle.clone(), lfs_max_ssh_transfers), - None => ssh_base, - }); - let xrpc_state = Arc::new(XrpcState { layout: layout.clone(), index: Arc::clone(&index), @@ -603,11 +603,27 @@ async fn main() -> anyhow::Result<()> { maintenance: maintenance_handle, appview: appview_endpoint, slots: slots.clone(), - events, - lfs: lfs_handle.map(|handle| knot_xrpc::LfsWeb::new(handle, lfs_max_http_downloads)), + firehose: Arc::new( + knot_events::EventLog::new(SystemClock, replay_bounds).with_floor(firehose_floor), + ), + lfs: lfs_handle + .clone() + .map(|handle| knot_xrpc::LfsWeb::new(handle, lfs_max_http_downloads)), catalog: Arc::clone(&catalog), }); - + let ssh_base = ssh_base.with_ref_surface(Arc::new(knot_xrpc::RefProjection { + state: Arc::clone(&xrpc_state), + })); + let ssh_state = Arc::new(match &lfs_handle { + Some(handle) => ssh_base.with_lfs(handle.clone(), lfs_max_ssh_transfers), + None => ssh_base, + }); + { + let state = Arc::clone(&xrpc_state); + tokio::task::spawn_blocking(move || knot_xrpc::materialize::enable_surfaces(&state)); + } + let seq_state = Arc::clone(&xrpc_state); + let seq_meta_path = seq_state.meta_path.clone(); let resolver: Arc = { let index = Arc::clone(&index); Arc::new(move |target: &knot_pack::RepoTarget| match target { @@ -641,13 +657,15 @@ async fn main() -> anyhow::Result<()> { ); knot_xrpc::legacy_admin::router(Arc::clone(&xrpc_state), secret) }); - let base_router = write_routes.merge(knot_xrpc::router(xrpc_state)).route( - "/.well-known/did.json", - get(move || { - let document = did_document.clone(); - async move { Json(document) } - }), - ); + let base_router = write_routes + .merge(knot_xrpc::router(Arc::clone(&xrpc_state))) + .route( + "/.well-known/did.json", + get(move || { + let document = did_document.clone(); + async move { Json(document) } + }), + ); let base_router = match legacy_admin_routes { Some(routes) => base_router.merge(routes), None => base_router, @@ -682,6 +700,29 @@ async fn main() -> anyhow::Result<()> { key_pace, shutdown.clone(), ); + { + let firehose = Arc::clone(&seq_state.firehose); + let seq_meta = seq_meta_path.clone(); + let seq_shutdown = shutdown.clone(); + tokio::spawn(async move { + let mut tick = tokio::time::interval(Duration::from_secs(SEQ_PERSIST_INTERVAL_SECS)); + loop { + tokio::select! { + _ = tick.tick() => {} + () = seq_shutdown.cancelled() => break, + } + let (firehose, seq_meta) = (Arc::clone(&firehose), seq_meta.clone()); + let _ = tokio::task::spawn_blocking(move || { + knot_xrpc::persist_seq(&firehose, &seq_meta) + }) + .await; + } + }); + } + tokio::spawn(knot_xrpc::plc::run_sweeper( + Arc::clone(&xrpc_state), + shutdown.clone(), + )); let mut edge_task = tokio::spawn(knot_edge::serve( edge_config, app, @@ -734,6 +775,12 @@ async fn main() -> anyhow::Result<()> { "aborting key fill drain after timeout" ); } + { + let firehose = Arc::clone(&seq_state.firehose); + let seq_meta = seq_meta_path.clone(); + let _ = + tokio::task::spawn_blocking(move || knot_xrpc::persist_seq(&firehose, &seq_meta)).await; + } if let Some(shutdown) = maintenance_shutdown { let _ = shutdown.send(true); } diff --git a/knot2/crates/knot-server/tests/invariants.rs b/knot2/crates/knot-server/tests/invariants.rs index 4bededef0..c2b3f01d0 100644 --- a/knot2/crates/knot-server/tests/invariants.rs +++ b/knot2/crates/knot-server/tests/invariants.rs @@ -211,7 +211,6 @@ fn the_shared_limit_defaults_match_the_config_defaults() { let server = ::Layer::default_values(); let bytes = ByteLimits::default(); let budgets = Budgets::default(); - let push_ms = budgets.languages_push.get().as_millis() as u64; let configured = [ ("body", xrpc.max_body_bytes), @@ -224,7 +223,6 @@ fn the_shared_limit_defaults_match_the_config_defaults() { ("tree_last_commit", xrpc.tree_last_commit_budget_ms), ("blob_last_commit", xrpc.blob_last_commit_budget_ms), ("languages", xrpc.languages_budget_ms), - ("languages_push", xrpc.languages_push_budget_ms), ] .map(|(name, value)| (name, value.expect("every limit has a config default"))); let shared = [ @@ -238,7 +236,6 @@ fn the_shared_limit_defaults_match_the_config_defaults() { ("tree_last_commit", ms(budgets.tree_last_commit.get())), ("blob_last_commit", ms(budgets.blob_last_commit.get())), ("languages", ms(budgets.languages.get())), - ("languages_push", push_ms), ]; assert_eq!( configured, shared, diff --git a/knot2/crates/knot-xrpc/src/cob.rs b/knot2/crates/knot-xrpc/src/cob.rs index 62d2b5367..379d9c0c4 100644 --- a/knot2/crates/knot-xrpc/src/cob.rs +++ b/knot2/crates/knot-xrpc/src/cob.rs @@ -1,5 +1,6 @@ -use knot_cob::{ChangePayload, Checkpoint, CobHome, CobStore, Evaluate}; -use knot_cobs::{Accept, Announcement, Effect, GrantChange, Offer, Roster, Standing}; +use knot_cob::{ChangePayload, Checkpoint, CobHome, CobStore, Evaluate, Materialized}; +use knot_cobs::{Accept, Effect, GrantChange, Offer, Roster, Standing}; +use knot_record::InviteRecord; use knot_runtime::Signer; use knot_types::{AccountDid, UnixSeconds}; @@ -49,32 +50,53 @@ impl Authorized { pub(crate) fn roster_apply( store: &CobStore, home: &CobHome, + collection: Option<&'static str>, Authorized(change): Authorized, signer: &dyn Signer, now: UnixSeconds, -) -> Result, XrpcError> +) -> Result<(), XrpcError> where E: Checkpoint + Evaluate, E::Change: ChangePayload + Clone + GrantChange, { - let announcement = change.effect().announcement(); match store.list::().map_err(XrpcError::from)?.as_slice() { [] => match refusal(&Roster::empty(), &change) { Some(error) => Err(error), - None if Roster::empty().would_change(&change) => store - .create(home, &change, signer, now) - .map(|_| announcement) - .map_err(XrpcError::from), - None => Ok(None), + None if Roster::empty().would_change(&change) => { + let record = record_for(&Roster::empty(), &change, collection)?; + store + .create_materialized( + home, + &Materialized { + record, + change: change.clone(), + }, + signer, + now, + ) + .map(|_| ()) + .map_err(XrpcError::from) + } + None => Ok(()), }, [object] => store - .update_maybe_checkpointed::(home, *object, signer, now, |roster| { - match refusal(roster, &change) { + .update_materialized_checkpointed::( + home, + *object, + signer, + now, + |roster| match refusal(roster, &change) { Some(error) => Err(error), - None => Ok(roster.would_change(&change).then(|| change.clone())), - } - }) - .map(|appended| appended.and(announcement)), + None => { + let record = record_for(roster, &change, collection)?; + Ok(roster.would_change(&change).then(|| Materialized { + record, + change: change.clone(), + })) + } + }, + ) + .map(|_| ()), many => Err(XrpcError::internal(format!( "{} collaborative objects of one type share namespace", many.len() @@ -82,6 +104,31 @@ where } } +fn record_for( + roster: &Roster, + change: &impl GrantChange, + collection: Option<&'static str>, +) -> Result>, XrpcError> { + match (change.effect(), collection) { + (Effect::Admit(Offer::Invited, grant), Some(collection)) => Ok(Some( + InviteRecord::new(collection, grant.created_at, grant.added_by.clone()) + .map_err(|error| XrpcError::internal(error.to_string()))? + .bytes(), + )), + (Effect::Accept(accept), Some(collection)) => Ok(roster + .get(&accept.subject) + .filter(|entry| entry.offer == Offer::Granted) + .map(|entry| InviteRecord::new(collection, entry.created_at, entry.added_by.clone())) + .transpose() + .map_err(|error| XrpcError::internal(error.to_string()))? + .map(|record| record.bytes())), + (Effect::Admit(Offer::Invited, _) | Effect::Accept(_), None) => Err(XrpcError::internal( + "a roster with invites is always materialized into a collection", + )), + (Effect::Admit(Offer::Granted, _) | Effect::Revoke(_), _) => Ok(None), + } +} + fn refusal(roster: &Roster, change: &impl GrantChange) -> Option { match change.effect() { Effect::Admit(Offer::Granted, grant) => { @@ -155,8 +202,7 @@ mod tests { fn a_grant_conflicts_with_an_outstanding_invite_and_the_other_three_effects_pass() { let invited = Roster::empty().apply_change(&Invite(offer())); assert!( - refusal(&invited, &offer()) - .is_some_and(|error| error.status() == StatusCode::CONFLICT), + refusal(&invited, &offer()).is_some_and(|error| error.status() == StatusCode::CONFLICT), "a Grant folds onto an outstanding invitation as a membership the account never \ signed for, so the locked path refuses it rather than swallowing it as a no-op, and \ by_operator plus knot-migrate's grandfathering leave this arm no caller at all" diff --git a/knot2/crates/knot-xrpc/src/collaborators.rs b/knot2/crates/knot-xrpc/src/collaborators.rs index e35314ead..6b0cc46b0 100644 --- a/knot2/crates/knot-xrpc/src/collaborators.rs +++ b/knot2/crates/knot-xrpc/src/collaborators.rs @@ -8,8 +8,7 @@ use serde::Deserialize; use knot_acl::{KnotAcl, can_manage_collaborators}; use knot_cob::{CobHome, CobStore}; -use knot_cobs::{Announcement, CollaboratorsChange, CollaboratorsCob, Grant, Invite, Removal}; -use knot_events::RepoCollaboratorUpdate; +use knot_cobs::{CollaboratorsChange, CollaboratorsCob, Grant, Invite, Removal}; use knot_index::Resolved; use knot_runtime::{Clock, HttpTransport}; use knot_types::{AccountDid, RepoDid, UnixSeconds}; @@ -89,33 +88,37 @@ async fn commit_collaborators( change: Authorized, now: UnixSeconds, ) -> Result { - let subject = change.subject().clone(); - let event_repo = repo.clone(); + let did = repo.as_str().to_owned(); let signer = state.secrets.signer(&state.knot_did)?; let layout = state.layout.clone(); let index = Arc::clone(&state.index); let cob_locks = Arc::clone(&state.cob_locks); - let events = Arc::clone(&state.events); + let firehose_state = Arc::clone(state); run_blocking(move || { let _guard = cob_locks.repo(&repo); owner_unmoved(&index, &repo, &owner)?; - if owner.is(&subject) { + if owner.is(change.subject()) { return Ok(()); } let git = layout.open(&repo)?; - let announced = roster_apply::( - &CobStore::new(&git), + let store = CobStore::new(&git); + let collection = crate::materialize::COLLABORATOR_INVITE_COLLECTION; + roster_apply::( + &store, &CobHome::from(&repo), + Some(collection), change, &signer, now, )?; - if let Some(update) = announced.map(|announcement| match announcement { - Announcement::Effective => RepoCollaboratorUpdate::added(subject, event_repo), - Announcement::Cleared => RepoCollaboratorUpdate::removed(subject, event_repo), - }) { - events.publish(&update); - } + crate::materialize::sync_and_publish::( + &firehose_state, + &git, + &did, + collection, + &signer, + now, + ); index.refresh_collaborators(&repo)?; Ok(()) }) diff --git a/knot2/crates/knot-xrpc/src/consent.rs b/knot2/crates/knot-xrpc/src/consent.rs index ed3729c81..53fe72986 100644 --- a/knot2/crates/knot-xrpc/src/consent.rs +++ b/knot2/crates/knot-xrpc/src/consent.rs @@ -54,9 +54,8 @@ impl AcceptanceRef { let rkey = uri .rkey() .ok_or_else(|| XrpcError::invalid_request("no record key in the acceptance at-uri"))?; - let subject = DidRkey::new(rkey.as_str()).map_err(|_| { - XrpcError::invalid_request("acceptance record key must be a bare DID") - })?; + let subject = DidRkey::new(rkey.as_str()) + .map_err(|_| XrpcError::invalid_request("acceptance record key must be a bare DID"))?; Ok(Self { author, subject, diff --git a/knot2/crates/knot-xrpc/src/lists.rs b/knot2/crates/knot-xrpc/src/lists.rs index 511b7577a..dc4faad6c 100644 --- a/knot2/crates/knot-xrpc/src/lists.rs +++ b/knot2/crates/knot-xrpc/src/lists.rs @@ -125,7 +125,9 @@ fn page(entries: Vec<(AccountDid, Entry)>, window: Window) -> Response { let mut rows: Vec<(UnixSeconds, AccountDid, Entry)> = entries .into_iter() .map(|(subject, entry)| { - let at = entry.effective_since().map_or(entry.created_at, EffectiveSince::seconds); + let at = entry + .effective_since() + .map_or(entry.created_at, EffectiveSince::seconds); (at, subject, entry) }) .collect(); @@ -187,9 +189,9 @@ pub(crate) async fn list_member_invites( Resolved::Ready(entries) => entries .into_iter() .map(|(subject, entry)| match state.index.is_blocked(&subject) { - Resolved::Warming => { - Err(XrpcError::warming("blocklist projection is still warming, retry shortly")) - } + Resolved::Warming => Err(XrpcError::warming( + "blocklist projection is still warming, retry shortly", + )), Resolved::Ready(true) => Ok(None), Resolved::Ready(false) => Ok(Some((subject, entry))), }) diff --git a/knot2/crates/knot-xrpc/src/members.rs b/knot2/crates/knot-xrpc/src/members.rs index a81c213fb..6ec536f36 100644 --- a/knot2/crates/knot-xrpc/src/members.rs +++ b/knot2/crates/knot-xrpc/src/members.rs @@ -8,8 +8,7 @@ use serde::Deserialize; use knot_acl::{KnotAcl, can_admin_knot}; use knot_cob::{CobHome, CobStore}; -use knot_cobs::{Announcement, Grant, Invite, MembersChange, MembersCob, Removal}; -use knot_events::KnotMemberUpdate; +use knot_cobs::{Grant, Invite, MembersChange, MembersCob, Removal}; use knot_git::Repo; use knot_runtime::{Clock, HttpTransport}; use knot_types::{AccountDid, UnixSeconds}; @@ -74,24 +73,27 @@ async fn commit_members( change: Authorized, now: UnixSeconds, ) -> Result { - let subject = change.subject().clone(); let signer = state.secrets.signer(&state.knot_did)?; let meta_path = state.meta_path.clone(); let index = Arc::clone(&state.index); let cob_locks = Arc::clone(&state.cob_locks); - let events = Arc::clone(&state.events); + let firehose_state = Arc::clone(state); + let knot_did = state.knot_did.as_str().to_owned(); let home = CobHome::from(&state.knot_did); run_blocking(move || { let _guard = cob_locks.meta(); let meta = Repo::open(&meta_path)?; - let announced = - roster_apply::(&CobStore::new(&meta), &home, change, &signer, now)?; - if let Some(update) = announced.map(|announcement| match announcement { - Announcement::Effective => KnotMemberUpdate::added(subject), - Announcement::Cleared => KnotMemberUpdate::removed(subject), - }) { - events.publish(&update); - } + let store = CobStore::new(&meta); + let collection = crate::materialize::MEMBER_INVITE_COLLECTION; + roster_apply::(&store, &home, Some(collection), change, &signer, now)?; + crate::materialize::sync_and_publish::( + &firehose_state, + &meta, + &knot_did, + collection, + &signer, + now, + ); index.refresh_members()?; Ok(()) }) diff --git a/knot2/crates/knot-xrpc/src/plc.rs b/knot2/crates/knot-xrpc/src/plc.rs new file mode 100644 index 000000000..d5951f778 --- /dev/null +++ b/knot2/crates/knot-xrpc/src/plc.rs @@ -0,0 +1,491 @@ +use std::collections::{BTreeSet, HashMap, HashSet}; +use std::path::Path; +use std::sync::Arc; +use std::time::Duration; + +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +use knot_atproto::AtprotoError; +use knot_atproto::RepoTarget; +use knot_git::Repo; +use knot_runtime::Clock; +use knot_runtime::HttpTransport; +use knot_types::RefName; +use knot_types::RepoDid; + +use crate::XrpcState; +use crate::error::XrpcError; + +const PENDING_REF: &str = "refs/atproto/plc/pending"; +const DELETED_REF: &str = "refs/atproto/deleted"; +pub const STATUS_DELETED: &str = "deleted"; +const PLC_TOMBSTONE_TYPE: &str = "plc_tombstone"; +const SWEEP_INTERVAL: Duration = Duration::from_secs(5); +const SUBMIT_GAP: Duration = Duration::from_millis(250); +const BACKOFF_BASE_MICROS: u64 = 30_000_000; +const BACKOFF_CAP_MICROS: u64 = 60 * 60 * 1_000_000; +pub(crate) const MAX_ATTEMPTS: u32 = 12; + +#[derive(Debug, Serialize, Deserialize)] +pub struct PendingGenesis { + pub did: String, + pub operation: serde_json::Value, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct DeletedAccount { + did: String, + deleted_at_micros: u64, +} + +fn load Deserialize<'de>>(repo: &Repo, name: &str) -> Result, XrpcError> { + let name = RefName::new(name).map_err(|error| XrpcError::internal(error.to_string()))?; + let Some(oid) = repo.find_ref(&name).map_err(XrpcError::from)? else { + return Ok(None); + }; + let object = repo + .git() + .find_object(oid.object_id()) + .map_err(|error| XrpcError::internal(error.to_string()))?; + serde_json::from_slice(&object.detach().data) + .map(Some) + .map_err(|error| XrpcError::internal(error.to_string())) +} + +fn store(repo: &Repo, name: &str, value: Option<&T>) -> Result<(), XrpcError> { + let name = RefName::new(name).map_err(|error| XrpcError::internal(error.to_string()))?; + let new = value + .map(|value| { + serde_json::to_vec(value) + .map_err(|error| XrpcError::internal(error.to_string())) + .and_then(|bytes| { + repo.git() + .write_blob(&bytes) + .map(|blob| knot_types::Oid::from(blob.detach())) + .map_err(|error| XrpcError::internal(error.to_string())) + }) + }) + .transpose()?; + let update = match (new, repo.find_ref(&name).map_err(XrpcError::from)?) { + (Some(new), Some(old)) => knot_git::RefUpdate::Update { name, old, new }, + (Some(new), None) => knot_git::RefUpdate::Create { name, new }, + (None, None) => return Ok(()), + (None, Some(old)) => knot_git::RefUpdate::Delete { name, old }, + }; + repo.update_ref(&update).map_err(XrpcError::from) +} + +fn edit_list(repo: &Repo, name: &str, edit: F) -> Result<(), XrpcError> +where + T: for<'de> Deserialize<'de> + Serialize + Default, + F: FnOnce(&mut T) -> bool, +{ + let mut entries = load::(repo, name)?.unwrap_or_default(); + if edit(&mut entries) { + store(repo, name, Some(&entries)) + } else { + Ok(()) + } +} + +pub fn pending_entries(meta: &Repo) -> Result, XrpcError> { + Ok(load(meta, PENDING_REF)?.unwrap_or_default()) +} + +pub fn enqueue_pending_genesis(meta: &Repo, entry: PendingGenesis) -> Result<(), XrpcError> { + edit_list( + meta, + PENDING_REF, + move |entries: &mut Vec| { + entries.retain(|existing| existing.did != entry.did); + entries.push(entry); + true + }, + ) +} + +pub fn drop_pending_genesis(meta: &Repo, did: &RepoDid) -> Result { + edit_list(meta, PENDING_REF, |entries: &mut Vec| { + let before = entries.len(); + entries.retain(|entry| entry.did != did.as_str()); + entries.len() != before + }) + .map(|()| true) +} + +pub fn record_deletion( + meta: &Repo, + did: &RepoDid, + deleted_at_micros: u64, +) -> Result<(), XrpcError> { + edit_list(meta, DELETED_REF, |accounts: &mut Vec| { + accounts.retain(|account| account.did != did.as_str()); + accounts.push(DeletedAccount { + did: did.as_str().to_owned(), + deleted_at_micros, + }); + true + }) +} + +pub fn deleted_dids(meta: &Repo) -> Result, XrpcError> { + Ok(load::>(meta, DELETED_REF)? + .unwrap_or_default() + .into_iter() + .map(|account| account.did) + .collect()) +} + +pub fn deleted_accounts(meta: &Repo) -> Result, XrpcError> { + Ok(load::>(meta, DELETED_REF)? + .unwrap_or_default() + .into_iter() + .map(|account| (account.did, account.deleted_at_micros)) + .collect()) +} + +async fn with_meta( + state: &XrpcState, + query: Q, +) -> Result +where + T: Send + 'static, + Q: FnOnce(&Repo) -> Result + Send + 'static, +{ + let cob_locks = Arc::clone(&state.cob_locks); + let meta_path = state.meta_path.clone(); + crate::run_blocking(move || { + let _guard = cob_locks.meta(); + query(&Repo::open(&meta_path)?) + }) + .await +} + +async fn clear_pending(state: &XrpcState, did: &RepoDid) { + let wanted = did.clone(); + let _ = with_meta(state, move |meta| drop_pending_genesis(meta, &wanted)).await; +} + +pub(crate) async fn is_deleted( + state: &XrpcState, + did: &str, +) -> Result { + let wanted = did.to_owned(); + with_meta(state, move |meta| { + deleted_dids(meta).map(|set| set.contains(&wanted)) + }) + .await +} + +#[derive(Clone, Copy, Default)] +pub(crate) struct Waiter { + attempts: u32, + due_at_micros: u64, +} + +#[derive(Default)] +pub(crate) struct SweepMemory { + pub(crate) settled: HashSet, + pub(crate) waiting: HashMap, +} + +impl SweepMemory { + fn is_due(&self, did: &str, now_micros: u64) -> bool { + self.waiting + .get(did) + .is_none_or(|waiter| now_micros >= waiter.due_at_micros) + } + + #[cfg(test)] + pub(crate) fn is_parked(&self, did: &str) -> bool { + self.waiting + .get(did) + .is_some_and(|waiter| waiter.due_at_micros == u64::MAX) + } + + fn schedule(&mut self, did: &str, now_micros: u64) -> u32 { + let waiter = self.waiting.entry(did.to_owned()).or_default(); + waiter.attempts = waiter.attempts.saturating_add(1); + waiter.due_at_micros = now_micros.saturating_add( + BACKOFF_BASE_MICROS + .saturating_mul(1u64 << waiter.attempts.saturating_sub(1).min(7)) + .min(BACKOFF_CAP_MICROS), + ); + let attempts = waiter.attempts; + if attempts >= MAX_ATTEMPTS { + tracing::error!( + repo = did, + "plc reconciler parked this repository after {MAX_ATTEMPTS} failed attempts" + ); + self.park(did); + } + attempts + } + fn park(&mut self, did: &str) { + self.waiting.insert( + did.to_owned(), + Waiter { + attempts: MAX_ATTEMPTS + 1, + due_at_micros: u64::MAX, + }, + ); + } + + fn settle(&mut self, did: &str) { + self.settled.insert(did.to_owned()); + self.waiting.remove(did); + } + + fn clear(&mut self, did: &str) { + self.waiting.remove(did); + } +} + +enum Outcome { + Settled, + Submitted, + Retriable(String), + Parked(String), +} + +pub async fn run_sweeper( + state: Arc>, + shutdown: tokio_util::sync::CancellationToken, +) { + let mut tick = tokio::time::interval(SWEEP_INTERVAL); + tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + let _first = tick.tick().await; + let mut memory = SweepMemory::default(); + loop { + tokio::select! { + () = shutdown.cancelled() => return, + _ = tick.tick() => sweep_tick(&state, &mut memory).await, + } + } +} + +pub(crate) async fn sweep_tick( + state: &XrpcState, + memory: &mut SweepMemory, +) { + let now_micros = state.atproto.now().get(); + let candidates: Vec = state + .index + .hosted_repos() + .into_iter() + .filter(|did| did.as_str().starts_with("did:plc:")) + .filter(|did| !memory.settled.contains(did.as_str())) + .filter(|did| memory.is_due(did.as_str(), now_micros)) + .collect(); + if candidates.is_empty() { + return; + } + let knot_public = match state.secrets.public_key(&state.knot_did) { + Ok(public) => public, + Err(error) => { + tracing::error!( + %error, + "the knot signer is unavailable, so this sweep ends before the first submission" + ); + return; + } + }; + let target = RepoTarget::new(&knot_public, &state.knot_service_url); + for did in candidates { + match reconcile(state, &target, &did).await { + Outcome::Settled => { + memory.settle(did.as_str()); + tracing::info!(repo = %did, "plc document matches the desired account shape"); + } + Outcome::Submitted => { + memory.clear(did.as_str()); + tracing::info!( + repo = %did, + "plc operation accepted; verification lands on the next pass" + ); + } + Outcome::Retriable(message) => { + if memory.schedule(did.as_str(), now_micros) <= 2 { + tracing::warn!(repo = %did, %message, "plc submission will retry"); + } + } + Outcome::Parked(message) => { + memory.park(did.as_str()); + tracing::error!(repo = %did, %message, "plc reconciler parked this repository"); + } + } + tokio::time::sleep(SUBMIT_GAP).await; + } +} + +async fn reconcile( + state: &XrpcState, + target: &RepoTarget, + did: &RepoDid, +) -> Outcome { + let last_op = match state.atproto.last_plc_operation(did).await { + Ok(Some(op)) => op, + Ok(None) => return submit_genesis(state, did).await, + Err(error) => return Outcome::Retriable(error.to_string()), + }; + if last_op.get("type").and_then(Value::as_str) == Some(PLC_TOMBSTONE_TYPE) { + return Outcome::Parked(format!( + "{did} was tombstoned in the directory; its history is closed" + )); + } + match target.satisfied_by(&last_op) { + Ok(true) => { + clear_pending(state, did).await; + Outcome::Settled + } + Ok(false) => { + let signer = match signer_for(state, target.knot_did_key(), &last_op, did) { + Ok(signer) => signer, + Err(outcome) => return outcome, + }; + match knot_atproto::update_operation(&signer, &last_op, target) { + Ok(update) => post(state, did, &update).await, + Err(error) => Outcome::Parked(error.to_string()), + } + } + Err(error) => Outcome::Parked(error.to_string()), + } +} + +async fn submit_genesis( + state: &XrpcState, + did: &RepoDid, +) -> Outcome { + let wanted = did.clone(); + match with_meta(state, move |meta| { + Ok(pending_entries(meta)? + .into_iter() + .find(|entry| entry.did == wanted.as_str())) + }) + .await + { + Err(error) => Outcome::Retriable(format!("couldn't read the pending queue: {error}")), + Ok(None) => { + Outcome::Parked("the directory holds no operations and no genesis is queued".to_owned()) + } + Ok(Some(entry)) => post(state, did, &entry.operation).await, + } +} + +async fn post( + state: &XrpcState, + did: &RepoDid, + operation: &Value, +) -> Outcome { + match state.atproto.submit_plc_operation(did, operation).await { + Ok(()) => Outcome::Submitted, + Err(AtprotoError::PlcSubmit { status, .. }) if !status.is_transient() => { + Outcome::Parked(format!( + "the directory rejected this operation with HTTP {status}; \ + resubmitting identical bytes can't converge" + )) + } + Err(error) => Outcome::Retriable(error.to_string()), + } +} + +fn signer_for( + state: &XrpcState, + knot_did_key: &str, + last_op: &Value, + did: &RepoDid, +) -> Result { + let rotation_names_knot = last_op + .get("rotationKeys") + .and_then(Value::as_array) + .is_some_and(|keys| keys.iter().any(|key| key.as_str() == Some(knot_did_key))); + if rotation_names_knot { + return state + .secrets + .signer(&state.knot_did) + .map_err(|_| Outcome::Parked("the sealed knot signer is missing".to_owned())); + } + knot_types::KnotId::new(did.as_str()) + .ok() + .and_then(|id| { + state + .secrets + .signer(knot_secrets::SealedKeyId::from(&id)) + .ok() + }) + .ok_or_else(|| { + Outcome::Parked(format!( + "{did} rotates on its migration per-repo key and none is sealed; \ + point [atproto].repo_signing_keys at the migration archive" + )) + }) +} + +#[derive(serde::Deserialize)] +struct ArchivedKey { + repo_did: String, + key_type: String, + secret_key_hex: String, +} + +pub fn ingest_legacy_keys( + secrets: &knot_secrets::SealedStore, + archive: &Path, +) -> Result { + let bytes = std::fs::read(archive) + .map_err(|error| XrpcError::internal(format!("reading {}: {error}", archive.display())))?; + let entries: Vec = + serde_json::from_slice(&bytes).map_err(|error| XrpcError::internal(error.to_string()))?; + let provided = entries + .into_iter() + .filter_map(|entry| match archived_scalar(&entry) { + Ok(pair) => Some(pair), + Err(reason) => { + tracing::warn!( + repo_did = %entry.repo_did, + %reason, + "skipped a row in the migration key archive" + ); + None + } + }); + secrets + .ingest(provided) + .map_err(|error| XrpcError::internal(error.to_string())) +} + +fn archived_scalar( + entry: &ArchivedKey, +) -> Result<(knot_secrets::SealedKeyId, [u8; 32]), &'static str> { + if entry.key_type != "k256" { + return Err("unsupported key type"); + } + let scalar = knot_types::decode_hex(&entry.secret_key_hex).ok_or("the secret key isn't hex")?; + let scalar: [u8; 32] = scalar + .try_into() + .map_err(|_| "the secret key isn't a 32-byte scalar")?; + let id = knot_types::KnotId::new(&entry.repo_did) + .map_err(|_| "the repo DID isn't a knot identity")?; + Ok((knot_secrets::SealedKeyId::from(&id), scalar)) +} + +#[cfg(test)] +mod backoff_tests { + use super::*; + + #[test] + fn backoff_doubles_to_3600_seconds_and_sweep_parks_on_twelfth_failure() { + let mut memory = SweepMemory::default(); + let did = "did:plc:nelpet"; + for want_seconds in [30u64, 60, 120, 240, 480, 960, 1920, 3600, 3600, 3600, 3600] { + memory.schedule(did, 0); + assert_eq!(memory.waiting[did].due_at_micros / 1_000_000, want_seconds); + } + memory.schedule(did, 0); + assert!( + memory.is_parked(did), + "the sweep parks the repository on the twelfth failure" + ); + } +} diff --git a/knot2/crates/knot-xrpc/src/repos.rs b/knot2/crates/knot-xrpc/src/repos.rs index b297a7d06..6a4b467cc 100644 --- a/knot2/crates/knot-xrpc/src/repos.rs +++ b/knot2/crates/knot-xrpc/src/repos.rs @@ -246,27 +246,14 @@ pub(crate) async fn create_repo( None => Vec::new(), }; - let submission = match &provisioning { - Provisioned::Minted(prepared) => state - .atproto - .submit_plc_operation(prepared) - .await - .map(|_| ()), - Provisioned::Reserved => Ok(()), - }; - if let Err(error) = submission { - let layout = state.layout.clone(); - let placed = repo_did.clone(); - let lfs = lfs_store.clone(); - let _ = run_blocking(move || { - rollback_local(&layout, lfs.as_deref(), &placed); - Ok(()) - }) - .await; - return Err(error.into()); - } - let reserved = matches!(provisioning, Provisioned::Reserved); + let genesis = match &provisioning { + Provisioned::Minted(prepared) => Some(crate::plc::PendingGenesis { + did: repo_did.as_str().to_owned(), + operation: prepared.operation().clone(), + }), + Provisioned::Reserved => None, + }; let layout = state.layout.clone(); let cob_locks = Arc::clone(&state.cob_locks); let meta_path = state.meta_path.clone(); @@ -274,11 +261,16 @@ pub(crate) async fn create_repo( let lfs = lfs_store.clone(); let home = CobHome::from(&state.knot_did); run_blocking(move || { - let outcome = { - let _guard = cob_locks.meta(); - Repo::open(&meta_path) - .map_err(XrpcError::from) - .and_then(|meta| { + let _guard = cob_locks.meta(); + let was_queued = genesis.is_some(); + let outcome = Repo::open(&meta_path) + .map_err(XrpcError::from) + .and_then(|meta| { + let queued = match genesis { + Some(entry) => crate::plc::enqueue_pending_genesis(&meta, entry), + None => Ok(()), + }; + let registered = queued.and_then(|()| { register( &CobStore::new(&meta), &home, @@ -286,16 +278,30 @@ pub(crate) async fn create_repo( &knot_signer, now, ) - }) - }; + }); + if registered.is_err() + && was_queued + && let Err(cleanup) = crate::plc::drop_pending_genesis(&meta, &placed) + { + tracing::warn!( + repo = %placed, + %cleanup, + "a dead genesis row stays queued beside the rolled back repository" + ); + } + registered + }); if outcome.is_err() { rollback_local(&layout, lfs.as_deref(), &placed); - if !reserved { - tracing::error!( - repo = %placed, - "registration failed after did:plc submitted to PLC directory" - ); - } + } + if let Err(error) = &outcome + && !reserved + { + tracing::error!( + repo = %placed, + %error, + "repository registration failed; the minted repository was rolled back" + ); } outcome }) @@ -316,6 +322,65 @@ pub(crate) async fn create_repo( }) .await?; + let knot_signer = state.secrets.signer(&state.knot_did)?; + let layout = state.layout.clone(); + let cob_locks = Arc::clone(&state.cob_locks); + let firehose = Arc::clone(&state.firehose); + let did = repo_did.clone(); + let initialized = repo_did.clone(); + let chunk = run_blocking(move || -> Result { + let identity = firehose.reserve(); + let account = firehose.reserve(); + let _guard = cob_locks.repo(&initialized); + let git = layout.open(&initialized)?; + let mut emit = crate::FirehoseEmit::new(&firehose, did.as_str(), now); + let tail = futures::executor::block_on(knot_record::chain::initialize( + &git, + did.as_str(), + &knot_signer, + &mut emit, + now, + ))?; + let drift = knot_record::git_ref::desired_records(&git, &[], None, &[]) + .inspect_err(|error| { + tracing::warn!(repo = %did, %error, "a fork's seeded ref set wouldn't build; the boot pass or its first push resyncs it") + }) + .ok() + .filter(|records| !records.is_empty()); + let tail = match drift { + Some(drift) => { + match futures::executor::block_on(knot_record::chain::sync( + &git, + did.as_str(), + knot_record::git_ref::GIT_REF_COLLECTION, + &drift, + &knot_signer, + &mut emit, + now, + )) { + Ok(chunks) => chunks.into_iter().last().unwrap_or(tail), + Err(error) => { + tracing::warn!(repo = %did, %error, "a fork's seeded refs wouldn't project at birth; the boot pass or its first push resyncs them"); + tail + } + } + } + None => tail, + }; + emit.finish(); + crate::firehose::fulfill_identity(identity, did.as_str(), now); + crate::firehose::fulfill_account(account, did.as_str(), true, None, now); + Ok(tail) + }) + .await?; + crate::firehose::publish_sync( + &state.firehose, + repo_did.as_str(), + &chunk.rev, + &chunk.commit, + now, + ); + Ok(( http::StatusCode::OK, Json(CreateOutput { @@ -402,9 +467,15 @@ async fn provision_repo_did( let knot_public = state.secrets.public_key(&state.knot_did)?; state .atproto - .verify_did_web_publishes_key(&did, &knot_public) + .verify_did_web_account_document(&did, &knot_public, &state.knot_service_url) .await - .map_err(XrpcError::from)?; + .map_err(|error| match error { + knot_atproto::AtprotoError::Resolve(_) => XrpcError::invalid_request( + "the did:web document must publish an #atproto verification method \ + holding this knot's key and an #atproto_pds service pointing at it", + ), + other => other.into(), + })?; Ok((did, Provisioned::Reserved)) } Some(_) => Err(XrpcError::invalid_request( @@ -509,6 +580,7 @@ pub(crate) async fn delete_repo( } let now = state.now(); + let now_micros = state.atproto.now().get(); let knot_signer = state.secrets.signer(&state.knot_did)?; let layout = state.layout.clone(); let index = Arc::clone(&state.index); @@ -523,6 +595,8 @@ pub(crate) async fn delete_repo( let _meta_guard = cob_locks.meta(); let meta = Repo::open(&meta_path)?; let store = CobStore::new(&meta); + crate::plc::record_deletion(&meta, &deleted, now_micros)?; + crate::plc::drop_pending_genesis(&meta, &deleted)?; deregister(&store, &home, target, &deleted, &knot_signer, now)?; let removal = layout.remove(&deleted); index.refresh_registry().map_err(|error| { @@ -547,6 +621,13 @@ pub(crate) async fn delete_repo( }) .await?; + crate::firehose::publish_account( + &state.firehose, + repo_did.as_str(), + false, + Some("deleted"), + state.now(), + ); Ok(ok_empty()) } diff --git a/knot2/example.toml b/knot2/example.toml index 67b660c5c..c0935ec3e 100644 --- a/knot2/example.toml +++ b/knot2/example.toml @@ -174,6 +174,9 @@ # Required! This value must be specified. #plc_directory = +# Can also be specified via environment variable `KNOT_REPO_SIGNING_KEYS`. +#repo_signing_keys = + [xrpc] # Can also be specified via environment variable `KNOT_XRPC_MAX_BODY_BYTES`. # Default value: 65536 @@ -206,10 +209,6 @@ # Default value: 1000 #languages_budget_ms = 1000 -# Can also be specified via environment variable `KNOT_XRPC_LANGUAGES_PUSH_BUDGET_MS`. -# Default value: 2000 -#languages_push_budget_ms = 2000 - # Body limit for the merge and mergeCheck procedures, whose patch payloads # routinely exceed the general XRPC body limit. # diff --git a/knot2/interop/atproto_firehose_test.go b/knot2/interop/atproto_firehose_test.go new file mode 100644 index 000000000..a09c1fe6e --- /dev/null +++ b/knot2/interop/atproto_firehose_test.go @@ -0,0 +1,210 @@ +package interop + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "strings" + "testing" + "time" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/gorilla/websocket" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + cbg "github.com/whyrusleeping/cbor-gen" +) + +func readStreamHeader(r io.Reader) (string, error) { + maj, count, err := cbg.CborReadHeader(r) + if err != nil { + return "", err + } + if maj != cbg.MajMap { + return "", fmt.Errorf("expected a cbor map header, got major type %d", maj) + } + msgType := "" + for range count { + key, err := cbg.ReadString(r) + if err != nil { + return "", err + } + switch key { + case "op": + if _, _, err := cbg.CborReadHeader(r); err != nil { + return "", err + } + case "t": + if msgType, err = cbg.ReadString(r); err != nil { + return "", err + } + default: + return "", fmt.Errorf("unknown header field %q", key) + } + } + return msgType, nil +} + +type firehoseFrame struct { + MsgType string + Commit *comatproto.SyncSubscribeRepos_Commit + Sync *comatproto.SyncSubscribeRepos_Sync + Identity *comatproto.SyncSubscribeRepos_Identity + Account *comatproto.SyncSubscribeRepos_Account + Info *comatproto.SyncSubscribeRepos_Info +} + +func nextFrame(t *testing.T, conn *websocket.Conn) firehoseFrame { + t.Helper() + require.NoError(t, conn.SetReadDeadline(time.Now().Add(15*time.Second))) + messageType, data, err := conn.ReadMessage() + require.NoError(t, err) + require.Equal(t, websocket.BinaryMessage, messageType) + reader := bytes.NewReader(data) + msgType, err := readStreamHeader(reader) + require.NoError(t, err, "frame % x", data) + frame := firehoseFrame{MsgType: msgType} + switch msgType { + case "#commit": + var evt comatproto.SyncSubscribeRepos_Commit + require.NoError(t, evt.UnmarshalCBOR(reader)) + frame.Commit = &evt + case "#sync": + var evt comatproto.SyncSubscribeRepos_Sync + require.NoError(t, evt.UnmarshalCBOR(reader)) + frame.Sync = &evt + case "#identity": + var evt comatproto.SyncSubscribeRepos_Identity + require.NoError(t, evt.UnmarshalCBOR(reader)) + frame.Identity = &evt + case "#account": + var evt comatproto.SyncSubscribeRepos_Account + require.NoError(t, evt.UnmarshalCBOR(reader)) + frame.Account = &evt + case "#info": + var evt comatproto.SyncSubscribeRepos_Info + require.NoError(t, evt.UnmarshalCBOR(reader)) + frame.Info = &evt + default: + t.Fatalf("unknown frame type %q", msgType) + } + require.Zero(t, reader.Len(), "a frame holds exactly two objects") + return frame +} + +func authedPost(t *testing.T, path, token string, body any) map[string]any { + t.Helper() + encoded, err := json.Marshal(body) + require.NoError(t, err) + request, err := http.NewRequest( + http.MethodPost, + os.Getenv("KNOT_TEST_BASE_URL")+path, + bytes.NewReader(encoded), + ) + require.NoError(t, err) + request.Header.Set("Content-Type", "application/json") + request.Header.Set("Authorization", "Bearer "+token) + response, err := http.DefaultClient.Do(request) + require.NoError(t, err) + defer response.Body.Close() + raw, err := io.ReadAll(response.Body) + require.NoError(t, err) + require.Equal( + t, + http.StatusOK, + response.StatusCode, + "%s answered %s: %s", + path, + response.Status, + string(raw), + ) + decoded := map[string]any{} + if len(raw) > 0 { + require.NoError(t, json.Unmarshal(raw, &decoded)) + } + return decoded +} + +func TestKnotServesSubscribeRepos(t *testing.T) { + base := os.Getenv("KNOT_TEST_BASE_URL") + require.NotEmpty(t, base, "KNOT_TEST_BASE_URL names the knot under test") + createToken := os.Getenv("KNOT_TEST_CREATE_TOKEN") + require.NotEmpty(t, createToken, "KNOT_TEST_CREATE_TOKEN holds the create-scope token") + deleteToken := os.Getenv("KNOT_TEST_DELETE_TOKEN") + require.NotEmpty(t, deleteToken, "KNOT_TEST_DELETE_TOKEN holds the delete-scope token") + endpoint := "ws://" + strings.TrimPrefix(base, "http://") + "/xrpc/com.atproto.sync.subscribeRepos" + conn, _, err := websocket.DefaultDialer.DialContext(context.Background(), endpoint, nil) + require.NoError(t, err) + defer conn.Close() + + created := authedPost(t, "/xrpc/sh.tangled.repo.create", createToken, map[string]any{ + "rkey": "limpetkey", + "name": "limpet", + }) + repoDid, _ := created["repoDid"].(string) + require.NotEmpty(t, repoDid, "create named the repository it minted: %v", created) + + identity := nextFrame(t, conn) + require.NotNil(t, identity.Identity) + assert.Equal(t, repoDid, identity.Identity.Did) + assert.Greater(t, identity.Identity.Seq, int64(0)) + assert.Nil(t, identity.Identity.Handle, "the knot names no handle it doesn't answer for") + + account := nextFrame(t, conn) + require.NotNil(t, account.Account) + assert.Equal(t, repoDid, account.Account.Did) + assert.True(t, account.Account.Active) + assert.Nil(t, account.Account.Status) + + commit := nextFrame(t, conn) + require.NotNil(t, commit.Commit) + genesis := commit.Commit + assert.Equal(t, repoDid, genesis.Repo) + assert.Empty(t, genesis.Ops, "a fresh account's first commit is empty") + assert.False(t, genesis.Rebase) + assert.False(t, genesis.TooBig) + assert.Nil(t, genesis.Since, "a genesis commit names no predecessor") + assert.Nil(t, genesis.PrevData, "a genesis commit omits prevData rather than nulling it") + assert.Empty(t, genesis.Blobs) + rev, err := syntax.ParseTID(genesis.Rev) + require.NoError(t, err) + assert.Equal( + t, + rev.Time().UnixMicro(), + genesis.Seq, + "the seq is the microsecond tick of the commit's own rev", + ) + assert.NotEmpty(t, genesis.Blocks, "the commit frame includes its CAR") + assert.NotEmpty(t, genesis.Commit.String()) + + sync := nextFrame(t, conn) + require.NotNil(t, sync.Sync) + assert.Equal(t, repoDid, sync.Sync.Did) + assert.Equal(t, genesis.Rev, sync.Sync.Rev, "the sync frame asserts the genesis commit") + assert.NotEmpty(t, sync.Sync.Blocks, "the sync frame includes the commit block") + + assert.Less(t, identity.Identity.Seq, account.Account.Seq, "identity precedes account") + assert.Less(t, account.Account.Seq, genesis.Seq, "account precedes the genesis commit") + assert.Less(t, genesis.Seq, sync.Sync.Seq, "the commit precedes its assertion") + + authedPost(t, "/xrpc/sh.tangled.repo.delete", deleteToken, map[string]any{ + "repo": repoDid, + "force": true, + }) + for range 10 { + frame := nextFrame(t, conn) + if frame.Account == nil || frame.Account.Did != repoDid { + continue + } + assert.False(t, frame.Account.Active) + require.NotNil(t, frame.Account.Status) + assert.Equal(t, "deleted", *frame.Account.Status) + return + } + t.Fatal("no account frame closed the lifecycle") +} diff --git a/knot2/interop/atproto_reads_test.go b/knot2/interop/atproto_reads_test.go new file mode 100644 index 000000000..ad32f05c4 --- /dev/null +++ b/knot2/interop/atproto_reads_test.go @@ -0,0 +1,165 @@ +package interop + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "testing" + "time" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/repo" + "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" + "github.com/samber/lo" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + cbg "github.com/whyrusleeping/cbor-gen" +) + +const memberInviteNSID = "sh.tangled.knot.memberInvite" + +type MemberInvite struct { + CreatedAt string `json:"createdAt"` + Editor string `json:"x-tngl-editor"` +} + +func (invite *MemberInvite) MarshalCBOR(w io.Writer) error { + return fmt.Errorf("the interop suite only reads") +} + + +func (invite *MemberInvite) UnmarshalCBOR(r io.Reader) error { + maj, count, err := cbg.CborReadHeader(r) + if err != nil { + return err + } + if maj != cbg.MajMap { + return fmt.Errorf("expected a cbor map, got major type %d", maj) + } + return invite.readFields(r, count) +} + +func (invite *MemberInvite) readFields(r io.Reader, remaining uint64) error { + if remaining == 0 { + return nil + } + key, err := cbg.ReadString(r) + if err != nil { + return err + } + value, err := cbg.ReadString(r) + if err != nil { + return err + } + switch key { + case "createdAt": + invite.CreatedAt = value + case "x-tngl-editor": + invite.Editor = value + } + return invite.readFields(r, remaining-1) +} + +func knotClient(t *testing.T) *indigoxrpc.Client { + t.Helper() + return &indigoxrpc.Client{ + Host: os.Getenv("KNOT_TEST_BASE_URL"), + Client: &http.Client{Timeout: 10 * time.Second}, + } +} + +func knotDirectory(t *testing.T) identity.Directory { + t.Helper() + var doc identity.DIDDocument + require.NoError(t, json.Unmarshal([]byte(os.Getenv("KNOT_TEST_DID_DOC")), &doc)) + require.Equal(t, os.Getenv("KNOT_TEST_DID"), doc.DID.String()) + dir := identity.NewMockDirectory() + dir.Insert(identity.ParseIdentity(&doc)) + return dir +} + +func TestKnotServesComAtprotoReads(t *testing.T) { + ctx := context.Background() + xc := knotClient(t) + did := os.Getenv("KNOT_TEST_DID") + collection := os.Getenv("KNOT_TEST_COLLECTION") + rkey := os.Getenv("KNOT_TEST_RKEY") + lexutil.RegisterType(memberInviteNSID, &MemberInvite{}) + + latest, err := comatproto.SyncGetLatestCommit(ctx, xc, did) + require.NoError(t, err, "sync.getLatestCommit answers an unmodified client") + assert.NotEmpty(t, latest.Rev) + + repoCar, err := comatproto.SyncGetRepo(ctx, xc, did, "") + require.NoError(t, err, "sync.getRepo answers an unmodified client") + + commit, err := repo.VerifyCommitSignatureFromCar(ctx, knotDirectory(t), repoCar) + require.NoError(t, err, "the signed commit verifies against the did document key") + assert.Equal(t, did, commit.DID) + assert.Equal(t, latest.Rev, commit.Rev) + + _, full, err := repo.LoadRepoFromCAR(ctx, bytes.NewReader(repoCar)) + require.NoError(t, err, "the whole tree loads from the getRepo car") + nsid, err := syntax.ParseNSID(collection) + require.NoError(t, err) + recordKey, err := syntax.ParseRecordKey(rkey) + require.NoError(t, err) + recordCID, err := full.GetRecordCID(ctx, nsid, recordKey) + require.NoError(t, err, "the invite's inclusion proof walks the mst under the signed commit") + require.NotNil(t, recordCID) + + proofCar, err := comatproto.SyncGetRecord(ctx, xc, collection, did, rkey) + require.NoError(t, err, "sync.getRecord answers an unmodified client") + proofCommit, _, err := repo.LoadCommitFromCAR(ctx, bytes.NewReader(proofCar)) + require.NoError(t, err, "the proof car roots at the signed commit") + assert.Equal(t, commit.Data.String(), proofCommit.Data.String()) + assert.Equal(t, commit.Sig, proofCommit.Sig, "the same commit rides both cars") + _, proofRepo, err := repo.LoadRepoFromCAR(ctx, bytes.NewReader(proofCar)) + require.NoError(t, err, "the inclusion proof carries the blocks the walk needs") + proofCID, err := proofRepo.GetRecordCID(ctx, nsid, recordKey) + require.NoError(t, err) + assert.Equal(t, recordCID.String(), proofCID.String()) + + blocksCar, err := comatproto.SyncGetBlocks(ctx, xc, []string{commit.Data.String()}, did) + require.NoError(t, err, "sync.getBlocks answers an unmodified client") + _, blocksRepo, err := repo.LoadRepoFromCAR(ctx, bytes.NewReader(blocksCar)) + require.NoError(t, err, "the mst root arrives by name") + namedCID, err := blocksRepo.GetRecordCID(ctx, nsid, recordKey) + require.NoError(t, err) + assert.Equal(t, recordCID.String(), namedCID.String()) + + listed, err := comatproto.SyncListRepos(ctx, xc, "", 100) + require.NoError(t, err, "sync.listRepos answers an unmodified client") + assert.True(t, lo.ContainsBy(listed.Repos, func(row *comatproto.SyncListRepos_Repo) bool { + return row.Did == did + })) + + records, err := comatproto.RepoListRecords(ctx, xc, collection, "", 50, did, false) + require.NoError(t, err, "repo.listRecords answers an unmodified client") + assert.Equal(t, "at://"+did+"/"+collection+"/"+rkey, records.Records[0].Uri) + invite, ok := records.Records[0].Value.Val.(*MemberInvite) + require.True(t, ok, "the value decodes as the registered lexicon type") + assert.Equal(t, "did:web:olaren.dev", invite.Editor) + assert.Equal(t, "1970-01-01T00:16:40Z", invite.CreatedAt) + + describe, err := comatproto.RepoDescribeRepo(ctx, xc, did) + require.NoError(t, err, "repo.describeRepo answers an unmodified client") + assert.Equal(t, did, describe.Did) + assert.Contains(t, lo.Map(describe.Collections, func(row string, _ int) string { return row }), collection) + + status, err := comatproto.SyncGetRepoStatus(ctx, xc, did) + require.NoError(t, err, "sync.getRepoStatus answers an unmodified client") + require.NotNil(t, status.Rev) + assert.Equal(t, latest.Rev, *status.Rev) + + resolved, err := comatproto.IdentityResolveHandle(ctx, xc, os.Getenv("KNOT_TEST_HANDLE")) + require.NoError(t, err, "identity.resolveHandle answers an unmodified client") + assert.Equal(t, did, resolved.Did) +} diff --git a/knot2/interop/knotfeed_consumer_test.go b/knot2/interop/knotfeed_consumer_test.go new file mode 100644 index 000000000..3593d7054 --- /dev/null +++ b/knot2/interop/knotfeed_consumer_test.go @@ -0,0 +1,80 @@ +package interop + +import ( + "context" + "log/slog" + "os" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tangled.org/core/knotfeed" +) + +func TestKnotfeedConsumerReadsRefRecords(t *testing.T) { + addr := os.Getenv("KNOT_TEST_ADDR") + require.NotEmpty(t, addr, "KNOT_TEST_ADDR names the knot's host:port") + repoDid := os.Getenv("KNOT_TEST_REPO_DID") + require.NotEmpty(t, repoDid, "KNOT_TEST_REPO_DID names the pushed repository") + editor := os.Getenv("KNOT_TEST_EDITOR") + require.NotEmpty(t, editor, "KNOT_TEST_EDITOR names the pushing account") + wantSha := os.Getenv("KNOT_TEST_SHA") + require.NotEmpty(t, wantSha, "KNOT_TEST_SHA names the pushed tip") + + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + defer cancel() + + records := make(chan knotfeed.RecordOp, 8) + var cursor atomic.Int64 + consumer := &knotfeed.Consumer{ + Host: addr, + NoTLS: true, + ReplayFromStart: true, + Logger: slog.Default(), + LoadCursor: func(context.Context) (int64, error) { + return cursor.Load(), nil + }, + StoreCursor: func(_ context.Context, seq int64) error { + cursor.Store(seq) + return nil + }, + Handle: func(_ context.Context, message knotfeed.Message) error { + if message.Type != knotfeed.TypeCommit || message.Commit == nil { + return nil + } + if message.Commit.Repo != repoDid { + return nil + } + for _, op := range message.Commit.Records { + if op.Collection == knotfeed.GitRefCollection && op.Action == "create" { + select { + case records <- op: + default: + } + } + } + return nil + }, + } + go consumer.Run(ctx) + + select { + case op := <-records: + refname, ok := knotfeed.UnescapeRkey(op.Rkey) + require.True(t, ok, "the record key %q decodes", op.Rkey) + assert.Equal(t, "refs/heads/main", refname) + rec, err := knotfeed.DecodeRefRecord(op.Bytes) + require.NoError(t, err, "the frozen record envelope decodes") + assert.Equal(t, wantSha, rec.Sha, "the record carries the pushed tip") + assert.Equal(t, editor, rec.Editor, "the record names the pusher") + require.Len(t, rec.PushOptions, 1) + assert.Equal(t, "skip-ci", rec.PushOptions[0], "the push options ride the undeclared field") + require.Eventually(t, func() bool { return cursor.Load() >= int64(1) }, + 5*time.Second, 10*time.Millisecond, "the consumer advanced its cursor") + case <-ctx.Done(): + t.Fatal("no ref record arrived before the deadline") + } +}