From b2ff1bbb373654efddec70b88e15ea38fbf06d78 Mon Sep 17 00:00:00 2001 From: Anirudh Oppiliappan Date: Thu, 13 Aug 2026 11:49:14 +0300 Subject: [PATCH] bobbin: filter notification recipients by subscription collection Subscriptions can name the collections they care about. The edge index now stores that filter alongside the source, ingest keeps it in sync as records land, and listRecipients drops subscribers whose filter excludes the requested collection. A subscriber with no filter still matches everything. Signed-off-by: Anirudh Oppiliappan --- bobbin/crates/edge-index/src/lib.rs | 95 +++++++++++++++++++++++++++++ bobbin/crates/ingest/src/lib.rs | 6 ++ bobbin/crates/types/src/edges.rs | 48 +++++++++++++++ bobbin/crates/xrpc/src/lib.rs | 7 ++- 4 files changed, 155 insertions(+), 1 deletion(-) 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()); -- 2.51.2