diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index cd372204c..d0667afcb 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -13,6 +13,7 @@ use jacquard_common::DefaultStr; use jacquard_common::deps::smol_str::SmolStr; 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; @@ -636,6 +637,7 @@ pub struct EdgeStore { source_collections: SccMap, RuntimeHasher>, issue_counts: StateViewIndex, pull_counts: StateViewIndex, + source_collections: SccMap, RuntimeHasher>, hasher: RuntimeHasher, writer: Mutex<()>, } @@ -654,6 +656,7 @@ impl EdgeStore { source_collections: 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(()), } @@ -741,7 +744,6 @@ impl EdgeStore { 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); @@ -3080,9 +3082,9 @@ mod tests { #[test] fn source_collections_set_and_clear() { let store = store(); - let source = at("at://did:plc:bob/sh.tangled.feed.subscription/r1"); + let source = at("at://did:plc:bob/org.tangled.feed.subscription/r1"); let edge = Edge { - kind: nsid("sh.tangled.feed.subscription"), + kind: nsid("org.tangled.feed.subscription"), subject: did_subj("did:plc:repo"), source: source.clone(), sort_micros: 1, @@ -3120,9 +3122,9 @@ mod tests { #[test] fn remove_source_cleans_up_collection_filter() { let store = store(); - let source = at("at://did:plc:bob/sh.tangled.feed.subscription/r1"); + let source = at("at://did:plc:bob/org.tangled.feed.subscription/r1"); let edge = Edge { - kind: nsid("sh.tangled.feed.subscription"), + kind: nsid("org.tangled.feed.subscription"), subject: did_subj("did:plc:repo"), source: source.clone(), sort_micros: 1, diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 15c4df15d..1999ab72c 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 c2b536fa1..d8e0345bb 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 2c0268cb4..1c37fe602 100644 --- a/bobbin/crates/types/src/edges.rs +++ b/bobbin/crates/types/src/edges.rs @@ -13,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; @@ -107,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), @@ -138,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", @@ -427,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()); }; @@ -437,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( @@ -873,38 +873,32 @@ mod tests { "createdAt": "2026-05-01T00:00:00Z", "subject": { "$type": "sh.tangled.feed.star#repo", "did": "did:plc:abalone" }, }), - ) - .unwrap(); + ).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"), + &nsid("org.tangled.feed.subscription"), json!({ - "$type": "sh.tangled.feed.subscription", + "$type": "org.tangled.feed.subscription", "createdAt": "2026-05-01T00:00:00Z", - "subject": { "$type": "sh.tangled.feed.subscription#repo", "did": "did:plc:repo" }, + "subject": { "$type": "org.tangled.feed.subscription#repo", "did": "did:plc:repo" }, }), - ) - .unwrap(); + ).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"), + &nsid("org.tangled.feed.subscription"), json!({ - "$type": "sh.tangled.feed.subscription", + "$type": "org.tangled.feed.subscription", "createdAt": "2026-05-01T00:00:00Z", - "subject": { "$type": "sh.tangled.feed.subscription#repo", "did": "did:plc:repo" }, + "subject": { "$type": "org.tangled.feed.subscription#repo", "did": "did:plc:repo" }, "collections": ["sh.tangled.repo.issue"], }), - ) - .unwrap(); + ).unwrap(); let cols = sub_filtered.subscription_collections(); - assert_eq!( - cols, - Some(Some(vec![SmolStr::new_static("sh.tangled.repo.issue")])) - ); + assert_eq!(cols, Some(Some(vec![SmolStr::new_static("sh.tangled.repo.issue")]))); } #[test] @@ -994,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 75ce18170..f8a1cb410 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 7606801dd..a17ca30c5 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, @@ -2983,6 +2997,52 @@ async fn await_record( }; Ok(Json(settled.into())) } +#[derive(Deserialize)] +#[allow(dead_code)] +struct ListRecipientsQuery { + subject: String, + collection: Option, +} + +#[derive(Serialize)] +struct ListRecipientsResponse { + dids: Vec, +} + +async fn list_recipients( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + use std::collections::BTreeSet; + + let subject_ref = if let Ok(uri) = AtUri::::new_owned(&q.subject) { + SubjectRef::Uri(uri) + } else if let Ok(did) = Did::::new_owned(&q.subject) { + SubjectRef::Did(did) + } else { + return Err(XrpcError::InvalidParams( + "subject must be an at-uri or did".into(), + )); + }; + + let key = EdgeKey::new(nsid_static("org.tangled.feed.subscription"), 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 &filtered { + if let Some(did) = owner_did_from_aturi(src) { + if seen.insert(did.to_string()) { + dids.push(did.to_string()); + } + } + } + + Ok(Json(ListRecipientsResponse { dids })) +} async fn search_query( State(state): State, XrpcQuery(q): XrpcQuery, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index 0a2b20603..b50725c26 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, @@ -3045,9 +3045,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 21856811b..6b7595da9 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": {