diff --git a/crates/tsunagi-cli/Cargo.toml b/crates/tsunagi-cli/Cargo.toml index d4f5ede..ea3fe6c 100644 --- a/crates/tsunagi-cli/Cargo.toml +++ b/crates/tsunagi-cli/Cargo.toml @@ -10,6 +10,10 @@ repository.workspace = true [[bin]] name = "tsunagi" path = "src/main.rs" +# The binary and the library share a name, which `cargo doc` cannot map to +# two directories. The command line is documented by `--help`, not rustdoc, +# so the library keeps the name and this opts out. +doc = false [dependencies] tsunagi = { path = "../tsunagi", version = "0.1.0", features = ["tun-device", "dns-publish"] } diff --git a/crates/tsunagi-wg-quic/src/plugin.rs b/crates/tsunagi-wg-quic/src/plugin.rs index c023614..f19da0c 100644 --- a/crates/tsunagi-wg-quic/src/plugin.rs +++ b/crates/tsunagi-wg-quic/src/plugin.rs @@ -45,6 +45,7 @@ use crate::device::{PeerSummary, WireguardDevice}; use crate::keys::{WgPublicKey, WgSecretKey}; use crate::store::WgKeyStore; use tsunagi::state::Ipv4Range; +use tsunagi::task::{TASK_GRACE, wind_down}; /// The protocol identifier this plugin announces. /// The protocol id of this plugin. @@ -795,11 +796,14 @@ impl IpPlugin for WireguardPlugin { Box::pin(async move { let (reply_tx, reply_rx) = oneshot::channel(); if self.commands.send(Command::Stop(reply_tx)).await.is_ok() { - let _ = reply_rx.await; + // Bounded here as well as by the agent's grace: that grace + // abandons this future, which would leave the runtime task + // running behind it. + let _ = tokio::time::timeout(TASK_GRACE, reply_rx).await; } let task = self.task.lock().ok().and_then(|mut guard| guard.take()); if let Some(task) = task { - let _ = task.await; + wind_down(task, TASK_GRACE, "wg-quic runtime").await; } }) } diff --git a/crates/tsunagi/src/agent/mod.rs b/crates/tsunagi/src/agent/mod.rs index eacf002..a8c7dda 100644 --- a/crates/tsunagi/src/agent/mod.rs +++ b/crates/tsunagi/src/agent/mod.rs @@ -51,6 +51,7 @@ use crate::proto::message::ControlMessage; use crate::state::Ipv4Range; use crate::storage::{CacheOutcome, Storage}; +use crate::task::{TASK_GRACE, wind_down}; use network::{InboundSession, NetCommand, NetworkHandle, RuntimeParams}; use shutdown::Shutdown; @@ -750,12 +751,22 @@ impl Agent { handle.stop().await; } - self.inner.adapter.close().await; + // iroh's own close waits for peers to acknowledge; a peer that has + // gone silent must not decide how long that takes. + if tokio::time::timeout(TASK_GRACE, self.inner.adapter.close()) + .await + .is_err() + { + tracing::warn!("the endpoint did not close in time; abandoning it"); + } - for handle in [&self.inner.accept_task, &self.inner.plugin_task] { + for (what, handle) in [ + ("inbound connections", &self.inner.accept_task), + ("plugin requests", &self.inner.plugin_task), + ] { let task = handle.lock().ok().and_then(|mut guard| guard.take()); if let Some(task) = task { - let _ = task.await; + wind_down(task, TASK_GRACE, what).await; } } diff --git a/crates/tsunagi/src/agent/network.rs b/crates/tsunagi/src/agent/network.rs index 7e6d32b..6e601c3 100644 --- a/crates/tsunagi/src/agent/network.rs +++ b/crates/tsunagi/src/agent/network.rs @@ -105,7 +105,7 @@ impl NetworkHandle { /// Stops the runtime and waits for its task to finish. pub(crate) async fn stop(self) { self.shutdown.trigger(); - let _ = self.task.await; + crate::task::wind_down(self.task, crate::task::TASK_GRACE, "network runtime").await; } } diff --git a/crates/tsunagi/src/agent/shutdown.rs b/crates/tsunagi/src/agent/shutdown.rs index f9bbd5b..5081ea6 100644 --- a/crates/tsunagi/src/agent/shutdown.rs +++ b/crates/tsunagi/src/agent/shutdown.rs @@ -51,3 +51,25 @@ impl Shutdown { // The sender is gone, which for our purposes means "stop". } } + +#[cfg(test)] +mod tests { + #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] + + use super::*; + + #[tokio::test] + async fn waiting_resolves_once_triggered_and_stays_resolved() { + let shutdown = Shutdown::new(); + assert!(!shutdown.is_triggered()); + let waiter = shutdown.clone(); + let waiting = tokio::spawn(async move { waiter.wait().await }); + shutdown.trigger(); + waiting.await.unwrap(); + assert!(shutdown.is_triggered()); + // Triggering twice is not an error, and a token cloned afterwards + // sees the trigger that already happened. + shutdown.trigger(); + shutdown.clone().wait().await; + } +} diff --git a/crates/tsunagi/src/error.rs b/crates/tsunagi/src/error.rs index f212097..47f9898 100644 --- a/crates/tsunagi/src/error.rs +++ b/crates/tsunagi/src/error.rs @@ -110,6 +110,16 @@ pub enum Error { #[error("agent is stopped")] Stopped, + /// Something that had to answer in bounded time did not. + /// + /// Reported instead of waited out: an agent that has wedged must not + /// leave a command hanging with nothing on screen and no way out. + #[error("{what} did not answer in time")] + Timeout { + /// What was asked, phrased so the message stands on its own. + what: String, + }, + /// Discovery backend failure. Never fatal for the agent. #[error("discovery error: {0}")] Discovery(String), diff --git a/crates/tsunagi/src/ipc/unix.rs b/crates/tsunagi/src/ipc/unix.rs index fbacc62..9440925 100644 --- a/crates/tsunagi/src/ipc/unix.rs +++ b/crates/tsunagi/src/ipc/unix.rs @@ -7,6 +7,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; +use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::{UnixListener, UnixStream}; @@ -146,7 +147,19 @@ async fn serve(listener: UnixListener, source: Arc) { } async fn handle(mut stream: UnixStream, source: Arc) -> Result<()> { - let request: Request = read_message(&mut stream).await?; + // Bounded, so a connection that sends nothing cannot hold a task open. + // Only the wait for the request: building the report afterwards takes as + // long as it takes, and cutting it off would answer a live client with a + // closed socket. + let request: Request = + match tokio::time::timeout(EXCHANGE_TIMEOUT, read_message(&mut stream)).await { + Ok(request) => request?, + Err(_) => { + return Err(Error::Timeout { + what: "a local control connection".to_string(), + }); + } + }; let response = match request { Request::Status => Response::Status(Box::new(source.report().await)), Request::SetHostname(hostname) => match source.set_hostname(hostname).await { @@ -157,17 +170,18 @@ async fn handle(mut stream: UnixStream, source: Arc) -> Result write_message(&mut stream, &response).await } +/// How long either end waits for the other. +/// +/// A local answer comes from memory, so anything this slow means the agent is +/// wedged rather than busy. Saying so beats waiting: unbounded, one wedged +/// runtime leaves `tsunagi status` hanging with nothing on screen and no way +/// out but Ctrl-C. +const EXCHANGE_TIMEOUT: Duration = Duration::from_secs(5); + /// Asks a running agent for its status. pub async fn request_status(path: impl AsRef) -> Result { let path = path.as_ref(); - let mut stream = UnixStream::connect(path) - .await - .map_err(|source| Error::Io { - path: path.to_path_buf(), - source, - })?; - write_message(&mut stream, &Request::Status).await?; - match read_message::(&mut stream).await? { + match exchange(path, &Request::Status, EXCHANGE_TIMEOUT).await? { Response::Status(report) => Ok(*report), Response::Error(reason) => Err(Error::Storage(reason)), other => Err(Error::Storage(format!("unexpected answer: {other:?}"))), @@ -180,20 +194,38 @@ pub async fn request_status(path: impl AsRef) -> Result { /// asked for and may differ from it. pub async fn set_hostname(path: impl AsRef, hostname: &str) -> Result { let path = path.as_ref(); - let mut stream = UnixStream::connect(path) - .await - .map_err(|source| Error::Io { - path: path.to_path_buf(), - source, - })?; - write_message(&mut stream, &Request::SetHostname(hostname.to_string())).await?; - match read_message::(&mut stream).await? { + let request = Request::SetHostname(hostname.to_string()); + match exchange(path, &request, EXCHANGE_TIMEOUT).await? { Response::Hostname(accepted) => Ok(accepted), Response::Error(reason) => Err(Error::Storage(reason)), other => Err(Error::Storage(format!("unexpected answer: {other:?}"))), } } +/// Connects, sends one request and reads the answer, all within `within`. +/// +/// The bound covers the whole exchange rather than each read: an agent that +/// answers the header and then stops is as stuck as one that never answers. +async fn exchange(path: &Path, request: &Request, within: Duration) -> Result { + let attempt = async { + let mut stream = UnixStream::connect(path) + .await + .map_err(|source| Error::Io { + path: path.to_path_buf(), + source, + })?; + write_message(&mut stream, request).await?; + read_message::(&mut stream).await + }; + + match tokio::time::timeout(within, attempt).await { + Ok(result) => result, + Err(_) => Err(Error::Timeout { + what: format!("the agent at {}", path.display()), + }), + } +} + /// Marks the wire format of the local control socket. /// /// `b"TSN"` followed by the version, so a mismatch is recognised as one @@ -258,3 +290,49 @@ fn io_error(source: std::io::Error) -> Error { source, } } + +#[cfg(test)] +mod tests { + #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] + + use super::*; + + #[tokio::test] + async fn a_silent_agent_is_reported_rather_than_waited_out() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("control.sock"); + // A socket that accepts and then says nothing, like a wedged agent: + // the connection succeeds, the answer never comes. Unbounded, this + // call would never return. + let _listener = UnixListener::bind(&path).unwrap(); + + let error = exchange(&path, &Request::Status, Duration::from_millis(100)) + .await + .expect_err("a silent agent cannot be reported as healthy"); + + match error { + Error::Timeout { what } => assert!( + what.contains("control.sock"), + "the message names what did not answer: {what}" + ), + other => panic!("expected a timeout, got {other:?}"), + } + } + + #[tokio::test] + async fn an_answered_request_is_not_affected_by_the_bound() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("control.sock"); + let source: Arc = Arc::new(|| -> BoxFuture<'static, StatusReport> { + Box::pin(async { StatusReport::default() }) + }); + let control = ControlSocket::bind(&path, source).await.unwrap(); + + let answer = exchange(&path, &Request::Status, EXCHANGE_TIMEOUT) + .await + .unwrap(); + assert!(matches!(answer, Response::Status(_))); + + control.shutdown().await; + } +} diff --git a/crates/tsunagi/src/lib.rs b/crates/tsunagi/src/lib.rs index 1619602..4444e5a 100644 --- a/crates/tsunagi/src/lib.rs +++ b/crates/tsunagi/src/lib.rs @@ -48,6 +48,7 @@ pub mod overlay; pub mod proto; pub mod state; pub mod storage; +pub mod task; #[cfg(feature = "testing")] pub mod testing; diff --git a/crates/tsunagi/src/overlay/interface.rs b/crates/tsunagi/src/overlay/interface.rs index a05768f..49a78cf 100644 --- a/crates/tsunagi/src/overlay/interface.rs +++ b/crates/tsunagi/src/overlay/interface.rs @@ -238,7 +238,17 @@ impl Interface { /// Removes the interface from the host. pub async fn remove(&self) { + // The packet loop holds the device open, and it ends only when the + // device reports end of stream — which a real interface never does + // while it exists. So it is stopped here, not left to notice. + let task = match self.task.lock() { + Ok(mut guard) => guard.take(), + Err(poisoned) => poisoned.into_inner().take(), + }; self.factory.destroy(self.device.name()).await; + if let Some(task) = task { + crate::task::wind_down(task, crate::task::TASK_GRACE, "overlay packet loop").await; + } } /// The counters as they stand. diff --git a/crates/tsunagi/src/task.rs b/crates/tsunagi/src/task.rs new file mode 100644 index 0000000..7e2beea --- /dev/null +++ b/crates/tsunagi/src/task.rs @@ -0,0 +1,74 @@ +//! Stopping the tasks an agent started. +//! +//! Shutdown is cooperative first and bounded always: every loop watches a +//! cancellation token, and anything still running when its grace period ends +//! is aborted. Nothing an agent owns may outlive the agent, and no peer may +//! hold shutdown open by refusing to read. +//! +//! Public because a protocol crate spawns tasks of its own and is held to +//! the same rule: the agent gives each plugin a grace period, but abandoning +//! the future it is awaiting does not stop the task behind it. + +use std::time::Duration; + +use tokio::task::JoinHandle; + +/// How long a task gets to notice cancellation before it is cut off. +/// +/// Winding down is cooperative first: a task in the middle of a write to a +/// peer gets a moment to finish it. After that it is aborted, because a peer +/// that stops reading must not be able to hold shutdown open. +pub const TASK_GRACE: Duration = Duration::from_secs(5); + +/// Waits for `task` to finish, aborting it if it outlasts `grace`. +/// +/// Returns whether it finished on its own. `what` names it in the warning, +/// which is the only signal that something was cut off rather than stopped. +pub async fn wind_down(task: JoinHandle<()>, grace: Duration, what: &str) -> bool { + let abort = task.abort_handle(); + if tokio::time::timeout(grace, task).await.is_err() { + tracing::warn!( + task = what, + grace_secs = grace.as_secs_f32(), + "task did not wind down in time; aborting it" + ); + abort.abort(); + return false; + } + true +} + +#[cfg(test)] +mod tests { + #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] + + use super::*; + + #[tokio::test] + async fn a_task_that_stops_on_its_own_is_not_aborted() { + let (tx, mut rx) = tokio::sync::mpsc::channel::<()>(1); + let task = tokio::spawn(async move { + tx.send(()).await.ok(); + }); + assert!(wind_down(task, Duration::from_secs(5), "cooperative").await); + assert_eq!(rx.recv().await, Some(()), "it ran to completion"); + } + + #[tokio::test] + async fn a_task_that_ignores_the_grace_is_aborted() { + // A stand-in for a task blocked writing to a peer that stopped + // reading: it will never finish by itself. + let (tx, mut rx) = tokio::sync::mpsc::channel::<()>(1); + let task = tokio::spawn(async move { + let _held = tx; + std::future::pending::<()>().await; + }); + + let finished = wind_down(task, Duration::from_millis(50), "wedged").await; + assert!(!finished, "it had to be cut off"); + + // Aborted for real: the task is gone, so what it held is dropped. + let closed = tokio::time::timeout(Duration::from_secs(5), rx.recv()).await; + assert_eq!(closed, Ok(None), "the aborted task dropped its sender"); + } +} diff --git a/crates/tsunagi/src/testing.rs b/crates/tsunagi/src/testing.rs index 86c0b8b..1baa9d8 100644 --- a/crates/tsunagi/src/testing.rs +++ b/crates/tsunagi/src/testing.rs @@ -155,14 +155,22 @@ where { let deadline = Instant::now() + DEADLINE; loop { - if let Some(value) = probe().await { - return value; - } + let remaining = deadline + .checked_duration_since(Instant::now()) + .unwrap_or_default(); assert!( - Instant::now() < deadline, + !remaining.is_zero(), "timed out waiting for condition: {what}" ); - tokio::time::sleep(POLL_INTERVAL).await; + + // The probe is bounded too. A deadline checked only between probes is + // no deadline at all: one call into a wedged agent that never answers + // would hang the test process for ever instead of failing it. + match tokio::time::timeout(remaining, probe()).await { + Ok(Some(value)) => return value, + Ok(None) => tokio::time::sleep(POLL_INTERVAL).await, + Err(_) => panic!("timed out waiting for condition: {what}"), + } } } @@ -190,3 +198,20 @@ pub async fn wait_for_peers( }) .await } + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test(start_paused = true)] + #[should_panic(expected = "timed out waiting for condition")] + async fn a_probe_that_never_answers_fails_the_deadline_instead_of_hanging() { + // Virtual time, so the deadline arrives at once rather than in thirty + // real seconds — and the deadline does arrive, which is the point: a + // probe that never returns used to hang the process instead. + wait_until::<(), _, _>("an answer that never comes", || async { + std::future::pending::>().await + }) + .await; + } +} diff --git a/docs/testing.md b/docs/testing.md index 79d7152..a08e732 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -16,9 +16,12 @@ rights and no changes to OS network settings: endpoints bind `127.0.0.1:0` and disabled, and net-report probing is reduced to its minimum. Synchronisation is always "wait for a specific event or condition under one -overall deadline" (`wait_event`, `wait_until`, 30 s). `settle()` exists only -for asserting that something did *not* happen. Ports are dynamic and -directories are isolated, so tests run in parallel. +overall deadline" (`wait_event`, `wait_until`, 30 s). The deadline covers the +probe as well as the gaps between probes: a call into a wedged agent that +never answers fails the test rather than hanging the process, which is the +difference between a red run and a test binary still burning a core the next +day. `settle()` exists only for asserting that something did *not* happen. +Ports are dynamic and directories are isolated, so tests run in parallel. Several library instances in one process is exactly that. It is **not** a test of several system processes, and is not presented as one.