From 69b74c77fa1fe7fcf5de9b4ec5c475e2c8e937bc Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Fri, 21 Aug 2026 02:12:28 +0900 Subject: [PATCH] Revert "bobbin: filter notification recipients by subscription collection" This reverts commit 37972939665e6e0d5e6fa8bef4ef6112db96b928, and partially 2c01b8d94a6c2ac5eb3a00003a78af2eb1bd0a6b. `2c01b8d9` accidentally reintroduced back the code reverted from `c6b02293e2c1357545fe9e8f9d473c0c18c1c758`. `collections` is not enough, say you want to subscribe to state change of specific ticket. State is not standalone record anymore. List of `org.tangled.event.*` NSIDs can solve this by defininig richer event definitions. Signed-off-by: Seongmin Lee --- bobbin/crates/edge-index/src/lib.rs | 95 ----------------------------- bobbin/crates/types/src/edges.rs | 48 --------------- 2 files changed, 143 deletions(-) diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index a5bb5d64f..c1b586294 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -12,7 +12,6 @@ use itertools::Itertools; use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; -use jacquard_common::deps::smol_str::SmolStr; use jacquard_common::types::string::AtUri; use lasso::{Key, Spur, ThreadedRodeo}; use scc::HashMap as SccMap; @@ -635,7 +634,6 @@ pub struct EdgeStore { reverse: SccMap, RuntimeHasher>, issue_counts: StateViewIndex, pull_counts: StateViewIndex, - source_collections: SccMap, RuntimeHasher>, hasher: RuntimeHasher, writer: Mutex<()>, } @@ -653,7 +651,6 @@ impl EdgeStore { reverse: SccMap::with_hasher(hasher.clone()), issue_counts: StateViewIndex::new(hasher.clone()), pull_counts: StateViewIndex::new(hasher.clone()), - source_collections: SccMap::with_hasher(hasher.clone()), hasher, writer: Mutex::new(()), } @@ -708,39 +705,6 @@ 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); @@ -787,7 +751,6 @@ 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; }; @@ -3075,62 +3038,4 @@ 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/org.tangled.feed.subscription/r1"); - let edge = Edge { - kind: nsid("org.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/org.tangled.feed.subscription/r1"); - let edge = Edge { - kind: nsid("org.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/types/src/edges.rs b/bobbin/crates/types/src/edges.rs index 1c37fe602..fd9eca694 100644 --- a/bobbin/crates/types/src/edges.rs +++ b/bobbin/crates/types/src/edges.rs @@ -1,7 +1,6 @@ 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, UriValue}; @@ -173,15 +172,6 @@ impl Record { Ok(out) } - /// 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(); @@ -863,44 +853,6 @@ 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("org.tangled.feed.subscription"), - json!({ - "$type": "org.tangled.feed.subscription", - "createdAt": "2026-05-01T00:00:00Z", - "subject": { "$type": "org.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("org.tangled.feed.subscription"), - json!({ - "$type": "org.tangled.feed.subscription", - "createdAt": "2026-05-01T00:00:00Z", - "subject": { "$type": "org.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( -- 2.51.2