node: drive inbound establish decision and net-new promote through the peer machine

Route XX inbound establishment through the per-peer machine: handle_msg3 hands
the establish decision to a transient machine and drives the net-new promote
(Promote + RestartThenPromote) through the executor's PromoteToActive action,
which transcribes the shared promote block (promote_connection plus the
msg2/tree-announce/bloom follow-ups). The four non-promote decisions (Reject,
ResendMsg2, CrossConnect, RekeyRespond) stay inline byte-for-byte; the two
session-swap arms are deferred to a later non-neutral step.

The decrypt-worker register stays inside promote_connection (not relocated),
no promotion-resolved feedback is fed back, and the transient machine is never
inserted into peer_machines (the persistent machine is created once inside
promote_connection). Behavior matches the inline path exactly, and operator
diagnostics are preserved: the relocated promoting-inbound log keeps its
handshake tracing target, and the epoch-mismatch restart breadcrumb is
re-emitted.
This commit is contained in:
Johnathan Corgan
2026-07-13 22:26:19 +00:00
parent 09d0d41639
commit cde1af2d1b
2 changed files with 226 additions and 166 deletions
+181 -36
View File
@@ -26,29 +26,34 @@
//! - `ReportLost` → `note_link_dead` (the reconciler loss reflex).
//! - `UnregisterDecryptSession` / `FreeIndex` index-plane cleanups.
//! - `RegisterDecryptSession` is a deliberate no-op — the decrypt-worker register
//! relocates into the driven register sites (`SwapSendState` here for the rekey
//! cutover; `PromoteToActive`'s deferred body for the establish promote).
//! for the rekey cutover relocates into the driven `SwapSendState` site here;
//! for the establish promote it stays INSIDE `promote_connection` (unlike the
//! IK lineage), so `PromoteToActive` does NOT re-register.
//!
//! Actions the M2 skeleton does not yet exercise are inert stubs carrying the
//! step that realizes them: the establish send-path (`SendHandshake` framing) and
//! the establish promote (`PromoteToActive`) both land when the inbound establish
//! path is wired; `SendRekey`/`SendLinkMessage`, the timers
//! Actions not yet exercised are inert stubs carrying the step that realizes
//! them: the establish send-path (`SendHandshake` framing) lands when the inbound
//! establish send is wired; `SendRekey`/`SendLinkMessage`, the timers
//! (`SetTimer`/`CancelTimer`), the outbound dial (`OpenTransport`), and the
//! connected-UDP plane are inert here.
//!
//! Three XX-semantic arms are explicit **deferred stubs** whose real bodies land
//! when the inbound establish path is wired: `PromoteToActive`,
//! `SwapToInboundSession`, and `RekeyRespondTrigger`. Each carries a
//! `debug_assert!(false, …)` guard + a benign no-op — they are unreachable now
//! (nothing drives the machine), so this is safe shadow.
//! `PromoteToActive` is LIVE: it drives the inbound establish promote
//! (`promote_connection` + the msg2/tree-announce/bloom follow-ups), transcribing
//! `handle_msg3`'s shared promote block. Two XX-semantic arms remain explicit
//! **deferred stubs** whose real bodies land when the inbound cross-connection /
//! rekey-responder paths are wired: `SwapToInboundSession` and
//! `RekeyRespondTrigger`. Each carries a `debug_assert!(false, …)` guard + a
//! benign no-op — they are unreachable now (those decisions stay inline in
//! `handle_msg3`).
use crate::node::Node;
use crate::node::reject::{HandshakeReject, RejectReason};
use crate::node::{Node, NodeError};
use crate::peer::machine::{PeerAction, PeerEvent};
use crate::proto::fmp::PromotionResult;
use crate::transport::{LinkId, TransportAddr, TransportId};
use crate::utils::index::SessionIndex;
use crate::{NodeAddr, PeerIdentity};
use std::collections::VecDeque;
use tracing::{info, trace};
use tracing::{debug, info, trace, warn};
/// Ambient shell facts a [`PeerAction`] executor needs that the machine's
/// runtime-agnostic action payloads deliberately omit (verified identity,
@@ -74,10 +79,11 @@ pub(in crate::node) struct PeerActionCtx {
pub(in crate::node) their_index: Option<SessionIndex>,
/// The wire timestamp driving this step (promotion ts / loss-report clock).
pub(in crate::node) now_ms: u64,
/// Establish direction for this exchange. Discriminates the `PromoteToActive`
/// failure cleanup (the inbound and outbound promote-Err arms differ). `false`
/// = inbound, `true` = outbound. Consumed by the deferred `PromoteToActive`
/// body.
/// Establish direction for this exchange. `false` = inbound, `true` =
/// outbound. The live `PromoteToActive` body drives the inbound establish
/// promote only (msg3 is always responder side), so it does not read this
/// yet; retained for the future outbound promote path, whose promote-Err
/// cleanup differs from the inbound one.
pub(in crate::node) is_outbound: bool,
}
@@ -144,17 +150,156 @@ impl Node {
// Encrypt + send a link-control frame (heartbeat / filter /
// tree / disconnect). Data-plane-owned; out of M2 scope.
}
PeerAction::PromoteToActive { .. } => {
// DEFERRED: crystallize identity + `promote_connection` +
// relocated decrypt-worker register + `PromotionResolved`
// feedback + cross-connection loser-link surgery. This is
// XX-semantic (identity at msg3) and is written, wired, and
// ci-local-validated when the inbound establish path is wired.
// Unreachable now (nothing drives the machine).
debug_assert!(
false,
"PromoteToActive executor body lands when inbound establish is wired"
PeerAction::PromoteToActive { link: promote_link } => {
// Establish promote, driven through the machine. Transcribes
// `handle_msg3`'s shared inbound promote block verbatim, adapted
// to the executor's ambient context. Two XX-specific choices vs
// the IK-lineage executor: (1) the decrypt-worker register stays
// INSIDE `promote_connection` (NOT relocated here) — re-
// registering would double-register; (2) NO `PromotionResolved`
// is fed back — the persistent machine is born `established()` by
// `promote_connection`, so feeding it would perturb that ctor
// state that the rekey/reap folds read.
// Capture msg2 BEFORE `promote_connection` removes the pending
// connection, so a duplicate msg1 can be answered with it.
let wire_msg2 = self
.connections
.get(&promote_link)
.and_then(|c| c.handshake_msg2().map(|m| m.to_vec()));
debug!(
// Relocated from `handlers/handshake.rs` (M6): pin the target
// so it stays visible under the harness's
// `fips::node::handlers::handshake=debug` filter.
target: "fips::node::handlers::handshake",
peer = %self.peer_display_name(ambient.verified_identity.node_addr()),
link_id = %promote_link,
our_index = ?ambient.our_index,
"handle_msg3: promoting inbound, peers_has_key={}",
self.peers.contains_key(ambient.verified_identity.node_addr()),
);
match self.promote_connection(
promote_link,
ambient.verified_identity,
ambient.now_ms,
) {
Ok(PromotionResult::Promoted(node_addr)) => {
// Store msg2 on peer for resend on duplicate msg1
if let (Some(peer), Some(msg2)) =
(self.peers.get_mut(&node_addr), wire_msg2)
{
peer.set_handshake_msg2(msg2);
}
// Promotion is logged once by `promote_connection`
// ("Connection promoted to active peer"); no separate
// inbound-path line.
// Send initial tree announce to new peer
if let Err(e) = self.send_tree_announce_to_peer(&node_addr).await {
debug!(peer = %self.peer_display_name(&node_addr), error = %e, "Failed to send initial TreeAnnounce");
}
// Schedule filter announce (sent on next tick via debounce)
self.bloom_state.mark_update_needed(node_addr);
self.reset_lookup_backoff();
}
Ok(PromotionResult::CrossConnectionWon {
loser_link_id,
node_addr,
}) => {
// UNREACHABLE on driven XX establish paths: `Promote`
// and `RestartThenPromote` (which removes the old peer
// first) both imply no existing peer at promote time, so
// `promote_connection` returns `Promoted`. Body kept
// byte-equivalent to next so a future path that drives a
// cross-connection through the executor trips the assert.
debug_assert!(
false,
"executor CrossConnectionWon is unreachable on driven \
XX inbound establish paths"
);
// Store msg2 on peer for resend on duplicate msg1
if let (Some(peer), Some(msg2)) =
(self.peers.get_mut(&node_addr), wire_msg2)
{
peer.set_handshake_msg2(msg2);
}
// Close the losing TCP connection (no-op for connectionless)
if let Some(loser_link) = self.links.get(&loser_link_id) {
let loser_tid = loser_link.transport_id();
let loser_addr = loser_link.remote_addr().clone();
if let Some(transport) = self.transports.get(&loser_tid) {
transport.close_connection(&loser_addr).await;
}
}
// Clean up the losing connection's link
self.remove_link(&loser_link_id);
debug!(
peer = %self.peer_display_name(&node_addr),
loser_link_id = %loser_link_id,
"Inbound cross-connection won, loser link cleaned up"
);
if let Err(e) = self.send_tree_announce_to_peer(&node_addr).await {
debug!(peer = %self.peer_display_name(&node_addr), error = %e, "Failed to send initial TreeAnnounce");
}
self.bloom_state.mark_update_needed(node_addr);
self.reset_lookup_backoff();
}
Ok(PromotionResult::CrossConnectionLost { winner_link_id }) => {
// UNREACHABLE on driven XX establish paths (see the Won
// arm). Body kept byte-equivalent to next; uses the
// ambient transport/addr in place of next's `packet.*`.
debug_assert!(
false,
"executor CrossConnectionLost is unreachable on driven \
XX inbound establish paths"
);
// Close the losing TCP connection (no-op for connectionless)
if let Some(transport) = self.transports.get(&ambient.transport_id) {
transport.close_connection(&ambient.remote_addr).await;
}
// This connection lost — clean up its link
self.remove_link(&promote_link);
// Restore addr_to_link for the winner's link
self.addr_to_link.insert(
(ambient.transport_id, ambient.remote_addr.clone()),
winner_link_id,
);
debug!(
winner_link_id = %winner_link_id,
"Inbound cross-connection lost, keeping existing"
);
}
Err(e) => {
// A max_peers rejection is expected policy, not a fault —
// log it at debug to avoid WARN spam when a cap'd node is
// under sustained inbound pressure. Other promotion
// failures remain at warn.
if matches!(e, NodeError::MaxPeersExceeded { .. }) {
debug!(
peer = %self.peer_display_name(ambient.verified_identity.node_addr()),
max = self.max_peers(),
"Rejecting inbound connection at max_peers cap (no promotion)"
);
} else {
warn!(
link_id = %promote_link,
error = %e,
"Failed to promote inbound connection"
);
}
// Clean up on promotion failure. promote_connection
// already freed our_index in its MaxPeersExceeded path;
// freeing again here is benign (IndexAllocator::free is a
// HashSet::remove, the second call returns Err(NotFound)
// and is ignored).
self.remove_link(&promote_link);
if let Some(idx) = ambient.our_index {
let _ = self.index_allocator.free(idx);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
}
}
}
PeerAction::SwapSendState { .. } => {
// Initiator cutover. Reproduces the `ConnAction::Cutover` body
@@ -278,15 +423,15 @@ impl Node {
}
PeerAction::RegisterDecryptSession { index } => {
let _ = index;
// No-op by design. The decrypt-worker registration relocates
// into the driven register sites: `SwapSendState` above (gated
// on `did_cutover`) for the rekey cutover, and the deferred
// `PromoteToActive` body for the establish promote. This
// machine-emitted action is redundant with those; kept as an
// inert no-op (rather than removing the emission) so the
// machine's action sequence and unit tests stay unchanged. The
// keyed-by-NodeAddr register does not need the machine's `index`
// payload.
// No-op by design. The rekey-cutover decrypt-worker register
// relocates into the driven `SwapSendState` site above (gated
// on `did_cutover`); the establish-promote register stays INSIDE
// `promote_connection`, so `PromoteToActive` does not re-register
// either. This machine-emitted action is redundant with both;
// kept as an inert no-op (rather than removing the emission) so
// the machine's action sequence and unit tests stay unchanged.
// The keyed-by-NodeAddr register does not need the machine's
// `index` payload.
}
PeerAction::UnregisterDecryptSession { index } => {
// Executor supplies `transport_id` from ambient; keyed by
+45 -130
View File
@@ -7,9 +7,10 @@
use crate::PeerIdentity;
use crate::node::acl::PeerAclContext;
use crate::node::dataplane::PeerActionCtx;
use crate::node::reject::{HandshakeReject, RejectReason};
use crate::node::{Node, NodeError};
use crate::peer::machine::PeerMachine;
use crate::peer::machine::{PeerEvent, PeerMachine};
use crate::peer::{ActivePeer, PeerConnection};
use crate::proto::fmp::wire::{Msg1Header, Msg2Header, Msg3Header, build_msg2, build_msg3};
use crate::proto::fmp::{
@@ -874,6 +875,14 @@ impl Node {
/// learn the initiator's identity and epoch, then performs identity-based
/// checks (restart detection, rekey detection, cross-connection resolution)
/// and promotes the connection to active peer.
//
// The four terminal `establish_inbound` arms (Reject/ResendMsg2/CrossConnect/
// RekeyRespond) keep the pre-refactor trailing `return;`, preserved byte-for-
// byte from before the promote block moved into the executor. With that block
// gone the `match` is now the function tail, so those returns read as
// `needless_return`; the allow keeps the arms untouched rather than editing
// preserved shell bodies to satisfy the lint.
#[allow(clippy::needless_return)]
pub(in crate::node) async fn handle_msg3(&mut self, packet: ReceivedPacket) {
// Parse header
let header = match Msg3Header::parse(&packet.data) {
@@ -1298,138 +1307,44 @@ impl Node {
);
return;
}
InboundDecision::RestartThenPromote { peer } => {
// Epoch mismatch — peer restarted. Tear down stale session, then
// fall through to promote the fresh connection in its place.
debug!(
peer = %self.peer_display_name(&peer),
"Peer restart detected (epoch mismatch), removing stale session"
);
self.remove_active_peer(&peer);
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
self.note_link_dead(peer, now_ms);
// Fall through to process as new connection.
}
InboundDecision::Promote => {
// Net-new inbound (or the post-restart re-promote): fall through
// to promote_connection, whose late max-peers cap and
// cross-connection won/lost handling stay shell-side.
}
}
// Promote the connection to active peer.
let wire_msg2 = self
.connections
.get(&link_id)
.and_then(|c| c.handshake_msg2().map(|m| m.to_vec()));
debug!(
peer = %self.peer_display_name(&peer_node_addr),
link_id = %link_id,
our_index = %our_index,
"handle_msg3: promoting inbound, peers_has_key={}",
self.peers.contains_key(&peer_node_addr),
);
match self.promote_connection(link_id, peer_identity, packet.timestamp_ms) {
Ok(result) => {
match result {
PromotionResult::Promoted(node_addr) => {
// Store msg2 on peer for resend on duplicate msg1
if let (Some(peer), Some(msg2)) =
(self.peers.get_mut(&node_addr), wire_msg2)
{
peer.set_handshake_msg2(msg2);
}
// Promotion is logged once by `promote_connection`
// ("Connection promoted to active peer"); no separate
// inbound-path line.
// Send initial tree announce to new peer
if let Err(e) = self.send_tree_announce_to_peer(&node_addr).await {
debug!(peer = %self.peer_display_name(&node_addr), error = %e, "Failed to send initial TreeAnnounce");
}
// Schedule filter announce (sent on next tick via debounce)
self.bloom_state.mark_update_needed(node_addr);
self.reset_lookup_backoff();
}
PromotionResult::CrossConnectionWon {
loser_link_id,
node_addr,
} => {
// Store msg2 on peer for resend on duplicate msg1
if let (Some(peer), Some(msg2)) =
(self.peers.get_mut(&node_addr), wire_msg2)
{
peer.set_handshake_msg2(msg2);
}
// Close the losing TCP connection (no-op for connectionless)
if let Some(loser_link) = self.links.get(&loser_link_id) {
let loser_tid = loser_link.transport_id();
let loser_addr = loser_link.remote_addr().clone();
if let Some(transport) = self.transports.get(&loser_tid) {
transport.close_connection(&loser_addr).await;
}
}
// Clean up the losing connection's link
self.remove_link(&loser_link_id);
debug!(
peer = %self.peer_display_name(&node_addr),
loser_link_id = %loser_link_id,
"Inbound cross-connection won, loser link cleaned up"
);
if let Err(e) = self.send_tree_announce_to_peer(&node_addr).await {
debug!(peer = %self.peer_display_name(&node_addr), error = %e, "Failed to send initial TreeAnnounce");
}
self.bloom_state.mark_update_needed(node_addr);
self.reset_lookup_backoff();
}
PromotionResult::CrossConnectionLost { winner_link_id } => {
// Close the losing TCP connection (no-op for connectionless)
if let Some(transport) = self.transports.get(&packet.transport_id) {
transport.close_connection(&packet.remote_addr).await;
}
// This connection lost — clean up its link
self.remove_link(&link_id);
// Restore addr_to_link for the winner's link
self.addr_to_link.insert(
(packet.transport_id, packet.remote_addr.clone()),
winner_link_id,
);
debug!(
winner_link_id = %winner_link_id,
"Inbound cross-connection lost, keeping existing"
);
}
}
}
Err(e) => {
// A max_peers rejection is expected policy, not a fault —
// log it at debug to avoid WARN spam when a cap'd node is
// under sustained inbound pressure. Other promotion failures
// remain at warn.
if matches!(e, NodeError::MaxPeersExceeded { .. }) {
decision @ (InboundDecision::RestartThenPromote { .. } | InboundDecision::Promote) => {
// Preserve next's epoch-mismatch restart breadcrumb — fires before
// the machine's teardown actions run, matching next's inline order
// (breadcrumb → remove_active_peer → note_link_dead → promote).
if let InboundDecision::RestartThenPromote { peer } = &decision {
debug!(
peer = %self.peer_display_name(&peer_node_addr),
max = self.max_peers(),
"Rejecting inbound connection at max_peers cap (no promotion)"
);
} else {
warn!(
link_id = %link_id,
error = %e,
"Failed to promote inbound connection"
peer = %self.peer_display_name(peer),
"Peer restart detected (epoch mismatch), removing stale session"
);
}
// Clean up on promotion failure. promote_connection already
// freed our_index in its MaxPeersExceeded path; freeing again
// here is benign (IndexAllocator::free is a HashSet::remove,
// the second call returns Err(NotFound) and is ignored).
self.remove_link(&link_id);
let _ = self.index_allocator.free(our_index);
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
// Machine-driven establish promote. A TRANSIENT inbound machine
// re-derives the decision and emits the action stream:
// `[PromoteToActive]` for `Promote`, or `[InvalidateSendState,
// ReportLost, PromoteToActive]` for `RestartThenPromote` (whose two
// teardown actions map to `remove_active_peer` / `note_link_dead`,
// in that order — the same order next ran them inline). The transient
// is never inserted into `peer_machines`; the persistent
// `established()` machine is created inside `promote_connection`
// (M3), so there is no double-insert. The relocated promote body,
// and the `RestartThenPromote` teardown ordering equivalence, live in
// the executor's `PromoteToActive` / `InvalidateSendState` /
// `ReportLost` arms.
let mut machine = PeerMachine::new_inbound(link_id, packet.timestamp_ms);
let actions = machine.step(
PeerEvent::InboundMsg3 { wire, est: snap },
packet.timestamp_ms,
&mut self.index_allocator,
);
let ambient = PeerActionCtx {
verified_identity: peer_identity,
transport_id: packet.transport_id,
remote_addr: packet.remote_addr.clone(),
our_index: Some(our_index),
their_index: Some(header.sender_idx),
now_ms: packet.timestamp_ms,
is_outbound: false,
};
self.execute_peer_actions(link_id, &ambient, actions).await;
}
}
}