Deliver a native datagram to an uncached, unrouted key after discovery

A native API flow to a node that was not a peer, had no session and had
not been resolved through DNS could not be reached when no route to it
was known. The session attempt failed for want of coordinates, so the
key the client supplied was never cached; discovery ran, but its answer
was dropped because the proof could not be verified without that key.
Even with the key cached, an accepted answer retried the session only
for queued TUN packets, so the held native datagram was never sent.

Cache the flow's key before trying the route, only when no key is
already cached for that destination, and retry the session after
discovery when native datagrams are held for the target as well as when
TUN packets are.
This commit is contained in:
Johnathan Corgan
2026-10-06 00:26:00 +00:00
parent f857278d4c
commit 55fb970d99
5 changed files with 234 additions and 8 deletions
+9 -5
View File
@@ -483,13 +483,17 @@ impl Node {
}
}
LookupAction::RetryQueuedPackets { target } => {
// If we have pending TUN packets for this target, retry session
// initiation. The coord_cache now has coords, so find_next_hop()
// should succeed.
if let Some(packets) = self.pending_tun_packets.get(&target) {
// If we hold TUN packets or native datagrams for this target,
// retry session initiation. The coord_cache now has coords, so
// find_next_hop() should succeed. Native datagrams are held in
// their own queue, and nothing else re-initiates for them.
let packets = self.pending_tun_packets.get(&target).map(|q| q.len());
let native = self.pending_native.get(&target).map(|q| q.len());
if packets.is_some() || native.is_some() {
debug!(
dest = %self.peer_display_name(&target),
queued_packets = packets.len(),
queued_packets = packets.unwrap_or(0),
queued_native = native.unwrap_or(0),
"Retrying queued packets after discovery"
);
self.retry_session_after_discovery(target).await;
+15
View File
@@ -157,6 +157,21 @@ impl Node {
// hashes the DH output x-only and forces even parity in the XK
// premessage, so both parities derive the same material.
let pubkey = peer.public_key(secp256k1::Parity::Even);
// Cache the key the client supplied before trying the route. If there
// is none, the lookup below verifies its answer against this cache, and
// `initiate_session` registers the key only after a send succeeds, so
// without this the answer is dropped. Only on a miss, the same rule
// `cache_session_identity` follows: when no route is known, a key
// already cached from DNS, a handshake or a configured peer keeps its
// own parity. The presence check refreshes that entry, so it is not
// the oldest one when the answer arrives.
let mut prefix = [0u8; 15];
prefix.copy_from_slice(&key.peer.as_bytes()[0..15]);
if self.lookup_by_fips_prefix(&prefix).is_none() {
self.register_identity(key.peer, pubkey);
}
if let Err(error) = self.initiate_session(key.peer, pubkey).await {
debug!(
peer = %self.peer_display_name(&key.peer),
+3 -2
View File
@@ -3383,8 +3383,9 @@ impl Node {
/// Retry session initiation after discovery provided coordinates.
///
/// Called when a LookupResponse arrives and we have pending TUN packets
/// for the discovered target. The coord_cache now has coords, so
/// `find_next_hop()` should succeed and the SessionSetup can be sent.
/// or native datagrams for the discovered target. The coord_cache now
/// has coords, so `find_next_hop()` should succeed and the SessionSetup
/// can be sent.
pub(in crate::node) async fn retry_session_after_discovery(&mut self, dest_addr: NodeAddr) {
// Look up the destination's public key from the identity cache
let mut prefix = [0u8; 15];
+205
View File
@@ -5,12 +5,15 @@
//! response routing.
use super::*;
use crate::native::link::NativeMessage;
use crate::native::registry::FlowKey;
use crate::proto::lookup::{LookupRequest, LookupResponse, RecentRequest};
use crate::proto::stp::TreeCoordinate;
use spanning_tree::{
cleanup_nodes, generate_random_edges, lock_large_network_test, process_available_packets,
run_tree_test, verify_tree_convergence,
};
use tokio::sync::{mpsc, oneshot};
// ============================================================================
// Unit Tests — LookupRequest Handler
@@ -2026,3 +2029,205 @@ async fn test_check_pending_lookups_default_sequence_unreachable() {
assert_eq!(icmp_frame[6], 58, "next_header must be IPPROTO_ICMPV6 (58)");
assert_eq!(icmp_frame[40], 1, "ICMPv6 type 1 = Destination Unreachable");
}
// ============================================================================
// Native API — first send to an uncached, unrouted key
// ============================================================================
/// The registry's port for the listener the native tests bind.
const NATIVE_PORT: u16 = 7000;
#[tokio::test]
async fn a_native_send_to_an_unrouted_key_caches_that_key_so_discovery_can_verify_the_answer() {
let mut node = make_node();
let dest = crate::Identity::generate();
let dest_addr = *dest.node_addr();
let key = FlowKey {
peer: dest_addr,
remote: NATIVE_PORT,
local: 7001,
};
assert!(
!node.has_cached_identity(&dest_addr),
"precondition: the destination's key must not be cached"
);
node.handle_native_outbound(key, dest.pubkey(), b"hello".to_vec())
.await;
assert!(
node.has_cached_identity(&dest_addr),
"a native send with no route must cache the flow's key, or the lookup it \
starts cannot verify its answer"
);
}
#[tokio::test]
async fn a_native_send_does_not_replace_a_key_already_cached_for_that_destination() {
let mut node = make_node();
// A fixed key with odd parity, so the native path's even-parity lift would
// be visible if it replaced the cached entry. Searched over a bounded set
// of fixed secrets so the test is deterministic.
let dest = (1u8..=64)
.map(|n| {
let mut secret = [0u8; 32];
secret[31] = n;
crate::Identity::from_secret_bytes(&secret).expect("a small non-zero secret is valid")
})
.find(|id| id.pubkey_full().x_only_public_key().1 == secp256k1::Parity::Odd)
.expect("one of 64 fixed secrets has an odd-parity public key");
let dest_addr = *dest.node_addr();
let original = dest.pubkey_full();
node.register_identity(dest_addr, original);
let key = FlowKey {
peer: dest_addr,
remote: NATIVE_PORT,
local: 7001,
};
node.handle_native_outbound(key, dest.pubkey(), b"hello".to_vec())
.await;
let mut prefix = [0u8; 15];
prefix.copy_from_slice(&dest_addr.as_bytes()[0..15]);
let (_, cached) = node
.lookup_by_fips_prefix(&prefix)
.expect("the destination's key must still be cached");
assert_eq!(
cached, original,
"a native send must not replace a cached key with its even-parity lift"
);
}
#[tokio::test]
async fn a_native_first_send_to_an_uncached_unrouted_key_is_delivered_after_discovery() {
// Topology: node0 — node1 — node2. Node 0 knows nothing of node 2: no
// peer link, no cached key, no cached coordinates, no session.
let edges = vec![(0, 1), (1, 2)];
let mut nodes = run_tree_test(3, &edges, false).await;
let node0_xonly = nodes[0].node.identity().pubkey();
let node2_addr = *nodes[2].node.node_addr();
let node2_xonly = nodes[2].node.identity().pubkey();
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
assert!(
nodes[0].node.get_peer(&node2_addr).is_none(),
"precondition: node 2 must not be a peer of node 0"
);
assert!(
!nodes[0].node.has_cached_identity(&node2_addr),
"precondition: node 0 must not have node 2's key cached"
);
assert!(
!nodes[0].node.coord_cache().contains(&node2_addr, now_ms),
"precondition: node 0 must not have node 2's coordinates cached"
);
assert!(
nodes[0].node.get_session(&node2_addr).is_none(),
"precondition: node 0 must have no session to node 2"
);
// Node 2 listens on the port node 0 will send to.
let (arrivals_tx, mut arrivals) = mpsc::channel(8);
let (reply_tx, mut reply_rx) = oneshot::channel();
nodes[2].node.handle_native(NativeMessage::Listen {
port: Some(NATIVE_PORT),
arrivals: arrivals_tx,
reply: reply_tx,
});
assert_eq!(
reply_rx
.try_recv()
.expect("Listen must answer synchronously"),
Ok(NATIVE_PORT),
"node 2 must hold the listener port"
);
let lookup = &nodes[0].node.metrics().lookup;
let initiated_before = lookup.req_initiated.get();
let accepted_before = lookup.resp_accepted.get();
let identity_miss_before = lookup.resp_identity_miss.get();
let payload = b"first native datagram".to_vec();
let flow = FlowKey {
peer: node2_addr,
remote: NATIVE_PORT,
local: 7001,
};
nodes[0]
.node
.handle_native_outbound(flow, node2_xonly, payload.clone())
.await;
// Bound the drive on the arrival itself: the initiator flushes the held
// datagram in the same call that sends msg3, and packets move one hop per
// pass, so the datagram is still in transit when node 0 first shows its
// session established. Never await `recv()` here: node 2's registry holds
// the sender, so a broken fix would hang rather than fail.
let mut arrival = None;
for _ in 0..40 {
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
if let Ok(seen) = arrivals.try_recv() {
arrival = Some(seen);
break;
}
}
let lookup = &nodes[0].node.metrics().lookup;
assert!(
lookup.req_initiated.get() > initiated_before,
"the unrouted native send must start discovery"
);
assert_eq!(
lookup.resp_identity_miss.get(),
identity_miss_before,
"discovery's answer must not be dropped for want of the target's key"
);
assert!(
lookup.resp_accepted.get() > accepted_before,
"discovery's answer must be verified and accepted"
);
assert!(
nodes[0]
.node
.get_session(&node2_addr)
.is_some_and(|entry| entry.is_established()),
"discovery must lead to an established session for the held native datagram"
);
assert_eq!(
nodes[0].node.metrics().native.sent_datagrams.get(),
1,
"node 0 must send the held native datagram once the session is up"
);
let arrival = arrival.expect("node 2's listener must see the native flow arrive");
assert_eq!(
arrival.pubkey, node0_xonly,
"the arrival must carry node 0's authenticated key"
);
let (sink_tx, _sink_rx) = mpsc::channel(8);
let (accept_tx, mut accept_rx) = oneshot::channel();
nodes[2].node.handle_native(NativeMessage::Accept {
flow: arrival.flow,
sink: sink_tx,
reply: accept_tx,
});
let accepted = accept_rx
.try_recv()
.expect("Accept must answer synchronously")
.expect("node 2 must accept the announced flow");
assert_eq!(
accepted.held,
vec![payload],
"the accepted flow must hold exactly the datagram node 0 sent"
);
cleanup_nodes(&mut nodes).await;
}
+2 -1
View File
@@ -53,7 +53,8 @@ pub(crate) enum LookupAction {
},
/// Reset the coords-warmup counter if an established session exists.
ResetWarmupIfEstablished { target: NodeAddr },
/// Retry queued TUN packets for the target if any are pending.
/// Retry queued TUN packets or native datagrams for the target if any are
/// pending.
RetryQueuedPackets { target: NodeAddr },
}