diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index b81240f36..657fd0c07 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -29,6 +29,7 @@ use bobbin_types::edges::Edge; use bobbin_types::ids::EdgeKey; use either::Either; use jacquard_common::DefaultStr; +use jacquard_common::deps::smol_str::SmolStr; use jacquard_common::types::string::AtUri; use lasso::{Key, Spur, ThreadedRodeo}; use scc::HashMap as SccMap; @@ -427,6 +428,7 @@ pub struct EdgeStore { next_key_id: AtomicU32, forward: SccMap, reverse: SccMap, RuntimeHasher>, + source_collections: SccMap, RuntimeHasher>, hasher: RuntimeHasher, writer: Mutex<()>, } @@ -441,6 +443,7 @@ impl EdgeStore { next_key_id: AtomicU32::new(0), forward: SccMap::with_hasher(hasher.clone()), reverse: SccMap::with_hasher(hasher.clone()), + source_collections: SccMap::with_hasher(hasher.clone()), hasher, writer: Mutex::new(()), } @@ -494,6 +497,39 @@ impl EdgeStore { self.clear_source_locked(source); } + /// Store collection filters for a source. None or empty = all collections (removes entry). + /// Call only after the source has been interned via add/upsert_source. + pub fn set_source_collections( + &self, + source: &AtUri, + collections: Option>, + ) { + let author = source_authority_did(source) + .and_then(|s| self.did_interner.get(s)) + .map(AuthorId::from_spur); + let Some(spur) = self.source_interner.get(self.source_key(source, author)) else { + return; + }; + let id = SourceId::from_spur(spur); + if let Some(cols) = collections + && !cols.is_empty() + { + let _ = self.source_collections.insert_sync(id, cols); + } else { + self.source_collections.remove_sync(&id); + } + } + + /// Returns None if no collection filter is stored (meaning all collections). + pub fn source_collections_for(&self, source: &AtUri) -> Option> { + let author = source_authority_did(source) + .and_then(|s| self.did_interner.get(s)) + .map(AuthorId::from_spur); + let spur = self.source_interner.get(self.source_key(source, author))?; + let id = SourceId::from_spur(spur); + self.source_collections.read_sync(&id, |_, cols| cols.clone()) + } + fn add_locked(&self, edge: Edge) { let author = self.intern_author(&edge.source); let source_key = self.source_key(&edge.source, author); @@ -540,6 +576,7 @@ impl EdgeStore { return; }; let id = SourceId::from_spur(source_spur); + self.source_collections.remove_sync(&id); let Some((_, entries)) = self.reverse.remove_sync(&id) else { return; }; @@ -1614,4 +1651,62 @@ mod tests { store.remove_source(&at("at://did:plc:nel/sh.tangled.feed.star/r99")); assert_eq!(store.count_distinct_authors(&key), 4, "nel fully removed"); } + + #[test] + fn source_collections_set_and_clear() { + let store = store(); + let source = at("at://did:plc:bob/sh.tangled.feed.subscription/r1"); + let edge = Edge { + kind: nsid("sh.tangled.feed.subscription"), + subject: did_subj("did:plc:repo"), + source: source.clone(), + sort_micros: 1, + }; + store.add(edge); + + // Initially no filter + assert_eq!(store.source_collections_for(&source), None); + + // Set a filter + store.set_source_collections(&source, Some(vec![ + SmolStr::new_static("sh.tangled.repo.issue"), + ])); + let stored = store.source_collections_for(&source).unwrap(); + assert_eq!(stored, vec![SmolStr::new_static("sh.tangled.repo.issue")]); + + // Update to empty = unrestricted + store.set_source_collections(&source, Some(vec![])); + assert_eq!(store.source_collections_for(&source), None); + + // Set a different filter + store.set_source_collections(&source, Some(vec![ + SmolStr::new_static("sh.tangled.repo.pull"), + ])); + let stored = store.source_collections_for(&source).unwrap(); + assert_eq!(stored, vec![SmolStr::new_static("sh.tangled.repo.pull")]); + + // Set to None = unrestricted + store.set_source_collections(&source, None); + assert_eq!(store.source_collections_for(&source), None); + } + + #[test] + fn remove_source_cleans_up_collection_filter() { + let store = store(); + let source = at("at://did:plc:bob/sh.tangled.feed.subscription/r1"); + let edge = Edge { + kind: nsid("sh.tangled.feed.subscription"), + subject: did_subj("did:plc:repo"), + source: source.clone(), + sort_micros: 1, + }; + store.add(edge); + store.set_source_collections(&source, Some(vec![ + SmolStr::new_static("sh.tangled.repo.issue"), + ])); + assert!(store.source_collections_for(&source).is_some()); + + store.remove_source(&source); + assert_eq!(store.source_collections_for(&source), None); + } } diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index abed14743..fe7bb1dce 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -1379,6 +1379,9 @@ async fn finalize_drained( remove_superseded_source(ctx.store, ctx.records, ctx.search, supersedes.as_ref()).await; cache_body(ctx.records, &source, cid, bytes); ctx.store.upsert_source(&source, edges); + if let Some(filter) = parsed.subscription_collections() { + ctx.store.set_source_collections(&source, filter); + } let outcome = apply_record_state(ctx.issue_states, ctx.pull_statuses, &source, &parsed); log_unknown_state_variant(outcome, &source); index_search(ctx.search, ctx.resolver, &source, parsed).await; @@ -1418,6 +1421,9 @@ async fn commit_pending( remove_superseded_source(store, records, search, supersedes.as_ref()).await; cache_body(records, &source, cid, bytes); store.upsert_source(&source, edges); + if let Some(filter) = parsed.subscription_collections() { + store.set_source_collections(&source, filter); + } let outcome = apply_record_state(issue_states, pull_statuses, &source, &parsed); log_unknown_state_variant(outcome, &source); index_search(search, resolver, &source, parsed).await; diff --git a/bobbin/crates/types/src/edges.rs b/bobbin/crates/types/src/edges.rs index 757aa5e28..001d4f932 100644 --- a/bobbin/crates/types/src/edges.rs +++ b/bobbin/crates/types/src/edges.rs @@ -1,6 +1,7 @@ use alloc::vec; use alloc::vec::Vec; +use jacquard_common::deps::smol_str::SmolStr; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::{AtStrError, AtUri, Datetime}; @@ -167,6 +168,15 @@ impl Record { Ok(append_mirror_edges(primary, source)) } + /// None = not a subscription. Some(None/Some([])) = unrestricted (all collections). + /// Some(Some(non_empty)) = filtered to those collection NSIDs. + pub fn subscription_collections(&self) -> Option>> { + match self { + Self::Subscription(r) => Some(r.collections.clone()), + _ => None, + } + } + pub fn sort_micros_for(&self, source: &AtUri) -> u64 { if let Some(dt) = self.created_at() { let micros = dt.timestamp_micros(); @@ -666,6 +676,44 @@ mod tests { assert!(edges.is_empty()); } + #[test] + fn subscription_collections_tri_state() { + // Non-subscription record (Star) → None + let star = Record::from_json_value( + &nsid("sh.tangled.feed.star"), + json!({ + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": { "$type": "sh.tangled.feed.star#repo", "did": "did:plc:abalone" }, + }), + ).unwrap(); + assert_eq!(star.subscription_collections(), None); + + // Unrestricted subscription (no collections field) → Some(None) + let sub_none = Record::from_json_value( + &nsid("sh.tangled.feed.subscription"), + json!({ + "$type": "sh.tangled.feed.subscription", + "createdAt": "2026-05-01T00:00:00Z", + "subject": { "$type": "sh.tangled.feed.subscription#repo", "did": "did:plc:repo" }, + }), + ).unwrap(); + assert_eq!(sub_none.subscription_collections(), Some(None)); + + // Filtered subscription → Some(Some(vec)) + let sub_filtered = Record::from_json_value( + &nsid("sh.tangled.feed.subscription"), + json!({ + "$type": "sh.tangled.feed.subscription", + "createdAt": "2026-05-01T00:00:00Z", + "subject": { "$type": "sh.tangled.feed.subscription#repo", "did": "did:plc:repo" }, + "collections": ["sh.tangled.repo.issue"], + }), + ).unwrap(); + let cols = sub_filtered.subscription_collections(); + assert_eq!(cols, Some(Some(vec![SmolStr::new_static("sh.tangled.repo.issue")]))); + } + #[test] fn follow_uses_subject_did() { let edges = extract( diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index d8959766b..3c12f4d81 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -2509,6 +2509,7 @@ async fn get_coverage(State(state): State) -> Json { #[allow(dead_code)] struct ListRecipientsQuery { subject: String, + collection: Option, } #[derive(Serialize)] @@ -2537,9 +2538,13 @@ async fn list_recipients( subject_ref, ); let sources = state.edges.sources_for(&key); + + // NOTE: collection filtering intentionally disabled. + // All subscribers are returned regardless of their stored collection filter. + let filtered = sources; let mut seen = BTreeSet::new(); let mut dids = Vec::new(); - for src in &sources { + for src in &filtered { if let Some(did) = owner_did_from_aturi(src) { if seen.insert(did.to_string()) { dids.push(did.to_string());