From 1a86edf3cfd36702476cf85adddf8626b0fa4100 Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 24 Jul 2026 23:49:51 +0300 Subject: [PATCH] bobbin: add sh.tangled.repo.countForks Signed-off-by: dawn --- bobbin/crates/ingest/src/lib.rs | 128 ++++++++++++++++++++++++ bobbin/crates/types/src/edges.rs | 100 +++++++++++++++++- bobbin/crates/xrpc/src/lib.rs | 19 ++++ bobbin/crates/xrpc/tests/aggregation.rs | 70 +++++++++++++ lexicons/repo/countForks.json | 39 ++++++++ 5 files changed, 354 insertions(+), 2 deletions(-) create mode 100644 lexicons/repo/countForks.json diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 28517c65..7c75812a 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -3038,6 +3038,134 @@ mod tests { ); } + #[tokio::test] + async fn fork_record_edges_against_its_source() { + let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + let fork: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": "abalone", + "source": "did:plc:abalone" + } + } + })); + handle_frame( + fork, + &store, + &issue_states, + &pull_statuses, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + &sys_clock(), + now(), + ) + .await; + let key = bobbin_types::ids::EdgeKey::new( + Nsid::new_static(bobbin_types::edges::REPO_SOURCE_EDGE_KIND).unwrap(), + did_subj("did:plc:abalone"), + ); + assert_eq!(store.count(&key), 1); + } + + fn fork_frame(source: &str) -> HydrantFrame { + parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": "abalone", + "source": source + } + } + })) + } + + fn fork_edge_key(subject: SubjectRef) -> bobbin_types::ids::EdgeKey { + bobbin_types::ids::EdgeKey::new( + Nsid::new_static(bobbin_types::edges::REPO_SOURCE_EDGE_KIND).unwrap(), + subject, + ) + } + + #[tokio::test] + async fn fork_by_at_uri_edges_against_the_source_did() { + let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + resolver + .observe( + Did::new_owned("did:plc:nel").unwrap(), + Rkey::new_owned("core").unwrap(), + Some(Did::new_owned("did:plc:abalone").unwrap()), + ) + .await; + handle_frame( + fork_frame("at://did:plc:nel/sh.tangled.repo/core"), + &store, + &issue_states, + &pull_statuses, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + &sys_clock(), + now(), + ) + .await; + assert_eq!(store.count(&fork_edge_key(did_subj("did:plc:abalone"))), 1); + assert_eq!( + store.count(&fork_edge_key(uri_subj( + "at://did:plc:nel/sh.tangled.repo/core" + ))), + 0, + "the rkey form would never match a bare-DID query", + ); + } + + #[tokio::test] + async fn fork_of_a_repo_without_a_did_has_nothing_to_count_against() { + let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + handle_frame( + fork_frame("at://did:plc:nel/sh.tangled.repo/core"), + &store, + &issue_states, + &pull_statuses, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + &sys_clock(), + now(), + ) + .await; + assert_eq!( + store.count(&fork_edge_key(uri_subj( + "at://did:plc:nel/sh.tangled.repo/core" + ))), + 0, + ); + } + #[tokio::test] async fn repo_record_without_a_name_indexes_its_rkey() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); diff --git a/bobbin/crates/types/src/edges.rs b/bobbin/crates/types/src/edges.rs index ea9d926b..4d831827 100644 --- a/bobbin/crates/types/src/edges.rs +++ b/bobbin/crates/types/src/edges.rs @@ -3,7 +3,7 @@ use alloc::vec::Vec; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; -use jacquard_common::types::string::{AtStrError, AtUri, Datetime}; +use jacquard_common::types::string::{AtStrError, AtUri, Datetime, UriValue}; use jacquard_common::types::tid::Tid; use jacquard_common::{BosStr, DefaultStr}; @@ -228,7 +228,7 @@ impl Record { Self::Knot(_) => Ok(owner_self_edges("sh.tangled.knot", source)), Self::LabelDefinition(_) => Ok(owner_self_edges("sh.tangled.label.definition", source)), Self::PublicKey(_) => Ok(owner_self_edges("sh.tangled.publicKey", source)), - Self::Repo(_) => Ok(owner_self_edges("sh.tangled.repo", source)), + Self::Repo(r) => Ok(repo_edges(source, r)), Self::Spindle(_) => Ok(owner_self_edges("sh.tangled.spindle", source)), Self::TangledString(_) => Ok(owner_self_edges("sh.tangled.string", source)), Self::Vouch(r) => vouch_edges(source, r), @@ -237,6 +237,29 @@ impl Record { } } +pub const REPO_SOURCE_EDGE_KIND: &str = "sh.tangled.repo.source"; + +fn repo_edges(source: &AtUri, record: &RepoRecord) -> Vec { + let mut edges = owner_self_edges("sh.tangled.repo", source); + if let Some(subject) = fork_source_subject(record) { + edges.push(Edge { + kind: nsid_static(REPO_SOURCE_EDGE_KIND), + subject, + source: source.clone(), + sort_micros: 0, + }); + } + edges +} + +fn fork_source_subject(record: &RepoRecord) -> Option { + match record.source.as_ref()? { + UriValue::Did(did) => Some(SubjectRef::Did(did.clone())), + UriValue::At(uri) => uri_subject_for_record(uri), + _ => None, + } +} + fn repo_subject( uri: &Option>, did: &Option>, @@ -580,6 +603,79 @@ mod tests { parsed.primary_edges(&at(source)).expect("extract") } + fn repo_body(source: Option<&str>) -> serde_json::Value { + let mut body = json!({ + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": "abalone" + }); + if let Some(source) = source { + body["source"] = json!(source); + } + body + } + + #[test] + fn plain_repo_only_makes_its_owner_edge() { + let edges = extract( + "sh.tangled.repo", + "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", + repo_body(None), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].kind, nsid("sh.tangled.repo")); + assert_eq!(edges[0].subject, did_subj("did:plc:nel")); + } + + #[test] + fn fork_of_a_repo_did_keys_on_the_did() { + let edges = extract( + "sh.tangled.repo", + "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", + repo_body(Some("did:plc:abalone")), + ); + let fork: Vec<&Edge> = edges + .iter() + .filter(|e| e.kind.as_ref() == REPO_SOURCE_EDGE_KIND) + .collect(); + assert_eq!(fork.len(), 1); + assert_eq!(fork[0].subject, did_subj("did:plc:abalone")); + } + + #[test] + fn fork_of_a_repo_without_a_did_keys_on_the_record() { + let edges = extract( + "sh.tangled.repo", + "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", + repo_body(Some("at://did:plc:olaren/sh.tangled.repo/core")), + ); + let fork: Vec<&Edge> = edges + .iter() + .filter(|e| e.kind.as_ref() == REPO_SOURCE_EDGE_KIND) + .collect(); + assert_eq!(fork.len(), 1); + assert_eq!( + fork[0].subject, + uri_subj("at://did:plc:olaren/sh.tangled.repo/core") + ); + } + + #[test] + fn fork_of_something_off_network_makes_no_fork_edge() { + let edges = extract( + "sh.tangled.repo", + "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", + repo_body(Some("https://github.com/90-008/awawa")), + ); + assert!( + edges + .iter() + .all(|e| e.kind.as_ref() != REPO_SOURCE_EDGE_KIND), + "a clone url is not a repo we can count against", + ); + } + #[test] fn star_repo_variant_keys_on_did() { let edges = extract( diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 338dff3c..9c1b5c0e 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -33,6 +33,7 @@ use bobbin_search::{ SearchCursor, SearchError, SearchFilters, SearchHit, SearchOffset, SearchReader, }; 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::knot_acl::{KnotOwnedSource, decode_knot_owned_source, knot_did_host}; use bobbin_types::record::RecordBody; @@ -239,6 +240,7 @@ pub fn router(state: AppState) -> Router { ) .route("/xrpc/sh.tangled.repo.listRepos", get(list_repos)) .route("/xrpc/sh.tangled.repo.countRepos", get(count_repos)) + .route("/xrpc/sh.tangled.repo.countForks", get(count_forks)) .route("/xrpc/sh.tangled.knot.listKnots", get(list_knots)) .route("/xrpc/sh.tangled.knot.countKnots", get(count_knots)) .route("/xrpc/sh.tangled.spindle.listSpindles", get(list_spindles)) @@ -1018,6 +1020,9 @@ edge_kinds! { "sh.tangled.string" => TangledStringRecord, SubjectShape::BareDid; } +// the fork edge is not a collection so it is not in the table above +const FORK_SUBJECT_SHAPE: SubjectShape = SubjectShape::BareDid; + fn parse_subject(raw: &SubjectQuery, shape: SubjectShape) -> Result { let uri = match raw { SubjectQuery::Did(did) => { @@ -2105,6 +2110,20 @@ async fn count_repos( count_for::(&state, q).map(Json) } +// forks are repo records pointing back at a repo, so they get their own edge +// kind instead of a collection +async fn count_forks( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + let subject = parse_subject(&q.subject, FORK_SUBJECT_SHAPE)?; + let key = EdgeKey::new(nsid_static(REPO_SOURCE_EDGE_KIND), subject); + Ok(Json(CountResponse { + count: state.edges.count(&key), + distinct_authors: state.edges.count_distinct_authors(&key), + })) +} + async fn list_knots( State(state): State, XrpcQuery(q): XrpcQuery>, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index 8c61e077..9f43c320 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -350,6 +350,76 @@ async fn count_distinct_authors_dedupes_per_author() { assert_eq!(body["distinctAuthors"], json!(2)); } +#[tokio::test] +async fn count_forks_counts_repos_pointing_at_the_source() { + let h = Harness::new().await; + let subject = at("at://did:plc:abalone"); + h.add_edge( + &nsid("sh.tangled.repo.source"), + &subject, + &at("at://did:plc:nel/sh.tangled.repo/f1"), + ); + h.add_edge( + &nsid("sh.tangled.repo.source"), + &subject, + &at("at://did:plc:olaren/sh.tangled.repo/f2"), + ); + // the source owner's other repos are not forks of it + h.add_edge( + &nsid("sh.tangled.repo"), + &subject, + &at("at://did:plc:nel/sh.tangled.repo/r1"), + ); + + let app = router(h.state.clone()); + let resp = app + .oneshot(list_request( + "sh.tangled.repo.countForks", + subject.as_ref(), + &[], + )) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["count"], json!(2)); + assert_eq!(body["distinctAuthors"], json!(2)); +} + +#[tokio::test] +async fn count_forks_of_an_unforked_repo_is_zero() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let resp = app + .oneshot(list_request( + "sh.tangled.repo.countForks", + "did:plc:abalone", + &[], + )) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["count"], json!(0)); +} + +#[tokio::test] +async fn count_forks_rejects_an_at_uri_subject() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let resp = app + .oneshot(list_request( + "sh.tangled.repo.countForks", + "at://did:plc:nel/sh.tangled.repo/core", + &[], + )) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "InvalidRequest"); +} + #[tokio::test] async fn list_items_stable_across_coverage_promotion() { let h = Harness::new().await; diff --git a/lexicons/repo/countForks.json b/lexicons/repo/countForks.json new file mode 100644 index 00000000..15c9f5e7 --- /dev/null +++ b/lexicons/repo/countForks.json @@ -0,0 +1,39 @@ +{ + "lexicon": 1, + "id": "sh.tangled.repo.countForks", + "defs": { + "main": { + "type": "query", + "parameters": { + "type": "params", + "required": ["subject"], + "properties": { + "subject": { + "type": "string", + "format": "did", + "description": "Repo DID to count forks of." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["count", "distinctAuthors"], + "properties": { + "count": { + "type": "integer", + "minimum": 0, + "description": "Total number of matching records." + }, + "distinctAuthors": { + "type": "integer", + "minimum": 0, + "description": "Number of distinct authors among the matching records." + } + } + } + } + } + } +} -- 2.51.2