diff --git a/CHANGELOG.md b/CHANGELOG.md index 87f52319..a03f3012 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,7 +26,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 `Degraded`, and starts carrying traffic the moment a port comes up. Whether an interface has carrier is reported separately as `interface.carrier` in `show_transports`, never acted on. - This closes the OpenWrt boot race (procd starts `fips` before wifi has created `fips-mesh0` / `fips-ap0`; both transports were skipped for the life of the process while the 802.11s peer link formed anyway, so the node looked @@ -75,6 +74,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 destroy-and-recreate, that an `optional` interface never moves node health, and that absence is logged once on the edge rather than once per retry. +#### Gateway + +- The gateway counts sessions on a kernel without `/proc/net/nf_conntrack`. + When the file is absent it dumps the conntrack table over netlink, as + `conntrack -L` does, so a mapping carrying traffic is pinned instead of + being reclaimed on its TTL and grace period alone. Kernels built without + `CONFIG_NF_CONNTRACK_PROCFS`, such as Ubuntu's, had session pinning off. +- The gateway says at startup whether it can read conntrack sessions. It + reads the table once, as each tick does, and logs either the source it read + or that no source is readable and session pinning is off. An operator on a + kernel with no readable source learned this only from a warning at the first + failed tick. + #### Node lifecycle - Transport-medium change detection, controlled by the new `node.netmon.*` @@ -157,6 +169,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 instead of them. The platform wiring (`registerNetworkCallback`, or `NWPathMonitor` on iOS) stays with the embedder, since the crate has no JNI layer. + - `TransportError::InterfaceUnavailable { interface }`. A missing interface and a typo'd interface name were previously the same flat `StartFailed(String)`; nothing downstream could branch on absence. @@ -211,7 +224,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 creating an interface — the daemon binds it on its own. `phy0-sta0` (`wwan`) is marked `optional: true` for the same reason: it only exists while a radio is in station mode. - **Upgrade note: an existing `/etc/fips/fips.yaml` is preserved and does not gain the new key.** It is a package conffile, so on a router where `fips-mesh-setup` or `fips-ap-setup` had already uncommented a block, that @@ -316,6 +328,41 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 binder tearing down and rebinding every second while teardown silently declined to abort anything. +#### Packaging (Debian) + +- An upgrade of the `.deb` now reapplies the firewall ruleset in place. Until + now an upgrade reloaded nothing, so a changed `/etc/fips/fips.nft` took + effect only at the next reboot or manual restart, and a restart deletes the + `fips` table and leaves the mesh interface unfiltered until the ruleset is + loaded again. `fips-firewall.service`, in both the Debian and the plain + systemd unit, gains a reload that replaces the ruleset in one transaction, + and the postinst reloads the unit only when it is already active, so an + upgrade never turns the firewall on for a host that has not opted in. A + reload that fails leaves the previous ruleset in place and is reported; the + upgrade goes on. + +#### Packaging (AUR) + +- The AUR publish on a release tag now waits until every package workflow of + that tag has succeeded. It used to push the new `pkgver` while the Linux, + macOS, Windows, OpenWrt and FreeBSD packages were still building; at v0.5.1 + the AUR was updated while the release had 15 of its 17 assets. Because the + AUR package pins the tag's source archive, withdrawing a bad release after + that point left the AUR package unbuildable. A failed or cancelled package + run now stops the publish, and one that has not finished within an hour + fails it. + +#### Dependencies + +- The lockfile moves `chacha20` from 0.10.1 to 0.10.2, because 0.10.1 is yanked. + It arrives through `rand`, a direct dependency, + so it sits on the built path rather than off to one side. The requirement in + `Cargo.toml` already admitted 0.10.2, so this is a lockfile change and no code + changed with it. **This is not a security fix**: `cargo audit` reports nothing + against `chacha20` at either version, and 0.10.1 was withdrawn by its + maintainer rather than flagged by an advisory. What it buys is that a fresh + checkout can resolve the lockfile without reaching for a yanked version. + ### Removed - **Source-breaking for consumers of the library crate**: `ActivePeer` no @@ -389,26 +436,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 update, or `fipsctl connect`), and a traversed link that goes quiet still releases the peer to every path. -- A heartbeat whose send failed no longer counts as one that was delivered. - The peer's "last heartbeat" timestamp was stamped before the send and left - alone whatever came back, so a failure suppressed the next attempt for a - full `node.heartbeat_interval_secs` even though the peer had heard nothing — - on a 10s interval against a 30s `link_dead_timeout_secs`, three failures in - a row were the whole budget. The timestamp now moves only on a send that - returned cleanly, and a separate record of the *attempt* spaces the retries - so a peer that keeps failing is retried in seconds rather than either - hammered every tick or left for a full interval. That retry spacing - applies to the failure path only: gating a healthy peer on it as well - would have floored `node.heartbeat_interval_secs` at two seconds, so a - configured value below that would silently not have been honoured. - - A peer that rotates its address no longer keeps sending from a socket aimed where it used to be. The authenticated-frame path updated the peer's address and discarded the flag saying it had changed, so the per-peer `connect()`-ed UDP socket stayed pinned to the old 5-tuple; the sibling path already cleared it. -#### Data plane +#### Data plane and transports - A peer that stops reading can no longer stall the node. TCP, Tor, Nym and BLE wrote to their links directly from the caller's task, and a write @@ -496,8 +530,42 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 socket is adopted, after its bind, so the traversal bind still receives a port no other socket holds. -#### Link rekey +#### Peering +- A heartbeat whose send failed no longer counts as one that was delivered. The + send was recorded before it was attempted, so a peer whose heartbeat could not + go out was treated as heartbeated and was not tried again for a whole + `heartbeat_interval_secs`, although it had heard nothing and its own link-dead + timer was running. The attempt and the delivery are now recorded separately: + the interval that paces a healthy peer advances only on a send that returned + cleanly, and a peer whose send failed is retried after a shorter fixed + interval instead. That retry interval gates only a peer whose last attempt + failed, so it cannot clamp a `heartbeat_interval_secs` configured below it. + +#### Session setup + +- A session whose last handshake message is lost no longer stays one-sided. + The initiator sent msg3 once and treated the session as established at once; + when that one datagram was lost, the responder kept waiting for it and + dropped every frame the initiator sent, and nothing sent msg3 again, because + the responder's repeated SessionAck was refused as arriving in the wrong + state. The session stayed that way until the next session rekey, or with + periodic rekey switched off, indefinitely. The initiator now keeps its msg3 + and resends it on the handshake resend interval, with backoff, until a frame + from the responder authenticates or `handshake_max_resends` resends have gone + out. The wire format is unchanged: the resend carries the same msg3, and a + responder that already completed the session refuses the duplicate as before. + +#### Link and session rekey + +- A session rekey this node started no longer stays in flight forever when its + setup or the peer's ack is lost. Nothing resends a rekey setup, and the only + expiry covered a rekey the peer started, so one lost datagram left the + rekey pending and blocked every later one: the session kept its current keys + and stopped rotating them. The rekey now expires on the handshake timeout, + timed from when this node sent its setup, and the next tick starts a fresh + one. A forged ack cannot extend it. Expiries are counted as + `rekey_unanswered`. The wire format is unchanged. - A forged rekey msg2 no longer takes the link down. The rekey initiator gave up its handshake before reading msg2 and abandoned the cycle when the read failed, although nothing authenticates a msg2 ahead of that read. Anyone on @@ -513,9 +581,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 forgery now costs the initiator the msg2 key agreement until the cycle ends, where before only the first one did; the msg1 resend budget bounds that. The wire format is unchanged. - -#### Session rekey - - A SessionAck that fails to read no longer ends a session rekey this node started. The handler took the rekey handshake off the session before reading the ack's msg2 and abandoned the rekey when the read failed, although nothing @@ -524,6 +589,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 state before the read, so the peer's genuine ack still completes the rekey, and the refusal is counted as `ack_handshake_failed`, as it already was for a first-contact session. The wire format is unchanged. +- A link rekey whose reply is lost no longer splits the link. The node that + answered a rekey used to switch to the new keys on its own next tick, before + the other side had them; when the reply was lost, frames from the answering + side were dropped until the link was torn down. The answering side now + switches only when a frame on the new keys arrives from the side that started + the rekey, and drops keys that were never adopted after a hold (120 s by + default) so the next rekey can proceed. #### Session coordinates @@ -554,7 +626,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 changes, while the row's `transport_id` and `remote_addr` stay those the link was created with. The counters cover authenticated link frames only, so they are not expected to match the transport totals in `show_transports`. - The response shape is unchanged. + The response shape is unchanged. Fixes #158. - `show_peers` (`fipsctl show peers`) now reports a peer that has gone quiet as `stale`. Its `connectivity` was read from a state that nothing outside the tests ever changed, so every peer read `connected` until it was @@ -566,7 +638,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 values the open-discovery tutorial described never occurred, and the tutorial no longer lists them. The response shape is unchanged. -#### Identity & config +#### Identity and config - A persistent node whose identity key path cannot be examined now refuses to start instead of coming up under a new identity. `Path::exists` reports false @@ -581,14 +653,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 #### Gateway -- A `.fips` query the gateway answers without an address no longer takes an - address from the pool. Every query type was allocated a mapping before the - code looked at what the client had asked for, and an A or HTTPS query was - then answered with NODATA, so any host that can reach the LAN resolver could - consume the pool one name at a time with a query type it is never given an - address for. Only AAAA and ANY allocate now. A non-AAAA query for a name that - already has a mapping still refreshes that mapping's TTL clock, so a client - querying both types does not lose half of its refresh. - Conntrack sessions are matched by address rather than by text, so live traffic pins a gateway mapping again. The session count searched each `/proc/net/nf_conntrack` line for `dst=` followed by the virtual IP in its @@ -608,19 +672,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 sessions for every mapping, as it always has, so reclamation keeps working rather than pinning the whole pool; but the first failure and each change of outcome after it are now logged, so an unreadable source is no longer - indistinguishable from an idle one. A kernel built without - `CONFIG_NF_CONNTRACK_PROCFS` has no `/proc/net/nf_conntrack` at all and fails - identically every tick, so a repeat is logged at debug rather than warn. -- The gateway says at startup whether it can read conntrack sessions. It - reads the table once, as each tick does, and logs either the source it read - or that no source is readable and session pinning is off. An operator on a - kernel with no readable source learned this only from a warning at the first - failed tick. -- The gateway counts sessions on a kernel without `/proc/net/nf_conntrack`. - When the file is absent it dumps the conntrack table over netlink, as - `conntrack -L` does, so a mapping carrying traffic is pinned instead of - being reclaimed on its TTL and grace period alone. Kernels built without - `CONFIG_NF_CONNTRACK_PROCFS`, such as Ubuntu's, had session pinning off. + indistinguishable from an idle one. A source that fails identically every + tick is logged at debug rather than warn on a repeat. - The NAT table is rebuilt in one netlink transaction. A rebuild deleted the `fips_gateway` table in a batch of its own, discarded that batch's result, and only then sent the batch that recreated the table, the chains, the @@ -643,6 +696,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 the batch and requests one acknowledgement per batch, and NAT errors now name the kernel errno. A rebuild that still fails is logged, and the next successful rebuild installs the mapping. +- A `.fips` query the gateway answers without an address no longer takes an + address from the pool. Every query type was allocated a mapping before the + code looked at what the client had asked for, and an A or HTTPS query was + then answered with NODATA, so any host that can reach the LAN resolver could + consume the pool one name at a time with a query type it is never given an + address for. Only AAAA and ANY allocate now. A non-AAAA query for a name that + already has a mapping still refreshes that mapping's TTL clock, so a client + querying both types does not lose half of its refresh. - The gateway's virtual-IP pool now limits how many mappings it holds and how fast it creates them. Any host that can reach the LAN resolver could ask for one new `.fips` name after another, and each got a mapping until the 65,535 @@ -653,6 +714,22 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 warning says which limit refused it. A name that already has a mapping is answered before either limit is consulted, so names in use keep resolving when the pool is full. The limits are compiled in, not configured. + +#### Nostr and NAT traversal + +- A node no longer publishes NIP-09 deletion requests signed with its routing + key after a NAT traversal attempt. Each request put the node's public + identity next to the ids of its offer and answer gift wraps on every relay it + reached, which the one-time signing keys on those wraps exist to prevent, and + most of the requests deleted nothing, since a relay deletes a gift wrap only + at its recipient's request. A relay that stores the wraps now keeps them + until their NIP-40 expiration; relays that do not store ephemeral events + never held them. The advertisement retraction still sends its deletion + request, since that names an event the routing key signed itself. The + discovery and traversal design documents describe the new behaviour. + +#### Packaging (OpenWrt) + - A new OpenWrt install no longer enables and starts `fips-gateway`. The generated postinst turned it on unconditionally, contradicting the init script's own header, the package README and the deployment tutorial, all of @@ -675,13 +752,30 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 port, add the LAN prefix and advertise the pool route, and only then start a daemon that exits immediately because the gateway is disabled, leaving `.fips` resolution pointed at a port nothing listens on. -- The four OpenWrt maintainer-script bodies now live in - `packaging/openwrt-ipk/scripts/` instead of inside heredocs in the two build - scripts, so the `.ipk` and `.apk` packages install the same bodies and the - scenarios in `testing/openwrt/` run what ships. +- The `.ipk` and `.apk` packages now install the same maintainer scripts. The + four script bodies live in `packaging/openwrt-ipk/scripts/` instead of inside + heredocs in the two build scripts, so the scenarios in `testing/openwrt/` run + what ships. -#### Packaging +#### Packaging (Debian) +- A `.deb` upgrade whose new daemon cannot start no longer hangs apt. The + postinst started `fips.service` and then `fips-dns.service` with blocking + calls, and because `fips-dns.service` requires the daemon, a daemon that + failed on every start left the second call, apt and everything queued behind + it waiting for ever with no message. Each start is now queued and waited on + for at most 60 seconds. A unit that does not come up has its status printed + and fails the configure step, so apt exits non-zero and names the unit; a + masked unit, or one whose condition is not met, is reported and skipped. +- The `.deb` maintainer scripts now manage `fips-gateway` with the rest of the + package's services. An upgrade stopped the daemon, which the gateway + requires, and never brought the gateway back, so an operator who had enabled + it lost it until the next reboot; removing or purging the package left the + gateway's enablement symlink behind, pointing at a unit file that no longer + exists. The gateway is now stopped before the daemon on upgrade and + restarted afterwards only when it is enabled and the daemon came up, and it + is stopped and disabled on remove and purge. A gateway that does not come + back is reported but does not fail the upgrade. - The `.deb` now declares `libgcc-s1 (>= 4.2)`. All four binaries link `libgcc_s.so.1`, but cargo-deb removes every libgcc entry from the dependencies it derives, so the package never said so. `libc6` depends on @@ -699,44 +793,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 outside the tree the build sees. The build image's tag now includes a hash of its Dockerfile, so a host with an older image cached builds a new one instead of reusing it. +- `packaging/debian/build-deb-container.sh` now returns the package it just + built. It picked the most recently modified `fips_*.deb` in the output + directory that sorted last by name, so a package with a higher version left + there by an earlier run was returned instead. -### Changed +#### Packaging (AUR) -- The lockfile moves `chacha20` from 0.10.1 to 0.10.2, because 0.10.1 is yanked. - It arrives through `rand`, a direct dependency, - so it sits on the built path rather than off to one side. The requirement in - `Cargo.toml` already admitted 0.10.2, so this is a lockfile change and no code - changed with it. **This is not a security fix**: `cargo audit` reports nothing - against `chacha20` at either version, and 0.10.1 was withdrawn by its - maintainer rather than flagged by an advisory. What it buys is that a fresh - checkout can resolve the lockfile without reaching for a yanked version. -- The dns-resolver test suite's end-to-end scenarios now run the `fips` and - `fips-gateway` binaries from the Debian package rather than compiling their - own. The suite used to build both in a Debian 12 image with whatever Rust was - current, a second release build on every CI run with no cache, and not the - toolchain or the build that ships. It now takes `--deb PATH`, and CI hands it - the package the install suite installs, so one package build serves both; run - on its own it builds the package through the same container script the - release uses. The GitHub leg moves to a job of its own that waits for the - package build, with its check name unchanged, and a local CI run builds the - package once for both suites. The suite now needs `dpkg-deb` on the host. -- The Linux release and CI package builds reuse their builder image across - GitHub runners instead of assembling it on every leg from apt, rustup and a - compile of `cargo-deb`. `build-deb-container.sh` gains `--print-image-tag` and - `--image-archive PATH`: the workflows key an Actions cache entry on the image - tag, load the image from it when present, and save it after a build. Only a - push to `maint`, `master` or `next` saves an entry; pull requests and topic - branches read the default branch's. A corrupt or mismatched archive is a - warning and a rebuild, never a failed build. A cached image is not refreshed - from apt or the base image until the base image name, the toolchain or - `Dockerfile.build` changes, as was already true of a developer's machine. -- CI now builds the arm64 `.deb` on an arm64 runner and installs it on Ubuntu - 22.04, the oldest supported distribution, starting the daemon, on every push - and pull request. Until now the arm64 package was floor-checked and never - installed anywhere in the pipeline. Its upgrade, purge and conffile paths - remain unexercised; those run on amd64 only. The parity check reads each - install leg's architecture, so the arm64 leg is reported as GitHub-only and - cannot stand in for a missing amd64 leg of the same distribution. +- The release `PKGBUILD` now lists `dbus` as a runtime dependency. The `fips` + binary links `libdbus-1`, and the `fips-git` package already declared it. ## [0.5.1] - 2026-09-06 @@ -4597,3 +4662,5 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Design documentation suite covering all protocol layers - CHANGELOG.md following Keep a Changelog format - Repository mirrored to [ngit](https://gitworkshop.dev/npub1y0gja7r4re0wyelmvdqa03qmjs62rwvcd8szzt4nf4t2hd43969qj000ly/relay.ngit.dev/fips) + + diff --git a/docs/design/fips-mesh-layer.md b/docs/design/fips-mesh-layer.md index 9de389af..b437fa22 100644 --- a/docs/design/fips-mesh-layer.md +++ b/docs/design/fips-mesh-layer.md @@ -388,10 +388,13 @@ the responder consumes msg1, builds msg2, and replies. After both sides have exchanged messages and finalised the new keys, traffic transitions from the old session to the new one. -Cutover is signalled in-band by the **K-bit** in the FMP flags byte. Each -side starts emitting frames under the new session with K set; on receipt -of the first K-marked frame the peer accepts the cutover and follows -suit. A new pair of session indices is allocated as part of the new +Cutover is signalled in-band by the **K-bit** in the FMP flags byte. The +initiator switches first, on its own schedule once it has read msg2, and +marks its frames under the new session with K; the responder follows on +the first K-marked frame that authenticates against its new session. A +responder drops a pending session the initiator never adopts after a hold +(the drain ceiling, 120 s at stock settings), so the next rekey can +proceed. A new pair of session indices is allocated as part of the new session, replacing the old indices on subsequent frames (see [Index Properties](#index-properties)). diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index 18d5e09b..a20454c4 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -716,9 +716,11 @@ impl Node { } } - // Store pending session on the existing peer. + // Store the new session as the responder's pending session. It + // is promoted by the initiator's first new-epoch frame, not by + // our own tick. if let Some(existing) = self.peers.get_mut(&peer) { - existing.set_pending_session(noise_session, our_new_index, wire.their_index); + existing.answer_rekey(noise_session, our_new_index, wire.their_index); existing.record_peer_rekey(); } @@ -1241,14 +1243,16 @@ impl Node { } // Nothing authenticated this msg2 before the read, and the // index it names travels in cleartext in our msg1, so it - // may be a forgery. The responder committed its new - // session when it answered that msg1 and cuts over on its - // own tick, so abandoning here would leave the two ends on - // different keys. The read rolled the handshake back: - // keep the cycle and its dispatch entry so the genuine - // msg2 can still complete it. If no readable msg2 ever - // arrives, the msg1 resend budget abandons the cycle as - // it would for a lost one. + // may be a forgery. The responder holds the session it + // answered with until our first new-epoch frame reaches + // it, so giving up the cycle here would throw away a + // cycle the genuine msg2 can still complete, and the + // responder would then hold that session until its + // retirement hold passes. The read rolled the handshake + // back: keep the cycle and its dispatch entry so the + // genuine msg2 can still complete it. If no readable msg2 + // ever arrives, the msg1 resend budget abandons the cycle + // as it would for a lost one. Err(e) if peer.awaits_msg2() => { debug!( peer = %display_name, diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index 1bdeb91b..a2fe719c 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -4,6 +4,7 @@ //! 1. Rekey trigger (time elapsed or send counter exceeded) //! 2. Drain window expiry (clean up previous session after cutover) //! 3. Initiator-side cutover (first send after handshake completion) +//! 4. Retirement of a responder's pending session the initiator never adopted use crate::NodeAddr; use crate::node::Node; @@ -59,24 +60,60 @@ const DRAIN_MAX_RETENTION_SECS: u64 = crate::proto::fsp::limits::DRAIN_WINDOW_SE /// the configured handshake timers imply, so the ceiling always clears /// the recovery it is supposed to leave room for. pub(in crate::node) fn drain_max_retention_ms(rate_limit: &crate::config::RateLimitConfig) -> u64 { - let mut ladder_ms: u64 = 0; - let mut interval = rate_limit.handshake_resend_interval_ms as f64; - for _ in 0..rate_limit.handshake_max_resends { - ladder_ms = ladder_ms.saturating_add(interval as u64); - interval *= rate_limit.handshake_resend_backoff; - } - let recovery_budget_ms = ladder_ms + let recovery_budget_ms = ladder_ms(rate_limit) .saturating_add(rate_limit.handshake_timeout_secs.saturating_mul(1000)) .saturating_add(crate::proto::fsp::limits::REKEY_DAMPENING_SECS * 1000); (DRAIN_MAX_RETENTION_SECS * 1000).max(recovery_budget_ms) } +/// Sum of the handshake resend intervals, in milliseconds: the initial +/// interval, multiplied by the backoff factor after each resend, over the +/// configured number of resends. +fn ladder_ms(rate_limit: &crate::config::RateLimitConfig) -> u64 { + let mut total: u64 = 0; + let mut interval = rate_limit.handshake_resend_interval_ms as f64; + for _ in 0..rate_limit.handshake_max_resends { + total = total.saturating_add(interval as u64); + interval *= rate_limit.handshake_resend_backoff; + } + total +} + +/// How long a rekey responder holds a pending session its initiator has +/// not adopted before retiring it. +/// +/// The initiator reads a msg2 only while it holds its handshake, so the +/// msg1 resend ladder, one tick to the abandon and one tick to the +/// cutover bound the in-flight age of a msg2 it can still adopt. That +/// assumes its sends succeed: a resend that fails to send does not +/// advance the ladder. After the cutover the hold must also outlast the +/// longer of the two ways the cutover reaches this node, or fails to: +/// the first frame on the new epoch on an idle link, up to one heartbeat +/// interval later, and the link-dead reap if every frame from the +/// initiator is lost. Both are allowed one more tick. Retiring before +/// either turns a late but legitimate adoption into a split. That floor +/// is 64 s at stock settings; the drain ceiling, which bounds residence +/// for the same recovering-peer reason, is 120 s and is used unless the +/// configured timers push the floor above it. +pub(in crate::node) fn pending_hold(node: &crate::config::NodeConfig) -> std::time::Duration { + let after_ms = node + .heartbeat_interval_secs + .max(node.link_dead_timeout_secs) + .saturating_mul(1000); + let floor_ms = ladder_ms(&node.rate_limit) + .saturating_add(node.tick_interval_secs.saturating_mul(3000)) + .saturating_add(after_ms); + std::time::Duration::from_millis(drain_max_retention_ms(&node.rate_limit).max(floor_ms)) +} + impl Node { /// Periodic rekey check. Called from the tick loop. /// /// For each active peer with a session: /// - If the initiator has a pending session, perform K-bit cutover /// - If the drain window has expired, clean up the previous session + /// - If a responder's pending session was never adopted by the initiator + /// within the hold, retire it /// - If the rekey timer/counter fires, initiate a new handshake pub(in crate::node) async fn check_rekey(&mut self) { if !self.config().node.rekey.enabled { @@ -122,6 +159,12 @@ impl Node { self.initiate_rekey(&node_addr).await; self.observe_rekey_initiated(&node_addr); } + // Retire a responder pending the initiator never adopted. Stays + // inline: the peer machine models only the initiator's pending + // (`on_rekey_msg2`), so there is nothing for it to consume. + ConnAction::RetirePending { peer: node_addr } => { + self.retire_unadopted(&node_addr); + } #[allow(unreachable_patterns)] _ => {} } @@ -250,6 +293,28 @@ impl Node { let _ = did_cutover; } + /// Retire `node_addr`'s responder-held pending session: unregister its + /// index and free it. The index was registered in `peers_by_index` when + /// the msg2 was sent and was never given to the decrypt worker, which + /// sees a session only once it is current. + fn retire_unadopted(&mut self, node_addr: &NodeAddr) { + let retired = self + .peers + .get_mut(node_addr) + .and_then(|peer| peer.retire_pending().map(|idx| (idx, peer.transport_id()))); + if let Some((idx, transport_id)) = retired { + if let Some(tid) = transport_id { + self.peers_by_index.remove(&(tid, idx.as_u32())); + } + let _ = self.index_allocator.free(idx); + debug!( + peer = %self.peer_display_name(node_addr), + index = %idx, + "Rekey pending session retired: the initiator never adopted it" + ); + } + } + /// Pre-refactor drain-completion body, retained as the release fallback for /// the (should-be-impossible) missing-machine case. Byte-identical to the old /// inline `ConnAction::Drain` arm and to the executor's `CompleteDrain` arm. @@ -300,6 +365,8 @@ impl Node { .map(|s| s.current_send_counter()) .unwrap_or(0), jitter_secs: peer.rekey_jitter_secs(), + pending_role: peer.pending_role(), + pending_expired: peer.pending_expired(pending_hold(&self.config().node)), }) .collect() } @@ -642,6 +709,9 @@ impl Node { /// out, abandon it (the handshake only — a completed rekey session /// is never discarded on a timer, see /// [`FspAction::AbandonHandshake`]) + /// - If a handshake this node initiated got no SessionAck within the + /// handshake timeout, drop it so the trigger can start a fresh one + /// (see [`FspAction::ExpireInitiation`]) /// - If the rekey timer/counter fires, initiate a new XK handshake /// (this last one only when `node.rekey.enabled`) /// @@ -707,6 +777,24 @@ impl Node { ); } } + FspAction::ExpireInitiation { addr } => { + // The handshake only: the trigger starts a fresh rekey + // on a later tick. + let age_ms = self + .sessions + .get(&addr) + .map(|entry| now_ms.saturating_sub(entry.initiated_ms())) + .unwrap_or(0); + if let Some(entry) = self.sessions.get_mut(&addr) { + entry.abandon_handshake(); + self.stats_mut().session.rekey_unanswered += 1; + info!( + peer = %self.peer_display_name(&addr), + age_ms, + "FSP rekey we initiated got no answer within the handshake timeout, retrying, session retained" + ); + } + } FspAction::InitiateRekey { addr } => { self.initiate_session_rekey(&addr).await; } @@ -726,6 +814,14 @@ impl Node { // finished, anchored on the peer's last accepted setup message, which // is the only stamp that path writes. A *completed* rekey has no such // bound and must not acquire one: see `FspAction::AbandonHandshake`. + // + // A handshake this node initiated has its own deadline, measured from + // the setup this node sent, and neither clock may stand in for the + // other: the peer's stamp says nothing about our setup, and ours + // says nothing about the peer's. Unlike the peer predicate, ours has + // no `!= 0` conjunct, deliberately: an initiator handshake armed + // without the stamp reads as expired, because dropping one costs a + // retry and keeping one stops rotation. let stale_handshake_ms = self.config().node.rate_limit.handshake_timeout_secs * 1000; // Absolute ceiling on `previous`-slot retention, measured from the // cutover. The sliding drain deadline is peer-progress-aware, so an @@ -749,6 +845,8 @@ impl Node { is_dampened: entry.is_rekey_dampened(now_ms, dampening_ms), armed_handshake_expired: entry.last_peer_rekey_ms() != 0 && now_ms.saturating_sub(entry.last_peer_rekey_ms()) > stale_handshake_ms, + initiation_expired: now_ms.saturating_sub(entry.initiated_ms()) + > stale_handshake_ms, elapsed_secs: now_ms.saturating_sub(entry.session_start_ms()) / 1000, counter: entry.send_counter(), jitter_secs: entry.rekey_jitter_secs(), @@ -759,7 +857,8 @@ impl Node { /// Initiate an FSP session rekey. /// /// Creates a new XK handshake as initiator, sends SessionSetup msg1 - /// through the mesh, and stores the handshake state on the existing entry. + /// through the mesh, and stores the handshake state on the existing entry, + /// stamping the deadline by which a SessionAck must complete it. async fn initiate_session_rekey(&mut self, dest_addr: &NodeAddr) { // Check route availability before paying crypto cost if self.find_next_hop(dest_addr).is_none() { @@ -826,9 +925,10 @@ impl Node { return; } - // Store rekey state on the existing session entry + // Store rekey state on the existing session entry. The deadline is + // stamped only now, so it runs from the setup actually on the wire. if let Some(entry) = self.sessions.get_mut(dest_addr) { - entry.set_rekey_state(handshake, true); + entry.begin_rekey(handshake, Self::now_ms()); } debug!( diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 25ac5b4d..15aad913 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -335,6 +335,12 @@ impl Node { } }; + // A frame that authenticates on this session, in any epoch slot, + // proves the peer completed the handshake, so a msg3 still held for + // resend has arrived. Only an established initiator holds one; for + // every other entry this is a no-op. + entry.clear_handshake_payload(); + // React to the epoch the frame decrypted against. The shell opened // the frame; the core classifies the post-decrypt reaction over the // plain-data slot + session flags, and the shell applies the @@ -893,13 +899,15 @@ impl Node { // 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 three sites become instances of the epoch - // discard the responder arm was fixed for. + // are the same action here. Arming as initiator has one caller, + // `initiate_session_rekey` through `begin_rekey`, and + // `set_rekey_state(_, true)` remains only at the restore below; + // `check_session_rekey` never reaches `initiate_session_rekey` 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 three + // 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, @@ -917,7 +925,9 @@ impl Node { // back to its pre-read state so it can still read the genuine // ack, and the refusal is counted. The rollback matters because // `read_xk_message_2` mixes the sender's ephemeral in before it - // authenticates. + // authenticates. The restore does not restamp the deadline, which + // runs from the setup this node sent, so an unreadable ack cannot + // hold the rekey open. if let Err(e) = handshake.try_read_xk_message_2(&ack.handshake_payload) { debug!(error = %e, "Failed to process rekey XK msg2, keeping the rekey"); entry.set_rekey_state(handshake, true); @@ -1044,7 +1054,7 @@ impl Node { let msg3_wire = SessionMsg3::new(msg3); let msg3_payload = msg3_wire.encode(); let my_addr = *self.node_addr(); - let mut datagram = SessionDatagram::new(my_addr, *src_addr, msg3_payload) + let mut datagram = SessionDatagram::new(my_addr, *src_addr, msg3_payload.clone()) .with_ttl(self.config().node.session.default_ttl); if let Err(e) = self.send_session_datagram(&mut datagram).await { @@ -1062,11 +1072,19 @@ impl Node { }; let now_ms = Self::now_ms(); + let resend_interval = self.config().node.rate_limit.handshake_resend_interval_ms; entry.set_state(EndToEndState::Established(session)); entry.set_coords_warmup_remaining(self.config().node.session.coords_warmup_packets); entry.mark_established(now_ms); entry.init_mmp(&self.config().node.session_mmp); - entry.clear_handshake_payload(); + // Keep msg3 for resend. This end is established once msg3 leaves, the + // responder only once it arrives, and nothing else repairs a lost + // msg3: the responder's resent SessionAck lands on the not-initiating + // arm above and is refused. `resend_pending_session_handshakes` + // resends it until a frame from the peer authenticates on this + // session or the resend budget is spent. The rekey arm keeps its + // msg3 for the same reason. + entry.set_handshake_payload(msg3_payload, now_ms + resend_interval); entry.touch(now_ms); self.sessions.insert(*src_addr, entry); self.insert_coord_hint(*src_addr, ack.src_coords.clone(), now_ms); diff --git a/src/node/handlers/timeout.rs b/src/node/handlers/timeout.rs index 7f62f9ac..aa833b2f 100644 --- a/src/node/handlers/timeout.rs +++ b/src/node/handlers/timeout.rs @@ -6,6 +6,7 @@ use crate::peer::machine::TimerKind; use crate::proto::fmp::{ ConnAction, ConnSnapshot, LifecycleView, PeerSnapshot, RekeyResendSnapshot, }; +use crate::proto::fsp::{FspAction, InitialMsg3ResendSnapshot}; use crate::transport::LinkId; use tracing::{debug, info, warn}; @@ -403,6 +404,10 @@ impl Node { /// - If the handshake has exceeded the timeout window, remove the session. /// - If a resend is due and under max resends, resend the stored payload /// wrapped in a fresh SessionDatagram (so routing can adapt). + /// + /// For an established initiator still holding its msg3, resend it until + /// the peer is heard from or the budget is spent (see + /// `resend_initial_msg3`). Established sessions are never removed here. pub(in crate::node) async fn resend_pending_session_handshakes(&mut self, now_ms: u64) { if self.sessions.is_empty() { return; @@ -474,6 +479,96 @@ impl Node { ); } } + + self.resend_initial_msg3(now_ms).await; + } + + /// Resend an established initiator's retained msg3 until the responder is + /// heard from, and stop retaining it once the resend budget is spent. + /// + /// The initiator is established the moment msg3 leaves; the responder only + /// once it arrives. Nothing else repairs a lost msg3: the responder's + /// resent SessionAck reaches an entry that is no longer initiating and is + /// refused. The payload is released by the first inbound frame that + /// authenticates on the session (`handle_encrypted_session_msg`) or, here, + /// when the budget is spent. Runs whether or not periodic rekey is enabled. + async fn resend_initial_msg3(&mut self, now_ms: u64) { + use crate::proto::link::SessionDatagram; + + let candidates = self.initial_msg3_resend_snapshots(now_ms); + if candidates.is_empty() { + return; + } + let max_resends = self.config().node.rate_limit.handshake_max_resends; + let interval_ms = self.config().node.rate_limit.handshake_resend_interval_ms; + let backoff = self.config().node.rate_limit.handshake_resend_backoff; + let ttl = self.config().node.session.default_ttl; + let my_addr = *self.node_addr(); + + for action in self.fsp.poll_initial_msg3_resends(candidates, max_resends) { + match action { + FspAction::ReleaseInitialMsg3 { addr } => { + if let Some(entry) = self.sessions.get_mut(&addr) { + entry.clear_handshake_payload(); + } + info!( + dest = %self.peer_display_name(&addr), + "Session msg3 unconfirmed after max resends, no longer resending" + ); + } + FspAction::ResendInitialMsg3 { addr } => { + let payload = match self.sessions.get(&addr).and_then(|e| e.handshake_payload()) + { + Some(p) => p.to_vec(), + None => continue, + }; + let mut datagram = SessionDatagram::new(my_addr, addr, payload).with_ttl(ttl); + let sent = match self.send_session_datagram(&mut datagram).await { + Ok(_) => true, + Err(e) => { + debug!( + dest = %self.peer_display_name(&addr), + error = %e, + "Session msg3 resend failed" + ); + false + } + }; + if sent && let Some(entry) = self.sessions.get_mut(&addr) { + let count = entry.resend_count() + 1; + let next = + now_ms + (interval_ms as f64 * backoff.powi(count as i32)) as u64; + entry.record_resend(next); + debug!( + dest = %self.peer_display_name(&addr), + resend = count, + "Resent session msg3" + ); + } + } + #[allow(unreachable_patterns)] + _ => {} + } + } + } + + /// Snapshot every established session still retaining its initial msg3, + /// pre-evaluating the resend-due predicate against `now_ms` so the core + /// reads no clock. + /// + /// `is_established()` partitions the shared handshake resend slot: a + /// non-established entry's SessionSetup or SessionAck belongs to the passes + /// above, an established entry's msg3 to this one. + fn initial_msg3_resend_snapshots(&self, now_ms: u64) -> Vec { + self.sessions + .iter() + .filter(|(_, entry)| entry.is_established() && entry.handshake_payload().is_some()) + .map(|(addr, entry)| InitialMsg3ResendSnapshot { + addr: *addr, + resend_count: entry.resend_count(), + resend_due: entry.next_resend_at_ms() != 0 && now_ms >= entry.next_resend_at_ms(), + }) + .collect() } /// Remove established sessions that have been idle too long. diff --git a/src/node/session/mod.rs b/src/node/session/mod.rs index 635c81f9..18c76b4c 100644 --- a/src/node/session/mod.rs +++ b/src/node/session/mod.rs @@ -123,8 +123,11 @@ pub(crate) struct SessionEntry { bytes_recv: u64, // === Handshake Resend === - /// Encoded session-layer payload for resend (SessionSetup or SessionAck). - /// Cleared on Established transition. + /// The last initial-handshake message this side sent and has no proof the + /// peer received: SessionSetup while initiating, SessionAck while awaiting + /// msg3, and on an established initiator its SessionMsg3 until an inbound + /// frame authenticates on the session or the resend budget is spent. The + /// state says which one it is, and which sweep resends it. handshake_payload: Option>, /// Number of resends performed. resend_count: u32, @@ -160,6 +163,14 @@ pub(crate) struct SessionEntry { rekey_initiator: bool, /// Dampening: last time peer sent us a rekey msg1 (Unix ms). last_peer_rekey_ms: u64, + /// When this node armed its current rekey handshake as initiator (Unix + /// ms), which bounds how long that handshake waits for a SessionAck. + /// + /// Written only by `begin_rekey`, so putting the handshake back after + /// an unreadable SessionAck does not move it. Read only while a + /// handshake is armed and `rekey_initiator` is set, which is why it is + /// never cleared. + initiated_ms: u64, /// When this side's FSP rekey handshake completed and produced the /// `pending` session (Unix ms): the initiator sending msg3, or the /// responder accepting it. Cleared on cutover. @@ -229,6 +240,7 @@ impl SessionEntry { pending_new_session: None, rekey_initiator: false, last_peer_rekey_ms: 0, + initiated_ms: 0, rekey_completed_ms: 0, rekey_msg3_payload: None, rekey_msg3_next_resend_ms: 0, @@ -408,6 +420,8 @@ impl SessionEntry { /// /// For initiators, this is the SessionSetup payload bytes. /// For responders, this is the SessionAck payload bytes. + /// For an initiator that has just become established, this is its + /// SessionMsg3 payload bytes, held until the responder is heard from. /// The payload is re-wrapped in a fresh SessionDatagram on each resend /// so routing can adapt to topology changes. pub(crate) fn set_handshake_payload(&mut self, payload: Vec, next_resend_at_ms: u64) { @@ -421,7 +435,8 @@ impl SessionEntry { self.handshake_payload.as_deref() } - /// Clear the stored handshake payload (called on Established transition). + /// Clear the stored handshake payload (an inbound frame authenticated on + /// the session, or the msg3 resend budget is spent). pub(crate) fn clear_handshake_payload(&mut self) { self.handshake_payload = None; self.next_resend_at_ms = 0; @@ -481,6 +496,15 @@ impl SessionEntry { self.last_peer_rekey_ms } + /// When this node armed its current rekey handshake as initiator, or 0 + /// if it never has. + /// + /// Bounds the age of a handshake this node armed, independently of any + /// rekey the peer started. + pub(crate) fn initiated_ms(&self) -> u64 { + self.initiated_ms + } + /// Record that the peer initiated a rekey (for dampening). pub(crate) fn record_peer_rekey(&mut self, now_ms: u64) { self.last_peer_rekey_ms = now_ms; @@ -674,6 +698,18 @@ impl SessionEntry { self.rekey_initiator = is_initiator; } + /// Arm a rekey handshake this node initiated, and stamp its deadline. + /// + /// The only writer of `initiated_ms`. The restore after an unreadable + /// SessionAck goes through `set_rekey_state` instead and must stay + /// there: a restamping restore would let a stream of forged acks hold + /// the rekey open for good. + pub(crate) fn begin_rekey(&mut self, state: HandshakeState, now_ms: u64) { + self.rekey_state = Some(state); + self.rekey_initiator = true; + self.initiated_ms = now_ms; + } + /// Take the rekey state for processing. pub(crate) fn take_rekey_state(&mut self) -> Option { self.rekey_state.take() @@ -825,10 +861,11 @@ impl SessionEntry { /// Abandon an in-progress rekey handshake, keeping any completed /// session already waiting for cut-over. /// - /// Used when a handshake the peer armed times out without its msg3. - /// The handshake holds no key material either endpoint can be using, - /// so dropping it costs nothing; a `pending` session alongside it is - /// the epoch the peer may already have moved to and must survive. + /// Used when an armed handshake times out: one the peer armed without + /// its msg3, or one this node armed without a SessionAck. The handshake + /// holds no key material either endpoint can be using, so dropping it + /// costs nothing; a `pending` session alongside it is the epoch the + /// peer may already have moved to and must survive. /// /// `rekey_initiator` is deliberately left alone: it describes whichever /// rekey artefact the entry still holds, and every reader gates on a diff --git a/src/node/stats.rs b/src/node/stats.rs index 0d82f62f..a638f131 100644 --- a/src/node/stats.rs +++ b/src/node/stats.rs @@ -64,6 +64,12 @@ pub struct SessionStats { /// the handshake timeout without a msg3 and was discarded. The /// established session is retained. pub rekey_expired: u64, + /// A rekey this node initiated got no readable SessionAck within the + /// handshake timeout, and its handshake was discarded so the trigger + /// can retry. The established session is retained. A sustained rate + /// means setups or acks to that peer are being lost, or the peer holds + /// a stuck handshake of its own and wins the tie-break. + pub rekey_unanswered: u64, /// A completed rekey session still waiting for the peer's cut-over was /// replaced by a newer one, completed from a msg3 carrying the same /// authenticated peer key. This drops key material the peer may @@ -101,6 +107,7 @@ impl SessionStats { rekey_yielded: self.rekey_yielded, rekey_pending: self.rekey_pending, rekey_expired: self.rekey_expired, + rekey_unanswered: self.rekey_unanswered, pending_replaced: self.pending_replaced, ack_handshake_failed: self.ack_handshake_failed, setup_rate_limited: self.setup_rate_limited, @@ -406,6 +413,7 @@ pub struct SessionStatsSnapshot { pub rekey_yielded: u64, pub rekey_pending: u64, pub rekey_expired: u64, + pub rekey_unanswered: u64, pub pending_replaced: u64, pub ack_handshake_failed: u64, pub setup_rate_limited: u64, diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index 78524156..f56f1d90 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -1395,6 +1395,20 @@ async fn pump_until_quiet(nodes: &mut [TestNode]) { } } +/// Drive node 0's rekey msg1 resend ladder on the synthetic clock from +/// `base_ms`: at the default 1 s interval and 2x backoff, resends at +1, +3, +/// +7, +15 and +31 s, then the abandon past the budget at +63 s, pumping +/// after each. +async fn walk_ladder(nodes: &mut [TestNode], base_ms: u64) { + for offset_s in [1u64, 3, 7, 15, 31, 63] { + nodes[0] + .node + .resend_pending_rekeys(base_ms + offset_s * 1000) + .await; + pump_until_quiet(nodes).await; + } +} + /// A forged rekey msg2 that carries the initiator's live rekey index must not /// split the link. /// @@ -1402,11 +1416,10 @@ async fn pump_until_quiet(nodes: &mut [TestNode]) { /// cleartext rekey index from it, and delivers a msg2 of the right size under /// that index ahead of the responder's real reply. Under IK the forgery cannot /// authenticate: only the responder's static key produces a msg2 the initiator -/// can read. The IK responder has already committed its new session when it -/// answered msg1 and cuts over on its own next rekey tick, so the initiator has -/// to keep the cycle through the forgery and complete it on the real msg2. An -/// initiator that gives the cycle up instead holds no session matching the one -/// the responder now sends on. +/// can read. The IK responder committed its new session as pending when it +/// answered msg1 and promotes it on the initiator's first new-epoch frame, so +/// the initiator must complete the cycle on the real msg2 for either direction +/// to survive. /// /// Deterministic, no wall-clock wait: both sessions are backdated past node 0's /// time trigger and node 1's rekey-acceptance floor, node 1 never initiates, @@ -1451,14 +1464,6 @@ async fn forged_rekey_msg2_does_not_split_the_link() { nodes[1].node.check_rekey().await; pump_until_quiet(&mut nodes).await; - // node 1 cut over to the session it committed at msg1. Without this the - // delivery checks below could pass because no rekey happened at all. - assert_ne!( - nodes[1].node.get_peer(&node0_addr).unwrap().our_index(), - node1_idx_before, - "node 1 must have cut over to its new session" - ); - let post_fwd = build_ipv6_packet(&fips0, &fips1, b"post-rekey 0 to 1"); let post_rev = build_ipv6_packet(&fips1, &fips0, b"post-rekey 1 to 0"); nodes[0].node.handle_tun_outbound(post_fwd.clone()).await; @@ -1472,12 +1477,21 @@ async fn forged_rekey_msg2_does_not_split_the_link() { "node 0 to node 1 must decode after the forged msg2" ); - // This is the assertion that tells the two outcomes apart; keep it. node 0 - // to node 1 passes either way inside this test, because node 1 keeps its - // previous session through the drain window and still decrypts node 0's - // old-session frames. node 1 to node 0 fails exactly when node 0 lost the - // cycle to the forgery: node 1 now sends on its new session, addressed to - // node 0's rekey index, and node 0 has no session registered under it. + // node 1 promotes its pending session when node 0's first frame on the new + // epoch authenticates against it. This is one of the two assertions that + // tell the outcomes apart (see below), and it also rules out a pass in + // which no rekey happened at all. + assert_ne!( + nodes[1].node.get_peer(&node0_addr).unwrap().our_index(), + node1_idx_before, + "node 1 must have promoted its pending session on node 0's first new-epoch frame" + ); + + // node 1 to node 0 must still decode, but it no longer tells the outcomes + // apart: if node 0 lost the cycle to the forgery, node 1 never promotes and + // both nodes stay on their original sessions, so this passes either way. + // The promotion assertion above and node 0's cutover assertion below are + // the ones that catch a lost cycle; keep both. let handshake = &nodes[0].node.stats().handshake; let (bad_state, unknown) = (handshake.bad_state, handshake.unknown_connection); let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); @@ -1505,26 +1519,23 @@ async fn forged_rekey_msg2_does_not_split_the_link() { /// A rekey msg2 lost in transit, with no attacker present, must not split the /// link. /// -/// node 1 answers node 0's rekey msg1, commits its new session at once, and -/// cuts over on its own next rekey tick. node 1's msg2 never arrives. node 0 -/// walks its whole msg1 resend ladder, which node 1 now meets holding a young -/// post-cutover session, then abandons the cycle at the budget, and its time -/// cadence fires again. The pair must end up able to carry data both ways. +/// node 1 answers node 0's rekey msg1 and holds its new session as pending. +/// node 1's msg2 never arrives, so node 0 never shows it holds the new keys, +/// and node 1's own rekey tick must not cut over to them. node 0 walks its +/// whole msg1 resend ladder and re-fires after the abandon; node 1 refuses +/// each of those msg1s while it holds the pending. Data must flow both ways on +/// the original sessions throughout. The assertion that tells the outcomes +/// apart is node 1 to node 0: a node 1 that cut over would seal to the rekey +/// index node 0 abandoned. /// /// What this does not model: the tick loop and the link-dead reap are not /// driven, only the rekey tick and the resend function, each called by hand. /// The resend ladder runs on a synthetic millisecond clock, while node 1's -/// session ages are real `Instant`s, so every resent msg1 lands inside node -/// 1's 30 s post-cutover window, as it would in real time for the first 30 s. -/// -/// It fails today at the node 1 to node 0 delivery: node 1 sends on the session -/// it cut over to, addressed to the rekey index node 0 abandoned, and node 0's -/// re-fired rekey msg1 is answered as a duplicate with a msg2 node 0 no longer -/// has a dispatch entry for. Ignored until the responder stops committing ahead -/// of the initiator. +/// pending install time is a real `Instant`, so node 1 is still inside its +/// hold when the ladder ends. The retirement of the pending and the rekey that +/// completes after it are covered by +/// `a_responder_retires_an_unadopted_rekey_and_the_next_rekey_completes`. #[tokio::test] -#[ignore = "known defect: a rekey msg2 lost in transit splits the link, because the \ - responder cuts over to keys the initiator never installs"] async fn dropped_rekey_msg2_does_not_split_the_link() { let HeldMsg2Pair { mut nodes, @@ -1543,25 +1554,22 @@ async fn dropped_rekey_msg2_does_not_split_the_link() { // this drop is of the real reply. drop(held_msg2); - // node 1 holds a pending session and no rekey in progress, so its tick - // cuts over to it. + // node 1 answered the rekey, so its tick must not commit to the new keys: + // node 0 has not shown it holds them. nodes[1].node.check_rekey().await; - assert_ne!( - nodes[1].node.get_peer(&node0_addr).unwrap().our_index(), + let node1_peer = nodes[1].node.get_peer(&node0_addr).unwrap(); + assert_eq!( + node1_peer.our_index(), node1_idx_before, - "node 1 must have cut over to the session it committed at msg1" + "node 1 must not cut over to a session node 0 has not adopted" ); + assert!( + node1_peer.pending_new_session().is_some(), + "node 1 must still hold the session it answered with" + ); + let node1_rejects_before = nodes[1].node.stats().handshake.bad_state; - // node 0's msg1 resend ladder at the default 1 s interval and 2x backoff: - // resends at +1, +3, +7, +15 and +31 s, then the abandon past the budget. - let base_ms = Node::now_ms(); - for offset_s in [1u64, 3, 7, 15, 31, 63] { - nodes[0] - .node - .resend_pending_rekeys(base_ms + offset_s * 1000) - .await; - pump_until_quiet(&mut nodes).await; - } + walk_ladder(&mut nodes, Node::now_ms()).await; assert!( !nodes[0] .node @@ -1578,6 +1586,11 @@ async fn dropped_rekey_msg2_does_not_split_the_link() { nodes[0].node.check_rekey().await; nodes[1].node.check_rekey().await; pump_until_quiet(&mut nodes).await; + assert_eq!( + nodes[1].node.stats().handshake.bad_state - node1_rejects_before, + 6, + "node 1 must refuse node 0's five resends and its re-fired msg1 while it holds the pending" + ); let post_fwd = build_ipv6_packet(&fips0, &fips1, b"post-loss 0 to 1"); let post_rev = build_ipv6_packet(&fips1, &fips0, b"post-loss 1 to 0"); @@ -1592,9 +1605,10 @@ async fn dropped_rekey_msg2_does_not_split_the_link() { "node 0 to node 1 must decode after the lost msg2" ); - // The discriminating assertion, as in the forged-msg2 test: node 1 sends - // on the session it cut over to, addressed to node 0's rekey index, and - // node 0 decodes it only if it holds a session registered under it. + // The discriminating assertion. node 1 still seals on its original session, + // to node 0's original index, which stays registered because node 0 never + // cut over. It decodes only because node 1 held back: had node 1 cut over + // on its own tick, it would seal to the rekey index node 0 abandoned. let handshake = &nodes[0].node.stats().handshake; let (bad_state, unknown) = (handshake.bad_state, handshake.unknown_connection); let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); @@ -1608,6 +1622,186 @@ async fn dropped_rekey_msg2_does_not_split_the_link() { cleanup_nodes(&mut nodes).await; } +/// A rekey responder whose msg2 was lost holds the pending session it +/// answered with until the hold passes, then retires it: the pending slot and +/// its role are emptied, its index is unregistered and freed, and the current +/// session is untouched. The initiator's next msg1 is then answered and the +/// rekey completes, with node 1 promoted by node 0's first new-epoch frame and +/// data flowing both ways. +#[tokio::test] +async fn a_responder_retires_an_unadopted_rekey_and_the_next_rekey_completes() { + use crate::node::handlers::rekey::pending_hold; + + let HeldMsg2Pair { + mut nodes, + node0_addr, + node1_addr, + fips0, + fips1, + tun0_rx, + tun1_rx, + node0_idx_before, + node1_idx_before, + held_msg2, + .. + } = rekey_pair_with_held_msg2().await; + + // 1. The msg2 is lost; node 1 holds the pending it answered with. + drop(held_msg2); + nodes[1].node.check_rekey().await; + let node1_peer = nodes[1].node.get_peer(&node0_addr).unwrap(); + assert_eq!( + node1_peer.our_index(), + node1_idx_before, + "node 1 must not cut over to a session node 0 has not adopted" + ); + let pending_idx = node1_peer + .pending_our_index() + .expect("node 1 must still hold the session it answered with"); + + // 2. node 0 walks its ladder to the abandon and re-fires; node 1 refuses + // the re-fired msg1 while it holds the pending. + walk_ladder(&mut nodes, Node::now_ms()).await; + nodes[0].node.check_rekey().await; + pump_until_quiet(&mut nodes).await; + assert!( + nodes[0] + .node + .get_peer(&node1_addr) + .unwrap() + .rekey_in_progress(), + "node 0 must have re-fired its rekey after the abandon" + ); + assert!( + nodes[1] + .node + .get_peer(&node0_addr) + .unwrap() + .pending_new_session() + .is_some(), + "node 1 must still hold its pending after refusing the re-fired msg1" + ); + + // 3. The hold passes. + let hold = pending_hold(&nodes[1].node.config().node); + nodes[1] + .node + .get_peer_mut(&node0_addr) + .unwrap() + .backdate_pending(hold + Duration::from_secs(1)); + + // 4. node 1's tick retires the pending. + nodes[1].node.check_rekey().await; + let node1_peer = nodes[1].node.get_peer(&node0_addr).unwrap(); + assert!( + node1_peer.pending_new_session().is_none(), + "node 1 must have retired the unadopted pending session" + ); + assert_eq!(node1_peer.pending_role(), None); + assert_eq!( + node1_peer.our_index(), + node1_idx_before, + "retirement must leave node 1's current session alone" + ); + assert!( + !nodes[1] + .node + .peers_by_index + .contains_key(&(nodes[1].transport_id, pending_idx.as_u32())), + "the retired pending index must be unregistered" + ); + assert!( + !nodes[1].node.index_allocator.is_allocated(pending_idx), + "the retired pending index must be freed" + ); + + // 5. node 0's next resend of the re-fired msg1 is answered now, and node 0 + // reads the msg2. + nodes[0] + .node + .resend_pending_rekeys(Node::now_ms() + 1_000) + .await; + pump_until_quiet(&mut nodes).await; + assert!( + nodes[0] + .node + .get_peer(&node1_addr) + .unwrap() + .pending_new_session() + .is_some(), + "node 0 must have completed the re-fired rekey on node 1's answer" + ); + + // 6. node 0 cuts over; node 1 holds until node 0's first new-epoch frame. + nodes[0].node.check_rekey().await; + nodes[1].node.check_rekey().await; + pump_until_quiet(&mut nodes).await; + + let post_fwd = build_ipv6_packet(&fips0, &fips1, b"post-retire 0 to 1"); + let post_rev = build_ipv6_packet(&fips1, &fips0, b"post-retire 1 to 0"); + nodes[0].node.handle_tun_outbound(post_fwd.clone()).await; + pump_until_quiet(&mut nodes).await; + nodes[1].node.handle_tun_outbound(post_rev.clone()).await; + pump_until_quiet(&mut nodes).await; + + let got: Vec> = std::iter::from_fn(|| tun1_rx.try_recv().ok()).collect(); + assert_eq!( + got, + vec![post_fwd], + "node 0 to node 1 must decode after the retry" + ); + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!( + got, + vec![post_rev], + "node 1 to node 0 must decode after the retry" + ); + assert_ne!( + nodes[0].node.get_peer(&node1_addr).unwrap().our_index(), + node0_idx_before, + "node 0 must have cut over on the retried rekey" + ); + assert_ne!( + nodes[1].node.get_peer(&node0_addr).unwrap().our_index(), + node1_idx_before, + "node 1 must have promoted on node 0's first new-epoch frame" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// The responder hold is the drain ceiling at stock settings, and a raised +/// link-dead timeout or heartbeat interval raises it past that ceiling, taking +/// the larger of the two rather than their sum. +#[test] +fn the_responder_hold_is_the_drain_ceiling_at_stock_settings_and_outlasts_a_raised_link_dead_timeout() + { + use crate::node::handlers::rekey::{drain_max_retention_ms, pending_hold}; + + let stock = crate::config::NodeConfig::default(); + assert_eq!(pending_hold(&stock), Duration::from_secs(120)); + assert_eq!( + pending_hold(&stock), + Duration::from_millis(drain_max_retention_ms(&stock.rate_limit)) + ); + + // 31 s of msg1 ladder (1+2+4+8+16), three 1 s ticks, and 200 s of + // link-dead timeout: 234 s. + let raised = crate::config::NodeConfig { + link_dead_timeout_secs: 200, + ..Default::default() + }; + assert_eq!(pending_hold(&raised), Duration::from_secs(234)); + + // The heartbeat term raised instead, link-dead at its default: the floor + // takes the larger of the two, so 234 s again, not 264 s. + let raised = crate::config::NodeConfig { + heartbeat_interval_secs: 200, + ..Default::default() + }; + assert_eq!(pending_hold(&raised), Duration::from_secs(234)); +} + #[tokio::test] async fn test_tun_outbound_triggers_session_initiation() { // Two connected nodes, no session yet. @@ -5090,6 +5284,302 @@ async fn test_forged_session_ack_leaves_the_initiation_able_to_complete_on_the_g cleanup_nodes(&mut nodes).await; } +// ============================================================================ +// Integration tests: a lost initial msg3 +// ============================================================================ + +/// Build a two-node pair where node 0's initial msg3 was sent and dropped, so +/// node 0 is established and node 1 is still waiting for msg3. +/// +/// Periodic rekey is off on both nodes: that is the configuration with no +/// other recovery, and it keeps the rekey drivers out of the picture. Every +/// step asserts its packet count, so a harness surprise fails loudly instead +/// of being read as the defect. +async fn pair_with_lost_initial_msg3() -> Vec { + use crate::proto::fmp::wire::{CommonPrefix, PHASE_ESTABLISHED}; + + 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(); + + nodes[0] + .node + .initiate_session(node1_addr, node1_pubkey) + .await + .expect("initiate_session failed"); + + assert_eq!( + process_available_packets(&mut nodes[1..]).await, + 1, + "node 1 must have exactly node 0's SessionSetup queued" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .is_awaiting_msg3(), + "node 1 must be awaiting msg3 after answering the SessionSetup" + ); + + assert_eq!( + process_available_packets(&mut nodes[..1]).await, + 1, + "node 0 must have exactly node 1's SessionAck queued" + ); + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .expect("initiator entry present") + .is_established(), + "node 0 must be established once it has sent msg3" + ); + + // Drop msg3: take it out of node 1's queue instead of processing it. + let dropped: Vec<_> = std::iter::from_fn(|| nodes[1].packet_rx.try_recv().ok()).collect(); + assert_eq!( + dropped.len(), + 1, + "node 1 must have only node 0's msg3 queued" + ); + assert_eq!( + CommonPrefix::parse(&dropped[0].data).map(|p| p.phase), + Some(PHASE_ESTABLISHED), + "the dropped packet must be a link data frame carrying the msg3" + ); + + pump_until_quiet(&mut nodes).await; + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .is_awaiting_msg3(), + "node 1 must still be awaiting msg3 after the drop" + ); + assert_eq!( + nodes[1].packet_rx.len(), + 0, + "nothing may be left queued at node 1" + ); + + nodes +} + +/// One lost initial msg3 must cost one resend, not the session, and the +/// recovery must not depend on periodic rekey. +#[tokio::test] +async fn a_lost_initial_msg3_is_resent_and_the_responder_completes_the_session() { + let mut nodes = pair_with_lost_initial_msg3().await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); + nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + let (tun1_tx, tun1_rx) = std::sync::mpsc::channel(); + nodes[1].node.supervisor.tun_tx = Some(tun1_tx); + let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); + let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); + + let interval_ms = nodes[0] + .node + .config() + .node + .rate_limit + .handshake_resend_interval_ms; + nodes[0] + .node + .resend_pending_session_handshakes(Node::now_ms() + interval_ms + 1) + .await; + pump_until_quiet(&mut nodes).await; + + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .is_established(), + "node 1 must complete the session once node 0 resends its msg3" + ); + + let fwd = build_ipv6_packet(&fips0, &fips1, b"after msg3 resend 0 to 1"); + let rev = build_ipv6_packet(&fips1, &fips0, b"after msg3 resend 1 to 0"); + nodes[0].node.handle_tun_outbound(fwd.clone()).await; + nodes[1].node.handle_tun_outbound(rev.clone()).await; + pump_until_quiet(&mut nodes).await; + let got: Vec> = std::iter::from_fn(|| tun1_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![fwd], "node 0 to node 1 must decode"); + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![rev], "node 1 to node 0 must decode"); + + cleanup_nodes(&mut nodes).await; +} + +/// Drain and count the packets queued at `node` without processing them. +fn drain_queued(node: &mut TestNode) -> usize { + std::iter::from_fn(|| node.packet_rx.try_recv().ok()).count() +} + +/// On a healthy session the retained msg3 is released by the responder's +/// first frame, so the resend window is one round trip wide and a later tick +/// sends nothing. +#[tokio::test] +async fn an_initiator_stops_resending_msg3_once_a_responder_frame_authenticates() { + let mut nodes = make_rekey_disabled_pair().await; + establish_pair_session(&mut nodes).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); + nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); + let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); + let interval_ms = nodes[0] + .node + .config() + .node + .rate_limit + .handshake_resend_interval_ms; + + // Control: before the responder has sent anything, the sweep does resend, + // so a silent sweep at the end is the release and not a dead driver. + let t1 = Node::now_ms() + interval_ms + 1; + nodes[0].node.resend_pending_session_handshakes(t1).await; + assert_eq!( + drain_queued(&mut nodes[1]), + 1, + "control: the sweep must resend msg3 while the responder is unheard" + ); + + let rev = build_ipv6_packet(&fips1, &fips0, b"responder's first frame"); + nodes[1].node.handle_tun_outbound(rev.clone()).await; + pump_until_quiet(&mut nodes).await; + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!( + got, + vec![rev], + "the responder's frame must authenticate at node 0" + ); + assert_eq!(nodes[1].packet_rx.len(), 0, "nothing may be left queued"); + + nodes[0] + .node + .resend_pending_session_handshakes(t1 + 64_000) + .await; + assert_eq!( + nodes[1].packet_rx.len(), + 0, + "a msg3 the responder has answered must not be resent" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// The resend is harmless to a responder that already completed: it is +/// refused as a bad-state reject and both directions keep decoding. This is +/// the wire-neutrality claim, observed rather than argued. +#[tokio::test] +async fn a_resent_msg3_reaching_an_established_responder_is_refused_and_the_session_keeps_working() +{ + let mut nodes = make_rekey_disabled_pair().await; + establish_pair_session(&mut nodes).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); + nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + let (tun1_tx, tun1_rx) = std::sync::mpsc::channel(); + nodes[1].node.supervisor.tun_tx = Some(tun1_tx); + let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); + let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); + let interval_ms = nodes[0] + .node + .config() + .node + .rate_limit + .handshake_resend_interval_ms; + let before = nodes[1].node.stats().session.bad_state; + + nodes[0] + .node + .resend_pending_session_handshakes(Node::now_ms() + interval_ms + 1) + .await; + assert_eq!( + nodes[1].packet_rx.len(), + 1, + "precondition: node 0 must have resent its msg3 to node 1" + ); + pump_until_quiet(&mut nodes).await; + + assert_eq!( + nodes[1].node.stats().session.bad_state, + before + 1, + "the duplicate msg3 must be refused as a bad-state reject" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder session present") + .is_established(), + "the duplicate msg3 must not disturb the responder's session" + ); + + let fwd = build_ipv6_packet(&fips0, &fips1, b"after duplicate msg3 0 to 1"); + let rev = build_ipv6_packet(&fips1, &fips0, b"after duplicate msg3 1 to 0"); + nodes[0].node.handle_tun_outbound(fwd.clone()).await; + nodes[1].node.handle_tun_outbound(rev.clone()).await; + pump_until_quiet(&mut nodes).await; + let got: Vec> = std::iter::from_fn(|| tun1_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![fwd], "node 0 to node 1 must decode"); + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![rev], "node 1 to node 0 must decode"); + + cleanup_nodes(&mut nodes).await; +} + +/// The resend ladder is bounded in count and in time: after +/// `handshake_max_resends` resends nothing more is sent and the payload is no +/// longer held, while the session itself is kept. +#[tokio::test] +async fn initial_msg3_resends_stop_at_the_budget_and_release_the_payload() { + let mut nodes = pair_with_lost_initial_msg3().await; + let node1_addr = *nodes[1].node.node_addr(); + let max_resends = nodes[0].node.config().node.rate_limit.handshake_max_resends; + + // 64 s steps pass any single backoff interval at stock settings, so each + // step is due. Node 0 holds only its established entry, so the sweep's + // timeout pass cannot remove anything on it. + let mut now = Node::now_ms(); + let mut sent = Vec::new(); + for _ in 0..(max_resends + 2) { + now += 64_000; + nodes[0].node.resend_pending_session_handshakes(now).await; + sent.push(drain_queued(&mut nodes[1])); + } + let mut expected = vec![1; max_resends as usize]; + expected.extend([0, 0]); + assert_eq!( + sent, expected, + "one resend per due tick up to the budget, then none" + ); + + let entry = nodes[0] + .node + .get_session(&node1_addr) + .expect("the release must not tear the session down"); + assert!( + entry.handshake_payload().is_none(), + "the msg3 must no longer be held once the budget is spent" + ); + assert!( + entry.is_established(), + "node 0's session must stay established after the release" + ); + + cleanup_nodes(&mut nodes).await; +} + /// A SessionAck that fails to read must not end an FSP rekey the node /// initiated. /// @@ -6177,6 +6667,493 @@ async fn test_fresh_peer_armed_rekey_is_not_expired() { ); } +// ============================================================================ +// Integration tests: a rekey this node initiated that is never answered +// ============================================================================ + +/// Build an established two-node pair, with the handshake timeout set to +/// `timeout_secs` on both nodes when given and left at its default otherwise. +/// +/// `rekeys[i]` says whether node `i` rekeys after a single sent message; +/// otherwise it never starts a rekey of its own. The session is opened from +/// node 0, which says nothing about which node later initiates a rekey. +async fn rekey_pair(rekeys: [bool; 2], timeout_secs: Option) -> Vec { + let configs = rekeys + .iter() + .map(|&rekeys| { + let mut config = Config::new(); + if let Some(secs) = timeout_secs { + config.node.rate_limit.handshake_timeout_secs = secs; + } + if rekeys { + config.node.rekey.after_messages = 1; + } else { + config.node.rekey.after_messages = u64::MAX; + config.node.rekey.after_secs = u64::MAX; + } + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + nodes +} + +/// Send one data frame from `nodes[from]` across its rekey trigger, deliver +/// it, then run `nodes[from]`'s tick so it sends a rekey SessionSetup, and +/// assert it now holds an initiated rekey. +async fn start_rekey(nodes: &mut [TestNode], from: usize) { + let peer = *nodes[1 - from].node.node_addr(); + nodes[from] + .node + .send_session_data(&peer, 0, 0, b"crosses the rekey trigger") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(nodes).await; + + nodes[from].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[from], &peer), + "node {from} must have initiated a rekey" + ); +} + +/// Whether `node` holds a rekey handshake toward `peer` that it initiated. +fn rekey_initiated(node: &TestNode, peer: &NodeAddr) -> bool { + node.node + .get_session(peer) + .is_some_and(|e| e.has_rekey_in_progress() && e.is_rekey_initiator()) +} + +/// Whether `node`'s session with `peer` holds a completed rekey session. +fn holds_pending(node: &TestNode, peer: &NodeAddr) -> bool { + node.node + .get_session(peer) + .is_some_and(|e| e.pending_new_session().is_some()) +} + +/// Discard every packet queued at `node` without processing it, as a lossy +/// link would, and return how many were dropped. +fn drop_queued(node: &mut TestNode) -> usize { + std::iter::from_fn(|| node.packet_rx.try_recv().ok()).count() +} + +/// Run three delivery passes over every node. +async fn pump_all(nodes: &mut [TestNode]) { + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(nodes).await; + } +} + +/// A rekey whose SessionSetup was lost is retired once the handshake timeout +/// has passed, and the trigger then starts a fresh rekey that completes. +#[tokio::test] +async fn test_a_rekey_whose_session_setup_was_lost_is_retired_after_the_handshake_timeout_and_retried() + { + let mut nodes = rekey_pair([true, false], Some(1)).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[1]) > 0, + "node 0's SessionSetup must have been queued at node 1 to be lost" + ); + assert_eq!(nodes[1].node.stats().session.rekey_armed, 0); + assert!( + !nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .has_rekey_in_progress(), + "the setup must really have been dropped before node 1 armed" + ); + + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[0].node.check_session_rekey().await; + assert!( + !nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .has_rekey_in_progress(), + "an unanswered rekey this node initiated must be retired after the \ + handshake timeout, or the trigger stays vetoed for good" + ); + + nodes[0].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the trigger must start a fresh rekey once the old one is retired" + ); + pump_all(&mut nodes).await; + assert_eq!( + nodes[1].node.stats().session.rekey_armed, + 1, + "the retried setup must reach node 1" + ); + assert!( + holds_pending(&nodes[0], &node1_addr) && holds_pending(&nodes[1], &node0_addr), + "the retried rekey must complete on both nodes" + ); + assert_eq!( + nodes[0].node.stats().session.rekey_unanswered, + 1, + "the retired handshake must be counted once" + ); + assert_eq!(nodes[1].node.stats().session.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + +/// A rekey whose SessionAck was lost is retired once the handshake timeout +/// has passed, beside the responder's own expiry, and the retry completes. +#[tokio::test] +async fn test_a_rekey_whose_session_ack_was_lost_is_retired_after_the_handshake_timeout_and_retried() + { + let mut nodes = rekey_pair([true, false], Some(1)).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes[1..]).await; + assert_eq!( + nodes[1].node.stats().session.rekey_armed, + 1, + "node 1 must have armed as the rekey responder" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[0]) > 0, + "node 1's SessionAck must have been queued at node 0 to be lost" + ); + + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[1].node.check_session_rekey().await; + assert_eq!( + nodes[1].node.stats().session.rekey_expired, + 1, + "node 1's own handshake must expire on the existing responder rule" + ); + assert!( + !nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .has_rekey_in_progress() + ); + + nodes[0].node.check_session_rekey().await; + assert!( + !nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .has_rekey_in_progress(), + "an unanswered rekey this node initiated must be retired after the \ + handshake timeout, or the trigger stays vetoed for good" + ); + + nodes[0].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the trigger must start a fresh rekey once the old one is retired" + ); + pump_all(&mut nodes).await; + assert_eq!( + nodes[1].node.stats().session.rekey_armed, + 2, + "the retried setup must reach node 1" + ); + assert!( + holds_pending(&nodes[0], &node1_addr) && holds_pending(&nodes[1], &node0_addr), + "the retried rekey must complete on both nodes" + ); + assert_eq!( + nodes[0].node.stats().session.rekey_unanswered, + 1, + "the retired handshake must be counted once" + ); + assert_eq!(nodes[1].node.stats().session.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + +/// Drive a lost SessionAck, then a retry that meets the responder's own +/// handshake from the first attempt, still armed because nothing has run the +/// responder's expiry yet. +/// +/// Both nodes would rekey after one message; the initiator is picked at run +/// time so that the responder holds the smaller address when +/// `responder_wins`, and the larger otherwise. Only the initiator's tick is +/// run before the responder has armed, after which the responder is +/// dampened and cannot start a rekey of its own inside the test. +/// +/// A responder that wins the tie-break drops the retry before arming +/// anything, so both handshakes expire and the next retry completes one +/// timeout later. A responder that loses yields and answers the retry at +/// once. +async fn lostack_retry(responder_wins: bool) { + let mut nodes = rekey_pair([true, true], Some(1)).await; + let node0_smaller = + crate::proto::fsp::initiation_winner(nodes[0].node.node_addr(), nodes[1].node.node_addr()); + let resp = if responder_wins == node0_smaller { + 0 + } else { + 1 + }; + let init = 1 - resp; + let init_addr = *nodes[init].node.node_addr(); + let resp_addr = *nodes[resp].node.node_addr(); + + // First attempt: the responder arms, and its SessionAck is lost. + start_rekey(&mut nodes, init).await; + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes[resp..=resp]).await; + assert_eq!( + nodes[resp].node.stats().session.rekey_armed, + 1, + "the responder must have armed" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[init]) > 0, + "the responder's SessionAck must have been queued to be lost" + ); + + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[init].node.check_session_rekey().await; + assert!( + !nodes[init] + .node + .get_session(&resp_addr) + .unwrap() + .has_rekey_in_progress(), + "an unanswered rekey this node initiated must be retired after the \ + handshake timeout, or the trigger stays vetoed for good" + ); + + // The retry reaches a responder still holding the first handshake. + nodes[init].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[init], &resp_addr), + "the trigger must start a fresh rekey once the old one is retired" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes[resp..=resp]).await; + let stats = &nodes[resp].node.stats().session; + let (tiebreak, yielded) = (stats.rekey_tiebreak, stats.rekey_yielded); + assert_eq!( + tiebreak + yielded, + 1, + "the retry must have met the responder's stale handshake" + ); + if responder_wins { + assert_eq!(tiebreak, 1, "the smaller responder must win"); + } else { + assert_eq!(yielded, 1, "the larger responder must yield"); + } + pump_all(&mut nodes).await; + + if responder_wins { + assert!( + !holds_pending(&nodes[init], &resp_addr) && !holds_pending(&nodes[resp], &init_addr), + "a retry the responder dropped completes nothing" + ); + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[resp].node.check_session_rekey().await; + assert_eq!( + nodes[resp].node.stats().session.rekey_expired, + 1, + "the responder's stale handshake must expire on its own rule" + ); + nodes[init].node.check_session_rekey().await; + assert!( + !nodes[init] + .node + .get_session(&resp_addr) + .unwrap() + .has_rekey_in_progress(), + "the retry the responder dropped must itself be retired" + ); + nodes[init].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[init], &resp_addr), + "the trigger must start a second retry" + ); + pump_all(&mut nodes).await; + } + + assert!( + holds_pending(&nodes[init], &resp_addr) && holds_pending(&nodes[resp], &init_addr), + "the retried rekey must complete on both nodes" + ); + // The first handshake always; the retry too when the responder dropped it. + assert_eq!( + nodes[init].node.stats().session.rekey_unanswered, + if responder_wins { 2 } else { 1 }, + "every retired handshake must be counted once" + ); + assert_eq!(nodes[resp].node.stats().session.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + +/// A retry dropped on the tie-break by a smaller responder still holding its +/// stale handshake completes once both handshakes have expired. +#[tokio::test] +async fn test_a_retry_dropped_by_a_smaller_responders_stale_handshake_completes_one_timeout_later() +{ + lostack_retry(true).await; +} + +/// A retry that a larger responder, still holding its stale handshake, +/// yields to completes at once. +#[tokio::test] +async fn test_a_retry_that_a_larger_responder_yields_to_completes_at_once() { + lostack_retry(false).await; +} + +/// A forged SessionAck arriving midway through an unanswered rekey must not +/// restart its deadline: the rekey is retired on the timeout measured from +/// the setup this node sent. +#[tokio::test] +async fn test_forged_session_acks_do_not_hold_an_unanswered_rekey_open_past_its_deadline() { + let mut nodes = rekey_pair([true, false], Some(1)).await; + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[1]) > 0, + "node 0's SessionSetup must have been queued at node 1 to be lost" + ); + + tokio::time::sleep(Duration::from_millis(600)).await; + let forged = forged_session_ack(&nodes[1]); + nodes[0] + .node + .handle_session_payload(&node1_addr, &node1_addr, &forged, 1280, false) + .await; + assert_eq!(nodes[0].node.stats().session.ack_handshake_failed, 1); + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the unreadable ack must have put the handshake back" + ); + + let forged_at = std::time::Instant::now(); + tokio::time::sleep(Duration::from_millis(600)).await; + nodes[0].node.check_session_rekey().await; + println!( + "forged ack to check: {} ms", + forged_at.elapsed().as_millis() + ); + assert!( + !nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .has_rekey_in_progress(), + "a forged SessionAck must not restart the deadline of the rekey this \ + node initiated" + ); + assert_eq!( + nodes[0].node.stats().session.rekey_unanswered, + 1, + "the retired handshake must be counted once" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// Putting a handshake back after an unreadable SessionAck leaves the stamp +/// its deadline runs from exactly where arming wrote it. +#[tokio::test] +async fn test_an_unreadable_session_ack_does_not_push_out_the_rekey_deadline() { + let mut nodes = rekey_pair([true, false], None).await; + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + let armed_at = nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(); + assert_ne!(armed_at, 0, "arming must stamp the deadline"); + + tokio::time::sleep(Duration::from_millis(5)).await; + let forged = forged_session_ack(&nodes[1]); + nodes[0] + .node + .handle_session_payload(&node1_addr, &node1_addr, &forged, 1280, false) + .await; + assert_eq!(nodes[0].node.stats().session.ack_handshake_failed, 1); + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the unreadable ack must have put the handshake back" + ); + assert_eq!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(), + armed_at, + "the restore must not restamp the deadline" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// A fresh rekey this node initiated is not retired because the peer's own +/// last rekey is older than the handshake timeout. +#[tokio::test] +async fn test_a_fresh_initiated_rekey_is_not_retired_on_the_peers_older_rekey_stamp() { + let mut nodes = rekey_pair([true, false], None).await; + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + // A peer rekey long past both the handshake timeout and the dampening. + nodes[0] + .node + .sessions + .get_mut(&node1_addr) + .unwrap() + .record_peer_rekey(wall_clock_ms() - 60_000); + let armed_at = nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(); + + nodes[0].node.check_session_rekey().await; + + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "a fresh rekey this node initiated must not be retired on the peer's clock" + ); + assert_eq!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(), + armed_at, + "the rekey must be the same one, not retired and re-armed" + ); + let stats = &nodes[0].node.stats().session; + assert_eq!(stats.rekey_expired, 0); + assert_eq!(stats.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + // ============================================================================ // Integration tests: a peer's cutover after a long silence // ============================================================================ diff --git a/src/peer/active.rs b/src/peer/active.rs index 81ba9b4e..0c7e4291 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -7,6 +7,7 @@ use crate::config::MmpConfig; use crate::node::REKEY_JITTER_SECS; use crate::noise::{HandshakeState as NoiseHandshakeState, NoiseError, NoiseSession}; use crate::proto::bloom::BloomFilter; +use crate::proto::fmp::RekeyRole; use crate::proto::mmp::MmpPeerState; use crate::proto::stp::{ParentDeclaration, TreeCoordinate}; use crate::transport::{LinkId, LinkStats, TransportAddr, TransportId}; @@ -15,7 +16,7 @@ use crate::{FipsAddress, NodeAddr, PeerIdentity}; use rand::RngExt; use secp256k1::XOnlyPublicKey; use std::fmt; -use std::time::Instant; +use std::time::{Duration, Instant}; /// Draw a fresh per-session rekey jitter from `[-REKEY_JITTER_SECS, +REKEY_JITTER_SECS]`. fn draw_rekey_jitter() -> i64 { @@ -268,6 +269,11 @@ pub struct ActivePeer { rekey_msg1_next_resend: u64, /// In-progress rekey: number of msg1 retransmissions performed so far. rekey_msg1_resend_count: u32, + /// Which side installed the pending session (`None` when no pending is + /// held). Set with the pending slot and cleared with it. + pending_role: Option, + /// When the pending session was installed, for the responder hold. + pending_since: Option, // === Published active-send-state (two-tier boundary) === /// The send-critical subset read (and, on roam/responder-cutover, written) @@ -309,6 +315,8 @@ impl ActivePeer { rekey_msg1: None, rekey_msg1_next_resend: 0, rekey_msg1_resend_count: 0, + pending_role: None, + pending_since: None, send: PeerSendState::new(link_id, now, authenticated_at), } } @@ -388,6 +396,8 @@ impl ActivePeer { rekey_msg1: None, rekey_msg1_next_resend: 0, rekey_msg1_resend_count: 0, + pending_role: None, + pending_since: None, send, } } @@ -917,6 +927,16 @@ impl ActivePeer { .unwrap_or_else(Instant::now); } + /// Test-only seam: backdate the pending session's install time so a test + /// can make the responder hold read as passed. Shifts only the private + /// timestamp; compiled out of release builds. + #[cfg(test)] + pub(crate) fn backdate_pending(&mut self, age: Duration) { + self.pending_since = self + .pending_since + .map(|t| t.checked_sub(age).unwrap_or_else(Instant::now)); + } + /// Test-only seam: install link-layer MMP state with a chosen operating /// mode on a peer that was constructed without a Noise session (the bare /// `new` constructor leaves `mmp` as `None`). This only attaches the same @@ -1024,19 +1044,62 @@ impl ActivePeer { self.send.pending_new_session.as_mut() } + /// Which side of the rekey handshake produced the pending session; `None` + /// when no pending session is held. + pub(crate) fn pending_role(&self) -> Option { + self.pending_role + } + + /// Check whether the pending session has been held for at least `hold` + /// since it was installed. False when no pending session is held. + pub(crate) fn pending_expired(&self, hold: Duration) -> bool { + self.pending_since.is_some_and(|t| t.elapsed() >= hold) + } + /// Store a completed rekey session and its indices. /// /// Called when the rekey handshake completes. The session is held /// as pending until the initiator flips the K-bit on the next outbound packet. + /// Records this node as the initiator; a pending session answered for the + /// peer is stored with [`answer_rekey`](Self::answer_rekey). pub fn set_pending_session( &mut self, session: NoiseSession, our_index: SessionIndex, their_index: SessionIndex, + ) { + self.install_pending(session, our_index, their_index, RekeyRole::Initiator); + } + + /// Store the session this node produced by answering the peer's rekey + /// msg1. It is held until a frame on the new epoch from the peer + /// authenticates against it ([`handle_peer_kbit_flip`](Self::handle_peer_kbit_flip)), + /// or until the responder hold passes and the node retires it + /// ([`retire_pending`](Self::retire_pending)); it is never cut over on this + /// node's own schedule. + pub(crate) fn answer_rekey( + &mut self, + session: NoiseSession, + our_index: SessionIndex, + their_index: SessionIndex, + ) { + self.install_pending(session, our_index, their_index, RekeyRole::Responder); + } + + /// Store a pending session with the role that produced it and the time it + /// was installed; the one writer that fills the pending slot. + fn install_pending( + &mut self, + session: NoiseSession, + our_index: SessionIndex, + their_index: SessionIndex, + role: RekeyRole, ) { self.send.pending_new_session = Some(session); self.send.pending_our_index = Some(our_index); self.send.pending_their_index = Some(their_index); + self.pending_role = Some(role); + self.pending_since = Some(Instant::now()); self.rekey_in_progress = false; // Clear initiator handshake state (index now lives in pending_our_index) self.rekey_our_index = None; @@ -1044,6 +1107,11 @@ impl ActivePeer { self.rekey_msg1 = None; self.rekey_msg1_next_resend = 0; self.rekey_msg1_resend_count = 0; + debug_assert_eq!( + self.pending_role.is_some(), + self.send.pending_new_session.is_some(), + "install_pending: pending role out of step with the pending slot" + ); } /// Cut over to the pending new session (initiator side). @@ -1055,6 +1123,8 @@ impl ActivePeer { let new_session = self.send.pending_new_session.take()?; let new_our_index = self.send.pending_our_index.take(); let new_their_index = self.send.pending_their_index.take(); + self.pending_role = None; + self.pending_since = None; // Demote current to previous self.send.previous_session = self.send.noise_session.take(); @@ -1081,6 +1151,11 @@ impl ActivePeer { mmp.reset_for_rekey(now_ms); } + debug_assert_eq!( + self.pending_role.is_some(), + self.send.pending_new_session.is_some(), + "cutover_to_new_session: pending role out of step with the pending slot" + ); self.send.previous_our_index } @@ -1092,6 +1167,8 @@ impl ActivePeer { let new_session = self.send.pending_new_session.take()?; let new_our_index = self.send.pending_our_index.take(); let new_their_index = self.send.pending_their_index.take(); + self.pending_role = None; + self.pending_since = None; // Demote current to previous self.send.previous_session = self.send.noise_session.take(); @@ -1118,6 +1195,11 @@ impl ActivePeer { mmp.reset_for_rekey(now_ms); } + debug_assert_eq!( + self.pending_role.is_some(), + self.send.pending_new_session.is_some(), + "handle_peer_kbit_flip: pending role out of step with the pending slot" + ); self.send.previous_our_index } @@ -1144,6 +1226,26 @@ impl ActivePeer { self.send.previous_our_index.take() } + /// Drop a pending session this node did not initiate, which the + /// initiator never adopted. Returns its index so the caller can + /// unregister and free it; `None` if no such pending is held. A pending + /// this node initiated is left alone: that one is cut over, not retired. + pub(crate) fn retire_pending(&mut self) -> Option { + if self.pending_role == Some(RekeyRole::Initiator) { + return None; + } + self.send.pending_new_session.take()?; + self.send.pending_their_index = None; + self.pending_role = None; + self.pending_since = None; + debug_assert_eq!( + self.pending_role.is_some(), + self.send.pending_new_session.is_some(), + "retire_pending: pending role out of step with the pending slot" + ); + self.send.pending_our_index.take() + } + /// Abandon an in-progress rekey. /// /// Returns the rekey our_index so the caller can free it. @@ -1156,11 +1258,19 @@ impl ActivePeer { self.rekey_msg1_resend_count = 0; self.rekey_in_progress = false; // Return whichever index needs freeing - self.rekey_our_index.take().or_else(|| { + let freed = self.rekey_our_index.take().or_else(|| { self.send.pending_new_session = None; self.send.pending_their_index = None; + self.pending_role = None; + self.pending_since = None; self.send.pending_our_index.take() - }) + }); + debug_assert_eq!( + self.pending_role.is_some(), + self.send.pending_new_session.is_some(), + "abandon_rekey: pending role out of step with the pending slot" + ); + freed } // === Rekey Handshake State (Initiator) === @@ -1729,4 +1839,67 @@ mod tests { }); assert_eq!(cur_pt.as_deref(), Some(&b"steady"[..])); } + + /// Retiring a pending session this node answered hands back its index for + /// the caller to free, empties the pending slot and its role, and leaves + /// the current session alone. + #[test] + fn retiring_an_answered_pending_returns_its_index_and_empties_the_slot() { + let (_cur_send, cur_recv) = ik_session_pair(); + let (_pend_send, pend_recv) = ik_session_pair(); + let mut peer = peer_with_current(cur_recv); + peer.answer_rekey(pend_recv, SessionIndex::new(3), SessionIndex::new(4)); + assert_eq!(peer.pending_role(), Some(RekeyRole::Responder)); + + assert_eq!(peer.retire_pending(), Some(SessionIndex::new(3))); + assert!(peer.pending_new_session().is_none()); + assert_eq!(peer.pending_role(), None); + assert_eq!(peer.pending_their_index(), None); + assert_eq!(peer.pending_our_index(), None); + assert!(peer.noise_session().is_some()); + assert_eq!(peer.our_index(), Some(SessionIndex::new(1))); + } + + /// Retirement leaves a pending session this node initiated in place: that + /// one is cut over on this node's schedule, never retired. + #[test] + fn retire_leaves_a_pending_this_node_initiated_in_place() { + let (_cur_send, cur_recv) = ik_session_pair(); + let (_pend_send, pend_recv) = ik_session_pair(); + let mut peer = peer_with_current(cur_recv); + peer.set_pending_session(pend_recv, SessionIndex::new(3), SessionIndex::new(4)); + assert_eq!(peer.pending_role(), Some(RekeyRole::Initiator)); + + assert_eq!(peer.retire_pending(), None); + assert!(peer.pending_new_session().is_some()); + assert_eq!(peer.pending_role(), Some(RekeyRole::Initiator)); + } + + /// Promoting an answered pending session on the peer's first new-epoch + /// frame clears its role and install time with the slot. + #[test] + fn promotion_on_the_peers_new_epoch_frame_clears_the_pending_role() { + let (_cur_send, cur_recv) = ik_session_pair(); + let (_pend_send, pend_recv) = ik_session_pair(); + let mut peer = peer_with_current(cur_recv); + peer.answer_rekey(pend_recv, SessionIndex::new(3), SessionIndex::new(4)); + + assert!(peer.handle_peer_kbit_flip().is_some()); + assert_eq!(peer.pending_role(), None); + assert!(!peer.pending_expired(Duration::ZERO)); + } + + /// The pending hold is measured from the install time: not expired inside + /// the hold, expired once the install time is older than it. + #[test] + fn pending_expiry_reads_the_install_time_against_the_hold() { + let (_cur_send, cur_recv) = ik_session_pair(); + let (_pend_send, pend_recv) = ik_session_pair(); + let mut peer = peer_with_current(cur_recv); + peer.answer_rekey(pend_recv, SessionIndex::new(3), SessionIndex::new(4)); + + assert!(!peer.pending_expired(Duration::from_secs(60))); + peer.backdate_pending(Duration::from_secs(61)); + assert!(peer.pending_expired(Duration::from_secs(60))); + } } diff --git a/src/peer/machine.rs b/src/peer/machine.rs index b812bc98..47030771 100644 --- a/src/peer/machine.rs +++ b/src/peer/machine.rs @@ -62,7 +62,7 @@ use crate::noise::{self, NoiseError, NoiseSession}; use crate::proto::fmp::{ ConnAction, ConnSnapshot, ConnectionState, EstablishSnapshot, Fmp, InboundDecision, OutboundDecision, OutboundSnapshot, PeerSnapshot, PromotionResult, RekeyCfg, - RekeyResendSnapshot, WireOutcome, + RekeyResendSnapshot, RekeyRole, WireOutcome, }; use crate::proto::link::LinkMessageType; use crate::transport::{LinkDirection, LinkId, LinkStats, TransportAddr, TransportId}; @@ -1836,6 +1836,9 @@ impl PeerMachine { elapsed_secs, counter: 0, jitter_secs: self.rekey_jitter_secs, + pending_role: (phase == Some(RekeyPhase::PendingCutover)) + .then_some(RekeyRole::Initiator), + pending_expired: false, } } } diff --git a/src/proto/fmp/core.rs b/src/proto/fmp/core.rs index 58598c0d..82b6a630 100644 --- a/src/proto/fmp/core.rs +++ b/src/proto/fmp/core.rs @@ -131,6 +131,19 @@ pub(crate) struct ConnSnapshot { pub msg1: Vec, } +/// Which side of a link rekey handshake produced a pending session. +/// +/// Only the side that initiated may commit to the new keys on its own +/// schedule: it holds proof that the peer derived them. The responder +/// learns that only when a frame sealed on the new epoch authenticates. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum RekeyRole { + /// This node sent the rekey msg1 and read the peer's msg2. + Initiator, + /// This node answered the peer's rekey msg1 with a msg2. + Responder, +} + /// A snapshot of one active peer's rekey-relevant state, taken by the shell. /// /// Every clock read is resolved shell-side into a plain `u64`/`bool` before the @@ -143,8 +156,15 @@ pub(crate) struct ConnSnapshot { pub(crate) struct PeerSnapshot { /// The peer's node address (cutover/drain/rekey target). pub addr: NodeAddr, - /// A pending post-rekey session is ready to cut over to. + /// A pending post-rekey session is held (cut over by the initiator, + /// promoted on the peer's first new-epoch frame by the responder). pub has_pending: bool, + /// Which role installed the pending session; `Some` exactly when + /// `has_pending`. + pub pending_role: Option, + /// The pending session has been held past the responder hold + /// (pre-evaluated shell-side against the hold, as `drain_expired` is). + pub pending_expired: bool, /// A rekey handshake is currently in flight. pub rekey_in_progress: bool, /// The peer is in its post-cutover drain window. @@ -293,6 +313,10 @@ pub(crate) enum ConnAction { /// Complete `peer`'s drain window: erase the previous session, free its /// index, and unregister its decrypt-worker entry. Drain { peer: NodeAddr }, + /// Retire `peer`'s pending session that this node did not initiate and the + /// initiator never adopted: drop it, unregister its index from + /// `peers_by_index`, and free the index. + RetirePending { peer: NodeAddr }, /// Initiate a fresh outbound rekey to `peer` (`initiate_rekey`: allocates a /// new index, builds and sends msg1, inserts `pending_outbound`). The msg1 /// construction is the establish leaf and stays shell-side; the action @@ -502,33 +526,47 @@ impl Fmp { /// snapshotted. Reproduces the pre-refactor priority and phase grouping /// exactly: /// - /// - **Cutover** takes precedence: a peer with a pending session and no - /// in-flight rekey cuts over and is considered for nothing else. - /// - Otherwise an expired drain window is completed, and — independently — - /// the rekey trigger fires when the peer is neither mid-rekey nor - /// dampened and its jittered time threshold or send counter is reached. - /// A draining peer can thus both drain and re-trigger in the same tick, - /// as before. + /// - **Cutover** takes precedence: a peer with a pending session this node + /// initiated and no in-flight rekey cuts over and is considered for + /// nothing else. + /// - A pending session this node answered is held, never cut over by the + /// tick: the peer's first frame on the new epoch promotes it. Once its + /// hold has passed it is retired. + /// - An expired drain window is completed, and — independently — the rekey + /// trigger fires when the peer is neither mid-rekey, dampened, nor + /// holding a pending session, and its jittered time threshold or send + /// counter is reached. A draining peer can thus both drain and + /// re-trigger in the same tick, as before. /// /// Actions are returned phase-grouped (all cutovers, then all drains, then - /// all rekey initiations) to preserve the pre-refactor global execution - /// order across peers, which the shared `index_allocator` observes. + /// all retirements, then all rekey initiations) to preserve the global + /// execution order across peers, which the shared `index_allocator` + /// observes: retirements free an index, so they run before initiations + /// allocate. pub(crate) fn poll_rekey(&self, peers: Vec, cfg: &RekeyCfg) -> Vec { let mut cutovers = Vec::new(); let mut drains = Vec::new(); + let mut retires = Vec::new(); let mut rekeys = Vec::new(); for p in peers { - // 1. Initiator-side cutover. - if p.has_pending && !p.rekey_in_progress { + let initiated = p.pending_role == Some(RekeyRole::Initiator); + // 1. Initiator-side cutover. A pending this node answered is promoted + // by the peer's first frame on the new epoch, never by this tick. + if p.has_pending && !p.rekey_in_progress && initiated { cutovers.push(ConnAction::Cutover { peer: p.addr }); continue; } + // 1b. A pending the initiator never adopted is retired at its hold. + if p.has_pending && !initiated && p.pending_expired { + retires.push(ConnAction::RetirePending { peer: p.addr }); + } // 2. Drain window expiry (does not preclude a trigger below). if p.is_draining && p.drain_expired { drains.push(ConnAction::Drain { peer: p.addr }); } - // 3. Rekey trigger. - if p.rekey_in_progress || p.is_dampened { + // 3. Rekey trigger. A held pending vetoes it: a new cycle's msg2 would + // overwrite the held slot. + if p.rekey_in_progress || p.is_dampened || p.has_pending { continue; } let effective_after = cfg.after_secs.saturating_add_signed(p.jitter_secs); @@ -537,6 +575,7 @@ impl Fmp { } } cutovers.extend(drains); + cutovers.extend(retires); cutovers.extend(rekeys); cutovers } diff --git a/src/proto/fmp/mod.rs b/src/proto/fmp/mod.rs index 8a970418..24262eb7 100644 --- a/src/proto/fmp/mod.rs +++ b/src/proto/fmp/mod.rs @@ -37,7 +37,7 @@ mod tests; pub(crate) use core::{ ConnAction, ConnSnapshot, EstablishSnapshot, EstablishView, InboundDecision, InboundReject, LifecycleView, OutboundDecision, OutboundSnapshot, PeerSnapshot, RekeyCfg, RekeyResendSnapshot, - WireOutcome, + RekeyRole, WireOutcome, }; pub use core::{PromotionResult, cross_connection_winner}; pub(crate) use limits::backoff_ms; diff --git a/src/proto/fmp/tests/core.rs b/src/proto/fmp/tests/core.rs index 90ad973a..ad81608e 100644 --- a/src/proto/fmp/tests/core.rs +++ b/src/proto/fmp/tests/core.rs @@ -7,7 +7,7 @@ use super::util::{ use crate::NodeAddr; use crate::proto::fmp::{ ConnAction, Fmp, InboundDecision, InboundReject, OutboundDecision, OutboundSnapshot, RekeyCfg, - cross_connection_winner, + RekeyRole, cross_connection_winner, }; use crate::testutil::make_node_addr; use crate::transport::LinkId; @@ -122,6 +122,7 @@ fn rekey_cutover_takes_precedence_over_trigger() { let fmp = Fmp::new(); let mut p = peer_snapshot(0x10); p.has_pending = true; + p.pending_role = Some(RekeyRole::Initiator); // Wildly over the time threshold, but cutover wins and nothing else fires. p.elapsed_secs = 10_000; p.counter = 10_000; @@ -229,6 +230,7 @@ fn rekey_actions_are_phase_grouped_across_peers() { a.elapsed_secs = 200; let mut b = peer_snapshot(0x02); b.has_pending = true; + b.pending_role = Some(RekeyRole::Initiator); let mut c = peer_snapshot(0x03); c.is_draining = true; c.drain_expired = true; @@ -607,3 +609,101 @@ fn test_cross_connection_symmetric() { // Exactly one survives assert!(a_outbound_wins != a_inbound_wins); } + +// --- rekey role gate: a pending session this node answered --- + +/// A pending session this node answered, held past the responder hold, is +/// retired, and the tick neither cuts it over nor starts a rekey of its own. +#[test] +fn a_responder_held_pending_past_its_hold_is_retired_and_nothing_else_fires() { + let fmp = Fmp::new(); + let mut p = peer_snapshot(0x20); + p.has_pending = true; + p.pending_role = Some(RekeyRole::Responder); + p.pending_expired = true; + p.elapsed_secs = 10_000; + p.counter = 10_000; + let actions = fmp.poll_rekey(vec![p], &cfg()); + assert_eq!(actions.len(), 1); + assert!( + matches!(actions[0], ConnAction::RetirePending { peer } if peer == make_node_addr(0x20)) + ); +} + +/// A pending session this node answered, still inside its hold, is left +/// alone: no cutover, no retirement, and no rekey trigger even with both +/// thresholds met, since a new cycle would overwrite the held slot. +#[test] +fn a_responder_held_pending_inside_its_hold_is_neither_cut_over_nor_allowed_to_trigger() { + let fmp = Fmp::new(); + let mut p = peer_snapshot(0x21); + p.has_pending = true; + p.pending_role = Some(RekeyRole::Responder); + p.pending_expired = false; + p.elapsed_secs = 10_000; + p.counter = 10_000; + assert!(fmp.poll_rekey(vec![p], &cfg()).is_empty()); +} + +/// A pending session this node initiated is cut over even when its hold +/// reads as passed: retirement never touches the initiator's pending. +#[test] +fn an_initiator_held_pending_cuts_over_even_when_its_hold_has_passed() { + let fmp = Fmp::new(); + let mut p = peer_snapshot(0x22); + p.has_pending = true; + p.pending_role = Some(RekeyRole::Initiator); + p.pending_expired = true; + let actions = fmp.poll_rekey(vec![p], &cfg()); + assert_eq!(actions.len(), 1); + assert!(matches!(actions[0], ConnAction::Cutover { peer } if peer == make_node_addr(0x22))); +} + +/// A pending session with no recorded role fails safe: it is held like a +/// responder's, never cut over by the tick, and vetoes the rekey trigger. +#[test] +fn a_pending_with_no_recorded_role_is_held_and_never_cut_over_by_the_tick() { + let fmp = Fmp::new(); + let mut p = peer_snapshot(0x23); + p.has_pending = true; + p.pending_role = None; + p.pending_expired = false; + p.elapsed_secs = 10_000; + p.counter = 10_000; + assert!(fmp.poll_rekey(vec![p], &cfg()).is_empty()); +} + +/// Retirements free an index, so they are grouped after drains (which free) +/// and before rekey initiations (which allocate): cutovers, drains, +/// retirements, rekeys. +#[test] +fn retirements_are_grouped_after_drains_and_before_rekey_initiations() { + let fmp = Fmp::new(); + // a: trigger only. b: initiator cutover. c: drain and trigger. d: retire. + let mut a = peer_snapshot(0x01); + a.elapsed_secs = 200; + let mut b = peer_snapshot(0x02); + b.has_pending = true; + b.pending_role = Some(RekeyRole::Initiator); + let mut c = peer_snapshot(0x03); + c.is_draining = true; + c.drain_expired = true; + c.counter = 5_000; + let mut d = peer_snapshot(0x04); + d.has_pending = true; + d.pending_role = Some(RekeyRole::Responder); + d.pending_expired = true; + let actions = fmp.poll_rekey(vec![a, b, c, d], &cfg()); + assert_eq!(actions.len(), 5); + assert!(matches!(actions[0], ConnAction::Cutover { peer } if peer == make_node_addr(0x02))); + assert!(matches!(actions[1], ConnAction::Drain { peer } if peer == make_node_addr(0x03))); + assert!( + matches!(actions[2], ConnAction::RetirePending { peer } if peer == make_node_addr(0x04)) + ); + assert!( + matches!(actions[3], ConnAction::InitiateRekey { peer } if peer == make_node_addr(0x01)) + ); + assert!( + matches!(actions[4], ConnAction::InitiateRekey { peer } if peer == make_node_addr(0x03)) + ); +} diff --git a/src/proto/fmp/tests/util.rs b/src/proto/fmp/tests/util.rs index 3684c76c..ebe5cd85 100644 --- a/src/proto/fmp/tests/util.rs +++ b/src/proto/fmp/tests/util.rs @@ -38,6 +38,8 @@ pub(super) fn peer_snapshot(addr_byte: u8) -> PeerSnapshot { elapsed_secs: 0, counter: 0, jitter_secs: 0, + pending_role: None, + pending_expired: false, } } diff --git a/src/proto/fsp/core.rs b/src/proto/fsp/core.rs index dc061cb0..33fc8b62 100644 --- a/src/proto/fsp/core.rs +++ b/src/proto/fsp/core.rs @@ -2,8 +2,8 @@ //! //! Pure, runtime-agnostic decisions for the FSP end-to-end session lifecycle: //! the per-tick rekey choreography (initiator cutover, drain completion, rekey -//! trigger), msg3 retransmission classification, and the post-decrypt epoch -//! reaction. The async I/O adapters in `node::handlers::{rekey,session}` build +//! trigger), rekey and initial-handshake msg3 resend classification, and the +//! post-decrypt epoch reaction. The async I/O adapters in `node::handlers::{rekey,session}` build //! the plain-data snapshots (pre-computing every clock read into `u64`/`bool`), //! call these decisions, and drive the returned effects — the sends, the //! `SessionEntry` mutations, metrics, and logging. No I/O, no clock, no crypto, @@ -65,9 +65,27 @@ pub(crate) enum FspAction { /// Discarding it would make every later frame from that peer /// undecryptable, so the handshake alone is dropped. AbandonHandshake { addr: NodeAddr }, + /// Drop only `addr`'s handshake that this node armed by sending a setup + /// message never answered within the handshake timeout + /// (`SessionEntry::abandon_handshake`). The trigger starts a fresh rekey + /// on a later tick. + /// + /// The handshake only, not [`AbandonRekey`](Self::AbandonRekey): an + /// entry whose handshake this node armed holds no pending session, and + /// dropping only the handshake keeps any such session safe even if that + /// ever stops holding. + ExpireInitiation { addr: NodeAddr }, /// Retransmit `addr`'s retained rekey msg3 (the shell re-reads the payload /// from the entry, sends it, then records the retransmission on success). ResendSessionMsg3 { addr: NodeAddr }, + /// Resend `addr`'s retained initial-handshake msg3 (the shell re-reads the + /// payload from the entry's handshake resend slot, sends it, then records + /// the resend on success). + ResendInitialMsg3 { addr: NodeAddr }, + /// Stop retaining `addr`'s initial-handshake msg3: the resend budget is + /// spent without an inbound frame showing the peer received it. The session + /// itself is kept; only the retained payload is dropped. + ReleaseInitialMsg3 { addr: NodeAddr }, /// Cache `coords` for `addr` in the shared coordinate cache /// (`coord_cache.insert`). CacheCoords { @@ -135,9 +153,14 @@ pub(crate) struct SessionSnapshot { /// A handshake armed by the peer's setup message has passed the handshake /// timeout without its msg3 (pre-evaluated: `last_peer_rekey_ms != 0 && /// now - last_peer_rekey_ms > handshake_timeout`). False when the entry - /// carries no peer-rekey stamp, so a handshake this side armed is never - /// aged out on the peer's clock. + /// carries no peer-rekey stamp. A handshake this side armed is never + /// aged out on the peer's clock; it has its own deadline, + /// [`initiation_expired`](Self::initiation_expired). pub armed_handshake_expired: bool, + /// This node's own armed handshake has passed the handshake timeout + /// measured from the setup message it sent (pre-evaluated: `now - + /// initiated_ms > handshake_timeout`). + pub initiation_expired: bool, /// Monotonic session age in seconds (`(now - session_start_ms) / 1000`). pub elapsed_secs: u64, /// Current Noise send counter. @@ -159,6 +182,18 @@ pub(crate) struct RekeyMsg3ResendSnapshot { pub resend_due: bool, } +/// A snapshot of one established session whose initiator still retains its +/// initial-handshake msg3, taken by the shell for the resend decision. +pub(crate) struct InitialMsg3ResendSnapshot { + /// The session's remote node address (release/resend target). + pub addr: NodeAddr, + /// How many msg3 resends have already happened. + pub resend_count: u32, + /// The retained msg3 is due as of the shell's `now_ms` (pre-evaluated: + /// `next_resend_at_ms != 0 && now_ms >= next_resend_at_ms`). + pub resend_due: bool, +} + /// Which key epoch a just-decrypted frame authenticated against — the shell-side /// [`EpochSlot`](crate::node::session::EpochSlot) mapped to a proto-local /// plain-data mirror so the core carries no `node` dependency. @@ -199,18 +234,22 @@ impl Fsp { /// in-flight rekey, and an elapsed liveness timer cuts over and is /// considered for nothing else. /// - Otherwise an expired drain window is completed, and — independently — - /// a handshake the peer armed and never finished is abandoned, which is - /// the last word on that session this tick. + /// a handshake the peer armed and never finished, or a handshake this + /// node armed and never got an answer to, is abandoned, which is the + /// last word on that session this tick. /// - Failing both, the rekey trigger fires when the session is neither /// mid-rekey, holding a pending session, retaining a msg3 payload, nor /// dampened, and its jittered time threshold or send counter is reached. - /// Only this last decision is gated on `cfg.enabled`: the other three - /// maintain state a peer's setup message can create with periodic rekey - /// switched off. + /// Only this last decision is gated on `cfg.enabled`: the cutover, the + /// drain and the abandon of a peer-armed handshake maintain state a + /// peer's setup message can create with periodic rekey switched off. + /// A handshake this node armed exists only when the trigger fired, but + /// its expiry sits above the gate too, so whether an existing + /// handshake is retired does not depend on whether new ones may start. /// /// Actions are returned phase-grouped (all cutovers, then all drains, then - /// all abandoned handshakes, then all rekey initiations) to preserve the - /// pre-refactor execution order. + /// all abandoned or expired handshakes, then all rekey initiations) to + /// preserve the pre-refactor execution order. pub(crate) fn poll_rekey( &self, sessions: Vec, @@ -244,6 +283,15 @@ impl Fsp { abandons.push(FspAction::AbandonHandshake { addr: s.addr }); continue; } + // 3b. Retire a handshake this node armed whose setup or + // SessionAck was lost. Arm 3 cannot: it runs on the peer's + // clock, which says nothing about our own setup. Only the + // handshake goes; the trigger below starts a fresh rekey on a + // later tick. + if s.is_rekey_initiator && s.rekey_in_progress && s.initiation_expired { + abandons.push(FspAction::ExpireInitiation { addr: s.addr }); + continue; + } // 4. Rekey trigger. if !cfg.enabled { continue; @@ -292,6 +340,33 @@ impl Fsp { abandons } + /// Decide the initial-handshake msg3 resends for the established sessions + /// the shell snapshotted as retaining one. A due candidate whose budget is + /// spent is released; an in-budget due candidate is resent; a candidate not + /// yet due gets nothing this tick. Releases come first, as abandons do in + /// [`poll_rekey_msg3_resends`](Self::poll_rekey_msg3_resends); the shell + /// commits a resend's count and reschedule only on a successful send. + pub(crate) fn poll_initial_msg3_resends( + &self, + candidates: Vec, + max_resends: u32, + ) -> Vec { + let mut releases = Vec::new(); + let mut resends = Vec::new(); + for c in candidates { + if !c.resend_due { + continue; + } + if c.resend_count >= max_resends { + releases.push(FspAction::ReleaseInitialMsg3 { addr: c.addr }); + continue; + } + resends.push(FspAction::ResendInitialMsg3 { addr: c.addr }); + } + releases.extend(resends); + releases + } + /// Classify the reaction to a frame that authenticated against `slot`. Pure /// over the slot and the two plain-data session flags; the shell applies the /// resulting `SessionEntry` mutation and observability. diff --git a/src/proto/fsp/mod.rs b/src/proto/fsp/mod.rs index 50c73e73..6ebb179c 100644 --- a/src/proto/fsp/mod.rs +++ b/src/proto/fsp/mod.rs @@ -12,7 +12,8 @@ //! address-only coordinate helpers downward from `crate::proto::stp`. //! //! - `core.rs` — the stateless [`Fsp`] anchor + [`FspAction`]: the pure rekey -//! choreography (`poll_rekey`/`poll_rekey_msg3_resends`), the post-decrypt +//! choreography (`poll_rekey`/`poll_rekey_msg3_resends`), the initial-handshake +//! msg3 resend decision (`poll_initial_msg3_resends`), the post-decrypt //! `classify_epoch`, the initiation tie-break, and the pure MTU-clamp / //! bounded-queue / ECN transforms. No clock/crypto/I/O/tracing. //! - `limits.rs` — the session-rekey timing constants. @@ -27,8 +28,9 @@ pub(crate) mod wire; mod tests; pub(crate) use core::{ - DecryptSlot, EpochReaction, Fsp, FspAction, RekeyCfg, RekeyMsg3ResendSnapshot, SessionSnapshot, - cutover_timer_elapsed, initiation_winner, mark_ipv6_ecn_ce, push_bounded_pending, + DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg, + RekeyMsg3ResendSnapshot, SessionSnapshot, cutover_timer_elapsed, initiation_winner, + mark_ipv6_ecn_ce, push_bounded_pending, }; pub use wire::{ FspInnerFlags, SessionAck, SessionFlags, SessionMessageType, SessionMsg3, SessionSetup, diff --git a/src/proto/fsp/tests/core.rs b/src/proto/fsp/tests/core.rs index f23d8ea7..0e80f3af 100644 --- a/src/proto/fsp/tests/core.rs +++ b/src/proto/fsp/tests/core.rs @@ -2,9 +2,9 @@ use crate::FipsAddress; use crate::proto::fsp::core::{ - DecryptSlot, EpochReaction, Fsp, FspAction, RekeyCfg, RekeyMsg3ResendSnapshot, SessionSnapshot, - cutover_timer_elapsed, initiation_winner, mark_ipv6_ecn_ce, push_bounded_pending, - should_apply_path_mtu, + DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg, + RekeyMsg3ResendSnapshot, SessionSnapshot, cutover_timer_elapsed, initiation_winner, + mark_ipv6_ecn_ce, push_bounded_pending, should_apply_path_mtu, }; use crate::proto::fsp::limits::FSP_CUTOVER_DELAY_MS; use crate::proto::stp::TreeCoordinate; @@ -16,7 +16,8 @@ fn coords(byte: u8) -> TreeCoordinate { } /// A quiescent established-session snapshot: no pending cutover, no drain, no -/// dampening, zero ages/counter/jitter. Tests set only the fields they exercise. +/// dampening, no expired initiation, zero ages/counter/jitter. Tests set only +/// the fields they exercise. fn session_snapshot(addr_byte: u8) -> SessionSnapshot { SessionSnapshot { addr: make_node_addr(addr_byte), @@ -29,6 +30,7 @@ fn session_snapshot(addr_byte: u8) -> SessionSnapshot { has_rekey_msg3_payload: false, is_dampened: false, armed_handshake_expired: false, + initiation_expired: false, elapsed_secs: 0, counter: 0, jitter_secs: 0, @@ -252,28 +254,135 @@ fn poll_rekey_abandons_the_handshake_without_touching_a_completed_pending() { ); } +/// A handshake this node armed whose setup or SessionAck was lost, past the +/// handshake timeout on its own clock and carrying no peer stamp, with the +/// send counter over the rekey threshold. +fn unanswered_initiation(addr_byte: u8) -> SessionSnapshot { + let mut s = session_snapshot(addr_byte); + s.rekey_in_progress = true; + s.is_rekey_initiator = true; + s.initiation_expired = true; + s.armed_handshake_expired = false; + s.counter = 5000; + s +} + +/// A rekey this node initiated whose setup or ack was lost is retired on its +/// own deadline, and once the handshake is gone the trigger fires again. #[test] -fn poll_rekey_does_not_abandon_a_fresh_or_locally_initiated_handshake() { +fn poll_rekey_retires_an_expired_handshake_this_node_initiated_so_the_trigger_fires_again() { let fsp = Fsp::new(); - // Still inside the handshake timeout: the peer's msg3 may be in flight. + let addr = make_node_addr(9); + + let mut s = unanswered_initiation(9); + assert_eq!( + fsp.poll_rekey(vec![unanswered_initiation(9)], &cfg(100, 1000)), + vec![FspAction::ExpireInitiation { addr }], + "the unanswered handshake must be retired, and the trigger must not \ + fire in the same tick" + ); + + // What the executor does for ExpireInitiation: `abandon_handshake` drops + // the handshake and leaves `rekey_initiator` as it was. + s.rekey_in_progress = false; + assert_eq!( + fsp.poll_rekey(vec![s], &cfg(100, 1000)), + vec![FspAction::InitiateRekey { addr }], + "with the handshake retired, the trigger must start a fresh rekey" + ); +} + +/// A handshake is retired only on the deadline of the side that armed it, +/// and nothing that is not an armed handshake is retired at all. +#[test] +fn poll_rekey_retires_an_initiated_handshake_only_on_its_own_deadline() { + let fsp = Fsp::new(); + let mut fresh = session_snapshot(9); fresh.rekey_in_progress = true; - fresh.armed_handshake_expired = false; - assert!(fsp.poll_rekey(vec![fresh], &cfg(100, 1000)).is_empty()); + fresh.is_rekey_initiator = true; + assert!( + fsp.poll_rekey(vec![fresh], &cfg(100, 1000)).is_empty(), + "a handshake this node armed inside the timeout may still be answered" + ); - // Expired, but this side is the initiator: the abandon-on-timeout rule is - // anchored on the peer's setup message and does not reach our own cycle, - // which the msg3 retransmission budget bounds instead. - let mut ours = session_snapshot(9); - ours.rekey_in_progress = true; - ours.is_rekey_initiator = true; - ours.armed_handshake_expired = true; - assert!(fsp.poll_rekey(vec![ours], &cfg(100, 1000)).is_empty()); + let mut peer_older = session_snapshot(9); + peer_older.rekey_in_progress = true; + peer_older.is_rekey_initiator = true; + peer_older.armed_handshake_expired = true; + assert!( + fsp.poll_rekey(vec![peer_older], &cfg(100, 1000)).is_empty(), + "the peer's older rekey stamp must not retire a fresh handshake this \ + node armed" + ); + + let mut theirs = session_snapshot(9); + theirs.rekey_in_progress = true; + theirs.initiation_expired = true; + assert!( + fsp.poll_rekey(vec![theirs], &cfg(100, 1000)).is_empty(), + "this node's stale initiation stamp must not retire a handshake the \ + peer armed" + ); + + let mut responder_fresh = session_snapshot(9); + responder_fresh.rekey_in_progress = true; + assert!( + fsp.poll_rekey(vec![responder_fresh], &cfg(100, 1000)) + .is_empty(), + "a handshake the peer armed inside the timeout may still get its msg3" + ); - // Expired stamp but no handshake left to abandon. let mut none = session_snapshot(9); none.armed_handshake_expired = true; - assert!(fsp.poll_rekey(vec![none], &cfg(100, 1000)).is_empty()); + none.initiation_expired = true; + assert!( + fsp.poll_rekey(vec![none], &cfg(100, 1000)).is_empty(), + "expired stamps with no handshake armed leave nothing to retire" + ); + + let mut completed = session_snapshot(9); + completed.is_rekey_initiator = true; + completed.has_pending = true; + completed.initiation_expired = true; + assert!( + fsp.poll_rekey(vec![completed], &cfg(100, 1000)).is_empty(), + "a completed cycle awaiting its cutover is not an expiring handshake" + ); +} + +/// The expiry of a handshake this node armed does not depend on whether +/// periodic rekey is enabled, and it is grouped with the other abandoned +/// handshakes in input order. +#[test] +fn poll_rekey_retires_an_expired_initiation_with_periodic_rekey_off() { + let fsp = Fsp::new(); + assert_eq!( + fsp.poll_rekey(vec![unanswered_initiation(9)], &cfg_rekey_off(100, 1000)), + vec![FspAction::ExpireInitiation { + addr: make_node_addr(9) + }], + "an existing handshake must be retired whatever the trigger policy" + ); + + let mut peer_armed = session_snapshot(7); + peer_armed.rekey_in_progress = true; + peer_armed.armed_handshake_expired = true; + assert_eq!( + fsp.poll_rekey( + vec![peer_armed, unanswered_initiation(9)], + &cfg_rekey_off(100, 1000) + ), + vec![ + FspAction::AbandonHandshake { + addr: make_node_addr(7) + }, + FspAction::ExpireInitiation { + addr: make_node_addr(9) + }, + ], + "both expiries sit in the abandon group, in input order" + ); } #[test] @@ -441,6 +550,85 @@ fn poll_msg3_abandons_first() { ); } +// ===== poll_initial_msg3_resends ===== + +/// Build an initial-handshake msg3 resend snapshot for the decision under test. +fn initial_msg3_snapshot( + addr_byte: u8, + resend_count: u32, + resend_due: bool, +) -> InitialMsg3ResendSnapshot { + InitialMsg3ResendSnapshot { + addr: make_node_addr(addr_byte), + resend_count, + resend_due, + } +} + +/// A candidate that is not yet due gets nothing, whether or not its budget is +/// spent. +#[test] +fn poll_initial_msg3_not_due_is_noop() { + let fsp = Fsp::new(); + assert!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(1, 0, false)], 3) + .is_empty() + ); + assert!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(1, 99, false)], 3) + .is_empty() + ); +} + +/// A due candidate within its budget is resent. +#[test] +fn poll_initial_msg3_resends_when_due_in_budget() { + let fsp = Fsp::new(); + assert_eq!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(2, 1, true)], 3), + vec![FspAction::ResendInitialMsg3 { + addr: make_node_addr(2) + }] + ); +} + +/// A due candidate whose budget is spent is released, not resent. +#[test] +fn poll_initial_msg3_releases_when_due_at_budget() { + let fsp = Fsp::new(); + assert_eq!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(3, 3, true)], 3), + vec![FspAction::ReleaseInitialMsg3 { + addr: make_node_addr(3) + }] + ); +} + +/// Releases are returned before resends. +#[test] +fn poll_initial_msg3_releases_before_resends() { + let fsp = Fsp::new(); + let actions = fsp.poll_initial_msg3_resends( + vec![ + initial_msg3_snapshot(1, 0, true), + initial_msg3_snapshot(2, 5, true), + ], + 3, + ); + assert_eq!( + actions, + vec![ + FspAction::ReleaseInitialMsg3 { + addr: make_node_addr(2) + }, + FspAction::ResendInitialMsg3 { + addr: make_node_addr(1) + }, + ], + "releases are grouped before resends" + ); +} + // ===== classify_epoch ===== #[test] diff --git a/testing/static/scripts/rekey-test.sh b/testing/static/scripts/rekey-test.sh index 2ca46a75..2d2030f1 100755 --- a/testing/static/scripts/rekey-test.sh +++ b/testing/static/scripts/rekey-test.sh @@ -494,10 +494,10 @@ echo "" # the log count is monotone over a container's lifetime, so a released # wait means the later count cannot lose the timing race. # -# What the count measures: the counted line is emitted by the cadence -# path on whichever side(s) flip first, so a completed link rekey yields -# one or two lines (a side promoted by receiving a flipped-K frame logs -# only a DEBUG on a target the test's log filter suppresses). With rekey +# What the count measures: the initiator logs the counted line once per +# completed link rekey; the responder is promoted by the initiator's first +# new-epoch frame and logs only a DEBUG on a target the test's log filter +# suppresses. With rekey # timers resetting at each cutover, the events landing inside this # test's window are first-cycle cutovers spread across the topology's # links, not a second cycle on one link.