diff --git a/CHANGELOG.md b/CHANGELOG.md index ffe9e4f7..ff0c67a1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,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 @@ -132,6 +141,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. + - A SessionDatagram carrying a truncated inner FSP payload no longer panics the forwarding path. The coordinate-cache warm path sliced the inner payload at the full 12-byte header offset while guarding only with the 4-byte common diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index 68ab80a4..305da063 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. @@ -932,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 1a44f38b..995ee4de 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -973,6 +973,25 @@ 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 diff --git a/src/config/node.rs b/src/config/node.rs index 36a573b2..eead8709 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.*`). diff --git a/src/node/handlers/forwarding.rs b/src/node/handlers/forwarding.rs index 0c93b318..c0b4dd42 100644 --- a/src/node/handlers/forwarding.rs +++ b/src/node/handlers/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, ) { @@ -58,6 +58,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, diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 83420569..2e892075 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -8,6 +8,7 @@ use crate::NodeAddr; use crate::mmp::report::ReceiverReport; use crate::mmp::{MAX_SESSION_REPORT_INTERVAL_MS, MIN_SESSION_REPORT_INTERVAL_MS}; +use crate::node::rate_limit::Msg1Class; use crate::node::reject::{RejectReason, SessionReject}; use crate::node::session::{EndToEndState, EpochSlot, SessionEntry}; use crate::node::session_wire::{ @@ -95,9 +96,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, @@ -117,7 +124,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; @@ -451,7 +458,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) => { @@ -469,6 +481,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() { @@ -519,8 +565,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(), @@ -541,13 +592,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 { @@ -703,6 +769,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, @@ -792,11 +870,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, @@ -878,6 +979,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, @@ -894,7 +1013,7 @@ impl Node { error = %e, "Failed to process rekey XK msg3" ); - entry.abandon_rekey(); + entry.abandon_handshake(); self.sessions.insert(*src_addr, entry); return; } @@ -906,8 +1025,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; } @@ -917,7 +1042,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)); @@ -928,8 +1053,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; } diff --git a/src/node/mod.rs b/src/node/mod.rs index 96d6663e..458fda71 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -32,7 +32,7 @@ mod tree; pub(crate) mod wire; use self::discovery_rate_limit::{DiscoveryBackoff, DiscoveryForwardRateLimiter}; -use self::rate_limit::HandshakeRateLimiter; +use self::rate_limit::{HandshakeRateLimiter, SessionSetupRateLimiter}; use self::reloadable::Reloadable; use self::routing_error_rate_limit::RoutingErrorRateLimiter; @@ -492,6 +492,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, /// Rate limiter for routing error signals (CoordsRequired / PathBroken). @@ -639,6 +641,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 { @@ -681,6 +705,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; @@ -764,6 +789,7 @@ impl Node { peers_by_index: HashMap::new(), pending_outbound: HashMap::new(), msg1_rate_limiter, + setup_rate_limiter, icmp_rate_limiter: IcmpRateLimiter::new(), routing_error_rate_limiter: RoutingErrorRateLimiter::new(), coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval( @@ -841,6 +867,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; @@ -921,6 +948,7 @@ impl Node { peers_by_index: HashMap::new(), pending_outbound: HashMap::new(), msg1_rate_limiter, + setup_rate_limiter, icmp_rate_limiter: IcmpRateLimiter::new(), routing_error_rate_limiter: RoutingErrorRateLimiter::new(), coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval( 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 d156b249..13c6cb58 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/stats.rs b/src/node/stats.rs index 0d69c4ce..ecaa830f 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/session.rs b/src/node/tests/session.rs index e1a38f18..44f800c7 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -9,6 +9,15 @@ use crate::node::tests::spanning_tree::{ }; use crate::protocol::{CoordsRequired, PathBroken, SessionAck, SessionDatagram, SessionMsg3}; +/// 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 @@ -3480,8 +3489,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(), @@ -3513,8 +3528,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 @@ -3548,8 +3569,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 @@ -3584,8 +3611,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()); @@ -3623,8 +3656,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!( @@ -3634,6 +3673,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 // ============================================================================ @@ -3731,7 +3877,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] @@ -3888,6 +4034,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 // ============================================================================ @@ -3901,6 +4292,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); @@ -3913,7 +4314,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. @@ -4015,7 +4628,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/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