From fd10f8f9f03a75674f330a4fcac4e45c39402652 Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 18 Aug 2026 11:53:31 +0300 Subject: [PATCH] knot2/crates/knot-xrpc: serve invite and acceptance endpoints Lewis: May this revision serve well! --- knot2/crates/knot-config/src/lib.rs | 28 +- knot2/crates/knot-server/src/main.rs | 20 +- knot2/crates/knot-sim/Cargo.toml | 1 + knot2/crates/knot-sim/src/harness.rs | 139 +++- knot2/crates/knot-sim/src/lib.rs | 7 + knot2/crates/knot-sim/src/trace.rs | 1 + knot2/crates/knot-sim/src/workload.rs | 390 +++++++--- knot2/crates/knot-sim/tests/reproducible.rs | 93 ++- knot2/crates/knot-ssh/Cargo.toml | 1 + knot2/crates/knot-ssh/src/exec.rs | 19 +- knot2/crates/knot-ssh/src/identity.rs | 4 + knot2/crates/knot-ssh/src/lib.rs | 7 +- knot2/crates/knot-ssh/src/server.rs | 9 + knot2/crates/knot-ssh/tests/ssh_push.rs | 207 +++++- knot2/crates/knot-xrpc/Cargo.toml | 1 + knot2/crates/knot-xrpc/src/blocklist.rs | 19 +- knot2/crates/knot-xrpc/src/cob.rs | 172 ++++- knot2/crates/knot-xrpc/src/collaborators.rs | 171 +++-- knot2/crates/knot-xrpc/src/forks.rs | 6 +- knot2/crates/knot-xrpc/src/legacy_admin.rs | 6 +- knot2/crates/knot-xrpc/src/lib.rs | 79 +- knot2/crates/knot-xrpc/src/lists.rs | 148 ++-- knot2/crates/knot-xrpc/src/members.rs | 117 +-- knot2/crates/knot-xrpc/src/reads.rs | 11 +- knot2/crates/knot-xrpc/src/repos.rs | 27 +- knot2/crates/knot-xrpc/src/tests.rs | 773 ++++++++++++++++---- knot2/crates/knot-xrpc/tests/common/mod.rs | 106 ++- knot2/crates/knot-xrpc/tests/reads.rs | 185 ++++- knot2/example.toml | 20 +- 29 files changed, 2202 insertions(+), 565 deletions(-) diff --git a/knot2/crates/knot-config/src/lib.rs b/knot2/crates/knot-config/src/lib.rs index 2c76fb7c5..1a972a371 100644 --- a/knot2/crates/knot-config/src/lib.rs +++ b/knot2/crates/knot-config/src/lib.rs @@ -60,6 +60,12 @@ pub struct AclConfig { #[config(env = "KNOT_LEGACY_ADMIN_SECRET_ENV")] pub legacy_admin_secret_env: Option, + + #[config(env = "KNOT_ACL_ACCEPTANCE_TTL_SECS", default = 60)] + pub acceptance_ttl_secs: u64, + + #[config(env = "KNOT_ACL_ACCEPTANCE_GRACE_SECS", default = 21_600)] + pub acceptance_grace_secs: u64, } #[derive(Debug, Config)] @@ -774,20 +780,32 @@ impl KnotConfig { self.xrpc.events_max_subscribers >= self.xrpc.events_max_per_peer, "xrpc.events_max_subscribers must be at least xrpc.events_max_per_peer", ), + check( + (1..=MAX_SPAN_SECS).contains(&self.acl.acceptance_ttl_secs), + "acl.acceptance_ttl_secs must be between one second and one year", + ), + check( + (1..=MAX_SPAN_SECS).contains(&self.acl.acceptance_grace_secs), + "acl.acceptance_grace_secs must be between one second and one year", + ), + check( + self.acl.acceptance_grace_secs >= self.acl.acceptance_ttl_secs, + "acl.acceptance_grace_secs must be at least acl.acceptance_ttl_secs", + ), check( (1..=MAX_KEYFILL_BUDGET_MIB).contains(&self.keyfill.key_budget_mib), "keyfill.key_budget_mib must be between one mebibyte and one tebibyte", ), check( - (1..=MAX_KEYFILL_SPAN_SECS).contains(&self.keyfill.ttl_secs), + (1..=MAX_SPAN_SECS).contains(&self.keyfill.ttl_secs), "keyfill.ttl_secs must be between one second and one year", ), check( - (1..=MAX_KEYFILL_SPAN_SECS).contains(&self.keyfill.reprieve_retry_secs), + (1..=MAX_SPAN_SECS).contains(&self.keyfill.reprieve_retry_secs), "keyfill.reprieve_retry_secs must be between one second and one year", ), check( - (1..=MAX_KEYFILL_SPAN_SECS).contains(&self.keyfill.reprieve_budget_secs), + (1..=MAX_SPAN_SECS).contains(&self.keyfill.reprieve_budget_secs), "keyfill.reprieve_budget_secs must be between one second and one year", ), check( @@ -1056,7 +1074,7 @@ pub enum EnvError { const MASTER_KEY_MIN_BYTES: usize = 32; -const MAX_KEYFILL_SPAN_SECS: u64 = 365 * 24 * 60 * 60; +const MAX_SPAN_SECS: u64 = 365 * 24 * 60 * 60; const MAX_KEYFILL_BUDGET_MIB: u32 = 1024 * 1024; @@ -1228,6 +1246,8 @@ mod tests { acl: AclConfig { admission: AdmissionPolicy::Closed, legacy_admin_secret_env: None, + acceptance_ttl_secs: 60, + acceptance_grace_secs: 21_600, }, repo: RepoConfig { scan_path: PathBuf::from("/srv/git"), diff --git a/knot2/crates/knot-server/src/main.rs b/knot2/crates/knot-server/src/main.rs index 3d7e82147..2f620b38d 100644 --- a/knot2/crates/knot-server/src/main.rs +++ b/knot2/crates/knot-server/src/main.rs @@ -224,12 +224,22 @@ async fn main() -> anyhow::Result<()> { .meta_path(&knot_did) .context("resolve meta-repo path")?; + knot_types::DidRkey::new(knot_did.as_str()).with_context(|| { + format!( + "{} can't be the record key of an acceptance record, so no member could ever accept \ + this knot: a did:web with a port or a path is a development form", + knot_did.as_str() + ) + })?; + let key_pace = keyfill_pace(&config.keyfill); - let index = Arc::new(Index::with_key_budget( - meta_path.clone(), - layout.clone(), - knot_index::KeyBudget::from_mib(config.keyfill.key_budget_mib as usize), - )); + let budget = knot_index::KeyBudget::from_mib(config.keyfill.key_budget_mib as usize); + let freshness = knot_index::Freshness::new( + knot_index::AcceptanceTtl::from_secs(config.acl.acceptance_ttl_secs), + knot_index::AcceptanceGrace::from_secs(config.acl.acceptance_grace_secs), + ); + let fresh = Index::with_key_budget(meta_path.clone(), layout.clone(), budget); + let index = Arc::new(fresh.with_freshness(freshness)); index.rebuild().context("rebuild index from meta-repo")?; tracing::info!(coverage = ?index.coverage(), "index ready"); diff --git a/knot2/crates/knot-sim/Cargo.toml b/knot2/crates/knot-sim/Cargo.toml index 0f22b57b2..b8153da02 100644 --- a/knot2/crates/knot-sim/Cargo.toml +++ b/knot2/crates/knot-sim/Cargo.toml @@ -19,6 +19,7 @@ knot-index = { workspace = true } knot-pack = { workspace = true } knot-resource = { workspace = true } knot-atproto = { workspace = true } +knot-consent = { workspace = true } knot-secrets = { workspace = true } knot-xrpc = { workspace = true } knot-events = { workspace = true } diff --git a/knot2/crates/knot-sim/src/harness.rs b/knot2/crates/knot-sim/src/harness.rs index bf727e9db..1b783e076 100644 --- a/knot2/crates/knot-sim/src/harness.rs +++ b/knot2/crates/knot-sim/src/harness.rs @@ -14,7 +14,7 @@ use url::Url; use knot_atproto::{Atproto, knot_did_document}; use knot_cob::{CobHome, CobStore}; -use knot_cobs::{Grant, Registration, RegistryChange}; +use knot_cobs::{Registration, RegistryChange}; use knot_events::{EventLog, GlobalSubscriberLimit, PerPeerSubscriberLimit, SubscriberGate}; use knot_git::{ EntryKind, Identity, Layout, NewCommit, RefUpdate, Repo, StagedAction, StagedChange, @@ -26,7 +26,7 @@ use knot_runtime::{ }; use knot_secrets::{MasterKey, SealedStore}; use knot_types::{ - AccountDid, AdmissionPolicy, AuthorName, BranchName, Email, KnotHostname, KnotId, + AccountDid, AdmissionPolicy, AuthorName, BranchName, DidRkey, Email, KnotHostname, KnotId, KnotServiceUrl, Oid, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, UnixSeconds, }; use knot_xrpc::{ @@ -94,13 +94,66 @@ impl Faults { } } +#[derive(Clone, Copy, PartialEq, Eq, Hash)] +pub(crate) enum AcceptanceKind { + Membership, + Collaboration, +} + +impl AcceptanceKind { + pub(crate) fn collection(self) -> &'static str { + match self { + AcceptanceKind::Membership => knot_consent::OF_MEMBERSHIP, + AcceptanceKind::Collaboration => knot_consent::OF_COLLABORATION, + } + } + + pub(crate) fn accept_method(self) -> &'static str { + match self { + AcceptanceKind::Membership => "sh.tangled.knot.acceptMembership", + AcceptanceKind::Collaboration => "sh.tangled.repo.acceptCollaboration", + } + } + + fn named_by(collection: &str) -> Option { + match collection { + knot_consent::OF_MEMBERSHIP => Some(AcceptanceKind::Membership), + knot_consent::OF_COLLABORATION => Some(AcceptanceKind::Collaboration), + _ => None, + } + } +} + +#[derive(Clone, PartialEq, Eq, Hash)] +pub(crate) struct AcceptanceRecord { + pub(crate) author: AccountDid, + pub(crate) kind: AcceptanceKind, + pub(crate) rkey: DidRkey, +} + +#[derive(Default)] +pub(crate) struct Published(Mutex>); + +impl Published { + pub(crate) fn write(&self, record: AcceptanceRecord) { + self.0.lock().expect("published records lock").insert(record); + } + + fn holds(&self, record: &AcceptanceRecord) -> bool { + self.0.lock().expect("published records lock").contains(record) + } +} + pub(crate) struct Harness { _dir: Option, pub clock: Arc, pub faults: Arc, + pub published: Arc, pub layout: Layout, pub index: Arc, pub knot_aud: KnotId, + pub knot_rkey: DidRkey, + pub plc_signer: K256Signer, router: Router, pub admin: Actor, pub strangers: Vec, @@ -156,9 +209,19 @@ impl Harness { let clock = Arc::new(ManualClock::new(UnixMicros::new(START_MICROS))); let pubkeys = pubkey_map(&admin, &strangers); let faults = Arc::new(Faults::default()); + let published = Arc::new(Published::default()); let plc = Url::parse("https://plc.directory/").expect("plc url"); - let responder = build_responder(Arc::clone(&faults), pubkeys, HashMap::new(), plc.clone()); + let plc_signer = plc_signer(); + let responder = build_responder( + Arc::clone(&faults), + Arc::clone(&published), + pubkeys, + HashMap::new(), + plc.clone(), + plc_signer.public_key(), + ); let knot_aud = knot.clone(); + let knot_rkey = DidRkey::new(knot.as_str()).expect("knot did is a record key"); let did_document = knot_did_document(&knot, &knot_pubkey, &knot_service_url); let router = assemble_router(StateParts { layout: layout.clone(), @@ -183,9 +246,12 @@ impl Harness { _dir: Some(dir), clock, faults, + published, layout, index, knot_aud, + knot_rkey, + plc_signer, router, admin, strangers, @@ -289,14 +355,19 @@ impl Harness { let clock = Arc::new(ManualClock::new(UnixMicros::new(START_MICROS))); let faults = Arc::new(Faults::default()); + let published = Arc::new(Published::default()); + let plc_signer = plc_signer(); let responder = build_responder( Arc::clone(&faults), + Arc::clone(&published), pubkey_map(&admin, &strangers), did_overrides, plc.clone(), + plc_signer.public_key(), ); let entropy = Arc::new(SeededEntropy::new(seed ^ 0x5eed_0050)); let knot_aud = knot.clone(); + let knot_rkey = DidRkey::new(knot.as_str()).expect("knot did is a record key"); let did_document = knot_did_document(&knot, &knot_pubkey, &knot_service_url); let router = assemble_router(StateParts { layout: layout.clone(), @@ -320,9 +391,12 @@ impl Harness { _dir: Some(scratch), clock, faults, + published, layout, index, knot_aud, + knot_rkey, + plc_signer, router, admin, strangers, @@ -373,7 +447,7 @@ impl Harness { let _ = self.index.ensure_collaborators(repo); RepoCollaborators { repo: repo.clone(), - subjects: sorted(grant_subjects(self.index.collaborator_entries(repo))), + subjects: subjects(self.index.effective_collaborators(repo)), } }) .collect(); @@ -385,8 +459,8 @@ impl Harness { Snapshot { round, clock_micros: self.clock.now_unix_micros(), - members: sorted(grant_subjects(self.index.member_entries())), - blocked: sorted(grant_subjects(self.index.blocked_entries())), + members: subjects(self.index.effective_members()), + blocked: subjects(self.index.blocked()), repos: repo_list, collaborators, } @@ -411,17 +485,19 @@ fn pubkey_map(admin: &Actor, strangers: &[Actor]) -> HashMap K256Signer { + K256Signer::generate(&SeededEntropy::new(7)) +} + fn build_responder( faults: Arc, + published: Arc, pubkeys: HashMap, did_overrides: HashMap, plc: Url, + plc_pubkey: PublicKeyBytes, ) -> Responder { - let plc_signer = K256Signer::generate(&SeededEntropy::new(7)); Box::new(move |request: &HttpRequest| { - if request.method == http::Method::POST { - return Ok(ok_body(Bytes::new())); - } let not_found = || HttpResponse { status: http::StatusCode::NOT_FOUND, headers: http::HeaderMap::new(), @@ -435,6 +511,17 @@ fn build_responder( "identity resolution dropped by simulation".to_string(), )); } + if request.method == http::Method::POST { + 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)); + return Ok(match written { + true => ok_body(Bytes::new()), + false => not_found(), + }); + } let is_plc = request.url.host() == plc.host(); let requested_did = if is_plc { request @@ -452,7 +539,7 @@ fn build_responder( return Ok(ok_body(did_doc(&requested_did, sec1))); } if is_plc { - return Ok(ok_body(did_doc(&requested_did, &plc_signer.public_key()))); + return Ok(ok_body(did_doc(&requested_did, &plc_pubkey))); } match pubkeys.get(&host) { Some(sec1) => Ok(ok_body(did_doc(&requested_did, sec1))), @@ -461,6 +548,19 @@ fn build_responder( }) } +fn query_value(url: &Url, key: &str) -> Option { + url.query_pairs() + .find_map(|(param, value)| (param == key).then(|| value.into_owned())) +} + +fn named_acceptance(url: &Url) -> Option { + Some(AcceptanceRecord { + author: AccountDid::new(query_value(url, "repo")?).ok()?, + kind: AcceptanceKind::named_by(&query_value(url, "collection")?)?, + rkey: DidRkey::new(query_value(url, "rkey")?).ok()?, + }) +} + fn did_doc(did: &AccountDid, sec1: &PublicKeyBytes) -> Bytes { let did = did.as_str(); let multikey = knot_types::crypto::multikey(0xe7, sec1.as_bytes()); @@ -583,26 +683,15 @@ fn register_seed_repo( .expect("register seed repo"); } -fn grant_subjects(resolved: Resolved>) -> Vec { +fn subjects(resolved: Resolved>) -> Vec { match resolved { - Resolved::Ready(grants) => grants - .into_iter() - .map(|grant| grant.subject.clone()) - .collect(), + Resolved::Ready(subjects) => subjects, Resolved::Warming => Vec::new(), } } -fn sorted(mut values: Vec) -> Vec { - values.sort(); - values -} - fn distinct_members(index: &Index) -> Vec { - let mut members = grant_subjects(index.member_entries()); - members.sort(); - members.dedup(); - members + subjects(index.effective_members()) } fn materialize_scratch(target: &Path, knot: &KnotId) -> Result { diff --git a/knot2/crates/knot-sim/src/lib.rs b/knot2/crates/knot-sim/src/lib.rs index 755b3a947..ba01c3441 100644 --- a/knot2/crates/knot-sim/src/lib.rs +++ b/knot2/crates/knot-sim/src/lib.rs @@ -20,6 +20,13 @@ pub async fn run(seed: u64, rounds: u32) -> Trace { workload::execute(harness, seed, plan).await } +pub fn error_digest(error: &str, message: &str) -> u64 { + workload::body_digest(&workload::encode_body(serde_json::json!({ + "error": error, + "message": message, + }))) +} + pub fn predict(seed: u64, rounds: u32) -> Projection { workload::predict(seed, rounds) } diff --git a/knot2/crates/knot-sim/src/trace.rs b/knot2/crates/knot-sim/src/trace.rs index 01804a1a6..1b4ea93dc 100644 --- a/knot2/crates/knot-sim/src/trace.rs +++ b/knot2/crates/knot-sim/src/trace.rs @@ -66,6 +66,7 @@ pub struct RepoCollaborators { pub struct Projection { pub members: Vec, pub blocked: Vec, + pub collaborators: Vec>, } #[derive(Debug, Clone, PartialEq, Eq, Serialize)] diff --git a/knot2/crates/knot-sim/src/workload.rs b/knot2/crates/knot-sim/src/workload.rs index 24add6cdc..2ad95f4ff 100644 --- a/knot2/crates/knot-sim/src/workload.rs +++ b/knot2/crates/knot-sim/src/workload.rs @@ -6,24 +6,24 @@ use futures::stream::StreamExt; use http::Method; use http::header::AUTHORIZATION; use serde_json::{Value, json}; -use std::collections::BTreeSet; +use std::collections::{BTreeMap, BTreeSet}; use std::net::SocketAddr; use std::sync::Arc; use tower::ServiceExt; use knot_runtime::{Entropy, K256Signer, SeededEntropy, Signer}; -use knot_types::{AccountDid, HttpStatus, KnotHostname, KnotId, RepoDid, UnixSeconds}; +use knot_types::{AccountDid, DidRkey, HttpStatus, KnotHostname, RepoDid, UnixSeconds}; -use crate::harness::{Harness, SUBJECT_DIDS}; +use crate::harness::{AcceptanceKind, AcceptanceRecord, Harness, SUBJECT_DIDS}; use crate::trace::{OperationIndex, Outcome, Projection, RoundNumber, Step, Trace, fnv1a}; const SKEW_BACKDATE_SECS: i64 = 600; const SKEW_LIFETIME_SECS: i64 = 60; -#[derive(Clone, Copy)] +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] struct RepoIndex(usize); -#[derive(Clone, Copy)] +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] struct SubjectIndex(usize); #[derive(Clone, Copy)] @@ -42,6 +42,7 @@ enum ReadOp { Tree(RepoIndex), Blob(RepoIndex), Languages(RepoIndex), + ListCollaborators(RepoIndex), } impl ReadOp { @@ -58,6 +59,7 @@ impl ReadOp { ReadOp::Tree(_) => "tree", ReadOp::Blob(_) => "blob", ReadOp::Languages(_) => "languages", + ReadOp::ListCollaborators(_) => "listCollaborators", } } } @@ -86,15 +88,85 @@ impl AdminOp { } #[derive(Clone, Copy)] -enum Planned { - Read { op: ReadOp, killed: bool }, +enum Mutation { Admin { op: AdminOp, skew: bool }, Probe { stranger: StrangerIndex, drop: bool }, Maintain { repo: RepoIndex }, } +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +enum Offered { + Membership(SubjectIndex), + Collaboration(RepoIndex, SubjectIndex), +} + +#[derive(Clone, Copy)] +struct Answer { + offer: Offered, + published: bool, + skew: bool, +} + +impl Answer { + fn name(self) -> &'static str { + match self.offer { + Offered::Membership(_) => "acceptMembership", + Offered::Collaboration(_, _) => "acceptCollaboration", + } + } + + fn fault(self) -> &'static str { + match (self.skew, self.published) { + (true, _) => "clock_skew", + (false, false) => "unpublished", + (false, true) => "none", + } + } +} + +#[derive(Clone, Copy)] +struct Observation { + op: ReadOp, + killed: bool, +} + +#[derive(Clone, Copy)] +enum Planned { + Mutate(Mutation), + Answer(Answer), + Observe(Observation), +} + +impl Planned { + fn creates_a_repo(self) -> bool { + matches!( + self, + Planned::Mutate(Mutation::Admin { + op: AdminOp::CreateRepo(_), + skew: false, + }) + ) + } +} + +enum Phase { + Mutate(Vec), + Answer(Vec), + Observe(Vec), +} + +impl Phase { + fn ops(&self) -> Vec { + match self { + Phase::Mutate(ops) => ops.iter().copied().map(Planned::Mutate).collect(), + Phase::Answer(ops) => ops.iter().copied().map(Planned::Answer).collect(), + Phase::Observe(ops) => ops.iter().copied().map(Planned::Observe).collect(), + } + } +} + pub(crate) struct Round { - ops: Vec, + phase: Phase, advance: std::time::Duration, } @@ -120,53 +192,108 @@ pub(crate) fn plan(seed: u64, rounds: u32, subjects: usize) -> Vec { let mut available: u32 = 1; let mut rkey: u32 = 0; let mut stranger: usize = 0; + let mut unanswered: BTreeSet = BTreeSet::new(); (0..rounds) .map(|round| { - let (ops, fresh_repos) = if round % 2 == 0 { - mutate_round(&rng, subjects, available, &mut rkey, &mut stranger) - } else { - (read_round(&rng, available), 0) + let (phase, fresh_repos) = match round % 3 { + 0 => { + let (ops, fresh_repos) = mutate_round( + &rng, + subjects, + available, + &mut rkey, + &mut stranger, + &mut unanswered, + ); + (Phase::Mutate(ops), fresh_repos) + } + 1 => ( + Phase::Answer(answer_round(&rng, subjects, &mut unanswered)), + 0, + ), + _ => (Phase::Observe(observe_round(&rng, available)), 0), }; let advance = std::time::Duration::from_micros(rng.below(3_000_000)); available += fresh_repos; - Round { ops, advance } + Round { phase, advance } }) .collect() } +#[derive(Default)] +struct Model { + members: BTreeSet, + blocked: BTreeSet, + collaborators: BTreeMap>, + offered: BTreeSet, + written: BTreeSet, +} + pub(crate) fn predict(seed: u64, rounds: u32) -> Projection { - let (members, blocked) = plan(seed, rounds, SUBJECT_DIDS.len()) - .iter() - .flat_map(|round| round.ops.iter()) - .fold( - (BTreeSet::::new(), BTreeSet::::new()), - |(mut members, mut blocked), planned| { - let subject_did = |subject: &SubjectIndex| { - AccountDid::new(SUBJECT_DIDS[subject.0]).expect("subject did") - }; - if let Planned::Admin { op, skew: false } = planned { - match op { - AdminOp::AddMember(subject) => { - members.insert(subject_did(subject)); - } - AdminOp::RemoveMember(subject) => { - members.remove(&subject_did(subject)); - } - AdminOp::Ban(subject) => { - blocked.insert(subject_did(subject)); - } - AdminOp::Unban(subject) => { - blocked.remove(&subject_did(subject)); - } - AdminOp::CreateRepo(_) | AdminOp::AddCollaborator(_, _) => {} - } + let subject_did = + |subject: SubjectIndex| AccountDid::new(SUBJECT_DIDS[subject.0]).expect("subject did"); + let plan = plan(seed, rounds, SUBJECT_DIDS.len()); + let ops = || plan.iter().flat_map(|round| round.phase.ops()); + let repos = 1 + ops().filter(|planned| planned.creates_a_repo()).count(); + let model = ops().fold(Model::default(), |mut model, planned| { + match planned { + Planned::Mutate(Mutation::Admin { op, skew: false }) => match op { + AdminOp::AddMember(subject) => { + model.offered.insert(Offered::Membership(subject)); } - (members, blocked) + AdminOp::RemoveMember(subject) => { + model.offered.remove(&Offered::Membership(subject)); + model.members.remove(&subject_did(subject)); + } + AdminOp::Ban(subject) => { + model.blocked.insert(subject_did(subject)); + } + AdminOp::Unban(subject) => { + model.blocked.remove(&subject_did(subject)); + } + AdminOp::AddCollaborator(repo, subject) => { + model.offered.insert(Offered::Collaboration(repo, subject)); + } + AdminOp::CreateRepo(_) => {} }, - ); + Planned::Answer(Answer { + offer, + published, + skew, + }) => { + if published { + model.written.insert(offer); + } + let admitted = + !skew && model.offered.contains(&offer) && model.written.contains(&offer); + match (admitted, offer) { + (false, _) => {} + (true, Offered::Membership(subject)) => { + model.members.insert(subject_did(subject)); + } + (true, Offered::Collaboration(repo, subject)) => { + model + .collaborators + .entry(repo.0) + .or_default() + .insert(subject_did(subject)); + } + } + } + Planned::Mutate(Mutation::Admin { skew: true, .. }) + | Planned::Mutate(Mutation::Probe { .. }) + | Planned::Mutate(Mutation::Maintain { .. }) + | Planned::Observe(_) => {} + } + model + }); + let mut rosters = model.collaborators; Projection { - members: members.into_iter().collect(), - blocked: blocked.into_iter().collect(), + 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()) + .collect(), } } @@ -176,8 +303,9 @@ fn mutate_round( available: u32, rkey: &mut u32, stranger: &mut usize, -) -> (Vec, u32) { - let mut ops: Vec = (0..subjects) + unanswered: &mut BTreeSet, +) -> (Vec, u32) { + let mut ops: Vec = (0..subjects) .filter(|_| rng.chance(2, 3)) .map(|subject| { let subject = SubjectIndex(subject); @@ -187,7 +315,7 @@ fn mutate_round( 2 => AdminOp::Ban(subject), _ => AdminOp::Unban(subject), }; - Planned::Admin { + Mutation::Admin { op, skew: rng.chance(1, 5), } @@ -197,16 +325,19 @@ fn mutate_round( (0..available) .filter(|_| rng.chance(1, 2)) .for_each(|repo| { + let repo = RepoIndex(repo as usize); let planned = if rng.chance(1, 2) { let subject = SubjectIndex(rng.below(subjects as u64) as usize); - Planned::Admin { - op: AdminOp::AddCollaborator(RepoIndex(repo as usize), subject), - skew: rng.chance(1, 6), + let skew = rng.chance(1, 6); + if !skew { + unanswered.insert(Offered::Collaboration(repo, subject)); } - } else { - Planned::Maintain { - repo: RepoIndex(repo as usize), + Mutation::Admin { + op: AdminOp::AddCollaborator(repo, subject), + skew, } + } else { + Mutation::Maintain { repo } }; ops.push(planned); }); @@ -219,7 +350,7 @@ fn mutate_round( if !skew { fresh_repos += 1; } - ops.push(Planned::Admin { + ops.push(Mutation::Admin { op: AdminOp::CreateRepo(key), skew, }); @@ -228,7 +359,7 @@ fn mutate_round( (0..rng.below(3)).for_each(|_| { let stranger_index = StrangerIndex(*stranger); *stranger += 1; - ops.push(Planned::Probe { + ops.push(Mutation::Probe { stranger: stranger_index, drop: rng.chance(1, 2), }); @@ -237,15 +368,38 @@ fn mutate_round( (ops, fresh_repos) } -fn read_round(rng: &Rng, available: u32) -> Vec { - let mut ops: Vec = [ +fn answer_round(rng: &Rng, subjects: usize, unanswered: &mut BTreeSet) -> Vec { + let answer = |offer| Answer { + offer, + published: rng.chance(3, 4), + skew: rng.chance(1, 6), + }; + let mut ops: Vec = (0..subjects) + .filter(|_| rng.chance(1, 2)) + .map(|subject| answer(Offered::Membership(SubjectIndex(subject)))) + .collect(); + + let answering: Vec = unanswered + .iter() + .copied() + .filter(|_| rng.chance(2, 3)) + .collect(); + answering.iter().for_each(|offer| { + unanswered.remove(offer); + }); + ops.extend(answering.into_iter().map(answer)); + ops +} + +fn observe_round(rng: &Rng, available: u32) -> Vec { + let mut ops: Vec = [ ReadOp::Version, ReadOp::Owner, ReadOp::ListMembers, ReadOp::DidJson, ] .into_iter() - .map(|op| Planned::Read { + .map(|op| Observation { op, killed: rng.chance(1, 5), }) @@ -254,16 +408,17 @@ fn read_round(rng: &Rng, available: u32) -> Vec { .filter(|_| rng.chance(2, 3)) .for_each(|repo| { let repo = RepoIndex(repo as usize); - let op = match rng.below(7) { + let op = match rng.below(8) { 0 => ReadOp::Branches(repo), 1 => ReadOp::Log(repo), 2 => ReadOp::DescribeRepo(repo), 3 => ReadOp::InfoRefs(repo), 4 => ReadOp::Tree(repo), 5 => ReadOp::Blob(repo), + 6 => ReadOp::ListCollaborators(repo), _ => ReadOp::Languages(repo), }; - ops.push(Planned::Read { + ops.push(Observation { op, killed: rng.chance(1, 5), }); @@ -290,10 +445,10 @@ pub(crate) async fn execute(harness: Arc, seed: u64, plan: Vec) let harness = Arc::clone(harness); async move { let round_no = RoundNumber::new(round_index as u32); - let drops = arm_drops(&harness, &round.ops); + let ops = round.phase.ops(); + let drops = arm_drops(&harness, &ops); let repos = Arc::new(repos); - let tasks = round - .ops + let tasks = ops .iter() .enumerate() .map(|(index, planned)| { @@ -325,19 +480,8 @@ pub(crate) async fn execute(harness: Arc, seed: u64, plan: Vec) .iter() .filter_map(|result| result.created.clone()) .collect(); - let planned_creates = round - .ops - .iter() - .filter(|planned| { - matches!( - planned, - Planned::Admin { - op: AdminOp::CreateRepo(_), - skew: false, - } - ) - }) - .count(); + let planned_creates = + ops.iter().filter(|planned| planned.creates_a_repo()).count(); assert_eq!( planned_creates, created.len(), @@ -380,10 +524,10 @@ fn arm_drops(harness: &Harness, ops: &[Planned]) -> Vec { let hosts: Vec = ops .iter() .filter_map(|planned| match planned { - Planned::Probe { + Planned::Mutate(Mutation::Probe { stranger, drop: true, - } => Some(harness.strangers[stranger.0].host.clone()), + }) => Some(harness.strangers[stranger.0].host.clone()), _ => None, }) .collect(); @@ -408,7 +552,7 @@ async fn run_op( }; match planned { - Planned::Maintain { repo } => { + Planned::Mutate(Mutation::Maintain { repo }) => { let repo = repo_at(repos, repo); let outcome = match harness.maintain(repo) { Ok(()) => Outcome::Answered { @@ -425,7 +569,7 @@ async fn run_op( created: None, } } - Planned::Read { op, killed } => { + Planned::Observe(Observation { op, killed }) => { let request = read_request(repos, op); if killed { drive_kill(harness.router(), request.method, &request.uri, request.body).await; @@ -455,7 +599,7 @@ async fn run_op( created: None, } } - Planned::Admin { op, skew } => { + Planned::Mutate(Mutation::Admin { op, skew }) => { let request = admin_request(harness, repos, op, skew, round, index); let (status, body) = http_call( harness.router(), @@ -482,7 +626,7 @@ async fn run_op( created, } } - Planned::Probe { stranger, drop } => { + Planned::Mutate(Mutation::Probe { stranger, drop }) => { let request = probe_request(harness, stranger, round, index); let (status, body) = http_call( harness.router(), @@ -505,6 +649,33 @@ async fn run_op( created: None, } } + Planned::Answer(answer) => { + let record = acceptance_of(harness, repos, answer); + if answer.published { + harness.published.write(record.clone()); + } + let request = answer_request(harness, &record, answer.skew, round, index); + let (status, body) = http_call( + harness.router(), + request.method, + &request.uri, + request.token.as_deref(), + request.body, + ) + .await; + OpResult { + step: make( + answer.name(), + request.actor, + answer.fault(), + Outcome::Answered { + status, + body: body_digest(&body), + }, + ), + created: None, + } + } } } @@ -539,6 +710,9 @@ fn read_request(repos: &[RepoDid], op: ReadOp) -> Request { ReadOp::DescribeRepo(repo) => repo_get("describeRepo", "repoDid", repo_at(repos, repo)), ReadOp::Tree(repo) => repo_get("tree", "repo", repo_at(repos, repo)), ReadOp::Languages(repo) => repo_get("languages", "repo", repo_at(repos, repo)), + ReadOp::ListCollaborators(repo) => { + repo_get("listCollaborators", "subject", repo_at(repos, repo)) + } ReadOp::Blob(repo) => Request { method: Method::GET, uri: format!( @@ -650,6 +824,58 @@ fn probe_request( } } +fn acceptance_of(harness: &Harness, repos: &[RepoDid], answer: Answer) -> AcceptanceRecord { + let (subject, kind, rkey) = match answer.offer { + Offered::Membership(subject) => ( + subject, + AcceptanceKind::Membership, + harness.knot_rkey.clone(), + ), + Offered::Collaboration(repo, subject) => ( + subject, + AcceptanceKind::Collaboration, + DidRkey::new(repo_at(repos, repo).as_str()).expect("repo did is a record key"), + ), + }; + AcceptanceRecord { + author: harness.subjects[subject.0].clone(), + kind, + rkey, + } +} + +fn answer_request( + harness: &Harness, + record: &AcceptanceRecord, + skew: bool, + round: RoundNumber, + index: OperationIndex, +) -> Request { + let (author, nsid) = (record.author.as_str(), record.kind.accept_method()); + let token = mint( + &harness.plc_signer, + &record.author, + &record.rkey, + nsid, + jwt_window(harness, skew), + round, + index, + ); + Request { + method: Method::POST, + uri: format!("/xrpc/{nsid}"), + token: Some(token), + body: encode_body(json!({ + "acceptance": format!( + "at://{author}/{}/{}", + record.kind.collection(), + record.rkey.as_str() + ) + })), + actor: author.to_string(), + } +} + fn get(path: &str) -> Request { Request { method: Method::GET, @@ -716,7 +942,7 @@ pub(crate) fn jwt_window(harness: &Harness, skew: bool) -> (UnixSeconds, UnixSec pub(crate) fn mint( signer: &K256Signer, issuer: &AccountDid, - aud: &KnotId, + aud: &impl AsRef, nsid: &str, window: (UnixSeconds, UnixSeconds), round: RoundNumber, @@ -726,7 +952,7 @@ pub(crate) fn mint( let payload = URL_SAFE_NO_PAD.encode( serde_json::to_vec(&json!({ "iss": issuer.as_str(), - "aud": aud.as_str(), + "aud": aud.as_ref(), "iat": window.0.get(), "exp": window.1.get(), "jti": format!("sim-{}-{}", round.get(), index.get()), diff --git a/knot2/crates/knot-sim/tests/reproducible.rs b/knot2/crates/knot-sim/tests/reproducible.rs index 1fa48a255..d983911f4 100644 --- a/knot2/crates/knot-sim/tests/reproducible.rs +++ b/knot2/crates/knot-sim/tests/reproducible.rs @@ -1,4 +1,4 @@ -use std::collections::BTreeSet; +use std::collections::{BTreeMap, BTreeSet}; use futures::future::join_all; use futures::stream::StreamExt; @@ -11,16 +11,24 @@ fn status_of(step: &Step) -> Option { } } +fn digest_of(step: &Step) -> Option { + match step.outcome { + Outcome::Answered { body, .. } => Some(body), + Outcome::Killed => None, + } +} + fn union_has(traces: &[Trace], predicate: impl Fn(&Step) -> bool) -> bool { traces .iter() .any(|trace| trace.steps.iter().any(&predicate)) } -const READ_OPS: [&str; 11] = [ +const READ_OPS: [&str; 12] = [ "version", "owner", "listMembers", + "listCollaborators", "didJson", "infoRefs", "branches", @@ -142,13 +150,34 @@ async fn the_injected_failures_are_correlated_with_their_observable_outcomes() { "resolved stranger with no fault must be denied 403 by access-control layer" ); + let absent = knot_sim::error_digest( + "Forbidden", + "no acceptance record in your repository, so write it first", + ); + let refused = traces + .iter() + .flat_map(|trace| trace.steps.iter()) + .filter(|step| { + step.fault == "unpublished" + && status_of(step) == Some(403) + && digest_of(step) == Some(absent) + }) + .count(); + assert!( + refused > 0, + "no 403 matched this digest, so the fault never fired or the literal drifted" + ); + [ "addMember", + "acceptMembership", + "acceptCollaboration", "addCollaborator", "createRepo", "maintain", "describeRepo", "listMembers", + "listCollaborators", "infoRefs", "didJson", "branches", @@ -190,6 +219,15 @@ async fn the_final_projection_matches_an_independent_model() { last.blocked, model.blocked, "seed {seed}: executed blocklist projection must equal the independent model" ); + assert_eq!( + last.collaborators + .iter() + .map(|roster| &roster.subjects) + .collect::>(), + model.collaborators.iter().collect::>(), + "seed {seed}: the run must host the repositories the model counts and open each \ + one to exactly the collaborators it predicts" + ); }); assert!( @@ -200,6 +238,12 @@ async fn the_final_projection_matches_an_independent_model() { predicted.iter().any(|model| !model.blocked.is_empty()), "oracle is vacuous unless at least one seed predicts a non-empty blocklist" ); + assert!( + predicted + .iter() + .any(|model| model.collaborators.iter().any(|repo| !repo.is_empty())), + "oracle is vacuous unless at least one seed predicts a repo with a collaborator" + ); } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] @@ -236,3 +280,48 @@ async fn the_final_projection_state_is_seed_stable() { last.repos ); } + +const OFFERS: [&str; 2] = ["addMember", "addCollaborator"]; +const ANSWERS: [&str; 2] = ["acceptMembership", "acceptCollaboration"]; +const ROSTER_READS: [&str; 2] = ["listMembers", "listCollaborators"]; + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn roster_read_and_roster_write_land_in_separate_rounds() { + let traces = join_all([1u64, 7, 13, 42, 99, 2026].map(|seed| knot_sim::run(seed, 18))).await; + let rounds: BTreeMap<(u64, u32), BTreeSet<&str>> = traces + .iter() + .flat_map(|trace| { + trace + .steps + .iter() + .map(move |step| ((trace.seed, step.round.get()), step.op)) + }) + .fold(BTreeMap::new(), |mut rounds, (round, op)| { + rounds.entry(round).or_default().insert(op); + rounds + }); + let holds = |ops: &BTreeSet<&str>, names: &[&str]| names.iter().any(|name| ops.contains(name)); + let writes = |ops: &BTreeSet<&str>| { + holds(ops, &OFFERS) || holds(ops, &ANSWERS) || ops.contains("removeMember") + }; + + rounds.iter().for_each(|((seed, round), ops)| { + assert!( + !(holds(ops, &OFFERS) && holds(ops, &ANSWERS)), + "seed {seed} round {round}: an offer and its acceptance ran in one round" + ); + assert!( + !(writes(ops) && holds(ops, &ROSTER_READS)), + "seed {seed} round {round}: a roster listing ran beside a roster write" + ); + }); + + [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-ssh/Cargo.toml b/knot2/crates/knot-ssh/Cargo.toml index 3601a5764..47cd14c3d 100644 --- a/knot2/crates/knot-ssh/Cargo.toml +++ b/knot2/crates/knot-ssh/Cargo.toml @@ -14,6 +14,7 @@ knot-lfs = { workspace = true } knot-pack = { workspace = true } knot-index = { workspace = true } knot-acl = { workspace = true } +knot-consent = { workspace = true } knot-atproto = { workspace = true } knot-cob = { workspace = true } knot-cobs = { workspace = true } diff --git a/knot2/crates/knot-ssh/src/exec.rs b/knot2/crates/knot-ssh/src/exec.rs index 4fcfc16b0..0a137e7e2 100644 --- a/knot2/crates/knot-ssh/src/exec.rs +++ b/knot2/crates/knot-ssh/src/exec.rs @@ -810,8 +810,12 @@ async fn authorize_push( ) -> PushAuth { match resolve_pusher(state, credential, repo, peer).await { PusherLookup::Matched(did) => { + let pds = knot_consent::Pds::reading(&state.atproto, &state.knot_did); + let consents = state.index.consents(); + let now = state.atproto.now().seconds(); + let consent = consents.of_collaboration(&pds, repo, &did, now).await; let acl = KnotAcl::new(&state.admins, state.admission, &state.index); - match can_push(&acl, &did, repo).is_allowed() { + match can_push(&acl, &consent).is_allowed() { true => PushAuth::Allowed(did), false => PushAuth::Refused { reason: "unauthorized", @@ -895,11 +899,10 @@ async fn resolve_pusher( let target = repo.clone(); let _ = tokio::task::spawn_blocking(move || index.ensure_collaborators(&target)).await; } - let collaborators = match state.index.collaborators_of(repo) { - Resolved::Ready(collaborators) => collaborators, - _ => Vec::new(), - }; - let candidates: Vec = owner.into_iter().chain(collaborators).collect(); + let collaborators = state.index.effective_collaborators(repo).ready_or_default(); + let invited = state.index.invited_collaborators(repo).ready_or_default(); + let grantable: Vec = owner.into_iter().chain(collaborators).collect(); + let candidates: Vec = grantable.iter().cloned().chain(invited).collect(); let now = state.atproto.now().seconds(); if let Some(publisher) = state.index.keys().publisher_among(&candidates, key, now) { return PusherLookup::Matched(publisher); @@ -917,7 +920,7 @@ async fn resolve_pusher( "push check has every candidate's keys on file, and the candidates don't publish \ the offered key" ); - return PusherLookup::Unmatched(candidates); + return PusherLookup::Unmatched(grantable); } let lease = state.key_ttl.lease_from(now); let unresolved = Arc::new(AtomicBool::new(false)); @@ -961,7 +964,7 @@ async fn resolve_pusher( match read.into_iter().flatten().next() { Some(did) => PusherLookup::Matched(did), None if unresolved.load(Ordering::Relaxed) => PusherLookup::Unavailable, - None => PusherLookup::Unmatched(candidates), + None => PusherLookup::Unmatched(grantable), } } diff --git a/knot2/crates/knot-ssh/src/identity.rs b/knot2/crates/knot-ssh/src/identity.rs index f2602f4b2..9fa1851b2 100644 --- a/knot2/crates/knot-ssh/src/identity.rs +++ b/knot2/crates/knot-ssh/src/identity.rs @@ -27,6 +27,7 @@ enum Claimed { Unreadable, } +#[derive(Clone)] pub(crate) enum Verdict { Identified(AccountDid), Offered, @@ -66,6 +67,9 @@ fn against_key_set( false => { if miss_worth_a_reread(state, peer) { state.index.keys().note_miss(); + if state.index.any_invited_collaborator() { + return Verdict::Offered; + } } Verdict::Refused } diff --git a/knot2/crates/knot-ssh/src/lib.rs b/knot2/crates/knot-ssh/src/lib.rs index 64897917d..5023a180e 100644 --- a/knot2/crates/knot-ssh/src/lib.rs +++ b/knot2/crates/knot-ssh/src/lib.rs @@ -17,7 +17,9 @@ 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}; +use knot_types::{ + AccountDid, ActorId, AdmissionPolicy, AppviewEndpoint, CiLogsAddr, KnotHostname, KnotId, +}; use russh::keys::ssh_key::rand_core; use russh::keys::{Algorithm, PrivateKey, ssh_key}; use russh::server::{Config, Server as _}; @@ -63,6 +65,7 @@ pub struct SshState { knot_actor: ActorId, events: Arc>, hostname: KnotHostname, + knot_did: KnotId, appview: AppviewEndpoint, admins: BTreeSet, admission: AdmissionPolicy, @@ -123,6 +126,7 @@ impl SshState { languages_push_budget, ci_logs, } = config; + let knot_did = hostname.knot_did(); Self { layout, index, @@ -130,6 +134,7 @@ impl SshState { knot_actor, events, hostname, + knot_did, appview, admins, admission, diff --git a/knot2/crates/knot-ssh/src/server.rs b/knot2/crates/knot-ssh/src/server.rs index 829ca390c..6c818c1cb 100644 --- a/knot2/crates/knot-ssh/src/server.rs +++ b/knot2/crates/knot-ssh/src/server.rs @@ -38,6 +38,7 @@ pub(crate) struct KnotSession { channels: HashMap>, protocols: HashSet, peer: Option, + decided: Option<(String, OfferedKey, Verdict)>, } fn reject() -> Auth { @@ -57,6 +58,7 @@ impl KnotSession { channels: HashMap::new(), protocols: HashSet::new(), peer, + decided: None, } } } @@ -68,6 +70,12 @@ impl KnotSession { public_key: &ssh_key::PublicKey, ) -> Option<(Verdict, OfferedKey)> { let key = OfferedKey::from_bytes(public_key.to_bytes().ok()?); + if let Some((asked_as, asked, verdict)) = &self.decided + && asked_as == user + && *asked == key + { + return Some((verdict.clone(), key)); + } let verdict = identity::verify( &self.state, OwnerRef::parse(user), @@ -76,6 +84,7 @@ impl KnotSession { &mut self.asserted, ) .await; + self.decided = Some((user.to_string(), key.clone(), verdict.clone())); Some((verdict, key)) } } diff --git a/knot2/crates/knot-ssh/tests/ssh_push.rs b/knot2/crates/knot-ssh/tests/ssh_push.rs index 64afdc14c..b61c094ce 100644 --- a/knot2/crates/knot-ssh/tests/ssh_push.rs +++ b/knot2/crates/knot-ssh/tests/ssh_push.rs @@ -171,6 +171,7 @@ fn fake_http(published_line: String) -> impl knot_runtime::HttpTransport { struct Accounts { identities: Arc>>>, unreachable: Arc>>, + acceptances: Arc>>, listings: Arc, } @@ -187,6 +188,14 @@ impl Accounts { self } + fn accepting(self, subject: &str, rkey: &str) -> Self { + self.acceptances + .lock() + .unwrap() + .insert(format!("{subject}/{rkey}")); + self + } + fn restore(&self, did: &str) { self.unreachable.lock().unwrap().remove(did); } @@ -229,6 +238,14 @@ fn multi_http(accounts: Accounts) -> impl knot_runtime::HttpTransport { .find(|(key, _)| key == "repo") .map(|(_, value)| value.into_owned()) .unwrap_or_default(); + if request.url.path().ends_with("com.atproto.repo.getRecord") { + let found = request.url.query_pairs().find(|(key, _)| key == "rkey"); + let rkey = found.map(|(_, rkey)| rkey.into_owned()).unwrap_or_default(); + let record = format!("{repo}/{rkey}"); + let seen = accounts.acceptances.lock().unwrap().contains(&record); + let body = seen.then(|| ok_body(b"{}".to_vec())); + return Ok(body.unwrap_or_else(not_found)); + } accounts.listings.fetch_add(1, Ordering::SeqCst); if accounts.unreachable.lock().unwrap().contains(&repo) { return Ok(server_error()); @@ -921,18 +938,25 @@ async fn ref_namespace_policy() { let fx = fixture().await; seed_work(&fx.work); - let (ok, out) = push( - &fx.work, - &fx.url, - &fx.key_path, - &["main:refs/hidden/feature/main"], - ) - .await; - assert!(!ok, "push to refs/hidden/* must be rejected:\n{out}"); - assert!( - ref_names(&fx.server).is_empty(), - "forbidden-ref push must land nothing" - ); + for (spec, names_itself) in [ + ("main:refs/hidden/feature/main", "server-side fork staging"), + ("main:refs/atproto/commit", "atproto commit chain"), + ( + "main:refs/cob-checkpoints/sh.tangled.knot.member/beef", + "knot-written collaborative-object snapshots", + ), + ] { + 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!( + out.contains(names_itself), + "git must show the pusher which namespace refused {spec}:\n{out}" + ); + assert!( + ref_names(&fx.server).is_empty(), + "a push refused by policy must leave the ref store empty" + ); + } let (ok, out) = push( &fx.work, @@ -1606,6 +1630,165 @@ async fn a_collaborator_pushes_its_repo_but_is_denied_on_a_repo_it_doesnt_collab ); } +const INVITE_COLLAB: &str = "did:plc:olaren"; +const INVITE_BYSTANDER: &str = "did:plc:bailey"; + +enum Invited { + Unanswered, + Answered, + AnsweredPdsDown, +} + +struct Invitation { + scratch: TempDir, + layout: Layout, + url: String, + head: Oid, + collab_key: String, + bystander_key: String, + accounts: Accounts, +} + +impl Invitation { + fn tip(&self) -> Option { + main_tip(&self.layout, &RepoDid::new(REPO_DID).unwrap()) + } + async fn push(&self, key: &str) -> (bool, String) { + let work = self.scratch.path().join("work"); + push(&work, &self.url, key, &["main"]).await + } +} + +async fn an_invited_collaborator(state: Invited) -> Invitation { + let scratch = tempfile::tempdir().unwrap(); + let dir = scratch.path(); + let (_owner_key, owner_line) = keygen(dir, "owner"); + let (collab_key, collab_line) = keygen(dir, "collab"); + let (bystander_key, bystander_line) = keygen(dir, "bystander"); + let (layout, repo, index) = registered_index(&scratch, knot_index::KeyBudget::DEFAULT); + let head = Oid::from_hex(&seed_work(&dir.join("work"))).unwrap(); + let now = UnixSeconds::new(1); + let invited = CollaboratorsChange::Invite(knot_cobs::Invite(Grant { + subject: AccountDid::new(INVITE_COLLAB).unwrap(), + added_by: AccountDid::new(OWNER_DID).unwrap(), + created_at: now, + })); + let signer = K256Signer::generate(&SeededEntropy::new(2)); + CobStore::new(&layout.open(&repo).unwrap()) + .create(&CobHome::from(&repo), &invited, &signer, now) + .unwrap(); + index.refresh_collaborators(&repo).unwrap(); + + let held = |line: &str| { + let key = russh::keys::ssh_key::PublicKey::from_openssh(line).unwrap(); + vec![knot_types::OfferedKey::from_bytes(key.to_bytes().unwrap())] + }; + let owner = AccountDid::new(OWNER_DID).unwrap(); + let bystander = AccountDid::new(INVITE_BYSTANDER).unwrap(); + let keys = index.keys(); + keys.record(&owner, held(&owner_line), forever()); + keys.record(&bystander, held(&bystander_line), forever()); + keys.mark_ready(index.generation()); + + let records = Accounts::publishing(HashMap::from([ + (OWNER_DID.to_string(), vec![owner_line]), + (INVITE_COLLAB.to_string(), vec![collab_line]), + ])); + let down = HashSet::from([INVITE_COLLAB.to_string()]); + let accounts = match state { + Invited::Unanswered => records, + Invited::Answered => records.accepting(INVITE_COLLAB, REPO_DID), + Invited::AnsweredPdsDown => records.accepting(INVITE_COLLAB, REPO_DID).unreachable(down), + }; + let port = launch(dir, layout.clone(), Arc::clone(&index), accounts.clone()).await; + Invitation { + url: format!("ssh://git@127.0.0.1:{port}/{REPO_DID}"), + scratch, + layout, + head, + collab_key, + bystander_key, + accounts, + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn an_invited_collaborator_without_an_acceptance_is_refused_after_key_identification() { + let invitation = an_invited_collaborator(Invited::Unanswered).await; + let (ok, out) = invitation.push(&invitation.collab_key).await; + assert!( + !ok, + "git reported success, so the acl let this push through with no acceptance:\n{out}" + ); + assert!( + out.contains("you aren't authorized to push to this repository"), + "the connection has to reach the acl for the acceptance to be resolved at all, so the \ + refusal must be authorization rather than key identification:\n{out}" + ); + assert!( + invitation.tip().is_none(), + "main moved on a repo the pusher never signed for" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn an_invited_collaborator_pushes_once_the_knot_reads_their_acceptance() { + let invitation = an_invited_collaborator(Invited::Answered).await; + let (ok, out) = invitation.push(&invitation.collab_key).await; + assert!( + ok, + "the acceptance is in the pusher's own repository and the push is what makes the knot \ + look, so no second call from the client is needed:\n{out}" + ); + assert_eq!(invitation.tip(), Some(invitation.head)); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn wide_door_admits_whichever_key_invitee_offers_first() { + let invitation = an_invited_collaborator(Invited::Answered).await; + let (decoy, _line) = keygen(invitation.scratch.path(), "decoy"); + let offered = format!("{decoy} -i {}", invitation.collab_key); + let (ok, out) = invitation.push(&offered).await; + assert!( + !ok, + "the ssh username is git and the repository arrives in the exec command, so at key \ + identification the knot can't know which invite to weigh a key against:\n{out}" + ); + assert!( + out.contains("IdentitiesOnly"), + "accepting the offer ends the client's key search, so the refusal has to tell an \ + invitee holding two keys how to offer the one they signed the acceptance under:\n{out}" + ); + assert!( + invitation.tip().is_none(), + "main moved, so the decoy key pushed" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn invitees_first_push_survives_a_strangers_failed_one() { + let invitation = an_invited_collaborator(Invited::AnsweredPdsDown).await; + let (ok, out) = invitation.push(&invitation.bystander_key).await; + assert!( + !ok, + "the bystander's key pushed, and no grant on this roster names it" + ); + assert!( + out.contains("couldn't read the account records"), + "the key is on file, so the session costs no miss reservation and the fan-out runs, \ + which is what spends the invitee's pacer while their pds is down:\n{out}" + ); + + invitation.accounts.restore(INVITE_COLLAB); + let (ok, out) = invitation.push(&invitation.collab_key).await; + assert!( + ok, + "an invitee with no keys on file is read on every attempt, so a pacer somebody else \ + spent can't turn their first push into an unregistered key:\n{out}" + ); + assert_eq!(invitation.tip(), Some(invitation.head)); +} + fn ssh_bare(key_path: &str, port: u16) -> (bool, String) { ssh_bare_as(key_path, "git", port) } diff --git a/knot2/crates/knot-xrpc/Cargo.toml b/knot2/crates/knot-xrpc/Cargo.toml index 430b81ad0..e69baba78 100644 --- a/knot2/crates/knot-xrpc/Cargo.toml +++ b/knot2/crates/knot-xrpc/Cargo.toml @@ -14,6 +14,7 @@ knot-pack = { workspace = true } knot-resource = { workspace = true } knot-index = { workspace = true } knot-acl = { workspace = true } +knot-consent = { workspace = true } knot-atproto = { workspace = true } knot-cob = { workspace = true } knot-cobs = { workspace = true } diff --git a/knot2/crates/knot-xrpc/src/blocklist.rs b/knot2/crates/knot-xrpc/src/blocklist.rs index 0d65698cd..0b00830b2 100644 --- a/knot2/crates/knot-xrpc/src/blocklist.rs +++ b/knot2/crates/knot-xrpc/src/blocklist.rs @@ -10,11 +10,10 @@ use knot_acl::{KnotAcl, can_admin_knot}; use knot_cob::{CobHome, CobStore}; use knot_cobs::{BlocklistChange, BlocklistCob, Grant, Removal}; use knot_git::Repo; -use knot_index::Resolved; use knot_runtime::{Clock, HttpTransport}; use knot_types::AccountDid; -use crate::cob::grant_set_apply; +use crate::cob::{Authorized, roster_apply}; use crate::error::XrpcError; use crate::{XrpcState, decode, ok_empty, run_blocking}; @@ -42,9 +41,6 @@ pub(crate) async fn ban( if state.admins.contains(&subject) { return Err(XrpcError::forbidden("admin cannot be banned")); } - if matches!(state.index.is_blocked(&subject), Resolved::Ready(true)) { - return Ok(ok_empty()); - } let now = state.now(); let grant = Grant { @@ -60,13 +56,12 @@ pub(crate) async fn ban( run_blocking(move || { let _guard = cob_locks.meta(); let meta = Repo::open(&meta_path)?; - grant_set_apply::( + let _ = roster_apply::( &CobStore::new(&meta), &home, - BlocklistChange::Add(grant), + Authorized::by_decree(BlocklistChange::Add(grant)), &signer, now, - true, )?; index.refresh_blocklist().map_err(XrpcError::from) }) @@ -88,9 +83,6 @@ pub(crate) async fn unban( } let SubjectInput { subject } = decode(&body)?; - if matches!(state.index.is_blocked(&subject), Resolved::Ready(false)) { - return Ok(ok_empty()); - } let now = state.now(); let removal = Removal { subject }; @@ -102,13 +94,12 @@ pub(crate) async fn unban( run_blocking(move || { let _guard = cob_locks.meta(); let meta = Repo::open(&meta_path)?; - grant_set_apply::( + let _ = roster_apply::( &CobStore::new(&meta), &home, - BlocklistChange::Remove(removal), + Authorized::by_decree(BlocklistChange::Remove(removal)), &signer, now, - false, )?; index.refresh_blocklist().map_err(XrpcError::from) }) diff --git a/knot2/crates/knot-xrpc/src/cob.rs b/knot2/crates/knot-xrpc/src/cob.rs index aa9505fea..62d2b5367 100644 --- a/knot2/crates/knot-xrpc/src/cob.rs +++ b/knot2/crates/knot-xrpc/src/cob.rs @@ -1,44 +1,172 @@ -use knot_cob::{ChangePayload, Checkpoint, CobError, CobHome, CobStore, Evaluate}; -use knot_cobs::{GrantChange, Roster}; +use knot_cob::{ChangePayload, Checkpoint, CobHome, CobStore, Evaluate}; +use knot_cobs::{Accept, Announcement, Effect, GrantChange, Offer, Roster, Standing}; use knot_runtime::Signer; -use knot_types::UnixSeconds; +use knot_types::{AccountDid, UnixSeconds}; use crate::error::XrpcError; -pub(crate) fn grant_set_apply( +const NO_OFFER: &str = "no offer on the roster for this account"; + +const NOT_YOURS: &str = "only the account itself may sign an acceptance"; + +const NOT_A_GRANT: &str = "an operator must invite, and the account must accept"; + +const INVITE_PENDING: &str = + "an invite for this account is outstanding, and a grant would skip its acceptance"; + +pub(crate) struct Authorized(C); + +impl Authorized { + pub(crate) fn by_operator(change: C) -> Result { + match change.effect() { + Effect::Accept(_) => Err(XrpcError::forbidden(NOT_YOURS)), + Effect::Admit(Offer::Granted, _) => Err(XrpcError::forbidden(NOT_A_GRANT)), + Effect::Admit(Offer::Invited, _) | Effect::Revoke(_) => Ok(()), + }?; + Ok(Self(change)) + } + + pub(crate) fn by_decree(change: C) -> Self { + Self(change) + } + + pub(crate) fn accepted_by( + accepting: impl FnOnce(Accept) -> C, + actor: &AccountDid, + verified_at: UnixSeconds, + ) -> Self { + Self(accepting(Accept { + subject: actor.clone(), + verified_at, + })) + } + + pub(crate) fn subject(&self) -> &AccountDid { + self.0.subject() + } +} + +pub(crate) fn roster_apply( store: &CobStore, home: &CobHome, - change: E::Change, + Authorized(change): Authorized, signer: &dyn Signer, now: UnixSeconds, - create_if_absent: bool, -) -> Result +) -> Result, XrpcError> where E: Checkpoint + Evaluate, E::Change: ChangePayload + Clone + GrantChange, { + let announcement = change.effect().announcement(); match store.list::().map_err(XrpcError::from)?.as_slice() { - [] if create_if_absent => { - store + [] => match refusal(&Roster::empty(), &change) { + Some(error) => Err(error), + None if Roster::empty().would_change(&change) => store .create(home, &change, signer, now) - .map_err(XrpcError::from)?; - Ok(true) - } - [] => Ok(false), + .map(|_| announcement) + .map_err(XrpcError::from), + None => Ok(None), + }, [object] => store - .update_maybe_checkpointed::(home, *object, signer, now, |roster| { - let redundant = change.adds() == roster.contains(change.subject()); - Ok(if redundant { - None - } else { - Some(change.clone()) - }) + .update_maybe_checkpointed::(home, *object, signer, now, |roster| { + match refusal(roster, &change) { + Some(error) => Err(error), + None => Ok(roster.would_change(&change).then(|| change.clone())), + } }) - .map(|change_id| change_id.is_some()) - .map_err(XrpcError::from), + .map(|appended| appended.and(announcement)), many => Err(XrpcError::internal(format!( "{} collaborative objects of one type share namespace", many.len() ))), } } + +fn refusal(roster: &Roster, change: &impl GrantChange) -> Option { + match change.effect() { + Effect::Admit(Offer::Granted, grant) => { + matches!(roster.standing(&grant.subject), Some(Standing::Invited)) + .then(|| XrpcError::conflict(INVITE_PENDING)) + } + Effect::Accept(accept) => roster + .get(&accept.subject) + .is_none() + .then(|| XrpcError::forbidden(NO_OFFER)), + Effect::Admit(Offer::Invited, _) | Effect::Revoke(_) => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use http::StatusCode; + use knot_cobs::{Grant, Invite, MembersChange, Removal}; + + fn did(suffix: &str) -> AccountDid { + AccountDid::new(format!("did:plc:{suffix}")).unwrap() + } + + fn offer() -> Grant { + Grant { + subject: did("nel"), + added_by: did("olaren"), + created_at: UnixSeconds::new(1), + } + } + + fn accept() -> Accept { + Accept { + subject: did("nel"), + verified_at: UnixSeconds::new(2), + } + } + + fn revoke() -> Removal { + Removal { + subject: did("nel"), + } + } + + #[test] + fn an_operator_invites_or_removes_and_the_account_signs_its_own_acceptance() { + let nel = did("nel"); + assert!( + Authorized::by_operator(MembersChange::Accept(accept())).is_err(), + "by_operator let an acceptance through, so an operator could sign for the account" + ); + assert!( + Authorized::by_operator(MembersChange::Add(offer())).is_err(), + "by_operator let a bare grant through, so an operator could add an account that \ + never signed" + ); + assert!( + Authorized::by_operator(MembersChange::Invite(Invite(offer()))).is_ok() + && Authorized::by_operator(MembersChange::Remove(revoke())).is_ok(), + "an invite and a removal need no signature from the account" + ); + assert_eq!( + Authorized::accepted_by(MembersChange::Accept, &nel, UnixSeconds::new(3)).subject(), + &nel, + "accepted_by takes its subject from the signing actor, never from the request body" + ); + } + + #[test] + fn a_grant_conflicts_with_an_outstanding_invite_and_the_other_three_effects_pass() { + let invited = Roster::empty().apply_change(&Invite(offer())); + assert!( + refusal(&invited, &offer()) + .is_some_and(|error| error.status() == StatusCode::CONFLICT), + "a Grant folds onto an outstanding invitation as a membership the account never \ + signed for, so the locked path refuses it rather than swallowing it as a no-op, and \ + by_operator plus knot-migrate's grandfathering leave this arm no caller at all" + ); + assert!( + refusal(&invited, &Invite(offer())).is_none() + && refusal(&invited, &accept()).is_none() + && refusal(&invited, &revoke()).is_none(), + "the offer again, the answer to it and its revocation are the invitation's other \ + three effects, and the append-time dedup turns each into its own idempotent answer" + ); + } +} diff --git a/knot2/crates/knot-xrpc/src/collaborators.rs b/knot2/crates/knot-xrpc/src/collaborators.rs index 6739de18b..e35314ead 100644 --- a/knot2/crates/knot-xrpc/src/collaborators.rs +++ b/knot2/crates/knot-xrpc/src/collaborators.rs @@ -8,18 +8,20 @@ use serde::Deserialize; use knot_acl::{KnotAcl, can_manage_collaborators}; use knot_cob::{CobHome, CobStore}; -use knot_cobs::{CollaboratorsChange, CollaboratorsCob, Grant, Removal}; +use knot_cobs::{Announcement, CollaboratorsChange, CollaboratorsCob, Grant, Invite, Removal}; use knot_events::RepoCollaboratorUpdate; use knot_index::Resolved; use knot_runtime::{Clock, HttpTransport}; -use knot_types::{AccountDid, RepoDid}; +use knot_types::{AccountDid, RepoDid, UnixSeconds}; -use crate::cob::grant_set_apply; +use crate::cob::{Authorized, roster_apply}; +use crate::consent::{AcceptanceInput, AcceptanceRef, Acceptances, Consent, consent_gate}; use crate::error::XrpcError; use crate::{XrpcState, decode, ok_empty, run_blocking}; pub(crate) const ADD_ROUTE: &str = "/xrpc/sh.tangled.repo.addCollaborator"; pub(crate) const REMOVE_ROUTE: &str = "/xrpc/sh.tangled.repo.removeCollaborator"; +pub(crate) const ACCEPT_ROUTE: &str = "/xrpc/sh.tangled.repo.acceptCollaboration"; #[derive(Deserialize)] struct CollaboratorInput { @@ -27,22 +29,25 @@ struct CollaboratorInput { subject: AccountDid, } +fn registered_owner( + state: &XrpcState, + repo: &RepoDid, +) -> Result { + match state.index.owner_of(repo) { + Resolved::Ready(Some(owner)) => Ok(owner), + Resolved::Ready(None) => Err(XrpcError::not_found( + "repository isn't registered on this knot", + )), + Resolved::Warming => Err(crate::reads::warming()), + } +} + fn require_owner( state: &XrpcState, actor: &AccountDid, repo: &RepoDid, ) -> Result { - let owner = match state.index.owner_of(repo) { - Resolved::Ready(Some(owner)) => owner, - Resolved::Ready(None) => { - return Err(XrpcError::not_found( - "repository isn't registered on this knot", - )); - } - Resolved::Warming => { - return Err(XrpcError::warming("registry projection is still warming")); - } - }; + let owner = registered_owner(state, repo)?; let acl = KnotAcl::new(&state.admins, state.admission, &state.index); if can_manage_collaborators(&acl, actor, repo).is_allowed() { Ok(owner) @@ -62,28 +67,31 @@ pub(crate) async fn add_collaborator( let actor = state.authenticate(&headers, &method).await?; let CollaboratorInput { repo, subject } = decode(&body)?; let owner = require_owner(&state, &actor, &repo)?; - - crate::fold_collaborators(&state, &repo).await; - if owner.is(&subject) - || matches!( - state.index.is_collaborator(&repo, &subject), - Resolved::Ready(true) - ) - { - return Ok(ok_empty()); - } - let now = state.now(); - let event_subject = subject.clone(); + commit_collaborators( + &state, + repo, + owner, + Authorized::by_operator(CollaboratorsChange::Invite(Invite(Grant { + subject, + added_by: actor, + created_at: now, + })))?, + now, + ) + .await +} + +async fn commit_collaborators( + state: &Arc>, + repo: RepoDid, + owner: knot_types::OwnerDid, + change: Authorized, + now: UnixSeconds, +) -> Result { + let subject = change.subject().clone(); let event_repo = repo.clone(); - let grant = Grant { - subject, - added_by: actor, - created_at: now, - }; - let signer = state.secrets.signer(&state.knot_did).map_err(|error| { - XrpcError::internal(format!("knot signing key is unavailable: {error}")) - })?; + let signer = state.secrets.signer(&state.knot_did)?; let layout = state.layout.clone(); let index = Arc::clone(&state.index); let cob_locks = Arc::clone(&state.cob_locks); @@ -91,19 +99,24 @@ pub(crate) async fn add_collaborator( run_blocking(move || { let _guard = cob_locks.repo(&repo); owner_unmoved(&index, &repo, &owner)?; + if owner.is(&subject) { + return Ok(()); + } let git = layout.open(&repo)?; - let changed = grant_set_apply::( + let announced = roster_apply::( &CobStore::new(&git), &CobHome::from(&repo), - CollaboratorsChange::Add(grant), + change, &signer, now, - true, )?; - index.refresh_collaborators(&repo)?; - if changed { - events.publish(&RepoCollaboratorUpdate::added(event_subject, event_repo)); + if let Some(update) = announced.map(|announcement| match announcement { + Announcement::Effective => RepoCollaboratorUpdate::added(subject, event_repo), + Announcement::Cleared => RepoCollaboratorUpdate::removed(subject, event_repo), + }) { + events.publish(&update); } + index.refresh_collaborators(&repo)?; Ok(()) }) .await?; @@ -134,44 +147,52 @@ pub(crate) async fn remove_collaborator( let CollaboratorInput { repo, subject } = decode(&body)?; let owner = require_owner(&state, &actor, &repo)?; - crate::fold_collaborators(&state, &repo).await; + commit_collaborators( + &state, + repo, + owner, + Authorized::by_operator(CollaboratorsChange::Remove(Removal { subject }))?, + state.now(), + ) + .await +} + +pub(crate) async fn accept_collaboration( + State(state): State>>, + headers: HeaderMap, + method: crate::Method, + body: Bytes, +) -> Result { + let AcceptanceInput { acceptance: uri } = decode(&body)?; + let reference = AcceptanceRef::parse(&uri, Acceptances::OfCollaboration)?; + let repo = RepoDid::from(reference.subject()); + + let actor = state + .authenticate_for_repo(&headers, &method, &repo) + .await?; + let acceptance = reference.written_by(&actor)?; + let owner = registered_owner(&state, &repo)?; + + crate::fold_collaborators(&state, &repo).await?; + let standing = state.index.collaborator_standing(&repo, &actor); if matches!( - state.index.is_collaborator(&repo, &subject), - Resolved::Ready(false) + consent_gate( + standing, + "no collaboration offer for you on this repository" + )?, + Consent::Recorded ) { return Ok(ok_empty()); } + acceptance.require_published(&state).await?; - let now = state.now(); - let event_subject = subject.clone(); - let event_repo = repo.clone(); - let removal = Removal { subject }; - let signer = state.secrets.signer(&state.knot_did).map_err(|error| { - XrpcError::internal(format!("knot signing key is unavailable: {error}")) - })?; - let layout = state.layout.clone(); - let index = Arc::clone(&state.index); - let cob_locks = Arc::clone(&state.cob_locks); - let events = Arc::clone(&state.events); - run_blocking(move || { - let _guard = cob_locks.repo(&repo); - owner_unmoved(&index, &repo, &owner)?; - let git = layout.open(&repo)?; - let changed = grant_set_apply::( - &CobStore::new(&git), - &CobHome::from(&repo), - CollaboratorsChange::Remove(removal), - &signer, - now, - false, - )?; - index.refresh_collaborators(&repo)?; - if changed { - events.publish(&RepoCollaboratorUpdate::removed(event_subject, event_repo)); - } - Ok(()) - }) - .await?; - - Ok(ok_empty()) + let verified_at = state.now(); + commit_collaborators( + &state, + repo, + owner, + Authorized::accepted_by(CollaboratorsChange::Accept, &actor, verified_at), + verified_at, + ) + .await } diff --git a/knot2/crates/knot-xrpc/src/forks.rs b/knot2/crates/knot-xrpc/src/forks.rs index 73eb1ef75..88c3ee992 100644 --- a/knot2/crates/knot-xrpc/src/forks.rs +++ b/knot2/crates/knot-xrpc/src/forks.rs @@ -9,7 +9,7 @@ use serde::{Deserialize, Serialize}; use url::Url; use knot_events::Reservation; -use knot_git::{Filter, GitError, Haves, RefUpdate, Repo, Staging, Wants}; +use knot_git::{Filter, GitError, Haves, RefClass, RefUpdate, Repo, Staging, Wants}; use knot_pack::{FetchError, HaveOids, PackLimits, UpstreamRefs, WantOids}; use knot_postreceive::{Actor, Ci}; use knot_runtime::{Clock, HttpTransport}; @@ -368,6 +368,10 @@ fn load_fork_state(repo: &Repo) -> Result { let haves = repo .references()? .into_iter() + .filter(|record| match RefClass::of(&record.name) { + RefClass::Public | RefClass::Hidden => true, + RefClass::Cob | RefClass::Checkpoint | RefClass::Atproto => false, + }) .map(|record| record.target) .collect(); Ok(ForkState { diff --git a/knot2/crates/knot-xrpc/src/legacy_admin.rs b/knot2/crates/knot-xrpc/src/legacy_admin.rs index d36ad9dd6..92fa3d019 100644 --- a/knot2/crates/knot-xrpc/src/legacy_admin.rs +++ b/knot2/crates/knot-xrpc/src/legacy_admin.rs @@ -14,7 +14,7 @@ use knot_cobs::Grant; use knot_runtime::{Clock, HttpTransport}; use crate::error::XrpcError; -use crate::members::{SubjectInput, grant_membership}; +use crate::members::{SubjectInput, offer_membership}; use crate::{XrpcState, basic_credentials, decode, enforce_pre_auth_limit}; pub const ADD_MEMBER_ROUTE: &str = "/admin/addMember"; @@ -80,9 +80,9 @@ async fn add_member( tracing::warn!( route = ADD_MEMBER_ROUTE, %subject, - "legacy admin route authorized a member grant" + "legacy shared-secret admin path is still in use" ); - grant_membership( + offer_membership( &admin.state, Grant { subject, diff --git a/knot2/crates/knot-xrpc/src/lib.rs b/knot2/crates/knot-xrpc/src/lib.rs index 36d2332d8..7f2f61051 100644 --- a/knot2/crates/knot-xrpc/src/lib.rs +++ b/knot2/crates/knot-xrpc/src/lib.rs @@ -3,6 +3,7 @@ mod body; mod branches; mod cob; mod collaborators; +mod consent; mod error; mod events; mod forks; @@ -213,6 +214,19 @@ impl XrpcState { .map_err(map_verify_error) } + pub(crate) async fn authenticate_for_repo( + &self, + headers: &HeaderMap, + method: &Method, + repo: &RepoDid, + ) -> Result { + let token = bearer(headers)?; + self.atproto + .verify_service_jwt_for_repo(&token, method.nsid(), repo) + .await + .map_err(map_verify_error) + } + pub(crate) async fn authenticate_push( &self, headers: &HeaderMap, @@ -277,6 +291,10 @@ pub fn router(state: Arc>) -> Router .merge(merge_routes) .route(members::ADD_ROUTE, post(members::add_member::)) .route(members::REMOVE_ROUTE, post(members::remove_member::)) + .route( + members::ACCEPT_ROUTE, + post(members::accept_membership::), + ) .route(blocklist::BAN_ROUTE, post(blocklist::ban::)) .route(blocklist::UNBAN_ROUTE, post(blocklist::unban::)) .route( @@ -287,6 +305,10 @@ pub fn router(state: Arc>) -> Router collaborators::REMOVE_ROUTE, post(collaborators::remove_collaborator::), ) + .route( + collaborators::ACCEPT_ROUTE, + post(collaborators::accept_collaboration::), + ) .route(repos::CREATE_ROUTE, post(repos::create_repo::)) .route(repos::DELETE_ROUTE, post(repos::delete_repo::)) .route(repos::RENAME_ROUTE, post(repos::rename_repo::)) @@ -328,6 +350,14 @@ pub fn router(state: Arc>) -> Router lists::LIST_COLLABORATORS_ROUTE, get(lists::list_collaborators::), ) + .route( + lists::LIST_MEMBER_INVITES_ROUTE, + get(lists::list_member_invites::), + ) + .route( + lists::LIST_COLLABORATOR_INVITES_ROUTE, + get(lists::list_collaborator_invites::), + ) .route(service::VERSION_ROUTE, get(service::version)) .route(service::OWNER_ROUTE, get(service::owner::)) .layer(DefaultBodyLimit::max(state.byte_limits.body.get())) @@ -475,10 +505,19 @@ pub(crate) fn current_owner( pub(crate) async fn fold_collaborators( state: &XrpcState, repo: &RepoDid, -) { +) -> Result<(), XrpcError> { let index = Arc::clone(&state.index); let target = repo.clone(); - let _ = run_blocking(move || Ok(index.ensure_collaborators(&target))).await; + run_blocking(move || index.ensure_collaborators(&target).map_err(XrpcError::from)).await +} + +pub(crate) async fn folded_member_standing( + state: &XrpcState, + actor: &AccountDid, +) -> Result>, XrpcError> { + let index = Arc::clone(&state.index); + run_blocking(move || index.refresh_members().map_err(XrpcError::from)).await?; + Ok(state.index.member_standing(actor)) } pub(crate) async fn authorize_push( @@ -487,15 +526,47 @@ pub(crate) async fn authorize_push( repo: &RepoDid, denied: &str, ) -> Result<(), XrpcError> { - fold_collaborators(state, repo).await; + fold_collaborators(state, repo).await?; + let consent = state + .index + .consents() + .of_collaboration(&pds(state), repo, actor, state.now()) + .await; let acl = knot_acl::KnotAcl::new(&state.admins, state.admission, &state.index); - if knot_acl::can_push(&acl, actor, repo).is_allowed() { + if knot_acl::can_push(&acl, &consent).is_allowed() { Ok(()) } else { Err(XrpcError::forbidden(denied)) } } +pub(crate) async fn authorize_create( + state: &XrpcState, + actor: &AccountDid, + denied: &str, +) -> Result<(), XrpcError> { + let acl = knot_acl::KnotAcl::new(&state.admins, state.admission, &state.index); + let consent = match knot_acl::needs_membership(&acl, actor) { + true => { + state + .index + .consents() + .of_membership(&pds(state), actor, state.now()) + .await + } + false => knot_index::MemberConsent::unread(actor), + }; + if knot_acl::can_create_repo(&acl, &consent).is_allowed() { + Ok(()) + } else { + Err(XrpcError::forbidden(denied)) + } +} + +fn pds(state: &XrpcState) -> knot_consent::Pds<'_, H, C> { + knot_consent::Pds::reading(&state.atproto, &state.knot_did) +} + pub(crate) async fn authenticate_and_authorize_push( state: &XrpcState, socket: SocketPeer, diff --git a/knot2/crates/knot-xrpc/src/lists.rs b/knot2/crates/knot-xrpc/src/lists.rs index fe1e40947..511b7577a 100644 --- a/knot2/crates/knot-xrpc/src/lists.rs +++ b/knot2/crates/knot-xrpc/src/lists.rs @@ -7,10 +7,9 @@ use http::StatusCode; use http::request::Parts; use serde::{Deserialize, Serialize}; -use knot_cobs::Grant; -use knot_index::Resolved; +use knot_index::{EffectiveSince, Entry, Resolved, Standing}; use knot_runtime::{Clock, HttpTransport}; -use knot_types::{AccountDid, RepoDid}; +use knot_types::{AccountDid, KnotId, RepoDid, UnixSeconds}; use crate::XrpcState; use crate::error::XrpcError; @@ -19,6 +18,9 @@ use crate::wire::rfc3339; pub(crate) const LIST_MEMBERS_ROUTE: &str = "/xrpc/sh.tangled.knot.listMembers"; pub(crate) const LIST_COLLABORATORS_ROUTE: &str = "/xrpc/sh.tangled.repo.listCollaborators"; +pub(crate) const LIST_MEMBER_INVITES_ROUTE: &str = "/xrpc/sh.tangled.knot.listMemberInvites"; +pub(crate) const LIST_COLLABORATOR_INVITES_ROUTE: &str = + "/xrpc/sh.tangled.repo.listCollaboratorInvites"; const DEFAULT_LIMIT: usize = 50; const MAX_LIMIT: usize = 1000; @@ -63,22 +65,23 @@ fn subject_param(parts: &Parts) -> Result { .ok_or_else(|| XrpcError::invalid_request("missing subject parameter")) } +fn unnameable(name: &'static str, message: &'static str) -> XrpcError { + XrpcError::named(StatusCode::BAD_REQUEST, name, message) +} + pub(crate) struct MemberSubject; -impl FromRequestParts for MemberSubject { +impl FromRequestParts>> for MemberSubject { type Rejection = XrpcError; - async fn from_request_parts(parts: &mut Parts, _state: &S) -> Result { - let subject = subject_param(parts)?; - AccountDid::new(subject) - .map(|_| MemberSubject) - .map_err(|_| { - XrpcError::named( - StatusCode::BAD_REQUEST, - "InvalidSubject", - "subject must be an account DID", - ) - }) + async fn from_request_parts( + parts: &mut Parts, + state: &Arc>, + ) -> Result { + KnotId::new(subject_param(parts)?) + .is_ok_and(|knot| knot == state.knot_did) + .then_some(MemberSubject) + .ok_or_else(|| unnameable("InvalidSubject", "subject must be this knot's DID")) } } @@ -88,14 +91,9 @@ impl FromRequestParts for CollaboratorRepo { type Rejection = XrpcError; async fn from_request_parts(parts: &mut Parts, _state: &S) -> Result { - let subject = subject_param(parts)?; - RepoDid::new(subject).map(CollaboratorRepo).map_err(|_| { - XrpcError::named( - StatusCode::BAD_REQUEST, - "InvalidRepo", - "subject must be a repo DID", - ) - }) + RepoDid::new(subject_param(parts)?) + .map(CollaboratorRepo) + .map_err(|_| unnameable("InvalidRepo", "subject must be a repo DID")) } } @@ -106,6 +104,10 @@ struct ItemWire { added_by: AccountDid, #[serde(rename = "createdAt")] created_at: String, + #[serde(rename = "verifiedAt", skip_serializing_if = "Option::is_none")] + verified_at: Option, + #[serde(rename = "effectiveSince", skip_serializing_if = "Option::is_none")] + effective_since: Option, } #[derive(Serialize)] @@ -115,24 +117,39 @@ struct PageWire { cursor: Option, } -fn respond(mut entries: Vec, window: Window) -> Response { - entries.sort_by(|a, b| { - let by_time = a.created_at.cmp(&b.created_at); - let by_time = match window.descending { - true => by_time.reverse(), - false => by_time, - }; - by_time.then_with(|| a.subject.cmp(&b.subject)) +fn invited(standing: Standing) -> bool { + matches!(standing, Standing::Invited) +} + +fn page(entries: Vec<(AccountDid, Entry)>, window: Window) -> Response { + let mut rows: Vec<(UnixSeconds, AccountDid, Entry)> = entries + .into_iter() + .map(|(subject, entry)| { + let at = entry.effective_since().map_or(entry.created_at, EffectiveSince::seconds); + (at, subject, entry) + }) + .collect(); + rows.sort_by(|(left, earlier, _), (right, later, _)| { + match window.descending { + true => right.cmp(left), + false => left.cmp(right), + } + .then_with(|| earlier.cmp(later)) }); - let total = entries.len(); - let items = entries + let total = rows.len(); + let items = rows .into_iter() .skip(window.offset.get()) .take(window.limit.get()) - .map(|grant| ItemWire { - subject: grant.subject, - added_by: grant.added_by, - created_at: rfc3339(grant.created_at.get(), 0), + .map(|(_, subject, entry)| { + let since = entry.effective_since(); + ItemWire { + subject, + added_by: entry.added_by, + created_at: rfc3339(entry.created_at.get(), 0), + verified_at: entry.verified_at.map(|at| rfc3339(at.get(), 0)), + effective_since: since.map(|since| rfc3339(since.seconds().get(), 0)), + } }) .collect(); let cursor = next_cursor(window.offset, window.limit, Total::new(total)); @@ -153,21 +170,44 @@ pub(crate) async fn list_members( ValidatedQuery(paging): ValidatedQuery, ) -> Result { let window = paging.window(); - match state.index.member_entries() { + match state.index.member_entries_where(Standing::is_effective) { Resolved::Warming => Err(members_warming()), - Resolved::Ready(entries) => Ok(respond(entries, window)), + Resolved::Ready(entries) => Ok(page(entries, window)), } } -pub(crate) async fn list_collaborators( +pub(crate) async fn list_member_invites( State(state): State>>, - CollaboratorRepo(repo): CollaboratorRepo, + _subject: MemberSubject, ValidatedQuery(paging): ValidatedQuery, ) -> Result { let window = paging.window(); - match state.index.owner_of(&repo) { + match state.index.member_entries_where(invited) { + Resolved::Warming => Err(members_warming()), + Resolved::Ready(entries) => entries + .into_iter() + .map(|(subject, entry)| match state.index.is_blocked(&subject) { + Resolved::Warming => { + Err(XrpcError::warming("blocklist projection is still warming, retry shortly")) + } + Resolved::Ready(true) => Ok(None), + Resolved::Ready(false) => Ok(Some((subject, entry))), + }) + .filter_map(Result::transpose) + .collect::, _>>() + .map(|open| page(open, window)), + } +} + +async fn collaborator_page( + state: &XrpcState, + repo: &RepoDid, + window: Window, + keep: fn(Standing) -> bool, +) -> Result { + match state.index.owner_of(repo) { Resolved::Warming => return Err(crate::reads::warming()), - Resolved::Ready(None) => return Ok(respond(Vec::new(), window)), + Resolved::Ready(None) => return Ok(page(Vec::new(), window)), Resolved::Ready(Some(_)) => {} } let index = Arc::clone(&state.index); @@ -176,11 +216,29 @@ pub(crate) async fn list_collaborators( index .ensure_collaborators(&target) .map_err(|_| collaborators_warming())?; - match index.collaborator_entries(&target) { + match index.collaborator_entries_where(&target, keep) { Resolved::Warming => Err(collaborators_warming()), Resolved::Ready(entries) => Ok(entries), } }) .await?; - Ok(respond(entries, window)) + Ok(page(entries, window)) +} + +pub(crate) async fn list_collaborators( + State(state): State>>, + CollaboratorRepo(repo): CollaboratorRepo, + ValidatedQuery(paging): ValidatedQuery, +) -> Result { + let window = paging.window(); + collaborator_page(&state, &repo, window, Standing::is_effective).await +} + +pub(crate) async fn list_collaborator_invites( + State(state): State>>, + CollaboratorRepo(repo): CollaboratorRepo, + ValidatedQuery(paging): ValidatedQuery, +) -> Result { + let window = paging.window(); + collaborator_page(&state, &repo, window, invited).await } diff --git a/knot2/crates/knot-xrpc/src/members.rs b/knot2/crates/knot-xrpc/src/members.rs index 01f9c508a..a81c213fb 100644 --- a/knot2/crates/knot-xrpc/src/members.rs +++ b/knot2/crates/knot-xrpc/src/members.rs @@ -8,19 +8,20 @@ use serde::Deserialize; use knot_acl::{KnotAcl, can_admin_knot}; use knot_cob::{CobHome, CobStore}; -use knot_cobs::{Grant, MembersChange, MembersCob, Removal}; +use knot_cobs::{Announcement, Grant, Invite, MembersChange, MembersCob, Removal}; use knot_events::KnotMemberUpdate; use knot_git::Repo; -use knot_index::Resolved; use knot_runtime::{Clock, HttpTransport}; -use knot_types::AccountDid; +use knot_types::{AccountDid, UnixSeconds}; -use crate::cob::grant_set_apply; +use crate::cob::{Authorized, roster_apply}; +use crate::consent::{AcceptanceInput, AcceptanceRef, Acceptances, Consent, consent_gate}; use crate::error::XrpcError; use crate::{XrpcState, decode, ok_empty, run_blocking}; pub(crate) const ADD_ROUTE: &str = "/xrpc/sh.tangled.knot.addMember"; pub(crate) const REMOVE_ROUTE: &str = "/xrpc/sh.tangled.knot.removeMember"; +pub(crate) const ACCEPT_ROUTE: &str = "/xrpc/sh.tangled.knot.acceptMembership"; #[derive(Deserialize)] pub(crate) struct SubjectInput { @@ -40,7 +41,7 @@ pub(crate) async fn add_member( } let SubjectInput { subject } = decode(&body)?; - grant_membership( + offer_membership( &state, Grant { subject, @@ -51,18 +52,29 @@ pub(crate) async fn add_member( .await } -pub(crate) async fn grant_membership( +pub(crate) async fn offer_membership( state: &Arc>, - grant: Grant, + offer: Grant, ) -> Result { - if state.admins.contains(&grant.subject) - || matches!(state.index.is_member(&grant.subject), Resolved::Ready(true)) - { + if state.admins.contains(&offer.subject) { return Ok(ok_empty()); } - let now = grant.created_at; - let event_subject = grant.subject.clone(); + let now = offer.created_at; + commit_members( + state, + Authorized::by_operator(MembersChange::Invite(Invite(offer)))?, + now, + ) + .await +} + +async fn commit_members( + state: &Arc>, + change: Authorized, + now: UnixSeconds, +) -> Result { + let subject = change.subject().clone(); let signer = state.secrets.signer(&state.knot_did)?; let meta_path = state.meta_path.clone(); let index = Arc::clone(&state.index); @@ -72,18 +84,15 @@ pub(crate) async fn grant_membership( run_blocking(move || { let _guard = cob_locks.meta(); let meta = Repo::open(&meta_path)?; - let changed = grant_set_apply::( - &CobStore::new(&meta), - &home, - MembersChange::Add(grant), - &signer, - now, - true, - )?; - index.refresh_members()?; - if changed { - events.publish(&KnotMemberUpdate::added(event_subject)); + let announced = + roster_apply::(&CobStore::new(&meta), &home, change, &signer, now)?; + if let Some(update) = announced.map(|announcement| match announcement { + Announcement::Effective => KnotMemberUpdate::added(subject), + Announcement::Cleared => KnotMemberUpdate::removed(subject), + }) { + events.publish(&update); } + index.refresh_members()?; Ok(()) }) .await?; @@ -104,37 +113,39 @@ pub(crate) async fn remove_member( } let SubjectInput { subject } = decode(&body)?; - if matches!(state.index.is_member(&subject), Resolved::Ready(false)) { + commit_members( + &state, + Authorized::by_operator(MembersChange::Remove(Removal { subject }))?, + state.now(), + ) + .await +} + +pub(crate) async fn accept_membership( + State(state): State>>, + headers: HeaderMap, + method: crate::Method, + body: Bytes, +) -> Result { + let actor = state.authenticate(&headers, &method).await?; + let AcceptanceInput { acceptance: uri } = decode(&body)?; + let acceptance = AcceptanceRef::parse(&uri, Acceptances::OfMembership)?.written_by(&actor)?; + acceptance.names_knot(&state.knot_did)?; + + let standing = crate::folded_member_standing(&state, &actor).await?; + if matches!( + consent_gate(standing, "no membership offer for you on this knot")?, + Consent::Recorded + ) { return Ok(ok_empty()); } + acceptance.require_published(&state).await?; - let now = state.now(); - let event_subject = subject.clone(); - let removal = Removal { subject }; - let signer = state.secrets.signer(&state.knot_did)?; - let meta_path = state.meta_path.clone(); - let index = Arc::clone(&state.index); - let cob_locks = Arc::clone(&state.cob_locks); - let events = Arc::clone(&state.events); - let home = CobHome::from(&state.knot_did); - run_blocking(move || { - let _guard = cob_locks.meta(); - let meta = Repo::open(&meta_path)?; - let changed = grant_set_apply::( - &CobStore::new(&meta), - &home, - MembersChange::Remove(removal), - &signer, - now, - false, - )?; - index.refresh_members()?; - if changed { - events.publish(&KnotMemberUpdate::removed(event_subject)); - } - Ok(()) - }) - .await?; - - Ok(ok_empty()) + let verified_at = state.now(); + commit_members( + &state, + Authorized::accepted_by(MembersChange::Accept, &actor, verified_at), + verified_at, + ) + .await } diff --git a/knot2/crates/knot-xrpc/src/reads.rs b/knot2/crates/knot-xrpc/src/reads.rs index 78a24d589..a530ef022 100644 --- a/knot2/crates/knot-xrpc/src/reads.rs +++ b/knot2/crates/knot-xrpc/src/reads.rs @@ -13,8 +13,8 @@ use tower_http::services::ServeFile; use knot_cobs::RepoRef; use knot_git::{ - ArchiveFormat, BinaryBudget, Commit, CommitRange, EntryKind, Layout, LogLimit, LogSkip, Repo, - SizedEntry, is_public_ref, screens_reserved, + ArchiveFormat, BinaryBudget, Commit, CommitRange, EntryKind, Layout, LogLimit, LogSkip, + RefClass, Repo, SizedEntry, }; use knot_index::{Coverage, Resolved}; use knot_runtime::{Clock, HttpTransport}; @@ -94,7 +94,10 @@ fn readme_serving_limit(response_limit: usize) -> u64 { } fn names_reserved(refspec: &str) -> bool { - screens_reserved(refspec) || screens_reserved(&format!("refs/{refspec}")) + match RefClass::of_revspec(refspec) { + RefClass::Cob | RefClass::Checkpoint | RefClass::Atproto => true, + RefClass::Public | RefClass::Hidden => false, + } } pub(crate) fn warming() -> XrpcError { @@ -1618,7 +1621,7 @@ pub(crate) async fn git_list_refs( let mut refs: Vec<_> = repo .references()? .into_iter() - .filter(|record| is_public_ref(&record.name)) + .filter(|record| RefClass::of(&record.name).is_public()) .collect(); refs.sort_by(|a, b| a.name.as_str().cmp(b.name.as_str())); let total = refs.len(); diff --git a/knot2/crates/knot-xrpc/src/repos.rs b/knot2/crates/knot-xrpc/src/repos.rs index 7cc81ed40..b297a7d06 100644 --- a/knot2/crates/knot-xrpc/src/repos.rs +++ b/knot2/crates/knot-xrpc/src/repos.rs @@ -7,7 +7,7 @@ use axum::response::{IntoResponse, Response}; use http::HeaderMap; use serde::{Deserialize, Serialize}; -use knot_acl::{KnotAcl, can_admin_knot, can_create_repo, can_delete_repo}; +use knot_acl::{KnotAcl, can_admin_knot, can_delete_repo}; use knot_atproto::{PreparedRepoDid, RecordPresence}; use knot_cob::{CobHome, CobStore}; use knot_cobs::{ @@ -85,13 +85,12 @@ pub(crate) async fn reserve_key( body: Bytes, ) -> Result { let actor = state.authenticate(&headers, &method).await?; - let acl = KnotAcl::new(&state.admins, state.admission, &state.index); - if !can_create_repo(&acl, &actor).is_allowed() { - return Err(XrpcError::forbidden( - "only knot admin or member may reserve a repository key", - )); - } - + crate::authorize_create( + &state, + &actor, + "only knot admin or member may reserve a repository key", + ) + .await?; let ReserveInput { repo_did } = decode(&body)?; if !repo_did.as_str().starts_with("did:web:") { return Err(XrpcError::invalid_request( @@ -156,12 +155,12 @@ pub(crate) async fn create_repo( body: Bytes, ) -> Result { let actor = state.authenticate(&headers, &method).await?; - let acl = KnotAcl::new(&state.admins, state.admission, &state.index); - if !can_create_repo(&acl, &actor).is_allowed() { - return Err(XrpcError::forbidden( - "only knot admin or member may create repositories", - )); - } + crate::authorize_create( + &state, + &actor, + "only knot admin or member may create repositories", + ) + .await?; let input: CreateInput = decode(&body)?; let head = input.default_branch.map(|branch| branch.head_ref()); diff --git a/knot2/crates/knot-xrpc/src/tests.rs b/knot2/crates/knot-xrpc/src/tests.rs index dd367f48d..8fe3c4d4b 100644 --- a/knot2/crates/knot-xrpc/src/tests.rs +++ b/knot2/crates/knot-xrpc/src/tests.rs @@ -1,6 +1,6 @@ -use std::collections::{BTreeSet, HashMap, HashSet}; +use std::collections::{BTreeSet, HashMap}; use std::path::PathBuf; -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use axum::body::Bytes; @@ -14,8 +14,9 @@ use serde_json::json; use tempfile::TempDir; use knot_atproto::Atproto; -use knot_git::Layout; -use knot_index::{Index, Resolved}; +use knot_cobs::{CollaboratorsChange, CollaboratorsCob, Grant, Invite, MembersChange, MembersCob}; +use knot_git::{Layout, Repo}; +use knot_index::{Index, Resolved, Standing}; use knot_runtime::{ FakeHttp, HttpRequest, HttpResponse, HttpTransport, K256Signer, ManualClock, NetworkError, OsEntropy, SeededEntropy, Signer, UnixMicros, @@ -23,10 +24,12 @@ use knot_runtime::{ use knot_secrets::{MasterKey, SealedStore}; use knot_types::{ AccountDid, AdmissionPolicy, AuthorName, Email, KnotHostname, KnotId, OriginUrl, OwnerDid, - RepoDid, RepoName, RepoRkey, + RepoDid, RepoName, RepoRkey, UnixSeconds, }; use crate::XrpcState; +use crate::collaborators::{add_collaborator, remove_collaborator}; +use crate::members::{add_member, remove_member}; const KNOT_HOST: &str = "knot.nel.pet"; const ADMIN_HOST: &str = "admin.nel.pet"; @@ -35,6 +38,10 @@ const STRANGER_HOST: &str = "stranger.nel.pet"; const ADD_MEMBER: &str = "sh.tangled.knot.addMember"; const REMOVE_MEMBER: &str = "sh.tangled.knot.removeMember"; +const ACCEPT_MEMBERSHIP: &str = "sh.tangled.knot.acceptMembership"; +const MEMBER_ACCEPTANCE: &str = "sh.tangled.knot.memberAcceptance"; +const ACCEPT_COLLABORATION: &str = "sh.tangled.repo.acceptCollaboration"; +const COLLABORATOR_ACCEPTANCE: &str = "sh.tangled.repo.collaboratorAcceptance"; const BAN: &str = "sh.tangled.knot.ban"; const UNBAN: &str = "sh.tangled.knot.unban"; const CREATE: &str = "sh.tangled.repo.create"; @@ -50,7 +57,7 @@ const FORK_SYNC: &str = "sh.tangled.repo.forkSync"; const HIDDEN_REF: &str = "sh.tangled.repo.hiddenRef"; type Responder = Box Result + Send + Sync>; -type SharedState = Arc, ManualClock>>; +type SharedState = Arc, Arc>>; static JTI: AtomicU64 = AtomicU64::new(0); @@ -116,12 +123,16 @@ fn repo_did_doc(did: &str, multikey: &str) -> Bytes { } fn mint(actor: &Actor, nsid: &str) -> String { + mint_for(actor, nsid, &format!("did:web:{KNOT_HOST}")) +} + +fn mint_for(actor: &Actor, nsid: &str, audience: &str) -> String { let jti = JTI.fetch_add(1, Ordering::Relaxed); let header = URL_SAFE_NO_PAD.encode(br#"{"alg":"ES256K","typ":"JWT"}"#); let payload = URL_SAFE_NO_PAD.encode( serde_json::to_vec(&json!({ "iss": actor.did.as_str(), - "aud": format!("did:web:{KNOT_HOST}"), + "aud": audience, "exp": 1_100, "iat": 999, "jti": format!("nonce-{jti}"), @@ -309,13 +320,14 @@ fn state_from( admission: AdmissionPolicy, reservations: Arc, git_http: Arc, -) -> SharedState { +) -> (SharedState, Arc) { let (layout, index, meta_path) = boot; let knot = knot_did(); let knot_url = knot_types::KnotServiceUrl::new(format!("https://{KNOT_HOST}")).unwrap(); + let clock = Arc::new(ManualClock::new(UnixMicros::new(1_000_000_000))); let atproto = Arc::new(Atproto::new( FakeHttp::new(responder), - ManualClock::new(UnixMicros::new(1_000_000_000)), + Arc::clone(&clock), knot.clone(), knot_atproto::PlcDirectory::new(url::Url::parse("https://plc.directory/").unwrap()) .unwrap(), @@ -329,7 +341,7 @@ fn state_from( .unwrap(), ); secrets.ensure(&knot).unwrap(); - Arc::new(XrpcState { + let state = Arc::new(XrpcState { layout, index, atproto, @@ -356,7 +368,7 @@ fn state_from( pack_limits: knot_pack::PackLimits::default(), service_owner: account(ADMIN_HOST), events: Arc::new(knot_events::EventLog::new( - ManualClock::new(UnixMicros::new(1_000_000_000)), + Arc::clone(&clock), knot_events::ReplayBounds::new( knot_events::ReplayEvents::new(1024).unwrap(), knot_events::ReplayBytes::new(16 << 20).unwrap(), @@ -371,13 +383,15 @@ fn state_from( slots: knot_resource::Slots::testing(8), lfs: None, catalog: Arc::new(knot_messages::Catalog::defaults()), - }) + }); + (state, clock) } fn world_responder( pubkeys: HashMap>, repo_docs: Arc>>, - pds_records: Arc>>, + pds_records: Arc>>, + reads: Arc, ) -> Responder { let doc_url = format!("https://{KNOT_HOST}"); Box::new(move |request: &HttpRequest| { @@ -389,24 +403,20 @@ fn world_responder( }); } if request.url.path().ends_with("com.atproto.repo.getRecord") { + reads.fetch_add(1, Ordering::Relaxed); let rkey = request .url .query_pairs() .find(|(key, _)| key == "rkey") .map(|(_, value)| value.into_owned()) .unwrap_or_default(); - let present = pds_records.lock().unwrap().contains(&rkey); + let answer = pds_records.lock().unwrap().get(&rkey).copied(); return Ok(HttpResponse { - status: if present { - StatusCode::OK - } else { - StatusCode::BAD_REQUEST - }, + status: answer.unwrap_or(StatusCode::BAD_REQUEST), headers: http::HeaderMap::new(), - body: if present { - Bytes::new() - } else { - Bytes::from_static(b"{\"error\":\"RecordNotFound\"}") + body: match answer { + Some(_) => Bytes::new(), + None => Bytes::from_static(b"{\"error\":\"RecordNotFound\"}"), }, }); } @@ -441,7 +451,9 @@ struct World { member: Actor, stranger: Actor, repo_docs: Arc>>, - pds_records: Arc>>, + pds_records: Arc>>, + clock: Arc, + reads: Arc, } impl World { @@ -511,7 +523,8 @@ impl World { let member = actor(2, MEMBER_HOST); let stranger = actor(3, STRANGER_HOST); let repo_docs: Arc>> = Arc::new(Mutex::new(HashMap::new())); - let pds_records: Arc>> = Arc::new(Mutex::new(HashSet::new())); + let pds_records: Arc>> = + Arc::new(Mutex::new(HashMap::new())); let pubkeys = HashMap::from([ ( ADMIN_HOST.to_string(), @@ -526,8 +539,14 @@ impl World { stranger.signer.public_key().as_bytes().to_vec(), ), ]); - let responder = world_responder(pubkeys, Arc::clone(&repo_docs), Arc::clone(&pds_records)); - let state = state_from( + let reads = Arc::new(AtomicUsize::new(0)); + let responder = world_responder( + pubkeys, + Arc::clone(&repo_docs), + Arc::clone(&pds_records), + Arc::clone(&reads), + ); + let (state, clock) = state_from( &dir, (layout.clone(), index, meta_path), responder, @@ -548,6 +567,8 @@ impl World { stranger, repo_docs, pds_records, + clock, + reads, } } @@ -563,14 +584,92 @@ impl World { } fn publish_pds_record(&self, rkey: &str) { - self.pds_records.lock().unwrap().insert(rkey.to_string()); + self.answer_pds_record(rkey, StatusCode::OK); + } + + fn answer_pds_record(&self, rkey: &str, status: StatusCode) { + self.pds_records + .lock() + .unwrap() + .insert(rkey.to_string(), status); + } + + fn advance(&self, secs: u64) { + self.clock.advance(std::time::Duration::from_secs(secs)); + } + + fn record_reads(&self) -> usize { + self.reads.load(Ordering::Relaxed) + } + + fn migrated(&self, host: &str) -> Grant { + Grant { + subject: account(host), + added_by: self.admin.did.clone(), + created_at: UnixSeconds::new(1_000), + } + } + + fn append(&self, repo: Option<&RepoDid>, change: &E::Change) + where + E: knot_cob::Evaluate, + { + let git = match repo { + None => Repo::open(&self.state.meta_path).unwrap(), + Some(repo) => self.layout.open(repo).unwrap(), + }; + let home = match repo { + None => knot_cob::CobHome::from(&self.state.knot_did), + Some(repo) => knot_cob::CobHome::from(repo), + }; + let store = knot_cob::CobStore::new(&git); + let signer = self.state.secrets.signer(&self.state.knot_did).unwrap(); + let at = UnixSeconds::new(1_000); + 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()), + } + } + + fn grandfather(&self, host: &str) { + self.append::(None, &MembersChange::Add(self.migrated(host))); + self.refresh(None); + } + + fn grandfather_collaborator(&self, repo: &RepoDid, host: &str) { + self.append::(Some(repo), &CollaboratorsChange::Add(self.migrated(host))); + self.refresh(Some(repo)); + } + + fn invite(&self, repo: Option<&RepoDid>, host: &str) { + let invite = Invite(self.migrated(host)); + match repo { + None => self.append::(repo, &MembersChange::Invite(invite)), + Some(_) => self.append::(repo, &CollaboratorsChange::Invite(invite)), + } + } + + fn refresh(&self, repo: Option<&RepoDid>) { + match repo { + None => self.state.index.refresh_members().unwrap(), + Some(repo) => self.state.index.refresh_collaborators(repo).unwrap(), + } + } + + fn standing(&self, repo: Option<&RepoDid>, subject: &AccountDid) -> Resolved> { + self.refresh(repo); + match repo { + None => self.state.index.member_standing(subject), + Some(repo) => self.state.index.collaborator_standing(repo, subject), + } } } fn build_state(responder: Responder, rebuild: bool) -> (TempDir, SharedState) { let dir = tempfile::tempdir().unwrap(); let (layout, index, meta_path) = bootstrap(&dir, rebuild, knot_types::ObjectFormat::SHA1); - let state = state_from( + let (state, _clock) = state_from( &dir, (layout, index, meta_path), responder, @@ -613,18 +712,56 @@ fn doc_responder( }) } -async fn add_member_helper(world: &World) { - assert_eq!( - as_admin( - world, - crate::members::add_member, - ADD_MEMBER, - json!({ "subject": format!("did:web:{MEMBER_HOST}") }) - ) - .await - .status(), - StatusCode::OK - ); +fn subject_body(repo: Option<&RepoDid>, subject: &AccountDid) -> serde_json::Value { + match repo { + None => json!({ "subject": subject.as_str() }), + Some(repo) => json!({ "repo": repo.as_str(), "subject": subject.as_str() }), + } +} + +async fn offer(world: &World, repo: Option<&RepoDid>, subject: &AccountDid) -> StatusCode { + let body = subject_body(repo, subject); + match repo { + None => as_admin(world, add_member, ADD_MEMBER, body).await, + Some(_) => as_member(world, add_collaborator, ADD_COLLAB, body).await, + } + .status() +} + +fn named_subject(world: &World, repo: Option<&RepoDid>) -> String { + repo.map_or_else(|| world.state.knot_did.as_str(), RepoDid::as_str) + .to_string() +} + +fn acceptance_uri(author: &AccountDid, repo: Option<&RepoDid>, subject: &str) -> serde_json::Value { + let collection = repo.map_or(MEMBER_ACCEPTANCE, |_| COLLABORATOR_ACCEPTANCE); + json!({ "acceptance": format!("at://{}/{collection}/{subject}", author.as_str()) }) +} + +async fn accept_call( + world: &World, + repo: Option<&RepoDid>, + actor: &Actor, + audience: &str, + uri: serde_json::Value, +) -> StatusCode { + let nsid = repo.map_or(ACCEPT_MEMBERSHIP, |_| ACCEPT_COLLABORATION); + let state = world.state(); + let headers = bearer(&mint_for(actor, nsid, audience)); + let method = crate::Method::from_nsid(nsid); + let body = body(uri); + into_response(match repo { + None => crate::members::accept_membership(state, headers, method, body).await, + Some(_) => crate::collaborators::accept_collaboration(state, headers, method, body).await, + }) + .status() +} + +async fn answer(world: &World, actor: &Actor, repo: Option<&RepoDid>) -> StatusCode { + let subject = named_subject(world, repo); + world.publish_pds_record(&subject); + let uri = acceptance_uri(&actor.did, repo, &subject); + accept_call(world, repo, actor, &subject, uri).await } async fn create_repo_helper(world: &World, name: &str) -> RepoDid { @@ -795,20 +932,23 @@ async fn blocklist_lifecycle() { async fn member_lifecycle() { let world = World::new(); let subject = json!({ "subject": format!("did:web:{MEMBER_HOST}") }); + let member = account(MEMBER_HOST); + assert_eq!(offer(&world, None, &member).await, StatusCode::OK); assert_eq!( - as_admin( - &world, - crate::members::add_member, - ADD_MEMBER, - subject.clone() - ) - .await - .status(), - StatusCode::OK + world.state.index.member_standing(&member), + Resolved::Ready(Some(Standing::Invited)), + "only the knot has signed this membership, so the account stands Invited" ); assert_eq!( - world.state.index.is_member(&account(MEMBER_HOST)), + 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" ); @@ -822,20 +962,14 @@ async fn member_lifecycle() { let baseline = event_count(&world); assert_eq!( - as_admin( - &world, - crate::members::add_member, - ADD_MEMBER, - subject.clone() - ) - .await - .status(), - StatusCode::OK + 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 add is a no-op and emits no event" + "a redundant offer appended a second add to the roster" ); assert_eq!( @@ -850,7 +984,7 @@ async fn member_lifecycle() { StatusCode::OK ); assert_eq!( - world.state.index.is_member(&account(MEMBER_HOST)), + world.state.index.effective_member(&account(MEMBER_HOST)), Resolved::Ready(false), "removed member is gone on the very next read" ); @@ -1004,14 +1138,14 @@ async fn concurrent_first_member_adds_converge_to_a_single_cob_object() { 1, "concurrent first adds serialize onto one singleton members COB, never splitting it" ); - assert_eq!( - world.state.index.is_member(&account("witchcraft.systems")), - Resolved::Ready(true) - ); - assert_eq!( - world.state.index.is_member(&account("isabelroses.com")), - Resolved::Ready(true) - ); + let hosts = ["witchcraft.systems", "isabelroses.com"]; + hosts.iter().for_each(|host| { + assert_eq!( + world.standing(None, &account(host)), + Resolved::Ready(Some(Standing::Invited)), + "both concurrent adds reached the one roster, and {host} stands Invited" + ); + }); } #[tokio::test] @@ -1071,7 +1205,7 @@ async fn re_adding_a_member_under_warming_appends_no_redundant_change() { #[tokio::test] async fn create_mints_a_did_plc_repo_and_refuses_a_duplicate_name() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "anemone").await; assert!( @@ -1112,7 +1246,7 @@ async fn create_mints_a_did_plc_repo_and_refuses_a_duplicate_name() { #[tokio::test] async fn create_and_reserve_reject_bad_identities() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); assert_eq!( create_status( @@ -1136,7 +1270,7 @@ async fn create_and_reserve_reject_bad_identities() { "the knot's own DID as a repoDid is a client error instead of a 500" ); assert_eq!( - world.state.index.is_member(&account(MEMBER_HOST)), + world.state.index.effective_member(&account(MEMBER_HOST)), Resolved::Ready(true), "the meta-repo is intact" ); @@ -1230,7 +1364,7 @@ async fn create_and_reserve_reject_bad_identities() { #[tokio::test] async fn a_byo_did_web_repo_is_accepted_and_its_key_is_returned() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let did = "did:web:nautilus.olaren.dev"; let reserved_key = reserve_repo_key(&world, did).await; @@ -1281,8 +1415,8 @@ async fn a_byo_did_web_repo_is_accepted_and_its_key_is_returned() { world .state .index - .is_collaborator(&repo_did, &account("witchcraft.systems")), - Resolved::Ready(true), + .collaborator_standing(&repo_did, &account("witchcraft.systems")), + Resolved::Ready(Some(Standing::Invited)), "the collaborator COB signed by the knot-held repo key lands and is visible" ); @@ -1462,14 +1596,14 @@ async fn resolve_by_name_matches_the_rkey_case_sensitively() { ); assert!( crate::merge::resolve_by_name(&*state, &owner, &RepoName::new("Anemone").unwrap()).is_err(), - "a differently-cased name must not resolve to a distinct rkey, atproto record keys are case-sensitive" + "rkeys are case-sensitive, so a differently-cased name mustn't resolve at all" ); } #[tokio::test] async fn reserve_key_refuses_once_the_pending_limit_is_reached() { let world = World::with_pending_limit(2); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); assert_eq!( reserve_status(&world, &world.member, "did:web:p0.olaren.dev").await, @@ -1491,7 +1625,7 @@ async fn reserve_key_refuses_once_the_pending_limit_is_reached() { #[tokio::test] async fn one_account_cannot_exhaust_the_global_reservation_budget() { let world = World::with_limits(256, 2); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let world_ref = &world; futures::stream::iter(["did:web:m0.olaren.dev", "did:web:m1.olaren.dev"]) @@ -1517,7 +1651,7 @@ async fn one_account_cannot_exhaust_the_global_reservation_budget() { #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn concurrent_reserve_key_calls_all_succeed() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let handles: Vec<_> = (0..24) .map(|i| { @@ -1554,34 +1688,38 @@ async fn concurrent_reserve_key_calls_all_succeed() { #[tokio::test] async fn collaborator_lifecycle() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "scallop").await; - let subject = || json!({ "repo": repo_did, "subject": "did:web:witchcraft.systems" }); + let collaborator = account(STRANGER_HOST); + let subject = || json!({ "repo": repo_did, "subject": collaborator.as_str() }); + let repo = Some(&repo_did); + assert_eq!(offer(&world, repo, &collaborator).await, StatusCode::OK); assert_eq!( - as_member( - &world, - crate::collaborators::add_collaborator, - ADD_COLLAB, - subject() - ) - .await - .status(), - StatusCode::OK + world + .state + .index + .collaborator_standing(&repo_did, &collaborator), + 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!( world .state .index - .is_collaborator(&repo_did, &account("witchcraft.systems")), + .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"], - account("witchcraft.systems").to_string() - ); + assert_eq!(added.payload["subject"], collaborator.to_string()); assert_eq!(added.payload["repo"], repo_did.to_string()); assert_eq!( @@ -1599,16 +1737,13 @@ async fn collaborator_lifecycle() { world .state .index - .is_collaborator(&repo_did, &account("witchcraft.systems")), + .effective_collaborator(&repo_did, &collaborator), 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"], - account("witchcraft.systems").to_string() - ); + assert_eq!(removed.payload["subject"], collaborator.to_string()); assert_eq!(removed.payload["repo"], repo_did.to_string()); let baseline = event_count(&world); @@ -1633,7 +1768,7 @@ async fn collaborator_lifecycle() { #[tokio::test] async fn repo_management_is_owner_or_collaborator_gated() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "squid").await; assert_eq!( @@ -1695,17 +1830,7 @@ async fn repo_management_is_owner_or_collaborator_gated() { "the canonical rkey is untouched by the rejected rename" ); - assert_eq!( - as_member( - &world, - crate::collaborators::add_collaborator, - ADD_COLLAB, - json!({ "repo": repo_did, "subject": format!("did:web:{STRANGER_HOST}") }) - ) - .await - .status(), - StatusCode::OK - ); + world.grandfather_collaborator(&repo_did, STRANGER_HOST); assert_eq!( rename_repo_as(&world, &world.stranger, &repo_did, "periwinkle").await, StatusCode::OK, @@ -1716,7 +1841,7 @@ async fn repo_management_is_owner_or_collaborator_gated() { #[tokio::test] async fn the_owner_sets_the_default_branch_by_repo_did() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "mussel").await; assert_eq!( @@ -1757,7 +1882,7 @@ async fn set_default_branch_rejections() { use knot_types::RefName; let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "mussel").await; let git = world.layout.open(&repo_did).unwrap(); @@ -1815,7 +1940,7 @@ async fn delete_branch_removes_a_branch_and_refuses_the_default() { use knot_types::RefName; let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "periwinkle").await; let git = world.layout.open(&repo_did).unwrap(); let now = world.state.now(); @@ -1899,7 +2024,7 @@ async fn delete_branch_removes_a_branch_and_refuses_the_default() { #[tokio::test] async fn delete_repo_lifecycle_and_guards() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let plain = create_repo_helper(&world, "whelk").await; assert_eq!( @@ -1962,7 +2087,7 @@ async fn delete_repo_lifecycle_and_guards() { #[tokio::test] async fn rename_alias_lifecycle() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_a = create_repo_helper(&world, "alpha").await; assert_eq!( @@ -2159,7 +2284,7 @@ async fn the_http_push_surface_sheds_a_bogus_credential_flood_from_one_peer() { use tower::ServiceExt; let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo = create_repo_helper(&world, "kelp").await; let app = crate::router(Arc::clone(&world.state)); let peer = SocketAddr::from(([203, 0, 113, 11], 5555)); @@ -2346,7 +2471,7 @@ mod merge_endpoints { #[tokio::test] async fn the_owner_merges_a_unified_patch_natively_and_cleans_up() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "kelp").await; let base = seed_main(&world, &repo_did, &[("reef.txt", "old line\n")]); @@ -2412,7 +2537,7 @@ mod merge_endpoints { #[tokio::test] async fn a_native_merge_advances_the_branch_without_a_pipeline_event() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "kelp").await; seed_main( &world, @@ -2458,7 +2583,7 @@ mod merge_endpoints { #[tokio::test] async fn a_format_patch_merge_creates_one_commit_per_patch_with_change_id() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "limpet").await; let base = seed_main(&world, &repo_did, &[("reef.txt", "one\n")]); @@ -2537,7 +2662,7 @@ mod merge_endpoints { #[tokio::test] async fn merge_rejections_move_nothing() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "scallop").await; let base = seed_main(&world, &repo_did, &[("reef.txt", "old line\n")]); @@ -2598,7 +2723,7 @@ mod merge_endpoints { #[tokio::test] async fn merge_check_is_open_and_reports_clean_conflicted_and_broken() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let repo_did = create_repo_helper(&world, "scallop").await; let base = seed_main(&world, &repo_did, &[("reef.txt", "old line\n")]); @@ -2653,7 +2778,7 @@ mod merge_endpoints { #[tokio::test] async fn a_repo_did_addresses_a_repo_whose_rkey_differs_from_its_name() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); assert_eq!( as_member( &world, @@ -2822,7 +2947,7 @@ mod fork_endpoints { async fn forked_world() -> ForkWorld { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let source_did = create_repo_helper(&world, "kelp").await; let source = world.layout.open(&source_did).unwrap(); advance(&source, &main_ref(), "reef.txt", "kelp forest\n", 1_000); @@ -2919,7 +3044,7 @@ mod fork_endpoints { #[tokio::test] async fn forking_a_source_this_knot_does_not_host_is_not_found() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); assert_eq!( as_member(&world, crate::repos::create_repo, CREATE, json!({ "rkey": "uni", "name": "uni", "source": format!("https://{KNOT_HOST}/did:plc:whelk/ghost") })).await.status(), StatusCode::NOT_FOUND @@ -2989,7 +3114,7 @@ mod fork_endpoints { #[tokio::test] async fn a_repo_did_addresses_a_fork_whose_rkey_differs_from_its_name() { let world = World::new(); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let source_did = create_repo_helper(&world, "kelp").await; let source = world.layout.open(&source_did).unwrap(); advance(&source, &main_ref(), "reef.txt", "kelp forest\n", 1_000); @@ -3198,7 +3323,7 @@ mod fork_endpoints { })); let world = World::with_git_http(git_http, knot_types::ObjectFormat::SHA256); - add_member_helper(&world).await; + world.grandfather(MEMBER_HOST); let plain = create_repo_helper(&world, "kelp").await; assert_eq!( world.layout.open(&plain).unwrap().object_format(), @@ -3325,7 +3450,7 @@ mod legacy_admin_route { StatusCode::PAYLOAD_TOO_LARGE ); assert_eq!( - world.state.index.is_member(&account(MEMBER_HOST)), + world.state.index.effective_member(&account(MEMBER_HOST)), Resolved::Ready(false), "a refused call grants nothing" ); @@ -3335,24 +3460,37 @@ mod legacy_admin_route { StatusCode::OK ); assert_eq!( - world.state.index.is_member(&account(MEMBER_HOST)), - Resolved::Ready(true) + world.state.index.member_standing(&account(MEMBER_HOST)), + 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 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 Resolved::Ready(members) = world.state.index.member_entries() else { panic!("the member roster is warm in this test"); }; assert_eq!( members .iter() - .find(|grant| grant.subject == account(MEMBER_HOST)) + .find(|(subject, _)| *subject == account(MEMBER_HOST)) .expect("the member is in the roster") + .1 .added_by, world.state.service_owner, - "the legacy grant records the service owner as the granter" + "the legacy route has no actor of its own, so `added_by` is the service owner" + ); + + assert_eq!(answer(&world, &world.member, None).await, StatusCode::OK); + assert_eq!( + 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!( @@ -3363,7 +3501,7 @@ mod legacy_admin_route { assert_eq!( event_count(&world), baseline, - "re-adding an existing member emits no event" + "the repeat call appended an add for a member already in the roster" ); } @@ -3429,3 +3567,362 @@ mod legacy_admin_route { ); } } + +mod rosters { + use super::*; + + const DENIED: &str = "you have no live collaboration on this repository"; + const TTL: u64 = 60; + const GRACE: u64 = 21_600; + + #[derive(Debug, Clone, Copy)] + enum Roster { + Knot, + Repository, + } + + #[derive(Debug, Clone, Copy)] + enum Seen { + Projected, + Unprojected, + } + + struct Offer { + world: World, + from: String, + repo: Option, + subject: AccountDid, + } + + impl Offer { + async fn opened(roster: Roster, seen: Seen) -> Self { + let world = World::new(); + let repo = match roster { + Roster::Knot => None, + Roster::Repository => { + world.grandfather(MEMBER_HOST); + Some(create_repo_helper(&world, "nautilus").await) + } + }; + let subject = account(STRANGER_HOST); + world.invite(repo.as_ref(), STRANGER_HOST); + let offer = Self { + world, + from: format!("{roster:?}/{seen:?}"), + repo, + subject, + }; + if matches!(seen, Seen::Projected) { + assert_eq!( + offer.standing(), + Resolved::Ready(Some(Standing::Invited)), + "{}: the fixture must open at Invited", + offer.from + ); + } + offer + } + + fn repo(&self) -> Option<&RepoDid> { + self.repo.as_ref() + } + + fn named(&self) -> String { + named_subject(&self.world, self.repo()) + } + + fn standing(&self) -> Resolved> { + self.world.standing(self.repo(), &self.subject) + } + + fn answered(&self, status: StatusCode) { + self.world.answer_pds_record(&self.named(), status); + } + + fn written_by(&self, author: &AccountDid) -> serde_json::Value { + acceptance_uri(author, self.repo(), &self.named()) + } + + fn own(&self) -> serde_json::Value { + self.written_by(&self.subject) + } + + async fn accept_to(&self, audience: &str, uri: serde_json::Value) -> StatusCode { + let actor = &self.world.stranger; + accept_call(&self.world, self.repo(), actor, audience, uri).await + } + + async fn accept(&self, uri: serde_json::Value) -> StatusCode { + self.accept_to(&self.named(), uri).await + } + + async fn accepted_by(&self, actor: &Actor) -> StatusCode { + let uri = self.written_by(&actor.did); + accept_call(&self.world, self.repo(), actor, &self.named(), uri).await + } + + async fn allowed(&self) -> bool { + let state = &self.world.state; + match self.repo() { + Some(repo) => crate::authorize_push(state, &self.subject, repo, DENIED).await, + None => crate::authorize_create(state, &self.subject, DENIED).await, + } + .is_ok() + } + } + + async fn every_roster(check: impl AsyncFn(Offer)) { + check(Offer::opened(Roster::Knot, Seen::Projected).await).await; + check(Offer::opened(Roster::Repository, Seen::Projected).await).await; + } + + async fn every_roster_and_projection(check: impl AsyncFn(Offer)) { + check(Offer::opened(Roster::Knot, Seen::Projected).await).await; + check(Offer::opened(Roster::Knot, Seen::Unprojected).await).await; + check(Offer::opened(Roster::Repository, Seen::Projected).await).await; + check(Offer::opened(Roster::Repository, Seen::Unprojected).await).await; + } + + #[tokio::test] + async fn only_the_accepters_own_published_record_moves_the_offer() { + every_roster(async |offer: Offer| { + let from = &offer.from; + let elsewhere = match offer.repo() { + None => "did:web:other.nel.pet", + Some(_) => "did:plc:3fwecdnvtcscjnrx2p4n7alz", + }; + + assert_eq!( + offer.accept(offer.own()).await, + StatusCode::FORBIDDEN, + "{from}: the accepter published no such record" + ); + + offer.answered(StatusCode::BAD_GATEWAY); + assert_eq!( + offer.accept(offer.own()).await, + StatusCode::SERVICE_UNAVAILABLE, + "{from}: a PDS the knot can't reach decides neither way" + ); + + offer.answered(StatusCode::OK); + assert_eq!( + offer.accept(offer.written_by(&account(MEMBER_HOST))).await, + StatusCode::FORBIDDEN, + "{from}: a record in another account's repository accepted for this one" + ); + + offer.world.publish_pds_record(elsewhere); + let audience = match offer.repo() { + None => offer.named(), + Some(_) => elsewhere.to_string(), + }; + let uri = acceptance_uri(&offer.subject, offer.repo(), elsewhere); + assert_eq!( + offer.accept_to(&audience, uri).await, + match offer.repo() { + None => StatusCode::FORBIDDEN, + Some(_) => StatusCode::NOT_FOUND, + }, + "{from}: an acceptance published against {elsewhere} answers no offer this \ + roster made" + ); + + assert_eq!( + offer.accepted_by(&offer.world.admin).await, + StatusCode::FORBIDDEN, + "{from}: an admin's authority is configuration, so no roster has a row to accept" + ); + assert_eq!( + offer.accepted_by(&offer.world.member).await, + StatusCode::FORBIDDEN, + "{from}: an account with no entry here, owner of the thing or not, has nothing \ + to accept" + ); + assert_eq!( + offer.standing(), + Resolved::Ready(Some(Standing::Invited)), + "{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, + StatusCode::OK, + "{from}: a client retrying through its own PDS outage mustn't lose an acceptance \ + the knot holds" + ); + assert_eq!( + offer.standing(), + 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!( + offer.accept_to(&fragment, offer.own()).await, + match offer.repo() { + None => StatusCode::OK, + Some(_) => StatusCode::UNAUTHORIZED, + }, + "{from}: each route takes the did of its own subject as the token audience, \ + service fragment and all" + ); + }) + .await; + } + + #[tokio::test] + async fn a_second_offer_leaves_the_invitation_and_the_admins_revocation_clears_it() { + every_roster_and_projection(async |invited: Offer| { + 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), + "{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" + ); + assert_eq!( + invited.standing(), + Resolved::Ready(Some(Standing::Invited)), + "{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, + Some(_) => as_member(world, remove_collaborator, REMOVE_COLLAB, body).await, + } + .status(); + assert_eq!( + (revoked, event_count(world) - before), + (StatusCode::OK, 1), + "{from}: a revocation is one of the two changes an acl consumer reads" + ); + assert_eq!( + invited.standing(), + Resolved::Ready(None), + "{from}: the revocation has to leave a Removal behind, or a later acceptance \ + still grants" + ); + }) + .await; + } + + #[tokio::test] + async fn grandfathered_collaborator_pushes_off_roster_alone() { + let world = World::new(); + world.grandfather(MEMBER_HOST); + let repo = create_repo_helper(&world, "nautilus").await; + let subject = account(STRANGER_HOST); + world.grandfather_collaborator(&repo, STRANGER_HOST); + let pushes = async || crate::authorize_push(&world.state, &subject, &repo, DENIED).await; + assert!( + pushes().await.is_ok(), + "the knot granted this entry itself, so no acceptance record is needed" + ); + world.advance(GRACE * 2); + assert_eq!( + (pushes().await.is_ok(), world.record_reads()), + (true, 0), + "a grandfathered entry is never resolved, so the grace passes with no PDS read" + ); + + let idle = Offer::opened(Roster::Repository, Seen::Projected).await; + idle.world.advance(GRACE * 2); + assert_eq!( + idle.world.record_reads(), + 0, + "time alone resolved this offer, and only a gated call may read the PDS" + ); + } + + #[tokio::test] + async fn the_next_gated_call_reads_the_acceptance_and_a_deletion_revokes_inside_the_ttl() { + every_roster(async |invited: Offer| { + let from = &invited.from; + assert_eq!( + (invited.allowed().await, invited.world.record_reads()), + (false, 1), + "{from}: an invitation with no acceptance must deny, and one read must settle it" + ); + + invited.answered(StatusCode::OK); + invited.world.advance(TTL + 1); + assert!( + invited.allowed().await, + "{from}: the acceptance is in the subject's own repository, and no second call \ + from them is needed for the knot to find it" + ); + + invited.answered(StatusCode::NOT_FOUND); + invited.world.advance(TTL + 1); + assert!( + !invited.allowed().await, + "{from}: deleting the acceptance is the whole revocation, and the first read past \ + the ttl denies" + ); + }) + .await; + } + + #[tokio::test] + async fn the_acceptance_route_takes_only_a_token_minted_for_its_own_method() { + use axum::body::Body; + use axum::extract::ConnectInfo; + use std::net::SocketAddr; + use tower::ServiceExt; + + let offer = Offer::opened(Roster::Knot, Seen::Projected).await; + offer.answered(StatusCode::OK); + let app = crate::router(Arc::clone(&offer.world.state)); + let peer = SocketAddr::from(([203, 0, 113, 41], 5555)); + + let post = |nsid: &str| { + let token = format!("Bearer {}", mint(&offer.world.stranger, nsid)); + let mut request = http::Request::builder() + .method("POST") + .uri(crate::members::ACCEPT_ROUTE) + .header(http::header::AUTHORIZATION, token) + .body(Body::from(serde_json::to_vec(&offer.own()).unwrap())) + .unwrap(); + request.extensions_mut().insert(ConnectInfo(peer)); + request + }; + + let sibling = post(ADD_MEMBER); + assert_eq!( + app.clone().oneshot(sibling).await.unwrap().status(), + StatusCode::UNAUTHORIZED, + "the acceptance route took a token minted for a sibling method" + ); + assert_eq!( + app.oneshot(post(ACCEPT_MEMBERSHIP)).await.unwrap().status(), + StatusCode::OK + ); + } +} diff --git a/knot2/crates/knot-xrpc/tests/common/mod.rs b/knot2/crates/knot-xrpc/tests/common/mod.rs index f14c396e7..c9234ecdd 100644 --- a/knot2/crates/knot-xrpc/tests/common/mod.rs +++ b/knot2/crates/knot-xrpc/tests/common/mod.rs @@ -17,10 +17,11 @@ use k256::ecdsa::{Signature, SigningKey}; use tower::ServiceExt; use knot_atproto::Atproto; -use knot_cob::{CobHome, CobStore}; +use knot_cob::{ChangePayload, CobHome, CobStore, Evaluate}; use knot_cobs::{ - CollaboratorsChange, CollaboratorsCob, Grant, MembersChange, MembersCob, Registration, - RegistryChange, RepoRegistryCob, register_repo, + Accept, BlocklistChange, BlocklistCob, CollaboratorsChange, CollaboratorsCob, Grant, Invite, + MembersChange, MembersCob, Registration, RegistryChange, Removal, RepoRegistryCob, + register_repo, }; use knot_git::{Layout, Repo}; use knot_index::Index; @@ -266,47 +267,69 @@ impl World { } pub fn add_member(&self, subject: &str, added_by: &str, at: i64) { - let meta = Repo::open(&self.state.meta_path).unwrap(); - let store = CobStore::new(&meta); - let home = CobHome::from(&self.state.knot_did); - let signer = self.state.secrets.signer(&self.state.knot_did).unwrap(); - let change = MembersChange::Add(grant(subject, added_by, at)); - match store.list::().unwrap().as_slice() { - [] => { - store - .create(&home, &change, &signer, UnixSeconds::new(at)) - .unwrap(); - } - [object] => { - store - .update(&home, *object, &change, &signer, UnixSeconds::new(at)) - .unwrap(); - } - many => panic!("{} members objects", many.len()), - } - self.state.index.refresh_members().unwrap(); + self.member(MembersChange::Add(grant(subject, added_by, at)), at); + } + + pub fn invite_member(&self, subject: &str, at: i64) { + self.member(MembersChange::Invite(Invite(grant(subject, OWNER, at))), at); + } + + pub fn accept_membership(&self, subject: &str, at: i64) { + self.member( + MembersChange::Accept(Accept { + subject: AccountDid::new(subject).unwrap(), + verified_at: UnixSeconds::new(at), + }), + at, + ); + } + + pub fn remove_member(&self, subject: &str, at: i64) { + let subject = AccountDid::new(subject).unwrap(); + self.member(MembersChange::Remove(Removal { subject }), at); + } + + pub fn ban(&self, subject: &str, at: i64) { + self.on_meta::(BlocklistChange::Add(grant(subject, OWNER, at)), at); + self.state.index.refresh_blocklist().unwrap(); } pub fn add_collaborator(&self, repo: &RepoDid, subject: &str, added_by: &str, at: i64) { + let change = CollaboratorsChange::Add(grant(subject, added_by, at)); + self.collaborator(repo, change, at); + } + + pub fn invite_collaborator(&self, repo: &RepoDid, subject: &str, at: i64) { + let change = CollaboratorsChange::Invite(Invite(grant(subject, OWNER, at))); + self.collaborator(repo, change, at); + } + + fn member(&self, change: MembersChange, at: i64) { + self.on_meta::(change, at); + self.state.index.refresh_members().unwrap(); + } + + fn collaborator(&self, repo: &RepoDid, change: CollaboratorsChange, at: i64) { let git = self.layout.open(repo).unwrap(); + self.append::(git, CobHome::from(repo), change, at); + self.state.index.refresh_collaborators(repo).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); + } + + fn append(&self, git: Repo, home: CobHome, change: E::Change, at: i64) { let store = CobStore::new(&git); - let home = CobHome::from(repo); let signer = self.state.secrets.signer(&self.state.knot_did).unwrap(); - let change = CollaboratorsChange::Add(grant(subject, added_by, at)); - match store.list::().unwrap().as_slice() { - [] => { - store - .create(&home, &change, &signer, UnixSeconds::new(at)) - .unwrap(); - } - [object] => { - store - .update(&home, *object, &change, &signer, UnixSeconds::new(at)) - .unwrap(); - } - many => panic!("{} collaborators objects", many.len()), + let kind = ::TYPE; + let at = UnixSeconds::new(at); + match store.list::().unwrap().as_slice() { + [] => _ = store.create(&home, &change, &signer, at).unwrap(), + [object] => _ = store.update(&home, *object, &change, &signer, at).unwrap(), + many => panic!("{} {kind} objects", many.len()), } - self.state.index.refresh_collaborators(repo).unwrap(); } } @@ -647,6 +670,15 @@ pub fn repo_dids(value: &serde_json::Value) -> Vec { .collect() } +pub fn column<'p>(page: &'p serde_json::Value, field: &str) -> Vec<&'p str> { + page["items"] + .as_array() + .unwrap() + .iter() + .map(|item| item[field].as_str().unwrap()) + .collect() +} + pub async fn archive_full(world: &World, did: &RepoDid) -> (String, String, Bytes) { let (status, headers, body) = get( world, diff --git a/knot2/crates/knot-xrpc/tests/reads.rs b/knot2/crates/knot-xrpc/tests/reads.rs index 01da15e98..1216f1f6c 100644 --- a/knot2/crates/knot-xrpc/tests/reads.rs +++ b/knot2/crates/knot-xrpc/tests/reads.rs @@ -10,13 +10,14 @@ 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_xrpc::{ArchiveLimit, ResponseLimit}; use common::{ OWNER, World, archive_full, assert_immutable_round_trip, assert_post_rejected, assert_warming, - commit_file, empty_repo, get, get_error, get_json, get_with_headers, git_run, post_authed, - post_json, ref_names, repo_dids, seeded, seeded_feature_branch, sh_git, sh_git_at, + column, commit_file, empty_repo, get, get_error, get_json, get_with_headers, git_run, + post_authed, post_json, ref_names, repo_dids, seeded, seeded_feature_branch, sh_git, sh_git_at, }; #[tokio::test] @@ -1263,6 +1264,8 @@ async fn every_projection_read_fails_closed_while_warming() { "/xrpc/sh.tangled.repo.listCollaborators?subject=did:plc:squid".to_string(), None, ), + (KNOT_OFFERS.to_string(), None), + (REPO_OFFERS.to_string(), None), ]; let w = &world; stream::iter(cases.iter()) @@ -1728,6 +1731,170 @@ async fn list_members_pages_in_the_wire_shape() { assert_eq!(ascending["items"][0]["subject"], "did:plc:limpet"); } +const KNOT_OFFERS: &str = "/xrpc/sh.tangled.knot.listMemberInvites?subject=did:web:knot.nel.pet"; +const REPO_OFFERS: &str = "/xrpc/sh.tangled.repo.listCollaboratorInvites?subject=did:plc:squid"; + +#[tokio::test] +async fn a_members_page_orders_on_when_each_member_started_granting() { + let world = World::new(); + world.add_member("did:plc:cowrie", OWNER, 500); + world.add_member("did:plc:limpet", OWNER, 1_000); + world.invite_member("did:plc:whelk", 1_000); + world.invite_member("did:plc:barnacle", 2_000); + world.accept_membership("did:plc:limpet", 9_000); + world.accept_membership("did:plc:whelk", 9_000); + + let page = get_json( + &world, + "/xrpc/sh.tangled.knot.listMembers?subject=did:web:knot.nel.pet&order=asc", + ) + .await; + let items = page["items"].as_array().unwrap(); + let stamps = |row: usize| { + ( + items[row]["createdAt"].as_str(), + items[row]["verifiedAt"].as_str(), + ) + }; + assert_eq!( + column(&page, "subject"), + ["did:plc:cowrie", "did:plc:limpet", "did:plc:whelk"], + "listMembers must leave barnacle's invite out; whelk was offered before limpet and \ + started granting after it, so ordering on the offer would shift the offsets under \ + anybody mid-walk" + ); + assert!( + items[0]["verifiedAt"].is_null(), + "a grandfathered row never had an Accept to timestamp" + ); + assert_eq!( + stamps(1), + stamps(2), + "a grant and an invite both accepted late have the same two timestamps, so neither can \ + be the sort key" + ); + assert_eq!( + column(&page, "effectiveSince"), + [ + "1970-01-01T00:08:20Z", + "1970-01-01T00:16:40Z", + "1970-01-01T02:30:00Z" + ], + "effectiveSince is the grant time for cowrie and limpet, and the acceptance time for \ + whelk, so a late acceptance never moves a grandfathered row" + ); + + let offers = get_json(&world, KNOT_OFFERS).await; + let offer = &offers["items"][0]; + assert_eq!( + column(&offers, "subject"), + ["did:plc:barnacle"], + "whelk answered and the granted rows were never offered anything, so barnacle alone is \ + outstanding, and the acl read above stays the effective set that knotacl takes a \ + membership from" + ); + assert_eq!(offer["addedBy"], OWNER); + assert_eq!(offer["createdAt"], "1970-01-01T00:33:20Z"); + assert!( + offer.get("effectiveSince").is_none() && offer.get("verifiedAt").is_none(), + "effectiveSince and verifiedAt arrive with the acceptance, and barnacle hasn't accepted" + ); + assert!( + offers.get("cursor").is_none(), + "one offer fits inside the default limit, so the page ends with no cursor" + ); +} + +#[tokio::test] +async fn either_roster_pages_its_outstanding_offers_and_drops_a_withdrawn_or_banned_one() { + let world = World::new(); + let squid = RepoDid::new("did:plc:squid").unwrap(); + world.layout.create(&squid).unwrap(); + world.register(&squid, "squid"); + let (other, _other) = seeded(&world, "limpet"); + world.add_member("did:plc:mussel", OWNER, 500); + world.add_collaborator(&squid, "did:plc:mussel", OWNER, 500); + world.invite_collaborator(&other, "did:plc:oyster", 2_000); + ["did:plc:whelk", "did:plc:barnacle", "did:plc:cowrie"] + .into_iter() + .zip([1_000, 2_000, 3_000]) + .for_each(|(subject, at)| { + world.invite_member(subject, at); + world.invite_collaborator(&squid, subject, at); + }); + + let w = &world; + stream::iter([KNOT_OFFERS, REPO_OFFERS]) + .for_each(|base| async move { + let page = get_json(w, &format!("{base}&limit=2")).await; + assert_eq!( + column(&page, "subject"), + ["did:plc:cowrie", "did:plc:barnacle"], + "the newest offer is served first without an order parameter, a granting row is \ + not an offer, and an offer belongs to the roster that made it, at {base}" + ); + let cursor = page["cursor"].as_str().unwrap(); + let next = get_json(w, &format!("{base}&limit=2&cursor={cursor}")).await; + assert_eq!( + column(&next, "subject"), + ["did:plc:whelk"], + "the cursor picks up at whelk with no row repeated, at {base}" + ); + assert!( + next.get("cursor").is_none(), + "whelk was the last of three, so a client stops walking here, at {base}" + ); + let ascending = get_json(w, &format!("{base}&order=asc&limit=1")).await; + assert_eq!( + column(&ascending, "subject"), + ["did:plc:whelk"], + "order=asc serves the oldest offer first, at {base}" + ); + }) + .await; + let unhosted = get_json(&world, &REPO_OFFERS.replace("squid", "unhosted")).await; + assert!( + unhosted["items"].as_array().unwrap().is_empty(), + "a client asking after an unhosted repository gets an empty page, not a 404" + ); + + world.remove_member("did:plc:barnacle", 4_000); + assert_eq!( + column(&get_json(&world, KNOT_OFFERS).await, "subject"), + ["did:plc:cowrie", "did:plc:whelk"], + "a removal clears barnacle's outstanding offer from the listing" + ); + world.ban("did:plc:whelk", 4_000); + assert_eq!( + column(&get_json(&world, KNOT_OFFERS).await, "subject"), + ["did:plc:cowrie"], + "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()), + Resolved::Ready(Some(Standing::Invited)), + "the blocklist is a roster of its own and neither change rewrites whelk's entry on the \ + members roster" + ); +} + +#[tokio::test] +async fn the_member_listings_reject_a_subject_that_names_another_knot() { + let world = World::new(); + let w = &world; + stream::iter(["listMemberInvites", "listMembers"]) + .for_each(|method| async move { + let path = format!("/xrpc/sh.tangled.knot.{method}?subject=did:web:oyster.cafe"); + let (status, error) = get_error(w, &path).await; + assert_eq!(status, StatusCode::BAD_REQUEST, "path {path}"); + assert_eq!( + error, "InvalidSubject", + "the subject must name this knot, whatever roster the method reads, at {path}" + ); + }) + .await; +} + #[tokio::test] async fn list_members_rejects_malformed_params_and_clamps_the_limit() { let world = World::new(); @@ -1805,15 +1972,6 @@ async fn the_subject_tie_break_stays_ascending_in_both_directions() { world.add_member("did:plc:whelk", OWNER, 1_000); world.add_member("did:plc:limpet", OWNER, 1_000); - let subjects = |page: &serde_json::Value| -> Vec { - page["items"] - .as_array() - .unwrap() - .iter() - .map(|item| item["subject"].as_str().unwrap().to_string()) - .collect() - }; - let asc = get_json( &world, "/xrpc/sh.tangled.knot.listMembers?subject=did:web:knot.nel.pet&order=asc", @@ -1824,10 +1982,9 @@ async fn the_subject_tie_break_stays_ascending_in_both_directions() { "/xrpc/sh.tangled.knot.listMembers?subject=did:web:knot.nel.pet&order=desc", ) .await; - assert_eq!(subjects(&asc), vec!["did:plc:limpet", "did:plc:whelk"]); assert_eq!( - subjects(&desc), - vec!["did:plc:limpet", "did:plc:whelk"], + [column(&asc, "subject"), column(&desc, "subject")], + [["did:plc:limpet", "did:plc:whelk"]; 2], "equal-createdAt entries keep an ascending subject tie-break regardless of sort direction" ); } diff --git a/knot2/example.toml b/knot2/example.toml index f2839f865..67b660c5c 100644 --- a/knot2/example.toml +++ b/knot2/example.toml @@ -110,6 +110,14 @@ # Can also be specified via environment variable `KNOT_LEGACY_ADMIN_SECRET_ENV`. #legacy_admin_secret_env = +# Can also be specified via environment variable `KNOT_ACL_ACCEPTANCE_TTL_SECS`. +# Default value: 60 +#acceptance_ttl_secs = 60 + +# Can also be specified via environment variable `KNOT_ACL_ACCEPTANCE_GRACE_SECS`. +# Default value: 21600 +#acceptance_grace_secs = 21600 + [repo] # Can also be specified via environment variable `KNOT_SCAN_PATH`. # Required! This value must be specified. @@ -494,8 +502,8 @@ [messages.reject] # Can also be specified via environment variable `KNOT_MESSAGES_REJECT_RESERVED_REFS`. -# Default value: "refs/cobs/* and refs/hidden/* are reserved and cannot be pushed" -#reserved_refs = "refs/cobs/* and refs/hidden/* are reserved and cannot be pushed" +# Default value: "refs/cobs/*, refs/cob-checkpoints/*, refs/hidden/* and refs/atproto/* are reserved and can't be pushed" +#reserved_refs = "refs/cobs/*, refs/cob-checkpoints/*, refs/hidden/* and refs/atproto/* are reserved and can't be pushed" # Can also be specified via environment variable `KNOT_MESSAGES_REJECT_COB_CREATE_ONLY`. # Default value: "existing refs/cobs/* object cannot be modified or deleted over the wire" @@ -509,6 +517,14 @@ # Default value: "refs/hidden/* is reserved for server-side fork staging and cannot be pushed" #hidden_reserved = "refs/hidden/* is reserved for server-side fork staging and cannot be pushed" +# Can also be specified via environment variable `KNOT_MESSAGES_REJECT_ATPROTO_RESERVED`. +# Default value: "refs/atproto/* is reserved for the knot's atproto commit chain and can't be pushed" +#atproto_reserved = "refs/atproto/* is reserved for the knot's atproto commit chain and can't be pushed" + +# Can also be specified via environment variable `KNOT_MESSAGES_REJECT_CHECKPOINT_RESERVED`. +# Default value: "refs/cob-checkpoints/* is reserved for knot-written collaborative-object snapshots and can't be pushed" +#checkpoint_reserved = "refs/cob-checkpoints/* is reserved for knot-written collaborative-object snapshots and can't be pushed" + # Can also be specified via environment variable `KNOT_MESSAGES_REJECT_COB_VERIFICATION`. # Default value: "collaborative-object verification failed: {error}" #cob_verification = "collaborative-object verification failed: {error}" -- 2.51.2