diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index c1b58629..a5bb5d64 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -12,6 +12,7 @@ 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; @@ -634,6 +635,7 @@ pub struct EdgeStore { reverse: SccMap, RuntimeHasher>, issue_counts: StateViewIndex, pull_counts: StateViewIndex, + source_collections: SccMap, RuntimeHasher>, hasher: RuntimeHasher, writer: Mutex<()>, } @@ -651,6 +653,7 @@ 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(()), } @@ -705,6 +708,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); @@ -751,6 +787,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; }; @@ -3038,4 +3075,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/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/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 15c4df15..1999ab72 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -49,6 +49,7 @@ pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot}; const TANGLED_PREFIX: &str = "sh.tangled."; +const ORG_TANGLED_PREFIX: &str = "org.tangled."; const RECONNECT_INITIAL_DELAY: Duration = Duration::from_millis(500); const RECONNECT_MAX_DELAY: Duration = Duration::from_secs(30); const PING_INTERVAL: Duration = Duration::from_secs(20); @@ -1089,7 +1090,7 @@ async fn prepare_record( debug!("record-typed frame missing payload, skipping"); return PendingOp::Noop; }; - if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { + if !record.collection.as_ref().starts_with(TANGLED_PREFIX) && !record.collection.as_ref().starts_with(ORG_TANGLED_PREFIX) { return PendingOp::Noop; } // hydrant only announces an identity when it changes, so nothing tells us diff --git a/bobbin/crates/types/Cargo.toml b/bobbin/crates/types/Cargo.toml index c2b536fa..d8e0345b 100644 --- a/bobbin/crates/types/Cargo.toml +++ b/bobbin/crates/types/Cargo.toml @@ -23,8 +23,9 @@ jacquard-lexicon = { workspace = true, features = ["codegen"] } walkdir = { workspace = true } [features] -default = ["sh_tangled"] +default = ["sh_tangled", "org_tangled"] com_atproto = [] com_bad_example = [] sh_tangled = ["com_atproto"] +org_tangled = ["com_atproto"] streaming = ["jacquard-common/websocket"] diff --git a/bobbin/crates/types/src/edges.rs b/bobbin/crates/types/src/edges.rs index a097e358..1c37fe60 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, UriValue}; @@ -12,7 +13,7 @@ use crate::sh_tangled::actor::profile::Profile; use crate::sh_tangled::feed::comment::Comment as FeedCommentRecord; use crate::sh_tangled::feed::reaction::Reaction; use crate::sh_tangled::feed::star::Star; -use crate::sh_tangled::feed::subscription::Subscription; +use crate::org_tangled::feed::subscription::Subscription; use crate::sh_tangled::git::ref_update::RefUpdate; use crate::sh_tangled::graph::follow::Follow; use crate::sh_tangled::graph::vouch::Vouch; @@ -106,7 +107,7 @@ impl Record { "sh.tangled.feed.comment" => parse!(FeedComment), "sh.tangled.feed.reaction" => parse!(Reaction), "sh.tangled.feed.star" => parse!(Star), - "sh.tangled.feed.subscription" => parse!(Subscription), + "org.tangled.feed.subscription" => parse!(Subscription), "sh.tangled.git.refUpdate" => parse!(RefUpdate), "sh.tangled.graph.follow" => parse!(Follow), "sh.tangled.graph.vouch" => parse!(Vouch), @@ -137,7 +138,7 @@ impl Record { Self::FeedComment(_) => "sh.tangled.feed.comment", Self::Reaction(_) => "sh.tangled.feed.reaction", Self::Star(_) => "sh.tangled.feed.star", - Self::Subscription(_) => "sh.tangled.feed.subscription", + Self::Subscription(_) => "org.tangled.feed.subscription", Self::RefUpdate(_) => "sh.tangled.git.refUpdate", Self::Follow(_) => "sh.tangled.graph.follow", Self::Vouch(_) => "sh.tangled.graph.vouch", @@ -172,6 +173,15 @@ 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(); @@ -417,9 +427,9 @@ fn subscription_edges( source: &AtUri, record: &Subscription, ) -> Result, ExtractError> { - use crate::sh_tangled::feed::subscription::SubscriptionSubject; + use crate::org_tangled::feed::subscription::SubscriptionSubject; let subject = match &record.subject { - SubscriptionSubject::Record(v) => { + SubscriptionSubject::Uri(v) => { let Some(subject) = uri_subject_for_record(&v.uri) else { return Ok(Vec::new()); }; @@ -427,7 +437,7 @@ fn subscription_edges( } SubscriptionSubject::Repo(v) => SubjectRef::Did(v.did.clone()), }; - Ok(one_edge("sh.tangled.feed.subscription", subject, source)) + Ok(one_edge("org.tangled.feed.subscription", subject, source)) } fn reaction_edges( @@ -853,6 +863,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("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( @@ -940,6 +988,7 @@ mod tests { "title": "feature", "createdAt": "2026-05-01T00:00:00Z", "rounds": [], + "versions": [], "target": {"repo": "did:plc:abalone", "branch": "main"} }), ); diff --git a/bobbin/crates/xrpc/src/enrich.rs b/bobbin/crates/xrpc/src/enrich.rs index 75ce1817..f8a1cb41 100644 --- a/bobbin/crates/xrpc/src/enrich.rs +++ b/bobbin/crates/xrpc/src/enrich.rs @@ -416,6 +416,8 @@ fn shape_accepts(shape: SubjectShape, reference: &SubjectRef) -> bool { .collection() .is_some_and(|c| allowed.contains(&c.as_ref())), (SubjectShape::BareDidOrOneOfCollections(_), SubjectRef::Did(_)) => true, + (SubjectShape::BareDidOrAnyAtUri, SubjectRef::Did(_)) => true, + (SubjectShape::BareDidOrAnyAtUri, SubjectRef::Uri(_)) => true, (SubjectShape::AnyAtUri, SubjectRef::Uri(_)) => true, _ => false, } diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index e5e05623..691bef4d 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -83,6 +83,8 @@ use bobbin_types::sh_tangled::spindle::{Spindle, SpindleRecord}; use bobbin_types::sh_tangled::string::{ TangledString, TangledStringGetRecordOutput, TangledStringRecord, }; +use bobbin_types::org_tangled::feed::subscription::SubscriptionRecord; + use futures::Stream; use futures::stream::{self, StreamExt, TryStreamExt}; use jacquard_axum::service_auth::{self, ExtractOptionalServiceAuth, ServiceAuthConfig}; @@ -538,6 +540,10 @@ pub fn router(state: AppState) -> Router { "/xrpc/org.tangled.temp.notification.listRecipients", get(list_recipients), ) + .route( + "/xrpc/org.tangled.feed.subscription.getForActor", + get(get_subscription_for_actor), + ) .route( "/xrpc/blue.microcosm.identity.resolveMiniDoc", get(resolve_mini_doc), @@ -1096,13 +1102,13 @@ fn parse_uri(raw: &str) -> Result, XrpcError> { AtUri::::new_owned(raw).map_err(|e| XrpcError::InvalidParams(format!("uri: {e}"))) } -#[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum SubjectShape { BareDid, Collection(&'static str), BareDidOrOneOfCollections(&'static [&'static str]), OneOfCollections(&'static [&'static str]), AnyAtUri, + BareDidOrAnyAtUri, } pub trait HasSubject { @@ -1171,6 +1177,7 @@ edge_kinds! { "sh.tangled.spindle" => SpindleRecord, SubjectShape::BareDid; "sh.tangled.spindle.member" => SpindleMemberRecord, SubjectShape::BareDid, mirror SpindleMemberBy; "sh.tangled.string" => TangledStringRecord, SubjectShape::BareDid; + "org.tangled.feed.subscription" => SubscriptionRecord, SubjectShape::BareDidOrAnyAtUri; } // the fork edge is not a collection so it is not in the table above @@ -1180,7 +1187,7 @@ fn parse_subject(raw: &SubjectQuery, shape: SubjectShape) -> Result { return match shape { - SubjectShape::BareDid | SubjectShape::BareDidOrOneOfCollections(_) => { + SubjectShape::BareDid | SubjectShape::BareDidOrOneOfCollections(_) | SubjectShape::BareDidOrAnyAtUri => { Ok(SubjectRef::Did(did.clone())) } SubjectShape::Collection(expected) => Err(XrpcError::InvalidParams(format!( @@ -1235,7 +1242,7 @@ fn parse_subject(raw: &SubjectQuery, shape: SubjectShape) -> Result// with nsid in [{}], got collection {c}", allowed.join(", "), ))), - SubjectShape::AnyAtUri => Ok(SubjectRef::Uri(uri.clone())), + SubjectShape::AnyAtUri | SubjectShape::BareDidOrAnyAtUri => Ok(SubjectRef::Uri(uri.clone())), } } @@ -2152,6 +2159,13 @@ async fn get_star( ) -> Result, XrpcError> { get_for::(&state, q).map(Json) } +async fn get_subscription_for_actor( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + get_for::(&state, q).map(Json) +} + async fn list_follows( State(state): State, @@ -2965,7 +2979,7 @@ async fn list_recipients( }; let key = EdgeKey::new( - nsid_static("sh.tangled.feed.subscription"), + nsid_static("org.tangled.feed.subscription"), subject_ref, ); let sources = state.edges.sources_for(&key); diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index 43cc1865..e65c8681 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -3020,9 +3020,9 @@ async fn knot_owned_collaborator_lists_by_subject_did() { async fn list_recipients_entity_subject_returns_subscriber_dids() { let h = Harness::new().await; let entity = at("at://did:plc:repo/sh.tangled.repo.issue/abc"); - let sub = at("at://did:plc:bob/sh.tangled.feed.subscription/rkey1"); + let sub = at("at://did:plc:bob/org.tangled.feed.subscription/rkey1"); h.edges.add(Edge { - kind: nsid("sh.tangled.feed.subscription"), + kind: nsid("org.tangled.feed.subscription"), subject: SubjectRef::Uri(entity.clone()), source: sub.clone(), sort_micros: 1, @@ -3047,9 +3047,9 @@ async fn list_recipients_entity_subject_returns_subscriber_dids() { #[tokio::test] async fn list_recipients_repo_subject_returns_repo_subscribers() { let h = Harness::new().await; - let repo_did_sub = at("at://did:plc:watcher/sh.tangled.feed.subscription/rkey1"); + let repo_did_sub = at("at://did:plc:watcher/org.tangled.feed.subscription/rkey1"); h.edges.add(Edge { - kind: nsid("sh.tangled.feed.subscription"), + kind: nsid("org.tangled.feed.subscription"), subject: SubjectRef::Did(did("did:plc:targetrepo")), source: repo_did_sub.clone(), sort_micros: 1, diff --git a/lexicons/feed/subscription.json b/lexicons/feed/subscription.json index 21856811..6b7595da 100644 --- a/lexicons/feed/subscription.json +++ b/lexicons/feed/subscription.json @@ -14,9 +14,14 @@ "properties": { "subject": { "type": "union", - "refs": ["#record", "#repo"], + "refs": ["#uri", "#repo"], "closed": true }, + "collections": { + "type": "array", + "items": { "type": "string" }, + "description": "Optional collection NSIDs to filter which notifications are sent. Empty or absent means all collections." + }, "createdAt": { "type": "string", "format": "datetime" @@ -24,7 +29,7 @@ } } }, - "record": { + "uri": { "type": "object", "required": ["uri"], "properties": {