diff --git a/src/mmp/mod.rs b/src/mmp/mod.rs index 6340a862..19232bfd 100644 --- a/src/mmp/mod.rs +++ b/src/mmp/mod.rs @@ -485,6 +485,35 @@ impl PathMtuState { // No change (equal or increase not yet confirmed) false } + + /// Forget the source-side path MTU after the path it described is gone. + /// + /// Called when a destination's path is declared broken. The tightened + /// value describes a path that no longer exists, and the increase ladder + /// in [`Self::apply_notification`] (three matching higher values spanning + /// two notification intervals) is far too slow to recover it on the + /// replacement path. Returning to the no-measurement state lets + /// [`Self::seed_source_mtu`] re-derive the value from the outbound + /// transport on the next send, exactly as a fresh session does. + /// + /// Returning to `u16::MAX` is not a licence to send oversized packets: + /// the TUN outbound path caps every packet at `effective_ipv6_mtu()` + /// before it consults the per-destination gate, and that gate is simply + /// inert at `u16::MAX` — the state [`Self::new`] already starts in. If + /// that earlier cap is ever removed or made conditional, this reset stops + /// being safe. + /// + /// Destination-side observation state (`last_observed_mtu`, + /// `observed_changed`, `last_notification_time`) is deliberately left + /// alone: it describes the reverse direction, which this event says + /// nothing about, and clearing it would suppress our notifications to the + /// peer until a fresh observation arrived. + pub fn reset_source_mtu(&mut self) { + self.current_mtu = u16::MAX; + self.consecutive_increase_count = 0; + self.first_increase_time = None; + self.pending_increase_mtu = 0; + } } impl Default for PathMtuState { @@ -626,4 +655,56 @@ owd_window_size: 48 "the increase is not yet due, so the effective MTU must be unchanged" ); } + + #[test] + fn test_reset_source_mtu_returns_to_the_no_measurement_state_and_clears_the_increase_ladder() { + // Two halves. The value half is trivially observable; the ladder half + // is not, because immediately after a reset every reported value is a + // decrease from u16::MAX and the decrease branch reads none of the + // increase counters. Observing them takes a later decrease followed by + // a full three-notification increase sequence: a stale + // `pending_increase_mtu` makes the first of those three take the + // "same value as pending" arm, which never sets `first_increase_time`, + // so the increase can never be accepted at all. + let t0 = Instant::now(); + let mut state = PathMtuState::new(); + + // Tighten, then part-build an increase sequence on top of it. + assert!(state.apply_notification(800, t0)); + state.apply_notification(1400, t0); + state.apply_notification(1400, t0 + Duration::from_secs(11)); + assert_eq!( + state.current_mtu(), + 800, + "precondition: the tightened value is in place and the increase is pending" + ); + + state.reset_source_mtu(); + + assert_eq!( + state.current_mtu(), + u16::MAX, + "the reset must return the source side to the no-measurement state" + ); + + // A fresh decrease, which zeroes the count and the first-increase time + // but would leave a stale pending value behind if the reset had not + // cleared it. + assert!(state.apply_notification(1000, t0 + Duration::from_secs(20))); + assert_eq!(state.current_mtu(), 1000); + + // Three matching higher values spanning two notification intervals. + // This is accepted only if the sequence starts from a cleared ladder. + let t1 = t0 + Duration::from_secs(30); + state.apply_notification(1400, t1); + state.apply_notification(1400, t1 + Duration::from_secs(11)); + state.apply_notification(1400, t1 + Duration::from_secs(21)); + + assert_eq!( + state.current_mtu(), + 1400, + "a reset that left the increase ladder behind strands the first of \ + the three notifications, so the increase is never accepted" + ); + } } diff --git a/src/node/handlers/discovery.rs b/src/node/handlers/discovery.rs index 42bc304e..24fe5703 100644 --- a/src/node/handlers/discovery.rs +++ b/src/node/handlers/discovery.rs @@ -266,21 +266,35 @@ impl Node { if path_mtu_actionable { match self.path_mtu_lookup.write() { Ok(mut map) => match map.get(&fips_addr).copied() { - Some(existing) if existing <= path_mtu => { + Some(existing) if existing.mtu <= path_mtu => { // Keep the tighter learned value; never loosen the // clamp. A reactive MtuExceeded or PathMtuNotification // tighten takes precedence over a looser discovery // estimate (cross-carrier keep-tighter). + // + // This arm deliberately leaves `learned_ms` alone. + // That is what bounds a replayed response: the + // replay of a value already stored takes this arm, + // so the entry still expires at first-write plus + // the TTL rather than being pushed out again on + // every injection. Refreshing the stamp here would + // read as a tidy-up and would silently restore + // indefinite pinning. debug!( target = %self.peer_display_name(&target), fips_addr = %fips_addr, path_mtu = path_mtu, - existing = existing, + existing = existing.mtu, "LookupResponse: keeping tighter existing path_mtu_lookup value" ); } other => { - map.insert(fips_addr, path_mtu); + // The one carrier with no release path, so this is + // the one write that carries a deadline. + map.insert( + fips_addr, + crate::upper::tun::PathMtuEntry::learned(path_mtu, now_ms), + ); debug!( target = %self.peer_display_name(&target), fips_addr = %fips_addr, @@ -749,18 +763,20 @@ impl Node { return; }; match map.get(&fips_addr).copied() { - Some(existing) if existing <= link_mtu => { + Some(existing) if existing.mtu <= link_mtu => { // Keep the tighter learned value; never loosen the clamp. debug!( peer = %self.peer_display_name(peer_addr), fips_addr = %fips_addr, link_mtu = link_mtu, - existing = existing, + existing = existing.mtu, "seed_path_mtu_for_link_peer: keeping tighter existing value" ); } other => { - map.insert(fips_addr, link_mtu); + // Held, not expiring: this describes a link this node can see + // for itself, and it is released when the link goes. + map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(link_mtu)); debug!( peer = %self.peer_display_name(peer_addr), fips_addr = %fips_addr, diff --git a/src/node/handlers/forwarding.rs b/src/node/handlers/forwarding.rs index d2cd9ae3..0c93b318 100644 --- a/src/node/handlers/forwarding.rs +++ b/src/node/handlers/forwarding.rs @@ -47,7 +47,7 @@ impl Node { // Coordinate cache warming from plaintext session-layer headers. // Runs ahead of both the delivery and the TTL decisions: the coords // a peer put on the wire are equally valid whichever way those go. - self.try_warm_coord_cache_ref(&datagram_ref); + self.try_warm_coord_cache_ref(&datagram_ref, payload.len()); // Local delivery: dispatch to session layer handlers without // materializing an owned SessionDatagram payload Vec. Delivery to @@ -181,7 +181,13 @@ impl Node { /// /// Decode failures are logged and silently ignored — they don't block /// forwarding. - fn try_warm_coord_cache_ref(&mut self, datagram: &SessionDatagramRef<'_>) { + /// + /// `outer_len` is the length of the msg_type-stripped `SessionDatagram` + /// buffer this view was decoded from. It is carried in rather than + /// reconstructed from the header size so the malformed-frame byte counter + /// measures the same population as its siblings — which are charged the + /// outer slice — instead of the inner FSP payload. + fn try_warm_coord_cache_ref(&mut self, datagram: &SessionDatagramRef<'_>, outer_len: usize) { let prefix = match FspCommonPrefix::parse(datagram.payload) { Some(p) => p, None => return, @@ -238,11 +244,10 @@ impl Node { // drill-down that separates a short frame from a bad version // or a U-flagged one. The level stays at debug: any peer past // the handshake can drive this at line rate. - self.metrics() - .forwarding - .record_warm_malformed(datagram.payload.len()); + self.metrics().forwarding.record_warm_malformed(outer_len); debug!( len = datagram.payload.len(), + outer_len, version = prefix.version, flags = prefix.flags, "Not a well-formed encrypted FSP message; not warming coords" diff --git a/src/node/handlers/rx_loop.rs b/src/node/handlers/rx_loop.rs index 518d8522..6ca0112e 100644 --- a/src/node/handlers/rx_loop.rs +++ b/src/node/handlers/rx_loop.rs @@ -259,7 +259,10 @@ impl Node { // distinct resources; the `path_mtu_lookup` cache and the // `nostr_discovery` subsystem are deliberately excluded // from `Reloadable` since neither reloads from a backing - // file (see `node::reloadable`). + // file (see `node::reloadable`). The `path_mtu_lookup` + // cache is nevertheless swept on this tick, by + // `purge_expired_path_mtu` below: that is expiry, not + // reload. self.reload_host_map().await; self.poll_pending_connects().await; self.poll_nostr_discovery().await; @@ -269,6 +272,7 @@ impl Node { self.resend_pending_session_handshakes(now_ms).await; self.resend_pending_session_msg3(now_ms).await; self.purge_idle_sessions(now_ms); + self.purge_expired_path_mtu(now_ms); self.process_pending_retries(now_ms).await; self.check_tree_state().await; self.check_bloom_state().await; diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 49b21d65..83420569 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -1201,17 +1201,24 @@ impl Node { let fips_addr = crate::FipsAddress::from_node_addr(src_addr); match self.path_mtu_lookup.write() { Ok(mut map) => match map.get(&fips_addr).copied() { - Some(existing) if existing <= new_mtu => { + Some(existing) if existing.mtu <= new_mtu => { debug!( dest = %peer_name, fips_addr = %fips_addr, new_mtu, - existing, + existing = existing.mtu, "PathMtuNotification: keeping tighter existing path_mtu_lookup value" ); } other => { - map.insert(fips_addr, new_mtu); + // Held, not expiring. This value arrives inside a session, + // and a session's teardown or a PathBroken naming it + // already releases the entry. A deadline here would + // instead recreate the gap this mirror exists to close: a + // peer repeating an identical value on a stable path takes + // the unchanged early-return above and never rewrites the + // entry, so an expiring one would vanish and stay gone. + map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(new_mtu)); debug!( dest = %peer_name, fips_addr = %fips_addr, @@ -1534,17 +1541,23 @@ impl Node { let fips_addr = crate::FipsAddress::from_node_addr(&msg.dest_addr); match self.path_mtu_lookup.write() { Ok(mut map) => match map.get(&fips_addr).copied() { - Some(existing) if existing <= msg.mtu => { + Some(existing) if existing.mtu <= msg.mtu => { debug!( dest = %peer_name, fips_addr = %fips_addr, bottleneck_mtu = msg.mtu, - existing, + existing = existing.mtu, "Reactive MtuExceeded: keeping tighter existing path_mtu_lookup value" ); } other => { - map.insert(fips_addr, msg.mtu); + // Held, not expiring. The admission gate above requires a + // session for the named destination, and that session's + // teardown releases this entry. Nothing re-sends the signal + // once traffic is sized to fit, so a deadline would drop a + // genuine persistent bottleneck and start the next flow at + // the conservative ceiling. + map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(msg.mtu)); debug!( dest = %peer_name, fips_addr = %fips_addr, diff --git a/src/node/handlers/timeout.rs b/src/node/handlers/timeout.rs index c03bc507..d8b11f24 100644 --- a/src/node/handlers/timeout.rs +++ b/src/node/handlers/timeout.rs @@ -4,7 +4,7 @@ use crate::node::Node; use crate::peer::HandshakeState; use crate::transport::LinkId; -use tracing::{debug, info}; +use tracing::{debug, info, warn}; impl Node { /// Check for timed-out handshake connections and clean them up. @@ -294,4 +294,92 @@ impl Node { ); } } + + /// Expire `path_mtu_lookup` entries that nothing else will ever release. + /// + /// The three callers of `path_mtu_lookup_release` all fire on session + /// state, so an entry written by the discovery `LookupResponse` carrier + /// for a destination this node never opens a session with has no release + /// path at all. Keep-tighter then makes one such response permanent: a + /// `path_mtu` of 256 pins that destination's SYN-time MSS clamp at 119 + /// bytes until the process restarts. Only those entries carry a + /// `learned_ms`, and only they are expired here. + /// + /// The deadline is the coordinate cache's own TTL, because the same + /// `LookupResponse` writes both stores: the clamp cannot outlive the + /// route it was learned with, and shortening `node.cache.coord_ttl_secs` + /// shortens this with it. A TTL of zero disables the pass, matching + /// `purge_idle_sessions`. + pub(in crate::node) fn purge_expired_path_mtu(&mut self, now_ms: u64) { + use crate::upper::tun::PathMtuEntry; + + let ttl_ms = self.config().node.cache.coord_ttl_secs * 1000; + if ttl_ms == 0 { + return; // disabled + } + let stale = |e: &PathMtuEntry| { + e.learned_ms + .is_some_and(|at| now_ms.saturating_sub(at) >= ttl_ms) + }; + + // Read-scan first: an ordinary tick expires nothing, and the TUN + // reader and writer take this lock on every packet. + let expired: Vec = match self.path_mtu_lookup.read() { + Ok(map) => map + .iter() + .filter(|(_, e)| stale(e)) + .map(|(a, _)| *a) + .collect(), + Err(e) => { + warn!(error = %e, "path_mtu_lookup read lock poisoned; expiry pass skipped"); + return; + } + }; + if expired.is_empty() { + return; + } + + match self.path_mtu_lookup.write() { + Ok(mut map) => { + for addr in &expired { + // Re-test under the write lock. The read guard was dropped + // before this one was taken, so a fresh value may have + // landed in between; without this the pass would delete a + // value that was just learned. + if map.get(addr).is_some_and(stale) { + map.remove(addr); + } + } + } + Err(e) => { + warn!(error = %e, "path_mtu_lookup write lock poisoned; entries not expired"); + return; + } + } + + // Restore what local configuration knows, the same way + // `path_mtu_lookup_release` does. A tighter remote claim overwrites a + // direct peer's link MTU under keep-tighter, so an expired entry may + // be sitting on top of a seed, and a bare removal would drop that peer + // to the conservative ceiling until its link re-handshakes. + let gone: std::collections::HashSet = expired.iter().copied().collect(); + let seeds: Vec<( + crate::NodeAddr, + crate::transport::TransportId, + crate::transport::TransportAddr, + )> = self + .peers + .iter() + .filter(|(addr, _)| gone.contains(&crate::FipsAddress::from_node_addr(addr))) + .filter_map(|(addr, p)| Some((*addr, p.transport_id()?, p.current_addr()?.clone()))) + .collect(); + for (addr, tid, taddr) in seeds { + self.seed_path_mtu_for_link_peer(&addr, tid, &taddr); + } + + debug!( + expired = expired.len(), + "Expired remote-learned path_mtu_lookup entries" + ); + } } diff --git a/src/node/metrics.rs b/src/node/metrics.rs index a7352b55..e81457b7 100644 --- a/src/node/metrics.rs +++ b/src/node/metrics.rs @@ -151,10 +151,16 @@ impl ForwardingMetrics { /// This is **not** a packet drop. The frame is still delivered or /// forwarded by the normal path; only the opportunistic warm attempt was /// abandoned, so this must never be folded into the rejection family or - /// rendered as dropped traffic. `bytes` is the payload size of the frame - /// whose warm attempt was abandoned, not volume dropped; it is carried so - /// the counter can be rendered as a packets-and-bytes pair like its - /// siblings. + /// rendered as dropped traffic. `bytes` is the outer `SessionDatagram` + /// payload of the frame whose warm attempt was abandoned — the buffer + /// after `dispatch_link_message` strips the msg_type byte, not the inner + /// FSP payload the warm path reads. That is the same basis as + /// [`Self::record_received`] and [`Self::record_reject_bytes`], and it has + /// to be: all three render through one `fwd_value` row on the fipstop + /// Routing tab, where a mixed basis reads as a smaller flood than the one + /// actually arriving. It is volume observed, not volume dropped, and is + /// carried so the counter can be rendered as a packets-and-bytes pair like + /// its siblings. #[inline] pub fn record_warm_malformed(&self, bytes: usize) { self.warm_malformed_packets.inc(); diff --git a/src/node/mod.rs b/src/node/mod.rs index 7a944c6c..96d6663e 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -370,7 +370,7 @@ pub struct Node { /// the TUN reader/writer threads at TCP MSS clamp time so the /// SYN/SYN-ACK clamp can use the smaller of the local-egress floor /// and the learned per-destination path MTU. - path_mtu_lookup: Arc>>, + path_mtu_lookup: crate::upper::tun::PathMtuLookup, // === Transports & Links === /// Active transports (owned by Node). @@ -2528,9 +2528,21 @@ impl Node { self.sessions.remove(remote) } - /// Read the path_mtu_lookup entry for a destination FipsAddress. + /// Read the path MTU stored for a destination FipsAddress. #[cfg(test)] pub(crate) fn path_mtu_lookup_get(&self, fips_addr: &crate::FipsAddress) -> Option { + self.path_mtu_lookup + .read() + .ok() + .and_then(|map| map.get(fips_addr).map(|e| e.mtu)) + } + + /// Read the whole path_mtu_lookup entry, including how it is released. + #[cfg(test)] + pub(crate) fn path_mtu_lookup_entry( + &self, + fips_addr: &crate::FipsAddress, + ) -> Option { self.path_mtu_lookup .read() .ok() @@ -2538,10 +2550,32 @@ impl Node { } /// Write a path_mtu_lookup entry directly (for tests that pre-seed the map). + /// + /// Writes a held entry, which is what a locally derived seed or a + /// session-carried value stores, so pre-seeding does not put a test at + /// the mercy of the expiry pass. Use `path_mtu_lookup_learn` for the + /// discovery-carrier shape. #[cfg(test)] pub(crate) fn path_mtu_lookup_insert(&self, fips_addr: crate::FipsAddress, mtu: u16) { if let Ok(mut map) = self.path_mtu_lookup.write() { - map.insert(fips_addr, mtu); + map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(mtu)); + } + } + + /// Write an expiring path_mtu_lookup entry directly, as the discovery + /// `LookupResponse` carrier does (for tests that drive the expiry pass). + #[cfg(test)] + pub(crate) fn path_mtu_lookup_learn( + &self, + fips_addr: crate::FipsAddress, + mtu: u16, + at_ms: u64, + ) { + if let Ok(mut map) = self.path_mtu_lookup.write() { + map.insert( + fips_addr, + crate::upper::tun::PathMtuEntry::learned(mtu, at_ms), + ); } } @@ -2558,7 +2592,26 @@ impl Node { /// link-peer seed keeps the second while discarding the first; a plain /// removal would silently drop a direct peer back to the conservative /// ceiling until its link re-handshakes. - fn path_mtu_lookup_release(&self, addr: &NodeAddr) { + /// + /// Two stores describe the same dead path, so this releases both: the + /// `FipsAddress`-keyed map the TCP MSS clamp reads, and the session's own + /// source-side path MTU estimate. + fn path_mtu_lookup_release(&mut self, addr: &NodeAddr) { + // The session's own source-side estimate described the same dead path, + // and the increase ladder is the only thing that would ever raise it + // again. Reset it here so the two halves of "this path is gone" stay + // together. The two timeout callers remove the session before calling + // this, so this arm is reached only from the PathBroken route, where + // the session survives the event. + // + // It runs first so the `&mut self.sessions` borrow ends before the + // shared `self.peers` borrow the reseed below takes. + if let Some(entry) = self.sessions.get_mut(addr) + && let Some(mmp) = entry.mmp_mut() + { + mmp.path_mtu.reset_source_mtu(); + } + let fips_addr = crate::FipsAddress::from_node_addr(addr); match self.path_mtu_lookup.write() { Ok(mut map) => { diff --git a/src/node/reloadable.rs b/src/node/reloadable.rs index 4d81acec..16d9b833 100644 --- a/src/node/reloadable.rs +++ b/src/node/reloadable.rs @@ -38,13 +38,19 @@ //! //! - `path_mtu_lookup` is an event-driven cache (`Arc>`) //! populated from observed path-MTU discovery traffic, not loaded from a -//! file. There is nothing to poll. Release is event-driven for the same -//! reason: an entry is dropped when the path it describes is declared +//! file. There is nothing to poll. Release is mostly event-driven for the +//! same reason: an entry is dropped when the path it describes is declared //! invalid (a `PathBroken` report, session idle expiry, or handshake -//! timeout) and the locally derived link MTU is reseeded in its place, so -//! there is no expiry sweep either. (Its read side could adopt the same -//! lock-free `ArcSwap` shape in the future, but that is an optimization, not -//! a reload.) +//! timeout) and the locally derived link MTU is reseeded in its place. All +//! three of those events read session state, which leaves one carrier +//! uncovered: a discovery `LookupResponse` writes an entry for a +//! destination this node may never open a session with. Those entries, and +//! only those, carry a learn time and are expired at the coordinate cache's +//! TTL by `purge_expired_path_mtu` on the same tick, which then reseeds any +//! direct peer whose entry went. Locally derived link MTUs and values +//! learned inside a session carry no deadline. That sweep is expiry, not a +//! reload. (The read side could adopt the same lock-free `ArcSwap` shape in +//! the future, but that is an optimization, not a reload.) //! - `nostr_discovery` is an async spawned subsystem, not a snapshot of disk //! state. //! diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index 45da022f..449aa155 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -1040,6 +1040,119 @@ async fn test_originator_lookup_response_keeps_tighter_path_mtu_lookup() { ); } +/// Build a verified LookupResponse for a fresh target and hand it to the +/// handler, returning the target and its FipsAddress. The identity is +/// registered so the proof verifies and the originator branch is taken. +fn make_verified_lookup_response( + node: &mut Node, + request_id: u64, + path_mtu: u16, +) -> (crate::NodeAddr, crate::FipsAddress, Vec) { + let target_identity = Identity::generate(); + let target = *target_identity.node_addr(); + let target_fips = crate::FipsAddress::from_node_addr(&target); + let root = make_node_addr(0xF0); + let coords = TreeCoordinate::from_addrs(vec![target, root]).unwrap(); + + node.register_identity(target, target_identity.pubkey_full()); + + let proof_data = LookupResponse::proof_bytes(request_id, &target, &coords); + let proof = target_identity.sign(&proof_data); + let mut response = LookupResponse::new(request_id, target, coords, proof); + response.path_mtu = path_mtu; + + (target, target_fips, response.encode()[1..].to_vec()) +} + +#[tokio::test] +async fn test_lookup_response_path_mtu_expires_without_a_session() { + // The discovery carrier writes an entry for a destination this node may + // never open a session with, and all three release callers fire on + // session state. Without a deadline, one response carrying 256 pins that + // destination's SYN-time MSS clamp at 119 bytes until the process + // restarts. + let mut node = make_node(); + let from = make_node_addr(0xAA); + + let (_target, target_fips, body) = make_verified_lookup_response(&mut node, 803, 256); + node.handle_lookup_response(&from, &body).await; + + let entry = node + .path_mtu_lookup_entry(&target_fips) + .expect("precondition: the response wrote an entry, or the rest observes nothing"); + assert_eq!(entry.mtu, 256, "precondition: the annotation was stored"); + let learned_ms = entry + .learned_ms + .expect("the discovery carrier has no release path, so its entry must carry a learn time"); + assert_eq!( + node.session_count(), + 0, + "precondition: no session exists, so nothing but the deadline would ever release this" + ); + + let ttl_ms = node.config().node.cache.coord_ttl_secs * 1000; + assert!(ttl_ms > 0, "precondition: the expiry pass is not disabled"); + + // The healthy half: an entry inside its lifetime must survive an + // ordinary tick, or the clamp loses a value it is entitled to. + node.purge_expired_path_mtu(learned_ms + ttl_ms - 1); + assert_eq!( + node.path_mtu_lookup_get(&target_fips), + Some(256), + "an entry inside its lifetime must survive the expiry pass" + ); + + node.purge_expired_path_mtu(learned_ms + ttl_ms + 1); + assert_eq!( + node.path_mtu_lookup_get(&target_fips), + None, + "past its deadline the entry must go, since nothing else will ever release it" + ); +} + +#[tokio::test] +async fn test_replayed_lookup_response_does_not_extend_the_path_mtu_deadline() { + // The response carries no replay dedupe, so a captured one can be + // re-injected indefinitely. What bounds the damage is that a replay of a + // value already stored takes the keep-tighter arm, which does not touch + // the learn time: each injection buys one TTL, not one per packet. + let mut node = make_node(); + let from = make_node_addr(0xAA); + + let (_target, target_fips, body) = make_verified_lookup_response(&mut node, 804, 256); + node.handle_lookup_response(&from, &body).await; + + let first = node + .path_mtu_lookup_entry(&target_fips) + .expect("precondition: the first response wrote an entry"); + let learned_ms = first + .learned_ms + .expect("precondition: the entry carries a learn time"); + + // Real elapsed wall-clock between replays. The handler stamps its own + // `Self::now_ms()`, so a version that refreshed the deadline would be + // indistinguishable from one that did not if all three landed in the + // same millisecond. + for _ in 0..2 { + std::thread::sleep(std::time::Duration::from_millis(5)); + node.handle_lookup_response(&from, &body).await; + } + + assert_eq!( + node.path_mtu_lookup_entry(&target_fips), + Some(first), + "a replay must leave the entry exactly as it was, deadline included" + ); + + let ttl_ms = node.config().node.cache.coord_ttl_secs * 1000; + node.purge_expired_path_mtu(learned_ms + ttl_ms + 1); + assert_eq!( + node.path_mtu_lookup_get(&target_fips), + None, + "replaying the same value must not push the deadline out" + ); +} + // ============================================================================ // Open-Discovery Sweep — cache-injection unit test // ============================================================================ diff --git a/src/node/tests/forwarding.rs b/src/node/tests/forwarding.rs index 3855c729..4ba0db08 100644 --- a/src/node/tests/forwarding.rs +++ b/src/node/tests/forwarding.rs @@ -330,6 +330,13 @@ async fn test_coord_cache_warming_encrypted_msg_with_coords() { node.coord_cache().get(&dest_addr, now_ms).is_some(), "dest coords not cached from encrypted message" ); + // Changing what the malformed counter charges is close enough to changing + // when it fires that the well-formed case is pinned in the same place. + assert_eq!( + node.metrics().forwarding.warm_malformed_packets.get(), + 0, + "a well-formed CP datagram must not be counted as an abandoned warm attempt" + ); } #[tokio::test] @@ -447,9 +454,29 @@ async fn test_coord_cache_warming_short_inner_payload_is_dropped_not_panic() { 24, "every datagram in both loops must be counted as an abandoned warm attempt" ); - assert!( - node.metrics().forwarding.warm_malformed_bytes.get() > 0, - "the byte counter must move alongside the packet counter" + // The byte counter shares a fipstop row with `received_bytes` and + // `decode_error_bytes`, so it has to measure the same population: the + // outer SessionDatagram payload, not the inner FSP one. Two assertions + // produced two different ways, because a single one cannot tell "the + // basis matches" from "two counters are wrong in the same direction". + // + // Self-derived: every one of the 24 datagrams reaches the warm guard, as + // the two packet counts above already pin, and `record_received` charges + // the identical outer slice. + assert_eq!( + node.metrics().forwarding.warm_malformed_bytes.get(), + node.metrics().forwarding.received_bytes.get(), + "the byte counter must be charged the same outer payload as its \ + siblings on the same row" + ); + // Literal cross-check. The first loop sends inner lengths 4..=11, so + // outer 39..=46, summing to 340; the second sends inner 12..=27, so + // outer 47..=62, summing to 872. Charging the inner payload instead + // reads 60 + 312 = 372, about 15% of the wire volume that arrived. + assert_eq!( + node.metrics().forwarding.warm_malformed_bytes.get(), + 1212, + "24 frames of 39..=46 and 47..=62 outer bytes sum to 1212" ); } diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index f1cabcc4..e1a38f18 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -2534,6 +2534,65 @@ async fn test_path_broken_releases_path_mtu_lookup_entry() { ); } +#[tokio::test] +async fn test_path_broken_resets_the_session_source_path_mtu() { + use crate::node::tests::spanning_tree::make_test_node; + use crate::protocol::PathBroken; + + // The other half of the same release. The map the SYN clamp reads is not + // the only store describing the dead path: the session's own source-side + // estimate gates every outbound packet, and the increase ladder is the + // only thing that would ever raise it again — three matching higher + // notifications spanning two notification intervals, which arrive only + // while the peer is still receiving our datagrams. + let mut tn = make_test_node().await; + + // An Established session, not an Initiating one: an Initiating entry + // carries no MMP state at all, which would make the assertion vacuous. + let remote = Identity::generate(); + install_established_session_with_mmp(&mut tn.node, &remote); + let dest = *remote.node_addr(); + let reporter = NodeAddr::from_bytes([0xBB; 16]); + + tn.node + .get_session_mut(&dest) + .expect("the session was just installed") + .mmp_mut() + .expect("install_established_session_with_mmp initialises MMP state") + .path_mtu + .apply_notification(800, std::time::Instant::now()); + assert_eq!( + tn.node + .get_session(&dest) + .and_then(|e| e.mmp()) + .map(|m| m.path_mtu.current_mtu()), + Some(800), + "precondition: the source-side estimate is tightened before the path dies" + ); + + // Same construction as the sibling test: encode() prepends a 4-byte FSP + // prefix and a msg_type byte, both already consumed by the dispatcher. + let encoded = PathBroken::new(dest, reporter).encode(); + let inner = &encoded[5..]; + assert!( + PathBroken::decode(inner).is_ok(), + "the test body must decode, or the handler returns early and the \ + assertion below observes nothing" + ); + + tn.node.handle_path_broken(&reporter, inner).await; + + assert_eq!( + tn.node + .get_session(&dest) + .and_then(|e| e.mmp()) + .map(|m| m.path_mtu.current_mtu()), + Some(u16::MAX), + "PathBroken must return the source-side estimate to the no-measurement \ + state, so the next send re-seeds it from the outbound transport" + ); +} + #[tokio::test] async fn test_idle_session_purge_keeps_link_peer_path_mtu_seed() { use crate::peer::ActivePeer; @@ -2609,6 +2668,169 @@ async fn test_idle_session_purge_keeps_link_peer_path_mtu_seed() { } } +/// A node with one UDP transport at `mtu`, and `path_mtu_lookup` seeded from +/// that transport's link MTU for a remote address. The remote is deliberately +/// *not* registered in `node.peers`: a test that wants the expiry pass to +/// reseed it must add the `ActivePeer` itself, so that the two tests below +/// can tell "restored by the reseed" apart from "never a candidate". +async fn node_with_link_seed( + mtu: u16, +) -> ( + Node, + crate::NodeAddr, + crate::FipsAddress, + TransportId, + TransportAddr, +) { + use crate::transport::udp::UdpTransport; + use crate::transport::{TransportHandle, packet_channel}; + + let mut node = make_node(); + let (packet_tx, packet_rx) = packet_channel(64); + node.packet_tx = Some(packet_tx); + node.packet_rx = Some(packet_rx); + + let (transport_packet_tx, _transport_packet_rx) = packet_channel(64); + let transport_id = TransportId::new(1); + let mut udp = UdpTransport::new( + transport_id, + Some("udp1".to_string()), + crate::config::UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(mtu), + ..Default::default() + }, + transport_packet_tx, + ); + udp.start_async().await.unwrap(); + node.transports + .insert(transport_id, TransportHandle::Udp(udp)); + + let remote = Identity::generate(); + let remote_addr = *remote.node_addr(); + let remote_fips = crate::FipsAddress::from_node_addr(&remote_addr); + let transport_addr = TransportAddr::from_string("127.0.0.1:2121"); + + node.seed_path_mtu_for_link_peer(&remote_addr, transport_id, &transport_addr); + + (node, remote_addr, remote_fips, transport_id, transport_addr) +} + +#[tokio::test] +async fn test_expired_path_mtu_keeps_the_link_peer_seed() { + use crate::peer::ActivePeer; + + // The same regression the release helper's reseed half exists to + // prevent, reproduced on the expiry path. A tighter discovery value + // overwrites a direct peer's link MTU under keep-tighter, so expiring it + // with a bare removal would silently drop that peer to the conservative + // ceiling until its link re-handshakes. + let (mut node, remote_addr, remote_fips, transport_id, transport_addr) = + node_with_link_seed(1452).await; + + let remote = Identity::generate(); + let peer_identity = PeerIdentity::from_pubkey_full(remote.pubkey_full()); + let mut peer = ActivePeer::new(peer_identity, LinkId::new(7), 0); + peer.set_current_addr(transport_id, transport_addr); + node.peers.insert(remote_addr, peer); + + assert_eq!( + node.path_mtu_lookup_get(&remote_fips), + Some(1452), + "precondition: the direct-link seed is in place" + ); + + let t0 = 5_000_000u64; + node.path_mtu_lookup_learn(remote_fips, 800, t0); + assert_eq!( + node.path_mtu_lookup_get(&remote_fips), + Some(800), + "precondition: a tighter remote-learned value is sitting on the seed" + ); + + let ttl_ms = node.config().node.cache.coord_ttl_secs * 1000; + node.purge_expired_path_mtu(t0 + ttl_ms + 1); + + assert_eq!( + node.path_mtu_lookup_get(&remote_fips), + Some(1452), + "expiring a remote value must restore the local link seed, not leave the \ + destination with no entry at all" + ); + + for transport in node.transports.values_mut() { + transport.stop().await.ok(); + } +} + +#[tokio::test] +async fn test_local_path_mtu_seed_never_expires() { + // Discriminating half of the test above, which on its own cannot tell + // "the seed was restored by the reseed sweep" from "the seed was never a + // candidate for expiry". Here the remote is not in `node.peers`, so there + // is no reseed to mask the difference: a seed that carried a deadline + // would be removed and stay removed. + let (mut node, _remote_addr, remote_fips, _tid, _taddr) = node_with_link_seed(1452).await; + assert_eq!( + node.path_mtu_lookup_get(&remote_fips), + Some(1452), + "precondition: the direct-link seed is in place" + ); + assert_eq!( + node.path_mtu_lookup_entry(&remote_fips) + .and_then(|e| e.learned_ms), + None, + "precondition: a locally derived seed carries no deadline" + ); + + let ttl_ms = node.config().node.cache.coord_ttl_secs * 1000; + node.purge_expired_path_mtu(10 * ttl_ms); + + assert_eq!( + node.path_mtu_lookup_get(&remote_fips), + Some(1452), + "a locally derived link MTU describes a link this node can still see, \ + so no amount of elapsed time may expire it" + ); + + for transport in node.transports.values_mut() { + transport.stop().await.ok(); + } +} + +#[tokio::test] +async fn test_mirrored_notification_path_mtu_survives_a_purge() { + // The proactive mirror exists because a peer repeating an identical value + // on a stable path never rewrites the entry: the handler returns early + // when the session-side MTU is unchanged. An entry from that carrier must + // therefore carry no deadline, or expiring it would permanently reopen + // the gap the mirror closed, for every long-lived multi-hop destination. + let mut node = make_node(); + let remote = Identity::generate(); + let remote_addr = *remote.node_addr(); + let remote_fips = crate::FipsAddress::from_node_addr(&remote_addr); + + install_established_session_with_mmp(&mut node, &remote); + + let body = build_path_mtu_notification_body(1280); + node.handle_session_path_mtu_notification(&remote_addr, &body); + assert_eq!( + node.path_mtu_lookup_get(&remote_fips), + Some(1280), + "precondition: the mirror wrote the notified value" + ); + + let ttl_ms = node.config().node.cache.coord_ttl_secs * 1000; + node.purge_expired_path_mtu(10 * ttl_ms); + + assert_eq!( + node.path_mtu_lookup_get(&remote_fips), + Some(1280), + "a value learned inside a session is released by the session, not by a \ + timer, and must survive any number of expiry passes" + ); +} + // ============================================================================ // Routing-signal admission: the named destination must be an address this // node bound itself, either by initiating toward it or by completing the diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 02d064f0..5467b57e 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -1509,7 +1509,7 @@ async fn test_seed_path_mtu_inserts_when_empty() { .read() .unwrap() .get(&fips_addr) - .copied(); + .map(|e| e.mtu); assert_eq!( stored, Some(1452), @@ -1549,7 +1549,7 @@ async fn test_seeded_narrow_link_mtu_reaches_the_clamp_as_a_tight_ceiling() { .read() .unwrap() .get(&fips_addr) - .copied(), + .map(|e| e.mtu), Some(240), "the seed stores a narrow link MTU unchanged" ); @@ -1587,7 +1587,7 @@ async fn test_seed_path_mtu_keeps_tighter_existing_value() { node.path_mtu_lookup .write() .unwrap() - .insert(fips_addr, 1280); + .insert(fips_addr, crate::upper::tun::PathMtuEntry::held(1280)); node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &transport_addr); @@ -1596,7 +1596,7 @@ async fn test_seed_path_mtu_keeps_tighter_existing_value() { .read() .unwrap() .get(&fips_addr) - .copied(); + .map(|e| e.mtu); assert_eq!( stored, Some(1280), @@ -1626,7 +1626,7 @@ async fn test_seed_path_mtu_tightens_looser_existing_value() { node.path_mtu_lookup .write() .unwrap() - .insert(fips_addr, 1452); + .insert(fips_addr, crate::upper::tun::PathMtuEntry::held(1452)); node.seed_path_mtu_for_link_peer(&peer_addr, TransportId::new(1), &transport_addr); @@ -1635,7 +1635,7 @@ async fn test_seed_path_mtu_tightens_looser_existing_value() { .read() .unwrap() .get(&fips_addr) - .copied(); + .map(|e| e.mtu); assert_eq!( stored, Some(1280), diff --git a/src/upper/tun.rs b/src/upper/tun.rs index b18982ca..970c4179 100644 --- a/src/upper/tun.rs +++ b/src/upper/tun.rs @@ -34,12 +34,50 @@ use tracing::{error, warn}; #[cfg(unix)] use tun::Layer; +/// One `path_mtu_lookup` entry: the MTU the TCP MSS clamp reads, plus how +/// the entry is released. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct PathMtuEntry { + /// Path MTU in bytes. + pub mtu: u16, + /// Unix ms at which a discovery `LookupResponse` supplied this value, or + /// `None` for an entry that some event releases instead of a timer. + /// + /// The discovery carrier is the one with no release path: it writes an + /// entry for a destination this node may never open a session with, and + /// all three callers of `path_mtu_lookup_release` fire on session state. + /// A link MTU this node derived from its own transport, and a value + /// learned inside a session, are both released by an event that says the + /// thing they describe is gone, so they carry no deadline. + pub learned_ms: Option, +} + +impl PathMtuEntry { + /// An entry released by an event rather than a timer: a locally derived + /// link MTU, or a value learned inside a session. + pub fn held(mtu: u16) -> Self { + Self { + mtu, + learned_ms: None, + } + } + + /// A remote party's claim stored at `at_ms` for a destination with no + /// other release path. Expires. + pub fn learned(mtu: u16, at_ms: u64) -> Self { + Self { + mtu, + learned_ms: Some(at_ms), + } + } +} + /// Read-only handle to the per-destination path MTU map. Populated by /// the discovery handler on `LookupResponse`; read by the TUN reader /// (outbound clamp) and writer (inbound clamp) at TCP MSS clamp time. /// Keyed by [`FipsAddress`] (16 bytes, the IPv6 form of a fips peer /// address). -pub type PathMtuLookup = Arc>>; +pub type PathMtuLookup = Arc>>; /// Compute the effective TCP MSS ceiling for a packet given its peer /// address bytes (a 16-byte IPv6 destination on outbound, source on @@ -106,7 +144,7 @@ pub(crate) fn per_flow_max_mss( ); return empty_lookup_ceiling; }; - let Some(&path_mtu) = map.get(&fips_addr) else { + let Some(entry) = map.get(&fips_addr).copied() else { trace!( fips_addr = %fips_addr, global_max_mss, @@ -116,6 +154,7 @@ pub(crate) fn per_flow_max_mss( ); return empty_lookup_ceiling; }; + let path_mtu = entry.mtu; let path_max_mss = mss_ceiling(path_mtu); // The actionable floor deliberately does not apply here. Every value a // remote party supplies is refused before it can reach this map, at the @@ -1511,7 +1550,10 @@ mod tests { // = min(1360, 1143) = 1143. let lookup = empty_lookup(); let addr = fips_addr_with_node_byte(0x42); - lookup.write().unwrap().insert(addr, 1280); + lookup + .write() + .unwrap() + .insert(addr, PathMtuEntry::held(1280)); assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1360), 1143); } @@ -1521,7 +1563,10 @@ mod tests { // global 1143 (the smaller of the two). let lookup = empty_lookup(); let addr = fips_addr_with_node_byte(0x42); - lookup.write().unwrap().insert(addr, 1452); + lookup + .write() + .unwrap() + .insert(addr, PathMtuEntry::held(1452)); // global=1143 (UDP-1280-derived); path_max = 1452-77-60 = 1315. assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1143), 1143); } @@ -1535,7 +1580,10 @@ mod tests { // learned value governs. let lookup = empty_lookup(); let addr = fips_addr_with_node_byte(0x42); - lookup.write().unwrap().insert(addr, 1452); + lookup + .write() + .unwrap() + .insert(addr, PathMtuEntry::held(1452)); // global=1360, path_max = 1452-77-60 = 1315; min(1360, 1315) = 1315. // 1315 > 1143, so the conservative ceiling did NOT clamp here. assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1360), 1315); @@ -1551,7 +1599,10 @@ mod tests { for stored in [0u16, 1, 100, 137] { let lookup = empty_lookup(); let addr = fips_addr_with_node_byte(0x42); - lookup.write().unwrap().insert(addr, stored); + lookup + .write() + .unwrap() + .insert(addr, PathMtuEntry::held(stored)); assert_eq!( per_flow_max_mss(&lookup, addr.as_bytes(), 1360), 1143, @@ -1585,7 +1636,10 @@ mod tests { .map(|&(stored, _)| { let lookup = empty_lookup(); let addr = fips_addr_with_node_byte(0x42); - lookup.write().unwrap().insert(addr, stored); + lookup + .write() + .unwrap() + .insert(addr, PathMtuEntry::held(stored)); (stored, per_flow_max_mss(&lookup, addr.as_bytes(), 1360)) }) .collect(); @@ -1605,10 +1659,10 @@ mod tests { // sub-floor table above. let lookup = empty_lookup(); let addr = fips_addr_with_node_byte(0x42); - lookup - .write() - .unwrap() - .insert(addr, super::super::icmp::MIN_ACTIONABLE_PATH_MTU); + lookup.write().unwrap().insert( + addr, + PathMtuEntry::held(super::super::icmp::MIN_ACTIONABLE_PATH_MTU), + ); assert_eq!(per_flow_max_mss(&lookup, addr.as_bytes(), 1360), 119); } @@ -1637,8 +1691,8 @@ mod tests { let lookup = empty_lookup(); let a = fips_addr_with_node_byte(0x10); let b = fips_addr_with_node_byte(0x20); - lookup.write().unwrap().insert(a, 1280); - lookup.write().unwrap().insert(b, 1452); + lookup.write().unwrap().insert(a, PathMtuEntry::held(1280)); + lookup.write().unwrap().insert(b, PathMtuEntry::held(1452)); assert_eq!(per_flow_max_mss(&lookup, a.as_bytes(), 1360), 1143); assert_eq!(per_flow_max_mss(&lookup, b.as_bytes(), 1360), 1315); }