mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 19:18:25 +00:00
Report a supervised child's death to the supervisor even when it panics
The DNS responder's exit report followed a loop that never returns, so it was unreachable, and a panic in the DNS task or either TUN thread unwound past its report. In each case the node stayed healthy with the child gone, and the DNS address stayed published for a socket nobody answered on. Each child body now runs inside a small wrapper that catches a panic, logs it, and reports the exit on either outcome. A deliberate stop still reports nothing, since aborting the task drops the wrapper before it can send. Tests cover both wrappers on panic and on return, the abort case, and the rx loop's handling of a DNS exit: the node degrades and retracts the responder's address.
This commit is contained in:
+101
-52
@@ -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<Output = ()>,
|
||||
tx: Option<tokio::sync::mpsc::Sender<Child>>,
|
||||
) {
|
||||
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<Child>>,
|
||||
) {
|
||||
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();
|
||||
|
||||
+191
-5
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user