From 75da9a8446d79a38b84ce1590a4ed6801ba08bf7 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 12 Sep 2026 14:48:16 +0000 Subject: [PATCH] test: preserve subprocess and recovery fixture listeners Releasing a reservation before spawn or recovery allowed unrelated tests to take the same port. Transfer the socket across exec through a private Unix test-only protocol, retaining a parent copy across Grasp restarts. Validate the inherited descriptor and keep it out of Git and SSH descendants. Offline identity fixtures retain an accept-and-close endpoint until recovery. Mock relay listeners and upgraded connections remain owned until shutdown. Use real HTTP or SDK readiness instead of grace sleeps; startup errors in the owned-listener path are not retried. Normal server startup and deployment configuration are unchanged. The capability probe permits older external harnesses to remain compatible. Validation: listener adoption unit tests passed; relay_identity passed with recovery coverage, MockRelay shutdown passed, and a real subprocess restart plus a cross-repository ngit harness smoke test passed. Full platform builds remain host validation. Assisted-by: Codex (GPT-6) --- src/main.rs | 16 +++ src/server.rs | 8 +- src/test_listener.rs | 93 +++++++++++++++++ tests/common/mock_relay.rs | 88 +++++++++++++--- tests/common/port.rs | 209 ++++++++++++++++++++++++------------- tests/common/relay.rs | 129 ++++++++--------------- tests/relay_identity.rs | 47 +++++++-- 7 files changed, 407 insertions(+), 183 deletions(-) create mode 100644 src/test_listener.rs diff --git a/src/main.rs b/src/main.rs index 33d7936..b1fcb2a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -15,6 +15,8 @@ use ngit_grasp::{ }; mod docs_export; +#[cfg(unix)] +mod test_listener; /// Top-level CLI dispatcher. /// @@ -45,6 +47,14 @@ enum Cli { #[tokio::main] async fn main() -> Result<()> { + // Harnesses probe this before handing a reserved Unix listener to the child. + if std::env::args_os().nth(1).as_deref() == Some(OsStr::new("--internal-test-listener-support")) + { + #[cfg(unix)] + println!("ngit-test-listener-v1"); + return Ok(()); + } + // Documentation builds need only the static command model. Dispatch before // dotenv, clap parsing, secret discovery, filesystem access, or relay // startup so exporting metadata is safe in an isolated build environment. @@ -123,6 +133,12 @@ async fn run_relay(config: Config) -> Result<()> { "Starting ngit-grasp" ); + #[cfg(unix)] + let server = match test_listener::from_env()? { + Some(listener) => RelayServer::start_with_listener(config, listener).await?, + None => RelayServer::start(config).await?, + }; + #[cfg(not(unix))] let server = RelayServer::start(config).await?; server.run_until(shutdown_signal()).await diff --git a/src/server.rs b/src/server.rs index d9ae464..3a6d85e 100644 --- a/src/server.rs +++ b/src/server.rs @@ -82,7 +82,7 @@ impl RelayServer { /// Expects `config.relay_owner_nsec` to be set (see [`Config::load`]) /// and does **not** install a tracing subscriber — that is the /// caller's concern. - pub async fn start(mut config: Config) -> Result { + pub async fn start(config: Config) -> Result { // Bind first so kernel-assigned ports are resolved before any // component captures the domain / bind address. let requested: SocketAddr = config @@ -92,6 +92,12 @@ impl RelayServer { let listener = TcpListener::bind(&requested) .await .with_context(|| format!("failed to bind {}", requested))?; + Self::start_with_listener(config, listener).await + } + + /// Start from an already bound listener without releasing its address. + /// Used by subprocess test fixtures to preserve their port reservation. + pub async fn start_with_listener(mut config: Config, listener: TcpListener) -> Result { let local_addr = listener.local_addr()?; config.bind_address = local_addr.to_string(); if config.domain.is_empty() { diff --git a/src/test_listener.rs b/src/test_listener.rs new file mode 100644 index 0000000..744999e --- /dev/null +++ b/src/test_listener.rs @@ -0,0 +1,93 @@ +//! Private listener handoff protocol for subprocess test fixtures. + +use std::os::fd::{FromRawFd, RawFd}; + +use anyhow::{bail, Context, Result}; +use tokio::net::TcpListener; + +pub(crate) fn from_env() -> Result> { + let Some(fd) = std::env::var_os("NGIT_TEST_LISTENER_FD") else { + return Ok(None); + }; + if std::env::var("NGIT_TEST").as_deref() != Ok("1") { + bail!("NGIT_TEST_LISTENER_FD is available only with NGIT_TEST=1"); + } + let fd: RawFd = fd + .to_str() + .context("listener descriptor is not UTF-8")? + .parse() + .context("listener descriptor is not an integer")?; + listener_from_fd(fd).map(Some) +} + +fn listener_from_fd(fd: RawFd) -> Result { + if fd < 3 { + bail!("test listener must not use a standard stream descriptor"); + } + // dup validates the descriptor and gives this function its own ownership. + // The inherited descriptor remains open until the test subprocess exits. + let owned = unsafe { libc::dup(fd) }; + if owned < 0 { + return Err(std::io::Error::last_os_error()).context("duplicate test listener"); + } + // SAFETY: dup returned a fresh owned descriptor, consumed exactly once. + let listener = unsafe { std::net::TcpListener::from_raw_fd(owned) }; + // Do not leak either copy into Git/SSH children spawned by the relay. + for descriptor in [fd, owned] { + if unsafe { libc::fcntl(descriptor, libc::F_SETFD, libc::FD_CLOEXEC) } < 0 { + return Err(std::io::Error::last_os_error()).context("protect test listener from exec"); + } + } + let mut accepting: libc::c_int = 0; + let mut length = std::mem::size_of_val(&accepting) as libc::socklen_t; + // SAFETY: both pointers reference live, correctly sized writable values. + let result = unsafe { + libc::getsockopt( + owned, + libc::SOL_SOCKET, + libc::SO_ACCEPTCONN, + (&mut accepting as *mut libc::c_int).cast(), + &mut length, + ) + }; + if result != 0 { + return Err(std::io::Error::last_os_error()).context("inspect test listener"); + } + if accepting == 0 || !listener.local_addr()?.ip().is_loopback() { + bail!("test listener must be a listening loopback TCP socket"); + } + listener.set_nonblocking(true)?; + TcpListener::from_std(listener).context("adopt test listener") +} + +#[cfg(test)] +mod tests { + use super::*; + use std::os::fd::AsRawFd; + + #[tokio::test] + async fn adoption_preserves_the_reserved_address() { + let reserved = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let address = reserved.local_addr().unwrap(); + let listener = listener_from_fd(reserved.as_raw_fd()).unwrap(); + drop(reserved); + assert_eq!(listener.local_addr().unwrap(), address); + assert_eq!( + std::net::TcpListener::bind(address).unwrap_err().kind(), + std::io::ErrorKind::AddrInUse + ); + let client = tokio::net::TcpStream::connect(address).await.unwrap(); + let (_, peer) = tokio::time::timeout(std::time::Duration::from_secs(5), listener.accept()) + .await + .unwrap() + .unwrap(); + assert_eq!(peer, client.local_addr().unwrap()); + } + + #[tokio::test] + async fn adoption_rejects_non_listener_descriptors() { + let file = std::fs::File::open("/dev/null").unwrap(); + assert!(listener_from_fd(file.as_raw_fd()).is_err()); + assert!(listener_from_fd(-1).is_err()); + } +} diff --git a/tests/common/mock_relay.rs b/tests/common/mock_relay.rs index ebbed89..2af1225 100644 --- a/tests/common/mock_relay.rs +++ b/tests/common/mock_relay.rs @@ -34,7 +34,7 @@ //! - Does NOT perform any GRASP validation (no purgatory, no git data checks) use std::net::SocketAddr; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; use http_body_util::Full; use hyper::body::Bytes; @@ -46,6 +46,7 @@ use hyper_util::rt::TokioIo; use nostr_sdk::prelude::*; use tokio::net::TcpListener; use tokio::sync::oneshot; +use tokio::task::JoinSet; /// Mock Nostr relay that accepts all events without validation. /// @@ -195,6 +196,25 @@ impl MockRelay { .await } + /// Recover on an owned listener, with initial events visible before accepts. + pub async fn start_on_listener(listener: std::net::TcpListener, events: Vec) -> Self { + let port = listener.local_addr().expect("mock listener address").port(); + listener + .set_nonblocking(true) + .expect("nonblocking mock listener"); + let listener = TcpListener::from_std(listener).expect("register mock listener"); + Self::start_with_listener( + listener, + port, + RateLimit::default(), + None, + None, + events, + None, + ) + .await + } + /// Internal method to start the relay with an existing listener. #[allow(clippy::too_many_arguments)] async fn start_with_listener( @@ -232,6 +252,8 @@ impl MockRelay { let server_relay = relay.clone(); let handle = tokio::spawn(async move { + let mut connections = JoinSet::new(); + let upgrades = Arc::new(Mutex::new(JoinSet::new())); loop { tokio::select! { accept_result = listener.accept() => { @@ -242,8 +264,10 @@ impl MockRelay { let custom_nip11 = custom_nip11.clone(); let io = TokioIo::new(stream); - tokio::spawn(async move { + let upgrades = upgrades.clone(); + connections.spawn(async move { let service = service_fn(move |req| { + let upgrades = upgrades.clone(); let relay = relay.clone(); let custom_nip11 = custom_nip11.clone(); async move { @@ -253,6 +277,7 @@ impl MockRelay { remote_addr, pagination, custom_nip11, + upgrades, ) .await } @@ -275,12 +300,15 @@ impl MockRelay { } } } - _ = &mut shutdown_rx => { - // Shutdown signal received - break; - } + _ = &mut shutdown_rx => break, + _ = connections.join_next(), if !connections.is_empty() => {} } } + // Stop HTTP services before draining upgrades so no service can + // create another WebSocket task after the snapshot is taken. + connections.shutdown().await; + let mut upgrades = std::mem::take(&mut *upgrades.lock().unwrap()); + upgrades.shutdown().await; }); let url = format!("ws://127.0.0.1:{}", port); @@ -321,13 +349,17 @@ impl MockRelay { // Wait for server task to complete if let Some(handle) = self.handle.take() { - let _ = handle.await; + tokio::time::timeout(std::time::Duration::from_secs(5), handle) + .await + .expect("mock relay tasks should stop") + .expect("mock relay server task should not panic"); } } } impl Drop for MockRelay { fn drop(&mut self) { + self.relay.shutdown(); // Send shutdown signal if not already sent if let Some(tx) = self.shutdown_tx.take() { let _ = tx.send(()); @@ -342,6 +374,7 @@ async fn handle_request( addr: SocketAddr, pagination: Option, custom_nip11: Option, + upgrades: Arc>>, ) -> Result>, hyper::Error> { // Check for WebSocket upgrade request let is_websocket = req @@ -361,8 +394,10 @@ async fn handle_request( if let Some(key) = key { let accept_key = derive_accept_key(key.as_bytes()); - // Spawn task to handle the upgraded connection - tokio::spawn(async move { + // Retain upgrade tasks until the fixture shuts down. + let mut upgrades = upgrades.lock().unwrap(); + while upgrades.try_join_next().is_some() {} + upgrades.spawn(async move { match hyper::upgrade::on(req).await { Ok(upgraded) => { if let Err(e) = relay.take_connection(TokioIo::new(upgraded), addr).await { @@ -471,6 +506,33 @@ mod tests { use nostr_sdk::prelude::*; use std::time::Duration; + #[tokio::test] + async fn stop_closes_active_websocket_sessions() { + use futures_util::{SinkExt, StreamExt}; + use tokio_tungstenite::tungstenite::Message; + tokio::time::timeout(Duration::from_secs(5), async { + let mock = MockRelay::start().await; + let (mut client, _) = tokio_tungstenite::connect_async(mock.url()).await.unwrap(); + client + .send(Message::Text(r#"["REQ","stop-probe",{"limit":1}]"#.into())) + .await + .unwrap(); + loop { + let message = client.next().await.unwrap().unwrap(); + if message.to_text().is_ok_and(|text| text.contains("EOSE")) { + break; + } + } + mock.stop().await; + assert!(matches!( + client.next().await, + None | Some(Err(_)) | Some(Ok(Message::Close(_))) + )); + }) + .await + .expect("mock stop must close active sessions"); + } + #[tokio::test] async fn test_mock_relay_starts_and_stops() { let mock = MockRelay::start().await; @@ -494,10 +556,10 @@ mod tests { .add_relay(mock.url()) .await .expect("Failed to add relay"); - client.connect().await; - - // Wait for connection - tokio::time::sleep(Duration::from_millis(500)).await; + client + .try_connect_relay(mock.url(), Duration::from_secs(5)) + .await + .expect("connect to mock relay"); // Create and send a simple event let event = EventBuilder::new(Kind::TextNote, "Test note from MockRelay test") diff --git a/tests/common/port.rs b/tests/common/port.rs index 042f6ec..5de7c4f 100644 --- a/tests/common/port.rs +++ b/tests/common/port.rs @@ -1,48 +1,13 @@ -//! Race-free port reservation for test fixtures. +//! Bound listener ownership for test fixtures. //! -//! ## The race +//! Reserving a port and dropping its listener before the server binds leaves +//! a race with every other concurrent test. Instead, retain the listener and +//! transfer it directly to an in-process fixture or through the subprocess +//! listener handoff protocol. Retain a clone across server restarts. //! -//! The naive pattern — bind `127.0.0.1:0`, read the kernel-assigned port, -//! drop the listener, hand the bare `u16` to whoever wants it — has a -//! TOCTOU window between drop and the consumer's actual `bind`. During -//! that window, anything else in the process (or, more rarely, another -//! process) can be handed the same port by the kernel. -//! -//! The race is rare on lightly-loaded hardware but has been observed in -//! CI and during development, and the failure mode -//! (`Address already in use (os error 98)`) is a hard test fail with no -//! useful information for the next debugger. The hazard is sharply worse -//! in patterns that reserve a port well in advance of binding (e.g. -//! pre-allocating a port to embed in an announcement event before -//! starting the relay that will host it). -//! -//! Our in-process fixtures ([`MockRelay`], [`SmartGitServer`]) avoid the -//! race entirely by **keeping the listener bound** and handing it -//! straight to their tokio accept loop. That trick doesn't work for the -//! [`TestRelay`] subprocess — `ngit-grasp` binds itself from -//! `NGIT_BIND_ADDRESS`, and inheriting the pre-bound fd would require -//! Unix-specific `pre_exec` plumbing we'd rather not own in the test -//! harness. -//! -//! ## The reservation pattern -//! -//! Instead, [`reserve_port`] returns a [`PortReservation`] that **holds the -//! bound `TcpListener`** until the caller is about to start the real -//! service. While any reservation is live, no other call to -//! `reserve_port` in this process can be handed the same port — the -//! kernel won't reissue a port that is currently bound. -//! -//! The caller drops the reservation immediately before the real bind, -//! shrinking the TOCTOU window from "however long the fixture takes to -//! spawn" (or, worse, "however long the test takes to build the -//! announcement event") to "a few microseconds inside the start -//! function". The retry loop in [`crate::common::relay::TestRelay`] -//! covers that residual window — defense-in-depth that has never been -//! observed to fire in local stress runs. -//! -//! [`MockRelay`]: crate::common::mock_relay::MockRelay -//! [`SmartGitServer`]: crate::common::git_server::SmartGitServer -//! [`TestRelay`]: crate::common::relay::TestRelay +//! [`UnavailableEndpoint`] accepts and closes connections while retaining the +//! address. Recovery stops that accept loop before handing the same listener +//! to the recovered service, so another test can never claim the port. use std::net::TcpListener; @@ -50,17 +15,9 @@ use std::net::TcpListener; /// a live `TcpListener` so that no other [`reserve_port`] call in this /// process can be handed the same number. /// -/// The reservation is released by: -/// -/// - calling [`PortReservation::release`] to consume the reservation and -/// return the port number (preferred — makes the release explicit at the -/// call site), or -/// - simply dropping the value (also fine, but the release point is then -/// tied to lexical scope). -/// -/// The caller should release **immediately** before the consuming service -/// performs its own `bind` so that the TOCTOU window between -/// reservation-release and service-bind is as small as possible. +/// Transfer it with [`PortReservation::into_std_listener`] to start a service +/// without releasing its address. Dropping it or calling `release` relinquishes +/// the address and must not be used before starting a replacement service. #[derive(Debug)] pub struct PortReservation { port: u16, @@ -71,10 +28,8 @@ pub struct PortReservation { /// Reserve a specific loopback port. /// -/// Used by restart flows that must reuse an address already embedded in -/// published events (e.g. a repository announcement naming the relay's -/// domain). Fails while the port is still bound — callers should wait -/// for the previous holder to exit and retry within a bounded deadline. +/// Fails while another listener owns the address. Restart flows should retain +/// and transfer their existing listener rather than binding again. pub fn reserve_specific(port: u16) -> std::io::Result { let listener = TcpListener::bind(("127.0.0.1", port))?; Ok(PortReservation { @@ -89,10 +44,23 @@ impl PortReservation { self.port } - /// Consume the reservation, dropping the underlying listener and - /// returning the port number. The port is now free for the caller's - /// real service to bind. Prefer this over relying on lexical drop — - /// it makes the release point explicit at the call site. + /// Transfer the bound listener without making its port available again. + pub fn into_std_listener(self) -> TcpListener { + self._listener + } + + /// Wrap a listener retained across a subprocess restart. + pub fn from_listener(listener: TcpListener) -> Self { + Self { + port: listener + .local_addr() + .expect("reserved listener address") + .port(), + _listener: listener, + } + } + + /// Relinquish the address. A later bind to this port is inherently racy. pub fn release(self) -> u16 { let port = self.port; // `self` is consumed; the listener inside is dropped here. @@ -103,7 +71,7 @@ impl PortReservation { /// Bind `127.0.0.1:0`, capture the assigned port, and **keep the listener /// bound** inside the returned [`PortReservation`] until the caller -/// releases it. +/// transfers or drops it. /// /// While the reservation is live, no other `reserve_port` call in this /// process will be handed the same port. See module docs for why this @@ -121,6 +89,79 @@ pub fn reserve_port() -> PortReservation { } } +/// A continuously reserved endpoint that rejects protocol connections until +/// its listener is transferred to a recovered service. +pub struct UnavailableEndpoint { + listener: Option, + task: Option>, +} + +impl Default for UnavailableEndpoint { + fn default() -> Self { + Self::new() + } +} + +impl UnavailableEndpoint { + pub fn new() -> Self { + Self::from_listener(reserve_port().into_std_listener()) + } + + pub fn from_listener(listener: TcpListener) -> Self { + listener + .set_nonblocking(true) + .expect("nonblocking unavailable endpoint"); + let accepting = tokio::net::TcpListener::from_std( + listener.try_clone().expect("retain unavailable listener"), + ) + .expect("register unavailable listener"); + let task = tokio::spawn(async move { + loop { + let (stream, _) = accepting + .accept() + .await + .expect("accept unavailable connection"); + drop(stream); + } + }); + Self { + listener: Some(listener), + task: Some(task), + } + } + + pub fn port(&self) -> u16 { + self.listener + .as_ref() + .expect("owned unavailable listener") + .local_addr() + .expect("unavailable address") + .port() + } + + pub async fn into_listener(mut self) -> TcpListener { + if let Some(task) = self.task.take() { + task.abort(); + match tokio::time::timeout(std::time::Duration::from_secs(5), task) + .await + .expect("unavailable accept loop should stop before recovery") + { + Err(error) if error.is_cancelled() => {} + result => panic!("unavailable accept loop ended unexpectedly: {result:?}"), + } + } + self.listener.take().expect("transfer unavailable listener") + } +} + +impl Drop for UnavailableEndpoint { + fn drop(&mut self) { + if let Some(task) = self.task.take() { + task.abort(); + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -139,14 +180,36 @@ mod tests { assert_ne!(a.port(), c.port()); } - // Note: there is intentionally no "released port is immediately - // bindable" unit test. Once released, the port re-enters the - // kernel's free pool, and under heavy parallel test load (where - // dozens of `reserve_port` / `TcpListener::bind("127.0.0.1:0")` - // calls are racing each other) another test can be handed that - // port number before this one rebinds. That race is exactly what - // `reserve_port` exists to suppress for the held-reservation - // window; once released, it is by design out of scope. The - // "bindable after release" property is implicitly exercised by - // every passing `TestRelay::start*` integration test. + #[tokio::test] + async fn unavailable_endpoint_preserves_address_through_recovery() { + use tokio::io::AsyncReadExt; + let endpoint = UnavailableEndpoint::new(); + let address = endpoint.listener.as_ref().unwrap().local_addr().unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(5), async { + let mut client = tokio::net::TcpStream::connect(address).await.unwrap(); + let mut byte = [0]; + match client.read(&mut byte).await { + Ok(0) => {} + Err(error) if error.kind() == std::io::ErrorKind::ConnectionReset => {} + result => panic!("unavailable endpoint should close clients: {result:?}"), + } + }) + .await + .expect("unavailable endpoint should reject connections"); + + let listener = endpoint.into_listener().await; + assert_eq!(listener.local_addr().unwrap(), address); + assert_eq!( + TcpListener::bind(address).unwrap_err().kind(), + std::io::ErrorKind::AddrInUse + ); + let listener = tokio::net::TcpListener::from_std(listener).unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(5), async { + let client = tokio::net::TcpStream::connect(address).await.unwrap(); + let (_, peer) = listener.accept().await.unwrap(); + assert_eq!(peer, client.local_addr().unwrap()); + }) + .await + .expect("recovered listener should own every new connection"); + } } diff --git a/tests/common/relay.rs b/tests/common/relay.rs index a097d36..cd3f909 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -16,7 +16,7 @@ use nostr_sdk::prelude::{Keys, ToBech32}; use std::path::PathBuf; use std::process::{Child, Command, Stdio}; use std::time::{Duration, Instant}; -use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::io::AsyncWriteExt; use tokio::time::sleep; use crate::common::port::{self, PortReservation}; @@ -26,31 +26,18 @@ use crate::common::port::{self, PortReservation}; const READY_TIMEOUT: Duration = Duration::from_secs(5); /// How often to retry the HTTP probe while waiting for readiness. const READY_POLL: Duration = Duration::from_millis(100); -/// Extra grace after the HTTP service responds before declaring the relay -/// ready. -const READY_GRACE: Duration = Duration::from_millis(100); /// Per-attempt timeout for the HTTP readiness probe once TCP connects. const READY_PROBE_TIMEOUT: Duration = Duration::from_secs(1); -/// How many fresh port reservations to attempt before giving up. The -/// subprocess binds itself from `NGIT_BIND_ADDRESS`, so there is a -/// microsecond-scale TOCTOU window between [`PortReservation::release`] -/// and the subprocess's own `bind`. If that window loses the race the -/// subprocess exits before its TCP listener accepts; the readiness -/// check picks that up via `try_wait` so we can retry on a fresh port -/// instead of hanging the full readiness timeout. -/// -/// In practice this loop has never been observed to fire in local -/// stress testing — kept as defense-in-depth for CI / loaded hardware. -const MAX_BIND_ATTEMPTS: usize = 5; /// Test relay fixture that manages relay lifecycle /// /// Automatically starts and stops the ngit-grasp relay for testing. -/// Uses a kernel-assigned port held open by a [`PortReservation`] until -/// just before subprocess spawn, eliminating the same-process port race -/// that plagued the older "bind, drop, return port" pattern. +/// Transfers a kernel-assigned listener to the child while retaining a +/// parent copy, so startup and same-address restart never release the port. pub struct TestRelay { process: Child, + /// Parent copy retains the address across a same-endpoint restart. + listener: std::net::TcpListener, url: String, port: u16, /// Relay-owner identity configured in the subprocess. @@ -725,49 +712,21 @@ impl TestRelay { /// Single entry point that drives the spawn+readiness loop. /// - /// Retries up to [`MAX_BIND_ATTEMPTS`] times if the subprocess exits - /// early — that's the signature of having lost the bind race in the - /// microseconds between [`PortReservation::release`] and the - /// subprocess's own `bind`. Each retry draws a brand-new - /// kernel-assigned port; two consecutive `AddrInUse` failures would - /// therefore require two independent races back to back. - async fn start_internal(initial_reservation: PortReservation, options: RelayOptions) -> Self { - let mut reservation = Some(initial_reservation); - for attempt in 1..=MAX_BIND_ATTEMPTS { - // Each attempt consumes the current reservation. On retry we - // re-acquire from the kernel — guaranteed to give us a port - // number different from any reservation currently held - // elsewhere in this process. - let r = reservation - .take() - .expect("reservation always present on attempt entry"); - match Self::try_start_once(r, &options).await { - StartOutcome::Ready(relay) => return relay, - StartOutcome::EarlyExit { status } if attempt < MAX_BIND_ATTEMPTS => { - eprintln!( - "[TestRelay] ngit-grasp exited early on attempt \ - {attempt}/{MAX_BIND_ATTEMPTS} (status: {status:?}); \ - likely a port-bind race — retrying with a fresh port", - ); - reservation = Some(port::reserve_port()); - continue; - } - StartOutcome::EarlyExit { status } => { - panic!( - "ngit-grasp subprocess exited early after {MAX_BIND_ATTEMPTS} attempts \ - (last exit status: {status:?}). If this is not a port-bind race, \ - check /tmp/relay-*.log for the relay's stdout." - ); - } - } + /// Transfer the reserved listener to the child. Startup errors are real + /// failures; no port is released and no address retry is needed. + async fn start_internal(reservation: PortReservation, options: RelayOptions) -> Self { + match Self::try_start_once(reservation, &options).await { + StartOutcome::Ready(relay) => relay, + StartOutcome::EarlyExit { status } => panic!( + "ngit-grasp exited before readiness (status: {status:?}); see /tmp/relay-*.log" + ), } - unreachable!("MAX_BIND_ATTEMPTS loop terminated without returning") } /// One attempt at spawning ngit-grasp on the given reservation and /// waiting for it to be ready. Returns [`StartOutcome::EarlyExit`] /// specifically when the subprocess died before the readiness probe - /// succeeded — the caller may retry in that case. + /// succeeded. async fn try_start_once(reservation: PortReservation, options: &RelayOptions) -> StartOutcome { let port = reservation.port(); let bind_address = format!("127.0.0.1:{}", port); @@ -960,16 +919,31 @@ impl TestRelay { ); } - // Release the port reservation immediately before spawning the - // subprocess that will bind it. Holding the reservation through - // env-var setup above is what keeps any concurrent - // `reserve_port()` calls from picking this same number. - let _ = reservation.release(); + // Transfer the reservation through exec and retain a parent copy + // for same-address restarts. The port is never released during startup. + let listener = reservation.into_std_listener(); + #[cfg(unix)] + { + use std::os::{fd::AsRawFd, unix::process::CommandExt}; + let fd = listener.as_raw_fd(); + cmd.env("NGIT_TEST_LISTENER_FD", fd.to_string()); + // SAFETY: fcntl is async-signal-safe; the parent owns the descriptor + // through spawn. The child adopts it before starting the server. + unsafe { + cmd.pre_exec(move || { + if libc::fcntl(fd, libc::F_SETFD, 0) < 0 { + return Err(std::io::Error::last_os_error()); + } + Ok(()) + }); + } + } let process = cmd.spawn().expect("Failed to start relay process"); let mut relay = Self { process, + listener, url, port, owner_keys: test_keys.clone(), @@ -1035,27 +1009,13 @@ impl TestRelay { pub async fn restart(mut self) -> Self { let port = self.port; let options = self.options.clone(); - + let reservation = port::PortReservation::from_listener( + self.listener.try_clone().expect("retain restart listener"), + ); let _ = self.process.kill(); let _ = self.process.wait(); drop(self); - // The kernel frees the port once the child is fully gone; retry - // the specific-port reservation within a bounded deadline. - let deadline = Instant::now() + Duration::from_secs(10); - let reservation = loop { - match port::reserve_specific(port) { - Ok(reservation) => break reservation, - Err(error) => { - assert!( - Instant::now() < deadline, - "port {port} not released by stopped relay within deadline: {error}" - ); - sleep(Duration::from_millis(50)).await; - } - } - }; - match Self::try_start_once(reservation, &options).await { StartOutcome::Ready(relay) => relay, StartOutcome::EarlyExit { status } => panic!( @@ -1102,7 +1062,6 @@ impl TestRelay { Ok(()) => { // HTTP service handled a request successfully, so the // accept loop, Hyper service, and relay wiring are ready. - sleep(READY_GRACE).await; return ReadyOutcome::Ready; } Err(_) if Instant::now() < deadline => { @@ -1131,8 +1090,10 @@ impl TestRelay { ); stream.write_all(request.as_bytes()).await?; - let mut response = [0_u8; 64]; - let read = stream.read(&mut response).await?; + let mut response = Vec::new(); + let mut reader = tokio::io::BufReader::new(stream); + let read = + tokio::io::AsyncBufReadExt::read_until(&mut reader, b'\n', &mut response).await?; if read == 0 { return Err(std::io::Error::new( std::io::ErrorKind::UnexpectedEof, @@ -1163,13 +1124,7 @@ impl TestRelay { /// Stop the relay pub async fn stop(mut self) { - // Kill the process (gracefully if possible) - let _ = self.process.kill(); - - // Wait a bit for graceful shutdown - sleep(Duration::from_millis(100)).await; - - // Force kill if still running + // kill() sends SIGKILL; reap directly instead of guessing a grace period. let _ = self.process.kill(); let _ = self.process.wait(); } diff --git a/tests/relay_identity.rs b/tests/relay_identity.rs index 1e920a2..b15db00 100644 --- a/tests/relay_identity.rs +++ b/tests/relay_identity.rs @@ -5,6 +5,7 @@ mod common; use std::collections::BTreeSet; use std::time::Duration; +use common::port::UnavailableEndpoint; use common::{reserve_port, wait_for_event_on_relay, MockRelay, TestClient, TestRelay}; use nostr::nips::nip65; use nostr_sdk::prelude::*; @@ -244,9 +245,9 @@ async fn wiped_relay_adopts_identity_from_user_index_instead_of_publishing() { #[tokio::test] async fn identity_publication_defers_until_a_user_index_relay_is_reachable() { - // Release the reservation so connection attempts fail fast with refused; - // the MockRelay rebinds the same port later in the test. - let port = reserve_port().release(); + // Reject protocol connections while retaining the recovery address. + let unavailable = UnavailableEndpoint::new(); + let port = unavailable.port(); let index_url = format!("ws://127.0.0.1:{port}"); let relay = TestRelay::start_with_sync(Some(index_url)).await; @@ -272,7 +273,7 @@ async fn identity_publication_defers_until_a_user_index_relay_is_reachable() { // Once an empty index becomes reachable, the generated identity is // released: seeded locally and published to the index. - let index = MockRelay::start_on_port(port).await; + let index = MockRelay::start_on_listener(unavailable.into_listener().await, Vec::new()).await; for url in [relay.url(), index.url()] { for kind in [Kind::Metadata, Kind::RelayList] { assert!( @@ -297,9 +298,15 @@ async fn stored_identity_is_not_pushed_to_recovering_index_holding_an_identity() // Index A stays unreachable until it "recovers" already holding the // operator's customized profile; index B is reachable so the phase-1 // index check can succeed without A. - let port_a = reserve_port().release(); - let port_b = reserve_port().release(); - let index_b = MockRelay::start_on_port(port_b).await; + let unavailable_a = UnavailableEndpoint::new(); + let port_a = unavailable_a.port(); + let listener_b = reserve_port().into_std_listener(); + let port_b = listener_b.local_addr().expect("index B address").port(); + let index_b = MockRelay::start_on_listener( + listener_b.try_clone().expect("retain index B listener"), + Vec::new(), + ) + .await; let git_data = tempfile::tempdir().expect("git data dir"); let relay_data = tempfile::tempdir().expect("relay data dir"); let relay = TestRelay::start_on_reservation_persistent_user_index_relays( @@ -338,8 +345,10 @@ async fn stored_identity_is_not_pushed_to_recovering_index_holding_an_identity() // identity events must wait for a successful index check before any // publication, then reach only relays confirmed to hold nothing. index_b.stop().await; + let unavailable_b = UnavailableEndpoint::from_listener(listener_b); let relay = relay.restart().await; - let index_b = MockRelay::start_on_port(port_b).await; + let index_b = + MockRelay::start_on_listener(unavailable_b.into_listener().await, Vec::new()).await; for kind in [Kind::Metadata, Kind::RelayList] { assert!( wait_for_event_on_relay( @@ -356,7 +365,11 @@ async fn stored_identity_is_not_pushed_to_recovering_index_holding_an_identity() // A recovers already holding the customized profile. The per-relay // re-check before every send must adopt it locally instead of pushing // the stored (formerly generated) profile over it. - let index_a = MockRelay::start_on_port_with_events(port_a, vec![customized.clone()]).await; + let index_a = MockRelay::start_on_listener( + unavailable_a.into_listener().await, + vec![customized.clone()], + ) + .await; assert!( wait_for_event_on_relay( relay.url(), @@ -549,3 +562,19 @@ fn assert_identity_shape(relay: &TestRelay, events: &[Event]) { assert_eq!(advertised[0].0.to_string(), relay.url()); assert_eq!(advertised[0].1, None, "unmarked means read and write"); } + +#[tokio::test] +async fn subprocess_restart_retains_the_reserved_endpoint() { + tokio::time::timeout(Duration::from_secs(20), async { + let relay = TestRelay::start().await; + let url = relay.url().to_string(); + let relay = relay.restart().await; + assert_eq!(relay.url(), url); + let (_connection, _) = tokio_tungstenite::connect_async(relay.url()) + .await + .expect("restarted subprocess must serve the inherited listener"); + relay.stop().await; + }) + .await + .expect("subprocess startup and restart must complete"); +}