diff --git a/docs/explanation/defensive-measures.md b/docs/explanation/defensive-measures.md index 08c3a54..34be641 100644 --- a/docs/explanation/defensive-measures.md +++ b/docs/explanation/defensive-measures.md @@ -120,7 +120,8 @@ outbound sink: dial, including reconnects, plus a registration-time gate so forbidden targets never enter the reconnect lifecycle. - **Git fetches** (`RealSyncContext::fetch_oids`): checked immediately before - spawning `git fetch`. + spawning the pass's git subprocesses (`git ls-remote` and `git fetch`), + with the vetted DNS answers pinned onto each of them. The policy enforces: diff --git a/docs/explanation/grasp-02-proactive-sync-purgatory-git-data.md b/docs/explanation/grasp-02-proactive-sync-purgatory-git-data.md index 8fb5798..0bd56b4 100644 --- a/docs/explanation/grasp-02-proactive-sync-purgatory-git-data.md +++ b/docs/explanation/grasp-02-proactive-sync-purgatory-git-data.md @@ -400,6 +400,41 @@ Soft expiry is the solution: the bare repo is deleted at 30 minutes (respecting --- +## Fetch Strategy: Advertised Tips First + +A sync pass for one URL (`RealSyncContext::fetch_oids`) runs three phases, +all through the same hardened/pinned git subprocess machinery: + +1. **Compare** — `git ls-remote` lists the remote's advertised refs. Most + needed OIDs are ref tips declared by state events, and PR tips appear + under `refs/nostr/`, so the advertisement reveals up front + which needed OIDs the remote can serve. +2. **Batch-fetch advertised tips** — one `git fetch …` for the + needed OIDs the remote advertises. Advertised OIDs are always valid + wants, so `not our ref` cannot occur and an OID the remote never had + cannot fail the batch. +3. **Residual OIDs one at a time** — needed OIDs that were neither + advertised nor discovered as ancestors of the fetched tips are + requested individually. Some servers refuse arbitrary-SHA1 wants, so a + per-OID failure costs exactly one small round trip and never aborts the + rest of the pass. + +**Why not just list every OID as a want?** An earlier implementation did +exactly that and dropped one OID from the batch on each `not our ref` +error. Against a state event declaring hundreds of tips that exist on no +reachable server (observed in production with the `market` repository), +that degenerated into a sorted oid-by-oid crawl: one failed upload-pack +round trip per missing tip, O(N²) want retransmission, pack data streamed +and discarded on every failure, and nothing fetched until the crawl +finished. + +Each pass emits one `Purgatory git fetch pass complete` INFO line +(needed / advertised tips / residual attempted / residual missing / +fetched) and increments `ngit_purgatory_git_fetch_passes_total` and +`ngit_purgatory_git_fetch_oids_total{kind}`. + +--- + ## Testability: Mock-Based Architecture A key design goal was **100% unit test coverage** without requiring real git servers or databases. @@ -584,13 +619,19 @@ Sync incomplete - applying backoff (identifier=test-repo, attempt_count=2, next_ Failed to fetch OIDs (url=https://server.com/repo.git, error=connection timeout) ``` -### Metrics (Future) +### Metrics -Planned Prometheus metrics for observability: +Implemented: + +- `ngit_purgatory_git_fetch_passes_total` - Outbound git fetch passes + (one ls-remote comparison plus fetches per pass) +- `ngit_purgatory_git_fetch_oids_total{kind}` - Per-pass OID outcomes + (`advertised_tip`, `residual_attempted`, `residual_missing`, `fetched`) + +Planned: - `purgatory_sync_queue_size` - Number of identifiers pending sync - `purgatory_sync_attempts_total{identifier}` - Total sync attempts per identifier -- `purgatory_sync_oids_fetched_total{identifier}` - OIDs successfully fetched - `purgatory_domain_in_flight{domain}` - Current in-flight requests per domain - `purgatory_domain_requests_total{domain}` - Total requests per domain @@ -670,8 +711,9 @@ End-to-end tests verify sync behavior with real relay instances: ### 3. Prometheus Metrics -**Current**: Structured logging only -**Future**: Export metrics for monitoring dashboards +**Current**: Structured logging plus per-pass fetch metrics +(`ngit_purgatory_git_fetch_*`, see [Observability](#observability)) +**Future**: Queue-depth and per-domain throttle metrics **Use case**: Operators want visibility into sync performance, throttle effectiveness, success rates. diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 578b041..f215a4e 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -176,6 +176,60 @@ lazy_static! { .expect("register manual ejection deleted metric"); metric }; + static ref PURGATORY_GIT_FETCH_PASSES_TOTAL: Counter = { + let metric = Counter::with_opts(Opts::new( + "ngit_purgatory_git_fetch_passes_total", + "Outbound purgatory git fetch passes (one ls-remote comparison plus fetches per pass)", + )) + .expect("build purgatory git fetch passes metric"); + REGISTRY + .register(Box::new(metric.clone())) + .expect("register purgatory git fetch passes metric"); + metric + }; + static ref PURGATORY_GIT_FETCH_OIDS_TOTAL: CounterVec = { + let metric = CounterVec::new( + Opts::new( + "ngit_purgatory_git_fetch_oids_total", + "Purgatory git fetch OID outcomes by kind \ + (advertised_tip, residual_attempted, residual_missing, fetched)", + ), + &["kind"], + ) + .expect("build purgatory git fetch oids metric"); + REGISTRY + .register(Box::new(metric.clone())) + .expect("register purgatory git fetch oids metric"); + metric + }; +} + +/// Record one completed outbound purgatory git fetch pass. +/// +/// Together with the per-pass summary log line this gives production the +/// client-side view of purgatory git fetching: how many needed OIDs were +/// advertised by the remote (batch-fetched), how many residual OIDs were +/// attempted individually, how many the remote did not have, and how many +/// OIDs actually arrived. +pub fn record_purgatory_git_fetch_pass( + advertised_tips: usize, + residual_attempted: usize, + residual_missing: usize, + fetched: usize, +) { + PURGATORY_GIT_FETCH_PASSES_TOTAL.inc(); + PURGATORY_GIT_FETCH_OIDS_TOTAL + .with_label_values(&["advertised_tip"]) + .inc_by(advertised_tips as f64); + PURGATORY_GIT_FETCH_OIDS_TOTAL + .with_label_values(&["residual_attempted"]) + .inc_by(residual_attempted as f64); + PURGATORY_GIT_FETCH_OIDS_TOTAL + .with_label_values(&["residual_missing"]) + .inc_by(residual_missing as f64); + PURGATORY_GIT_FETCH_OIDS_TOTAL + .with_label_values(&["fetched"]) + .inc_by(fetched as f64); } pub fn record_blacklist_deletion_attempt(phase: &str) { diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index 075d266..e8ddb3d 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -300,7 +300,8 @@ fn resolve_pin_entry(resolved: &ResolvedTarget) -> Option { )) } -/// Build a `git fetch` command that cannot escape the outbound target policy. +/// Build a git command (fetch / ls-remote) that cannot escape the outbound +/// target policy. /// /// The URL was authorized immediately before this call; these controls keep /// the subprocess pointed at the vetted target: @@ -313,11 +314,7 @@ fn resolve_pin_entry(resolved: &ResolvedTarget) -> Option { /// of requests to event-directed servers; /// - `http.curloptResolve` pins the vetted DNS answers so the fetch cannot be /// re-bound to a different address between authorization and connection. -fn hardened_git_fetch_command( - repo_path: &Path, - resolve_pin: Option<&str>, - args: &[String], -) -> Command { +fn hardened_git_command(repo_path: &Path, resolve_pin: Option<&str>, args: &[String]) -> Command { let mut command = Command::new("git"); command .arg("-c") @@ -343,6 +340,66 @@ fn hardened_git_fetch_command( command } +/// Parse `git ls-remote` output into the set of advertised object ids. +/// +/// Each line has the form `\t`; peeled tag lines +/// (`refs/tags/x^{}`) are included since their OIDs are fetchable too. +fn parse_advertised_oids(stdout: &str) -> HashSet { + stdout + .lines() + .filter_map(|line| line.split_whitespace().next()) + .filter(|oid| oid.len() == 40 && oid.chars().all(|c| c.is_ascii_hexdigit())) + .map(|oid| oid.to_string()) + .collect() +} + +/// Whether a git fetch failure means the remote simply does not have (or +/// will not serve) the requested object — an expected per-OID outcome, not +/// a remote malfunction. +fn is_object_missing_error(stderr: &str) -> bool { + // Server-side rejection of an unknown want, relayed by the client as + // "fatal: remote error: upload-pack: not our ref ". + stderr.contains("not our ref") + // Client-side refusal when the server does not advertise the + // object and does not allow unadvertised wants (protocol v0). + || stderr.contains("Server does not allow request for unadvertised object") +} + +/// Record a remote failure with the naughty list tracker (only error +/// categories the tracker classifies as persistent are recorded). +fn record_remote_failure( + naughty_list: &NaughtyListTracker, + url: &str, + stderr: &str, + operation: &str, +) { + let Some(domain) = extract_domain(url) else { + return; + }; + let Some(category) = NaughtyListTracker::classify_error(stderr) else { + return; + }; + let is_new = naughty_list.record(&domain, category, stderr.to_string()); + + if is_new { + tracing::warn!( + domain = %domain, + category = %category, + operation = %operation, + error = %stderr, + "Git remote domain added to naughty list" + ); + } else { + debug!( + domain = %domain, + category = %category, + operation = %operation, + error = %stderr, + "Git remote operation failed (domain on naughty list)" + ); + } +} + #[async_trait] impl SyncContext for RealSyncContext { fn collect_pr_clone_urls(&self, identifier: &str) -> HashSet { @@ -489,110 +546,76 @@ impl SyncContext for RealSyncContext { let naughty_list = self.git_naughty_list.clone(); tokio::task::spawn_blocking(move || -> Result> { - let mut remaining_oids = missing_oids.clone(); - let mut missing_from_remote: Vec = Vec::new(); - - // Retry loop: keep fetching until success or no OIDs left - loop { - if remaining_oids.is_empty() { - // All OIDs were missing from remote - debug!( - url = %url, - missing_count = missing_from_remote.len(), - "All requested OIDs missing from remote" - ); - return Ok(vec![]); + // Phase 1: compare the remote's advertised refs against our + // needs. Most needed OIDs are ref tips declared by state events + // (PR tips appear under `refs/nostr/`), so the + // advertisement tells us up front which OIDs the remote can + // serve — instead of discovering each missing one through a + // failed "not our ref" upload-pack round trip. + let ls_remote_args = vec!["ls-remote".to_string(), url.clone()]; + let advertised = match hardened_git_command( + &repo_path, + resolve_pin.as_deref(), + &ls_remote_args, + ) + .output() + { + Ok(result) if result.status.success() => { + parse_advertised_oids(&String::from_utf8_lossy(&result.stdout)) } + Ok(result) => { + let stderr = String::from_utf8_lossy(&result.stderr); + record_remote_failure(&naughty_list, &url, &stderr, "ls-remote"); + return Err(anyhow::anyhow!( + "git ls-remote failed for {}: {}", + url, + stderr + )); + } + Err(e) => { + return Err(anyhow::anyhow!( + "git ls-remote command error for {}: {}", + url, + e + )) + } + }; - // git fetch ... - fetch all OIDs with full history + // Phase 2: batch-fetch the needed OIDs the remote actually + // advertises. Advertised OIDs are always valid wants, so + // "not our ref" cannot occur here and one batch cannot be + // failed by an OID the remote never had. + let advertised_tips: Vec = missing_oids + .iter() + .filter(|oid| advertised.contains(*oid)) + .cloned() + .collect(); + if !advertised_tips.is_empty() { let mut args = vec!["fetch".to_string(), url.clone()]; - args.extend(remaining_oids.iter().cloned()); + args.extend(advertised_tips.iter().cloned()); - let output = - hardened_git_fetch_command(&repo_path, resolve_pin.as_deref(), &args).output(); - - match output { - Ok(result) if result.status.success() => { - // Fetch succeeded - count how many OIDs we now have - let fetched: Vec = missing_oids - .iter() - .filter(|oid| crate::git::oid_exists(&repo_path, oid)) - .cloned() - .collect(); - - if !missing_from_remote.is_empty() { - debug!( - url = %url, - fetched_count = fetched.len(), - missing_count = missing_from_remote.len(), - missing_oids = ?missing_from_remote, - "Fetch completed after retries - some OIDs were missing from remote" - ); - } else { - debug!(url = %url, fetched_count = fetched.len(), "Successfully fetched OIDs"); - } - - return Ok(fetched); - } + match hardened_git_command(&repo_path, resolve_pin.as_deref(), &args).output() { + Ok(result) if result.status.success() => {} Ok(result) => { let stderr = String::from_utf8_lossy(&result.stderr); - - // Check for "not our ref" error - this is retryable - if stderr.contains("upload-pack: not our ref") { - // Parse out the missing OID from stderr - let missing_oid = stderr.lines().find_map(|line| { - if line.contains("not our ref") { - // Extract the OID from lines like: - // "fatal: remote error: upload-pack: not our ref " - line.split("not our ref") - .nth(1) - .map(|s| s.trim().to_string()) - } else { - None - } - }); - - if let Some(ref oid) = missing_oid { - // Remove the missing OID and retry with remaining - remaining_oids.retain(|o| o != oid); - missing_from_remote.push(oid.clone()); - - debug!( - url = %url, - missing_oid = %oid, - remaining_count = remaining_oids.len(), - "OID not found on remote, retrying with remaining OIDs" - ); - - continue; // Retry with remaining OIDs - } + if is_object_missing_error(&stderr) { + // The remote's refs changed between ls-remote and + // fetch. Fall through: every still-missing OID is + // retried individually below. + debug!( + url = %url, + error = %stderr, + "Advertised tip vanished between ls-remote and \ + fetch - falling back to per-OID fetches" + ); + } else { + record_remote_failure(&naughty_list, &url, &stderr, "fetch"); + return Err(anyhow::anyhow!( + "git fetch failed for {}: {}", + url, + stderr + )); } - - // Non-retryable error - record to naughty list and return error - if let Some(domain) = extract_domain(&url) { - if let Some(category) = NaughtyListTracker::classify_error(&stderr) { - let is_new = - naughty_list.record(&domain, category, stderr.to_string()); - - if is_new { - tracing::warn!( - domain = %domain, - category = %category, - error = %stderr, - "Git remote domain added to naughty list" - ); - } else { - debug!( - domain = %domain, - category = %category, - error = %stderr, - "Git fetch failed (domain on naughty list)" - ); - } - } - } - - return Err(anyhow::anyhow!("git fetch failed for {}: {}", url, stderr)); } Err(e) => { return Err(anyhow::anyhow!( @@ -603,6 +626,74 @@ impl SyncContext for RealSyncContext { } } } + + // Phase 3: request residual OIDs that were not advertised and + // did not arrive as ancestors of the fetched tips, one at a + // time. Arbitrary-SHA1 wants may still be refused by some + // servers, so a per-OID failure must cost exactly one small + // round trip and must not abort the rest. + let mut residual_attempted = 0usize; + let mut residual_missing = 0usize; + for oid in &missing_oids { + if crate::git::oid_exists(&repo_path, oid) { + continue; + } + residual_attempted += 1; + + let args = vec!["fetch".to_string(), url.clone(), oid.clone()]; + match hardened_git_command(&repo_path, resolve_pin.as_deref(), &args).output() { + Ok(result) if result.status.success() => {} + Ok(result) => { + let stderr = String::from_utf8_lossy(&result.stderr); + if is_object_missing_error(&stderr) { + residual_missing += 1; + debug!( + url = %url, + oid = %oid, + "OID not available on remote" + ); + } else { + // A genuine remote failure: record it and stop + // hammering the server, but keep what the batch + // already fetched. + record_remote_failure(&naughty_list, &url, &stderr, "fetch"); + break; + } + } + Err(e) => { + debug!(url = %url, oid = %oid, error = %e, "git fetch command error"); + break; + } + } + } + + let fetched: Vec = missing_oids + .iter() + .filter(|oid| crate::git::oid_exists(&repo_path, oid)) + .cloned() + .collect(); + + // One summary line per outbound fetch pass: this is the + // client-side visibility the drop-one-and-retry loop lacked + // (it logged only at debug level, so a relay crawling a remote + // was invisible in its own production logs). + tracing::info!( + url = %url, + needed = missing_oids.len(), + advertised_tips = advertised_tips.len(), + residual_attempted = residual_attempted, + residual_missing = residual_missing, + fetched = fetched.len(), + "Purgatory git fetch pass complete" + ); + crate::metrics::record_purgatory_git_fetch_pass( + advertised_tips.len(), + residual_attempted, + residual_missing, + fetched.len(), + ); + + Ok(fetched) }) .await .map_err(|e| anyhow::anyhow!("Failed to spawn blocking task: {}", e))? @@ -674,6 +765,43 @@ impl SyncContext for RealSyncContext { } } +#[cfg(test)] +mod fetch_helper_tests { + use super::*; + + #[test] + fn parse_advertised_oids_collects_ref_tips_and_peeled_tags() { + let stdout = "6451bf9d83e64647ba714740e9a15e5327d28be1\tHEAD\n\ + 6451bf9d83e64647ba714740e9a15e5327d28be1\trefs/heads/main\n\ + ace0a97f55e3b5343f312b1340717d34e2b4d2c6\trefs/tags/v1.0\n\ + beef000000000000000000000000000000000001\trefs/tags/v1.0^{}\n"; + let advertised = parse_advertised_oids(stdout); + assert_eq!(advertised.len(), 3); + assert!(advertised.contains("6451bf9d83e64647ba714740e9a15e5327d28be1")); + assert!(advertised.contains("ace0a97f55e3b5343f312b1340717d34e2b4d2c6")); + assert!(advertised.contains("beef000000000000000000000000000000000001")); + } + + #[test] + fn parse_advertised_oids_ignores_malformed_lines() { + let stdout = "warning: something\nnot-an-oid\trefs/heads/x\n\n"; + assert!(parse_advertised_oids(stdout).is_empty()); + } + + #[test] + fn object_missing_errors_are_recognised() { + assert!(is_object_missing_error( + "fatal: remote error: upload-pack: not our ref beef0000" + )); + assert!(is_object_missing_error( + "error: Server does not allow request for unadvertised object beef0000" + )); + assert!(!is_object_missing_error( + "fatal: unable to access 'https://x/': SSL certificate problem" + )); + } +} + // ============================================================================= // Mock Implementation for Testing // ============================================================================= diff --git a/tests/common/mod.rs b/tests/common/mod.rs index 1259145..6af269a 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -8,6 +8,7 @@ pub mod mock_relay; pub mod neg_limiting_proxy; pub mod nip09_helpers; pub mod req_limiting_proxy; +pub mod upload_pack_counting_proxy; pub mod port; pub mod purgatory_helpers; pub mod relay; diff --git a/tests/common/upload_pack_counting_proxy.rs b/tests/common/upload_pack_counting_proxy.rs new file mode 100644 index 0000000..8f38e30 --- /dev/null +++ b/tests/common/upload_pack_counting_proxy.rs @@ -0,0 +1,322 @@ +//! Counting Git Smart-HTTP Proxy for Purgatory Sync Tests +//! +//! A transparent HTTP proxy that sits between a syncing relay's outbound +//! purgatory git fetches and the relay serving the repository. It records +//! every `POST /git-upload-pack` exchange so tests can assert *how* the +//! syncing relay asks for missing git data, not just whether the data +//! eventually arrives: +//! +//! - the `want ` lines contained in each request body; +//! - whether the response carried an upload-pack `not our ref` error +//! (the signature a git server produces when a want names an object it +//! does not have); +//! - the order of exchanges, so tests can assert that available data is +//! fetched before (and independently of) requests destined to fail. +//! +//! `GET .../info/refs?service=git-upload-pack` requests (ref +//! advertisements, used by both `git fetch` and `git ls-remote`) are +//! counted but not recorded in detail. +//! +//! Bodies are buffered, not streamed: test repositories are tiny and the +//! proxy needs complete request/response bytes to classify the exchange. + +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; + +use http_body_util::{BodyExt, Full}; +use hyper::body::Bytes; +use hyper::server::conn::http1; +use hyper::service::service_fn; +use hyper::{Request, Response, StatusCode}; +use hyper_util::rt::TokioIo; +use tokio::net::TcpListener; +use tokio::sync::oneshot; + +/// One observed `POST /git-upload-pack` exchange. +#[derive(Debug, Clone)] +pub struct UploadPackExchange { + /// Object ids named in `want` lines of the request body, in order. + pub wants: Vec, + /// The response contained an upload-pack "not our ref" error. + pub not_our_ref: bool, + /// The exchange succeeded: HTTP 200 and no "not our ref" in the body. + pub ok: bool, + /// The request body carried a `Content-Encoding` the proxy cannot + /// inspect (git only compresses very large bodies; tests assert this + /// never happens so `wants` is trustworthy). + pub opaque_body: bool, +} + +/// Transparent counting proxy for a git smart-HTTP endpoint. +pub struct UploadPackCountingProxy { + url: String, + exchanges: Arc>>, + info_refs: Arc, + shutdown_tx: Option>, + handle: Option>, +} + +impl UploadPackCountingProxy { + /// Start the proxy on a random loopback port, forwarding every request + /// to `backend_url` (e.g. `http://127.0.0.1:`). + pub async fn start(backend_url: &str) -> Self { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("UploadPackCountingProxy failed to bind"); + let port = listener + .local_addr() + .expect("UploadPackCountingProxy local_addr") + .port(); + + let exchanges = Arc::new(Mutex::new(Vec::new())); + let info_refs = Arc::new(AtomicUsize::new(0)); + let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>(); + + let backend_url = backend_url.trim_end_matches('/').to_string(); + let accept_exchanges = exchanges.clone(); + let accept_info_refs = info_refs.clone(); + + let handle = tokio::spawn(async move { + // One shared upstream client; keep it plain so request and + // response bodies pass through byte-for-byte. + let client = reqwest::Client::builder() + .no_proxy() + .build() + .expect("build reqwest client"); + + loop { + tokio::select! { + accepted = listener.accept() => { + let Ok((stream, _)) = accepted else { break }; + let backend_url = backend_url.clone(); + let client = client.clone(); + let exchanges = accept_exchanges.clone(); + let info_refs = accept_info_refs.clone(); + tokio::spawn(async move { + let io = TokioIo::new(stream); + let service = service_fn(move |req| { + let backend_url = backend_url.clone(); + let client = client.clone(); + let exchanges = exchanges.clone(); + let info_refs = info_refs.clone(); + async move { + forward(req, &backend_url, &client, &exchanges, &info_refs) + .await + } + }); + if let Err(error) = http1::Builder::new() + .serve_connection(io, service) + .await + { + // Client disconnects mid-test are expected. + if !error.to_string().contains("connection") { + eprintln!( + "UploadPackCountingProxy connection error: {error}" + ); + } + } + }); + } + _ = &mut shutdown_rx => break, + } + } + }); + + Self { + url: format!("http://127.0.0.1:{port}"), + exchanges, + info_refs, + shutdown_tx: Some(shutdown_tx), + handle: Some(handle), + } + } + + /// Base URL of the proxy (`http://127.0.0.1:`); append the + /// `/{npub}/{identifier}.git` path when building clone URLs. + pub fn url(&self) -> &str { + &self.url + } + + /// Snapshot of all recorded `POST /git-upload-pack` exchanges, oldest + /// first. + pub fn exchanges(&self) -> Vec { + self.exchanges.lock().unwrap().clone() + } + + /// Number of ref-advertisement requests + /// (`GET .../info/refs?service=git-upload-pack`) observed. + pub fn info_refs_count(&self) -> usize { + self.info_refs.load(Ordering::Relaxed) + } + + /// Stop the proxy. + pub async fn stop(mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + if let Some(handle) = self.handle.take() { + let _ = handle.await; + } + } +} + +impl Drop for UploadPackCountingProxy { + fn drop(&mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + } +} + +/// Forward one request to the backend, recording upload-pack exchanges. +async fn forward( + req: Request, + backend_url: &str, + client: &reqwest::Client, + exchanges: &Mutex>, + info_refs: &AtomicUsize, +) -> Result>, hyper::Error> { + let method = req.method().clone(); + let path_and_query = req + .uri() + .path_and_query() + .map(|pq| pq.as_str().to_string()) + .unwrap_or_else(|| req.uri().path().to_string()); + let headers = req.headers().clone(); + let body_bytes = req.collect().await?.to_bytes(); + + let is_upload_pack_post = + method == hyper::Method::POST && path_and_query.ends_with("/git-upload-pack"); + if method == hyper::Method::GET + && path_and_query.contains("/info/refs") + && path_and_query.contains("service=git-upload-pack") + { + info_refs.fetch_add(1, Ordering::Relaxed); + } + + let opaque_body = headers.contains_key(hyper::header::CONTENT_ENCODING); + let wants = if is_upload_pack_post && !opaque_body { + parse_wants(&body_bytes) + } else { + Vec::new() + }; + + // Forward with original headers; the upstream client recomputes + // Host and Content-Length itself. + let mut upstream = client + .request( + reqwest::Method::from_bytes(method.as_str().as_bytes()).expect("valid method"), + format!("{backend_url}{path_and_query}"), + ) + .body(body_bytes.to_vec()); + for (name, value) in headers.iter() { + let name = name.as_str(); + if name.eq_ignore_ascii_case("host") || name.eq_ignore_ascii_case("content-length") { + continue; + } + upstream = upstream.header(name, value.as_bytes()); + } + + let upstream_response = match upstream.send().await { + Ok(response) => response, + Err(error) => { + eprintln!("UploadPackCountingProxy upstream error: {error}"); + return Ok(Response::builder() + .status(StatusCode::BAD_GATEWAY) + .body(Full::new(Bytes::from("upstream error"))) + .unwrap()); + } + }; + + let status = upstream_response.status(); + let response_headers = upstream_response.headers().clone(); + let response_bytes = upstream_response.bytes().await.unwrap_or_default(); + + if is_upload_pack_post { + let not_our_ref = contains(&response_bytes, b"not our ref"); + exchanges.lock().unwrap().push(UploadPackExchange { + wants, + not_our_ref, + ok: status.is_success() && !not_our_ref, + opaque_body, + }); + } + + let mut builder = Response::builder() + .status(hyper::StatusCode::from_u16(status.as_u16()).expect("valid status")); + for (name, value) in response_headers.iter() { + let name = name.as_str(); + // hyper recomputes framing headers for the buffered body. + if name.eq_ignore_ascii_case("transfer-encoding") + || name.eq_ignore_ascii_case("content-length") + { + continue; + } + builder = builder.header(name, value.as_bytes()); + } + Ok(builder.body(Full::new(response_bytes)).unwrap()) +} + +/// Extract the object ids named in `want` pkt-lines of a smart-HTTP +/// request body (works for protocol v0 and v2: both carry plain-text +/// `want <40-hex-oid>` sequences inside pkt-line framing). +fn parse_wants(body: &[u8]) -> Vec { + const NEEDLE: &[u8] = b"want "; + const OID_LEN: usize = 40; + let mut wants = Vec::new(); + let mut index = 0; + while index + NEEDLE.len() + OID_LEN <= body.len() { + if &body[index..index + NEEDLE.len()] == NEEDLE { + let oid = &body[index + NEEDLE.len()..index + NEEDLE.len() + OID_LEN]; + if oid.iter().all(u8::is_ascii_hexdigit) { + wants.push(String::from_utf8_lossy(oid).to_string()); + index += NEEDLE.len() + OID_LEN; + continue; + } + } + index += 1; + } + wants +} + +/// Byte-level substring search (`response_bytes` is protocol data, not +/// guaranteed to be valid UTF-8). +fn contains(haystack: &[u8], needle: &[u8]) -> bool { + haystack + .windows(needle.len()) + .any(|window| window == needle) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parse_wants_extracts_oids_from_pkt_lines() { + let body = b"0032want aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\n\ + 0054want bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb multi_ack side-band-64k\n\ + 0032have cccccccccccccccccccccccccccccccccccccccc\n0000"; + let wants = parse_wants(body); + assert_eq!( + wants, + vec![ + "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".to_string(), + "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".to_string(), + ] + ); + } + + #[test] + fn parse_wants_ignores_short_or_non_hex_suffixes() { + assert!(parse_wants(b"0009want zz\n").is_empty()); + assert!(parse_wants(b"want deadbeef").is_empty()); + } + + #[test] + fn contains_finds_error_text_in_binary_bodies() { + let mut body = vec![0u8, 1, 2]; + body.extend_from_slice(b"ERR upload-pack: not our ref beef"); + assert!(contains(&body, b"not our ref")); + assert!(!contains(&body, b"unrelated")); + } +} diff --git a/tests/sync.rs b/tests/sync.rs index f9316a9..b761032 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -39,6 +39,7 @@ mod sync { pub mod maintainer_reprocessing; pub mod metrics; pub mod neg_concurrency; + pub mod purgatory_fetch; pub mod req_concurrency; pub mod tag_variations; } diff --git a/tests/sync/mod.rs b/tests/sync/mod.rs index a351a62..af0a8cf 100644 --- a/tests/sync/mod.rs +++ b/tests/sync/mod.rs @@ -136,5 +136,6 @@ pub mod live_sync; pub mod maintainer_reprocessing; pub mod metrics; pub mod neg_concurrency; +pub mod purgatory_fetch; pub mod req_concurrency; pub mod tag_variations; \ No newline at end of file diff --git a/tests/sync/purgatory_fetch.rs b/tests/sync/purgatory_fetch.rs new file mode 100644 index 0000000..83d715c --- /dev/null +++ b/tests/sync/purgatory_fetch.rs @@ -0,0 +1,339 @@ +//! Purgatory Git Fetch Strategy Tests +//! +//! Regression coverage for a production failure observed on gitnostr.com +//! (2026-08-05, 07:56–08:41 UTC window): purgatory git sync requested every +//! needed commit id as an explicit want in a single +//! `git fetch …`. When the remote's upload-pack rejected +//! one want with "not our ref", the retry loop parsed that single oid out of +//! stderr, removed it, and re-sent the entire remaining batch. Against the +//! `market` repository — whose state event declares ~470 ref tips that exist +//! on no reachable server — this produced a sorted oid-by-oid crawl: +//! 4,521 `Git upload-pack failed after streaming stdout: … not our ref` +//! errors in 45 minutes on the serving side, one failed upload-pack round +//! trip per missing tip with O(N²) want retransmission, and nothing fetched +//! until the loop had crawled through every missing oid. +//! +//! The agreed fix compares the remote's advertised ref list (`git ls-remote` +//! through the same hardened subprocess machinery) against the needed oids +//! first, batch-fetches only advertised tips (always valid wants, so +//! "not our ref" cannot occur for them), and only then requests residual +//! oids one at a time so a single missing object cannot fail a batch. +//! +//! The scenario drives the real purgatory sync path end to end: a genuine +//! ngit-grasp relay serves a repository with two real branch tips behind a +//! counting proxy, and the relay under test holds a state event declaring +//! those two tips plus eight tips that exist nowhere (mirroring `market`). +//! The proxy records every upload-pack POST with its `want` lines, so the +//! test can assert the shape of the outbound fetching, not just the result. + +use std::path::Path; +use std::time::Duration; + +use nostr_sdk::prelude::*; + +use crate::common::purgatory_helpers::{ + add_commit_to_repo, create_branch, create_state_event, create_test_repo_with_commit, + push_to_relay, verify_event_not_served, wait_for_event_served, CommitVariant, +}; +use crate::common::upload_pack_counting_proxy::UploadPackCountingProxy; +use crate::common::{port, MockRelay, TestRelay}; + +/// Declared ref tips that exist on no reachable server, mirroring the +/// `market` repository's unfetchable state event. +const MISSING_TIP_COUNT: usize = 8; + +/// Wait until every given oid exists in the bare repository at `repo_path`. +async fn wait_for_oids_in_repo(repo_path: &Path, oids: &[&str], deadline: Duration) -> bool { + let end = tokio::time::Instant::now() + deadline; + loop { + let all_present = repo_path.exists() + && oids.iter().all(|oid| { + grasp_audit::git_command() + .args(["cat-file", "-e", oid]) + .current_dir(repo_path) + .output() + .map(|output| output.status.success()) + .unwrap_or(false) + }); + if all_present { + return true; + } + if tokio::time::Instant::now() >= end { + return false; + } + tokio::time::sleep(Duration::from_millis(200)).await; + } +} + +/// Wait until the proxy has observed at least `min_total` upload-pack +/// exchanges and the count has been stable for `stable_for`. +async fn wait_for_upload_pack_quiescence( + proxy: &UploadPackCountingProxy, + min_total: usize, + stable_for: Duration, + deadline: Duration, +) -> bool { + let end = tokio::time::Instant::now() + deadline; + let mut last_total = 0usize; + let mut stable_since = tokio::time::Instant::now(); + loop { + let total = proxy.exchanges().len(); + if total != last_total { + last_total = total; + stable_since = tokio::time::Instant::now(); + } else if total >= min_total && stable_since.elapsed() >= stable_for { + return true; + } + if tokio::time::Instant::now() >= end { + return false; + } + tokio::time::sleep(Duration::from_millis(200)).await; + } +} + +/// Scenario: +/// 1. A source relay hosts one repository with two distinct branch tips +/// (`main`, `feature`). Its own announcement/state clone URLs point only +/// at itself, so the source never fetches outbound. +/// 2. A counting proxy fronts the source's git smart-HTTP endpoint. +/// 3. The relay under test bootstrap-syncs from a MockRelay serving a +/// later announcement whose clone tag points at the proxy, plus a state +/// event declaring the two real tips and eight tips that exist nowhere. +/// Both events sit in purgatory; the purgatory sync loop fetches +/// through the proxy. +/// 4. The two real tips must arrive, and the outbound request shape must +/// not degrade into the production oid crawl: +/// - the first upload-pack request must succeed (available tips are +/// batch-fetched first, not after crawling through every missing oid); +/// - every failed upload-pack request must carry at most one want (a +/// single missing object must never fail a batch). +#[tokio::test] +async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() { + // 1. Source relay + counting proxy in front of its git endpoint, and + // the MockRelay that will carry the events for the relay under test. + let source = TestRelay::start().await; + let proxy = UploadPackCountingProxy::start(&format!("http://{}", source.domain())).await; + let mock = MockRelay::start().await; + + // 2. Repository with two distinct branch tips: main → commit_b, + // feature → commit_a (commit_a is commit_b's parent). + let git_dir = tempfile::tempdir().expect("create git repo dir"); + let commit_a = create_test_repo_with_commit(git_dir.path(), CommitVariant::StateTest) + .expect("create first commit"); + let commit_b = add_commit_to_repo(git_dir.path(), CommitVariant::SecondCommit) + .expect("create second commit"); + create_branch(git_dir.path(), "feature", Some(&commit_a)).expect("create feature branch"); + + let keys = Keys::generate(); + let npub = keys.public_key().to_bech32().expect("npub"); + let identifier = "fetch-strategy-repo"; + + // 3. Source-side events reference only the source itself, so the + // source's own purgatory sync has no external URL to fetch from and + // the proxy sees exclusively the relay under test. + let source_clone_url = format!("http://{}/{}/{}.git", source.domain(), npub, identifier); + let source_relay_url = format!("ws://{}", source.domain()); + let source_announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "Fetch strategy repo") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom("clone", vec![source_clone_url.clone()]), + Tag::custom("relays", vec![source_relay_url.clone()]), + ]) + .finalize(&keys) + .expect("sign source announcement"); + let source_state = create_state_event( + &keys, + identifier, + &[("main", &commit_b), ("feature", &commit_a)], + &[], + &[&source_clone_url], + &[&source_relay_url], + ) + .expect("create source state event"); + + let source_client = Client::builder() + .authenticator(SignerAuthenticator::new(keys.clone())) + .build(); + source_client + .add_relay(source.url()) + .await + .expect("add source relay"); + source_client.connect().await; + tokio::time::sleep(Duration::from_millis(500)).await; + source_client + .send_event(&source_announcement) + .await + .expect("send announcement to source"); + source_client + .send_event(&source_state) + .await + .expect("send state event to source"); + + // The state event in purgatory authorizes this push; the push releases + // it, after which the source serves both branches over git HTTP. + push_to_relay(git_dir.path(), &source.domain(), &npub, identifier) + .expect("push git data to source relay"); + wait_for_event_served(source.url(), &source_state.id, Duration::from_secs(15)) + .await + .expect("source state event should be released after push"); + + // 4. Relay under test: a later announcement whose clone tag points at + // the proxy, plus a state event declaring the two real tips and + // MISSING_TIP_COUNT tips that exist nowhere. Both are served by a + // MockRelay (no validation, no purgatory, no outbound fetching of + // its own) configured as the bootstrap relay: events arriving via + // sync take the immediate purgatory-sync path instead of the + // 3-minute wait-for-push delay applied to direct submissions. + let syncing_reservation = port::reserve_port(); + let syncing_domain = format!("127.0.0.1:{}", syncing_reservation.port()); + let proxy_clone_url = format!("{}/{}/{}.git", proxy.url(), npub, identifier); + let syncing_clone_url = format!("http://{}/{}/{}.git", syncing_domain, npub, identifier); + let syncing_relay_url = format!("ws://{}", syncing_domain); + + // The relays tag must list the relay under test (so it accepts the + // announcement) and the MockRelay (so the per-repo state subscription + // targets the MockRelay and delivers the state event via sync). + let syncing_announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "Fetch strategy repo") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom( + "clone", + vec![proxy_clone_url.clone(), syncing_clone_url.clone()], + ), + Tag::custom( + "relays", + vec![syncing_relay_url.clone(), mock.url().to_string()], + ), + ]) + .finalize(&keys) + .expect("sign syncing announcement"); + + let missing_tips: Vec = (0..MISSING_TIP_COUNT) + .map(|index| format!("beef{index:036x}")) + .collect(); + let missing_branch_names: Vec = (0..MISSING_TIP_COUNT) + .map(|index| format!("missing-{index}")) + .collect(); + let mut branches: Vec<(&str, &str)> = + vec![("main", commit_b.as_str()), ("feature", commit_a.as_str())]; + for (name, tip) in missing_branch_names.iter().zip(missing_tips.iter()) { + branches.push((name.as_str(), tip.as_str())); + } + let syncing_state = create_state_event( + &keys, + identifier, + &branches, + &[], + &[&proxy_clone_url, &syncing_clone_url], + &[&syncing_relay_url, mock.url()], + ) + .expect("create syncing state event"); + + let mock_client = Client::builder() + .authenticator(SignerAuthenticator::new(keys.clone())) + .build(); + mock_client + .add_relay(mock.url()) + .await + .expect("add mock relay"); + mock_client.connect().await; + tokio::time::sleep(Duration::from_millis(500)).await; + mock_client + .send_event(&syncing_announcement) + .await + .expect("send announcement to mock relay"); + mock_client + .send_event(&syncing_state) + .await + .expect("send state event to mock relay"); + + // Negentropy is disabled because MockRelay does not support NIP-77. + let syncing = TestRelay::start_on_reservation_with_options( + syncing_reservation, + Some(mock.url().to_string()), + true, + ) + .await; + + // 5. The two real tips must arrive through purgatory sync. + let repo_path = syncing + .git_data_path() + .join(&npub) + .join(format!("{identifier}.git")); + assert!( + wait_for_oids_in_repo( + &repo_path, + &[commit_a.as_str(), commit_b.as_str()], + Duration::from_secs(90), + ) + .await, + "available tips should be fetched into {} despite the unfetchable \ + tips declared alongside them (proxy exchanges: {:?})", + repo_path.display(), + proxy.exchanges(), + ); + + // 6. Let the sync pass finish so requests aimed at the missing tips + // (however the client shapes them) are all recorded. + wait_for_upload_pack_quiescence( + &proxy, + 1, + Duration::from_secs(3), + Duration::from_secs(60), + ) + .await; + + // 7. Regression assertions on the outbound request shape. + let exchanges = proxy.exchanges(); + assert!( + !exchanges.is_empty(), + "purgatory sync should fetch through the proxy" + ); + assert!( + exchanges.iter().all(|exchange| !exchange.opaque_body), + "upload-pack request bodies should be inspectable (no compression)" + ); + // Protocol v2 sends a want-less `ls-refs` POST before each fetch; the + // ordering guarantee is about the first request that names wants. + let first_want_exchange = exchanges + .iter() + .find(|exchange| !exchange.wants.is_empty()) + .expect("at least one upload-pack request should carry wants"); + assert!( + first_want_exchange.ok, + "the first want-carrying upload-pack request must batch-fetch the \ + available tips and succeed; instead it {} with wants {:?} — the \ + client crawled instead of consulting the advertised refs first", + if first_want_exchange.not_our_ref { + "failed with 'not our ref'" + } else { + "failed" + }, + first_want_exchange.wants, + ); + for (index, exchange) in exchanges.iter().enumerate() { + if !exchange.ok { + assert!( + exchange.wants.len() <= 1, + "failed upload-pack request #{index} carried {} wants {:?}; \ + a missing object must cost one single-want round trip and \ + must never fail a batch", + exchange.wants.len(), + exchange.wants, + ); + } + } + + // 8. The state event stays in purgatory — its declared tips are + // unfetchable, exactly like `market` in production. + verify_event_not_served(syncing.url(), &syncing_state.id, Duration::from_secs(1)) + .await + .expect("state event with unfetchable tips must stay in purgatory"); + + source_client.disconnect().await; + mock_client.disconnect().await; + syncing.stop().await; + mock.stop().await; + proxy.stop().await; + source.stop().await; +}