From f5b01370faa8cdd77350a55ad673fba079c6ec3a Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Tue, 11 Aug 2026 06:21:31 +0000 Subject: [PATCH] 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. --- src/node/handlers/handshake.rs | 115 +++- src/node/handlers/rekey.rs | 33 +- src/node/mod.rs | 61 ++ src/node/tests/acl.rs | 62 +- src/node/tests/handshake.rs | 1045 +++++++++++++++++++++++++++++++- src/utils/index.rs | 5 +- 6 files changed, 1308 insertions(+), 13 deletions(-) diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index 020137e3..d4d1570c 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -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; diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index 12ad27a0..40195c2d 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -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), diff --git a/src/node/mod.rs b/src/node/mod.rs index b025514b..160748d1 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -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`, 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. /// diff --git a/src/node/tests/acl.rs b/src/node/tests/acl.rs index a8fc8921..902469df 100644 --- a/src/node/tests/acl.rs +++ b/src/node/tests/acl.rs @@ -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" + ); } diff --git a/src/node/tests/handshake.rs b/src/node/tests/handshake.rs index d50260e3..73626bd9 100644 --- a/src/node/tests/handshake.rs +++ b/src/node/tests/handshake.rs @@ -2761,24 +2761,72 @@ async fn test_anonymous_msg2_crystallizes_identity_and_promotes() { stop_hs(&mut responder).await; } +/// The msg2 self-connect drop must release all three things the outbound leg +/// holds: its session index, its link, and its `pending_outbound` entry. +/// +/// The three assertions are independent by construction — each names a +/// different registry — so deleting any one of the three releases from the arm +/// reds exactly one of them. +/// +/// The node handshakes with itself, so it holds TWO legs here: the outbound one +/// being dropped and the inbound one its own `handle_msg1` created. That is why +/// the allocator goes to `baseline - 1` rather than to zero, and why the +/// inbound leg's index staying allocated is a useful control that the right +/// index was freed. #[tokio::test] async fn test_anonymous_self_connect_drop_disposes_machine() { let mut node = make_hs_node(Config::new()).await; let self_addr = node.addr.clone(); + let transport_id = node.transport_id; // Anonymously dial our own bound address (a shared-media beacon can echo // ourselves back at us). node.node - .initiate_connection(node.transport_id, self_addr, None) + .initiate_connection(node.transport_id, self_addr.clone(), None) .await .expect("anonymous self dial"); let leg_link = node.node.connections().next().unwrap().1.link_id(); assert_eq!(node.node.peer_machines.len(), 1); + let outbound_idx = node + .node + .peer_machines + .get(&leg_link) + .unwrap() + .our_index() + .expect("the dial allocated an index for the outbound leg"); + assert!( + node.node + .pending_outbound + .contains_key(&(transport_id, outbound_idx.as_u32())), + "control: the dial registered the leg for msg2 dispatch" + ); // We answer our own msg1, then our msg2 processing discovers the learned // identity is our own and drops the leg — machine included. let msg1_pkt = recv_phase(&mut node.packet_rx, 1, "msg1").await; node.node.handle_msg1(msg1_pkt).await; + + // The inbound leg msg1 just parked: the one machine that is not the dial's. + let inbound_link = *node + .node + .peer_machines + .keys() + .find(|k| **k != leg_link) + .expect("msg1 parks an inbound leg of its own"); + let inbound_idx = node + .node + .peer_machines + .get(&inbound_link) + .unwrap() + .our_index() + .expect("msg1 allocates an index for the inbound leg"); + let baseline = node.node.index_allocator.count(); + assert_eq!(baseline, 2, "one index per leg, and there are two legs"); + assert!( + node.node.index_allocator.is_allocated(outbound_idx), + "control: the outbound leg's index is live before the drop" + ); + let msg2_pkt = recv_phase(&mut node.packet_rx, 2, "msg2").await; node.node.handle_msg2(msg2_pkt).await; @@ -2791,6 +2839,50 @@ async fn test_anonymous_self_connect_drop_disposes_machine() { !node.node.peer_machines.contains_key(&leg_link), "self-connect drop disposes the outbound leg's machine" ); + + // Leak 1: the session index. + assert!( + !node.node.index_allocator.is_allocated(outbound_idx), + "the self-connect drop must return the outbound leg's index" + ); + assert_eq!( + node.node.index_allocator.count(), + baseline - 1, + "and must free exactly one" + ); + assert!( + node.node.index_allocator.is_allocated(inbound_idx), + "the surviving inbound leg keeps its own index" + ); + + // Leak 2: the link. + assert!( + !node.node.links.contains_key(&leg_link), + "the self-connect drop must remove the outbound leg's link" + ); + + // Leak 3: the pending_outbound entry. + assert!( + !node + .node + .pending_outbound + .contains_key(&(transport_id, outbound_idx.as_u32())), + "the self-connect drop must clear the leg's msg2 dispatch entry" + ); + + // Not a leak assertion: a guard against the wrong fix. `handle_msg1` + // overwrote the reverse mapping for this address with the INBOUND leg's + // link, and `remove_link` declines to clear an entry that no longer points + // at the link being removed. An implementer replacing `remove_link` with a + // bare `links.remove` plus `addr_to_link.remove` would take down the live + // inbound leg's routing, and this is where that shows up. + assert_eq!( + node.node.addr_to_link.get(&(transport_id, self_addr)), + Some(&inbound_link), + "the surviving inbound leg keeps the address mapping; the dropped \ + outbound leg must not take it down" + ); + node.node.debug_assert_peer_maps_coherent(); stop_hs(&mut node).await; @@ -3314,3 +3406,954 @@ async fn test_anonymous_dial_msg2_promotes_whoever_answers() { stop_hs(&mut initiator).await; stop_hs(&mut responder).await; } + +// =========================================================================== +// Reject-arm resource release on the XX handshake path +// +// Every arm below disposes a leg and must return what that leg allocated. The +// discriminating assertion is always allocator-level and names the index: +// `is_allocated(idx)` before, `!is_allocated(idx)` after, plus a `count()` +// limb that catches a fix which frees an extra index. A `peers_by_index` +// assertion is worthless at all of these sites — nothing inserts there ahead +// of the msg3 gates, so its pass value and its failure value coincide. +// =========================================================================== + +/// Hand-carry one XX exchange up to the responder's parked msg3 wait, with no +/// transports involved. Returns the responder's link, its msg1-allocated index, +/// and the initiator leg (which still owes a msg3). +async fn park_inbound_leg( + node_b: &mut Node, + initiator_keypair: secp256k1::Keypair, + initiator_epoch: [u8; 8], + transport_id: TransportId, + remote_addr: &TransportAddr, + sender_idx: SessionIndex, +) -> (LinkId, SessionIndex, PeerMachine) { + use crate::proto::fmp::wire::build_msg1; + + let peer_b_identity = PeerIdentity::from_pubkey_full(node_b.identity().pubkey_full()); + let mut conn_a = outbound_leg(LinkId::new(1), peer_b_identity, 1000); + let noise_msg1 = conn_a + .start_handshake(initiator_keypair, initiator_epoch, 1000) + .unwrap(); + + node_b + .handle_msg1(ReceivedPacket::with_timestamp( + transport_id, + remote_addr.clone(), + build_msg1(sender_idx, &noise_msg1), + 1000, + )) + .await; + + let link_id_b = *node_b + .peer_machines + .keys() + .next() + .expect("msg1 parks a machine on the responder"); + let our_index_b = node_b + .peer_machines + .get(&link_id_b) + .unwrap() + .our_index() + .expect("msg1 allocates the responder's session index"); + (link_id_b, our_index_b, conn_a) +} + +/// Finish the initiator's half: read the framed msg2 off the responder's +/// carrier and produce the noise msg3, optionally carrying a negotiation +/// payload. +fn answer_with_msg3( + node_b: &Node, + link_id_b: LinkId, + conn_a: &mut PeerMachine, + negotiation: Option<&[u8]>, +) -> Vec { + use crate::proto::fmp::wire::Msg2Header; + + let wire_msg2 = node_b + .peer_machines + .get(&link_id_b) + .unwrap() + .conn_handshake_msg2() + .expect("msg1 stores the framed msg2 on the carrier") + .to_vec(); + let msg2_header = Msg2Header::parse(&wire_msg2).unwrap(); + let (noise_msg3, _) = conn_a + .complete_handshake(msg2_header.noise_msg2(&wire_msg2), negotiation, 1100) + .unwrap(); + noise_msg3 +} + +/// S3/T1 — the orphaned pending-inbound arm frees the index it can only read +/// off the wire, and disposes the link with it. +/// +/// The planted state is a `pending_inbound` entry whose link has lost its +/// machine. It is produced through `remove_peer_machine`, the disposal +/// choke-point, so the peer's timer store goes with it exactly as a real +/// teardown would leave things. +/// +/// `unknown_connection` is this arm's own counter, which separates it from +/// every `bad_state` arm. It does NOT separate it from the +/// `pending_inbound`-miss arm above, which records the same counter — the +/// `is_allocated` transition is what does, since that arm has no index to free. +#[tokio::test] +async fn test_msg3_orphaned_pending_inbound_frees_index() { + use crate::proto::fmp::wire::build_msg3; + + let mut node_b = make_node(); + let node_a = make_node(); + let transport_id = TransportId::new(1); + let remote_addr = TransportAddr::from_string("127.0.0.1:5000"); + let sender_idx = SessionIndex::new(7); + + let (link_id_b, our_index_b, mut conn_a) = park_inbound_leg( + &mut node_b, + node_a.identity().keypair(), + node_a.startup_epoch(), + transport_id, + &remote_addr, + sender_idx, + ) + .await; + let noise_msg3 = answer_with_msg3(&node_b, link_id_b, &mut conn_a, None); + + // The unaudited disposal this arm exists to survive: the machine goes, the + // pending-inbound entry does not. + node_b.remove_peer_machine(link_id_b); + assert!( + node_b + .pending_inbound + .contains_key(&(transport_id, our_index_b.as_u32())), + "control: the pending-inbound entry outlived its machine" + ); + + let baseline = node_b.index_allocator.count(); + assert!( + node_b.index_allocator.is_allocated(our_index_b), + "control: the leg allocated an index" + ); + assert!( + node_b.links.contains_key(&link_id_b), + "control: the orphaned leg still holds a link" + ); + + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr, + build_msg3(sender_idx, our_index_b, &noise_msg3), + 1200, + )) + .await; + + assert!( + !node_b.index_allocator.is_allocated(our_index_b), + "the orphaned pending-inbound arm must return the index it allocated" + ); + assert_eq!( + node_b.index_allocator.count(), + baseline - 1, + "and must free exactly one" + ); + assert!( + !node_b.links.contains_key(&link_id_b), + "the link goes with the orphaned entry" + ); + assert_eq!( + node_b.stats().handshake.unknown_connection, + 1, + "this arm records its own reject counter" + ); +} + +/// S3/T2a — a hostile msg3 naming a PROMOTED peer's index must not free it. +/// +/// `receiver_idx` is chosen by the sender. Freeing it the way the other arms +/// free would turn a slow leak into a remote session-teardown primitive: name +/// an index belonging to an unrelated live session and the node hands it back +/// to the allocator. The ownership check is what stops that. +#[tokio::test] +async fn test_msg3_orphaned_entry_naming_live_peer_index_does_not_free() { + use crate::proto::fmp::wire::build_msg3; + + let mut node_b = make_node(); + let node_a = make_node(); + let transport_id = TransportId::new(1); + let remote_addr = TransportAddr::from_string("127.0.0.1:5000"); + let sender_idx = SessionIndex::new(7); + + let (link_id_b, our_index_b, mut conn_a) = park_inbound_leg( + &mut node_b, + node_a.identity().keypair(), + node_a.startup_epoch(), + transport_id, + &remote_addr, + sender_idx, + ) + .await; + let noise_msg3 = answer_with_msg3(&node_b, link_id_b, &mut conn_a, None); + + // First msg3 promotes: the responder now holds a live session index. + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr.clone(), + build_msg3(sender_idx, our_index_b, &noise_msg3), + 1200, + )) + .await; + assert_eq!(node_b.peer_count(), 1, "control: the first msg3 promoted"); + let peer_addr = *PeerIdentity::from_pubkey_full(node_a.identity().pubkey_full()).node_addr(); + let victim_idx = node_b + .get_peer(&peer_addr) + .unwrap() + .our_index() + .expect("a promoted peer holds a session index"); + let victim_link = node_b.get_peer(&peer_addr).unwrap().link_id(); + assert!( + node_b + .peers_by_index + .contains_key(&(transport_id, victim_idx.as_u32())), + "control: the live session index is registered for dispatch" + ); + + // The aliasing state: a stale pending-inbound entry naming the LIVE index, + // pointing at a link that has no machine. + let dead_link = LinkId::new(4242); + assert!(!node_b.peer_machines.contains_key(&dead_link)); + node_b + .pending_inbound + .insert((transport_id, victim_idx.as_u32()), dead_link); + + let baseline = node_b.index_allocator.count(); + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr, + build_msg3(sender_idx, victim_idx, &noise_msg3), + 1300, + )) + .await; + + assert!( + node_b.index_allocator.is_allocated(victim_idx), + "a live peer's session index must never be freed by a wire-named \ + orphan cleanup" + ); + assert_eq!( + node_b.index_allocator.count(), + baseline, + "nothing at all was freed" + ); + assert!( + node_b + .peers_by_index + .contains_key(&(transport_id, victim_idx.as_u32())), + "the victim's dispatch entry survives" + ); + assert_eq!(node_b.peer_count(), 1, "the victim peer survives"); + assert!(node_b.links.contains_key(&victim_link)); + assert!(node_b.peer_machines.contains_key(&victim_link)); +} + +/// S3/T2b — the same hostile shape against a victim `peers_by_index` cannot +/// see: a live PENDING INBOUND leg, whose index exists only on its machine. +/// +/// T2a alone proves only the predicate's first limb, since a promoted peer +/// short-circuits it. This is the population that motivates the rest of the +/// scan: deleting the `peer_machines` limb must red this test and leave T2a +/// green. +#[tokio::test] +async fn test_msg3_orphaned_entry_naming_pending_leg_index_does_not_free() { + use crate::proto::fmp::wire::build_msg3; + + let mut node_b = make_node(); + let node_a = make_node(); + let transport_id = TransportId::new(1); + let remote_addr = TransportAddr::from_string("127.0.0.1:5000"); + let sender_idx = SessionIndex::new(7); + + let (victim_link, victim_idx, mut conn_a) = park_inbound_leg( + &mut node_b, + node_a.identity().keypair(), + node_a.startup_epoch(), + transport_id, + &remote_addr, + sender_idx, + ) + .await; + let noise_msg3 = answer_with_msg3(&node_b, victim_link, &mut conn_a, None); + + // The victim's index is held by a parked machine and by nothing else: no + // peer exists yet, and the msg3 gate consumes the pending-inbound key + // before the ownership check runs. + assert_eq!(node_b.peer_count(), 0); + assert!(node_b.peers_by_index.is_empty()); + assert!(node_b.pending_outbound.is_empty()); + + // Overwrite the victim's own pending-inbound entry with a machine-less + // link. The point is a link with no machine, which is what selects the arm. + let dead_link = LinkId::new(4242); + assert!(!node_b.peer_machines.contains_key(&dead_link)); + node_b + .pending_inbound + .insert((transport_id, victim_idx.as_u32()), dead_link); + + let baseline = node_b.index_allocator.count(); + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr, + build_msg3(sender_idx, victim_idx, &noise_msg3), + 1200, + )) + .await; + + assert!( + node_b.index_allocator.is_allocated(victim_idx), + "an index held by a live pending leg must never be freed by a \ + wire-named orphan cleanup" + ); + assert_eq!(node_b.index_allocator.count(), baseline); + assert!( + node_b.peer_machines.contains_key(&victim_link), + "the victim leg's machine survives" + ); + assert!(node_b.links.contains_key(&victim_link)); +} + +/// S3/T2c — the same hostile shape delivered on a DIFFERENT transport from the +/// one the victim's index is registered under. +/// +/// This is the discriminating test for the one regression the ownership check +/// has: re-scoping `session_index_is_claimed` to the transport the msg3 arrived +/// on. `index_allocator` is a single node-global `HashSet` while every +/// registry it is checked against is keyed `(TransportId, u32)`, so a +/// transport-filtered scan reports "unclaimed" for an index that is very much +/// claimed, on another transport — and the arm hands a live session's index +/// back to the allocator, which is exactly the teardown primitive the check +/// exists to prevent. +/// +/// The victim has to be chosen with care, and T2a's victim will not do. A +/// promoted peer's `our_index` also sits on that peer's surviving `PeerMachine`, +/// and the `peer_machines` limb has no transport in its key to filter by, so it +/// masks the mutation and the test reports green either way. An in-flight +/// REKEY index is the one live index held by no machine: it lives in +/// `pending_outbound` under the peer's transport and on the `ActivePeer`, both +/// of which a transport-scoped scan filters out. The two controls below pin +/// that. +#[tokio::test] +async fn test_msg3_orphan_on_other_transport_does_not_free_live_index() { + use crate::proto::fmp::wire::build_msg3; + + let make_config = || { + let mut c = Config::new(); + c.node.rekey.enabled = true; + c.node.rekey.after_secs = 30; + c + }; + + let mut initiator = make_hs_node(make_config()).await; + let mut responder = make_hs_node(make_config()).await; + + let responder_addr = + *PeerIdentity::from_pubkey_full(responder.node.identity().pubkey_full()).node_addr(); + + let msg3 = drive_to_msg3(&mut initiator, &mut responder, 1000).await; + responder.node.handle_msg3(msg3).await; + assert_eq!(initiator.node.peer_count(), 1, "control: the session is up"); + + // Start a real rekey so the victim index is produced exactly as production + // produces it. + initiator + .node + .get_peer_mut(&responder_addr) + .unwrap() + .test_backdate_session_established(std::time::Duration::from_secs(120)); + initiator.node.check_rekey().await; + + let victim_idx = initiator + .node + .get_peer(&responder_addr) + .unwrap() + .rekey_our_index() + .expect("check_rekey started a cycle and allocated its index"); + + // Control 1: the victim is registered under the peer's OWN transport only. + assert!( + initiator + .node + .pending_outbound + .contains_key(&(initiator.transport_id, victim_idx.as_u32())), + "control: the rekey index is registered under the peer's transport" + ); + // Control 2: no machine holds it. This is what makes the test able to see a + // transport-scoped scan at all — the `peer_machines` limb is unfiltered, so + // a machine-held victim would mask the regression. + assert!( + !initiator + .node + .peer_machines + .values() + .any(|m| m.our_index() == Some(victim_idx)), + "control: no PeerMachine holds the rekey index, so only the \ + transport-keyed limbs can answer for it" + ); + + // The hostile msg3: a different transport, an orphaned pending-inbound entry + // naming the live rekey index, and a link with no machine. + let other_transport = TransportId::new(99); + assert_ne!(other_transport, initiator.transport_id); + let dead_link = LinkId::new(4242); + assert!(!initiator.node.peer_machines.contains_key(&dead_link)); + initiator + .node + .pending_inbound + .insert((other_transport, victim_idx.as_u32()), dead_link); + + let baseline = initiator.node.index_allocator.count(); + initiator + .node + .handle_msg3(ReceivedPacket::with_timestamp( + other_transport, + TransportAddr::from_string("127.0.0.1:5999"), + build_msg3( + SessionIndex::new(7), + victim_idx, + &[0u8; crate::noise::HANDSHAKE_MSG3_SIZE], + ), + 1300, + )) + .await; + + assert!( + initiator.node.index_allocator.is_allocated(victim_idx), + "an index claimed on ANOTHER transport must not be freed: the \ + allocator is node-global, so the ownership scan must be too" + ); + assert_eq!( + initiator.node.index_allocator.count(), + baseline, + "nothing at all was freed" + ); + assert!( + initiator + .node + .pending_outbound + .contains_key(&(initiator.transport_id, victim_idx.as_u32())), + "the in-flight rekey's dispatch entry survives" + ); + assert_eq!(initiator.node.peer_count(), 1, "the victim peer survives"); +} + +/// S3/T2d — the `ActivePeer` limb of the ownership scan, pinned directly. +/// +/// **The state below is fabricated and is not reachable in production.** Every +/// `self.peers` insertion goes through `ActivePeer::with_session` with a +/// concrete `TransportId` and an unconditional `peers_by_index` insert beside +/// it, so today a promoted peer's `our_index` is always answerable from the +/// maps and from its surviving machine. T2a and T2b ride those two limbs, and +/// deleting the `peers` scan entirely leaves the whole suite green — which is +/// the reason this test exists rather than a reason it should not. +/// +/// The limb is defence-in-depth against holders the transport-keyed maps cannot +/// answer for: `remove_active_peer` gates its index frees on `if let Some(tid)`, +/// and the `peers_by_index` insert for a rekeying peer's `pending_our_index` is +/// gated the same way, so an index can in principle be held by an `ActivePeer` +/// and by nothing else. This calls the predicate directly rather than dressing +/// that state up as an arrival, because pretending it is reachable would be the +/// worse dishonesty. +#[tokio::test] +async fn test_session_index_is_claimed_answers_from_the_peer_when_maps_cannot() { + use crate::proto::fmp::wire::build_msg3; + + let mut node_b = make_node(); + let node_a = make_node(); + let transport_id = TransportId::new(1); + let remote_addr = TransportAddr::from_string("127.0.0.1:5000"); + let sender_idx = SessionIndex::new(7); + + let (link_id_b, our_index_b, mut conn_a) = park_inbound_leg( + &mut node_b, + node_a.identity().keypair(), + node_a.startup_epoch(), + transport_id, + &remote_addr, + sender_idx, + ) + .await; + let noise_msg3 = answer_with_msg3(&node_b, link_id_b, &mut conn_a, None); + + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr, + build_msg3(sender_idx, our_index_b, &noise_msg3), + 1200, + )) + .await; + assert_eq!(node_b.peer_count(), 1, "control: the peer promoted"); + + let peer_addr = *PeerIdentity::from_pubkey_full(node_a.identity().pubkey_full()).node_addr(); + let victim_idx = node_b + .get_peer(&peer_addr) + .unwrap() + .our_index() + .expect("a promoted peer holds a session index"); + let victim_link = node_b.get_peer(&peer_addr).unwrap().link_id(); + + assert!( + node_b.session_index_is_claimed(victim_idx), + "control: with every registry intact the index reads as claimed" + ); + + // Strip the two limbs that answer today, leaving the `ActivePeer` as the + // only holder. Fabricated, per the doc comment above. + node_b + .peers_by_index + .remove(&(transport_id, victim_idx.as_u32())); + node_b.remove_peer_machine(victim_link); + assert!( + !node_b + .peers_by_index + .keys() + .any(|(_, i)| *i == victim_idx.as_u32()), + "control: no transport-keyed registry names the index any more" + ); + assert!( + !node_b + .peer_machines + .values() + .any(|m| m.our_index() == Some(victim_idx)), + "control: no machine names it either" + ); + + assert!( + node_b.session_index_is_claimed(victim_idx), + "the peer itself still holds the index, so the orphan arm must not \ + free it" + ); +} + +/// S4 — the msg3 FMP-negotiation reject arm must free the msg1-allocated index. +/// +/// The payload must be non-empty or the responder's split yields `None` and the +/// negotiation block is skipped entirely. Three bytes of `0xff` decrypt fine +/// and then fail `NegotiationPayload::decode`'s length gate, which is the arm +/// this test is aimed at. +/// +/// This arm sits ahead of the inbound ACL gate, so any peer able to complete a +/// Noise msg3 reaches it without being authorized. +/// +/// The counters alone do not identify the arm: the msg3-processing-failure arm +/// 25 lines above records the identical `bad_state` + empty registries AND +/// already frees the index, so a harness slip that makes the msg3 fail to +/// DECRYPT rather than fail to DECODE would pass every assertion here with or +/// without the fix. `test_msg3_negotiation_wellformed_promotes` is the control +/// for that: it runs the identical construction with a well-formed payload and +/// requires promotion, which can only happen if the appended bytes decrypt. +#[tokio::test] +async fn test_msg3_negotiation_failure_frees_index() { + use crate::proto::fmp::wire::build_msg3; + + let mut node_b = make_node(); + let node_a = make_node(); + let transport_id = TransportId::new(1); + let remote_addr = TransportAddr::from_string("127.0.0.1:5000"); + let sender_idx = SessionIndex::new(7); + + let (link_id_b, our_index_b, mut conn_a) = park_inbound_leg( + &mut node_b, + node_a.identity().keypair(), + node_a.startup_epoch(), + transport_id, + &remote_addr, + sender_idx, + ) + .await; + + let bad_neg = vec![0xffu8; 3]; + let noise_msg3 = answer_with_msg3(&node_b, link_id_b, &mut conn_a, Some(&bad_neg)); + + let baseline = node_b.index_allocator.count(); + assert!( + node_b.index_allocator.is_allocated(our_index_b), + "control: msg1 allocated the responder's index" + ); + + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr, + build_msg3(sender_idx, our_index_b, &noise_msg3), + 1200, + )) + .await; + + assert!( + !node_b.index_allocator.is_allocated(our_index_b), + "the negotiation-failure arm must return the msg1-allocated index" + ); + assert_eq!( + node_b.index_allocator.count(), + baseline - 1, + "and must free exactly one" + ); + assert_eq!( + node_b.peer_count(), + 0, + "a failed negotiation never promotes" + ); + assert!(node_b.peer_machines.is_empty()); + assert_eq!(node_b.link_count(), 0); + assert_eq!(node_b.stats().handshake.bad_state, 1); +} + +/// The positive sibling of `test_msg3_negotiation_failure_frees_index`, and the +/// only thing that proves the appended payload reaches the responder's +/// negotiation step at all rather than failing earlier in the Noise decrypt. +#[tokio::test] +async fn test_msg3_negotiation_wellformed_promotes() { + use crate::proto::fmp::NegotiationPayload; + use crate::proto::fmp::wire::build_msg3; + + let mut node_b = make_node(); + let node_a = make_node(); + let transport_id = TransportId::new(1); + let remote_addr = TransportAddr::from_string("127.0.0.1:5000"); + let sender_idx = SessionIndex::new(7); + + let (link_id_b, our_index_b, mut conn_a) = park_inbound_leg( + &mut node_b, + node_a.identity().keypair(), + node_a.startup_epoch(), + transport_id, + &remote_addr, + sender_idx, + ) + .await; + + let good_neg = NegotiationPayload::fmp(1, 1, node_a.node_profile()).encode(); + let noise_msg3 = answer_with_msg3(&node_b, link_id_b, &mut conn_a, Some(&good_neg)); + + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr, + build_msg3(sender_idx, our_index_b, &noise_msg3), + 1200, + )) + .await; + + assert_eq!( + node_b.peer_count(), + 1, + "a well-formed negotiation payload on the same construction promotes, \ + which is what proves the appended bytes decrypt on the responder" + ); + assert_eq!(node_b.stats().handshake.bad_state, 0); +} + +/// S7 — the msg3 self-connect drop must free the msg1-allocated index. +/// +/// A real self-dial cannot reach this arm: the msg2 self-connect drop takes the +/// leg down first. It is constructed directly by handing the initiator leg +/// node_b's OWN keypair, so the static in msg3 is node_b's own. The inbound ACL +/// gate runs first and default-allows our own identity, so the arm is reached. +/// +/// Limitation, stated rather than fixed: the assertions do not by themselves +/// identify this arm. The msg3-processing-failure arm above records the same +/// `bad_state == 1` with the same empty registries and already frees the index, +/// so a harness slip that made the hand-built msg3 fail to DECRYPT would leave +/// every assertion here green with or without S7's fix. That the test lands on +/// S7 today was established by mutation — deleting S7's free reds it — not by +/// anything the assertions can distinguish. Unlike S4, no positive sibling +/// control is available here: a well-formed self-connect msg3 cannot promote, +/// which is the point of the arm. +#[tokio::test] +async fn test_msg3_self_connect_frees_index() { + use crate::proto::fmp::wire::build_msg3; + + let mut node_b = make_node(); + let transport_id = TransportId::new(1); + let remote_addr = TransportAddr::from_string("127.0.0.1:5000"); + let sender_idx = SessionIndex::new(7); + + let own_keypair = node_b.identity().keypair(); + let own_epoch = node_b.startup_epoch(); + let (link_id_b, our_index_b, mut conn_a) = park_inbound_leg( + &mut node_b, + own_keypair, + own_epoch, + transport_id, + &remote_addr, + sender_idx, + ) + .await; + let noise_msg3 = answer_with_msg3(&node_b, link_id_b, &mut conn_a, None); + + let baseline = node_b.index_allocator.count(); + assert!( + node_b.index_allocator.is_allocated(our_index_b), + "control: msg1 allocated the responder's index" + ); + + node_b + .handle_msg3(ReceivedPacket::with_timestamp( + transport_id, + remote_addr, + build_msg3(sender_idx, our_index_b, &noise_msg3), + 1200, + )) + .await; + + assert!( + !node_b.index_allocator.is_allocated(our_index_b), + "the msg3 self-connect drop must return the msg1-allocated index" + ); + assert_eq!( + node_b.index_allocator.count(), + baseline - 1, + "and must free exactly one" + ); + assert_eq!(node_b.peer_count(), 0, "we never promote ourselves"); + assert_eq!(node_b.link_count(), 0); + assert_eq!(node_b.stats().handshake.bad_state, 1); +} + +/// S1b — the outbound ACL reject at msg2 must put the dial back on the retry +/// schedule, or a configured peer is dialed once at startup and never again. +/// +/// Ordering matters and is a vacuous-pass trap: `initiate_connection` consults +/// the ACL itself and returns `AccessDenied`, so denying BEFORE the dial means +/// no msg1 is ever sent, the arm is never reached, and the `retry_pending` +/// assertion passes for the wrong reason. The deny is written and reloaded +/// AFTER the dial, which is also the runtime-reload scenario that motivates +/// the fix. +/// +/// The emptiness control is anchored twice — after the dial and again after the +/// reload — because three sites insert into `retry_pending`, so a single +/// pre-dial check would not attribute the entry to this arm. +/// +/// Not asserted here: that the reschedule names the DIALED peer rather than the +/// learned static. In this scenario the two are the same node, so the limb +/// cannot fail. The sibling gate's own test covers that with two distinct +/// identities. +#[tokio::test] +async fn test_outbound_msg2_acl_reject_reschedules_dial() { + use crate::node::acl::PeerAclReloader; + + let mut responder = make_hs_node(Config::new()).await; + let responder_identity = + PeerIdentity::from_pubkey_full(responder.node.identity().pubkey_full()); + let responder_addr = *responder_identity.node_addr(); + + // The dialed peer is a configured auto-connect peer: the only shape the + // retry machinery seeds a schedule entry for, and the shape the missing + // reschedule strands. + let mut config = Config::new(); + config.peers.push(crate::config::PeerConfig::new( + responder.node.npub(), + "udp", + "10.0.0.2:2121", + )); + let mut initiator = make_hs_node(config).await; + + let dir = tempfile::tempdir().unwrap(); + initiator.node.peer_acl = PeerAclReloader::with_paths( + dir.path().join("peers.allow"), + dir.path().join("peers.deny"), + ); + + initiator + .node + .initiate_connection( + initiator.transport_id, + responder.addr.clone(), + Some(responder_identity), + ) + .await + .expect("named dial"); + assert!( + initiator.node.peering.reconciler.retry_pending.is_empty(), + "nothing is scheduled by the dial itself" + ); + + // Deny only now, so the dial really happened and the msg2 arm is the gate + // that turns the peer away. + std::fs::write( + dir.path().join("peers.deny"), + format!("{}\n", responder.node.npub()), + ) + .unwrap(); + assert!(initiator.node.reload_peer_acl().await); + assert!( + initiator.node.peering.reconciler.retry_pending.is_empty(), + "the reload itself schedules nothing" + ); + + let msg1 = recv_phase(&mut responder.packet_rx, 1, "msg1").await; + responder.node.handle_msg1(msg1).await; + let msg2 = recv_phase(&mut initiator.packet_rx, 2, "msg2").await; + initiator.node.handle_msg2(msg2).await; + + // Arm identity: without these, an ACL that failed to load produces a + // promoted peer and an empty schedule, and the assertion below would be + // testing nothing. + assert_eq!( + initiator.node.peer_count(), + 0, + "the denied peer never promotes" + ); + assert_eq!( + initiator.node.stats().handshake.bad_state, + 1, + "the denial is attributed to the handshake state-machine counter" + ); + assert!( + initiator.node.peer_machines.is_empty(), + "the leg is disposed" + ); + + assert!( + initiator + .node + .peering + .reconciler + .retry_pending + .contains_key(&responder_addr), + "the ACL-rejected dial must leave the configured peer scheduled for \ + retry, or it is never dialed again for the life of the process" + ); + + stop_hs(&mut initiator).await; + stop_hs(&mut responder).await; +} + +/// S8 — abandoning a rekey cycle must return the rekey index AND clear the +/// `pending_outbound` entry seeded with it. +/// +/// `handshake_max_resends = 0` fires the abandon on the first +/// `resend_pending_rekeys` call, because the classifier tests +/// `resend_count >= max_resends` BEFORE it consults the resend-due predicate. +/// That takes the wall clock out of the path entirely: the rekey deadline is +/// seeded from `SystemTime::now()`, so a synthetic `now_ms` can never make the +/// msg1 look due and the budget would never be spent. +/// +/// The rekey msg2 is deliberately never delivered — that is what makes this the +/// abandon path. +/// +/// No assertion is made on `peers_by_index` for the rekey index: nothing +/// inserts there until the rekey msg2 install, so it would be inert here. +#[tokio::test] +async fn test_abandon_rekey_frees_index_and_pending_outbound() { + let make_config = || { + let mut c = Config::new(); + c.node.rekey.enabled = true; + c.node.rekey.after_secs = 30; + c.node.rate_limit.handshake_max_resends = 0; + c + }; + + let mut initiator = make_hs_node(make_config()).await; + let mut responder = make_hs_node(make_config()).await; + + let responder_addr = + *PeerIdentity::from_pubkey_full(responder.node.identity().pubkey_full()).node_addr(); + + let msg3 = drive_to_msg3(&mut initiator, &mut responder, 1000).await; + responder.node.handle_msg3(msg3).await; + assert_eq!(initiator.node.peer_count(), 1); + + let session_idx = initiator + .node + .get_peer(&responder_addr) + .unwrap() + .our_index() + .expect("the established peer holds a session index"); + + // Age the session past the (jittered) trigger and let the real cadence fire + // so the rekey index and its pending_outbound entry are produced exactly as + // in production. + initiator + .node + .get_peer_mut(&responder_addr) + .unwrap() + .test_backdate_session_established(std::time::Duration::from_secs(120)); + initiator.node.check_rekey().await; + + let rekey_idx = initiator + .node + .get_peer(&responder_addr) + .unwrap() + .rekey_our_index() + .expect("check_rekey started a cycle and allocated its index"); + let baseline = initiator.node.index_allocator.count(); + assert!( + initiator.node.index_allocator.is_allocated(rekey_idx), + "control: the rekey allocated an index" + ); + assert!( + initiator + .node + .pending_outbound + .contains_key(&(initiator.transport_id, rekey_idx.as_u32())), + "control: the rekey registered its index for msg2 dispatch" + ); + + // The rekey msg2 never arrives; the first poll spends the (zero) budget. + initiator.node.resend_pending_rekeys(2000).await; + + assert!( + !initiator + .node + .get_peer(&responder_addr) + .unwrap() + .rekey_in_progress(), + "control: the abandon actually fired" + ); + assert!( + initiator + .node + .get_peer(&responder_addr) + .unwrap() + .rekey_our_index() + .is_none() + ); + + assert!( + !initiator.node.index_allocator.is_allocated(rekey_idx), + "the abandoned cycle must return its index" + ); + assert_eq!( + initiator.node.index_allocator.count(), + baseline - 1, + "and must free exactly one" + ); + assert!( + !initiator + .node + .pending_outbound + .contains_key(&(initiator.transport_id, rekey_idx.as_u32())), + "the abandoned cycle must clear its msg2 dispatch entry" + ); + + // The limb that catches a fix binding the wrong index: the live session must + // be untouched. + assert_eq!(initiator.node.peer_count(), 1, "the session survives"); + assert!( + initiator.node.index_allocator.is_allocated(session_idx), + "the established session's own index must not be freed" + ); + assert!( + initiator + .node + .peers_by_index + .contains_key(&(initiator.transport_id, session_idx.as_u32())), + "the established session stays registered for dispatch" + ); + + stop_hs(&mut initiator).await; + stop_hs(&mut responder).await; +} diff --git a/src/utils/index.rs b/src/utils/index.rs index 94a94636..29d638ec 100644 --- a/src/utils/index.rs +++ b/src/utils/index.rs @@ -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.