diff --git a/src/config/mod.rs b/src/config/mod.rs index d2ea8f91..5a3a775b 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -1601,58 +1601,9 @@ node: assert_eq!(metadata.mode() & 0o777, 0o644); } - /// Collect formatted tracing events on the current thread. - /// - /// `resolve_identity` reports its identity-loss conditions only in the - /// log, so the log is what the tests have to assert on. Installed with - /// `tracing::subscriber::with_default`, which is thread-local, so parallel - /// tests do not see each other's events. - #[derive(Clone, Default)] - struct LogCapture(std::sync::Arc>>); - - impl LogCapture { - fn warnings(&self) -> Vec { - self.0 - .lock() - .unwrap() - .iter() - .filter(|line| line.starts_with("WARN")) - .cloned() - .collect() - } - } - - impl tracing_subscriber::Layer for LogCapture { - fn on_event( - &self, - event: &tracing::Event<'_>, - _ctx: tracing_subscriber::layer::Context<'_, S>, - ) { - struct Fields(String); - impl tracing::field::Visit for Fields { - fn record_debug( - &mut self, - field: &tracing::field::Field, - value: &dyn std::fmt::Debug, - ) { - self.0.push_str(&format!(" {}={:?}", field.name(), value)); - } - } - - let mut fields = Fields(event.metadata().level().to_string()); - event.record(&mut fields); - self.0.lock().unwrap().push(fields.0); - } - } - - fn capture_logs(f: impl FnOnce() -> T) -> (T, LogCapture) { - use tracing_subscriber::layer::SubscriberExt; - - let capture = LogCapture::default(); - let subscriber = tracing_subscriber::registry().with(capture.clone()); - let out = tracing::subscriber::with_default(subscriber, f); - (out, capture) - } + // `resolve_identity` reports its identity-loss conditions only in the log, + // so the log is what these tests assert on. + use crate::testutil::capture_logs; #[cfg(unix)] #[test] diff --git a/src/node/handlers/lookup.rs b/src/node/handlers/lookup.rs index eacc4570..9685db00 100644 --- a/src/node/handlers/lookup.rs +++ b/src/node/handlers/lookup.rs @@ -253,6 +253,7 @@ impl Node { // cross-subsystem effects for us to drive. let actions = crate::proto::lookup::on_response_accepted( &mut self.lookup, + response.request_id, &target, response.target_coords, now_ms, @@ -270,6 +271,7 @@ impl Node { for action in actions { match action { LookupAction::CacheCoords { + request_id, target, coords, now_ms, @@ -283,6 +285,7 @@ impl Node { // response at all. if path_mtu < crate::upper::icmp::MIN_ACTIONABLE_PATH_MTU { warn!( + request_id = request_id, target = %self.peer_display_name(&target), path_mtu = path_mtu, floor = crate::upper::icmp::MIN_ACTIONABLE_PATH_MTU, diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index f60dd43d..83656d1d 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -939,7 +939,9 @@ async fn test_originator_ignores_sub_floor_path_mtu_but_still_caches_coords() { response.path_mtu = 64; let payload = &response.encode()[1..]; + let (logs, guard) = crate::testutil::capture_logs_scoped(); node.handle_lookup_response(&from, payload).await; + drop(guard); let now_ms = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) @@ -965,6 +967,21 @@ async fn test_originator_ignores_sub_floor_path_mtu_but_still_caches_coords() { 1, "refusing the annotation must be visible on a counter, not only in a log" ); + + // The counter says a refusal happened; only the correlator says which + // exchange it happened in. The warning is the sole place that pairing + // exists, so the field an operator greps on is asserted here rather than + // left to survive on the strength of compiling. + let floor_warning = logs + .warnings() + .into_iter() + .find(|line| line.contains("below the actionable floor")) + .expect("the sub-floor refusal must be logged at WARN"); + assert!( + floor_warning.contains("request_id=801"), + "the sub-floor warning must carry the correlator of the response it \ + refused, got: {floor_warning}" + ); } #[tokio::test] diff --git a/src/proto/lookup/core.rs b/src/proto/lookup/core.rs index 2eeb5e57..2602fac7 100644 --- a/src/proto/lookup/core.rs +++ b/src/proto/lookup/core.rs @@ -29,7 +29,13 @@ pub(crate) enum LookupAction { /// `bytes` is `Arc`-shared so a fan-out encodes once. SendLink { peer: NodeAddr, bytes: Arc<[u8]> }, /// Cache the verified destination coordinates + path MTU (coord_cache). + /// + /// `request_id` is the correlator of the response this write came from. The + /// shell logs a warning here when `path_mtu` is below the actionable floor, + /// and that line is only correlatable with the rest of the exchange if the + /// action carries the id — by then the response itself is out of scope. CacheCoords { + request_id: u64, target: NodeAddr, coords: crate::TreeCoordinate, now_ms: u64, @@ -262,6 +268,7 @@ pub(crate) fn plan_response_route(lookup: &Lookup, request_id: u64) -> ResponseR /// Verification is the shell's job — this runs only after the proof checked out. pub(crate) fn on_response_accepted( lookup: &mut Lookup, + request_id: u64, target: &NodeAddr, coords: crate::TreeCoordinate, now_ms: u64, @@ -271,6 +278,7 @@ pub(crate) fn on_response_accepted( lookup.pending_lookups.remove(target); vec![ LookupAction::CacheCoords { + request_id, target: *target, coords, now_ms, diff --git a/src/proto/lookup/tests/core.rs b/src/proto/lookup/tests/core.rs index e402e9e0..9bb2e19a 100644 --- a/src/proto/lookup/tests/core.rs +++ b/src/proto/lookup/tests/core.rs @@ -203,7 +203,7 @@ fn on_response_accepted_clears_state_and_emits_effects() { let coords = TreeCoordinate::root(target); let now_ms = 12_345u64; let path_mtu = 1400u16; - let actions = on_response_accepted(&mut lookup, &target, coords, now_ms, path_mtu); + let actions = on_response_accepted(&mut lookup, 7, &target, coords, now_ms, path_mtu); // Success state must be cleared. assert!( @@ -256,6 +256,36 @@ fn on_response_accepted_clears_state_and_emits_effects() { } } +#[test] +fn cache_coords_action_carries_the_request_id_the_response_arrived_with() { + let target = make_node_addr(0x5A); + let mut lookup = empty_lookup(); + + // Distinct from now_ms and path_mtu below, so substituting either for the + // correlator fails rather than coincidentally matching. + let request_id = 0xDEAD_BEEF_0000_0001u64; + let now_ms = 12_345u64; + let path_mtu = 1400u16; + let actions = on_response_accepted( + &mut lookup, + request_id, + &target, + TreeCoordinate::root(target), + now_ms, + path_mtu, + ); + + match &actions[0] { + LookupAction::CacheCoords { request_id: id, .. } => assert_eq!( + *id, request_id, + "the cache-coords action must carry the correlator of the response it \ + was planned from, or the shell's sub-floor path-MTU warning cannot \ + name the exchange its sibling warnings do" + ), + _ => panic!("action[0] must be CacheCoords"), + } +} + #[test] fn poll_pending_no_action_before_first_deadline() { let target = make_node_addr(0x30); diff --git a/src/testutil.rs b/src/testutil.rs index 4765983e..65874819 100644 --- a/src/testutil.rs +++ b/src/testutil.rs @@ -8,3 +8,68 @@ pub(crate) fn make_node_addr(val: u8) -> NodeAddr { bytes[0] = val; NodeAddr::from_bytes(bytes) } + +/// Collects emitted tracing events so a test can assert on a log line. +/// +/// Some behaviour is reported only in the log: a structured field an operator +/// greps on is part of the contract even when no counter or return value +/// carries it. Installed with `tracing::subscriber::with_default`, which is +/// thread-local, so tests running in parallel do not see each other's events. +#[derive(Clone, Default)] +pub(crate) struct LogCapture(std::sync::Arc>>); + +impl LogCapture { + /// Only the captured lines emitted at WARN. + pub(crate) fn warnings(&self) -> Vec { + self.0 + .lock() + .unwrap() + .iter() + .filter(|line| line.starts_with("WARN")) + .cloned() + .collect() + } +} + +impl tracing_subscriber::Layer for LogCapture { + fn on_event( + &self, + event: &tracing::Event<'_>, + _ctx: tracing_subscriber::layer::Context<'_, S>, + ) { + struct Fields(String); + impl tracing::field::Visit for Fields { + fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) { + self.0.push_str(&format!(" {}={:?}", field.name(), value)); + } + } + + let mut fields = Fields(event.metadata().level().to_string()); + event.record(&mut fields); + self.0.lock().unwrap().push(fields.0); + } +} + +/// Run `f` with a capturing subscriber installed, returning its value and the capture. +pub(crate) fn capture_logs(f: impl FnOnce() -> T) -> (T, LogCapture) { + use tracing_subscriber::layer::SubscriberExt; + + let capture = LogCapture::default(); + let subscriber = tracing_subscriber::registry().with(capture.clone()); + let out = tracing::subscriber::with_default(subscriber, f); + (out, capture) +} + +/// Install a capturing subscriber for the rest of the current scope. +/// +/// The async counterpart of [`capture_logs`]: an `async` test cannot wrap its +/// awaits in a closure, so it holds this guard instead and reads the capture +/// once the awaited work has run. +pub(crate) fn capture_logs_scoped() -> (LogCapture, tracing::subscriber::DefaultGuard) { + use tracing_subscriber::layer::SubscriberExt; + + let capture = LogCapture::default(); + let subscriber = tracing_subscriber::registry().with(capture.clone()); + let guard = tracing::subscriber::set_default(subscriber); + (capture, guard) +}