Release the session index and link at every XX handshake reject arm

The msg2 self-connect drop and the msg3 reject arms disposed a leg without
returning what that leg had allocated. Each site now releases the session
index and the link it holds, and captures the index before disposal, since
reading it afterwards yields None.

One arm is not like the others and must not be fixed like them. Its index
comes from the receiver field of the incoming header, which a peer supplies,
so freeing it unconditionally would let a hostile peer release an index
belonging to an unrelated live session: a memory leak traded for a remote
session teardown. That arm now frees only after establishing that the index
is not claimed elsewhere, by a predicate that scans the pending maps, the
peer machines and the active peers.

The predicate is deliberately transport-blind. Scoping it to the transport
the packet arrived on would let an orphaned entry on one transport free an
index live on another, and that regression was invisible to the whole suite
until the test added here: with the scoping applied, 1806 tests passed.

Each of the predicate's limbs is now decided by exactly one test, checked by
mutating the production code and confirming the intended test fails alone.
The outbound ACL-reject arm also regains the reschedule call its dial-gate
sibling makes, so a configured peer no longer drops off the dial schedule.
This commit is contained in:
Johnathan Corgan
2026-08-11 06:21:31 +00:00
parent c7e6d00333
commit f5b01370fa
6 changed files with 1308 additions and 13 deletions
+113 -2
View File
@@ -933,6 +933,35 @@ impl Node {
// and its pending connection is dropped with it.
self.remove_peer_machine(link_id);
self.remove_link(&link_id);
// `our_index` is the index WE allocated at msg1 preparation, read
// back off the machine above before any disposal — never the
// `receiver_idx` the msg2 header supplied. Note what makes that
// true: `link_id` came from `pending_outbound[(tid,
// header.receiver_idx)]`, an attacker-chosen KEY, and on the rekey
// path that same map names a PROMOTED peer's link, whose machine
// carries the peer's live session index rather than a rekey index.
// What keeps this arm off that state is the leg gate above
// (`peer_machines.get(&link_id).is_none_or(|m| m.leg().is_none())`):
// every sub-branch under it returns, so only a leg with a live
// handshake carrier — a fresh dial — reaches here. If that gate is
// ever relaxed, this free returns a live peer's current index to
// the pool.
if let Some(idx) = our_index {
let _ = self.index_allocator.free(idx);
}
// Put the dial back on the retry schedule. The disposal above takes
// this leg out of both reapers, so the stuck-leg sweep that normally
// reaches `note_handshake_timeout` never runs for it, and that
// reflex is the only thing that seeds `retry_pending` for a
// configured peer.
//
// Targets `dialed_peer`, the dial-time expectation, NOT
// `peer_identity` — the completion above already overwrote the
// machine's expected identity with the answering static, so reading
// that here would reschedule against whoever answered.
if let Some(peer) = dialed_peer {
self.note_handshake_timeout(peer, packet.timestamp_ms);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
return;
@@ -946,10 +975,26 @@ impl Node {
// arrives here — the dial-identity gate above catches it first,
// and only a dialed == learned == us leg gets this far. This leg
// never promotes; its machine goes with it (dropping the embedded
// pending connection). The index, link, and `pending_outbound`
// entry are deliberately NOT freed here (pre-existing shape).
// pending connection), and everything else the leg holds — the
// index, the link, and the `pending_outbound` entry — goes with it.
//
// No reschedule fires, unlike the dial-identity gate above: the dial
// named ourselves, so there is nothing to retry.
//
// Link disposal goes through `remove_link`, never a bare
// `links.remove` plus `addr_to_link.remove`. We have answered our own
// msg1, so `handle_msg1` has already overwritten
// `addr_to_link[(tid, self_addr)]` with the INBOUND leg's link;
// `remove_link` clears the reverse entry only when it still points at
// the link being removed, so it correctly leaves the live inbound
// leg's mapping alone. A hand-rolled removal would destroy it.
debug!(link_id = %link_id, "Discovered self via shared-media beacon, dropping");
self.pending_outbound.remove(&key);
self.remove_peer_machine(link_id);
self.remove_link(&link_id);
if let Some(idx) = our_index {
let _ = self.index_allocator.free(idx);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
return;
@@ -1301,10 +1346,43 @@ impl Node {
let machine = match self.peer_machines.get_mut(&link_id) {
Some(m) => m,
None => {
// The pending-inbound entry outlived its machine. `key` was
// written by us at msg1 (`pending_inbound` has one insertion
// site, and the index it carries came straight from
// `index_allocator.allocate()`), so the index named here is
// ours to reclaim — UNLESS this entry is stale enough that
// the index has already been freed and re-drawn for something
// live, in which case freeing it would tear down an unrelated
// session. No path to this arm has been found; the check
// below is unconditional defence-in-depth, not a fix for a
// demonstrated state. Leaving one index leaked on a
// should-not-happen path is strictly better than a remote
// teardown primitive, since `receiver_idx` is a wire field
// the sender chooses.
//
// The link goes regardless: it has no machine on either
// branch, its pending-inbound key is already gone, and
// `LinkId`s come from a monotonic counter
// (`allocate_link_id`) so a stale value can never name a
// different live connection. `remove_link` drops the
// `addr_to_link` reverse entry only if it still points here,
// so a newer leg on the same address is untouched.
debug!(
link_id = %link_id,
"No pending connection for msg3"
);
self.remove_link(&link_id);
let orphan_index = header.receiver_idx;
if self.session_index_is_claimed(orphan_index) {
warn!(
link_id = %link_id,
receiver_idx = %orphan_index,
"Orphaned pending-inbound entry names an index claimed by live \
state; leaving it allocated"
);
} else {
let _ = self.index_allocator.free(orphan_index);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::UnknownConnection));
return;
@@ -1345,8 +1423,19 @@ impl Node {
Ok(()) => {}
Err(e) => {
warn!(link_id = %link_id, our_profile = %our_profile, error = %e, "FMP negotiation failed");
// Capture the msg1-allocated index before disposing the
// machine that carries it; a read placed after the
// disposal yields None and the free is silently skipped.
// This arm sits ahead of the inbound ACL gate, so any
// peer able to complete a Noise msg3 reaches it without
// being authorized.
let our_idx_to_free =
self.peer_machines.get(&link_id).and_then(|m| m.our_index());
self.remove_link(&link_id);
self.remove_peer_machine(link_id);
if let Some(idx) = our_idx_to_free {
let _ = self.index_allocator.free(idx);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
return;
@@ -1359,8 +1448,18 @@ impl Node {
Some(id) => *id,
None => {
warn!("Identity not learned from msg3");
// Same capture-before-dispose shape as the negotiation arm
// above. Defensive: `complete_handshake_msg3` sets the
// identity on success, and no path reaching here with `None`
// has been constructed, so this free is by symmetry with its
// siblings rather than by test.
let our_idx_to_free =
self.peer_machines.get(&link_id).and_then(|m| m.our_index());
self.remove_link(&link_id);
self.remove_peer_machine(link_id);
if let Some(idx) = our_idx_to_free {
let _ = self.index_allocator.free(idx);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
return;
@@ -1446,6 +1545,14 @@ impl Node {
}
self.remove_link(&link_id);
self.remove_peer_machine(link_id);
// The msg1-allocated index, captured above before any disposal and
// carried across the `Disconnect` send (`SessionIndex` is `Copy`).
// The free deliberately follows the send: the reject signal is built
// from state on the machine and is this arm's user-visible
// behaviour.
if let Some(idx) = our_index {
let _ = self.index_allocator.free(idx);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
return;
@@ -1455,6 +1562,10 @@ impl Node {
debug!(link_id = %link_id, "Received msg3 from self, dropping");
self.remove_link(&link_id);
self.remove_peer_machine(link_id);
// Same msg1-allocated index, same capture point above.
if let Some(idx) = our_index {
let _ = self.index_allocator.free(idx);
}
self.stats_mut()
.record_reject(RejectReason::Handshake(HandshakeReject::BadState));
return;
+31 -2
View File
@@ -389,8 +389,37 @@ impl Node {
match action {
// Abandon rekey cycles that exhausted their retransmission budget.
ConnAction::AbandonRekey { peer: node_addr } => {
if let Some(peer) = self.peers.get_mut(&node_addr) {
peer.abandon_rekey();
// `abandon_rekey` hands back whichever index the abandoned
// cycle owned; dropping the return value orphans it, and the
// `pending_outbound` entry seeded at rekey msg1 with it. Same
// shape as the msg3 tie-break loser in `peer_actions`.
//
// The `peers_by_index` removal is shape-parity with those two
// model arms and is a no-op at THIS call site: `AbandonRekey`
// is emitted only from the msg1-resend-budget classification,
// so no rekey msg2 ever arrived and nothing was inserted.
//
// Known exposure, kept for parity with the model arms rather
// than closed here: `transport_id()` is RE-READ, while the
// `pending_outbound` entry was keyed by whatever it was when
// rekey msg1 went out. A roam in between (`set_current_addr`
// overwrites `send.transport_id`) makes the removal miss, so
// the index is freed with a stale entry still pointing at the
// peer's live link. Walked to its end: a later msg2 naming
// that index on the old transport resolves the stale link,
// finds the promoted peer's machine leg-less, finds no peer
// with a matching `rekey_our_index` (this arm cleared it),
// and takes the "not a rekey" arm, which removes the stale
// entry and records a reject. No teardown, no wrong-peer
// effect. Removing by index VALUE would close it outright.
if let Some(peer) = self.peers.get_mut(&node_addr)
&& let Some(idx) = peer.abandon_rekey()
{
if let Some(tid) = peer.transport_id() {
self.peers_by_index.remove(&(tid, idx.as_u32()));
self.pending_outbound.remove(&(tid, idx.as_u32()));
}
let _ = self.index_allocator.free(idx);
}
debug!(
peer = %self.peer_display_name(&node_addr),
+61
View File
@@ -2368,6 +2368,67 @@ impl Node {
self.peer_timers.remove(&link);
}
/// Does anything live claim this session index?
///
/// The orphaned-pending-inbound arm in `handle_msg3` is the one place that
/// must free an index it cannot read off a machine — the machine's absence
/// is what selects the arm — so the only value available is the msg3
/// header's `receiver_idx`, a field the sender chooses. The
/// `pending_inbound` hit that precedes the arm proves the index was one WE
/// allocated, but not that it is still that leg's: every free site is
/// expected to drop the map entry along with the index, and that is a
/// convention rather than an invariant. No path to that arm has been found;
/// this is unconditional defence-in-depth, not a fix for a demonstrated
/// state. If any live holder claims the index, we do not free it.
///
/// **Deliberately transport-blind.** `index_allocator` is one node-global
/// `HashSet<u32>`, so an index value is unique across every transport, while
/// every registry scanned below is keyed `(TransportId, u32)`. Filtering
/// this scan by the transport the msg3 arrived on would let an index held by
/// a peer on a DIFFERENT transport read as unclaimed — the one direction
/// that turns this guard into the teardown primitive it exists to prevent.
/// The `ActivePeer` limb must likewise not be filtered by `transport_id()`.
/// No production path builds a peer with `transport_id() == None` — both
/// `self.peers` insertion sites go through `ActivePeer::with_session`, which
/// takes a concrete `TransportId` — but `remove_active_peer` gates its index
/// frees on `if let Some(tid)`, so if one ever arose its indices would stay
/// allocated forever and are exactly what must not be handed back here.
/// Scan by index value only.
///
/// The `ActivePeer` limb also has to carry slots the maps do not: the
/// `peers_by_index` insert for a rekeying peer's `pending_our_index` is
/// itself gated on `if let Some(tid)` in `handle_msg2`'s rekey branch.
/// `rekey_responder_our_index` is listed for completeness — its only writer,
/// `set_rekey_responder_state`, has no callers today, so the slot is always
/// `None` and nothing exercises it; whoever wires the XX rekey-responder
/// path must not have to remember to add it here.
///
/// Should-not-happen-path guard, not a general-purpose lookup: it is O(live
/// sessions) and has exactly one call site.
fn session_index_is_claimed(&self, idx: crate::utils::index::SessionIndex) -> bool {
let raw = idx.as_u32();
if self.peers_by_index.keys().any(|(_, i)| *i == raw)
|| self.pending_outbound.keys().any(|(_, i)| *i == raw)
|| self.pending_inbound.keys().any(|(_, i)| *i == raw)
{
return true;
}
if self
.peer_machines
.values()
.any(|m| m.our_index() == Some(idx))
{
return true;
}
self.peers.values().any(|p| {
p.our_index() == Some(idx)
|| p.rekey_our_index() == Some(idx)
|| p.rekey_responder_our_index() == Some(idx)
|| p.pending_our_index() == Some(idx)
|| p.previous_our_index() == Some(idx)
})
}
/// Debug-build coherence sweep over the peer-lifecycle maps, run once per
/// rx-loop tick and invoked directly by unit tests.
///
+55 -7
View File
@@ -53,6 +53,20 @@ async fn test_outbound_connect_denied_by_denylist() {
// The equivalent gate on this line runs once msg3 has revealed the initiator's
// static key; `test_inbound_msg3_denied_by_acl` covers it there.
/// The outbound ACL reject at msg2 must return the dial-time session index.
///
/// The index asserted on is `our_index_a`, allocated by hand at the top of this
/// test and read back off the machine by the arm before it disposes the leg.
///
/// The arm is not identified by a stats counter: `bad_state == 1` with
/// `peer_count() == 0` is produced identically by the dial-identity gate 45
/// lines above, which also frees the index. What establishes that the ACL arm
/// is the one reached is structural rather than asserted — the seed below names
/// `peer_b_identity` as the dialed identity and the msg2 is built from node_b's
/// own static, so the dial-identity gate returns `Accept` and falls through.
/// The behavioural evidence that the assertion is attached to the right arm is
/// the mutation that deletes the ACL arm's free specifically, which reds
/// `is_allocated` here and nothing else.
#[tokio::test]
async fn test_outbound_msg2_denied_after_acl_reload() {
let (dir, mut node_a) = make_acl_node();
@@ -112,6 +126,15 @@ async fn test_outbound_msg2_denied_after_acl_reload() {
std::fs::write(deny_path(&dir), format!("{}\n", node_b.npub())).unwrap();
assert!(node_a.reload_peer_acl().await);
// Control: the leg holds exactly the index allocated above, so a run that
// never reached the arm cannot report green on the post-condition.
let baseline = node_a.index_allocator.count();
assert_eq!(baseline, 1, "the dial allocated exactly one index");
assert!(
node_a.index_allocator.is_allocated(our_index_a),
"control: the leg allocated an index"
);
let packet = ReceivedPacket::with_timestamp(transport_id, remote_addr, wire_msg2, 1100);
node_a.handle_msg2(packet).await;
@@ -119,6 +142,15 @@ async fn test_outbound_msg2_denied_after_acl_reload() {
assert_eq!(node_a.connection_count(), 0);
assert_eq!(node_a.link_count(), 0);
assert!(node_a.pending_outbound.is_empty());
assert!(
!node_a.index_allocator.is_allocated(our_index_a),
"the outbound ACL reject must return the index it allocated"
);
assert_eq!(
node_a.index_allocator.count(),
baseline - 1,
"and must free exactly one"
);
}
/// Inbound rejection at msg3 must also cut down the initiator.
@@ -342,13 +374,15 @@ async fn test_outbound_connect_not_denied_by_allowlist_miss() {
/// peering-budget slot forever, so that is what this pins, along with the
/// link and the connection count.
///
/// What it deliberately does NOT pin is session-index hygiene. This arm does
/// not return the index to the allocator today, unlike the sibling bad-state
/// arm just above it; that gap is tracked and is left exactly as it is here,
/// since this change adds tests only. Asserting on `peers_by_index` would
/// look like coverage of it and would be worthless: nothing maps an index
/// until promotion, which is downstream of this gate, so such an assertion
/// holds no matter what the denial cleans up.
/// Session-index hygiene is pinned at the ALLOCATOR, by name, and not through
/// `peers_by_index`. The registry assertion that looks like coverage here is
/// worthless and must not be reintroduced: nothing maps an index until
/// promotion, which is downstream of this gate, so
/// `assert!(!peers_by_index.contains_key(..))` holds whether the denial frees
/// the index or leaks it — its pass value and its failure value coincide.
/// `debug_assert_peer_maps_coherent()` is no better; it walks `peer_machines`
/// for carriers and says nothing about the allocator. Naming the index and
/// asserting `is_allocated` before and after is what discriminates.
#[tokio::test]
async fn test_acl_rejected_msg3_leaves_no_registry_trace() {
use crate::proto::fmp::wire::{Msg2Header, build_msg1, build_msg3};
@@ -395,6 +429,11 @@ async fn test_acl_rejected_msg3_leaves_no_registry_trace() {
.expect("msg1 stores the framed msg2 on the carrier")
.to_vec();
assert_eq!(node_b.link_count(), 1);
let baseline = node_b.index_allocator.count();
assert!(
node_b.index_allocator.is_allocated(our_index_b),
"control: msg1 allocated the responder's index"
);
// A completes on msg2 and answers msg3, which is where its static key —
// and therefore the denial — first reaches B.
@@ -424,4 +463,13 @@ async fn test_acl_rejected_msg3_leaves_no_registry_trace() {
1,
"the denial is attributed to the handshake state-machine counter"
);
assert!(
!node_b.index_allocator.is_allocated(our_index_b),
"the inbound ACL reject must return the index msg1 allocated"
);
assert_eq!(
node_b.index_allocator.count(),
baseline - 1,
"and must free exactly one"
);
}
File diff suppressed because it is too large Load Diff
+4 -1
View File
@@ -70,10 +70,13 @@ impl std::fmt::Display for SessionIndex {
}
}
/// Allocator for session indices within a single transport.
/// Allocator for session indices, node-global across every transport.
///
/// Manages a pool of random 32-bit indices, tracking which are in use
/// to prevent collision. Thread-safe for single-threaded async use.
///
/// There is exactly one of these per node and its whole state is the `in_use`
/// set below, so an index value is unique across the node, not per transport.
#[derive(Debug)]
pub struct IndexAllocator {
/// Set of currently allocated indices.