diff --git a/CHANGELOG.md b/CHANGELOG.md index f15ac91..49ad3a8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -72,6 +72,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 `handshake.rs:1114` is unchanged. Introduces a shared `Node::outbound_admission_check()` helper so the invariant is grep-able and unit-testable. +- Inbound `handle_msg1` now silent-drops at `node.limits.max_peers` + saturation *before* building/sending Msg2, instead of replying with + Msg2 and then rejecting at `promote_connection`. Adds an early cap + check positioned after identity verification (so the + reconnect / cross-connection bypass for known peers still fires) and + before index allocation + Msg2 wire send. The late cap check inside + `promote_connection` is intentionally retained as + defense-in-depth. Wire savings observed in a 45 s ops tcpdump at + saturation: ~3.6 cap-denials/s × Msg2 (~104 B + AEAD compute) each. + Bigger win is cleaner peer-side semantics — no fake-completed + handshake whose subsequent data frames fail decryption on this side. - Mesh-size estimator (`compute_mesh_size`) no longer double-counts the parent's bloom cardinality during the transient cache window after a local parent-switch. Symptom: `fipsctl show status` / fipstop displayed diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index ae2cc45..0f4dc26 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -207,6 +207,35 @@ impl Node { possible_restart = true; } + // Early cap check: at max_peers and this is a net-new identity? + // Bypass for known peers (reconnect / cross-connection) — admitting + // them doesn't grow peers.len(). This silent-drops the Msg1 before + // the Msg2 build/send and index allocation, avoiding wasted wire + // bytes and giving the peer cleaner semantics (no fake-completed + // handshake whose data frames subsequently fail decryption here). + // The late cap check inside promote_connection() is intentionally + // left in place as defense-in-depth. + if self.max_peers > 0 && self.peers.len() >= self.max_peers { + let is_known_active = self.peers.contains_key(&peer_node_addr); + let is_pending_outbound = self.connections.iter().any(|(_, conn)| { + conn.expected_identity() + .map(|id| *id.node_addr() == peer_node_addr) + .unwrap_or(false) + }); + if !is_known_active && !is_pending_outbound { + debug!( + peer = %self.peer_display_name(&peer_node_addr), + max = self.max_peers, + "Silent-dropping Msg1 at max_peers cap (early gate; no Msg2 sent)" + ); + // `link_id` was allocated above but `conn` is still a local + // (not yet inserted into self.connections / self.links / + // self.addr_to_link), so the local drop suffices. + self.msg1_rate_limiter.complete_handshake(); + return; + } + } + // Epoch-based restart detection and duplicate msg1 handling. // // If we fell through from the addr_to_link check above with diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 7822727..4671625 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -1497,3 +1497,227 @@ fn nostr_discovery_outbound_admission_atomic_roundtrip() { "after recovery store: traversal initiator/responder must see true" ); } + +/// Sender-side helper: build a wire-format Msg1 from a fresh peer +/// identity targeting `node_b`, *and* send it on the wire over `socket_a` +/// to `addr_b`. Returns the sender's NodeAddr so the test can assert on +/// identity-keyed maps. +/// +/// Uses the same outbound-PeerConnection->Noise IK pattern as the +/// integration handshake tests, but inlined and unit-scoped. +async fn craft_and_send_msg1( + node_b: &Node, + sender_identity: &Identity, + socket_a: &tokio::net::UdpSocket, + addr_b: std::net::SocketAddr, + timestamp_ms: u64, +) -> NodeAddr { + use crate::node::wire::build_msg1; + use crate::utils::index::SessionIndex; + + let peer_b_identity = PeerIdentity::from_pubkey_full(node_b.identity.pubkey_full()); + let sender_pubkey_id = PeerIdentity::from_pubkey_full(sender_identity.pubkey_full()); + let sender_node_addr = *sender_pubkey_id.node_addr(); + + let link_id = LinkId::new(0xDEAD_BEEF); + let mut conn = PeerConnection::outbound(link_id, peer_b_identity, timestamp_ms); + + let sender_keypair = sender_identity.keypair(); + let mut startup_epoch = [0u8; 8]; + rand::Rng::fill_bytes(&mut rand::rng(), &mut startup_epoch); + let noise_msg1 = conn + .start_handshake(sender_keypair, startup_epoch, timestamp_ms) + .expect("start_handshake should produce noise msg1"); + + let sender_index = SessionIndex::new(0x5151); + let wire_msg1 = build_msg1(sender_index, &noise_msg1); + + socket_a + .send_to(&wire_msg1, addr_b) + .await + .expect("sender_socket.send_to"); + sender_node_addr +} + +/// Helper: deliver a packet from `node`'s registered UDP transport to +/// `node.handle_msg1`. Returns Ok(()) on success or Err if the packet +/// was not received within `timeout`. +async fn pump_one_msg1_into_node( + node: &mut Node, + packet_rx: &mut crate::transport::PacketRx, + timeout_ms: u64, +) -> Result<(), &'static str> { + use tokio::time::{Duration, timeout}; + let packet = timeout(Duration::from_millis(timeout_ms), packet_rx.recv()) + .await + .map_err(|_| "timed out waiting for msg1 on packet_rx")? + .ok_or("packet channel closed")?; + node.handle_msg1(packet).await; + Ok(()) +} + +/// Verifies the early max_peers cap check in `handle_msg1` silent-drops +/// a Msg1 from a brand-new identity at saturation: no peer is admitted, +/// no Msg2 response goes back on the wire, and the msg1 rate-limiter +/// pending_count returns to baseline. +/// +/// Wire-observable Msg2 absence is the load-bearing discriminator. With +/// the early cap gate removed (stash-verify), the late gate inside +/// `promote_connection` still rejects the new identity — but only +/// *after* `handle_msg1` has already built the Msg2 frame and +/// `transport.send(...wire_msg2)` has put it on the wire. The +/// post-call wire-side poll catches that Msg2 (FAIL pre-fix; the +/// silent timeout is the PASS post-fix). +#[tokio::test] +async fn handle_msg1_silent_drops_at_cap_for_new_peer() { + use crate::config::UdpConfig; + use tokio::time::{Duration, timeout}; + + let mut node = make_node(); + node.set_max_peers(2); + inject_dummy_peers(&mut node, 2); + assert_eq!(node.peer_count(), 2, "precondition: at cap"); + + // === UDP transport setup for node_b (the unit under test) === + let transport_id_b = TransportId::new(1); + let udp_config = UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(1280), + ..Default::default() + }; + let (packet_tx_b, mut packet_rx_b) = packet_channel(64); + let mut transport_b = UdpTransport::new(transport_id_b, None, udp_config, packet_tx_b); + transport_b.start_async().await.unwrap(); + let addr_b = transport_b.local_addr().unwrap(); + node.transports + .insert(transport_id_b, TransportHandle::Udp(transport_b)); + + // === Sender-side socket === + let socket_a = tokio::net::UdpSocket::bind("127.0.0.1:0") + .await + .expect("bind sender socket"); + + let before_peers = node.peer_count(); + let before_pending = node.msg1_rate_limiter.pending_count(); + + // Fresh sender identity — never seen by `node`. + let sender = Identity::generate(); + let sender_node_addr = craft_and_send_msg1(&node, &sender, &socket_a, addr_b, 1000).await; + + // Sanity: new identity is not currently a peer. + assert!( + !node.peers.contains_key(&sender_node_addr), + "precondition: new sender not yet a peer" + ); + + // Pump the wire-arrived Msg1 into the node's handler. + pump_one_msg1_into_node(&mut node, &mut packet_rx_b, 1000) + .await + .expect("msg1 must reach packet_rx_b"); + + // Post-call state checks. + assert_eq!( + node.peer_count(), + before_peers, + "early cap gate must not adopt a new peer at saturation" + ); + assert!( + !node.peers.contains_key(&sender_node_addr), + "new sender must not appear in peers map" + ); + assert_eq!( + node.msg1_rate_limiter.pending_count(), + before_pending, + "rate limiter must rebalance: start_handshake() then \ + complete_handshake() before silent-drop return" + ); + + // Wire-observable discriminator: with the early gate in place, no + // Msg2 should come back. With the gate removed, Msg2 IS sent + // before promote_connection rejects. + let mut buf = [0u8; 2048]; + let recv = timeout(Duration::from_millis(300), socket_a.recv_from(&mut buf)).await; + let received_bytes = recv.ok().and_then(|inner| inner.ok()).map(|(n, _)| n); + assert!( + received_bytes.is_none(), + "Msg2 must NOT be sent in response when at max_peers cap; \ + observed {received_bytes:?} wire bytes — the fingerprint of \ + the late-gate path replying with Msg2 before rejecting" + ); +} + +/// Verifies the bypass: at saturation, an inbound Msg1 from an +/// *existing* peer's identity is not silent-dropped by the early cap +/// check (the gate would otherwise wedge legitimate +/// reconnect/restart/rekey traffic against an at-cap node). +/// +/// The cap-gate's `is_known_active = self.peers.contains_key(&peer_node_addr)` +/// branch admits this case; the downstream handling (restart-detect or +/// duplicate-msg1 resend) then runs per existing semantics. The +/// observable assertion here is the existing peer's continued +/// presence — the rate-limiter rebalance is the same in +/// bypass-admit and silent-drop, so this test isn't a discriminator +/// against the no-gate (stash) build; it's a regression check that the +/// gate doesn't accidentally evict known peers. +#[tokio::test] +async fn handle_msg1_admits_existing_peer_at_cap() { + use crate::config::UdpConfig; + + let mut node = make_node(); + node.set_max_peers(2); + + inject_dummy_peers(&mut node, 1); + + let existing_sender = Identity::generate(); + let existing_pid = PeerIdentity::from_pubkey_full(existing_sender.pubkey_full()); + let existing_node_addr = *existing_pid.node_addr(); + let existing_link_id = LinkId::new(7777); + { + use crate::peer::ActivePeer; + let peer = ActivePeer::new(existing_pid, existing_link_id, 0); + node.peers.insert(existing_node_addr, peer); + } + assert_eq!(node.peer_count(), 2, "precondition: at cap"); + + let transport_id_b = TransportId::new(1); + let udp_config = UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + mtu: Some(1280), + ..Default::default() + }; + let (packet_tx_b, mut packet_rx_b) = packet_channel(64); + let mut transport_b = UdpTransport::new(transport_id_b, None, udp_config, packet_tx_b); + transport_b.start_async().await.unwrap(); + let addr_b = transport_b.local_addr().unwrap(); + node.transports + .insert(transport_id_b, TransportHandle::Udp(transport_b)); + + let socket_a = tokio::net::UdpSocket::bind("127.0.0.1:0") + .await + .expect("bind sender socket"); + + let before_pending = node.msg1_rate_limiter.pending_count(); + + let sender_node_addr = + craft_and_send_msg1(&node, &existing_sender, &socket_a, addr_b, 2000).await; + assert_eq!( + sender_node_addr, existing_node_addr, + "sanity: crafted msg1 carries the existing peer's NodeAddr" + ); + + pump_one_msg1_into_node(&mut node, &mut packet_rx_b, 1000) + .await + .expect("msg1 must reach packet_rx_b"); + + // Bypass must not evict the existing peer or grow peer count. + assert_eq!(node.peer_count(), 2, "peer count unchanged"); + assert!( + node.peers.contains_key(&existing_node_addr), + "existing peer must still be present after bypass-admitted msg1" + ); + assert_eq!( + node.msg1_rate_limiter.pending_count(), + before_pending, + "rate limiter must rebalance after the (bypass-admitted) handler returns" + ); +}