Act on a lookup response only when it answers a lookup we issued

The originator path took any LookupResponse whose request_id was not in
the transit dedup map, so an admitted peer could harvest one genuine
signed response for a target and re-inject it whenever it liked. Each
injection cleared our in-flight lookup, recorded a reachability success
for a target that might be unreachable, refreshed the cached coordinates
for a further full TTL, and flushed our queued packets onto that route at
a moment the sender chose. The signature verify ran before any check that
the response was wanted, so it was the first cost gate on the path.

Record the request_id of every lookup request we send on that target's
pending entry, and drop a response unless it names a target with a lookup
outstanding and carries one of the ids issued for it. The id is fresh
64-bit randomness drawn per attempt and the target signs over it, so a
harvested response is bound to the request it answered and cannot be
redirected or replayed. The check runs before the identity-cache resolve
and before the verify, so a response nobody asked for costs nothing.

initiate_lookup now establishes the pending entry itself rather than
relying on its callers, which keeps "if a request went out, its id is
recorded" true everywhere. The recorded set is capped at eight ids and
evicts the oldest rather than refusing the newest, so a retry ladder
longer than the cap cannot discard the attempt most likely to be
answered; replies to earlier attempts of an outstanding lookup are still
accepted, which is the common case on a link whose round trip exceeds the
first rung.

