proto: rename the mesh discovery subsystem module from discovery to lookup

Rename src/proto/discovery to src/proto/lookup and bring the module's naming
onto the lookup stem, matching its concept: a mesh lookup of a node's
coordinates from its pubkey, sent into the mesh as a bloom-filter-guided
multicast request that returns a unicast response with the coordinates.

- the module directory and its declaration (proto::discovery -> proto::lookup)
- exported types: DiscoveryAction -> LookupAction, DiscoveryBackoff ->
  LookupBackoff, DiscoveryForwardRateLimiter -> LookupForwardRateLimiter,
  MAX_RECENT_DISCOVERY_REQUESTS -> MAX_RECENT_LOOKUP_REQUESTS, and the Discovery
  state struct -> Lookup
- internal terminology: doc comments, the "overlay-lookup" phrasing, the
  empty_discovery/suppressing_discovery test helpers, and the disc
  parameter/variable all take the lookup names
- the already-lookup-named wire types (LookupRequest/LookupResponse) are unchanged

Behavior-neutral: no wire bytes or decision logic change. Two references are
intentionally kept as "discovery": the still-named node::handlers::discovery
shell module and the node.discovery.* config keys, which belong to the broader
disambiguation of the shell, config, and metric surfaces still to come.
This commit is contained in:
Johnathan Corgan
2026-07-08 21:54:24 +00:00
parent 2b009196b5
commit e3e03f6a5d
15 changed files with 295 additions and 275 deletions
+2 -2
View File
@@ -73,8 +73,8 @@ pub use proto::stp::TreeAnnounce;
// Re-export bloom wire types (relocated from protocol:: to proto::bloom) // Re-export bloom wire types (relocated from protocol:: to proto::bloom)
pub use proto::bloom::FilterAnnounce; pub use proto::bloom::FilterAnnounce;
// Re-export discovery wire types (relocated from protocol:: to proto::discovery) // Re-export discovery wire types (relocated from protocol:: to proto::lookup)
pub use proto::discovery::{LookupRequest, LookupResponse}; pub use proto::lookup::{LookupRequest, LookupResponse};
// Re-export routing wire types (relocated from protocol:: to proto::routing) // Re-export routing wire types (relocated from protocol:: to proto::routing)
pub use proto::routing::{ pub use proto::routing::{
+31 -31
View File
@@ -7,8 +7,8 @@
use crate::node::Node; use crate::node::Node;
use crate::node::reject::DiscoveryReject; use crate::node::reject::DiscoveryReject;
use crate::proto::discovery::{ use crate::proto::lookup::{
DiscoveryAction, LookupRequest, LookupResponse, MAX_RECENT_DISCOVERY_REQUESTS, LookupAction, LookupRequest, LookupResponse, MAX_RECENT_LOOKUP_REQUESTS,
}; };
use crate::transport::{TransportAddr, TransportId}; use crate::transport::{TransportAddr, TransportId};
use crate::{NodeAddr, PeerIdentity}; use crate::{NodeAddr, PeerIdentity};
@@ -26,7 +26,7 @@ struct NodeRoutingView<'a> {
node: &'a Node, node: &'a Node,
} }
impl crate::proto::discovery::RoutingView for NodeRoutingView<'_> { impl crate::proto::lookup::RoutingView for NodeRoutingView<'_> {
fn is_tree_peer(&self, addr: &NodeAddr) -> bool { fn is_tree_peer(&self, addr: &NodeAddr) -> bool {
self.node.is_tree_peer(addr) self.node.is_tree_peer(addr)
} }
@@ -67,15 +67,15 @@ impl Node {
let now_ms = Self::now_ms(); let now_ms = Self::now_ms();
let recent_expiry_ms = self.config().node.discovery.recent_expiry_secs * 1000; let recent_expiry_ms = self.config().node.discovery.recent_expiry_secs * 1000;
let my_addr = *self.node_addr(); let my_addr = *self.node_addr();
use crate::proto::discovery::RequestOutcome; use crate::proto::lookup::RequestOutcome;
match crate::proto::discovery::classify_request( match crate::proto::lookup::classify_request(
&mut self.discovery, &mut self.discovery,
&request, &request,
from, from,
&my_addr, &my_addr,
now_ms, now_ms,
recent_expiry_ms, recent_expiry_ms,
MAX_RECENT_DISCOVERY_REQUESTS, MAX_RECENT_LOOKUP_REQUESTS,
) { ) {
RequestOutcome::Duplicate => { RequestOutcome::Duplicate => {
self.metrics() self.metrics()
@@ -95,7 +95,7 @@ impl Node {
request_id = request.request_id, request_id = request.request_id,
from = %self.peer_display_name(from), from = %self.peer_display_name(from),
recent_requests = len, recent_requests = len,
max_recent_requests = MAX_RECENT_DISCOVERY_REQUESTS, max_recent_requests = MAX_RECENT_LOOKUP_REQUESTS,
"Discovery request dedup cache full, dropping LookupRequest" "Discovery request dedup cache full, dropping LookupRequest"
); );
} }
@@ -161,8 +161,8 @@ impl Node {
let now_ms = Self::now_ms(); let now_ms = Self::now_ms();
// Check if we forwarded this request (transit node) or originated it // Check if we forwarded this request (transit node) or originated it
match crate::proto::discovery::classify_response(&mut self.discovery, response.request_id) { match crate::proto::lookup::classify_response(&mut self.discovery, response.request_id) {
crate::proto::discovery::ResponseRoute::AlreadyForwarded => { crate::proto::lookup::ResponseRoute::AlreadyForwarded => {
// Already forwarded a response for this request — drop to // Already forwarded a response for this request — drop to
// prevent response routing loops. // prevent response routing loops.
debug!( debug!(
@@ -171,7 +171,7 @@ impl Node {
"Response already forwarded for this request, dropping" "Response already forwarded for this request, dropping"
); );
} }
crate::proto::discovery::ResponseRoute::Transit { from_peer } => { crate::proto::lookup::ResponseRoute::Transit { from_peer } => {
// Transit node: reverse-path forward // Transit node: reverse-path forward
self.metrics().discovery.resp_forwarded.inc(); self.metrics().discovery.resp_forwarded.inc();
@@ -195,7 +195,7 @@ impl Node {
); );
} }
} }
crate::proto::discovery::ResponseRoute::Originator => { crate::proto::lookup::ResponseRoute::Originator => {
// We originated this request — verify proof before caching // We originated this request — verify proof before caching
let target = response.target; let target = response.target;
let path_mtu = response.path_mtu; let path_mtu = response.path_mtu;
@@ -251,7 +251,7 @@ impl Node {
// Apply the accept-side effects: the core clears the success // Apply the accept-side effects: the core clears the success
// state (backoff + pending lookup) and returns the // state (backoff + pending lookup) and returns the
// cross-subsystem effects for us to drive. // cross-subsystem effects for us to drive.
let actions = crate::proto::discovery::on_response_accepted( let actions = crate::proto::lookup::on_response_accepted(
&mut self.discovery, &mut self.discovery,
&target, &target,
response.target_coords, response.target_coords,
@@ -266,10 +266,10 @@ impl Node {
/// Drive the cross-subsystem effects returned by the discovery core's /// Drive the cross-subsystem effects returned by the discovery core's
/// accept-side planning. Each arm reproduces the original inline effect /// accept-side planning. Each arm reproduces the original inline effect
/// exactly (same metrics/logs/writes, same order). /// exactly (same metrics/logs/writes, same order).
async fn drive_response_actions(&mut self, actions: Vec<DiscoveryAction>) { async fn drive_response_actions(&mut self, actions: Vec<LookupAction>) {
for action in actions { for action in actions {
match action { match action {
DiscoveryAction::CacheCoords { LookupAction::CacheCoords {
target, target,
coords, coords,
now_ms, now_ms,
@@ -278,7 +278,7 @@ impl Node {
self.coord_cache self.coord_cache
.insert_with_path_mtu(target, coords, now_ms, path_mtu); .insert_with_path_mtu(target, coords, now_ms, path_mtu);
} }
DiscoveryAction::WritePathMtu { target, path_mtu } => { LookupAction::WritePathMtu { target, path_mtu } => {
// Mirror path_mtu into the FipsAddress-keyed read-only lookup // Mirror path_mtu into the FipsAddress-keyed read-only lookup
// map used by the TUN reader/writer at TCP MSS clamp time. // map used by the TUN reader/writer at TCP MSS clamp time.
let fips_addr = crate::FipsAddress::from_node_addr(&target); let fips_addr = crate::FipsAddress::from_node_addr(&target);
@@ -305,7 +305,7 @@ impl Node {
} }
} }
} }
DiscoveryAction::ResetWarmupIfEstablished { target } => { LookupAction::ResetWarmupIfEstablished { target } => {
// If an established session exists, reset the warmup counter. // If an established session exists, reset the warmup counter.
let n = self.config().node.session.coords_warmup_packets; let n = self.config().node.session.coords_warmup_packets;
if let Some(entry) = self.sessions.get_mut(&target) if let Some(entry) = self.sessions.get_mut(&target)
@@ -319,7 +319,7 @@ impl Node {
); );
} }
} }
DiscoveryAction::RetryQueuedPackets { target } => { LookupAction::RetryQueuedPackets { target } => {
// If we have pending TUN packets for this target, retry session // If we have pending TUN packets for this target, retry session
// initiation. The coord_cache now has coords, so find_next_hop() // initiation. The coord_cache now has coords, so find_next_hop()
// should succeed. // should succeed.
@@ -332,7 +332,7 @@ impl Node {
self.retry_session_after_discovery(target).await; self.retry_session_after_discovery(target).await;
} }
} }
DiscoveryAction::SendLink { peer, bytes } => { LookupAction::SendLink { peer, bytes } => {
if let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await { if let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await {
debug!( debug!(
peer = %self.peer_display_name(&peer), peer = %self.peer_display_name(&peer),
@@ -360,8 +360,8 @@ impl Node {
// Route toward origin. The reverse-path decision (the peer the request // Route toward origin. The reverse-path decision (the peer the request
// arrived from, recorded in recent_requests) is the sans-IO core's; the // 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. // greedy tree-route fallback is a &mut coord-cache op kept in the shell.
use crate::proto::discovery::ResponseRouteDecision; use crate::proto::lookup::ResponseRouteDecision;
let next_hop_addr = match crate::proto::discovery::plan_response_route( let next_hop_addr = match crate::proto::lookup::plan_response_route(
&self.discovery, &self.discovery,
request.request_id, request.request_id,
) { ) {
@@ -423,18 +423,18 @@ impl Node {
// fan-out; the shell keeps all metrics/logging and drives the sends. // fan-out; the shell keeps all metrics/logging and drives the sends.
let outcome = { let outcome = {
let rv = NodeRoutingView { node: self }; let rv = NodeRoutingView { node: self };
crate::proto::discovery::plan_forward(&mut request, &rv) crate::proto::lookup::plan_forward(&mut request, &rv)
}; };
match outcome { match outcome {
crate::proto::discovery::ForwardOutcome::TtlExhausted => {} crate::proto::lookup::ForwardOutcome::TtlExhausted => {}
crate::proto::discovery::ForwardOutcome::NoPeers => { crate::proto::lookup::ForwardOutcome::NoPeers => {
self.metrics().discovery.req_no_tree_peer.inc(); self.metrics().discovery.req_no_tree_peer.inc();
trace!( trace!(
request_id = request.request_id, request_id = request.request_id,
"No eligible peers to forward LookupRequest" "No eligible peers to forward LookupRequest"
); );
} }
crate::proto::discovery::ForwardOutcome::Forward { crate::proto::lookup::ForwardOutcome::Forward {
actions, actions,
used_fallback, used_fallback,
} => { } => {
@@ -458,7 +458,7 @@ impl Node {
); );
} }
for action in actions { for action in actions {
if let DiscoveryAction::SendLink { peer, bytes } = action if let LookupAction::SendLink { peer, bytes } = action
&& let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await && let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await
{ {
debug!( debug!(
@@ -494,7 +494,7 @@ impl Node {
// the shell drives the sends and keeps all metrics/logging. // the shell drives the sends and keeps all metrics/logging.
let actions = { let actions = {
let rv = NodeRoutingView { node: self }; let rv = NodeRoutingView { node: self };
crate::proto::discovery::plan_initiate(&request, &rv) crate::proto::lookup::plan_initiate(&request, &rv)
}; };
let peer_count = actions.len(); let peer_count = actions.len();
@@ -509,7 +509,7 @@ impl Node {
); );
for action in actions { for action in actions {
if let DiscoveryAction::SendLink { peer, bytes } = action if let LookupAction::SendLink { peer, bytes } = action
&& let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await && let Err(e) = self.send_encrypted_link_message(&peer, &bytes).await
{ {
debug!( debug!(
@@ -539,8 +539,8 @@ impl Node {
// overlapping the immutable peer-table read. // overlapping the immutable peer-table read.
let reachable = self.peers.values().any(|peer| peer.may_reach(dest)); let reachable = self.peers.values().any(|peer| peer.may_reach(dest));
use crate::proto::discovery::InitiateDecision; use crate::proto::lookup::InitiateDecision;
match crate::proto::discovery::initiate_gate(&mut self.discovery, dest, now_ms, reachable) { match crate::proto::lookup::initiate_gate(&mut self.discovery, dest, now_ms, reachable) {
InitiateDecision::Deduplicated => { InitiateDecision::Deduplicated => {
self.metrics().discovery.req_deduplicated.inc(); self.metrics().discovery.req_deduplicated.inc();
debug!( debug!(
@@ -569,7 +569,7 @@ impl Node {
// If no tree peers had the target, fail immediately // If no tree peers had the target, fail immediately
if sent == 0 { if sent == 0 {
crate::proto::discovery::initiate_failed(&mut self.discovery, dest, now_ms); crate::proto::lookup::initiate_failed(&mut self.discovery, dest, now_ms);
debug!( debug!(
target_node = %self.peer_display_name(dest), target_node = %self.peer_display_name(dest),
"Discovery failed, no tree peers with bloom match" "Discovery failed, no tree peers with bloom match"
@@ -592,7 +592,7 @@ impl Node {
pub(in crate::node) async fn check_pending_lookups(&mut self, now_ms: u64) { pub(in crate::node) async fn check_pending_lookups(&mut self, now_ms: u64) {
let attempt_timeouts = self.config().node.discovery.attempt_timeouts_secs.clone(); let attempt_timeouts = self.config().node.discovery.attempt_timeouts_secs.clone();
let outcome = let outcome =
crate::proto::discovery::poll_pending(&mut self.discovery, now_ms, &attempt_timeouts); crate::proto::lookup::poll_pending(&mut self.discovery, now_ms, &attempt_timeouts);
for (target, attempt) in outcome.retries { for (target, attempt) in outcome.retries {
let ttl = self.config().node.discovery.ttl; let ttl = self.config().node.discovery.ttl;
+7 -7
View File
@@ -43,9 +43,9 @@ use crate::cache::CoordCache;
use crate::node::session::SessionEntry; use crate::node::session::SessionEntry;
use crate::peer::{ActivePeer, PeerConnection}; use crate::peer::{ActivePeer, PeerConnection};
use crate::proto::bloom::{BloomFilter, BloomState}; use crate::proto::bloom::{BloomFilter, BloomState};
use crate::proto::discovery::{Discovery, DiscoveryBackoff, DiscoveryForwardRateLimiter};
use crate::proto::fmp::Fmp; use crate::proto::fmp::Fmp;
use crate::proto::fsp::Fsp; use crate::proto::fsp::Fsp;
use crate::proto::lookup::{Lookup, LookupBackoff, LookupForwardRateLimiter};
use crate::proto::mmp::Mmp; use crate::proto::mmp::Mmp;
use crate::proto::routing::{self, Router, RoutingErrorRateLimiter}; use crate::proto::routing::{self, Router, RoutingErrorRateLimiter};
use crate::proto::stp::TreeState; use crate::proto::stp::TreeState;
@@ -337,7 +337,7 @@ pub struct Node {
// === Discovery === // === Discovery ===
/// Discovery-subsystem state: recent-request dedup cache, in-flight /// Discovery-subsystem state: recent-request dedup cache, in-flight
/// lookups, originator-side backoff, and transit-side forward limiter. /// lookups, originator-side backoff, and transit-side forward limiter.
discovery: Discovery, discovery: Lookup,
// === Counters === // === Counters ===
/// Next link ID to allocate. /// Next link ID to allocate.
@@ -678,9 +678,9 @@ impl Node {
coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval_ms( coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval_ms(
coords_response_interval_ms, coords_response_interval_ms,
), ),
discovery: Discovery::new( discovery: Lookup::new(
DiscoveryBackoff::with_params(backoff_base_secs, backoff_max_secs), LookupBackoff::with_params(backoff_base_secs, backoff_max_secs),
DiscoveryForwardRateLimiter::with_interval_ms(forward_min_interval_secs * 1000), LookupForwardRateLimiter::with_interval_ms(forward_min_interval_secs * 1000),
), ),
pending_connects: Vec::new(), pending_connects: Vec::new(),
retry_pending: HashMap::new(), retry_pending: HashMap::new(),
@@ -843,7 +843,7 @@ impl Node {
coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval_ms( coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval_ms(
coords_response_interval_ms, coords_response_interval_ms,
), ),
discovery: Discovery::new(DiscoveryBackoff::new(), DiscoveryForwardRateLimiter::new()), discovery: Lookup::new(LookupBackoff::new(), LookupForwardRateLimiter::new()),
pending_connects: Vec::new(), pending_connects: Vec::new(),
retry_pending: HashMap::new(), retry_pending: HashMap::new(),
nostr_discovery: None, nostr_discovery: None,
@@ -2541,7 +2541,7 @@ impl Node {
/// Iterate over pending discovery lookups for diagnostics. /// Iterate over pending discovery lookups for diagnostics.
pub fn pending_lookups_iter( pub fn pending_lookups_iter(
&self, &self,
) -> impl Iterator<Item = (&NodeAddr, &crate::proto::discovery::PendingLookup)> { ) -> impl Iterator<Item = (&NodeAddr, &crate::proto::lookup::PendingLookup)> {
self.discovery.pending_lookups.iter() self.discovery.pending_lookups.iter()
} }
+2 -2
View File
@@ -5,7 +5,7 @@
//! response routing. //! response routing.
use super::*; use super::*;
use crate::proto::discovery::{LookupRequest, LookupResponse, RecentRequest}; use crate::proto::lookup::{LookupRequest, LookupResponse, RecentRequest};
use crate::proto::stp::TreeCoordinate; use crate::proto::stp::TreeCoordinate;
use spanning_tree::{ use spanning_tree::{
cleanup_nodes, generate_random_edges, lock_large_network_test, process_available_packets, cleanup_nodes, generate_random_edges, lock_large_network_test, process_available_packets,
@@ -1040,7 +1040,7 @@ async fn test_open_discovery_sweep_queues_eligible_skips_filtered() {
async fn test_check_pending_lookups_default_sequence_unreachable() { async fn test_check_pending_lookups_default_sequence_unreachable() {
use crate::peer::ActivePeer; use crate::peer::ActivePeer;
use crate::proto::bloom::BloomFilter; use crate::proto::bloom::BloomFilter;
use crate::proto::discovery::PendingLookup; use crate::proto::lookup::PendingLookup;
use crate::transport::LinkId; use crate::transport::LinkId;
use std::sync::mpsc; use std::sync::mpsc;
@@ -1,17 +1,17 @@
//! Sans-IO discovery decision core. //! Sans-IO mesh lookup decision core.
//! //!
//! Pure, runtime-agnostic decision logic for the discovery protocol. The //! Pure, runtime-agnostic decision logic for the mesh lookup protocol. The
//! async I/O adapter in `node::handlers::discovery` decodes wire bytes, //! async I/O adapter in `node::handlers::discovery` decodes wire bytes,
//! calls into this core, and drives the returned actions (the actual //! calls into this core, and drives the returned actions (the actual
//! encrypted sends). No I/O, no clock, no metrics, no logging here. //! encrypted sends). No I/O, no clock, no metrics, no logging here.
use alloc::sync::Arc; use alloc::sync::Arc;
use super::state::{Discovery, PendingLookup, RecentRequest}; use super::state::{Lookup, PendingLookup, RecentRequest};
use super::wire::LookupRequest; use super::wire::LookupRequest;
use crate::NodeAddr; use crate::NodeAddr;
/// Read-only view of routing state the discovery core needs. /// Read-only view of routing state the lookup core needs.
/// ///
/// The core defines this interface; the async shell (`node`) implements it /// The core defines this interface; the async shell (`node`) implements it
/// over the live peer/tree tables. Keeping it a trait keeps `proto` free of /// over the live peer/tree tables. Keeping it a trait keeps `proto` free of
@@ -24,8 +24,8 @@ pub(crate) trait RoutingView {
} }
/// An I/O action the async shell performs on the core's behalf. /// An I/O action the async shell performs on the core's behalf.
pub(crate) enum DiscoveryAction { pub(crate) enum LookupAction {
/// Send an encoded discovery PDU to a peer as an encrypted link message. /// Send an encoded lookup PDU to a peer as an encrypted link message.
/// `bytes` is `Arc`-shared so a fan-out encodes once. /// `bytes` is `Arc`-shared so a fan-out encodes once.
SendLink { peer: NodeAddr, bytes: Arc<[u8]> }, SendLink { peer: NodeAddr, bytes: Arc<[u8]> },
/// Cache the verified destination coordinates + path MTU (coord_cache). /// Cache the verified destination coordinates + path MTU (coord_cache).
@@ -52,7 +52,7 @@ pub(crate) enum ForwardOutcome {
/// Forward: one SendLink per selected peer. `used_fallback` is true when /// Forward: one SendLink per selected peer. `used_fallback` is true when
/// the non-tree bloom-match fallback set was used (no tree peer matched). /// the non-tree bloom-match fallback set was used (no tree peer matched).
Forward { Forward {
actions: Vec<DiscoveryAction>, actions: Vec<LookupAction>,
used_fallback: bool, used_fallback: bool,
}, },
} }
@@ -88,7 +88,7 @@ pub(crate) fn plan_forward(request: &mut LookupRequest, rv: &impl RoutingView) -
let bytes: Arc<[u8]> = Arc::from(request.encode()); let bytes: Arc<[u8]> = Arc::from(request.encode());
let actions = targets let actions = targets
.into_iter() .into_iter()
.map(|peer| DiscoveryAction::SendLink { .map(|peer| LookupAction::SendLink {
peer, peer,
bytes: bytes.clone(), bytes: bytes.clone(),
}) })
@@ -111,10 +111,7 @@ pub(crate) fn plan_forward(request: &mut LookupRequest, rv: &impl RoutingView) -
/// the pre-sans-IO `initiate_lookup` to keep this extraction behavior-neutral; /// the pre-sans-IO `initiate_lookup` to keep this extraction behavior-neutral;
/// it is a known origination gap (ISSUE-2026-0059) whose fix adds the fallback /// it is a known origination gap (ISSUE-2026-0059) whose fix adds the fallback
/// branch as a separate, behavior-changing change. /// branch as a separate, behavior-changing change.
pub(crate) fn plan_initiate( pub(crate) fn plan_initiate(request: &LookupRequest, rv: &impl RoutingView) -> Vec<LookupAction> {
request: &LookupRequest,
rv: &impl RoutingView,
) -> Vec<DiscoveryAction> {
let targets: Vec<NodeAddr> = rv let targets: Vec<NodeAddr> = rv
.peers_reaching(&request.target) .peers_reaching(&request.target)
.into_iter() .into_iter()
@@ -126,14 +123,14 @@ pub(crate) fn plan_initiate(
let bytes: Arc<[u8]> = Arc::from(request.encode()); let bytes: Arc<[u8]> = Arc::from(request.encode());
targets targets
.into_iter() .into_iter()
.map(|peer| DiscoveryAction::SendLink { .map(|peer| LookupAction::SendLink {
peer, peer,
bytes: bytes.clone(), bytes: bytes.clone(),
}) })
.collect() .collect()
} }
/// Classification of an inbound LookupRequest, decided from Discovery state. /// Classification of an inbound LookupRequest, decided from Lookup state.
pub(crate) enum RequestOutcome { pub(crate) enum RequestOutcome {
/// request_id already in the dedup cache — drop. /// request_id already in the dedup cache — drop.
Duplicate, Duplicate,
@@ -152,9 +149,9 @@ pub(crate) enum RequestOutcome {
/// Classify an inbound LookupRequest against the recent-request dedup cache and /// Classify an inbound LookupRequest against the recent-request dedup cache and
/// the transit forward rate limiter. Purges expired dedup entries, records the /// the transit forward rate limiter. Purges expired dedup entries, records the
/// request for reverse-path forwarding on the non-drop paths, and decides the /// request for reverse-path forwarding on the non-drop paths, and decides the
/// route. Pure over Discovery state + node addr + injected clock; no I/O, no view. /// route. Pure over Lookup state + node addr + injected clock; no I/O, no view.
pub(crate) fn classify_request( pub(crate) fn classify_request(
disc: &mut Discovery, lookup: &mut Lookup,
request: &LookupRequest, request: &LookupRequest,
from: &NodeAddr, from: &NodeAddr,
my_addr: &NodeAddr, my_addr: &NodeAddr,
@@ -163,25 +160,30 @@ pub(crate) fn classify_request(
max_recent: usize, max_recent: usize,
) -> RequestOutcome { ) -> RequestOutcome {
// Purge expired dedup entries (was purge_expired_requests). // Purge expired dedup entries (was purge_expired_requests).
disc.recent_requests lookup
.recent_requests
.retain(|_, entry| !entry.is_expired(now_ms, recent_expiry_ms)); .retain(|_, entry| !entry.is_expired(now_ms, recent_expiry_ms));
if disc.recent_requests.contains_key(&request.request_id) { if lookup.recent_requests.contains_key(&request.request_id) {
return RequestOutcome::Duplicate; return RequestOutcome::Duplicate;
} }
if disc.recent_requests.len() >= max_recent { if lookup.recent_requests.len() >= max_recent {
return RequestOutcome::DedupCacheFull { return RequestOutcome::DedupCacheFull {
len: disc.recent_requests.len(), len: lookup.recent_requests.len(),
}; };
} }
disc.recent_requests lookup
.recent_requests
.insert(request.request_id, RecentRequest::new(*from, now_ms)); .insert(request.request_id, RecentRequest::new(*from, now_ms));
if request.target == *my_addr { if request.target == *my_addr {
return RequestOutcome::RespondAsTarget; return RequestOutcome::RespondAsTarget;
} }
if request.can_forward() { if request.can_forward() {
if disc.forward_limiter.should_forward(&request.target, now_ms) { if lookup
.forward_limiter
.should_forward(&request.target, now_ms)
{
RequestOutcome::Forward RequestOutcome::Forward
} else { } else {
RequestOutcome::ForwardRateLimited RequestOutcome::ForwardRateLimited
@@ -205,10 +207,10 @@ pub(crate) enum ResponseRoute {
/// Classify an inbound LookupResponse against the recent-request dedup cache. /// Classify an inbound LookupResponse against the recent-request dedup cache.
/// ///
/// Pure decision over `Discovery` state: sets `response_forwarded` when this is /// Pure decision over `Lookup` state: sets `response_forwarded` when this is
/// the first response we transit for the request. No I/O, no view, no metrics. /// the first response we transit for the request. No I/O, no view, no metrics.
pub(crate) fn classify_response(disc: &mut Discovery, request_id: u64) -> ResponseRoute { pub(crate) fn classify_response(lookup: &mut Lookup, request_id: u64) -> ResponseRoute {
match disc.recent_requests.get_mut(&request_id) { match lookup.recent_requests.get_mut(&request_id) {
Some(recent) => { Some(recent) => {
if recent.response_forwarded { if recent.response_forwarded {
ResponseRoute::AlreadyForwarded ResponseRoute::AlreadyForwarded
@@ -233,13 +235,13 @@ pub(crate) enum ResponseRouteDecision {
} }
/// Decide the first hop for a LookupResponse we originate as the target, from /// Decide the first hop for a LookupResponse we originate as the target, from
/// the recent-request reverse-path record. Pure over `Discovery` state. /// the recent-request reverse-path record. Pure over `Lookup` state.
/// ///
/// Only the reverse-path decision is pure. The `NeedsTreeRoute` fallback (greedy /// Only the reverse-path decision is pure. The `NeedsTreeRoute` fallback (greedy
/// tree routing toward the origin) is a `&mut Node` coord-cache operation with a /// tree routing toward the origin) is a `&mut Node` coord-cache operation with a
/// TTL-touch side effect, so it stays in the shell rather than moving here. /// TTL-touch side effect, so it stays in the shell rather than moving here.
pub(crate) fn plan_response_route(disc: &Discovery, request_id: u64) -> ResponseRouteDecision { pub(crate) fn plan_response_route(lookup: &Lookup, request_id: u64) -> ResponseRouteDecision {
match disc.recent_requests.get(&request_id) { match lookup.recent_requests.get(&request_id) {
Some(recent) => ResponseRouteDecision::ReversePath(recent.from_peer), Some(recent) => ResponseRouteDecision::ReversePath(recent.from_peer),
None => ResponseRouteDecision::NeedsTreeRoute, None => ResponseRouteDecision::NeedsTreeRoute,
} }
@@ -247,31 +249,31 @@ pub(crate) fn plan_response_route(disc: &Discovery, request_id: u64) -> Response
/// Apply the accept-side effects of a verified LookupResponse we originated. /// Apply the accept-side effects of a verified LookupResponse we originated.
/// ///
/// Mutates the Discovery success state (clears backoff, drops the pending /// Mutates the Lookup success state (clears backoff, drops the pending
/// lookup) and returns the cross-subsystem effects for the shell to drive. /// lookup) and returns the cross-subsystem effects for the shell to drive.
/// Verification is the shell's job — this runs only after the proof checked out. /// Verification is the shell's job — this runs only after the proof checked out.
pub(crate) fn on_response_accepted( pub(crate) fn on_response_accepted(
disc: &mut Discovery, lookup: &mut Lookup,
target: &NodeAddr, target: &NodeAddr,
coords: crate::TreeCoordinate, coords: crate::TreeCoordinate,
now_ms: u64, now_ms: u64,
path_mtu: u16, path_mtu: u16,
) -> Vec<DiscoveryAction> { ) -> Vec<LookupAction> {
disc.backoff.record_success(target); lookup.backoff.record_success(target);
disc.pending_lookups.remove(target); lookup.pending_lookups.remove(target);
vec![ vec![
DiscoveryAction::CacheCoords { LookupAction::CacheCoords {
target: *target, target: *target,
coords, coords,
now_ms, now_ms,
path_mtu, path_mtu,
}, },
DiscoveryAction::WritePathMtu { LookupAction::WritePathMtu {
target: *target, target: *target,
path_mtu, path_mtu,
}, },
DiscoveryAction::ResetWarmupIfEstablished { target: *target }, LookupAction::ResetWarmupIfEstablished { target: *target },
DiscoveryAction::RetryQueuedPackets { target: *target }, LookupAction::RetryQueuedPackets { target: *target },
] ]
} }
@@ -287,11 +289,11 @@ pub(crate) struct PollOutcome {
pub timeouts: Vec<(NodeAddr, u32)>, pub timeouts: Vec<(NodeAddr, u32)>,
} }
/// Advance the pending-lookup retry ladder. Pure over `Discovery` state + /// Advance the pending-lookup retry ladder. Pure over `Lookup` state +
/// injected clock: partitions due entries into retries (attempt bumped) and /// injected clock: partitions due entries into retries (attempt bumped) and
/// final timeouts (removed + backoff failure recorded). No I/O, no view. /// final timeouts (removed + backoff failure recorded). No I/O, no view.
pub(crate) fn poll_pending( pub(crate) fn poll_pending(
disc: &mut Discovery, lookup: &mut Lookup,
now_ms: u64, now_ms: u64,
attempt_timeouts_secs: &[u64], attempt_timeouts_secs: &[u64],
) -> PollOutcome { ) -> PollOutcome {
@@ -301,7 +303,7 @@ pub(crate) fn poll_pending(
let mut retry_targets: Vec<NodeAddr> = Vec::new(); let mut retry_targets: Vec<NodeAddr> = Vec::new();
let mut timeout_targets: Vec<NodeAddr> = Vec::new(); let mut timeout_targets: Vec<NodeAddr> = Vec::new();
for (&target, entry) in &disc.pending_lookups { for (&target, entry) in &lookup.pending_lookups {
let idx = (entry.attempt as usize).saturating_sub(1); let idx = (entry.attempt as usize).saturating_sub(1);
let to_ms = attempt_timeouts_secs.get(idx).copied().unwrap_or(0) * 1000; let to_ms = attempt_timeouts_secs.get(idx).copied().unwrap_or(0) * 1000;
if now_ms.saturating_sub(entry.last_sent_ms) >= to_ms { if now_ms.saturating_sub(entry.last_sent_ms) >= to_ms {
@@ -315,7 +317,7 @@ pub(crate) fn poll_pending(
let mut retries: Vec<(NodeAddr, u8)> = Vec::new(); let mut retries: Vec<(NodeAddr, u8)> = Vec::new();
for target in retry_targets { for target in retry_targets {
if let Some(entry) = disc.pending_lookups.get_mut(&target) { if let Some(entry) = lookup.pending_lookups.get_mut(&target) {
entry.attempt += 1; entry.attempt += 1;
entry.last_sent_ms = now_ms; entry.last_sent_ms = now_ms;
retries.push((target, entry.attempt)); retries.push((target, entry.attempt));
@@ -324,16 +326,16 @@ pub(crate) fn poll_pending(
let mut timeouts: Vec<(NodeAddr, u32)> = Vec::new(); let mut timeouts: Vec<(NodeAddr, u32)> = Vec::new();
for target in timeout_targets { for target in timeout_targets {
disc.pending_lookups.remove(&target); lookup.pending_lookups.remove(&target);
disc.backoff.record_failure(&target, now_ms); lookup.backoff.record_failure(&target, now_ms);
let failures = disc.backoff.failure_count(&target); let failures = lookup.backoff.failure_count(&target);
timeouts.push((target, failures)); timeouts.push((target, failures));
} }
PollOutcome { retries, timeouts } PollOutcome { retries, timeouts }
} }
/// Decision for whether/how to initiate a discovery lookup for a target. /// Decision for whether/how to initiate a lookup for a target.
pub(crate) enum InitiateDecision { pub(crate) enum InitiateDecision {
/// A lookup is already pending for this target — skip. /// A lookup is already pending for this target — skip.
Deduplicated, Deduplicated,
@@ -345,36 +347,37 @@ pub(crate) enum InitiateDecision {
Proceed, Proceed,
} }
/// Gate a discovery-lookup initiation against pending-dedup, backoff /// Gate a lookup initiation against pending-dedup, backoff
/// suppression, and bloom reachability (passed in — the shell reads the peer /// suppression, and bloom reachability (passed in — the shell reads the peer
/// filters). On BloomMiss records a failure; on Proceed inserts the pending /// filters). On BloomMiss records a failure; on Proceed inserts the pending
/// lookup. Pure over Discovery state + injected clock; no I/O, no view. /// lookup. Pure over Lookup state + injected clock; no I/O, no view.
pub(crate) fn initiate_gate( pub(crate) fn initiate_gate(
disc: &mut Discovery, lookup: &mut Lookup,
dest: &NodeAddr, dest: &NodeAddr,
now_ms: u64, now_ms: u64,
reachable: bool, reachable: bool,
) -> InitiateDecision { ) -> InitiateDecision {
if disc.pending_lookups.contains_key(dest) { if lookup.pending_lookups.contains_key(dest) {
return InitiateDecision::Deduplicated; return InitiateDecision::Deduplicated;
} }
if disc.backoff.is_suppressed(dest, now_ms) { if lookup.backoff.is_suppressed(dest, now_ms) {
return InitiateDecision::Suppressed { return InitiateDecision::Suppressed {
failures: disc.backoff.failure_count(dest), failures: lookup.backoff.failure_count(dest),
}; };
} }
if !reachable { if !reachable {
disc.backoff.record_failure(dest, now_ms); lookup.backoff.record_failure(dest, now_ms);
return InitiateDecision::BloomMiss; return InitiateDecision::BloomMiss;
} }
disc.pending_lookups lookup
.pending_lookups
.insert(*dest, PendingLookup::new(now_ms)); .insert(*dest, PendingLookup::new(now_ms));
InitiateDecision::Proceed InitiateDecision::Proceed
} }
/// Roll back a lookup whose first attempt reached no tree peers (sent == 0): /// Roll back a lookup whose first attempt reached no tree peers (sent == 0):
/// drop the pending entry and record a backoff failure. /// drop the pending entry and record a backoff failure.
pub(crate) fn initiate_failed(disc: &mut Discovery, dest: &NodeAddr, now_ms: u64) { pub(crate) fn initiate_failed(lookup: &mut Lookup, dest: &NodeAddr, now_ms: u64) {
disc.pending_lookups.remove(dest); lookup.pending_lookups.remove(dest);
disc.backoff.record_failure(dest, now_ms); lookup.backoff.record_failure(dest, now_ms);
} }
@@ -1,15 +1,15 @@
//! Discovery protocol rate limiting and backoff. //! Mesh lookup protocol rate limiting and backoff.
//! //!
//! Two complementary mechanisms: //! Two complementary mechanisms:
//! //!
//! - **`DiscoveryBackoff`** (originator-side, optional): Exponential //! - **`LookupBackoff`** (originator-side, optional): Exponential
//! suppression of fresh lookups after the per-attempt sequence in //! suppression of fresh lookups after the per-attempt sequence in
//! `node.discovery.attempt_timeouts_secs` has been exhausted. //! `node.discovery.attempt_timeouts_secs` has been exhausted.
//! **Disabled by default** (base/cap = 0); the per-attempt sequence //! **Disabled by default** (base/cap = 0); the per-attempt sequence
//! is the only retry pacing in the standard configuration. Reset on //! is the only retry pacing in the standard configuration. Reset on
//! topology changes (parent change, new peer, first RTT, reconnection). //! topology changes (parent change, new peer, first RTT, reconnection).
//! //!
//! - **`DiscoveryForwardRateLimiter`** (transit-side): Per-target minimum //! - **`LookupForwardRateLimiter`** (transit-side): Per-target minimum
//! interval for forwarded requests. Defense-in-depth against misbehaving //! interval for forwarded requests. Defense-in-depth against misbehaving
//! nodes generating fresh request_ids at high rate. //! nodes generating fresh request_ids at high rate.
@@ -23,10 +23,10 @@ use alloc::collections::BTreeMap;
/// Maximum number of recent LookupRequests retained for dedup and /// Maximum number of recent LookupRequests retained for dedup and
/// reverse-path routing before the cache is treated as full. /// reverse-path routing before the cache is treated as full.
pub(crate) const MAX_RECENT_DISCOVERY_REQUESTS: usize = 4096; pub(crate) const MAX_RECENT_LOOKUP_REQUESTS: usize = 4096;
// ============================================================================ // ============================================================================
// Originator-side: Discovery Backoff // Originator-side: Lookup Backoff
// ============================================================================ // ============================================================================
/// Default base backoff after first lookup failure. `0` = disabled. /// Default base backoff after first lookup failure. `0` = disabled.
@@ -35,11 +35,11 @@ const DEFAULT_BACKOFF_BASE_SECS: u64 = 0;
/// Default maximum backoff cap. `0` = disabled. /// Default maximum backoff cap. `0` = disabled.
const DEFAULT_BACKOFF_MAX_SECS: u64 = 0; const DEFAULT_BACKOFF_MAX_SECS: u64 = 0;
/// Exponential backoff for failed discovery lookups. /// Exponential backoff for failed lookups.
/// ///
/// Tracks targets whose lookups have timed out and suppresses /// Tracks targets whose lookups have timed out and suppresses
/// re-initiation with increasing delays. Cleared on topology changes. /// re-initiation with increasing delays. Cleared on topology changes.
pub struct DiscoveryBackoff { pub struct LookupBackoff {
/// Maps target → (suppress_until, consecutive_failures). /// Maps target → (suppress_until, consecutive_failures).
pub(crate) entries: BTreeMap<NodeAddr, BackoffEntry>, pub(crate) entries: BTreeMap<NodeAddr, BackoffEntry>,
/// Base backoff in milliseconds (first failure). /// Base backoff in milliseconds (first failure).
@@ -55,7 +55,7 @@ pub(crate) struct BackoffEntry {
failures: u32, failures: u32,
} }
impl DiscoveryBackoff { impl LookupBackoff {
/// Create with default parameters (disabled — base/cap = 0). /// Create with default parameters (disabled — base/cap = 0).
pub fn new() -> Self { pub fn new() -> Self {
Self::with_params(DEFAULT_BACKOFF_BASE_SECS, DEFAULT_BACKOFF_MAX_SECS) Self::with_params(DEFAULT_BACKOFF_BASE_SECS, DEFAULT_BACKOFF_MAX_SECS)
@@ -138,14 +138,14 @@ impl DiscoveryBackoff {
} }
} }
impl Default for DiscoveryBackoff { impl Default for LookupBackoff {
fn default() -> Self { fn default() -> Self {
Self::new() Self::new()
} }
} }
// ============================================================================ // ============================================================================
// Transit-side: Discovery Forward Rate Limiter // Transit-side: Lookup Forward Rate Limiter
// ============================================================================ // ============================================================================
/// Default minimum interval between forwarded lookups for the same target. /// Default minimum interval between forwarded lookups for the same target.
@@ -154,14 +154,14 @@ const DEFAULT_FORWARD_MIN_INTERVAL_MS: u64 = 2_000;
/// Maximum age of entries before cleanup. /// Maximum age of entries before cleanup.
const FORWARD_MAX_AGE_MS: u64 = 60_000; const FORWARD_MAX_AGE_MS: u64 = 60_000;
/// Rate limiter for forwarded discovery requests. /// Rate limiter for forwarded lookup requests.
/// ///
/// Tracks the last time a LookupRequest was forwarded for each target /// Tracks the last time a LookupRequest was forwarded for each target
/// and enforces a minimum interval to prevent floods from misbehaving /// and enforces a minimum interval to prevent floods from misbehaving
/// nodes generating fresh request_ids. /// nodes generating fresh request_ids.
pub struct DiscoveryForwardRateLimiter(PerAddrRateLimiter); pub struct LookupForwardRateLimiter(PerAddrRateLimiter);
impl DiscoveryForwardRateLimiter { impl LookupForwardRateLimiter {
/// Create with default parameters (2s interval). /// Create with default parameters (2s interval).
pub fn new() -> Self { pub fn new() -> Self {
Self(PerAddrRateLimiter::new( Self(PerAddrRateLimiter::new(
@@ -201,7 +201,7 @@ impl DiscoveryForwardRateLimiter {
} }
} }
impl Default for DiscoveryForwardRateLimiter { impl Default for LookupForwardRateLimiter {
fn default() -> Self { fn default() -> Self {
Self::new() Self::new()
} }
@@ -1,14 +1,14 @@
//! Sans-IO discovery protocol state. //! Sans-IO mesh lookup protocol state.
//! //!
//! Pure, runtime-agnostic discovery state and rate limiting, migrated out //! Pure, runtime-agnostic mesh lookup state and rate limiting, migrated out
//! of the async node shell. The async I/O handlers remain in //! of the async node shell. The async I/O handlers remain in
//! `node::handlers::discovery`. The discovery wire codec now lives here in //! `node::handlers::discovery`. The mesh lookup wire codec now lives here in
//! `wire.rs` (the `LookupRequest` / `LookupResponse` structs), per the //! `wire.rs` (the `LookupRequest` / `LookupResponse` structs), per the
//! wire-migrates-with-subsystem policy. //! wire-migrates-with-subsystem policy.
//! //!
//! The sans-IO decision core lives in `core.rs`: it defines the `RoutingView` //! The sans-IO decision core lives in `core.rs`: it defines the `RoutingView`
//! read-seam trait plus the `plan_forward` / `plan_initiate` LookupRequest //! read-seam trait plus the `plan_forward` / `plan_initiate` LookupRequest
//! planners and their `DiscoveryAction` / `ForwardOutcome` types. The async //! planners and their `LookupAction` / `ForwardOutcome` types. The async
//! shell decodes wire //! shell decodes wire
//! bytes, calls the planner, and drives the returned actions. //! bytes, calls the planner, and drives the returned actions.
@@ -21,15 +21,13 @@ mod wire;
mod tests; mod tests;
pub(crate) use core::{ pub(crate) use core::{
DiscoveryAction, ForwardOutcome, InitiateDecision, RequestOutcome, ResponseRoute, ForwardOutcome, InitiateDecision, LookupAction, RequestOutcome, ResponseRoute,
ResponseRouteDecision, RoutingView, classify_request, classify_response, initiate_failed, ResponseRouteDecision, RoutingView, classify_request, classify_response, initiate_failed,
initiate_gate, on_response_accepted, plan_forward, plan_initiate, plan_response_route, initiate_gate, on_response_accepted, plan_forward, plan_initiate, plan_response_route,
poll_pending, poll_pending,
}; };
pub(crate) use limits::{ pub(crate) use limits::{LookupBackoff, LookupForwardRateLimiter, MAX_RECENT_LOOKUP_REQUESTS};
DiscoveryBackoff, DiscoveryForwardRateLimiter, MAX_RECENT_DISCOVERY_REQUESTS,
};
#[cfg(test)] #[cfg(test)]
pub(crate) use state::RecentRequest; pub(crate) use state::RecentRequest;
pub(crate) use state::{Discovery, PendingLookup}; pub(crate) use state::{Lookup, PendingLookup};
pub use wire::{LookupRequest, LookupResponse}; pub use wire::{LookupRequest, LookupResponse};
@@ -1,14 +1,14 @@
//! Discovery-subsystem state owned by [`Node`](crate::node::Node). //! Mesh lookup subsystem state owned by [`Node`](crate::node::Node).
//! //!
//! Groups the four discovery-related state fields (recent-request dedup //! Groups the four lookup-related state fields (recent-request dedup
//! cache, in-flight lookups, originator-side backoff, transit-side forward //! cache, in-flight lookups, originator-side backoff, transit-side forward
//! rate limiter) behind a single struct so the discovery handlers can //! rate limiter) behind a single struct so the lookup handlers can
//! evolve toward a sans-IO core without threading four fields through //! evolve toward a sans-IO core without threading four fields through
//! `Node`. //! `Node`.
use alloc::collections::BTreeMap; use alloc::collections::BTreeMap;
use super::limits::{DiscoveryBackoff, DiscoveryForwardRateLimiter}; use super::limits::{LookupBackoff, LookupForwardRateLimiter};
use crate::NodeAddr; use crate::NodeAddr;
/// Recent request tracking for dedup and reverse-path forwarding. /// Recent request tracking for dedup and reverse-path forwarding.
@@ -44,7 +44,7 @@ impl RecentRequest {
} }
} }
/// Tracks a pending discovery lookup with retry state. /// Tracks a pending lookup with retry state.
pub struct PendingLookup { pub struct PendingLookup {
/// When the lookup was first initiated. /// When the lookup was first initiated.
pub initiated_ms: u64, pub initiated_ms: u64,
@@ -64,30 +64,27 @@ impl PendingLookup {
} }
} }
/// Discovery-subsystem state. /// Mesh lookup subsystem state.
pub(crate) struct Discovery { pub(crate) struct Lookup {
/// Recent discovery requests (dedup + reverse-path forwarding). /// Recent lookup requests (dedup + reverse-path forwarding).
/// Maps request_id → RecentRequest. /// Maps request_id → RecentRequest.
pub(crate) recent_requests: BTreeMap<u64, RecentRequest>, pub(crate) recent_requests: BTreeMap<u64, RecentRequest>,
/// Tracks in-flight discovery lookups. Maps target NodeAddr to the /// Tracks in-flight lookups. Maps target NodeAddr to the
/// initiation timestamp (Unix ms). Prevents duplicate flood queries. /// initiation timestamp (Unix ms). Prevents duplicate flood queries.
pub(crate) pending_lookups: BTreeMap<NodeAddr, PendingLookup>, pub(crate) pending_lookups: BTreeMap<NodeAddr, PendingLookup>,
/// Backoff for failed discovery lookups (originator-side). /// Backoff for failed lookups (originator-side).
pub(crate) backoff: DiscoveryBackoff, pub(crate) backoff: LookupBackoff,
/// Rate limiter for forwarded discovery requests (transit-side). /// Rate limiter for forwarded lookup requests (transit-side).
pub(crate) forward_limiter: DiscoveryForwardRateLimiter, pub(crate) forward_limiter: LookupForwardRateLimiter,
} }
impl Discovery { impl Lookup {
/// Create discovery state with the given backoff and forward limiter. /// Create mesh lookup state with the given backoff and forward limiter.
/// ///
/// The two limiters are constructed by the caller so each `Node` /// The two limiters are constructed by the caller so each `Node`
/// constructor can supply its own configured/default variant, matching /// constructor can supply its own configured/default variant, matching
/// the pre-refactor initialization exactly. /// the pre-refactor initialization exactly.
pub(crate) fn new( pub(crate) fn new(backoff: LookupBackoff, forward_limiter: LookupForwardRateLimiter) -> Self {
backoff: DiscoveryBackoff,
forward_limiter: DiscoveryForwardRateLimiter,
) -> Self {
Self { Self {
recent_requests: BTreeMap::new(), recent_requests: BTreeMap::new(),
pending_lookups: BTreeMap::new(), pending_lookups: BTreeMap::new(),
@@ -96,7 +93,7 @@ impl Discovery {
} }
} }
/// Reset discovery backoff on topology changes. Returns the number of /// Reset lookup backoff on topology changes. Returns the number of
/// entries cleared (0 if already empty) so the shell can log the reset — /// entries cleared (0 if already empty) so the shell can log the reset —
/// observability stays out of the pure core. /// observability stays out of the pure core.
pub(crate) fn reset_backoff(&mut self) -> usize { pub(crate) fn reset_backoff(&mut self) -> usize {
@@ -1,11 +1,10 @@
//! Tests for the sans-IO discovery decision core. //! Tests for the sans-IO lookup decision core.
use super::util::{ use super::util::{
MockRoutingView, action_peers, empty_discovery, make_request, make_request_id, MockRoutingView, action_peers, empty_lookup, make_request, make_request_id, suppressing_lookup,
suppressing_discovery,
}; };
use crate::TreeCoordinate; use crate::TreeCoordinate;
use crate::proto::discovery::*; use crate::proto::lookup::*;
use crate::testutil::make_node_addr; use crate::testutil::make_node_addr;
#[test] #[test]
@@ -125,10 +124,12 @@ fn initiate_returns_empty_when_nothing_reaches_target() {
#[test] #[test]
fn response_route_uses_recorded_reverse_path() { fn response_route_uses_recorded_reverse_path() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(9); let from = make_node_addr(9);
disc.recent_requests.insert(42, RecentRequest::new(from, 0)); lookup
match plan_response_route(&disc, 42) { .recent_requests
.insert(42, RecentRequest::new(from, 0));
match plan_response_route(&lookup, 42) {
ResponseRouteDecision::ReversePath(peer) => assert_eq!(peer, from), ResponseRouteDecision::ReversePath(peer) => assert_eq!(peer, from),
ResponseRouteDecision::NeedsTreeRoute => panic!("expected ReversePath"), ResponseRouteDecision::NeedsTreeRoute => panic!("expected ReversePath"),
} }
@@ -136,9 +137,9 @@ fn response_route_uses_recorded_reverse_path() {
#[test] #[test]
fn response_route_needs_tree_route_when_no_record() { fn response_route_needs_tree_route_when_no_record() {
let disc = empty_discovery(); let lookup = empty_lookup();
assert!(matches!( assert!(matches!(
plan_response_route(&disc, 42), plan_response_route(&lookup, 42),
ResponseRouteDecision::NeedsTreeRoute ResponseRouteDecision::NeedsTreeRoute
)); ));
} }
@@ -146,40 +147,42 @@ fn response_route_needs_tree_route_when_no_record() {
#[test] #[test]
fn classify_response_transit_on_fresh_forwarded_request() { fn classify_response_transit_on_fresh_forwarded_request() {
let from_peer = make_node_addr(0x11); let from_peer = make_node_addr(0x11);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
disc.recent_requests lookup
.recent_requests
.insert(42, RecentRequest::new(from_peer, 1000)); .insert(42, RecentRequest::new(from_peer, 1000));
match classify_response(&mut disc, 42) { match classify_response(&mut lookup, 42) {
ResponseRoute::Transit { from_peer: peer } => assert_eq!(peer, from_peer), ResponseRoute::Transit { from_peer: peer } => assert_eq!(peer, from_peer),
_ => panic!("expected Transit"), _ => panic!("expected Transit"),
} }
// The dedup flag must flip after the first transit. // The dedup flag must flip after the first transit.
assert!(disc.recent_requests.get(&42).unwrap().response_forwarded); assert!(lookup.recent_requests.get(&42).unwrap().response_forwarded);
} }
#[test] #[test]
fn classify_response_already_forwarded_on_second_call() { fn classify_response_already_forwarded_on_second_call() {
let from_peer = make_node_addr(0x22); let from_peer = make_node_addr(0x22);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
disc.recent_requests lookup
.recent_requests
.insert(7, RecentRequest::new(from_peer, 1000)); .insert(7, RecentRequest::new(from_peer, 1000));
assert!(matches!( assert!(matches!(
classify_response(&mut disc, 7), classify_response(&mut lookup, 7),
ResponseRoute::Transit { .. } ResponseRoute::Transit { .. }
)); ));
assert!(matches!( assert!(matches!(
classify_response(&mut disc, 7), classify_response(&mut lookup, 7),
ResponseRoute::AlreadyForwarded ResponseRoute::AlreadyForwarded
)); ));
} }
#[test] #[test]
fn classify_response_originator_when_request_absent() { fn classify_response_originator_when_request_absent() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
assert!(matches!( assert!(matches!(
classify_response(&mut disc, 999), classify_response(&mut lookup, 999),
ResponseRoute::Originator ResponseRoute::Originator
)); ));
} }
@@ -187,34 +190,35 @@ fn classify_response_originator_when_request_absent() {
#[test] #[test]
fn on_response_accepted_clears_state_and_emits_effects() { fn on_response_accepted_clears_state_and_emits_effects() {
let target = make_node_addr(0x5A); let target = make_node_addr(0x5A);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
// Seed a backoff entry and a pending lookup for the target. // Seed a backoff entry and a pending lookup for the target.
disc.backoff.record_failure(&target, 1000); lookup.backoff.record_failure(&target, 1000);
assert!(!disc.backoff.is_empty(), "precondition: backoff seeded"); assert!(!lookup.backoff.is_empty(), "precondition: backoff seeded");
disc.pending_lookups lookup
.pending_lookups
.insert(target, PendingLookup::new(1000)); .insert(target, PendingLookup::new(1000));
assert!(disc.pending_lookups.contains_key(&target)); assert!(lookup.pending_lookups.contains_key(&target));
let coords = TreeCoordinate::root(target); let coords = TreeCoordinate::root(target);
let now_ms = 12_345u64; let now_ms = 12_345u64;
let path_mtu = 1400u16; let path_mtu = 1400u16;
let actions = on_response_accepted(&mut disc, &target, coords, now_ms, path_mtu); let actions = on_response_accepted(&mut lookup, &target, coords, now_ms, path_mtu);
// Success state must be cleared. // Success state must be cleared.
assert!( assert!(
disc.backoff.is_empty(), lookup.backoff.is_empty(),
"backoff entry must clear on success" "backoff entry must clear on success"
); );
assert!( assert!(
!disc.pending_lookups.contains_key(&target), !lookup.pending_lookups.contains_key(&target),
"pending lookup must be dropped" "pending lookup must be dropped"
); );
// Exactly the four effect actions, in order. // Exactly the four effect actions, in order.
assert_eq!(actions.len(), 4, "expected four effect actions"); assert_eq!(actions.len(), 4, "expected four effect actions");
match &actions[0] { match &actions[0] {
DiscoveryAction::CacheCoords { LookupAction::CacheCoords {
target: t, target: t,
now_ms: n, now_ms: n,
path_mtu: p, path_mtu: p,
@@ -227,7 +231,7 @@ fn on_response_accepted_clears_state_and_emits_effects() {
_ => panic!("action[0] must be CacheCoords"), _ => panic!("action[0] must be CacheCoords"),
} }
match &actions[1] { match &actions[1] {
DiscoveryAction::WritePathMtu { LookupAction::WritePathMtu {
target: t, target: t,
path_mtu: p, path_mtu: p,
} => { } => {
@@ -237,11 +241,11 @@ fn on_response_accepted_clears_state_and_emits_effects() {
_ => panic!("action[1] must be WritePathMtu"), _ => panic!("action[1] must be WritePathMtu"),
} }
match &actions[2] { match &actions[2] {
DiscoveryAction::ResetWarmupIfEstablished { target: t } => assert_eq!(*t, target), LookupAction::ResetWarmupIfEstablished { target: t } => assert_eq!(*t, target),
_ => panic!("action[2] must be ResetWarmupIfEstablished"), _ => panic!("action[2] must be ResetWarmupIfEstablished"),
} }
match &actions[3] { match &actions[3] {
DiscoveryAction::RetryQueuedPackets { target: t } => assert_eq!(*t, target), LookupAction::RetryQueuedPackets { target: t } => assert_eq!(*t, target),
_ => panic!("action[3] must be RetryQueuedPackets"), _ => panic!("action[3] must be RetryQueuedPackets"),
} }
} }
@@ -249,17 +253,19 @@ fn on_response_accepted_clears_state_and_emits_effects() {
#[test] #[test]
fn poll_pending_no_action_before_first_deadline() { fn poll_pending_no_action_before_first_deadline() {
let target = make_node_addr(0x30); let target = make_node_addr(0x30);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let t0 = 10_000u64; let t0 = 10_000u64;
disc.pending_lookups.insert(target, PendingLookup::new(t0)); lookup
.pending_lookups
.insert(target, PendingLookup::new(t0));
// Just before the attempt-1 deadline (1s): nothing fires. // Just before the attempt-1 deadline (1s): nothing fires.
let outcome = poll_pending(&mut disc, t0 + 999, &[1, 2, 4, 8]); let outcome = poll_pending(&mut lookup, t0 + 999, &[1, 2, 4, 8]);
assert!(outcome.retries.is_empty(), "no retry before deadline"); assert!(outcome.retries.is_empty(), "no retry before deadline");
assert!(outcome.timeouts.is_empty(), "no timeout before deadline"); assert!(outcome.timeouts.is_empty(), "no timeout before deadline");
// Entry unchanged. // Entry unchanged.
let entry = disc.pending_lookups.get(&target).unwrap(); let entry = lookup.pending_lookups.get(&target).unwrap();
assert_eq!(entry.attempt, 1); assert_eq!(entry.attempt, 1);
assert_eq!(entry.last_sent_ms, t0); assert_eq!(entry.last_sent_ms, t0);
} }
@@ -267,17 +273,19 @@ fn poll_pending_no_action_before_first_deadline() {
#[test] #[test]
fn poll_pending_retries_at_first_deadline() { fn poll_pending_retries_at_first_deadline() {
let target = make_node_addr(0x31); let target = make_node_addr(0x31);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let t0 = 10_000u64; let t0 = 10_000u64;
disc.pending_lookups.insert(target, PendingLookup::new(t0)); lookup
.pending_lookups
.insert(target, PendingLookup::new(t0));
// At the attempt-1 deadline (t0 + 1000): one retry to attempt 2. // At the attempt-1 deadline (t0 + 1000): one retry to attempt 2.
let outcome = poll_pending(&mut disc, t0 + 1000, &[1, 2, 4, 8]); let outcome = poll_pending(&mut lookup, t0 + 1000, &[1, 2, 4, 8]);
assert_eq!(outcome.retries, vec![(target, 2)]); assert_eq!(outcome.retries, vec![(target, 2)]);
assert!(outcome.timeouts.is_empty()); assert!(outcome.timeouts.is_empty());
// Entry mutated: attempt bumped, last_sent refreshed. // Entry mutated: attempt bumped, last_sent refreshed.
let entry = disc.pending_lookups.get(&target).unwrap(); let entry = lookup.pending_lookups.get(&target).unwrap();
assert_eq!(entry.attempt, 2); assert_eq!(entry.attempt, 2);
assert_eq!(entry.last_sent_ms, t0 + 1000); assert_eq!(entry.last_sent_ms, t0 + 1000);
} }
@@ -285,17 +293,17 @@ fn poll_pending_retries_at_first_deadline() {
#[test] #[test]
fn poll_pending_final_timeout_at_max_attempt() { fn poll_pending_final_timeout_at_max_attempt() {
let target = make_node_addr(0x32); let target = make_node_addr(0x32);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
// Drive the entry to attempt == max (4) with a known last_sent. // Drive the entry to attempt == max (4) with a known last_sent.
let tn = 50_000u64; let tn = 50_000u64;
let mut entry = PendingLookup::new(tn); let mut entry = PendingLookup::new(tn);
entry.attempt = 4; entry.attempt = 4;
entry.last_sent_ms = tn; entry.last_sent_ms = tn;
disc.pending_lookups.insert(target, entry); lookup.pending_lookups.insert(target, entry);
// attempt_timeouts_secs[3] == 8 → deadline at tn + 8000. // attempt_timeouts_secs[3] == 8 → deadline at tn + 8000.
let outcome = poll_pending(&mut disc, tn + 8000, &[1, 2, 4, 8]); let outcome = poll_pending(&mut lookup, tn + 8000, &[1, 2, 4, 8]);
assert!(outcome.retries.is_empty(), "max attempt cannot retry"); assert!(outcome.retries.is_empty(), "max attempt cannot retry");
assert_eq!( assert_eq!(
outcome.timeouts, outcome.timeouts,
@@ -305,89 +313,98 @@ fn poll_pending_final_timeout_at_max_attempt() {
// Entry removed and a backoff failure recorded. // Entry removed and a backoff failure recorded.
assert!( assert!(
!disc.pending_lookups.contains_key(&target), !lookup.pending_lookups.contains_key(&target),
"timed-out entry must be removed" "timed-out entry must be removed"
); );
assert_eq!(disc.backoff.failure_count(&target), 1); assert_eq!(lookup.backoff.failure_count(&target), 1);
} }
// --- classify_request tests --- // --- classify_request tests ---
#[test] #[test]
fn classify_request_forwards_fresh_and_records_it() { fn classify_request_forwards_fresh_and_records_it() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(0x01); let from = make_node_addr(0x01);
let my_addr = make_node_addr(0x99); let my_addr = make_node_addr(0x99);
let target = make_node_addr(0xAA); let target = make_node_addr(0xAA);
let request = make_request_id(1, target, 3); let request = make_request_id(1, target, 3);
let outcome = classify_request(&mut disc, &request, &from, &my_addr, 1000, 5000, 4096); let outcome = classify_request(&mut lookup, &request, &from, &my_addr, 1000, 5000, 4096);
assert!(matches!(outcome, RequestOutcome::Forward)); assert!(matches!(outcome, RequestOutcome::Forward));
// Recorded for reverse-path forwarding. // Recorded for reverse-path forwarding.
assert!(disc.recent_requests.contains_key(&1)); assert!(lookup.recent_requests.contains_key(&1));
assert_eq!(disc.recent_requests.get(&1).unwrap().from_peer, from); assert_eq!(lookup.recent_requests.get(&1).unwrap().from_peer, from);
} }
#[test] #[test]
fn classify_request_duplicate_on_second_call() { fn classify_request_duplicate_on_second_call() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(0x01); let from = make_node_addr(0x01);
let my_addr = make_node_addr(0x99); let my_addr = make_node_addr(0x99);
let target = make_node_addr(0xAA); let target = make_node_addr(0xAA);
let request = make_request_id(1, target, 3); let request = make_request_id(1, target, 3);
assert!(matches!( assert!(matches!(
classify_request(&mut disc, &request, &from, &my_addr, 1000, 5000, 4096), classify_request(&mut lookup, &request, &from, &my_addr, 1000, 5000, 4096),
RequestOutcome::Forward RequestOutcome::Forward
)); ));
assert!(matches!( assert!(matches!(
classify_request(&mut disc, &request, &from, &my_addr, 1000, 5000, 4096), classify_request(&mut lookup, &request, &from, &my_addr, 1000, 5000, 4096),
RequestOutcome::Duplicate RequestOutcome::Duplicate
)); ));
} }
#[test] #[test]
fn classify_request_dedup_cache_full() { fn classify_request_dedup_cache_full() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(0x01); let from = make_node_addr(0x01);
let my_addr = make_node_addr(0x99); let my_addr = make_node_addr(0x99);
let target = make_node_addr(0xAA); let target = make_node_addr(0xAA);
// Fill the cache to max_recent with distinct request_ids. // Fill the cache to max_recent with distinct request_ids.
let max_recent = 3usize; let max_recent = 3usize;
for id in 100..(100 + max_recent as u64) { for id in 100..(100 + max_recent as u64) {
disc.recent_requests lookup
.recent_requests
.insert(id, RecentRequest::new(from, 1000)); .insert(id, RecentRequest::new(from, 1000));
} }
assert_eq!(disc.recent_requests.len(), max_recent); assert_eq!(lookup.recent_requests.len(), max_recent);
let request = make_request_id(1, target, 3); let request = make_request_id(1, target, 3);
match classify_request(&mut disc, &request, &from, &my_addr, 1000, 5000, max_recent) { match classify_request(
&mut lookup,
&request,
&from,
&my_addr,
1000,
5000,
max_recent,
) {
RequestOutcome::DedupCacheFull { len } => assert_eq!(len, max_recent), RequestOutcome::DedupCacheFull { len } => assert_eq!(len, max_recent),
_ => panic!("expected DedupCacheFull"), _ => panic!("expected DedupCacheFull"),
} }
// The new request must not have been recorded on the drop path. // The new request must not have been recorded on the drop path.
assert!(!disc.recent_requests.contains_key(&1)); assert!(!lookup.recent_requests.contains_key(&1));
} }
#[test] #[test]
fn classify_request_respond_as_target() { fn classify_request_respond_as_target() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(0x01); let from = make_node_addr(0x01);
let my_addr = make_node_addr(0xAA); let my_addr = make_node_addr(0xAA);
// target == my_addr // target == my_addr
let request = make_request_id(1, my_addr, 3); let request = make_request_id(1, my_addr, 3);
assert!(matches!( assert!(matches!(
classify_request(&mut disc, &request, &from, &my_addr, 1000, 5000, 4096), classify_request(&mut lookup, &request, &from, &my_addr, 1000, 5000, 4096),
RequestOutcome::RespondAsTarget RequestOutcome::RespondAsTarget
)); ));
// Recorded before the target decision. // Recorded before the target decision.
assert!(disc.recent_requests.contains_key(&1)); assert!(lookup.recent_requests.contains_key(&1));
} }
#[test] #[test]
fn classify_request_ttl_exhausted_for_non_target() { fn classify_request_ttl_exhausted_for_non_target() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(0x01); let from = make_node_addr(0x01);
let my_addr = make_node_addr(0x99); let my_addr = make_node_addr(0x99);
let target = make_node_addr(0xAA); let target = make_node_addr(0xAA);
@@ -395,71 +412,74 @@ fn classify_request_ttl_exhausted_for_non_target() {
let request = make_request_id(1, target, 0); let request = make_request_id(1, target, 0);
assert!(matches!( assert!(matches!(
classify_request(&mut disc, &request, &from, &my_addr, 1000, 5000, 4096), classify_request(&mut lookup, &request, &from, &my_addr, 1000, 5000, 4096),
RequestOutcome::TtlExhausted RequestOutcome::TtlExhausted
)); ));
} }
#[test] #[test]
fn classify_request_forward_rate_limited() { fn classify_request_forward_rate_limited() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(0x01); let from = make_node_addr(0x01);
let my_addr = make_node_addr(0x99); let my_addr = make_node_addr(0x99);
let target = make_node_addr(0xAA); let target = make_node_addr(0xAA);
// Pre-seed the forward limiter so should_forward(target) returns false // Pre-seed the forward limiter so should_forward(target) returns false
// on the next call within the (default 2s) min interval. // on the next call within the (default 2s) min interval.
assert!(disc.forward_limiter.should_forward(&target, 1000)); assert!(lookup.forward_limiter.should_forward(&target, 1000));
let request = make_request_id(1, target, 3); let request = make_request_id(1, target, 3);
assert!(matches!( assert!(matches!(
classify_request(&mut disc, &request, &from, &my_addr, 1000, 5000, 4096), classify_request(&mut lookup, &request, &from, &my_addr, 1000, 5000, 4096),
RequestOutcome::ForwardRateLimited RequestOutcome::ForwardRateLimited
)); ));
} }
#[test] #[test]
fn classify_request_purges_expired_entries() { fn classify_request_purges_expired_entries() {
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let from = make_node_addr(0x01); let from = make_node_addr(0x01);
let my_addr = make_node_addr(0x99); let my_addr = make_node_addr(0x99);
let target = make_node_addr(0xAA); let target = make_node_addr(0xAA);
// Seed an entry that is expired at now_ms with the given expiry window. // Seed an entry that is expired at now_ms with the given expiry window.
// is_expired: now - timestamp > expiry_ms → expired. // is_expired: now - timestamp > expiry_ms → expired.
disc.recent_requests lookup
.recent_requests
.insert(55, RecentRequest::new(from, 1000)); .insert(55, RecentRequest::new(from, 1000));
// now_ms = 10_000, expiry_ms = 5000 → 9000 > 5000 → expired. // now_ms = 10_000, expiry_ms = 5000 → 9000 > 5000 → expired.
let request = make_request_id(1, target, 3); let request = make_request_id(1, target, 3);
let outcome = classify_request(&mut disc, &request, &from, &my_addr, 10_000, 5000, 4096); let outcome = classify_request(&mut lookup, &request, &from, &my_addr, 10_000, 5000, 4096);
assert!(matches!(outcome, RequestOutcome::Forward)); assert!(matches!(outcome, RequestOutcome::Forward));
// The expired entry (55) must have been purged. // The expired entry (55) must have been purged.
assert!(!disc.recent_requests.contains_key(&55)); assert!(!lookup.recent_requests.contains_key(&55));
// The fresh request is recorded. // The fresh request is recorded.
assert!(disc.recent_requests.contains_key(&1)); assert!(lookup.recent_requests.contains_key(&1));
} }
#[test] #[test]
fn poll_pending_full_ladder_end_to_end() { fn poll_pending_full_ladder_end_to_end() {
let target = make_node_addr(0x33); let target = make_node_addr(0x33);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let t0 = 0u64; let t0 = 0u64;
disc.pending_lookups.insert(target, PendingLookup::new(t0)); lookup
.pending_lookups
.insert(target, PendingLookup::new(t0));
let ladder = [1u64, 2, 4, 8]; let ladder = [1u64, 2, 4, 8];
// attempt 1 → 2 at deadline 1s // attempt 1 → 2 at deadline 1s
let o = poll_pending(&mut disc, t0 + 1000, &ladder); let o = poll_pending(&mut lookup, t0 + 1000, &ladder);
assert_eq!(o.retries, vec![(target, 2)]); assert_eq!(o.retries, vec![(target, 2)]);
// attempt 2 → 3 at deadline 2s after last send // attempt 2 → 3 at deadline 2s after last send
let o = poll_pending(&mut disc, t0 + 1000 + 2000, &ladder); let o = poll_pending(&mut lookup, t0 + 1000 + 2000, &ladder);
assert_eq!(o.retries, vec![(target, 3)]); assert_eq!(o.retries, vec![(target, 3)]);
// attempt 3 → 4 at deadline 4s after last send // attempt 3 → 4 at deadline 4s after last send
let o = poll_pending(&mut disc, t0 + 1000 + 2000 + 4000, &ladder); let o = poll_pending(&mut lookup, t0 + 1000 + 2000 + 4000, &ladder);
assert_eq!(o.retries, vec![(target, 4)]); assert_eq!(o.retries, vec![(target, 4)]);
// attempt 4 is max → final timeout at deadline 8s after last send // attempt 4 is max → final timeout at deadline 8s after last send
let last = t0 + 1000 + 2000 + 4000; let last = t0 + 1000 + 2000 + 4000;
let o = poll_pending(&mut disc, last + 8000, &ladder); let o = poll_pending(&mut lookup, last + 8000, &ladder);
assert!(o.retries.is_empty()); assert!(o.retries.is_empty());
assert_eq!(o.timeouts, vec![(target, 1)]); assert_eq!(o.timeouts, vec![(target, 1)]);
assert!(!disc.pending_lookups.contains_key(&target)); assert!(!lookup.pending_lookups.contains_key(&target));
} }
// --- initiate_gate / initiate_failed tests --- // --- initiate_gate / initiate_failed tests ---
@@ -467,12 +487,12 @@ fn poll_pending_full_ladder_end_to_end() {
#[test] #[test]
fn initiate_gate_deduplicated_when_pending() { fn initiate_gate_deduplicated_when_pending() {
let dest = make_node_addr(0x40); let dest = make_node_addr(0x40);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
disc.pending_lookups.insert(dest, PendingLookup::new(500)); lookup.pending_lookups.insert(dest, PendingLookup::new(500));
// reachable=true would otherwise Proceed, but the pending entry wins. // reachable=true would otherwise Proceed, but the pending entry wins.
assert!(matches!( assert!(matches!(
initiate_gate(&mut disc, &dest, 1000, true), initiate_gate(&mut lookup, &dest, 1000, true),
InitiateDecision::Deduplicated InitiateDecision::Deduplicated
)); ));
} }
@@ -480,48 +500,48 @@ fn initiate_gate_deduplicated_when_pending() {
#[test] #[test]
fn initiate_gate_suppressed_by_backoff() { fn initiate_gate_suppressed_by_backoff() {
let dest = make_node_addr(0x41); let dest = make_node_addr(0x41);
let mut disc = suppressing_discovery(); let mut lookup = suppressing_lookup();
// One failure arms suppression under with_params(30, 300). // One failure arms suppression under with_params(30, 300).
disc.backoff.record_failure(&dest, 1000); lookup.backoff.record_failure(&dest, 1000);
assert!( assert!(
disc.backoff.is_suppressed(&dest, 1000), lookup.backoff.is_suppressed(&dest, 1000),
"precondition: suppressed" "precondition: suppressed"
); );
match initiate_gate(&mut disc, &dest, 1000, true) { match initiate_gate(&mut lookup, &dest, 1000, true) {
InitiateDecision::Suppressed { failures } => assert_eq!(failures, 1), InitiateDecision::Suppressed { failures } => assert_eq!(failures, 1),
_ => panic!("expected Suppressed"), _ => panic!("expected Suppressed"),
} }
// No pending entry was inserted on the suppress path. // No pending entry was inserted on the suppress path.
assert!(!disc.pending_lookups.contains_key(&dest)); assert!(!lookup.pending_lookups.contains_key(&dest));
} }
#[test] #[test]
fn initiate_gate_bloom_miss_records_failure() { fn initiate_gate_bloom_miss_records_failure() {
let dest = make_node_addr(0x42); let dest = make_node_addr(0x42);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
assert!(matches!( assert!(matches!(
initiate_gate(&mut disc, &dest, 1000, false), initiate_gate(&mut lookup, &dest, 1000, false),
InitiateDecision::BloomMiss InitiateDecision::BloomMiss
)); ));
// A backoff failure was recorded, and no pending entry created. // A backoff failure was recorded, and no pending entry created.
assert_eq!(disc.backoff.failure_count(&dest), 1); assert_eq!(lookup.backoff.failure_count(&dest), 1);
assert!(!disc.pending_lookups.contains_key(&dest)); assert!(!lookup.pending_lookups.contains_key(&dest));
} }
#[test] #[test]
fn initiate_gate_proceed_inserts_pending() { fn initiate_gate_proceed_inserts_pending() {
let dest = make_node_addr(0x43); let dest = make_node_addr(0x43);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
let now_ms = 7_777u64; let now_ms = 7_777u64;
assert!(matches!( assert!(matches!(
initiate_gate(&mut disc, &dest, now_ms, true), initiate_gate(&mut lookup, &dest, now_ms, true),
InitiateDecision::Proceed InitiateDecision::Proceed
)); ));
// The pending entry now exists, stamped with now_ms. // The pending entry now exists, stamped with now_ms.
let entry = disc let entry = lookup
.pending_lookups .pending_lookups
.get(&dest) .get(&dest)
.expect("Proceed must insert a pending lookup"); .expect("Proceed must insert a pending lookup");
@@ -532,13 +552,15 @@ fn initiate_gate_proceed_inserts_pending() {
#[test] #[test]
fn initiate_failed_drops_pending_and_records_failure() { fn initiate_failed_drops_pending_and_records_failure() {
let dest = make_node_addr(0x44); let dest = make_node_addr(0x44);
let mut disc = empty_discovery(); let mut lookup = empty_lookup();
disc.pending_lookups.insert(dest, PendingLookup::new(1000)); lookup
.pending_lookups
.insert(dest, PendingLookup::new(1000));
initiate_failed(&mut disc, &dest, 1000); initiate_failed(&mut lookup, &dest, 1000);
assert!( assert!(
!disc.pending_lookups.contains_key(&dest), !lookup.pending_lookups.contains_key(&dest),
"pending entry must be dropped" "pending entry must be dropped"
); );
assert_eq!(disc.backoff.failure_count(&dest), 1); assert_eq!(lookup.backoff.failure_count(&dest), 1);
} }
@@ -1,13 +1,13 @@
//! Tests for discovery rate limiting and backoff. //! Tests for lookup rate limiting and backoff.
use crate::proto::discovery::{DiscoveryBackoff, DiscoveryForwardRateLimiter}; use crate::proto::lookup::{LookupBackoff, LookupForwardRateLimiter};
use crate::testutil::make_node_addr as addr; use crate::testutil::make_node_addr as addr;
// --- DiscoveryBackoff tests --- // --- LookupBackoff tests ---
#[test] #[test]
fn test_backoff_not_suppressed_initially() { fn test_backoff_not_suppressed_initially() {
let backoff = DiscoveryBackoff::new(); let backoff = LookupBackoff::new();
assert!(!backoff.is_suppressed(&addr(1), 0)); assert!(!backoff.is_suppressed(&addr(1), 0));
} }
@@ -15,7 +15,7 @@ fn test_backoff_not_suppressed_initially() {
fn test_backoff_suppressed_after_failure() { fn test_backoff_suppressed_after_failure() {
// Backoff is opt-in; exercise the suppression path with explicit params. // Backoff is opt-in; exercise the suppression path with explicit params.
let now = 1_000; let now = 1_000;
let mut backoff = DiscoveryBackoff::with_params(30, 300); let mut backoff = LookupBackoff::with_params(30, 300);
backoff.record_failure(&addr(1), now); backoff.record_failure(&addr(1), now);
assert!(backoff.is_suppressed(&addr(1), now)); assert!(backoff.is_suppressed(&addr(1), now));
// Different target not affected // Different target not affected
@@ -25,7 +25,7 @@ fn test_backoff_suppressed_after_failure() {
#[test] #[test]
fn test_backoff_cleared_on_success() { fn test_backoff_cleared_on_success() {
let now = 1_000; let now = 1_000;
let mut backoff = DiscoveryBackoff::with_params(30, 300); let mut backoff = LookupBackoff::with_params(30, 300);
backoff.record_failure(&addr(1), now); backoff.record_failure(&addr(1), now);
assert!(backoff.is_suppressed(&addr(1), now)); assert!(backoff.is_suppressed(&addr(1), now));
@@ -36,7 +36,7 @@ fn test_backoff_cleared_on_success() {
#[test] #[test]
fn test_backoff_reset_all() { fn test_backoff_reset_all() {
let now = 1_000; let now = 1_000;
let mut backoff = DiscoveryBackoff::new(); let mut backoff = LookupBackoff::new();
backoff.record_failure(&addr(1), now); backoff.record_failure(&addr(1), now);
backoff.record_failure(&addr(2), now); backoff.record_failure(&addr(2), now);
assert_eq!(backoff.len(), 2); assert_eq!(backoff.len(), 2);
@@ -49,7 +49,7 @@ fn test_backoff_reset_all() {
#[test] #[test]
fn test_backoff_exponential() { fn test_backoff_exponential() {
let now = 1_000; let now = 1_000;
let mut backoff = DiscoveryBackoff::with_params(1, 300); let mut backoff = LookupBackoff::with_params(1, 300);
// First failure: 1s backoff // First failure: 1s backoff
backoff.record_failure(&addr(1), now); backoff.record_failure(&addr(1), now);
@@ -67,7 +67,7 @@ fn test_backoff_exponential() {
#[test] #[test]
fn test_backoff_expires() { fn test_backoff_expires() {
let now = 1_000; let now = 1_000;
let mut backoff = DiscoveryBackoff::with_params(0, 0); let mut backoff = LookupBackoff::with_params(0, 0);
backoff.record_failure(&addr(1), now); backoff.record_failure(&addr(1), now);
// With 0s backoff, should not be suppressed // With 0s backoff, should not be suppressed
assert!(!backoff.is_suppressed(&addr(1), now)); assert!(!backoff.is_suppressed(&addr(1), now));
@@ -76,7 +76,7 @@ fn test_backoff_expires() {
#[test] #[test]
fn test_backoff_capped() { fn test_backoff_capped() {
let now = 1_000; let now = 1_000;
let mut backoff = DiscoveryBackoff::with_params(1, 10); let mut backoff = LookupBackoff::with_params(1, 10);
// Record many failures // Record many failures
for _ in 0..20 { for _ in 0..20 {
@@ -89,18 +89,18 @@ fn test_backoff_capped() {
assert!(remaining <= 11_000); assert!(remaining <= 11_000);
} }
// --- DiscoveryForwardRateLimiter tests --- // --- LookupForwardRateLimiter tests ---
#[test] #[test]
fn test_forward_first_allowed() { fn test_forward_first_allowed() {
let mut limiter = DiscoveryForwardRateLimiter::new(); let mut limiter = LookupForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1), 0)); assert!(limiter.should_forward(&addr(1), 0));
} }
#[test] #[test]
fn test_forward_rapid_rate_limited() { fn test_forward_rapid_rate_limited() {
let now = 1_000; let now = 1_000;
let mut limiter = DiscoveryForwardRateLimiter::new(); let mut limiter = LookupForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1), now)); assert!(limiter.should_forward(&addr(1), now));
assert!(!limiter.should_forward(&addr(1), now)); assert!(!limiter.should_forward(&addr(1), now));
assert!(!limiter.should_forward(&addr(1), now)); assert!(!limiter.should_forward(&addr(1), now));
@@ -109,7 +109,7 @@ fn test_forward_rapid_rate_limited() {
#[test] #[test]
fn test_forward_different_targets_independent() { fn test_forward_different_targets_independent() {
let now = 1_000; let now = 1_000;
let mut limiter = DiscoveryForwardRateLimiter::new(); let mut limiter = LookupForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1), now)); assert!(limiter.should_forward(&addr(1), now));
assert!(limiter.should_forward(&addr(2), now)); assert!(limiter.should_forward(&addr(2), now));
assert!(!limiter.should_forward(&addr(1), now)); assert!(!limiter.should_forward(&addr(1), now));
@@ -119,7 +119,7 @@ fn test_forward_different_targets_independent() {
#[test] #[test]
fn test_forward_allowed_after_interval() { fn test_forward_allowed_after_interval() {
let now = 1_000; let now = 1_000;
let mut limiter = DiscoveryForwardRateLimiter::with_interval_ms(100); let mut limiter = LookupForwardRateLimiter::with_interval_ms(100);
assert!(limiter.should_forward(&addr(1), now)); assert!(limiter.should_forward(&addr(1), now));
// Advance past the minimum interval. // Advance past the minimum interval.
@@ -129,7 +129,7 @@ fn test_forward_allowed_after_interval() {
#[test] #[test]
fn test_forward_cleanup_removes_old() { fn test_forward_cleanup_removes_old() {
let now = 1_000; let now = 1_000;
let mut limiter = DiscoveryForwardRateLimiter::new(); let mut limiter = LookupForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1), now)); assert!(limiter.should_forward(&addr(1), now));
assert!(limiter.should_forward(&addr(2), now)); assert!(limiter.should_forward(&addr(2), now));
assert_eq!(limiter.len(), 2); assert_eq!(limiter.len(), 2);
@@ -142,7 +142,7 @@ fn test_forward_cleanup_removes_old() {
#[test] #[test]
fn test_forward_cleanup_preserves_recent() { fn test_forward_cleanup_preserves_recent() {
let now = 1_000; let now = 1_000;
let mut limiter = DiscoveryForwardRateLimiter::new(); let mut limiter = LookupForwardRateLimiter::new();
assert!(limiter.should_forward(&addr(1), now)); assert!(limiter.should_forward(&addr(1), now));
assert_eq!(limiter.len(), 1); assert_eq!(limiter.len(), 1);
@@ -1,4 +1,4 @@
//! Discovery subsystem unit tests, extracted from the co-located `#[cfg(test)]` //! Mesh lookup subsystem unit tests, extracted from the co-located `#[cfg(test)]`
//! blocks in the sibling source modules. Shared helpers live in `util`. //! blocks in the sibling source modules. Shared helpers live in `util`.
mod core; mod core;
@@ -1,10 +1,10 @@
//! Shared test helpers for the discovery subsystem unit tests. //! Shared test helpers for the lookup subsystem unit tests.
use sha2::Digest; use sha2::Digest;
use crate::proto::discovery::{ use crate::proto::lookup::{
Discovery, DiscoveryAction, DiscoveryBackoff, DiscoveryForwardRateLimiter, LookupRequest, Lookup, LookupAction, LookupBackoff, LookupForwardRateLimiter, LookupRequest, LookupResponse,
LookupResponse, RoutingView, RoutingView,
}; };
use crate::testutil::make_node_addr; use crate::testutil::make_node_addr;
use crate::{NodeAddr, TreeCoordinate}; use crate::{NodeAddr, TreeCoordinate};
@@ -49,29 +49,29 @@ pub(super) fn make_coords(ids: &[u8]) -> TreeCoordinate {
TreeCoordinate::from_addrs(ids.iter().map(|&v| make_node_addr(v)).collect()).unwrap() TreeCoordinate::from_addrs(ids.iter().map(|&v| make_node_addr(v)).collect()).unwrap()
} }
pub(super) fn action_peers(actions: &[DiscoveryAction]) -> Vec<NodeAddr> { pub(super) fn action_peers(actions: &[LookupAction]) -> Vec<NodeAddr> {
actions actions
.iter() .iter()
.map(|action| match action { .map(|action| match action {
DiscoveryAction::SendLink { peer, .. } => *peer, LookupAction::SendLink { peer, .. } => *peer,
_ => panic!("expected SendLink, got a different action variant"), _ => panic!("expected SendLink, got a different action variant"),
}) })
.collect() .collect()
} }
pub(super) fn empty_discovery() -> Discovery { pub(super) fn empty_lookup() -> Lookup {
Discovery::new( Lookup::new(
DiscoveryBackoff::default(), LookupBackoff::default(),
DiscoveryForwardRateLimiter::default(), LookupForwardRateLimiter::default(),
) )
} }
/// A Discovery whose backoff is armed (non-zero base/cap) so that a single /// A Lookup whose backoff is armed (non-zero base/cap) so that a single
/// recorded failure suppresses the target — the default backoff is inert. /// recorded failure suppresses the target — the default backoff is inert.
pub(super) fn suppressing_discovery() -> Discovery { pub(super) fn suppressing_lookup() -> Lookup {
Discovery::new( Lookup::new(
DiscoveryBackoff::with_params(30, 300), LookupBackoff::with_params(30, 300),
DiscoveryForwardRateLimiter::default(), LookupForwardRateLimiter::default(),
) )
} }
@@ -1,7 +1,7 @@
//! Tests for the discovery wire codec (`LookupRequest` / `LookupResponse`). //! Tests for the lookup wire codec (`LookupRequest` / `LookupResponse`).
use super::util::{make_coords, signed_response}; use super::util::{make_coords, signed_response};
use crate::proto::discovery::{LookupRequest, LookupResponse}; use crate::proto::lookup::{LookupRequest, LookupResponse};
use crate::testutil::make_node_addr; use crate::testutil::make_node_addr;
#[test] #[test]
@@ -1,4 +1,4 @@
//! Discovery messages: LookupRequest and LookupResponse. //! Mesh lookup messages: LookupRequest and LookupResponse.
use crate::NodeAddr; use crate::NodeAddr;
use crate::proto::Error; use crate::proto::Error;
@@ -7,7 +7,7 @@ use crate::proto::stp::TreeCoordinate;
use crate::proto::stp::{decode_coords, encode_coords}; use crate::proto::stp::{decode_coords, encode_coords};
use secp256k1::schnorr::Signature; use secp256k1::schnorr::Signature;
/// Request to discover a node's coordinates. /// Request to look up a node's coordinates.
/// ///
/// Routed through the spanning tree via bloom-filter-guided forwarding. /// Routed through the spanning tree via bloom-filter-guided forwarding.
/// Each transit node forwards only to tree peers whose bloom filter /// Each transit node forwards only to tree peers whose bloom filter
+1 -1
View File
@@ -9,10 +9,10 @@ pub use error::Error;
pub(crate) mod bloom; pub(crate) mod bloom;
pub(crate) mod codec; pub(crate) mod codec;
pub(crate) mod coord; pub(crate) mod coord;
pub(crate) mod discovery;
pub(crate) mod fmp; pub(crate) mod fmp;
pub(crate) mod fsp; pub(crate) mod fsp;
pub(crate) mod link; pub(crate) mod link;
pub(crate) mod lookup;
pub(crate) mod math; pub(crate) mod math;
pub(crate) mod mmp; pub(crate) mod mmp;
pub(crate) mod rate_limit; pub(crate) mod rate_limit;