diff --git a/tests/git_response_streaming.rs b/tests/git_response_streaming.rs index f669493..bd64639 100644 --- a/tests/git_response_streaming.rs +++ b/tests/git_response_streaming.rs @@ -8,7 +8,7 @@ use std::time::Duration; use async_trait::async_trait; use clap::Parser; use http_body_util::BodyExt; -use hyper::body::{Bytes, Frame}; +use hyper::body::{Body, Bytes, Frame}; use ngit_grasp::config::Config; use ngit_grasp::git::handlers::handle_receive_pack; use ngit_grasp::git::sync::PurgatoryPromotionHooks; @@ -24,6 +24,8 @@ use ngit_grasp::purgatory::Purgatory; use ngit_grasp::sync::rejected_index::RejectedEventsIndex; use nostr_sdk::prelude::LocalRelayBuilder; use nostr_sdk::prelude::*; +use tokio::io::AsyncWriteExt; +use tokio::net::TcpListener; use tokio::sync::Semaphore; use tokio::time::timeout; @@ -61,7 +63,8 @@ async fn receive_pack_response_streams_stdout_before_subprocess_exit() { let _env_lock = PATH_ENV_LOCK.lock().await; let fake_bin = tempfile::tempdir().expect("fake git bin tempdir"); - write_fake_git(fake_bin.path()); + let git_gate = GitGate::new().await; + write_fake_git_script(fake_bin.path(), false, git_gate.port()); let _path = PathOverride::prepend(fake_bin.path()); let repo = tempfile::tempdir().expect("repo tempdir"); @@ -91,39 +94,10 @@ async fn receive_pack_response_streams_stdout_before_subprocess_exit() { let mut body = response.into_body(); - let first = timeout(Duration::from_secs(1), body.frame()) - .await - .expect("first stdout chunk should arrive before fake git exits") - .expect("body should still be open") - .expect("first frame should not be an HTTP body error"); - let first = frame_data(first); - assert!( - !first.is_empty(), - "first progress frame should contain data" - ); - - let no_second_yet = timeout(Duration::from_millis(250), body.frame()).await; - assert!( - no_second_yet.is_err(), - "body produced another frame while fake git was still sleeping; \ - this test needs the first frame to be observed before subprocess EOF" - ); - - let mut streamed = first.to_vec(); - loop { - let frame = timeout(Duration::from_secs(3), body.frame()) - .await - .expect("remaining stdout should arrive after fake git wakes"); - let Some(frame) = frame else { - break; - }; - streamed.extend_from_slice( - &frame - .expect("remaining frame should not be an HTTP body error") - .into_data() - .expect("remaining frame should contain data"), - ); - } + let mut streamed = read_first_progress(&mut body).await; + // 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"); } @@ -150,7 +124,8 @@ async fn receive_pack_terminal_flush_waits_for_purgatory_promotion() { let _env_lock = PATH_ENV_LOCK.lock().await; let fake_bin = tempfile::tempdir().expect("fake git bin tempdir"); - write_fake_git_with_terminal_flush(fake_bin.path()); + let git_gate = GitGate::new().await; + write_fake_git_script(fake_bin.path(), true, git_gate.port()); let _path = PathOverride::prepend(fake_bin.path()); let keys = Keys::generate(); @@ -203,12 +178,8 @@ async fn receive_pack_terminal_flush_waits_for_purgatory_promotion() { .expect("receive-pack handler should start fake subprocess"); let mut body = response.into_body(); - let first = timeout(Duration::from_secs(1), body.frame()) - .await - .expect("receive-pack progress should stream before promotion") - .expect("body should contain progress") - .expect("progress frame should not be an HTTP body error"); - let mut streamed = frame_data(first).to_vec(); + let mut streamed = read_first_progress(&mut body).await; + git_gate.release().await; timeout(Duration::from_secs(3), entered.acquire()) .await @@ -225,56 +196,31 @@ async fn receive_pack_terminal_flush_waits_for_purgatory_promotion() { "announcement must not be queryable while promotion is blocked" ); - while let Ok(Some(frame)) = timeout(Duration::from_millis(25), body.frame()).await { - streamed.extend_from_slice( - &frame - .expect("progress frame should not be an HTTP body error") - .into_data() - .expect("progress frame should contain data"), - ); - } - assert!( - !streamed.ends_with(b"0000"), - "receive-pack terminal flush must remain hidden while promotion is blocked" - ); - - let keepalive = timeout(Duration::from_secs(6), body.frame()) - .await - .expect("sideband keepalive should arrive during blocked promotion") - .expect("body should remain open while promotion is blocked") - .expect("keepalive should not be an HTTP body error"); - streamed.extend_from_slice( - &keepalive - .into_data() - .expect("keepalive frame should contain data"), - ); - assert!( - streamed + // Accumulate across arbitrary frame boundaries until the finalization + // keepalive is visible. Promotion remains blocked throughout this read. + timeout(Duration::from_secs(6), async { + while !streamed .windows(b"GRASP is finalizing the push\n".len()) - .any(|window| window == b"GRASP is finalizing the push\n"), - "blocked post-push processing should emit sideband progress" - ); - assert!( - !streamed.ends_with(b"0000"), - "keepalive must not expose the terminal flush" - ); + .any(|window| window == b"GRASP is finalizing the push\n") + { + let frame = body + .frame() + .await + .expect("body should remain open while promotion is blocked") + .expect("progress frame should not be an HTTP body error"); + streamed.extend_from_slice(&frame_data(frame)); + assert!( + !streamed.ends_with(b"0000"), + "terminal flush must remain hidden while promotion is blocked" + ); + } + }) + .await + .expect("sideband keepalive should arrive during blocked promotion"); release.add_permits(1); - loop { - let frame = timeout(Duration::from_secs(1), body.frame()) - .await - .expect("response should finish after promotion is released"); - let Some(frame) = frame else { - break; - }; - streamed.extend_from_slice( - &frame - .expect("terminal frame should not be an HTTP body error") - .into_data() - .expect("terminal frame should contain data"), - ); - } + finish_body(&mut body, &mut streamed).await; assert!( streamed.starts_with(b"first-progress\nsecond-progress\n"), @@ -300,7 +246,8 @@ async fn prs_receive_pack_streams_stdout_before_cleanup_removes_empty_repo() { let _env_lock = PATH_ENV_LOCK.lock().await; let fake_bin = tempfile::tempdir().expect("fake git bin tempdir"); - write_fake_git_with_terminal_flush(fake_bin.path()); + let git_gate = GitGate::new().await; + write_fake_git_script(fake_bin.path(), true, git_gate.port()); let _path = PathOverride::prepend(fake_bin.path()); let keys = Keys::generate(); @@ -350,61 +297,92 @@ async fn prs_receive_pack_streams_stdout_before_cleanup_removes_empty_repo() { let mut body = response.into_body(); - let first = timeout(Duration::from_secs(1), body.frame()) - .await - .expect("first /prs/ stdout chunk should arrive before fake git exits") - .expect("/prs/ body should still be open") - .expect("first /prs/ frame should not be an HTTP body error"); - let mut streamed = frame_data(first).to_vec(); - assert!( - b"first-progress\n".starts_with(&streamed), - "the four-byte terminal look-behind may split the first progress chunk" - ); - + let mut streamed = read_first_progress(&mut body).await; assert!( repo_path.exists(), "/prs/ cleanup must not remove the repo before receive-pack exits" ); - - let no_second_yet = timeout(Duration::from_millis(250), body.frame()).await; - assert!( - no_second_yet.is_err(), - "/prs/ body produced another frame while fake git was still sleeping; \ - this test needs the first frame to be observed before subprocess EOF" - ); - - let second = timeout(Duration::from_secs(3), body.frame()) - .await - .expect("second /prs/ stdout chunk should arrive after fake git wakes") - .expect("/prs/ body should still be open for second chunk") - .expect("second /prs/ frame should not be an HTTP body error"); - streamed.extend_from_slice(&frame_data(second)); - - // Family ref retention happens behind this boundary. Keep the deadline - // bounded but allow finalization the same scheduling headroom as the - // subprocess wake above. - let terminal = timeout(Duration::from_secs(3), body.frame()) - .await - .expect("/prs/ terminal flush should arrive after cleanup") - .expect("/prs/ body should contain the terminal flush") - .expect("/prs/ terminal frame should not be an HTTP body error"); - streamed.extend_from_slice(&frame_data(terminal)); + git_gate.release().await; + finish_body(&mut body, &mut streamed).await; assert_eq!(streamed, b"first-progress\nsecond-progress\n0000"); - let eof = timeout(Duration::from_secs(1), body.frame()) - .await - .expect("/prs/ body should close after cleanup"); - assert!( - eof.is_none(), - "/prs/ body should be closed after fake git exits" - ); - assert!( !repo_path.exists(), "/prs/ cleanup should remove zero-ref repo after receive-pack exits" ); } +// Keep the listener bound until Git has observed the parent's release. +struct GitGate(TcpListener); + +impl GitGate { + async fn new() -> Self { + Self( + TcpListener::bind("127.0.0.1:0") + .await + .expect("bind fake git gate"), + ) + } + + fn port(&self) -> u16 { + self.0.local_addr().expect("fake git gate address").port() + } + + async fn release(self) { + timeout(Duration::from_secs(3), async { + let (mut stream, _) = self.0.accept().await.expect("accept fake git gate"); + stream + .write_all(b"release\n") + .await + .expect("release fake git"); + }) + .await + .expect("fake git should connect to its gate"); + } +} + +async fn read_first_progress(body: &mut B) -> Vec +where + B: Body + Unpin, + B::Error: std::fmt::Debug, +{ + const FIRST: &[u8] = b"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 { + let frame = body + .frame() + .await + .expect("progress before fake git exit") + .expect("progress body frame"); + streamed.extend_from_slice(&frame_data(frame)); + assert!( + FIRST.starts_with(&streamed), + "unexpected initial git progress" + ); + } + streamed + }) + .await + .expect("first progress should stream before git is released") +} + +async fn finish_body(body: &mut B, streamed: &mut Vec) +where + B: Body + Unpin, + B::Error: std::fmt::Debug, +{ + timeout(Duration::from_secs(3), async { + while let Some(frame) = body.frame().await { + streamed.extend_from_slice(&frame_data(frame.expect("remaining body frame"))); + } + }) + .await + .expect("response should finish after fixture release"); +} + fn frame_data(frame: Frame) -> Bytes { frame.into_data().expect("frame should contain data") } @@ -452,15 +430,7 @@ fn test_write_policy( ) } -fn write_fake_git(bin_dir: &Path) { - write_fake_git_script(bin_dir, false); -} - -fn write_fake_git_with_terminal_flush(bin_dir: &Path) { - write_fake_git_script(bin_dir, true); -} - -fn write_fake_git_script(bin_dir: &Path, terminal_flush: bool) { +fn write_fake_git_script(bin_dir: &Path, terminal_flush: bool, gate_port: u16) { let git_path = bin_dir.join("git"); let terminal_flush = if terminal_flush { "printf '0000'\n" @@ -481,8 +451,11 @@ done if [ "$is_receive_pack" = "1" ]; then cat >/dev/null +exec 3<>/dev/tcp/127.0.0.1/__GATE_PORT__ printf 'first-progress\n' -sleep 2 +IFS= read -r -t 10 release <&3 +[ "$release" = release ] +exec 3<&- printf 'second-progress\n' __TERMINAL_FLUSH__exit 0 fi @@ -513,7 +486,8 @@ fi echo "unsupported fake git invocation: $*" >&2 exit 1 "# - .replace("__TERMINAL_FLUSH__", terminal_flush), + .replace("__TERMINAL_FLUSH__", terminal_flush) + .replace("__GATE_PORT__", &gate_port.to_string()), ) .expect("write fake git executable");