Merge maint into master, carrying the rekey and session-setup fixes and the unreleased changelog rewrite

This commit is contained in:
Johnathan Corgan
2026-09-22 21:45:08 +00:00
19 changed files with 2124 additions and 233 deletions
+151 -84
View File
@@ -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)
<!-- markdownlint-configure-file { "MD024": { "siblings_only": true } } -->
+7 -4
View File
@@ -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)).
+14 -10
View File
@@ -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,
+110 -10
View File
@@ -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!(
+28 -10
View File
@@ -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);
+95
View File
@@ -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<InitialMsg3ResendSnapshot> {
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.
+44 -7
View File
@@ -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<Vec<u8>>,
/// 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<u8>, 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<HandshakeState> {
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
+8
View File
@@ -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,
+1029 -52
View File
File diff suppressed because it is too large Load Diff
+176 -3
View File
@@ -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<RekeyRole>,
/// When the pending session was installed, for the responder hold.
pending_since: Option<Instant>,
// === 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<RekeyRole> {
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<SessionIndex> {
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)));
}
}
+4 -1
View File
@@ -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,
}
}
}
+53 -14
View File
@@ -131,6 +131,19 @@ pub(crate) struct ConnSnapshot {
pub msg1: Vec<u8>,
}
/// 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<RekeyRole>,
/// 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<PeerSnapshot>, cfg: &RekeyCfg) -> Vec<ConnAction> {
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
}
+1 -1
View File
@@ -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;
+101 -1
View File
@@ -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))
);
}
+2
View File
@@ -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,
}
}
+86 -11
View File
@@ -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<SessionSnapshot>,
@@ -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<InitialMsg3ResendSnapshot>,
max_resends: u32,
) -> Vec<FspAction> {
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.
+5 -3
View File
@@ -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,
+206 -18
View File
@@ -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]
+4 -4
View File
@@ -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.