From a0d381381751c86d4ed4b8ef46b8b7c6e42d8659 Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 16 Jun 2026 22:29:31 +0300 Subject: [PATCH] bobbin/knot-ingest: handle knot115 knot-owned data Lewis: May this revision serve well! --- Cargo.toml | 1 + bobbin/crates/bobbin-sim/src/runtime.rs | 2 + bobbin/crates/bobbin/Cargo.toml | 1 + bobbin/crates/bobbin/src/main.rs | 49 +- bobbin/crates/edge-index/src/lib.rs | 42 +- bobbin/crates/ingest/Cargo.toml | 1 + bobbin/crates/ingest/examples/smoke.rs | 2 + bobbin/crates/ingest/src/lib.rs | 171 ++++++ bobbin/crates/knot-ingest/Cargo.toml | 30 ++ bobbin/crates/knot-ingest/src/client.rs | 497 ++++++++++++++++++ bobbin/crates/knot-ingest/src/gate.rs | 255 +++++++++ bobbin/crates/knot-ingest/src/lib.rs | 13 + bobbin/crates/knot-ingest/src/orchestrator.rs | 463 ++++++++++++++++ bobbin/crates/knot-ingest/src/registry.rs | 197 +++++++ bobbin/crates/knot-ingest/src/roster.rs | 445 ++++++++++++++++ bobbin/crates/knot-ingest/src/stream.rs | 446 ++++++++++++++++ bobbin/crates/knot-proxy/src/dns.rs | 4 +- bobbin/crates/knot-proxy/src/lib.rs | 3 +- bobbin/crates/runtime/src/lib.rs | 6 +- bobbin/crates/runtime/src/network.rs | 49 ++ bobbin/crates/types/src/knot_acl.rs | 289 ++++++++++ bobbin/crates/types/src/lib.rs | 1 + bobbin/crates/xrpc/src/lib.rs | 72 ++- bobbin/crates/xrpc/tests/aggregation.rs | 76 +++ 24 files changed, 3097 insertions(+), 18 deletions(-) create mode 100644 bobbin/crates/knot-ingest/Cargo.toml create mode 100644 bobbin/crates/knot-ingest/src/client.rs create mode 100644 bobbin/crates/knot-ingest/src/gate.rs create mode 100644 bobbin/crates/knot-ingest/src/lib.rs create mode 100644 bobbin/crates/knot-ingest/src/orchestrator.rs create mode 100644 bobbin/crates/knot-ingest/src/registry.rs create mode 100644 bobbin/crates/knot-ingest/src/roster.rs create mode 100644 bobbin/crates/knot-ingest/src/stream.rs create mode 100644 bobbin/crates/types/src/knot_acl.rs diff --git a/Cargo.toml b/Cargo.toml index 52c26b59..123a783f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -34,6 +34,7 @@ bobbin-resolver = { path = "bobbin/crates/resolver" } bobbin-slingshot-client = { path = "bobbin/crates/slingshot-client" } bobbin-record-lru = { path = "bobbin/crates/record-lru" } bobbin-knot-proxy = { path = "bobbin/crates/knot-proxy" } +bobbin-knot-ingest = { path = "bobbin/crates/knot-ingest" } bobbin-runtime = { path = "bobbin/crates/runtime" } bobbin-search = { path = "bobbin/crates/search" } bobbin-sim = { path = "bobbin/crates/bobbin-sim" } diff --git a/bobbin/crates/bobbin-sim/src/runtime.rs b/bobbin/crates/bobbin-sim/src/runtime.rs index d74e59be..fd229047 100644 --- a/bobbin/crates/bobbin-sim/src/runtime.rs +++ b/bobbin/crates/bobbin-sim/src/runtime.rs @@ -122,6 +122,8 @@ impl Sim { disconnects: Some(disconnects.clone()), warming_shadow: Some(warming_shadow.clone()), warming_buffer: warming_buffer_enabled.then(|| warming_buffer.clone()), + knot_registry: None, + knot_gate: None, }; let ingest_config = IngestConfig { hydrant_base, diff --git a/bobbin/crates/bobbin/Cargo.toml b/bobbin/crates/bobbin/Cargo.toml index fc040eb4..6ee7112b 100644 --- a/bobbin/crates/bobbin/Cargo.toml +++ b/bobbin/crates/bobbin/Cargo.toml @@ -12,6 +12,7 @@ path = "src/main.rs" [dependencies] bobbin-edge-index = { workspace = true } bobbin-ingest = { workspace = true } +bobbin-knot-ingest = { workspace = true } bobbin-knot-proxy = { workspace = true } bobbin-record-lru = { workspace = true } bobbin-runtime = { workspace = true } diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index 49ace415..bdb082dd 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -10,9 +10,13 @@ use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor, StateIndex}; use bobbin_ingest::{ IngestConfig, IngestRuntime, RepoIdResolver, WarmingBuffer, run as run_ingest, }; -use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; +use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry, Orchestrator}; +use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig, classify_ip}; use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; -use bobbin_runtime::{Clock, MemoryBudget, OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; +use bobbin_runtime::{ + Clock, GuardedWs, MemoryBudget, NetworkError, OsEntropy, RuntimeHasher, SystemClock, + TungsteniteWs, WsTransport, +}; use bobbin_search::{SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; use bobbin_xrpc::{ @@ -186,6 +190,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { let pull_statuses = Arc::new(StateIndex::new(hasher.clone())); let coverage = Arc::new(CoverageWatch::new()); let warming_buffer = Arc::new(WarmingBuffer::new(hasher.clone())); + let knot_registry = Arc::new(KnotRegistry::new()); let knots = Arc::new(KnotProxy::new( KnotProxyConfig { allow_private_hosts: cfg.knot.allow_private, @@ -215,6 +220,29 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { }; let cancel = CancellationToken::new(); let ingest_coverage = coverage.clone(); + + let knot_acl_dev = !cfg.knot.require_https; + let knot_allow_private = cfg.knot.allow_private; + let knot_client = KnotClient::with_default_http(knot_allow_private)?; + let knot_gate = Arc::new(CapabilityGate::new( + knot_client.clone(), + clock.clone(), + knot_acl_dev, + knot_allow_private, + )); + let knot_ws: Arc = if knot_allow_private { + ws.clone() + } else { + GuardedWs::shared(Arc::new(|addrs: &[SocketAddr]| { + match addrs.iter().find_map(|sa| classify_ip(&sa.ip())) { + Some(reason) => Err(NetworkError::Connect(format!( + "knot eventstream resolves to {reason} address space" + ))), + None => Ok(()), + } + })) + }; + let ingest_runtime = IngestRuntime { store: edges.clone(), issue_states: issue_states.clone(), @@ -225,14 +253,29 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { resolver: resolver.clone(), clock: clock.clone(), entropy, - ws, + ws: ws.clone(), cancel: cancel.clone(), disconnects: None, warming_shadow: None, warming_buffer: Some(warming_buffer), + knot_registry: Some(knot_registry.clone()), + knot_gate: Some(knot_gate.clone()), }; let mut ingest_handle = tokio::spawn(run_ingest(ingest_cfg, ingest_runtime)); + let knot_orchestrator = Orchestrator { + client: Arc::new(knot_client), + gate: knot_gate, + registry: knot_registry, + store: edges.clone(), + ws: knot_ws, + clock: clock.clone(), + dev: knot_acl_dev, + allow_private: knot_allow_private, + cancel: cancel.clone(), + }; + let _knot_acl_handle = tokio::spawn(knot_orchestrator.run()); + let _adaptive_watcher = budget.zip(limiter.as_ref()).map(|(b, l)| { mem::spawn_adaptive_watcher( l.clone(), diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index d26a35b8..b81240f3 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -227,10 +227,22 @@ impl PageLimit { #[derive(Debug)] pub struct EdgePage { - pub items: Vec>, + pub items: Vec, pub next: Option, } +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct EdgeItem { + pub uri: AtUri, + pub sort_micros: u64, +} + +impl AsRef for EdgeItem { + fn as_ref(&self) -> &str { + self.uri.as_ref() + } +} + #[derive(Clone, Copy, Debug, Default)] pub struct EdgeMemReport { pub key_count: u64, @@ -639,6 +651,23 @@ impl EdgeStore { .unwrap_or(0) } + pub fn sources_for(&self, key: &EdgeKey) -> Vec> { + self.lookup_key(key) + .and_then(|id| { + self.forward.read_sync(&id, |_, sources| { + sources + .directed(PageCursor::Start, SortDir::Desc) + .filter_map(|bucket| { + let spur = bucket.source.to_spur()?; + let stored = self.source_interner.try_resolve(&spur)?; + AtUri::new_owned(self.decode_source(stored)?).ok() + }) + .collect::>() + }) + }) + .unwrap_or_default() + } + pub fn list( &self, key: &EdgeKey, @@ -659,7 +688,11 @@ impl EdgeStore { .filter_map(|&key| { let spur = key.source.to_spur()?; let stored = self.source_interner.try_resolve(&spur)?; - AtUri::new_owned(self.decode_source(stored)?).ok() + let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; + Some(EdgeItem { + uri, + sort_micros: key.micros.0, + }) }) .collect(); let next = has_more @@ -732,7 +765,10 @@ impl EdgeStore { .matched .into_iter() .take(visible_len) - .map(|(_, uri)| uri) + .map(|(key, uri)| EdgeItem { + uri, + sort_micros: key.micros.0, + }) .collect(); EdgePage { items, next } }) diff --git a/bobbin/crates/ingest/Cargo.toml b/bobbin/crates/ingest/Cargo.toml index 6f00621b..3506424e 100644 --- a/bobbin/crates/ingest/Cargo.toml +++ b/bobbin/crates/ingest/Cargo.toml @@ -8,6 +8,7 @@ rust-version.workspace = true [dependencies] bobbin-types = { workspace = true } bobbin-edge-index = { workspace = true } +bobbin-knot-ingest = { workspace = true } bobbin-record-lru = { workspace = true } bobbin-resolver = { workspace = true } bobbin-runtime = { workspace = true } diff --git a/bobbin/crates/ingest/examples/smoke.rs b/bobbin/crates/ingest/examples/smoke.rs index 96d7f8a7..7a105c7f 100644 --- a/bobbin/crates/ingest/examples/smoke.rs +++ b/bobbin/crates/ingest/examples/smoke.rs @@ -46,6 +46,8 @@ async fn main() { disconnects: None, warming_shadow: None, warming_buffer: None, + knot_registry: None, + knot_gate: None, }; let task = tokio::spawn(async move { let _ = run(cfg, runtime).await; diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index b0e413ce..5ec1278c 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -7,6 +7,7 @@ use bobbin_edge_index::{ ApplyOutcome, Coverage, CoverageWatch, EdgeStore, HydrantCursor, IssueStateKind, PromotionSignal, PullStatusKind, StateIndex, apply_record_state, }; +use bobbin_knot_ingest::{CapabilityGate, KnotRegistry}; use bobbin_record_lru::RecordStore; use bobbin_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at}; use bobbin_runtime::{ @@ -15,6 +16,7 @@ use bobbin_runtime::{ }; use bobbin_types::edges::{Edge, ExtractError, Record}; use bobbin_types::ids::{RepoIdent, SubjectRef}; +use bobbin_types::knot_acl::KnotHostKey; use bobbin_types::record::RecordBody; use bobbin_types::search::{SearchSink, SearchableRecord}; use bytes::Bytes; @@ -215,6 +217,8 @@ pub struct IngestRuntime { pub disconnects: Option>, pub warming_shadow: Option>, pub warming_buffer: Option>, + pub knot_registry: Option>, + pub knot_gate: Option>, } impl Clone for IngestRuntime { @@ -234,6 +238,8 @@ impl Clone for IngestRuntime { disconnects: self.disconnects.clone(), warming_shadow: self.warming_shadow.clone(), warming_buffer: self.warming_buffer.clone(), + knot_registry: self.knot_registry.clone(), + knot_gate: self.knot_gate.clone(), } } } @@ -250,6 +256,8 @@ impl IngestRuntime { search: &self.search, shadow: self.warming_shadow.as_deref(), buffer: self.warming_buffer.as_deref(), + knot_registry: self.knot_registry.as_deref(), + knot_gate: self.knot_gate.as_deref(), } } } @@ -264,6 +272,8 @@ struct PipelineCtx<'a, S: SearchSink + 'static> { search: &'a S, shadow: Option<&'a WarmingShadowBuffer>, buffer: Option<&'a WarmingBuffer>, + knot_registry: Option<&'a KnotRegistry>, + knot_gate: Option<&'a CapabilityGate>, } pub async fn run( @@ -1062,6 +1072,28 @@ async fn prepare_record( finalize_drained(ctx, drained).await; } } + if let Some(registry) = ctx.knot_registry { + let host = KnotHostKey::new(repo.knot.as_ref()); + match repo.repo_did.clone() { + Some(repo_did) => registry.observe_repo(&host, repo_did), + None => registry.observe_host(&host), + } + } + } + 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::Delete { source, nsid }; + } + 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) { Ok(es) => es, @@ -1085,6 +1117,11 @@ async fn prepare_record( if nsid.as_ref() == "sh.tangled.repo" { ctx.resolver.forget(&record.did, &record.rkey).await; } + if nsid.as_ref() == "sh.tangled.knot.member" + && let Some(registry) = ctx.knot_registry + { + registry.forget_legacy_member(&source); + } PendingOp::Delete { source, nsid } } RecordAction::Other => { @@ -1094,6 +1131,43 @@ async fn prepare_record( } } +enum AclDisposition { + Other, + NativeSkip, + LegacyMember { host: KnotHostKey }, +} + +fn acl_disposition( + parsed: &Record, + gate: Option<&CapabilityGate>, + registry: Option<&KnotRegistry>, +) -> AclDisposition { + let Some(gate) = gate else { + 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)) + .is_some_and(|host| gate.is_native(&host)); + if native { + AclDisposition::NativeSkip + } else { + AclDisposition::Other + } + } + _ => AclDisposition::Other, + } +} + fn fallback_rfc3339( rkey: &Rkey, rev: &jacquard_common::types::tid::Tid, @@ -1375,6 +1449,8 @@ async fn handle_frame( search, shadow: None, buffer: None, + knot_registry: None, + knot_gate: None, }; let pending = prepare_frame(frame, &ctx, now).await; let pending = resolve_pending(pending, &ctx).await; @@ -1613,6 +1689,95 @@ mod tests { assert_eq!(cov.snapshot().events_processed(), 1); } + #[tokio::test] + async fn native_knot_member_skipped_legacy_indexed() { + use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry}; + use wiremock::matchers::{method, path}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "version": "1.1.0", + "capabilities": ["knot-acl"] + }))) + .mount(&server) + .await; + let url = url::Url::parse(&server.uri()).unwrap(); + let native_host = format!("{}:{}", url.host_str().unwrap(), url.port().unwrap()); + + let gate = CapabilityGate::new( + KnotClient::with_default_http(true).unwrap(), + Arc::new(SystemClock::new()), + true, + true, + ); + assert!(gate.has_knot_acl(&KnotHostKey::new(&native_host)).await); + + let registry = KnotRegistry::new(); + let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + let ctx = PipelineCtx { + resolver: &resolver, + store: &store, + issue_states: &issue_states, + pull_statuses: &pull_statuses, + coverage: &cov, + records: &NoopRecordStore, + search: &NoopSearchSink, + shadow: None, + buffer: None, + knot_registry: Some(®istry), + knot_gate: Some(&gate), + }; + + let member_frame = |id: u64, rkey: &str, domain: &str| { + parse_frame(json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": "did:plc:akshay", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.knot.member", + "rkey": rkey, + "action": "create", + "record": { + "$type": "sh.tangled.knot.member", + "subject": "did:plc:boltless", + "domain": domain, + "createdAt": "2026-06-01T00:00:00Z" + } + } + })) + }; + + let native = + prepare_frame(member_frame(1, "aaaaaaaaaaaaz", &native_host), &ctx, now()).await; + assert!( + matches!(native.op, PendingOp::Delete { .. }), + "member record for a native knot must be dropped" + ); + + let legacy = + prepare_frame(member_frame(2, "bbbbbbbbbbbbz", "legacy.knot"), &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" + ); + } + #[tokio::test] async fn create_then_delete_round_trips_a_star() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); @@ -1752,6 +1917,8 @@ mod tests { search: &search, shadow: None, buffer: None, + knot_registry: None, + knot_gate: None, }; let mk = |live: bool| -> HydrantFrame { parse_frame(json!({ @@ -2418,6 +2585,8 @@ mod tests { disconnects: None, warming_shadow: None, warming_buffer: None, + knot_registry: None, + knot_gate: None, } } @@ -3046,6 +3215,8 @@ mod tests { disconnects: None, warming_shadow: None, warming_buffer: None, + knot_registry: None, + knot_gate: None, }; let parallelism = 4usize; diff --git a/bobbin/crates/knot-ingest/Cargo.toml b/bobbin/crates/knot-ingest/Cargo.toml new file mode 100644 index 00000000..36253c89 --- /dev/null +++ b/bobbin/crates/knot-ingest/Cargo.toml @@ -0,0 +1,30 @@ +[package] +name = "bobbin-knot-ingest" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +bobbin-edge-index = { workspace = true } +bobbin-knot-proxy = { workspace = true } +bobbin-runtime = { workspace = true } +bobbin-types = { workspace = true } +jacquard-common = { workspace = true } + +bytes = { workspace = true } +chrono = { workspace = true } +futures = { workspace = true } +http = { workspace = true } +reqwest = { workspace = true } +serde = { workspace = true } +serde_json = { workspace = true } +thiserror = { workspace = true } +tokio = { workspace = true } +tokio-util = { workspace = true } +tracing = { workspace = true } +url = { workspace = true } + +[dev-dependencies] +tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } +wiremock = { workspace = true } diff --git a/bobbin/crates/knot-ingest/src/client.rs b/bobbin/crates/knot-ingest/src/client.rs new file mode 100644 index 00000000..773366b1 --- /dev/null +++ b/bobbin/crates/knot-ingest/src/client.rs @@ -0,0 +1,497 @@ +use std::future::Future; +use std::pin::Pin; +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 bytes::{Bytes, BytesMut}; +use chrono::{DateTime, Utc}; +use futures::TryStreamExt; +use http::{HeaderMap, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; +use serde::Deserialize; +use thiserror::Error; +use url::Url; + +const USER_AGENT: &str = concat!("bobbin/", env!("CARGO_PKG_VERSION")); +const REQUEST_TIMEOUT: Duration = Duration::from_secs(10); +const CONNECT_TIMEOUT: Duration = Duration::from_secs(5); +const MAX_BODY_BYTES: u64 = 4 * 1024 * 1024; +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"; + +#[derive(Clone)] +pub struct KnotClient { + http: Arc, +} + +#[derive(Clone, Debug, Eq, PartialEq, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct AclEntry { + pub subject: Did, + pub created_at: DateTime, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum Completeness { + Complete, + Truncated, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AclListing { + pub entries: Vec, + pub completeness: Completeness, +} + +#[derive(Debug, Error)] +pub enum KnotClientError { + #[error("knot host: {0}")] + Host(#[from] KnotHostError), + #[error("blocked: knot {host} resolves to {reason} address space")] + PrivateHost { + host: String, + reason: PrivateHostReason, + }, + #[error("http client build: {0}")] + Build(String), + #[error("network: {0}")] + Network(#[from] NetworkError), + #[error("xrpc not found")] + NotFound, + #[error("upstream returned status {0}")] + Upstream(StatusCode), + #[error("response body exceeded {limit} bytes")] + BodyTooLarge { limit: u64 }, + #[error("decode response: {0}")] + Decode(#[from] serde_json::Error), +} + +pub fn knot_endpoint( + host: &str, + dev: bool, + allow_private: bool, +) -> Result { + let raw = if dev { + format!("http://{host}") + } else { + host.to_owned() + }; + let knot = KnotHost::parse(&raw)?; + if !allow_private && let Some(reason) = knot.private_literal_reason() { + return Err(KnotClientError::PrivateHost { + host: host.to_owned(), + reason, + }); + } + Ok(knot) +} + +fn default_http_client(allow_private: bool) -> Result { + reqwest::Client::builder() + .user_agent(USER_AGENT) + .timeout(REQUEST_TIMEOUT) + .connect_timeout(CONNECT_TIMEOUT) + .redirect(reqwest::redirect::Policy::none()) + .dns_resolver(Arc::new(PrivateAddressFilter::new(allow_private))) + .build() +} + +fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).expect("static nsid literal must validate") +} + +pub(crate) fn authority(host: &KnotHost) -> String { + let url = host.url(); + match (url.host_str(), url.port()) { + (Some(h), Some(p)) => format!("{h}:{p}"), + (Some(h), None) => h.to_owned(), + (None, _) => String::new(), + } +} + +impl KnotClient { + pub fn new(http: Arc) -> Self { + Self { http } + } + + pub fn with_default_http(allow_private: bool) -> Result { + let client = default_http_client(allow_private) + .map_err(|e| KnotClientError::Build(e.to_string()))?; + Ok(Self::new(ReqwestHttp::shared(client))) + } + + pub async fn capabilities(&self, host: &KnotHost) -> Result, KnotClientError> { + let mut url = host.xrpc_url(&nsid(VERSION_NSID)); + url.set_query(None); + let bytes = self.get_json(url).await?; + let resp: VersionWire = serde_json::from_slice(&bytes)?; + Ok(resp.capabilities.unwrap_or_default()) + } + + pub async fn list_members(&self, host: &KnotHost) -> Result { + let subject = authority(host); + self.drain(host, LIST_MEMBERS_NSID, subject, None, 0, Vec::new()) + .await + } + + pub async fn list_collaborators( + &self, + host: &KnotHost, + repo: &Did, + ) -> Result { + self.drain( + host, + LIST_COLLABORATORS_NSID, + repo.as_ref().to_owned(), + None, + 0, + Vec::new(), + ) + .await + } + + fn drain<'a>( + &'a self, + host: &'a KnotHost, + endpoint: &'static str, + subject: String, + cursor: Option, + page: usize, + mut acc: Vec, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + if page >= MAX_LIST_PAGES { + tracing::warn!( + host = %authority(host), + endpoint, + pages = page, + "knot list truncated at page cap" + ); + return Ok(AclListing { + entries: acc, + completeness: Completeness::Truncated, + }); + } + let resp = self + .fetch_page(host, endpoint, &subject, cursor.as_deref()) + .await?; + acc.extend(resp.items); + match resp.cursor.filter(|c| !c.is_empty()) { + Some(next) => { + self.drain(host, endpoint, subject, Some(next), page + 1, acc) + .await + } + None => Ok(AclListing { + entries: acc, + completeness: Completeness::Complete, + }), + } + }) + } + + async fn fetch_page( + &self, + host: &KnotHost, + endpoint: &'static str, + subject: &str, + cursor: Option<&str>, + ) -> Result { + let mut url = host.xrpc_url(&nsid(endpoint)); + { + let mut q = url.query_pairs_mut(); + q.clear(); + q.append_pair("subject", subject); + q.append_pair("limit", &LIST_PAGE_LIMIT.to_string()); + if let Some(c) = cursor { + q.append_pair("cursor", c); + } + } + let bytes = self.get_json(url).await?; + Ok(serde_json::from_slice(&bytes)?) + } + + async fn get_json(&self, url: Url) -> Result { + let resp = self + .http + .execute(HttpRequest { + url, + headers: HeaderMap::new(), + }) + .await?; + match resp.status { + StatusCode::OK => read_bounded(resp).await, + StatusCode::NOT_FOUND => Err(KnotClientError::NotFound), + other => Err(KnotClientError::Upstream(other)), + } + } +} + +#[derive(Deserialize)] +struct VersionWire { + #[serde(default)] + capabilities: Option>, +} + +#[derive(Deserialize)] +struct ListWire { + #[serde(default)] + items: Vec, + #[serde(default)] + cursor: Option, +} + +async fn read_bounded(resp: HttpResponseHead) -> Result { + if resp.content_length.is_some_and(|len| len > MAX_BODY_BYTES) { + return Err(KnotClientError::BodyTooLarge { + limit: MAX_BODY_BYTES, + }); + } + let buf = resp + .body + .map_err(KnotClientError::Network) + .try_fold(BytesMut::new(), |mut acc, chunk| async move { + if (acc.len() as u64).saturating_add(chunk.len() as u64) > MAX_BODY_BYTES { + return Err(KnotClientError::BodyTooLarge { + limit: MAX_BODY_BYTES, + }); + } + acc.extend_from_slice(&chunk); + Ok(acc) + }) + .await?; + Ok(buf.freeze()) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + use wiremock::matchers::{method, path, query_param}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + + fn client() -> KnotClient { + KnotClient::new(ReqwestHttp::shared(default_http_client(true).unwrap())) + } + + fn endpoint(server: &MockServer) -> KnotHost { + KnotHost::parse(&server.uri()).unwrap() + } + + #[tokio::test] + async fn capabilities_returns_declared_tokens() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "version": "1.0.0 (deadbeef)", + "capabilities": ["knot-acl"] + }))) + .mount(&server) + .await; + + let caps = client().capabilities(&endpoint(&server)).await.unwrap(); + assert_eq!(caps, vec!["knot-acl".to_owned()]); + } + + #[tokio::test] + async fn capabilities_empty_when_field_absent() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with( + ResponseTemplate::new(200).set_body_json(json!({ "version": "1.0.0 (cafe)" })), + ) + .mount(&server) + .await; + + let caps = client().capabilities(&endpoint(&server)).await.unwrap(); + assert!(caps.is_empty()); + } + + #[tokio::test] + async fn list_members_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", authority(&host).as_str())) + .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-02T12:00:00Z"} + ] + }))) + .mount(&server) + .await; + + let listing = client().list_members(&host).await.unwrap(); + assert_eq!(listing.completeness, Completeness::Complete); + assert_eq!( + listing.entries, + vec![ + AclEntry { + subject: did("did:plc:boltless"), + created_at: "2026-06-01T00:00:00Z".parse().unwrap(), + }, + AclEntry { + subject: did("did:plc:akshay"), + created_at: "2026-06-02T12:00:00Z".parse().unwrap(), + }, + ] + ); + } + + #[tokio::test] + async fn list_members_drains_multiple_pages() { + let server = MockServer::start().await; + let host = endpoint(&server); + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.listMembers")) + .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"}] + }))) + .mount(&server) + .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"}], + "cursor": "p2" + }))) + .mount(&server) + .await; + + let listing = client().list_members(&host).await.unwrap(); + assert_eq!(listing.completeness, Completeness::Complete); + let subjects: Vec<_> = listing.entries.into_iter().map(|m| m.subject).collect(); + assert_eq!( + subjects, + vec![did("did:plc:boltless"), did("did:plc:akshay")] + ); + } + + #[tokio::test] + async fn list_collaborators_uses_repo_did_subject() { + let server = MockServer::start().await; + let host = endpoint(&server); + Mock::given(method("GET")) + .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:olaren", "addedBy": "did:plc:boltless", "createdAt": "2026-06-03T00:00:00Z"}] + }))) + .mount(&server) + .await; + + let listing = client() + .list_collaborators(&host, &did("did:plc:scallop")) + .await + .unwrap(); + assert_eq!(listing.completeness, Completeness::Complete); + assert_eq!(listing.entries.len(), 1); + assert_eq!(listing.entries[0].subject, did("did:plc:olaren")); + } + + #[tokio::test] + async fn drain_reports_truncation_at_page_cap() { + let server = MockServer::start().await; + let host = endpoint(&server); + 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"}], + "cursor": "more" + }))) + .mount(&server) + .await; + + let listing = client().list_members(&host).await.unwrap(); + assert_eq!( + listing.completeness, + Completeness::Truncated, + "a never-terminating cursor must surface as a truncated listing" + ); + assert_eq!(listing.entries.len(), MAX_LIST_PAGES); + } + + #[tokio::test] + async fn maps_404_to_not_found() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .respond_with(ResponseTemplate::new(404)) + .mount(&server) + .await; + + let err = client() + .capabilities(&endpoint(&server)) + .await + .expect_err("404 must surface"); + assert!(matches!(err, KnotClientError::NotFound)); + } + + #[tokio::test] + async fn maps_5xx_to_upstream() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .respond_with(ResponseTemplate::new(503)) + .mount(&server) + .await; + + let err = client() + .capabilities(&endpoint(&server)) + .await + .expect_err("5xx must surface"); + match err { + KnotClientError::Upstream(s) => assert_eq!(s.as_u16(), 503), + other => panic!("wrong variant: {other:?}"), + } + } + + #[tokio::test] + async fn does_not_follow_redirects() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with(ResponseTemplate::new(301).insert_header( + "location", + "https://kt.tngl.oyster.cafe/xrpc/sh.tangled.knot.version", + )) + .mount(&server) + .await; + + let err = client() + .capabilities(&endpoint(&server)) + .await + .expect_err("a redirect must surface as an error, not be followed to another knot"); + match err { + KnotClientError::Upstream(s) => assert_eq!(s.as_u16(), 301), + other => panic!("expected Upstream(301), got {other:?}"), + } + } + + #[test] + fn knot_endpoint_rejects_private_host() { + let err = + knot_endpoint("127.0.0.1:9", true, false).expect_err("private host must be refused"); + assert!(matches!(err, KnotClientError::PrivateHost { .. })); + } + + #[test] + fn knot_endpoint_allows_private_when_permitted() { + let knot = knot_endpoint("127.0.0.1:9", true, true).expect("private host allowed"); + assert_eq!(knot.url().scheme(), "http"); + } +} diff --git a/bobbin/crates/knot-ingest/src/gate.rs b/bobbin/crates/knot-ingest/src/gate.rs new file mode 100644 index 00000000..9bf6447e --- /dev/null +++ b/bobbin/crates/knot-ingest/src/gate.rs @@ -0,0 +1,255 @@ +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; +use std::sync::Mutex; +use std::time::Duration; + +use bobbin_runtime::Clock; +use bobbin_types::knot_acl::KnotHostKey; +use tokio::time::Instant; + +use crate::client::{KnotClient, knot_endpoint}; + +const KNOT_ACL_CAPABILITY: &str = "knot-acl"; +const LEGACY_REPROBE_INTERVAL: Duration = Duration::from_secs(300); +const ERROR_REPROBE_INTERVAL: Duration = Duration::from_secs(60); + +struct ProbeRecord { + at: Instant, + retry_after: Duration, +} + +pub struct CapabilityGate { + client: KnotClient, + clock: Arc, + dev: bool, + allow_private: bool, + native: Mutex>, + last_probe: Mutex>, +} + +impl CapabilityGate { + pub fn new(client: KnotClient, clock: Arc, dev: bool, allow_private: bool) -> Self { + Self { + client, + clock, + dev, + allow_private, + native: Mutex::new(HashSet::new()), + last_probe: Mutex::new(HashMap::new()), + } + } + + pub fn is_native(&self, host: &KnotHostKey) -> bool { + self.native.lock().unwrap().contains(host) + } + + pub async fn has_knot_acl(&self, host: &KnotHostKey) -> bool { + if self.is_native(host) { + return true; + } + let now = self.clock.now_instant(); + if self.throttled(host, now) { + return false; + } + match self.probe(host).await { + Ok(true) => { + self.native.lock().unwrap().insert(host.clone()); + true + } + Ok(false) => { + self.mark(host, now, LEGACY_REPROBE_INTERVAL); + false + } + Err(err) => { + tracing::warn!(host = %host, error = %err, "knot capability probe failed"); + self.mark(host, now, ERROR_REPROBE_INTERVAL); + false + } + } + } + + async fn probe(&self, host: &KnotHostKey) -> Result { + let endpoint = knot_endpoint(host.as_str(), self.dev, self.allow_private)?; + let caps = self.client.capabilities(&endpoint).await?; + Ok(caps.iter().any(|cap| cap == KNOT_ACL_CAPABILITY)) + } + + fn throttled(&self, host: &KnotHostKey, now: Instant) -> bool { + self.last_probe + .lock() + .unwrap() + .get(host) + .is_some_and(|rec| now.saturating_duration_since(rec.at) < rec.retry_after) + } + + fn mark(&self, host: &KnotHostKey, now: Instant, retry_after: Duration) { + self.last_probe.lock().unwrap().insert( + host.clone(), + ProbeRecord { + at: now, + retry_after, + }, + ); + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::atomic::{AtomicU64, Ordering}; + + use bobbin_runtime::{ReqwestHttp, SleepFuture, UnixMicros}; + use serde_json::json; + use wiremock::matchers::{method, path}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + struct ManualClock { + base: Instant, + offset_micros: AtomicU64, + } + + impl ManualClock { + fn new() -> Self { + Self { + base: Instant::now(), + offset_micros: AtomicU64::new(0), + } + } + + fn advance(&self, by: Duration) { + self.offset_micros + .fetch_add(by.as_micros() as u64, Ordering::SeqCst); + } + } + + impl Clock for ManualClock { + fn now_unix_micros(&self) -> UnixMicros { + UnixMicros::new(self.offset_micros.load(Ordering::SeqCst)) + } + fn now_instant(&self) -> Instant { + self.base + Duration::from_micros(self.offset_micros.load(Ordering::SeqCst)) + } + fn sleep(&self, _: Duration) -> SleepFuture { + Box::pin(async {}) + } + fn sleep_until(&self, _: Instant) -> SleepFuture { + Box::pin(async {}) + } + } + + fn gate(server: &MockServer, clock: Arc) -> (CapabilityGate, KnotHostKey) { + let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); + let url = url::Url::parse(&server.uri()).unwrap(); + let host = format!("{}:{}", url.host_str().unwrap(), url.port().unwrap()); + ( + CapabilityGate::new(client, clock, true, true), + KnotHostKey::new(&host), + ) + } + + async fn mount_version(server: &MockServer, caps: serde_json::Value, expect: u64) { + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with( + ResponseTemplate::new(200) + .set_body_json(json!({ "version": "1.0.0 (cafe)", "capabilities": caps })), + ) + .expect(expect) + .mount(server) + .await; + } + + #[tokio::test] + async fn declares_knot_acl() { + let server = MockServer::start().await; + mount_version(&server, json!(["knot-acl"]), 1).await; + let (gate, host) = gate(&server, Arc::new(ManualClock::new())); + assert!(gate.has_knot_acl(&host).await); + assert!(gate.is_native(&host)); + } + + #[tokio::test] + async fn legacy_knot_without_capability() { + let server = MockServer::start().await; + mount_version(&server, json!([]), 1).await; + let (gate, host) = gate(&server, Arc::new(ManualClock::new())); + assert!(!gate.has_knot_acl(&host).await); + assert!(!gate.is_native(&host)); + } + + #[tokio::test] + async fn native_is_latched_and_survives_probe_error() { + let server = MockServer::start().await; + mount_version(&server, json!(["knot-acl"]), 1).await; + let clock = Arc::new(ManualClock::new()); + let (gate, host) = gate(&server, clock.clone()); + assert!(gate.has_knot_acl(&host).await); + + server.reset().await; + clock.advance(LEGACY_REPROBE_INTERVAL + Duration::from_secs(1)); + assert!( + gate.has_knot_acl(&host).await, + "latched native never re-probes" + ); + assert!(gate.is_native(&host)); + } + + #[tokio::test] + async fn legacy_throttled_then_reprobed_after_interval() { + let server = MockServer::start().await; + mount_version(&server, json!([]), 2).await; + let clock = Arc::new(ManualClock::new()); + let (gate, host) = gate(&server, clock.clone()); + assert!(!gate.has_knot_acl(&host).await); + clock.advance(Duration::from_secs(60)); + assert!(!gate.has_knot_acl(&host).await); + clock.advance(LEGACY_REPROBE_INTERVAL); + assert!(!gate.has_knot_acl(&host).await); + } + + #[tokio::test] + async fn legacy_upgrade_is_detected_on_reprobe() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "version": "1.0.0" }))) + .up_to_n_times(1) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with( + ResponseTemplate::new(200) + .set_body_json(json!({ "version": "1.1.0", "capabilities": ["knot-acl"] })), + ) + .mount(&server) + .await; + let clock = Arc::new(ManualClock::new()); + let (gate, host) = gate(&server, clock.clone()); + assert!(!gate.has_knot_acl(&host).await); + clock.advance(LEGACY_REPROBE_INTERVAL + Duration::from_secs(1)); + assert!(gate.has_knot_acl(&host).await); + assert!(gate.is_native(&host)); + } + + #[tokio::test] + async fn probe_error_throttled_briefly_then_reprobed() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.knot.version")) + .respond_with(ResponseTemplate::new(503)) + .expect(2) + .mount(&server) + .await; + let clock = Arc::new(ManualClock::new()); + let (gate, host) = gate(&server, clock.clone()); + assert!(!gate.has_knot_acl(&host).await); + clock.advance(Duration::from_secs(1)); + assert!( + !gate.has_knot_acl(&host).await, + "error reprobe is throttled" + ); + clock.advance(ERROR_REPROBE_INTERVAL); + assert!(!gate.has_knot_acl(&host).await); + } +} diff --git a/bobbin/crates/knot-ingest/src/lib.rs b/bobbin/crates/knot-ingest/src/lib.rs new file mode 100644 index 00000000..0bc6a1af --- /dev/null +++ b/bobbin/crates/knot-ingest/src/lib.rs @@ -0,0 +1,13 @@ +pub mod client; +pub mod gate; +pub mod orchestrator; +pub mod registry; +pub mod roster; +pub mod stream; + +pub use client::{AclEntry, AclListing, Completeness, KnotClient, KnotClientError, knot_endpoint}; +pub use gate::CapabilityGate; +pub use orchestrator::Orchestrator; +pub use registry::KnotRegistry; +pub use roster::{AclOp, Cursor, Roster}; +pub use stream::{StreamConfig, run_stream}; diff --git a/bobbin/crates/knot-ingest/src/orchestrator.rs b/bobbin/crates/knot-ingest/src/orchestrator.rs new file mode 100644 index 00000000..2144bb7b --- /dev/null +++ b/bobbin/crates/knot-ingest/src/orchestrator.rs @@ -0,0 +1,463 @@ +use std::collections::{HashMap, HashSet}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use bobbin_edge_index::EdgeStore; +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; +use tokio_util::sync::CancellationToken; + +use crate::client::{AclListing, Completeness, KnotClient, knot_endpoint}; +use crate::gate::CapabilityGate; +use crate::registry::KnotRegistry; +use crate::roster::{AclOp, Cursor, Roster}; +use crate::stream::{StreamConfig, run_stream}; + +const POLL_INTERVAL: Duration = Duration::from_secs(30); +const RECONCILE_INTERVAL: Duration = Duration::from_secs(300); +const PROBE_CONCURRENCY: usize = 16; + +pub struct Orchestrator { + pub client: Arc, + pub gate: Arc, + pub registry: Arc, + pub store: Arc, + pub ws: Arc, + pub clock: Arc, + pub dev: bool, + pub allow_private: bool, + pub cancel: CancellationToken, +} + +impl Orchestrator { + pub async fn run(self) { + let mut subscribed: HashMap = HashMap::new(); + let mut unspawnable: HashSet = HashSet::new(); + loop { + if self.cancel.is_cancelled() { + break; + } + self.discover(&mut subscribed, &mut unspawnable).await; + tokio::select! { + _ = self.cancel.cancelled() => break, + _ = self.clock.sleep(POLL_INTERVAL) => {} + } + } + } + + async fn discover( + &self, + subscribed: &mut HashMap, + unspawnable: &mut HashSet, + ) { + let candidates: Vec = self + .registry + .hosts() + .into_iter() + .filter(|host| !subscribed.contains_key(host) && !unspawnable.contains(host)) + .collect(); + let approved: Vec = futures::stream::iter(candidates) + .map(|host| async move { self.gate.has_knot_acl(&host).await.then_some(host) }) + .buffer_unordered(PROBE_CONCURRENCY) + .filter_map(|approved| async move { approved }) + .collect() + .await; + approved + .into_iter() + .for_each(|host| match self.spawn(&host) { + Some(token) => { + subscribed.insert(host, token); + } + None => { + tracing::warn!(host = %host, "skipping unspawnable knot endpoint"); + unspawnable.insert(host); + } + }); + } + + fn spawn(&self, host: &KnotHostKey) -> Option { + let endpoint = knot_endpoint(host.as_str(), self.dev, self.allow_private).ok()?; + let knot_did = host_to_knot_did(host.as_str())?; + let roster = Arc::new(Mutex::new(Roster::new( + self.store.clone(), + knot_did, + self.registry.clone(), + host.clone(), + ))); + let token = self.cancel.child_token(); + + let stream_cfg = StreamConfig { + ws: self.ws.clone(), + clock: self.clock.clone(), + cancel: token.clone(), + }; + let stream_roster = roster.clone(); + let stream_endpoint = endpoint.clone(); + tokio::spawn(async move { + run_stream(&stream_cfg, &stream_endpoint, &stream_roster, 0).await; + }); + + let client = self.client.clone(); + let registry = self.registry.clone(); + let clock = self.clock.clone(); + let host_owned = host.clone(); + let reconcile_token = token.clone(); + tokio::spawn(async move { + reconcile_loop( + &client, + &endpoint, + &host_owned, + ®istry, + &roster, + &*clock, + &reconcile_token, + ) + .await; + }); + + Some(token) + } +} + +async fn reconcile_loop( + client: &KnotClient, + endpoint: &KnotHost, + host: &KnotHostKey, + registry: &KnotRegistry, + roster: &Mutex, + clock: &dyn Clock, + cancel: &CancellationToken, +) { + loop { + if cancel.is_cancelled() { + return; + } + reconcile_once(client, endpoint, host, registry, roster).await; + tokio::select! { + _ = cancel.cancelled() => return, + _ = clock.sleep(RECONCILE_INTERVAL) => {} + } + } +} + +async fn reconcile_once( + client: &KnotClient, + endpoint: &KnotHost, + host: &KnotHostKey, + registry: &KnotRegistry, + roster: &Mutex, +) { + let horizon = roster.lock().unwrap().max_cursor(); + + match client.list_members(endpoint).await { + Ok(AclListing { + entries, + completeness, + }) => { + let present: HashSet> = + entries.iter().map(|entry| entry.subject.clone()).collect(); + let mut guard = roster.lock().unwrap(); + entries.into_iter().for_each(|entry| { + guard.apply_member(AclOp::Add, entry.subject, Cursor(nanos(entry.created_at))); + }); + if completeness == Completeness::Complete { + guard.reap_members(&present, horizon); + } + } + Err(err) => tracing::warn!(host = %host, error = %err, "knot member reconcile failed"), + } + + futures::stream::iter(registry.repos(host)) + .for_each(|repo| async move { + match client.list_collaborators(endpoint, &repo).await { + Ok(AclListing { + entries, + completeness, + }) => { + let present: HashSet> = + entries.iter().map(|entry| entry.subject.clone()).collect(); + let mut guard = roster.lock().unwrap(); + entries.into_iter().for_each(|entry| { + guard.apply_collaborator( + AclOp::Add, + repo.clone(), + entry.subject, + Cursor(nanos(entry.created_at)), + ); + }); + if completeness == Completeness::Complete { + guard.reap_collaborators(&repo, &present, horizon); + } + } + Err(err) => { + tracing::warn!(host = %host, repo = %repo.as_ref(), error = %err, "knot collaborator reconcile failed") + } + } + }) + .await; + + roster.lock().unwrap().purge_legacy(); +} + +fn nanos(timestamp: DateTime) -> i64 { + timestamp.timestamp_nanos_opt().unwrap_or(0) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::client::authority; + use bobbin_runtime::{ReqwestHttp, RuntimeHasher}; + use bobbin_types::edges::Edge; + use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; + use jacquard_common::types::string::AtUri; + use serde_json::json; + use wiremock::matchers::{method, path, query_param}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + fn did(s: &str) -> Did { + 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"), + SubjectRef::Did(did(repo)), + )) + } + + #[tokio::test] + async fn reconcile_backfills_members_and_collaborators() { + let server = MockServer::start().await; + let endpoint = KnotHost::parse(&server.uri()).unwrap(); + let host = KnotHostKey::new(&authority(&endpoint)); + + 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(&server) + .await; + Mock::given(method("GET")) + .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:olaren", "addedBy": "did:plc:boltless", "createdAt": "2026-06-03T00:00:00Z"} + ] + }))) + .mount(&server) + .await; + + let registry = Arc::new(KnotRegistry::new()); + registry.observe_repo(&host, did("did:plc:scallop")); + + let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); + let knot_did = host_to_knot_did(host.as_str()).unwrap(); + let roster = Mutex::new(Roster::new( + store.clone(), + knot_did, + registry.clone(), + host.clone(), + )); + let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); + + reconcile_once(&client, &endpoint, &host, ®istry, &roster).await; + + assert_eq!(member_count(&store, "did:plc:boltless"), 1); + assert_eq!(member_count(&store, "did:plc:akshay"), 1); + assert_eq!(collaborator_count(&store, "did:plc:scallop"), 1); + } + + #[tokio::test] + async fn reconcile_reaps_departed_member() { + let server = MockServer::start().await; + let endpoint = KnotHost::parse(&server.uri()).unwrap(); + let host = KnotHostKey::new(&authority(&endpoint)); + + let registry = Arc::new(KnotRegistry::new()); + let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); + let knot_did = host_to_knot_did(host.as_str()).unwrap(); + let roster = Mutex::new(Roster::new( + store.clone(), + knot_did, + registry.clone(), + host.clone(), + )); + let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); + + 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(&server) + .await; + reconcile_once(&client, &endpoint, &host, ®istry, &roster).await; + assert_eq!(member_count(&store, "did:plc:boltless"), 1); + assert_eq!(member_count(&store, "did:plc:akshay"), 1); + + server.reset().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:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"} + ] + }))) + .mount(&server) + .await; + reconcile_once(&client, &endpoint, &host, ®istry, &roster).await; + + assert_eq!( + member_count(&store, "did:plc:boltless"), + 0, + "a member dropped from the authoritative snapshot is reaped on reconcile" + ); + assert_eq!(member_count(&store, "did:plc:akshay"), 1); + } + + #[tokio::test] + async fn reconcile_skips_reap_when_member_list_truncated() { + let server = MockServer::start().await; + let endpoint = KnotHost::parse(&server.uri()).unwrap(); + let host = KnotHostKey::new(&authority(&endpoint)); + + 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"}], + "cursor": "more" + }))) + .mount(&server) + .await; + + let registry = Arc::new(KnotRegistry::new()); + let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); + let knot_did = host_to_knot_did(host.as_str()).unwrap(); + let roster = Mutex::new(Roster::new( + store.clone(), + knot_did, + registry.clone(), + host.clone(), + )); + let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); + + let stayed = did("did:plc:akshay"); + roster + .lock() + .unwrap() + .apply_member(AclOp::Add, stayed, Cursor(1_000_000)); + assert_eq!(member_count(&store, "did:plc:akshay"), 1); + + reconcile_once(&client, &endpoint, &host, ®istry, &roster).await; + + assert_eq!( + member_count(&store, "did:plc:akshay"), + 1, + "a truncated member snapshot must not reap members it could not enumerate" + ); + assert_eq!(member_count(&store, "did:plc:boltless"), 1); + } + + fn seed_legacy_edge(store: &EdgeStore, kind: &'static str, subject: &str, source: &str) { + let source = AtUri::new_owned(source).unwrap(); + store.upsert_source( + &source, + vec![Edge { + kind: nsid_static(kind), + subject: SubjectRef::Did(did(subject)), + source: source.clone(), + sort_micros: 0, + }], + ); + } + + #[tokio::test] + async fn reconcile_purges_seeded_legacy_acl() { + let server = MockServer::start().await; + let endpoint = KnotHost::parse(&server.uri()).unwrap(); + let host = KnotHostKey::new(&authority(&endpoint)); + + 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(&server) + .await; + Mock::given(method("GET")) + .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:olaren", "addedBy": "did:plc:boltless", "createdAt": "2026-06-03T00:00:00Z"} + ] + }))) + .mount(&server) + .await; + + let registry = Arc::new(KnotRegistry::new()); + let repo = did("did:plc:scallop"); + registry.observe_repo(&host, repo.clone()); + + let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); + let member_source = "at://did:plc:akshay/sh.tangled.knot.member/r1"; + seed_legacy_edge( + &store, + "sh.tangled.knot.member", + "did:plc:boltless", + member_source, + ); + registry.note_legacy_member(AtUri::new_owned(member_source).unwrap(), &host); + seed_legacy_edge( + &store, + "sh.tangled.repo.collaborator", + "did:plc:scallop", + "at://did:plc:akshay/sh.tangled.repo.collaborator/r2", + ); + + let knot_did = host_to_knot_did(host.as_str()).unwrap(); + let roster = Mutex::new(Roster::new( + store.clone(), + knot_did, + registry.clone(), + host.clone(), + )); + let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); + + reconcile_once(&client, &endpoint, &host, ®istry, &roster).await; + + assert_eq!( + member_count(&store, "did:plc:boltless"), + 1, + "stale legacy member purged, knot-owned member kept" + ); + assert_eq!( + collaborator_count(&store, "did:plc:scallop"), + 1, + "stale legacy collaborator purged, knot-owned collaborator kept" + ); + } +} diff --git a/bobbin/crates/knot-ingest/src/registry.rs b/bobbin/crates/knot-ingest/src/registry.rs new file mode 100644 index 00000000..67f731f3 --- /dev/null +++ b/bobbin/crates/knot-ingest/src/registry.rs @@ -0,0 +1,197 @@ +use std::collections::{HashMap, HashSet}; +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)] +pub struct KnotRegistry { + inner: Mutex, +} + +impl KnotRegistry { + pub fn new() -> Self { + Self::default() + } + + pub fn observe_host(&self, host: &KnotHostKey) { + self.inner.lock().unwrap().hosts.insert(host.clone()); + } + + pub fn observe_repo(&self, host: &KnotHostKey, repo: Did) { + let mut inner = self.inner.lock().unwrap(); + inner.hosts.insert(host.clone()); + inner + .repos + .entry(host.clone()) + .or_default() + .insert(repo.clone()); + inner.repo_host.insert(repo, host.clone()); + } + + pub fn hosts(&self) -> Vec { + self.inner.lock().unwrap().hosts.iter().cloned().collect() + } + + pub fn repos(&self, host: &KnotHostKey) -> Vec> { + self.inner + .lock() + .unwrap() + .repos + .get(host) + .map(|set| set.iter().cloned().collect()) + .unwrap_or_default() + } + + pub fn repo_on_host(&self, host: &KnotHostKey, repo: &Did) -> bool { + self.inner + .lock() + .unwrap() + .repos + .get(host) + .is_some_and(|set| set.contains(repo)) + } + + pub fn host_of_repo(&self, repo: &Did) -> Option { + 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)] +mod tests { + use super::*; + + fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + + fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() + } + + fn host(s: &str) -> KnotHostKey { + KnotHostKey::new(s) + } + + #[test] + fn observe_repo_registers_host_and_repo() { + let registry = KnotRegistry::new(); + registry.observe_repo(&host("oyster.cafe"), did("did:plc:scallop")); + registry.observe_repo(&host("oyster.cafe"), did("did:plc:limpet")); + + assert_eq!(registry.hosts(), vec![host("oyster.cafe")]); + let mut repos = registry.repos(&host("oyster.cafe")); + repos.sort_by(|a, b| a.as_ref().cmp(b.as_ref())); + assert_eq!(repos, vec![did("did:plc:limpet"), did("did:plc:scallop")]); + } + + #[test] + fn observe_repo_dedups() { + let registry = KnotRegistry::new(); + registry.observe_repo(&host("nel.pet"), did("did:plc:whelk")); + registry.observe_repo(&host("nel.pet"), did("did:plc:whelk")); + assert_eq!(registry.repos(&host("nel.pet")), vec![did("did:plc:whelk")]); + } + + #[test] + fn observe_host_without_repos() { + let registry = KnotRegistry::new(); + registry.observe_host(&host("oyster.cafe")); + assert_eq!(registry.hosts(), vec![host("oyster.cafe")]); + assert!(registry.repos(&host("oyster.cafe")).is_empty()); + } + + #[test] + fn host_lookups_are_case_insensitive() { + let registry = KnotRegistry::new(); + registry.observe_repo(&host("KT.Oyster.Cafe"), did("did:plc:scallop")); + assert!(registry.repo_on_host(&host("kt.oyster.cafe"), &did("did:plc:scallop"))); + assert_eq!( + registry.host_of_repo(&did("did:plc:scallop")), + Some(host("kt.oyster.cafe")) + ); + } + + #[test] + fn host_of_repo_resolves_owning_knot() { + let registry = KnotRegistry::new(); + registry.observe_repo(&host("oyster.cafe"), did("did:plc:scallop")); + assert_eq!( + registry.host_of_repo(&did("did:plc:scallop")), + Some(host("oyster.cafe")) + ); + 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 new file mode 100644 index 00000000..72d8ce5c --- /dev/null +++ b/bobbin/crates/knot-ingest/src/roster.rs @@ -0,0 +1,445 @@ +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; + +use bobbin_edge_index::EdgeStore; +use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; +use bobbin_types::knot_acl::{self, KnotHostKey}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use serde::Deserialize; + +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); + +impl Cursor { + fn micros(self) -> u64 { + self.0.max(0) as u64 / 1000 + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum AclOp { + Add, + Remove, +} + +#[derive(Clone, Eq, Hash, PartialEq)] +enum DedupKey { + Member(Did), + Collaborator(Did, Did), +} + +struct SeenState { + cursor: Cursor, + present: bool, +} + +pub struct Roster { + store: Arc, + knot: Did, + registry: Arc, + host: KnotHostKey, + seen: HashMap, +} + +impl Roster { + pub fn new( + store: Arc, + knot: Did, + registry: Arc, + host: KnotHostKey, + ) -> Self { + Self { + store, + knot, + registry, + host, + seen: HashMap::new(), + } + } + + pub fn apply_member(&mut self, op: AclOp, subject: Did, cursor: Cursor) { + if !self.advance( + DedupKey::Member(subject.clone()), + cursor, + matches!(op, AclOp::Add), + ) { + return; + } + match op { + AclOp::Add => { + if let Some((source, edges)) = + knot_acl::member_upsert(&self.knot, &subject, cursor.micros()) + { + self.store.upsert_source(&source, edges); + } + } + AclOp::Remove => { + if let Some(source) = knot_acl::member_source(&self.knot, &subject) { + self.store.remove_source(&source); + } + } + } + } + + pub fn apply_collaborator( + &mut self, + op: AclOp, + repo: Did, + subject: Did, + cursor: Cursor, + ) { + if !self.registry.repo_on_host(&self.host, &repo) { + return; + } + if !self.advance( + DedupKey::Collaborator(repo.clone(), subject.clone()), + cursor, + matches!(op, AclOp::Add), + ) { + return; + } + match op { + AclOp::Add => { + if let Some((source, edges)) = + knot_acl::collaborator_upsert(&repo, &subject, cursor.micros()) + { + self.store.upsert_source(&source, edges); + } + } + AclOp::Remove => { + if let Some(source) = knot_acl::collaborator_source(&repo, &subject) { + self.store.remove_source(&source); + } + } + } + } + + pub fn max_cursor(&self) -> Cursor { + self.seen + .values() + .map(|state| state.cursor) + .max() + .unwrap_or(Cursor(0)) + } + + pub fn reap_members(&mut self, present: &HashSet>, horizon: Cursor) { + let stale: Vec> = self + .seen + .iter() + .filter_map(|(key, state)| match key { + DedupKey::Member(subject) + if state.present && state.cursor <= horizon && !present.contains(subject) => + { + Some(subject.clone()) + } + _ => None, + }) + .collect(); + stale + .into_iter() + .for_each(|subject| self.retire_member(subject)); + } + + pub fn reap_collaborators( + &mut self, + repo: &Did, + present: &HashSet>, + horizon: Cursor, + ) { + let stale: Vec> = self + .seen + .iter() + .filter_map(|(key, state)| match key { + DedupKey::Collaborator(edge_repo, subject) + if edge_repo == repo + && state.present + && state.cursor <= horizon + && !present.contains(subject) => + { + Some(subject.clone()) + } + _ => None, + }) + .collect(); + stale + .into_iter() + .for_each(|subject| self.retire_collaborator(repo.clone(), subject)); + } + + 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() + .for_each(|repo| { + let key = EdgeKey::new(nsid_static(REPO_COLLABORATOR_KIND), SubjectRef::Did(repo)); + self.store + .sources_for(&key) + .into_iter() + .filter(|source| knot_acl::decode_knot_owned_source(source).is_none()) + .for_each(|source| self.store.remove_source(&source)); + }); + } + + fn retire_member(&mut self, subject: Did) { + if let Some(source) = knot_acl::member_source(&self.knot, &subject) { + self.store.remove_source(&source); + } + if let Some(state) = self.seen.get_mut(&DedupKey::Member(subject)) { + state.present = false; + } + } + + fn retire_collaborator(&mut self, repo: Did, subject: Did) { + if let Some(source) = knot_acl::collaborator_source(&repo, &subject) { + self.store.remove_source(&source); + } + if let Some(state) = self.seen.get_mut(&DedupKey::Collaborator(repo, subject)) { + state.present = false; + } + } + + 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 + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use bobbin_runtime::RuntimeHasher; + use bobbin_types::edges::Edge; + use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; + use jacquard_common::types::string::AtUri; + + fn store() -> Arc { + Arc::new(EdgeStore::new(RuntimeHasher::default())) + } + + fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + + fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() + } + + fn host() -> KnotHostKey { + KnotHostKey::new("oyster.cafe") + } + + fn add_legacy_edge( + store: &EdgeStore, + kind: &'static str, + subject: &Did, + source: &AtUri, + ) { + store.upsert_source( + source, + vec![Edge { + kind: nsid_static(kind), + subject: SubjectRef::Did(subject.clone()), + source: source.clone(), + sort_micros: 0, + }], + ); + } + + fn knot() -> Did { + 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"), + SubjectRef::Did(repo.clone()), + )) + } + + fn empty_registry() -> Arc { + Arc::new(KnotRegistry::new()) + } + + fn member_roster(store: Arc) -> Roster { + Roster::new(store, knot(), empty_registry(), host()) + } + + fn subjects(items: &[&Did]) -> HashSet> { + items.iter().map(|d| (*d).clone()).collect() + } + + #[test] + fn member_add_then_remove() { + let store = store(); + let mut roster = member_roster(store.clone()); + let m = did("did:plc:boltless"); + roster.apply_member(AclOp::Add, m.clone(), Cursor(1_000_000)); + assert_eq!(member_count(&store, &m), 1); + roster.apply_member(AclOp::Remove, m.clone(), Cursor(2_000_000)); + assert_eq!(member_count(&store, &m), 0); + } + + #[test] + fn stale_add_cannot_resurrect_removed_member() { + let store = store(); + let mut roster = member_roster(store.clone()); + let m = did("did:plc:boltless"); + roster.apply_member(AclOp::Remove, m.clone(), Cursor(5)); + roster.apply_member(AclOp::Add, m.clone(), Cursor(1)); + assert_eq!(member_count(&store, &m), 0); + } + + #[test] + fn duplicate_cursor_is_idempotent() { + let store = store(); + let mut roster = member_roster(store.clone()); + let m = did("did:plc:akshay"); + roster.apply_member(AclOp::Add, m.clone(), Cursor(10)); + roster.apply_member(AclOp::Add, m.clone(), Cursor(10)); + assert_eq!(member_count(&store, &m), 1); + } + + #[test] + fn collaborator_keyed_on_repo_when_hosted() { + 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(AclOp::Add, repo.clone(), subject.clone(), Cursor(7)); + assert_eq!(collaborator_count(&store, &repo), 1); + roster.apply_collaborator(AclOp::Remove, repo.clone(), subject.clone(), Cursor(8)); + assert_eq!(collaborator_count(&store, &repo), 0); + } + + #[test] + fn reap_removes_departed_member() { + let store = store(); + let mut roster = member_roster(store.clone()); + let stayed = did("did:plc:akshay"); + let left = did("did:plc:boltless"); + roster.apply_member(AclOp::Add, stayed.clone(), Cursor(10)); + roster.apply_member(AclOp::Add, left.clone(), Cursor(20)); + + let horizon = roster.max_cursor(); + roster.reap_members(&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" + ); + } + + #[test] + fn reap_skips_member_added_after_horizon() { + let store = store(); + let mut roster = member_roster(store.clone()); + let m = did("did:plc:boltless"); + let horizon = roster.max_cursor(); + roster.apply_member(AclOp::Add, m.clone(), Cursor(100)); + + roster.reap_members(&subjects(&[]), horizon); + + assert_eq!( + member_count(&store, &m), + 1, + "a member 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 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(AclOp::Add, repo.clone(), left.clone(), Cursor(7)); + + let horizon = roster.max_cursor(); + roster.reap_collaborators(&repo, &subjects(&[]), horizon); + + assert_eq!(collaborator_count(&store, &repo), 0); + } + + #[test] + 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()); + + roster.apply_collaborator(AclOp::Add, repo.clone(), did("did:plc:olaren"), Cursor(7)); + add_legacy_edge( + &store, + "sh.tangled.repo.collaborator", + &repo, + &at("at://did:plc:akshay/sh.tangled.repo.collaborator/r1"), + ); + assert_eq!(collaborator_count(&store, &repo), 2); + + roster.purge_legacy(); + assert_eq!( + collaborator_count(&store, &repo), + 1, + "only the knot-owned collaborator survives the purge" + ); + } + + #[test] + 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 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 = Roster::new(store.clone(), knot(), empty_registry(), host()); + roster.apply_collaborator(AclOp::Add, repo.clone(), did("did:plc:olaren"), Cursor(7)); + assert_eq!( + collaborator_count(&store, &repo), + 0, + "a knot cannot assert collaborators on a repo it does not host" + ); + } +} diff --git a/bobbin/crates/knot-ingest/src/stream.rs b/bobbin/crates/knot-ingest/src/stream.rs new file mode 100644 index 00000000..f1be6f85 --- /dev/null +++ b/bobbin/crates/knot-ingest/src/stream.rs @@ -0,0 +1,446 @@ +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use bobbin_knot_proxy::KnotHost; +use bobbin_runtime::{Clock, WsConn, WsMessage, WsTransport}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use serde::Deserialize; +use serde_json::value::RawValue; +use tokio_util::sync::CancellationToken; +use url::Url; + +use crate::client::authority; +use crate::roster::{AclOp, Cursor, Roster}; + +const KNOT_MEMBER_UPDATE_NSID: &str = "sh.tangled.knot.memberUpdate"; +const REPO_COLLABORATOR_UPDATE_NSID: &str = "sh.tangled.repo.collaboratorUpdate"; +const RECONNECT_INITIAL: Duration = Duration::from_secs(1); +const RECONNECT_MAX: Duration = Duration::from_secs(60); +const HEALTHY_SESSION_MIN: Duration = Duration::from_secs(15); + +#[derive(Clone)] +pub struct StreamConfig { + pub ws: Arc, + pub clock: Arc, + pub cancel: CancellationToken, +} + +enum SessionEnd { + Cancelled, + Closed { progressed: bool }, + ConnectFailed, +} + +pub async fn run_stream( + cfg: &StreamConfig, + host: &KnotHost, + roster: &Mutex, + initial_cursor: i64, +) { + let mut cursor = initial_cursor; + let mut backoff = RECONNECT_INITIAL; + loop { + if cfg.cancel.is_cancelled() { + return; + } + let started = cfg.clock.now_instant(); + let end = run_session(cfg, host, roster, &mut cursor).await; + if matches!(&end, SessionEnd::Cancelled) { + return; + } + let elapsed = cfg.clock.now_instant().saturating_duration_since(started); + if matches!(&end, SessionEnd::Closed { progressed: false }) && cursor != 0 { + cursor = 0; + } + if session_was_healthy(&end, elapsed) { + backoff = RECONNECT_INITIAL; + } + let delay = jitter_delay(backoff, cfg.clock.now_unix_micros().raw()); + tokio::select! { + _ = cfg.cancel.cancelled() => return, + _ = cfg.clock.sleep(delay) => {} + } + backoff = (backoff * 2).min(RECONNECT_MAX); + } +} + +fn session_was_healthy(end: &SessionEnd, elapsed: Duration) -> bool { + matches!(end, SessionEnd::Closed { progressed: true }) && elapsed >= HEALTHY_SESSION_MIN +} + +fn jitter_delay(base: Duration, entropy: u64) -> Duration { + let frac = (entropy % 1024) as f64 / 1024.0; + base.mul_f64(0.5 + 0.5 * frac) +} + +async fn run_session( + cfg: &StreamConfig, + host: &KnotHost, + roster: &Mutex, + cursor: &mut i64, +) -> SessionEnd { + let Some(url) = events_url(host, *cursor) else { + return SessionEnd::ConnectFailed; + }; + let conn = tokio::select! { + _ = cfg.cancel.cancelled() => return SessionEnd::Cancelled, + res = cfg.ws.connect(url) => match res { + Ok(conn) => conn, + Err(err) => { + tracing::warn!(host = %authority(host), error = %err, "knot eventstream connect failed"); + return SessionEnd::ConnectFailed; + } + }, + }; + let WsConn { + mut sink, + mut stream, + } = conn; + let mut progressed = false; + loop { + let message = tokio::select! { + _ = cfg.cancel.cancelled() => return SessionEnd::Cancelled, + message = stream.next() => message, + }; + match message { + None => return SessionEnd::Closed { progressed }, + Some(Ok(WsMessage::Text(text))) => { + progressed = true; + process_frame(&text, roster, cursor); + } + Some(Ok(WsMessage::Ping(payload))) => { + progressed = true; + let _ = sink.send(WsMessage::Pong(payload)).await; + } + Some(Ok(WsMessage::Close { .. })) => return SessionEnd::Closed { progressed }, + Some(Ok(_)) => {} + Some(Err(err)) => { + tracing::warn!(host = %authority(host), error = %err, "knot eventstream read error"); + return SessionEnd::Closed { progressed }; + } + } + } +} + +fn process_frame(text: &str, roster: &Mutex, cursor: &mut i64) { + let Ok(frame) = serde_json::from_str::(text) else { + return; + }; + *cursor = frame.created; + match frame.nsid.as_str() { + KNOT_MEMBER_UPDATE_NSID => { + if let Ok(update) = serde_json::from_str::(frame.event.get()) { + roster.lock().unwrap().apply_member( + update.op, + update.subject, + Cursor(frame.created), + ); + } + } + REPO_COLLABORATOR_UPDATE_NSID => { + if let Ok(update) = serde_json::from_str::(frame.event.get()) { + roster.lock().unwrap().apply_collaborator( + update.op, + update.repo, + update.subject, + Cursor(frame.created), + ); + } + } + _ => {} + } +} + +fn events_url(host: &KnotHost, cursor: i64) -> Option { + let scheme = if host.url().scheme() == "https" { + "wss" + } else { + "ws" + }; + let base = format!("{scheme}://{}/events", authority(host)); + let full = if cursor != 0 { + format!("{base}?cursor={cursor}") + } else { + base + }; + Url::parse(&full).ok() +} + +#[derive(Deserialize)] +struct FrameWire { + nsid: String, + event: Box, + created: i64, +} + +#[derive(Deserialize)] +struct MemberUpdate { + op: AclOp, + subject: Did, +} + +#[derive(Deserialize)] +struct CollaboratorUpdate { + op: AclOp, + subject: Did, + repo: Did, +} + +#[cfg(test)] +mod tests { + use super::*; + use std::collections::VecDeque; + use std::sync::Mutex; + + use bobbin_edge_index::EdgeStore; + use bobbin_runtime::{ + NetworkError, SystemClock, WsConnectFuture, WsMessageFuture, WsSendFuture, WsSink, WsStream, + }; + use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; + use bobbin_types::knot_acl; + use bytes::Bytes; + + use crate::registry::KnotRegistry; + + struct ScriptStream { + msgs: VecDeque, + } + impl WsStream for ScriptStream { + fn next<'a>(&'a mut self) -> WsMessageFuture<'a> { + Box::pin(async move { self.msgs.pop_front().map(Ok) }) + } + } + + struct RecordSink { + sent: Arc>>, + } + impl WsSink for RecordSink { + fn send<'a>(&'a mut self, message: WsMessage) -> WsSendFuture<'a> { + let sent = self.sent.clone(); + Box::pin(async move { + sent.lock().unwrap().push(message); + Ok(()) + }) + } + } + + struct ScriptWs { + msgs: Mutex>>, + sent: Arc>>, + fail: bool, + } + impl WsTransport for ScriptWs { + fn connect(&self, _url: Url) -> WsConnectFuture { + if self.fail { + return Box::pin(async { + Err(NetworkError::Connect("scripted failure".to_owned())) + }); + } + let msgs = self.msgs.lock().unwrap().take().unwrap_or_default(); + let sent = self.sent.clone(); + Box::pin(async move { + Ok(WsConn { + sink: Box::new(RecordSink { sent }), + stream: Box::new(ScriptStream { msgs }), + }) + }) + } + } + + fn member_frame(op: &str, subject: &str, created: i64) -> String { + format!( + r#"{{"rkey":"r{created}","nsid":"sh.tangled.knot.memberUpdate","event":{{"op":"{op}","subject":"{subject}"}},"created":{created}}}"# + ) + } + + fn collab_frame(op: &str, subject: &str, repo: &str, created: i64) -> String { + format!( + r#"{{"rkey":"r{created}","nsid":"sh.tangled.repo.collaboratorUpdate","event":{{"op":"{op}","subject":"{subject}","repo":"{repo}"}},"created":{created}}}"# + ) + } + + fn did(s: &str) -> Did { + 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"), + SubjectRef::Did(did(repo)), + )) + } + + fn cfg(ws: Arc, cancel: CancellationToken) -> StreamConfig { + StreamConfig { + ws, + clock: Arc::new(SystemClock::new()), + cancel, + } + } + + #[tokio::test] + async fn session_dispatches_deltas_pongs_and_advances_cursor() { + let store = Arc::new(EdgeStore::new(bobbin_runtime::RuntimeHasher::default())); + let knot = knot_acl::host_to_knot_did("oyster.cafe").unwrap(); + let registry = Arc::new(KnotRegistry::new()); + registry.observe_repo( + &knot_acl::KnotHostKey::new("oyster.cafe"), + Did::new_owned("did:plc:scallop").unwrap(), + ); + let roster = Mutex::new(Roster::new( + store.clone(), + knot, + registry, + knot_acl::KnotHostKey::new("oyster.cafe"), + )); + let frames = VecDeque::from(vec![ + WsMessage::Text(member_frame("add", "did:plc:boltless", 100)), + WsMessage::Text(collab_frame( + "add", + "did:plc:olaren", + "did:plc:scallop", + 200, + )), + WsMessage::Ping(Bytes::from_static(b"ka")), + WsMessage::Text(member_frame("remove", "did:plc:boltless", 300)), + WsMessage::Close { + code: 1000, + reason: String::new(), + }, + ]); + let sent = Arc::new(Mutex::new(Vec::new())); + let ws: Arc = Arc::new(ScriptWs { + msgs: Mutex::new(Some(frames)), + sent: sent.clone(), + fail: false, + }); + let config = cfg(ws, CancellationToken::new()); + let host = KnotHost::parse("http://oyster.cafe").unwrap(); + let mut cursor = 0i64; + + let end = run_session(&config, &host, &roster, &mut cursor).await; + + assert!(matches!(end, SessionEnd::Closed { progressed: true })); + assert_eq!(cursor, 300); + assert_eq!(member_count(&store, "did:plc:boltless"), 0); + assert_eq!(collaborator_count(&store, "did:plc:scallop"), 1); + let sent = sent.lock().unwrap(); + assert_eq!(sent.len(), 1); + assert!(matches!(&sent[0], WsMessage::Pong(p) if p.as_ref() == b"ka")); + } + + #[tokio::test] + async fn session_reports_no_progress_on_immediate_close() { + let store = Arc::new(EdgeStore::new(bobbin_runtime::RuntimeHasher::default())); + let knot = knot_acl::host_to_knot_did("oyster.cafe").unwrap(); + let roster = Mutex::new(Roster::new( + store, + knot, + Arc::new(KnotRegistry::new()), + knot_acl::KnotHostKey::new("oyster.cafe"), + )); + let frames = VecDeque::from(vec![WsMessage::Close { + code: 1000, + reason: String::new(), + }]); + let ws: Arc = Arc::new(ScriptWs { + msgs: Mutex::new(Some(frames)), + sent: Arc::new(Mutex::new(Vec::new())), + fail: false, + }); + let config = cfg(ws, CancellationToken::new()); + let host = KnotHost::parse("http://oyster.cafe").unwrap(); + let mut cursor = 99i64; + + let end = run_session(&config, &host, &roster, &mut cursor).await; + + assert!( + matches!(end, SessionEnd::Closed { progressed: false }), + "a session that delivers no frames before closing reports no progress" + ); + } + + #[tokio::test] + async fn pre_cancelled_stream_returns_without_connecting() { + let store = Arc::new(EdgeStore::new(bobbin_runtime::RuntimeHasher::default())); + let knot = knot_acl::host_to_knot_did("oyster.cafe").unwrap(); + let roster = Mutex::new(Roster::new( + store, + knot, + Arc::new(KnotRegistry::new()), + knot_acl::KnotHostKey::new("oyster.cafe"), + )); + let ws: Arc = Arc::new(ScriptWs { + msgs: Mutex::new(None), + sent: Arc::new(Mutex::new(Vec::new())), + fail: true, + }); + let cancel = CancellationToken::new(); + cancel.cancel(); + let config = cfg(ws, cancel); + let host = KnotHost::parse("http://oyster.cafe").unwrap(); + + run_stream(&config, &host, &roster, 0).await; + } + + #[test] + fn events_url_carries_scheme_and_cursor() { + let host = KnotHost::parse("http://oyster.cafe").unwrap(); + assert_eq!( + events_url(&host, 0).unwrap().as_str(), + "ws://oyster.cafe/events" + ); + assert_eq!( + events_url(&host, 42).unwrap().as_str(), + "ws://oyster.cafe/events?cursor=42" + ); + let secure = KnotHost::parse("https://nel.pet").unwrap(); + assert_eq!( + events_url(&secure, 7).unwrap().as_str(), + "wss://nel.pet/events?cursor=7" + ); + } + + #[test] + fn only_a_long_lived_session_resets_backoff() { + assert!(session_was_healthy( + &SessionEnd::Closed { progressed: true }, + HEALTHY_SESSION_MIN + )); + assert!( + !session_was_healthy( + &SessionEnd::Closed { progressed: true }, + HEALTHY_SESSION_MIN - Duration::from_millis(1) + ), + "a knot that flaps faster than the healthy floor must keep backing off" + ); + assert!(!session_was_healthy( + &SessionEnd::Closed { progressed: false }, + Duration::from_secs(3600) + )); + assert!(!session_was_healthy( + &SessionEnd::ConnectFailed, + Duration::from_secs(3600) + )); + } + + #[test] + fn jitter_delay_stays_within_half_to_full_window() { + let base = Duration::from_secs(8); + [0u64, 1, 511, 512, 1023, 1024, u64::MAX] + .into_iter() + .for_each(|entropy| { + let delay = jitter_delay(base, entropy); + assert!(delay >= base / 2, "delay {delay:?} below half of {base:?}"); + assert!(delay <= base, "delay {delay:?} above {base:?}"); + }); + } +} diff --git a/bobbin/crates/knot-proxy/src/dns.rs b/bobbin/crates/knot-proxy/src/dns.rs index 5d3ea648..f2ac808c 100644 --- a/bobbin/crates/knot-proxy/src/dns.rs +++ b/bobbin/crates/knot-proxy/src/dns.rs @@ -7,12 +7,12 @@ use reqwest::dns::{Addrs, Name, Resolve, Resolving}; use crate::host::{PrivateHostReason, classify_ip}; -pub(crate) struct PrivateAddressFilter { +pub struct PrivateAddressFilter { allow_private: bool, } impl PrivateAddressFilter { - pub(crate) fn new(allow_private: bool) -> Self { + pub fn new(allow_private: bool) -> Self { Self { allow_private } } } diff --git a/bobbin/crates/knot-proxy/src/lib.rs b/bobbin/crates/knot-proxy/src/lib.rs index fc3131c8..0f1e5c10 100644 --- a/bobbin/crates/knot-proxy/src/lib.rs +++ b/bobbin/crates/knot-proxy/src/lib.rs @@ -22,7 +22,8 @@ mod dns; mod host; pub use breaker::{Breaker, BreakerPermit, CircuitOpen, FailureThreshold, ThresholdError}; -pub use host::{KnotHost, KnotHostError, PrivateHostReason, RepoSlug, RepoSlugError}; +pub use dns::PrivateAddressFilter; +pub use host::{KnotHost, KnotHostError, PrivateHostReason, RepoSlug, RepoSlugError, classify_ip}; const USER_AGENT: &str = concat!("bobbin/", env!("CARGO_PKG_VERSION")); const HTTPS_SCHEME: &str = "https"; diff --git a/bobbin/crates/runtime/src/lib.rs b/bobbin/crates/runtime/src/lib.rs index f5466e17..c8cb3809 100644 --- a/bobbin/crates/runtime/src/lib.rs +++ b/bobbin/crates/runtime/src/lib.rs @@ -14,7 +14,7 @@ pub use mem_network::{ MemWsResponder, MemWsServerFuture, MemWsTransport, }; pub use network::{ - BodyStream, HttpRequest, HttpResponseFuture, HttpResponseHead, HttpResult, HttpTransport, - NetworkError, ReqwestHttp, TungsteniteWs, WsConn, WsConnectFuture, WsMessage, WsMessageFuture, - WsSendFuture, WsSink, WsStream, WsTransport, + AddrGuard, BodyStream, GuardedWs, HttpRequest, HttpResponseFuture, HttpResponseHead, + HttpResult, HttpTransport, NetworkError, ReqwestHttp, TungsteniteWs, WsConn, WsConnectFuture, + WsMessage, WsMessageFuture, WsSendFuture, WsSink, WsStream, WsTransport, }; diff --git a/bobbin/crates/runtime/src/network.rs b/bobbin/crates/runtime/src/network.rs index 65c2efe3..80afcee7 100644 --- a/bobbin/crates/runtime/src/network.rs +++ b/bobbin/crates/runtime/src/network.rs @@ -1,4 +1,5 @@ use std::future::Future; +use std::net::SocketAddr; use std::pin::Pin; use std::sync::Arc; @@ -137,6 +138,8 @@ pub trait WsTransport: Send + Sync + 'static { fn connect(&self, url: Url) -> WsConnectFuture; } +pub type AddrGuard = Arc Result<(), NetworkError> + Send + Sync>; + #[derive(Clone, Copy, Debug, Default)] pub struct TungsteniteWs; @@ -163,6 +166,52 @@ impl WsTransport for TungsteniteWs { } } +pub struct GuardedWs { + guard: AddrGuard, +} + +impl GuardedWs { + pub fn shared(guard: AddrGuard) -> Arc { + Arc::new(Self { guard }) + } +} + +impl WsTransport for GuardedWs { + fn connect(&self, url: Url) -> WsConnectFuture { + let guard = self.guard.clone(); + Box::pin(async move { + let host = url + .host_str() + .ok_or_else(|| NetworkError::Connect("ws url missing host".to_owned()))? + .to_owned(); + let port = url + .port_or_known_default() + .ok_or_else(|| NetworkError::Connect("ws url missing port".to_owned()))?; + let addrs: Vec = tokio::net::lookup_host((host.as_str(), port)) + .await + .map_err(|e| NetworkError::Connect(e.to_string()))? + .collect(); + guard(&addrs)?; + let addr = addrs + .into_iter() + .next() + .ok_or_else(|| NetworkError::Connect(format!("no addresses for {host}")))?; + let tcp = tokio::net::TcpStream::connect(addr) + .await + .map_err(|e| NetworkError::Connect(e.to_string()))?; + let (ws, _resp) = tokio_tungstenite::client_async_tls(url.as_str(), tcp) + .await + .map_err(|e| NetworkError::Connect(e.to_string()))?; + let (sink_inner, stream_inner) = futures::StreamExt::split(ws); + let sink: Box = Box::new(TungsteniteSink { inner: sink_inner }); + let stream: Box = Box::new(TungsteniteStream { + inner: stream_inner, + }); + Ok(WsConn { sink, stream }) + }) + } +} + type TungsteniteWsStream = tokio_tungstenite::WebSocketStream>; diff --git a/bobbin/crates/types/src/knot_acl.rs b/bobbin/crates/types/src/knot_acl.rs new file mode 100644 index 00000000..c571d375 --- /dev/null +++ b/bobbin/crates/types/src/knot_acl.rs @@ -0,0 +1,289 @@ +use core::fmt; + +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::ident::AtIdentifier; +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:"; + +#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct KnotHostKey(String); + +impl KnotHostKey { + pub fn new(host: &str) -> Self { + Self(host.trim_end_matches('.').to_ascii_lowercase()) + } + + pub fn as_str(&self) -> &str { + &self.0 + } +} + +impl AsRef for KnotHostKey { + fn as_ref(&self) -> &str { + &self.0 + } +} + +impl fmt::Display for KnotHostKey { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(&self.0) + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum KnotOwnedSource { + Member { + knot: Did, + subject: Did, + }, + Collaborator { + repo: Did, + subject: Did, + }, +} + +pub fn host_to_knot_did(host: &str) -> Option> { + if host.contains('/') { + return None; + } + let normalized = KnotHostKey::new(host); + let host = normalized.as_str(); + if host.is_empty() { + return None; + } + 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, +) -> Option> { + build_source( + repo.as_ref(), + KNOT_COLLABORATOR_COLLECTION, + subject.as_ref(), + ) +} + +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, vec![edge])) +} + +pub fn collaborator_upsert( + repo: &Did, + subject: &Did, + created_micros: u64, +) -> Option<(AtUri, Vec)> { + let source = collaborator_source(repo, subject)?; + let edge = Edge { + kind: nsid_static(REPO_COLLABORATOR_KIND), + subject: SubjectRef::Did(repo.clone()), + source: source.clone(), + sort_micros: created_micros, + }; + Some((source, vec![edge])) +} + +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 { + 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, + } +} + +fn build_source(authority: &str, collection: &str, rkey: &str) -> Option> { + AtUri::new_owned(format!("at://{authority}/{collection}/{rkey}")).ok() +} + +#[cfg(test)] +mod tests { + use super::*; + + fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + + fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() + } + + #[test] + fn host_did_round_trips() { + 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() { + 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] + fn host_rejects_slash_and_empty() { + assert_eq!(host_to_knot_did("oyster.cafe/evil"), None); + assert_eq!(host_to_knot_did(""), None); + assert_eq!(host_to_knot_did("."), None); + } + + #[test] + fn host_key_normalizes_case_and_trailing_dot() { + assert_eq!( + KnotHostKey::new("Kt.Oyster.Cafe.").as_str(), + "kt.oyster.cafe" + ); + assert_eq!( + KnotHostKey::new("KT.OYSTER.CAFE:3000").as_str(), + "kt.oyster.cafe:3000" + ); + assert_eq!( + KnotHostKey::new("kt.oyster.cafe"), + KnotHostKey::new("KT.OYSTER.CAFE") + ); + } + + #[test] + fn host_to_knot_did_normalizes_before_encoding() { + assert_eq!( + host_to_knot_did("KT.Oyster.Cafe").unwrap().as_ref(), + "did:web:kt.oyster.cafe" + ); + assert_eq!( + host_to_knot_did("KT.Oyster.Cafe:3000").unwrap().as_ref(), + "did:web:kt.oyster.cafe%3A3000" + ); + } + + #[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"); + let subject = did("did:plc:olaren"); + let source = collaborator_source(&repo, &subject).expect("build collaborator source"); + assert_eq!( + source.as_ref(), + "at://did:plc:scallop/sh.tangled.bobbin.knotCollaborator/did:plc:olaren" + ); + assert_eq!( + decode_knot_owned_source(&source), + Some(KnotOwnedSource::Collaborator { repo, subject }) + ); + } + + #[test] + fn decode_ignores_legacy_pds_member_record() { + let source = at("at://did:plc:nel/sh.tangled.knot.member/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"); + assert_eq!(decode_knot_owned_source(&source), None); + } + + #[test] + fn member_upsert_builds_decodable_primary_edge() { + 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(), 1); + let edge = &edges[0]; + assert_eq!(edge.kind.as_ref(), "sh.tangled.knot.member"); + assert_eq!(edge.subject, SubjectRef::Did(subject.clone())); + assert_eq!(edge.source, source); + assert_eq!(edge.sort_micros, 1_700_000_000_000_000); + assert_eq!( + decode_knot_owned_source(&source), + Some(KnotOwnedSource::Member { knot, subject }) + ); + } + + #[test] + fn collaborator_upsert_keys_on_repo_did() { + let repo = did("did:plc:scallop"); + let subject = did("did:plc:olaren"); + let (source, edges) = + collaborator_upsert(&repo, &subject, 42).expect("build collaborator upsert"); + assert_eq!(edges.len(), 1); + let edge = &edges[0]; + assert_eq!(edge.kind.as_ref(), "sh.tangled.repo.collaborator"); + assert_eq!(edge.subject, SubjectRef::Did(repo.clone())); + assert_eq!(edge.source, source); + assert_eq!(edge.sort_micros, 42); + assert_eq!( + decode_knot_owned_source(&source), + Some(KnotOwnedSource::Collaborator { repo, subject }) + ); + } +} diff --git a/bobbin/crates/types/src/lib.rs b/bobbin/crates/types/src/lib.rs index 52d0dc01..3a61600e 100644 --- a/bobbin/crates/types/src/lib.rs +++ b/bobbin/crates/types/src/lib.rs @@ -19,6 +19,7 @@ pub use _lex::*; pub mod edges; pub mod ids; +pub mod knot_acl; pub mod legacy; pub mod record; pub mod search; diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index c8585bb3..370859eb 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -23,8 +23,8 @@ use axum::{ routing::get, }; use bobbin_edge_index::{ - Coverage, CoverageWatch, CursorParseError, EdgePage, EdgeStore, IssueStateKind, PageCursor, - PageLimit, PageToken, PullStatusKind, SortDir, StateIndex, StateKind, + Coverage, CoverageWatch, CursorParseError, EdgeItem, EdgePage, EdgeStore, IssueStateKind, + PageCursor, PageLimit, PageToken, PullStatusKind, SortDir, StateIndex, StateKind, }; use bobbin_knot_proxy::{KnotHost, KnotProxy, KnotProxyError, ProxyResponse, RepoSlug}; use bobbin_record_lru::RecordStore; @@ -34,6 +34,7 @@ use bobbin_search::{ }; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; +use bobbin_types::knot_acl::{KnotOwnedSource, decode_knot_owned_source, knot_did_host}; use bobbin_types::record::RecordBody; use bobbin_types::search::SearchableRecord; use bobbin_types::sh_tangled::actor::profile::{Profile, ProfileGetRecordOutput, ProfileRecord}; @@ -1434,10 +1435,14 @@ async fn hydrate_record_view( state: &AppState, nsid: &Nsid, uri: AtUri, + sort_micros: u64, ) -> Result>, XrpcError> where V: serde::de::DeserializeOwned + NormalizeRepoRefs, { + if let Some(source) = decode_knot_owned_source(&uri) { + return synthesize_knot_owned_view::(state, uri, source, sort_micros).await; + } let body = resolve_for_view(state, nsid, uri).await?; let value: V = deserialize_or_upgrade::(state, nsid, &body.value).await?; let Some(value) = value.normalize(&state.resolver).await else { @@ -1450,6 +1455,53 @@ where })) } +async fn synthesize_knot_owned_view( + state: &AppState, + uri: AtUri, + source: KnotOwnedSource, + sort_micros: u64, +) -> Result>, XrpcError> +where + V: serde::de::DeserializeOwned + NormalizeRepoRefs, +{ + let Some(body) = synth_knot_owned_value(source, sort_micros) else { + return Ok(None); + }; + let Ok(value) = serde_json::from_value::(body) else { + return Ok(None); + }; + let Some(value) = value.normalize(&state.resolver).await else { + return Ok(None); + }; + Ok(Some(RecordView { + uri, + cid: None, + value, + })) +} + +fn synth_knot_owned_value(source: KnotOwnedSource, sort_micros: u64) -> Option { + let created_at = micros_to_rfc3339(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(), + "createdAt": created_at, + })), + } +} + +fn micros_to_rfc3339(micros: u64) -> Option { + let micros = i64::try_from(micros).ok()?; + chrono::DateTime::from_timestamp_micros(micros) + .map(|dt| dt.to_rfc3339_opts(chrono::SecondsFormat::Micros, true)) +} + #[derive(Clone, Copy)] enum HitProvenance { ClientSupplied, @@ -1533,18 +1585,19 @@ where fn hydrate_record_stream( state: &AppState, nsid: Nsid, - uris: Vec>, + items: Vec, provenance: HitProvenance, ) -> impl Stream, XrpcError>> + Send + 'static where V: serde::de::DeserializeOwned + Serialize + NormalizeRepoRefs + Send + 'static, { let owned = state.clone(); - hydrate_stream(uris, move |uri| { + hydrate_stream(items, move |item| { let owned = owned.clone(); let nsid = nsid.clone(); async move { - let result = hydrate_record_view::(&owned, &nsid, uri.clone()).await; + let EdgeItem { uri, sort_micros } = item; + let result = hydrate_record_view::(&owned, &nsid, uri.clone(), sort_micros).await; if matches!(provenance, HitProvenance::Indexed) && let Err(err) = &result && is_index_evictable(err) @@ -1673,7 +1726,14 @@ where ))); } let permit = state.heavy_permit()?; - let views = hydrate_record_stream::(state, nsid, parsed, HitProvenance::ClientSupplied); + let items = parsed + .into_iter() + .map(|uri| EdgeItem { + uri, + sort_micros: 0, + }) + .collect(); + let views = hydrate_record_stream::(state, nsid, items, HitProvenance::ClientSupplied); Ok(json_stream::, _>( "items", views, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index f5e70308..a0b1924d 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -2061,3 +2061,79 @@ 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]["cid"].is_null()); + 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_collaborator_is_synthesized_without_slingshot() { + let harness = Harness::new().await; + let repo = did("did:plc:scallop"); + let subject = did("did:plc:olaren"); + let created = chrono::DateTime::parse_from_rfc3339("2026-06-03T12:00:00Z").unwrap(); + let micros = created.timestamp_micros() as u64; + let (source, edges) = + bobbin_types::knot_acl::collaborator_upsert(&repo, &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.repo.listCollaborators", + repo.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]["cid"].is_null()); + assert_eq!(items[0]["value"]["repo"], json!("did:plc:scallop")); + assert_eq!(items[0]["value"]["subject"], json!("did:plc:olaren")); +} -- 2.51.2