diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index f30cb3d8..6d16427f 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -25,7 +25,7 @@ use std::collections::{HashMap, HashSet}; use std::net::SocketAddr; use std::thread; use std::time::Duration; -use tracing::{debug, info, warn}; +use tracing::{debug, error, info, warn}; const OPEN_DISCOVERY_RETRY_LIFETIME_MULTIPLIER: u64 = 2; const MAX_PARALLEL_PATH_CANDIDATES_PER_PEER: usize = 4; @@ -38,6 +38,48 @@ fn socket_addr_families_compatible(local: SocketAddr, remote: SocketAddr) -> boo ) } +/// Run a supervised child's async body, then report the child's exit on `tx`, +/// whether the body returned or panicked. +/// +/// A panic would otherwise unwind past the report and leave the node healthy +/// with the child gone. Aborting the task still reports nothing: the abort +/// drops this whole future, so a deliberate stop does not read as a death. +pub(in crate::node) async fn report_exit( + child: Child, + body: impl std::future::Future, + tx: Option>, +) { + use futures::FutureExt; + + if std::panic::AssertUnwindSafe(body) + .catch_unwind() + .await + .is_err() + { + // The panic hook has already printed the payload. + error!(child = ?child, "Supervised child panicked; reporting its exit"); + } + if let Some(tx) = tx { + let _ = tx.send(child).await; + } +} + +/// Run a supervised child's thread body, then report the child's exit on `tx`, +/// whether the body returned or panicked. +pub(in crate::node) fn report_thread( + child: Child, + body: impl FnOnce(), + tx: Option<&tokio::sync::mpsc::Sender>, +) { + if std::panic::catch_unwind(std::panic::AssertUnwindSafe(body)).is_err() { + // The panic hook has already printed the payload. + error!(child = ?child, "Supervised child panicked; reporting its exit"); + } + if let Some(tx) = tx { + let _ = tx.blocking_send(child); + } +} + impl Node { /// Replace the runtime peer list. /// @@ -1674,16 +1716,18 @@ impl Node { let (writer, tun_tx) = device.create_writer(max_mss, self.path_mtu_lookup.clone())?; - // Spawn writer thread. On exit it self-reports - // `Child::Tun` (sync context → `blocking_send`); TUN - // is one compound child, so both threads reporting is - // fine (the FSM de-dups via `up.remove`). + // Spawn writer thread. On exit, including a panic, + // it self-reports `Child::Tun` (sync context → + // `blocking_send`); TUN is one compound child, so + // both threads reporting is fine (the FSM de-dups + // via `up.remove`). let writer_child_tx = self.child_exit_tx.clone(); let writer_handle = thread::spawn(move || { - writer.run(); - if let Some(tx) = &writer_child_tx { - let _ = tx.blocking_send(Child::Tun); - } + report_thread( + Child::Tun, + move || writer.run(), + writer_child_tx.as_ref(), + ); }); // Clone tun_tx for the reader @@ -1695,42 +1739,49 @@ impl Node { tokio::sync::mpsc::channel(tun_channel_size); // Spawn reader thread. Like the writer, it - // self-reports `Child::Tun` on exit (sync context → - // `blocking_send`). Exactly one cfg variant compiles, - // so the single clone is moved into that closure. + // self-reports `Child::Tun` on exit or panic (sync + // context → `blocking_send`). Exactly one cfg + // variant compiles, so the single clone is moved + // into that closure. let transport_mtu = self.transport_mtu(); let path_mtu_lookup = self.path_mtu_lookup.clone(); let reader_child_tx = self.child_exit_tx.clone(); #[cfg(any(target_os = "macos", target_os = "freebsd"))] let reader_handle = thread::spawn(move || { - run_tun_reader( - device, - mtu, - our_addr, - reader_tun_tx, - outbound_tx, - transport_mtu, - path_mtu_lookup, - shutdown_read_fd, + report_thread( + Child::Tun, + move || { + run_tun_reader( + device, + mtu, + our_addr, + reader_tun_tx, + outbound_tx, + transport_mtu, + path_mtu_lookup, + shutdown_read_fd, + ) + }, + reader_child_tx.as_ref(), ); - if let Some(tx) = &reader_child_tx { - let _ = tx.blocking_send(Child::Tun); - } }); #[cfg(not(any(target_os = "macos", target_os = "freebsd")))] let reader_handle = thread::spawn(move || { - run_tun_reader( - device, - mtu, - our_addr, - reader_tun_tx, - outbound_tx, - transport_mtu, - path_mtu_lookup, + report_thread( + Child::Tun, + move || { + run_tun_reader( + device, + mtu, + our_addr, + reader_tun_tx, + outbound_tx, + transport_mtu, + path_mtu_lookup, + ) + }, + reader_child_tx.as_ref(), ); - if let Some(tx) = &reader_child_tx { - let _ = tx.blocking_send(Child::Tun); - } }); self.tun_state = TunState::Active; @@ -1820,23 +1871,24 @@ impl Node { ); // Self-report on exit so the supervisor FSM // routes health when the DNS task dies at - // runtime. On a deliberate stop the task is - // `.abort()`ed before this send; even if it - // fired, the FSM ignores it outside `Running`. + // runtime. The responder never returns, so + // in practice that is a panic. On a + // deliberate stop the task is `.abort()`ed, + // which drops the report with it; even if + // one fired, the FSM ignores it outside + // `Running`. let dns_child_tx = self.child_exit_tx.clone(); - let handle = tokio::spawn(async move { + let handle = tokio::spawn(report_exit( + Child::Dns, crate::upper::dns::run_dns_responder( socket, identity_tx, dns_ttl, reloader, mesh_ifindex, - ) - .await; - if let Some(tx) = dns_child_tx { - let _ = tx.send(Child::Dns).await; - } - }); + ), + dns_child_tx, + )); self.supervisor.dns_identity_rx = Some(identity_rx); self.supervisor.dns_task = Some(handle); self.supervisor.dns_local_addr = Some(local_addr); @@ -2224,13 +2276,10 @@ impl Node { /// reads it to rebuild the teardown set, and aborting an already-finished /// handle there is harmless. /// - /// **Dormant for `Dns` as written.** `run_dns_responder` is an unconditional - /// loop whose every failure arm continues, so it never returns and the - /// `Child::Dns` send that follows it is unreachable — nothing produces the - /// event this consumes. The consumer side is correct and lands here so the - /// producer fix does not have to rediscover it. A responder that *panics* is - /// not covered either way, since the unwind goes past the send rather than - /// through it; that is true of every child producer, not just this one. + /// For `Dns` the event comes only from a panic. `run_dns_responder` is an + /// unconditional loop whose every failure arm continues, so it has no + /// ordinary exit; [`report_exit`] catches a panic in it and reports + /// `Child::Dns`, which is what reaches this. pub(in crate::node) fn retract_child_publications(&mut self, child: Child) { if matches!(child, Child::Dns) { self.supervisor.dns_local_addr.take(); diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 500547e9..b20cc4e5 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -3835,12 +3835,12 @@ async fn dns_responder_serves_a_proxying_embedder() { /// /// Scoped to the helper deliberately, and named for that rather than for the /// scenario: no responder dies here, and deleting the `run_rx_loop` call site -/// leaves this green. Driving a real exit through the loop needs the node moved -/// into a task, which puts `dns_local_addr()` out of reach — and the producer -/// side cannot deliver `Child::Dns` today regardless, since `run_dns_responder` -/// never returns. +/// leaves this green. The rx-loop wiring is pinned separately, by +/// `a_dns_exit_on_the_child_channel_degrades_the_node_and_retracts_its_address`. +/// The producer side reports `Child::Dns` only when the responder panics, since +/// `run_dns_responder` has no ordinary exit. /// -/// What it does pin is the behavior the eventual wiring depends on: the FSM's +/// What it does pin is the behavior the wiring depends on: the FSM's /// `ChildExited` handling only republishes node health, so without this /// retraction `dns_local_addr()` would keep naming a socket nobody is listening /// on, and a proxying embedder would see `.fips` queries silently time out @@ -3871,6 +3871,192 @@ async fn retract_child_publications_clears_the_dns_address() { node.stop().await.unwrap(); } +/// A supervised task that panics is still reported to the supervisor. +/// +/// Without the report the panic ends the task silently: the `JoinError` sits +/// on a handle nobody reads until `stop()`, and the node stays `Running` with +/// the child gone. +#[tokio::test] +async fn report_exit_reports_a_body_that_panics() { + use crate::node::lifecycle::report_exit; + use crate::node::lifecycle::supervisor::Child; + + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + let handle = tokio::spawn(report_exit( + Child::Dns, + async { panic!("the supervised body died") }, + Some(tx), + )); + let joined = handle.await; + + assert_eq!( + rx.recv().await, + Some(Child::Dns), + "a panicking child must still report its exit", + ); + assert!( + joined.is_ok(), + "the panic must be contained by the reporter" + ); +} + +/// A supervised task whose body returns is reported, so a future ordinary +/// exit from the DNS responder needs no further wiring. +#[tokio::test] +async fn report_exit_reports_a_body_that_returns() { + use crate::node::lifecycle::report_exit; + use crate::node::lifecycle::supervisor::Child; + + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + tokio::spawn(report_exit(Child::Dns, async {}, Some(tx))) + .await + .unwrap(); + + assert_eq!(rx.recv().await, Some(Child::Dns)); +} + +/// A deliberate teardown aborts a running child, and that must not read as a +/// death: the abort drops the whole reporting future before any send. +#[tokio::test] +async fn report_exit_stays_silent_when_the_task_is_aborted() { + use crate::node::lifecycle::report_exit; + use crate::node::lifecycle::supervisor::Child; + + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + let (started_tx, started_rx) = tokio::sync::oneshot::channel(); + let handle = tokio::spawn(report_exit( + Child::Dns, + async move { + let _ = started_tx.send(()); + std::future::pending::<()>().await; + }, + Some(tx), + )); + // The body is running, as a live responder is when `stop()` aborts it. + started_rx.await.unwrap(); + handle.abort(); + + assert!(handle.await.unwrap_err().is_cancelled()); + assert_eq!(rx.recv().await, None, "an aborted child must not report"); +} + +/// A supervised thread that panics is still reported to the supervisor. +#[test] +fn report_thread_reports_a_body_that_panics() { + use crate::node::lifecycle::report_thread; + use crate::node::lifecycle::supervisor::Child; + + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + let joined = std::thread::spawn(move || { + report_thread( + Child::Tun, + || panic!("the supervised thread died"), + Some(&tx), + ) + }) + .join(); + + assert_eq!( + rx.try_recv(), + Ok(Child::Tun), + "a panicking thread must still report its exit", + ); + assert!( + joined.is_ok(), + "the panic must be contained by the reporter" + ); +} + +/// A supervised thread whose body returns is reported. +#[test] +fn report_thread_reports_a_body_that_returns() { + use crate::node::lifecycle::report_thread; + use crate::node::lifecycle::supervisor::Child; + + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + std::thread::spawn(move || report_thread(Child::Tun, || {}, Some(&tx))) + .join() + .unwrap(); + + assert_eq!(rx.try_recv(), Ok(Child::Tun)); +} + +/// A node with a UDP transport and a DNS responder on an ephemeral port, and +/// no control socket for the rx loop to bind. +fn dns_config() -> crate::Config { + let mut config = crate::Config::new(); + config.node.control.enabled = false; + config.transports.udp = crate::config::TransportInstances::Single(crate::config::UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + ..Default::default() + }); + config.dns.enabled = true; + config.dns.bind_addr = Some("::1".to_string()); + config.dns.port = Some(0); + config +} + +/// A DNS exit reaching the child channel degrades the node and retracts the +/// responder's published address, through the real rx loop. +/// +/// The exit is queued before the loop starts and the loop is driven in place +/// under a bound, so the node's state is readable afterwards. The shutdown +/// future never fires, which keeps the loop out of the drain path. +#[tokio::test] +async fn a_dns_exit_on_the_child_channel_degrades_the_node_and_retracts_its_address() { + let mut node = make_node_with(dns_config()); + node.start().await.unwrap(); + assert_eq!(node.state(), NodeState::Running); + assert!(node.dns_local_addr().is_some(), "responder came up"); + + node.child_exit_tx + .clone() + .unwrap() + .send(crate::node::lifecycle::supervisor::Child::Dns) + .await + .unwrap(); + let drive = tokio::time::timeout( + Duration::from_millis(500), + node.run_rx_loop_with_shutdown(std::future::pending()), + ) + .await; + assert!( + drive.is_err(), + "the rx loop must still be running: {drive:?}" + ); + + assert_eq!(node.state(), NodeState::Degraded); + assert!( + node.dns_local_addr().is_none(), + "a dead responder must not keep publishing an address to dial", + ); + + node.stop().await.unwrap(); +} + +/// The healthy side of the wiring test: the same drive with no exit queued +/// leaves the node `Running` and its responder's address published. +#[tokio::test] +async fn an_rx_loop_with_no_child_exit_leaves_the_node_running_and_its_dns_address_published() { + let mut node = make_node_with(dns_config()); + node.start().await.unwrap(); + + let drive = tokio::time::timeout( + Duration::from_millis(500), + node.run_rx_loop_with_shutdown(std::future::pending()), + ) + .await; + assert!( + drive.is_err(), + "the rx loop must still be running: {drive:?}" + ); + + assert_eq!(node.state(), NodeState::Running); + assert!(node.dns_local_addr().is_some()); + + node.stop().await.unwrap(); +} + /// `dns.enabled` with a bind that fails must report `None`, not an address. /// /// This is the third state an embedder has to tell apart, and the one that