Send the XX handshake replies and remaining tick sends without dialing

The merge from master brought the non-dialing send and switched the call
sites whose code next shares with master: the rekey msg1 and its resend,
the msg1 resend on an outbound handshake, the executor's stored msg1
send, the encrypted link send and the executor's msg2 arm (unreachable
on next, since handle_msg1 sends its msg2 inline). The XX handshake's
own sites still called the dialing transport send, so a reply to a
connection that had closed held the rx loop for the whole connect
timeout, as in #176.

The replies now use send_existing and never dial: the msg2 to a msg1,
the msg2 resent for a duplicate msg1, the msg3 to a dial's msg2 and to a
rekey msg2, the msg2 resent for a duplicate handshake at msg3, and the
Disconnect sent to a peer the access list refuses at msg3. The rekey
msg3 resend and the anonymous dial's inline msg1 go through send_nowait,
which starts a background connect only toward an address this node
dialed, and the tests cover that for a dialed peer and an inbound one.

process_packet and start_handshake are visible to the node tests, which
call them directly. The tests in src/node/tests/rx_stall.rs are the XX
counterparts of master's.
This commit is contained in:
Johnathan Corgan
2026-10-04 22:29:04 +00:00
parent aa753714bd
commit e65950bc5e
7 changed files with 1459 additions and 17 deletions
+1 -1
View File
@@ -619,7 +619,7 @@ impl Node {
/// Process a single received packet.
///
/// Dispatches based on the phase field in the 4-byte common prefix.
async fn process_packet(&mut self, packet: ReceivedPacket) {
pub(in crate::node) async fn process_packet(&mut self, packet: ReceivedPacket) {
if packet.data.len() < COMMON_PREFIX_SIZE {
return; // Drop packets too short for common prefix
}
+29 -9
View File
@@ -329,11 +329,15 @@ impl Node {
"Msg1 differs from the one the pending handshake at this address answered; starting a new handshake"
);
} else {
// Genuinely pending handshake — resend msg2
// Genuinely pending handshake — resend msg2. Like every
// reply on the rx loop it never dials: with the msg1's
// connection gone, a dial to its address (an inbound
// connection's ephemeral port) would hold the loop for up
// to the connect timeout.
let msg2_bytes = self.find_stored_msg2(existing_link_id);
if let Some(msg2) = msg2_bytes {
if let Some(transport) = self.transports.get(&packet.transport_id) {
match transport.send(&packet.remote_addr, &msg2).await {
match transport.send_existing(&packet.remote_addr, &msg2).await {
Ok(_) => debug!(
remote_addr = %packet.remote_addr,
"Resent msg2 for duplicate msg1"
@@ -469,8 +473,15 @@ impl Node {
machine.set_conn_handshake_msg1(packet.data.clone(), 0);
self.peer_machines.insert(link_id, machine);
// The msg2 goes back on the msg1's connection and never dials: if that
// connection has closed, a dial to its address (for an inbound
// connection, the initiator's ephemeral port) would hold the rx loop
// for up to the connect timeout. The failure tears the leg down below.
if let Some(transport) = self.transports.get(&packet.transport_id) {
match transport.send(&packet.remote_addr, &wire_msg2).await {
match transport
.send_existing(&packet.remote_addr, &wire_msg2)
.await
{
Ok(bytes) => {
debug!(
link_id = %link_id,
@@ -686,14 +697,18 @@ impl Node {
peer.set_remote_epoch(remote_epoch);
}
// Send msg3 before setting pending session
// Send msg3 before setting pending session.
// A reply on the rx loop, so it never dials;
// with the link's connection gone the send
// fails at once and the rekey is abandoned
// below, as on any send failure.
let wire_msg3 =
build_msg3(our_index, header.sender_idx, &msg3_bytes);
let msg3_sent = if let (Some(tid), Some(addr)) =
(transport_id, &remote_addr)
&& let Some(transport) = self.transports.get(&tid)
{
match transport.send(addr, &wire_msg3).await {
match transport.send_existing(addr, &wire_msg3).await {
Ok(_) => {
debug!(
peer = %display_name,
@@ -1117,12 +1132,17 @@ impl Node {
return;
}
// Build and send msg3
// Build and send msg3, on the msg2's connection. A reply on the rx
// loop, so it never dials: with that connection gone the send fails
// at once and the sweep reclaims the handshake, as on any send failure.
let our_index = our_index.unwrap_or(header.receiver_idx);
let wire_msg3 = build_msg3(our_index, header.sender_idx, &msg3_bytes);
if let Some(transport) = self.transports.get(&packet.transport_id) {
match transport.send(&packet.remote_addr, &wire_msg3).await {
match transport
.send_existing(&packet.remote_addr, &wire_msg3)
.await
{
Ok(bytes) => {
debug!(
peer = %self.peer_display_name(&peer_node_addr),
@@ -1886,11 +1906,11 @@ impl Node {
// Not a rekey — duplicate handshake from same epoch. Resend the
// stored msg2 bytes as-is (a driver mechanism: replaying the
// stored frame, not rebuilding it), leaving the active peer
// untouched.
// untouched. On the msg3's connection, never dialing.
if let Some(msg2) = msg2
&& let Some(transport) = self.transports.get(&packet.transport_id)
{
match transport.send(&packet.remote_addr, &msg2).await {
match transport.send_existing(&packet.remote_addr, &msg2).await {
Ok(_) => debug!(
peer = %self.peer_display_name(&peer_node_addr),
"Resent msg2 for duplicate handshake (same epoch)"
+7 -3
View File
@@ -650,16 +650,20 @@ impl Node {
}
for (node_addr, payload) in to_resend {
let (transport_id, remote_addr) = match self.peers.get(&node_addr) {
let (link_id, transport_id, remote_addr) = match self.peers.get(&node_addr) {
Some(p) => match (p.transport_id(), p.current_addr()) {
(Some(tid), Some(addr)) => (tid, addr.clone()),
(Some(tid), Some(addr)) => (p.link_id(), tid, addr.clone()),
_ => continue,
},
None => continue,
};
// A failed send records no resend, so the msg3 stays due and is
// retried next tick, over any connection the failed send started.
let sent = if let Some(transport) = self.transports.get(&transport_id) {
transport.send(&remote_addr, &payload).await.is_ok()
self.send_nowait(transport, link_id, &remote_addr, &payload)
.await
.is_ok()
} else {
false
};
+9 -3
View File
@@ -782,7 +782,7 @@ impl Node {
/// connect-resolution). Anonymous discovery (no `peer_identity`) leaves
/// identity to be learned from the XX msg2, which crystallizes it onto the
/// leg-born machine.
pub(super) async fn start_handshake(
pub(in crate::node) async fn start_handshake(
&mut self,
link_id: LinkId,
transport_id: TransportId,
@@ -877,9 +877,15 @@ impl Node {
// already carry the msg1-prep provenance.
self.peer_machines.insert(link_id, machine);
// Send the wire format handshake message
// Send the wire format handshake message. It never dials: if the
// connection the dial resolved to has gone, the send fails at once
// into the failure path below, after starting a background connect
// to the dial address.
if let Some(transport) = self.transports.get(&transport_id) {
match transport.send(&remote_addr, &wire_msg1).await {
match self
.send_nowait(transport, link_id, &remote_addr, &wire_msg1)
.await
{
Ok(bytes) => {
debug!(
link_id = %link_id,
+4 -1
View File
@@ -4182,8 +4182,11 @@ impl Node {
.get(&transport_id)
.ok_or(NodeError::TransportNotFound(transport_id))?;
// The one caller answers a msg3 on the rx loop, so this never dials:
// a dial to the msg3's address, with its connection gone, would hold
// the loop for up to the connect timeout.
transport
.send(remote_addr, &wire_packet)
.send_existing(remote_addr, &wire_packet)
.await
.map(|_| ())
.map_err(|e| match e {
+1
View File
@@ -24,6 +24,7 @@ mod mmp_chartests;
mod netmon;
mod probe;
mod routing;
mod rx_stall;
mod session;
mod spanning_tree;
mod tcp;
File diff suppressed because it is too large Load Diff