mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
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)
This commit is contained in:
+16
@@ -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
|
||||
|
||||
+7
-1
@@ -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<Self> {
|
||||
pub async fn start(config: Config) -> Result<Self> {
|
||||
// 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<Self> {
|
||||
let local_addr = listener.local_addr()?;
|
||||
config.bind_address = local_addr.to_string();
|
||||
if config.domain.is_empty() {
|
||||
|
||||
@@ -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<Option<TcpListener>> {
|
||||
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<TcpListener> {
|
||||
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());
|
||||
}
|
||||
}
|
||||
+75
-13
@@ -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<Event>) -> 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<PaginationConfig>,
|
||||
custom_nip11: Option<serde_json::Value>,
|
||||
upgrades: Arc<Mutex<JoinSet<()>>>,
|
||||
) -> Result<Response<Full<Bytes>>, 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")
|
||||
|
||||
+136
-73
@@ -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<PortReservation> {
|
||||
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<TcpListener>,
|
||||
task: Option<tokio::task::JoinHandle<()>>,
|
||||
}
|
||||
|
||||
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");
|
||||
}
|
||||
}
|
||||
|
||||
+42
-87
@@ -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();
|
||||
}
|
||||
|
||||
+38
-9
@@ -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");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user