Keep a link rekey responder from switching to keys the initiator never installed

The node that answers a link rekey stored its new session as pending and
cut over to it on its own next tick, before the initiator had read the
reply. When that reply was lost, the initiator abandoned the cycle and
freed the index the responder was now sealing to, so frames from the
responder were dropped at the initiator until the link was torn down.

The pending session now records which side of the handshake produced it.
Only a pending session this node initiated is cut over by the rekey tick.
One it answered is promoted when a frame on the new keys from the
initiator authenticates against it, as the data path already does, and a
held pending session also stops this node from starting a rekey of its
own that would overwrite it. If the initiator never adopts the keys, the
responder drops the pending session after a hold that outlasts the
initiator's resend ladder and cutover, its first heartbeat and the
link-dead timeout, and frees its index, so the next rekey can be
answered.

The lost-reply test now runs, and new tests cover the hold, the
retirement and the rekey that completes after it.
This commit is contained in:
Johnathan Corgan
2026-09-22 21:43:28 +00:00
parent c8a907826a
commit 5db6382788
12 changed files with 689 additions and 97 deletions
+7
View File
@@ -148,6 +148,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
+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
@@ -717,9 +717,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();
}
@@ -1225,14 +1227,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,
+74 -7
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()
}
+246 -52
View File
@@ -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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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.
+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 {
@@ -287,6 +288,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)
@@ -329,6 +335,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),
}
}
@@ -409,6 +417,8 @@ impl ActivePeer {
rekey_msg1: None,
rekey_msg1_next_resend: 0,
rekey_msg1_resend_count: 0,
pending_role: None,
pending_since: None,
send,
}
}
@@ -977,6 +987,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
@@ -1084,19 +1104,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;
@@ -1104,6 +1167,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).
@@ -1115,6 +1183,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();
@@ -1141,6 +1211,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
}
@@ -1152,6 +1227,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();
@@ -1178,6 +1255,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
}
@@ -1204,6 +1286,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.
@@ -1216,11 +1318,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) ===
@@ -1824,4 +1934,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.
@@ -295,6 +315,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
@@ -504,33 +528,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);
@@ -539,6 +577,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;
@@ -614,3 +616,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,
}
}
+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.