diff --git a/CHANGELOG.md b/CHANGELOG.md index 5604bc1..3f99c49 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -57,8 +57,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 promotion after ref advertisement, while retaining authorization for changed targets and Git's protection against concurrent ref updates. Validate PR refs independently so mixed pushes do not require state authorization for unchanged - branches. Complete fully satisfied pushes containing deletions with a - verify-only ref transaction, without risking deletion of a recreated ref. + branches. Verify already-completed deletions without writing refs, including + when mixed with new updates. Preserve per-ref outcomes for ordinary pushes + and hold verified refs locked across atomic pushes. Keep rejected refs from + being immediately overwritten by post-push state promotion. - Reject PR and PR-update events missing commit metadata as terminal invalid events, avoiding repeated recovery attempts for immutable malformed data. - Omit raw rejected relay payloads from SDK diagnostics while retaining URLs and diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index c7b36d0..e6fc3d4 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -448,14 +448,21 @@ pub struct Purgatory { and signed commands and `refs/nostr/` commands are not rewritten. - State no-op checks consider only ordinary refs. PR refs retain their own validation and cannot make an unchanged branch require a pending state. - - Fully satisfied unsigned pushes containing ordinary-ref deletions pass - normal authorization and then verify every target in one Git ref - transaction. This checks both existing targets and absent refs without - writing refs or unpacking the redundant upload. A changed or recreated - ref rejects verification. The response honors report-status (including v2) - and sideband negotiation. Commands making changes still use receive-pack; - zero-to-zero deletions are never synthesized for that process because Git - can interpret them as unconditional deletes. + - Unsigned pushes containing already-completed ordinary-ref deletions verify + satisfied commands without writing refs, and forward remaining commands + and the original pack/options to receive-pack. Non-atomic verifications + succeed or fail independently. Atomic requests prepare one verify-only + transaction and hold its ref locks until receive-pack finishes; failed + verification prevents any Git updates. Never send synthesized zero-to-zero + deletions to receive-pack: Git can treat them as unconditional deletes. + - Native Git ref statuses are combined with verification outcomes in request + order, preserving report-status/v2 and sideband progress. Normal signed-state + authorization remains a whole-request check; ref execution follows Git's + negotiated atomic/non-atomic behavior. Ordinary unsigned pushes also inspect + per-ref statuses: receive-pack can exit successfully after rejecting a ref. + Such requests retain successful tips but defer immediate state promotion, + so the request's follow-up work cannot override Git's rejection. Normal + background recovery can still reconcile the signed state independently. - Helper functions in [`helpers.rs`](../../src/purgatory/helpers.rs) handle ref extraction 4. **Bidirectional Waiting**: Either side can arrive first diff --git a/src/git/authorization.rs b/src/git/authorization.rs index 1411292..3eed02a 100644 --- a/src/git/authorization.rs +++ b/src/git/authorization.rs @@ -898,6 +898,28 @@ pub fn validate_push_refs( for (old_oid, new_oid, ref_name) in pushed_refs { debug!("Validating push: {} {} -> {}", ref_name, old_oid, new_oid); + // A deletion is authorized by absence from the signed state. Whole + // state matching above already checks all remaining branches and tags. + if new_oid == "0000000000000000000000000000000000000000" + && (ref_name.starts_with("refs/heads/") || ref_name.starts_with("refs/tags/")) + { + let declared = ref_name + .strip_prefix("refs/heads/") + .and_then(|name| state.get_branch_commit(name)) + .or_else(|| { + ref_name + .strip_prefix("refs/tags/") + .and_then(|name| state.get_tag_commit(name)) + }); + if declared.is_some() { + return Err(anyhow!( + "Cannot delete ref still declared in state: {}", + ref_name + )); + } + continue; + } + // Handle branch updates if let Some(branch_name) = ref_name.strip_prefix("refs/heads/") { if let Some(expected_commit) = state.get_branch_commit(branch_name) { @@ -2133,6 +2155,33 @@ mod tests { assert!(!membership.invited.contains(&hex(&bob))); } + #[test] + fn deletion_requires_absence_from_the_authorizing_state() { + let keys = create_test_keys(); + let oid = "a".repeat(40); + let state = RepositoryState::from_event( + EventBuilder::new(Kind::RepoState, "") + .tags([ + Tag::identifier("deletion"), + Tag::custom("refs/heads/main", [&oid]), + Tag::custom("refs/tags/v1", [&oid]), + ]) + .finalize(&keys) + .unwrap(), + ) + .unwrap(); + for name in ["refs/heads/main", "refs/tags/v1"] { + assert!( + validate_push_refs(&state, &[(oid.clone(), "0".repeat(40), name.into())]).is_err() + ); + } + for name in ["refs/heads/gone", "refs/tags/gone"] { + assert!( + validate_push_refs(&state, &[(oid.clone(), "0".repeat(40), name.into())]).is_ok() + ); + } + } + #[test] fn test_validate_push_refs_success() { let alice = create_test_keys(); diff --git a/src/git/completed_push.rs b/src/git/completed_push.rs deleted file mode 100644 index 0d9b91a..0000000 --- a/src/git/completed_push.rs +++ /dev/null @@ -1,272 +0,0 @@ -//! Complete fully satisfied pushes containing deletions without issuing writes. -//! A zero-to-zero receive-pack deletion is not a compare-and-swap: Git can -//! delete a concurrently recreated ref. Verify all targets in one transaction -//! instead, then report the push as satisfied at that point in time. - -use super::{authorization::parse_ref_line, protocol::PktLine}; -use hyper::body::Bytes; -use std::{ - collections::{HashMap, HashSet}, - io::{self, Write}, - path::Path, - process::{Command, Stdio}, -}; - -const ZERO: &str = "0000000000000000000000000000000000000000"; - -pub(super) struct CompletedPush { - /// Only for authorization. Never pass these zero-to-zero commands to Git. - pub authorization_body: Bytes, - targets: Vec<(String, String)>, - report_status: bool, - sideband: bool, -} - -impl CompletedPush { - /// Only unsigned, well-formed commands with every target already satisfied - /// qualify. Pushes making any change continue through receive-pack. - pub fn prepare(request: &[u8], local_refs: &HashMap) -> Option { - let mut input = request; - let mut targets = Vec::new(); - let mut seen = HashSet::new(); - let mut authorization = Vec::new(); - let mut has_deletion = false; - let mut report_status = false; - let mut sideband = false; - loop { - let (packet, remaining) = PktLine::parse(input).ok()?; - input = remaining; - let PktLine::Data(payload) = packet else { - break; - }; - // Certificates and malformed commands cannot parse as a ref line. - let (old, new, name) = parse_ref_line(&payload)?; - let line = std::str::from_utf8(&payload).ok()?.trim_end_matches('\n'); - let (command, caps) = line.split_once('\0').unwrap_or((line, "")); - if command != format!("{old} {new} {name}") - || !name.starts_with("refs/") - || !seen.insert(name.clone()) - || (!targets.is_empty() && line.contains('\0')) - { - return None; - } - if targets.is_empty() { - report_status = caps - .split_whitespace() - .any(|c| matches!(c, "report-status" | "report-status-v2")); - sideband = caps.split_whitespace().any(|c| c == "side-band-64k"); - } - let satisfied = if new == ZERO { - !local_refs.contains_key(&name) - } else { - local_refs.get(&name) == Some(&new) - }; - if !satisfied { - return None; - } - let ordinary = !name.starts_with("refs/nostr/"); - has_deletion |= ordinary && new == ZERO; - let auth_old = if ordinary { &new } else { &old }; - authorization - .extend(PktLine::data(format!("{auth_old} {new} {name}\n").into_bytes()).encode()); - targets.push((name, new)); - } - if !has_deletion { - return None; - } - authorization.extend(PktLine::flush().encode()); - Some(Self { - authorization_body: authorization.into(), - targets, - report_status, - sideband, - }) - } - - /// Git locks every ref during prepare. A ref recreated or changed since - /// inspection rejects the whole verification; no update or delete is sent. - pub fn verify(&self, repo: &Path) -> io::Result<()> { - let mut commands = b"start\0option no-deref\0".to_vec(); - for (name, oid) in &self.targets { - commands.extend(format!("verify {name}\0{oid}\0").as_bytes()); - } - commands.extend(b"prepare\0commit\0"); - let mut child = Command::new("git") - .current_dir(repo) - .args(["update-ref", "--stdin", "-z"]) - .stdin(Stdio::piped()) - .stdout(Stdio::null()) - .stderr(Stdio::piped()) - .spawn()?; - let write_result = child - .stdin - .take() - .expect("piped stdin") - .write_all(&commands); - let output = child.wait_with_output()?; - write_result?; - if !output.status.success() { - return Err(io::Error::other( - String::from_utf8_lossy(&output.stderr).trim().to_owned(), - )); - } - Ok(()) - } - - pub fn response(&self) -> Vec { - let mut report = Vec::new(); - if self.report_status { - report.extend(PktLine::data(b"unpack ok\n".to_vec()).encode()); - for (name, _) in &self.targets { - report.extend(PktLine::data(format!("ok {name}\n").into_bytes()).encode()); - } - report.extend(PktLine::flush().encode()); - } - if !self.sideband { - return report; - } - let mut response = Vec::new(); - for chunk in report.chunks(65515) { - let mut payload = vec![1]; - payload.extend(chunk); - response.extend(PktLine::data(payload).encode()); - } - response.extend(PktLine::flush().encode()); - response - } -} - -#[cfg(test)] -mod tests { - use super::*; - - fn git(repo: &Path, args: &[&str]) -> String { - let output = Command::new("git") - .current_dir(repo) - .args(["-c", "user.name=Test", "-c", "user.email=test@example.com"]) - .args(args) - .output() - .unwrap(); - assert!( - output.status.success(), - "{}", - String::from_utf8_lossy(&output.stderr) - ); - String::from_utf8(output.stdout).unwrap().trim().to_owned() - } - - fn request(oid: &str, caps: &str) -> Vec { - let mut bytes = - PktLine::data(format!("{oid} {ZERO} refs/heads/gone\0{caps}\n").into_bytes()).encode(); - bytes.extend(PktLine::flush().encode()); - bytes - } - - #[test] - fn verification_rejects_recreated_ref_without_deleting_it() { - let temp = tempfile::tempdir().unwrap(); - git(temp.path(), &["init", "--bare"]); - let tree = git(temp.path(), &["mktree"]); - let oid = git( - temp.path(), - &["commit-tree", &tree, "-m", "existing object"], - ); - let push = - CompletedPush::prepare(&request(&oid, "report-status"), &HashMap::new()).unwrap(); - push.verify(temp.path()).unwrap(); - assert!(git(temp.path(), &["for-each-ref"]).is_empty()); - // Recreate exactly the stale client's old target between inspection - // and verification. It must survive, not become an authorized delete. - git(temp.path(), &["update-ref", "refs/heads/gone", &oid]); - assert!(push.verify(temp.path()).is_err()); - assert_eq!(git(temp.path(), &["rev-parse", "refs/heads/gone"]), oid); - assert!(CompletedPush::prepare( - &request(&oid, "report-status"), - &HashMap::from([("refs/heads/gone".into(), oid)]) - ) - .is_none()); - } - - #[test] - fn verification_checks_companion_refs_without_reverting_them() { - let temp = tempfile::tempdir().unwrap(); - git(temp.path(), &["init", "--bare"]); - let tree = git(temp.path(), &["mktree"]); - let first = git(temp.path(), &["commit-tree", &tree, "-m", "first"]); - let second = git(temp.path(), &["commit-tree", &tree, "-m", "second"]); - git(temp.path(), &["update-ref", "refs/heads/main", &first]); - let mut bytes = request(&first, "report-status"); - bytes.truncate(bytes.len() - 4); - bytes.extend( - PktLine::data(format!("{ZERO} {first} refs/heads/main\n").into_bytes()).encode(), - ); - bytes.extend(PktLine::flush().encode()); - let push = - CompletedPush::prepare(&bytes, &HashMap::from([("refs/heads/main".into(), first)])) - .unwrap(); - push.verify(temp.path()).unwrap(); - git(temp.path(), &["update-ref", "refs/heads/main", &second]); - assert!(push.verify(temp.path()).is_err()); - assert_eq!(git(temp.path(), &["rev-parse", "refs/heads/main"]), second); - assert!(git(temp.path(), &["for-each-ref", "refs/heads/gone"]).is_empty()); - } - - #[test] - fn response_respects_status_capabilities_and_sideband_framing() { - for caps in [ - "", - "report-status", - "report-status-v2", - "report-status side-band-64k", - "report-status-v2 side-band-64k", - "side-band-64k", - ] { - let push = - CompletedPush::prepare(&request(&"1".repeat(40), caps), &HashMap::new()).unwrap(); - let wire = push.response(); - let mut report = Vec::new(); - if caps.contains("side-band-64k") { - let mut remaining = wire.as_slice(); - loop { - let (packet, tail) = PktLine::parse(remaining).unwrap(); - remaining = tail; - match packet { - PktLine::Data(payload) => { - assert_eq!(payload[0], 1); - report.extend(&payload[1..]); - } - PktLine::Flush => { - assert!(remaining.is_empty()); - break; - } - } - } - } else { - report = wire; - } - let expected = if caps.contains("report-status") { - [ - PktLine::data(b"unpack ok\n".to_vec()).encode(), - PktLine::data(b"ok refs/heads/gone\n".to_vec()).encode(), - PktLine::flush().encode(), - ] - .concat() - } else { - Vec::new() - }; - assert_eq!(report, expected, "{caps}"); - } - } - - #[test] - fn signed_malformed_and_duplicate_commands_do_not_take_shortcut() { - let line = request(&"1".repeat(40), "report-status"); - let mut signed = PktLine::data(b"push-cert\0report-status\n".to_vec()).encode(); - signed.extend(&line); - assert!(CompletedPush::prepare(&signed, &HashMap::new()).is_none()); - assert!(CompletedPush::prepare(&line[..line.len() - 4], &HashMap::new()).is_none()); - let mut duplicate = line[..line.len() - 4].to_vec(); - duplicate.extend(&line); - assert!(CompletedPush::prepare(&duplicate, &HashMap::new()).is_none()); - } -} diff --git a/src/git/handlers.rs b/src/git/handlers.rs index 79d215f..5706f1d 100644 --- a/src/git/handlers.rs +++ b/src/git/handlers.rs @@ -16,8 +16,8 @@ use tokio::sync::mpsc; use tokio::time::MissedTickBehavior; use tracing::{debug, error, info, warn, Instrument}; -use super::completed_push::CompletedPush; use super::protocol::{GitService, PktLine}; +use super::receive_pack_plan::ReceivePackPlan; use super::storage::{FamilyKey, FamilyWriteLease, LocalGitStorage}; use super::subprocess::GitSubprocess; use super::{full_body, GitResponseBody}; @@ -689,18 +689,16 @@ pub async fn handle_receive_pack( // advertised and removed its state from purgatory. Normalize only commands // already satisfied by this view; keep ordinary authorization for changes // and Git's compare-and-swap protection against a later writer. - let (request_body, completed) = match crate::git::list_refs(&repo_path) { + let (request_body, push_plan) = match crate::git::list_refs(&repo_path) { Ok(refs) => { let refs = refs.into_iter().collect(); - let completed = CompletedPush::prepare(&request_body, &refs); - ( - normalize_applied_ref_updates(&request_body, &refs), - completed, - ) + let normalized = normalize_applied_ref_updates(&request_body, &refs); + let push_plan = ReceivePackPlan::prepare(&normalized, &refs); + (normalized, push_plan) } Err(_) => (request_body, None), }; - let authorization_body = completed + let authorization_body = push_plan .as_ref() .map(|push| &push.authorization_body) .unwrap_or(&request_body); @@ -752,42 +750,43 @@ pub async fn handle_receive_pack( } }; - if let Some(completed) = completed { + let mut push_plan = if let Some(mut push_plan) = push_plan { let verify_path = repo_path.clone(); - let result = tokio::task::spawn_blocking(move || { - completed.verify(&verify_path)?; - Ok::<_, io::Error>(completed.response()) - }) - .await - .map_err(|error| GitError::Storage(error.to_string()))?; - return match result { - Ok(body) => { - info!( - identifier, - "Completed already-applied push with a verify-only ref transaction" - ); - record_git_operation(&metrics, "push", "success"); - Ok(Response::builder() - .status(StatusCode::OK) - .header( - "content-type", - GitService::ReceivePack.result_content_type(), - ) - .header("cache-control", "no-cache") - .body(full_body(body)) - .unwrap()) - } - Err(error) => { - warn!(identifier, %error, "Already-applied push failed ref verification"); - record_git_operation(&metrics, "push", "error"); - Ok(build_git_protocol_error_response( - GitService::ReceivePack, - "already-applied push no longer matches local refs", - Some(&request_body), - )) - } - }; + Some( + tokio::task::spawn_blocking(move || { + push_plan.verify(&verify_path); + push_plan + }) + .await + .map_err(|error| GitError::Storage(error.to_string()))?, + ) + } else { + None + }; + if let Some(plan) = push_plan.as_mut() { + if !plan.needs_git() { + let body = plan.response(&[]); + plan.release_locks(); + record_git_operation( + &metrics, + "push", + if plan.failed { "error" } else { "success" }, + ); + return Ok(Response::builder() + .status(StatusCode::OK) + .header( + "content-type", + GitService::ReceivePack.result_content_type(), + ) + .header("cache-control", "no-cache") + .body(full_body(body)) + .unwrap()); + } } + let forwarded_body = push_plan + .as_ref() + .map(|plan| &plan.forwarded_body) + .unwrap_or(&request_body); // Ref updates remain in the selected view; receive-pack's quarantine and // final objects are installed directly in the shared family inventory. @@ -802,14 +801,8 @@ pub async fn handle_receive_pack( ) .map_err(GitError::ProcessSpawnFailed)?; - // Write request to git's stdin - if let Some(mut stdin) = git.take_stdin() { - stdin.write_all(&request_body).await.map_err(|e| { - error!("Failed to write to git receive-pack stdin: {}", e); - GitError::IoError(e) - })?; - drop(stdin); - } + let stdin = git.take_stdin(); + let forwarded_body = forwarded_body.clone(); let stdout = git.take_stdout().ok_or_else(|| { GitError::IoError(io::Error::new( @@ -835,6 +828,23 @@ pub async fn handle_receive_pack( let git_data_path = git_data_path.to_owned(); tokio::spawn(async move { + // Own the subprocess and verification locks together before the first + // upload await. Cancellation must stop Git before those locks release. + if let Some(mut stdin) = stdin { + let written = tokio::select! { + result = stdin.write_all(&forwarded_body) => result, + _ = tx.closed() => Err(io::Error::new(io::ErrorKind::BrokenPipe, "push client disconnected")), + }; + drop(stdin); + if let Err(error) = written { + let _ = git.kill().await; + let _ = git.wait().await; + record_git_operation(&metrics, "push", "error"); + let _ = tx.send(Err(error)).await; + return; + } + } + stream_receive_pack_output( git, stdout, @@ -855,6 +865,7 @@ pub async fn handle_receive_pack( family_key, family_lease, pushed_refs, + push_plan, ) .await; }); @@ -883,12 +894,17 @@ async fn stream_receive_pack_output( family_key: FamilyKey, family_lease: Option, pushed_refs: Vec<(String, String, String)>, + mut push_plan: Option, ) where S: tokio::io::AsyncRead + Unpin + Send + 'static, E: tokio::io::AsyncRead + Unpin + Send + 'static, { let stderr_task = stderr.map(|stderr| tokio::spawn(read_stderr_to_end(stderr))); - let (pump_result, terminal_flush) = pump_receive_pack_stdout_to_channel(stdout, &tx).await; + let (pump_result, terminal_flush) = if let Some(plan) = push_plan.as_mut() { + plan.pump(stdout, &tx).await + } else { + pump_receive_pack_stdout_to_channel(stdout, &tx).await + }; if !matches!(pump_result, PumpResult::Eof { .. }) { let _ = git.kill().await; @@ -904,6 +920,10 @@ async fn stream_receive_pack_output( } }; + if let Some(plan) = push_plan.as_mut() { + plan.release_locks(); + } + let stderr_output = match stderr_task { Some(task) => task.await.unwrap_or_default(), None => Vec::new(), @@ -986,6 +1006,17 @@ async fn stream_receive_pack_output( // the subprocess boundary would deadlock. drop(repo_lifecycle_guard); + // A successful receive-pack exit can still contain rejected refs. Do not + // immediately apply the complete pending state over Git's partial result. + // The normal background state-recovery path remains responsible for it. + if push_plan.as_ref().is_some_and(|plan| plan.failed) { + record_git_operation(&metrics, "push", "error"); + if let Some(report) = terminal_flush { + let _ = send_body_bytes(&tx, report).await; + } + return; + } + // Git's receive-pack progress has already been streamed verbatim, but its // final flush remains withheld. Run GRASP follow-up work before releasing // that client-visible success boundary so push completion implies local diff --git a/src/git/mod.rs b/src/git/mod.rs index 7aa3279..ef01271 100644 --- a/src/git/mod.rs +++ b/src/git/mod.rs @@ -19,12 +19,12 @@ pub mod authorization; pub mod authorization_integrity; -mod completed_push; pub mod handlers; pub mod integrity; pub mod migration; pub mod process; pub mod protocol; +mod receive_pack_plan; pub mod storage; pub mod subprocess; pub mod sync; diff --git a/src/git/receive_pack_plan.rs b/src/git/receive_pack_plan.rs new file mode 100644 index 0000000..52bab63 --- /dev/null +++ b/src/git/receive_pack_plan.rs @@ -0,0 +1,704 @@ +//! Keep already-completed deletions out of receive-pack (zero-to-zero deletes +//! are unconditional there). Verify them separately and merge their per-ref +//! statuses into Git's report. Atomic pushes hold verification locks until Git exits. + +use super::{authorization::parse_ref_line, handlers::PumpResult, protocol::PktLine}; +use hyper::body::{Bytes, Frame}; +use std::{ + collections::{HashMap, HashSet}, + io::{self, BufRead, BufReader, Write}, + path::Path, + process::{Child, Command, Stdio}, +}; +use tokio::{ + io::{AsyncRead, AsyncReadExt}, + sync::mpsc, +}; + +const ZERO: &str = "0000000000000000000000000000000000000000"; + +struct CommandRef { + old: String, + new: String, + name: String, + completed: bool, + verified: Option, +} + +pub(super) struct ReceivePackPlan { + /// Authorization only: completed deletions appear as no-op commands. + pub authorization_body: Bytes, + pub forwarded_body: Bytes, + commands: Vec, + report_status: bool, + sideband: bool, + atomic: bool, + abort: bool, + pub failed: bool, + locks: Option, +} + +/// A prepared verify-only transaction. Closing stdin makes Git release its +/// locks without changing refs, including on early returns or cancellation. +struct RefLocks(Child); +impl RefLocks { + fn prepare(repo: &Path, targets: &[(&str, &str)]) -> io::Result { + let child = Command::new("git") + .current_dir(repo) + .args(["update-ref", "--stdin", "-z"]) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .spawn()?; + let mut guard = Self(child); + let mut commands = b"start\0".to_vec(); + for (name, oid) in targets { + commands.extend(format!("option no-deref\0verify {name}\0{oid}\0").as_bytes()); + } + commands.extend(b"prepare\0"); + guard + .0 + .stdin + .as_mut() + .expect("piped stdin") + .write_all(&commands)?; + let mut output = BufReader::new(guard.0.stdout.take().expect("piped stdout")); + for expected in ["start: ok\n", "prepare: ok\n"] { + let mut line = String::new(); + output.read_line(&mut line)?; + if line != expected { + return Err(io::Error::other("ref absence verification failed")); + } + } + Ok(guard) + } +} +impl Drop for RefLocks { + fn drop(&mut self) { + // No commit is needed: only verification was prepared, never a write. + self.0.stdin.take(); + let _ = self.0.wait(); + } +} + +impl ReceivePackPlan { + /// Only unsigned canonical commands qualify; retain pack/options bytes. + /// Non-deletion commands are already normalized by the ordinary push path. + pub fn prepare(request: &[u8], local_refs: &HashMap) -> Option { + let mut input = request; + let mut commands = Vec::new(); + let mut seen = HashSet::new(); + let mut caps = String::new(); + loop { + let (packet, remaining) = PktLine::parse(input).ok()?; + input = remaining; + let PktLine::Data(payload) = packet else { + break; + }; + let (old, new, name) = parse_ref_line(&payload)?; + let line = std::str::from_utf8(&payload).ok()?.trim_end_matches('\n'); + let (command, capabilities) = line.split_once('\0').unwrap_or((line, "")); + if command != format!("{old} {new} {name}") + || !name.starts_with("refs/") + || !seen.insert(name.clone()) + || (!commands.is_empty() && line.contains('\0')) + { + return None; + } + if commands.is_empty() { + caps = capabilities.to_owned(); + } + let completed = !name.starts_with("refs/nostr/") + && if new == ZERO { + !local_refs.contains_key(&name) + } else { + local_refs.get(&name) == Some(&new) + }; + commands.push(CommandRef { + old, + new, + name, + completed, + verified: None, + }); + } + if commands.is_empty() { + return None; + } + // Ordinary pushes remain entirely in receive-pack. Only its unsafe + // already-absent deletion case needs separate verification. + if !commands.iter().any(|c| c.completed && c.new == ZERO) { + for command in &mut commands { + command.completed = false; + } + } + let report_status = caps + .split_whitespace() + .any(|c| matches!(c, "report-status" | "report-status-v2")); + let sideband = caps.split_whitespace().any(|c| c == "side-band-64k"); + let atomic = caps.split_whitespace().any(|c| c == "atomic"); + // We need Git's individual outcomes even if the client omitted reports. + if !report_status { + caps.push_str(" report-status"); + } + let mut authorization = Vec::new(); + let mut forwarded = Vec::new(); + for command in &commands { + let old = if command.completed { + &command.new + } else { + &command.old + }; + authorization.extend( + PktLine::data(format!("{old} {} {}\n", command.new, command.name).into_bytes()) + .encode(), + ); + if !command.completed { + let suffix = if forwarded.is_empty() { + format!("\0{caps}") + } else { + String::new() + }; + forwarded.extend( + PktLine::data( + format!("{} {} {}{suffix}\n", command.old, command.new, command.name) + .into_bytes(), + ) + .encode(), + ); + } + } + authorization.extend(PktLine::flush().encode()); + forwarded.extend(PktLine::flush().encode()); + forwarded.extend(input); + Some(Self { + authorization_body: authorization.into(), + forwarded_body: forwarded.into(), + commands, + report_status, + sideband, + atomic, + abort: false, + failed: false, + locks: None, + }) + } + + /// Non-atomic no-ops complete independently. Atomic no-ops retain every + /// ref lock across receive-pack's transaction on the other refs. + pub fn verify(&mut self, repo: &Path) { + if self.atomic { + let names: Vec<_> = self + .commands + .iter() + .filter(|c| c.completed) + .map(|c| (c.name.as_str(), c.new.as_str())) + .collect(); + if names.is_empty() { + return; + } + match RefLocks::prepare(repo, &names) { + Ok(locks) => self.locks = Some(locks), + Err(_) => self.abort = true, + } + for command in &mut self.commands { + if command.completed { + command.verified = Some(!self.abort); + } + } + } else { + for command in &mut self.commands { + if command.completed { + command.verified = + Some(RefLocks::prepare(repo, &[(&command.name, &command.new)]).is_ok()); + } + } + } + } + + pub fn needs_git(&self) -> bool { + !self.abort && self.commands.iter().any(|c| !c.completed) + } + + pub fn release_locks(&mut self) { + self.locks.take(); + } + + /// Merge native Git outcomes with verified deletion outcomes. Preserve v2 + /// option lines and request ordering; only atomic requests couple failures. + pub fn response(&mut self, git_report: &[u8]) -> Vec { + let mut input = git_report; + let mut unpack = if self.needs_git() { + "unpack missing Git status\n".to_owned() + } else { + "unpack ok\n".to_owned() + }; + let mut statuses: HashMap>> = HashMap::new(); + let mut last_ref = None; + while let Ok((packet, tail)) = PktLine::parse(input) { + input = tail; + let PktLine::Data(payload) = packet else { + break; + }; + let line = String::from_utf8_lossy(&payload); + if line.starts_with("unpack ") { + unpack = line.into_owned(); + } else if line.starts_with("ok ") || line.starts_with("ng ") { + if let Some(name) = line.split_whitespace().nth(1) { + last_ref = Some(name.to_owned()); + statuses.insert(name.to_owned(), vec![payload]); + } + } else if line.starts_with("option ") { + if let Some(lines) = last_ref.as_ref().and_then(|name| statuses.get_mut(name)) { + lines.push(payload); + } + } + } + let mut outcomes = Vec::new(); + for command in &self.commands { + let lines = if command.completed { + let status = if command.verified == Some(true) { + format!("ok {}\n", command.name) + } else { + format!("ng {} ref absence verification failed\n", command.name) + }; + vec![status.into_bytes()] + } else { + statuses.remove(&command.name).unwrap_or_else(|| { + vec![format!("ng {} missing Git status\n", command.name).into_bytes()] + }) + }; + outcomes.push(lines); + } + self.failed = self.abort + || unpack != "unpack ok\n" + || outcomes.iter().any(|lines| !lines[0].starts_with(b"ok ")); + let atomic_failed = self.atomic + && (self.abort + || unpack != "unpack ok\n" + || outcomes.iter().any(|lines| !lines[0].starts_with(b"ok "))); + let mut report = Vec::new(); + if self.report_status { + report.extend(PktLine::data(unpack.as_bytes().to_vec()).encode()); + for (command, lines) in self.commands.iter().zip(outcomes) { + // Native statuses describe actual Git effects. Never turn a + // native success into a claimed rollback; only synthesize + // atomic failures for omitted commands or a pre-Git abort. + if self.abort || (command.completed && atomic_failed) || unpack != "unpack ok\n" { + let reason = if atomic_failed { + "atomic push failure" + } else { + "unpacker error" + }; + report.extend( + PktLine::data(format!("ng {} {reason}\n", command.name).into_bytes()) + .encode(), + ); + } else { + for line in lines { + report.extend(PktLine::data(line).encode()); + } + } + } + report.extend(PktLine::flush().encode()); + } + if !self.sideband { + return report; + } + let mut response = Vec::new(); + for chunk in report.chunks(65515) { + let mut payload = vec![1]; + payload.extend(chunk); + response.extend(PktLine::data(payload).encode()); + } + response.extend(PktLine::flush().encode()); + response + } + + /// Forward sideband progress immediately; only the small status report is + /// retained for merging. Native receive-pack streams no pack bytes here. + pub async fn pump( + &mut self, + mut stdout: R, + tx: &mpsc::Sender, io::Error>>, + ) -> (PumpResult, Option>) { + let mut report = Vec::new(); + let mut sent_stdout = false; + let reading = async { + loop { + let mut header = [0; 4]; + let first = stdout.read(&mut header[..1]).await?; + if first == 0 { + break; + } + stdout.read_exact(&mut header[1..]).await?; + let len = std::str::from_utf8(&header) + .ok() + .and_then(|s| usize::from_str_radix(s, 16).ok()) + .ok_or_else(|| io::Error::other("invalid Git status framing"))?; + if len == 0 { + break; + } + if !(5..=65520).contains(&len) { + return Err(io::Error::other("invalid Git status packet length")); + } + let mut payload = vec![0; len - 4]; + stdout.read_exact(&mut payload).await?; + if self.sideband { + if payload[0] == 1 { + report.extend(&payload[1..]); + } else { + sent_stdout = true; + tx.send(Ok(Frame::data(PktLine::data(payload).encode().into()))) + .await + .map_err(|_| { + io::Error::new( + io::ErrorKind::BrokenPipe, + "push client disconnected", + ) + })?; + } + } else { + report.extend(PktLine::data(payload).encode()); + } + if report.len() + > self + .commands + .len() + .saturating_mul(4096) + .saturating_add(65536) + { + return Err(io::Error::other("Git status report exceeds command budget")); + } + } + Ok(()) + }; + let result: io::Result<()> = tokio::select! { + result = reading => result, + _ = tx.closed() => return (PumpResult::ClientDisconnected, None), + }; + match result { + Ok(()) => ( + PumpResult::Eof { sent_stdout }, + Some(self.response(&report)), + ), + Err(error) => { + let _ = tx.send(Err(error)).await; + (PumpResult::ReadError, None) + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use tokio::time::{timeout, Duration}; + + fn git(repo: &Path, args: &[&str]) -> String { + let output = Command::new("git") + .current_dir(repo) + .args(["-c", "user.name=Test", "-c", "user.email=test@example.com"]) + .args(args) + .output() + .unwrap(); + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8(output.stdout).unwrap().trim().into() + } + fn fixture() -> (tempfile::TempDir, String) { + let dir = tempfile::tempdir().unwrap(); + git(dir.path(), &["init", "--bare"]); + let tree = git(dir.path(), &["mktree"]); + let oid = git(dir.path(), &["commit-tree", &tree, "-m", "available"]); + (dir, oid) + } + fn request(oid: &str, caps: &str, second_delete: bool) -> Vec { + let lines = [ + format!("{oid} {ZERO} refs/heads/gone\0{caps}\n"), + if second_delete { + format!("{oid} {ZERO} refs/heads/other\n") + } else { + format!("{ZERO} {oid} refs/heads/feature\n") + }, + ]; + let mut bytes = Vec::new(); + for line in lines { + bytes.extend(PktLine::data(line.into_bytes()).encode()); + } + bytes.extend(b"0000PACK-tail-preserved"); + bytes + } + fn native_report(status: &str) -> Vec { + [ + PktLine::data(b"unpack ok\n".to_vec()).encode(), + PktLine::data(status.as_bytes().to_vec()).encode(), + PktLine::flush().encode(), + ] + .concat() + } + + #[test] + fn non_atomic_deletions_verify_independently_and_preserve_recreated_refs() { + let (dir, oid) = fixture(); + let mut plan = + ReceivePackPlan::prepare(&request(&oid, "report-status", true), &HashMap::new()) + .unwrap(); + git(dir.path(), &["update-ref", "refs/heads/gone", &oid]); + plan.verify(dir.path()); + assert!(!plan.needs_git()); + let response = String::from_utf8(plan.response(&[])).unwrap(); + assert!(response.contains("ng refs/heads/gone ref absence verification failed")); + assert!(response.contains("ok refs/heads/other")); + assert_eq!(git(dir.path(), &["rev-parse", "refs/heads/gone"]), oid); + } + + #[test] + fn atomic_verification_failure_aborts_other_commands() { + let (dir, oid) = fixture(); + let mut plan = ReceivePackPlan::prepare( + &request(&oid, "report-status atomic", false), + &HashMap::new(), + ) + .unwrap(); + git(dir.path(), &["update-ref", "refs/heads/gone", &oid]); + plan.verify(dir.path()); + assert!(!plan.needs_git()); + let response = String::from_utf8(plan.response(&[])).unwrap(); + assert!(response.contains("ng refs/heads/gone atomic push failure")); + assert!(response.contains("ng refs/heads/feature atomic push failure")); + assert!(git(dir.path(), &["for-each-ref", "refs/heads/feature"]).is_empty()); + assert_eq!(git(dir.path(), &["rev-parse", "refs/heads/gone"]), oid); + } + + #[test] + fn atomic_absence_lock_survives_until_git_finishes_and_cleans_up() { + let (dir, oid) = fixture(); + let mut plan = ReceivePackPlan::prepare( + &request(&oid, "report-status atomic", false), + &HashMap::new(), + ) + .unwrap(); + plan.verify(dir.path()); + assert!(plan.needs_git()); + let create = || { + Command::new("git") + .current_dir(dir.path()) + .args([ + "-c", + "core.filesRefLockTimeout=0", + "update-ref", + "refs/heads/gone", + &oid, + ]) + .output() + .unwrap() + }; + assert!( + !create().status.success(), + "absence lock must prevent recreation during atomic push" + ); + let response = + String::from_utf8(plan.response(&native_report("ok refs/heads/feature\n"))).unwrap(); + assert!(response.contains("ok refs/heads/gone")); + plan.release_locks(); + assert!( + create().status.success(), + "release must not leave a stale ref lock" + ); + } + + #[test] + fn forwarding_preserves_pack_and_moves_capabilities_without_zero_deletes() { + let request = request( + &"1".repeat(40), + "report-status-v2 side-band-64k atomic", + false, + ); + let plan = ReceivePackPlan::prepare(&request, &HashMap::new()).unwrap(); + let (packet, tail) = PktLine::parse(&plan.forwarded_body).unwrap(); + let PktLine::Data(payload) = packet else { + panic!("command missing") + }; + assert!(String::from_utf8(payload) + .unwrap() + .contains("refs/heads/feature\0report-status-v2 side-band-64k atomic")); + assert_eq!(tail, b"0000PACK-tail-preserved"); + let mut signed = PktLine::data(b"push-cert\0report-status\n".to_vec()).encode(); + signed.extend(&request); + assert!(ReceivePackPlan::prepare(&signed, &HashMap::new()).is_none()); + assert!(ReceivePackPlan::prepare(b"0009bad!!", &HashMap::new()).is_none()); + } + + #[test] + fn atomic_native_failure_changes_every_status_but_non_atomic_does_not() { + for atomic in [false, true] { + let (dir, oid) = fixture(); + let caps = if atomic { + "report-status atomic" + } else { + "report-status" + }; + let mut plan = + ReceivePackPlan::prepare(&request(&oid, caps, false), &HashMap::new()).unwrap(); + plan.verify(dir.path()); + let response = String::from_utf8( + plan.response(&native_report("ng refs/heads/feature rejected\n")), + ) + .unwrap(); + assert!( + response.contains(&format!( + "{} refs/heads/gone", + if atomic { "ng" } else { "ok" } + )), + "{response}" + ); + assert!(response.contains("ng refs/heads/feature")); + } + } + + #[test] + fn changed_companion_noop_does_not_fail_an_independent_deletion() { + let (dir, oid) = fixture(); + git(dir.path(), &["update-ref", "refs/heads/feature", &oid]); + let refs = HashMap::from([("refs/heads/feature".into(), oid.clone())]); + let mut plan = + ReceivePackPlan::prepare(&request(&oid, "report-status", false), &refs).unwrap(); + let tree = git(dir.path(), &["mktree"]); + let later = git(dir.path(), &["commit-tree", &tree, "-m", "later"]); + git(dir.path(), &["update-ref", "refs/heads/feature", &later]); + plan.verify(dir.path()); + assert!(!plan.needs_git()); + let response = String::from_utf8(plan.response(&[])).unwrap(); + assert!(response.contains("ok refs/heads/gone")); + assert!(response.contains("ng refs/heads/feature")); + assert_eq!(git(dir.path(), &["rev-parse", "refs/heads/feature"]), later); + } + + #[test] + fn internal_status_negotiation_is_not_exposed_to_clients_without_reports() { + for caps in ["", "side-band-64k"] { + let (dir, oid) = fixture(); + let mut plan = + ReceivePackPlan::prepare(&request(&oid, caps, false), &HashMap::new()).unwrap(); + assert!(String::from_utf8_lossy(&plan.forwarded_body).contains("report-status")); + plan.verify(dir.path()); + let response = plan.response(&native_report("ok refs/heads/feature\n")); + assert_eq!( + response, + if caps.is_empty() { + Vec::new() + } else { + b"0000".to_vec() + } + ); + assert!(!plan.failed); + } + } + + #[test] + fn native_success_is_never_relabelled_as_a_rollback() { + let (dir, oid) = fixture(); + let mut plan = ReceivePackPlan::prepare( + &request(&oid, "report-status atomic", false), + &HashMap::new(), + ) + .unwrap(); + plan.commands.push(CommandRef { + old: ZERO.into(), + new: oid, + name: "refs/heads/refused".into(), + completed: false, + verified: None, + }); + plan.verify(dir.path()); + let mut report = native_report("ok refs/heads/feature\n"); + report.truncate(report.len() - 4); + report.extend(PktLine::data(b"ng refs/heads/refused rejected\n".to_vec()).encode()); + report.extend(b"0000"); + let response = String::from_utf8(plan.response(&report)).unwrap(); + assert!(response.contains("ok refs/heads/feature")); + assert!(response.contains("ng refs/heads/refused")); + assert!(response.contains("ng refs/heads/gone")); + assert!(plan.failed); + } + + #[tokio::test] + async fn disconnected_client_stops_waiting_for_git_and_releases_verification_locks() { + let (dir, oid) = fixture(); + let mut plan = ReceivePackPlan::prepare( + &request(&oid, "report-status atomic", false), + &HashMap::new(), + ) + .unwrap(); + plan.verify(dir.path()); + assert!(plan.locks.is_some()); + let (_writer, reader) = tokio::io::duplex(32); + let (tx, rx) = mpsc::channel(1); + let task = tokio::spawn(async move { plan.pump(reader, &tx).await }); + drop(rx); + let (outcome, report) = timeout(Duration::from_secs(5), task) + .await + .unwrap() + .unwrap(); + assert_eq!(outcome, PumpResult::ClientDisconnected); + assert!(report.is_none()); + // This would fail immediately if the abandoned transaction kept a lock. + git( + dir.path(), + &[ + "-c", + "core.filesRefLockTimeout=0", + "update-ref", + "refs/heads/gone", + &oid, + ], + ); + } + + #[tokio::test] + async fn sideband_progress_streams_before_final_status_and_status_v2_survives() { + let (dir, oid) = fixture(); + let mut plan = ReceivePackPlan::prepare( + &request(&oid, "report-status-v2 side-band-64k", false), + &HashMap::new(), + ) + .unwrap(); + plan.verify(dir.path()); + let (mut writer, reader) = tokio::io::duplex(4096); + let (tx, mut rx) = mpsc::channel(4); + let task = tokio::spawn(async move { plan.pump(reader, &tx).await }); + use tokio::io::AsyncWriteExt; + let progress = PktLine::data(b"\x02counting objects\n".to_vec()).encode(); + writer.write_all(&progress).await.unwrap(); + let frame = timeout(Duration::from_secs(5), rx.recv()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!(frame.into_data().unwrap().as_ref(), progress.as_slice()); + let mut native = native_report("ok refs/heads/feature\n"); + native.truncate(native.len() - 4); + native.extend(PktLine::data(b"option refname refs/heads/feature\n".to_vec()).encode()); + native.extend(b"0000"); + let mut band = vec![1]; + band.extend(native); + writer + .write_all(&PktLine::data(band).encode()) + .await + .unwrap(); + writer.write_all(b"0000").await.unwrap(); + drop(writer); + let (_, final_report) = timeout(Duration::from_secs(5), task) + .await + .unwrap() + .unwrap(); + let report = String::from_utf8(final_report.unwrap()).unwrap(); + assert!(report.contains("ok refs/heads/gone")); + assert!(report.contains("ok refs/heads/feature")); + assert!(report.contains("option refname refs/heads/feature")); + } +} diff --git a/tests/git_push_promotion_race.rs b/tests/git_push_promotion_race.rs index 1959b58..13c497e 100644 --- a/tests/git_push_promotion_race.rs +++ b/tests/git_push_promotion_race.rs @@ -225,3 +225,168 @@ async fn promoted_push(with_pr: bool, deletion_caps: Option<&str>) { ); relay.shutdown(); } + +#[tokio::test] +async fn completed_deletion_and_new_branch_succeed_together() { + mixed_deletion_push(false, false, true).await; +} + +#[tokio::test] +async fn atomic_completed_deletion_and_new_branch_succeed_together() { + mixed_deletion_push(true, false, true).await; +} + +#[tokio::test] +async fn non_atomic_push_reports_success_and_rejection_per_ref() { + mixed_deletion_push(false, true, true).await; +} + +#[tokio::test] +async fn atomic_push_rejects_every_ref_when_one_update_fails() { + mixed_deletion_push(true, true, true).await; +} + +#[tokio::test] +async fn ordinary_non_atomic_push_preserves_partial_result() { + mixed_deletion_push(false, true, false).await; +} + +#[tokio::test] +async fn ordinary_atomic_push_preserves_rejection() { + mixed_deletion_push(true, true, false).await; +} + +async fn mixed_deletion_push(atomic: bool, reject_branch: bool, include_deletion: bool) { + let dir = tempfile::tempdir().unwrap(); + let owner = Keys::generate(); + let identifier = "mixed-push"; + let repo = dir + .path() + .join(owner.public_key().to_bech32().unwrap()) + .join("mixed-push.git"); + let storage = LocalGitStorage::new(dir.path()); + let key = FamilyKey::sha1(identifier).unwrap(); + storage.create_thin_view(&key, &repo).unwrap(); + let family = storage.ensure_family(&key).unwrap(); + let tree = String::from_utf8(git(&family, &["mktree"])).unwrap(); + let oid = String::from_utf8(git( + &family, + &["commit-tree", tree.trim(), "-m", "available"], + )) + .unwrap() + .trim() + .to_owned(); + let bad = oid.clone(); + if reject_branch { + use std::os::unix::fs::PermissionsExt; + let hook = repo.join("hooks/update"); + std::fs::write(&hook, b"#!/bin/sh\n[ \"$1\" != refs/heads/unavailable ]\n").unwrap(); + std::fs::set_permissions(hook, std::fs::Permissions::from_mode(0o755)).unwrap(); + } + let db: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded()); + db.save_event( + &EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags([ + Tag::identifier(identifier), + Tag::custom("clone", ["https://service.example/mixed-push.git"]), + ]) + .finalize(&owner) + .unwrap(), + ) + .await + .unwrap(); + let mut tags = vec![ + Tag::identifier(identifier), + Tag::custom("refs/heads/feature", [&oid]), + ]; + if reject_branch { + tags.push(Tag::custom("refs/heads/unavailable", [&bad])); + } + let state = EventBuilder::new(Kind::RepoState, "") + .tags(tags) + .finalize(&owner) + .unwrap(); + let purgatory = Arc::new(Purgatory::new(dir.path().to_path_buf())); + purgatory.add_state(state, identifier.into(), owner.public_key(), false); + let relay = LocalRelayBuilder::default().database(db.clone()).build(); + let caps = if atomic { + "report-status-v2 side-band-64k atomic" + } else { + "report-status" + }; + // Capabilities begin on the omitted deletion and must move to the first + // actual Git update. Its pack must remain intact. + let mut commands = if include_deletion { + vec![ + format!("{oid} {} refs/heads/gone\0{caps}\n", "0".repeat(40)), + format!("{} {oid} refs/heads/feature\n", "0".repeat(40)), + ] + } else { + vec![format!( + "{} {oid} refs/heads/feature\0{caps}\n", + "0".repeat(40) + )] + }; + if reject_branch { + commands.push(format!("{} {bad} refs/heads/unavailable\n", "0".repeat(40))); + } + let mut raw = Vec::new(); + for command in commands { + raw.extend(format!("{:04x}{command}", command.len() + 4).as_bytes()); + } + raw.extend(b"0000"); + raw.extend(git(&repo, &["pack-objects", "--stdout"])); + let response = tokio::time::timeout( + Duration::from_secs(10), + handle_receive_pack( + repo.clone(), + raw.into(), + db, + relay.clone(), + identifier, + &owner.public_key().to_hex(), + purgatory, + dir.path().to_str().unwrap(), + None, + None, + None, + None, + ), + ) + .await + .unwrap() + .unwrap(); + let body = tokio::time::timeout(Duration::from_secs(10), response.into_body().collect()) + .await + .unwrap() + .unwrap() + .to_bytes(); + let report = String::from_utf8_lossy(&body); + let rejected_all = atomic && reject_branch; + let names = if include_deletion { + vec!["gone", "feature"] + } else { + vec!["feature"] + }; + for name in names { + assert!( + report.contains(&format!( + "{} refs/heads/{name}", + if rejected_all { "ng" } else { "ok" } + )), + "{report}" + ); + } + if reject_branch { + assert!(report.contains("ng refs/heads/unavailable"), "{report}"); + } + let refs = String::from_utf8(git( + &repo, + &["for-each-ref", "--format=%(refname) %(objectname)"], + )) + .unwrap(); + assert_eq!(refs.contains("refs/heads/feature"), !rejected_all, "{refs}"); + assert!(!refs.contains("refs/heads/gone"), "{refs}"); + assert!(!refs.contains("refs/heads/unavailable"), "{refs}"); + relay.shutdown(); +} diff --git a/tests/git_response_streaming.rs b/tests/git_response_streaming.rs index bd64639..6b0ea1b 100644 --- a/tests/git_response_streaming.rs +++ b/tests/git_response_streaming.rs @@ -11,6 +11,7 @@ use http_body_util::BodyExt; use hyper::body::{Body, Bytes, Frame}; use ngit_grasp::config::Config; use ngit_grasp::git::handlers::handle_receive_pack; +use ngit_grasp::git::protocol::PktLine; use ngit_grasp::git::sync::PurgatoryPromotionHooks; use ngit_grasp::grasp06::endpoint::PrsUrl; use ngit_grasp::grasp06::paths::prs_repo_path; @@ -64,7 +65,7 @@ async fn receive_pack_response_streams_stdout_before_subprocess_exit() { let fake_bin = tempfile::tempdir().expect("fake git bin tempdir"); let git_gate = GitGate::new().await; - write_fake_git_script(fake_bin.path(), false, git_gate.port()); + write_fake_git_script(fake_bin.path(), git_gate.port()); let _path = PathOverride::prepend(fake_bin.path()); let repo = tempfile::tempdir().expect("repo tempdir"); @@ -98,7 +99,7 @@ async fn receive_pack_response_streams_stdout_before_subprocess_exit() { // Git cannot exit until this test releases it, regardless of scheduling. git_gate.release().await; finish_body(&mut body, &mut streamed).await; - assert_eq!(streamed, b"first-progress\nsecond-progress\n"); + assert_eq!(streamed, expected_git_output()); } #[derive(Clone)] @@ -125,7 +126,7 @@ async fn receive_pack_terminal_flush_waits_for_purgatory_promotion() { let fake_bin = tempfile::tempdir().expect("fake git bin tempdir"); let git_gate = GitGate::new().await; - write_fake_git_script(fake_bin.path(), true, git_gate.port()); + write_fake_git_script(fake_bin.path(), git_gate.port()); let _path = PathOverride::prepend(fake_bin.path()); let keys = Keys::generate(); @@ -223,7 +224,13 @@ async fn receive_pack_terminal_flush_waits_for_purgatory_promotion() { finish_body(&mut body, &mut streamed).await; assert!( - streamed.starts_with(b"first-progress\nsecond-progress\n"), + streamed.starts_with( + &[ + progress_packet("first-progress\n"), + progress_packet("second-progress\n") + ] + .concat() + ), "git progress should remain at the start of the response" ); assert!( @@ -247,7 +254,7 @@ async fn prs_receive_pack_streams_stdout_before_cleanup_removes_empty_repo() { let fake_bin = tempfile::tempdir().expect("fake git bin tempdir"); let git_gate = GitGate::new().await; - write_fake_git_script(fake_bin.path(), true, git_gate.port()); + write_fake_git_script(fake_bin.path(), git_gate.port()); let _path = PathOverride::prepend(fake_bin.path()); let keys = Keys::generate(); @@ -304,7 +311,7 @@ async fn prs_receive_pack_streams_stdout_before_cleanup_removes_empty_repo() { ); git_gate.release().await; finish_body(&mut body, &mut streamed).await; - assert_eq!(streamed, b"first-progress\nsecond-progress\n0000"); + assert_eq!(streamed, expected_git_output()); assert!( !repo_path.exists(), @@ -346,12 +353,12 @@ where B: Body + Unpin, B::Error: std::fmt::Debug, { - const FIRST: &[u8] = b"first-progress\n"; + let first = progress_packet("first-progress\n"); timeout(Duration::from_secs(3), async { let mut streamed = Vec::new(); // Receive-pack can retain the last four bytes until further stdout // arrives, and pipe reads need not match writes or HTTP frames. - while streamed.len() < FIRST.len() - 4 { + while streamed.len() < first.len() - 4 { let frame = body .frame() .await @@ -359,7 +366,7 @@ where .expect("progress body frame"); streamed.extend_from_slice(&frame_data(frame)); assert!( - FIRST.starts_with(&streamed), + first.starts_with(&streamed), "unexpected initial git progress" ); } @@ -430,13 +437,36 @@ fn test_write_policy( ) } -fn write_fake_git_script(bin_dir: &Path, terminal_flush: bool, gate_port: u16) { +fn progress_packet(message: &str) -> Vec { + let mut payload = vec![2]; + payload.extend(message.as_bytes()); + PktLine::data(payload).encode() +} + +fn git_status_packet() -> Vec { + let mut payload = vec![1]; + payload.extend(PktLine::data(b"unpack ok\n".to_vec()).encode()); + payload + .extend(PktLine::data(format!("ok refs/nostr/{}\n", "a".repeat(64)).into_bytes()).encode()); + payload.extend(b"0000"); + [PktLine::data(payload).encode(), b"0000".to_vec()].concat() +} + +fn expected_git_output() -> Vec { + [ + progress_packet("first-progress\n"), + progress_packet("second-progress\n"), + git_status_packet(), + ] + .concat() +} + +fn shell_bytes(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("\\{byte:03o}")).collect() +} + +fn write_fake_git_script(bin_dir: &Path, gate_port: u16) { let git_path = bin_dir.join("git"); - let terminal_flush = if terminal_flush { - "printf '0000'\n" - } else { - "" - }; std::fs::write( &git_path, r#"#!/usr/bin/env bash @@ -452,12 +482,13 @@ done if [ "$is_receive_pack" = "1" ]; then cat >/dev/null exec 3<>/dev/tcp/127.0.0.1/__GATE_PORT__ -printf 'first-progress\n' +printf '%b' '__FIRST__' IFS= read -r -t 10 release <&3 [ "$release" = release ] exec 3<&- -printf 'second-progress\n' -__TERMINAL_FLUSH__exit 0 +printf '%b' '__SECOND__' +printf '%b' '__STATUS__' +exit 0 fi if [ "$1" = "init" ]; then @@ -486,7 +517,15 @@ fi echo "unsupported fake git invocation: $*" >&2 exit 1 "# - .replace("__TERMINAL_FLUSH__", terminal_flush) + .replace( + "__FIRST__", + &shell_bytes(&progress_packet("first-progress\n")), + ) + .replace( + "__SECOND__", + &shell_bytes(&progress_packet("second-progress\n")), + ) + .replace("__STATUS__", &shell_bytes(&git_status_packet())) .replace("__GATE_PORT__", &gate_port.to_string()), ) .expect("write fake git executable");