diff --git a/CHANGELOG.md b/CHANGELOG.md index b5b3a691..2c273402 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -102,6 +102,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 reporting channel at all, so someone with a finding had to guess at an address or open a public issue. +- `node.rate_limit.session_setup_burst` (64) and + `node.rate_limit.session_setup_rate` (16.0), the parameters of the new + per-link-peer session-setup limiter. Setup messages naming a peer this node + is already established with are metered on a second per-link bucket derived + from `node.limits.max_peers`, `node.rekey.after_secs` and + `node.rate_limit.handshake_max_resends`, so raising the peer limit sizes it + automatically. A zero burst or a non-positive rate is rejected at config + validation rather than silently refusing every session. + - `node.rate_limit.established_handshake_burst` and `node.rate_limit.established_handshake_rate`, the parameters of the new established-link msg1 token bucket. Both are optional; omitting them (the @@ -123,6 +132,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 refused together with `--no-build`, which would stamp the marker onto binaries the features never reached. +- `node.rendezvous.nostr.max_concurrent_offers_per_npub`, defaulting to 4, which + bounds how many inbound traversal offers one sender npub may have in flight + at once. It sits inside `max_concurrent_incoming_offers`, which remains the + outer bound, so a value above that is inert; zero is rejected at config + validation, since it refuses every inbound offer rather than disabling the + limit. Existing configurations parse unchanged, the key being optional. + ### Changed - Node health is determined at start completion instead of unconditionally @@ -175,6 +191,42 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 folded into the new tables with a one-time deprecation warning; migrate your `fips.yaml` to the new keys. + +- Inbound traversal offers are now admitted against a per-sender allowance as + well as the global pool. The intake path previously took a permit from a + single semaphore before any identity check, with the sender's npub used only + as a log field, so one sender could hold every slot and deny traversal + onboarding to every other peer for as long as it kept offering. Admission now + takes a per-npub permit and a global permit together. A sender over its own + allowance is refused at debug rather than warn, because the party tripping it + is by definition sending faster than the node wants and a record per + rejection would turn the spam into log volume; the global bound being reached + keeps its warn, which is the operator's signal that the node is genuinely + saturated. **This does not make the pool inexhaustible.** Nostr identities + cost nothing to generate and the signal subscription carries no author + restriction, so an attacker running four throwaway npubs still saturates the + shipped 16-slot pool at an unchanged total offer rate. What the change buys + is that one identity can no longer do it alone, and that the two refusals are + distinguishable in the log. The permit is still held across the whole + attempt; that duration remains inferred from the attempt timeout rather than + measured. + +- Config validation now rejects a `node.rendezvous.nostr.signal_ttl_secs` that + is too large for the configured `replay_window_secs`. A traversal signal is + acceptable over its TTL plus 60s of clock-skew grace on each side, and that + span has to stay strictly inside the replay window, or a session id evicted + from the replay cache on expiry is still fresh enough to be accepted a second + time. The relation was documented but unenforced, so raising the TTL past + 180s silently voided it. The bound is derived from the skew constant rather + than restated, and is checked whether or not nostr discovery is enabled, for + the same reason the rekey rules are. The shipped defaults (120s against 300s) + are unaffected, but a configuration that had widened the TTL or narrowed the + replay window now fails to load, with an error naming the concrete floor for + `replay_window_secs`. The NAT lab's config generator was one such + configuration and its generated `replay_window_secs` moves from 60 to 180. + Note that this covers eviction on expiry only: `seen_sessions_max_entries` + remains a separate capacity-eviction route that no config relation bounds. + - Config validation now rejects two `node.rekey` settings that appear to disable the trigger and in fact fire it continuously. `after_messages` of zero makes the message-count arm true on every poll, because the trigger @@ -252,6 +304,66 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Inbound session-setup messages are now rate limited, keyed on the + authenticated link peer the datagram arrived over. The setup path allocated + a session entry and sent a routed SessionAck for every well-formed message + naming an address it had no entry for, and that address is an envelope field + the sender picks, so one neighbour could grow the session table at whatever + rate it could transmit and buy an ack per entry to a destination of its + choosing. The limiter sits ahead of every send and both handshake + constructions in the handler, so a refused message emits nothing and costs + no cryptography. The key is the link peer rather than the claimed source + address, which is what makes it a limit at all: keying on the source would + hand a single sender a fresh full bucket per forged message. + + Two consequences worth stating rather than discovering. The limiter bounds + each neighbour's contribution and makes a flood attributable; it does not + give the node an absolute ceiling, which stays at roughly + `peers * rate * handshake_timeout_secs`. And a legitimate peer reaching this + node over the *same* link as an attacker shares that attacker's bucket, so + establishment behind a flooded neighbour is refused until it refills. Rekey + and restart traffic is deliberately not subject to that: setup messages + naming an already-established peer draw on a separate per-link bucket, + because suppressed key rotation is silent — nothing errors and no session + drops — and would have shown up only as a flat `rekey_armed`. + +- A forged SessionAck no longer destroys an in-flight session initiation. The + handler removed the session entry to take ownership of the handshake state + and, when the XK msg2 read failed, returned without putting it back. Nothing + in that message is authenticated — the only thing tying it to the initiation + is the datagram's source address, which the sender chooses — so any node able + to reach the victim could cancel any initiation with 57 bytes of the right + length, and hold establishment down by repeating it. The entry is now kept. + Keeping it is not enough on its own, and the second half is the part worth + naming: the msg2 read mixes the sender's ephemeral into the symmetric state + before it authenticates anything, so an entry put back as the failed read + left it holds a handshake that can never read the genuine msg2, which trades + a one-round-trip denial for one lasting the full handshake timeout. The read + is therefore rolled back to its pre-read state before the entry goes back. + The three later failure paths in the same handler still drop the entry: each + is downstream of a msg2 that authenticated, so it is a local failure rather + than a possible forgery. The entry's activity stamp is deliberately not + refreshed on the failure path, so a spray cannot hold a dead initiation past + its original sweep deadline, and a new `ack_handshake_failed` counter makes + the refusals visible at the default log level. Only the XK handshake on this + branch is covered; the additional drop sites in the XX handshake on the + development branch are not. + +- An unauthenticated session msg3 no longer discards a completed key epoch. + The four failure paths in the responder-side rekey arm of the msg3 handler + abandoned the whole rekey, which nulls a `pending` session sitting beside the + handshake, when only the handshake had failed. That `pending` session is the + epoch the real peer may already have cut over to, so discarding it kills the + reverse direction until the session idles out. Two unauthenticated messages + reached it: a forged setup message arms a handshake beside a completed rekey + once that rekey has waited a full idle timeout for a peer that never appeared + on the new epoch, and any garbage msg3 of the right length then finishes the + job. All four paths now abandon only the handshake. The neighbouring comment + in the setup handler, which claimed no path in that handler discards the + pending keys, has been corrected rather than left to be trusted: the + dual-initiation arm still calls the destructive form on an unauthenticated + msg1, and that is now named at the site as the one path that does. + - The packaged macOS daemon now recreates and binds its control socket at `/var/run/fips/control.sock` instead of falling through to the shared `/tmp/fips-control.sock` path after boot. The previous fallback also made the diff --git a/docs/design/fips-nostr-discovery.md b/docs/design/fips-nostr-discovery.md index 5b7c93fc..528cffc4 100644 --- a/docs/design/fips-nostr-discovery.md +++ b/docs/design/fips-nostr-discovery.md @@ -212,11 +212,20 @@ NIP-17 DM relay list (kind 10050), and falls back to `dm_relays` if the inbox-relays fetch fails. Each side also publishes its own inbox relay list on startup so dialers can discover it. -On the receiving side, an inbound semaphore bounds concurrent offer -processing at `max_concurrent_incoming_offers`. When the semaphore is -full, the offer is dropped with a warn log; this is the primary guard -against offer-spam from a misbehaving or compromised relay. A -`sessionId` replay cache (bounded by `seen_sessions_max_entries`, with +On the receiving side, admission is a pair of bounds taken together: a +per-sender allowance of `max_concurrent_offers_per_npub`, keyed on the +npub that signed the gift wrap, nested inside a global +`max_concurrent_incoming_offers`. A sender over its own allowance is +refused at debug, since by definition it is sending faster than the node +wants and a record per rejection would turn the spam into log volume; the +global bound being reached is the operator-visible warn, because that one +says the node is genuinely saturated. Together they keep one identity +from holding the whole pool. They do not make the pool inexhaustible: +nostr identities are free to generate, so an attacker running +`ceil(max_concurrent_incoming_offers / max_concurrent_offers_per_npub)` +throwaway npubs still saturates it at the same total offer rate. Raising +the attacker's cost beyond keypairs would mean pricing the offer itself. +A `sessionId` replay cache (bounded by `seen_sessions_max_entries`, with entries valid for `replay_window_secs`) rejects duplicates. The responder runs its own STUN query and replies with a @@ -306,6 +315,7 @@ machinery: | Mechanism | Default | What it prevents | Behavior at limit | | --- | --- | --- | --- | | Offer semaphore (`max_concurrent_incoming_offers`) | 16 | CPU and memory exhaustion from offer spam on DM relays. | Warn log, offer dropped. | +| Per-npub offer allowance (`max_concurrent_offers_per_npub`) | 4 | One sender identity holding every offer slot and denying traversal onboarding to everyone else. Does not prevent the same denial from several throwaway npubs. | Debug log, offer dropped. | | Advert cache (`advert_cache_max_entries`) | 2048 | Memory growth from ambient advert traffic under `policy: open`. | LRU-by-expiry eviction. | | Seen-sessions (`seen_sessions_max_entries`) | 2048 | Replay of stale `sessionId` values. | Oldest entry evicted. | | Signal TTL (`signal_ttl_secs`) | 120 s | Indefinite in-flight offers on relays. | Expired offers rejected at validation. | diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index c5333b33..765efbf4 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -136,6 +136,8 @@ Handshake rate limiting protects against DoS on the Noise IK handshake path. | `node.rate_limit.handshake_max_resends` | u32 | `5` | Max resends per handshake attempt | | `node.rate_limit.established_handshake_burst` | u32 | derived | Burst capacity of the established-link bucket. Derived default is `node.limits.max_peers` (128) | | `node.rate_limit.established_handshake_rate` | f64 | derived | Refill rate of that bucket. Derived default is `(max_peers / max(node.rekey.after_secs, 1)) * (1 + handshake_max_resends)`, floored at 1.0/s — 6.4/s at shipped defaults | +| `node.rate_limit.session_setup_burst` | u32 | `64` | Per-link-peer burst for inbound session-setup messages that would open a new session | +| `node.rate_limit.session_setup_rate` | f64 | `16.0` | Per-link-peer refill rate for those messages, in tokens per second | Msg1 whose source matches an established link (rekey and restart maintenance traffic) draws on a second bucket rather than competing with @@ -150,6 +152,25 @@ burst and 16.4/s at shipped defaults, of which the established half is reachable only by a source that already matches a live link. Size against the sum when budgeting handshake crypto load for a host. +The `session_setup_*` pair is a separate limiter on the session layer, not +the link layer. It is keyed on the authenticated link peer a session datagram +arrived over, so each neighbour gets its own budget and a flood is +attributable. Setup messages naming a peer this node is already established +with (inbound rekey and restart traffic) draw on a second per-link bucket +derived from `max_peers`, `node.rekey.after_secs` and `handshake_max_resends`, +exactly as `established_handshake_*` is, so a stranger flood cannot suppress +rekey traffic sharing the link. + +At the defaults one neighbour can force at most +`session_setup_rate * handshake_timeout_secs` half-open entries (480) and +`session_setup_rate * (1 + handshake_max_resends)` acks per second (96). The +node-wide ceiling is still that times the peer count, since the limiter bounds +each neighbour rather than the aggregate. A legitimate peer whose traffic +reaches this node over the *same* link as an attacker's shares that +attacker's stranger bucket, so establishment behind a flooded neighbour is +refused until the bucket refills; the initiator's own resend schedule (1s, 2s, +4s, 8s, 16s) covers a short drain. + ### Retry / Backoff (`node.retry.*`) Connection retry with exponential backoff. @@ -207,6 +228,7 @@ inert otherwise. | `node.discovery.nostr.policy` | string | `"configured_only"` | Advert discovery policy: `disabled`, `configured_only`, `open` | | `node.discovery.nostr.open_discovery_max_pending` | usize | `64` | Max open-discovery peers queued in outbound retry/connection state at once | | `node.discovery.nostr.max_concurrent_incoming_offers` | usize | `16` | Max concurrent inbound traversal offers processed at once (rate limit against offer spam) | +| `node.discovery.nostr.max_concurrent_offers_per_npub` | usize | `4` | Max concurrent inbound traversal offers accepted from any one sender npub, so a single identity cannot hold the whole pool. Sits inside `max_concurrent_incoming_offers`, which stays the outer bound; a larger value is inert. Zero is rejected, since it refuses every inbound offer rather than disabling the limit | | `node.discovery.nostr.advert_cache_max_entries` | usize | `2048` | Max cached overlay adverts retained from relay traffic | | `node.discovery.nostr.seen_sessions_max_entries` | usize | `2048` | Max seen-session IDs retained for replay detection | | `node.discovery.nostr.advertise` | bool | `true` | Publish local endpoint adverts | @@ -931,6 +953,8 @@ node: handshake_resend_interval_ms: 1000 handshake_resend_backoff: 2.0 handshake_max_resends: 5 + session_setup_burst: 64 + session_setup_rate: 16.0 retry: max_retries: 5 base_interval_secs: 5 diff --git a/src/config/mod.rs b/src/config/mod.rs index 0eb4b0f5..d2ea8f91 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -26,6 +26,7 @@ mod peer; mod transport; use crate::node::REKEY_JITTER_SECS; +use crate::nostr::FRESHNESS_SKEW_TOLERANCE_MS; use crate::upper::config::{DnsConfig, TunConfig}; use crate::{Identity, IdentityError}; use serde::{Deserialize, Serialize}; @@ -1113,6 +1114,67 @@ impl Config { ))); } + // The per-link session-setup bucket. The same trap as above in a + // non-optional field: zero does not disable the limiter, it refuses + // every inbound setup message and so refuses every session. + if rl.session_setup_burst == 0 { + return Err(ConfigError::Validation( + "`node.rate_limit.session_setup_burst` is 0, which refuses every inbound SessionSetup rather than disabling the limit. \ + Set a positive burst; a very large value effectively disables it." + .to_string(), + )); + } + + let setup_rate = rl.session_setup_rate; + if !(setup_rate.is_finite() && setup_rate > 0.0) { + return Err(ConfigError::Validation(format!( + "`node.rate_limit.session_setup_rate` is {setup_rate}, but must be a finite value greater than 0; \ + a non-positive or non-finite refill rate never replenishes a link's setup bucket, so that link stops establishing sessions once its initial burst is spent." + ))); + } + + // The freshness window backstops session-id replay protection: an + // offer evicted from the replay cache must already be too old to pass + // the freshness check, or it can be accepted a second time. A signal + // is acceptable over `signal_ttl_secs` plus the skew tolerance on each + // side, so that span has to stay strictly inside `replay_window_secs`. + // Checked regardless of `nostr.enabled`, for the reason the rekey + // block above gives: enabling the feature later must not surface a + // config error at a surprising moment. + let skew_secs = FRESHNESS_SKEW_TOLERANCE_MS / 1000; + let freshness_window_secs = nostr.signal_ttl_secs.saturating_add(2 * skew_secs); + if freshness_window_secs >= nostr.replay_window_secs { + return Err(ConfigError::Validation(format!( + "`node.rendezvous.nostr.signal_ttl_secs` is {}, which with {skew_secs}s of clock-skew grace on each side makes a traversal signal acceptable over a {freshness_window_secs}s window, \ + but `node.rendezvous.nostr.replay_window_secs` is {}. \ + The freshness window must be strictly narrower than the replay window, or a session id evicted from the replay cache is still fresh enough to be accepted a second time. \ + Raise `replay_window_secs` above {freshness_window_secs}, or lower `signal_ttl_secs`.", + nostr.signal_ttl_secs, nostr.replay_window_secs + ))); + } + + // Zero here is the same trap as `established_handshake_burst`: it + // reads as "no limit" and in fact refuses every inbound offer. The + // upper bound exists because the per-npub semaphore is built lazily + // inside the intake path rather than at startup, so an oversized + // value would panic there instead of failing loudly at load. + if nostr.max_concurrent_offers_per_npub == 0 { + return Err(ConfigError::Validation( + "`node.rendezvous.nostr.max_concurrent_offers_per_npub` is 0, which refuses every inbound traversal offer rather than disabling the per-sender limit. \ + Omit the key for the default, or set a positive allowance; `max_concurrent_incoming_offers` remains the outer bound." + .to_string(), + )); + } + + if nostr.max_concurrent_offers_per_npub > tokio::sync::Semaphore::MAX_PERMITS { + return Err(ConfigError::Validation(format!( + "`node.rendezvous.nostr.max_concurrent_offers_per_npub` is {}, which exceeds the maximum {} permits a semaphore can hold. \ + Use a value at or below `max_concurrent_incoming_offers`, which is the outer bound anything larger is inert against.", + nostr.max_concurrent_offers_per_npub, + tokio::sync::Semaphore::MAX_PERMITS + ))); + } + Ok(()) } @@ -2443,6 +2505,74 @@ node: assert!(err.to_string().contains("after_secs")); } + #[test] + fn test_validate_signal_ttl_at_or_above_the_replay_window_margin_rejected() { + // 180 is the boundary: 180 + 2 * 60 = 300, which is not strictly less + // than the default 300s replay window. + for signal_ttl_secs in [180, 181, 3600, u64::MAX] { + let mut config = Config::default(); + config.node.rendezvous.nostr.signal_ttl_secs = signal_ttl_secs; + + match config.validate() { + Err(e) => { + let msg = e.to_string(); + assert!(msg.contains("signal_ttl_secs"), "got: {msg}"); + assert!(msg.contains("replay_window_secs"), "got: {msg}"); + } + Ok(()) => panic!("signal_ttl_secs = {signal_ttl_secs} should be rejected"), + } + } + } + + #[test] + fn test_validate_signal_ttl_just_inside_the_replay_window_margin_accepted() { + let mut config = Config::default(); + config.node.rendezvous.nostr.signal_ttl_secs = 179; + + config + .validate() + .expect("179 + 2 * 60 = 299 leaves the freshness window inside the 300s replay window"); + } + + #[test] + fn test_validate_per_npub_offer_allowance_of_zero_rejected() { + let mut config = Config::default(); + config.node.rendezvous.nostr.max_concurrent_offers_per_npub = 0; + + let err = config.validate().expect_err("validation should fail"); + assert!( + err.to_string().contains("max_concurrent_offers_per_npub"), + "got: {err}" + ); + } + + #[test] + fn test_validate_per_npub_offer_allowance_of_one_accepted() { + let mut config = Config::default(); + config.node.rendezvous.nostr.max_concurrent_offers_per_npub = 1; + + config + .validate() + .expect("an allowance of one offer per sender is restrictive but well defined"); + } + + #[test] + fn test_validate_shipped_defaults_satisfy_the_freshness_invariant() { + Config::default() + .validate() + .expect("the shipped defaults must satisfy every validation rule"); + + // Stated against the constant rather than a literal, so this reds if + // anyone changes FRESHNESS_SKEW_TOLERANCE_MS or either default without + // re-checking the relation they jointly have to satisfy. + let defaults = Config::default(); + let nostr = &defaults.node.rendezvous.nostr; + assert!( + nostr.signal_ttl_secs + 2 * (FRESHNESS_SKEW_TOLERANCE_MS / 1000) + < nostr.replay_window_secs + ); + } + #[test] fn test_outbound_only_forces_ephemeral_bind() { let cfg = UdpConfig { diff --git a/src/config/node.rs b/src/config/node.rs index c3a283bc..2847e217 100644 --- a/src/config/node.rs +++ b/src/config/node.rs @@ -91,6 +91,27 @@ pub struct RateLimitConfig { /// `node.rekey.after_secs` and `handshake_max_resends`. #[serde(default)] pub established_handshake_rate: Option, + /// Per-link-peer burst capacity for inbound FSP SessionSetup messages + /// that would open a new session (`node.rate_limit.session_setup_burst`). + /// + /// 64 absorbs a legitimate reconnect burst arriving behind one + /// neighbour. It bounds nothing on its own; `session_setup_rate` is what + /// bounds the sustained cost. + #[serde(default = "RateLimitConfig::default_session_setup_burst")] + pub session_setup_burst: u32, + /// Per-link-peer refill rate for those messages, in tokens per second + /// (`node.rate_limit.session_setup_rate`). + /// + /// 16/s caps one neighbour's forced half-open occupancy at + /// `rate * handshake_timeout_secs` (480 entries at defaults) and its ack + /// amplification at `rate * (1 + handshake_max_resends)` (96 acks/s). + /// + /// Setup messages naming a peer this node is already established with + /// are metered on a separate per-link bucket, derived from + /// `node.limits.max_peers` exactly as `established_handshake_*` is, so a + /// stranger flood cannot suppress rekey traffic sharing the link. + #[serde(default = "RateLimitConfig::default_session_setup_rate")] + pub session_setup_rate: f64, } impl Default for RateLimitConfig { @@ -104,6 +125,8 @@ impl Default for RateLimitConfig { handshake_max_resends: 5, established_handshake_burst: None, established_handshake_rate: None, + session_setup_burst: 64, + session_setup_rate: 16.0, } } } @@ -127,6 +150,12 @@ impl RateLimitConfig { fn default_handshake_max_resends() -> u32 { 5 } + fn default_session_setup_burst() -> u32 { + 64 + } + fn default_session_setup_rate() -> f64 { + 16.0 + } } /// Retry/backoff configuration (`node.retry.*`). @@ -374,6 +403,11 @@ pub struct NostrRendezvousConfig { /// Acts as a rate limit against offer spam from relays. #[serde(default = "NostrRendezvousConfig::default_max_concurrent_incoming_offers")] pub max_concurrent_incoming_offers: usize, + /// Max concurrent inbound traversal offers accepted from any one sender + /// npub. Sits inside `max_concurrent_incoming_offers`, which remains the + /// outer bound. + #[serde(default = "NostrRendezvousConfig::default_max_concurrent_offers_per_npub")] + pub max_concurrent_offers_per_npub: usize, /// Max cached overlay adverts retained from relay traffic. /// Bounds memory under ambient advert volume. #[serde(default = "NostrRendezvousConfig::default_advert_cache_max_entries")] @@ -462,6 +496,7 @@ impl Default for NostrRendezvousConfig { policy: NostrRendezvousPolicy::default(), open_discovery_max_pending: Self::default_open_discovery_max_pending(), max_concurrent_incoming_offers: Self::default_max_concurrent_incoming_offers(), + max_concurrent_offers_per_npub: Self::default_max_concurrent_offers_per_npub(), advert_cache_max_entries: Self::default_advert_cache_max_entries(), seen_sessions_max_entries: Self::default_seen_sessions_max_entries(), attempt_timeout_secs: Self::default_attempt_timeout_secs(), @@ -527,6 +562,20 @@ impl NostrRendezvousConfig { 16 } + /// Four, derived rather than picked. The initiator side already admits at + /// most one in-flight traversal per peer npub, so one concurrent offer per + /// peer is the honest steady state. An offer is published to and consumed + /// from the whole DM relay set, three URLs by default, and whether the + /// notification stream deduplicates one event delivered by three relays is + /// not established here — if it does not, one honest offer can present as + /// three near-simultaneous admissions before the replay check rejects the + /// duplicates. Four is that worst-case fan-out plus one, so a retry + /// overlapping a still-timing-out attempt is still admitted, and it is a + /// quarter of the default global bound. + fn default_max_concurrent_offers_per_npub() -> usize { + 4 + } + fn default_advert_cache_max_entries() -> usize { 2048 } diff --git a/src/instr/recorder.rs b/src/instr/recorder.rs index 1aef0fa5..7addb7f1 100644 --- a/src/instr/recorder.rs +++ b/src/instr/recorder.rs @@ -57,6 +57,7 @@ pub(crate) enum Step { ResendPendingSessionHandshakes, ResendPendingSessionMsg3, PurgeIdleSessions, + PurgeExpiredPathMtu, ProcessPendingRetries, CheckTreeState, CheckBloomState, @@ -94,6 +95,7 @@ pub(crate) const STEPS: [Step; N_STEPS] = [ Step::ResendPendingSessionHandshakes, Step::ResendPendingSessionMsg3, Step::PurgeIdleSessions, + Step::PurgeExpiredPathMtu, Step::ProcessPendingRetries, Step::CheckTreeState, Step::CheckBloomState, @@ -126,6 +128,7 @@ impl Step { Step::ResendPendingSessionHandshakes => "resend_pending_session_handshakes", Step::ResendPendingSessionMsg3 => "resend_pending_session_msg3", Step::PurgeIdleSessions => "purge_idle_sessions", + Step::PurgeExpiredPathMtu => "purge_expired_path_mtu", Step::ProcessPendingRetries => "process_pending_retries", Step::CheckTreeState => "check_tree_state", Step::CheckBloomState => "check_bloom_state", diff --git a/src/node/dataplane/forwarding.rs b/src/node/dataplane/forwarding.rs index 75731fbd..7c514f86 100644 --- a/src/node/dataplane/forwarding.rs +++ b/src/node/dataplane/forwarding.rs @@ -27,7 +27,7 @@ impl Node { /// has already had its msg_type byte stripped by dispatch. pub(in crate::node) async fn handle_session_datagram( &mut self, - _from: &NodeAddr, + from: &NodeAddr, payload: &[u8], incoming_ce: bool, ) { @@ -51,7 +51,7 @@ impl Node { // coords a peer put on the wire are equally valid whichever way those // go, and the only arrivals this newly warms from are those with an // exhausted TTL, whose every insert is already achievable at TTL 1. - self.try_warm_coord_cache_ref(&datagram_ref); + self.try_warm_coord_cache_ref(&datagram_ref, payload.len()); // Pre-resolve the next hop only for datagrams the core can actually // forward: not locally destined, and carrying a TTL that survives the @@ -107,6 +107,7 @@ impl Node { self.metrics().forwarding.record_delivered(payload.len()); self.handle_session_payload( &datagram_ref.src_addr, + from, datagram_ref.payload, datagram_ref.path_mtu, incoming_ce, @@ -219,7 +220,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, @@ -276,11 +283,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/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index 75ec0797..78bf3e8c 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -355,7 +355,10 @@ impl Node { // distinct resources; the `path_mtu_lookup` cache and the // `nostr_rendezvous` 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. instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ReloadHostMap, self.reload_host_map().await); instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PollPendingConnects, @@ -374,6 +377,8 @@ impl Node { self.resend_pending_session_msg3(now_ms).await); instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PurgeIdleSessions, self.purge_idle_sessions(now_ms)); + instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::PurgeExpiredPathMtu, + self.purge_expired_path_mtu(now_ms)); instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::ProcessPendingRetries, self.process_pending_retries(now_ms).await); instr_step!(instr_on, crate::instr::Domain::Tick, crate::instr::Step::CheckTreeState, diff --git a/src/node/handlers/lookup.rs b/src/node/handlers/lookup.rs index 556ec18b..eacc4570 100644 --- a/src/node/handlers/lookup.rs +++ b/src/node/handlers/lookup.rs @@ -296,7 +296,11 @@ impl Node { .insert_with_path_mtu(target, coords, now_ms, path_mtu); } } - LookupAction::WritePathMtu { target, path_mtu } => { + LookupAction::WritePathMtu { + target, + now_ms, + path_mtu, + } => { // Refused as absent on the CacheCoords arm the core always // pairs with this one, so there is nothing to mirror; the // warning and the counter are emitted there, once. @@ -308,22 +312,37 @@ impl Node { let fips_addr = crate::FipsAddress::from_node_addr(&target); 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, @@ -753,18 +772,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/session.rs b/src/node/handlers/session.rs index a2d0c607..3b7155ef 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -7,6 +7,7 @@ use crate::NodeAddr; use crate::node::handlers::mmp::format_throughput; +use crate::node::rate_limit::Msg1Class; use crate::node::reject::{RejectReason, SessionReject}; use crate::node::session::{EndToEndState, EpochSlot, SessionEntry}; use crate::node::{Node, NodeError}; @@ -100,9 +101,15 @@ impl Node { /// - Phase 0x3 → SessionMsg3 (XK handshake msg3) /// - Phase 0x0 + U flag → plaintext error signal (CoordsRequired/PathBroken) /// - Phase 0x0 + !U → encrypted session message (data, reports, etc.) + /// + /// `src_addr` is the datagram's claimed source, an envelope field the + /// sender chooses. `link_peer` is the authenticated FMP peer the datagram + /// arrived over, and is the only identity on this path worth keying a + /// limiter on. pub(in crate::node) async fn handle_session_payload( &mut self, src_addr: &NodeAddr, + link_peer: &NodeAddr, payload: &[u8], path_mtu: u16, ce_flag: bool, @@ -122,7 +129,7 @@ impl Node { match prefix.phase { FSP_PHASE_MSG1 => { - self.handle_session_setup(src_addr, inner).await; + self.handle_session_setup(src_addr, link_peer, inner).await; } FSP_PHASE_MSG2 => { self.handle_session_ack(src_addr, inner).await; @@ -470,7 +477,12 @@ impl Node { /// The remote node wants to establish an end-to-end session with us. /// We create an XK responder handshake, process msg1, send SessionAck with msg2, /// and transition to AwaitingMsg3. - async fn handle_session_setup(&mut self, src_addr: &NodeAddr, inner: &[u8]) { + async fn handle_session_setup( + &mut self, + src_addr: &NodeAddr, + link_peer: &NodeAddr, + inner: &[u8], + ) { let setup = match SessionSetup::decode(inner) { Ok(s) => s, Err(e) => { @@ -488,6 +500,40 @@ impl Node { return; } + // Meter the setup before anything is spent on it. This sits ahead of + // every `send_session_datagram` call in the handler — the duplicate + // ack resend, the rekey ack and the fresh-setup ack alike — so a + // refused msg1 emits nothing at all, which is what bounds the ack + // amplification. It also precedes both responder handshake + // constructions, so a refusal costs no crypto. Moving it below the + // existing-entry lookup would leave the duplicate resend unmetered. + // + // The class is read from the session table but the *key* is the link + // peer: `src_addr` is chosen by the sender, so keying on it would be + // no limit at all. A setup naming an established peer cannot grow the + // table and is metered separately, so that a stranger flood over a + // shared link cannot stop that peer's rekey from arming. + let class = if self + .sessions + .get(src_addr) + .is_some_and(|e| e.is_established()) + { + Msg1Class::EstablishedLink + } else { + Msg1Class::Stranger + }; + if !self.setup_rate_limiter.try_admit(link_peer, class) { + debug!( + link_peer = %self.peer_display_name(link_peer), + src = %self.peer_display_name(src_addr), + ?class, + "SessionSetup rate limited" + ); + self.stats_mut() + .record_reject(RejectReason::Session(SessionReject::SetupRateLimited)); + return; + } + // Check for existing session with this remote if let Some(existing) = self.sessions.get(src_addr) { if existing.is_initiating() { @@ -538,8 +584,13 @@ impl Node { // on the new epoch, it stops vetoing: a peer that restarted, // or one whose own cycle lapsed, would otherwise be refused // for as long as our own sends kept the session from idling - // out. The pending keys are not discarded here either way — - // only an authenticated msg3 replaces them. + // out. The pending keys survive both outcomes of this test: + // the veto returns without touching them, and the + // fall-through only arms a handshake beside them. Adopting + // them is still an authenticated msg3's job alone. No arm + // below drops one either: the dual-initiation arm did until + // it was narrowed to abandon only the handshake, for the + // reason recorded at that site. let pending_outranks = existing.pending_new_session().is_some() && !existing.pending_stale( Self::now_ms(), @@ -560,13 +611,28 @@ impl Node { .record_reject(RejectReason::Session(SessionReject::RekeyTiebreak)); return; } - // We lose — abandon our rekey, become responder below. + // We lose — abandon the armed handshake, become responder + // below. + // + // `abandon_handshake`, not `abandon_rekey`: the gate + // above is `has_rekey_in_progress`, which says only that + // *some* handshake is armed, not that we armed it. A + // handshake the peer armed carries `rekey_initiator == + // false` and can sit beside a completed epoch that a + // stale `pending_outranks` no longer vetoes, so a + // stranger reaches this line with two unauthenticated + // setup messages: one to arm the handshake, one to lose + // the tie-break against it. Dropping the pending session + // there kills the epoch the peer may already have cut + // over to. Only the handshake is ours to discard, and + // discarding it costs nothing, since an armed handshake + // holds no key material either endpoint is using. debug!( src = %self.peer_display_name(src_addr), "Dual FSP rekey initiation: we lose (larger addr), abandoning ours" ); let entry = self.sessions.get_mut(src_addr).unwrap(); - entry.abandon_rekey(); + entry.abandon_handshake(); self.stats_mut() .record_reject(RejectReason::Session(SessionReject::RekeyYielded)); } else if pending_outranks { @@ -722,6 +788,18 @@ impl Node { }; // Rekey path: entry is Established with rekey_state + // + // `abandon_rekey` below, not `abandon_handshake` as in the responder + // arm of `handle_session_msg3`, and the difference rests on an + // invariant rather than on a different judgement: an entry with + // `rekey_initiator` set holds no pending session, so the two calls + // are the same action here. `set_rekey_state(_, true)` has one + // caller, `initiate_session_rekey`, which `check_session_rekey` + // never reaches for an entry holding a pending session; and + // `set_pending_session` clears `rekey_state`, so a completed + // initiator cycle leaves at most one of the two set. If that ever + // stops holding, these four sites become instances of the epoch + // discard the responder arm was fixed for. if entry.is_established() && entry.has_rekey_in_progress() && entry.is_rekey_initiator() { let mut handshake = match entry.take_rekey_state() { Some(hs) => hs, @@ -811,11 +889,34 @@ impl Node { }; // Process XK msg2: read_xk_message_2 (extracts responder's epoch) - if let Err(e) = handshake.read_xk_message_2(&ack.handshake_payload) { + // + // Nothing here has been authenticated: the only thing tying this + // message to the initiation is the datagram's source address, an + // envelope field the sender chooses. Dropping the entry would let + // anyone able to reach us cancel any initiation in flight, so the + // entry goes back with the handshake rolled back to its pre-read + // state and the stored msg1 still scheduled for resend. The rollback + // is what makes the reinsert worth anything: `read_xk_message_2` + // mixes the sender's ephemeral in before it authenticates, so a + // handshake put back as it was left could never read the genuine + // msg2. A real peer's corrupt ack is covered by the same path — the + // responder resends its stored msg2, and the handshake sweep reaps + // the entry on its original deadline if none arrives. `touch()` is + // deliberately not called, so a spray cannot push that deadline out. + if let Err(e) = handshake.try_read_xk_message_2(&ack.handshake_payload) { debug!(error = %e, "Failed to process Noise XK msg2 in SessionAck"); - return; // Entry was already removed, don't put back a broken session + entry.set_state(EndToEndState::Initiating(handshake)); + self.sessions.insert(*src_addr, entry); + self.stats_mut() + .record_reject(RejectReason::Session(SessionReject::AckHandshakeFailed)); + return; } + // The three drops below stay drops. Each is downstream of a msg2 that + // already authenticated, so they are local failures rather than + // possible forgeries, and a handshake left at `Message2Done` cannot + // be re-driven from a resent msg1 anyway. + // Generate XK msg3: write_xk_message_3 (sends encrypted static + epoch) let msg3 = match handshake.write_xk_message_3() { Ok(m) => m, @@ -897,6 +998,24 @@ impl Node { }; // Rekey path: entry is Established with rekey_state (responder side) + // + // Every failure below abandons only the handshake. Nothing in a msg3 + // is authenticated until `read_xk_message_3` has both succeeded and + // produced a static key matching this session's peer, so a failure + // here proves nothing about the sender and must not cost the entry + // anything it would miss. A `pending_new_session` beside the + // handshake is the epoch the real peer may already have cut over to, + // and dropping it kills the reverse direction on two unauthenticated + // messages: a forged msg1 to arm the handshake, then any garbage + // msg3. `abandon_handshake` keeps it; `abandon_rekey` does not. + // + // What `abandon_handshake` leaves behind, and why each is safe here: + // `rekey_completed_ms` must survive, since `pending_stale` reads it + // and a zeroed stamp reads as freshly completed. A stranded + // `rekey_msg3_payload` belongs to an initiator cycle and clears + // itself once `resend_pending_session_msg3` exhausts its budget. + // `peer_new_epoch_confirmed` only stops that retransmission, and + // `rekey_initiator` is false throughout this arm by its own gate. if entry.is_established() && entry.has_rekey_in_progress() && !entry.is_rekey_initiator() { let mut handshake = match entry.take_rekey_state() { Some(hs) => hs, @@ -913,7 +1032,7 @@ impl Node { error = %e, "Failed to process rekey XK msg3" ); - entry.abandon_rekey(); + entry.abandon_handshake(); self.sessions.insert(*src_addr, entry); return; } @@ -925,8 +1044,14 @@ impl Node { let rekey_pubkey = match handshake.remote_static() { Some(pk) => *pk, None => { + // Not independently exercised by any test: a successful + // `read_xk_message_3` always sets the remote static, so + // reaching this needs fault injection into + // `HandshakeState`. Changed with its three siblings so + // the arm has one rule rather than three plus an + // exception. debug!("No remote static key after processing rekey XK msg3"); - entry.abandon_rekey(); + entry.abandon_handshake(); self.sessions.insert(*src_addr, entry); return; } @@ -936,7 +1061,7 @@ impl Node { src = %self.peer_display_name(src_addr), "FSP rekey: initiator static key differs from the established peer key" ); - entry.abandon_rekey(); + entry.abandon_handshake(); self.sessions.insert(*src_addr, entry); self.stats_mut() .record_reject(RejectReason::Session(SessionReject::RekeyKeyMismatch)); @@ -947,8 +1072,11 @@ impl Node { let session = match handshake.into_session() { Ok(s) => s, Err(e) => { + // Also not independently exercised, for the same reason + // as the missing-static arm above: a handshake that read + // msg3 successfully always converts. debug!(error = %e, "Failed to create session from rekey XK msg3"); - entry.abandon_rekey(); + entry.abandon_handshake(); self.sessions.insert(*src_addr, entry); return; } @@ -1441,19 +1569,29 @@ impl Node { // Read existing, decide, and apply the write under one guard so // the keep-tighter update stays atomic. let prior = map.get(&fips_addr).copied(); - let actions = self.fsp.plan_path_mtu_tighten(fips_addr, prior, new_mtu); + let actions = + self.fsp + .plan_path_mtu_tighten(fips_addr, prior.map(|e| e.mtu), new_mtu); if actions.is_empty() { debug!( dest = %peer_name, fips_addr = %fips_addr, new_mtu, - existing = prior.unwrap_or(new_mtu), + existing = prior.map(|e| e.mtu).unwrap_or(new_mtu), "PathMtuNotification: keeping tighter existing path_mtu_lookup value" ); } for action in actions { if let FspAction::TightenPathMtuLookup { fips_addr, mtu } = action { - map.insert(fips_addr, 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(mtu)); debug!( dest = %peer_name, fips_addr = %fips_addr, @@ -1814,19 +1952,27 @@ impl Node { // Read existing, decide, and apply the write under one guard so // the keep-tighter update stays atomic. let prior = map.get(&fips_addr).copied(); - let actions = self.fsp.plan_path_mtu_tighten(fips_addr, prior, msg.mtu); + let actions = + self.fsp + .plan_path_mtu_tighten(fips_addr, prior.map(|e| e.mtu), msg.mtu); if actions.is_empty() { debug!( dest = %peer_name, fips_addr = %fips_addr, bottleneck_mtu = msg.mtu, - existing = prior.unwrap_or(msg.mtu), + existing = prior.map(|e| e.mtu).unwrap_or(msg.mtu), "Reactive MtuExceeded: keeping tighter existing path_mtu_lookup value" ); } for action in actions { if let FspAction::TightenPathMtuLookup { fips_addr, mtu } = action { - map.insert(fips_addr, 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(mtu)); debug!( dest = %peer_name, fips_addr = %fips_addr, diff --git a/src/node/handlers/timeout.rs b/src/node/handlers/timeout.rs index 3fd5319d..c8d51420 100644 --- a/src/node/handlers/timeout.rs +++ b/src/node/handlers/timeout.rs @@ -7,7 +7,7 @@ use crate::proto::fmp::{ ConnAction, ConnSnapshot, LifecycleView, PeerSnapshot, RekeyResendSnapshot, }; use crate::transport::LinkId; -use tracing::{debug, info}; +use tracing::{debug, info, warn}; impl LifecycleView for Node { fn stale_connections(&self, now_ms: u64, timeout_ms: u64) -> Vec { @@ -515,4 +515,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 e2f323ad..963f5b83 100644 --- a/src/node/metrics.rs +++ b/src/node/metrics.rs @@ -118,10 +118,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 00fe03a4..fa4b93f9 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -28,7 +28,7 @@ pub(crate) mod stats_history; mod tests; mod tree; -use self::rate_limit::HandshakeRateLimiter; +use self::rate_limit::{HandshakeRateLimiter, SessionSetupRateLimiter}; use self::reloadable::Reloadable; /// Half-range of the symmetric jitter applied to the per-session rekey timer. @@ -338,7 +338,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). @@ -474,6 +474,8 @@ pub struct Node { // === Rate Limiting === /// Rate limiter for msg1 processing (DoS protection). msg1_rate_limiter: HandshakeRateLimiter, + /// Rate limiter for inbound FSP SessionSetup, keyed on the link peer. + setup_rate_limiter: SessionSetupRateLimiter, /// Rate limiter for ICMP Packet Too Big messages. icmp_rate_limiter: IcmpRateLimiter, /// Routing-subsystem state (routing error-signal rate limiter). @@ -585,6 +587,28 @@ fn build_msg1_rate_limiter(config: &Config) -> HandshakeRateLimiter { ) } +/// Build the per-link session-setup limiter's two bucket sizes. +/// +/// The stranger bucket is configured directly. The established bucket is +/// derived exactly as the FMP limiter's is, from `max_peers`, the rekey +/// period and the resend budget: a hub neighbour can legitimately carry the +/// rekey traffic of every session this node holds, so that is the population +/// the per-link bucket has to cover. +fn build_setup_rate_limiter(config: &Config) -> SessionSetupRateLimiter { + let rl = &config.node.rate_limit; + let established = rate_limit::derive_established_bucket( + config.node.limits.max_peers, + config.node.rekey.after_secs, + rl.handshake_max_resends, + rl.session_setup_burst, + rl.session_setup_rate, + ); + SessionSetupRateLimiter::with_params( + (rl.session_setup_burst, rl.session_setup_rate), + established, + ) +} + impl Node { /// Create a new node from configuration. pub fn new(config: Config) -> Result { @@ -630,6 +654,7 @@ impl Node { config.node.cache.coord_ttl_secs * 1000, ); let msg1_rate_limiter = build_msg1_rate_limiter(&config); + let setup_rate_limiter = build_setup_rate_limiter(&config); let max_connections = config.node.limits.max_connections; let max_peers = config.node.limits.max_peers; @@ -705,6 +730,7 @@ impl Node { peers_by_index: HashMap::new(), pending_outbound: HashMap::new(), msg1_rate_limiter, + setup_rate_limiter, icmp_rate_limiter: IcmpRateLimiter::new(), routing: Router::new(), fmp: Fmp::new(), @@ -777,6 +803,7 @@ impl Node { config.node.cache.coord_ttl_secs * 1000, ); let msg1_rate_limiter = build_msg1_rate_limiter(&config); + let setup_rate_limiter = build_setup_rate_limiter(&config); let max_connections = config.node.limits.max_connections; let max_peers = config.node.limits.max_peers; @@ -849,6 +876,7 @@ impl Node { peers_by_index: HashMap::new(), pending_outbound: HashMap::new(), msg1_rate_limiter, + setup_rate_limiter, icmp_rate_limiter: IcmpRateLimiter::new(), routing: Router::new(), fmp: Fmp::new(), @@ -2592,9 +2620,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() @@ -2602,10 +2642,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), + ); } } @@ -2622,7 +2684,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/rate_limit.rs b/src/node/rate_limit.rs index 32a581cc..0e577b7e 100644 --- a/src/node/rate_limit.rs +++ b/src/node/rate_limit.rs @@ -36,10 +36,23 @@ //! The concurrency limb (`max_pending`) is deliberately *not* split, so no //! equivalent inflation happens there: one counter bounds simultaneous //! in-flight handshake state whoever holds the slot. +//! +//! ## Why the session-setup limiter *is* keyed, when the msg1 limiter is not +//! +//! "Not per-source, since UDP sources are spoofable" is about the FMP link +//! layer, where the source is a transport address on an unauthenticated +//! datagram. [`SessionSetupRateLimiter`] sits a layer up and keys on +//! something different: the FMP link peer the datagram arrived over, which +//! the hop-by-hop Noise AEAD authenticates and whose population is bounded by +//! the peer table. Keying on the FSP `src_addr` instead would be the mistake +//! that sentence warns about, since that field is chosen by the sender and a +//! single sender can mint an unbounded number of distinct values. +use crate::NodeAddr; +use std::collections::HashMap; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; -use std::time::Instant; +use std::time::{Duration, Instant}; /// Default burst capacity (max tokens). pub const DEFAULT_BURST_CAPACITY: u32 = 100; @@ -398,6 +411,104 @@ impl HandshakeRateLimiter { } } +/// How long a link peer's buckets are kept after its last setup message. +const SETUP_BUCKET_IDLE: Duration = Duration::from_secs(300); + +/// One link peer's pair of session-setup buckets. +struct LinkBuckets { + /// Setup messages that would create a new half-open session entry. + stranger: TokenBucket, + /// Setup messages naming a peer this node already has a session with: + /// rekey and restart traffic, which creates no new entry. + established: TokenBucket, + /// Last time this peer was charged, for idle pruning. + seen: Instant, +} + +/// Rate limiter for inbound FSP SessionSetup messages, keyed on the link peer. +/// +/// The setup path allocates a `SessionEntry` and sends a SessionAck for every +/// well-formed msg1 naming an address it has no entry for, and the address is +/// an envelope field the sender picks. Without a limiter one neighbour can +/// grow the session table at whatever rate it can transmit, and buy a routed +/// ack per entry to a destination it chooses. +/// +/// **The key is the FMP link peer the datagram arrived over, never the FSP +/// `src_addr`.** The link peer is authenticated by the hop-by-hop Noise AEAD +/// and its population is bounded by the peer table; `src_addr` is chosen by +/// the sender, so keying on it would let one sender mint a fresh full bucket +/// per forged message. See the module doc for how this squares with the msg1 +/// limiter being unkeyed. +/// +/// Two buckets per link, for the same reason [`HandshakeRateLimiter`] has +/// two: a drained stranger bucket must not also stop an established peer's +/// rekey msg1 from arming. Suppressed rekey is quiet — nothing errors and no +/// session drops — so folding both classes into one bucket would let a +/// sprayer one hop away hold forward-secrecy rotation off for everything +/// behind that link with no signal but a flat `rekey_armed`. +/// +/// Idle links are pruned lazily on the admit path. Pruning only ever relaxes +/// the limit, and it cannot be farmed: earning a fresh bucket costs a full +/// [`SETUP_BUCKET_IDLE`] of silence on that link, which at any sane sizing is +/// a far lower sustained rate than simply waiting for the bucket to refill. +pub struct SessionSetupRateLimiter { + /// Per-link-peer buckets, created on first use. + buckets: HashMap, + /// Burst and refill rate for a new link's stranger bucket. + stranger: (u32, f64), + /// Burst and refill rate for a new link's established bucket. + established: (u32, f64), +} + +impl SessionSetupRateLimiter { + /// Create a limiter whose per-link buckets take the given parameters. + /// + /// Each pair is `(burst, tokens per second)`. + pub fn with_params(stranger: (u32, f64), established: (u32, f64)) -> Self { + Self { + buckets: HashMap::new(), + stranger, + established, + } + } + + /// Charge one setup message of `class` to `link_peer`. + /// + /// Returns `false` when the class's bucket for that link is empty, in + /// which case the caller must drop the message before doing any work. + pub fn try_admit(&mut self, link_peer: &NodeAddr, class: Msg1Class) -> bool { + let now = Instant::now(); + let stranger = self.stranger; + let established = self.established; + let link = self + .buckets + .entry(*link_peer) + .or_insert_with(|| LinkBuckets { + stranger: TokenBucket::with_params(stranger.0, stranger.1), + established: TokenBucket::with_params(established.0, established.1), + seen: now, + }); + link.seen = now; + + let admitted = match class { + Msg1Class::Stranger => link.stranger.try_acquire(), + Msg1Class::EstablishedLink => link.established.try_acquire(), + }; + + if admitted { + self.buckets + .retain(|_, link| now.duration_since(link.seen) < SETUP_BUCKET_IDLE); + } + admitted + } + + /// Number of link peers currently holding buckets. + #[cfg(test)] + pub fn len(&self) -> usize { + self.buckets.len() + } +} + #[cfg(test)] mod tests { use super::*; @@ -724,4 +835,42 @@ mod tests { assert_eq!(burst, 100); assert_eq!(rate, 10.0); } + + fn addr(byte: u8) -> NodeAddr { + NodeAddr::from_bytes([byte; 16]) + } + + #[test] + fn setup_limiter_draining_one_link_leaves_another_links_budget_untouched() { + let mut limiter = SessionSetupRateLimiter::with_params((2, 0.001), (2, 0.001)); + let noisy = addr(0x01); + let quiet = addr(0x02); + + assert!(limiter.try_admit(&noisy, Msg1Class::Stranger)); + assert!(limiter.try_admit(&noisy, Msg1Class::Stranger)); + assert!( + !limiter.try_admit(&noisy, Msg1Class::Stranger), + "the noisy link's own bucket must run out" + ); + assert!( + limiter.try_admit(&quiet, Msg1Class::Stranger), + "a second link peer must not share the first one's budget" + ); + assert_eq!(limiter.len(), 2); + } + + #[test] + fn setup_limiter_draining_the_stranger_bucket_still_admits_established_peer_setups() { + let mut limiter = SessionSetupRateLimiter::with_params((1, 0.001), (1, 0.001)); + let link = addr(0x01); + + assert!(limiter.try_admit(&link, Msg1Class::Stranger)); + assert!(!limiter.try_admit(&link, Msg1Class::Stranger)); + assert!( + limiter.try_admit(&link, Msg1Class::EstablishedLink), + "rekey traffic must not be starved by a stranger flood on the \ + same link; suppressed rotation is silent and would show only as \ + a flat rekey_armed counter" + ); + } } diff --git a/src/node/reject.rs b/src/node/reject.rs index 616e411b..382323a4 100644 --- a/src/node/reject.rs +++ b/src/node/reject.rs @@ -226,6 +226,18 @@ pub enum SessionReject { /// rather than arming a second handshake. Tracked via /// [`SessionStats::rekey_pending`](crate::node::stats::SessionStats). RekeyPending, + /// An inbound SessionAck naming a session we are initiating failed the + /// XK msg2 read. The message carries no authenticator tying it to the + /// initiation, so the entry is kept and the handshake rolled back + /// rather than discarded; a sustained rate here is either a broken path + /// to the responder or somebody spraying forged acks to hold + /// establishment down. Tracked via + /// [`SessionStats::ack_handshake_failed`](crate::node::stats::SessionStats). + AckHandshakeFailed, + /// A setup message was refused by the per-link-peer setup limiter + /// before any handshake state was created or any ack sent. Tracked via + /// [`SessionStats::setup_rate_limited`](crate::node::stats::SessionStats). + SetupRateLimited, } /// MMP rejection reasons. diff --git a/src/node/reloadable.rs b/src/node/reloadable.rs index bf2d4a1c..3fc58b2c 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 lookup `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_rendezvous` is an async spawned subsystem, not a snapshot of disk //! state. //! diff --git a/src/node/stats.rs b/src/node/stats.rs index 49d06750..96503f43 100644 --- a/src/node/stats.rs +++ b/src/node/stats.rs @@ -70,6 +70,15 @@ pub struct SessionStats { /// already have adopted, so a sustained rate means one side keeps /// rekeying while the other never appears on the new epoch. pub pending_replaced: u64, + /// An inbound SessionAck failed the XK msg2 read against a session we + /// are initiating. The entry is kept and the handshake rolled back, + /// since the message authenticates nothing; a sustained rate is either + /// a broken path to the responder or forged acks holding establishment + /// down. + pub ack_handshake_failed: u64, + /// A setup message was refused by the per-link-peer setup limiter, + /// before any handshake state was created or any ack sent. + pub setup_rate_limited: u64, } impl SessionStats { @@ -85,6 +94,8 @@ impl SessionStats { rekey_pending: self.rekey_pending, rekey_expired: self.rekey_expired, pending_replaced: self.pending_replaced, + ack_handshake_failed: self.ack_handshake_failed, + setup_rate_limited: self.setup_rate_limited, } } @@ -97,6 +108,8 @@ impl SessionStats { SessionReject::RekeyTiebreak => self.rekey_tiebreak += 1, SessionReject::RekeyYielded => self.rekey_yielded += 1, SessionReject::RekeyPending => self.rekey_pending += 1, + SessionReject::AckHandshakeFailed => self.ack_handshake_failed += 1, + SessionReject::SetupRateLimited => self.setup_rate_limited += 1, } } } @@ -364,6 +377,8 @@ pub struct SessionStatsSnapshot { pub rekey_pending: u64, pub rekey_expired: u64, pub pending_replaced: u64, + pub ack_handshake_failed: u64, + pub setup_rate_limited: u64, } #[derive(Clone, Debug, Default, Serialize)] diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index 523eeb0a..f60dd43d 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -1047,6 +1047,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 d66f039b..ea84f04e 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] @@ -1191,8 +1198,28 @@ 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 3ad70b81..a281f7ef 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -10,6 +10,15 @@ use crate::node::tests::spanning_tree::{ use crate::proto::fsp::{SessionAck, SessionMsg3}; use crate::proto::link::SessionDatagram; +/// A stand-in for the authenticated FMP link peer a datagram arrived over. +/// +/// Tests that call `handle_session_payload` directly have no link underneath +/// them. The setup limiter keys on this address, so a test wanting to drain a +/// bucket has to drive `handle_session_datagram` instead. +fn stub_link_peer() -> NodeAddr { + make_node_addr(0xFE) +} + /// Populate all nodes' coordinate caches with each other's coords. /// /// This enables routing between non-adjacent nodes (bloom filter + tree @@ -2681,6 +2690,228 @@ 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::proto::routing::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, 1_000); + 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" + ); +} + +/// 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.supervisor.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 @@ -2945,9 +3176,11 @@ async fn test_coords_required_naming_a_dest_with_no_session_is_counted_as_an_unk node.handle_coords_required(&reporter, &encoded[5..]).await; assert_eq!(node.stats().session.unknown_session, 1); - // The gate runs ahead of the response rate limiter, so a second - // identical signal is rejected the same way rather than being - // absorbed by rate-limiter state keyed on an attacker-chosen address. + // A second identical signal is refused the same way. This does not pin + // the gate's position relative to the response rate limiter: should_send + // returning false would not short-circuit the handler, so this counter + // reaches 2 either way. The ordering is pinned by + // test_coords_required_for_an_unbound_dest_never_reaches_the_response_rate_limiter. node.handle_coords_required(&reporter, &encoded[5..]).await; assert_eq!(node.stats().session.unknown_session, 2); @@ -2979,6 +3212,55 @@ async fn test_coords_required_naming_a_dest_with_no_session_is_counted_as_an_unk ); } +#[tokio::test] +async fn test_coords_required_for_an_unbound_dest_never_reaches_the_response_rate_limiter() { + use crate::proto::routing::CoordsRequired; + + let mut node = make_node(); + + let dest = NodeAddr::from_bytes([0xCC; 16]); + let reporter = NodeAddr::from_bytes([0xBB; 16]); + + assert_eq!( + node.coords_response_rate_limiter.len(), + 0, + "precondition: the response rate limiter holds nothing before the signal" + ); + + let encoded = CoordsRequired::new(dest, reporter).encode(); + node.handle_coords_required(&reporter, &encoded[5..]).await; + + assert_eq!(node.stats().session.unknown_session, 1); + assert_eq!( + node.coords_response_rate_limiter.len(), + 0, + "an inadmissible signal must be refused before should_send can insert \ + the attacker-chosen address into last_sent" + ); +} + +#[tokio::test] +async fn test_coords_required_for_a_bound_dest_does_reach_the_response_rate_limiter() { + use crate::proto::routing::CoordsRequired; + + let mut node = make_node(); + + let remote = Identity::generate(); + install_initiating(&mut node, &remote); + let dest = *remote.node_addr(); + let reporter = NodeAddr::from_bytes([0xBB; 16]); + + let encoded = CoordsRequired::new(dest, reporter).encode(); + node.handle_coords_required(&reporter, &encoded[5..]).await; + + assert_eq!(node.stats().session.unknown_session, 0); + assert_eq!( + node.coords_response_rate_limiter.len(), + 1, + "an admitted signal must still consult the response rate limiter" + ); +} + #[tokio::test] async fn test_mtu_exceeded_whose_claimed_source_is_the_destination_it_names_is_dropped() { let mut node = make_node(); @@ -3368,8 +3650,14 @@ async fn test_session_msg3_rejects_spoofed_source_address() { ); node.sessions.insert(victim_addr, entry); - node.handle_session_payload(&victim_addr, &SessionMsg3::new(msg3).encode(), 1280, false) - .await; + node.handle_session_payload( + &victim_addr, + &stub_link_peer(), + &SessionMsg3::new(msg3).encode(), + 1280, + false, + ) + .await; assert_eq!( node.session_count(), @@ -3401,8 +3689,14 @@ async fn test_session_msg3_accepts_matching_source_address() { ); node.sessions.insert(peer_addr, entry); - node.handle_session_payload(&peer_addr, &SessionMsg3::new(msg3).encode(), 1280, false) - .await; + node.handle_session_payload( + &peer_addr, + &stub_link_peer(), + &SessionMsg3::new(msg3).encode(), + 1280, + false, + ) + .await; assert!( node.sessions @@ -3436,8 +3730,14 @@ async fn test_rekey_msg3_rejects_different_static_key() { entry.set_rekey_state(responder, false); node.sessions.insert(peer_addr, entry); - node.handle_session_payload(&peer_addr, &SessionMsg3::new(msg3).encode(), 1280, false) - .await; + node.handle_session_payload( + &peer_addr, + &stub_link_peer(), + &SessionMsg3::new(msg3).encode(), + 1280, + false, + ) + .await; let entry = node .sessions @@ -3472,8 +3772,14 @@ async fn test_rekey_msg3_accepts_established_peer_key() { entry.set_rekey_state(responder, false); node.sessions.insert(peer_addr, entry); - node.handle_session_payload(&peer_addr, &SessionMsg3::new(msg3).encode(), 1280, false) - .await; + node.handle_session_payload( + &peer_addr, + &stub_link_peer(), + &SessionMsg3::new(msg3).encode(), + 1280, + false, + ) + .await; let entry = node.sessions.get(&peer_addr).expect("session present"); assert!(entry.pending_new_session().is_some()); @@ -3511,8 +3817,14 @@ async fn test_rekey_msg3_accepts_odd_parity_peer_stored_as_even() { entry.set_rekey_state(responder, false); node.sessions.insert(peer_addr, entry); - node.handle_session_payload(&peer_addr, &SessionMsg3::new(msg3).encode(), 1280, false) - .await; + node.handle_session_payload( + &peer_addr, + &stub_link_peer(), + &SessionMsg3::new(msg3).encode(), + 1280, + false, + ) + .await; let entry = node.sessions.get(&peer_addr).expect("session present"); assert!( @@ -3522,6 +3834,113 @@ async fn test_rekey_msg3_accepts_odd_parity_peer_stored_as_even() { assert_eq!(node.stats().session.rekey_key_mismatch, 0); } +/// Install the shape the msg3 epoch-discard defect needs: an established +/// entry holding a completed rekey the peer has not yet cut over to, stamped +/// stale, with a second handshake armed beside it by a stranger's setup. +/// +/// Returns the node, the peer's address and the cryptographically valid msg3 +/// the stranger would send to finish the handshake it armed. +fn install_stale_pending_beside_a_stranger_armed_handshake( + peer: &Identity, + stranger: &Identity, +) -> (Node, crate::NodeAddr, Vec) { + let (mut node, peer_addr) = make_node_with_established_peer(false, peer); + let msg3 = arm_stranger_handshake_beside_stale_pending(&mut node, &peer_addr, peer, stranger); + (node, peer_addr, msg3) +} + +/// Put a stale completed rekey and a stranger-armed handshake on an entry +/// that is already established, and return the msg3 that finishes the +/// stranger's handshake. +/// +/// This is what a forged setup leaves behind once `pending_stale` has +/// lapsed: the veto no longer fires, so the fall-through arms a responder +/// handshake beside pending keys it does not touch. `set_pending_session` +/// clears `rekey_state`, so the arming has to follow it, as it does in the +/// handler. +fn arm_stranger_handshake_beside_stale_pending( + node: &mut Node, + peer_addr: &crate::NodeAddr, + peer: &Identity, + stranger: &Identity, +) -> Vec { + let pending = make_noise_session(node.identity(), peer); + let (responder, msg3) = drive_xk_to_msg3(stranger, node.identity()); + + let idle_ms = node.config().node.session.idle_timeout_secs * 1000; + let now_ms = wall_clock_ms(); + let entry = node.sessions.get_mut(peer_addr).unwrap(); + entry.set_pending_session(pending); + // Stale enough that `pending_stale` is true, which is what lets a forged + // setup arm the handshake this state starts from. + entry.set_rekey_completed_ms(now_ms - idle_ms - 60_000); + entry.set_rekey_state(responder, false); + entry.record_peer_rekey(now_ms); + + msg3 +} + +#[tokio::test] +async fn test_forged_msg3_against_a_peer_armed_handshake_leaves_the_completed_epoch_intact() { + let peer = Identity::generate(); + let stranger = Identity::generate(); + let (mut node, peer_addr, _valid_msg3) = + install_stale_pending_beside_a_stranger_armed_handshake(&peer, &stranger); + + // Garbage of the right length: `read_xk_message_3` fails on the AEAD. + let forged = SessionMsg3::new(vec![0u8; crate::noise::XK_HANDSHAKE_MSG3_SIZE]).encode(); + node.handle_session_payload(&peer_addr, &stub_link_peer(), &forged, 1280, false) + .await; + + let entry = node.sessions.get(&peer_addr).expect("session present"); + assert!( + entry.pending_new_session().is_some(), + "an unauthenticated msg3 must not discard the key epoch the peer may \ + already have cut over to; only the handshake it failed belongs to it" + ); + assert!( + entry.is_established(), + "the running session must be left intact alongside the pending one" + ); + assert!( + !entry.has_rekey_in_progress(), + "the handshake the msg3 failed against must still be abandoned" + ); +} + +#[tokio::test] +async fn test_rekey_msg3_from_a_different_static_key_leaves_the_completed_epoch_intact() { + let peer = Identity::generate(); + let stranger = Identity::generate(); + let (mut node, peer_addr, valid_msg3) = + install_stale_pending_beside_a_stranger_armed_handshake(&peer, &stranger); + + // Cryptographically valid for the handshake the stranger armed, so + // `read_xk_message_3` succeeds and the key-mismatch branch decides. + node.handle_session_payload( + &peer_addr, + &stub_link_peer(), + &SessionMsg3::new(valid_msg3).encode(), + 1280, + false, + ) + .await; + + let entry = node.sessions.get(&peer_addr).expect("session present"); + assert!( + entry.pending_new_session().is_some(), + "a msg3 whose static key is not this session's peer must not discard \ + the completed epoch either" + ); + assert!(entry.is_established()); + assert_eq!( + node.stats().session.rekey_key_mismatch, + 1, + "the key mismatch must still be counted, so this test also pins that \ + the refusal itself did not move" + ); +} + // ============================================================================ // Integration tests: a setup message naming an established peer // ============================================================================ @@ -3619,7 +4038,7 @@ async fn test_forged_setup_naming_established_peer_leaves_session_carrying_traff let forged = forge_setup_from_stranger(&nodes); nodes[1] .node - .handle_session_payload(&node0_addr, &forged, 1280, false) + .handle_session_payload(&node0_addr, &node0_addr, &forged, 1280, false) .await; let entry = nodes[1] @@ -3776,6 +4195,251 @@ async fn test_genuine_peer_restart_reestablishes_session_with_rekey_enabled() { cleanup_nodes(&mut nodes).await; } +// ============================================================================ +// Integration tests: the per-link-peer session-setup limiter +// ============================================================================ + +/// Build a two-node routable mesh with the setup limiter sized for a test. +async fn make_setup_limited_pair(burst: u32, rate: f64) -> Vec { + let configs = (0..2) + .map(|_| { + let mut config = Config::new(); + config.node.rekey.enabled = false; + config.node.rate_limit.session_setup_burst = burst; + config.node.rate_limit.session_setup_rate = rate; + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + nodes +} + +/// Deliver one forged SessionSetup to `nodes[1]` over the link from +/// `nodes[0]`, naming a fresh source address nobody has seen. +/// +/// Driven through `handle_session_datagram` rather than +/// `handle_session_payload` for two reasons: it is the only path that binds +/// the link peer the limiter keys on, and its coordinate-cache warming is +/// what gives the forged address a route, without which the ack send fails +/// and the entry is never inserted even in unlimited code. +async fn deliver_forged_setup_over_link(nodes: &mut [TestNode]) { + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let forged_src = *Identity::generate().node_addr(); + + let setup = forge_setup_from_stranger(nodes); + let datagram = SessionDatagram::new(forged_src, node1_addr, setup).with_ttl(64); + let encoded = datagram.encode(); + + nodes[1] + .node + .handle_session_datagram(&node0_addr, &encoded[1..], false) + .await; +} + +#[tokio::test] +async fn test_forged_setups_from_one_link_peer_stop_creating_session_entries_once_the_bucket_is_drained() + { + const BURST: u32 = 4; + // Slow enough that nothing refills during the test. + let mut nodes = make_setup_limited_pair(BURST, 0.5).await; + + let before = nodes[1].node.sessions.len(); + for _ in 0..BURST { + deliver_forged_setup_over_link(&mut nodes).await; + } + assert_eq!( + nodes[1].node.sessions.len(), + before + BURST as usize, + "the burst must be admitted, or this test would pass for the wrong reason" + ); + assert_eq!(nodes[1].node.stats().session.setup_rate_limited, 0); + + // Every SessionAck the handler emits goes out through + // `send_session_datagram`, which is the only thing that bumps this + // counter on a node with no transit traffic. A refused setup must not + // move it: that is the ack amplification bound, measured rather than + // argued from where the check sits. + let originated = nodes[1].node.metrics().forwarding.originated_packets.get(); + + for _ in 0..3 { + deliver_forged_setup_over_link(&mut nodes).await; + } + + assert_eq!( + nodes[1].node.sessions.len(), + before + BURST as usize, + "a drained bucket must stop the session table growing" + ); + assert_eq!( + nodes[1].node.stats().session.setup_rate_limited, + 3, + "each refusal must be counted; the DEBUG line is invisible by default" + ); + assert_eq!( + nodes[1].node.metrics().forwarding.originated_packets.get(), + originated, + "a refused setup must emit nothing at all, so it buys the sender no \ + packet to an address it chose" + ); + + cleanup_nodes(&mut nodes).await; +} + +#[tokio::test] +async fn test_a_drained_setup_bucket_refills_and_admits_the_next_legitimate_setup() { + // Fast refill: the point is that the denial is transient, and that the + // initiator's own msg1 resend schedule covers a window this short. + let mut nodes = make_setup_limited_pair(2, 50.0).await; + + for _ in 0..3 { + deliver_forged_setup_over_link(&mut nodes).await; + } + assert!( + nodes[1].node.stats().session.setup_rate_limited > 0, + "the bucket must actually be drained before the refill is tested" + ); + + tokio::time::sleep(Duration::from_millis(100)).await; + establish_pair_session(&mut nodes).await; + + cleanup_nodes(&mut nodes).await; +} + +#[tokio::test] +async fn test_a_drained_stranger_bucket_still_admits_a_setup_naming_an_established_peer() { + // Burst 2: one token for the genuine msg1 that establishes the pair, one + // for a forged stranger setup, and the third stranger setup is refused. + let mut nodes = make_setup_limited_pair(2, 0.5).await; + establish_pair_session(&mut nodes).await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + deliver_forged_setup_over_link(&mut nodes).await; + deliver_forged_setup_over_link(&mut nodes).await; + assert!( + nodes[1].node.stats().session.setup_rate_limited > 0, + "the stranger bucket must be drained before the established class is tested" + ); + + // The same message, but naming the established peer: this is the shape an + // inbound rekey arrives in. It creates no new table entry, so it draws on + // its own bucket rather than competing with stranger admission. + let setup = forge_setup_from_stranger(&nodes); + let datagram = SessionDatagram::new(node0_addr, node1_addr, setup).with_ttl(64); + let encoded = datagram.encode(); + nodes[1] + .node + .handle_session_datagram(&node0_addr, &encoded[1..], false) + .await; + + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("the established session must still be there") + .has_rekey_in_progress(), + "a drained stranger bucket must not stop an established peer's rekey \ + arming: suppressed rotation is silent, and the operator's only \ + signal would be a flat rekey_armed" + ); + assert_eq!(nodes[1].node.stats().session.rekey_armed, 1); + + cleanup_nodes(&mut nodes).await; +} + +// ============================================================================ +// Integration tests: a forged SessionAck against an in-flight initiation +// ============================================================================ + +#[tokio::test] +async fn test_forged_session_ack_leaves_the_initiation_able_to_complete_on_the_genuine_ack() { + let mut nodes = make_rekey_disabled_pair().await; + + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node1_pubkey = nodes[1].node.identity().pubkey_full(); + + // Initiate but do not pump: node 0 sits in Initiating with its msg1 in + // flight, which is the state the forgery targets. + nodes[0] + .node + .initiate_session(node1_addr, node1_pubkey) + .await + .expect("initiate_session failed"); + let activity_before = nodes[0] + .node + .get_session(&node1_addr) + .expect("initiating entry present") + .last_activity(); + + // A forged ack of exactly the right length. The leading 33 bytes are a + // valid compressed point, which is the point of the test: random bytes + // usually fail `PublicKey::from_slice` before anything has been mixed + // into the symmetric state, so they would not discriminate the rollback. + let mut payload = Identity::generate().pubkey_full().serialize().to_vec(); + payload.extend_from_slice(&[0u8; crate::noise::EPOCH_ENCRYPTED_SIZE]); + assert_eq!(payload.len(), crate::noise::XK_HANDSHAKE_MSG2_SIZE); + let coords = nodes[1].node.tree_state().my_coords().clone(); + let forged = SessionAck::new(coords.clone(), coords) + .with_handshake(payload) + .encode(); + + nodes[0] + .node + .handle_session_payload(&node1_addr, &node1_addr, &forged, 1280, false) + .await; + + let entry = nodes[0] + .node + .get_session(&node1_addr) + .expect("an unauthenticated ack must not destroy the initiation"); + assert!( + entry.is_initiating(), + "the entry must still be the initiation it was, not a broken one" + ); + assert_eq!( + entry.last_activity(), + activity_before, + "the reinsert must not push the handshake sweep's deadline out, or a \ + spray would keep a dead entry alive" + ); + assert_eq!( + nodes[0].node.stats().session.ack_handshake_failed, + 1, + "the refusal must be counted; its DEBUG line is invisible at the \ + default log level" + ); + + // The genuine exchange now runs to completion over the same handshake. + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes).await; + } + + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .expect("initiator session present") + .is_established(), + "the initiation must still complete when the genuine ack arrives" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder session present") + .is_established(), + "and the responder must reach Established too" + ); + + cleanup_nodes(&mut nodes).await; +} + // ============================================================================ // Tick-loop maintenance with periodic rekey disabled // ============================================================================ @@ -3789,6 +4453,16 @@ fn make_node_with_established_peer( let mut config = Config::new(); config.node.rekey.enabled = rekey_enabled; let mut node = make_node_with(config); + let peer_addr = install_established_peer(&mut node, peer); + (node, peer_addr) +} + +/// Install one established session with `peer` on an existing node. +/// +/// Split out of `make_node_with_established_peer` for the tests that must +/// choose the peer identity relative to the node's own address, which needs +/// the node to exist first. +fn install_established_peer(node: &mut Node, peer: &Identity) -> crate::NodeAddr { let peer_addr = *peer.node_addr(); let session = make_noise_session(node.identity(), peer); @@ -3801,7 +4475,219 @@ fn make_node_with_established_peer( ); entry.mark_established(1000); node.sessions.insert(peer_addr, entry); - (node, peer_addr) + peer_addr +} + +/// Generate an identity whose address sorts strictly above `node_addr`. +/// +/// The dual-initiation tie-break compares the two addresses directly, so a +/// test that wants a specific side of it has to pick the peer to match. +/// Roughly two draws on average, as with `generate_odd_parity_identity`. +fn peer_identity_sorting_above(node_addr: &crate::NodeAddr) -> Identity { + loop { + let id = Identity::generate(); + if id.node_addr() > node_addr { + return id; + } + } +} + +/// Build the initiator-side XK handshake `initiate_session_rekey` would +/// leave on the entry, without needing a route to send its msg1 over. +fn our_rekey_initiator_handshake(node: &Node, peer: &Identity) -> crate::noise::HandshakeState { + let mut handshake = crate::noise::HandshakeState::new_xk_initiator( + node.identity().keypair(), + peer.pubkey_full(), + ); + handshake.set_local_epoch([0x11; 8]); + handshake + .write_xk_message_1() + .expect("our own msg1 must build"); + handshake +} + +/// Generate an identity whose address sorts strictly below `node_addr`. +fn peer_identity_sorting_below(node_addr: &crate::NodeAddr) -> Identity { + loop { + let id = Identity::generate(); + if id.node_addr() < node_addr { + return id; + } + } +} + +#[tokio::test] +async fn test_setup_naming_a_peer_whose_address_sorts_above_ours_keeps_our_rekey_and_counts_the_tiebreak() + { + let mut config = Config::new(); + config.node.rekey.enabled = false; + let mut node = make_node_with(config); + + // Our address sorts smaller, so the tie-break keeps us as initiator. + let peer = peer_identity_sorting_above(node.node_addr()); + let peer_addr = install_established_peer(&mut node, &peer); + + // Our own rekey is in flight as initiator. + let our_handshake = our_rekey_initiator_handshake(&node, &peer); + node.sessions + .get_mut(&peer_addr) + .unwrap() + .set_rekey_state(our_handshake, true); + + let forged = forge_setup_for(&node); + node.handle_session_payload(&peer_addr, &stub_link_peer(), &forged, 1280, false) + .await; + + assert_eq!( + node.stats().session.rekey_tiebreak, + 1, + "winning the dual-initiation tie-break must be counted; its DEBUG line \ + is invisible at the default log level" + ); + assert_eq!(node.stats().session.rekey_yielded, 0); + assert_eq!( + node.stats().session.rekey_armed, + 0, + "we won, so nothing may have been armed for the sender" + ); + assert!( + node.sessions + .get(&peer_addr) + .unwrap() + .has_rekey_in_progress(), + "our own rekey must survive, which is the behaviour the counter reports" + ); +} + +#[tokio::test] +async fn test_setup_naming_a_peer_whose_address_sorts_below_ours_yields_our_rekey_and_counts_it() { + let mut config = Config::new(); + config.node.rekey.enabled = false; + let mut node = make_node_with(config); + + // Our address sorts larger, so the tie-break makes us the responder. + let peer = peer_identity_sorting_below(node.node_addr()); + let peer_addr = install_established_peer(&mut node, &peer); + + let our_handshake = our_rekey_initiator_handshake(&node, &peer); + node.sessions + .get_mut(&peer_addr) + .unwrap() + .set_rekey_state(our_handshake, true); + + let forged = forge_setup_for(&node); + node.handle_session_payload(&peer_addr, &stub_link_peer(), &forged, 1280, false) + .await; + + assert_eq!( + node.stats().session.rekey_yielded, + 1, + "yielding our own rekey to an unauthenticated setup message must be \ + counted; a sustained rate here is local key rotation being suppressed" + ); + assert_eq!(node.stats().session.rekey_tiebreak, 0); + // The yield counter is recorded before the SessionAck send, so this + // assertion needs no routing. The two below depend on the send failing: + // a standalone node has no peers and an empty coord cache, so + // `send_session_datagram` returns and the responder arming never runs. + assert_eq!( + node.stats().session.rekey_armed, + 0, + "no route, so the handler returns before arming the responder side" + ); + assert!( + !node + .sessions + .get(&peer_addr) + .unwrap() + .has_rekey_in_progress(), + "our rekey was abandoned by the yield" + ); +} + +#[tokio::test] +async fn test_losing_the_tiebreak_against_a_peer_armed_handshake_keeps_the_completed_epoch() { + let mut config = Config::new(); + config.node.rekey.enabled = false; + let mut node = make_node_with(config); + + // Our address sorts larger, so the second setup loses the tie-break. + // Which side of it a given pair lands on is fixed by the two addresses, + // not chosen by the sender, so this is half of all peers rather than + // something an attacker selects. + let peer = peer_identity_sorting_below(node.node_addr()); + let peer_addr = install_established_peer(&mut node, &peer); + + // What a first forged setup leaves: a handshake the *stranger* armed, + // beside a completed epoch too stale for `pending_outranks` to veto. The + // tie-break arm gates on `has_rekey_in_progress`, which this satisfies, + // so a second forged setup reaches the yield with a pending session + // present. Nothing here required us to be the rekey initiator. + let stranger = Identity::generate(); + arm_stranger_handshake_beside_stale_pending(&mut node, &peer_addr, &peer, &stranger); + assert!( + !node.sessions.get(&peer_addr).unwrap().is_rekey_initiator(), + "the state under test is a handshake we did not arm" + ); + + let forged = forge_setup_for(&node); + node.handle_session_payload(&peer_addr, &stub_link_peer(), &forged, 1280, false) + .await; + + let entry = node.sessions.get(&peer_addr).expect("session present"); + assert_eq!( + node.stats().session.rekey_yielded, + 1, + "the test must actually reach the yield arm, or it proves nothing" + ); + assert!( + entry.pending_new_session().is_some(), + "yielding a tie-break to an unauthenticated setup must not discard \ + the key epoch the peer may already have cut over to; two forged \ + setups would otherwise kill the reverse direction" + ); + assert!( + entry.is_established(), + "the running session must be left intact alongside the pending one" + ); + assert!( + !entry.has_rekey_in_progress(), + "the handshake we yielded must still be abandoned" + ); +} + +#[tokio::test] +async fn test_a_responder_handshake_with_no_peer_rekey_stamp_is_not_expired_by_the_tick_loop() { + let peer = Identity::generate(); + let stranger = Identity::generate(); + let (mut node, peer_addr) = make_node_with_established_peer(false, &peer); + + // Arm a responder-side handshake but leave `last_peer_rekey_ms` at zero. + // The expiry predicate's `!= 0` conjunct is what stops that unstamped + // zero being read as an age of the whole Unix epoch. This pins a + // defence-in-depth guard: the state is unreachable in production, since + // the only responder arming stamps the field on the adjacent line. + let (responder, _msg3) = drive_xk_to_msg3(&stranger, node.identity()); + node.sessions + .get_mut(&peer_addr) + .unwrap() + .set_rekey_state(responder, false); + assert_eq!( + node.sessions.get(&peer_addr).unwrap().last_peer_rekey_ms(), + 0, + "test fixture must actually leave the stamp unset" + ); + + node.check_session_rekey().await; + + assert!( + node.sessions + .get(&peer_addr) + .unwrap() + .has_rekey_in_progress(), + "an unstamped handshake must not be read as infinitely old" + ); + assert_eq!(node.stats().session.rekey_expired, 0); } /// Wall-clock milliseconds, matching the clock the tick loop reads. @@ -3903,7 +4789,7 @@ async fn test_setup_naming_peer_with_pending_session_is_dropped_and_counted() { .set_pending_session(pending); let forged = forge_setup_for(&node); - node.handle_session_payload(&peer_addr, &forged, 1280, false) + node.handle_session_payload(&peer_addr, &stub_link_peer(), &forged, 1280, false) .await; let entry = node.sessions.get(&peer_addr).unwrap(); diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 1d0f535b..63c25f1d 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -1787,7 +1787,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), @@ -1827,7 +1827,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" ); @@ -1865,7 +1865,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); @@ -1874,7 +1874,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), @@ -1904,7 +1904,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); @@ -1913,7 +1913,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/noise/handshake.rs b/src/noise/handshake.rs index 737c0f02..a7de40bc 100644 --- a/src/noise/handshake.rs +++ b/src/noise/handshake.rs @@ -13,6 +13,11 @@ use std::fmt; /// Symmetric state during handshake. /// /// Maintains the chaining key (ck), handshake hash (h), and current cipher. +/// +/// `Clone` exists for [`HandshakeState::try_read_xk_message_2`], which has to +/// put the pre-read state back after a message that mixed material in before +/// failing to authenticate. +#[derive(Clone)] struct SymmetricState { /// Chaining key for key derivation. ck: [u8; 32], @@ -743,6 +748,39 @@ impl HandshakeState { Ok(()) } + /// Read XK message 2, leaving the handshake untouched when the message + /// does not authenticate. + /// + /// `read_xk_message_2` mixes the sender's ephemeral into the symmetric + /// state before it authenticates the encrypted epoch, so a message that + /// fails partway leaves a handshake that can never read the genuine msg2 + /// afterwards. A caller that keeps its session entry across a failed read + /// — because the message may be a forgery rather than a real peer's + /// corrupt reply — needs the pre-read state back. + /// + /// The saved set is exactly what `read_xk_message_2` writes: + /// `symmetric`, `remote_ephemeral`, `remote_epoch` and `progress`. **That + /// mirror is manual.** A later edit that adds a write to + /// `read_xk_message_2` without adding it here silently reintroduces the + /// poisoning, and no caller can detect it. + pub fn try_read_xk_message_2(&mut self, message: &[u8]) -> Result<(), NoiseError> { + let symmetric = self.symmetric.clone(); + let remote_ephemeral = self.remote_ephemeral; + let remote_epoch = self.remote_epoch; + let progress = self.progress; + + match self.read_xk_message_2(message) { + Ok(()) => Ok(()), + Err(e) => { + self.symmetric = symmetric; + self.remote_ephemeral = remote_ephemeral; + self.remote_epoch = remote_epoch; + self.progress = progress; + Err(e) + } + } + } + /// Write XK message 3 (initiator only). /// /// XK msg3: `-> s, se` + encrypted epoch diff --git a/src/nostr/mod.rs b/src/nostr/mod.rs index af300adf..b7f5a5d3 100644 --- a/src/nostr/mod.rs +++ b/src/nostr/mod.rs @@ -2,6 +2,7 @@ mod advert; mod driver; mod failure_state; mod handoff; +mod offer_admission; mod runtime; mod signal; mod stun; @@ -12,6 +13,8 @@ mod types; #[cfg(test)] mod tests; +pub(crate) use signal::FRESHNESS_SKEW_TOLERANCE_MS; + pub use driver::{AdvertTransportSnapshot, RendezvousDriver}; pub use handoff::{BootstrapHandoffResult, EstablishedTraversal, is_punch_packet}; pub use runtime::NostrRendezvous; diff --git a/src/nostr/offer_admission.rs b/src/nostr/offer_admission.rs new file mode 100644 index 00000000..521b2a96 --- /dev/null +++ b/src/nostr/offer_admission.rs @@ -0,0 +1,202 @@ +//! Per-npub admission control for inbound traversal offers. +//! +//! The intake path used to hold a single global semaphore sized by +//! `max_concurrent_incoming_offers`, acquired before any identity check, so +//! one sender could hold every slot and deny traversal onboarding to every +//! other peer. This type takes a per-npub permit and a global permit +//! together, so a single npub can occupy at most its own allowance. +//! +//! **What this does not fix.** The global pool stays exhaustible. Nostr +//! identities are free to generate and the signal subscription carries no +//! author restriction, so an attacker willing to run +//! `ceil(global_limit / per_sender_limit)` throwaway npubs still saturates +//! the pool at an unchanged total offer rate, and an honest peer still gets +//! the same drop. What the per-npub allowance buys is that one identity can +//! no longer do it alone, and that the refusal a spamming sender receives is +//! distinguishable in the log from genuine saturation. A defence that raises +//! the attacker's cost in more than keypairs would have to price the offer +//! itself, which is a larger change than this one. +//! +//! **Lock discipline, which the correctness of the map bound rests on.** +//! `try_admit` does the prune, the map lookup and both permit acquisitions +//! under one `std::sync::Mutex` and never awaits inside it. The prune uses +//! `Arc::strong_count` as an exact idle test: an `OwnedSemaphorePermit` owns +//! an `Arc`, so a count of 1 means the map holds the only +//! reference and no permit for that npub is outstanding. Moving +//! `try_acquire_owned` outside the lock breaks this — a concurrent prune +//! could evict an entry whose permit was still live, and the next offer from +//! that npub would build a fresh semaphore and over-admit. No test in this +//! module catches that rearrangement, so it has to be preserved by reading. + +use std::collections::HashMap; +use std::sync::{Arc, Mutex}; + +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; + +/// Why an inbound offer was refused a slot. +/// +/// The two classes stay apart because they are different operator stories: +/// one says the node is saturated, the other says a single sender is over its +/// allowance and everyone else is unaffected. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum AdmissionReject { + /// The sender already holds its whole per-npub allowance. + SenderFull, + /// The node is at `max_concurrent_incoming_offers` across all senders. + GlobalFull, +} + +/// A granted admission. Both permits release on drop, so the caller only has +/// to keep this alive for as long as the offer is being processed. +pub(super) struct OfferPermit { + _sender: OwnedSemaphorePermit, + _global: OwnedSemaphorePermit, +} + +/// Admits inbound offers against a per-sender allowance nested inside a +/// global bound. +pub(super) struct OfferAdmission { + global: Arc, + per_sender: Mutex>>, + per_sender_limit: usize, +} + +impl OfferAdmission { + /// Build an admission gate bounded globally by `global_limit` and per + /// sender npub by `per_sender_limit`. + pub(super) fn new(global_limit: usize, per_sender_limit: usize) -> Self { + Self { + global: Arc::new(Semaphore::new(global_limit)), + per_sender: Mutex::new(HashMap::new()), + per_sender_limit, + } + } + + /// Try to take one slot for `npub`, or say which bound refused it. + /// + /// The per-sender check runs first so a spamming sender is turned away + /// without ever touching the global pool, and therefore causes no churn + /// there. + pub(super) fn try_admit(&self, npub: &str) -> Result { + let mut map = self + .per_sender + .lock() + .expect("offer-admission mutex poisoned"); + + // Drop entries with no outstanding permit. Every live per-sender + // permit implies a live global permit, so at most `global_limit` + // senders survive a prune and the map is bounded by + // `global_limit + 1` after the insert below. A benign race with a + // permit dropping concurrently can only retain an idle entry for one + // more round, or drop an entry all of whose permits are already free; + // neither over-admits. + map.retain(|_, sem| Arc::strong_count(sem) > 1); + + let sem = map + .entry(npub.to_string()) + .or_insert_with(|| Arc::new(Semaphore::new(self.per_sender_limit))) + .clone(); + let sender = sem + .try_acquire_owned() + .map_err(|_| AdmissionReject::SenderFull)?; + let global = self + .global + .clone() + .try_acquire_owned() + .map_err(|_| AdmissionReject::GlobalFull)?; + + Ok(OfferPermit { + _sender: sender, + _global: global, + }) + } + + /// How many sender npubs the map currently holds. + #[cfg(test)] + pub(super) fn tracked_senders(&self) -> usize { + self.per_sender + .lock() + .expect("offer-admission mutex poisoned") + .len() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn one_sender_cannot_take_more_than_its_allowance_or_deny_another_sender() { + let admission = OfferAdmission::new(16, 4); + + // The attacker pushes for the whole global pool, which is what makes + // the victim's assertion below discriminating: without per-npub + // accounting it takes all sixteen and the victim is refused. + let mut held = Vec::new(); + let mut refusals = Vec::new(); + for _ in 0..16 { + match admission.try_admit("npub1attacker") { + Ok(permit) => held.push(permit), + Err(reject) => refusals.push(reject), + } + } + + assert_eq!( + held.len(), + 4, + "one npub must not hold more than its own allowance" + ); + assert!( + refusals + .iter() + .all(|reject| *reject == AdmissionReject::SenderFull), + "a sender over its allowance is refused on its own account, not on the pool's: {refusals:?}" + ); + assert!( + admission.try_admit("npub1victim").is_ok(), + "a second sender must still be admitted while the pool has slots" + ); + } + + #[test] + fn several_peers_bootstrapping_at_once_are_all_admitted() { + let admission = OfferAdmission::new(16, 4); + + let mut held = Vec::new(); + for peer in 0..8 { + held.push( + admission + .try_admit(&format!("npub1peer{peer}")) + .unwrap_or_else(|e| panic!("peer {peer} should be admitted, got {e:?}")), + ); + } + } + + #[test] + fn an_idle_sender_is_reclaimed_so_the_map_cannot_grow_without_bound() { + let admission = OfferAdmission::new(16, 4); + + // Three senders stay busy for the whole test, so the map bound and + // not a collapse to a single entry is what the assertion measures. + let mut held = Vec::new(); + for peer in 0..3 { + held.push( + admission + .try_admit(&format!("npub1busy{peer}")) + .expect("a busy peer is inside both bounds"), + ); + } + + for sender in 0..100 { + drop( + admission + .try_admit(&format!("npub1transient{sender}")) + .expect("a transient sender is inside both bounds"), + ); + } + + // Three still-held entries, plus the last transient sender, which was + // inserted after the prune that would have removed it. + assert_eq!(admission.tracked_senders(), 4); + } +} diff --git a/src/nostr/runtime.rs b/src/nostr/runtime.rs index c34933cc..cfa0f0c6 100644 --- a/src/nostr/runtime.rs +++ b/src/nostr/runtime.rs @@ -12,13 +12,14 @@ use nostr::prelude::{ }; use nostr_sdk::{Client, ClientOptions, prelude::RelayPoolNotification}; use serde::Serialize; -use tokio::sync::{Mutex, Notify, RwLock, Semaphore, broadcast, mpsc, oneshot}; +use tokio::sync::{Mutex, Notify, RwLock, broadcast, mpsc, oneshot}; use tokio::task::JoinHandle; use tracing::{debug, info, trace, warn}; use super::advert::{AdvertMachine, PublishPlan}; use super::failure_state::FailureState; use super::handoff::EstablishedTraversal; +use super::offer_admission::{AdmissionReject, OfferAdmission}; use super::signal::{ FreshnessOutcome, SignalEnvelope, build_signal_event, create_traversal_answer, create_traversal_offer, estimate_clock_skew, unwrap_signal_event, validate_offer_freshness, @@ -204,7 +205,7 @@ pub struct NostrRendezvous { advert: AdvertMachine, traversal: TraversalMachine, pending_answers: Mutex>>>, - offer_slots: Arc, + admission: OfferAdmission, event_tx: mpsc::UnboundedSender, event_rx: Mutex>, connect_task: Mutex>>, @@ -304,7 +305,10 @@ impl NostrRendezvous { let pubkey = keys.public_key(); let npub = crate::encode_npub(&identity.pubkey()); let (event_tx, event_rx) = mpsc::unbounded_channel(); - let offer_slots = Arc::new(Semaphore::new(config.max_concurrent_incoming_offers)); + let admission = OfferAdmission::new( + config.max_concurrent_incoming_offers, + config.max_concurrent_offers_per_npub, + ); let failure_state = FailureState::new( config.failure_streak_threshold, @@ -332,7 +336,7 @@ impl NostrRendezvous { advert, traversal, pending_answers: Mutex::new(HashMap::new()), - offer_slots, + admission, event_tx, event_rx: Mutex::new(event_rx), connect_task: Mutex::new(None), @@ -801,13 +805,30 @@ impl NostrRendezvous { && offer.message_type == "offer" && offer.recipient_npub == self.npub { - let Ok(permit) = self.offer_slots.clone().try_acquire_owned() else { - warn!( - sender_npub = %sender_npub, - limit = self.config.max_concurrent_incoming_offers, - "rate-limited inbound traversal offer (max_concurrent_incoming_offers reached); offer dropped" - ); - continue; + let permit = match self.admission.try_admit(&sender_npub) { + Ok(permit) => permit, + Err(AdmissionReject::GlobalFull) => { + warn!( + sender_npub = %sender_npub, + limit = self.config.max_concurrent_incoming_offers, + "rate-limited inbound traversal offer (max_concurrent_incoming_offers reached); offer dropped" + ); + continue; + } + // Debug, not warn: the party that trips this is by + // definition sending faster than the node wants, so + // a record per rejection turns the spam into log + // volume. The global-full arm above stays at warn + // and remains the operator's signal that the node + // is actually saturated. + Err(AdmissionReject::SenderFull) => { + debug!( + sender_npub = %sender_npub, + limit = self.config.max_concurrent_offers_per_npub, + "inbound traversal offer refused: sender is at its per-npub offer allowance" + ); + continue; + } }; let runtime = Arc::clone(&self); let peer_short = short_npub(&sender_npub); @@ -1791,7 +1812,10 @@ impl NostrRendezvous { .opts(ClientOptions::new().autoconnect(false)) .build(); let config = NostrRendezvousConfig::default(); - let offer_slots = Arc::new(Semaphore::new(config.max_concurrent_incoming_offers)); + let admission = OfferAdmission::new( + config.max_concurrent_incoming_offers, + config.max_concurrent_offers_per_npub, + ); let (event_tx, event_rx) = mpsc::unbounded_channel(); let failure_state = FailureState::new( config.failure_streak_threshold, @@ -1818,7 +1842,7 @@ impl NostrRendezvous { advert, traversal, pending_answers: Mutex::new(HashMap::new()), - offer_slots, + admission, event_tx, event_rx: Mutex::new(event_rx), connect_task: Mutex::new(None), diff --git a/src/nostr/signal.rs b/src/nostr/signal.rs index 29b8b855..6b2306e5 100644 --- a/src/nostr/signal.rs +++ b/src/nostr/signal.rs @@ -11,7 +11,11 @@ use super::types::{BootstrapError, PunchHint, SIGNAL_KIND, TraversalAnswer, Trav /// past ~minutes erodes the freshness guarantee that backstops session-id /// replay protection. Tightening it below the size of a typical un-NTP'd /// drift defeats the purpose. 60s sits comfortably between those. -pub(super) const FRESHNESS_SKEW_TOLERANCE_MS: u64 = 60_000; +/// +/// `pub(crate)` because `Config::validate` derives the freshness window from +/// it rather than restating the number, so changing it here moves the +/// validation boundary with it. +pub(crate) const FRESHNESS_SKEW_TOLERANCE_MS: u64 = 60_000; pub(super) struct SignalEnvelope { pub(super) payload: T, diff --git a/src/nostr/tests.rs b/src/nostr/tests.rs index ec260797..c7bcdee4 100644 --- a/src/nostr/tests.rs +++ b/src/nostr/tests.rs @@ -523,15 +523,76 @@ fn planned_remote_endpoints_cap_targets_from_an_oversized_candidate_list() { ) .expect("endpoint planning should succeed"); + // The exact figures are determined by the inputs: the plan is one + // reflexive-to-reflexive pair plus 300 reflexive-to-candidate pairs, all + // distinct, so the cap lands on exactly MAX_PUNCH_TARGETS. An inequality + // here could not tell the cap from a collapse to a single target. assert!(tally.capped > 0, "the cap should have discarded targets"); assert!(tally.suspicious()); - assert!( - endpoints.len() <= 8, - "expected at most 8 endpoints, got {}", - endpoints.len() + assert_eq!(tally.admitted, 8); + assert_eq!(endpoints.len(), 8); +} + +/// The IPv6 arm of `is_never_punchable_ip` is exercised as a pure predicate +/// elsewhere; this drives it end to end through the planner, which is the +/// path the reflector attack actually uses. +/// +/// The local address is a ULA so `lan_refs` is non-empty and `same_subnet_24` +/// is genuinely called with two IPv6 strings. It splits on `.` and requires +/// four parts, so no IPv6 candidate can ever satisfy the /24 gate and +/// `fd00::1` is refused off-subnet rather than admitted. +#[test] +fn planned_remote_endpoints_reject_never_punchable_ipv6_candidates() { + let (endpoints, _tally) = planned_remote_endpoints( + &[addr("fd00::2", 62000)], + Some(&addr("203.0.113.10", 62000)), + &[ + addr("::1", 63000), + addr("::", 63000), + addr("ff02::1", 63000), + addr("fe80::1", 63000), + addr("fd00::1", 63000), + ], + Some(&addr("198.51.100.20", 63000)), + ) + .expect("endpoint planning should succeed"); + + assert_eq!( + endpoints, + vec!["198.51.100.20:63000".parse::().unwrap()] ); } +/// The cap test above uses public candidates so the cap and not the filter is +/// what bounds the output. The attack shape is the opposite: several hundred +/// attacker-chosen unroutable addresses. This is the only test that pins the +/// filter and the cap acting together on that input, and its failure would +/// mean the reflector is back. +#[test] +fn planned_remote_endpoints_bound_an_oversized_list_of_unroutable_candidates() { + let mut remotes = Vec::new(); + for host in 1..=150u8 { + remotes.push(addr(&format!("127.0.0.{host}"), 63000)); + remotes.push(addr(&format!("224.0.0.{host}"), 63000)); + } + + let (endpoints, tally) = planned_remote_endpoints( + &[], + Some(&addr("203.0.113.10", 62000)), + &remotes, + Some(&addr("198.51.100.20", 63000)), + ) + .expect("endpoint planning should succeed"); + + assert_eq!( + endpoints, + vec!["198.51.100.20:63000".parse::().unwrap()] + ); + assert_eq!(tally.unroutable, 300); + assert_eq!(tally.admitted, 1); + assert!(tally.suspicious()); +} + /// Guards the deployment whose STUN server sits inside the private network, /// so the observed reflexive address is itself private. Applying the /24 gate /// to a peer's reflexive address would drop it and remove the only branch diff --git a/src/nostr/traversal.rs b/src/nostr/traversal.rs index 8521e98c..b5309a97 100644 --- a/src/nostr/traversal.rs +++ b/src/nostr/traversal.rs @@ -469,14 +469,26 @@ pub(super) fn nonce() -> String { /// step, which is what a resume produces, saturates the punch start delay to /// zero so punching begins immediately; the attempt's own bounds are monotonic /// `Instant` deadlines, so its length is unaffected. A backward step lengthens -/// that delay instead and can cost a single punch attempt, which retries. Early -/// eviction from the replay window cannot admit a replay under the shipped -/// defaults, because the freshness window a replayed offer would also have to -/// satisfy (`signal_ttl_secs` plus `FRESHNESS_SKEW_TOLERANCE_MS` on each side, -/// 240s under the shipped defaults) is strictly narrower than the replay window -/// itself (`replay_window_secs`, 300s). The margin holds while -/// `signal_ttl_secs + 120 < replay_window_secs`; nothing in config validation -/// enforces that relation today. +/// that delay instead and can cost a single punch attempt, which retries. +/// Expiry-driven eviction from the replay window cannot admit a replay, because +/// the freshness window a replayed offer would also have to satisfy +/// (`signal_ttl_secs` plus `FRESHNESS_SKEW_TOLERANCE_MS` on each side, 240s +/// under the shipped defaults) is strictly narrower than the replay window +/// itself (`replay_window_secs`, 300s). `Config::validate` enforces +/// `signal_ttl_secs + 2 * FRESHNESS_SKEW_TOLERANCE_MS < replay_window_secs`, so +/// that margin can no longer be configured away. +/// +/// That covers the expiry route only. `mark_session_seen` also evicts on +/// capacity, dropping the entries nearest expiry once the cache exceeds +/// `seen_sessions_max_entries`, and a session id dropped that way stays +/// replayable for the rest of its freshness window whatever the time relation +/// is. Nothing bounds that route: whether it is reachable depends on the +/// admitted-offer rate against the cache size. The shipped defaults leave a +/// wide margin — filling 2048 entries inside 300s needs about 6.8 admitted +/// offers per second, against roughly 1.2/s from a 16-slot pool whose permits +/// are held for the length of an attempt — but raising +/// `max_concurrent_incoming_offers` narrows it. That admission rate is inferred +/// from the attempt timeout, not measured. pub(super) fn now_ms() -> u64 { SystemTime::now() .duration_since(UNIX_EPOCH) diff --git a/src/proto/lookup/core.rs b/src/proto/lookup/core.rs index 291ba8e9..2eeb5e57 100644 --- a/src/proto/lookup/core.rs +++ b/src/proto/lookup/core.rs @@ -36,7 +36,15 @@ pub(crate) enum LookupAction { path_mtu: u16, }, /// Mirror path_mtu into the FipsAddress-keyed TUN-shared lookup map. - WritePathMtu { target: NodeAddr, path_mtu: u16 }, + /// + /// `now_ms` stamps the entry's expiry deadline. It is the same instant + /// [`Self::CacheCoords`] carries, because the clamp this write feeds must + /// not outlive the route it was learned with. + WritePathMtu { + target: NodeAddr, + now_ms: u64, + path_mtu: u16, + }, /// Reset the coords-warmup counter if an established session exists. ResetWarmupIfEstablished { target: NodeAddr }, /// Retry queued TUN packets for the target if any are pending. @@ -270,6 +278,7 @@ pub(crate) fn on_response_accepted( }, LookupAction::WritePathMtu { target: *target, + now_ms, path_mtu, }, LookupAction::ResetWarmupIfEstablished { target: *target }, diff --git a/src/proto/lookup/tests/core.rs b/src/proto/lookup/tests/core.rs index a783c3b3..e402e9e0 100644 --- a/src/proto/lookup/tests/core.rs +++ b/src/proto/lookup/tests/core.rs @@ -233,9 +233,15 @@ fn on_response_accepted_clears_state_and_emits_effects() { match &actions[1] { LookupAction::WritePathMtu { target: t, + now_ms: n, path_mtu: p, } => { assert_eq!(*t, target); + assert_eq!( + *n, now_ms, + "the path-MTU entry must be stamped with the same instant the \ + coordinates are cached at, or the clamp outlives its route" + ); assert_eq!(*p, path_mtu); } _ => panic!("action[1] must be WritePathMtu"), diff --git a/src/proto/mmp/path_mtu.rs b/src/proto/mmp/path_mtu.rs index 3c7fdca6..8407de97 100644 --- a/src/proto/mmp/path_mtu.rs +++ b/src/proto/mmp/path_mtu.rs @@ -179,6 +179,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_ms`) 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_ms = None; + self.pending_increase_mtu = 0; + } } impl Default for PathMtuState { diff --git a/src/proto/mmp/tests/path_mtu.rs b/src/proto/mmp/tests/path_mtu.rs index 3ba46efe..c0899478 100644 --- a/src/proto/mmp/tests/path_mtu.rs +++ b/src/proto/mmp/tests/path_mtu.rs @@ -63,3 +63,54 @@ fn apply_notification_increase_counter_does_not_overflow_on_a_repeated_value() { "the increase is not yet due, so the effective MTU must be unchanged" ); } + +#[test] +fn 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_ms`, so the + // increase can never be accepted at all. + let t0 = 100_000u64; + 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 + 11_000); + 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 + 20_000)); + 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 + 30_000; + state.apply_notification(1400, t1); + state.apply_notification(1400, t1 + 11_000); + state.apply_notification(1400, t1 + 21_000); + + 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/proto/stp/tests/limits.rs b/src/proto/stp/tests/limits.rs index f0324e8a..097b6f3d 100644 --- a/src/proto/stp/tests/limits.rs +++ b/src/proto/stp/tests/limits.rs @@ -442,3 +442,77 @@ fn test_handle_parent_lost_keeps_its_published_bool_return() { assert!(published(&mut state, &BTreeMap::new(), 2000, 2000)); assert!(state.is_root()); } + +#[test] +fn test_flap_dampening_counter_resets_when_the_episode_lapses_not_when_it_engages() { + // Retirement clears the switch counter when a lapsed episode is retired, + // not when the episode is armed, so switches taken *during* an episode do + // not carry into the next window. Telling those two apart needs an episode + // that is genuinely live for at least one switch, which the injected clock + // supplies without any real elapsed time. + let my_node = make_node_addr(5); + let mut state = TreeState::new(my_node, 1000); + state.set_flap_dampening(3, 60, 2); + state.set_hold_down(0); + + let peer_a = make_node_addr(1); + let peer_b = make_node_addr(2); + let root = make_node_addr(0); + + state.update_peer( + ParentDeclaration::new(peer_a, root, 1, 1000), + make_coords(&[1, 0]), + ); + state.update_peer( + ParentDeclaration::new(peer_b, root, 1, 1000), + make_coords(&[2, 0]), + ); + + assert!(!state.set_parent(peer_a, 1, 1000, 1000)); + state.recompute_coords(); + assert!(!state.set_parent(peer_b, 2, 2000, 2000)); + state.recompute_coords(); + assert!( + state.set_parent(peer_a, 3, 3000, 3000), + "the third switch inside the window must arm an episode" + ); + state.recompute_coords(); + assert!( + state.is_flap_dampened(3000), + "a two-second episode must be live immediately after it is armed" + ); + + // In-episode accumulation: in production a mandatory switch bypasses the + // veto and still feeds the counter. `set_parent` is the same entry point + // those switches reach, and it must not re-arm the running episode. + assert!( + !state.set_parent(peer_b, 4, 4000, 4000), + "a switch inside a live episode must not re-arm it" + ); + state.recompute_coords(); + assert!(state.is_flap_dampened(4000)); + + assert!( + !state.is_flap_dampened(6000), + "the episode must have lapsed after its two seconds" + ); + + // The lapsed episode is retired on the next switch, taking the counter + // with it: the in-episode switch must not count toward the next episode. + assert!( + !state.set_parent(peer_a, 5, 5000, 6000), + "retiring a lapsed episode must clear the switch counter" + ); + state.recompute_coords(); + assert!( + !state.set_parent(peer_b, 6, 6000, 7000), + "a second episode must cost a fresh threshold of switches" + ); + state.recompute_coords(); + assert!( + state.set_parent(peer_a, 7, 7000, 8000), + "a second episode must arm once the fresh threshold is reached" + ); + state.recompute_coords(); + assert!(state.is_flap_dampened(8000)); +} diff --git a/src/upper/tun.rs b/src/upper/tun.rs index 29b6f74c..45262aeb 100644 --- a/src/upper/tun.rs +++ b/src/upper/tun.rs @@ -35,12 +35,50 @@ use tracing::{error, warn}; #[cfg(any(target_os = "linux", target_os = "macos", target_os = "freebsd"))] 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 @@ -107,7 +145,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, @@ -117,6 +155,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 @@ -1631,7 +1670,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); } @@ -1641,7 +1683,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); } @@ -1655,7 +1700,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); @@ -1671,7 +1719,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, @@ -1705,7 +1756,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(); @@ -1725,10 +1779,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); } @@ -1757,8 +1811,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); } diff --git a/testing/nat/scripts/generate-configs.sh b/testing/nat/scripts/generate-configs.sh index a2d07fbc..409a2efa 100755 --- a/testing/nat/scripts/generate-configs.sh +++ b/testing/nat/scripts/generate-configs.sh @@ -102,7 +102,10 @@ node: - "$stun_addr" signal_ttl_secs: 30 attempt_timeout_secs: 6 - replay_window_secs: 60 + # Must stay above signal_ttl_secs plus 60s of clock-skew grace on each + # side, or config validation refuses to start the node. 180 keeps a 30s + # margin over the 150s freshness window the 30s TTL implies. + replay_window_secs: 180 punch_start_delay_ms: 500 punch_interval_ms: 100 punch_duration_ms: 2500