From 0ebd1b44c01f6c249a7a1f22fc0898443d5c32b6 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Mon, 13 Jul 2026 14:51:11 +0000 Subject: [PATCH] node: add the initiator-cutover and drain executor actions (unwired) Prepare the establish executor to drive FMP rekey by activating the SwapSendState action (initiator K-bit cutover via cutover_to_new_session, with the gated decrypt-worker re-registration) and adding a CompleteDrain action (erase the drained previous session: free its index, drop its peers_by_index entry, unregister its decrypt session). Both reproduce the current inline rekey.rs cutover and drain bodies. The machine's drain mapping now emits CompleteDrain, using the real drained index rather than a shadow copy. Unwired: nothing drives the machine's rekey path yet (live rekey still runs inline), so these arms are unreachable and the change is behavior-neutral. The cadence fold that routes cutover and drain through the machine follows. --- src/node/dataplane/peer_actions.rs | 71 ++++++++++++++++++++++++++++-- src/peer/machine.rs | 33 +++++++++++--- 2 files changed, 95 insertions(+), 9 deletions(-) diff --git a/src/node/dataplane/peer_actions.rs b/src/node/dataplane/peer_actions.rs index 65097da..b879ce1 100644 --- a/src/node/dataplane/peer_actions.rs +++ b/src/node/dataplane/peer_actions.rs @@ -28,7 +28,7 @@ use crate::transport::{LinkId, TransportAddr, TransportId}; use crate::utils::index::SessionIndex; use crate::{NodeAddr, PeerIdentity}; use std::collections::VecDeque; -use tracing::warn; +use tracing::{debug, trace, warn}; /// Ambient shell facts a [`PeerAction`] executor needs that the machine's /// runtime-agnostic action payloads deliberately omit (verified identity, @@ -323,8 +323,73 @@ impl Node { } } PeerAction::SwapSendState { .. } => { - // C4: initiator cutover (`active.rs:1033` - // `cutover_to_new_session`). + // C4-0: initiator cutover. Reproduces the `ConnAction::Cutover` + // body in `handlers/rekey.rs:53-88` EXACTLY. `addr` is resolved + // from the ambient verified identity (as `InvalidateSendState` + // does). The decrypt re-register folds HERE, gated on + // `did_cutover` — the generic `RegisterDecryptSession` arm stays a + // no-op so a promote never double-registers. Shadow-only until the + // cadence fold routes here (C4-1). + let node_addr = *ambient.verified_identity.node_addr(); + let did_cutover = if let Some(peer) = self.peers.get_mut(&node_addr) { + if let Some(_old_our_index) = peer.cutover_to_new_session() { + // New index was pre-registered in peers_by_index + // during msg2 handling (handshake.rs). + debug_assert!( + peer.transport_id().is_some() + && peer.our_index().is_some() + && self.peers_by_index.contains_key(&( + peer.transport_id().unwrap(), + peer.our_index().unwrap().as_u32() + )), + "peers_by_index should contain pre-registered new index after cutover" + ); + debug!( + peer = %self.peer_display_name(&node_addr), + "Rekey cutover complete (initiator), K-bit flipped" + ); + true + } else { + false + } + } else { + false + }; + // Re-register the new session with the decrypt worker — the + // cache_key (transport_id, our_index) just changed, so the + // old worker entry is stale and every packet on the new + // session would miss the worker's HashMap lookup. + #[cfg(unix)] + if did_cutover { + self.register_decrypt_worker_session(&node_addr); + } + #[cfg(not(unix))] + let _ = did_cutover; + } + PeerAction::CompleteDrain { peer: node_addr } => { + // C4-0: initiator drain completion. Reproduces the + // `ConnAction::Drain` body in `handlers/rekey.rs:90-111` EXACTLY. + // Extract the real previous index + transport_id under the peer + // borrow, drop the borrow, then run the cache_key cleanup (which + // takes &mut self for unregister_decrypt_worker_session). + // Shadow-only until the cadence fold routes here (C4-1). + let drained = self.peers.get_mut(&node_addr).and_then(|peer| { + peer.complete_drain().map(|idx| (idx, peer.transport_id())) + }); + if let Some((old_our_index, transport_id)) = drained { + if let Some(tid) = transport_id { + let cache_key = (tid, old_our_index.as_u32()); + self.peers_by_index.remove(&cache_key); + #[cfg(unix)] + self.unregister_decrypt_worker_session(cache_key); + } + let _ = self.index_allocator.free(old_our_index); + trace!( + peer = %self.peer_display_name(&node_addr), + old_index = %old_our_index, + "Drain complete, previous session erased" + ); + } } PeerAction::InvalidateSendState => { // GAP-4 (biggest): the FULL teardown. `remove_active_peer` diff --git a/src/peer/machine.rs b/src/peer/machine.rs index 2b5e885..61b358a 100644 --- a/src/peer/machine.rs +++ b/src/peer/machine.rs @@ -268,6 +268,11 @@ pub(crate) enum PeerAction { /// Initiator-side rekey cutover: swap the published send-state to the pending /// epoch. SwapSendState { epoch: [u8; 8] }, + /// Complete an initiator-side rekey drain: retire the previous session slot + /// (drop its `peers_by_index`/decrypt-worker entry, free its index). The + /// executor reads the REAL previous index from `ActivePeer::complete_drain` + /// (not a machine-shadow index, which can drift). + CompleteDrain { peer: NodeAddr }, /// Invalidate the published send-state (close/loss). InvalidateSendState, /// Register a decrypt-worker entry for `index`. @@ -829,13 +834,12 @@ impl PeerMachine { actions } ConnAction::Drain { peer } => { + // The executor reads the real previous_our_index from + // `ActivePeer::complete_drain` and does the peers_by_index / + // decrypt-worker / index-free cleanup, replacing the old + // shadow-index emission (which could drift from the real index). self.state = PeerState::Active { addr: peer }; - let mut actions = Vec::new(); - if let Some(idx) = self.draining_index.take() { - actions.push(PeerAction::UnregisterDecryptSession { index: idx }); - actions.push(PeerAction::FreeIndex { index: idx }); - } - actions + vec![PeerAction::CompleteDrain { peer }] } ConnAction::InitiateRekey { peer } => { // Fresh outbound rekey: allocate our new index, send msg1 (Noise @@ -1264,6 +1268,23 @@ mod tests { kind: MaintainKind::Rekey(RekeyPhase::Draining) } ); + + // A second cadence tick from the (expired) drain window completes the + // drain: the machine now emits the single `CompleteDrain` send-state + // write (executor reads the real previous index) instead of the old + // shadow-index `[UnregisterDecryptSession, FreeIndex]` pair. + let drain_actions = m.step( + PeerEvent::Timeout { + kind: TimerKind::RekeyCadence, + }, + 20_000, + &mut alloc, + ); + assert_eq!( + drain_actions, + vec![PeerAction::CompleteDrain { peer: addr }] + ); + assert_eq!(m.state(), PeerState::Active { addr }); } // ---- Test 2: responder cutover (data-plane owned) ---------------------