diff --git a/knot2/crates/knot-sim/src/harness.rs b/knot2/crates/knot-sim/src/harness.rs index 1b783e076..e7434fb2e 100644 --- a/knot2/crates/knot-sim/src/harness.rs +++ b/knot2/crates/knot-sim/src/harness.rs @@ -2,7 +2,6 @@ use std::collections::{BTreeSet, HashMap, HashSet}; use std::fs; use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex}; -use std::time::Duration; use axum::{Json, Router, routing::get}; use base64::Engine; @@ -31,9 +30,9 @@ use knot_types::{ }; use knot_xrpc::{ BlobReadBudget, BodyLimit, Budgets, ByteLimits, CobLocks, Committer, GlobalInflight, - GlobalQuota, LanguagesPushBudget, LanguagesReadBudget, LimitConfig, MaxWireBytes, - PerActorQuota, PerPeerInflight, PreAuthLimiter, ReadBudget, ReservationTtl, Reservations, - ResponseLimit, TreeReadBudget, XrpcState, + GlobalQuota, LanguagesReadBudget, LimitConfig, MaxWireBytes, PerActorQuota, PerPeerInflight, + PreAuthLimiter, ReadBudget, ReservationTtl, Reservations, ResponseLimit, TreeReadBudget, + XrpcState, }; use crate::realdata::{ @@ -136,11 +135,17 @@ pub(crate) struct Published(Mutex>); impl Published { pub(crate) fn write(&self, record: AcceptanceRecord) { - self.0.lock().expect("published records lock").insert(record); + self.0 + .lock() + .expect("published records lock") + .insert(record); } fn holds(&self, record: &AcceptanceRecord) -> bool { - self.0.lock().expect("published records lock").contains(record) + self.0 + .lock() + .expect("published records lock") + .contains(record) } } @@ -515,8 +520,8 @@ fn build_responder( return Ok(ok_body(Bytes::new())); } if request.url.path().ends_with("com.atproto.repo.getRecord") { - let written = named_acceptance(&request.url) - .is_some_and(|record| published.holds(&record)); + let written = + named_acceptance(&request.url).is_some_and(|record| published.holds(&record)); return Ok(match written { true => ok_body(Bytes::new()), false => not_found(), @@ -836,7 +841,6 @@ fn assemble_router(parts: StateParts) -> Router { tree_last_commit: TreeReadBudget::new(ReadBudget::Unbounded), blob_last_commit: BlobReadBudget::new(ReadBudget::Unbounded), languages: LanguagesReadBudget::new(ReadBudget::Unbounded), - languages_push: LanguagesPushBudget::new(Duration::from_secs(120)), }, git_http: Arc::new(FakeHttp::new(|_request: &HttpRequest| { Err(NetworkError::Connect( @@ -845,7 +849,7 @@ fn assemble_router(parts: StateParts) -> Router { })), pack_limits: knot_pack::PackLimits::default(), service_owner, - events: Arc::new(EventLog::new( + firehose: Arc::new(EventLog::new( Arc::clone(&clock), knot_events::ReplayBounds::new( knot_events::ReplayEvents::new(4096).expect("replay event maximum is nonzero"), diff --git a/knot2/crates/knot-sim/src/workload.rs b/knot2/crates/knot-sim/src/workload.rs index 2ad95f4ff..903a9c9dd 100644 --- a/knot2/crates/knot-sim/src/workload.rs +++ b/knot2/crates/knot-sim/src/workload.rs @@ -292,7 +292,13 @@ pub(crate) fn predict(seed: u64, rounds: u32) -> Projection { members: model.members.into_iter().collect(), blocked: model.blocked.into_iter().collect(), collaborators: (0..repos) - .map(|repo| rosters.remove(&repo).unwrap_or_default().into_iter().collect()) + .map(|repo| { + rosters + .remove(&repo) + .unwrap_or_default() + .into_iter() + .collect() + }) .collect(), } } @@ -480,8 +486,10 @@ pub(crate) async fn execute(harness: Arc, seed: u64, plan: Vec) .iter() .filter_map(|result| result.created.clone()) .collect(); - let planned_creates = - ops.iter().filter(|planned| planned.creates_a_repo()).count(); + let planned_creates = ops + .iter() + .filter(|planned| planned.creates_a_repo()) + .count(); assert_eq!( planned_creates, created.len(), diff --git a/knot2/crates/knot-sim/tests/lfs_roundtrip.rs b/knot2/crates/knot-sim/tests/lfs_roundtrip.rs index ca93ba675..5f6c3b2b0 100644 --- a/knot2/crates/knot-sim/tests/lfs_roundtrip.rs +++ b/knot2/crates/knot-sim/tests/lfs_roundtrip.rs @@ -5,7 +5,6 @@ use std::io::Write; use std::path::{Path, PathBuf}; use std::process::{Command, Stdio}; use std::sync::Arc; -use std::time::Duration; use base64::Engine; use base64::engine::general_purpose::URL_SAFE_NO_PAD; @@ -300,27 +299,18 @@ async fn spawn(published_line: String, with_h3: bool) -> World { let key_dir = scan.path().join("hostkey"); std::fs::create_dir_all(&key_dir).unwrap(); let host_key = knot_ssh::load_or_create_host_key(&key_dir.join("host")).unwrap(); - let events = Arc::new(knot_events::EventLog::new( - ManualClock::new(UnixMicros::new(1_000_000_000)), - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(64).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - )); let ssh_state = Arc::new( knot_ssh::SshState::new(knot_ssh::SshConfig { layout: layout.clone(), index: Arc::clone(&index), atproto: Arc::clone(&atproto), knot_actor: knot_types::ActorId::from_secp256k1(actor_signer().public_key().as_bytes()), - events: Arc::clone(&events), hostname: KnotHostname::new("nel.pet").unwrap(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), admins: BTreeSet::from([AccountDid::new(OWNER_DID).unwrap()]), admission: AdmissionPolicy::Closed, max_pack_bytes: knot_xrpc::MaxWireBytes::new(1 << 30), archive_limit: knot_git::ArchiveLimit::default(), - languages_push_budget: knot_xrpc::LanguagesPushBudget::new(Duration::from_secs(2)), ci_logs: None, }) .with_lfs(lfs.clone(), 16), @@ -374,7 +364,13 @@ async fn spawn(published_line: String, with_h3: bool) -> World { })), pack_limits: knot_pack::PackLimits::default(), service_owner: AccountDid::new(OWNER_DID).unwrap(), - events, + firehose: Arc::new(knot_events::EventLog::new( + ManualClock::new(UnixMicros::new(1_000_000_000)), + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(64).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), + )), subscriber_gate: Arc::new(knot_events::SubscriberGate::new( knot_events::GlobalSubscriberLimit::new(16), knot_events::PerPeerSubscriberLimit::new(4), diff --git a/knot2/crates/knot-sim/tests/reproducible.rs b/knot2/crates/knot-sim/tests/reproducible.rs index d983911f4..c6bcd1036 100644 --- a/knot2/crates/knot-sim/tests/reproducible.rs +++ b/knot2/crates/knot-sim/tests/reproducible.rs @@ -316,12 +316,16 @@ async fn roster_read_and_roster_write_land_in_separate_rounds() { ); }); - [OFFERS.as_slice(), ANSWERS.as_slice(), ROSTER_READS.as_slice()] - .into_iter() - .for_each(|names| { - assert!( - rounds.values().any(|ops| holds(ops, names)), - "the rule is vacuous unless {names:?} run in some round" - ); - }); + [ + OFFERS.as_slice(), + ANSWERS.as_slice(), + ROSTER_READS.as_slice(), + ] + .into_iter() + .for_each(|names| { + assert!( + rounds.values().any(|ops| holds(ops, names)), + "the rule is vacuous unless {names:?} run in some round" + ); + }); } diff --git a/knot2/crates/knot-sim/tests/ssh.rs b/knot2/crates/knot-sim/tests/ssh.rs index 389d32cb0..bd8ca9a0e 100644 --- a/knot2/crates/knot-sim/tests/ssh.rs +++ b/knot2/crates/knot-sim/tests/ssh.rs @@ -2,7 +2,6 @@ use std::path::Path; use std::process::Command; use std::sync::Arc; -use futures::stream::StreamExt; use knot_atproto::Atproto; use knot_cob::{CobHome, CobStore}; use knot_cobs::{Registration, RegistryChange}; @@ -136,7 +135,6 @@ struct Server { layout: Layout, repo_did: RepoDid, port: u16, - events: Arc>, } async fn spawn(published_line: String) -> Server { @@ -176,28 +174,17 @@ async fn spawn(published_line: String) -> Server { let key_dir = scan.path().join("hostkey"); std::fs::create_dir_all(&key_dir).unwrap(); let host_key = knot_ssh::load_or_create_host_key(&key_dir.join("host")).unwrap(); - let events = Arc::new(knot_events::EventLog::new( - ManualClock::new(UnixMicros::new(1_000_000_000)), - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(64).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - )); let state = Arc::new(knot_ssh::SshState::new(knot_ssh::SshConfig { layout: layout.clone(), index, atproto, knot_actor: actor_for_seed(1), - events: Arc::clone(&events), hostname: knot_types::KnotHostname::new("knot.test").unwrap(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), admins: std::collections::BTreeSet::new(), admission: knot_types::AdmissionPolicy::Closed, max_pack_bytes: knot_xrpc::MaxWireBytes::new(1 << 30), archive_limit: knot_git::ArchiveLimit::default(), - languages_push_budget: knot_xrpc::LanguagesPushBudget::new(std::time::Duration::from_secs( - 2, - )), ci_logs: None, })); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); @@ -210,7 +197,6 @@ async fn spawn(published_line: String) -> Server { layout, repo_did, port, - events, } } @@ -232,7 +218,7 @@ fn seed_work(work: &Path) -> String { head.trim().to_string() } -async fn push_once(scratch: &Path) -> (String, Option, serde_json::Value) { +async fn push_once(scratch: &Path) -> (String, Option) { let (key_path, public_line) = keygen(scratch); let server = spawn(public_line).await; let url = format!( @@ -260,46 +246,17 @@ async fn push_once(scratch: &Path) -> (String, Option, serde_js .unwrap() .find_ref(&knot_types::RefName::new("refs/heads/main").unwrap()) .unwrap(); - let event = poll_for_event(&server.events).await; drop(server); - (head, stored, event) -} - -fn replay_bounds() -> knot_events::ReplayBounds { - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(32).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ) -} - -async fn poll_for_event(events: &knot_events::EventLog) -> serde_json::Value { - futures::stream::iter(0..100) - .then(|_| async { - let hit = events - .replay(knot_events::EventCursor::START, replay_bounds()) - .events - .into_iter() - .find(|event| event.nsid == "sh.tangled.git.refUpdate") - .map(|event| serde_json::to_value(&*event).unwrap()["event"].clone()); - if hit.is_none() { - tokio::time::sleep(std::time::Duration::from_millis(20)).await; - } - hit - }) - .filter_map(|hit| async move { hit }) - .boxed() - .next() - .await - .expect("refUpdate event must be published within polling window") + (head, stored) } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn the_simulated_ssh_write_path_lands_a_seed_deterministic_tip() { let first_dir = tempfile::tempdir().unwrap(); - let (first_head, first_stored, first_event) = push_once(first_dir.path()).await; + let (first_head, first_stored) = push_once(first_dir.path()).await; let second_dir = tempfile::tempdir().unwrap(); - let (second_head, second_stored, second_event) = push_once(second_dir.path()).await; + let (second_head, second_stored) = push_once(second_dir.path()).await; let tip = knot_types::Oid::from_hex(&first_head).unwrap(); assert_eq!( @@ -315,11 +272,4 @@ async fn the_simulated_ssh_write_path_lands_a_seed_deterministic_tip() { first_stored, second_stored, "assembled-against-doubles ssh write path is logically reproducible" ); - - assert_eq!(first_event["ref"], "refs/heads/main"); - assert_eq!(first_event["newSha"], first_head); - assert_eq!( - first_event["newSha"], second_event["newSha"], - "ref-update event the push emits is seed-stable too" - ); } diff --git a/knot2/crates/knot-ssh/Cargo.toml b/knot2/crates/knot-ssh/Cargo.toml index 47cd14c3d..4442ef936 100644 --- a/knot2/crates/knot-ssh/Cargo.toml +++ b/knot2/crates/knot-ssh/Cargo.toml @@ -18,7 +18,6 @@ knot-consent = { workspace = true } knot-atproto = { workspace = true } knot-cob = { workspace = true } knot-cobs = { workspace = true } -knot-events = { workspace = true } knot-maintenance = { workspace = true } knot-postreceive = { workspace = true } knot-receive = { workspace = true } @@ -39,5 +38,10 @@ http = { workspace = true } bytes = { workspace = true } url = { workspace = true } russh = "0.61" +knot-events = { workspace = true } +serde = { workspace = true } +knot-secrets = { workspace = true } +knot-xrpc = { workspace = true } +serde_ipld_dagcbor = { workspace = true } tikv-jemallocator = { workspace = true } tikv-jemalloc-ctl = { workspace = true } diff --git a/knot2/crates/knot-ssh/examples/ephemeral_knot.rs b/knot2/crates/knot-ssh/examples/ephemeral_knot.rs index 7424a5f12..bc2ffe27c 100644 --- a/knot2/crates/knot-ssh/examples/ephemeral_knot.rs +++ b/knot2/crates/knot-ssh/examples/ephemeral_knot.rs @@ -193,13 +193,6 @@ async fn main() { std::fs::create_dir_all(&key_dir).unwrap(); let host_key = knot_ssh::load_or_create_host_key(&key_dir.join("host")).unwrap(); - let events = Arc::new(knot_events::EventLog::new( - ManualClock::new(UnixMicros::new(1_000_000_000)), - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(64).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - )); let actor = knot_types::ActorId::from_secp256k1( K256Signer::generate(&SeededEntropy::new(1)) .public_key() @@ -210,14 +203,12 @@ async fn main() { index, atproto, knot_actor: actor, - events, hostname: KnotHostname::new("knot.test").unwrap(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), admins: BTreeSet::new(), admission: AdmissionPolicy::Closed, max_pack_bytes: knot_pack::MaxWireBytes::new(1 << 34), archive_limit: knot_git::ArchiveLimit::default(), - languages_push_budget: knot_postreceive::LanguagesPushBudget::new(Duration::from_secs(2)), ci_logs: None, })); diff --git a/knot2/crates/knot-ssh/src/exec.rs b/knot2/crates/knot-ssh/src/exec.rs index 0a137e7e2..48ef0d004 100644 --- a/knot2/crates/knot-ssh/src/exec.rs +++ b/knot2/crates/knot-ssh/src/exec.rs @@ -712,16 +712,15 @@ async fn serve_receive( limits: state.limits, knot_actor: state.knot_actor.clone(), committer, - events: Arc::clone(&state.events), index: &state.index, atproto: &state.atproto, resolve_slots: &state.slots.resolve, appview: &state.appview, maintenance: &state.maintenance, hostname: &state.hostname, - languages_push_budget: state.languages_push_budget, catalog: Arc::clone(&state.catalog), ci_logs: state.ci_logs.clone(), + ref_surface: state.ref_surface.clone(), }) .await; match landed { diff --git a/knot2/crates/knot-ssh/src/lib.rs b/knot2/crates/knot-ssh/src/lib.rs index 5023a180e..125e97201 100644 --- a/knot2/crates/knot-ssh/src/lib.rs +++ b/knot2/crates/knot-ssh/src/lib.rs @@ -10,12 +10,10 @@ use std::sync::Arc; use std::time::Duration; use knot_atproto::Atproto; -use knot_events::EventLog; use knot_git::{ArchiveLimit, Layout}; use knot_index::{Index, KeyTtl}; use knot_maintenance::MaintenanceHandle; use knot_pack::{MaxWireBytes, PackLimits}; -use knot_postreceive::LanguagesPushBudget; use knot_runtime::{Clock, Entropy, HttpTransport, OsEntropy}; use knot_types::{ AccountDid, ActorId, AdmissionPolicy, AppviewEndpoint, CiLogsAddr, KnotHostname, KnotId, @@ -63,7 +61,6 @@ pub struct SshState { index: Arc, atproto: Arc>, knot_actor: ActorId, - events: Arc>, hostname: KnotHostname, knot_did: KnotId, appview: AppviewEndpoint, @@ -72,7 +69,6 @@ pub struct SshState { limits: PackLimits, max_pack_bytes: MaxWireBytes, archive_limit: ArchiveLimit, - languages_push_budget: LanguagesPushBudget, ci_logs: Option, slots: Slots, lookup_slots: ResolveSlots, @@ -84,6 +80,7 @@ pub struct SshState { maintenance: MaintenanceHandle, lfs: Option, catalog: Arc, + ref_surface: Option>, } #[derive(Clone)] @@ -98,14 +95,12 @@ pub struct SshConfig { pub index: Arc, pub atproto: Arc>, pub knot_actor: ActorId, - pub events: Arc>, pub hostname: KnotHostname, pub appview: AppviewEndpoint, pub admins: BTreeSet, pub admission: AdmissionPolicy, pub max_pack_bytes: MaxWireBytes, pub archive_limit: ArchiveLimit, - pub languages_push_budget: LanguagesPushBudget, pub ci_logs: Option, } @@ -116,14 +111,12 @@ impl SshState { index, atproto, knot_actor, - events, hostname, appview, admins, admission, max_pack_bytes, archive_limit, - languages_push_budget, ci_logs, } = config; let knot_did = hostname.knot_did(); @@ -132,7 +125,6 @@ impl SshState { index, atproto, knot_actor, - events, hostname, knot_did, appview, @@ -141,7 +133,6 @@ impl SshState { limits: PackLimits::default(), max_pack_bytes, archive_limit, - languages_push_budget, ci_logs, slots: Slots::for_machine(), lookup_slots: ResolveSlots::new(MAX_INFLIGHT_LOOKUPS), @@ -168,6 +159,7 @@ impl SshState { maintenance: MaintenanceHandle::disabled(), lfs: None, catalog: Arc::new(knot_messages::Catalog::defaults()), + ref_surface: None, } } @@ -176,6 +168,11 @@ impl SshState { self } + pub fn with_ref_surface(mut self, surface: Arc) -> Self { + self.ref_surface = Some(surface); + self + } + pub fn with_slots(mut self, slots: Slots) -> Self { self.slots = slots; self diff --git a/knot2/crates/knot-ssh/tests/ssh_push.rs b/knot2/crates/knot-ssh/tests/ssh_push.rs index b61c094ce..db38e2f59 100644 --- a/knot2/crates/knot-ssh/tests/ssh_push.rs +++ b/knot2/crates/knot-ssh/tests/ssh_push.rs @@ -11,7 +11,6 @@ use knot_cobs::{CollaboratorsChange, Grant, MembersChange, Registration, Registr use knot_git::{ArchiveLimit, Layout, Repo}; use knot_index::{Index, Resolved}; use knot_pack::MaxWireBytes; -use knot_postreceive::LanguagesPushBudget; use knot_runtime::{ FakeDns, FakeHttp, HttpResponse, K256Signer, ManualClock, SeededEntropy, Signer, UnixMicros, }; @@ -271,7 +270,6 @@ struct Server { layout: Layout, repo_did: RepoDid, port: u16, - events: Arc>, } async fn spawn_server( @@ -286,12 +284,13 @@ async fn spawn_server_with( max_pack_bytes: MaxWireBytes, warm: bool, ) -> (Server, Arc) { - let (server, index, _, _) = spawn_server_core( + let (server, index, _, _, _) = spawn_server_core( published_line, max_pack_bytes, ArchiveLimit::default(), warm, None, + false, ) .await; (server, index) @@ -303,11 +302,13 @@ async fn spawn_server_core( archive_limit: ArchiveLimit, warm: bool, lfs: Option, + projection: bool, ) -> ( Server, Arc, tokio_util::sync::CancellationToken, tokio::task::JoinHandle<()>, + Option>>, ) { let scan = tempfile::tempdir().unwrap(); let meta_path = scan.path().join("meta"); @@ -373,28 +374,91 @@ async fn spawn_server_core( std::fs::create_dir_all(&key_dir).unwrap(); let host_key = knot_ssh::load_or_create_host_key(&key_dir.join("host")).unwrap(); - let events = Arc::new(knot_events::EventLog::new( - ManualClock::new(UnixMicros::new(1_000_000_000)), - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(64).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - )); + let (xrpc_state, firehose) = if projection { + let knot = KnotId::new("did:web:nel.pet").unwrap(); + let secrets = Arc::new( + knot_secrets::SealedStore::open( + scan.path().join("keys.sealed"), + &knot_secrets::MasterKey::new([7u8; 32]).unwrap(), + Box::new(knot_runtime::OsEntropy), + ) + .unwrap(), + ); + secrets.ensure(&knot).unwrap(); + let state = Arc::new(knot_xrpc::XrpcState { + layout: layout.clone(), + index: Arc::clone(&index), + atproto: Arc::clone(&atproto), + secrets, + entropy: Arc::new(knot_runtime::OsEntropy), + ci_logs: None, + admins: std::collections::BTreeSet::new(), + admission: knot_types::AdmissionPolicy::Closed, + knot_did: knot, + knot_hostname: knot_types::KnotHostname::new("knot.test").unwrap(), + meta_path: scan.path().join("meta"), + knot_service_url: knot_types::KnotServiceUrl::new("https://knot.test".to_string()) + .unwrap(), + limiter: Arc::new(knot_xrpc::PreAuthLimiter::with_config( + knot_xrpc::LimitConfig::default(), + )), + cob_locks: Arc::new(knot_xrpc::CobLocks::default()), + reservations: Arc::new(knot_xrpc::Reservations::new( + knot_xrpc::ReservationTtl::new(1_000_000), + knot_xrpc::PerActorQuota::new(16), + knot_xrpc::GlobalQuota::new(16), + )), + proxy_trust: knot_types::ProxyTrust::default(), + committer: knot_xrpc::Committer { + name: knot_types::AuthorName::new("Tangled"), + email: knot_types::Email::new("noreply@tangled.sh"), + }, + byte_limits: knot_xrpc::ByteLimits::default(), + budgets: knot_xrpc::Budgets::default(), + git_http: Arc::new(knot_runtime::FakeHttp::new( + |_request: &knot_runtime::HttpRequest| Ok(not_found()), + )), + pack_limits: knot_pack::PackLimits::default(), + service_owner: AccountDid::new(OWNER_DID).unwrap(), + subscriber_gate: Arc::new(knot_events::SubscriberGate::new( + knot_events::GlobalSubscriberLimit::new(16), + knot_events::PerPeerSubscriberLimit::new(4), + )), + maintenance: knot_maintenance::MaintenanceHandle::disabled(), + appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), + slots: knot_resource::Slots::testing(8), + lfs: None, + catalog: Arc::new(knot_messages::Catalog::defaults()), + firehose: Arc::new(knot_events::EventLog::new( + ManualClock::new(UnixMicros::new(1_000_000_000)), + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(4096).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), + )), + }); + let firehose = Arc::clone(&state.firehose); + (Some(state), Some(firehose)) + } else { + (None, None) + }; let base = knot_ssh::SshState::new(knot_ssh::SshConfig { layout: layout.clone(), index: Arc::clone(&index), atproto, knot_actor: actor_for_seed(1), - events: Arc::clone(&events), hostname: knot_types::KnotHostname::new("knot.test").unwrap(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), admins: std::collections::BTreeSet::new(), admission: knot_types::AdmissionPolicy::Closed, max_pack_bytes, archive_limit, - languages_push_budget: LanguagesPushBudget::new(std::time::Duration::from_secs(2)), ci_logs: None, }); + let base = match xrpc_state { + Some(xrpc) => base.with_ref_surface(Arc::new(knot_xrpc::RefProjection { state: xrpc })), + None => base, + }; let state = Arc::new(match lfs { Some(handle) => base.with_lfs(handle, 2), None => base, @@ -416,11 +480,11 @@ async fn spawn_server_core( layout, repo_did, port, - events, }, index, shutdown, serve_task, + firehose, ) } @@ -524,32 +588,6 @@ fn ref_names(server: &Server) -> Vec { .collect() } -fn replay_bounds() -> knot_events::ReplayBounds { - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(32).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ) -} - -async fn poll_for_event( - events: &knot_events::EventLog, - nsid: &str, -) -> serde_json::Value { - for _ in 0..50 { - if let Some(payload) = events - .replay(knot_events::EventCursor::START, replay_bounds()) - .events - .into_iter() - .find(|event| event.nsid == nsid) - .map(|event| serde_json::to_value(&*event).unwrap()["event"].clone()) - { - return payload; - } - tokio::time::sleep(std::time::Duration::from_millis(20)).await; - } - panic!("no {nsid} event was published within the polling window"); -} - struct Fixture { scratch: TempDir, server: Server, @@ -566,12 +604,13 @@ async fn fixture() -> Fixture { async fn fixture_with_archive_limit(archive_limit: ArchiveLimit) -> Fixture { let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "client"); - let (server, index, _, _) = spawn_server_core( + let (server, index, _, _, _) = spawn_server_core( public_line, MaxWireBytes::new(1 << 30), archive_limit, true, None, + false, ) .await; let url = format!( @@ -806,22 +845,6 @@ async fn an_authorized_push_over_ssh_succeeds_and_a_clone_reads_it_back() { ); } -#[tokio::test(flavor = "multi_thread", worker_threads = 4)] -async fn an_authorized_push_emits_a_ref_update_event() { - let fx = fixture().await; - let head = seed_work(&fx.work); - - let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &["main"]).await; - assert!(ok, "authorized push over ssh must succeed:\n{out}"); - - let event = poll_for_event(&fx.server.events, "sh.tangled.git.refUpdate").await; - assert_eq!(event["ref"], "refs/heads/main"); - assert_eq!(event["newSha"], head); - assert_eq!(event["committerDid"], OWNER_DID); - assert_eq!(event["ownerDid"], OWNER_DID); - assert_eq!(event["meta"]["isDefaultRef"], true); -} - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn a_client_requesting_ssh_compression_clones_an_incompressible_pack() { let fx = fixture().await; @@ -947,7 +970,10 @@ async fn ref_namespace_policy() { ), ] { let (ok, out) = push(&fx.work, &fx.url, &fx.key_path, &[spec]).await; - assert!(!ok, "the ssh path screens namespaces too, so {spec} must be refused:\n{out}"); + assert!( + !ok, + "the ssh path screens namespaces too, so {spec} must be refused:\n{out}" + ); assert!( out.contains(names_itself), "git must show the pusher which namespace refused {spec}:\n{out}" @@ -1255,26 +1281,17 @@ async fn launch(host_key_dir: &Path, layout: Layout, index: Arc, accounts )); std::fs::create_dir_all(host_key_dir).unwrap(); let host_key = knot_ssh::load_or_create_host_key(&host_key_dir.join("host")).unwrap(); - let events = Arc::new(knot_events::EventLog::new( - ManualClock::new(UnixMicros::new(1_000_000_000)), - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(64).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - )); let state = Arc::new(knot_ssh::SshState::new(knot_ssh::SshConfig { layout, index, atproto, knot_actor: actor_for_seed(77), - events, hostname: knot_types::KnotHostname::new("knot.test").unwrap(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), admins: std::collections::BTreeSet::new(), admission: knot_types::AdmissionPolicy::Closed, max_pack_bytes: MaxWireBytes::new(1 << 30), archive_limit: ArchiveLimit::default(), - languages_push_budget: LanguagesPushBudget::new(std::time::Duration::from_secs(2)), ci_logs: None, })); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); @@ -2077,12 +2094,13 @@ async fn shutdown_drains_an_in_flight_lfs_transfer_before_exit() { knot_lfs::FreeSpaceFloor::new(0), ) .unwrap(); - let (server, _index, shutdown, serve_task) = spawn_server_core( + let (server, _index, shutdown, serve_task, _) = spawn_server_core( public_line, MaxWireBytes::new(1 << 20), ArchiveLimit::default(), true, Some(handle.clone()), + false, ) .await; @@ -2152,3 +2170,77 @@ async fn shutdown_drains_an_in_flight_lfs_transfer_before_exit() { "a connection after shutdown must be refused, the drain only covers in-flight work" ); } + +#[derive(serde::Deserialize)] +struct OpProbe { + path: String, +} + +#[derive(serde::Deserialize)] +struct CommitProbe { + ops: Vec, +} + +fn frame_kind_and_ops(frame: &[u8]) -> (&'static str, Vec) { + let mut read = frame; + let header: serde_json::Value = serde_ipld_dagcbor::de::from_reader_once(&mut read).unwrap(); + match header["t"].as_str().unwrap_or_default() { + "#commit" => { + let payload: CommitProbe = serde_ipld_dagcbor::de::from_reader_once(&mut read).unwrap(); + let ops = payload.ops.into_iter().map(|op| op.path).collect(); + ("#commit", ops) + } + "#sync" => ("#sync", Vec::new()), + _ => ("#ignored", Vec::new()), + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn an_authorized_ssh_push_projects_the_ref_records_the_http_push_projects() { + let scratch = tempfile::tempdir().unwrap(); + let (key_path, public_line) = keygen(scratch.path(), "client"); + let (server, _index, _shutdown, _serve_task, firehose) = spawn_server_core( + public_line, + MaxWireBytes::new(1 << 30), + ArchiveLimit::default(), + true, + None, + true, + ) + .await; + let firehose = firehose.expect("the projection fixture serves a firehose"); + + let url = format!( + "ssh://git@127.0.0.1:{}/{OWNER_DID}/{REPO_NAME}", + server.port + ); + let work = scratch.path().join("work"); + let head = seed_work(&work); + + let (ok, out) = push(&work, &url, &key_path, &["main"]).await; + assert!(ok, "authorized ssh push must succeed:\n{out}"); + assert_eq!( + main_tip(&server.layout, &server.repo_did), + Some(Oid::from_hex(&head).unwrap()), + "pushed commit must be the repository's main tip" + ); + + let ref_ops: Vec = firehose + .replay( + knot_events::EventCursor::START, + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(4096).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), + ) + .events + .iter() + .flat_map(|event| frame_kind_and_ops(&event.frame).1) + .filter(|path| path.starts_with("sh.tangled.git.ref/")) + .collect(); + assert_eq!( + ref_ops, + vec!["sh.tangled.git.ref/refs~2fheads~2fmain".to_string()], + "the ssh push projected its ref into an atproto record, exactly as the http push does" + ); +} diff --git a/knot2/crates/knot-xrpc/src/atproto.rs b/knot2/crates/knot-xrpc/src/atproto.rs new file mode 100644 index 000000000..ea1e2eb13 --- /dev/null +++ b/knot2/crates/knot-xrpc/src/atproto.rs @@ -0,0 +1,685 @@ +use std::collections::{BTreeMap, BTreeSet}; +use std::str::FromStr; +use std::sync::Arc; + +use axum::body::Body; +use axum::extract::State; +use axum::http::{HeaderValue, StatusCode, header}; +use axum::response::{IntoResponse, Response}; +use bytes::Bytes; +use cid::Cid; +use jacquard_common::types::nsid; +use jacquard_common::types::recordkey::Rkey; +use jacquard_repo::RepoError; +use jacquard_repo::car::write_car_bytes; +use jacquard_repo::storage::BlockStore; +use knot_git::Repo; +use knot_index::Resolved; +use knot_record::blocks::GitBlocks; +use knot_record::chain::{self, ChainTip}; +use knot_runtime::Signer as _; +use knot_runtime::{Clock, HttpTransport}; +use knot_types::{Handle, KnotHostname, RepoDid}; +use serde::Deserialize; +use serde_json::Value; + +use crate::XrpcState; +use crate::error::XrpcError; +use crate::query::{Limit, Offset}; +use crate::run_blocking; + +const CAR_CONTENT_TYPE: &str = "application/vnd.ipld.car"; +const LIST_RECORDS_DEFAULT: usize = 50; +const LIST_RECORDS_MAX: usize = 100; +const LIST_REPOS_DEFAULT: usize = 500; +const LIST_REPOS_MAX: usize = 1000; +const GET_BLOCKS_CIDS_MAX: usize = 1000; + +fn repo_not_found(did: &str) -> XrpcError { + XrpcError::named( + StatusCode::BAD_REQUEST, + XrpcError::REPO_NOT_FOUND, + format!("couldn't find repo: {did}"), + ) +} + +fn record_not_found(uri: &str) -> XrpcError { + XrpcError::named( + StatusCode::BAD_REQUEST, + "RecordNotFound", + format!("couldn't locate record: {uri}"), + ) +} + +fn handle_not_found(handle: &str) -> XrpcError { + XrpcError::named( + StatusCode::BAD_REQUEST, + "HandleNotFound", + format!("couldn't resolve handle to a repo here: {handle}"), + ) +} + +fn storage(error: RepoError) -> XrpcError { + XrpcError::internal(error.to_string()) +} + +/// One atproto repository this knot can serve: +/// the knot's own meta-repo under its `did:web`, +/// or a hosted repository under its repo DID. +struct AtprotoRepo { + did: String, + repo: Repo, +} + +fn resolve( + state: &XrpcState, + repo: &str, +) -> Result { + if repo == state.knot_did.as_str() || repo == state.knot_hostname.as_str() { + return Ok(AtprotoRepo { + did: state.knot_did.as_str().to_owned(), + repo: Repo::open(&state.meta_path)?, + }); + } + let did = RepoDid::new(repo).map_err(|_| repo_not_found(repo))?; + match state.index.ownership_of(&did) { + Resolved::Ready(Some(_)) => Ok(AtprotoRepo { + did: did.as_str().to_owned(), + repo: state.layout.open(&did)?, + }), + Resolved::Ready(None) => Err(repo_not_found(repo)), + Resolved::Warming => Err(XrpcError::warming("registry projection is still warming")), + } +} + +fn tip_of(repo: &Repo, did: &str) -> Result { + chain::tip(repo)?.ok_or_else(|| repo_not_found(did)) +} + +fn car(response_limit: usize, bytes: Vec) -> Result { + if bytes.len() > response_limit { + return Err(XrpcError::request_too_large(format!( + "the repository export is {} bytes over the {} byte response limit", + bytes.len(), + response_limit + ))); + } + let mut response = Response::new(Body::from(bytes)); + response.headers_mut().insert( + header::CONTENT_TYPE, + HeaderValue::from_static(CAR_CONTENT_TYPE), + ); + Ok(response) +} + +fn json(value: impl serde::Serialize, limit: usize) -> Result { + let bytes = + serde_json::to_vec(&value).map_err(|error| XrpcError::internal(error.to_string()))?; + if bytes.len() > limit { + return Err(XrpcError::request_too_large(format!( + "the response is {} bytes over the {} byte limit", + bytes.len(), + limit + ))); + } + Ok(([(header::CONTENT_TYPE, "application/json")], bytes).into_response()) +} + +async fn serve_car( + found: AtprotoRepo, + response_limit: usize, + collect: R, +) -> Result +where + R: FnOnce(&Repo, &ChainTip) -> Result, XrpcError> + Send + 'static, +{ + run_blocking(move || { + let tip = tip_of(&found.repo, &found.did)?; + let commit_bytes = chain::commit_bytes(&found.repo, &tip)?; + let commit_cid = chain::commit_cid(&commit_bytes)?; + let mut blocks = collect(&found.repo, &tip)?; + blocks.insert(commit_cid, Bytes::from(commit_bytes)); + let bytes = futures::executor::block_on(write_car_bytes(commit_cid, blocks)) + .map_err(|error| XrpcError::internal(error.to_string()))?; + car(response_limit, bytes) + }) + .await +} + +#[derive(Deserialize)] +pub(crate) struct DidParams { + pub(crate) did: String, +} + +pub(crate) async fn sync_get_repo( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + let found = resolve(&state, ¶ms.did)?; + serve_car(found, state.byte_limits.response.get(), |repo, tip| { + let tree = chain::load_tree(repo, tip); + let store = GitBlocks::at(repo, tip.blocks); + let root = tip.data; + futures::executor::block_on(async move { + let mut cids = tree.collect_node_cids().await?; + cids.push(root); + cids.extend(tree.leaves().await?.into_iter().map(|(_, cid)| cid)); + let mut blocks = BTreeMap::new(); + for cid in cids { + insert_or_corrupt(&store, cid, &mut blocks).await?; + } + Ok::, RepoError>(blocks) + }) + .map_err(|error| XrpcError::internal(error.to_string())) + }) + .await +} + +pub(crate) async fn sync_get_latest_commit( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + let found = resolve(&state, ¶ms.did)?; + let response_limit = state.byte_limits.response.get(); + run_blocking(move || { + let tip = tip_of(&found.repo, &found.did)?; + let cid = chain::commit_cid(&chain::commit_bytes(&found.repo, &tip)?)?; + json( + serde_json::json!({ "cid": cid.to_string(), "rev": tip.rev.to_string() }), + response_limit, + ) + }) + .await +} + +#[derive(Deserialize)] +pub(crate) struct SyncRecordParams { + did: String, + collection: String, + rkey: String, +} + +pub(crate) async fn sync_get_record( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + checked_collection(¶ms.collection)?; + checked_rkey(¶ms.rkey)?; + let found = resolve(&state, ¶ms.did)?; + let key = format!("{}/{}", params.collection, params.rkey); + serve_car(found, state.byte_limits.response.get(), |repo, tip| { + let tree = chain::load_tree(repo, tip); + let store = GitBlocks::at(repo, tip.blocks); + futures::executor::block_on(async move { + let leaves = tree.leaves().await?; + let mut keys: Vec<&str> = leaves.iter().map(|(key, _)| key.as_str()).collect(); + keys.sort_unstable(); + let mut wanted = neighbors_of(&keys, &key); + wanted.push(key.clone()); + let mut blocks = BTreeMap::new(); + for probe in wanted { + tree.blocks_for_path(&probe, &mut blocks).await?; + } + if let Some(cid) = tree.get(&key).await? { + insert_or_corrupt(&store, cid, &mut blocks).await?; + } + Ok::, RepoError>(blocks) + }) + .map_err(|error| XrpcError::internal(error.to_string())) + }) + .await +} + +fn neighbors_of(sorted: &[&str], key: &str) -> Vec { + let position = sorted.partition_point(|existing| *existing < key); + [ + position + .checked_sub(1) + .and_then(|before| sorted.get(before)), + sorted.get(position), + ] + .into_iter() + .flatten() + .map(|neighbor| (*neighbor).to_owned()) + .collect() +} + +fn checked_collection(collection: &str) -> Result<(), XrpcError> { + nsid::validate_nsid(collection) + .map_err(|_| XrpcError::invalid_request(format!("invalid collection nsid: {collection}"))) +} + +fn checked_rkey(rkey: &str) -> Result<(), XrpcError> { + Rkey::::new_owned(rkey) + .map(|_| ()) + .map_err(|_| XrpcError::invalid_request(format!("invalid record key: {rkey}"))) +} + +async fn insert_or_corrupt( + store: &GitBlocks, + cid: Cid, + blocks: &mut BTreeMap, +) -> jacquard_repo::error::Result<()> { + let bytes = store.get(&cid).await?.ok_or_else(|| { + RepoError::storage(std::io::Error::other(format!( + "the chain names block {cid} and the block store doesn't hold it" + ))) + })?; + blocks.insert(cid, bytes); + Ok(()) +} + +pub(crate) struct BlocksQuery { + pub(crate) did: String, + pub(crate) cids: Vec, +} + +impl axum::extract::FromRequestParts for BlocksQuery { + type Rejection = XrpcError; + + async fn from_request_parts( + parts: &mut http::request::Parts, + _state: &S, + ) -> Result { + let pairs = url::form_urlencoded::parse(parts.uri.query().unwrap_or_default().as_bytes()); + let did = pairs + .clone() + .find(|(key, _)| key == "did") + .map(|(_, value)| value.to_string()) + .ok_or_else(|| XrpcError::invalid_request("did is required"))?; + let cids: Vec = pairs + .filter(|(key, _)| key == "cids") + .map(|(_, value)| value.to_string()) + .collect(); + if cids.len() > GET_BLOCKS_CIDS_MAX { + return Err(XrpcError::invalid_request(format!( + "at most {GET_BLOCKS_CIDS_MAX} cids may be requested per call" + ))); + } + Ok(Self { did, cids }) + } +} + +pub(crate) async fn sync_get_blocks( + State(state): State>>, + params: BlocksQuery, +) -> Result { + let found = resolve(&state, ¶ms.did)?; + let wanted: Vec = params + .cids + .iter() + .map(|cid| Cid::from_str(cid)) + .collect::>() + .map_err(|_| XrpcError::invalid_request("a requested cid doesn't parse"))?; + serve_car(found, state.byte_limits.response.get(), |repo, tip| { + let store = GitBlocks::at(repo, tip.blocks); + futures::executor::block_on(async move { + let mut blocks = BTreeMap::new(); + for cid in &wanted { + if let Some(bytes) = store.get(cid).await? { + blocks.insert(*cid, bytes); + } + } + Ok::, RepoError>(blocks) + }) + .map_err(|error| XrpcError::internal(error.to_string())) + }) + .await +} + +#[derive(Deserialize)] +pub(crate) struct ListReposParams { + #[serde(default)] + limit: Limit, + #[serde(default)] + cursor: Offset, +} + +pub(crate) async fn sync_list_repos( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + if matches!( + state.index.coverage().registry, + knot_index::Coverage::Warming + ) { + return Err(XrpcError::warming("registry projection is still warming")); + } + let dids: Vec = std::iter::once(state.knot_did.as_str().to_owned()) + .chain( + state + .index + .hosted_repos() + .into_iter() + .map(|did| did.as_str().to_owned()), + ) + .collect(); + let cursor = next_offset(params.cursor.get(), params.limit.get(), dids.len()); + let page: Vec = dids + .into_iter() + .skip(params.cursor.get()) + .take(params.limit.get()) + .collect(); + let response_limit = state.byte_limits.response.get(); + run_blocking(move || { + let repos: Vec = page + .iter() + .filter_map(|did| { + resolve(&state, did) + .and_then(|found| { + tip_of(&found.repo, &found.did).and_then(|tip| { + let head = chain::commit_cid(&chain::commit_bytes(&found.repo, &tip)?)?; + Ok(serde_json::json!({ + "did": did, + "head": head.to_string(), + "rev": tip.rev.to_string(), + "active": true, + })) + }) + }) + .map_err(|error| { + tracing::warn!( + did, + %error, + "listRepos left a repo out: it failed to open or read" + ); + error + }) + .ok() + }) + .collect(); + json( + ListReposOut { + repos, + cursor: cursor.map(|taken| taken.to_string()), + }, + response_limit, + ) + }) + .await +} + +#[derive(serde::Serialize)] +struct ListReposOut { + repos: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + cursor: Option, +} + +#[derive(serde::Serialize)] +struct RepoStatusOut { + did: String, + active: bool, + #[serde(skip_serializing_if = "Option::is_none")] + status: Option<&'static str>, + #[serde(skip_serializing_if = "Option::is_none")] + rev: Option, +} + +fn next_offset(offset: usize, limit: usize, total: usize) -> Option { + let taken = offset + limit; + (taken < total).then_some(taken) +} + +pub(crate) async fn sync_get_repo_status( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + let response_limit = state.byte_limits.response.get(); + match resolve(&state, ¶ms.did) { + Ok(found) => { + run_blocking(move || { + let rev = chain::tip(&found.repo)?.map(|tip| tip.rev.to_string()); + json( + RepoStatusOut { + did: found.did, + active: true, + status: None, + rev, + }, + response_limit, + ) + }) + .await + } + Err(error) if error.is_repo_not_found() => { + match crate::plc::is_deleted(&state, ¶ms.did).await { + Ok(true) => json( + RepoStatusOut { + did: params.did, + active: false, + status: Some(crate::plc::STATUS_DELETED), + rev: None, + }, + response_limit, + ), + Ok(false) => Err(error), + Err(tombstone_error) => Err(tombstone_error), + } + } + Err(error) => Err(error), + } +} + +#[derive(Deserialize)] +pub(crate) struct RepoRecordParams { + repo: String, + collection: String, + rkey: String, +} + +pub(crate) async fn repo_get_record( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + checked_collection(¶ms.collection)?; + checked_rkey(¶ms.rkey)?; + let found = resolve(&state, ¶ms.repo)?; + let key = format!("{}/{}", params.collection, params.rkey); + let uri = format!("at://{}/{key}", found.did); + let response_limit = state.byte_limits.response.get(); + run_blocking(move || { + let tip = tip_of(&found.repo, &found.did)?; + let tree = chain::load_tree(&found.repo, &tip); + let store = GitBlocks::at(&found.repo, tip.blocks); + let (cid, value) = futures::executor::block_on(async { + let cid = tree + .get(&key) + .await + .map_err(storage)? + .ok_or_else(|| record_not_found(&uri))?; + let bytes = store.get(&cid).await.map_err(storage)?.ok_or_else(|| { + XrpcError::internal(format!( + "the chain names {uri} and the block store doesn't hold its block" + )) + })?; + let value: Value = serde_ipld_dagcbor::from_slice(&bytes) + .map_err(|error| XrpcError::internal(error.to_string()))?; + Result::<(Cid, Value), XrpcError>::Ok((cid, value)) + })?; + json( + serde_json::json!({ + "uri": uri, + "cid": cid.to_string(), + "value": value, + }), + response_limit, + ) + }) + .await +} + +#[derive(Deserialize)] +pub(crate) struct ListRecordsParams { + repo: String, + collection: String, + #[serde(default)] + limit: Limit, + cursor: Option, + #[serde(default)] + reverse: bool, +} + +pub(crate) async fn repo_list_records( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + checked_collection(¶ms.collection)?; + let found = resolve(&state, ¶ms.repo)?; + let prefix = format!("{}/", params.collection); + let limit = params.limit; + let reverse = params.reverse; + let response_limit = state.byte_limits.response.get(); + run_blocking(move || { + let tip = tip_of(&found.repo, &found.did)?; + let tree = chain::load_tree(&found.repo, &tip); + let store = GitBlocks::at(&found.repo, tip.blocks); + let records = futures::executor::block_on(async { + let mut matching: Vec<(String, Cid)> = tree + .leaves() + .await + .map_err(storage)? + .into_iter() + .filter_map(|(key, cid)| { + Some((key.as_str().strip_prefix(&prefix)?.to_owned(), cid)) + }) + .collect(); + matching.sort_by(|a, b| a.0.cmp(&b.0)); + if reverse { + matching.reverse(); + } + let start = params.cursor.as_ref().map_or(0, |after| { + matching.partition_point(|(rkey, _)| { + if reverse { + rkey.as_str() >= after.as_str() + } else { + rkey.as_str() <= after.as_str() + } + }) + }); + let mut page: Vec<(String, Cid, Value)> = Vec::new(); + for (rkey, cid) in matching.into_iter().skip(start).take(limit.get() + 1) { + let bytes = store.get(&cid).await.map_err(storage)?.ok_or_else(|| { + XrpcError::internal(format!( + "the chain names record {}{rkey} and the block store doesn't hold it", + prefix + )) + })?; + let value: Value = serde_ipld_dagcbor::from_slice(&bytes) + .map_err(|error| XrpcError::internal(error.to_string()))?; + page.push((rkey, cid, value)); + } + Result::, XrpcError>::Ok(page) + })?; + let next = records + .get(limit.get() - 1) + .filter(|_| records.len() > limit.get()) + .map(|(rkey, _, _)| rkey.clone()); + let items: Vec = records + .into_iter() + .take(limit.get()) + .map(|(rkey, cid, value)| { + serde_json::json!({ + "uri": format!("at://{}/{prefix}{rkey}", found.did), + "cid": cid.to_string(), + "value": value, + }) + }) + .collect(); + json( + serde_json::json!({ "records": items, "cursor": next }), + response_limit, + ) + }) + .await +} + +#[derive(Deserialize)] +pub(crate) struct RepoQuery { + repo: String, +} + +pub(crate) async fn repo_describe_repo( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + let found = resolve(&state, ¶ms.repo)?; + let knot = found.did == state.knot_did.as_str(); + let handle = if knot { + state.knot_hostname.as_str().to_owned() + } else { + String::new() + }; + let did_doc = knot + .then(|| { + let key = state.secrets.signer(&state.knot_did)?.public_key(); + Ok::<_, XrpcError>(knot_atproto::knot_did_document( + &state.knot_did, + &key, + &state.knot_service_url, + )) + }) + .transpose()?; + let response_limit = state.byte_limits.response.get(); + run_blocking(move || { + let collections = chain::tip(&found.repo)? + .map(|tip| chain::load_tree(&found.repo, &tip)) + .map(|tree| { + futures::executor::block_on(async { + let names: BTreeSet = tree + .leaves() + .await? + .iter() + .filter_map(|(key, _)| key.as_str().split('/').next().map(String::from)) + .collect(); + Ok::, RepoError>(names.into_iter().collect()) + }) + }) + .transpose() + .map_err(storage)? + .unwrap_or_default(); + let mut body = serde_json::json!({ + "did": found.did, + "handle": handle, + "handleIsCorrect": false, + "collections": collections, + }); + if let Some(document) = did_doc { + body["didDoc"] = document; + } + json(body, response_limit) + }) + .await +} + +pub(crate) async fn describe_server( + State(state): State>>, +) -> Result { + json( + serde_json::json!({ + "did": state.knot_did.as_str(), + "availableUserDomains": [], + }), + state.byte_limits.response.get(), + ) +} + +#[derive(Deserialize)] +pub(crate) struct HandleParams { + handle: String, +} + +pub(crate) async fn resolve_handle( + State(state): State>>, + axum::extract::Query(params): axum::extract::Query, +) -> Result { + let text = params.handle; + let handle = Handle::new(text.clone()).map_err(|_| handle_not_found(&text))?; + if handle.as_str() == KnotHostname::as_str(&state.knot_hostname) { + return json( + serde_json::json!({ "did": state.knot_did.as_str() }), + state.byte_limits.response.get(), + ); + } + Err(handle_not_found(handle.as_str())) +} diff --git a/knot2/crates/knot-xrpc/src/error.rs b/knot2/crates/knot-xrpc/src/error.rs index 0652aa496..5c08b84f2 100644 --- a/knot2/crates/knot-xrpc/src/error.rs +++ b/knot2/crates/knot-xrpc/src/error.rs @@ -18,6 +18,11 @@ impl XrpcError { message: message.into(), } } + pub(crate) const REPO_NOT_FOUND: &str = "RepoNotFound"; + + pub(crate) fn is_repo_not_found(&self) -> bool { + self.error == Self::REPO_NOT_FOUND + } pub fn invalid_request(message: impl Into) -> Self { Self::new(StatusCode::BAD_REQUEST, "InvalidRequest", message) @@ -147,6 +152,24 @@ impl From for XrpcError { } } +impl From for XrpcError { + fn from(error: knot_record::chain::ChainError) -> Self { + use knot_record::chain::ChainError; + let message = error.to_string(); + match error { + ChainError::RecordTooLarge(_) | ChainError::ChunkTooLarge => { + Self::request_too_large(message) + } + ChainError::NotEnabled => Self::not_found(message), + ChainError::AlreadyMaterialized => Self::conflict(message), + ChainError::Record(_) + | ChainError::Block(_) + | ChainError::Git(_) + | ChainError::Malformed(_) => Self::internal(message), + } + } +} + impl From for XrpcError { fn from(error: knot_git::ApplyError) -> Self { use knot_git::ApplyError; @@ -157,6 +180,15 @@ impl From for XrpcError { } } +impl From for XrpcError { + fn from(error: crate::materialize::EnableError) -> Self { + match error { + crate::materialize::EnableError::Chain(chain) => chain.into(), + other => Self::internal(other.to_string()), + } + } +} + impl From for XrpcError { fn from(error: knot_cob::CobError) -> Self { use knot_cob::CobError; diff --git a/knot2/crates/knot-xrpc/src/tests.rs b/knot2/crates/knot-xrpc/src/tests.rs index 8fe3c4d4b..f81ae9b04 100644 --- a/knot2/crates/knot-xrpc/src/tests.rs +++ b/knot2/crates/knot-xrpc/src/tests.rs @@ -1,7 +1,9 @@ +use parking_lot::Mutex; +use parking_lot::RwLock; use std::collections::{BTreeSet, HashMap}; use std::path::PathBuf; +use std::sync::Arc; use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; -use std::sync::{Arc, Mutex}; use axum::body::Bytes; use axum::extract::State; @@ -106,16 +108,20 @@ fn did_web_doc(did: &str, sec1: &[u8], pds: &str) -> Bytes { .unwrap(), ) } - -fn repo_did_doc(did: &str, multikey: &str) -> Bytes { +fn repo_account_doc(did: &str, multikey: &str, pds: &str) -> Bytes { Bytes::from( serde_json::to_vec(&json!({ "id": did, "verificationMethod": [{ - "id": format!("{did}#repo"), + "id": format!("{did}#atproto"), "type": "Multikey", "controller": did, "publicKeyMultibase": multikey + }], + "service": [{ + "id": "#atproto_pds", + "type": "AtprotoPersonalDataServer", + "serviceEndpoint": pds }] })) .unwrap(), @@ -175,6 +181,20 @@ async fn json_of(response: Response) -> serde_json::Value { serde_json::from_slice(&bytes).unwrap() } +async fn repo_status(world: &World, did: &RepoDid) -> serde_json::Value { + json_of( + crate::atproto::sync_get_repo_status( + State(Arc::clone(&world.state)), + axum::extract::Query(crate::atproto::DidParams { + did: did.as_str().to_owned(), + }), + ) + .await + .unwrap(), + ) + .await +} + async fn call( world: &World, handler: F, @@ -248,50 +268,6 @@ fn resolve(world: &World, rkey: &str) -> Resolved> { .resolve_repo(&member_owner(), &RepoRkey::new(rkey).unwrap()) } -fn replay(world: &World) -> Vec> { - world - .state - .events - .replay( - knot_events::EventCursor::START, - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(64).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - ) - .events -} - -fn event_count(world: &World) -> usize { - replay(world).len() -} - -fn last_event(world: &World, nsid: &str) -> std::sync::Arc { - replay(world) - .into_iter() - .rev() - .find(|event| event.nsid == nsid) - .unwrap_or_else(|| panic!("{nsid} event is emitted")) -} - -fn git_events(world: &World) -> Vec> { - replay(world) - .into_iter() - .filter(|event| { - !matches!( - event.nsid, - "sh.tangled.knot.memberUpdate" | "sh.tangled.repo.collaboratorUpdate" - ) - }) - .collect() -} - -fn only_git_event(world: &World) -> std::sync::Arc { - let mut events = git_events(world); - assert_eq!(events.len(), 1, "expected exactly one non-acl event"); - events.remove(0) -} - fn bootstrap( dir: &TempDir, rebuild: bool, @@ -367,13 +343,6 @@ fn state_from( git_http, pack_limits: knot_pack::PackLimits::default(), service_owner: account(ADMIN_HOST), - events: Arc::new(knot_events::EventLog::new( - Arc::clone(&clock), - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(1024).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - )), subscriber_gate: Arc::new(knot_events::SubscriberGate::new( knot_events::GlobalSubscriberLimit::new(16), knot_events::PerPeerSubscriberLimit::new(8), @@ -383,6 +352,13 @@ fn state_from( slots: knot_resource::Slots::testing(8), lfs: None, catalog: Arc::new(knot_messages::Catalog::defaults()), + firehose: Arc::new(knot_events::EventLog::new( + Arc::clone(&clock), + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(1024).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), + )), }); (state, clock) } @@ -410,7 +386,7 @@ fn world_responder( .find(|(key, _)| key == "rkey") .map(|(_, value)| value.into_owned()) .unwrap_or_default(); - let answer = pds_records.lock().unwrap().get(&rkey).copied(); + let answer = pds_records.lock().get(&rkey).copied(); return Ok(HttpResponse { status: answer.unwrap_or(StatusCode::BAD_REQUEST), headers: http::HeaderMap::new(), @@ -421,11 +397,15 @@ fn world_responder( }); } let host = request.url.host_str().unwrap_or_default(); - if let Some(multikey) = repo_docs.lock().unwrap().get(host).cloned() { + if let Some(multikey) = repo_docs.lock().get(host).cloned() { return Ok(HttpResponse { status: StatusCode::OK, headers: http::HeaderMap::new(), - body: repo_did_doc(&format!("did:web:{host}"), &multikey), + body: repo_account_doc( + &format!("did:web:{host}"), + &multikey, + &format!("https://{KNOT_HOST}"), + ), }); } match pubkeys.get(host) { @@ -579,7 +559,6 @@ impl World { fn publish_repo_doc(&self, host: &str, multikey: &str) { self.repo_docs .lock() - .unwrap() .insert(host.to_string(), multikey.to_string()); } @@ -588,10 +567,7 @@ impl World { } fn answer_pds_record(&self, rkey: &str, status: StatusCode) { - self.pds_records - .lock() - .unwrap() - .insert(rkey.to_string(), status); + self.pds_records.lock().insert(rkey.to_string(), status); } fn advance(&self, secs: u64) { @@ -628,7 +604,10 @@ impl World { match store.list::().unwrap().as_slice() { [] => _ = store.create(&home, change, &signer, at).unwrap(), [object] => _ = store.update(&home, *object, change, &signer, at).unwrap(), - many => panic!("{} roster objects under a namespace that must hold one", many.len()), + many => panic!( + "{} roster objects under a namespace that must hold one", + many.len() + ), } } @@ -667,11 +646,19 @@ impl World { } fn build_state(responder: Responder, rebuild: bool) -> (TempDir, SharedState) { + let (dir, state, _clock) = build_state_with_clock(responder, rebuild); + (dir, state) +} + +fn build_state_with_clock( + responder: Responder, + rebuild: bool, +) -> (TempDir, SharedState, Arc) { let dir = tempfile::tempdir().unwrap(); - let (layout, index, meta_path) = bootstrap(&dir, rebuild, knot_types::ObjectFormat::SHA1); - let (state, _clock) = state_from( + let boot = bootstrap(&dir, rebuild, knot_types::ObjectFormat::SHA1); + let (state, clock) = state_from( &dir, - (layout, index, meta_path), + boot, responder, AdmissionPolicy::Closed, Arc::new(crate::Reservations::new( @@ -681,7 +668,7 @@ fn build_state(responder: Responder, rebuild: bool) -> (TempDir, SharedState) { )), no_git_upstream(), ); - (dir, state) + (dir, state, clock) } fn no_git_upstream() -> Arc { @@ -712,6 +699,128 @@ fn doc_responder( }) } +#[derive(Default)] +struct PlcFake { + reject_posts: bool, + ops: Vec<(String, serde_json::Value)>, + post_count: usize, + tombstones: std::collections::HashSet, + last_override: Option, +} + +fn reply(status: StatusCode, body: Bytes) -> HttpResponse { + HttpResponse { + status, + headers: http::HeaderMap::new(), + body, + } +} + +fn plc_stage(fake: Arc>, sec1: Vec) -> Responder { + let pds = format!("https://{KNOT_HOST}"); + Box::new(move |request: &HttpRequest| { + let host = request.url.host_str().unwrap_or_default(); + if host != "plc.directory" { + return Ok(reply( + StatusCode::OK, + did_web_doc(&format!("did:web:{host}"), &sec1, &pds), + )); + } + let mut fake = fake.lock(); + if request.method == http::Method::POST { + fake.post_count += 1; + if fake.reject_posts { + return Ok(reply(StatusCode::BAD_REQUEST, Bytes::new())); + } + let did = request.url.path().trim_start_matches('/').to_owned(); + let operation = + serde_json::from_slice(request.body.as_ref().expect("post body")).unwrap(); + fake.ops.push((did, operation)); + return Ok(reply(StatusCode::OK, Bytes::new())); + } + let path = request.url.path().to_owned(); + match path + .strip_suffix("/log/last") + .map(|did| did.trim_start_matches('/')) + { + None => Ok(reply(StatusCode::NOT_FOUND, Bytes::new())), + Some(did) => { + if fake.tombstones.contains(did) { + return Ok(reply( + StatusCode::OK, + Bytes::from( + serde_json::to_vec(&json!({ + "type": "plc_tombstone", + "sig": "AAAA", + })) + .unwrap(), + ), + )); + } + let last = match &fake.last_override { + Some(op) => Some(op.clone()), + None => fake + .ops + .iter() + .rev() + .find(|(entry_did, _)| entry_did == did) + .map(|(_, op)| op.clone()), + }; + match last { + None => Ok(reply(StatusCode::NOT_FOUND, Bytes::new())), + Some(op) => Ok(reply( + StatusCode::OK, + Bytes::from(serde_json::to_vec(&op).unwrap()), + )), + } + } + } + }) +} + +struct SweptRepo { + _dir: TempDir, + state: SharedState, + clock: Arc, + did: RepoDid, + fake: Arc>, + create_status: StatusCode, +} + +async fn minted_repo(fake: PlcFake) -> SweptRepo { + let admin = actor(1, ADMIN_HOST); + let key = admin.signer.public_key().as_bytes().to_vec(); + let fake = Arc::new(Mutex::new(fake)); + let responder = plc_stage(Arc::clone(&fake), key); + let (_dir, state, clock) = build_state_with_clock(responder, true); + let create_status = into_response( + crate::repos::create_repo( + State(Arc::clone(&state)), + bearer(&mint(&admin, CREATE)), + crate::Method::from_nsid(CREATE), + body(json!({ "rkey": "anemone", "name": "anemone" })), + ) + .await, + ) + .status(); + let owner = OwnerDid::new(format!("did:web:{ADMIN_HOST}")).unwrap(); + let did = match state + .index + .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()) + { + Resolved::Ready(Some(did)) => did, + other => panic!("repo wasn't registered despite the queue: {other:?}"), + }; + SweptRepo { + _dir, + state, + clock, + did, + fake, + create_status, + } +} + fn subject_body(repo: Option<&RepoDid>, subject: &AccountDid) -> serde_json::Value { match repo { None => json!({ "subject": subject.as_str() }), @@ -940,37 +1049,17 @@ async fn member_lifecycle() { Resolved::Ready(Some(Standing::Invited)), "only the knot has signed this membership, so the account stands Invited" ); - assert_eq!( - event_count(&world), - 0, - "an unsigned offer emitted an event that acl consumers read as a grant" - ); - assert_eq!(answer(&world, &world.member, None).await, StatusCode::OK); assert_eq!( world.state.index.effective_member(&account(MEMBER_HOST)), Resolved::Ready(true), - "member is effective on the very next read, with no firehose" + "the member accepted, so the account stands effective" ); - let events = replay(&world); - let [added] = events.as_slice() else { - panic!("expected exactly one event, got {}", events.len()); - }; - assert_eq!(added.nsid, "sh.tangled.knot.memberUpdate"); - assert_eq!(added.payload["op"], "add"); - assert_eq!(added.payload["subject"], account(MEMBER_HOST).to_string()); - - let baseline = event_count(&world); assert_eq!( offer(&world, None, &member).await, StatusCode::OK, "the knot refused a second offer for a member it already added" ); - assert_eq!( - event_count(&world), - baseline, - "a redundant offer appended a second add to the roster" - ); assert_eq!( as_admin( @@ -988,11 +1077,6 @@ async fn member_lifecycle() { Resolved::Ready(false), "removed member is gone on the very next read" ); - let removed = last_event(&world, "sh.tangled.knot.memberUpdate"); - assert_eq!(removed.payload["op"], "remove"); - assert_eq!(removed.payload["subject"], account(MEMBER_HOST).to_string()); - - let baseline = event_count(&world); assert_eq!( as_admin( &world, @@ -1004,11 +1088,6 @@ async fn member_lifecycle() { .status(), StatusCode::OK ); - assert_eq!( - event_count(&world), - baseline, - "removing a non-member is a no-op and emits no event" - ); } #[tokio::test] @@ -1440,134 +1519,66 @@ async fn a_byo_did_web_repo_is_accepted_and_its_key_is_returned() { } #[tokio::test] -async fn a_rejected_plc_submission_is_a_bad_gateway() { - let admin = actor(1, ADMIN_HOST); - let key = admin.signer.public_key().as_bytes().to_vec(); - let responder = doc_responder(key, || StatusCode::BAD_REQUEST); - let (_dir, state) = build_state(responder, true); - let token = mint(&admin, CREATE); - let status = into_response( - crate::repos::create_repo( - State(Arc::clone(&state)), - bearer(&token), - crate::Method::from_nsid(CREATE), - body(json!({ "rkey": "conch", "name": "conch" })), - ) - .await, - ) - .status(); - assert_eq!( - status, - StatusCode::BAD_GATEWAY, - "non-transient PLC rejection surfaces as 502, distinct from an internal 500" - ); - assert_eq!( - state.secrets.len(), - 1, - "rejected PLC submission rolls back the minted repo key, leaving only the knot's own. Unpublished did:plc orphans nothing" - ); - assert_eq!( - state.index.resolve_repo( - &OwnerDid::new(format!("did:web:{ADMIN_HOST}")).unwrap(), - &RepoRkey::new("conch").unwrap() - ), - Resolved::Ready(None), - "repo isn't registered after a rejected PLC submission" - ); -} - -#[tokio::test] -async fn a_rejected_plc_submission_never_touches_the_registry() { - let admin = actor(1, ADMIN_HOST); - let key = admin.signer.public_key().as_bytes().to_vec(); - let reject_posts = Arc::new(std::sync::atomic::AtomicBool::new(false)); - let reject = Arc::clone(&reject_posts); - let responder = doc_responder(key, move || { - if reject.load(Ordering::Relaxed) { - StatusCode::BAD_REQUEST - } else { - StatusCode::OK - } - }); - let (_dir, state) = build_state(responder, true); - let owner = OwnerDid::new(format!("did:web:{ADMIN_HOST}")).unwrap(); - - let mint_create = || mint(&admin, CREATE); - let anemone = || body(json!({ "rkey": "anemone", "name": "anemone" })); +async fn a_rejected_genesis_stays_queued_and_the_queue_lands_both_operations() { + let repo = minted_repo(PlcFake { + reject_posts: true, + ..PlcFake::default() + }) + .await; assert_eq!( - into_response( - crate::repos::create_repo( - State(Arc::clone(&state)), - bearer(&mint_create()), - crate::Method::from_nsid(CREATE), - anemone(), - ) - .await - ) - .status(), - StatusCode::OK + repo.create_status, + StatusCode::OK, + "a directory outage mustn't fail the create; submission is queued" ); - let victim = match state - .index - .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()) { - Resolved::Ready(Some(did)) => did, - other => panic!("victim repo wasn't registered: {other:?}"), - }; + let meta = Repo::open(&repo.state.meta_path).unwrap(); + let pending = crate::plc::pending_entries(&meta).unwrap(); + assert_eq!(pending.len(), 1, "the genesis op sits in the pending queue"); + assert_eq!(pending[0].did, repo.did.as_str()); + } - let token = mint(&admin, RENAME); - assert_eq!( - into_response( - crate::repos::rename_repo( - State(Arc::clone(&state)), - bearer(&token), - crate::Method::from_nsid(RENAME), - body(json!({ "repo": victim.as_str(), "rkey": "barnacle", "name": "barnacle" })), - ) - .await - ) - .status(), - StatusCode::OK - ); + let mut memory = crate::plc::SweepMemory::default(); + repo.fake.lock().reject_posts = false; + for _ in 0..4 { + repo.clock.advance(std::time::Duration::from_secs(31)); + crate::plc::sweep_tick(&repo.state, &mut memory).await; + } - reject_posts.store(true, Ordering::Relaxed); + let snapshot = repo.fake.lock(); assert_eq!( - into_response( - crate::repos::create_repo( - State(Arc::clone(&state)), - bearer(&mint_create()), - crate::Method::from_nsid(CREATE), - anemone(), - ) - .await - ) - .status(), - StatusCode::BAD_GATEWAY - ); + snapshot.post_count, 2, + "exactly genesis and update should POST; ops: {:?}", + snapshot.ops + ); + drop(snapshot); + assert_eq!(memory.settled.len(), 1, "the document settles"); + let meta = Repo::open(&repo.state.meta_path).unwrap(); + let pending = crate::plc::pending_entries(&meta).unwrap(); + assert!(pending.is_empty(), "queue must drain; entries: {pending:?}"); +} +#[tokio::test] +async fn deleting_a_repo_records_the_deactivation_for_get_repo_status() { + let world = World::new(); + world.grandfather(MEMBER_HOST); + let repo = create_repo_helper(&world, "reef").await; - assert_eq!( - state - .index - .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()), - Resolved::Ready(Some(victim.clone())), - "stale alias still resolves to its prior holder because the failed create never registered" - ); + let live = repo_status(&world, &repo).await; + assert_eq!(live["active"], json!(true)); + assert!(live.get("status").is_none()); - state.index.refresh_registry().unwrap(); - assert_eq!( - state - .index - .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()), - Resolved::Ready(Some(victim.clone())), - "durable registry COB has no trace of the failed create" - ); - assert_eq!( - state.index.rkey_of(&victim), - Resolved::Ready(Some(RepoRkey::new("barnacle").unwrap())), - "victim's canonical rkey is unmoved" - ); -} + as_member( + &world, + crate::repos::delete_repo, + DELETE, + json!({ "repo": repo.as_str() }), + ) + .await; + let deleted = repo_status(&world, &repo).await; + assert_eq!(deleted["did"], json!(repo.as_str())); + assert_eq!(deleted["active"], json!(false)); + assert_eq!(deleted["status"], json!("deleted")); +} #[tokio::test] async fn resolve_by_name_matches_the_rkey_case_sensitively() { let admin = actor(1, ADMIN_HOST); @@ -1703,11 +1714,6 @@ async fn collaborator_lifecycle() { Resolved::Ready(Some(Standing::Invited)), "only the repository has signed this collaboration, so the account stands Invited" ); - assert_eq!( - event_count(&world), - 0, - "spindle would mint an rbac grant off this event, and the collaborator never signed" - ); assert_eq!(answer(&world, &world.stranger, repo).await, StatusCode::OK); assert_eq!( @@ -1717,11 +1723,6 @@ async fn collaborator_lifecycle() { .effective_collaborator(&repo_did, &collaborator), Resolved::Ready(true) ); - let added = last_event(&world, "sh.tangled.repo.collaboratorUpdate"); - assert_eq!(added.payload["op"], "add"); - assert_eq!(added.payload["subject"], collaborator.to_string()); - assert_eq!(added.payload["repo"], repo_did.to_string()); - assert_eq!( as_member( &world, @@ -1741,12 +1742,6 @@ async fn collaborator_lifecycle() { Resolved::Ready(false), "removed collaborator is gone on the very next read" ); - let removed = last_event(&world, "sh.tangled.repo.collaboratorUpdate"); - assert_eq!(removed.payload["op"], "remove"); - assert_eq!(removed.payload["subject"], collaborator.to_string()); - assert_eq!(removed.payload["repo"], repo_did.to_string()); - - let baseline = event_count(&world); assert_eq!( as_member( &world, @@ -1758,11 +1753,6 @@ async fn collaborator_lifecycle() { .status(), StatusCode::OK ); - assert_eq!( - event_count(&world), - baseline, - "removing a non-collaborator is a no-op and emits no event" - ); } #[tokio::test] @@ -1862,16 +1852,6 @@ async fn the_owner_sets_the_default_branch_by_repo_did() { Some("refs/heads/trunk".to_string()), "HEAD now points at the requested default branch" ); - - let event = only_git_event(&world); - assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); - assert_eq!(event.payload["repo"], repo_did.to_string()); - assert_eq!(event.payload["ownerDid"], account(MEMBER_HOST).to_string()); - assert_eq!( - event.payload["committerDid"], - account(MEMBER_HOST).to_string(), - "the actor who set the default branch is the committer on the wire" - ); } #[tokio::test] @@ -2000,25 +1980,6 @@ async fn delete_branch_removes_a_branch_and_refuses_the_default() { StatusCode::BAD_REQUEST, "the current default branch cannot be deleted" ); - - let event = only_git_event(&world); - assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); - assert_eq!(event.payload["repo"], repo_did.to_string()); - assert_eq!(event.payload["ref"], "refs/heads/trunk"); - assert_eq!( - event.payload["oldSha"], - oid.to_string(), - "the deletion event includes the branch's old tip" - ); - assert_eq!( - event.payload["newSha"], - git.object_format().null_oid().to_string(), - "deletion reports the null oid as the new sha" - ); - assert_eq!( - event.payload["committerDid"], - account(MEMBER_HOST).to_string() - ); } #[tokio::test] @@ -2507,17 +2468,6 @@ mod merge_endpoints { assert_eq!(commit.message, "Merge tide\n\nbody text\n"); assert_eq!(blob_at(&repo, tip, "reef.txt"), b"new line\n"); - let event = only_git_event(&world); - assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); - assert_eq!(event.payload["repo"], repo_did.to_string()); - assert_eq!(event.payload["ref"], "refs/heads/main"); - assert_eq!(event.payload["oldSha"], base.to_string()); - assert_eq!(event.payload["newSha"], tip.to_string()); - assert_eq!( - event.payload["committerDid"], - account(MEMBER_HOST).to_string() - ); - let staging = std::fs::read_dir(repo.path()) .unwrap() .flatten() @@ -2569,15 +2519,6 @@ mod merge_endpoints { .status(), StatusCode::OK ); - - let update = last_event(&world, "sh.tangled.git.refUpdate"); - assert_eq!(update.payload["ref"], "refs/heads/main"); - assert!( - replay(&world) - .iter() - .all(|event| event.nsid != "sh.tangled.pipeline"), - "the knot emits no pipeline record" - ); } #[tokio::test] @@ -3071,26 +3012,16 @@ mod fork_endpoints { let fork = setup.world.layout.open(&setup.fork_did).unwrap(); assert_eq!(fork.find_ref(&main_ref()).unwrap(), Some(new_tip)); - let event = only_git_event(&setup.world); - assert_eq!(event.nsid, "sh.tangled.git.refUpdate"); - assert_eq!(event.payload["repo"], setup.fork_did.to_string()); - assert_eq!(event.payload["ref"], "refs/heads/main"); - assert_eq!(event.payload["oldSha"], setup.tip.to_string()); - assert_eq!(event.payload["newSha"], new_tip.to_string()); - assert_eq!( - event.payload["committerDid"], - account(MEMBER_HOST).to_string() - ); - assert_eq!( sync_fork(&setup.world, &setup.world.member, &setup.fork_did, "main").await, StatusCode::OK, "an up-to-date sync is a no-op" ); + let fork_after = setup.world.layout.open(&setup.fork_did).unwrap(); assert_eq!( - git_events(&setup.world).len(), - 1, - "an up-to-date sync emits no further event" + fork_after.find_ref(&main_ref()).unwrap(), + Some(new_tip), + "an up-to-date sync moves nothing" ); assert_eq!( @@ -3304,11 +3235,11 @@ mod fork_endpoints { advance(&upstream, &main_ref(), "reef.txt", "kelp forest\n", 1_000); let tip = advance(&upstream, &main_ref(), "tide.txt", "rock pool\n", 1_001); - let served = Arc::new(std::sync::RwLock::new(upstream_path.clone())); + let served = Arc::new(RwLock::new(upstream_path.clone())); let path = Arc::clone(&served); let git_http: Arc = Arc::new(FakeHttp::new(move |request: &HttpRequest| { - let repo = Repo::open(path.read().unwrap().clone()).unwrap(); + let repo = Repo::open(path.read().clone()).unwrap(); let body = if request.url.path().ends_with("/info/refs") { knot_pack::advertise_upload(&repo).unwrap() } else { @@ -3377,7 +3308,7 @@ mod fork_endpoints { let sha256 = Repo::create_with_format(&replaced, knot_types::ObjectFormat::SHA256).unwrap(); sha256.set_head(&main_ref()).unwrap(); advance(&sha256, &main_ref(), "reef.txt", "kelp forest\n", 1_000); - *served.write().unwrap() = replaced; + *served.write() = replaced; assert_eq!( sync_fork(&world, &world.member, &fork_did, "main").await, StatusCode::CONFLICT, @@ -3464,11 +3395,6 @@ mod legacy_admin_route { Resolved::Ready(Some(Standing::Invited)), "the legacy route can only invite, and the acceptance stays with the account" ); - assert_eq!( - event_count(&world), - 0, - "the legacy route published an event ahead of any acceptance" - ); let Resolved::Ready(members) = world.state.index.member_entries() else { panic!("the member roster is warm in this test"); }; @@ -3488,21 +3414,11 @@ mod legacy_admin_route { world.state.index.effective_member(&account(MEMBER_HOST)), Resolved::Ready(true) ); - let added = last_event(&world, "sh.tangled.knot.memberUpdate"); - assert_eq!(added.payload["op"], "add"); - assert_eq!(added.payload["subject"], account(MEMBER_HOST).to_string()); - - let baseline = event_count(&world); assert_eq!( call(&router, "admin", SECRET, &subject).await, StatusCode::OK, "the legacy route is idempotent, matching the Go knot" ); - assert_eq!( - event_count(&world), - baseline, - "the repeat call appended an add for a member already in the roster" - ); } #[tokio::test] @@ -3745,19 +3661,12 @@ mod rosters { "{from}: a refused acceptance must leave the standing at Invited" ); - let before = event_count(&offer.world); assert_eq!(offer.accept(offer.own()).await, StatusCode::OK, "{from}"); let accepted = offer.standing(); assert!( matches!(accepted, Resolved::Ready(Some(Standing::Accepted { .. }))), "{from}: standing must be Accepted with the time the knot verified it" ); - assert_eq!( - event_count(&offer.world), - before + 1, - "{from}: no event for the acceptance, so spindle's ACL stays stale until half B" - ); - offer.answered(StatusCode::BAD_GATEWAY); assert_eq!( offer.accept(offer.own()).await, @@ -3770,11 +3679,6 @@ mod rosters { accepted, "{from}: the retry re-verified and moved the standing" ); - assert_eq!( - event_count(&offer.world), - before + 1, - "{from}: a retry that changed no row must publish no event" - ); let fragment = format!("did:web:{KNOT_HOST}#tangled_knot"); assert_eq!( @@ -3796,14 +3700,12 @@ mod rosters { let from = &invited.from; let world = &invited.world; - let before = event_count(world); let repeated = offer(world, invited.repo(), &invited.subject).await; assert_eq!( - (repeated, event_count(world) - before), - (StatusCode::OK, 0), + repeated, + StatusCode::OK, "{from}: an operator offering again holds the same offer open, appending \ - nothing and announcing nothing an acl consumer would read as a subject \ - that grants" + nothing" ); assert_eq!( invited.standing(), @@ -3811,7 +3713,6 @@ mod rosters { "{from}: the second offer moved the standing off Invited" ); - let before = event_count(world); let body = subject_body(invited.repo(), &invited.subject); let revoked = match invited.repo() { None => as_admin(world, remove_member, REMOVE_MEMBER, body).await, @@ -3819,8 +3720,8 @@ mod rosters { } .status(); assert_eq!( - (revoked, event_count(world) - before), - (StatusCode::OK, 1), + revoked, + StatusCode::OK, "{from}: a revocation is one of the two changes an acl consumer reads" ); assert_eq!( @@ -3925,4 +3826,395 @@ mod rosters { StatusCode::OK ); } + + mod materialization { + use super::*; + use knot_cob::{ChangePayload, CobStore}; + use knot_record::InviteRecord; + + fn members_changes(world: &World) -> Vec { + let meta = Repo::open(world.state.meta_path.clone()).unwrap(); + let store = CobStore::new(&meta); + let object = store.list::().unwrap()[0]; + store + .changes_since::(object, None) + .unwrap() + .changes + } + + fn record_carrying(world: &World) -> knot_cob::Change { + members_changes(world) + .into_iter() + .find(|change| change.record.is_some()) + .expect("a change carries its materialized record") + } + + fn roster_history( + world: &World, + changes: Vec<(MembersChange, UnixSeconds)>, + ) -> knot_cob::CobId { + let meta = Repo::open(world.state.meta_path.clone()).unwrap(); + let store = CobStore::new(&meta); + let home = knot_cob::CobHome::from(&world.state.knot_did); + let signer = world.state.secrets.signer(&world.state.knot_did).unwrap(); + let mut history = changes.into_iter(); + let first = history.next().unwrap(); + let object = store + .create(&home, &first.0, &signer, first.1) + .unwrap() + .object; + let remaining: Vec<(MembersChange, UnixSeconds)> = history.collect(); + if !remaining.is_empty() { + store + .extend( + &home, + object, + remaining.iter().map(|(change, at)| (change, *at)), + &signer, + ) + .unwrap() + .unwrap(); + } + object + } + + #[tokio::test] + async fn an_invite_change_stores_its_record_bytes_in_the_tree() { + let world = World::new(); + assert_eq!( + offer(&world, None, &account(STRANGER_HOST)).await, + StatusCode::OK + ); + let carrying = record_carrying(&world); + let MembersChange::Invite(invite) = MembersChange::decode(carrying.payload()).unwrap() + else { + panic!("the record-carrying change is the invite"); + }; + assert_eq!( + carrying.record.unwrap().as_bytes(), + InviteRecord::new( + knot_record::MEMBER_INVITE_COLLECTION, + invite.0.created_at, + invite.0.added_by.clone(), + ) + .unwrap() + .bytes(), + "stored bytes equal the frozen encoder's output, and reads never re-derive" + ); + } + + #[tokio::test] + async fn a_grandfathered_acceptance_mints_the_lazy_invite_in_the_same_change() { + let world = World::new(); + world.grandfather(MEMBER_HOST); + assert_eq!(answer(&world, &world.member, None).await, StatusCode::OK); + let carrying = record_carrying(&world); + let grant = world.migrated(MEMBER_HOST); + assert_eq!( + carrying.record.as_ref().unwrap().as_bytes(), + InviteRecord::new( + knot_record::MEMBER_INVITE_COLLECTION, + grant.created_at, + grant.added_by.clone(), + ) + .unwrap() + .bytes(), + "the minted-late invite names the admin whose offer it was, at the offer's time" + ); + let MembersChange::Accept(accept) = MembersChange::decode(carrying.payload()).unwrap() + else { + panic!("the record-carrying change is the acceptance"); + }; + assert_eq!(accept.subject, world.member.did); + } + + async fn enable_at( + world: &World, + object: knot_cob::CobId, + ) -> Result, crate::materialize::EnableError> { + let meta = Repo::open(world.state.meta_path.clone()).unwrap(); + let store = CobStore::new(&meta); + let signer = world.state.secrets.signer(&world.state.knot_did).unwrap(); + crate::materialize::upkeep::( + &store, + object, + world.state.knot_did.as_str(), + crate::materialize::MEMBER_INVITE_COLLECTION, + &signer, + &mut jacquard_common::types::tid::Ticker::new(), + UnixSeconds::new(2_000), + ) + .await + } + + fn framed_events(world: &World) -> Vec> { + world + .state + .firehose + .replay( + knot_events::EventCursor::START, + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(64).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), + ) + .events + } + + #[tokio::test] + async fn boot_pass_enables_roster_written_before_chain_existed() { + let world = World::new(); + world.invite(None, "limpet.nel.pet"); + world.invite(None, "whelk.nel.pet"); + let meta = Repo::open(world.state.meta_path.clone()).unwrap(); + assert!( + knot_record::chain::tip(&meta).unwrap().is_none(), + "the roster exists with no chain until something enables it" + ); + crate::materialize::enable_surfaces(&world.state); + let tip = knot_record::chain::tip(&meta) + .unwrap() + .expect("the boot pass materialized the roster"); + let leaves = knot_record::chain::load_tree(&meta, &tip) + .leaves() + .await + .unwrap(); + assert_eq!(leaves.len(), 2); + assert_eq!( + framed_events(&world).len(), + 2, + "the enabled chunk and its closing sync frame reach the firehose without any roster write" + ); + crate::materialize::enable_surfaces(&world.state); + let again = knot_record::chain::tip(&meta).unwrap().unwrap(); + assert_eq!(again.git, tip.git, "a second boot changes nothing"); + } + + #[tokio::test] + async fn a_failed_chain_sync_still_announces_and_refreshes() { + use axum::body::Body; + use axum::extract::ConnectInfo; + use std::net::SocketAddr; + use tower::ServiceExt; + + let world = World::new(); + let app = crate::router(Arc::clone(&world.state)); + let post = |nsid: &'static str, subject: &str| { + let mut request = http::Request::builder() + .method("POST") + .uri(match nsid { + ADD_MEMBER => crate::members::ADD_ROUTE, + _ => crate::members::REMOVE_ROUTE, + }) + .header( + http::header::AUTHORIZATION, + format!("Bearer {}", mint(&world.admin, nsid)), + ) + .body(Body::from( + serde_json::to_vec(&json!({ "subject": subject })).unwrap(), + )) + .unwrap(); + request + .extensions_mut() + .insert(ConnectInfo(SocketAddr::from(([203, 0, 113, 31], 5555)))); + app.clone().oneshot(request) + }; + + let holds = + |world: &World, subject: &str| match world.state.index.member_entries_where( + |standing| matches!(standing, knot_cobs::Standing::Invited), + ) { + knot_index::Resolved::Ready(entries) => { + entries.iter().any(|(did, _)| did.as_str() == subject) + } + knot_index::Resolved::Warming => panic!("the members projection is warming"), + }; + let response = post(ADD_MEMBER, "did:plc:limpet").await.unwrap(); + assert_eq!( + response.status(), + StatusCode::OK, + "the healthy invite lands" + ); + assert!(holds(&world, "did:plc:limpet")); + + let meta = Repo::open(world.state.meta_path.clone()).unwrap(); + let old = knot_record::chain::tip(&meta) + .unwrap() + .expect("the healthy invite materialized the chain") + .git; + let junk = meta.git().write_blob(b"not a commit").unwrap().detach(); + meta.update_ref(&knot_git::RefUpdate::Update { + name: knot_types::RefName::new(knot_record::chain::CHAIN_REF).unwrap(), + old, + new: knot_types::Oid::from(junk), + }) + .unwrap(); + + let response = post(ADD_MEMBER, "did:plc:whelk").await.unwrap(); + assert_eq!( + response.status(), + StatusCode::OK, + "an invite stands even when the chain can't sync" + ); + assert!( + holds(&world, "did:plc:whelk"), + "the projection refreshes even though the chain sync failed" + ); + + let response = post(REMOVE_MEMBER, "did:plc:whelk").await.unwrap(); + assert_eq!( + response.status(), + StatusCode::OK, + "the removal stands even when the chain can't sync" + ); + assert!( + !holds(&world, "did:plc:whelk"), + "the projection drops the removed invite" + ); + } + #[tokio::test] + async fn migration_of_grandfathered_entries_alone_leaves_chain_empty() { + let world = World::new(); + let object = roster_history( + &world, + vec![( + MembersChange::Add(world.migrated("old0.nel.pet")), + UnixSeconds::new(1_000), + )], + ); + let chunks = enable_at(&world, object).await.unwrap(); + assert!( + chunks.is_empty(), + "grandfathering writes no record, so a migration emits no commit at all" + ); + } + + #[tokio::test] + async fn a_frameless_tip_from_before_the_upgrade_gets_a_sync_frame_at_boot() { + let world = World::new(); + world.invite(None, "limpet.nel.pet"); + let meta = Repo::open(world.state.meta_path.clone()).unwrap(); + let store = CobStore::new(&meta); + let [object] = store.list::().unwrap()[..] else { + panic!("one roster object exists"); + }; + enable_at(&world, object).await.unwrap(); + let tip = knot_record::chain::tip(&meta).unwrap().unwrap(); + assert!( + knot_record::chain::unframed_tip(&meta).unwrap().is_some(), + "the enable helper writes its chain with no frame blobs" + ); + crate::materialize::enable_surfaces(&world.state); + let frames = framed_events(&world); + let [event] = &frames[..] else { + panic!("a single sync frame covers the frameless tip"); + }; + let mut read = &event.frame[..]; + let header: serde_json::Value = + serde_ipld_dagcbor::de::from_reader_once(&mut read).unwrap(); + assert_eq!(header["t"], "#sync"); + #[derive(serde::Deserialize)] + struct SyncProbe { + did: String, + rev: String, + #[serde(with = "serde_bytes")] + blocks: Vec, + } + let payload: SyncProbe = serde_ipld_dagcbor::de::from_reader_once(&mut read).unwrap(); + assert_eq!( + payload.rev, + tip.rev.to_string(), + "the sync frame asserts the tip the upgrade inherited" + ); + assert_eq!(payload.did, world.state.knot_did.as_str()); + assert!( + !payload.blocks.is_empty(), + "the sync frame carries the commit block in its CAR" + ); + crate::materialize::enable_surfaces(&world.state); + assert_eq!( + framed_events(&world).len(), + 2, + "each boot reasserts an unframed tip, and no commit accompanies it" + ); + } + } +} + +#[tokio::test] +async fn sweep_parks_tombstoned_repository_before_first_post() { + let repo = minted_repo(PlcFake::default()).await; + repo.fake + .lock() + .tombstones + .insert(repo.did.as_str().to_owned()); + + let mut memory = crate::plc::SweepMemory::default(); + repo.clock.advance(std::time::Duration::from_secs(31)); + crate::plc::sweep_tick(&repo.state, &mut memory).await; + + assert_eq!( + repo.fake.lock().post_count, + 0, + "the sweep reads the tombstone and stops before posting" + ); + assert_eq!(memory.settled.len(), 0); + assert!(memory.is_parked(repo.did.as_str())); +} + +#[tokio::test] +async fn sweep_parks_legacy_rotation_missing_from_sealed_archive() { + let foreign_rotation = "did:key:z6MkhaXgBZDvotDkL5257faiztiGiC2QtKLGpbnnEGta2doK"; + let legacy = json!({ + "type": "plc_operation", + "rotationKeys": [foreign_rotation], + "verificationMethods": {"atproto": foreign_rotation}, + "alsoKnownAs": [], + "services": {"atproto_pds": { + "type": "AtprotoPersonalDataServer", + "endpoint": format!("https://{KNOT_HOST}"), + }}, + "prev": null, + "sig": "AAAA", + }); + let repo = minted_repo(PlcFake { + last_override: Some(legacy), + ..PlcFake::default() + }) + .await; + + let mut memory = crate::plc::SweepMemory::default(); + repo.clock.advance(std::time::Duration::from_secs(31)); + crate::plc::sweep_tick(&repo.state, &mut memory).await; + + assert_eq!( + repo.fake.lock().post_count, + 0, + "the sweep signs only from the sealed archive" + ); + assert_eq!(memory.settled.len(), 0); + assert!(memory.is_parked(repo.did.as_str())); +} + +#[tokio::test] +async fn sweep_parks_permanently_refused_repository_after_one_post() { + let repo = minted_repo(PlcFake { + reject_posts: true, + ..PlcFake::default() + }) + .await; + + let mut memory = crate::plc::SweepMemory::default(); + for _ in 0..2 { + repo.clock.advance(std::time::Duration::from_secs(31)); + crate::plc::sweep_tick(&repo.state, &mut memory).await; + } + + assert_eq!( + repo.fake.lock().post_count, + 1, + "the sweep parks at the first permanent refusal, and the second due pass skips it" + ); + assert_eq!(memory.settled.len(), 0); + assert!(memory.is_parked(repo.did.as_str())); } diff --git a/knot2/crates/knot-xrpc/tests/atproto.rs b/knot2/crates/knot-xrpc/tests/atproto.rs new file mode 100644 index 000000000..dcd431209 --- /dev/null +++ b/knot2/crates/knot-xrpc/tests/atproto.rs @@ -0,0 +1,575 @@ +mod common; + +use std::collections::BTreeMap; +use std::sync::Arc; +use std::time::Duration; + +use axum::body::Body; +use cid::Cid; +use common::{KNOT_HOST, OWNER, World, percent_encode, post_authed}; +use futures::StreamExt; +use http::{HeaderValue, StatusCode, header}; +use jacquard_repo::Mst; +use jacquard_repo::car::parse_car_bytes; +use jacquard_repo::storage::MemoryBlockStore; +use knot_git::Repo; +use tower::ServiceExt; + +const KNOT_DID: &str = "did:web:knot.nel.pet"; +const COLLECTION: &str = "sh.tangled.knot.memberInvite"; +const FIVE_SECONDS: Duration = Duration::from_secs(5); + +async fn get(world: &World, path_and_query: &str) -> (StatusCode, HeaderValue, Vec) { + let request = http::Request::builder() + .method("GET") + .uri(path_and_query) + .body(Body::empty()) + .unwrap(); + let response = world.router.clone().oneshot(request).await.unwrap(); + let status = response.status(); + let content_type = response + .headers() + .get(header::CONTENT_TYPE) + .cloned() + .unwrap_or_else(|| HeaderValue::from_static("")); + let body = axum::body::to_bytes(response.into_body(), usize::MAX) + .await + .unwrap() + .to_vec(); + (status, content_type, body) +} + +async fn get_json(world: &World, path_and_query: &str) -> serde_json::Value { + let (status, content_type, body) = get(world, path_and_query).await; + assert_eq!(status, StatusCode::OK, "{path_and_query}"); + assert!( + content_type + .to_str() + .unwrap() + .starts_with("application/json"), + "{path_and_query} serves json" + ); + serde_json::from_slice(&body).unwrap() +} + +async fn error_of(world: &World, path_and_query: &str) -> (StatusCode, serde_json::Value) { + let (status, _content_type, body) = get(world, path_and_query).await; + assert_eq!(status, StatusCode::BAD_REQUEST, "{path_and_query}"); + (status, serde_json::from_slice(&body).unwrap()) +} + +fn invited_world() -> World { + let world = World::new(); + world.invite_member("did:plc:limpet", 1_000); + world.invite_member("did:plc:whelk", 1_100); + world +} + +async fn cold_tree(car: &[u8]) -> (Cid, BTreeMap) { + let parsed = parse_car_bytes(car).await.unwrap(); + (parsed.root, parsed.blocks) +} + +async fn cold_loaded_tree(car: &[u8]) -> Mst { + let (head, blocks) = cold_tree(car).await; + let commit = jacquard_repo::commit::Commit::::from_cbor( + blocks.get(&head).expect("the commit block rides the car"), + ) + .unwrap(); + Mst::load( + Arc::new(MemoryBlockStore::new_from_blocks(blocks)), + *commit.data(), + None, + ) +} + +async fn root_data_of(world: &World) -> Cid { + let meta = Repo::open(&world.state.meta_path).unwrap(); + knot_record::chain::tip(&meta).unwrap().unwrap().data +} + +type Ws = + tokio_tungstenite::WebSocketStream>; +#[derive(serde::Deserialize)] +struct CommitProbe { + seq: u64, + rebase: bool, + repo: String, + rev: String, + ops: Vec, +} + +#[derive(serde::Deserialize)] +struct SyncProbe { + seq: u64, + did: String, + rev: String, + #[serde(with = "serde_bytes")] + blocks: Vec, +} + +fn frame_header(frame: &[u8]) -> (serde_json::Value, &[u8]) { + let mut read = frame; + let header: serde_json::Value = serde_ipld_dagcbor::de::from_reader_once(&mut read).unwrap(); + (header, read) +} + +fn frame_payload(rest: &[u8]) -> T { + serde_ipld_dagcbor::from_slice(rest).unwrap() +} + +async fn next_binary(ws: &mut Ws) -> Vec { + let received = tokio::time::timeout(FIVE_SECONDS, ws.next()) + .await + .expect("a frame arrives within the timeout") + .expect("the stream stays open") + .expect("the frame is readable"); + match received { + tokio_tungstenite::tungstenite::Message::Binary(bytes) => bytes.to_vec(), + other => panic!("the frame is binary, not {other:?}"), + } +} + +async fn drain_binaries(ws: &mut Ws) -> Vec> { + let mut seen = Vec::new(); + while let Ok(Some(Ok(frame))) = tokio::time::timeout(FIVE_SECONDS, ws.next()).await { + if let tokio_tungstenite::tungstenite::Message::Binary(bytes) = frame { + seen.push(bytes.to_vec()); + } + } + seen +} + +#[tokio::test] +async fn an_unmodified_cold_client_verifies_an_invite_off_the_proof_car() { + let world = invited_world(); + let key = format!("{COLLECTION}/did:plc:limpet"); + let (status, content_type, body) = get( + &world, + &format!( + "/xrpc/com.atproto.sync.getRecord?did={KNOT_DID}&collection={COLLECTION}&rkey=did%3Aplc%3Alimpet" + ), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(content_type, "application/vnd.ipld.car"); + let latest = get_json( + &world, + &format!("/xrpc/com.atproto.sync.getLatestCommit?did={KNOT_DID}"), + ) + .await; + let head: Cid = latest["cid"].as_str().unwrap().parse().unwrap(); + let (root, blocks) = cold_tree(&body).await; + assert_eq!(root, head, "the car roots at the signed commit"); + let commit = jacquard_repo::commit::Commit::::from_cbor( + blocks.get(&head).expect("the commit block rides the car"), + ) + .unwrap(); + assert_eq!( + commit.data(), + &root_data_of(&world).await, + "the commit's data cid names the tree the car proves" + ); + let found = cold_loaded_tree(&body) + .await + .get(&key) + .await + .unwrap() + .expect("the record is present"); + let record: serde_json::Value = + serde_ipld_dagcbor::from_slice(blocks.get(&found).unwrap()).unwrap(); + assert_eq!(record["createdAt"], "1970-01-01T00:16:40Z"); + assert_eq!(record["x-tngl-editor"], OWNER); + assert_eq!( + record["$type"], COLLECTION, + "the stored bytes carry their own $type, so a cold car reader can type the record" + ); + + let (_status, _type, absent) = get( + &world, + &format!( + "/xrpc/com.atproto.sync.getRecord?did={KNOT_DID}&collection={COLLECTION}&rkey=did%3Aplc%3Aabsent" + ), + ) + .await; + assert!( + cold_loaded_tree(&absent) + .await + .get(&format!("{COLLECTION}/did:plc:absent")) + .await + .unwrap() + .is_none(), + "non-inclusion is a proof, not an error" + ); + let (_, error) = error_of( + &world, + &format!( + "/xrpc/com.atproto.repo.getRecord?repo={KNOT_DID}&collection={COLLECTION}&rkey=did%3Aplc%3Aabsent" + ), + ) + .await; + assert_eq!(error["error"], "RecordNotFound"); + + let described = get_json( + &world, + &format!("/xrpc/com.atproto.repo.describeRepo?repo={KNOT_DID}"), + ) + .await; + assert_eq!( + described["handle"], KNOT_HOST, + "the knot's own handle is its hostname" + ); + assert!( + described["didDoc"].is_object(), + "the knot's own repo carries a did doc" + ); + assert_eq!(described["collections"], serde_json::json!([COLLECTION])); + + let described = get_json(&world, "/xrpc/com.atproto.server.describeServer").await; + assert_eq!(described["did"], KNOT_DID); + let resolved = get_json( + &world, + &format!("/xrpc/com.atproto.identity.resolveHandle?handle={KNOT_HOST}"), + ) + .await; + assert_eq!(resolved["did"], KNOT_DID); + let (_, error) = error_of( + &world, + "/xrpc/com.atproto.identity.resolveHandle?handle=stranger.example", + ) + .await; + assert_eq!(error["error"], "HandleNotFound"); +} + +#[tokio::test] +async fn get_repo_serves_the_whole_tree_and_get_blocks_serves_named_cids() { + let world = invited_world(); + let latest = get_json( + &world, + &format!("/xrpc/com.atproto.sync.getLatestCommit?did={KNOT_DID}"), + ) + .await; + let head: Cid = latest["cid"].as_str().unwrap().parse().unwrap(); + + let (status, _content_type, body) = get( + &world, + &format!("/xrpc/com.atproto.sync.getRepo?did={KNOT_DID}"), + ) + .await; + assert_eq!(status, StatusCode::OK); + let (root, _blocks) = cold_tree(&body).await; + assert_eq!(root, head); + assert_eq!( + cold_loaded_tree(&body).await.leaves().await.unwrap().len(), + 2, + "the whole roster exports" + ); + for wanted in [head, "bafkqaaa".parse::().unwrap()] { + let (status, _content_type, body) = get( + &world, + &format!("/xrpc/com.atproto.sync.getBlocks?did={KNOT_DID}&cids={wanted}"), + ) + .await; + assert_eq!(status, StatusCode::OK); + let (root, blocks) = cold_tree(&body).await; + assert_eq!(root, head, "the car still roots at the commit"); + assert_eq!( + blocks.contains_key(&wanted), + wanted == head, + "a named block is served and a cid too short to shard reads as absent, not a panic" + ); + } +} + +#[tokio::test] +async fn a_removal_takes_the_record_off_the_surface() { + let world = invited_world(); + world.remove_member("did:plc:limpet", 1_200); + let listed = get_json( + &world, + &format!("/xrpc/com.atproto.repo.listRecords?repo={KNOT_DID}&collection={COLLECTION}"), + ) + .await; + assert_eq!( + listed["records"][0]["uri"], + format!("at://{KNOT_DID}/{COLLECTION}/did:plc:whelk"), + "the removed invite is the only row left" + ); +} + +#[tokio::test] +async fn hosted_repositories_read_over_com_atproto_until_a_chain_exists() { + let world = World::new(); + let (repo, _name, _kept) = common::empty_repo(&world, "squid"); + world.invite_collaborator(&repo, "did:plc:limpet", 1_000); + let listed = get_json( + &world, + &format!( + "/xrpc/com.atproto.repo.listRecords?repo={}&collection=sh.tangled.repo.collaboratorInvite", + percent_encode(repo.as_str()) + ), + ) + .await; + assert_eq!( + listed["records"][0]["uri"], + format!( + "at://{}/sh.tangled.repo.collaboratorInvite/did:plc:limpet", + repo.as_str() + ), + "a repo with a chain lists its own collaborator invites" + ); + let repos = get_json(&world, "/xrpc/com.atproto.sync.listRepos").await; + assert!( + repos["repos"] + .as_array() + .unwrap() + .iter() + .any(|row| row["did"] == repo.as_str()), + "a repo with a chain is listed" + ); + + let (lonely, _name, _kept) = common::empty_repo(&world, "lonely"); + let did = percent_encode(lonely.as_str()); + let (_, error) = error_of( + &world, + &format!("/xrpc/com.atproto.sync.getLatestCommit?did={did}"), + ) + .await; + assert_eq!(error["error"], "RepoNotFound"); + let repos = get_json(&world, "/xrpc/com.atproto.sync.listRepos").await; + assert!( + !repos["repos"] + .as_array() + .unwrap() + .iter() + .any(|row| row["did"] == lonely.as_str()), + "a registered repo with no chain has no head to name" + ); + let described = get_json( + &world, + &format!("/xrpc/com.atproto.sync.getRepoStatus?did={did}"), + ) + .await; + assert_eq!(described["active"], true); + assert!( + described.get("rev").is_none(), + "a repo with no chain carries no rev rather than a null one" + ); +} + +#[tokio::test] +async fn an_evicted_cursor_answers_an_out_of_cursor_info_frame() { + let world = invited_world(); + let addr = common::serve(&world).await; + let (mut ws, _response) = tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos?cursor=0" + )) + .await + .unwrap(); + let first = next_binary(&mut ws).await; + let (header, rest) = frame_header(&first); + let payload: CommitProbe = frame_payload(rest); + assert_eq!(header["op"], 1); + assert_eq!( + header["t"], "#commit", + "a zero cursor against a window that lost nothing replays without an info frame" + ); + assert_eq!(payload.repo, KNOT_DID); + assert!(!payload.rebase); + assert!(payload.seq > 0); + let second = next_binary(&mut ws).await; + let (header, _rest) = frame_header(&second); + assert_eq!(header["t"], "#commit", "both invites reach the stream"); + + (0..10).for_each(|n| world.invite_member(&format!("did:plc:subject{n}"), 2_000 + n)); + let (mut ws, _response) = tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos?cursor=0" + )) + .await + .unwrap(); + let notified = next_binary(&mut ws).await; + let (header, rest) = frame_header(¬ified); + assert_eq!(header["t"], "#info"); + let info: serde_json::Value = frame_payload(rest); + assert_eq!(info["name"], "AgedNonCommitCursor"); + let seqs = drain_binaries(&mut ws) + .await + .iter() + .map(|bytes| frame_payload::(frame_header(bytes).1).seq) + .collect::>(); + assert_eq!( + seqs.len(), + 12, + "every commit the ring evicted still replays off the chains" + ); + assert!( + seqs.windows(2).all(|pair| pair[0] < pair[1]), + "the replay arrives in seq order" + ); +} + +#[tokio::test] +async fn account_creation_announces_identity_account_and_an_empty_commit() { + let world = World::administered(); + world.add_member(OWNER, OWNER, 1_000); + let addr = common::serve(&world).await; + let (mut ws, _response) = tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos" + )) + .await + .unwrap(); + let (status, body) = post_authed( + &world, + "/xrpc/sh.tangled.repo.create", + OWNER, + serde_json::json!({ "rkey": "limpetkey", "name": "limpet" }), + ) + .await; + assert_eq!(status, StatusCode::OK, "create failed: {body}"); + let repo_did = body["repoDid"].as_str().unwrap().to_owned(); + let identity = next_binary(&mut ws).await; + let (header, rest) = frame_header(&identity); + assert_eq!(header["t"], "#identity"); + let payload: serde_json::Value = frame_payload(rest); + assert_eq!(payload["did"], repo_did); + let identity_seq = payload["seq"].as_u64().unwrap(); + + let account = next_binary(&mut ws).await; + let (header, rest) = frame_header(&account); + assert_eq!(header["t"], "#account"); + let payload: serde_json::Value = frame_payload(rest); + assert_eq!(payload["did"], repo_did); + assert_eq!(payload["active"], true); + let account_seq = payload["seq"].as_u64().unwrap(); + + let commit = next_binary(&mut ws).await; + let (header, rest) = frame_header(&commit); + assert_eq!(header["t"], "#commit"); + let payload: CommitProbe = frame_payload(rest); + assert_eq!(payload.repo, repo_did); + assert_eq!( + payload.ops.len(), + 0, + "a fresh account's first commit is empty" + ); + assert!( + identity_seq < account_seq && account_seq < payload.seq, + "creation goes identity, account, commit" + ); + let genesis_rev = payload.rev.clone(); + let commit_seq = payload.seq; + let sync = next_binary(&mut ws).await; + let (header, rest) = frame_header(&sync); + assert_eq!(header["t"], "#sync"); + let payload: SyncProbe = frame_payload(rest); + assert!( + commit_seq < payload.seq, + "the sync assertion follows the commit it names" + ); + assert_eq!(payload.did, repo_did); + assert_eq!( + payload.rev, genesis_rev, + "the sync frame asserts the genesis commit just announced" + ); + assert!( + !payload.blocks.is_empty(), + "the sync frame carries the genesis commit in its CAR" + ); + + let (status, body) = post_authed( + &world, + "/xrpc/sh.tangled.repo.delete", + OWNER, + serde_json::json!({ "repo": repo_did, "force": true }), + ) + .await; + assert_eq!(status, StatusCode::OK, "delete failed: {body}"); + let closed = drain_binaries(&mut ws) + .await + .iter() + .filter(|bytes| frame_header(bytes).0["t"] == "#account") + .map(|bytes| frame_payload::(frame_header(bytes).1)) + .next_back() + .expect("an account frame closed the lifecycle"); + assert_eq!(closed["did"], repo_did); + assert_eq!(closed["active"], false); + assert_eq!(closed["status"], "deleted"); +} +#[tokio::test] +async fn a_non_integer_cursor_is_rejected_in_the_xrpc_error_shape() { + let addr = common::serve(&invited_world()).await; + let error = tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos?cursor=soon" + )) + .await + .unwrap_err(); + let response = match error { + tokio_tungstenite::tungstenite::Error::Http(response) => response, + other => panic!("the handshake answers http, not {other:?}"), + }; + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + let body = response.into_body().unwrap_or_default(); + let error: serde_json::Value = serde_json::from_slice(&body).unwrap(); + assert_eq!(error["error"], "InvalidRequest"); + assert_eq!(error["message"], "cursor must be an integer"); +} + +#[tokio::test] +async fn a_future_cursor_answers_an_error_frame_and_a_try_again_close() { + let addr = common::serve(&invited_world()).await; + let (mut ws, _response) = tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos?cursor=9000000000000000000" + )) + .await + .unwrap(); + let refused = next_binary(&mut ws).await; + let (header, rest) = frame_header(&refused); + let refusal: serde_json::Value = frame_payload(rest); + assert_eq!(header["op"], -1); + assert!( + header.get("t").is_none(), + "an error frame carries no type field" + ); + assert_eq!(refusal["error"], "FutureCursor"); + let closed = tokio::time::timeout(FIVE_SECONDS, ws.next()) + .await + .expect("the close arrives within the timeout") + .expect("the stream ends") + .expect("the close is readable"); + match closed { + tokio_tungstenite::tungstenite::Message::Close(Some(frame)) => assert_eq!( + u16::from(frame.code), + 1013, + "the client is told to try again later" + ), + other => panic!("the stream closes with a code, not {other:?}"), + } +} + +async fn refusal(world: &World, query: String) { + let (_, error) = error_of(world, &query).await; + assert_eq!(error["error"], "InvalidRequest", "{query}"); +} + +#[tokio::test] +async fn malformed_collections_rkeys_and_cid_lists_answer_invalid_request() { + let world = invited_world(); + [ + format!( + "/xrpc/com.atproto.sync.getRecord?did={KNOT_DID}&collection=not%20an%20nsid&rkey=did%3Aplc%3Alimpet" + ), + format!( + "/xrpc/com.atproto.sync.getRecord?did={KNOT_DID}&collection={COLLECTION}&rkey=has%20spaces" + ), + format!( + "/xrpc/com.atproto.repo.getRecord?repo={KNOT_DID}&collection=all+wrong&rkey=did%3Aplc%3Alimpet" + ), + format!("/xrpc/com.atproto.repo.listRecords?repo={KNOT_DID}&collection=also+wrong"), + format!( + "/xrpc/com.atproto.sync.getBlocks?did={KNOT_DID}&{}", + (0..1001).map(|n| format!("cids=bafkqaa{n}")).collect::>().join("&") + ), + ] + .into_iter() + .for_each(|query| futures::executor::block_on(refusal(&world, query))); +} diff --git a/knot2/crates/knot-xrpc/tests/atproto_interop.rs b/knot2/crates/knot-xrpc/tests/atproto_interop.rs new file mode 100644 index 000000000..1482bb7a0 --- /dev/null +++ b/knot2/crates/knot-xrpc/tests/atproto_interop.rs @@ -0,0 +1,122 @@ +mod common; + +use common::{OWNER, World, commit_file, empty_repo, http_push, serve_pushes, sh_git}; + +const KNOT_DID: &str = "did:web:knot.nel.pet"; +const COLLECTION: &str = "sh.tangled.knot.memberInvite"; +const RKEY: &str = "did:plc:limpet"; + +#[tokio::test] +async fn an_indigo_client_reads_and_verifies_the_knot() { + let world = World::new(); + world.invite_member(RKEY, 1_000); + world.invite_member("did:plc:whelk", 1_100); + + let signer = world.state.secrets.signer(&world.state.knot_did).unwrap(); + let doc = knot_atproto::knot_did_document( + &world.state.knot_did, + &knot_runtime::Signer::public_key(&signer), + &world.state.knot_service_url, + ); + let addr = common::serve(&world).await; + + let repo_root = env!("CARGO_MANIFEST_DIR").to_owned() + "/../../.."; + let go = tokio::process::Command::new("go") + .arg("test") + .arg("./knot2/interop/") + .arg("-run") + .arg("TestKnotServesComAtprotoReads") + .arg("-count=1") + .arg("-v") + .current_dir(&repo_root) + .env("KNOT_TEST_BASE_URL", format!("http://{addr}")) + .env("KNOT_TEST_DID", KNOT_DID) + .env("KNOT_TEST_HANDLE", common::KNOT_HOST) + .env("KNOT_TEST_DID_DOC", doc.to_string()) + .env("KNOT_TEST_COLLECTION", COLLECTION) + .env("KNOT_TEST_RKEY", RKEY) + .output() + .await + .expect("go test launches"); + assert!( + go.status.success(), + "the unmodified indigo client failed:\n{}{}", + String::from_utf8_lossy(&go.stdout), + String::from_utf8_lossy(&go.stderr), + ); +} + +#[tokio::test] +async fn an_indigo_client_consumes_the_firehose() { + let world = World::administered(); + world.add_member(common::OWNER, common::OWNER, 1_000); + let addr = common::serve(&world).await; + + let repo_root = env!("CARGO_MANIFEST_DIR").to_owned() + "/../../.."; + let go = tokio::process::Command::new("go") + .arg("test") + .arg("./knot2/interop/") + .arg("-run") + .arg("TestKnotServesSubscribeRepos") + .arg("-count=1") + .arg("-v") + .current_dir(&repo_root) + .env("KNOT_TEST_BASE_URL", format!("http://{addr}")) + .env( + "KNOT_TEST_CREATE_TOKEN", + common::service_jwt("sh.tangled.repo.create", common::OWNER), + ) + .env( + "KNOT_TEST_DELETE_TOKEN", + common::service_jwt("sh.tangled.repo.delete", common::OWNER), + ) + .output() + .await + .expect("go test launches"); + assert!( + go.status.success(), + "the unmodified indigo firehose client failed:\n{}{}", + String::from_utf8_lossy(&go.stdout), + String::from_utf8_lossy(&go.stderr), + ); +} + +#[tokio::test] +async fn a_knotfeed_consumer_reads_the_ref_records() { + let world = World::new(); + let (did, _bare, work) = empty_repo(&world, "conch"); + commit_file( + work.path(), + "README.md", + b"hello\n", + "first", + "2026-06-01T12:30:00+02:00", + ); + let addr = serve_pushes(&world).await; + http_push(work.path(), addr, &did, OWNER, &["-o", "skip-ci", "main"]).await; + let sha = sh_git(work.path(), &["rev-parse", "HEAD"]); + let sha = sha.trim().to_string(); + + let repo_root = env!("CARGO_MANIFEST_DIR").to_owned() + "/../../.."; + let go = tokio::process::Command::new("go") + .arg("test") + .arg("./knot2/interop/") + .arg("-run") + .arg("TestKnotfeedConsumerReadsRefRecords") + .arg("-count=1") + .arg("-v") + .current_dir(&repo_root) + .env("KNOT_TEST_ADDR", format!("127.0.0.1:{}", addr.port())) + .env("KNOT_TEST_REPO_DID", did.as_str()) + .env("KNOT_TEST_EDITOR", OWNER) + .env("KNOT_TEST_SHA", &sha) + .output() + .await + .expect("go test launches"); + assert!( + go.status.success(), + "the knotfeed consumer failed:\n{}{}", + String::from_utf8_lossy(&go.stdout), + String::from_utf8_lossy(&go.stderr), + ); +} diff --git a/knot2/crates/knot-xrpc/tests/common/mod.rs b/knot2/crates/knot-xrpc/tests/common/mod.rs index c9234ecdd..3433be094 100644 --- a/knot2/crates/knot-xrpc/tests/common/mod.rs +++ b/knot2/crates/knot-xrpc/tests/common/mod.rs @@ -1,5 +1,112 @@ #![allow(dead_code)] +pub async fn serve(world: &World) -> std::net::SocketAddr { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let router = world.router.clone(); + tokio::spawn(async move { + axum::serve( + listener, + router.into_make_service_with_connect_info::(), + ) + .await + .unwrap(); + }); + addr +} + +pub async fn serve_pushes(world: &World) -> std::net::SocketAddr { + let resolver: Arc = { + let index = Arc::clone(&world.state.index); + Arc::new(move |target: &knot_pack::RepoTarget| match target { + knot_pack::RepoTarget::Did(did) => { + knot_pack::RepoLookup::from_resolved(index.owner_of(did), |_| did.clone()) + } + knot_pack::RepoTarget::OwnerPath(owner, path) => knot_pack::RepoLookup::from_resolved( + index.resolve_clone_path(owner, path), + |found| found, + ), + }) + }; + let (_, advertisement) = knot_pack::edge_routes(knot_pack::EdgeConfig { + receive: Some(knot_xrpc::receive_advertiser(Arc::clone(&world.state))), + ..knot_pack::EdgeConfig::serving( + world.layout.clone(), + resolver, + Arc::new(knot_runtime::SystemClock), + ) + }); + let router = advertisement.into_router().merge(world.router.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + axum::serve( + listener, + router.into_make_service_with_connect_info::(), + ) + .await + .unwrap(); + }); + addr +} + +pub async fn http_push( + work: &Path, + addr: std::net::SocketAddr, + did: &knot_types::RepoDid, + actor: &str, + refspecs: &[&str], +) -> (bool, String) { + let header = format!( + "http.extraHeader=Authorization: Basic {}", + base64::engine::general_purpose::STANDARD.encode(format!( + "x-tangled-token:{}", + service_jwt("sh.tangled.repo.push", actor) + )) + ); + let url = format!("http://{addr}/{did}"); + let mut args = vec![ + "-c".to_string(), + header, + "push".to_string(), + "-q".to_string(), + ]; + args.push(url); + args.extend(refspecs.iter().map(|spec| spec.to_string())); + let work = work.to_path_buf(); + tokio::task::spawn_blocking(move || { + let output = std::process::Command::new("git") + .args(&args) + .env("GIT_TERMINAL_PROMPT", "0") + .current_dir(&work) + .output() + .expect("git is available"); + ( + output.status.success(), + format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ), + ) + }) + .await + .unwrap() +} + +pub fn percent_encode(value: &str) -> String { + value + .as_bytes() + .iter() + .fold(String::with_capacity(value.len()), |mut out, byte| { + match byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'.' | b'_' | b'~') { + true => out.push(*byte as char), + false => out.push_str(&format!("%{byte:02X}")), + } + out + }) +} + use std::collections::BTreeSet; use std::os::unix::fs::PermissionsExt; use std::path::Path; @@ -95,6 +202,16 @@ impl World { ) } + pub fn administered() -> Self { + Self::assemble( + true, + ByteLimits::default(), + ObjectFormat::SHA1, + LimitConfig::default(), + BTreeSet::from([AccountDid::new(OWNER).unwrap()]), + ) + } + fn build(rebuilt: bool, byte_limits: ByteLimits, object_format: ObjectFormat) -> Self { Self::build_with_limits(rebuilt, byte_limits, object_format, LimitConfig::default()) } @@ -104,6 +221,16 @@ impl World { byte_limits: ByteLimits, object_format: ObjectFormat, limits: LimitConfig, + ) -> Self { + Self::assemble(rebuilt, byte_limits, object_format, limits, BTreeSet::new()) + } + + fn assemble( + rebuilt: bool, + byte_limits: ByteLimits, + object_format: ObjectFormat, + limits: LimitConfig, + admins: BTreeSet, ) -> Self { let dir = tempfile::tempdir().unwrap(); let scan_path = dir.path().join("repos"); @@ -172,7 +299,7 @@ impl World { secrets, entropy: Arc::new(OsEntropy), ci_logs: None, - admins: BTreeSet::new(), + admins, admission: knot_types::AdmissionPolicy::Closed, knot_did: knot, knot_hostname: KnotHostname::new(KNOT_HOST).unwrap(), @@ -200,13 +327,6 @@ impl World { })), pack_limits: knot_pack::PackLimits::default(), service_owner: AccountDid::new(OWNER).unwrap(), - events: Arc::new(knot_events::EventLog::new( - ManualClock::new(UnixMicros::new(1_000_000_000)), - knot_events::ReplayBounds::new( - knot_events::ReplayEvents::new(1024).unwrap(), - knot_events::ReplayBytes::new(16 << 20).unwrap(), - ), - )), subscriber_gate: Arc::new(knot_events::SubscriberGate::new( knot_events::GlobalSubscriberLimit::new(16), knot_events::PerPeerSubscriberLimit::new(4), @@ -216,6 +336,13 @@ impl World { slots: knot_resource::Slots::testing(8), lfs: Some(knot_xrpc::LfsWeb::new(lfs_handle, 8)), catalog: Arc::new(knot_messages::Catalog::defaults()), + firehose: Arc::new(knot_events::EventLog::new( + ManualClock::new(UnixMicros::new(1_000_000_000)), + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(8).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), + )), }); let router = knot_xrpc::router(Arc::clone(&state)); Self { @@ -306,15 +433,76 @@ impl World { fn member(&self, change: MembersChange, at: i64) { self.on_meta::(change, at); + self.upkeep_members(); self.state.index.refresh_members().unwrap(); } fn collaborator(&self, repo: &RepoDid, change: CollaboratorsChange, at: i64) { + self.append::( + self.layout.open(repo).unwrap(), + CobHome::from(repo), + change, + at, + ); let git = self.layout.open(repo).unwrap(); - self.append::(git, CobHome::from(repo), change, at); + self.upkeep_chain::( + &git, + repo.as_str(), + knot_xrpc::materialize::COLLABORATOR_INVITE_COLLECTION, + ); self.state.index.refresh_collaborators(repo).unwrap(); } + fn upkeep_members(&self) { + let meta = Repo::open(&self.state.meta_path).unwrap(); + self.upkeep_chain::( + &meta, + self.state.knot_did.as_str(), + knot_xrpc::materialize::MEMBER_INVITE_COLLECTION, + ); + } + + fn upkeep_chain(&self, git: &Repo, did: &str, collection: &'static str) + where + E: knot_cob::Checkpoint + knot_cob::Evaluate, + E::Change: ChangePayload + knot_cobs::GrantChange + Clone, + { + let store = CobStore::new(git); + if let [object] = store.list::().unwrap().as_slice() { + let path = git.path().to_path_buf(); + std::thread::scope(|scope| { + scope + .spawn(move || { + let git = Repo::open(&path).unwrap(); + let store = CobStore::new(&git); + let signer = self.state.secrets.signer(&self.state.knot_did).unwrap(); + let mut emit = knot_xrpc::FirehoseEmit::new( + &self.state.firehose, + did, + UnixSeconds::new(1_000), + ); + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(knot_xrpc::materialize::upkeep::( + &store, + *object, + did, + collection, + &signer, + &mut emit, + UnixSeconds::new(1_000), + )) + .unwrap(); + emit.finish(); + }) + .join() + .unwrap(); + }); + } + } + fn on_meta(&self, change: E::Change, at: i64) { let meta = Repo::open(&self.state.meta_path).unwrap(); self.append::(meta, CobHome::from(&self.state.knot_did), change, at); @@ -372,7 +560,7 @@ fn did_doc_for(did: &str) -> Bytes { static JTI: AtomicU64 = AtomicU64::new(0); -fn service_jwt(nsid: &str, actor: &str) -> String { +pub fn service_jwt(nsid: &str, actor: &str) -> String { let nonce = JTI.fetch_add(1, Ordering::SeqCst); let header = URL_SAFE_NO_PAD.encode(br#"{"alg":"ES256K","typ":"JWT"}"#); let claims = serde_json::json!({ diff --git a/knot2/crates/knot-xrpc/tests/reads.rs b/knot2/crates/knot-xrpc/tests/reads.rs index 1216f1f6c..45e5595b2 100644 --- a/knot2/crates/knot-xrpc/tests/reads.rs +++ b/knot2/crates/knot-xrpc/tests/reads.rs @@ -1,7 +1,5 @@ mod common; -use std::future::Future; -use std::pin::Pin; use std::time::{Duration, Instant}; use futures::StreamExt; @@ -9,9 +7,8 @@ use futures::stream; use http::{HeaderMap, StatusCode, header}; use tokio_tungstenite::tungstenite; -use knot_events::{EventCursor, GitRefUpdate}; use knot_index::{Resolved, Standing}; -use knot_types::{AccountDid, Oid, OwnerDid, RepoDid}; +use knot_types::{AccountDid, Oid, RepoDid}; use knot_xrpc::{ArchiveLimit, ResponseLimit}; use common::{ @@ -1871,7 +1868,10 @@ async fn either_roster_pages_its_outstanding_offers_and_drops_a_withdrawn_or_ban "the blocklist filters whelk out of the listing, since a banned account can't accept" ); assert_eq!( - world.state.index.member_standing(&AccountDid::new("did:plc:whelk").unwrap()), + world + .state + .index + .member_standing(&AccountDid::new("did:plc:whelk").unwrap()), Resolved::Ready(Some(Standing::Invited)), "the blocklist is a roster of its own and neither change rewrites whelk's entry on the \ members roster" @@ -2022,15 +2022,7 @@ async fn service_metadata_endpoints_answer() { assert_eq!(owner["owner"], OWNER); } -fn publish_update(world: &World, repo: &str) -> EventCursor { - world.state.events.publish(&GitRefUpdate::new( - RepoDid::new(repo).unwrap(), - Some(OwnerDid::new(OWNER).unwrap()), - AccountDid::new("did:plc:nel").unwrap(), - )) -} - -async fn serve_events(world: &World) -> std::net::SocketAddr { +async fn serve_router(world: &World) -> std::net::SocketAddr { let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); let router = world.router.clone(); @@ -2045,70 +2037,6 @@ async fn serve_events(world: &World) -> std::net::SocketAddr { addr } -type Ws = - tokio_tungstenite::WebSocketStream>; - -fn next_event(ws: &mut Ws) -> Pin + '_>> { - Box::pin(async move { - let received = tokio::time::timeout(std::time::Duration::from_secs(5), ws.next()) - .await - .expect("an event arrives within the timeout") - .expect("the stream stays open") - .expect("the frame is readable"); - match received { - tungstenite::Message::Text(text) => serde_json::from_str(text.as_str()).unwrap(), - _ => next_event(ws).await, - } - }) -} - -#[tokio::test] -async fn the_events_stream_replays_resumes_and_rejects_bad_cursors() { - let world = World::new(); - let first = publish_update(&world, "did:plc:squid"); - publish_update(&world, "did:plc:anemone"); - let addr = serve_events(&world).await; - - let (mut ws, _) = tokio_tungstenite::connect_async(format!("ws://{addr}/events")) - .await - .unwrap(); - let replayed_first = next_event(&mut ws).await; - let replayed_second = next_event(&mut ws).await; - assert_eq!(replayed_first["nsid"], "sh.tangled.git.refUpdate"); - assert_eq!(replayed_first["event"]["repo"], "did:plc:squid"); - assert_eq!(replayed_first["event"]["ownerDid"], OWNER); - assert_eq!(replayed_first["event"]["committerDid"], "did:plc:nel"); - assert_eq!(replayed_first["rkey"].as_str().unwrap().len(), 13); - assert_eq!(replayed_second["event"]["repo"], "did:plc:anemone"); - assert!( - replayed_first["created"].as_i64().unwrap() < replayed_second["created"].as_i64().unwrap() - ); - - publish_update(&world, "did:plc:whelk"); - let live = next_event(&mut ws).await; - assert_eq!(live["event"]["repo"], "did:plc:whelk"); - - let (mut resumed, _) = - tokio_tungstenite::connect_async(format!("ws://{addr}/events?cursor={}", first.get())) - .await - .unwrap(); - let resumed_event = next_event(&mut resumed).await; - assert_eq!( - resumed_event["event"]["repo"], "did:plc:anemone", - "a cursor resumes past the event it names" - ); - - let (mut garbled, _) = - tokio_tungstenite::connect_async(format!("ws://{addr}/events?cursor=banana")) - .await - .unwrap(); - let replayed = next_event(&mut garbled).await; - assert_eq!( - replayed["event"]["repo"], "did:plc:squid", - "a garbled cursor replays from the start" - ); -} - fn refused(result: Result) { match result { Err(tungstenite::Error::Http(response)) => { @@ -2120,9 +2048,9 @@ fn refused(result: Result) { } #[tokio::test] -async fn a_subscriber_beyond_the_events_limit_is_refused() { +async fn a_subscriber_beyond_the_firehose_limit_is_refused() { let world = World::new(); - let addr = serve_events(&world).await; + let addr = serve_router(&world).await; let saturated: Vec<_> = (0..16u8) .map(|octet| { world @@ -2134,24 +2062,36 @@ async fn a_subscriber_beyond_the_events_limit_is_refused() { .expect("distinct peers fill the global limit") }) .collect(); - refused(tokio_tungstenite::connect_async(format!("ws://{addr}/events")).await); + refused( + tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos" + )) + .await, + ); drop(saturated); } #[tokio::test] -async fn a_single_peer_cannot_monopolize_the_events_stream() { +async fn one_peer_opens_four_subscriptions_and_fifth_is_refused() { let world = World::new(); - let addr = serve_events(&world).await; + let addr = serve_router(&world).await; let held: Vec<_> = stream::iter(0..4) .then(|_| async { - tokio_tungstenite::connect_async(format!("ws://{addr}/events")) - .await - .expect("a connection within the per-peer limit is admitted") - .0 + tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos" + )) + .await + .expect("a connection within the per-peer limit is admitted") + .0 }) .collect() .await; - refused(tokio_tungstenite::connect_async(format!("ws://{addr}/events")).await); + refused( + tokio_tungstenite::connect_async(format!( + "ws://{addr}/xrpc/com.atproto.sync.subscribeRepos" + )) + .await, + ); drop(held); }