From c9581d7b782ae1fcb549eec9a451c664cb04550e Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Thu, 17 Sep 2026 05:14:13 +0900 Subject: [PATCH] bobbin,web: remove knot member index from bobbin Signed-off-by: Seongmin Lee --- bobbin/crates/edge-index/src/invites.rs | 53 +-- bobbin/crates/ingest/src/lib.rs | 74 ++--- bobbin/crates/knot-ingest/src/client.rs | 73 +---- bobbin/crates/knot-ingest/src/firehose.rs | 47 +-- bobbin/crates/knot-ingest/src/legacy.rs | 64 ++-- bobbin/crates/knot-ingest/src/orchestrator.rs | 192 +++++------ bobbin/crates/knot-ingest/src/registry.rs | 66 ---- bobbin/crates/knot-ingest/src/roster.rs | 304 +++++++----------- bobbin/crates/knot-ingest/src/stream.rs | 58 ++-- bobbin/crates/resolver/src/legacy_upgrade.rs | 33 +- bobbin/crates/resolver/src/normalize.rs | 2 - bobbin/crates/types/src/edges.rs | 35 -- bobbin/crates/types/src/knot_acl.rs | 110 +------ bobbin/crates/types/src/legacy.rs | 17 - bobbin/crates/types/src/search.rs | 1 - bobbin/crates/xrpc/src/lib.rs | 107 +----- bobbin/crates/xrpc/tests/aggregation.rs | 124 +------ bobbin/crates/xrpc/tests/extended.rs | 50 --- web/src/lib/api/count.ts | 2 - web/src/lib/api/repoCreationTargets.test.ts | 59 ++-- web/src/lib/api/repoCreationTargets.ts | 59 +--- 21 files changed, 386 insertions(+), 1144 deletions(-) diff --git a/bobbin/crates/edge-index/src/invites.rs b/bobbin/crates/edge-index/src/invites.rs index 8a5ac6c7e..f2fe38b54 100644 --- a/bobbin/crates/edge-index/src/invites.rs +++ b/bobbin/crates/edge-index/src/invites.rs @@ -1,3 +1,5 @@ +//! repository collaborator invite + use std::collections::{HashMap, HashSet}; use bobbin_runtime::UnixMicros; @@ -7,27 +9,9 @@ use jacquard_common::types::did::Did; use parking_lot::RwLock; #[derive(Clone, Debug, Eq, Hash, PartialEq)] -pub enum InviteScope { - Knot(Did), - Repo { - knot: Did, - repo: Did, - }, -} - -impl InviteScope { - pub fn knot(&self) -> &Did { - match self { - Self::Knot(knot) | Self::Repo { knot, .. } => knot, - } - } - - pub fn subject(&self) -> &Did { - match self { - Self::Knot(knot) => knot, - Self::Repo { repo, .. } => repo, - } - } +pub struct InviteScope { + pub knot: Did, + pub repo: Did, } #[derive(Clone, Debug, Eq, PartialEq)] @@ -107,10 +91,7 @@ impl InviteIndex { .map(|(scope, offer)| (scope.clone(), offer.clone())) .sorted_by(|(left_scope, left), (right_scope, right)| { right.offered_at.cmp(&left.offered_at).then_with(|| { - left_scope - .subject() - .as_ref() - .cmp(right_scope.subject().as_ref()) + left_scope.repo.as_ref().cmp(right_scope.repo.as_ref()) }) }) .collect() @@ -127,12 +108,8 @@ mod tests { Did::new_owned(s).unwrap() } - fn knot_scope() -> InviteScope { - InviteScope::Knot(did("did:web:knot.oyster.cafe")) - } - fn repo_scope(repo: &str) -> InviteScope { - InviteScope::Repo { + InviteScope { knot: did("did:web:knot.oyster.cafe"), repo: did(repo), } @@ -152,7 +129,7 @@ mod tests { #[test] fn an_invitee_reads_their_outstanding_offers_newest_first() { let index = InviteIndex::new(); - offer(&index, knot_scope(), "did:plc:akshay", 10); + offer(&index, repo_scope("did:plc:whelk"), "did:plc:akshay", 10); offer(&index, repo_scope("did:plc:scallop"), "did:plc:akshay", 30); offer(&index, repo_scope("did:plc:limpet"), "did:plc:olaren", 20); @@ -167,25 +144,25 @@ mod tests { assert_eq!(outstanding[0].0, repo_scope("did:plc:scallop")); assert_eq!(outstanding[0].1.invited_by, did("did:plc:akshay")); assert_eq!( - outstanding[0].0.knot(), - &did("did:web:knot.oyster.cafe"), - "repo-scoped offer still reports knot, since pending is keyed on knot" + outstanding[0].0.knot, + did("did:web:knot.oyster.cafe"), + "an offer still reports its knot, since pending is keyed on knot" ); - assert_eq!(outstanding[0].0.subject(), &did("did:plc:scallop")); + assert_eq!(outstanding[0].0.repo, did("did:plc:scallop")); } #[test] fn re_offer_replaces_original_and_withdrawal_clears_it() { let index = InviteIndex::new(); - offer(&index, knot_scope(), "did:plc:akshay", 10); - offer(&index, knot_scope(), "did:plc:olaren", 40); + offer(&index, repo_scope("did:plc:scallop"), "did:plc:akshay", 10); + offer(&index, repo_scope("did:plc:scallop"), "did:plc:olaren", 40); let outstanding = index.outstanding(&did(INVITEE)); assert_eq!(outstanding.len(), 1); assert_eq!(outstanding[0].1.invited_by, did("did:plc:olaren")); assert_eq!(outstanding[0].1.offered_at, UnixMicros::new(40)); - index.withdraw(&did(INVITEE), &knot_scope()); + index.withdraw(&did(INVITEE), &repo_scope("did:plc:scallop")); assert!(index.outstanding(&did(INVITEE)).is_empty()); } diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 84b018975..189a67d7e 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -1219,9 +1219,6 @@ async fn prepare_record( }; match acl_disposition(&parsed, ctx.knot_gate, ctx.knot_registry) { AclDisposition::NativeSkip => { - if let Some(registry) = ctx.knot_registry { - registry.forget_legacy_member(&source); - } return PendingOp::NotIndexed { source, nsid, @@ -1229,12 +1226,6 @@ async fn prepare_record( detail: KNOT_AUTHORITATIVE, }; } - AclDisposition::LegacyMember { host } => { - if let Some(registry) = ctx.knot_registry { - registry.observe_host(&host); - registry.note_legacy_member(source.clone(), &host); - } - } AclDisposition::Other => {} } let edges = match parsed.extract_edges(&source) { @@ -1255,14 +1246,7 @@ async fn prepare_record( supersedes: None, })) } - RecordAction::Delete => { - if nsid.as_ref() == "sh.tangled.knot.member" - && let Some(registry) = ctx.knot_registry - { - registry.forget_legacy_member(&source); - } - PendingOp::Delete { source, nsid } - } + RecordAction::Delete => PendingOp::Delete { source, nsid }, RecordAction::Other => { debug!(collection = %nsid, "ignoring unknown record action"); PendingOp::Noop @@ -1286,7 +1270,6 @@ fn rejected( enum AclDisposition { Other, NativeSkip, - LegacyMember { host: KnotHostKey }, } fn acl_disposition( @@ -1298,14 +1281,6 @@ fn acl_disposition( return AclDisposition::Other; }; match parsed { - Record::KnotMember(member) => { - let host = KnotHostKey::new(member.domain.as_ref()); - if gate.is_native(&host) { - AclDisposition::NativeSkip - } else { - AclDisposition::LegacyMember { host } - } - } Record::Collaborator(collaborator) => { let native = registry .and_then(|registry| registry.host_of_repo(&collaborator.repo)) @@ -1989,7 +1964,7 @@ mod tests { } #[tokio::test] - async fn native_knot_member_skipped_legacy_indexed() { + async fn native_knot_collaborator_skipped_legacy_indexed() { use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry}; use wiremock::matchers::{method, path}; use wiremock::{Mock, MockServer, ResponseTemplate}; @@ -2033,7 +2008,7 @@ mod tests { settlements: None, }; - let member_frame = |id: u64, rkey: &str, domain: &str| { + let collaborator_frame = |id: u64, rkey: &str, repo: &str| { parse_frame(json!({ "id": id, "type": "record", @@ -2041,42 +2016,45 @@ mod tests { "live": false, "did": "did:plc:akshay", "rev": fresh_tid().as_str(), - "collection": "sh.tangled.knot.member", + "collection": "sh.tangled.repo.collaborator", "rkey": rkey, "action": "create", "record": { - "$type": "sh.tangled.knot.member", + "$type": "sh.tangled.repo.collaborator", "subject": "did:plc:boltless", - "domain": domain, + "repo": repo, "createdAt": "2026-06-01T00:00:00Z" } } })) }; - let native = - prepare_frame(member_frame(1, "aaaaaaaaaaaaz", &native_host), &ctx, now()).await; + registry.observe_repo(&KnotHostKey::new(&native_host), Did::new_owned("did:plc:scallop").unwrap()); + registry.observe_repo( + &KnotHostKey::new("legacy.knot"), + Did::new_owned("did:plc:whelk").unwrap(), + ); + + let native = prepare_frame( + collaborator_frame(1, "aaaaaaaaaaaaz", "did:plc:scallop"), + &ctx, + now(), + ) + .await; assert!( matches!(native.op, PendingOp::NotIndexed { .. }), - "member record for a native knot must be dropped" + "collaborator record for a repo on a native knot must be dropped" ); - let legacy = - prepare_frame(member_frame(2, "bbbbbbbbbbbbz", "legacy.knot"), &ctx, now()).await; + let legacy = prepare_frame( + collaborator_frame(2, "bbbbbbbbbbbbz", "did:plc:whelk"), + &ctx, + now(), + ) + .await; assert!( matches!(legacy.op, PendingOp::Upsert(_)), - "member record for a legacy knot must be ingested" - ); - assert!( - registry.hosts().contains(&KnotHostKey::new("legacy.knot")), - "a member record seeds host discovery even before any repo is seen" - ); - assert_eq!( - registry - .drain_legacy_members(&KnotHostKey::new("legacy.knot")) - .len(), - 1, - "legacy member edge is indexed for later purge once the knot upgrades" + "collaborator record for a repo on a legacy knot must be ingested" ); } diff --git a/bobbin/crates/knot-ingest/src/client.rs b/bobbin/crates/knot-ingest/src/client.rs index eed6384ea..91295bfde 100644 --- a/bobbin/crates/knot-ingest/src/client.rs +++ b/bobbin/crates/knot-ingest/src/client.rs @@ -28,9 +28,7 @@ const LIST_PAGE_LIMIT: i64 = 1000; const MAX_LIST_PAGES: usize = 256; const VERSION_NSID: &str = "sh.tangled.knot.version"; -const LIST_MEMBERS_NSID: &str = "sh.tangled.knot.listMembers"; const LIST_COLLABORATORS_NSID: &str = "sh.tangled.repo.listCollaborators"; -const LIST_MEMBER_INVITES_NSID: &str = "sh.tangled.knot.listMemberInvites"; const LIST_COLLABORATOR_INVITES_NSID: &str = "sh.tangled.repo.listCollaboratorInvites"; #[derive(Clone)] @@ -179,14 +177,6 @@ impl KnotClient { .collect()) } - pub async fn list_members( - &self, - host: &KnotHost, - knot: &Did, - ) -> Result { - self.list(host, LIST_MEMBERS_NSID, knot).await - } - pub async fn list_collaborators( &self, host: &KnotHost, @@ -195,14 +185,6 @@ impl KnotClient { self.list(host, LIST_COLLABORATORS_NSID, repo).await } - pub async fn list_member_invites( - &self, - host: &KnotHost, - knot: &Did, - ) -> Result { - self.list(host, LIST_MEMBER_INVITES_NSID, knot).await - } - pub async fn list_collaborator_invites( &self, host: &KnotHost, @@ -404,12 +386,12 @@ mod tests { } #[tokio::test] - async fn list_members_drains_single_page() { + async fn list_collaborators_drains_single_page() { let server = MockServer::start().await; let host = endpoint(&server); Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) - .and(query_param("subject", "did:web:knot.oyster.cafe")) + .and(path("/xrpc/sh.tangled.repo.listCollaborators")) + .and(query_param("subject", "did:plc:scallop")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "items": [ {"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}, @@ -419,8 +401,8 @@ mod tests { .mount(&server) .await; - let knot = did("did:web:knot.oyster.cafe"); - let listing = client().list_members(&host, &knot).await.unwrap(); + let repo = did("did:plc:scallop"); + let listing = client().list_collaborators(&host, &repo).await.unwrap(); assert_eq!(listing.completeness, Completeness::Complete); assert_eq!( listing.entries, @@ -438,11 +420,11 @@ mod tests { } #[tokio::test] - async fn list_members_drains_multiple_pages() { + async fn list_collaborators_drains_multiple_pages() { let server = MockServer::start().await; let host = endpoint(&server); Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) + .and(path("/xrpc/sh.tangled.repo.listCollaborators")) .and(query_param("cursor", "p2")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "items": [{"subject": "did:plc:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"}] @@ -450,7 +432,7 @@ mod tests { .mount(&server) .await; Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) + .and(path("/xrpc/sh.tangled.repo.listCollaborators")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "items": [{"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}], "cursor": "p2" @@ -458,8 +440,8 @@ mod tests { .mount(&server) .await; - let knot = did("did:web:knot.oyster.cafe"); - let listing = client().list_members(&host, &knot).await.unwrap(); + let repo = did("did:plc:scallop"); + let listing = client().list_collaborators(&host, &repo).await.unwrap(); assert_eq!(listing.completeness, Completeness::Complete); let subjects: Vec<_> = listing.entries.into_iter().map(|m| m.subject).collect(); assert_eq!( @@ -491,17 +473,9 @@ mod tests { } #[tokio::test] - async fn an_invite_listing_decodes_both_facets_at_their_own_subject() { + async fn an_invite_listing_decodes_at_the_repo_subject() { let server = MockServer::start().await; let host = endpoint(&server); - Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMemberInvites")) - .and(query_param("subject", "did:web:knot.oyster.cafe")) - .respond_with(ResponseTemplate::new(200).set_body_json(json!({ - "items": [{"subject": "did:plc:limpet", "addedBy": "did:plc:akshay", "createdAt": "2026-06-04T00:00:00Z"}] - }))) - .mount(&server) - .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.listCollaboratorInvites")) .and(query_param("subject", "did:plc:scallop")) @@ -511,24 +485,11 @@ mod tests { .mount(&server) .await; - let members = client() - .list_member_invites(&host, &did("did:web:knot.oyster.cafe")) - .await - .unwrap(); - assert_eq!(members.completeness, Completeness::Complete); - assert_eq!( - members.entries, - vec![InviteEntry { - subject: did("did:plc:limpet"), - added_by: did("did:plc:akshay"), - created_at: stamp("2026-06-04T00:00:00Z"), - }] - ); - let collaborators = client() .list_collaborator_invites(&host, &did("did:plc:scallop")) .await .unwrap(); + assert_eq!(collaborators.completeness, Completeness::Complete); assert_eq!( collaborators.entries, vec![InviteEntry { @@ -544,7 +505,7 @@ mod tests { let server = MockServer::start().await; let host = endpoint(&server); Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMemberInvites")) + .and(path("/xrpc/sh.tangled.repo.listCollaboratorInvites")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "items": [{"subject": "did:plc:limpet", "createdAt": "2026-06-04T00:00:00Z"}] }))) @@ -552,7 +513,7 @@ mod tests { .await; let rejected = client() - .list_member_invites(&host, &did("did:web:knot.oyster.cafe")) + .list_collaborator_invites(&host, &did("did:plc:scallop")) .await .is_err(); assert!( @@ -566,7 +527,7 @@ mod tests { let server = MockServer::start().await; let host = endpoint(&server); Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) + .and(path("/xrpc/sh.tangled.repo.listCollaborators")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "items": [{"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}], "cursor": "more" @@ -574,8 +535,8 @@ mod tests { .mount(&server) .await; - let knot = did("did:web:knot.oyster.cafe"); - let listing = client().list_members(&host, &knot).await.unwrap(); + let repo = did("did:plc:scallop"); + let listing = client().list_collaborators(&host, &repo).await.unwrap(); assert_eq!( listing.completeness, Completeness::Truncated, diff --git a/bobbin/crates/knot-ingest/src/firehose.rs b/bobbin/crates/knot-ingest/src/firehose.rs index 89afc81dc..2706c1ab3 100644 --- a/bobbin/crates/knot-ingest/src/firehose.rs +++ b/bobbin/crates/knot-ingest/src/firehose.rs @@ -74,38 +74,15 @@ impl OpAction { } } -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub enum InviteEvent { - KnotMember, - RepoCollaborator, -} - -impl InviteEvent { - const KINDS: [Self; 2] = [Self::KnotMember, Self::RepoCollaborator]; - - pub const fn collection(self) -> &'static str { - match self { - Self::KnotMember => "sh.tangled.knot.memberInvite", - Self::RepoCollaborator => "sh.tangled.repo.collaboratorInvite", - } - } +pub const COLLABORATOR_INVITE_COLLECTION: &str = "sh.tangled.repo.collaboratorInvite"; - pub fn of(collection: &Nsid) -> Option { - Self::KINDS - .into_iter() - .find(|kind| kind.collection() == collection.as_str()) - } -} - -#[derive(Clone, Debug, Eq, PartialEq)] -pub enum InviteTarget { - Knot, - Repo(Did), +pub fn is_collaborator_invite(collection: &Nsid) -> bool { + collection.as_str() == COLLABORATOR_INVITE_COLLECTION } #[derive(Clone, Debug, PartialEq)] pub struct InviteOp { - pub target: InviteTarget, + pub repo: Did, pub invitee: Did, pub offer: Option, } @@ -118,11 +95,7 @@ struct InviteWire { editor: Did, } -pub fn invite_op_of(op: &RecordOp, kind: InviteEvent, repo: &Did) -> Option { - let target = match kind { - InviteEvent::KnotMember => InviteTarget::Knot, - InviteEvent::RepoCollaborator => InviteTarget::Repo(repo.clone()), - }; +pub fn invite_op_of(op: &RecordOp, repo: &Did) -> Option { let invitee = Did::::new_owned(op.rkey.as_str()).ok()?; let offer = if op.action.is_delete() { None @@ -134,7 +107,7 @@ pub fn invite_op_of(op: &RecordOp, kind: InviteEvent, repo: &Did) -> }) }; Some(InviteOp { - target, + repo: repo.clone(), invitee, offer, }) @@ -415,7 +388,7 @@ mod tests { ("collection that is not an nsid", "not an nsid", "3mug"), ( "record key with a slash in it", - InviteEvent::KnotMember.collection(), + COLLABORATOR_INVITE_COLLECTION, "did:plc:limpet/extra", ), ]; @@ -487,9 +460,9 @@ mod tests { #[test] fn malformed_invite_never_decodes_to_offer() { - let collection = InviteEvent::KnotMember.collection(); + let collection = COLLABORATOR_INVITE_COLLECTION; let full = invite_record(collection, "2026-06-01T00:00:00Z", "did:plc:akshay"); - let knot = Did::::new_owned("did:web:knot.oyster.cafe").unwrap(); + let repo = Did::::new_owned("did:plc:scallop").unwrap(); let cases = [ ("record key that isn't a did", "3lkm2xqbolt2s", full.clone()), ( @@ -507,7 +480,7 @@ mod tests { record: Some(Bytes::from(body)), }; assert!( - invite_op_of(&op, InviteEvent::KnotMember, &knot).is_none(), + invite_op_of(&op, &repo).is_none(), "invite with {case} still decoded" ); }); diff --git a/bobbin/crates/knot-ingest/src/legacy.rs b/bobbin/crates/knot-ingest/src/legacy.rs index 3afcd4efb..95b61453c 100644 --- a/bobbin/crates/knot-ingest/src/legacy.rs +++ b/bobbin/crates/knot-ingest/src/legacy.rs @@ -5,28 +5,7 @@ use serde::Deserialize; use crate::stream::{Cursor, FeedHandler, Outcome}; -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -enum RosterUpdate { - KnotMember, - RepoCollaborator, -} - -impl RosterUpdate { - const KINDS: [Self; 2] = [Self::KnotMember, Self::RepoCollaborator]; - - const fn nsid(self) -> &'static str { - match self { - Self::KnotMember => "sh.tangled.knot.memberUpdate", - Self::RepoCollaborator => "sh.tangled.repo.collaboratorUpdate", - } - } - - fn of(nsid: &Nsid) -> Option { - Self::KINDS - .into_iter() - .find(|kind| kind.nsid() == nsid.as_str()) - } -} +const COLLABORATOR_UPDATE_NSID: &str = "sh.tangled.repo.collaboratorUpdate"; #[derive(Deserialize)] struct FrameWire { @@ -48,10 +27,10 @@ pub(crate) fn dispatch(text: &str, handler: &dyn FeedHandler, cursor: &mut Curso return Outcome::Idle; } cursor.advance(frame.created); - match (RosterUpdate::of(&frame.nsid), frame.event.repo) { - (Some(RosterUpdate::KnotMember), _) => handler.roster_changed(), - (Some(RosterUpdate::RepoCollaborator), Some(repo)) => handler.repo_roster_changed(&repo), - _ => {} + if frame.nsid.as_str() == COLLABORATOR_UPDATE_NSID + && let Some(repo) = frame.event.repo + { + handler.repo_roster_changed(&repo); } Outcome::Progressed } @@ -81,43 +60,50 @@ mod tests { } } - const MEMBER: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.knot.memberUpdate", - "event":{"op":"add","subject":"did:plc:limpet"},"created":6521339806120}"#; - const MEMBER_OLD: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.knot.memberUpdate", - "event":{"op":"add","subject":"did:plc:limpet"},"created":7}"#; - const MEMBER_UNDATED: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.knot.memberUpdate", - "event":{"op":"add","subject":"did:plc:limpet"},"created":0}"#; const COLLABORATOR: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.repo.collaboratorUpdate", + "event":{"op":"remove","subject":"did:plc:limpet","repo":"did:plc:scallop"},"created":6521339806120}"#; + const COLLABORATOR_OLD: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.repo.collaboratorUpdate", "event":{"op":"remove","subject":"did:plc:limpet","repo":"did:plc:scallop"},"created":7}"#; + const COLLABORATOR_UNDATED: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.repo.collaboratorUpdate", + "event":{"op":"remove","subject":"did:plc:limpet","repo":"did:plc:scallop"},"created":0}"#; + const COLLABORATOR_NO_REPO: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.repo.collaboratorUpdate", + "event":{"op":"remove","subject":"did:plc:limpet"},"created":7}"#; const REF: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.git.refUpdate", "event":{"ref":"refs/heads/main"},"created":9}"#; - const GARBLED: &str = r#"{"nsid":"sh.tangled.knot.memberUpdate"}"#; + const GARBLED: &str = r#"{"nsid":"sh.tangled.repo.collaboratorUpdate"}"#; const NOT_AN_NSID: &str = r#"{"rkey":"3mug","nsid":"not an nsid at all", "event":{"op":"add","subject":"did:plc:limpet"},"created":7}"#; #[test] fn each_legacy_frame_moves_cursor_and_calls_handler() { let cases: [(&str, &str, i64, bool, i64, &[&str]); 7] = [ - ("member update", MEMBER, 0, true, 6521339806120, &["roster"]), ( "collaborator update", COLLABORATOR, 0, true, - 7, + 6521339806120, &["did:plc:scallop"], ), ("unrelated ref update", REF, 4, true, 9, &[]), ( - "older member update", - MEMBER_OLD, + "older collaborator update", + COLLABORATOR_OLD, 100, true, 100, - &["roster"], + &["did:plc:scallop"], + ), + ( + "collaborator update naming no repo", + COLLABORATOR_NO_REPO, + 0, + true, + 7, + &[], ), ("garbled frame", GARBLED, 3, false, 3, &[]), - ("undated frame", MEMBER_UNDATED, 5, false, 5, &[]), + ("undated frame", COLLABORATOR_UNDATED, 5, false, 5, &[]), ("frame with a junk nsid", NOT_AN_NSID, 3, false, 3, &[]), ]; diff --git a/bobbin/crates/knot-ingest/src/orchestrator.rs b/bobbin/crates/knot-ingest/src/orchestrator.rs index 20705d22b..3a2a8a623 100644 --- a/bobbin/crates/knot-ingest/src/orchestrator.rs +++ b/bobbin/crates/knot-ingest/src/orchestrator.rs @@ -17,10 +17,10 @@ use tokio::sync::mpsc; use crate::client::{ AclEntry, Completeness, InviteEntry, KnotClient, KnotClientError, Listing, knot_endpoint, }; -use crate::firehose::{InviteOp, InviteTarget}; +use crate::firehose::InviteOp; use crate::gate::CapabilityGate; use crate::registry::KnotRegistry; -use crate::roster::{AclKind, Applied, Assertion, Roster, Scope, Source}; +use crate::roster::{AclKind, Applied, Assertion, Roster, Source}; use crate::stream::{Feed, FeedHandler, StreamConfig, run_stream}; const POLL_INTERVAL: Duration = Duration::from_secs(30); @@ -187,26 +187,22 @@ impl FeedHandler for KnotHandlers { fn invite_op(&self, op: InviteOp) { let InviteOp { - target, + repo, invitee, offer, } = op; - let scope = match &target { - InviteTarget::Knot => Scope::Knot, - InviteTarget::Repo(repo) => Scope::Repo(repo), - }; let mut roster = self.roster.lock(); match offer { Some(offer) => roster.apply( Assertion::Invite { - scope, + repo: &repo, invited_by: offer.invited_by, }, invitee, offer.offered_at, Source::Firehose, ), - None => roster.retire(AclKind::invite(scope), invitee), + None => roster.retire(AclKind::invite(&repo), invitee), } } } @@ -304,12 +300,6 @@ async fn apply_batch(ctx: &ReconcileCtx<'_>, batch: NudgeBatch) { async fn reconcile_once(ctx: &ReconcileCtx<'_>) { let horizon = ctx.roster.lock().applied(); - let members = ctx.client.list_members(ctx.endpoint, ctx.knot).await; - reconcile_listing(ctx, Scope::Knot, members, horizon); - - let invited = ctx.client.list_member_invites(ctx.endpoint, ctx.knot).await; - let offers = reconcile_listing(ctx, Scope::Knot, invited, horizon); - let repo_offers = futures::stream::iter(ctx.registry.repos(ctx.host)) .then(|repo| async move { reconcile_repo_once(ctx, &repo, horizon).await }) .fold( @@ -318,7 +308,7 @@ async fn reconcile_once(ctx: &ReconcileCtx<'_>) { ) .await; - if offers.and(repo_offers) == Completeness::Complete { + if repo_offers == Completeness::Complete { ctx.invites.backfilled(ctx.knot); } @@ -331,48 +321,48 @@ async fn reconcile_repo_once( horizon: Applied, ) -> Completeness { let collaborators = ctx.client.list_collaborators(ctx.endpoint, repo).await; - reconcile_listing(ctx, Scope::Repo(repo), collaborators, horizon); + reconcile_listing(ctx, repo, collaborators, horizon); let invited = ctx .client .list_collaborator_invites(ctx.endpoint, repo) .await; - reconcile_listing(ctx, Scope::Repo(repo), invited, horizon) + reconcile_listing(ctx, repo, invited, horizon) } trait Listed { - fn kind(scope: Scope<'_>) -> AclKind<'_>; + fn kind(repo: &Did) -> AclKind<'_>; fn subject(&self) -> &Did; - fn apply_to(self, roster: &mut Roster, scope: Scope<'_>, source: Source); + fn apply_to(self, roster: &mut Roster, repo: &Did, source: Source); } impl Listed for AclEntry { - fn kind(scope: Scope<'_>) -> AclKind<'_> { - AclKind::acl(scope) + fn kind(repo: &Did) -> AclKind<'_> { + AclKind::acl(repo) } fn subject(&self) -> &Did { &self.subject } - fn apply_to(self, roster: &mut Roster, scope: Scope<'_>, source: Source) { - roster.apply(Assertion::Acl(scope), self.subject, self.created_at, source); + fn apply_to(self, roster: &mut Roster, repo: &Did, source: Source) { + roster.apply(Assertion::Acl(repo), self.subject, self.created_at, source); } } impl Listed for InviteEntry { - fn kind(scope: Scope<'_>) -> AclKind<'_> { - AclKind::invite(scope) + fn kind(repo: &Did) -> AclKind<'_> { + AclKind::invite(repo) } fn subject(&self) -> &Did { &self.subject } - fn apply_to(self, roster: &mut Roster, scope: Scope<'_>, source: Source) { + fn apply_to(self, roster: &mut Roster, repo: &Did, source: Source) { roster.apply( Assertion::Invite { - scope, + repo, invited_by: self.added_by, }, self.subject, @@ -384,11 +374,11 @@ impl Listed for InviteEntry { fn reconcile_listing( ctx: &ReconcileCtx<'_>, - scope: Scope<'_>, + repo: &Did, listing: Result, KnotClientError>, horizon: Applied, ) -> Completeness { - let kind = T::kind(scope); + let kind = T::kind(repo); let Listing { entries, completeness, @@ -407,7 +397,7 @@ fn reconcile_listing( let mut guard = ctx.roster.lock(); entries .into_iter() - .for_each(|entry| entry.apply_to(&mut guard, scope, Source::Listing(horizon))); + .for_each(|entry| entry.apply_to(&mut guard, repo, Source::Listing(horizon))); if completeness == Completeness::Complete { guard.reap(kind, &present, horizon); } @@ -417,7 +407,7 @@ fn reconcile_listing( fn warn_unreadable(ctx: &ReconcileCtx<'_>, kind: AclKind<'_>, err: &KnotClientError) { tracing::warn!( host = %ctx.host, - repo = ?kind.scope.repo().map(|repo| repo.as_ref()), + repo = %kind.repo.as_ref(), listing = kind.listing(), error = %err, "knot reconcile failed" @@ -441,13 +431,6 @@ mod tests { Did::new_owned(s).unwrap() } - fn member_count(store: &EdgeStore, subject: &str) -> u64 { - store.count(&EdgeKey::new( - nsid_static("sh.tangled.knot.member"), - SubjectRef::Did(did(subject)), - )) - } - fn collaborator_count(store: &EdgeStore, repo: &str) -> u64 { store.count(&EdgeKey::new( nsid_static("sh.tangled.repo.collaborator"), @@ -510,7 +493,8 @@ mod tests { } } - const MEMBER_INVITES: &str = "sh.tangled.knot.listMemberInvites"; + const COLLABORATOR_INVITES: &str = "sh.tangled.repo.listCollaboratorInvites"; + const COLLABORATORS: &str = "sh.tangled.repo.listCollaborators"; async fn serve(h: &Harness, endpoint: &str, response: ResponseTemplate) { Mock::given(method("GET")) @@ -534,19 +518,9 @@ mod tests { } #[tokio::test] - async fn reconcile_backfills_members_and_collaborators() { + async fn reconcile_backfills_collaborators() { let h = harness().await; - Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) - .respond_with(ResponseTemplate::new(200).set_body_json(json!({ - "items": [ - {"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}, - {"subject": "did:plc:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"} - ] - }))) - .mount(&h.server) - .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.listCollaborators")) .and(query_param("subject", "did:plc:scallop")) @@ -562,17 +536,16 @@ mod tests { reconcile_once(&h.ctx()).await; - assert_eq!(member_count(&h.store, "did:plc:boltless"), 1); - assert_eq!(member_count(&h.store, "did:plc:akshay"), 1); assert_eq!(collaborator_count(&h.store, "did:plc:scallop"), 1); } #[tokio::test] - async fn reconcile_backfills_outstanding_member_invites() { + async fn reconcile_backfills_outstanding_collaborator_invites() { let h = harness().await; + h.registry.observe_repo(&h.host, did("did:plc:scallop")); serve( &h, - MEMBER_INVITES, + COLLABORATOR_INVITES, ResponseTemplate::new(200).set_body_json(offers(&["did:plc:limpet", "did:plc:olaren"])), ) .await; @@ -588,13 +561,14 @@ mod tests { #[tokio::test] async fn completed_backfill_stops_knot_being_pending() { let h = harness().await; + h.registry.observe_repo(&h.host, did("did:plc:scallop")); h.invites.awaiting(h.knot.clone()); h.invites.all_knots_discovered(); assert_eq!(h.invites.pending(), vec![h.knot.clone()]); serve( &h, - MEMBER_INVITES, + COLLABORATOR_INVITES, ResponseTemplate::new(200).set_body_json(offers(&[])), ) .await; @@ -606,10 +580,11 @@ mod tests { #[tokio::test] async fn knot_without_invite_listings_stops_being_pending() { let h = harness().await; + h.registry.observe_repo(&h.host, did("did:plc:scallop")); h.invites.awaiting(h.knot.clone()); h.invites.all_knots_discovered(); - serve(&h, MEMBER_INVITES, ResponseTemplate::new(404)).await; + serve(&h, COLLABORATOR_INVITES, ResponseTemplate::new(404)).await; reconcile_once(&h.ctx()).await; assert!( @@ -621,12 +596,13 @@ mod tests { #[tokio::test] async fn failed_backfill_leaves_only_its_own_knot_pending() { let h = harness().await; + h.registry.observe_repo(&h.host, did("did:plc:scallop")); h.invites.awaiting(h.knot.clone()); h.invites.awaiting(did("did:web:knot.nel.pet")); h.invites.backfilled(&did("did:web:knot.nel.pet")); h.invites.all_knots_discovered(); - serve(&h, MEMBER_INVITES, ResponseTemplate::new(500)).await; + serve(&h, COLLABORATOR_INVITES, ResponseTemplate::new(500)).await; reconcile_once(&h.ctx()).await; assert_eq!( @@ -637,77 +613,80 @@ mod tests { } #[tokio::test] - async fn reconcile_reaps_departed_member() { + async fn reconcile_reaps_departed_collaborator() { let h = harness().await; + let repo = did("did:plc:scallop"); + h.registry.observe_repo(&h.host, repo.clone()); - Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) - .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + serve( + &h, + COLLABORATORS, + ResponseTemplate::new(200).set_body_json(json!({ "items": [ {"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}, {"subject": "did:plc:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"} ] - }))) - .mount(&h.server) - .await; + })), + ) + .await; reconcile_once(&h.ctx()).await; - assert_eq!(member_count(&h.store, "did:plc:boltless"), 1); - assert_eq!(member_count(&h.store, "did:plc:akshay"), 1); + assert_eq!(collaborator_count(&h.store, "did:plc:scallop"), 2); h.server.reset().await; - Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) - .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + serve( + &h, + COLLABORATORS, + ResponseTemplate::new(200).set_body_json(json!({ "items": [ {"subject": "did:plc:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"} ] - }))) - .mount(&h.server) - .await; + })), + ) + .await; reconcile_once(&h.ctx()).await; assert_eq!( - member_count(&h.store, "did:plc:boltless"), - 0, - "a member dropped from the authoritative snapshot is reaped on reconcile" + collaborator_count(&h.store, "did:plc:scallop"), + 1, + "a collaborator dropped from the authoritative snapshot is reaped on reconcile" ); - assert_eq!(member_count(&h.store, "did:plc:akshay"), 1); } #[tokio::test] - async fn reconcile_skips_reap_when_member_list_truncated() { + async fn reconcile_skips_reap_when_collaborator_list_truncated() { let h = harness().await; + let repo = did("did:plc:scallop"); + h.registry.observe_repo(&h.host, repo.clone()); - Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) - .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + serve( + &h, + COLLABORATORS, + ResponseTemplate::new(200).set_body_json(json!({ "items": [{"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}], "cursor": "more" - }))) - .mount(&h.server) - .await; + })), + ) + .await; - let stayed = did("did:plc:akshay"); { let mut roster = h.roster.lock(); let taken = roster.applied(); roster.apply( - Assertion::Acl(Scope::Knot), - stayed, + Assertion::Acl(&repo), + did("did:plc:akshay"), UnixMicros::new(1_000_000), Source::Listing(taken), ); } - assert_eq!(member_count(&h.store, "did:plc:akshay"), 1); + assert_eq!(collaborator_count(&h.store, "did:plc:scallop"), 1); reconcile_once(&h.ctx()).await; assert_eq!( - member_count(&h.store, "did:plc:akshay"), - 1, - "a truncated member snapshot must not reap members it could not enumerate" + collaborator_count(&h.store, "did:plc:scallop"), + 2, + "a truncated snapshot must not reap collaborators it could not enumerate" ); - assert_eq!(member_count(&h.store, "did:plc:boltless"), 1); } fn seed_legacy_edge(store: &EdgeStore, kind: &'static str, subject: &str, source: &str) { @@ -727,15 +706,6 @@ mod tests { async fn reconcile_purges_seeded_legacy_acl() { let h = harness().await; - Mock::given(method("GET")) - .and(path("/xrpc/sh.tangled.knot.listMembers")) - .respond_with(ResponseTemplate::new(200).set_body_json(json!({ - "items": [ - {"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"} - ] - }))) - .mount(&h.server) - .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.listCollaborators")) .and(query_param("subject", "did:plc:scallop")) @@ -750,15 +720,6 @@ mod tests { let repo = did("did:plc:scallop"); h.registry.observe_repo(&h.host, repo.clone()); - let member_source = "at://did:plc:akshay/sh.tangled.knot.member/r1"; - seed_legacy_edge( - &h.store, - "sh.tangled.knot.member", - "did:plc:boltless", - member_source, - ); - h.registry - .note_legacy_member(AtUri::new_owned(member_source).unwrap(), &h.host); seed_legacy_edge( &h.store, "sh.tangled.repo.collaborator", @@ -768,11 +729,6 @@ mod tests { reconcile_once(&h.ctx()).await; - assert_eq!( - member_count(&h.store, "did:plc:boltless"), - 1, - "stale legacy member purged, knot-owned member kept" - ); assert_eq!( collaborator_count(&h.store, "did:plc:scallop"), 1, @@ -812,12 +768,14 @@ mod tests { #[tokio::test] async fn handlers_apply_firehose_invite_straight_to_roster() { let h = harness().await; + let repo = did("did:plc:scallop"); + h.registry.observe_repo(&h.host, repo.clone()); let (tx, _rx) = mpsc::channel(4); let handlers = test_handlers(h.roster.clone(), tx); let invitee = did("did:plc:limpet"); handlers.invite_op(InviteOp { - target: InviteTarget::Knot, + repo: repo.clone(), invitee: invitee.clone(), offer: Some(Offer { invited_by: did("did:plc:akshay"), @@ -827,7 +785,7 @@ mod tests { assert_eq!(h.invites.outstanding(&invitee).len(), 1); handlers.invite_op(InviteOp { - target: InviteTarget::Knot, + repo: repo.clone(), invitee: invitee.clone(), offer: None, }); diff --git a/bobbin/crates/knot-ingest/src/registry.rs b/bobbin/crates/knot-ingest/src/registry.rs index 67f731f3d..c37d856b3 100644 --- a/bobbin/crates/knot-ingest/src/registry.rs +++ b/bobbin/crates/knot-ingest/src/registry.rs @@ -4,14 +4,12 @@ use std::sync::Mutex; use bobbin_types::knot_acl::KnotHostKey; use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; -use jacquard_common::types::string::AtUri; #[derive(Default)] struct Inner { hosts: HashSet, repos: HashMap>>, repo_host: HashMap, KnotHostKey>, - legacy_members: HashMap, KnotHostKey>, } #[derive(Default)] @@ -66,31 +64,6 @@ impl KnotRegistry { self.inner.lock().unwrap().repo_host.get(repo).cloned() } - pub fn note_legacy_member(&self, source: AtUri, host: &KnotHostKey) { - self.inner - .lock() - .unwrap() - .legacy_members - .insert(source, host.clone()); - } - - pub fn forget_legacy_member(&self, source: &AtUri) { - self.inner.lock().unwrap().legacy_members.remove(source); - } - - pub fn drain_legacy_members(&self, host: &KnotHostKey) -> Vec> { - let mut inner = self.inner.lock().unwrap(); - let matched: Vec> = inner - .legacy_members - .iter() - .filter(|(_, member_host)| member_host.as_str() == host.as_str()) - .map(|(source, _)| source.clone()) - .collect(); - matched.iter().for_each(|source| { - inner.legacy_members.remove(source); - }); - matched - } } #[cfg(test)] @@ -101,10 +74,6 @@ mod tests { Did::new_owned(s).unwrap() } - fn at(s: &str) -> AtUri { - AtUri::new_owned(s).unwrap() - } - fn host(s: &str) -> KnotHostKey { KnotHostKey::new(s) } @@ -159,39 +128,4 @@ mod tests { assert_eq!(registry.host_of_repo(&did("did:plc:limpet")), None); } - #[test] - fn drain_legacy_members_returns_only_matching_host() { - let registry = KnotRegistry::new(); - let here = at("at://did:plc:akshay/sh.tangled.knot.member/r1"); - let elsewhere = at("at://did:plc:akshay/sh.tangled.knot.member/r2"); - registry.note_legacy_member(here.clone(), &host("oyster.cafe")); - registry.note_legacy_member(elsewhere.clone(), &host("nel.pet")); - - assert_eq!( - registry.drain_legacy_members(&host("oyster.cafe")), - vec![here] - ); - assert!( - registry - .drain_legacy_members(&host("oyster.cafe")) - .is_empty() - ); - assert_eq!( - registry.drain_legacy_members(&host("nel.pet")), - vec![elsewhere] - ); - } - - #[test] - fn forget_legacy_member_drops_source() { - let registry = KnotRegistry::new(); - let source = at("at://did:plc:akshay/sh.tangled.knot.member/r1"); - registry.note_legacy_member(source.clone(), &host("oyster.cafe")); - registry.forget_legacy_member(&source); - assert!( - registry - .drain_legacy_members(&host("oyster.cafe")) - .is_empty() - ); - } } diff --git a/bobbin/crates/knot-ingest/src/roster.rs b/bobbin/crates/knot-ingest/src/roster.rs index 0a3dfc362..b359bd848 100644 --- a/bobbin/crates/knot-ingest/src/roster.rs +++ b/bobbin/crates/knot-ingest/src/roster.rs @@ -18,46 +18,31 @@ enum Facet { Invite, } -#[derive(Clone, Copy, Eq, Hash, PartialEq)] -pub enum Scope<'a> { - Knot, - Repo(&'a Did), -} - -impl<'a> Scope<'a> { - pub(crate) fn repo(self) -> Option<&'a Did> { - match self { - Self::Knot => None, - Self::Repo(repo) => Some(repo), - } - } -} - #[derive(Clone, Eq, Hash, PartialEq)] struct DedupKey { facet: Facet, - repo: Option>, + repo: Did, subject: Did, } #[derive(Clone, Copy)] pub struct AclKind<'a> { facet: Facet, - pub scope: Scope<'a>, + pub repo: &'a Did, } impl<'a> AclKind<'a> { - pub(crate) fn acl(scope: Scope<'a>) -> Self { + pub(crate) fn acl(repo: &'a Did) -> Self { Self { facet: Facet::Acl, - scope, + repo, } } - pub(crate) fn invite(scope: Scope<'a>) -> Self { + pub(crate) fn invite(repo: &'a Did) -> Self { Self { facet: Facet::Invite, - scope, + repo, } } @@ -66,27 +51,25 @@ impl<'a> AclKind<'a> { } pub(crate) fn listing(self) -> &'static str { - match (self.facet, self.scope) { - (Facet::Acl, Scope::Knot) => "members", - (Facet::Acl, Scope::Repo(_)) => "collaborators", - (Facet::Invite, Scope::Knot) => "member invites", - (Facet::Invite, Scope::Repo(_)) => "collaborator invites", + match self.facet { + Facet::Acl => "collaborators", + Facet::Invite => "collaborator invites", } } fn dedup(self, subject: Did) -> DedupKey { DedupKey { facet: self.facet, - repo: self.scope.repo().cloned(), + repo: self.repo.clone(), subject, } } } pub enum Assertion<'a> { - Acl(Scope<'a>), + Acl(&'a Did), Invite { - scope: Scope<'a>, + repo: &'a Did, invited_by: Did, }, } @@ -94,8 +77,8 @@ pub enum Assertion<'a> { impl<'a> Assertion<'a> { fn kind(&self) -> AclKind<'a> { match self { - Self::Acl(scope) => AclKind::acl(*scope), - Self::Invite { scope, .. } => AclKind::invite(*scope), + Self::Acl(repo) => AclKind::acl(repo), + Self::Invite { repo, .. } => AclKind::invite(repo), } } } @@ -144,13 +127,10 @@ impl Roster { } } - fn invite_scope(&self, scope: Scope<'_>) -> InviteScope { - match scope { - Scope::Knot => InviteScope::Knot(self.knot.clone()), - Scope::Repo(repo) => InviteScope::Repo { - knot: self.knot.clone(), - repo: repo.clone(), - }, + fn invite_scope(&self, repo: &Did) -> InviteScope { + InviteScope { + knot: self.knot.clone(), + repo: repo.clone(), } } @@ -162,29 +142,19 @@ impl Roster { source: Source, ) { let kind = assertion.kind(); - if !kind - .scope - .repo() - .is_none_or(|repo| self.registry.repo_on_host(&self.host, repo)) - { + if !self.registry.repo_on_host(&self.host, kind.repo) { return; } if !self.advance(kind.dedup(subject.clone()), at, source) { return; } match assertion { - Assertion::Acl(scope) => { - let upserted = match scope { - Scope::Knot => knot_acl::member_upsert(&self.knot, &subject, at.raw()), - Scope::Repo(repo) => knot_acl::collaborator_upsert(repo, &subject, at.raw()), - }; - upserted - .into_iter() - .for_each(|(record, edges)| self.store.upsert_source(&record, edges)); - } - Assertion::Invite { scope, invited_by } => self.invites.offer( + Assertion::Acl(repo) => knot_acl::collaborator_upsert(repo, &subject, at.raw()) + .into_iter() + .for_each(|(record, edges)| self.store.upsert_source(&record, edges)), + Assertion::Invite { repo, invited_by } => self.invites.offer( subject, - self.invite_scope(scope), + self.invite_scope(repo), Offer { invited_by, offered_at: at, @@ -213,7 +183,7 @@ impl Roster { .iter() .filter(|(key, state)| { key.facet == kind.facet - && key.repo.as_ref() == kind.scope.repo() + && &key.repo == kind.repo && state.present && state.applied <= horizon && !present.contains(&key.subject) @@ -226,11 +196,6 @@ impl Roster { } pub fn purge_legacy(&self) { - self.registry - .drain_legacy_members(&self.host) - .iter() - .for_each(|source| self.store.remove_source(source)); - self.registry .repos(&self.host) .into_iter() @@ -246,18 +211,12 @@ impl Roster { pub fn retire(&mut self, kind: AclKind<'_>, subject: Did) { match kind.facet { - Facet::Acl => { - let source = match kind.scope { - Scope::Knot => knot_acl::member_source(&self.knot, &subject), - Scope::Repo(repo) => knot_acl::collaborator_source(repo, &subject), - }; - source - .iter() - .for_each(|record| self.store.remove_source(record)); - } + Facet::Acl => knot_acl::collaborator_source(kind.repo, &subject) + .iter() + .for_each(|record| self.store.remove_source(record)), Facet::Invite => self .invites - .withdraw(&subject, &self.invite_scope(kind.scope)), + .withdraw(&subject, &self.invite_scope(kind.repo)), } let applied = self.step(); if let Some(state) = self.seen.get_mut(&kind.dedup(subject)) { @@ -339,13 +298,6 @@ mod tests { knot_acl::host_to_knot_did("oyster.cafe").unwrap() } - fn member_count(store: &EdgeStore, subject: &Did) -> u64 { - store.count(&EdgeKey::new( - nsid_static("sh.tangled.knot.member"), - SubjectRef::Did(subject.clone()), - )) - } - fn collaborator_count(store: &EdgeStore, repo: &Did) -> u64 { store.count(&EdgeKey::new( nsid_static("sh.tangled.repo.collaborator"), @@ -361,7 +313,7 @@ mod tests { Arc::new(InviteIndex::new()) } - fn member_roster(store: Arc) -> Roster { + fn bare_roster(store: Arc) -> Roster { Roster::new(store, invites(), knot(), empty_registry(), host()) } @@ -375,30 +327,21 @@ mod tests { Roster::new(store, index, knot(), registry, host()) } - fn indexed_roster(index: Arc) -> Roster { - Roster::new(store(), index, knot(), empty_registry(), host()) + fn indexed_roster(index: Arc, repo: &Did) -> Roster { + repo_roster(store(), index, repo) } fn subjects(items: &[&Did]) -> HashSet> { items.iter().map(|d| (*d).clone()).collect() } - const MEMBER: Assertion<'static> = Assertion::Acl(Scope::Knot); - fn collab(repo: &Did) -> Assertion<'_> { - Assertion::Acl(Scope::Repo(repo)) - } - - fn offer_of(admin: &str) -> Assertion<'static> { - Assertion::Invite { - scope: Scope::Knot, - invited_by: did(admin), - } + Assertion::Acl(repo) } - fn repo_offer_of<'a>(repo: &'a Did, admin: &str) -> Assertion<'a> { + fn offer_of<'a>(repo: &'a Did, admin: &str) -> Assertion<'a> { Assertion::Invite { - scope: Scope::Repo(repo), + repo, invited_by: did(admin), } } @@ -420,23 +363,25 @@ mod tests { } #[test] - fn late_add_leaves_newer_member_row_standing() { + fn late_add_leaves_newer_row_standing() { let store = store(); - let mut roster = member_roster(store.clone()); - let m = did("did:plc:boltless"); - roster.listed(MEMBER, &m, 5); - roster.listed(MEMBER, &m, 1); - assert_eq!(member_count(&store, &m), 1); + let repo = did("did:plc:scallop"); + let mut roster = repo_roster(store.clone(), invites(), &repo); + let who = did("did:plc:boltless"); + roster.listed(collab(&repo), &who, 5); + roster.listed(collab(&repo), &who, 1); + assert_eq!(collaborator_count(&store, &repo), 1); } #[test] fn duplicate_cursor_is_idempotent() { let store = store(); - let mut roster = member_roster(store.clone()); - let m = did("did:plc:akshay"); - roster.listed(MEMBER, &m, 10); - roster.listed(MEMBER, &m, 10); - assert_eq!(member_count(&store, &m), 1); + let repo = did("did:plc:scallop"); + let mut roster = repo_roster(store.clone(), invites(), &repo); + let who = did("did:plc:akshay"); + roster.listed(collab(&repo), &who, 10); + roster.listed(collab(&repo), &who, 10); + assert_eq!(collaborator_count(&store, &repo), 1); } #[test] @@ -449,55 +394,43 @@ mod tests { } #[test] - fn reap_removes_departed_member() { + fn reap_removes_departed_collaborator() { let store = store(); - let mut roster = member_roster(store.clone()); + let repo = did("did:plc:scallop"); + let mut roster = repo_roster(store.clone(), invites(), &repo); let stayed = did("did:plc:akshay"); let left = did("did:plc:boltless"); - roster.listed(MEMBER, &stayed, 10); - roster.listed(MEMBER, &left, 20); + roster.listed(collab(&repo), &stayed, 10); + roster.listed(collab(&repo), &left, 20); + assert_eq!(collaborator_count(&store, &repo), 2); let horizon = roster.applied(); - roster.reap(AclKind::acl(Scope::Knot), &subjects(&[&stayed]), horizon); + roster.reap(AclKind::acl(&repo), &subjects(&[&stayed]), horizon); - assert_eq!(member_count(&store, &stayed), 1); assert_eq!( - member_count(&store, &left), - 0, - "a member absent from the authoritative snapshot is reaped" + collaborator_count(&store, &repo), + 1, + "a collaborator absent from the authoritative snapshot is reaped" ); } #[test] - fn reap_skips_member_added_after_horizon() { + fn reap_skips_collaborator_added_after_horizon() { let store = store(); - let mut roster = member_roster(store.clone()); - let m = did("did:plc:boltless"); + let repo = did("did:plc:scallop"); + let mut roster = repo_roster(store.clone(), invites(), &repo); let horizon = roster.applied(); - roster.listed(MEMBER, &m, 100); + roster.listed(collab(&repo), &did("did:plc:boltless"), 100); - roster.reap(AclKind::acl(Scope::Knot), &subjects(&[]), horizon); + roster.reap(AclKind::acl(&repo), &subjects(&[]), horizon); assert_eq!( - member_count(&store, &m), + collaborator_count(&store, &repo), 1, - "a member added after the snapshot horizon must survive the reap" + "a collaborator added after the snapshot horizon must survive the reap" ); } - #[test] - fn reap_removes_departed_collaborator() { - let store = store(); - let repo = did("did:plc:scallop"); - let mut roster = repo_roster(store.clone(), invites(), &repo); - roster.listed(collab(&repo), &did("did:plc:olaren"), 7); - - let horizon = roster.applied(); - roster.reap(AclKind::acl(Scope::Repo(&repo)), &subjects(&[]), horizon); - - assert_eq!(collaborator_count(&store, &repo), 0); - } - #[test] fn purge_legacy_strips_pds_collaborator_keeps_knot_owned() { let store = store(); @@ -521,27 +454,11 @@ mod tests { ); } - #[test] - fn purge_legacy_removes_indexed_member_edges() { - let store = store(); - let registry = empty_registry(); - let roster = Roster::new(store.clone(), invites(), knot(), registry.clone(), host()); - let member = did("did:plc:boltless"); - let source = at("at://did:plc:akshay/sh.tangled.knot.member/r1"); - - add_legacy_edge(&store, "sh.tangled.knot.member", &member, &source); - registry.note_legacy_member(source.clone(), &host()); - assert_eq!(member_count(&store, &member), 1); - - roster.purge_legacy(); - assert_eq!(member_count(&store, &member), 0); - } - #[test] fn collaborator_for_unhosted_repo_is_dropped() { let store = store(); let repo = did("did:plc:scallop"); - let mut roster = member_roster(store.clone()); + let mut roster = bare_roster(store.clone()); roster.listed(collab(&repo), &did("did:plc:olaren"), 7); assert_eq!( collaborator_count(&store, &repo), @@ -551,55 +468,55 @@ mod tests { } #[test] - fn offer_reaches_index_scoped_to_its_source_and_leaves_when_retired() { + fn offer_reaches_index_scoped_to_its_repo_and_leaves_when_retired() { let index = invites(); - let repo = did("did:plc:scallop"); - let mut roster = repo_roster(store(), index.clone(), &repo); - let member = did("did:plc:limpet"); - let collaborator = did("did:plc:olaren"); - - roster.listed(offer_of("did:plc:akshay"), &member, 7_000_000); - roster.listed( - repo_offer_of(&repo, "did:plc:boltless"), - &collaborator, - 9_000_000, - ); + let store = store(); + let one = did("did:plc:scallop"); + let two = did("did:plc:whelk"); + let registry = empty_registry(); + registry.observe_repo(&host(), one.clone()); + registry.observe_repo(&host(), two.clone()); + let mut roster = Roster::new(store, index.clone(), knot(), registry, host()); + let here = did("did:plc:limpet"); + let there = did("did:plc:olaren"); - let offered = index.outstanding(&member); + roster.listed(offer_of(&one, "did:plc:akshay"), &here, 7_000_000); + roster.listed(offer_of(&two, "did:plc:boltless"), &there, 9_000_000); + + let offered = index.outstanding(&here); assert_eq!(offered.len(), 1); - assert_eq!(offered[0].0, InviteScope::Knot(knot())); - assert_eq!(offered[0].1.invited_by, did("did:plc:akshay")); - assert_eq!(offered[0].1.offered_at, UnixMicros::new(7_000_000)); assert_eq!( - index.outstanding(&collaborator)[0].0, - InviteScope::Repo { + offered[0].0, + InviteScope { knot: knot(), - repo: repo.clone() - }, - "repo's offer is scoped to repo and never to knot" + repo: one.clone() + } ); + assert_eq!(offered[0].1.invited_by, did("did:plc:akshay")); + assert_eq!(offered[0].1.offered_at, UnixMicros::new(7_000_000)); - roster.retire(AclKind::invite(Scope::Knot), member.clone()); - assert!(index.outstanding(&member).is_empty()); + roster.retire(AclKind::invite(&one), here.clone()); + assert!(index.outstanding(&here).is_empty()); assert_eq!( - index.outstanding(&collaborator).len(), + index.outstanding(&there).len(), 1, - "withdrawing knot's offer leaves repo's standing" + "withdrawing one repo's offer leaves the other standing" ); } #[test] fn reap_drops_offer_knot_stopped_listing_and_takes_it_back_when_it_returns() { let index = invites(); - let mut roster = indexed_roster(index.clone()); + let repo = did("did:plc:scallop"); + let mut roster = indexed_roster(index.clone(), &repo); let stayed = did("did:plc:limpet"); let withdrawn = did("did:plc:boltless"); let moment = 8_000_000; - roster.listed(offer_of("did:plc:akshay"), &stayed, 7_000_000); - roster.listed(offer_of("did:plc:akshay"), &withdrawn, moment); + roster.listed(offer_of(&repo, "did:plc:akshay"), &stayed, 7_000_000); + roster.listed(offer_of(&repo, "did:plc:akshay"), &withdrawn, moment); let horizon = roster.applied(); - roster.reap(AclKind::invite(Scope::Knot), &subjects(&[&stayed]), horizon); + roster.reap(AclKind::invite(&repo), &subjects(&[&stayed]), horizon); assert_eq!(index.outstanding(&stayed).len(), 1); assert!( @@ -607,27 +524,28 @@ mod tests { "knot's filtered listing wins over firehose record it no longer reports" ); - roster.listed(offer_of("did:plc:akshay"), &withdrawn, moment); + roster.listed(offer_of(&repo, "did:plc:akshay"), &withdrawn, moment); assert_eq!( index.outstanding(&withdrawn).len(), 1, - "listMemberInvites hides blocked invitees, and lifting the block returns offers at timestamps they always had" + "listCollaboratorInvites hides blocked invitees, and lifting the block returns offers at timestamps they always had" ); } #[test] fn listing_fetched_before_withdrawal_leaves_offer_withdrawn() { let index = invites(); - let mut roster = indexed_roster(index.clone()); + let repo = did("did:plc:scallop"); + let mut roster = indexed_roster(index.clone(), &repo); let invitee = did("did:plc:limpet"); let moment = 8_000_000; - roster.listed(offer_of("did:plc:akshay"), &invitee, moment); + roster.listed(offer_of(&repo, "did:plc:akshay"), &invitee, moment); let fetched = roster.applied(); - roster.retire(AclKind::invite(Scope::Knot), invitee.clone()); + roster.retire(AclKind::invite(&repo), invitee.clone()); roster.apply( - offer_of("did:plc:akshay"), + offer_of(&repo, "did:plc:akshay"), invitee.clone(), UnixMicros::new(moment), Source::Listing(fetched), @@ -642,13 +560,14 @@ mod tests { #[test] fn firehose_reoffer_in_same_second_reaches_invitee() { let index = invites(); - let mut roster = indexed_roster(index.clone()); + let repo = did("did:plc:scallop"); + let mut roster = indexed_roster(index.clone(), &repo); let invitee = did("did:plc:limpet"); let moment = 8_000_000; - roster.firehosed(offer_of("did:plc:akshay"), &invitee, moment); - roster.retire(AclKind::invite(Scope::Knot), invitee.clone()); - roster.firehosed(offer_of("did:plc:olaren"), &invitee, moment); + roster.firehosed(offer_of(&repo, "did:plc:akshay"), &invitee, moment); + roster.retire(AclKind::invite(&repo), invitee.clone()); + roster.firehosed(offer_of(&repo, "did:plc:olaren"), &invitee, moment); let outstanding = index.outstanding(&invitee); assert_eq!( @@ -662,16 +581,17 @@ mod tests { #[test] fn offer_applied_mid_fetch_survives_snapshot_that_missed_it() { let index = invites(); - let mut roster = indexed_roster(index.clone()); + let repo = did("did:plc:scallop"); + let mut roster = indexed_roster(index.clone(), &repo); let listed = did("did:plc:limpet"); let raced = did("did:plc:olaren"); let moment = 8_000_000; - roster.listed(offer_of("did:plc:akshay"), &listed, moment); + roster.listed(offer_of(&repo, "did:plc:akshay"), &listed, moment); let horizon = roster.applied(); - roster.firehosed(offer_of("did:plc:akshay"), &raced, moment); - roster.reap(AclKind::invite(Scope::Knot), &subjects(&[&listed]), horizon); + roster.firehosed(offer_of(&repo, "did:plc:akshay"), &raced, moment); + roster.reap(AclKind::invite(&repo), &subjects(&[&listed]), horizon); assert_eq!( index.outstanding(&raced).len(), diff --git a/bobbin/crates/knot-ingest/src/stream.rs b/bobbin/crates/knot-ingest/src/stream.rs index e27f3b2d8..66b1a85c2 100644 --- a/bobbin/crates/knot-ingest/src/stream.rs +++ b/bobbin/crates/knot-ingest/src/stream.rs @@ -10,7 +10,7 @@ use tokio_util::sync::CancellationToken; use url::Url; use crate::client::authority; -use crate::firehose::{self, Frame, InviteEvent, InviteOp}; +use crate::firehose::{self, Frame, InviteOp}; use crate::legacy; const SUBSCRIBE_REPOS_PATH: &str = "xrpc/com.atproto.sync.subscribeRepos"; @@ -302,13 +302,10 @@ fn dispatch_commit(commit: &firehose::Commit, handler: &dyn FeedHandler) { commit .records .iter() - .filter_map(|op| InviteEvent::of(&op.collection).map(|kind| (op, kind))) - .for_each(|(op, kind)| { - match kind { - InviteEvent::KnotMember => handler.roster_changed(), - InviteEvent::RepoCollaborator => handler.repo_roster_changed(&commit.repo), - } - match firehose::invite_op_of(op, kind, &commit.repo) { + .filter(|op| firehose::is_collaborator_invite(&op.collection)) + .for_each(|op| { + handler.repo_roster_changed(&commit.repo); + match firehose::invite_op_of(op, &commit.repo) { Some(invite) => handler.invite_op(invite), None => tracing::warn!( repo = %commit.repo.as_ref(), @@ -351,7 +348,7 @@ mod tests { use super::*; use crate::firehose::testcbor::*; - use crate::firehose::{InviteTarget, OpAction}; + use crate::firehose::OpAction; #[derive(Clone, Debug, PartialEq)] enum Entry { @@ -496,8 +493,8 @@ mod tests { const FIREHOSE_PATH: &str = "xrpc/com.atproto.sync.subscribeRepos"; const EVENTS_URL: &str = "ws://oyster.cafe/events"; - const MEMBER_UPDATE: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.knot.memberUpdate", - "event":{"op":"add","subject":"did:plc:limpet"},"created":11}"#; + const COLLABORATOR_UPDATE: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.repo.collaboratorUpdate", + "event":{"op":"add","subject":"did:plc:limpet","repo":"did:plc:scallop"},"created":11}"#; const STRAY_REF: &str = r#"{"nsid":"sh.tangled.git.refUpdate","created":9}"#; fn text(json: &str) -> WsMessage { @@ -526,13 +523,13 @@ mod tests { invite_commit( "did:plc:limpet", 100, - InviteEvent::KnotMember.collection(), + firehose::COLLABORATOR_INVITE_COLLECTION, "did:plc:boltless", ), invite_commit( repo, 200, - InviteEvent::RepoCollaborator.collection(), + firehose::COLLABORATOR_INVITE_COLLECTION, "did:plc:olaren", ), WsMessage::Ping(Bytes::from_static(b"ka")), @@ -574,7 +571,10 @@ mod tests { assert!(matches!(end, SessionEnd::Closed { progressed: true })); assert_eq!(cursor.seq(), 400); - assert_eq!(handler.entries(), vec![Entry::Full, Entry::Repo(did(repo))]); + assert_eq!( + handler.entries(), + vec![Entry::Repo(did("did:plc:limpet")), Entry::Repo(did(repo))] + ); let sent = sent.lock(); assert_eq!(sent.len(), 1); assert!(matches!(&sent[0], WsMessage::Pong(p) if p.as_ref() == b"ka")); @@ -603,9 +603,9 @@ mod tests { #[test] fn invite_create_both_nudges_roster_and_delivers_offer() { let handler = RecordingHandler::default(); - let collection = InviteEvent::KnotMember.collection(); + let collection = firehose::COLLABORATOR_INVITE_COLLECTION; let commit = invite_record_commit( - "did:web:knot.oyster.cafe", + "did:plc:scallop", OpAction::Create, collection, "did:plc:limpet", @@ -621,9 +621,9 @@ mod tests { assert_eq!( handler.entries(), vec![ - Entry::Full, + Entry::Repo(did("did:plc:scallop")), Entry::Offer(InviteOp { - target: InviteTarget::Knot, + repo: did("did:plc:scallop"), invitee: did("did:plc:limpet"), offer: Some(Offer { invited_by: did("did:plc:akshay"), @@ -641,7 +641,7 @@ mod tests { let commit = invite_record_commit( "did:plc:scallop", OpAction::Delete, - InviteEvent::RepoCollaborator.collection(), + firehose::COLLABORATOR_INVITE_COLLECTION, "did:plc:olaren", None, ); @@ -653,7 +653,7 @@ mod tests { vec![ Entry::Repo(did("did:plc:scallop")), Entry::Offer(InviteOp { - target: InviteTarget::Repo(did("did:plc:scallop")), + repo: did("did:plc:scallop"), invitee: did("did:plc:olaren"), offer: None, }), @@ -719,7 +719,7 @@ mod tests { invite_commit( "did:plc:limpet", 500, - InviteEvent::KnotMember.collection(), + firehose::COLLABORATOR_INVITE_COLLECTION, "did:plc:boltless", ), WsMessage::Close { @@ -735,7 +735,7 @@ mod tests { assert!(matches!(end, SessionEnd::Closed { progressed: true })); assert_eq!(cursor.seq(), 500); - assert_eq!(handler.entries(), vec![Entry::Full]); + assert_eq!(handler.entries(), vec![Entry::Repo(did("did:plc:limpet"))]); } #[tokio::test(start_paused = true)] @@ -751,7 +751,7 @@ mod tests { session(vec![invite_commit( "did:plc:limpet", 700, - InviteEvent::KnotMember.collection(), + firehose::COLLABORATOR_INVITE_COLLECTION, "did:plc:boltless", )]), session(vec![]), @@ -780,7 +780,7 @@ mod tests { assert_eq!( handler.entries(), - vec![Entry::Full, Entry::Full], + vec![Entry::Full, Entry::Repo(did("did:plc:limpet"))], "an outdated cursor resyncs to live and keeps delivering after the reconnect" ); assert!( @@ -791,13 +791,13 @@ mod tests { #[tokio::test(start_paused = true)] async fn refused_firehose_falls_back_to_knot_event_stream() { - let sessions = vec![session(vec![text(MEMBER_UPDATE)]), session(vec![])]; + let sessions = vec![session(vec![text(COLLABORATOR_UPDATE)]), session(vec![])]; let (entries, urls) = drive(1, sessions).await; assert_eq!( entries, - vec![Entry::Full], + vec![Entry::Repo(did("did:plc:scallop"))], "legacy frame arrives after the refusal" ); assert!(urls[0].ends_with(FIREHOSE_PATH), "firehose is tried first"); @@ -812,13 +812,13 @@ mod tests { #[tokio::test(start_paused = true)] async fn text_frame_on_firehose_switches_to_legacy_feed() { let stray = session(vec![text(STRAY_REF)]); - let member = session(vec![text(MEMBER_UPDATE)]); + let collaborator = session(vec![text(COLLABORATOR_UPDATE)]); - let (entries, urls) = drive(0, vec![stray, member]).await; + let (entries, urls) = drive(0, vec![stray, collaborator]).await; assert_eq!( entries, - vec![Entry::Full], + vec![Entry::Repo(did("did:plc:scallop"))], "legacy frame shows up once we switch" ); assert!( diff --git a/bobbin/crates/resolver/src/legacy_upgrade.rs b/bobbin/crates/resolver/src/legacy_upgrade.rs index 0771c7715..2834a7949 100644 --- a/bobbin/crates/resolver/src/legacy_upgrade.rs +++ b/bobbin/crates/resolver/src/legacy_upgrade.rs @@ -3,12 +3,11 @@ use bobbin_types::com_atproto::repo::strong_ref::StrongRef; use bobbin_types::edges::{ExtractError, Record}; use bobbin_types::legacy::{ LEGACY_COMMENT_SENTINEL_CID, LegacyCollaborator, LegacyIssue, LegacyIssueComment, - LegacyKnotMember, LegacyPublicKey, LegacyPull, LegacyPullComment, LegacyRecord, LegacyRepo, + LegacyPublicKey, LegacyPull, LegacyPullComment, LegacyRecord, LegacyRepo, LegacySource, LegacyStar, LegacyTarget, }; use bobbin_types::sh_tangled::feed::comment::Comment as FeedComment; use bobbin_types::sh_tangled::feed::star::{Repo as StarRepo, Star, StarString, StarSubject}; -use bobbin_types::sh_tangled::knot::member::Member as KnotMember; use bobbin_types::sh_tangled::markup::markdown::Markdown; use bobbin_types::sh_tangled::public_key::PublicKey; use bobbin_types::sh_tangled::repo::Repo; @@ -217,9 +216,8 @@ fn serialize_canon_variant(record: &Record) -> Result, serde Record::Star(r) => serde_json::to_vec(r), Record::PublicKey(r) => serde_json::to_vec(r), Record::Repo(r) => serde_json::to_vec(r), - Record::KnotMember(r) => serde_json::to_vec(r), _ => unreachable!( - "upgrade only produces FeedComment/Issue/Pull/Collaborator/Star/PublicKey/Repo/KnotMember" + "upgrade only produces FeedComment/Issue/Pull/Collaborator/Star/PublicKey/Repo" ), } } @@ -236,7 +234,6 @@ pub async fn upgrade(legacy: LegacyRecord, resolver: &RepoIdResolver) -> Option< LegacyRecord::Star(l) => upgrade_star(l, resolver).await.map(Record::Star), LegacyRecord::PublicKey(l) => Some(Record::PublicKey(upgrade_public_key(l))), LegacyRecord::Repo(l) => Some(Record::Repo(upgrade_repo(l))), - LegacyRecord::KnotMember(l) => Some(Record::KnotMember(upgrade_knot_member(l))), } } @@ -420,15 +417,6 @@ fn upgrade_repo(l: LegacyRepo) -> Repo { } } -fn upgrade_knot_member(l: LegacyKnotMember) -> KnotMember { - KnotMember { - created_at: l.added_at, - domain: l.domain, - subject: l.member, - extra_data: l.extra_data, - } -} - async fn upgrade_collaborator( l: LegacyCollaborator, resolver: &RepoIdResolver, @@ -928,23 +916,6 @@ mod tests { } } - #[tokio::test] - async fn upgrade_knot_member_renames_added_at_and_member() { - let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let json = br#"{"$type":"sh.tangled.knot.member","addedAt":"2025-03-31T05:14:09Z","domain":"knot.example","member":"did:plc:nel"}"#; - let legacy = - LegacyRecord::from_json_bytes(&nsid("sh.tangled.knot.member"), json).expect("decode"); - let canon = upgrade(legacy, &resolver).await.expect("upgrade"); - match canon { - Record::KnotMember(m) => { - assert_eq!(m.created_at.as_str(), "2025-03-31T05:14:09Z"); - assert_eq!(m.domain.as_str(), "knot.example"); - assert_eq!(m.subject, did("did:plc:nel")); - } - other => panic!("expected canon knot.member, got {other:?}"), - } - } - #[tokio::test] async fn legacy_pull_target_with_empty_repo_did_treats_as_none() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); diff --git a/bobbin/crates/resolver/src/normalize.rs b/bobbin/crates/resolver/src/normalize.rs index d2d1fc663..0ea7fd816 100644 --- a/bobbin/crates/resolver/src/normalize.rs +++ b/bobbin/crates/resolver/src/normalize.rs @@ -111,7 +111,6 @@ use bobbin_types::sh_tangled::feed::star::Star; use bobbin_types::sh_tangled::graph::follow::Follow; use bobbin_types::sh_tangled::graph::vouch::Vouch; use bobbin_types::sh_tangled::knot::Knot; -use bobbin_types::sh_tangled::knot::member::Member as KnotMember; use bobbin_types::sh_tangled::label::definition::Definition as LabelDefinition; use bobbin_types::sh_tangled::label::op::Op as LabelOp; use bobbin_types::sh_tangled::pipeline::status::Status as PipelineStatus; @@ -134,7 +133,6 @@ identity_normalize!( Follow, Vouch, Knot, - KnotMember, LabelDefinition, LabelOp, PipelineStatus, diff --git a/bobbin/crates/types/src/edges.rs b/bobbin/crates/types/src/edges.rs index a12d4ef7a..d504c9d87 100644 --- a/bobbin/crates/types/src/edges.rs +++ b/bobbin/crates/types/src/edges.rs @@ -16,7 +16,6 @@ use crate::sh_tangled::feed::star::Star; use crate::sh_tangled::graph::follow::Follow; use crate::sh_tangled::graph::vouch::Vouch; use crate::sh_tangled::knot::Knot; -use crate::sh_tangled::knot::member::Member as KnotMemberRecord; use crate::sh_tangled::label::definition::Definition as LabelDefinitionRecord; use crate::sh_tangled::label::op::Op as LabelOpRecord; use crate::sh_tangled::pipeline::Pipeline; @@ -63,7 +62,6 @@ pub enum Record { Follow(Follow), Vouch(Vouch), Knot(Knot), - KnotMember(KnotMemberRecord), LabelDefinition(LabelDefinitionRecord), LabelOp(LabelOpRecord), Pipeline(Pipeline), @@ -108,7 +106,6 @@ impl Record { "sh.tangled.graph.follow" => parse!(Follow), "sh.tangled.graph.vouch" => parse!(Vouch), "sh.tangled.knot" => parse!(Knot), - "sh.tangled.knot.member" => parse!(KnotMember), "sh.tangled.label.definition" => parse!(LabelDefinition), "sh.tangled.label.op" => parse!(LabelOp), "sh.tangled.pipeline" => parse!(Pipeline), @@ -138,7 +135,6 @@ impl Record { Self::Follow(_) => "sh.tangled.graph.follow", Self::Vouch(_) => "sh.tangled.graph.vouch", Self::Knot(_) => "sh.tangled.knot", - Self::KnotMember(_) => "sh.tangled.knot.member", Self::LabelDefinition(_) => "sh.tangled.label.definition", Self::LabelOp(_) => "sh.tangled.label.op", Self::Pipeline(_) => "sh.tangled.pipeline", @@ -193,7 +189,6 @@ impl Record { Self::Follow(r) => Some(&r.created_at), Self::Vouch(r) => Some(&r.created_at), Self::Knot(r) => Some(&r.created_at), - Self::KnotMember(r) => Some(&r.created_at), Self::LabelDefinition(r) => Some(&r.created_at), Self::LabelOp(r) => Some(&r.performed_at), Self::Pipeline(_) => None, @@ -219,7 +214,6 @@ impl Record { Self::FeedComment(r) => feed_comment_edges(source, r), Self::Reaction(r) => reaction_edges(source, r), Self::Follow(r) => follow_edges(source, r), - Self::KnotMember(r) => knot_member_edges(source, r), Self::LabelOp(r) => label_op_edges(source, r), Self::PipelineStatus(r) => pipeline_status_edges(source, r), Self::Artifact(r) => artifact_edges(source, r), @@ -317,7 +311,6 @@ const MIRROR_KINDS: &[(&str, &str)] = &[ ("sh.tangled.feed.reaction", "sh.tangled.feed.reaction.by"), ("sh.tangled.graph.follow", "sh.tangled.graph.follow.by"), ("sh.tangled.graph.vouch", "sh.tangled.graph.vouch.by"), - ("sh.tangled.knot.member", "sh.tangled.knot.member.by"), ("sh.tangled.label.op", "sh.tangled.label.op.by"), ("sh.tangled.pipeline", "sh.tangled.pipeline.by"), ( @@ -444,17 +437,6 @@ fn follow_edges( )) } -fn knot_member_edges( - source: &AtUri, - record: &KnotMemberRecord, -) -> Result, ExtractError> { - Ok(one_edge( - "sh.tangled.knot.member", - SubjectRef::Did(record.subject.clone()), - source, - )) -} - fn label_op_edges( source: &AtUri, record: &LabelOpRecord, @@ -1165,23 +1147,6 @@ mod tests { assert_eq!(edges[0].subject, uri_subj(pipeline_uri)); } - #[test] - fn knot_member_keys_on_subject_did() { - let edges = extract( - "sh.tangled.knot.member", - "at://did:plc:teq/sh.tangled.knot.member/abcabcabcabcz", - json!({ - "$type": "sh.tangled.knot.member", - "createdAt": "2026-05-01T00:00:00Z", - "subject": "did:plc:nel", - "domain": "oyster.cafe" - }), - ); - assert_eq!(edges.len(), 1); - assert_eq!(edges[0].kind, nsid("sh.tangled.knot.member")); - assert_eq!(edges[0].subject, did_subj("did:plc:nel")); - } - #[test] fn spindle_member_keys_on_subject_did() { let edges = extract( diff --git a/bobbin/crates/types/src/knot_acl.rs b/bobbin/crates/types/src/knot_acl.rs index 78cdc5fe0..a47b0bc96 100644 --- a/bobbin/crates/types/src/knot_acl.rs +++ b/bobbin/crates/types/src/knot_acl.rs @@ -8,10 +8,8 @@ use jacquard_common::types::string::AtUri; use crate::edges::Edge; use crate::ids::{SubjectRef, nsid_static}; -pub const KNOT_MEMBER_COLLECTION: &str = "sh.tangled.bobbin.knotMember"; pub const KNOT_COLLABORATOR_COLLECTION: &str = "sh.tangled.bobbin.knotCollaborator"; -const KNOT_MEMBER_KIND: &str = "sh.tangled.knot.member"; const REPO_COLLABORATOR_KIND: &str = "sh.tangled.repo.collaborator"; const DID_WEB_PREFIX: &str = "did:web:"; @@ -43,10 +41,6 @@ impl fmt::Display for KnotHostKey { #[derive(Clone, Debug, Eq, PartialEq)] pub enum KnotOwnedSource { - Member { - knot: Did, - subject: Did, - }, Collaborator { repo: Did, subject: Did, @@ -65,21 +59,6 @@ pub fn host_to_knot_did(host: &str) -> Option> { Did::new_owned(format!("{DID_WEB_PREFIX}{}", host.replace(':', "%3A"))).ok() } -pub fn knot_did_host(knot: &Did) -> Option { - let encoded = knot.as_ref().strip_prefix(DID_WEB_PREFIX)?; - if encoded.is_empty() { - return None; - } - Some(encoded.replace("%3A", ":")) -} - -pub fn member_source( - knot: &Did, - subject: &Did, -) -> Option> { - build_source(knot.as_ref(), KNOT_MEMBER_COLLECTION, subject.as_ref()) -} - pub fn collaborator_source( repo: &Did, subject: &Did, @@ -91,24 +70,6 @@ pub fn collaborator_source( ) } -pub fn member_upsert( - knot: &Did, - subject: &Did, - created_micros: u64, -) -> Option<(AtUri, Vec)> { - let source = member_source(knot, subject)?; - let edge = Edge { - kind: nsid_static(KNOT_MEMBER_KIND), - subject: SubjectRef::Did(subject.clone()), - source: source.clone(), - sort_micros: created_micros, - }; - Some(( - source, - crate::edges::subject_keyed_mirror(edge, SubjectRef::Did(knot.clone())), - )) -} - pub fn collaborator_upsert( repo: &Did, subject: &Did, @@ -129,25 +90,17 @@ pub fn collaborator_upsert( pub fn decode_knot_owned_source(source: &AtUri) -> Option { let collection = source.collection()?; - let collection = collection.as_ref(); - if collection != KNOT_MEMBER_COLLECTION && collection != KNOT_COLLABORATOR_COLLECTION { + if collection.as_ref() != KNOT_COLLABORATOR_COLLECTION { return None; } let AtIdentifier::Did(authority) = source.authority() else { return None; }; let subject = Did::new_owned(source.rkey()?.as_ref()).ok()?; - match collection { - KNOT_MEMBER_COLLECTION => Some(KnotOwnedSource::Member { - knot: Did::new_owned(authority.as_ref()).ok()?, - subject, - }), - KNOT_COLLABORATOR_COLLECTION => Some(KnotOwnedSource::Collaborator { - repo: Did::new_owned(authority.as_ref()).ok()?, - subject, - }), - _ => None, - } + Some(KnotOwnedSource::Collaborator { + repo: Did::new_owned(authority.as_ref()).ok()?, + subject, + }) } fn build_source(authority: &str, collection: &str, rkey: &str) -> Option> { @@ -167,17 +120,15 @@ mod tests { } #[test] - fn host_did_round_trips() { + fn host_becomes_a_did_web() { let knot = host_to_knot_did("oyster.cafe").unwrap(); assert_eq!(knot.as_ref(), "did:web:oyster.cafe"); - assert_eq!(knot_did_host(&knot), Some("oyster.cafe".to_owned())); } #[test] - fn host_did_round_trips_with_port() { + fn host_with_a_port_becomes_a_did_web() { let knot = host_to_knot_did("oyster.cafe:3000").unwrap(); assert_eq!(knot.as_ref(), "did:web:oyster.cafe%3A3000"); - assert_eq!(knot_did_host(&knot), Some("oyster.cafe:3000".to_owned())); } #[test] @@ -215,21 +166,6 @@ mod tests { ); } - #[test] - fn member_source_round_trips() { - let knot = host_to_knot_did("oyster.cafe").unwrap(); - let subject = did("did:plc:nel"); - let source = member_source(&knot, &subject).expect("build member source"); - assert_eq!( - source.as_ref(), - "at://did:web:oyster.cafe/sh.tangled.bobbin.knotMember/did:plc:nel" - ); - assert_eq!( - decode_knot_owned_source(&source), - Some(KnotOwnedSource::Member { knot, subject }) - ); - } - #[test] fn collaborator_source_round_trips() { let repo = did("did:plc:scallop"); @@ -246,43 +182,17 @@ mod tests { } #[test] - fn decode_ignores_legacy_pds_member_record() { - let source = at("at://did:plc:nel/sh.tangled.knot.member/abcabcabcabcz"); + fn decode_ignores_pds_mirrored_record() { + let source = at("at://did:plc:nel/sh.tangled.repo.collaborator/abcabcabcabcz"); assert_eq!(decode_knot_owned_source(&source), None); } #[test] fn decode_ignores_handle_authority() { - let source = at("at://oyster.cafe/sh.tangled.bobbin.knotMember/did:plc:nel"); + let source = at("at://oyster.cafe/sh.tangled.bobbin.knotCollaborator/did:plc:nel"); assert_eq!(decode_knot_owned_source(&source), None); } - #[test] - fn member_upsert_builds_primary_and_knot_mirror_edges() { - let knot = host_to_knot_did("oyster.cafe").unwrap(); - let subject = did("did:plc:nel"); - let (source, edges) = - member_upsert(&knot, &subject, 1_700_000_000_000_000).expect("build member upsert"); - assert_eq!(edges.len(), 2); - - let primary = &edges[0]; - assert_eq!(primary.kind.as_ref(), "sh.tangled.knot.member"); - assert_eq!(primary.subject, SubjectRef::Did(subject.clone())); - assert_eq!(primary.source, source); - assert_eq!(primary.sort_micros, 1_700_000_000_000_000); - - let mirror = &edges[1]; - assert_eq!(mirror.kind.as_ref(), "sh.tangled.knot.member.by"); - assert_eq!(mirror.subject, SubjectRef::Did(knot.clone())); - assert_eq!(mirror.source, source); - assert_eq!(mirror.sort_micros, 1_700_000_000_000_000); - - assert_eq!( - decode_knot_owned_source(&source), - Some(KnotOwnedSource::Member { knot, subject }) - ); - } - #[test] fn collaborator_upsert_builds_primary_and_subject_mirror_edges() { let repo = did("did:plc:scallop"); diff --git a/bobbin/crates/types/src/legacy.rs b/bobbin/crates/types/src/legacy.rs index 85933304f..d9edbd9be 100644 --- a/bobbin/crates/types/src/legacy.rs +++ b/bobbin/crates/types/src/legacy.rs @@ -226,21 +226,6 @@ pub struct LegacyRepo { pub extra_data: Option>>, } -#[derive(Debug, Deserialize)] -#[serde( - rename_all = "camelCase", - rename = "sh.tangled.knot.member", - tag = "$type", - bound(deserialize = "S: Deserialize<'de> + BosStr") -)] -pub struct LegacyKnotMember { - pub added_at: Datetime, - pub domain: S, - pub member: Did, - #[serde(flatten, default, skip_serializing_if = "Option::is_none")] - pub extra_data: Option>>, -} - #[derive(Debug)] pub enum LegacyRecord { Issue(LegacyIssue), @@ -251,7 +236,6 @@ pub enum LegacyRecord { Star(LegacyStar), PublicKey(LegacyPublicKey), Repo(LegacyRepo), - KnotMember(LegacyKnotMember), } impl LegacyRecord { @@ -272,7 +256,6 @@ impl LegacyRecord { "sh.tangled.feed.star" => Ok(Self::Star(serde_json::from_slice(bytes)?)), "sh.tangled.publicKey" => Ok(Self::PublicKey(serde_json::from_slice(bytes)?)), "sh.tangled.repo" => Ok(Self::Repo(serde_json::from_slice(bytes)?)), - "sh.tangled.knot.member" => Ok(Self::KnotMember(serde_json::from_slice(bytes)?)), other => Err(ExtractError::UnknownCollection(other.into())), } } diff --git a/bobbin/crates/types/src/search.rs b/bobbin/crates/types/src/search.rs index c5fc22056..40b132237 100644 --- a/bobbin/crates/types/src/search.rs +++ b/bobbin/crates/types/src/search.rs @@ -69,7 +69,6 @@ impl SearchableRecord { | Record::Follow(_) | Record::Vouch(_) | Record::Knot(_) - | Record::KnotMember(_) | Record::LabelOp(_) | Record::Pipeline(_) | Record::PipelineStatus(_) diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 32f8150e5..f47b3f175 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -42,7 +42,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, owner_did_from_aturi}; -use bobbin_types::knot_acl::{KnotOwnedSource, decode_knot_owned_source, knot_did_host}; +use bobbin_types::knot_acl::{KnotOwnedSource, decode_knot_owned_source}; use bobbin_types::org_tangled::feed::subscription::SubscriptionRecord; use bobbin_types::record::RecordBody; use bobbin_types::search::SearchableRecord; @@ -54,9 +54,6 @@ use bobbin_types::sh_tangled::feed::reaction::{Reaction, ReactionRecord}; use bobbin_types::sh_tangled::feed::star::{Star, StarRecord}; use bobbin_types::sh_tangled::graph::follow::{Follow, FollowRecord}; use bobbin_types::sh_tangled::graph::vouch::{Vouch, VouchRecord}; -use bobbin_types::sh_tangled::knot::member::{ - Member as KnotMember, MemberRecord as KnotMemberRecord, -}; use bobbin_types::sh_tangled::knot::{Knot, KnotRecord}; use bobbin_types::sh_tangled::label::definition::{ Definition as LabelDefinition, DefinitionRecord as LabelDefinitionRecord, @@ -437,22 +434,10 @@ pub fn router(state: AppState) -> Router { "/xrpc/sh.tangled.graph.countFollowsBy", get(count_follows_by), ) - .route( - "/xrpc/sh.tangled.knot.listMembersBy", - get(list_knot_members_by), - ) - .route( - "/xrpc/sh.tangled.knot.listMemberInvitesBy", - get(list_knot_member_invites_by), - ) .route( "/xrpc/sh.tangled.repo.listCollaboratorInvitesBy", get(list_collaborator_invites_by), ) - .route( - "/xrpc/sh.tangled.knot.countMembersBy", - get(count_knot_members_by), - ) .route("/xrpc/sh.tangled.label.listOpsBy", get(list_label_ops_by)) .route("/xrpc/sh.tangled.label.countOpsBy", get(count_label_ops_by)) .route( @@ -556,11 +541,6 @@ pub fn router(state: AppState) -> Router { ) .route("/xrpc/sh.tangled.repo.listArtifacts", get(list_artifacts)) .route("/xrpc/sh.tangled.repo.countArtifacts", get(count_artifacts)) - .route("/xrpc/sh.tangled.knot.listMembers", get(list_knot_members)) - .route( - "/xrpc/sh.tangled.knot.countMembers", - get(count_knot_members), - ) .route( "/xrpc/sh.tangled.spindle.listMembers", get(list_spindle_members), @@ -1289,7 +1269,6 @@ edge_kinds! { "sh.tangled.graph.follow" => FollowRecord, SubjectShape::BareDid, mirror FollowBy; "sh.tangled.graph.vouch" => VouchRecord, SubjectShape::BareDid, mirror VouchBy; "sh.tangled.knot" => KnotRecord, SubjectShape::BareDid; - "sh.tangled.knot.member" => KnotMemberRecord, SubjectShape::BareDid, mirror KnotMemberBy; "sh.tangled.label.definition" => LabelDefinitionRecord, SubjectShape::BareDid; "sh.tangled.label.op" => LabelOpRecord, SubjectShape::OneOfCollections(&["sh.tangled.repo.issue", "sh.tangled.repo.pull"]), mirror LabelOpBy; "sh.tangled.pipeline" => PipelineRecord, SubjectShape::BareDid, mirror PipelineBy; @@ -1821,11 +1800,6 @@ where fn synth_knot_owned_value(source: KnotOwnedSource, sort_micros: u64) -> Option { let created_at = datetime_of(UnixMicros::new(sort_micros))?; match source { - KnotOwnedSource::Member { knot, subject } => Some(serde_json::json!({ - "domain": knot_did_host(&knot)?, - "subject": subject.as_ref(), - "createdAt": created_at, - })), KnotOwnedSource::Collaborator { repo, subject } => Some(serde_json::json!({ "repo": repo.as_ref(), "subject": subject.as_ref(), @@ -3317,18 +3291,6 @@ async fn count_follows_by( count_mirror::(&state, q).map(Json) } -async fn list_knot_members_by( - State(state): State, - XrpcQuery(q): XrpcQuery>, -) -> Result { - list_mirror::, _>(&state, q).await -} -async fn count_knot_members_by( - State(state): State, - XrpcQuery(q): XrpcQuery, -) -> Result, XrpcError> { - count_mirror::(&state, q).map(Json) -} #[derive(serde::Deserialize)] struct InvitesByQuery { @@ -3340,8 +3302,7 @@ struct InvitesByQuery { struct InviteItem { uri: AtUri, knot: Did, - #[serde(skip_serializing_if = "Option::is_none")] - repo: Option>, + repo: Did, added_by: Did, created_at: Datetime, } @@ -3353,44 +3314,18 @@ struct InvitesByOutput { truncated: bool, } -#[derive(Clone, Copy)] -enum InviteFacet { - Knot, - Repo, -} - -impl InviteFacet { - fn collection(self) -> &'static str { - match self { - Self::Knot => "sh.tangled.knot.memberInvite", - Self::Repo => "sh.tangled.repo.collaboratorInvite", - } - } - - fn accepts(self, scope: &InviteScope) -> bool { - matches!( - (self, scope), - (Self::Knot, InviteScope::Knot(_)) | (Self::Repo, InviteScope::Repo { .. }) - ) - } -} +const COLLABORATOR_INVITE_COLLECTION: &str = "sh.tangled.repo.collaboratorInvite"; fn list_invites_by( state: &AppState, q: InvitesByQuery, - facet: InviteFacet, ) -> Result, XrpcError> { let subject = q.subject; if !state.invites.discovered() { return Err(XrpcError::Warming); } - let offered: Vec<(InviteScope, Offer)> = state - .invites - .outstanding(&subject) - .into_iter() - .filter(|(scope, _)| facet.accepts(scope)) - .collect(); + let offered: Vec<(InviteScope, Offer)> = state.invites.outstanding(&subject); let truncated = offered.len() > MAX_OUTSTANDING_OFFERS; let items = offered @@ -3406,8 +3341,8 @@ fn list_invites_by( })?; let uri = AtUri::::new_owned(format!( "at://{}/{}/{}", - scope.subject().as_ref(), - facet.collection(), + scope.repo.as_ref(), + COLLABORATOR_INVITE_COLLECTION, subject.as_ref() )) .map_err(|e| { @@ -3418,11 +3353,8 @@ fn list_invites_by( })?; Ok(InviteItem { uri, - knot: scope.knot().clone(), - repo: match &scope { - InviteScope::Knot(_) => None, - InviteScope::Repo { repo, .. } => Some(repo.clone()), - }, + knot: scope.knot, + repo: scope.repo, added_by: offer.invited_by, created_at, }) @@ -3436,18 +3368,11 @@ fn list_invites_by( })) } -async fn list_knot_member_invites_by( - State(state): State, - XrpcQuery(q): XrpcQuery, -) -> Result, XrpcError> { - list_invites_by(&state, q, InviteFacet::Knot) -} - async fn list_collaborator_invites_by( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result, XrpcError> { - list_invites_by(&state, q, InviteFacet::Repo) + list_invites_by(&state, q) } async fn list_label_ops_by( @@ -3690,20 +3615,6 @@ async fn count_artifacts( count_for::(&state, q).map(Json) } -async fn list_knot_members( - State(state): State, - XrpcQuery(q): XrpcQuery>, -) -> Result { - list_records::, _>(&state, q).await -} - -async fn count_knot_members( - State(state): State, - XrpcQuery(q): XrpcQuery, -) -> Result, XrpcError> { - count_for::(&state, q).map(Json) -} - async fn list_spindle_members( State(state): State, XrpcQuery(q): XrpcQuery>, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index d46b72323..0b6d21459 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -3167,81 +3167,6 @@ async fn list_issues_by_state_filter_narrows_results() { assert_eq!(items[0]["uri"], json!(closed_uri.as_ref())); } -#[tokio::test] -async fn knot_owned_member_is_synthesized_without_slingshot() { - let harness = Harness::new().await; - let knot = bobbin_types::knot_acl::host_to_knot_did("kt.oyster.cafe").unwrap(); - let subject = did("did:plc:boltless"); - let created = chrono::DateTime::parse_from_rfc3339("2026-06-01T00:00:00Z").unwrap(); - let micros = created.timestamp_micros() as u64; - let (source, edges) = bobbin_types::knot_acl::member_upsert(&knot, &subject, micros).unwrap(); - harness.edges.upsert_source(&source, edges); - harness.promote_ready(1, 1); - - let (status, body) = json_response( - router(harness.state.clone()) - .oneshot(list_request( - "sh.tangled.knot.listMembers", - subject.as_ref(), - &[], - )) - .await - .unwrap(), - ) - .await; - - assert_eq!(status, StatusCode::OK); - let items = body["items"].as_array().expect("items array"); - assert_eq!( - items.len(), - 1, - "synthesized member must hydrate with no slingshot mock mounted" - ); - assert_eq!(items[0]["uri"], json!(source.as_ref())); - assert!(items[0].get("cid").is_none()); - assert_eq!(items[0]["value"]["domain"], json!("kt.oyster.cafe")); - assert_eq!(items[0]["value"]["subject"], json!("did:plc:boltless")); - let got = chrono::DateTime::parse_from_rfc3339( - items[0]["value"]["createdAt"] - .as_str() - .expect("createdAt string"), - ) - .unwrap(); - assert_eq!(got.timestamp_micros(), micros as i64); -} - -#[tokio::test] -async fn knot_owned_member_lists_by_knot_did() { - let harness = Harness::new().await; - let knot = bobbin_types::knot_acl::host_to_knot_did("kt.oyster.cafe").unwrap(); - let subject = did("did:plc:boltless"); - let created = chrono::DateTime::parse_from_rfc3339("2026-06-01T00:00:00Z").unwrap(); - let micros = created.timestamp_micros() as u64; - let (source, edges) = bobbin_types::knot_acl::member_upsert(&knot, &subject, micros).unwrap(); - harness.edges.upsert_source(&source, edges); - harness.promote_ready(1, 1); - - let (status, body) = json_response( - router(harness.state.clone()) - .oneshot(list_request( - "sh.tangled.knot.listMembersBy", - knot.as_ref(), - &[], - )) - .await - .unwrap(), - ) - .await; - - assert_eq!(status, StatusCode::OK); - let items = body["items"].as_array().expect("items array"); - assert_eq!(items.len(), 1); - assert_eq!(items[0]["uri"], json!(source.as_ref())); - assert!(items[0].get("cid").is_none()); - assert_eq!(items[0]["value"]["domain"], json!("kt.oyster.cafe")); - assert_eq!(items[0]["value"]["subject"], json!("did:plc:boltless")); -} - #[tokio::test] async fn knot_owned_collaborator_is_synthesized_without_slingshot() { let harness = Harness::new().await; @@ -3392,20 +3317,16 @@ async fn list_recipients_empty_for_unknown_subject() { } const KNOT: &str = "did:web:knot.oyster.cafe"; -const MEMBER_INVITES_BY: &str = "sh.tangled.knot.listMemberInvitesBy"; const COLLABORATOR_INVITES_BY: &str = "sh.tangled.repo.listCollaboratorInvitesBy"; const JUNE_FIRST: u64 = 1_780_272_000_000_000; const JUNE_FOURTH: u64 = 1_780_531_200_000_000; -fn offer(h: &Harness, repo: Option<&str>, invitee: &str, added_by: &str, micros: u64) { +fn offer(h: &Harness, repo: &str, invitee: &str, added_by: &str, micros: u64) { h.invites.offer( did(invitee), - match repo { - None => InviteScope::Knot(did(KNOT)), - Some(repo) => InviteScope::Repo { - knot: did(KNOT), - repo: did(repo), - }, + InviteScope { + knot: did(KNOT), + repo: did(repo), }, Offer { invited_by: did(added_by), @@ -3431,34 +3352,17 @@ async fn invites_body(h: &Harness, endpoint: &str, subject: &str) -> (StatusCode } #[tokio::test] -async fn each_listing_serves_invitee_one_facet_of_their_offers() { +async fn the_listing_serves_an_invitee_their_standing_offers() { let h = Harness::new().await; h.invites.all_knots_discovered(); - offer(&h, None, "did:plc:limpet", "did:plc:akshay", JUNE_FIRST); offer( &h, - Some("did:plc:scallop"), + "did:plc:scallop", "did:plc:limpet", "did:plc:boltless", JUNE_FOURTH, ); - let (status, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; - assert_eq!(status, StatusCode::OK); - assert_eq!( - body, - json!({ - "items": [{ - "uri": "at://did:web:knot.oyster.cafe/sh.tangled.knot.memberInvite/did:plc:limpet", - "knot": KNOT, - "addedBy": "did:plc:akshay", - "createdAt": "2026-06-01T00:00:00.000000Z", - }], - "pending": [], - "truncated": false, - }) - ); - let (status, body) = invites_body(&h, COLLABORATOR_INVITES_BY, "did:plc:limpet").await; assert_eq!(status, StatusCode::OK); assert_eq!( @@ -3478,7 +3382,7 @@ async fn each_listing_serves_invitee_one_facet_of_their_offers() { async fn index_still_finding_its_knots_answers_warming() { let h = Harness::new().await; - let (status, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; + let (status, body) = invites_body(&h, COLLABORATOR_INVITES_BY, "did:plc:limpet").await; assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); assert_eq!( @@ -3492,9 +3396,15 @@ async fn unread_knot_goes_into_pending_and_answer_still_serves() { let h = Harness::new().await; h.invites.all_knots_discovered(); h.invites.awaiting(did(KNOT)); - offer(&h, None, "did:plc:limpet", "did:plc:akshay", JUNE_FIRST); + offer( + &h, + "did:plc:scallop", + "did:plc:limpet", + "did:plc:akshay", + JUNE_FIRST, + ); - let (status, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; + let (status, body) = invites_body(&h, COLLABORATOR_INVITES_BY, "did:plc:limpet").await; assert_eq!(status, StatusCode::OK); assert_eq!( body["pending"], @@ -3504,7 +3414,7 @@ async fn unread_knot_goes_into_pending_and_answer_still_serves() { assert_eq!(body["items"].as_array().unwrap().len(), 1); h.invites.backfilled(&did(KNOT)); - let (_, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; + let (_, body) = invites_body(&h, COLLABORATOR_INVITES_BY, "did:plc:limpet").await; assert_eq!(body["pending"], json!([])); } @@ -3515,7 +3425,7 @@ async fn flood_of_offers_is_limited_and_says_so() { (0..200u64).for_each(|n| { offer( &h, - Some(&format!("did:plc:repo{n:0>3}")), + &format!("did:plc:repo{n:0>3}"), "did:plc:limpet", "did:plc:boltless", JUNE_FOURTH + n, diff --git a/bobbin/crates/xrpc/tests/extended.rs b/bobbin/crates/xrpc/tests/extended.rs index 49536fd00..694fa6b33 100644 --- a/bobbin/crates/xrpc/tests/extended.rs +++ b/bobbin/crates/xrpc/tests/extended.rs @@ -253,14 +253,6 @@ fn artifact_body(repo_did: &Did, name: &str) -> Value { }) } -fn knot_member_body(subject_did: &Did) -> Value { - json!({ - "$type": "sh.tangled.knot.member", - "createdAt": "2026-05-01T00:00:00Z", - "subject": subject_did.as_ref(), - "domain": "oyster.cafe" - }) -} fn spindle_member_body(subject_did: &Did) -> Value { json!({ @@ -696,48 +688,6 @@ async fn list_artifacts_keys_on_repo_did() { assert_eq!(items[0]["value"]["repoDid"], json!(repo_did.as_ref())); } -#[tokio::test] -async fn list_knot_members_keys_on_subject_did() { - let h = Harness::new().await; - let subject_did = did("did:plc:nel"); - let subject = at(&format!("at://{}", subject_did.as_ref())); - let admin = did("did:plc:teq"); - let rk = rkey("m1"); - h.add_edge( - &nsid("sh.tangled.knot.member"), - &subject, - &at(&format!( - "at://{}/sh.tangled.knot.member/{}", - admin.as_ref(), - rk.as_ref() - )), - ); - h.mount( - &admin, - &nsid("sh.tangled.knot.member"), - &rk, - knot_member_body(&subject_did), - ) - .await; - - let app = router(h.state.clone()); - let (status, body) = json_response( - app.oneshot(list_request( - "sh.tangled.knot.listMembers", - subject.as_ref(), - &[], - )) - .await - .unwrap(), - ) - .await; - assert_eq!(status, StatusCode::OK); - let items = body["items"].as_array().unwrap(); - assert_eq!(items.len(), 1); - assert_eq!(items[0]["value"]["subject"], json!(subject_did.as_ref())); - assert_eq!(items[0]["value"]["domain"], json!("oyster.cafe")); -} - #[tokio::test] async fn list_spindle_members_keys_on_subject_did() { let h = Harness::new().await; diff --git a/web/src/lib/api/count.ts b/web/src/lib/api/count.ts index 622262b76..3f8355df7 100644 --- a/web/src/lib/api/count.ts +++ b/web/src/lib/api/count.ts @@ -14,8 +14,6 @@ export type CountName = | "sh.tangled.graph.countVouches" | "sh.tangled.graph.countVouchesBy" | "sh.tangled.knot.countKnots" - | "sh.tangled.knot.countMembers" - | "sh.tangled.knot.countMembersBy" | "sh.tangled.label.countDefinitions" | "sh.tangled.label.countOps" | "sh.tangled.label.countOpsBy" diff --git a/web/src/lib/api/repoCreationTargets.test.ts b/web/src/lib/api/repoCreationTargets.test.ts index 8cff13a3d..ad547cdd8 100644 --- a/web/src/lib/api/repoCreationTargets.test.ts +++ b/web/src/lib/api/repoCreationTargets.test.ts @@ -57,36 +57,20 @@ describe("availableSpindles", () => { }); describe("availableKnots", () => { - it("merges owned knots and memberships", async () => { - const pages: Record = { - "/xrpc/sh.tangled.knot.listKnots": { - items: [{ uri: "at://did:plc:alice/sh.tangled.knot/owned.example", value: {} }] - }, - "/xrpc/sh.tangled.knot.listMembers": { - items: [ - { - uri: "at://did:plc:owner/sh.tangled.knot.member/one", - value: { domain: "member.example" } - }, - { - uri: "at://did:web:member.example/sh.tangled.bobbin.knotMember/did:plc:alice", - value: { - subject: "did:plc:alice", - addedBy: "did:plc:owner", - createdAt: "2026-06-01T00:00:00.000Z" - } - }, - { - uri: "at://did:plc:owner/sh.tangled.knot.member/two", - value: { domain: "owned.example" } - }, - { uri: "at://did:plc:owner/sh.tangled.knot.member/bad", value: {} } - ] - } - }; + it("offers the shared knot alongside the knots this account owns", async () => { const fetch = vi.fn(async (input) => { const asked = new URL(String(input)); - const page = pages[asked.pathname]; + const page = + asked.pathname === "/xrpc/sh.tangled.knot.listKnots" + ? { + items: [ + { + uri: "at://did:plc:alice/sh.tangled.knot/owned.example", + value: {} + } + ] + } + : null; return new Response(JSON.stringify(page ?? { error: "MethodNotImplemented" }), { status: page ? 200 : 501, headers: { "content-type": "application/json" } @@ -95,15 +79,24 @@ describe("availableKnots", () => { const ctx = createBobbinClient({ serviceUrl: "https://bobbin.oyster.cafe", fetch }); await expect(availableKnots(ctx, did)).resolves.toEqual([ - "member.example", + "knot2.tngl.boltless.dev", "owned.example" ]); expect(fetch.mock.calls.map(([input]) => new URL(String(input)).pathname)).toEqual([ - "/xrpc/sh.tangled.knot.listKnots", - "/xrpc/sh.tangled.knot.listMembers" + "/xrpc/sh.tangled.knot.listKnots" ]); - expect(String(fetch.mock.calls[1][0])).toBe( - "https://bobbin.oyster.cafe/xrpc/sh.tangled.knot.listMembers?subject=did%3Aplc%3Aalice&limit=1000" + }); + + it("still offers the shared knot when this account owns none", async () => { + const fetch = vi.fn( + async () => + new Response(JSON.stringify({ items: [] }), { + status: 200, + headers: { "content-type": "application/json" } + }) ); + const ctx = createBobbinClient({ serviceUrl: "https://bobbin.oyster.cafe", fetch }); + + await expect(availableKnots(ctx, did)).resolves.toEqual(["knot2.tngl.boltless.dev"]); }); }); diff --git a/web/src/lib/api/repoCreationTargets.ts b/web/src/lib/api/repoCreationTargets.ts index 0e5207d02..6bf3b1bc9 100644 --- a/web/src/lib/api/repoCreationTargets.ts +++ b/web/src/lib/api/repoCreationTargets.ts @@ -1,12 +1,10 @@ import { ok } from "@atcute/client"; import type { Did } from "@atcute/lexicons/syntax"; import type { BobbinContext } from "$lib/api/client"; -import { jsonGet } from "$lib/api/_request"; import { mainSchema as listKnotsSchema } from "$lib/api/lexicons/types/sh/tangled/knot/listKnots"; import { mainSchema as listSpindleMembersSchema } from "$lib/api/lexicons/types/sh/tangled/spindle/listMembers"; import { mainSchema as listSpindlesSchema } from "$lib/api/lexicons/types/sh/tangled/spindle/listSpindles"; import { parseResourceUri } from "@atcute/lexicons/syntax"; -import { hostForServiceDid } from "$lib/auth/agent"; interface Page { items: T[]; @@ -29,20 +27,6 @@ const allPages = async (load: (cursor?: string) => Promise>): Promise return items; }; -const knotDomain = (value: unknown): string | null => { - if (typeof value !== "object" || value === null || !("domain" in value)) return null; - return typeof value.domain === "string" && value.domain ? value.domain : null; -}; - -const repoFromUri = (uri: unknown): string | null => { - if (typeof uri !== "string") return null; - try { - return parseResourceUri(uri).repo; - } catch { - return null; - } -}; - const recordName = (uri: unknown): string | null => { if (typeof uri !== "string") return null; try { @@ -52,30 +36,18 @@ const recordName = (uri: unknown): string | null => { } }; -const knotMembershipHost = (item: unknown): string | null => { - if (typeof item !== "object" || item === null) return null; - const record = item as { uri?: unknown; value?: unknown }; - const repo = repoFromUri(record.uri); - const host = repo ? hostForServiceDid(repo) : null; - if (host) return host; - if (record.value !== undefined) { - const domain = knotDomain(record.value); - if (domain) return domain; - } - return null; -}; - const spindleInstance = (value: unknown): string | null => { if (typeof value !== "object" || value === null || !("instance" in value)) return null; return typeof value.instance === "string" && value.instance ? value.instance : null; }; const targetNames = ( + seed: string[], owned: RecordItem[], - memberships: RecordItem[], - memberName: (item: RecordItem) => string | null + memberships: RecordItem[] = [], + memberName: (item: RecordItem) => string | null = () => null ): string[] => { - const names = new Set(); + const names = new Set(seed); for (const item of owned) { const name = recordName(item.uri); if (name) names.add(name); @@ -87,25 +59,20 @@ const targetNames = ( return [...names].sort((a, b) => a.localeCompare(b)); }; -export const availableKnots = async (ctx: BobbinContext, did: string): Promise => { - const [owned, memberships] = await Promise.all([ - allPages((cursor) => +// TODO(boltless): list bookmarked repos instead. +const SHARED_KNOTS = ["knot2.tngl.boltless.dev"]; + +export const availableKnots = async (ctx: BobbinContext, did: string): Promise => + targetNames( + SHARED_KNOTS, + await allPages((cursor) => ok( ctx.xrpc.call(listKnotsSchema, { params: { subject: did as Did, limit: 1_000, cursor } }) ) - ), - allPages((cursor) => - jsonGet>(ctx, "sh.tangled.knot.listMembers", { - subject: did, - limit: 1_000, - cursor - }) ) - ]); - return targetNames(owned, memberships, knotMembershipHost); -}; + ); export const availableSpindles = async (ctx: BobbinContext, did: string): Promise => { const [owned, memberships] = await Promise.all([ @@ -124,5 +91,5 @@ export const availableSpindles = async (ctx: BobbinContext, did: string): Promis ) ) ]); - return targetNames(owned, memberships, (item) => spindleInstance(item.value)); + return targetNames([], owned, memberships, (item) => spindleInstance(item.value)); }; -- 2.51.2