diff --git a/Cargo.lock b/Cargo.lock index 6a31d4fbf..c4f5b8b8b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4994,6 +4994,7 @@ dependencies = [ "k256", "knot-cache", "knot-lexicons", + "knot-record", "knot-runtime", "knot-types", "proptest", diff --git a/knot2/crates/knot-atproto/Cargo.toml b/knot2/crates/knot-atproto/Cargo.toml index 9f967b783..c6df4a1de 100644 --- a/knot2/crates/knot-atproto/Cargo.toml +++ b/knot2/crates/knot-atproto/Cargo.toml @@ -7,6 +7,7 @@ license.workspace = true [dependencies] knot-types = { workspace = true } +knot-record = { workspace = true } knot-runtime = { workspace = true } knot-cache = { workspace = true } cid = { workspace = true } diff --git a/knot2/crates/knot-atproto/src/lib.rs b/knot2/crates/knot-atproto/src/lib.rs index 85b655a77..2806c5789 100644 --- a/knot2/crates/knot-atproto/src/lib.rs +++ b/knot2/crates/knot-atproto/src/lib.rs @@ -26,11 +26,11 @@ pub mod fuzz { } pub fn did_document(data: &[u8]) { - let did = knot_types::AccountDid::new("did:plc:nel").expect("constant test did is valid"); + 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 repo_did = knot_types::RepoDid::new("did:plc:nel").expect("Constant test DID is valid"); let pds = knot_types::KnotServiceUrl::new("https://knot.oyster.cafe") - .expect("constant test url is valid"); + .expect("Constant test URL is valid"); let _ = crate::resolve::document_publishes_account_shape( &repo_did, data, @@ -66,8 +66,8 @@ use knot_runtime::{ UnixMicros, }; use knot_types::{ - AccountDid, Collection, DidRkey, Handle, HttpStatus, KnotId, Nsid, OfferedKey, RepoDid, - RepoRkey, Rkey, UnixSeconds, + AccountDid, Collection, DidRkey, Handle, HttpStatus, KnotId, Nsid, OfferedKey, RecordRkey, + RepoDid, RepoRkey, Rkey, UnixSeconds, }; use pubkeys::Cursor; use serde::{Deserialize, Serialize}; @@ -78,11 +78,15 @@ const NEGATIVE_TTL: Duration = Duration::from_secs(30); const STALE_TTL: Duration = Duration::from_secs(30); const PUBKEY_PAGE_LIMIT: u16 = 100; fn repo_collection() -> Nsid { - Nsid::new_static("sh.tangled.repo").expect("literal nsid parses") + Nsid::new_static("sh.tangled.repo").expect("Literal NSID parses") } const PUBKEY_MAX_PAGES: usize = 8; const PUBKEY_TTL: Duration = Duration::from_secs(30); const MAX_PUBKEY_CACHE_BYTES: u64 = 1 << 20; +const DEF_TTL: Duration = Duration::from_secs(30); +const MAX_DEF_CACHE_BYTES: u64 = 1 << 20; +const MAX_REMOTE_DEF_BYTES: usize = 64 * 1024; +const DEF_CACHE_ENTRY_OVERHEAD: u64 = 128; const PUBKEY_CACHE_ENTRY_OVERHEAD: u64 = 128; const PUBKEY_CACHE_KEY_OVERHEAD: u64 = 192; const MAX_IDENTITY_CACHE: usize = 4096; @@ -96,7 +100,7 @@ pub enum AtprotoError { Resolve(#[from] ResolveError), #[error(transparent)] Jwt(#[from] JwtError), - #[error("network failure: {0}")] + #[error("Network failure: {0}")] Network(#[from] NetworkError), #[error("listRecords for {did} returned HTTP {status}")] ListRecords { did: AccountDid, status: HttpStatus }, @@ -104,17 +108,17 @@ pub enum AtprotoError { MalformedRecords(String), #[error("PDS endpoint {pds:?} isn't a usable base URL")] BadPdsEndpoint { pds: String }, - #[error("service token {jti:?} has already been presented")] + #[error("Service token {jti:?} has already been presented")] Replay { jti: JwtNonce }, - #[error("replay-protection store is full and cannot accept another nonce")] + #[error("Replay-protection store is full and can't accept another nonce")] ReplayStoreSaturated, - #[error("issuer {issuer} has too many live replay nonces")] + #[error("Issuer {issuer} has too many live replay nonces")] ReplayShareExhausted { issuer: AccountDid }, #[error(transparent)] Identity(#[from] IdentityError), - #[error("plc submission for {did} returned HTTP {status}")] + #[error("PLC submission for {did} returned HTTP {status}")] PlcSubmit { did: RepoDid, status: HttpStatus }, - #[error("plc log fetch for {did} returned HTTP {status}")] + #[error("PLC log fetch for {did} returned HTTP {status}")] PlcFetch { did: RepoDid, status: HttpStatus }, #[error("putRecord for {subject} returned HTTP {status}")] PutRecord { @@ -126,7 +130,16 @@ pub enum AtprotoError { owner: AccountDid, status: HttpStatus, }, - #[error("pointer record couldn't be encoded: {0}")] + #[error( + "getRecord for label def {rkey} of {owner} returned a body that doesn't parse as a def: \ + {reason}" + )] + MalformedDef { + owner: AccountDid, + rkey: RecordRkey, + reason: String, + }, + #[error("Pointer record couldn't be encoded: {0}")] PointerEncode(String), #[error("putRecord response isn't a valid receipt: {0}")] MalformedReceipt(String), @@ -164,6 +177,7 @@ impl AtprotoError { | AtprotoError::PlcFetch { .. } | AtprotoError::PutRecord { .. } | AtprotoError::GetRecord { .. } + | AtprotoError::MalformedDef { .. } | AtprotoError::PointerEncode(_) | AtprotoError::MalformedReceipt(_) => false, } @@ -196,6 +210,15 @@ pub enum ClaimedKeys { Unread(Arc), } +#[derive(Clone)] +pub enum RemoteDef { + Read { + def: knot_record::label::Def, + bytes: usize, + }, + Unread(Arc), +} + fn claimed_weight(cached: &Cached) -> Weight { Weight::new(match &cached.resolution { ClaimedKeys::Published(keys) => keys @@ -206,6 +229,12 @@ fn claimed_weight(cached: &Cached) -> Weight { ClaimedKeys::Unread(_) => PUBKEY_CACHE_ENTRY_OVERHEAD, }) } +fn def_weight(cached: &Cached) -> Weight { + Weight::new(match &cached.resolution { + RemoteDef::Read { bytes, .. } => *bytes as u64 + DEF_CACHE_ENTRY_OVERHEAD, + RemoteDef::Unread(_) => DEF_CACHE_ENTRY_OVERHEAD, + }) +} pub struct Atproto { http: H, @@ -216,6 +245,7 @@ pub struct Atproto { identities: MokaFuture>, handles: MokaFuture>, claimed: MokaFuture>, + defs: MokaFuture>, seen_jti: Expiring<(AccountDid, JwtNonce), AccountDid, ()>, } @@ -230,6 +260,10 @@ impl Atproto { identities: MokaFuture::by_count(EntryCount::new(MAX_IDENTITY_CACHE as u64)), handles: MokaFuture::by_count(EntryCount::new(MAX_IDENTITY_CACHE as u64)), claimed: MokaFuture::by_weight(Weight::new(MAX_PUBKEY_CACHE_BYTES), claimed_weight), + defs: MokaFuture::by_weight( + Weight::new(MAX_DEF_CACHE_BYTES), + def_weight, + ), seen_jti: Expiring::new(Quotas { per_group: GroupQuota::new(MAX_JTI_PER_ISSUER), total: TotalQuota::new(MAX_SEEN_JTI), @@ -568,6 +602,104 @@ impl Atproto { } } + pub async fn remote_def( + &self, + wanted: &knot_record::label::DefRef, + ) -> RemoteDef { + let now = self.clock.now_unix_micros(); + let filled = self + .defs + .get_or_fill_if( + wanted.clone(), + |cached: &Cached| cached.expires_at <= now, + self.fill_def(wanted, now), + ) + .await; + filled.value.resolution + } + + async fn fill_def( + &self, + wanted: &knot_record::label::DefRef, + now: UnixMicros, + ) -> Cached { + match self.fetch_def(wanted).await { + Ok((value, bytes)) => match knot_record::label::Def::parse(value) { + Ok(def) => Cached { + resolution: RemoteDef::Read { def, bytes }, + expires_at: expires(now, DEF_TTL), + }, + Err(reason) => Cached { + resolution: RemoteDef::Unread(Arc::new( + AtprotoError::MalformedDef { + owner: AccountDid::from(wanted.repo().clone()), + rkey: wanted.rkey().clone(), + reason: reason.to_string(), + }, + )), + expires_at: expires(now, NEGATIVE_TTL), + }, + }, + Err(error) => Cached { + expires_at: if error.is_transient() { + now + } else { + expires(now, NEGATIVE_TTL) + }, + resolution: RemoteDef::Unread(Arc::new(error)), + }, + } + } + + async fn fetch_def( + &self, + wanted: &knot_record::label::DefRef, + ) -> Result<(serde_json::Value, usize), AtprotoError> { + let owner = AccountDid::from(wanted.repo().clone()); + let rkey = wanted.rkey().clone(); + let cached = self.identity(&owner).await?; + let collection: Nsid = Nsid::new_static(knot_record::label::LABEL_DEF_COLLECTION) + .expect("Literal NSID parses"); + let url = get_record_url(&cached.value.pds, &owner, &collection, rkey.as_str())?; + resolve::guard_fetch_url(&url)?; + let response = self.http.execute(HttpRequest::get(url)).await?; + if !response.status.is_success() { + return Err(AtprotoError::GetRecord { + owner: owner.clone(), + status: HttpStatus::from(response.status), + }); + } + if response.body.len() > MAX_REMOTE_DEF_BYTES { + return Err(AtprotoError::MalformedDef { + owner: owner.clone(), + rkey: rkey.clone(), + reason: format!( + "the body is {} bytes, over the {MAX_REMOTE_DEF_BYTES} byte limit for \ + a label def", + response.body.len() + ), + }); + } + let bytes = response.body.len(); + let decoded: serde_json::Value = + serde_json::from_slice(response.body.as_ref()).map_err(|error| { + AtprotoError::MalformedDef { + owner: owner.clone(), + rkey: rkey.clone(), + reason: error.to_string(), + } + })?; + decoded + .get("value") + .cloned() + .map(|value| (value, bytes)) + .ok_or_else(|| AtprotoError::MalformedDef { + owner: owner.clone(), + rkey: rkey.clone(), + reason: "the response is missing a value field".to_owned(), + }) + } + async fn read_record( &self, identity: &Identity, @@ -826,8 +958,8 @@ fn list_records_url( cursor: Option<&Cursor>, ) -> Result { let method: Nsid = - Nsid::new_static("com.atproto.repo.listRecords").expect("literal nsid parses"); - let collection: Nsid = Nsid::new_static("sh.tangled.publicKey").expect("literal nsid parses"); + Nsid::new_static("com.atproto.repo.listRecords").expect("Literal NSID parses"); + let collection: Nsid = Nsid::new_static("sh.tangled.publicKey").expect("Literal NSID parses"); let mut url = xrpc_url(pds, &method)?; url.query_pairs_mut() .append_pair("repo", did.as_str()) @@ -855,7 +987,7 @@ fn get_record_url( collection: &Nsid, rkey: &str, ) -> Result { - let method: Nsid = Nsid::new_static("com.atproto.repo.getRecord").expect("literal nsid parses"); + let method: Nsid = Nsid::new_static("com.atproto.repo.getRecord").expect("Literal NSID parses"); let mut url = xrpc_url(pds, &method)?; url.query_pairs_mut() .append_pair("repo", owner.as_str()) @@ -976,7 +1108,7 @@ mod tests { assert_eq!( dns_hits.load(Ordering::SeqCst), 1, - "a resolved handle is served from cache" + "A resolved handle is served from cache" ); } @@ -989,7 +1121,7 @@ mod tests { Ok(ok(Bytes::from_static(b"did:plc:squid\n"))) } "https://plc.directory/did:plc:squid" => Ok(ok(squid_doc(&signing))), - other => panic!("unexpected url {other}"), + other => panic!("Unexpected URL {other}"), }) }; let empty = resolver(FakeDns::new(|_| Ok(Vec::new())), well_known()); @@ -1018,7 +1150,7 @@ mod tests { "did=did:plc:limpet".to_string(), ]) }), - FakeHttp::new(|_| panic!("resolution must stop before any fetch")), + FakeHttp::new(|_| panic!("Resolution must stop before any fetch")), ); assert!(matches!( ambiguous @@ -1072,7 +1204,7 @@ mod tests { assert_eq!( hits.load(Ordering::SeqCst), 1, - "an unresolvable handle is resolved once then served from the negative cache" + "An unresolvable handle is resolved once, then served from the negative cache" ); let hits = Arc::new(AtomicUsize::new(0)); @@ -1089,7 +1221,7 @@ mod tests { assert_eq!( hits.load(Ordering::SeqCst), 2, - "a transient well-known status mustn't be negatively cached" + "A transient well-known status mustn't be negatively cached" ); } @@ -1450,7 +1582,128 @@ mod tests { assert_eq!( listings.load(Ordering::SeqCst), 2, - "a key published after the last read will be read once the cache turn is over" + "A key published after the last read will be read once the cache turn is over" + ); + } + + #[tokio::test] + async fn remote_def_costs_one_getrecord_per_ttl() { + let signing = signer(11); + let reads = Arc::new(AtomicUsize::new(0)); + let counter = reads.clone(); + let def = serde_json::json!({ + "uri": "at://did:plc:squid/sh.tangled.label.definition/wontfix", + "cid": "bafyreihdwdhef9h1o6xtmxpgpkwp0shftvtjhu2suxcbhelperphpriljku", + "value": { + "$type": "sh.tangled.label.definition", + "name": "wontfix", + "scope": ["sh.tangled.repo.issue"], + "multiple": false, + "createdAt": "2026-09-01T00:00:00Z", + "valueType": {"type": "null", "format": "any"} + } + }); + let http = FakeHttp::new(move |request| { + if request.url.host_str() == Some("plc.directory") { + return Ok(ok(squid_doc(&signing))); + } + assert!(request.url.path().ends_with("com.atproto.repo.getRecord")); + assert!(request.url.query().unwrap().contains("rkey=wontfix")); + counter.fetch_add(1, Ordering::SeqCst); + Ok(ok(Bytes::from(serde_json::to_vec(&def).unwrap()))) + }); + let atproto = Atproto::new(http, clock(), knot_did(KNOT), plc()); + + let wontfix = wanted_def("wontfix"); + let first = atproto.remote_def(&wontfix).await; + assert!( + matches!(&first, RemoteDef::Read { def, .. } if def.name().as_str() == "wontfix"), + "Def comes back read" + ); + let second = atproto.remote_def(&wontfix).await; + assert!(matches!(&second, RemoteDef::Read { .. })); + assert_eq!( + reads.load(Ordering::SeqCst), + 1, + "several operands on one def share a single getRecord" + ); + + atproto + .clock + .advance(DEF_TTL + Duration::from_micros(1)); + let refreshed = atproto.remote_def(&wontfix).await; + assert!(matches!(&refreshed, RemoteDef::Read { .. })); + assert_eq!(reads.load(Ordering::SeqCst), 2); + } + + #[tokio::test] + async fn unreadable_def_caches_unread_for_the_negative_ttl() { + let signing = signer(12); + let reads = Arc::new(AtomicUsize::new(0)); + let counter = reads.clone(); + let absent = FakeHttp::new(move |request| { + if request.url.host_str() == Some("plc.directory") { + return Ok(ok(squid_doc(&signing))); + } + counter.fetch_add(1, Ordering::SeqCst); + Ok(status(StatusCode::NOT_FOUND, Bytes::new())) + }); + let atproto = Atproto::new(absent, clock(), knot_did(KNOT), plc()); + let gone = wanted_def("gone"); + let first = atproto.remote_def(&gone).await; + assert!(matches!(&first, RemoteDef::Unread(_))); + let second = atproto.remote_def(&gone).await; + assert!(matches!(&second, RemoteDef::Unread(_))); + assert_eq!( + reads.load(Ordering::SeqCst), + 1, + "an absent def caches unread for the negative TTL" + ); + + let bloat = serde_json::json!({ + "uri": "at://did:plc:squid/sh.tangled.label.definition/bloat", + "cid": "bafyreihdwdhef9h1o6xtmxpgpkwp0shftvtjhu2suxcbhelperphpriljku", + "value": { + "$type": "sh.tangled.label.definition", + "name": "bloat", + "scope": ["sh.tangled.repo.issue"], + "multiple": false, + "createdAt": "2026-09-01T00:00:00Z", + "valueType": { + "type": "string", + "format": "any", + "enum": (0..8192) + .map(|index| format!("wontfix-{index}")) + .collect::>() + } + } + }); + let signing = signer(13); + let reads = Arc::new(AtomicUsize::new(0)); + let counter = reads.clone(); + let serving = FakeHttp::new(move |request| { + if request.url.host_str() == Some("plc.directory") { + return Ok(ok(squid_doc(&signing))); + } + counter.fetch_add(1, Ordering::SeqCst); + Ok(ok(Bytes::from(serde_json::to_vec(&bloat).unwrap()))) + }); + let atproto = Atproto::new(serving, clock(), knot_did(KNOT), plc()); + let bloat = wanted_def("bloat"); + let first = atproto.remote_def(&bloat).await; + let RemoteDef::Unread(error) = &first else { + panic!("An oversized def comes back unread"); + }; + assert!( + matches!(&**error, AtprotoError::MalformedDef { .. }), + "the unread error for an oversized body is MalformedDef" + ); + let second = atproto.remote_def(&bloat).await; + assert!(matches!(&second, RemoteDef::Unread(_))); + assert_eq!( + reads.load(Ordering::SeqCst), + 1, + "the unread error caches for the negative TTL" ); } @@ -1970,7 +2223,7 @@ mod tests { atproto .verify_service_jwt(&mint(&signing, &bystander), &method) .await - .expect("unrelated issuer is unaffected by the hog"); + .expect("Unrelated issuer is unaffected by the hog"); atproto.clock.advance(Duration::from_secs(62)); let after_expiry = serde_json::json!({ @@ -1980,7 +2233,7 @@ mod tests { atproto .verify_service_jwt(&mint(&signing, &after_expiry), &method) .await - .expect("hog recovers once its nonces expire from the store"); + .expect("Hog recovers once its nonces expire from the store"); } #[tokio::test] @@ -2042,7 +2295,7 @@ mod tests { assert_eq!( hot_hits.load(Ordering::SeqCst), before, - "frequently resolved identity is retained through the flood" + "Frequently resolved identity is retained through the flood" ); let novel = did("did:web:sentinel.oyster.cafe"); atproto.resolve_identity(&novel).await.unwrap(); @@ -2050,12 +2303,12 @@ mod tests { assert_eq!( sentinel_hits.load(Ordering::SeqCst), 1, - "saturated cache still admits a new entry and reuses it without a re-fetch" + "Saturated cache still accepts a new entry and reuses it without a re-fetch" ); atproto.identities.run_pending_tasks().await; assert!( atproto.identities.entry_count().get() <= MAX_IDENTITY_CACHE as u64, - "cache stays bounded after eviction" + "Cache stays bounded after eviction" ); } @@ -2120,7 +2373,7 @@ mod tests { #[tokio::test] async fn a_concurrent_429_wave_gets_the_real_error_and_never_poisons_the_cache() { let (results, calls) = gated_wave(StatusCode::TOO_MANY_REQUESTS, Bytes::new()).await; - assert_eq!(calls, 1, "failing wave coalesces into one outbound fetch"); + assert_eq!(calls, 1, "Failing wave coalesces into one outbound fetch"); assert!( results.iter().all(|outcome| matches!( outcome, @@ -2145,8 +2398,8 @@ mod tests { .count(); assert_eq!( gone, 2000, - "coalescing mustn't decide which callers learn the account is missing, since the \ - callers served from the cache act on the answer the same way" + "Coalescing mustn't decide which callers learn the account is missing; \ + every caller sees the same answer" ); } @@ -2298,7 +2551,7 @@ mod tests { .get(http::header::AUTHORIZATION) .and_then(|value| value.to_str().ok()) .and_then(|value| value.strip_prefix("Bearer ")) - .expect("request includes a bearer service token"); + .expect("Request includes a bearer service token"); let parsed = knot_types::service_auth::parse_jwt(bearer).unwrap(); assert_eq!(parsed.claims().iss.as_str(), KNOT); assert_eq!(parsed.claims().aud.as_str(), "did:web:pds.oyster.cafe"); @@ -2310,7 +2563,7 @@ mod tests { let key = knot_types::service_auth::PublicKey::from_k256_bytes(knot_public.as_bytes()) .unwrap(); knot_types::service_auth::verify_signature(&parsed, &key) - .expect("token is signed by the knot key"); + .expect("Token is signed by the knot key"); let body: serde_json::Value = serde_json::from_slice(request.body.as_ref().unwrap()).unwrap(); assert_eq!(body["repo"], SQUID); diff --git a/knot2/crates/knot-atproto/src/test_support.rs b/knot2/crates/knot-atproto/src/test_support.rs index 3e6e03101..b8345e1db 100644 --- a/knot2/crates/knot-atproto/src/test_support.rs +++ b/knot2/crates/knot-atproto/src/test_support.rs @@ -27,6 +27,14 @@ pub(crate) fn repo_did(value: &str) -> RepoDid { RepoDid::new(value).unwrap() } +pub(crate) fn wanted_def(rkey: &str) -> knot_record::label::DefRef { + knot_record::label::DefRef::parse( + &knot_types::AtUri::new_owned(format!("at://{SQUID}/sh.tangled.label.definition/{rkey}")) + .unwrap(), + ) + .unwrap() +} + pub(crate) fn owner_did(value: &str) -> OwnerDid { OwnerDid::new(value).unwrap() } @@ -163,7 +171,7 @@ pub(crate) fn mint_with_header(signing: &SigningKey, header: &[u8], claims: &Val "{signing_input}.{}", URL_SAFE_NO_PAD.encode(signature.to_bytes()) )) - .expect("minted token is structurally a JWT") + .expect("Minted token is structurally a JWT") } pub(crate) fn runtime_signer(seed: u64) -> K256Signer { diff --git a/knot2/crates/knot-xrpc/src/social.rs b/knot2/crates/knot-xrpc/src/social.rs index 122c14850..9a631d17b 100644 --- a/knot2/crates/knot-xrpc/src/social.rs +++ b/knot2/crates/knot-xrpc/src/social.rs @@ -5,6 +5,8 @@ use axum::Json; use axum::body::Bytes; use axum::extract::State; use axum::response::{IntoResponse, Response}; +use futures::StreamExt; +use futures::TryStreamExt; use http::HeaderMap; use jacquard_common::types::tid::Ticker; use knot_acl::{KnotAcl, RepoAccess}; @@ -13,8 +15,8 @@ use knot_cob::{ }; use knot_cobs::{ Comment, CommentChange, Creation, Editing, Erasure, Issue, IssueChange, IssueStateChange, - IssueStateCob, Reaction, ReactionKind, SocialBody, SocialChange, SocialCob, SocialObject, - Transition, + IssueStateCob, LabelDef, LabelDefChange, LabelOp, Opening, Reaction, + ReactionKind, SocialBody, SocialChange, SocialCob, SocialObject, Transition, }; use knot_git::Repo; use knot_index::{Resolved, SocialEntry}; @@ -26,6 +28,10 @@ use knot_record::issue::{ IssueProse, IssueRecord, IssueState, IssueStateRecord, IssueTitle, MarkdownBody, issue_collection, issue_state_collection, }; +use knot_record::label::{ + Def, DefRef, LabelDefRecord, LabelOpRecord, LabelOperand, + label_def_collection, label_op_collection, +}; use knot_record::reaction::{ReactionRecord, reaction_collection}; use knot_resource::StructureBytes; use knot_runtime::{Clock, HttpTransport, Signer}; @@ -45,9 +51,9 @@ pub(crate) const PUT_RECORD_ROUTE: &str = "/xrpc/com.atproto.repo.putRecord"; pub(crate) const DELETE_RECORD_ROUTE: &str = "/xrpc/com.atproto.repo.deleteRecord"; pub(crate) const LIST_COB_VERSIONS_ROUTE: &str = "/xrpc/sh.tangled.repo.listCobVersions"; -const NOT_STATE_AUTHORITY: &str = "state changes belong to the issue's author or a maintainer"; +const NOT_STATE_AUTHORITY: &str = "State changes belong to the issue's author or a maintainer"; -const COLLABORATE_REMEDY: &str = "ask the repository owner to add you as a collaborator, or to open the contribution policy to members or anyone"; +const COLLABORATE_REMEDY: &str = "Ask the repository owner to add you as a collaborator, or to open the contribution policy to members or anyone"; pub struct RecordKeys(Mutex); @@ -135,7 +141,7 @@ fn typed( collection: &RecordCollection, ) -> Result { serde_json::from_value(record).map_err(|error| { - XrpcError::invalid_request(format!("record for {collection} is unusable: {error}")) + XrpcError::invalid_request(format!("Record for {collection} is unusable: {error}")) }) } @@ -247,7 +253,7 @@ impl MarkupInput { } let unusable = |what: &'static str| { move |error| { - XrpcError::invalid_request(format!("comment body's {what} is unusable: {error}")) + XrpcError::invalid_request(format!("Comment body's {what} is unusable: {error}")) } }; let text = MarkdownBody::new(self.text).map_err(unusable("text"))?; @@ -342,6 +348,39 @@ impl ReactionDraft { } } +fn parsed_def(record: serde_json::Value) -> Result { + Def::parse_written(record).map_err(|error| XrpcError::invalid_request(error.to_string())) +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct LabelOpDraft { + #[serde(rename = "$type", default)] + kind: Option, + subject: AtUri, + #[serde(default)] + add: Vec, + #[serde(default)] + delete: Vec, + #[serde(rename = "performedAt")] + _performed_at: serde::de::IgnoredAny, + #[serde(rename = "x-tngl-editor", default)] + _editor: serde::de::IgnoredAny, +} + +impl LabelOpDraft { + fn parse( + self, + written_to: &RepoDid, + collection: &RecordCollection, + ) -> Result<(RecordAddress, Vec, Vec), XrpcError> { + kind_matches(self.kind.as_ref(), collection)?; + let subject = any_target(&self.subject)?; + on_this_knot(&subject.repo, written_to, "subject")?; + Ok((subject.address, self.add, self.delete)) + } +} + fn on_this_knot(repo: &RepoDid, written_to: &RepoDid, what: &'static str) -> Result<(), XrpcError> { if repo == written_to { return Ok(()); @@ -423,7 +462,7 @@ async fn permission( repo: &RepoDid, ) -> Result { let policy = ready(state.index.contribution_policy(repo), "registry")? - .ok_or_else(|| XrpcError::not_found("no such repository on this knot"))?; + .ok_or_else(|| XrpcError::not_found("No such repository on this knot"))?; project_collaborators(state, repo).await?; let pds = crate::pds(state); let consents = state.index.consents(); @@ -432,7 +471,7 @@ async fn permission( .await; let membership = consents.of_membership(&pds, actor, state.now()).await; let access = RepoAccess::of(&collaboration, &membership, &policy) - .ok_or_else(|| XrpcError::internal("consents answered with mismatched subjects"))?; + .ok_or_else(|| XrpcError::internal("Consents answered with mismatched subjects"))?; let acl = KnotAcl::new(&state.admins, state.admission, &state.index); ready(knot_acl::permission_of(&acl, &access), "roll") } @@ -452,7 +491,7 @@ fn next_number( collection: &RecordCollection, ) -> Result { ready(state.index.next_subject_number(repo, collection), "social")? - .ok_or_else(|| XrpcError::internal("this repository has numbered past what fits")) + .ok_or_else(|| XrpcError::internal("This repository has numbered past what fits")) } fn object_at( @@ -544,7 +583,7 @@ impl Ledger<'_, H, C> { repo = %self.repo, uri = %address.at_uri(self.repo), %drift, - "the chain stood where the projection didn't expect it, so the write reconciled the chain to the projection before answering" + "Chain stood where the projection didn't expect it, so the write reconciled the chain to the projection before answering" ); reconcile(self.state, self.git, self.repo, self.signer, self.now) .and_then(|_| self.head()) @@ -565,7 +604,7 @@ impl Ledger<'_, H, C> { fn head(&self) -> Result { let tip = chain::tip(self.git)?.ok_or_else(|| { - XrpcError::internal("a repository that took a social write has a chain tip") + XrpcError::internal("A repository that took a social write has a chain tip") })?; Ok(CommitHead { cid: chain::commit_cid(&chain::commit_bytes(self.git, &tip)?)?.into(), @@ -600,14 +639,14 @@ impl Ledger<'_, H, C> { repo = %self.repo, %cause, commits = reconciled.commits, - "a social write failed after its change was committed, so the chain was \ + "A social write failed after its change was committed, so the chain was \ reconciled to the projection" ), Err(error) => tracing::error!( repo = %self.repo, %cause, %error, - "a social write failed after its change was committed and the reconciler failed too; the next restart reconciles" + "A social write failed after its change was committed and the reconciler failed too; the next restart reconciles" ), } } @@ -684,10 +723,7 @@ impl Ledger<'_, H, C> { ) .map_err(|error| match (verb, error) { (Some(verb), CobError::AppendRejected { .. }) => { - XrpcError::forbidden(format!( - "sorry, only its author or a moderator can {verb} this {}", - noun_of(address.collection()) - )) + append_denial(B::OPENING, verb, address.collection()) } (_, other) => other.into(), })?, @@ -784,7 +820,7 @@ pub(crate) fn reconcile( repo = %repo, uri = %entry.address.at_uri(repo), %version, - "a live social record is missing its body in the change tree; the chain keeps its served leaf and leaves the body out" + "A live social record is missing its body in the change tree; the chain keeps its served leaf and leaves the body out" ); (leaf.and_then(RecordCid::new), None) } @@ -799,7 +835,7 @@ pub(crate) fn reconcile( tracing::error!( repo = %repo, records = serving.len(), - "social records exist on a repository whose chain was never enabled, so they stay off the atproto surface" + "Social records exist on a repository whose chain was never enabled, so they stay off the atproto surface" ); } return Ok(Reconciled { commits: 0 }); @@ -915,6 +951,7 @@ pub(crate) async fn create_record( issue_collection(), comment_collection(), reaction_collection(), + label_op_collection(), ] .contains(&collection) { @@ -930,6 +967,10 @@ pub(crate) async fn create_record( comment_open(&state, &headers, &method, repo, record).await } else if collection == reaction_collection() { reaction_open(&state, &headers, &method, repo, record).await + } else if collection == label_def_collection() { + label_def_open(&state, &headers, &method, repo, rkey, record).await + } else if collection == label_op_collection() { + label_op_open(&state, &headers, &method, repo, record).await } else { Err(rejected_collection(&collection)) } @@ -944,6 +985,35 @@ async fn open( body: B, admit: F, ) -> Result +where + B: SocialBody + Send + 'static, + H: HttpTransport, + C: Clock, + F: FnOnce(&Ledger<'_, H, C>, &WriteContext<'_>) -> Result, XrpcError> + + Send + + 'static, +{ + open_at( + state, + session, + RecordAddress::new(collection.clone(), state.rkeys.mint()), + parent, + record, + body, + admit, + ) + .await +} + +async fn open_at( + state: &Arc>, + session: Session, + address: RecordAddress, + parent: Option, + record: RecordBody, + body: B, + admit: F, +) -> Result where B: SocialBody + Send + 'static, H: HttpTransport, @@ -953,11 +1023,12 @@ where + 'static, { let produces = record.cid(); - let address = RecordAddress::new(collection.clone(), state.rkeys.mint()); + let collection = address.collection().clone(); on_repo(state, session, move |ledger, cx| { if let Some(answered) = admit(ledger, cx)? { return Ok(answered); } + let change = SocialChange::Open { creation: Creation { address: address.clone(), @@ -1050,10 +1121,10 @@ fn standing_version( ) -> Result<(), XrpcError> { let serving = live_record_at(ledger, reference.address())? .version() - .ok_or_else(|| XrpcError::internal(format!("the current {what} is missing its version")))?; + .ok_or_else(|| XrpcError::internal(format!("The current {what} is missing its version")))?; if serving != reference.cid() { return Err(XrpcError::invalid_request(format!( - "the {what}'s strongRef cid is {}, but the record at {} now has version {serving}", + "The {what}'s strongRef CID is {}, but the record at {} now has version {serving}", reference.cid(), reference.address().at_uri(ledger.repo), ))); @@ -1105,10 +1176,10 @@ async fn reaction_open( "social", )? .ok_or_else(|| { - XrpcError::internal("the matching reaction is missing its address") + XrpcError::internal("The matching reaction is missing its address") })?; let version = live_record_at(ledger, &address)?.version().ok_or_else(|| { - XrpcError::internal("the matching reaction should have a version and doesn't") + XrpcError::internal("The matching reaction should have a version and doesn't") })?; return Ok(Some(wrote(ledger.repo, &address, version, &ledger.head()?))); } @@ -1122,10 +1193,11 @@ fn noun_of(collection: &RecordCollection) -> &'static str { knot_record::issue::ISSUE_COLLECTION => "issue", knot_record::comment::COMMENT_COLLECTION => "comment", knot_record::reaction::REACTION_COLLECTION => "reaction", + knot_record::label::LABEL_DEF_COLLECTION => "label def", + knot_record::label::LABEL_OP_COLLECTION => "label op", _ => "record", } } - fn standing( ledger: &Ledger<'_, H, C>, address: &RecordAddress, @@ -1189,10 +1261,18 @@ fn materialized( ledger: &Ledger<'_, H, C>, address: &RecordAddress, object: CobId, +) -> Result<(knot_cobs::SocialState, ChangeId), XrpcError> { + folded_at(ledger.git, address, object) +} + +fn folded_at( + git: &Repo, + address: &RecordAddress, + object: CobId, ) -> Result<(knot_cobs::SocialState, ChangeId), XrpcError> { let type_name = TypeName::new(address.collection().as_str()) .expect("Record collection is a valid type name"); - knot_cobs::social_state(&CobStore::new(ledger.git), &type_name, object) + knot_cobs::social_state(&CobStore::new(git), &type_name, object) .map_err(XrpcError::from)? .ok_or_else(|| { XrpcError::internal(format!( @@ -1212,11 +1292,11 @@ fn opened_comment( .changes .into_iter() .next() - .ok_or_else(|| XrpcError::internal("a comment object's root change went missing"))?; + .ok_or_else(|| XrpcError::internal("A comment object's root change went missing"))?; match CommentChange::decode(root.payload()) { Ok(SocialChange::Open { body, .. }) => Ok(body), Ok(_) => Err(XrpcError::internal( - "a comment object's root change isn't the change that created it", + "A comment object's root change isn't the change that created it", )), Err(error) => Err(XrpcError::internal(error.to_string())), } @@ -1231,10 +1311,14 @@ pub(crate) async fn put_record( let SwapInput { collection, .. } = &swap; match collection.as_str() { knot_record::issue::ISSUE_STATE_COLLECTION => Err(XrpcError::invalid_request( - "a state transition is created, never edited", + "A state transition is created, never edited", )), knot_record::reaction::REACTION_COLLECTION => Err(XrpcError::invalid_request( - "createRecord makes a reaction and deleteRecord renoves it; a reaction is never edited", + "createRecord makes a reaction and deleteRecord removes it; a reaction is never edited", + )), + knot_record::label::LABEL_OP_COLLECTION => Err(XrpcError::invalid_request( + "createRecord writes a label op and deleteRecord retracts the label op; putRecord \ + can't edit label ops, and the current label set is the fold of the ops that stand", )), knot_record::issue::ISSUE_COLLECTION => { issue_put(&state, &headers, &method, swap, record).await @@ -1242,6 +1326,9 @@ pub(crate) async fn put_record( knot_record::comment::COMMENT_COLLECTION => { comment_put(&state, &headers, &method, swap, record).await } + knot_record::label::LABEL_DEF_COLLECTION => { + label_def_put(&state, &headers, &method, swap, record).await + } _ => Err(rejected_collection(&swap.collection)), } } @@ -1393,6 +1480,12 @@ pub(crate) async fn delete_record( knot_record::reaction::REACTION_COLLECTION => { erase_of::(&state, &headers, &method, swap).await } + knot_record::label::LABEL_DEF_COLLECTION => { + erase_of::(&state, &headers, &method, swap).await + } + knot_record::label::LABEL_OP_COLLECTION => { + erase_of::(&state, &headers, &method, swap).await + } _ => Err(rejected_collection(&swap.collection)), } } @@ -1432,11 +1525,11 @@ fn current_state(store: &CobStore, object: CobId) -> Result Ok(body.state), SocialChange::Edit { .. } | SocialChange::Erase { .. } => Err(XrpcError::internal( - "a state object opens with a transition", + "A state object opens with a transition", )), } } @@ -1459,7 +1552,7 @@ fn state_of( ledger.state.index.social_object_at(ledger.repo, at), "social", )? - .ok_or_else(|| XrpcError::internal("the latest transition has its object")) + .ok_or_else(|| XrpcError::internal("The latest transition has its object")) .and_then(|object| current_state(&CobStore::new(ledger.git), object)) }) .transpose() @@ -1519,6 +1612,280 @@ async fn state_open( .await } +const NOT_LABEL_AUTHORITY: &str = "label defs belong to the repository's owner alone"; + +const NOT_TRIAGE_AUTHORITY: &str = "labels belong to the repository's owner and collaborators"; + +const TRIAGE_REMEDY: &str = "ask the repository owner to add you as a collaborator"; + +fn label_denial(opening: Opening, permission: ContributionPermission) -> XrpcError { + match (opening, permission.moderates().is_allowed()) { + (Opening::Ownership, _) => XrpcError::forbidden(NOT_LABEL_AUTHORITY), + (Opening::Triage, true) => XrpcError::forbidden(NOT_TRIAGE_AUTHORITY), + (Opening::Triage, false) => { + XrpcError::forbidden(format!("{NOT_TRIAGE_AUTHORITY}. {TRIAGE_REMEDY}")) + } + (Opening::Contribution | Opening::Moderation, _) => { + XrpcError::internal("A label denial arrived under an opening without label rules") + } + } +} + +fn require_label_permission( + opening: Opening, + permission: ContributionPermission, +) -> Result<(), XrpcError> { + if opening.decide(permission).is_allowed() { + Ok(()) + } else { + Err(label_denial(opening, permission)) + } +} + +fn append_denial(opening: Opening, verb: &str, collection: &RecordCollection) -> XrpcError { + match opening { + Opening::Ownership => XrpcError::forbidden(NOT_LABEL_AUTHORITY), + _ => XrpcError::forbidden(format!( + "Sorry, only the author or a moderator can {verb} this {}", + noun_of(collection) + )), + } +} + +async fn label_def_open( + state: &Arc>, + headers: &HeaderMap, + method: &crate::Method, + repo: RepoDid, + rkey: Option, + record: serde_json::Value, +) -> Result { + let collection = label_def_collection(); + let slug = rkey + .and_then(|rkey| rkey.as_label_key().cloned()) + .ok_or_else(|| { + XrpcError::invalid_request( + "a label def's record key is the slug the owner chooses; send that slug \ + as the rkey, and a minted TID is refused", + ) + })?; + let def = parsed_def(record)?; + let session = session(state, headers, method, repo).await?; + require_label_permission(Opening::Ownership, session.permission)?; + let record = + LabelDefRecord::new(def, session.now, session.actor.clone())?.body()?; + let address = RecordAddress::new(collection, RecordRkey::composed(slug)); + open_at( + state, + session, + address.clone(), + None, + record, + LabelDef, + move |ledger, _| match object_at(ledger.state, ledger.repo, &address) { + Ok(object) => match serving_object(ledger, &address, object) { + Ok(Some(_)) => Err(XrpcError::conflict(format!( + "a label def already stands at {}", + address.at_uri(ledger.repo) + ))), + Ok(None) => Err(XrpcError::conflict(format!( + "a label def was erased at {}, and an erased def's slug stays \ + retired while the versions stay in the change graph", + address.at_uri(ledger.repo) + ))), + Err(interrupted) => Err(interrupted), + }, + Err(missing) if missing.status() == http::StatusCode::NOT_FOUND => Ok(None), + Err(other) => Err(other), + }, + ) + .await +} + +async fn label_def_put( + state: &Arc>, + headers: &HeaderMap, + method: &crate::Method, + swap: SwapInput, + record: serde_json::Value, +) -> Result { + let collection = label_def_collection(); + let (repo, rkey, from) = swap.edit_of(&collection)?; + let def = parsed_def(record)?; + let address = RecordAddress::new(collection, rkey); + let session = session(state, headers, method, repo).await?; + on_repo(state, session, move |ledger, cx| { + let (object, live) = standing(ledger, &address)?; + let record = LabelDefRecord::new(def, live.created_at(), cx.actor().clone())? + .body()?; + let produces = record.cid(); + if live.version() == Some(produces) && from == produces { + return Ok(wrote(ledger.repo, &address, produces, &ledger.head()?)); + } + let change: LabelDefChange = SocialChange::Edit { + editing: Editing { + author: cx.actor().clone(), + from, + produces, + edited_at: ledger.now, + }, + body: LabelDef, + }; + let (_, _, head) = ledger.write( + cx, + &address, + change, + Doing::Edit { + object, + from, + body: &record, + }, + )?; + Ok(wrote(ledger.repo, &address, produces, &head)) + }) + .await +} + +const DEF_FETCH_WIDTH: usize = 8; + +async fn remote_defs( + state: &XrpcState, + wanted: &[DefRef], +) -> Result, XrpcError> { + futures::stream::iter(wanted.to_vec()) + .map(|wanted| async move { + let def = match state.atproto.remote_def(&wanted).await { + knot_atproto::RemoteDef::Read { def, .. } => def.clone(), + knot_atproto::RemoteDef::Unread(error) => { + return Err(XrpcError::bad_gateway(match &*error { + knot_atproto::AtprotoError::MalformedDef { .. } => format!( + "the def at {wanted} doesn't parse, so the knot won't store \ + the citing op: {error}" + ), + _ => format!( + "the PDS didn't serve the label def at {wanted}, so the knot \ + won't store the citing op: {error}" + ), + })); + } + }; + Ok((wanted, def)) + }) + .buffered(DEF_FETCH_WIDTH) + .try_collect::>() + .await +} + +fn local_def( + index: &Arc, + layout: &knot_git::Layout, + wanted: &DefRef, +) -> Result { + let repo = wanted.repo(); + let address = RecordAddress::new(label_def_collection(), wanted.rkey().clone()); + index.ensure_social(repo).map_err(XrpcError::from)?; + let object = ready(index.social_object_at(repo, &address), "social")?.ok_or_else(|| { + XrpcError::invalid_request(format!( + "this repository doesn't have a label def at {wanted}" + )) + })?; + let git = layout.open(repo)?; + let store = CobStore::new(&git); + let live = folded_at(&git, &address, object)? + .0 + .object() + .filter(|live| !live.is_erased()) + .cloned() + .ok_or_else(|| { + XrpcError::invalid_request(format!( + "the label def at {wanted} was deleted from this repository" + )) + })?; + let version = live + .version() + .ok_or_else(|| XrpcError::internal("A live def is missing a version"))?; + let body = tree_body(&store, address.collection(), object, version)?.ok_or_else(|| { + XrpcError::internal("A live def is missing the record bytes in the change tree") + })?; + let def = Def::decode(body.as_bytes()).map_err(|error| { + XrpcError::invalid_request(format!("The def at {wanted} doesn't parse: {error}")) + })?; + Ok(def) +} + +async fn label_op_open( + state: &Arc>, + headers: &HeaderMap, + method: &crate::Method, + repo: RepoDid, + record: serde_json::Value, +) -> Result { + let collection = label_op_collection(); + let (subject, adds, deletes) = + typed::(record, &collection)?.parse(&repo, &collection)?; + let session = session(state, headers, method, repo).await?; + require_label_permission(Opening::Triage, session.permission)?; + warm_social(state, &session.repo).await?; + ready( + state.index.social_object_at(&session.repo, &subject), + "social", + )? + .ok_or_else(|| { + XrpcError::invalid_request(format!( + "{} isn't a record on this knot", + subject.at_uri(&session.repo) + )) + })?; + let mut wanted: Vec = adds + .iter() + .chain(deletes.iter()) + .map(|operand| operand.key().clone()) + .collect(); + wanted.sort_unstable(); + wanted.dedup(); + let index = state.index.clone(); + let layout = state.layout.clone(); + let (mut defs, remote) = run_blocking(move || { + let (local, remote): (Vec, Vec) = wanted + .into_iter() + .partition(|wanted| layout.open(wanted.repo()).is_ok()); + let mut defs = BTreeMap::new(); + local.iter().try_for_each(|wanted| { + let def = local_def(&index, &layout, wanted)?; + defs.insert(wanted.clone(), def); + Ok::<(), XrpcError>(()) + })?; + Ok((defs, remote)) + }) + .await?; + defs.extend(remote_defs(state, &remote).await?); + let op = LabelOpRecord::new( + session.repo.clone(), + subject.clone(), + adds, + deletes, + session.now, + session.actor.clone(), + ) + .map_err(|error| XrpcError::invalid_request(error.to_string()))?; + let record = op.body()?; + let body = LabelOp::applications(&op, &defs) + .map_err(|error| XrpcError::invalid_request(error.to_string()))?; + open( + state, + session, + collection, + Some(subject.clone()), + record, + body, + move |ledger, _| { + live_record_at(ledger, &subject)?; + Ok(None) + }, + ) + .await +} + const COB_VERSIONS_DEFAULT: usize = 100; const COB_VERSIONS_MAX: usize = 500; @@ -1563,7 +1930,7 @@ pub(crate) async fn list_cob_versions( ))); } ready(state.index.ownership_of(&repo), "registry")? - .ok_or_else(|| XrpcError::not_found("no such repository on this knot"))?; + .ok_or_else(|| XrpcError::not_found("No such repository on this knot"))?; warm_social(&state, &repo).await?; let object = object_at(&state, &repo, &address)?; let index_state = Arc::clone(&state); @@ -1571,7 +1938,7 @@ pub(crate) async fn list_cob_versions( let response_limit = state.byte_limits.response.get(); let (limit, cursor) = (query.limit, query.cursor); let type_name = knot_types::TypeName::new(collection.as_str()) - .map_err(|_| XrpcError::invalid_request("the uri's collection isn't a type name"))?; + .map_err(|_| XrpcError::invalid_request("The URI's collection isn't a type name"))?; let (versions, total) = run_blocking(move || { let git = index_state.layout.open(&repo)?; let store = CobStore::new(&git);