From 59be433ea676992f96a12191f2822dd5e71428c3 Mon Sep 17 00:00:00 2001 From: Lewis Date: Mon, 21 Sep 2026 12:45:42 +0300 Subject: [PATCH] knot-xrpc: label tests, merge handler test-only Lewis: May this revision serve well! --- knot2/crates/knot-migrate/src/emit.rs | 70 +- knot2/crates/knot-migrate/tests/migrate.rs | 4 +- knot2/crates/knot-xrpc/src/lib.rs | 42 +- knot2/crates/knot-xrpc/src/merge.rs | 80 +- knot2/crates/knot-xrpc/src/out.rs | 3 +- knot2/crates/knot-xrpc/src/receive.rs | 16 +- knot2/crates/knot-xrpc/src/tests.rs | 1089 ++++++++++++++++++-- 7 files changed, 1156 insertions(+), 148 deletions(-) diff --git a/knot2/crates/knot-migrate/src/emit.rs b/knot2/crates/knot-migrate/src/emit.rs index 0be4f4f53..822c5ebcf 100644 --- a/knot2/crates/knot-migrate/src/emit.rs +++ b/knot2/crates/knot-migrate/src/emit.rs @@ -18,44 +18,47 @@ use crate::mapping::{AdoptRepo, MappedGrant, Mapping}; #[derive(Debug, thiserror::Error)] pub enum EmitError { - #[error("meta-repo bootstrap: {0}")] + #[error("Meta-repo bootstrap: {0}")] Meta(#[from] GitError), - #[error("open adopted repo {repo}: {source}")] + #[error("Open adopted repo {repo}: {source}")] OpenRepo { repo: RepoDid, source: GitError }, - #[error("registry record for {repo} has name {existing:?} and the mapping has {name:?}")] + #[error("Registry record for {repo} has name {existing:?} and the mapping has {name:?}")] RegistryNameChanged { repo: RepoDid, existing: RepoName, name: RepoName, }, #[error("{cob} write failed: {source}")] - Cob { cob: &'static str, source: CobError }, - #[error("registry write failed: {0}")] - Registry(#[from] RegistryError), + Cob { + cob: &'static str, + source: Box, + }, + #[error("Registry write failed: {0}")] + Registry(Box), #[error("{cob} is split across {count} objects")] SplitObject { cob: &'static str, count: usize }, - #[error("write {path}: {source}")] + #[error("Write {path}: {source}")] Io { path: PathBuf, source: std::io::Error, }, - #[error("read host key {path}: {source}")] + #[error("Read host key {path}: {source}")] ReadHostKey { path: PathBuf, source: std::io::Error, }, #[error( - "host key {path} is the public half of a key pair. Point --host-key at the private key, usually the same path without the .pub." + "Host key {path} is the public half of a key pair. Point --host-key at the private key, usually the same path without the .pub." )] PublicHostKey { path: PathBuf }, - #[error("host key {path} doesn't parse as an OpenSSH private key: {source}")] + #[error("Host key {path} doesn't parse as an OpenSSH private key: {source}")] HostKey { path: PathBuf, source: ssh_key::Error, }, - #[error("knot cannot load the passphrase-protected host key {path} unattended")] + #[error("Knot can't load the passphrase-protected host key {path} unattended")] EncryptedHostKey { path: PathBuf }, - #[error("config template has no line for {section}.{key}")] + #[error("Config template has no line for {section}.{key}")] TemplateDrift { section: &'static str, key: &'static str, @@ -191,22 +194,22 @@ fn write_registry( registry.record_of(&repo.did), registry.resolve(&repo.owner, &repo.rkey), ) { - (Some(record), _) if record.owner != repo.owner => { - Err(EmitError::Registry(RegistryError::AlreadyRegistered { + (Some(record), _) if record.owner != repo.owner => Err(EmitError::Registry( + Box::new(RegistryError::AlreadyRegistered { repo: repo.did.clone(), owner: record.owner.clone(), rkey: record.rkey.clone(), field: RegistrationField::Owner, - })) - } - (Some(record), _) if record.rkey != repo.rkey => { - Err(EmitError::Registry(RegistryError::AlreadyRegistered { + }), + )), + (Some(record), _) if record.rkey != repo.rkey => Err(EmitError::Registry( + Box::new(RegistryError::AlreadyRegistered { repo: repo.did.clone(), owner: record.owner.clone(), rkey: record.rkey.clone(), field: RegistrationField::Rkey, - })) - } + }), + )), (Some(record), _) if record.name != repo.name => { Err(EmitError::RegistryNameChanged { repo: repo.did.clone(), @@ -215,11 +218,11 @@ fn write_registry( }) } (None, Some(holder)) if holder != &repo.did => { - Err(EmitError::Registry(RegistryError::RkeyTaken { + Err(EmitError::Registry(Box::new(RegistryError::RkeyTaken { owner: repo.owner.clone(), rkey: repo.rkey.clone(), existing: holder.clone(), - })) + }))) } _ => Ok(()), } @@ -245,7 +248,10 @@ where E::State: Serialize + DeserializeOwned, E::Change: KnotAuthored + Clone, { - let fail = |source: CobError| EmitError::Cob { cob, source }; + let fail = |source: CobError| EmitError::Cob { + cob, + source: Box::new(source), + }; let objects = store.list::().map_err(fail)?; let (object, state, created) = match (objects.as_slice(), items) { (_, []) => return Ok(GrantSetOutcome::default()), @@ -352,7 +358,7 @@ pub fn write_key_archive(path: &Path, repos: &[AdoptRepo]) -> Result<(), EmitErr }) .collect(); let body = zeroize::Zeroizing::new( - serde_json::to_string_pretty(&keys).expect("key archive serializes"), + serde_json::to_string_pretty(&keys).expect("Key archive serializes"), ); write_private(path, body.as_bytes()) } @@ -391,7 +397,7 @@ pub struct Fingerprints { #[derive(Debug, thiserror::Error)] pub enum HostKeyConflict { #[error( - "the migration will replace the different host key at {path}, whose fingerprint {} your users already trust, with {}", + "The migration will replace the different host key at {path}, whose fingerprint {} your users already trust, with {}", .fingerprints.found, .fingerprints.importing )] @@ -399,10 +405,10 @@ pub enum HostKeyConflict { path: PathBuf, fingerprints: Box, }, - #[error("the migration won't write over the key already at the target: {source}")] + #[error("The migration won't write over the key already at the target: {source}")] Unparsable { source: EmitError }, #[error( - "read {path}, which this process can't examine and the old knot might still be serving: {source}" + "Read {path}, which this process can't examine and the old knot might still be serving: {source}" )] Unreadable { path: PathBuf, @@ -410,13 +416,13 @@ pub enum HostKeyConflict { }, #[error("{path} isn't a regular file, so the migration won't write the host key over it")] NotAFile { path: PathBuf }, - #[error("examine {path}: {source}")] + #[error("Examine {path}: {source}")] Unexaminable { path: PathBuf, source: std::io::Error, }, #[error( - "the migration will write the host key over {path}, which this process can't write: {source}" + "The migration will write the host key over {path}, which this process can't write: {source}" )] Unwritable { path: PathBuf, @@ -630,11 +636,11 @@ impl MasterKeyEnv { #[derive(Debug, thiserror::Error)] pub enum MasterKeyError { - #[error("master key env var {0} isn't set")] + #[error("Master key env var {0} isn't set")] Unset(MasterKeyEnv), - #[error("master key env var {0} is set to bytes that aren't text")] + #[error("Master key env var {0} is set to bytes that aren't text")] NotText(MasterKeyEnv), - #[error("master key env var {0} isn't base64")] + #[error("Master key env var {0} isn't base64")] NotBase64(MasterKeyEnv), #[error("{source}, decoded from master key env var {env}")] Weak { diff --git a/knot2/crates/knot-migrate/tests/migrate.rs b/knot2/crates/knot-migrate/tests/migrate.rs index 36d9a143e..6bd83f0b4 100644 --- a/knot2/crates/knot-migrate/tests/migrate.rs +++ b/knot2/crates/knot-migrate/tests/migrate.rs @@ -1606,7 +1606,7 @@ fn a_rehearsal_report_names_what_the_key_at_the_target_costs() { }))); assert!( unwritable.contains( - "host key: the migration will write the host key over \ + "host key: The migration will write the host key over \ /srv/knot/ssh_host_key, which this process can't write" ), "{unwritable}" @@ -1654,7 +1654,7 @@ fn a_rehearsal_states_the_inputs_that_the_real_run_still_needs() { "{rendered}" ); assert!( - rendered.contains("master key env var KNOT_MASTER_KEY_THAT_NOBODY_SETS isn't set"), + rendered.contains("Master key env var KNOT_MASTER_KEY_THAT_NOBODY_SETS isn't set"), "{rendered}" ); } diff --git a/knot2/crates/knot-xrpc/src/lib.rs b/knot2/crates/knot-xrpc/src/lib.rs index c3c475e7c..69508262e 100644 --- a/knot2/crates/knot-xrpc/src/lib.rs +++ b/knot2/crates/knot-xrpc/src/lib.rs @@ -264,7 +264,7 @@ impl XrpcState { headers: &HeaderMap, ) -> Result { let token = push_credential(headers)?; - let method = Nsid::new_owned(PUSH_NSID).expect("push nsid is always a valid nsid"); + let method = Nsid::new_owned(PUSH_NSID).expect("Push NSID is always a valid NSID"); self.atproto .verify_service_jwt_guarded( &token, @@ -293,7 +293,7 @@ impl Method { #[cfg(test)] pub(crate) fn from_nsid(nsid: &str) -> Self { - Self(Nsid::new_owned(nsid).expect("test route nsid parses")) + Self(Nsid::new_owned(nsid).expect("Test route NSID parses")) } } @@ -303,14 +303,14 @@ impl FromRequestParts for Method { async fn from_request_parts(parts: &mut Parts, state: &S) -> Result { let matched = MatchedPath::from_request_parts(parts, state) .await - .map_err(|_| XrpcError::internal("xrpc handler reached without a matched route"))?; + .map_err(|_| XrpcError::internal("XRPC handler reached without a matched route"))?; let nsid = matched .as_str() .strip_prefix("/xrpc/") - .ok_or_else(|| XrpcError::internal("xrpc route paths are prefixed with /xrpc/"))?; + .ok_or_else(|| XrpcError::internal("XRPC route paths are prefixed with /xrpc/"))?; Nsid::new_owned(nsid) .map(Self) - .map_err(|_| XrpcError::internal("route nsid is always a valid nsid")) + .map_err(|_| XrpcError::internal("Route NSID is always a valid NSID")) } } @@ -498,10 +498,10 @@ pub(crate) fn admit_pre_auth( .map_err(|rejection| match rejection { Rejection::RateLimited => XrpcError::rate_limited( knot_resource::Bucket::Peer, - "too many pre-authentication requests, please retry shortly", + "Too many pre-authentication requests, please retry shortly", ), Rejection::Saturated => { - XrpcError::overloaded("knot is shedding pre-authentication load, retry shortly") + XrpcError::overloaded("Knot is shedding pre-authentication load, retry shortly") } }) } @@ -516,10 +516,10 @@ pub(crate) fn admit_write( .map_err(|rejection| match rejection { Rejection::RateLimited => XrpcError::rate_limited( knot_resource::Bucket::ActorWrites, - "this account is writing faster than knot takes writes, please retry shortly", + "This account is writing faster than knot takes writes, please retry shortly", ), Rejection::Saturated => { - XrpcError::overloaded("knot already tracks as many writers as it can take, please retry shortly, thank you for your patience") + XrpcError::overloaded("Knot already tracks as many writers as it can take, please retry shortly, thank you for your patience") } }) } @@ -553,7 +553,7 @@ fn bearer(headers: &HeaderMap) -> Result { .and_then(strip_bearer) .map(str::trim) .and_then(|token| ServiceJwt::new(token).ok()) - .ok_or_else(|| XrpcError::auth_required("missing or malformed Bearer authorization header")) + .ok_or_else(|| XrpcError::auth_required("Missing or malformed Bearer authorization header")) } pub(crate) struct BasicUser(String); @@ -603,20 +603,20 @@ fn push_credential(headers: &HeaderMap) -> Result { let value = headers .get(AUTHORIZATION) .and_then(|value| value.to_str().ok()) - .ok_or_else(|| XrpcError::auth_required("missing authorization header"))?; + .ok_or_else(|| XrpcError::auth_required("Missing authorization header"))?; strip_bearer(value) .map(str::trim) .map(str::to_string) .or_else(|| strip_basic(value)) .and_then(|token| ServiceJwt::new(token).ok()) .ok_or_else(|| { - XrpcError::auth_required("authorization isn't a bearer token or basic credential") + XrpcError::auth_required("Authorization isn't a bearer token or basic credential") }) } pub(crate) fn decode(body: &Bytes) -> Result { serde_json::from_slice(body) - .map_err(|error| XrpcError::invalid_request(format!("invalid request body: {error}"))) + .map_err(|error| XrpcError::invalid_request(format!("Invalid request body: {error}"))) } pub(crate) fn ok_empty() -> Response { @@ -712,7 +712,7 @@ where { match tokio::task::spawn_blocking(task).await { Ok(result) => result, - Err(_) => Err(XrpcError::internal("blocking task failed to complete")), + Err(_) => Err(XrpcError::internal("Blocking task failed to complete")), } } @@ -758,12 +758,12 @@ pub(crate) fn resolve_repo_did( ) -> Result { let raw = segment.as_str(); let trimmed = raw.strip_suffix(".git").unwrap_or(raw); - let did = RepoDid::new(trimmed).map_err(|_| XrpcError::not_found("repository not found"))?; + let did = RepoDid::new(trimmed).map_err(|_| XrpcError::not_found("Repository not found"))?; match state.index.owner_of(&did) { Resolved::Ready(Some(_)) => Ok(did), - Resolved::Ready(None) => Err(XrpcError::not_found("repository not found")), + Resolved::Ready(None) => Err(XrpcError::not_found("Repository not found")), Resolved::Warming => Err(XrpcError::warming( - "registry projection is still warming, retry shortly", + "Registry projection is still warming, retry shortly", )), } } @@ -775,12 +775,12 @@ pub(crate) async fn resolve_repo_named( ) -> Result { let owner = resolve_owner_segment(state, owner).await?; let path = ClonePath::parse(name.as_str()) - .ok_or_else(|| XrpcError::not_found("repository not found"))?; + .ok_or_else(|| XrpcError::not_found("Repository not found"))?; match state.index.resolve_clone_path(&owner, &path) { Resolved::Ready(Some(did)) => Ok(did), - Resolved::Ready(None) => Err(XrpcError::not_found("repository not found")), + Resolved::Ready(None) => Err(XrpcError::not_found("Repository not found")), Resolved::Warming => Err(XrpcError::warming( - "registry projection is still warming, retry shortly", + "Registry projection is still warming, retry shortly", )), } } @@ -789,7 +789,7 @@ async fn resolve_owner_segment( state: &XrpcState, owner: &OwnerSegment, ) -> Result { - let not_found = || XrpcError::not_found("repository not found"); + let not_found = || XrpcError::not_found("Repository not found"); match OwnerRef::parse(owner.as_str()).ok_or_else(not_found)? { OwnerRef::Did(did) => Ok(did), OwnerRef::Handle(handle) => state diff --git a/knot2/crates/knot-xrpc/src/merge.rs b/knot2/crates/knot-xrpc/src/merge.rs index 0fb0c25e7..7add01efc 100644 --- a/knot2/crates/knot-xrpc/src/merge.rs +++ b/knot2/crates/knot-xrpc/src/merge.rs @@ -4,29 +4,39 @@ use axum::Json; use axum::body::Bytes; use axum::extract::State; use axum::response::{IntoResponse, Response}; -use http::{HeaderMap, StatusCode}; +#[cfg(test)] +use http::HeaderMap; +use http::StatusCode; use serde::{Deserialize, Serialize}; use knot_git::{ - ApplyError, ApplyOutcome, Conflict, Identity, NewCommit, ParsedFile, PatchApplier, - PatchParseError, RefUpdate, Repo, StagedChange, Staging, is_format_patch, - parse_mailbox_bounded, parse_patch_bounded, + ApplyError, ApplyOutcome, Conflict, ParsedFile, PatchApplier, PatchParseError, Repo, + StagedChange, is_format_patch, parse_mailbox_bounded, parse_patch_bounded, }; +#[cfg(test)] +use knot_git::{Identity, NewCommit, RefUpdate, Staging}; use knot_runtime::{Clock, HttpTransport}; +#[cfg(test)] +use knot_types::UnixSeconds; use knot_types::{ - AuthorName, BranchName, Email, Oid, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, UnixSeconds, + AuthorName, BranchName, Email, Oid, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, }; -use crate::body::{CommitBody, CommitMessage, Patch}; +use crate::body::Patch; +#[cfg(test)] +use crate::body::{CommitBody, CommitMessage}; use crate::error::XrpcError; +#[cfg(test)] +use crate::ok_empty; use crate::query::RepoArg; use crate::reads::{HostedRepo, open, repo_not_found, require_hosted, resolve_repo}; -use crate::{XrpcState, decode, ok_empty, run_blocking}; +use crate::{XrpcState, decode, run_blocking}; -pub(crate) const MERGE_ROUTE: &str = "/xrpc/sh.tangled.repo.merge"; pub(crate) const MERGE_CHECK_ROUTE: &str = "/xrpc/sh.tangled.repo.mergeCheck"; +const CONFLICT_MESSAGE: &str = "patch can't be applied cleanly"; + +#[cfg(test)] const MERGE_RETRIES: u32 = 3; -const CONFLICT_MESSAGE: &str = "patch cannot be applied cleanly"; #[derive(Debug, Clone, PartialEq, Eq)] pub struct Committer { @@ -34,6 +44,7 @@ pub struct Committer { pub email: Email, } +#[cfg(test)] #[derive(Deserialize)] #[serde(rename_all = "camelCase")] struct MergeInput { @@ -107,6 +118,7 @@ impl MergeCheckOutput { } } +#[cfg(test)] struct CommitSpec { files: Vec, author: Option, @@ -114,12 +126,25 @@ struct CommitSpec { change_id: Option, } +#[cfg(test)] struct MailAuthor { name: AuthorName, email: Email, date: String, } +fn parse_patches(patch: &str, max_bytes: u64) -> Result>, PatchParseError> { + if is_format_patch(patch) { + Ok(parse_mailbox_bounded(patch, max_bytes)? + .into_iter() + .map(|mail| mail.files) + .collect()) + } else { + Ok(vec![parse_patch_bounded(patch, max_bytes)?]) + } +} + +#[cfg(test)] fn parse_specs( patch: &str, message: String, @@ -161,9 +186,10 @@ pub(crate) fn resolve_by_name( fn branch_tip(repo: &Repo, refname: &RefName) -> Result { repo.find_ref(refname)? - .ok_or_else(|| XrpcError::invalid_request("no such branch to merge into")) + .ok_or_else(|| XrpcError::invalid_request("No such branch to merge into")) } +#[cfg(test)] fn mail_time(date: &str, now: UnixSeconds) -> (UnixSeconds, i32) { let trimmed = date.trim(); chrono::DateTime::parse_from_rfc2822(trimmed) @@ -177,6 +203,7 @@ fn mail_time(date: &str, now: UnixSeconds) -> (UnixSeconds, i32) { .unwrap_or((now, 0)) } +#[cfg(test)] fn spec_identities( spec: &CommitSpec, fallback: &Identity, @@ -215,19 +242,19 @@ enum StageStop { fn stage_all( repo: &Repo, tip: Oid, - specs: &[CommitSpec], + files: &[Vec], ) -> Result>, Vec>, ApplyError> { let mut applier = PatchApplier::new(repo, tip); - let staged = specs.iter().try_fold(Vec::new(), |mut clean, spec| { - match applier.step(&spec.files) { + let staged = files + .iter() + .try_fold(Vec::new(), |mut clean, files| match applier.step(files) { Ok(ApplyOutcome::Clean(staged)) => { clean.push(staged); Ok(clean) } Ok(ApplyOutcome::Conflicted(conflicts)) => Err(StageStop::Conflict(conflicts)), Err(error) => Err(StageStop::Apply(error)), - } - }); + }); match staged { Ok(clean) => Ok(Ok(clean)), Err(StageStop::Conflict(conflicts)) => Ok(Err(conflicts)), @@ -235,17 +262,20 @@ fn stage_all( } } +#[cfg(test)] enum MergeAttempt { Done { old: Oid, new: Oid }, Conflicted(Vec), Raced, } +#[cfg(test)] enum Merged { Done { old: Oid, new: Oid }, Conflicted(Vec), } +#[cfg(test)] fn attempt_merge( repo: &Repo, refname: &RefName, @@ -254,10 +284,11 @@ fn attempt_merge( now: UnixSeconds, ) -> Result { if specs.iter().any(|spec| spec.message.trim().is_empty()) { - return Err(XrpcError::invalid_request("commit message is required")); + return Err(XrpcError::invalid_request("Commit message is required")); } let tip = branch_tip(repo, refname)?; - let staged = match stage_all(repo, tip, specs).map_err(XrpcError::from)? { + let files: Vec> = specs.iter().map(|spec| spec.files.clone()).collect(); + let staged = match stage_all(repo, tip, &files).map_err(XrpcError::from)? { Ok(staged) => staged, Err(conflicts) => return Ok(MergeAttempt::Conflicted(conflicts)), }; @@ -311,6 +342,7 @@ fn attempt_merge( } } +#[cfg(test)] fn merge_with_retry( repo: &Repo, refname: &RefName, @@ -329,6 +361,7 @@ fn merge_with_retry( } } +#[cfg(test)] fn merge_conflict(conflicts: &[Conflict]) -> XrpcError { let detail = conflicts .first() @@ -347,6 +380,7 @@ fn merge_conflict(conflicts: &[Conflict]) -> XrpcError { ) } +#[cfg(test)] pub(crate) async fn merge( State(state): State>>, headers: HeaderMap, @@ -360,7 +394,7 @@ pub(crate) async fn merge( &state, &actor, &repo_did, - "only repository owner or a collaborator may merge", + "Only repository owner or a collaborator may merge", ) .await?; let refname = input.branch.head_ref(); @@ -398,13 +432,14 @@ pub(crate) async fn merge( ) .await { - tracing::warn!(repo = %repo_label, %error, "post-receive after merge failed"); + tracing::warn!(repo = %repo_label, %error, "Post-receive after merge failed"); } Ok(ok_empty()) } } } +#[cfg(test)] fn unified_message(input: &MergeInput) -> String { let message = input .commit_message @@ -422,6 +457,7 @@ fn unified_message(input: &MergeInput) -> String { } } +#[cfg(test)] fn unified_author(input: &MergeInput) -> Option { match (input.author_name.as_ref(), input.author_email.as_ref()) { (Some(name), Some(email)) if !name.as_str().is_empty() && !email.as_str().is_empty() => { @@ -446,13 +482,13 @@ pub(crate) async fn merge_check( let max_patch_bytes = state.byte_limits.patch_decompressed.get(); let output = run_blocking(move || { - let specs = match parse_specs(input.patch.as_str(), String::new(), None, max_patch_bytes) { - Ok(specs) => specs, + let files = match parse_patches(input.patch.as_str(), max_patch_bytes) { + Ok(files) => files, Err(error) => return Ok(MergeCheckOutput::broken(error.to_string())), }; let repo = open(&layout, &repo_did)?; let tip = branch_tip(&repo, &refname)?; - match stage_all(&repo, tip, &specs) { + match stage_all(&repo, tip, &files) { Ok(Ok(_)) => Ok(MergeCheckOutput::clean()), Ok(Err(conflicts)) => Ok(MergeCheckOutput::conflicted(conflicts)), Err(ApplyError::TooLarge) => { diff --git a/knot2/crates/knot-xrpc/src/out.rs b/knot2/crates/knot-xrpc/src/out.rs index 129d4b631..7cf853b15 100644 --- a/knot2/crates/knot-xrpc/src/out.rs +++ b/knot2/crates/knot-xrpc/src/out.rs @@ -24,7 +24,7 @@ fn display_opt( fn zoned(seconds: i64, offset_seconds: i32) -> chrono::DateTime { let offset = chrono::FixedOffset::east_opt(offset_seconds) - .unwrap_or_else(|| chrono::FixedOffset::east_opt(0).expect("zero offset is valid")); + .unwrap_or_else(|| chrono::FixedOffset::east_opt(0).expect("Zero offset is valid")); chrono::DateTime::from_timestamp(seconds, 0) .unwrap_or_default() .with_timezone(&offset) @@ -50,6 +50,7 @@ pub(crate) fn datetime_of(at: UnixSeconds) -> Datetime { pub(crate) struct CidString(cid::Cid); impl CidString { + #[cfg(test)] pub(crate) fn cid(&self) -> cid::Cid { self.0 } diff --git a/knot2/crates/knot-xrpc/src/receive.rs b/knot2/crates/knot-xrpc/src/receive.rs index bf01dd06e..0419349f3 100644 --- a/knot2/crates/knot-xrpc/src/receive.rs +++ b/knot2/crates/knot-xrpc/src/receive.rs @@ -61,7 +61,7 @@ async fn serve_advertisement( headers: &HeaderMap, ) -> Response { if let Err(response) = authorized_pusher(state, peer, headers, &repo).await { - return response; + return *response; } let layout = state.layout.clone(); let target = repo.clone(); @@ -116,7 +116,7 @@ async fn serve_receive( ) -> Response { let pusher = match authorized_pusher(state, peer, headers, &repo_did).await { Ok(pusher) => pusher, - Err(response) => return response, + Err(response) => return *response, }; let knot_actor = match state.secrets.public_key(&state.knot_did) { @@ -192,7 +192,7 @@ async fn serve_receive( match landed { Ok(framed) => git_response(RESULT, framed), Err(error) => { - tracing::warn!(repo = repo_did.as_str(), %error, "http receive-pack failed"); + tracing::warn!(repo = repo_did.as_str(), %error, "HTTP receive-pack failed"); error.into_response() } } @@ -207,12 +207,12 @@ fn drain_pack( messages: &HttpMessages, ) -> Result { let mut receiver = PackReceiver::new(scratch, limit, limits, format.kind()) - .map_err(|error| XrpcError::internal(format!("receive staging failed: {error}")))?; + .map_err(|error| XrpcError::internal(format!("Receive staging failed: {error}")))?; let mut buffer = [0u8; READ_CHUNK]; loop { let read = reader .read(&mut buffer) - .map_err(|error| XrpcError::invalid_request(format!("receive read error: {error}")))?; + .map_err(|error| XrpcError::invalid_request(format!("Receive read error: {error}")))?; if read == 0 { break; } @@ -233,11 +233,11 @@ async fn authorized_pusher( peer: SocketPeer, headers: &HeaderMap, repo: &RepoDid, -) -> Result { +) -> Result> { let denied = state.catalog.http.push_denied.text(); authenticate_and_authorize_push(state, peer, headers, repo, &denied) .await - .map_err(challenge) + .map_err(|error| Box::new(challenge(error))) } fn challenge(error: XrpcError) -> Response { @@ -286,7 +286,7 @@ fn map_receive_read(error: knot_pack::ReceiveReadError, messages: &HttpMessages) .line(|ErrorKey::Error| error.to_string()), ), knot_pack::ReceiveReadError::Io(error) => { - XrpcError::internal(format!("receive io error: {error}")) + XrpcError::internal(format!("Receive I/O error: {error}")) } } } diff --git a/knot2/crates/knot-xrpc/src/tests.rs b/knot2/crates/knot-xrpc/src/tests.rs index 99e75c352..8a028840f 100644 --- a/knot2/crates/knot-xrpc/src/tests.rs +++ b/knot2/crates/knot-xrpc/src/tests.rs @@ -390,6 +390,7 @@ fn world_responder( pubkeys: HashMap>, repo_docs: Arc>>, pds_records: Arc>>, + pds_bodies: Arc>>, reads: Arc, ) -> Responder { let doc_url = format!("https://{KNOT_HOST}"); @@ -409,6 +410,13 @@ fn world_responder( .find(|(key, _)| key == "rkey") .map(|(_, value)| value.into_owned()) .unwrap_or_default(); + if let Some(body) = pds_bodies.lock().get(&rkey).cloned() { + return Ok(HttpResponse { + status: StatusCode::OK, + headers: http::HeaderMap::new(), + body, + }); + } let answer = pds_records.lock().get(&rkey).copied(); return Ok(HttpResponse { status: answer.unwrap_or(StatusCode::BAD_REQUEST), @@ -456,6 +464,7 @@ struct World { passerby: Actor, repo_docs: Arc>>, pds_records: Arc>>, + pds_bodies: Arc>>, clock: Arc, reads: Arc, } @@ -596,11 +605,13 @@ impl World { ) }) .collect(); + let pds_bodies: Arc>> = Arc::new(Mutex::new(HashMap::new())); let reads = Arc::new(AtomicUsize::new(0)); let responder = world_responder( pubkeys, Arc::clone(&repo_docs), Arc::clone(&pds_records), + Arc::clone(&pds_bodies), Arc::clone(&reads), ); let (state, clock) = state_from( @@ -629,6 +640,7 @@ impl World { passerby, repo_docs, pds_records, + pds_bodies, clock, reads, } @@ -638,18 +650,31 @@ impl World { State(Arc::clone(&self.state)) } + fn publish_pds_record(&self, rkey: &str) { + self.respond_pds_record(rkey, StatusCode::OK); + } + + fn respond_pds_record(&self, rkey: &str, status: StatusCode) { + self.pds_records.lock().insert(rkey.to_string(), status); + } fn publish_repo_doc(&self, host: &str, multikey: &str) { self.repo_docs .lock() .insert(host.to_string(), multikey.to_string()); } - fn publish_pds_record(&self, rkey: &str) { - self.answer_pds_record(rkey, StatusCode::OK); - } - - fn answer_pds_record(&self, rkey: &str, status: StatusCode) { - self.pds_records.lock().insert(rkey.to_string(), status); + fn publish_remote_def(&self, rkey: &str, value: serde_json::Value) { + self.pds_bodies.lock().insert( + rkey.to_string(), + Bytes::from( + serde_json::to_vec(&json!({ + "uri": format!("at://did/sh.tangled.label.definition/{rkey}"), + "cid": "bafyreihdwdhef9h1o6xtmxpgpkwp0shftvtjhu2suxcbhelperphpriljku", + "value": value + })) + .unwrap(), + ), + ); } fn advance(&self, secs: u64) { @@ -825,7 +850,7 @@ fn plc_stage(fake: Arc>, sec1: Vec) -> Responder { } let did = request.url.path().trim_start_matches('/').to_owned(); let operation = - serde_json::from_slice(request.body.as_ref().expect("post body")).unwrap(); + serde_json::from_slice(request.body.as_ref().expect("Post body")).unwrap(); fake.ops.push((did, operation)); return Ok(reply(StatusCode::OK, Bytes::new())); } @@ -900,7 +925,7 @@ async fn minted_repo(fake: PlcFake) -> SweptRepo { .resolve_repo(&owner, &RepoRkey::new("anemone").unwrap()) { Resolved::Ready(Some(did)) => did, - other => panic!("repo wasn't registered despite the queue: {other:?}"), + other => panic!("Repo wasn't registered despite the queue: {other:?}"), }; SweptRepo { _dir, @@ -978,7 +1003,7 @@ async fn create_repo_helper(world: &World, name: &str) -> RepoDid { ); match resolve(world, name) { Resolved::Ready(Some(did)) => did, - other => panic!("repo {name} wasn't registered: {other:?}"), + other => panic!("Repo {name} wasn't registered: {other:?}"), } } @@ -1101,7 +1126,7 @@ async fn a_repository_takes_the_knots_default_policy_and_its_owner_or_the_operat assert_eq!( set_policy_as(&world, &world.admin, &repo, ContributionPolicy::Anyone).await, StatusCode::OK, - "operator answers for text this knot serves, policy is the operator's to move" + "Operator answers for text this knot serves; policy is the operator's to move" ); assert_eq!( policy_now(), @@ -1145,7 +1170,7 @@ async fn a_second_write_inside_one_account_s_burst_answers_429_named_by_bucket() ) .await, StatusCode::OK, - "the limit is the writer's, leaving another account's token untouched" + "The limit is the writer's, leaving another account's token untouched" ); } @@ -1359,7 +1384,7 @@ async fn delete_issue( fn rkey_of(uri: &AtUri) -> String { uri.rkey() - .expect("an at-uri of a record has its record key") + .expect("An at-uri of a record has its record key") .as_ref() .to_owned() } @@ -1373,7 +1398,7 @@ fn record_of( repo: crate::atproto::RepoSubject::Did(repo.clone()), collection: knot_types::RecordCollection::new( uri.collection() - .expect("an at-uri of a record has its collection") + .expect("An at-uri of a record has its collection") .as_str(), ) .unwrap(), @@ -1492,7 +1517,7 @@ fn message_contains(body: &serde_json::Value, fragments: &[&str]) { let message = body["message"].as_str().unwrap_or_default(); assert!( fragments.iter().all(|fragment| message.contains(fragment)), - "the answer oughtn't miss any of {fragments:?}: {body}" + "The answer oughtn't miss any of {fragments:?}: {body}" ); } @@ -1537,7 +1562,7 @@ async fn stranger_opens_issues_until_policy_closes_and_reads_collaborators_remed .as_str() .unwrap() .starts_with("bafyrei"), - "the answer oughtta include the signed commit stock clients expect" + "The answer oughtta include the signed commit stock clients expect" ); assert!(!written["commit"]["rev"].as_str().unwrap().is_empty()); assert_eq!(written["validationStatus"], "valid"); @@ -1601,24 +1626,26 @@ async fn stock_routes_dispatch_on_collection_and_reject_foreign_collections() { ); message_contains(&minted, &["mints the record key"]); - futures::stream::iter([Some("3lubrptx57d22".to_owned()), None]).for_each(|rkey| { - let (world, repo) = (&world, &repo); - async move { - let rejected = failed( - create_as( - world, - &world.stranger, - repo, - "com.example.gripe", - rkey, - retype(issue_record(repo, "no", ""), "com.example.gripe"), - ) - .await, - StatusCode::BAD_REQUEST, - ); - message_contains(&rejected, &["isn't a collection this knot takes"]); - } - }); + futures::stream::iter([Some("3lubrptx57d22".to_owned()), None]) + .for_each(|rkey| { + let (world, repo) = (&world, &repo); + async move { + let rejected = failed( + create_as( + world, + &world.stranger, + repo, + "com.example.gripe", + rkey, + retype(issue_record(repo, "no", ""), "com.example.gripe"), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&rejected, &["isn't a collection this knot takes"]); + } + }) + .await; let wrong_type = failed( create_as( @@ -1695,7 +1722,7 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { ); assert!( replayed_blocks(&world).contains(&v1.cid()), - "before the delete, replay off the ring oughtta serve the leaf" + "Before the delete, replay off the ring oughtta serve the leaf" ); let second = ok(edit_issue(&world, &world.stranger, &repo, &uri, &v1, "kelp, again").await); @@ -1731,7 +1758,7 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { assert_eq!( same["cid"], v2.to_string(), - "without writing, an edit that doesn't change anything oughtta answer the standing version" + "Without writing, an edit that doesn't change anything oughtta answer the standing version" ); assert_eq!(same["commit"], second["commit"]); @@ -1751,13 +1778,13 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { assert_eq!( record["value"]["x-tngl-editor"], world.member.did.as_str(), - "the version ought to answer with the maintainer who wrote it" + "The version ought to answer with the maintainer who wrote it" ); let requested = ok(record_at(&world, &repo, &uri, Some(&v3)).await); assert_eq!( requested["cid"], v3.to_string(), - "a cid matching the standing version should serve it" + "A CID matching the standing version should serve it" ); let superseded = failed( record_at(&world, &repo, &uri, Some(&v2)).await, @@ -1765,7 +1792,7 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { ); assert_eq!( superseded["error"], "RecordNotFound", - "a cid from a superseded version must answer as a record that left" + "A CID from a superseded version must answer as a record that left" ); failed( @@ -1783,7 +1810,7 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { let deleted = ok(delete_issue(&world, &world.member, &repo, &uri, &v3).await); assert!( deleted.get("uri").is_none() && deleted.get("cid").is_none(), - "the tombstone answer oughtta be the stock commit alone" + "Tombstone answer oughtta be the stock commit alone" ); assert!( deleted["commit"]["cid"] @@ -1798,7 +1825,7 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { assert!(records_in(&world, &repo, ISSUE).await.is_empty()); assert!( block_served(&world, &repo, &v1).await, - "the block tree still keeps the prose git can't delete" + "The block tree still keeps the prose git can't delete" ); let versions = versions_page(&world, &uri, None).await; let chain = versions["versions"].as_array().unwrap(); @@ -1811,7 +1838,7 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { assert_eq!( chain[0]["editor"], world.stranger.did.as_str(), - "the chain lists the opener first" + "Chain lists the opener first" ); assert_eq!(chain[1]["cid"], v2.to_string()); assert_eq!(chain[1]["editor"], world.stranger.did.as_str()); @@ -1841,7 +1868,7 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { assert_eq!(again["message"], format!("issue at {uri} was deleted")); assert!( replayed_blocks(&world).contains(&v1.cid()), - "replay serves the stored frame, whose create commit still includes the leaf" + "Replay serves the stored frame, whose create commit still includes the leaf" ); } @@ -1908,7 +1935,7 @@ async fn state_transition_belongs_to_author_or_maintainer_and_repeat_answers_409 ); assert_eq!( rejected["message"], - "state changes belong to the issue's author or a maintainer" + "State changes belong to the issue's author or a maintainer" ); } @@ -1998,14 +2025,14 @@ async fn social_budget_429s_over_serving_bytes_and_edits_or_deletes_free_them() assert_eq!(rejected["error"], "RateLimitExceeded"); assert_eq!( rejected["bucket"], "repo-social", - "the edited version still occupies the budget until it is freed" + "The edited version still occupies the budget until it is freed" ); let message = rejected["message"].as_str().unwrap(); assert!( message.contains("serving record text") && message.contains("frees the bytes of the version it replaces") && !message.contains("pool"), - "the 429 oughtta explain how to free the bytes: {rejected}" + "The 429 oughtta explain how to free the bytes: {rejected}" ); ok(delete_issue(&world, &world.member, &repo, &uri, &v2).await); @@ -2015,7 +2042,7 @@ async fn social_budget_429s_over_serving_bytes_and_edits_or_deletes_free_them() .as_str() .unwrap() .contains("/sh.tangled.repo.issue/"), - "erasure frees the bytes of the version it retracts" + "Erasure frees the bytes of the version it retracts" ); } #[tokio::test] @@ -2025,7 +2052,7 @@ async fn admission_screens_repo_creation() { assert_eq!( create_status(&closed, &closed.stranger, make()).await, StatusCode::FORBIDDEN, - "a closed knot denies a stranger" + "A closed knot denies a stranger" ); let open = World::open(); @@ -2041,7 +2068,7 @@ async fn admission_screens_repo_creation() { ), Resolved::Ready(Some(_)) ), - "an open knot registers the stranger's repo without membership" + "An open knot registers the stranger's repo without membership" ); } @@ -3779,7 +3806,7 @@ mod merge_endpoints { assert_eq!(conflicted["is_conflicted"], true); assert_eq!(conflicted["conflicts"][0]["filename"], "reef.txt"); assert_eq!(conflicted["conflicts"][0]["reason"], "patch doesn't apply"); - assert_eq!(conflicted["message"], "patch cannot be applied cleanly"); + assert_eq!(conflicted["message"], "patch can't be applied cleanly"); let broken = json_of( crate::merge::merge_check(world.state(), input("hello world\n")) @@ -3815,7 +3842,7 @@ mod merge_endpoints { ); let repo_did = match resolve(&world, "3mjmslfzgwb22") { Resolved::Ready(Some(did)) => did, - other => panic!("repo wasn't registered: {other:?}"), + other => panic!("Repo wasn't registered: {other:?}"), }; let base = seed_main(&world, &repo_did, &[("reef.txt", "old line\n")]); @@ -3957,7 +3984,7 @@ mod fork_endpoints { .resolve_repo(&member_did(), &RepoRkey::new(rkey).unwrap()) { Resolved::Ready(Some(did)) => did, - other => panic!("fork {rkey} wasn't registered: {other:?}"), + other => panic!("Fork {rkey} wasn't registered: {other:?}"), } } @@ -4249,7 +4276,7 @@ mod fork_endpoints { assert_eq!( track_hidden(&setup.world, &setup.fork_did, "feature", "main").await, StatusCode::NOT_FOUND, - "the knot reports not found for a file origin with an unknown repo did" + "The knot reports not found for a file origin with an unknown repo DID" ); fork.set_origin_url(&OriginUrl::new("ssh://knot.nel.pet/did:plc:whelk/ghost")) @@ -4634,7 +4661,7 @@ mod rolls { } fn answered(&self, status: StatusCode) { - self.world.answer_pds_record(&self.named(), status); + self.world.respond_pds_record(&self.named(), status); } fn written_by(&self, author: &AccountDid) -> serde_json::Value { @@ -4976,7 +5003,7 @@ mod rolls { members_changes(world) .into_iter() .find(|change| change.record.bytes().is_some()) - .expect("a change carries its materialized record") + .expect("A change includes its materialized record") } fn roll_history( @@ -5019,7 +5046,7 @@ mod rolls { let MembersChange::Invite(invite) = MembersChange::decode(with_bytes.payload()).unwrap() else { - panic!("invite is change that stores record"); + panic!("Invite is change that stores record"); }; assert_eq!( with_bytes.record.bytes().unwrap(), @@ -5055,7 +5082,7 @@ mod rolls { let MembersChange::Accept(accept) = MembersChange::decode(with_bytes.payload()).unwrap() else { - panic!("acceptance is change that stores record"); + panic!("Acceptance is change that stores record"); }; assert_eq!(accept.subject, world.member.did); } @@ -5160,7 +5187,7 @@ mod rolls { knot_index::Resolved::Ready(entries) => { entries.iter().any(|(did, _)| did.as_str() == subject) } - knot_index::Resolved::Warming => panic!("the members projection is warming"), + knot_index::Resolved::Warming => panic!("The members projection is warming"), }; let response = post(ADD_MEMBER, "did:plc:limpet").await.unwrap(); assert_eq!( @@ -5173,7 +5200,7 @@ mod rolls { let meta = Repo::open(world.state.meta_path.clone()).unwrap(); let old = knot_record::chain::tip(&meta) .unwrap() - .expect("the healthy invite materialized the chain") + .expect("The healthy invite materialized the chain") .git; let junk = meta.git().write_blob(b"not a commit").unwrap().detach(); meta.update_ref(&knot_git::RefUpdate::Update { @@ -5229,7 +5256,7 @@ mod rolls { let meta = Repo::open(world.state.meta_path.clone()).unwrap(); let store = CobStore::new(&meta); let [object] = store.list::().unwrap()[..] else { - panic!("one roll object exists"); + panic!("One roll object exists"); }; enable_at(&world, object).await.unwrap(); let tip = knot_record::chain::tip(&meta).unwrap().unwrap(); @@ -5240,7 +5267,7 @@ mod rolls { crate::materialize::enable_surfaces(&world.state); let frames = framed_events(&world); let [event] = &frames[..] else { - panic!("a single sync frame covers the frameless tip"); + panic!("A single sync frame covers the frameless tip"); }; let mut read = &event.frame[..]; let header: serde_json::Value = @@ -5679,7 +5706,7 @@ async fn comment_rejects_subjects_this_knot_doesnt_serve() { .await, StatusCode::BAD_REQUEST, ); - message_contains(&stale_subject, &["strongRef cid is", "now has version"]); + message_contains(&stale_subject, &["strongRef CID is", "now has version"]); assert!( stale_subject["message"] .as_str() @@ -6095,7 +6122,7 @@ async fn comment_subject_and_reply_are_fixed_at_open() { .await, StatusCode::BAD_REQUEST, ); - message_contains(&stale, &["reply's strongRef cid is", "now has version"]); + message_contains(&stale, &["reply's strongRef CID is", "now has version"]); assert!( stale["message"] .as_str() @@ -6249,3 +6276,941 @@ async fn reaction_dedup_wont_answer_when_policy_rejects() { "the rejected second react shouldn't create a duplicate" ); } + +const LABEL_DEF: &str = "sh.tangled.label.definition"; +const LABEL_OP: &str = "sh.tangled.label.op"; + +fn def_body(name: &str, scope: &[&str], value_type: serde_json::Value) -> serde_json::Value { + json!({ + "$type": LABEL_DEF, + "name": name, + "scope": scope, + "multiple": false, + "createdAt": "2026-09-01T00:00:00Z", + "valueType": value_type, + }) +} + +fn null_type() -> serde_json::Value { + json!({"type": "null", "format": "any"}) +} + +fn assignee_type() -> serde_json::Value { + json!({"type": "string", "format": "did"}) +} + +fn op_body( + subject: &AtUri, + add: serde_json::Value, + delete: serde_json::Value, +) -> serde_json::Value { + json!({ + "$type": LABEL_OP, + "subject": subject.as_str(), + "add": add, + "delete": delete, + "performedAt": "2026-09-01T00:00:00Z", + }) +} + +fn operand(key: &str, value: &str) -> serde_json::Value { + json!({ "key": key, "value": value }) +} + +async fn define_label( + world: &World, + repo: &RepoDid, + slug: &str, + name: &str, + scope: &[&str], + value_type: serde_json::Value, + multiple: bool, +) -> (AtUri, crate::out::CidString) { + let mut body = def_body(name, scope, value_type); + body["multiple"] = json!(multiple); + let written = ok(create_as( + world, + &world.member, + repo, + LABEL_DEF, + Some(slug.to_owned()), + body, + ) + .await); + (uri_of(&written), cid_of(&written)) +} + +async fn wontfix(world: &World, repo: &RepoDid) -> (AtUri, crate::out::CidString) { + define_label( + world, + repo, + "wontfix", + "wontfix", + &["sh.tangled.repo.issue"], + null_type(), + false, + ) + .await +} + +fn remote_triage(world: &World, rkey: &str) -> String { + world.publish_remote_def( + rkey, + json!({ + "name": rkey, + "scope": ["sh.tangled.repo.issue"], + "multiple": true, + "createdAt": "2026-09-01T00:00:00Z", + "valueType": {"type": "string", "format": "any"} + }), + ); + format!( + "at://{}/{LABEL_DEF}/{rkey}", + world.stranger.did.as_str() + ) +} + +#[tokio::test] +async fn label_defs_belong_to_the_owner_and_go_by_their_slug() { + let (world, repo) = issue_world().await; + + let refused = failed( + create_as( + &world, + &world.stranger, + &repo, + LABEL_DEF, + Some("wontfix".to_owned()), + def_body("wontfix", &["sh.tangled.repo.issue"], null_type()), + ) + .await, + StatusCode::FORBIDDEN, + ); + message_contains( + &refused, + &["label defs belong to the repository's owner"], + ); + + let refused_policy = failed( + create_as( + &world, + &world.admin, + &repo, + LABEL_DEF, + Some("wontfix".to_owned()), + def_body("wontfix", &["sh.tangled.repo.issue"], null_type()), + ) + .await, + StatusCode::FORBIDDEN, + ); + message_contains( + &refused_policy, + &["label defs belong to the repository's owner"], + ); + + let no_slug = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_DEF, + None, + def_body("wontfix", &["sh.tangled.repo.issue"], null_type()), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&no_slug, &["slug"]); + + let (uri, cid) = wontfix(&world, &repo).await; + assert_eq!( + uri.as_str(), + format!("at://{repo}/{LABEL_DEF}/wontfix"), + "the def's record key is the slug the owner sent" + ); + assert_eq!( + records_in(&world, &repo, LABEL_DEF).await, + vec![uri.as_str().to_owned()] + ); + let served = ok(record_at(&world, &repo, &uri, None).await); + assert_eq!(served["value"]["name"], "wontfix"); + assert_eq!(served["value"]["x-tngl-editor"], world.member.did.as_str()); + + let duplicate = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_DEF, + Some("wontfix".to_owned()), + def_body("wontfix2", &["sh.tangled.repo.issue"], null_type()), + ) + .await, + StatusCode::CONFLICT, + ); + message_contains(&duplicate, &["already stands"]); + + let edited = ok(put_as( + &world, + &world.member, + &repo, + LABEL_DEF, + "wontfix".to_owned(), + json!({ "swapRecord": cid }), + def_body("never_happening", &["sh.tangled.repo.issue"], null_type()), + ) + .await); + assert_eq!(uri_of(&edited).as_str(), uri.as_str()); + let renamed = ok(record_at(&world, &repo, &uri, None).await); + assert_eq!(renamed["value"]["name"], "never_happening"); + + let outsider = failed( + put_as( + &world, + &world.stranger, + &repo, + LABEL_DEF, + "wontfix".to_owned(), + json!({ "swapRecord": cid_of(&edited) }), + def_body("taken_over", &["sh.tangled.repo.issue"], null_type()), + ) + .await, + StatusCode::FORBIDDEN, + ); + message_contains( + &outsider, + &["label defs belong to the repository's owner alone"], + ); + + ok(delete_as( + &world, + &world.member, + &repo, + LABEL_DEF, + "wontfix".to_owned(), + json!({ "swapRecord": cid_of(&edited) }), + ) + .await); + assert_eq!( + records_in(&world, &repo, LABEL_DEF).await, + Vec::::new(), + "the erased def stops serving, the way every social record does" + ); +} + +#[tokio::test] +async fn label_ops_require_triage_and_vet_their_operands() { + let (world, repo) = issue_world().await; + let (issue, _) = opened(&world, &world.stranger, &repo, "kelp").await; + let (assignee, _) = define_label( + &world, + &repo, + "assignee", + "assignee", + &["sh.tangled.repo.issue", "sh.tangled.repo.pull"], + assignee_type(), + true, + ) + .await; + let (wontfix, _) = wontfix(&world, &repo).await; + + let refused = failed( + create_as( + &world, + &world.stranger, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([operand(wontfix.as_str(), "null")]), + json!([]), + ), + ) + .await, + StatusCode::FORBIDDEN, + ); + message_contains( + &refused, + &[ + "labels belong to the repository's owner and collaborators", + "add you as a collaborator", + ], + ); + + let wrong_value = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([operand(assignee.as_str(), "squid")]), + json!([]), + ), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&wrong_value, &["did-typed label's value is a DID"]); + + let missing_subject = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &AtUri::new_owned(format!("at://{repo}/{ISSUE}/3lubrptx5zzzz")).unwrap(), + json!([operand(wontfix.as_str(), "null")]), + json!([]), + ), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&missing_subject, &["isn't a record on this knot"]); + + let gone_def = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([operand( + &format!("at://{repo}/{LABEL_DEF}/absent"), + "null" + )]), + json!([]), + ), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&gone_def, &["doesn't have a label def at"]); + + let first = ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([ + operand(assignee.as_str(), "did:plc:abalone"), + operand(wontfix.as_str(), "null"), + ]), + json!([]), + ), + ) + .await); + let op_uri = uri_of(&first); + assert_eq!( + op_uri.as_str(), + format!("at://{repo}/{LABEL_OP}/{}", rkey_of(&op_uri)), + "the op's record key is the minted TID, like every social record but the def" + ); + + let second = ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([operand(assignee.as_str(), "did:plc:conch")]), + json!([]), + ), + ) + .await); + let third = ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([]), + json!([operand(assignee.as_str(), "did:plc:abalone")]), + ), + ) + .await); + assert_ne!(cid_of(&second), cid_of(&third)); + + let labels = labels_on(&world, &repo, &issue); + assert_eq!( + labels, + std::collections::BTreeMap::from([ + ( + format!("at://{repo}/{LABEL_DEF}/assignee"), + vec!["did:plc:conch".to_owned()] + ), + ( + format!("at://{repo}/{LABEL_DEF}/wontfix"), + vec!["null".to_owned()] + ), + ]), + "multiple accumulates then a delete removes one value, and single keeps one standing value" + ); + + let not_edited = failed( + put_as( + &world, + &world.member, + &repo, + LABEL_OP, + rkey_of(&op_uri), + json!({ "swapRecord": cid_of(&first) }), + op_body( + &issue, + json!([operand(wontfix.as_str(), "null")]), + json!([]), + ), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(¬_edited, &["putRecord can't edit"]); + + ok(delete_as( + &world, + &world.member, + &repo, + LABEL_OP, + rkey_of(&op_uri), + json!({ "swapRecord": cid_of(&first) }), + ) + .await); + let after_erase = labels_on(&world, &repo, &issue); + assert_eq!( + after_erase, + std::collections::BTreeMap::from([( + format!("at://{repo}/{LABEL_DEF}/assignee"), + vec!["did:plc:conch".to_owned()] + ),]), + "erasing the first op removes the two applications, and the standing set is what the \ + surviving ops fold to" + ); + + let mut echoed = op_body( + &issue, + json!([operand(wontfix.as_str(), "null")]), + json!([]), + ); + echoed["x-tngl-editor"] = json!("did:plc:system"); + let response = create_as(&world, &world.member, &repo, LABEL_OP, None, echoed).await; + assert_eq!( + response.0, + StatusCode::OK, + "the server stamps the editor, and a client echoing the field back is heard and \ + overwritten: {response:?}" + ); +} + +#[tokio::test] +async fn remote_def_reads_through_getrecord_and_an_unreadable_def_rejects_the_op() { + let (world, repo) = issue_world().await; + let (issue, _) = opened(&world, &world.stranger, &repo, "kelp").await; + + let remote = remote_triage(&world, "triage"); + + ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([ + operand(&remote, "needs-a-look"), + operand(&remote, "next-up") + ]), + json!([]), + ), + ) + .await); + let reads_after_first = world.record_reads(); + ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body(&issue, json!([]), json!([operand(&remote, "needs-a-look")])), + ) + .await); + assert!( + world.record_reads() <= reads_after_first + 1, + "the def caches the way public keys do, so a second op costs at most one more \ + getRecord" + ); + + let labels = labels_on(&world, &repo, &issue); + assert_eq!( + labels, + std::collections::BTreeMap::from([(remote.clone(), vec!["next-up".to_owned()])]) + ); + + let unreadable = format!( + "at://{}/{LABEL_DEF}/gone", + world.stranger.did.as_str() + ); + let refused = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body(&issue, json!([operand(&unreadable, "null")]), json!([])), + ) + .await, + StatusCode::BAD_GATEWAY, + ); + message_contains(&refused, &["didn't serve", "won't store"]); + assert_eq!( + records_in(&world, &repo, LABEL_OP).await.len(), + 2, + "the refused op doesn't leave a record behind" + ); +} + +#[tokio::test] +async fn refused_label_ops_leave_the_defs_unread() { + let (world, repo) = issue_world().await; + let (issue, _) = opened(&world, &world.stranger, &repo, "kelp").await; + + let remote = remote_triage(&world, "triage"); + let before = world.record_reads(); + let refused = failed( + create_as( + &world, + &world.stranger, + &repo, + LABEL_OP, + None, + op_body(&issue, json!([operand(&remote, "needs-a-look")]), json!([])), + ) + .await, + StatusCode::FORBIDDEN, + ); + message_contains( + &refused, + &["labels belong to the repository's owner and collaborators"], + ); + assert_eq!( + world.record_reads(), + before, + "the triage rule rejects the write before any def is read, so a stranger's write \ + costs zero getRecord calls" + ); + + let bogus = AtUri::new_owned(format!("at://{repo}/{ISSUE}/3lubrptx5zzzz")).unwrap(); + let before = world.record_reads(); + let refused = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body(&bogus, json!([operand(&remote, "needs-a-look")]), json!([])), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&refused, &["isn't a record on this knot"]); + assert_eq!( + world.record_reads(), + before, + "the knot rejects an unknown subject before fetching any def" + ); +} + +#[tokio::test] +async fn an_op_citing_an_erased_def_is_refused() { + let (world, repo) = issue_world().await; + let (issue, _) = opened(&world, &world.stranger, &repo, "kelp").await; + let (wontfix, cid) = wontfix(&world, &repo).await; + ok(delete_as( + &world, + &world.member, + &repo, + LABEL_DEF, + "wontfix".to_owned(), + json!({ "swapRecord": cid }), + ) + .await); + let refused = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([operand(wontfix.as_str(), "null")]), + json!([]), + ), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&refused, &["was deleted from this repository"]); + assert_eq!( + records_in(&world, &repo, LABEL_OP).await, + Vec::::new(), + "an op citing a deleted def is rejected without a record" + ); +} + +#[tokio::test] +async fn remote_def_with_unknown_fields_still_validates() { + let (world, repo) = issue_world().await; + let (issue, _) = opened(&world, &world.stranger, &repo, "kelp").await; + + let triage = format!( + "at://{}/{LABEL_DEF}/triage", + world.stranger.did.as_str() + ); + let blocker = format!( + "at://{}/{LABEL_DEF}/blocker", + world.stranger.did.as_str() + ); + world.publish_remote_def( + "triage", + json!({ + "name": "triage", + "scope": ["sh.tangled.repo.issue"], + "multiple": false, + "createdAt": "2026-09-01T00:00:00Z", + "futureField": {"note": "a lexicon revision this knot doesn't know"}, + "valueType": {"type": "string", "format": "any", "anotherFutureField": true} + }), + ); + world.publish_remote_def( + "blocker", + json!({ + "name": "blocker", + "scope": ["sh.tangled.repo.issue"], + "multiple": false, + "createdAt": "2026-09-01T00:00:00Z", + "valueType": {"type": "null", "format": "any"} + }), + ); + ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([operand(&triage, "needs-a-look"), operand(&blocker, "null")]), + json!([]), + ), + ) + .await); + let labels = labels_on(&world, &repo, &issue); + assert_eq!( + labels, + std::collections::BTreeMap::from([ + (triage, vec!["needs-a-look".to_owned()]), + (blocker, vec!["null".to_owned()]), + ]), + "a PDS def with fields this knot doesn't know still validates the op, the way \ + every atproto consumer tolerates a lexicon revision that the consumer hasn't met" + ); +} + +fn labels_on( + world: &World, + repo: &RepoDid, + subject: &AtUri, +) -> std::collections::BTreeMap> { + match world.state.index.social_labels_of( + repo, + &match knot_types::RecordAddress::parse_at_uri(subject) { + Ok((_, address)) => address, + Err(_) => panic!("The subject URI parses"), + }, + ) { + knot_index::Resolved::Ready(labels) => labels + .into_iter() + .map(|(def, values)| { + ( + def.to_string(), + values + .into_iter() + .map(|value| value.as_str().to_owned()) + .collect(), + ) + }) + .collect(), + knot_index::Resolved::Warming => panic!("The projection was warmed by the writes"), + } +} + +#[tokio::test] +async fn null_def_with_a_declared_format_is_refused_locally_and_remotely() { + let (world, repo) = issue_world().await; + let (issue, _) = opened(&world, &world.stranger, &repo, "kelp").await; + + let refused = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_DEF, + Some("bent".to_owned()), + def_body( + "bent", + &["sh.tangled.repo.issue"], + json!({"type": "null", "format": "did"}), + ), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&refused, &["format"]); + + let remote = format!( + "at://{}/{LABEL_DEF}/bent", + world.stranger.did.as_str() + ); + world.publish_remote_def( + "bent", + json!({ + "name": "bent", + "scope": ["sh.tangled.repo.issue"], + "multiple": false, + "createdAt": "2026-09-01T00:00:00Z", + "valueType": {"type": "null", "format": "did"} + }), + ); + let refused_remote = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body(&issue, json!([operand(&remote, "null")]), json!([])), + ) + .await, + StatusCode::BAD_GATEWAY, + ); + message_contains(&refused_remote, &["doesn't parse"]); + assert_eq!( + records_in(&world, &repo, LABEL_OP).await, + Vec::::new(), + "the knot doesn't store an op citing a malformed def" + ); +} + +#[tokio::test] +async fn erased_def_slug_stays_retired() { + let (world, repo) = issue_world().await; + let (_, cid) = wontfix(&world, &repo).await; + ok(delete_as( + &world, + &world.member, + &repo, + LABEL_DEF, + "wontfix".to_owned(), + json!({ "swapRecord": cid }), + ) + .await); + let refused = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_DEF, + Some("wontfix".to_owned()), + def_body("wontfix", &["sh.tangled.repo.issue"], null_type()), + ) + .await, + StatusCode::CONFLICT, + ); + message_contains(&refused, &["was erased at", "slug stays retired"]); + assert_eq!( + records_in(&world, &repo, LABEL_DEF).await, + Vec::::new(), + "the erased def doesn't serve, and the slug stays retired: one record address is \ + one object for life" + ); +} + +#[tokio::test] +async fn an_out_of_scope_operand_rejects_the_op() { + let (world, repo) = issue_world().await; + let (issue, cid) = opened(&world, &world.stranger, &repo, "kelp").await; + let comment = ok(create_as( + &world, + &world.member, + &repo, + COMMENT, + None, + comment_record(&issue, &cid, "kelp"), + ) + .await); + let (issue_only, _) = wontfix(&world, &repo).await; + let refused = failed( + create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &uri_of(&comment), + json!([operand(issue_only.as_str(), "null")]), + json!([]), + ), + ) + .await, + StatusCode::BAD_REQUEST, + ); + message_contains(&refused, &["outside this def's scope"]); +} + +#[tokio::test] +async fn a_def_flip_leaves_each_op_under_the_flag_in_force_at_write_time() { + let (world, repo) = issue_world().await; + let (second, _) = opened(&world, &world.stranger, &repo, "urchin").await; + let text_type = json!({"type": "string", "format": "any"}); + let (reviewer, cid) = define_label( + &world, + &repo, + "reviewer", + "reviewer", + &["sh.tangled.repo.issue"], + text_type.clone(), + true, + ) + .await; + ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &second, + json!([ + operand(reviewer.as_str(), "conch"), + operand(reviewer.as_str(), "limpet") + ]), + json!([]), + ), + ) + .await); + assert_eq!( + labels_on(&world, &repo, &second), + std::collections::BTreeMap::from([( + reviewer.as_str().to_owned(), + vec!["conch".to_owned(), "limpet".to_owned()] + )]), + "under multiple the two adds accumulate" + ); + + ok(put_as( + &world, + &world.member, + &repo, + LABEL_DEF, + "reviewer".to_owned(), + json!({ "swapRecord": cid }), + def_body("reviewer", &["sh.tangled.repo.issue"], text_type), + ) + .await); + ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &second, + json!([]), + json!([operand(reviewer.as_str(), "conch")]), + ), + ) + .await); + assert_eq!( + labels_on(&world, &repo, &second), + std::collections::BTreeMap::new(), + "a single delete written after the flip spells the standing first value and clears \ + limpet too, which was added under the multiple the def doesn't declare anymore" + ); +} + +#[tokio::test] +async fn a_def_on_another_repo_reads_from_that_repo_s_store() { + let (world, repo) = issue_world().await; + let (issue, _) = opened(&world, &world.stranger, &repo, "kelp").await; + let other = create_repo_helper(&world, "barnacle").await; + let (reviewer, _) = define_label( + &world, + &other, + "reviewer", + "reviewer", + &["sh.tangled.repo.issue"], + json!({"type": "string", "format": "any"}), + false, + ) + .await; + + let before = world.record_reads(); + ok(create_as( + &world, + &world.member, + &repo, + LABEL_OP, + None, + op_body( + &issue, + json!([operand(reviewer.as_str(), "nel")]), + json!([]), + ), + ) + .await); + assert_eq!( + world.record_reads(), + before, + "a def this knot hosts is read from that repo's store, without a getRecord \ + against the knot's public endpoint" + ); +} + +#[tokio::test] +async fn a_stranger_probing_a_taken_slug_hears_ownership_first() { + let (world, repo) = issue_world().await; + wontfix(&world, &repo).await; + let response = create_as( + &world, + &world.stranger, + &repo, + LABEL_DEF, + Some("wontfix".to_owned()), + def_body("wontfix", &["sh.tangled.repo.issue"], null_type()), + ) + .await; + assert_eq!(response.0, StatusCode::FORBIDDEN); + message_contains( + &response.1, + &["label defs belong to the repository's owner alone"], + ); +} -- 2.51.2