diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index 2194ad7b3..ebe99c940 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -6,7 +6,9 @@ use std::sync::Arc; use std::time::Duration; use anyhow::{Context, anyhow}; -use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor, Settlements, StateIndex}; +use bobbin_edge_index::{ + CoverageWatch, EdgeStore, HydrantCursor, InviteIndex, Settlements, StateIndex, +}; use bobbin_ingest::{ IngestConfig, IngestRuntime, RepoIdResolver, WarmingBuffer, run as run_ingest, }; @@ -218,6 +220,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { hasher.clone(), )); let edges = Arc::new(EdgeStore::new(hasher.clone())); + let invites = Arc::new(InviteIndex::new()); let issue_states = Arc::new(StateIndex::new(hasher.clone())); let pull_statuses = Arc::new(StateIndex::new(hasher.clone())); let coverage = Arc::new(CoverageWatch::new()); @@ -367,6 +370,8 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { gate: knot_gate, registry: knot_registry, store: edges.clone(), + invites: invites.clone(), + coverage: coverage.clone(), ws: knot_ws, clock: clock.clone(), dev: knot_acl_dev, @@ -429,7 +434,8 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { .with_mirror_v2(mirror_v2) .with_proxies(trusted_proxies) .with_service_did(cfg.service_auth.did.clone()) - .with_settlements(settlements); + .with_settlements(settlements) + .with_invites(invites); let app = router(state); let _debug_server = match (debug_bind, mem_probe) { diff --git a/bobbin/crates/knot-ingest/src/client.rs b/bobbin/crates/knot-ingest/src/client.rs index 2ca8b7fbe..eed6384ea 100644 --- a/bobbin/crates/knot-ingest/src/client.rs +++ b/bobbin/crates/knot-ingest/src/client.rs @@ -5,16 +5,18 @@ use std::sync::Arc; use std::time::Duration; use bobbin_knot_proxy::{KnotHost, KnotHostError, PrivateAddressFilter, PrivateHostReason}; -use bobbin_runtime::{HttpRequest, HttpResponseHead, HttpTransport, NetworkError, ReqwestHttp}; +use bobbin_runtime::{ + HttpRequest, HttpResponseHead, HttpTransport, NetworkError, ReqwestHttp, UnixMicros, +}; use bytes::{Bytes, BytesMut}; -use chrono::{DateTime, Utc}; use futures::TryStreamExt; use http::{HeaderMap, StatusCode}; use jacquard_common::DefaultStr; +use jacquard_common::types::datetime::Datetime; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use knot_capability::Capability; -use serde::Deserialize; +use serde::{Deserialize, Deserializer, de::DeserializeOwned}; use thiserror::Error; use url::Url; @@ -28,17 +30,36 @@ 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)] pub struct KnotClient { http: Arc, } +pub(crate) fn knot_stamp<'de, D: Deserializer<'de>>(de: D) -> Result { + let at = Datetime::deserialize(de)?; + u64::try_from(at.timestamp_micros()) + .map(UnixMicros::new) + .map_err(|_| serde::de::Error::custom("stamp before 1970 won't fit unix micros")) +} + #[derive(Clone, Debug, Eq, PartialEq, Deserialize)] #[serde(rename_all = "camelCase")] pub struct AclEntry { pub subject: Did, - pub created_at: DateTime, + #[serde(deserialize_with = "knot_stamp")] + pub created_at: UnixMicros, +} + +#[derive(Clone, Debug, Eq, PartialEq, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct InviteEntry { + pub subject: Did, + pub added_by: Did, + #[serde(deserialize_with = "knot_stamp")] + pub created_at: UnixMicros, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -47,12 +68,24 @@ pub enum Completeness { Truncated, } +impl Completeness { + pub fn and(self, other: Self) -> Self { + match (self, other) { + (Self::Complete, Self::Complete) => Self::Complete, + _ => Self::Truncated, + } + } +} + #[derive(Clone, Debug, Eq, PartialEq)] -pub struct AclListing { - pub entries: Vec, +pub struct Listing { + pub entries: Vec, pub completeness: Completeness, } +pub type AclListing = Listing; +pub type InviteListing = Listing; + #[derive(Debug, Error)] pub enum KnotClientError { #[error("knot host: {0}")] @@ -151,15 +184,7 @@ impl KnotClient { host: &KnotHost, knot: &Did, ) -> Result { - self.drain( - host, - LIST_MEMBERS_NSID, - knot.as_ref().to_owned(), - None, - 0, - Vec::new(), - ) - .await + self.list(host, LIST_MEMBERS_NSID, knot).await } pub async fn list_collaborators( @@ -167,10 +192,35 @@ impl KnotClient { host: &KnotHost, repo: &Did, ) -> Result { + 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, + repo: &Did, + ) -> Result { + self.list(host, LIST_COLLABORATOR_INVITES_NSID, repo).await + } + + async fn list( + &self, + host: &KnotHost, + endpoint: &'static str, + subject: &Did, + ) -> Result, KnotClientError> { self.drain( host, - LIST_COLLABORATORS_NSID, - repo.as_ref().to_owned(), + endpoint, + subject.as_ref().to_owned(), None, 0, Vec::new(), @@ -178,15 +228,15 @@ impl KnotClient { .await } - fn drain<'a>( + fn drain<'a, T: DeserializeOwned + Send + 'a>( &'a self, host: &'a KnotHost, endpoint: &'static str, subject: String, cursor: Option, page: usize, - mut acc: Vec, - ) -> Pin> + Send + 'a>> { + mut acc: Vec, + ) -> Pin, KnotClientError>> + Send + 'a>> { Box::pin(async move { if page >= MAX_LIST_PAGES { tracing::warn!( @@ -195,7 +245,7 @@ impl KnotClient { pages = page, "knot list truncated at page cap" ); - return Ok(AclListing { + return Ok(Listing { entries: acc, completeness: Completeness::Truncated, }); @@ -209,7 +259,7 @@ impl KnotClient { self.drain(host, endpoint, subject, Some(next), page + 1, acc) .await } - None => Ok(AclListing { + None => Ok(Listing { entries: acc, completeness: Completeness::Complete, }), @@ -217,13 +267,13 @@ impl KnotClient { }) } - async fn fetch_page( + async fn fetch_page( &self, host: &KnotHost, endpoint: &'static str, subject: &str, cursor: Option<&str>, - ) -> Result { + ) -> Result, KnotClientError> { let mut url = host.xrpc_url(&nsid(endpoint)); { let mut q = url.query_pairs_mut(); @@ -261,9 +311,10 @@ struct VersionWire { } #[derive(Deserialize)] -struct ListWire { - #[serde(default)] - items: Vec, +#[serde(bound = "T: DeserializeOwned")] +struct ListWire { + #[serde(default = "Vec::new")] + items: Vec, #[serde(default)] cursor: Option, } @@ -301,6 +352,10 @@ mod tests { Did::new_owned(s).unwrap() } + fn stamp(iso: &str) -> UnixMicros { + knot_stamp(serde_json::Value::String(iso.to_owned())).unwrap() + } + fn client() -> KnotClient { KnotClient::new(ReqwestHttp::shared(default_http_client(true).unwrap())) } @@ -372,11 +427,11 @@ mod tests { vec![ AclEntry { subject: did("did:plc:boltless"), - created_at: "2026-06-01T00:00:00Z".parse().unwrap(), + created_at: stamp("2026-06-01T00:00:00Z"), }, AclEntry { subject: did("did:plc:akshay"), - created_at: "2026-06-02T12:00:00Z".parse().unwrap(), + created_at: stamp("2026-06-02T12:00:00Z"), }, ] ); @@ -435,6 +490,77 @@ mod tests { assert_eq!(listing.entries[0].subject, did("did:plc:olaren")); } + #[tokio::test] + async fn an_invite_listing_decodes_both_facets_at_their_own_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")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "items": [{"subject": "did:plc:olaren", "addedBy": "did:plc:boltless", "createdAt": "2026-06-05T00:00:00Z"}] + }))) + .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.entries, + vec![InviteEntry { + subject: did("did:plc:olaren"), + added_by: did("did:plc:boltless"), + created_at: stamp("2026-06-05T00:00:00Z"), + }] + ); + } + + #[tokio::test] + async fn invite_listing_rejects_offer_missing_its_inviter() { + let server = MockServer::start().await; + let host = endpoint(&server); + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.listMemberInvites")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "items": [{"subject": "did:plc:limpet", "createdAt": "2026-06-04T00:00:00Z"}] + }))) + .mount(&server) + .await; + + let rejected = client() + .list_member_invites(&host, &did("did:web:knot.oyster.cafe")) + .await + .is_err(); + assert!( + rejected, + "reject here, since rows can't record who invited without this field" + ); + } + #[tokio::test] async fn drain_reports_truncation_at_page_cap() { let server = MockServer::start().await; diff --git a/bobbin/crates/knot-ingest/src/firehose.rs b/bobbin/crates/knot-ingest/src/firehose.rs index ba64ed78e..9c3f4263d 100644 --- a/bobbin/crates/knot-ingest/src/firehose.rs +++ b/bobbin/crates/knot-ingest/src/firehose.rs @@ -2,6 +2,8 @@ use std::collections::HashMap; use std::convert::Infallible; use std::io::Cursor as IoCursor; +use bobbin_edge_index::Offer; +use bobbin_runtime::UnixMicros; use bytes::Bytes; use jacquard_api::com_atproto::sync::subscribe_repos; use jacquard_common::DefaultStr; @@ -95,6 +97,49 @@ impl InviteEvent { } } +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum InviteTarget { + Knot, + Repo(Did), +} + +#[derive(Clone, Debug, PartialEq)] +pub struct InviteOp { + pub target: InviteTarget, + pub invitee: Did, + pub offer: Option, +} + +#[derive(Deserialize)] +struct InviteWire { + #[serde(rename = "createdAt", deserialize_with = "crate::client::knot_stamp")] + created_at: UnixMicros, + #[serde(rename = "x-tngl-editor")] + 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()), + }; + let invitee = Did::::new_owned(op.rkey.as_str()).ok()?; + let offer = if op.action.is_delete() { + None + } else { + let wire: InviteWire = serde_ipld_dagcbor::from_slice(op.record.as_deref()?).ok()?; + Some(Offer { + invited_by: wire.editor, + offered_at: wire.created_at, + }) + }; + Some(InviteOp { + target, + invitee, + offer, + }) +} + #[derive(Deserialize)] struct Header { op: Option, @@ -439,6 +484,38 @@ mod tests { bytes.extend(encode_one(&nest(MAX_CBOR_DEPTH - 4))); assert_eq!(decode_frame(&bytes).unwrap(), Frame::Other); } + + #[test] + fn malformed_invite_never_decodes_to_offer() { + let collection = InviteEvent::KnotMember.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 cases = [ + ( + "record key that isn't a did", + "3lkm2xqbolt2s", + full.clone(), + ), + ( + "body cut short of its last field", + "did:plc:limpet", + full[..full.len() - 4].to_vec(), + ), + ]; + + cases.into_iter().for_each(|(case, rkey, body)| { + let op = RecordOp { + action: OpAction::Create, + collection: Nsid::::new_owned(collection).unwrap(), + rkey: Rkey::::new_owned(rkey).unwrap(), + record: Some(Bytes::from(body)), + }; + assert!( + invite_op_of(&op, InviteEvent::KnotMember, &knot).is_none(), + "invite with {case} still decoded" + ); + }); + } } #[cfg(test)] @@ -602,4 +679,12 @@ pub(crate) mod testcbor { } Val::Map(pairs) } + + pub(crate) fn invite_record(collection: &str, created_at: &str, editor: &str) -> Vec { + encode_one(&Val::Map(vec![ + (t("$type"), t(collection)), + (t("createdAt"), t(created_at)), + (t("x-tngl-editor"), t(editor)), + ])) + } } diff --git a/bobbin/crates/knot-ingest/src/legacy.rs b/bobbin/crates/knot-ingest/src/legacy.rs index 836b41d6a..3afcd4efb 100644 --- a/bobbin/crates/knot-ingest/src/legacy.rs +++ b/bobbin/crates/knot-ingest/src/legacy.rs @@ -61,6 +61,7 @@ mod tests { use parking_lot::Mutex; use super::*; + use crate::firehose::InviteOp; use crate::stream::Feed; #[derive(Default)] @@ -74,6 +75,10 @@ mod tests { fn repo_roster_changed(&self, repo: &Did) { self.0.lock().push(repo.as_ref().to_owned()); } + + fn invite_op(&self, op: InviteOp) { + self.0.lock().push(op.invitee.as_ref().to_owned()); + } } const MEMBER: &str = r#"{"rkey":"3mug","nsid":"sh.tangled.knot.memberUpdate", diff --git a/bobbin/crates/knot-ingest/src/lib.rs b/bobbin/crates/knot-ingest/src/lib.rs index fda2e157f..bf864b3e0 100644 --- a/bobbin/crates/knot-ingest/src/lib.rs +++ b/bobbin/crates/knot-ingest/src/lib.rs @@ -11,5 +11,5 @@ pub use client::{AclEntry, AclListing, Completeness, KnotClient, KnotClientError pub use gate::CapabilityGate; pub use orchestrator::Orchestrator; pub use registry::KnotRegistry; -pub use roster::{Cursor, Roster}; +pub use roster::Roster; pub use stream::{Feed, StreamConfig, run_stream}; diff --git a/bobbin/crates/knot-ingest/src/orchestrator.rs b/bobbin/crates/knot-ingest/src/orchestrator.rs index bfe8ed7fb..20705d22b 100644 --- a/bobbin/crates/knot-ingest/src/orchestrator.rs +++ b/bobbin/crates/knot-ingest/src/orchestrator.rs @@ -3,11 +3,10 @@ use std::collections::{HashMap, HashSet}; use std::sync::Arc; use std::time::Duration; -use bobbin_edge_index::EdgeStore; +use bobbin_edge_index::{CoverageWatch, EdgeStore, InviteIndex}; use bobbin_knot_proxy::KnotHost; use bobbin_runtime::{Clock, WsTransport}; use bobbin_types::knot_acl::{KnotHostKey, host_to_knot_did}; -use chrono::{DateTime, Utc}; use futures::StreamExt; use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; @@ -15,10 +14,13 @@ use tokio_util::sync::CancellationToken; use tokio::sync::mpsc; -use crate::client::{AclListing, Completeness, KnotClient, KnotClientError, knot_endpoint}; +use crate::client::{ + AclEntry, Completeness, InviteEntry, KnotClient, KnotClientError, Listing, knot_endpoint, +}; +use crate::firehose::{InviteOp, InviteTarget}; use crate::gate::CapabilityGate; use crate::registry::KnotRegistry; -use crate::roster::{AclKind, Cursor, Roster}; +use crate::roster::{AclKind, Applied, Assertion, Roster, Scope, Source}; use crate::stream::{Feed, FeedHandler, StreamConfig, run_stream}; const POLL_INTERVAL: Duration = Duration::from_secs(30); @@ -32,6 +34,8 @@ pub struct Orchestrator { pub gate: Arc, pub registry: Arc, pub store: Arc, + pub invites: Arc, + pub coverage: Arc, pub ws: Arc, pub clock: Arc, pub dev: bool, @@ -47,7 +51,11 @@ impl Orchestrator { if self.cancel.is_cancelled() { break; } + let discovered = self.coverage.snapshot().is_ready(); self.discover(&mut subscribed, &mut unspawnable).await; + if discovered { + self.invites.all_knots_discovered(); + } tokio::select! { _ = self.cancel.cancelled() => break, _ = self.clock.sleep(POLL_INTERVAL) => {} @@ -66,6 +74,10 @@ impl Orchestrator { .into_iter() .filter(|host| !subscribed.contains_key(host) && !unspawnable.contains(host)) .collect(); + candidates + .iter() + .filter_map(|host| host_to_knot_did(host.as_str())) + .for_each(|knot| self.invites.awaiting(knot)); let approved: Vec<(KnotHostKey, Feed)> = futures::stream::iter(candidates) .map(|host| async move { self.gate.admit(&host).await.map(|feed| (host, feed)) }) .buffer_unordered(PROBE_CONCURRENCY) @@ -90,6 +102,7 @@ impl Orchestrator { let knot_did = host_to_knot_did(host.as_str())?; let roster = Arc::new(Mutex::new(Roster::new( self.store.clone(), + self.invites.clone(), knot_did.clone(), self.registry.clone(), host.clone(), @@ -105,6 +118,7 @@ impl Orchestrator { let handlers = KnotHandlers { host: host.clone(), acl: nudge_tx, + roster: roster.clone(), }; let stream_endpoint = endpoint.clone(); let stream_cursor = crate::stream::Cursor::start(feed); @@ -114,6 +128,7 @@ impl Orchestrator { let client = self.client.clone(); let registry = self.registry.clone(); + let invites = self.invites.clone(); let clock = self.clock.clone(); let host_owned = host.clone(); let reconcile_token = token.clone(); @@ -124,6 +139,7 @@ impl Orchestrator { knot: &knot_did, host: &host_owned, registry: ®istry, + invites: &invites, roster: &roster, }; reconcile_loop(&ctx, &*clock, &reconcile_token, nudge_rx).await; @@ -142,6 +158,7 @@ enum Nudge { struct KnotHandlers { host: KnotHostKey, acl: mpsc::Sender, + roster: Arc>, } impl KnotHandlers { @@ -167,6 +184,31 @@ impl FeedHandler for KnotHandlers { fn repo_roster_changed(&self, repo: &Did) { self.nudge(Nudge::Repo(repo.clone())); } + + fn invite_op(&self, op: InviteOp) { + let InviteOp { + target, + 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, + invited_by: offer.invited_by, + }, + invitee, + offer.offered_at, + Source::Firehose, + ), + None => roster.retire(AclKind::invite(scope), invitee), + } + } } struct ReconcileCtx<'a> { @@ -175,6 +217,7 @@ struct ReconcileCtx<'a> { knot: &'a Did, host: &'a KnotHostKey, registry: &'a KnotRegistry, + invites: &'a InviteIndex, roster: &'a Mutex, } @@ -250,77 +293,143 @@ async fn apply_batch(ctx: &ReconcileCtx<'_>, batch: NudgeBatch) { reconcile_once(ctx).await; return; } - let horizon = ctx.roster.lock().max_cursor(); - for repo in batch.repos { - reconcile_repo_once(ctx, &repo, horizon).await; - } + let horizon = ctx.roster.lock().applied(); + futures::stream::iter(batch.repos) + .for_each(|repo| async move { + reconcile_repo_once(ctx, &repo, horizon).await; + }) + .await; } + async fn reconcile_once(ctx: &ReconcileCtx<'_>) { - let horizon = ctx.roster.lock().max_cursor(); + let horizon = ctx.roster.lock().applied(); let members = ctx.client.list_members(ctx.endpoint, ctx.knot).await; - reconcile_listing(ctx, AclKind::Member, members, horizon); + 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); - futures::stream::iter(ctx.registry.repos(ctx.host)) - .for_each(|repo| async move { reconcile_repo_once(ctx, &repo, horizon).await }) + let repo_offers = futures::stream::iter(ctx.registry.repos(ctx.host)) + .then(|repo| async move { reconcile_repo_once(ctx, &repo, horizon).await }) + .fold( + Completeness::Complete, + |all, one| async move { all.and(one) }, + ) .await; + if offers.and(repo_offers) == Completeness::Complete { + ctx.invites.backfilled(ctx.knot); + } + ctx.roster.lock().purge_legacy(); } -async fn reconcile_repo_once(ctx: &ReconcileCtx<'_>, repo: &Did, horizon: Cursor) { +async fn reconcile_repo_once( + ctx: &ReconcileCtx<'_>, + repo: &Did, + horizon: Applied, +) -> Completeness { let collaborators = ctx.client.list_collaborators(ctx.endpoint, repo).await; - reconcile_listing(ctx, AclKind::Collaborator(repo), collaborators, horizon); + reconcile_listing(ctx, Scope::Repo(repo), collaborators, horizon); + + let invited = ctx + .client + .list_collaborator_invites(ctx.endpoint, repo) + .await; + reconcile_listing(ctx, Scope::Repo(repo), invited, horizon) +} + +trait Listed { + fn kind(scope: Scope<'_>) -> AclKind<'_>; + fn subject(&self) -> &Did; + fn apply_to(self, roster: &mut Roster, scope: Scope<'_>, source: Source); +} + +impl Listed for AclEntry { + fn kind(scope: Scope<'_>) -> AclKind<'_> { + AclKind::acl(scope) + } + + 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); + } +} + +impl Listed for InviteEntry { + fn kind(scope: Scope<'_>) -> AclKind<'_> { + AclKind::invite(scope) + } + + fn subject(&self) -> &Did { + &self.subject + } + + fn apply_to(self, roster: &mut Roster, scope: Scope<'_>, source: Source) { + roster.apply( + Assertion::Invite { + scope, + invited_by: self.added_by, + }, + self.subject, + self.created_at, + source, + ); + } } -fn reconcile_listing( +fn reconcile_listing( ctx: &ReconcileCtx<'_>, - kind: AclKind<'_>, - listing: Result, - horizon: Cursor, -) { - match listing { - Ok(AclListing { - entries, - completeness, - }) => { - let present: HashSet> = - entries.iter().map(|entry| entry.subject.clone()).collect(); - let mut guard = ctx.roster.lock(); - entries.into_iter().for_each(|entry| match kind { - AclKind::Member => { - guard.apply_member(entry.subject, Cursor(nanos(entry.created_at))) - } - AclKind::Collaborator(repo) => guard.apply_collaborator( - repo.clone(), - entry.subject, - Cursor(nanos(entry.created_at)), - ), - }); - if completeness == Completeness::Complete { - guard.reap(kind, &present, horizon); - } + scope: Scope<'_>, + listing: Result, KnotClientError>, + horizon: Applied, +) -> Completeness { + let kind = T::kind(scope); + let Listing { + entries, + completeness, + } = match listing { + Ok(listing) => listing, + Err(KnotClientError::NotFound) if kind.is_invite() => return Completeness::Complete, + Err(err) => { + warn_unreadable(ctx, kind, &err); + return Completeness::Truncated; } - Err(err) => match kind { - AclKind::Member => { - tracing::warn!(host = %ctx.host, error = %err, "knot member reconcile failed") - } - AclKind::Collaborator(repo) => { - tracing::warn!(host = %ctx.host, repo = %repo.as_ref(), error = %err, "knot collaborator reconcile failed") - } - }, + }; + let present: HashSet> = entries + .iter() + .map(|entry| entry.subject().clone()) + .collect(); + let mut guard = ctx.roster.lock(); + entries + .into_iter() + .for_each(|entry| entry.apply_to(&mut guard, scope, Source::Listing(horizon))); + if completeness == Completeness::Complete { + guard.reap(kind, &present, horizon); } + completeness } -fn nanos(timestamp: DateTime) -> i64 { - timestamp.timestamp_nanos_opt().unwrap_or(0) +fn warn_unreadable(ctx: &ReconcileCtx<'_>, kind: AclKind<'_>, err: &KnotClientError) { + tracing::warn!( + host = %ctx.host, + repo = ?kind.scope.repo().map(|repo| repo.as_ref()), + listing = kind.listing(), + error = %err, + "knot reconcile failed" + ); } #[cfg(test)] mod tests { use super::*; use crate::client::authority; - use bobbin_runtime::{ReqwestHttp, RuntimeHasher, SystemClock}; + use bobbin_edge_index::Offer; + use bobbin_runtime::{ReqwestHttp, RuntimeHasher, SystemClock, UnixMicros}; use bobbin_types::edges::Edge; use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; use jacquard_common::types::string::AtUri; @@ -353,6 +462,7 @@ mod tests { host: KnotHostKey, registry: Arc, store: Arc, + invites: Arc, roster: Arc>, client: KnotClient, } @@ -365,6 +475,7 @@ mod tests { knot: &self.knot, host: &self.host, registry: &self.registry, + invites: &self.invites, roster: &self.roster, } } @@ -376,9 +487,11 @@ mod tests { let host = KnotHostKey::new(&authority(&endpoint)); let registry = Arc::new(KnotRegistry::new()); let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); + let invites = Arc::new(InviteIndex::new()); let knot_did = host_to_knot_did(host.as_str()).unwrap(); let roster = Arc::new(Mutex::new(Roster::new( store.clone(), + invites.clone(), knot_did.clone(), registry.clone(), host.clone(), @@ -391,11 +504,35 @@ mod tests { host, registry, store, + invites, roster, client, } } + const MEMBER_INVITES: &str = "sh.tangled.knot.listMemberInvites"; + + async fn serve(h: &Harness, endpoint: &str, response: ResponseTemplate) { + Mock::given(method("GET")) + .and(path(format!("/xrpc/{endpoint}"))) + .respond_with(response) + .mount(&h.server) + .await; + } + + fn offers(invitees: &[&str]) -> serde_json::Value { + json!({ + "items": invitees + .iter() + .map(|invitee| json!({ + "subject": invitee, + "addedBy": "did:plc:akshay", + "createdAt": "2026-06-04T00:00:00Z" + })) + .collect::>() + }) + } + #[tokio::test] async fn reconcile_backfills_members_and_collaborators() { let h = harness().await; @@ -430,6 +567,75 @@ mod tests { assert_eq!(collaborator_count(&h.store, "did:plc:scallop"), 1); } + #[tokio::test] + async fn reconcile_backfills_outstanding_member_invites() { + let h = harness().await; + serve( + &h, + MEMBER_INVITES, + ResponseTemplate::new(200).set_body_json(offers(&["did:plc:limpet", "did:plc:olaren"])), + ) + .await; + + reconcile_once(&h.ctx()).await; + + assert_eq!(h.invites.outstanding(&did("did:plc:olaren")).len(), 1); + let offer = h.invites.outstanding(&did("did:plc:limpet")).remove(0).1; + assert_eq!(offer.invited_by, did("did:plc:akshay")); + assert_eq!(offer.offered_at, UnixMicros::new(1_780_531_200_000_000)); + } + + #[tokio::test] + async fn completed_backfill_stops_knot_being_pending() { + let h = harness().await; + h.invites.awaiting(h.knot.clone()); + h.invites.all_knots_discovered(); + assert_eq!(h.invites.pending(), vec![h.knot.clone()]); + + serve( + &h, + MEMBER_INVITES, + ResponseTemplate::new(200).set_body_json(offers(&[])), + ) + .await; + reconcile_once(&h.ctx()).await; + + assert!(h.invites.pending().is_empty()); + } + + #[tokio::test] + async fn knot_without_invite_listings_stops_being_pending() { + let h = harness().await; + h.invites.awaiting(h.knot.clone()); + h.invites.all_knots_discovered(); + + serve(&h, MEMBER_INVITES, ResponseTemplate::new(404)).await; + reconcile_once(&h.ctx()).await; + + assert!( + h.invites.pending().is_empty(), + "knots predating invite listings won't list offers or stream them, so bobbin has read all it can" + ); + } + + #[tokio::test] + async fn failed_backfill_leaves_only_its_own_knot_pending() { + let h = harness().await; + 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; + reconcile_once(&h.ctx()).await; + + assert_eq!( + h.invites.pending(), + vec![h.knot.clone()], + "bobbin puts unreadable knots into pending, and answers anyway" + ); + } + #[tokio::test] async fn reconcile_reaps_departed_member() { let h = harness().await; @@ -482,7 +688,16 @@ mod tests { .await; let stayed = did("did:plc:akshay"); - h.roster.lock().apply_member(stayed, Cursor(1_000_000)); + { + let mut roster = h.roster.lock(); + let taken = roster.applied(); + roster.apply( + Assertion::Acl(Scope::Knot), + stayed, + UnixMicros::new(1_000_000), + Source::Listing(taken), + ); + } assert_eq!(member_count(&h.store, "did:plc:akshay"), 1); reconcile_once(&h.ctx()).await; @@ -565,13 +780,19 @@ mod tests { ); } - #[test] - fn handlers_nudge_full_for_roster_changes() { - let (tx, mut rx) = mpsc::channel(4); - let handlers = KnotHandlers { + fn test_handlers(roster: Arc>, acl: mpsc::Sender) -> KnotHandlers { + KnotHandlers { host: KnotHostKey::new("oyster.cafe"), - acl: tx, - }; + acl, + roster, + } + } + + #[tokio::test] + async fn handlers_nudge_full_for_roster_changes() { + let h = harness().await; + let (tx, mut rx) = mpsc::channel(4); + let handlers = test_handlers(h.roster.clone(), tx); handlers.roster_changed(); handlers.roster_changed(); assert!(matches!(rx.try_recv(), Ok(Nudge::Full))); @@ -579,17 +800,43 @@ mod tests { assert!(rx.try_recv().is_err()); } - #[test] - fn handlers_nudge_repo_for_repo_roster_changes() { + #[tokio::test] + async fn handlers_nudge_repo_for_repo_roster_changes() { + let h = harness().await; let (tx, mut rx) = mpsc::channel(4); - let handlers = KnotHandlers { - host: KnotHostKey::new("oyster.cafe"), - acl: tx, - }; + let handlers = test_handlers(h.roster.clone(), tx); handlers.repo_roster_changed(&did("did:plc:scallop")); assert_eq!(rx.try_recv(), Ok(Nudge::Repo(did("did:plc:scallop"))),); } + #[tokio::test] + async fn handlers_apply_firehose_invite_straight_to_roster() { + let h = harness().await; + 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, + invitee: invitee.clone(), + offer: Some(Offer { + invited_by: did("did:plc:akshay"), + offered_at: UnixMicros::new(1_780_272_000_000_000), + }), + }); + assert_eq!(h.invites.outstanding(&invitee).len(), 1); + + handlers.invite_op(InviteOp { + target: InviteTarget::Knot, + invitee: invitee.clone(), + offer: None, + }); + assert!( + h.invites.outstanding(&invitee).is_empty(), + "firehose delete withdraws offer it created" + ); + } + #[tokio::test(start_paused = true)] async fn collaborator_nudge_reconciles_only_that_repo() { let Harness { @@ -599,6 +846,7 @@ mod tests { host, registry, store, + invites, roster, client, } = harness().await; @@ -623,6 +871,7 @@ mod tests { let loop_host = host.clone(); let loop_registry = registry.clone(); let loop_roster = roster.clone(); + let loop_invites = invites.clone(); let loop_cancel = cancel.clone(); let task = tokio::spawn(async move { let ctx = ReconcileCtx { @@ -631,6 +880,7 @@ mod tests { knot: &knot, host: &loop_host, registry: &loop_registry, + invites: &loop_invites, roster: &loop_roster, }; reconcile_loop(&ctx, clock.as_ref(), &loop_cancel, rx).await; diff --git a/bobbin/crates/knot-ingest/src/roster.rs b/bobbin/crates/knot-ingest/src/roster.rs index a1d23f273..0a3dfc362 100644 --- a/bobbin/crates/knot-ingest/src/roster.rs +++ b/bobbin/crates/knot-ingest/src/roster.rs @@ -1,7 +1,8 @@ use std::collections::{HashMap, HashSet}; use std::sync::Arc; -use bobbin_edge_index::EdgeStore; +use bobbin_edge_index::{EdgeStore, InviteIndex, InviteScope, Offer}; +use bobbin_runtime::UnixMicros; use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; use bobbin_types::knot_acl::{self, KnotHostKey}; use jacquard_common::DefaultStr; @@ -11,115 +12,214 @@ use crate::registry::KnotRegistry; const REPO_COLLABORATOR_KIND: &str = "sh.tangled.repo.collaborator"; -#[derive(Clone, Copy, Debug, Eq, PartialEq, Ord, PartialOrd)] -pub struct Cursor(pub i64); +#[derive(Clone, Copy, Eq, Hash, PartialEq)] +enum Facet { + Acl, + Invite, +} + +#[derive(Clone, Copy, Eq, Hash, PartialEq)] +pub enum Scope<'a> { + Knot, + Repo(&'a Did), +} -impl Cursor { - fn micros(self) -> u64 { - self.0.max(0) as u64 / 1000 +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)] -enum DedupKey { - Member(Did), - Collaborator(Did, Did), +struct DedupKey { + facet: Facet, + repo: Option>, + subject: Did, } #[derive(Clone, Copy)] -pub enum AclKind<'a> { - Member, - Collaborator(&'a Did), +pub struct AclKind<'a> { + facet: Facet, + pub scope: Scope<'a>, +} + +impl<'a> AclKind<'a> { + pub(crate) fn acl(scope: Scope<'a>) -> Self { + Self { + facet: Facet::Acl, + scope, + } + } + + pub(crate) fn invite(scope: Scope<'a>) -> Self { + Self { + facet: Facet::Invite, + scope, + } + } + + pub(crate) fn is_invite(self) -> bool { + self.facet == Facet::Invite + } + + 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", + } + } + + fn dedup(self, subject: Did) -> DedupKey { + DedupKey { + facet: self.facet, + repo: self.scope.repo().cloned(), + subject, + } + } +} + +pub enum Assertion<'a> { + Acl(Scope<'a>), + Invite { + scope: Scope<'a>, + invited_by: Did, + }, +} + +impl<'a> Assertion<'a> { + fn kind(&self) -> AclKind<'a> { + match self { + Self::Acl(scope) => AclKind::acl(*scope), + Self::Invite { scope, .. } => AclKind::invite(*scope), + } + } +} + +#[derive(Clone, Copy, Debug, Default, Eq, Ord, PartialEq, PartialOrd)] +pub struct Applied(u64); + +#[derive(Clone, Copy, Eq, PartialEq)] +pub enum Source { + Firehose, + Listing(Applied), } struct SeenState { - cursor: Cursor, + at: UnixMicros, present: bool, + applied: Applied, } pub struct Roster { store: Arc, + invites: Arc, knot: Did, registry: Arc, host: KnotHostKey, seen: HashMap, + applied: Applied, } impl Roster { pub fn new( store: Arc, + invites: Arc, knot: Did, registry: Arc, host: KnotHostKey, ) -> Self { Self { store, + invites, knot, registry, host, seen: HashMap::new(), + applied: Applied::default(), } } - pub fn apply_member(&mut self, subject: Did, cursor: Cursor) { - if !self.advance(DedupKey::Member(subject.clone()), cursor, true) { - return; - } - if let Some((source, edges)) = - knot_acl::member_upsert(&self.knot, &subject, cursor.micros()) - { - self.store.upsert_source(&source, edges); + 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(), + }, } } - pub fn apply_collaborator( + pub fn apply( &mut self, - repo: Did, + assertion: Assertion<'_>, subject: Did, - cursor: Cursor, + at: UnixMicros, + source: Source, ) { - if !self.registry.repo_on_host(&self.host, &repo) { + let kind = assertion.kind(); + if !kind + .scope + .repo() + .is_none_or(|repo| self.registry.repo_on_host(&self.host, repo)) + { return; } - if !self.advance( - DedupKey::Collaborator(repo.clone(), subject.clone()), - cursor, - true, - ) { + if !self.advance(kind.dedup(subject.clone()), at, source) { return; } - if let Some((source, edges)) = - knot_acl::collaborator_upsert(&repo, &subject, cursor.micros()) - { - self.store.upsert_source(&source, edges); + 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( + subject, + self.invite_scope(scope), + Offer { + invited_by, + offered_at: at, + }, + ), } } - pub fn max_cursor(&self) -> Cursor { - self.seen - .values() - .map(|state| state.cursor) - .max() - .unwrap_or(Cursor(0)) - } - - pub fn reap(&mut self, kind: AclKind<'_>, present: &HashSet>, horizon: Cursor) { - let stale: Vec> = - self.seen - .iter() - .filter_map(|(key, state)| { - let subject = match (key, kind) { - (DedupKey::Member(subject), AclKind::Member) => subject, - ( - DedupKey::Collaborator(edge_repo, subject), - AclKind::Collaborator(repo), - ) if edge_repo == repo => subject, - _ => return None, - }; - (state.present && state.cursor <= horizon && !present.contains(subject)) - .then(|| subject.clone()) - }) - .collect(); + pub fn applied(&self) -> Applied { + self.applied + } + + fn step(&mut self) -> Applied { + self.applied = Applied(self.applied.0.saturating_add(1)); + self.applied + } + + pub fn reap( + &mut self, + kind: AclKind<'_>, + present: &HashSet>, + horizon: Applied, + ) { + let stale: Vec> = self + .seen + .iter() + .filter(|(key, state)| { + key.facet == kind.facet + && key.repo.as_ref() == kind.scope.repo() + && state.present + && state.applied <= horizon + && !present.contains(&key.subject) + }) + .map(|(key, _)| key.subject.clone()) + .collect(); stale .into_iter() .for_each(|subject| self.retire(kind, subject)); @@ -144,31 +244,53 @@ impl Roster { }); } - fn retire(&mut self, kind: AclKind<'_>, subject: Did) { - let source = match kind { - AclKind::Member => knot_acl::member_source(&self.knot, &subject), - AclKind::Collaborator(repo) => knot_acl::collaborator_source(repo, &subject), - }; - if let Some(source) = source { - self.store.remove_source(&source); + 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::Invite => self + .invites + .withdraw(&subject, &self.invite_scope(kind.scope)), } - let key = match kind { - AclKind::Member => DedupKey::Member(subject), - AclKind::Collaborator(repo) => DedupKey::Collaborator(repo.clone(), subject), - }; - if let Some(state) = self.seen.get_mut(&key) { + let applied = self.step(); + if let Some(state) = self.seen.get_mut(&kind.dedup(subject)) { state.present = false; + state.applied = applied; } } - fn advance(&mut self, key: DedupKey, cursor: Cursor, present: bool) -> bool { - match self.seen.get(&key) { - Some(state) if cursor <= state.cursor => false, - _ => { - self.seen.insert(key, SeenState { cursor, present }); - true + fn advance(&mut self, key: DedupKey, at: UnixMicros, source: Source) -> bool { + let seen = self + .seen + .get(&key) + .map(|state| (state.at, state.present, state.applied)); + let stale = match (seen, source) { + (Some((previous, present, applied)), Source::Listing(taken)) => { + applied > taken || (present && at <= previous) } + _ => false, + }; + if stale { + return false; } + let applied = self.step(); + let at = seen.map_or(at, |(previous, _, _)| at.max(previous)); + self.seen.insert( + key, + SeenState { + at, + present: true, + applied, + }, + ); + true } } @@ -235,21 +357,75 @@ mod tests { Arc::new(KnotRegistry::new()) } + fn invites() -> Arc { + Arc::new(InviteIndex::new()) + } + fn member_roster(store: Arc) -> Roster { - Roster::new(store, knot(), empty_registry(), host()) + Roster::new(store, invites(), knot(), empty_registry(), host()) + } + + fn repo_roster( + store: Arc, + index: Arc, + repo: &Did, + ) -> Roster { + let registry = empty_registry(); + registry.observe_repo(&host(), repo.clone()); + Roster::new(store, index, knot(), registry, host()) + } + + fn indexed_roster(index: Arc) -> Roster { + Roster::new(store(), index, knot(), empty_registry(), host()) } 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), + } + } + + fn repo_offer_of<'a>(repo: &'a Did, admin: &str) -> Assertion<'a> { + Assertion::Invite { + scope: Scope::Repo(repo), + invited_by: did(admin), + } + } + + impl Roster { + fn listed(&mut self, assertion: Assertion<'_>, subject: &Did, micros: u64) { + let source = Source::Listing(self.applied()); + self.apply(assertion, subject.clone(), UnixMicros::new(micros), source); + } + + fn firehosed(&mut self, assertion: Assertion<'_>, subject: &Did, micros: u64) { + self.apply( + assertion, + subject.clone(), + UnixMicros::new(micros), + Source::Firehose, + ); + } + } + #[test] fn late_add_leaves_newer_member_row_standing() { let store = store(); let mut roster = member_roster(store.clone()); let m = did("did:plc:boltless"); - roster.apply_member(m.clone(), Cursor(5)); - roster.apply_member(m.clone(), Cursor(1)); + roster.listed(MEMBER, &m, 5); + roster.listed(MEMBER, &m, 1); assert_eq!(member_count(&store, &m), 1); } @@ -258,8 +434,8 @@ mod tests { let store = store(); let mut roster = member_roster(store.clone()); let m = did("did:plc:akshay"); - roster.apply_member(m.clone(), Cursor(10)); - roster.apply_member(m.clone(), Cursor(10)); + roster.listed(MEMBER, &m, 10); + roster.listed(MEMBER, &m, 10); assert_eq!(member_count(&store, &m), 1); } @@ -267,11 +443,8 @@ mod tests { fn collaborator_on_hosted_repo_is_indexed() { let store = store(); let repo = did("did:plc:scallop"); - let registry = empty_registry(); - registry.observe_repo(&host(), repo.clone()); - let mut roster = Roster::new(store.clone(), knot(), registry, host()); - let subject = did("did:plc:olaren"); - roster.apply_collaborator(repo.clone(), subject.clone(), Cursor(7)); + let mut roster = repo_roster(store.clone(), invites(), &repo); + roster.listed(collab(&repo), &did("did:plc:olaren"), 7); assert_eq!(collaborator_count(&store, &repo), 1); } @@ -281,11 +454,11 @@ mod tests { let mut roster = member_roster(store.clone()); let stayed = did("did:plc:akshay"); let left = did("did:plc:boltless"); - roster.apply_member(stayed.clone(), Cursor(10)); - roster.apply_member(left.clone(), Cursor(20)); + roster.listed(MEMBER, &stayed, 10); + roster.listed(MEMBER, &left, 20); - let horizon = roster.max_cursor(); - roster.reap(AclKind::Member, &subjects(&[&stayed]), horizon); + let horizon = roster.applied(); + roster.reap(AclKind::acl(Scope::Knot), &subjects(&[&stayed]), horizon); assert_eq!(member_count(&store, &stayed), 1); assert_eq!( @@ -300,10 +473,10 @@ mod tests { let store = store(); let mut roster = member_roster(store.clone()); let m = did("did:plc:boltless"); - let horizon = roster.max_cursor(); - roster.apply_member(m.clone(), Cursor(100)); + let horizon = roster.applied(); + roster.listed(MEMBER, &m, 100); - roster.reap(AclKind::Member, &subjects(&[]), horizon); + roster.reap(AclKind::acl(Scope::Knot), &subjects(&[]), horizon); assert_eq!( member_count(&store, &m), @@ -316,14 +489,11 @@ mod tests { fn reap_removes_departed_collaborator() { let store = store(); let repo = did("did:plc:scallop"); - let registry = empty_registry(); - registry.observe_repo(&host(), repo.clone()); - let mut roster = Roster::new(store.clone(), knot(), registry, host()); - let left = did("did:plc:olaren"); - roster.apply_collaborator(repo.clone(), left.clone(), Cursor(7)); + let mut roster = repo_roster(store.clone(), invites(), &repo); + roster.listed(collab(&repo), &did("did:plc:olaren"), 7); - let horizon = roster.max_cursor(); - roster.reap(AclKind::Collaborator(&repo), &subjects(&[]), horizon); + let horizon = roster.applied(); + roster.reap(AclKind::acl(Scope::Repo(&repo)), &subjects(&[]), horizon); assert_eq!(collaborator_count(&store, &repo), 0); } @@ -332,11 +502,9 @@ mod tests { fn purge_legacy_strips_pds_collaborator_keeps_knot_owned() { let store = store(); let repo = did("did:plc:scallop"); - let registry = empty_registry(); - registry.observe_repo(&host(), repo.clone()); - let mut roster = Roster::new(store.clone(), knot(), registry, host()); + let mut roster = repo_roster(store.clone(), invites(), &repo); - roster.apply_collaborator(repo.clone(), did("did:plc:olaren"), Cursor(7)); + roster.listed(collab(&repo), &did("did:plc:olaren"), 7); add_legacy_edge( &store, "sh.tangled.repo.collaborator", @@ -357,7 +525,7 @@ mod tests { fn purge_legacy_removes_indexed_member_edges() { let store = store(); let registry = empty_registry(); - let roster = Roster::new(store.clone(), knot(), registry.clone(), host()); + 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"); @@ -373,12 +541,142 @@ mod tests { fn collaborator_for_unhosted_repo_is_dropped() { let store = store(); let repo = did("did:plc:scallop"); - let mut roster = Roster::new(store.clone(), knot(), empty_registry(), host()); - roster.apply_collaborator(repo.clone(), did("did:plc:olaren"), Cursor(7)); + let mut roster = member_roster(store.clone()); + roster.listed(collab(&repo), &did("did:plc:olaren"), 7); assert_eq!( collaborator_count(&store, &repo), 0, "a knot cannot assert collaborators on a repo it does not host" ); } + + #[test] + fn offer_reaches_index_scoped_to_its_source_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 offered = index.outstanding(&member); + 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 { + knot: knot(), + repo: repo.clone() + }, + "repo's offer is scoped to repo and never to knot" + ); + + roster.retire(AclKind::invite(Scope::Knot), member.clone()); + assert!(index.outstanding(&member).is_empty()); + assert_eq!( + index.outstanding(&collaborator).len(), + 1, + "withdrawing knot's offer leaves repo's 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 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); + + let horizon = roster.applied(); + roster.reap(AclKind::invite(Scope::Knot), &subjects(&[&stayed]), horizon); + + assert_eq!(index.outstanding(&stayed).len(), 1); + assert!( + index.outstanding(&withdrawn).is_empty(), + "knot's filtered listing wins over firehose record it no longer reports" + ); + + roster.listed(offer_of("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" + ); + } + + #[test] + fn listing_fetched_before_withdrawal_leaves_offer_withdrawn() { + let index = invites(); + let mut roster = indexed_roster(index.clone()); + let invitee = did("did:plc:limpet"); + let moment = 8_000_000; + roster.listed(offer_of("did:plc:akshay"), &invitee, moment); + let fetched = roster.applied(); + roster.retire(AclKind::invite(Scope::Knot), invitee.clone()); + + roster.apply( + offer_of("did:plc:akshay"), + invitee.clone(), + UnixMicros::new(moment), + Source::Listing(fetched), + ); + + assert!( + index.outstanding(&invitee).is_empty(), + "firehose delivered delete after this page was fetched, so page shows knot that has already moved on" + ); + } + + #[test] + fn firehose_reoffer_in_same_second_reaches_invitee() { + let index = invites(); + let mut roster = indexed_roster(index.clone()); + 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); + + let outstanding = index.outstanding(&invitee); + assert_eq!( + outstanding.len(), + 1, + "knots stamp offers to the second, so withdraw and re-offer inside one second spell one cursor; firehose order settles which won" + ); + assert_eq!(outstanding[0].1.invited_by, did("did:plc:olaren")); + } + + #[test] + fn offer_applied_mid_fetch_survives_snapshot_that_missed_it() { + let index = invites(); + let mut roster = indexed_roster(index.clone()); + 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); + let horizon = roster.applied(); + + roster.firehosed(offer_of("did:plc:akshay"), &raced, moment); + roster.reap(AclKind::invite(Scope::Knot), &subjects(&[&listed]), horizon); + + assert_eq!( + index.outstanding(&raced).len(), + 1, + "admins inviting two people inside one second stamp both offers alike, and the second reached bobbin after this page was fetched" + ); + } } diff --git a/bobbin/crates/knot-ingest/src/stream.rs b/bobbin/crates/knot-ingest/src/stream.rs index d20be6712..e27f3b2d8 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}; +use crate::firehose::{self, Frame, InviteEvent, InviteOp}; use crate::legacy; const SUBSCRIBE_REPOS_PATH: &str = "xrpc/com.atproto.sync.subscribeRepos"; @@ -105,6 +105,7 @@ pub struct StreamConfig { pub trait FeedHandler: Send + Sync + 'static { fn roster_changed(&self); fn repo_roster_changed(&self, repo: &Did); + fn invite_op(&self, op: InviteOp); } #[derive(Debug)] @@ -301,10 +302,20 @@ fn dispatch_commit(commit: &firehose::Commit, handler: &dyn FeedHandler) { commit .records .iter() - .filter_map(|op| InviteEvent::of(&op.collection)) - .for_each(|invite| match invite { - InviteEvent::KnotMember => handler.roster_changed(), - InviteEvent::RepoCollaborator => handler.repo_roster_changed(&commit.repo), + .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) { + Some(invite) => handler.invite_op(invite), + None => tracing::warn!( + repo = %commit.repo.as_ref(), + rkey = %op.rkey, + "invite record didn't decode, periodic reconcile will cover it" + ), + } }); } @@ -328,20 +339,25 @@ mod tests { use std::collections::VecDeque; use std::sync::atomic::{AtomicUsize, Ordering}; + use bobbin_edge_index::Offer; use bobbin_runtime::{ NetworkError, SimClock, SystemClock, UnixMicros, WsConnectFuture, WsMessageFuture, WsSendFuture, WsSink, WsStream, }; use bytes::Bytes; + use jacquard_common::types::nsid::Nsid; + use jacquard_common::types::recordkey::Rkey; use parking_lot::Mutex; use super::*; use crate::firehose::testcbor::*; + use crate::firehose::{InviteTarget, OpAction}; #[derive(Clone, Debug, PartialEq)] enum Entry { Full, Repo(Did), + Offer(InviteOp), } #[derive(Default)] @@ -367,6 +383,10 @@ mod tests { fn repo_roster_changed(&self, repo: &Did) { self.push(Entry::Repo(repo.clone())); } + + fn invite_op(&self, op: InviteOp) { + self.push(Entry::Offer(op)); + } } struct ScriptStream { @@ -560,6 +580,87 @@ mod tests { assert!(matches!(&sent[0], WsMessage::Pong(p) if p.as_ref() == b"ka")); } + fn invite_record_commit( + repo: &str, + action: OpAction, + collection: &str, + rkey: &str, + record: Option>, + ) -> firehose::Commit { + firehose::Commit { + repo: did(repo), + seq: 12, + rev: "3zzzzzzzzzzzz".to_owned(), + records: vec![firehose::RecordOp { + action, + collection: Nsid::::new_owned(collection).unwrap(), + rkey: Rkey::::new_owned(rkey).unwrap(), + record: record.map(Bytes::from), + }], + } + } + + #[test] + fn invite_create_both_nudges_roster_and_delivers_offer() { + let handler = RecordingHandler::default(); + let collection = InviteEvent::KnotMember.collection(); + let commit = invite_record_commit( + "did:web:knot.oyster.cafe", + OpAction::Create, + collection, + "did:plc:limpet", + Some(invite_record( + collection, + "2026-06-01T00:00:00Z", + "did:plc:akshay", + )), + ); + + dispatch_commit(&commit, &handler); + + assert_eq!( + handler.entries(), + vec![ + Entry::Full, + Entry::Offer(InviteOp { + target: InviteTarget::Knot, + invitee: did("did:plc:limpet"), + offer: Some(Offer { + invited_by: did("did:plc:akshay"), + offered_at: UnixMicros::new(1_780_272_000_000_000), + }), + }), + ], + "one record is both roster signal and standing offer" + ); + } + + #[test] + fn invite_delete_delivers_op_without_offer() { + let handler = RecordingHandler::default(); + let commit = invite_record_commit( + "did:plc:scallop", + OpAction::Delete, + InviteEvent::RepoCollaborator.collection(), + "did:plc:olaren", + None, + ); + + dispatch_commit(&commit, &handler); + + assert_eq!( + handler.entries(), + vec![ + Entry::Repo(did("did:plc:scallop")), + Entry::Offer(InviteOp { + target: InviteTarget::Repo(did("did:plc:scallop")), + invitee: did("did:plc:olaren"), + offer: None, + }), + ] + ); + } + #[tokio::test] async fn outdated_cursor_info_ends_the_session_in_resync() { let handler = RecordingHandler::default();