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"; 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())), ); 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": "distinctAuthors" }, { "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["stats"]; 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"]["distinctAuthors"], 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["stats"]["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}" ); } #[tokio::test] async fn sources_scope_which_refs_get_enriched() { 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" }], "sources": ["items[].value.repoDid"] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); let stats = &body["stats"]; assert_eq!( stats["did:plc:limpet"]["sh.tangled.feed.star:subject"]["count"], json!(3) ); assert_eq!(stats.as_object().unwrap().len(), 1, "{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" }], "sources": ["items[].value.nope"] }))) .await .unwrap(), ) .await; assert_eq!(status, StatusCode::OK, "{body}"); assert_eq!(body["stats"], 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; let 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.uri", "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", "type": "count" }], "sources": ["items["] }), // no defaults, a source without a colon is rejected json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star", "type": "count" }] }), ]; for case in 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}"); } // serde rejects a missing type before the handler sees it // that returns plain-text 422, not our 400 json body let app = router(h.state.clone()); let response = app .oneshot(enrich_request(json!({ "xrpc": "sh.tangled.repo.countRepos", "params": { "subject": "did:plc:nel" }, "enrich": [{ "source": "sh.tangled.feed.star:subject" }] }))) .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["stats"][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["stats"][repo_did.as_str()]["sh.tangled.feed.star:subject"]; assert_eq!(stats["viewer"], Value::Null); }