Drops are counted as resp_unsolicited, in show routing, show metrics and
the fipstop routing pane. The counter has a nonzero floor in healthy
operation: a request is flooded to every qualifying tree peer, so the
duplicate replies land there once the first has been accepted.
This commit is contained in:
Johnathan Corgan
2026-08-23 11:46:15 +01:00
parent 42622c8efe
commit cbe35f1cac
8 changed files with 392 additions and 1 deletions
+25
View File
@@ -1052,6 +1052,31 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
claimed source and destination pairing no honest forwarder could produce.
The drop log line now carries the signal type and the refusal class.
- A discovery lookup response is now acted on only when it answers a lookup
this node actually has outstanding. The originator path took any response
whose `request_id` was not in the transit dedup map, so an admitted peer
could harvest one genuine signed response for a target and re-inject it at
will: each injection cleared the victim's in-flight lookup, recorded a
reachability success for a target that might be unreachable, refreshed the
cached coordinates for a further full TTL, and flushed the victim's queued
packets onto a route at a moment the sender chose. It also reached the
signature verify before any check that the response was wanted, so the
verify was the first cost gate on the path. The node now records the
`request_id` of every lookup request it sends on that target's pending
entry, and a response is dropped unless it names a target with a lookup
outstanding and carries one of the ids issued for it. Because the id is
fresh 64-bit randomness drawn per attempt and the target signs over it, a
harvested response is bound to the request it answered and cannot be
redirected or replayed. The check runs before the identity-cache resolve
and before the signature verify, so a response nobody asked for costs
nothing. Replies to earlier attempts of a still-outstanding lookup are
still accepted, which is the common case on a link whose round trip
exceeds the first rung of the retry ladder. Drops are counted as
`resp_unsolicited`, visible through `show routing`, `show metrics` and the
fipstop routing pane; the counter has a nonzero floor in healthy operation,
because a request is flooded to every qualifying tree peer and the
duplicate replies land there once the first has been accepted.
#### Admission / peer caps
- The Ethernet transport's discovery buffer is now bounded and no longer costs
+1
View File
@@ -206,6 +206,7 @@ fn draw_routing_stats(
("Timed Out", disc("resp_timed_out")),
("Identity Miss", disc("resp_identity_miss")),
("Proof Failed", disc("resp_proof_failed")),
("Unsolicited", disc("resp_unsolicited")),
("Decode Error", disc("resp_decode_error")),
],
));
+2 -1
View File
@@ -29,7 +29,8 @@
"resp_no_route": 0,
"resp_proof_failed": 0,
"resp_received": 0,
"resp_timed_out": 0
"resp_timed_out": 0,
"resp_unsolicited": 0
},
"error_signals": {
"coords_required": 0,
+71
View File
@@ -185,6 +185,32 @@ impl Node {
let target = response.target;
let path_mtu = response.path_mtu;
// Correlate against our own outstanding lookups first. The
// request_id is fresh 64-bit randomness we drew per attempt and
// the target signs over it, so requiring the response to carry
// one we issued for this target is what makes this path
// solicited: an unsolicited or replayed response is dropped
// here, before the identity resolve and before the signature
// verify, and so cannot clear pending state, record a backoff
// success, refresh the coordinate cache, or flush queued
// packets. A duplicate of a response already accepted lands
// here too, which is ordinary and is why this is debug level.
let solicited = self
.pending_lookups
.get(&target)
.is_some_and(|pending| pending.matches(response.request_id));
if !solicited {
self.metrics()
.discovery
.record_reject(DiscoveryReject::RespUnsolicited);
debug!(
request_id = response.request_id,
target = %self.peer_display_name(&target),
"LookupResponse does not match an outstanding request, dropping"
);
return;
}
// Look up the target's public key from identity_cache
let mut prefix = [0u8; 15];
prefix.copy_from_slice(&target.as_bytes()[0..15]);
@@ -487,6 +513,10 @@ impl Node {
/// filters contain the target. Returns the number of peers sent to.
/// The originator does NOT record the request_id in recent_requests,
/// so when the response arrives, it's recognized as "our request".
/// It records the id on the target's pending entry instead, which is
/// what the response path correlates against; recording it here rather
/// than in the callers keeps "if a request went out, its id is
/// recorded" true for every caller.
pub(in crate::node) async fn initiate_lookup(&mut self, target: &NodeAddr, ttl: u8) -> usize {
self.metrics().discovery.req_initiated.inc();
@@ -494,6 +524,12 @@ impl Node {
let origin_coords = self.tree_state().my_coords().clone();
let request = LookupRequest::generate(*target, origin, origin_coords, ttl, 0);
let now_ms = Self::now_ms();
self.pending_lookups
.entry(*target)
.or_insert_with(|| PendingLookup::new(now_ms))
.record(request.request_id);
// Send only to tree peers whose bloom filter contains the target
let peer_addrs: Vec<NodeAddr> = self
.peers
@@ -790,6 +826,18 @@ impl Node {
}
}
/// How many outstanding `request_id`s one pending lookup remembers.
///
/// Bounds the per-target correlator at eight u64s. The retry ladder
/// (`node.discovery.attempt_timeouts_secs`) is operator configuration and
/// can be longer than this, so the recorder evicts the oldest id rather
/// than refusing the newest: dropping the newest would discard the id most
/// likely to be answered and fail a healthy lookup. Raising this costs
/// eight bytes per extra attempt on every pending target and widens the
/// set of ids a late response may still match; lowering it means a reply
/// to an early attempt on a long ladder is dropped as unsolicited.
const MAX_RECORDED_IDS: usize = 8;
/// Tracks a pending discovery lookup with retry state.
pub struct PendingLookup {
/// When the lookup was first initiated.
@@ -798,6 +846,12 @@ pub struct PendingLookup {
pub last_sent_ms: u64,
/// Current attempt number (1 = initial, 2 = first retry, ...).
pub attempt: u8,
/// `request_id`s issued for this target, oldest first, capped at
/// [`MAX_RECORDED_IDS`]. A response is only acted on when it carries
/// one of these, which is what makes the accept path solicited. The
/// entry itself is dropped at ladder timeout, so this set needs no
/// expiry of its own.
pub ids: Vec<u64>,
}
impl PendingLookup {
@@ -806,6 +860,23 @@ impl PendingLookup {
initiated_ms: now_ms,
last_sent_ms: now_ms,
attempt: 1,
ids: Vec::new(),
}
}
/// Remember a `request_id` we just put on the wire for this target.
pub fn record(&mut self, request_id: u64) {
if self.ids.contains(&request_id) {
return;
}
if self.ids.len() >= MAX_RECORDED_IDS {
self.ids.remove(0);
}
self.ids.push(request_id);
}
/// Whether `request_id` is one this node issued for this target.
pub fn matches(&self, request_id: u64) -> bool {
self.ids.contains(&request_id)
}
}
+3
View File
@@ -281,6 +281,7 @@ pub struct DiscoveryMetrics {
pub resp_forwarded: Counter,
pub resp_identity_miss: Counter,
pub resp_proof_failed: Counter,
pub resp_unsolicited: Counter,
pub resp_no_route: Counter,
pub resp_accepted: Counter,
pub resp_timed_out: Counter,
@@ -299,6 +300,7 @@ impl DiscoveryMetrics {
DiscoveryReject::RespDecodeError => self.resp_decode_error.inc(),
DiscoveryReject::RespIdentityMiss => self.resp_identity_miss.inc(),
DiscoveryReject::RespProofFailed => self.resp_proof_failed.inc(),
DiscoveryReject::RespUnsolicited => self.resp_unsolicited.inc(),
DiscoveryReject::RespNoRoute => self.resp_no_route.inc(),
}
}
@@ -325,6 +327,7 @@ impl DiscoveryMetrics {
resp_forwarded: self.resp_forwarded.get(),
resp_identity_miss: self.resp_identity_miss.get(),
resp_proof_failed: self.resp_proof_failed.get(),
resp_unsolicited: self.resp_unsolicited.get(),
resp_no_route: self.resp_no_route.get(),
resp_accepted: self.resp_accepted.get(),
resp_timed_out: self.resp_timed_out.get(),
+9
View File
@@ -136,6 +136,14 @@ pub enum DiscoveryReject {
/// Response proof signature failed verification. Tracked via
/// [`DiscoveryStats::resp_proof_failed`](crate::node::stats::DiscoveryStats).
RespProofFailed,
/// Response arrived on the originator path but carries no
/// `request_id` this node has outstanding for the named target, so
/// it answers no lookup of ours. Expected to be nonzero in healthy
/// operation: the request is flooded to every qualifying tree peer,
/// so duplicate replies land here after the first is accepted.
/// Tracked via
/// [`DiscoveryStats::resp_unsolicited`](crate::node::stats::DiscoveryStats).
RespUnsolicited,
/// Response could not be routed toward the origin: no reverse-path
/// entry for the `request_id` and no greedy tree route to the
/// origin. Tracked via
@@ -371,6 +379,7 @@ mod tests {
DiscoveryReject::RespDecodeError,
DiscoveryReject::RespIdentityMiss,
DiscoveryReject::RespProofFailed,
DiscoveryReject::RespUnsolicited,
];
for v in variants {
let r = RejectReason::Discovery(v);
+1
View File
@@ -334,6 +334,7 @@ pub struct DiscoveryStatsSnapshot {
pub resp_forwarded: u64,
pub resp_identity_miss: u64,
pub resp_proof_failed: u64,
pub resp_unsolicited: u64,
pub resp_no_route: u64,
pub resp_accepted: u64,
pub resp_timed_out: u64,
+280
View File
@@ -84,6 +84,18 @@ async fn test_request_ttl_zero_not_forwarded() {
// Unit Tests — LookupResponse Handler
// ============================================================================
/// Record `request_id` as outstanding for `target`, exactly as
/// `initiate_lookup` does when it puts a request on the wire. The response
/// handler correlates against this, so a unit test that hands the handler a
/// response without it is testing the correlation gate rather than whatever
/// it names.
fn seed_pending_lookup(node: &mut Node, target: crate::NodeAddr, request_id: u64) {
node.pending_lookups
.entry(target)
.or_insert_with(|| handlers::discovery::PendingLookup::new(Node::now_ms()))
.record(request_id);
}
#[tokio::test]
async fn test_response_decode_error() {
let mut node = make_node();
@@ -107,6 +119,8 @@ async fn test_response_originator_caches_route() {
// Register target identity in cache so verification can find it
node.register_identity(target, target_identity.pubkey_full());
seed_pending_lookup(&mut node, target, 555);
// Create a valid response with a real proof signature (includes coords)
let proof_data = LookupResponse::proof_bytes(555, &target, &coords);
let proof = target_identity.sign(&proof_data);
@@ -183,6 +197,8 @@ async fn test_response_proof_verification_success() {
// Register target in identity_cache
node.register_identity(target, target_identity.pubkey_full());
seed_pending_lookup(&mut node, target, 700);
// Sign with correct proof_bytes (including coords)
let proof_data = LookupResponse::proof_bytes(700, &target, &coords);
let proof = target_identity.sign(&proof_data);
@@ -218,6 +234,8 @@ async fn test_response_proof_verification_failure() {
node.register_identity(target, target_identity.pubkey_full());
// Sign with a DIFFERENT identity (wrong key)
seed_pending_lookup(&mut node, target, 701);
let wrong_identity = Identity::generate();
let proof_data = LookupResponse::proof_bytes(701, &target, &coords);
let proof = wrong_identity.sign(&proof_data);
@@ -251,6 +269,8 @@ async fn test_response_identity_cache_miss() {
// Do NOT register target in identity_cache
seed_pending_lookup(&mut node, target, 702);
let proof_data = LookupResponse::proof_bytes(702, &target, &coords);
let proof = target_identity.sign(&proof_data);
@@ -285,6 +305,8 @@ async fn test_response_coord_substitution_detected() {
// Register target in identity_cache
node.register_identity(target, target_identity.pubkey_full());
seed_pending_lookup(&mut node, target, 703);
// Sign proof with real coords
let proof_data = LookupResponse::proof_bytes(703, &target, &real_coords);
let proof = target_identity.sign(&proof_data);
@@ -305,6 +327,255 @@ async fn test_response_coord_substitution_detected() {
);
}
/// Build a signed LookupResponse body for `target_identity` over
/// `request_id`, ready to hand to `handle_lookup_response`.
fn signed_response_body(
target_identity: &Identity,
request_id: u64,
coords: &TreeCoordinate,
) -> Vec<u8> {
let target = *target_identity.node_addr();
let proof_data = LookupResponse::proof_bytes(request_id, &target, coords);
let proof = target_identity.sign(&proof_data);
LookupResponse::new(request_id, target, coords.clone(), proof).encode()[1..].to_vec()
}
/// Register `target_identity` and return its address and a plausible
/// coordinate for it, the shared preamble of the correlation tests.
fn register_lookup_target(node: &mut Node, target_identity: &Identity) -> TreeCoordinate {
let target = *target_identity.node_addr();
node.register_identity(target, target_identity.pubkey_full());
TreeCoordinate::from_addrs(vec![target, make_node_addr(0xF0)]).unwrap()
}
fn wall_clock_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
#[tokio::test]
async fn test_unsolicited_lookup_response_is_dropped_before_proof_verification() {
// Any admitted peer can hand us a correctly signed response for a target
// we never asked about. Accepting it lets that peer clear our pending
// state, refresh a cache entry's TTL and flush our queued packets at a
// moment it picks, so the response must not be acted on at all.
let mut node = make_node();
let from = make_node_addr(0xAA);
let target_identity = Identity::generate();
let target = *target_identity.node_addr();
let coords = register_lookup_target(&mut node, &target_identity);
let body = signed_response_body(&target_identity, 900, &coords);
assert!(
node.pending_lookups.is_empty(),
"precondition: this node has no lookup outstanding for anything"
);
node.handle_lookup_response(&from, &body).await;
assert!(
!node.coord_cache().contains(&target, wall_clock_ms()),
"a response answering no request of ours must not reach the coordinate cache"
);
assert_eq!(
node.metrics().discovery.resp_accepted.get(),
0,
"an unsolicited response must not count as accepted"
);
assert_eq!(
node.metrics().discovery.resp_unsolicited.get(),
1,
"the drop must be visible on a counter, not only in a log"
);
assert_eq!(
node.metrics().discovery.resp_proof_failed.get(),
0,
"the drop must happen before the signature verify, so the verify is not a cost gate"
);
}
#[tokio::test]
async fn test_lookup_response_with_a_request_id_we_never_issued_is_dropped() {
// Correlating on the target alone would leave the attack open: there is
// no inbound limiter on responses, so a peer can spray a harvested one
// and land inside any window in which we happen to be looking that
// target up. The id must match too.
let mut node = make_node();
let from = make_node_addr(0xAA);
let target_identity = Identity::generate();
let target = *target_identity.node_addr();
let coords = register_lookup_target(&mut node, &target_identity);
node.initiate_lookup(&target, 5).await;
let issued = node.pending_lookups.get(&target).unwrap().ids.clone();
assert_eq!(issued.len(), 1, "precondition: one attempt went out");
let body = signed_response_body(&target_identity, issued[0] ^ 1, &coords);
node.handle_lookup_response(&from, &body).await;
assert!(
!node.coord_cache().contains(&target, wall_clock_ms()),
"a response bearing an id we never issued must not reach the coordinate cache"
);
assert!(
node.pending_lookups.contains_key(&target),
"it must not cancel the lookup that is genuinely outstanding"
);
assert_eq!(node.metrics().discovery.resp_unsolicited.get(), 1);
}
#[tokio::test]
async fn test_a_response_matching_a_pending_attempt_is_accepted_and_clears_the_pending_lookup() {
// The healthy path. A fix that reds a legitimate lookup is no use, and
// this is the test that catches it.
let mut node = make_node();
let from = make_node_addr(0xAA);
let target_identity = Identity::generate();
let target = *target_identity.node_addr();
let coords = register_lookup_target(&mut node, &target_identity);
node.initiate_lookup(&target, 5).await;
let issued = node.pending_lookups.get(&target).unwrap().ids[0];
let body = signed_response_body(&target_identity, issued, &coords);
node.handle_lookup_response(&from, &body).await;
assert_eq!(
node.coord_cache().get(&target, wall_clock_ms()),
Some(&coords),
"a response to our own outstanding request must be cached"
);
assert!(
!node.pending_lookups.contains_key(&target),
"accepting it must clear the pending lookup"
);
assert_eq!(node.metrics().discovery.resp_accepted.get(), 1);
assert_eq!(node.metrics().discovery.resp_unsolicited.get(), 0);
}
#[tokio::test]
async fn test_a_late_response_for_an_earlier_retry_attempt_is_still_accepted() {
// Each retry draws a fresh id, and on any link with more than a second
// of round trip the reply to an earlier attempt is the common case. A
// correlator that remembered only the newest id would drop it.
let mut node = make_node();
let from = make_node_addr(0xAA);
let target_identity = Identity::generate();
let target = *target_identity.node_addr();
let coords = register_lookup_target(&mut node, &target_identity);
node.initiate_lookup(&target, 5).await;
node.initiate_lookup(&target, 5).await;
let issued = node.pending_lookups.get(&target).unwrap().ids.clone();
assert_eq!(issued.len(), 2, "precondition: two attempts, two ids");
let body = signed_response_body(&target_identity, issued[0], &coords);
node.handle_lookup_response(&from, &body).await;
assert_eq!(
node.coord_cache().get(&target, wall_clock_ms()),
Some(&coords),
"the first attempt's id is still ours and its answer must be accepted"
);
}
#[tokio::test]
async fn test_a_second_genuine_response_after_the_first_is_accepted_is_dropped() {
// The request is flooded to every qualifying tree peer, so duplicate
// replies are routine. They are dropped at the correlation gate, which
// gives the unsolicited counter a nonzero floor in healthy operation:
// it is not by itself a sign of attack traffic.
let mut node = make_node();
let from = make_node_addr(0xAA);
let target_identity = Identity::generate();
let target = *target_identity.node_addr();
let coords = register_lookup_target(&mut node, &target_identity);
node.initiate_lookup(&target, 5).await;
let issued = node.pending_lookups.get(&target).unwrap().ids[0];
let body = signed_response_body(&target_identity, issued, &coords);
node.handle_lookup_response(&from, &body).await;
node.handle_lookup_response(&from, &body).await;
assert_eq!(
node.coord_cache().get(&target, wall_clock_ms()),
Some(&coords),
"the value written by the first response must still be there"
);
assert_eq!(
node.metrics().discovery.resp_accepted.get(),
1,
"only the first of the two answers our request"
);
assert_eq!(
node.metrics().discovery.resp_unsolicited.get(),
1,
"the duplicate is counted, which is why the counter has a healthy floor"
);
}
#[tokio::test]
async fn test_a_validly_signed_response_for_a_retired_lookup_is_dropped() {
// The pending entry's lifetime is what bounds how stale an accepted
// coordinate can be. Once the retry ladder is exhausted and the entry
// goes, a transit node holding the genuine reply can no longer deliver
// it late.
let mut node = make_node();
let from = make_node_addr(0xAA);
let target_identity = Identity::generate();
let target = *target_identity.node_addr();
let coords = register_lookup_target(&mut node, &target_identity);
node.initiate_lookup(&target, 5).await;
let issued = node.pending_lookups.get(&target).unwrap().ids[0];
// Drive the whole ladder: three retries, then the final timeout.
let mut now_ms = Node::now_ms();
for _ in 0..4 {
now_ms += 100_000;
node.check_pending_lookups(now_ms).await;
}
assert!(
!node.pending_lookups.contains_key(&target),
"precondition: the ladder retired the lookup"
);
let body = signed_response_body(&target_identity, issued, &coords);
node.handle_lookup_response(&from, &body).await;
assert!(
!node.coord_cache().contains(&target, wall_clock_ms()),
"a reply to a retired lookup must not install a coordinate"
);
assert_eq!(node.metrics().discovery.resp_unsolicited.get(), 1);
}
#[test]
fn pending_lookup_id_set_evicts_the_oldest_id_rather_than_refusing_the_newest() {
// The retry ladder is operator configuration and can be longer than the
// recorded-id cap. Refusing the newest id would discard the attempt most
// likely to be answered and fail a healthy lookup.
let mut pending = handlers::discovery::PendingLookup::new(0);
for id in 0..12u64 {
pending.record(id);
}
assert!(
pending.matches(11),
"the newest attempt's id must always be remembered"
);
assert!(!pending.matches(0), "the oldest id is the one evicted");
assert_eq!(pending.ids.len(), 8, "the set stays bounded");
}
// ============================================================================
// Unit Tests — RecentRequest Expiry
// ============================================================================
@@ -886,6 +1157,8 @@ async fn test_originator_stores_path_mtu_in_cache() {
node.register_identity(target, target_identity.pubkey_full());
seed_pending_lookup(&mut node, target, 800);
let proof_data = LookupResponse::proof_bytes(800, &target, &coords);
let proof = target_identity.sign(&proof_data);
@@ -930,6 +1203,8 @@ async fn test_originator_ignores_sub_floor_path_mtu_but_still_caches_coords() {
node.register_identity(target, target_identity.pubkey_full());
seed_pending_lookup(&mut node, target, 801);
let proof_data = LookupResponse::proof_bytes(801, &target, &coords);
let proof = target_identity.sign(&proof_data);
@@ -980,6 +1255,8 @@ async fn test_actionable_lookup_response_path_mtu_does_not_bump_below_floor_coun
node.register_identity(target, target_identity.pubkey_full());
seed_pending_lookup(&mut node, target, 802);
let proof_data = LookupResponse::proof_bytes(802, &target, &coords);
let proof = target_identity.sign(&proof_data);
@@ -1021,6 +1298,8 @@ async fn test_originator_lookup_response_keeps_tighter_path_mtu_lookup() {
let target_fips = crate::FipsAddress::from_node_addr(&target);
node.path_mtu_lookup_insert(target_fips, 1280);
seed_pending_lookup(&mut node, target, 800);
let proof_data = LookupResponse::proof_bytes(800, &target, &coords);
let proof = target_identity.sign(&proof_data);
@@ -1055,6 +1334,7 @@ fn make_verified_lookup_response(
let coords = TreeCoordinate::from_addrs(vec![target, root]).unwrap();
node.register_identity(target, target_identity.pubkey_full());
seed_pending_lookup(node, target, request_id);
let proof_data = LookupResponse::proof_bytes(request_id, &target, &coords);
let proof = target_identity.sign(&proof_data);