Merge refactor-sans-io: reconcile discovery sans-IO core with the FMP v2 delta

Forward-merge the master-side sans-IO refactor (discovery migration + no_std
reductions, collapsed to one commit) into the next-side branch. The discovery
decision core meets next's FMP profile/LookupRequest delta here: next's four
transit-forward predicates (Leaf no-forward, Full-profile, min_mtu, tree/fallback)
are folded into the pure core planners via an extended RoutingView seam
(node_is_leaf / peer_is_full / peer_meets_mtu, new ForwardOutcome::LeafNoForward),
and next's v2 LookupRequest API delta (origin_coords removed, tlv_entries added)
is reconciled across the discovery test tree, with next's four TLV wire tests
ported into the relocated wire test module.

The master-side branch now also carries the no_std+alloc reductions (BTreeMap,
alloc::sync::Arc, backoff-reset log moved to the shell, extern crate alloc),
which come through cleanly on the pilot-only core files.

Validated: cargo fmt / clippy (-D warnings) / test --lib all green, including
next's own transit MTU-pruning integration tests driving the refactored core
through the full shell path.
This commit is contained in:
Johnathan Corgan
2026-07-05 22:01:06 +00:00
18 changed files with 2357 additions and 1139 deletions
-376
View File
@@ -1,376 +0,0 @@
//! Discovery protocol rate limiting and backoff.
//!
//! Two complementary mechanisms:
//!
//! - **`DiscoveryBackoff`** (originator-side, optional): Exponential
//! suppression of fresh lookups after the per-attempt sequence in
//! `node.discovery.attempt_timeouts_secs` has been exhausted.
//! **Disabled by default** (base/cap = 0); the per-attempt sequence
//! is the only retry pacing in the standard configuration. Reset on
//! topology changes (parent change, new peer, first RTT, reconnection).
//!
//! - **`DiscoveryForwardRateLimiter`** (transit-side): Per-target minimum
//! interval for forwarded requests. Defense-in-depth against misbehaving
//! nodes generating fresh request_ids at high rate.
use crate::NodeAddr;
use std::collections::HashMap;
use std::time::{Duration, Instant};
// ============================================================================
// Originator-side: Discovery Backoff
// ============================================================================
/// Default base backoff after first lookup failure. `0` = disabled.
const DEFAULT_BACKOFF_BASE_SECS: u64 = 0;
/// Default maximum backoff cap. `0` = disabled.
const DEFAULT_BACKOFF_MAX_SECS: u64 = 0;
/// Backoff multiplier per consecutive failure.
const BACKOFF_MULTIPLIER: u64 = 2;
/// Exponential backoff for failed discovery lookups.
///
/// Tracks targets whose lookups have timed out and suppresses
/// re-initiation with increasing delays. Cleared on topology changes.
pub struct DiscoveryBackoff {
/// Maps target → (suppress_until, consecutive_failures).
entries: HashMap<NodeAddr, BackoffEntry>,
/// Base backoff duration (first failure).
base: Duration,
/// Maximum backoff cap.
max: Duration,
}
struct BackoffEntry {
/// Don't re-initiate until this instant.
suppress_until: Instant,
/// Consecutive failures (drives exponential backoff).
failures: u32,
}
impl DiscoveryBackoff {
/// Create with default parameters (disabled — base/cap = 0).
pub fn new() -> Self {
Self::with_params(DEFAULT_BACKOFF_BASE_SECS, DEFAULT_BACKOFF_MAX_SECS)
}
/// Create with custom base and max backoff in seconds.
pub fn with_params(base_secs: u64, max_secs: u64) -> Self {
Self {
entries: HashMap::new(),
base: Duration::from_secs(base_secs),
max: Duration::from_secs(max_secs),
}
}
/// Check if a lookup for this target is suppressed.
///
/// Returns true if the target is in backoff and should not be
/// looked up yet.
pub fn is_suppressed(&self, target: &NodeAddr) -> bool {
if let Some(entry) = self.entries.get(target) {
Instant::now() < entry.suppress_until
} else {
false
}
}
/// Record a lookup failure (timeout) for a target.
///
/// Increments the failure count and sets the next suppression
/// window using exponential backoff.
pub fn record_failure(&mut self, target: &NodeAddr) {
let now = Instant::now();
let failures = self.entries.get(target).map_or(0, |e| e.failures) + 1;
let backoff_secs = self
.base
.as_secs()
.saturating_mul(BACKOFF_MULTIPLIER.saturating_pow(failures.saturating_sub(1)));
let backoff = Duration::from_secs(backoff_secs.min(self.max.as_secs()));
self.entries.insert(
*target,
BackoffEntry {
suppress_until: now + backoff,
failures,
},
);
}
/// Record a successful lookup — remove backoff for this target.
pub fn record_success(&mut self, target: &NodeAddr) {
self.entries.remove(target);
}
/// Clear all backoff entries.
///
/// Called on topology changes that might make previously-unreachable
/// targets reachable (parent change, new peer, first RTT, reconnection).
pub fn reset_all(&mut self) {
self.entries.clear();
}
/// Whether any entries exist.
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
/// Current number of entries.
pub fn entry_count(&self) -> usize {
self.entries.len()
}
/// Get the failure count for a target (for logging).
pub fn failure_count(&self, target: &NodeAddr) -> u32 {
self.entries.get(target).map_or(0, |e| e.failures)
}
#[cfg(test)]
pub fn len(&self) -> usize {
self.entries.len()
}
}
impl Default for DiscoveryBackoff {
fn default() -> Self {
Self::new()
}
}
// ============================================================================
// Transit-side: Discovery Forward Rate Limiter
// ============================================================================
/// Default minimum interval between forwarded lookups for the same target.
const DEFAULT_FORWARD_MIN_INTERVAL: Duration = Duration::from_secs(2);
/// Maximum age of entries before cleanup.
const FORWARD_MAX_AGE: Duration = Duration::from_secs(60);
/// Rate limiter for forwarded discovery requests.
///
/// Tracks the last time a LookupRequest was forwarded for each target
/// and enforces a minimum interval to prevent floods from misbehaving
/// nodes generating fresh request_ids.
pub struct DiscoveryForwardRateLimiter {
last_forwarded: HashMap<NodeAddr, Instant>,
min_interval: Duration,
max_age: Duration,
}
impl DiscoveryForwardRateLimiter {
/// Create with default parameters (2s interval).
pub fn new() -> Self {
Self {
last_forwarded: HashMap::new(),
min_interval: DEFAULT_FORWARD_MIN_INTERVAL,
max_age: FORWARD_MAX_AGE,
}
}
/// Create with a custom minimum interval.
pub fn with_interval(min_interval: Duration) -> Self {
Self {
last_forwarded: HashMap::new(),
min_interval,
max_age: FORWARD_MAX_AGE,
}
}
/// Check if we should forward a lookup for this target.
///
/// Returns true if enough time has passed since the last forward
/// for this target. Updates internal state when returning true.
pub fn should_forward(&mut self, target: &NodeAddr) -> bool {
let now = Instant::now();
if let Some(&last) = self.last_forwarded.get(target)
&& now.duration_since(last) < self.min_interval
{
return false;
}
self.last_forwarded.insert(*target, now);
self.cleanup(now);
true
}
/// Replace the minimum interval (e.g., set to zero to disable).
#[cfg(test)]
pub fn set_interval(&mut self, interval: Duration) {
self.min_interval = interval;
}
/// Remove entries older than max_age.
fn cleanup(&mut self, now: Instant) {
self.last_forwarded
.retain(|_, &mut last| now.duration_since(last) < self.max_age);
}
#[cfg(test)]
pub fn len(&self) -> usize {
self.last_forwarded.len()
}
}
impl Default for DiscoveryForwardRateLimiter {
fn default() -> Self {
Self::new()
}
}
// ============================================================================
// Tests
// ============================================================================
#[cfg(test)]
mod tests {
use super::*;
use std::thread;
fn addr(val: u8) -> NodeAddr {
let mut bytes = [0u8; 16];
bytes[0] = val;
NodeAddr::from_bytes(bytes)
}
// --- DiscoveryBackoff tests ---
#[test]
fn test_backoff_not_suppressed_initially() {
let backoff = DiscoveryBackoff::new();
assert!(!backoff.is_suppressed(&addr(1)));
}
#[test]
fn test_backoff_suppressed_after_failure() {
// Backoff is opt-in; exercise the suppression path with explicit params.
let mut backoff = DiscoveryBackoff::with_params(30, 300);
backoff.record_failure(&addr(1));
assert!(backoff.is_suppressed(&addr(1)));
// Different target not affected
assert!(!backoff.is_suppressed(&addr(2)));
}
#[test]
fn test_backoff_cleared_on_success() {
let mut backoff = DiscoveryBackoff::with_params(30, 300);
backoff.record_failure(&addr(1));
assert!(backoff.is_suppressed(&addr(1)));
backoff.record_success(&addr(1));
assert!(!backoff.is_suppressed(&addr(1)));
}
#[test]
fn test_backoff_reset_all() {
let mut backoff = DiscoveryBackoff::new();
backoff.record_failure(&addr(1));
backoff.record_failure(&addr(2));
assert_eq!(backoff.len(), 2);
backoff.reset_all();
assert_eq!(backoff.len(), 0);
assert!(!backoff.is_suppressed(&addr(1)));
}
#[test]
fn test_backoff_exponential() {
let mut backoff = DiscoveryBackoff::with_params(1, 300);
// First failure: 1s backoff
backoff.record_failure(&addr(1));
assert_eq!(backoff.failure_count(&addr(1)), 1);
// Second failure: 2s backoff
backoff.record_failure(&addr(1));
assert_eq!(backoff.failure_count(&addr(1)), 2);
// Third failure: 4s backoff
backoff.record_failure(&addr(1));
assert_eq!(backoff.failure_count(&addr(1)), 3);
}
#[test]
fn test_backoff_expires() {
let mut backoff = DiscoveryBackoff::with_params(0, 0);
backoff.record_failure(&addr(1));
// With 0s backoff, should not be suppressed
assert!(!backoff.is_suppressed(&addr(1)));
}
#[test]
fn test_backoff_capped() {
let mut backoff = DiscoveryBackoff::with_params(1, 10);
// Record many failures
for _ in 0..20 {
backoff.record_failure(&addr(1));
}
// Backoff should be capped at max (10s), not overflow
let entry = backoff.entries.get(&addr(1)).unwrap();
let remaining = entry.suppress_until.duration_since(Instant::now());
assert!(remaining <= Duration::from_secs(11));
}
// --- DiscoveryForwardRateLimiter tests ---
#[test]
fn test_forward_first_allowed() {
let mut limiter = DiscoveryForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1)));
}
#[test]
fn test_forward_rapid_rate_limited() {
let mut limiter = DiscoveryForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1)));
assert!(!limiter.should_forward(&addr(1)));
assert!(!limiter.should_forward(&addr(1)));
}
#[test]
fn test_forward_different_targets_independent() {
let mut limiter = DiscoveryForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1)));
assert!(limiter.should_forward(&addr(2)));
assert!(!limiter.should_forward(&addr(1)));
assert!(!limiter.should_forward(&addr(2)));
}
#[test]
fn test_forward_allowed_after_interval() {
let mut limiter = DiscoveryForwardRateLimiter::with_interval(Duration::from_millis(100));
assert!(limiter.should_forward(&addr(1)));
thread::sleep(Duration::from_millis(110));
assert!(limiter.should_forward(&addr(1)));
}
#[test]
fn test_forward_cleanup_removes_old() {
let mut limiter = DiscoveryForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1)));
assert!(limiter.should_forward(&addr(2)));
assert_eq!(limiter.len(), 2);
let future = Instant::now() + Duration::from_secs(61);
limiter.cleanup(future);
assert_eq!(limiter.len(), 0);
}
#[test]
fn test_forward_cleanup_preserves_recent() {
let mut limiter = DiscoveryForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1)));
assert_eq!(limiter.len(), 1);
limiter.cleanup(Instant::now());
assert_eq!(limiter.len(), 1);
}
}
+380 -402
View File
@@ -5,15 +5,56 @@
//! bloom filter contains the target. TTL and request_id dedup provide
//! safety bounds.
use crate::node::Node;
use crate::node::reject::DiscoveryReject;
use crate::node::{Node, RecentRequest};
use crate::protocol::{LookupRequest, LookupResponse};
use crate::proto::discovery::{DiscoveryAction, LookupRequest, LookupResponse};
use crate::transport::{TransportAddr, TransportId};
use crate::{NodeAddr, PeerIdentity};
use tracing::{debug, info, trace, warn};
const MAX_RECENT_DISCOVERY_REQUESTS: usize = 4096;
/// Shell adapter exposing the live routing tables to the sans-IO discovery
/// core's `RoutingView` read seam. Lives in `node` so it can read `Node`'s
/// private `peers` map and call the crate-private tree/bloom predicates.
///
/// Holding `&Node` whole is fine for the forward path because it does not
/// also need `&mut self.discovery` concurrently. A later commit whose core
/// step needs `&mut discovery` while reading routing state should narrow this
/// to borrow only `peers` + `tree_state` instead of the whole node.
struct NodeRoutingView<'a> {
node: &'a Node,
}
impl crate::proto::discovery::RoutingView for NodeRoutingView<'_> {
fn is_tree_peer(&self, addr: &NodeAddr) -> bool {
self.node.is_tree_peer(addr)
}
fn peers_reaching(&self, target: &NodeAddr) -> Vec<NodeAddr> {
self.node
.peers
.iter()
.filter(|(_, peer)| peer.may_reach(target))
.map(|(addr, _)| *addr)
.collect()
}
fn node_is_leaf(&self) -> bool {
self.node.node_profile() == crate::protocol::NodeProfile::Leaf
}
fn peer_is_full(&self, addr: &NodeAddr) -> bool {
self.node
.peers
.get(addr)
.is_some_and(|peer| peer.peer_profile() == crate::protocol::NodeProfile::Full)
}
fn peer_meets_mtu(&self, addr: &NodeAddr, min_mtu: u16) -> bool {
self.node
.peers
.get(addr)
.is_some_and(|peer| self.node.peer_meets_mtu(peer, min_mtu))
}
}
impl Node {
/// Handle an incoming LookupRequest from a peer.
///
@@ -39,80 +80,71 @@ impl Node {
};
let now_ms = Self::now_ms();
self.purge_expired_requests(now_ms);
// Dedup: drop if we've already seen this request_id.
// Also serves as loop protection — tree routing is loop-free,
// but request_id dedup catches edge cases during tree restructuring.
if self.recent_requests.contains_key(&request.request_id) {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqDuplicate);
debug!(
request_id = request.request_id,
from = %self.peer_display_name(from),
"Duplicate LookupRequest, dropping"
);
return;
}
if self.recent_requests.len() >= MAX_RECENT_DISCOVERY_REQUESTS {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqDedupCacheFull);
debug!(
request_id = request.request_id,
from = %self.peer_display_name(from),
recent_requests = self.recent_requests.len(),
max_recent_requests = MAX_RECENT_DISCOVERY_REQUESTS,
"Discovery request dedup cache full, dropping LookupRequest"
);
return;
}
// Record for reverse-path forwarding and dedup
self.recent_requests
.insert(request.request_id, RecentRequest::new(*from, now_ms));
// Are we the target?
if request.target == *self.node_addr() {
self.metrics().discovery.req_target_is_us.inc();
debug!(
request_id = request.request_id,
origin = %self.peer_display_name(&request.origin),
"We are the lookup target, generating response"
);
self.send_lookup_response(&request).await;
return;
}
// Forward if TTL permits
if request.can_forward() {
// Transit-side rate limit: collapse rapid-fire lookups for the
// same target from misbehaving nodes generating fresh request_ids.
if !self
.discovery_forward_limiter
.should_forward(&request.target)
{
let recent_expiry_ms = self.config().node.discovery.recent_expiry_secs * 1000;
let my_addr = *self.node_addr();
use crate::proto::discovery::RequestOutcome;
match crate::proto::discovery::classify_request(
&mut self.discovery,
&request,
from,
&my_addr,
now_ms,
recent_expiry_ms,
MAX_RECENT_DISCOVERY_REQUESTS,
) {
RequestOutcome::Duplicate => {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqDuplicate);
debug!(
request_id = request.request_id,
from = %self.peer_display_name(from),
"Duplicate LookupRequest, dropping"
);
}
RequestOutcome::DedupCacheFull { len } => {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqDedupCacheFull);
debug!(
request_id = request.request_id,
from = %self.peer_display_name(from),
recent_requests = len,
max_recent_requests = MAX_RECENT_DISCOVERY_REQUESTS,
"Discovery request dedup cache full, dropping LookupRequest"
);
}
RequestOutcome::RespondAsTarget => {
self.metrics().discovery.req_target_is_us.inc();
debug!(
request_id = request.request_id,
origin = %self.peer_display_name(&request.origin),
"We are the lookup target, generating response"
);
self.send_lookup_response(&request).await;
}
RequestOutcome::Forward => {
self.metrics().discovery.req_forwarded.inc();
self.forward_lookup_request(request).await;
}
RequestOutcome::ForwardRateLimited => {
self.metrics().discovery.req_forward_rate_limited.inc();
debug!(
request_id = request.request_id,
target = %self.peer_display_name(&request.target),
"Forward rate limited, suppressing LookupRequest"
);
return;
}
self.metrics().discovery.req_forwarded.inc();
self.forward_lookup_request(request).await;
} else {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqTtlExhausted);
debug!(
request_id = request.request_id,
target = %self.peer_display_name(&request.target),
"LookupRequest TTL exhausted"
);
RequestOutcome::TtlExhausted => {
self.metrics()
.discovery
.record_reject(DiscoveryReject::ReqTtlExhausted);
debug!(
request_id = request.request_id,
target = %self.peer_display_name(&request.target),
"LookupRequest TTL exhausted"
);
}
}
}
@@ -144,150 +176,186 @@ impl Node {
let now_ms = Self::now_ms();
// Check if we forwarded this request (transit node) or originated it
if let Some(recent) = self.recent_requests.get_mut(&response.request_id) {
// Already forwarded a response for this request — drop to
// prevent response routing loops.
if recent.response_forwarded {
match crate::proto::discovery::classify_response(&mut self.discovery, response.request_id) {
crate::proto::discovery::ResponseRoute::AlreadyForwarded => {
// Already forwarded a response for this request — drop to
// prevent response routing loops.
debug!(
request_id = response.request_id,
target = %self.peer_display_name(&response.target),
"Response already forwarded for this request, dropping"
);
return;
}
recent.response_forwarded = true;
crate::proto::discovery::ResponseRoute::Transit { from_peer } => {
// Transit node: reverse-path forward
self.metrics().discovery.resp_forwarded.inc();
// Transit node: reverse-path forward
let from_peer = recent.from_peer;
self.metrics().discovery.resp_forwarded.inc();
// Apply path_mtu min() from the outgoing link's transport MTU
self.apply_outgoing_link_mtu_to_response(&mut response, &from_peer);
// Apply path_mtu min() from the outgoing link's transport MTU
self.apply_outgoing_link_mtu_to_response(&mut response, &from_peer);
debug!(
request_id = response.request_id,
target = %self.peer_display_name(&response.target),
next_hop = %self.peer_display_name(&from_peer),
path_mtu = response.path_mtu,
"Reverse-path forwarding LookupResponse"
);
let encoded = response.encode();
if let Err(e) = self.send_encrypted_link_message(&from_peer, &encoded).await {
debug!(
request_id = response.request_id,
target = %self.peer_display_name(&response.target),
next_hop = %self.peer_display_name(&from_peer),
error = %e,
"Failed to forward LookupResponse"
path_mtu = response.path_mtu,
"Reverse-path forwarding LookupResponse"
);
}
} else {
// We originated this request — verify proof before caching
let target = response.target;
let path_mtu = response.path_mtu;
// Look up the target's public key from identity_cache
let mut prefix = [0u8; 15];
prefix.copy_from_slice(&target.as_bytes()[0..15]);
let target_pubkey = match self.lookup_by_fips_prefix(&prefix) {
Some((_addr, pubkey)) => pubkey,
None => {
let encoded = response.encode();
if let Err(e) = self.send_encrypted_link_message(&from_peer, &encoded).await {
debug!(
next_hop = %self.peer_display_name(&from_peer),
error = %e,
"Failed to forward LookupResponse"
);
}
}
crate::proto::discovery::ResponseRoute::Originator => {
// We originated this request — verify proof before caching
let target = response.target;
let path_mtu = response.path_mtu;
// Look up the target's public key from identity_cache
let mut prefix = [0u8; 15];
prefix.copy_from_slice(&target.as_bytes()[0..15]);
let target_pubkey = match self.lookup_by_fips_prefix(&prefix) {
Some((_addr, pubkey)) => pubkey,
None => {
self.metrics()
.discovery
.record_reject(DiscoveryReject::RespIdentityMiss);
warn!(
request_id = response.request_id,
target = %self.peer_display_name(&target),
"identity_cache miss for lookup target, cannot verify proof"
);
return;
}
};
// Verify the proof signature
let (xonly, _parity) = target_pubkey.x_only_public_key();
let peer_id = PeerIdentity::from_pubkey(xonly);
let proof_data = LookupResponse::proof_bytes(
response.request_id,
&target,
&response.target_coords,
);
if !peer_id.verify(&proof_data, &response.proof) {
self.metrics()
.discovery
.record_reject(DiscoveryReject::RespIdentityMiss);
.record_reject(DiscoveryReject::RespProofFailed);
warn!(
request_id = response.request_id,
target = %self.peer_display_name(&target),
"identity_cache miss for lookup target, cannot verify proof"
"LookupResponse proof verification failed, discarding"
);
return;
}
};
// Verify the proof signature
let (xonly, _parity) = target_pubkey.x_only_public_key();
let peer_id = PeerIdentity::from_pubkey(xonly);
let proof_data =
LookupResponse::proof_bytes(response.request_id, &target, &response.target_coords);
if !peer_id.verify(&proof_data, &response.proof) {
self.metrics()
.discovery
.record_reject(DiscoveryReject::RespProofFailed);
warn!(
self.metrics().discovery.resp_accepted.inc();
info!(
request_id = response.request_id,
target = %self.peer_display_name(&target),
"LookupResponse proof verification failed, discarding"
depth = response.target_coords.depth(),
path_mtu = path_mtu,
"Discovery succeeded, proof verified, route cached"
);
return;
// Apply the accept-side effects: the core clears the success
// state (backoff + pending lookup) and returns the
// cross-subsystem effects for us to drive.
let actions = crate::proto::discovery::on_response_accepted(
&mut self.discovery,
&target,
response.target_coords,
now_ms,
path_mtu,
);
self.drive_response_actions(actions).await;
}
}
}
self.metrics().discovery.resp_accepted.inc();
// Clear backoff on success — target is reachable
self.discovery_backoff.record_success(&target);
info!(
request_id = response.request_id,
target = %self.peer_display_name(&target),
depth = response.target_coords.depth(),
path_mtu = path_mtu,
"Discovery succeeded, proof verified, route cached"
);
self.coord_cache
.insert_with_path_mtu(target, response.target_coords, now_ms, path_mtu);
// Mirror path_mtu into the FipsAddress-keyed read-only lookup
// map used by the TUN reader/writer at TCP MSS clamp time.
let fips_addr = crate::FipsAddress::from_node_addr(&target);
match self.path_mtu_lookup.write() {
Ok(mut map) => {
let prior = map.insert(fips_addr, path_mtu);
debug!(
target = %self.peer_display_name(&target),
fips_addr = %fips_addr,
path_mtu = path_mtu,
prior = ?prior,
map_len = map.len(),
"Wrote path_mtu_lookup from discovery LookupResponse"
);
/// Drive the cross-subsystem effects returned by the discovery core's
/// accept-side planning. Each arm reproduces the original inline effect
/// exactly (same metrics/logs/writes, same order).
async fn drive_response_actions(&mut self, actions: Vec<DiscoveryAction>) {
for action in actions {
match action {
DiscoveryAction::CacheCoords {
target,
coords,
now_ms,
path_mtu,
} => {
self.coord_cache
.insert_with_path_mtu(target, coords, now_ms, path_mtu);
}
Err(e) => {
warn!(
target = %self.peer_display_name(&target),
fips_addr = %fips_addr,
path_mtu = path_mtu,
error = %e,
"path_mtu_lookup write lock poisoned; clamp will not see this update"
);
DiscoveryAction::WritePathMtu { target, path_mtu } => {
// Mirror path_mtu into the FipsAddress-keyed read-only lookup
// map used by the TUN reader/writer at TCP MSS clamp time.
let fips_addr = crate::FipsAddress::from_node_addr(&target);
match self.path_mtu_lookup.write() {
Ok(mut map) => {
let prior = map.insert(fips_addr, path_mtu);
debug!(
target = %self.peer_display_name(&target),
fips_addr = %fips_addr,
path_mtu = path_mtu,
prior = ?prior,
map_len = map.len(),
"Wrote path_mtu_lookup from discovery LookupResponse"
);
}
Err(e) => {
warn!(
target = %self.peer_display_name(&target),
fips_addr = %fips_addr,
path_mtu = path_mtu,
error = %e,
"path_mtu_lookup write lock poisoned; clamp will not see this update"
);
}
}
}
DiscoveryAction::ResetWarmupIfEstablished { target } => {
// If an established session exists, reset the warmup counter.
let n = self.config().node.session.coords_warmup_packets;
if let Some(entry) = self.sessions.get_mut(&target)
&& entry.is_established()
{
entry.set_coords_warmup_remaining(n);
debug!(
dest = %self.peer_display_name(&target),
warmup_packets = n,
"Reset coords warmup after discovery for existing session"
);
}
}
DiscoveryAction::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) {
debug!(
dest = %self.peer_display_name(&target),
queued_packets = packets.len(),
"Retrying queued packets after discovery"
);
self.retry_session_after_discovery(target).await;
}
}
DiscoveryAction::SendLink { peer, bytes } => {
if let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await {
debug!(
peer = %self.peer_display_name(&peer),
error = %e,
"Failed to send discovery link message"
);
}
}
}
// Clean up pending lookup tracking
self.pending_lookups.remove(&target);
// If an established session exists, reset the warmup counter.
let n = self.config().node.session.coords_warmup_packets;
if let Some(entry) = self.sessions.get_mut(&target)
&& entry.is_established()
{
entry.set_coords_warmup_remaining(n);
debug!(
dest = %self.peer_display_name(&target),
warmup_packets = n,
"Reset coords warmup after discovery for existing session"
);
}
// 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) {
debug!(
dest = %self.peer_display_name(&target),
queued_packets = packets.len(),
"Retrying queued packets after discovery"
);
self.retry_session_after_discovery(target).await;
}
}
}
@@ -304,12 +372,16 @@ impl Node {
let mut response =
LookupResponse::new(request.request_id, request.target, our_coords, proof);
// Route toward origin via reverse path.
let next_hop_addr = if let Some(recent) = self.recent_requests.get(&request.request_id) {
recent.from_peer
} else {
// Fallback: try greedy tree routing toward origin
match self.find_next_hop(&request.origin) {
// Route toward origin. The reverse-path decision (the peer the request
// arrived from, recorded in recent_requests) is the sans-IO core's; the
// greedy tree-route fallback is a &mut coord-cache op kept in the shell.
use crate::proto::discovery::ResponseRouteDecision;
let next_hop_addr = match crate::proto::discovery::plan_response_route(
&self.discovery,
request.request_id,
) {
ResponseRouteDecision::ReversePath(peer) => peer,
ResponseRouteDecision::NeedsTreeRoute => match self.find_next_hop(&request.origin) {
Some(peer) => *peer.node_addr(),
None => {
debug!(
@@ -321,7 +393,7 @@ impl Node {
.record_reject(DiscoveryReject::RespNoRoute);
return;
}
}
},
};
// Fold our outgoing-link MTU into path_mtu so the target-edge link
@@ -361,83 +433,58 @@ impl Node {
/// bloom contains the target. This recovers from dead ends caused by
/// stale bloom filters, tree restructuring, or transit node failures.
async fn forward_lookup_request(&mut self, mut request: LookupRequest) {
if !request.forward() {
return;
}
// Leaf nodes don't forward discovery requests
if self.node_profile() == crate::protocol::NodeProfile::Leaf {
return;
}
// Collect full tree peers whose bloom filter contains the target
let min_mtu = request.min_mtu;
let forward_to: Vec<NodeAddr> = self
.peers
.iter()
.filter(|(addr, peer)| {
peer.peer_profile() == crate::protocol::NodeProfile::Full
&& self.is_tree_peer(addr)
&& peer.may_reach(&request.target)
&& self.peer_meets_mtu(peer, min_mtu)
})
.map(|(addr, _)| *addr)
.collect();
// Fallback: if no tree peer matches, try non-tree full bloom-matching peers
let (forward_to, used_fallback) = if forward_to.is_empty() {
let fallback: Vec<NodeAddr> = self
.peers
.iter()
.filter(|(addr, peer)| {
peer.peer_profile() == crate::protocol::NodeProfile::Full
&& !self.is_tree_peer(addr)
&& peer.may_reach(&request.target)
&& self.peer_meets_mtu(peer, min_mtu)
})
.map(|(addr, _)| *addr)
.collect();
if fallback.is_empty() {
// Plan the forward with the sans-IO decision core. The core owns the
// TTL decrement, Leaf suppression, Full+MTU eligibility, tree/fallback
// peer selection, and single-encode fan-out; the shell keeps all
// metrics/logging and drives the sends.
let outcome = {
let rv = NodeRoutingView { node: self };
crate::proto::discovery::plan_forward(&mut request, &rv)
};
match outcome {
crate::proto::discovery::ForwardOutcome::TtlExhausted => {}
crate::proto::discovery::ForwardOutcome::LeafNoForward => {}
crate::proto::discovery::ForwardOutcome::NoPeers => {
self.metrics().discovery.req_no_tree_peer.inc();
trace!(
request_id = request.request_id,
"No eligible peers to forward LookupRequest"
);
return;
}
(fallback, true)
} else {
(forward_to, false)
};
if used_fallback {
self.metrics().discovery.req_fallback_forwarded.inc();
debug!(
request_id = request.request_id,
target = %self.peer_display_name(&request.target),
ttl = request.ttl,
peer_count = forward_to.len(),
"Forwarding LookupRequest via non-tree fallback"
);
} else {
debug!(
request_id = request.request_id,
target = %self.peer_display_name(&request.target),
ttl = request.ttl,
peer_count = forward_to.len(),
"Forwarding LookupRequest"
);
}
let encoded = request.encode();
for peer_addr in forward_to {
if let Err(e) = self.send_encrypted_link_message(&peer_addr, &encoded).await {
debug!(
peer = %self.peer_display_name(&peer_addr),
error = %e,
"Failed to forward LookupRequest to peer"
);
crate::proto::discovery::ForwardOutcome::Forward {
actions,
used_fallback,
} => {
let peer_count = actions.len();
if used_fallback {
self.metrics().discovery.req_fallback_forwarded.inc();
debug!(
request_id = request.request_id,
target = %self.peer_display_name(&request.target),
ttl = request.ttl,
peer_count,
"Forwarding LookupRequest via non-tree fallback"
);
} else {
debug!(
request_id = request.request_id,
target = %self.peer_display_name(&request.target),
ttl = request.ttl,
peer_count,
"Forwarding LookupRequest"
);
}
for action in actions {
if let DiscoveryAction::SendLink { peer, bytes } = action
&& let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await
{
debug!(
peer = %self.peer_display_name(&peer),
error = %e,
"Failed to forward LookupRequest to peer"
);
}
}
}
}
}
@@ -455,20 +502,16 @@ impl Node {
let min_mtu = self.config().tun.mtu();
let request = LookupRequest::generate(*target, origin, ttl, min_mtu);
// Send only to full tree peers whose bloom filter contains the target
let peer_addrs: Vec<NodeAddr> = self
.peers
.iter()
.filter(|(addr, peer)| {
peer.peer_profile() == crate::protocol::NodeProfile::Full
&& self.is_tree_peer(addr)
&& peer.may_reach(target)
&& self.peer_meets_mtu(peer, request.min_mtu)
})
.map(|(addr, _)| *addr)
.collect();
// Tree-peer selection restricted to Full peers meeting min_mtu, plus the
// single encode, live in the sans-IO core. The core keeps the tree-only
// (no non-tree fallback) behavior; the shell drives the sends and keeps
// all metrics/logging.
let actions = {
let rv = NodeRoutingView { node: self };
crate::proto::discovery::plan_initiate(&request, &rv)
};
let peer_count = peer_addrs.len();
let peer_count = actions.len();
debug!(
request_id = request.request_id,
@@ -479,16 +522,12 @@ impl Node {
"Discovery lookup initiated"
);
if peer_count == 0 {
return 0;
}
let encoded = request.encode();
for peer_addr in peer_addrs {
if let Err(e) = self.send_encrypted_link_message(&peer_addr, &encoded).await {
for action in actions {
if let DiscoveryAction::SendLink { peer, bytes } = action
&& let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await
{
debug!(
peer = %self.peer_display_name(&peer_addr),
peer = %self.peer_display_name(&peer),
error = %e,
"Failed to send LookupRequest to peer"
);
@@ -508,54 +547,49 @@ impl Node {
pub(in crate::node) async fn maybe_initiate_lookup(&mut self, dest: &NodeAddr) {
let now_ms = Self::now_ms();
// Dedup: any pending lookup means we are already trying.
if self.pending_lookups.contains_key(dest) {
self.metrics().discovery.req_deduplicated.inc();
debug!(
target_node = %self.peer_display_name(dest),
"Discovery lookup deduplicated, already pending"
);
return;
}
// Optional post-failure suppression. Defaults are 0/0 (inert);
// operators can opt in by setting `node.discovery.backoff_*_secs`.
if self.discovery_backoff.is_suppressed(dest) {
self.metrics().discovery.req_backoff_suppressed.inc();
debug!(
target_node = %self.peer_display_name(dest),
failures = self.discovery_backoff.failure_count(dest),
"Discovery lookup suppressed by backoff"
);
return;
}
// Bloom filter pre-check: if no peer's filter contains the target,
// it's not in the mesh — skip the lookup and record as failure.
// Bloom filter pre-check (view read) BEFORE the core call: if no peer's
// filter contains the target, it's not in the mesh. Reading `self.peers`
// here keeps the `&mut self.discovery` borrow in `initiate_gate` from
// overlapping the immutable peer-table read.
let reachable = self.peers.values().any(|peer| peer.may_reach(dest));
if !reachable {
self.metrics().discovery.req_bloom_miss.inc();
self.discovery_backoff.record_failure(dest);
debug!(
target_node = %self.peer_display_name(dest),
"Discovery skipped, target not in any peer bloom filter"
);
return;
}
self.pending_lookups
.insert(*dest, PendingLookup::new(now_ms));
let ttl = self.config().node.discovery.ttl;
let sent = self.initiate_lookup(dest, ttl).await;
use crate::proto::discovery::InitiateDecision;
match crate::proto::discovery::initiate_gate(&mut self.discovery, dest, now_ms, reachable) {
InitiateDecision::Deduplicated => {
self.metrics().discovery.req_deduplicated.inc();
debug!(
target_node = %self.peer_display_name(dest),
"Discovery lookup deduplicated, already pending"
);
}
InitiateDecision::Suppressed { failures } => {
self.metrics().discovery.req_backoff_suppressed.inc();
debug!(
target_node = %self.peer_display_name(dest),
failures = failures,
"Discovery lookup suppressed by backoff"
);
}
InitiateDecision::BloomMiss => {
self.metrics().discovery.req_bloom_miss.inc();
debug!(
target_node = %self.peer_display_name(dest),
"Discovery skipped, target not in any peer bloom filter"
);
}
InitiateDecision::Proceed => {
let ttl = self.config().node.discovery.ttl;
let sent = self.initiate_lookup(dest, ttl).await;
// If no tree peers had the target, fail immediately
if sent == 0 {
self.pending_lookups.remove(dest);
self.discovery_backoff.record_failure(dest);
debug!(
target_node = %self.peer_display_name(dest),
"Discovery failed, no tree peers with bloom match"
);
// If no tree peers had the target, fail immediately
if sent == 0 {
crate::proto::discovery::initiate_failed(&mut self.discovery, dest, now_ms);
debug!(
target_node = %self.peer_display_name(dest),
"Discovery failed, no tree peers with bloom match"
);
}
}
}
}
@@ -570,53 +604,24 @@ impl Node {
/// - Otherwise: declare the destination unreachable, drop queued packets,
/// and emit ICMPv6 destination-unreachable for each.
pub(in crate::node) async fn check_pending_lookups(&mut self, now_ms: u64) {
let timeouts = self.config().node.discovery.attempt_timeouts_secs.clone();
let max_attempts = timeouts.len() as u8;
let attempt_timeouts = self.config().node.discovery.attempt_timeouts_secs.clone();
let outcome =
crate::proto::discovery::poll_pending(&mut self.discovery, now_ms, &attempt_timeouts);
// Collect targets needing action
let mut to_retry: Vec<NodeAddr> = Vec::new();
let mut to_timeout: Vec<NodeAddr> = Vec::new();
for (&target, entry) in &self.pending_lookups {
let attempt_idx = (entry.attempt as usize).saturating_sub(1);
let attempt_timeout_ms = timeouts.get(attempt_idx).copied().unwrap_or(0) * 1000;
if now_ms.saturating_sub(entry.last_sent_ms) >= attempt_timeout_ms {
if entry.attempt >= max_attempts {
to_timeout.push(target);
} else {
to_retry.push(target);
}
for (target, attempt) in outcome.retries {
let ttl = self.config().node.discovery.ttl;
let sent = self.initiate_lookup(&target, ttl).await;
if sent > 0 {
debug!(
target_node = %self.peer_display_name(&target),
attempt = attempt,
"Discovery retry sent"
);
}
}
// Process retries
for target in to_retry {
if let Some(entry) = self.pending_lookups.get_mut(&target) {
entry.attempt += 1;
entry.last_sent_ms = now_ms;
let attempt = entry.attempt;
let ttl = self.config().node.discovery.ttl;
let sent = self.initiate_lookup(&target, ttl).await;
if sent > 0 {
debug!(
target_node = %self.peer_display_name(&target),
attempt = attempt,
"Discovery retry sent"
);
}
}
}
// Process timeouts
for addr in to_timeout {
for (addr, failures) in outcome.timeouts {
self.metrics().discovery.resp_timed_out.inc();
self.pending_lookups.remove(&addr);
// Record failure for optional backoff
self.discovery_backoff.record_failure(&addr);
let failures = self.discovery_backoff.failure_count(&addr);
let queued = self.pending_tun_packets.remove(&addr);
let pkt_count = queued.as_ref().map_or(0, |p| p.len());
info!(
@@ -635,12 +640,12 @@ impl Node {
/// Reset discovery backoff on topology changes.
pub(in crate::node) fn reset_discovery_backoff(&mut self) {
if !self.discovery_backoff.is_empty() {
let cleared = self.discovery.reset_backoff();
if cleared > 0 {
debug!(
entries = self.discovery_backoff.entry_count(),
entries = cleared,
"Resetting discovery backoff on topology change"
);
self.discovery_backoff.reset_all();
}
}
@@ -666,13 +671,6 @@ impl Node {
}
}
/// Remove expired entries from the recent_requests cache.
fn purge_expired_requests(&mut self, current_time_ms: u64) {
let expiry_ms = self.config().node.discovery.recent_expiry_secs * 1000;
self.recent_requests
.retain(|_, entry| !entry.is_expired(current_time_ms, expiry_ms));
}
/// Min-fold our outgoing-link MTU into a LookupResponse's `path_mtu`.
///
/// Used at both transit-side reverse-path forward and at the target's
@@ -759,23 +757,3 @@ impl Node {
}
}
}
/// Tracks a pending discovery lookup with retry state.
pub struct PendingLookup {
/// When the lookup was first initiated.
pub initiated_ms: u64,
/// When the last attempt was sent.
pub last_sent_ms: u64,
/// Current attempt number (1 = initial, 2 = first retry, ...).
pub attempt: u8,
}
impl PendingLookup {
pub fn new(now_ms: u64) -> Self {
Self {
initiated_ms: now_ms,
last_sent_ms: now_ms,
attempt: 1,
}
}
}
+15 -61
View File
@@ -9,7 +9,6 @@ mod bloom;
pub(crate) mod context;
#[cfg(unix)]
pub(crate) mod decrypt_worker;
mod discovery_rate_limit;
#[cfg(unix)]
pub(crate) mod encrypt_worker;
mod handlers;
@@ -29,7 +28,6 @@ mod tests;
mod tree;
pub(crate) mod wire;
use self::discovery_rate_limit::{DiscoveryBackoff, DiscoveryForwardRateLimiter};
use self::rate_limit::HandshakeRateLimiter;
use self::reloadable::Reloadable;
use self::routing_error_rate_limit::RoutingErrorRateLimiter;
@@ -61,6 +59,7 @@ use crate::bloom::{BloomFilter, BloomState};
use crate::cache::CoordCache;
use crate::node::session::SessionEntry;
use crate::peer::{ActivePeer, PeerConnection};
use crate::proto::discovery::{Discovery, DiscoveryBackoff, DiscoveryForwardRateLimiter};
use crate::protocol::NodeProfile;
#[cfg(unix)]
use crate::transport::ethernet::EthernetTransport;
@@ -232,39 +231,6 @@ pub struct UpdatePeersOutcome {
pub unchanged: usize,
}
/// Recent request tracking for dedup and reverse-path forwarding.
///
/// When a LookupRequest is forwarded through a node, the node stores the
/// request_id and which peer sent it. When the corresponding LookupResponse
/// arrives, it's forwarded back to that peer (reverse-path forwarding).
/// The `response_forwarded` flag prevents response routing loops.
#[derive(Clone, Debug)]
pub(crate) struct RecentRequest {
/// The peer who sent this request to us.
pub(crate) from_peer: NodeAddr,
/// When we received this request (Unix milliseconds).
pub(crate) timestamp_ms: u64,
/// Whether we've already forwarded a response for this request.
/// Prevents response routing loops when convergent request paths
/// create bidirectional entries in recent_requests.
pub(crate) response_forwarded: bool,
}
impl RecentRequest {
pub(crate) fn new(from_peer: NodeAddr, timestamp_ms: u64) -> Self {
Self {
from_peer,
timestamp_ms,
response_forwarded: false,
}
}
/// Check if this entry has expired (older than expiry_ms).
pub(crate) fn is_expired(&self, current_time_ms: u64, expiry_ms: u64) -> bool {
current_time_ms.saturating_sub(self.timestamp_ms) > expiry_ms
}
}
/// Key for addr_to_link reverse lookup.
type AddrKey = (TransportId, TransportAddr);
@@ -335,9 +301,6 @@ pub struct Node {
// === Routing ===
/// Address -> coordinates cache (from session setup and discovery).
coord_cache: CoordCache,
/// Recent discovery requests (dedup + reverse-path forwarding).
/// Maps request_id → RecentRequest.
recent_requests: HashMap<u64, RecentRequest>,
/// Per-destination path MTU lookup, keyed by FipsAddress (mirrors
/// `coord_cache.entries[*].path_mtu`). Sync read-only access from
/// the TUN reader/writer threads at TCP MSS clamp time so the
@@ -385,10 +348,11 @@ pub struct Node {
/// Packets queued while waiting for session establishment.
/// Keyed by destination NodeAddr, bounded per-dest and total.
pending_tun_packets: HashMap<NodeAddr, VecDeque<Vec<u8>>>,
// === Pending Discovery Lookups ===
/// Tracks in-flight discovery lookups. Maps target NodeAddr to the
/// initiation timestamp (Unix ms). Prevents duplicate flood queries.
pending_lookups: HashMap<NodeAddr, handlers::discovery::PendingLookup>,
// === Discovery ===
/// Discovery-subsystem state: recent-request dedup cache, in-flight
/// lookups, originator-side backoff, and transit-side forward limiter.
discovery: Discovery,
// === Counters ===
/// Next link ID to allocate.
@@ -475,10 +439,6 @@ pub struct Node {
routing_error_rate_limiter: RoutingErrorRateLimiter,
/// Rate limiter for source-side CoordsRequired/PathBroken responses.
coords_response_rate_limiter: RoutingErrorRateLimiter,
/// Backoff for failed discovery lookups (originator-side).
discovery_backoff: DiscoveryBackoff,
/// Rate limiter for forwarded discovery requests (transit-side).
discovery_forward_limiter: DiscoveryForwardRateLimiter,
// === Pending Transport Connects ===
/// Links waiting for transport-level connection establishment before
@@ -680,7 +640,6 @@ impl Node {
tree_state,
bloom_state,
coord_cache,
recent_requests: HashMap::new(),
transports: HashMap::new(),
transport_drops: HashMap::new(),
links: HashMap::new(),
@@ -692,7 +651,6 @@ impl Node {
sessions: HashMap::new(),
identity_cache: HashMap::new(),
pending_tun_packets: HashMap::new(),
pending_lookups: HashMap::new(),
next_link_id: 1,
next_transport_id: 1,
stats: stats::NodeStats::new(),
@@ -727,9 +685,9 @@ impl Node {
coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval(
std::time::Duration::from_millis(coords_response_interval_ms),
),
discovery_backoff: DiscoveryBackoff::with_params(backoff_base_secs, backoff_max_secs),
discovery_forward_limiter: DiscoveryForwardRateLimiter::with_interval(
std::time::Duration::from_secs(forward_min_interval_secs),
discovery: Discovery::new(
DiscoveryBackoff::with_params(backoff_base_secs, backoff_max_secs),
DiscoveryForwardRateLimiter::with_interval_ms(forward_min_interval_secs * 1000),
),
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
@@ -843,7 +801,6 @@ impl Node {
tree_state,
bloom_state,
coord_cache,
recent_requests: HashMap::new(),
transports: HashMap::new(),
transport_drops: HashMap::new(),
links: HashMap::new(),
@@ -855,7 +812,6 @@ impl Node {
sessions: HashMap::new(),
identity_cache: HashMap::new(),
pending_tun_packets: HashMap::new(),
pending_lookups: HashMap::new(),
next_link_id: 1,
next_transport_id: 1,
stats: stats::NodeStats::new(),
@@ -890,8 +846,7 @@ impl Node {
coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval(
std::time::Duration::from_millis(coords_response_interval_ms),
),
discovery_backoff: DiscoveryBackoff::new(),
discovery_forward_limiter: DiscoveryForwardRateLimiter::new(),
discovery: Discovery::new(DiscoveryBackoff::new(), DiscoveryForwardRateLimiter::new()),
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
nostr_discovery: None,
@@ -2485,8 +2440,7 @@ impl Node {
/// Disable the discovery forward rate limiter (for tests).
#[cfg(test)]
pub(crate) fn disable_discovery_forward_rate_limit(&mut self) {
self.discovery_forward_limiter
.set_interval(std::time::Duration::ZERO);
self.discovery.forward_limiter.set_interval_ms(0);
}
#[cfg(test)]
@@ -2598,19 +2552,19 @@ impl Node {
/// Number of pending discovery lookups.
pub fn pending_lookup_count(&self) -> usize {
self.pending_lookups.len()
self.discovery.pending_lookups.len()
}
/// Iterate over pending discovery lookups for diagnostics.
pub fn pending_lookups_iter(
&self,
) -> impl Iterator<Item = (&NodeAddr, &handlers::discovery::PendingLookup)> {
self.pending_lookups.iter()
) -> impl Iterator<Item = (&NodeAddr, &crate::proto::discovery::PendingLookup)> {
self.discovery.pending_lookups.iter()
}
/// Number of recent discovery requests tracked.
pub fn recent_request_count(&self) -> usize {
self.recent_requests.len()
self.discovery.recent_requests.len()
}
/// Count of destinations with queued TUN packets awaiting session setup.
+29 -23
View File
@@ -5,8 +5,7 @@
//! response routing.
use super::*;
use crate::node::RecentRequest;
use crate::protocol::{LookupRequest, LookupResponse};
use crate::proto::discovery::{LookupRequest, LookupResponse, RecentRequest};
use crate::tree::TreeCoordinate;
use spanning_tree::{
cleanup_nodes, generate_random_edges, lock_large_network_test, process_available_packets,
@@ -23,7 +22,7 @@ async fn test_request_decode_error() {
let from = make_node_addr(0xAA);
// Too-short payload: should log error and return without panic
node.handle_lookup_request(&from, &[0x00; 5]).await;
assert!(node.recent_requests.is_empty());
assert!(node.discovery.recent_requests.is_empty());
}
#[tokio::test]
@@ -38,11 +37,11 @@ async fn test_request_dedup() {
// First request: accepted
node.handle_lookup_request(&from, payload).await;
assert_eq!(node.recent_requests.len(), 1);
assert_eq!(node.discovery.recent_requests.len(), 1);
// Duplicate request: dropped
node.handle_lookup_request(&from, payload).await;
assert_eq!(node.recent_requests.len(), 1);
assert_eq!(node.discovery.recent_requests.len(), 1);
}
#[tokio::test]
@@ -59,7 +58,7 @@ async fn test_request_target_is_self() {
// Should succeed without panic (response send will fail silently
// since we have no peers to route toward origin)
node.handle_lookup_request(&from, payload).await;
assert!(node.recent_requests.contains_key(&777));
assert!(node.discovery.recent_requests.contains_key(&777));
}
#[tokio::test]
@@ -74,7 +73,7 @@ async fn test_request_ttl_zero_not_forwarded() {
node.handle_lookup_request(&from, payload).await;
// Request recorded, but not forwarded (TTL=0, and no peers anyway)
assert!(node.recent_requests.contains_key(&666));
assert!(node.discovery.recent_requests.contains_key(&666));
}
// ============================================================================
@@ -112,7 +111,7 @@ async fn test_response_originator_caches_route() {
let payload = &response.encode()[1..]; // skip msg_type
// No entry in recent_requests for 555 → we're the originator
assert!(!node.recent_requests.contains_key(&555));
assert!(!node.discovery.recent_requests.contains_key(&555));
node.handle_lookup_response(&from, payload).await;
@@ -146,7 +145,8 @@ async fn test_response_transit_needs_recent_request() {
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
node.recent_requests
node.discovery
.recent_requests
.insert(444, RecentRequest::new(make_node_addr(0xDD), now_ms));
// Handle response — should try to reverse-path forward to 0xDD
@@ -316,14 +316,16 @@ async fn test_recent_request_expiry() {
.as_millis() as u64;
// Insert an old request (11 seconds ago)
node.recent_requests
node.discovery
.recent_requests
.insert(123, RecentRequest::new(make_node_addr(1), now_ms - 11_000));
// Insert a recent request
node.recent_requests
node.discovery
.recent_requests
.insert(456, RecentRequest::new(make_node_addr(2), now_ms));
assert_eq!(node.recent_requests.len(), 2);
assert_eq!(node.discovery.recent_requests.len(), 2);
// Trigger purge via a new lookup request
let target = make_node_addr(0xBB);
@@ -334,9 +336,9 @@ async fn test_recent_request_expiry() {
.await;
// Old entry (123) should be purged, recent entry (456) and new entry (789) kept
assert!(!node.recent_requests.contains_key(&123));
assert!(node.recent_requests.contains_key(&456));
assert!(node.recent_requests.contains_key(&789));
assert!(!node.discovery.recent_requests.contains_key(&123));
assert!(node.discovery.recent_requests.contains_key(&456));
assert!(node.discovery.recent_requests.contains_key(&789));
}
// ============================================================================
@@ -373,7 +375,7 @@ async fn test_request_forwarding_two_node() {
// Node1 should have recorded the request
assert!(
nodes[1].node.recent_requests.contains_key(&42),
nodes[1].node.discovery.recent_requests.contains_key(&42),
"Node 1 should have recorded the forwarded request"
);
@@ -441,13 +443,13 @@ async fn test_request_three_node_chain() {
// Node1 should have been a transit node (has the request_id in recent_requests)
assert!(
!nodes[1].node.recent_requests.is_empty(),
!nodes[1].node.discovery.recent_requests.is_empty(),
"Node 1 should have recorded the forwarded request"
);
// Node2 should have received the request (it's the target)
assert!(
!nodes[2].node.recent_requests.is_empty(),
!nodes[2].node.discovery.recent_requests.is_empty(),
"Node 2 should have received the request"
);
@@ -494,7 +496,7 @@ async fn test_request_dedup_convergent_paths() {
// Node2 (the target) must have received the request
assert!(
nodes[2].node.recent_requests.contains_key(&300),
nodes[2].node.discovery.recent_requests.contains_key(&300),
"Node 2 (target) should have received the request"
);
@@ -1159,8 +1161,8 @@ async fn test_open_discovery_sweep_queues_eligible_skips_filtered() {
#[tokio::test]
async fn test_check_pending_lookups_default_sequence_unreachable() {
use crate::bloom::BloomFilter;
use crate::node::handlers::discovery::PendingLookup;
use crate::peer::ActivePeer;
use crate::proto::discovery::PendingLookup;
use crate::transport::LinkId;
use std::sync::mpsc;
@@ -1228,7 +1230,8 @@ async fn test_check_pending_lookups_default_sequence_unreachable() {
// Inject a PendingLookup directly: attempt=1, last_sent_ms=0. This
// mirrors the post-condition of a successful `maybe_initiate_lookup`
// at t=0 without depending on wall-clock-derived `Self::now_ms()`.
node.pending_lookups
node.discovery
.pending_lookups
.insert(target_addr, PendingLookup::new(0));
let baseline_initiated = node.metrics().discovery.req_initiated.get();
@@ -1238,6 +1241,7 @@ async fn test_check_pending_lookups_default_sequence_unreachable() {
node.check_pending_lookups(1100).await;
{
let entry = node
.discovery
.pending_lookups
.get(&target_addr)
.expect("still pending");
@@ -1254,6 +1258,7 @@ async fn test_check_pending_lookups_default_sequence_unreachable() {
node.check_pending_lookups(3100).await;
{
let entry = node
.discovery
.pending_lookups
.get(&target_addr)
.expect("still pending");
@@ -1270,6 +1275,7 @@ async fn test_check_pending_lookups_default_sequence_unreachable() {
node.check_pending_lookups(7100).await;
{
let entry = node
.discovery
.pending_lookups
.get(&target_addr)
.expect("still pending");
@@ -1285,7 +1291,7 @@ async fn test_check_pending_lookups_default_sequence_unreachable() {
// --- Just-before-final: at t=15099ms the 8s window is not yet reached ---
node.check_pending_lookups(15_099).await;
assert!(
node.pending_lookups.contains_key(&target_addr),
node.discovery.pending_lookups.contains_key(&target_addr),
"8s window not yet expired: pending_lookup must persist"
);
assert_eq!(
@@ -1308,7 +1314,7 @@ async fn test_check_pending_lookups_default_sequence_unreachable() {
// Pending lookup is dropped.
assert!(
!node.pending_lookups.contains_key(&target_addr),
!node.discovery.pending_lookups.contains_key(&target_addr),
"final timeout must remove the pending_lookups entry"
);
// resp_timed_out counter ticked.