diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index b9f894812..e5e056236 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -38,7 +38,7 @@ use bobbin_search::{ }; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::edges::REPO_SOURCE_EDGE_KIND; -use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; +use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static, owner_did_from_aturi}; use bobbin_types::knot_acl::{KnotOwnedSource, decode_knot_owned_source, knot_did_host}; use bobbin_types::record::RecordBody; use bobbin_types::search::SearchableRecord; @@ -534,6 +534,10 @@ pub fn router(state: AppState) -> Router { ) .route("/xrpc/sh.tangled.bobbin.getCoverage", get(get_coverage)) .route("/xrpc/sh.tangled.bobbin.awaitRecord", get(await_record)) + .route( + "/xrpc/org.tangled.temp.notification.listRecipients", + get(list_recipients), + ) .route( "/xrpc/blue.microcosm.identity.resolveMiniDoc", get(resolve_mini_doc), @@ -2933,6 +2937,51 @@ async fn await_record( Ok(Json(settled.into())) } +#[derive(Deserialize)] +#[allow(dead_code)] +struct ListRecipientsQuery { + subject: String, +} + +#[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("sh.tangled.feed.subscription"), + subject_ref, + ); + let sources = state.edges.sources_for(&key); + let mut seen = BTreeSet::new(); + let mut dids = Vec::new(); + for src in &sources { + 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 3af207e5b..43cc18655 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -3015,3 +3015,78 @@ async fn knot_owned_collaborator_lists_by_subject_did() { assert_eq!(items[0]["value"]["repo"], json!("did:plc:scallop")); assert_eq!(items[0]["value"]["subject"], json!("did:plc:olaren")); } + +#[tokio::test] +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"); + h.edges.add(Edge { + kind: nsid("sh.tangled.feed.subscription"), + subject: SubjectRef::Uri(entity.clone()), + source: sub.clone(), + sort_micros: 1, + }); + let resp = router(h.state.clone()) + .oneshot( + Request::builder() + .uri("/xrpc/org.tangled.temp.notification.listRecipients?subject=at://did:plc:repo/sh.tangled.repo.issue/abc") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let body: Value = serde_json::from_slice( + &to_bytes(resp.into_body(), 1 << 20).await.unwrap(), + ) + .unwrap(); + assert_eq!(body["dids"], json!(["did:plc:bob"])); +} + +#[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"); + h.edges.add(Edge { + kind: nsid("sh.tangled.feed.subscription"), + subject: SubjectRef::Did(did("did:plc:targetrepo")), + source: repo_did_sub.clone(), + sort_micros: 1, + }); + let resp = router(h.state.clone()) + .oneshot( + Request::builder() + .uri("/xrpc/org.tangled.temp.notification.listRecipients?subject=did:plc:targetrepo") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let body: Value = serde_json::from_slice( + &to_bytes(resp.into_body(), 1 << 20).await.unwrap(), + ) + .unwrap(); + assert_eq!(body["dids"], json!(["did:plc:watcher"])); +} + +#[tokio::test] +async fn list_recipients_empty_for_unknown_subject() { + let h = Harness::new().await; + let resp = router(h.state.clone()) + .oneshot( + Request::builder() + .uri("/xrpc/org.tangled.temp.notification.listRecipients?subject=at://did:plc:repo/sh.tangled.repo.issue/nonexistent") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let body: Value = serde_json::from_slice( + &to_bytes(resp.into_body(), 1 << 20).await.unwrap(), + ) + .unwrap(); + assert_eq!(body["dids"], json!([])); +}