use std::sync::Arc; use axum::body::{Body, to_bytes}; use bobbin_edge_index::{CoverageWatch, EdgeStore, IssueStateKind, PullStatusKind, StateIndex}; use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; use bobbin_record_lru::{CacheCapacity, LruRecordStore}; use bobbin_resolver::RepoIdResolver; use bobbin_runtime::{RuntimeHasher, SystemClock}; use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; use bobbin_types::edges::Edge; use bobbin_types::ids::SubjectRef; use bobbin_xrpc::{AppState, router}; use http::{Request, StatusCode}; use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::AtUri; use serde_json::{Value, json}; use tower::ServiceExt; use url::Url; use wiremock::matchers::{method, path, query_param}; use wiremock::{Mock, MockServer, ResponseTemplate}; const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; const COUNT: &str = "sh.tangled.query.enrichResponse#count"; const DISTINCT_AUTHORS: &str = "sh.tangled.query.enrichResponse#distinctAuthors"; const VIEWER: &str = "sh.tangled.query.enrichResponse#viewer"; const MINIDOC: &str = "com.bad-example.identity.miniDoc"; fn at(s: &str) -> AtUri { AtUri::new_owned(s).unwrap() } fn did(s: &str) -> Did { Did::new_owned(s).unwrap() } fn nsid(s: &'static str) -> Nsid { Nsid::new_static(s).unwrap() } static EDGE_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1); fn next_sort_micros() -> u64 { EDGE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed) } struct Harness { server: MockServer, edges: Arc, state: AppState, } impl Harness { async fn new() -> Self { let server = MockServer::start().await; let edges = Arc::new(EdgeStore::new(RuntimeHasher::default())); let coverage = Arc::new(CoverageWatch::new()); let state = AppState::new( Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), edges.clone(), Arc::new(StateIndex::::new(RuntimeHasher::default())), Arc::new(StateIndex::::new(RuntimeHasher::default())), coverage.clone(), Arc::new( KnotProxy::new( KnotProxyConfig::default(), KnotHttpConfig::default(), Arc::new(SystemClock::new()), RuntimeHasher::default(), ) .unwrap(), ), Arc::new( SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), ) as Arc, Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), Arc::new(bobbin_xrpc::default_directory()), ); Self { server, edges, state, } } fn add_edge(&self, kind: &'static str, subject: SubjectRef, source: &AtUri) { self.edges.add(Edge { kind: nsid(kind), subject, source: source.clone(), sort_micros: next_sort_micros(), }); } async fn mount(&self, did: &Did, collection: &str, rkey: &str, value: Value) { let uri = format!("at://{}/{}/{}", did.as_ref(), collection, rkey); let body = json!({ "uri": uri, "cid": CID, "value": value }); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) .and(query_param("repo", did.as_ref())) .and(query_param("collection", collection)) .and(query_param("rkey", rkey)) .respond_with(ResponseTemplate::new(200).set_body_json(body)) .mount(&self.server) .await; } } fn enrich_request(body: Value) -> Request { Request::builder() .method("POST") .uri("/xrpc/sh.tangled.query.enrichResponse") .header("content-type", "application/json") .body(Body::from(serde_json::to_vec(&body).unwrap())) .unwrap() } async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { let status = resp.status(); let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); (status, parsed) } fn repo_body(name: &str, repo_did: &Did) -> Value { json!({ "$type": "sh.tangled.repo", "name": name, "knot": "oyster.cafe", "repoDid": repo_did.as_ref(), "createdAt": "2026-05-01T00:00:00Z" }) } fn follow_body(subject: &Did) -> Value { json!({ "$type": "sh.tangled.graph.follow", "createdAt": "2026-05-01T00:00:00Z", "subject": subject.as_ref() }) } /// one repo owned by `owner`, with `stars`/`issues` counts against its repo did async fn repo_fixture(h: &Harness, owner: &Did, repo_did: &Did) { let repo_uri = at(&format!("at://{}/sh.tangled.repo/reef", owner.as_ref())); h.add_edge("sh.tangled.repo", SubjectRef::Did(owner.clone()), &repo_uri); h.mount( owner, "sh.tangled.repo", "reef", repo_body("reef", repo_did), ) .await; for (i, stargazer) in ["did:plc:a", "did:plc:b", "did:plc:a"].iter().enumerate() { h.add_edge( "sh.tangled.feed.star", SubjectRef::Did(repo_did.clone()), &at(&format!("at://{stargazer}/sh.tangled.feed.star/s{i}")), ); } h.add_edge( "sh.tangled.repo.issue", SubjectRef::Did(repo_did.clone()), &at("at://did:plc:a/sh.tangled.repo.issue/i0"), ); } #[tokio::test] async fn zero_config_counts_stars_and_issues_for_repo_did() { let h = Harness::new().await; let owner = did("did:plc:nel"); let repo_did = did("did:plc:limpet"); repo_fixture(&h, &owner, &repo_did).await; let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [ { "source": "sh.tangled.feed.star:subject", "type": COUNT }, { "source": "sh.tangled.feed.star:subject", "type": DISTINCT_AUTHORS }, { "source": "sh.tangled.repo.issue:subject", "type": COUNT } ] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); assert_eq!(body["output"]["items"].as_array().unwrap().len(), 1); let stats = &body["data"]; assert_eq!( stats["did:plc:limpet"]["sh.tangled.feed.star:subject"][COUNT], json!(3) ); assert_eq!( stats["did:plc:limpet"]["sh.tangled.feed.star:subject"][DISTINCT_AUTHORS], json!(2) ); assert_eq!( stats["did:plc:limpet"]["sh.tangled.repo.issue:subject"][COUNT], json!(1) ); assert!(stats["at://did:plc:nel/sh.tangled.repo/reef"].is_null()); } #[tokio::test] async fn follow_counts_cover_both_directions() { let h = Harness::new().await; let owner = did("did:plc:nel"); // followers, edges pointing at owner for (i, fan) in ["did:plc:a", "did:plc:b"].iter().enumerate() { h.add_edge( "sh.tangled.graph.follow", SubjectRef::Did(owner.clone()), &at(&format!("at://{fan}/sh.tangled.graph.follow/f{i}")), ); h.mount( &did(fan), "sh.tangled.graph.follow", &format!("f{i}"), follow_body(&owner), ) .await; } // following, via the .by mirror edge since owner is the author here h.add_edge( "sh.tangled.graph.follow.by", SubjectRef::Did(owner.clone()), &at("at://did:plc:nel/sh.tangled.graph.follow/f0"), ); let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.graph.listFollows", "params": { "subject": owner.as_ref() }, "enrich": [{ "source": "sh.tangled.graph.follow:subject", "type": COUNT }, { "source": "sh.tangled.graph.follow:.repo", "type": COUNT }] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); let nel = &body["data"]["did:plc:nel"]; assert_eq!( nel["sh.tangled.graph.follow:subject"][COUNT], json!(2), "{body}" ); assert_eq!( nel["sh.tangled.graph.follow:.repo"][COUNT], json!(1), "{body}" ); } // at-uri authorities join the ref set, so record authors get stats keyed by // their bare did without appearing as a value anywhere in the response #[tokio::test] async fn authorities_of_record_uris_become_refs() { let h = Harness::new().await; let owner = did("did:plc:nel"); for (i, fan) in ["did:plc:a", "did:plc:b"].iter().enumerate() { h.add_edge( "sh.tangled.graph.follow", SubjectRef::Did(owner.clone()), &at(&format!("at://{fan}/sh.tangled.graph.follow/f{i}")), ); h.mount( &did(fan), "sh.tangled.graph.follow", &format!("f{i}"), follow_body(&owner), ) .await; } // each fan also follows one other person for (i, fan) in ["did:plc:a", "did:plc:b"].iter().enumerate() { h.add_edge( "sh.tangled.graph.follow.by", SubjectRef::Did(did(fan)), &at(&format!("at://{fan}/sh.tangled.graph.follow/g{i}")), ); } let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.graph.listFollows", "params": { "subject": owner.as_ref() }, "enrich": [ { "source": "sh.tangled.graph.follow:subject", "type": COUNT }, { "source": "sh.tangled.graph.follow:.repo", "type": COUNT } ] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); for fan in ["did:plc:a", "did:plc:b"] { let entry = &body["data"][fan]; assert_eq!( entry["sh.tangled.graph.follow:subject"][COUNT], json!(0), "{body}" ); assert_eq!( entry["sh.tangled.graph.follow:.repo"][COUNT], json!(1), "{body}" ); } } #[tokio::test] async fn targets_scope_each_payload_independently() { let h = Harness::new().await; let owner = did("did:plc:nel"); let repo_did = did("did:plc:limpet"); repo_fixture(&h, &owner, &repo_did).await; let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [ { "source": "sh.tangled.feed.star:subject", "type": COUNT, "targets": ["items[].value.repoDid"] }, { "source": "sh.tangled.repo.issue:subject", "type": COUNT, "targets": ["items[].uri"] } ] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); let stats = &body["data"]; assert_eq!( stats["did:plc:limpet"]["sh.tangled.feed.star:subject"][COUNT], json!(3) ); assert_eq!(stats.as_object().unwrap().len(), 2, "{body}"); assert_eq!( stats[owner.as_str()]["sh.tangled.repo.issue:subject"][COUNT], json!(0) ); assert!( stats[repo_did.as_str()]["sh.tangled.repo.issue:subject"].is_null(), "{body}" ); assert!( stats[owner.as_str()]["sh.tangled.feed.star:subject"].is_null(), "{body}" ); // a path matching nothing is empty stats, not an error, since selection is vector-matched let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": COUNT, "targets": ["items[].value.nope"] }] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); assert_eq!(body["data"], json!({})); } #[tokio::test] async fn inner_record_miss_passes_through_as_404() { let h = Harness::new().await; let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.repo.getRepo", "params": { "repo": "at://did:plc:nel/sh.tangled.repo/absent" }, "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": COUNT }] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::NOT_FOUND, "{body}"); assert_eq!(body["error"], json!("RecordNotFound")); } #[tokio::test] async fn rejects_bad_requests() { let h = Harness::new().await; // semantic rejections: the descriptor parses, the handler refuses it let handler_cases = [ json!({ "xrpc": "sh.tangled.nope.nope", "enrich": [] }), json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.nope:subject", "type": COUNT }] }), json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "sh.tangled.query.enrichResponse#bogus" }] }), json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": COUNT, "targets": ["items["] }] }), ]; for case in handler_cases { let app = router(h.state.clone()); let (status, body) = json_response(app.oneshot(enrich_request(case)).await.unwrap()).await; assert_eq!(status, StatusCode::BAD_REQUEST, "{body}"); assert_eq!(body["error"], json!("InvalidRequest"), "{body}"); } // structural rejections: serde refuses the descriptor before the handler // sees it, which is plain-text 422 rather than our 400 json body let serde_cases = [ json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star:subject" }] }), json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star", "type": COUNT }] }), json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star:.rkey", "type": COUNT }] }), json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star:subject.uri", "type": COUNT }] }), ]; for case in serde_cases { let app = router(h.state.clone()); let response = app.oneshot(enrich_request(case)).await.unwrap(); assert_eq!(response.status(), StatusCode::UNPROCESSABLE_ENTITY); } } #[tokio::test] async fn viewer_aggregation_uses_explicit_viewer_param() { let h = Harness::new().await; let owner = did("did:plc:abc"); let repo_did = did("did:plc:limpet"); repo_fixture(&h, &owner, &repo_did).await; // the viewer already starred this repo, for the checks below let subject = SubjectRef::Did(repo_did.clone()); h.state.edges.add(Edge { kind: nsid("sh.tangled.feed.star"), subject, source: at("at://did:plc:nel/sh.tangled.feed.star/r99"), sort_micros: 99, }); let app = router(h.state.clone()); // missing viewer param is a 400, viewer descriptors require it let no_viewer = json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": VIEWER }] }); let (status, _) = json_response( app.clone() .oneshot(enrich_request(no_viewer)) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::BAD_REQUEST); // viewer who starred it gets their own star uri back let starred_viewer = json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": VIEWER }], "viewer": "did:plc:nel" }); let (status, resp) = json_response( app.clone() .oneshot(enrich_request(starred_viewer)) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK); let stats = &resp["data"][repo_did.as_str()]["sh.tangled.feed.star:subject"]; assert_eq!( stats[VIEWER], json!("at://did:plc:nel/sh.tangled.feed.star/r99") ); // a viewer who never starred it gets an explicit null, not absent let other_viewer = json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": VIEWER }], "viewer": "did:plc:someoneelse" }); let (status, resp) = json_response(app.oneshot(enrich_request(other_viewer)).await.unwrap()).await; assert_eq!(status, StatusCode::OK); let stats = &resp["data"][repo_did.as_str()]["sh.tangled.feed.star:subject"]; assert_eq!(stats[VIEWER], Value::Null); } #[tokio::test] async fn minidoc_payloads_resolve_record_authors() { let h = Harness::new().await; let owner = did("did:plc:nel"); for (i, fan) in ["did:plc:a", "did:plc:b"].iter().enumerate() { h.add_edge( "sh.tangled.graph.follow", SubjectRef::Did(owner.clone()), &at(&format!("at://{fan}/sh.tangled.graph.follow/f{i}")), ); h.mount( &did(fan), "sh.tangled.graph.follow", &format!("f{i}"), follow_body(&owner), ) .await; } Mock::given(method("GET")) .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) .and(query_param("identifier", "did:plc:a")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "did": "did:plc:a", "handle": "a.example.com", "pds": "https://pds.example.com" }))) .expect(1) .mount(&h.server) .await; Mock::given(method("GET")) .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) .and(query_param("identifier", "did:plc:b")) .respond_with(ResponseTemplate::new(404)) .expect(1) .mount(&h.server) .await; let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.graph.listFollows", "params": { "subject": owner.as_ref() }, "enrich": [ { "source": "sh.tangled.graph.follow:.repo", "type": MINIDOC }, { "source": "sh.tangled.graph.follow:.repo", "type": MINIDOC }, { "source": "sh.tangled.feed.star:.repo", "type": MINIDOC } ] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); assert_eq!( body["data"]["did:plc:a"]["sh.tangled.graph.follow:.repo"][MINIDOC]["handle"], json!("a.example.com") ); assert_eq!( body["data"]["did:plc:a"]["sh.tangled.feed.star:.repo"][MINIDOC]["handle"], json!("a.example.com") ); // resolution failures are dropped, the client falls back for misses assert!(body["data"]["did:plc:b"].is_null(), "{body}"); // the profile owner authored nothing here, so it earns no minidoc assert!(body["data"]["did:plc:nel"].is_null(), "{body}"); } #[tokio::test] async fn minidoc_repo_sources_skip_the_author_index() { let h = Harness::new().await; let owner = did("did:plc:nel"); let repo_did = did("did:plc:limpet"); repo_fixture(&h, &owner, &repo_did).await; Mock::given(method("GET")) .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) .and(query_param("identifier", "did:plc:nel")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "did": "did:plc:nel", "handle": "nel.example.com" }))) .mount(&h.server) .await; let app = router(h.state.clone()); // sh.tangled.repo has no author mirror; stats would 400, minidocs must not let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [{ "source": "sh.tangled.repo:.repo", "type": MINIDOC }] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); assert_eq!( body["data"]["did:plc:nel"]["sh.tangled.repo:.repo"][MINIDOC]["handle"], json!("nel.example.com") ); let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(enrich_request(json!({ "xrpc": "sh.tangled.repo.listRepos", "params": { "subject": owner.as_ref() }, "enrich": [{ "source": "sh.tangled.repo:.repo", "type": COUNT }] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::BAD_REQUEST, "{body}"); }