diff --git a/src/lib.rs b/src/lib.rs index 5d38a36..8a62bb4 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -73,8 +73,8 @@ pub use proto::stp::TreeAnnounce; // Re-export bloom wire types (relocated from protocol:: to proto::bloom) pub use proto::bloom::FilterAnnounce; -// Re-export discovery wire types (relocated from protocol:: to proto::discovery) -pub use proto::discovery::{LookupRequest, LookupResponse}; +// Re-export discovery wire types (relocated from protocol:: to proto::lookup) +pub use proto::lookup::{LookupRequest, LookupResponse}; // Re-export routing wire types (relocated from protocol:: to proto::routing) pub use proto::routing::{ diff --git a/src/node/handlers/discovery.rs b/src/node/handlers/discovery.rs index 3df02a6..8d3c152 100644 --- a/src/node/handlers/discovery.rs +++ b/src/node/handlers/discovery.rs @@ -7,8 +7,8 @@ use crate::node::Node; use crate::node::reject::DiscoveryReject; -use crate::proto::discovery::{ - DiscoveryAction, LookupRequest, LookupResponse, MAX_RECENT_DISCOVERY_REQUESTS, +use crate::proto::lookup::{ + LookupAction, LookupRequest, LookupResponse, MAX_RECENT_LOOKUP_REQUESTS, }; use crate::transport::{TransportAddr, TransportId}; use crate::{NodeAddr, PeerIdentity}; @@ -26,7 +26,7 @@ struct NodeRoutingView<'a> { 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 { self.node.is_tree_peer(addr) } @@ -67,15 +67,15 @@ impl Node { let now_ms = Self::now_ms(); 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( + use crate::proto::lookup::RequestOutcome; + match crate::proto::lookup::classify_request( &mut self.discovery, &request, from, &my_addr, now_ms, recent_expiry_ms, - MAX_RECENT_DISCOVERY_REQUESTS, + MAX_RECENT_LOOKUP_REQUESTS, ) { RequestOutcome::Duplicate => { self.metrics() @@ -95,7 +95,7 @@ impl Node { request_id = request.request_id, from = %self.peer_display_name(from), recent_requests = len, - max_recent_requests = MAX_RECENT_DISCOVERY_REQUESTS, + max_recent_requests = MAX_RECENT_LOOKUP_REQUESTS, "Discovery request dedup cache full, dropping LookupRequest" ); } @@ -161,8 +161,8 @@ impl Node { let now_ms = Self::now_ms(); // Check if we forwarded this request (transit node) or originated it - match crate::proto::discovery::classify_response(&mut self.discovery, response.request_id) { - crate::proto::discovery::ResponseRoute::AlreadyForwarded => { + match crate::proto::lookup::classify_response(&mut self.discovery, response.request_id) { + crate::proto::lookup::ResponseRoute::AlreadyForwarded => { // Already forwarded a response for this request — drop to // prevent response routing loops. debug!( @@ -171,7 +171,7 @@ impl Node { "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 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 let target = response.target; let path_mtu = response.path_mtu; @@ -251,7 +251,7 @@ impl Node { // 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( + let actions = crate::proto::lookup::on_response_accepted( &mut self.discovery, &target, response.target_coords, @@ -266,10 +266,10 @@ impl Node { /// 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) { + async fn drive_response_actions(&mut self, actions: Vec) { for action in actions { match action { - DiscoveryAction::CacheCoords { + LookupAction::CacheCoords { target, coords, now_ms, @@ -278,7 +278,7 @@ impl Node { self.coord_cache .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 // map used by the TUN reader/writer at TCP MSS clamp time. 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. let n = self.config().node.session.coords_warmup_packets; 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 // initiation. The coord_cache now has coords, so find_next_hop() // should succeed. @@ -332,7 +332,7 @@ impl Node { 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 { debug!( peer = %self.peer_display_name(&peer), @@ -360,8 +360,8 @@ impl Node { // 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( + use crate::proto::lookup::ResponseRouteDecision; + let next_hop_addr = match crate::proto::lookup::plan_response_route( &self.discovery, request.request_id, ) { @@ -423,18 +423,18 @@ impl Node { // 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) + crate::proto::lookup::plan_forward(&mut request, &rv) }; match outcome { - crate::proto::discovery::ForwardOutcome::TtlExhausted => {} - crate::proto::discovery::ForwardOutcome::NoPeers => { + crate::proto::lookup::ForwardOutcome::TtlExhausted => {} + crate::proto::lookup::ForwardOutcome::NoPeers => { self.metrics().discovery.req_no_tree_peer.inc(); trace!( request_id = request.request_id, "No eligible peers to forward LookupRequest" ); } - crate::proto::discovery::ForwardOutcome::Forward { + crate::proto::lookup::ForwardOutcome::Forward { actions, used_fallback, } => { @@ -458,7 +458,7 @@ impl Node { ); } 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 { debug!( @@ -494,7 +494,7 @@ impl Node { // the shell drives the sends and keeps all metrics/logging. let actions = { let rv = NodeRoutingView { node: self }; - crate::proto::discovery::plan_initiate(&request, &rv) + crate::proto::lookup::plan_initiate(&request, &rv) }; let peer_count = actions.len(); @@ -509,7 +509,7 @@ impl Node { ); 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 { debug!( @@ -539,8 +539,8 @@ impl Node { // overlapping the immutable peer-table read. let reachable = self.peers.values().any(|peer| peer.may_reach(dest)); - use crate::proto::discovery::InitiateDecision; - match crate::proto::discovery::initiate_gate(&mut self.discovery, dest, now_ms, reachable) { + use crate::proto::lookup::InitiateDecision; + match crate::proto::lookup::initiate_gate(&mut self.discovery, dest, now_ms, reachable) { InitiateDecision::Deduplicated => { self.metrics().discovery.req_deduplicated.inc(); debug!( @@ -569,7 +569,7 @@ impl Node { // If no tree peers had the target, fail immediately 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!( target_node = %self.peer_display_name(dest), "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) { 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); + crate::proto::lookup::poll_pending(&mut self.discovery, now_ms, &attempt_timeouts); for (target, attempt) in outcome.retries { let ttl = self.config().node.discovery.ttl; diff --git a/src/node/mod.rs b/src/node/mod.rs index c5aa8e5..dfced01 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -43,9 +43,9 @@ use crate::cache::CoordCache; use crate::node::session::SessionEntry; use crate::peer::{ActivePeer, PeerConnection}; use crate::proto::bloom::{BloomFilter, BloomState}; -use crate::proto::discovery::{Discovery, DiscoveryBackoff, DiscoveryForwardRateLimiter}; use crate::proto::fmp::Fmp; use crate::proto::fsp::Fsp; +use crate::proto::lookup::{Lookup, LookupBackoff, LookupForwardRateLimiter}; use crate::proto::mmp::Mmp; use crate::proto::routing::{self, Router, RoutingErrorRateLimiter}; use crate::proto::stp::TreeState; @@ -337,7 +337,7 @@ pub struct Node { // === Discovery === /// Discovery-subsystem state: recent-request dedup cache, in-flight /// lookups, originator-side backoff, and transit-side forward limiter. - discovery: Discovery, + discovery: Lookup, // === Counters === /// Next link ID to allocate. @@ -678,9 +678,9 @@ impl Node { coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval_ms( coords_response_interval_ms, ), - discovery: Discovery::new( - DiscoveryBackoff::with_params(backoff_base_secs, backoff_max_secs), - DiscoveryForwardRateLimiter::with_interval_ms(forward_min_interval_secs * 1000), + discovery: Lookup::new( + LookupBackoff::with_params(backoff_base_secs, backoff_max_secs), + LookupForwardRateLimiter::with_interval_ms(forward_min_interval_secs * 1000), ), pending_connects: Vec::new(), retry_pending: HashMap::new(), @@ -843,7 +843,7 @@ impl Node { coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval_ms( coords_response_interval_ms, ), - discovery: Discovery::new(DiscoveryBackoff::new(), DiscoveryForwardRateLimiter::new()), + discovery: Lookup::new(LookupBackoff::new(), LookupForwardRateLimiter::new()), pending_connects: Vec::new(), retry_pending: HashMap::new(), nostr_discovery: None, @@ -2541,7 +2541,7 @@ impl Node { /// Iterate over pending discovery lookups for diagnostics. pub fn pending_lookups_iter( &self, - ) -> impl Iterator { + ) -> impl Iterator { self.discovery.pending_lookups.iter() } diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index 21eb6fd..f1e0de2 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -5,7 +5,7 @@ //! response routing. use super::*; -use crate::proto::discovery::{LookupRequest, LookupResponse, RecentRequest}; +use crate::proto::lookup::{LookupRequest, LookupResponse, RecentRequest}; use crate::proto::stp::TreeCoordinate; use spanning_tree::{ cleanup_nodes, generate_random_edges, lock_large_network_test, process_available_packets, @@ -1040,7 +1040,7 @@ async fn test_open_discovery_sweep_queues_eligible_skips_filtered() { async fn test_check_pending_lookups_default_sequence_unreachable() { use crate::peer::ActivePeer; use crate::proto::bloom::BloomFilter; - use crate::proto::discovery::PendingLookup; + use crate::proto::lookup::PendingLookup; use crate::transport::LinkId; use std::sync::mpsc; diff --git a/src/proto/discovery/core.rs b/src/proto/lookup/core.rs similarity index 79% rename from src/proto/discovery/core.rs rename to src/proto/lookup/core.rs index 0fe72c4..9af173a 100644 --- a/src/proto/discovery/core.rs +++ b/src/proto/lookup/core.rs @@ -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, //! calls into this core, and drives the returned actions (the actual //! encrypted sends). No I/O, no clock, no metrics, no logging here. use alloc::sync::Arc; -use super::state::{Discovery, PendingLookup, RecentRequest}; +use super::state::{Lookup, PendingLookup, RecentRequest}; use super::wire::LookupRequest; 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 /// 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. -pub(crate) enum DiscoveryAction { - /// Send an encoded discovery PDU to a peer as an encrypted link message. +pub(crate) enum LookupAction { + /// Send an encoded lookup PDU to a peer as an encrypted link message. /// `bytes` is `Arc`-shared so a fan-out encodes once. SendLink { peer: NodeAddr, bytes: Arc<[u8]> }, /// 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 /// the non-tree bloom-match fallback set was used (no tree peer matched). Forward { - actions: Vec, + actions: Vec, 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 actions = targets .into_iter() - .map(|peer| DiscoveryAction::SendLink { + .map(|peer| LookupAction::SendLink { peer, 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; /// it is a known origination gap (ISSUE-2026-0059) whose fix adds the fallback /// branch as a separate, behavior-changing change. -pub(crate) fn plan_initiate( - request: &LookupRequest, - rv: &impl RoutingView, -) -> Vec { +pub(crate) fn plan_initiate(request: &LookupRequest, rv: &impl RoutingView) -> Vec { let targets: Vec = rv .peers_reaching(&request.target) .into_iter() @@ -126,14 +123,14 @@ pub(crate) fn plan_initiate( let bytes: Arc<[u8]> = Arc::from(request.encode()); targets .into_iter() - .map(|peer| DiscoveryAction::SendLink { + .map(|peer| LookupAction::SendLink { peer, bytes: bytes.clone(), }) .collect() } -/// Classification of an inbound LookupRequest, decided from Discovery state. +/// Classification of an inbound LookupRequest, decided from Lookup state. pub(crate) enum RequestOutcome { /// request_id already in the dedup cache — drop. Duplicate, @@ -152,9 +149,9 @@ pub(crate) enum RequestOutcome { /// Classify an inbound LookupRequest against the recent-request dedup cache and /// the transit forward rate limiter. Purges expired dedup entries, records 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( - disc: &mut Discovery, + lookup: &mut Lookup, request: &LookupRequest, from: &NodeAddr, my_addr: &NodeAddr, @@ -163,25 +160,30 @@ pub(crate) fn classify_request( max_recent: usize, ) -> RequestOutcome { // Purge expired dedup entries (was purge_expired_requests). - disc.recent_requests + lookup + .recent_requests .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; } - if disc.recent_requests.len() >= max_recent { + if lookup.recent_requests.len() >= max_recent { 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)); if request.target == *my_addr { return RequestOutcome::RespondAsTarget; } 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 } else { RequestOutcome::ForwardRateLimited @@ -205,10 +207,10 @@ pub(crate) enum ResponseRoute { /// 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. -pub(crate) fn classify_response(disc: &mut Discovery, request_id: u64) -> ResponseRoute { - match disc.recent_requests.get_mut(&request_id) { +pub(crate) fn classify_response(lookup: &mut Lookup, request_id: u64) -> ResponseRoute { + match lookup.recent_requests.get_mut(&request_id) { Some(recent) => { if recent.response_forwarded { ResponseRoute::AlreadyForwarded @@ -233,13 +235,13 @@ pub(crate) enum ResponseRouteDecision { } /// 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 /// 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. -pub(crate) fn plan_response_route(disc: &Discovery, request_id: u64) -> ResponseRouteDecision { - match disc.recent_requests.get(&request_id) { +pub(crate) fn plan_response_route(lookup: &Lookup, request_id: u64) -> ResponseRouteDecision { + match lookup.recent_requests.get(&request_id) { Some(recent) => ResponseRouteDecision::ReversePath(recent.from_peer), 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. /// -/// 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. /// Verification is the shell's job — this runs only after the proof checked out. pub(crate) fn on_response_accepted( - disc: &mut Discovery, + lookup: &mut Lookup, target: &NodeAddr, coords: crate::TreeCoordinate, now_ms: u64, path_mtu: u16, -) -> Vec { - disc.backoff.record_success(target); - disc.pending_lookups.remove(target); +) -> Vec { + lookup.backoff.record_success(target); + lookup.pending_lookups.remove(target); vec![ - DiscoveryAction::CacheCoords { + LookupAction::CacheCoords { target: *target, coords, now_ms, path_mtu, }, - DiscoveryAction::WritePathMtu { + LookupAction::WritePathMtu { target: *target, path_mtu, }, - DiscoveryAction::ResetWarmupIfEstablished { target: *target }, - DiscoveryAction::RetryQueuedPackets { target: *target }, + LookupAction::ResetWarmupIfEstablished { target: *target }, + LookupAction::RetryQueuedPackets { target: *target }, ] } @@ -287,11 +289,11 @@ pub(crate) struct PollOutcome { 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 /// final timeouts (removed + backoff failure recorded). No I/O, no view. pub(crate) fn poll_pending( - disc: &mut Discovery, + lookup: &mut Lookup, now_ms: u64, attempt_timeouts_secs: &[u64], ) -> PollOutcome { @@ -301,7 +303,7 @@ pub(crate) fn poll_pending( let mut retry_targets: Vec = Vec::new(); let mut timeout_targets: Vec = 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 to_ms = attempt_timeouts_secs.get(idx).copied().unwrap_or(0) * 1000; 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(); 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.last_sent_ms = now_ms; retries.push((target, entry.attempt)); @@ -324,16 +326,16 @@ pub(crate) fn poll_pending( let mut timeouts: Vec<(NodeAddr, u32)> = Vec::new(); for target in timeout_targets { - disc.pending_lookups.remove(&target); - disc.backoff.record_failure(&target, now_ms); - let failures = disc.backoff.failure_count(&target); + lookup.pending_lookups.remove(&target); + lookup.backoff.record_failure(&target, now_ms); + let failures = lookup.backoff.failure_count(&target); timeouts.push((target, failures)); } 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 { /// A lookup is already pending for this target — skip. Deduplicated, @@ -345,36 +347,37 @@ pub(crate) enum InitiateDecision { 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 /// 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( - disc: &mut Discovery, + lookup: &mut Lookup, dest: &NodeAddr, now_ms: u64, reachable: bool, ) -> InitiateDecision { - if disc.pending_lookups.contains_key(dest) { + if lookup.pending_lookups.contains_key(dest) { return InitiateDecision::Deduplicated; } - if disc.backoff.is_suppressed(dest, now_ms) { + if lookup.backoff.is_suppressed(dest, now_ms) { return InitiateDecision::Suppressed { - failures: disc.backoff.failure_count(dest), + failures: lookup.backoff.failure_count(dest), }; } if !reachable { - disc.backoff.record_failure(dest, now_ms); + lookup.backoff.record_failure(dest, now_ms); return InitiateDecision::BloomMiss; } - disc.pending_lookups + lookup + .pending_lookups .insert(*dest, PendingLookup::new(now_ms)); InitiateDecision::Proceed } /// Roll back a lookup whose first attempt reached no tree peers (sent == 0): /// drop the pending entry and record a backoff failure. -pub(crate) fn initiate_failed(disc: &mut Discovery, dest: &NodeAddr, now_ms: u64) { - disc.pending_lookups.remove(dest); - disc.backoff.record_failure(dest, now_ms); +pub(crate) fn initiate_failed(lookup: &mut Lookup, dest: &NodeAddr, now_ms: u64) { + lookup.pending_lookups.remove(dest); + lookup.backoff.record_failure(dest, now_ms); } diff --git a/src/proto/discovery/limits.rs b/src/proto/lookup/limits.rs similarity index 90% rename from src/proto/discovery/limits.rs rename to src/proto/lookup/limits.rs index b92a4b1..b4cb26f 100644 --- a/src/proto/discovery/limits.rs +++ b/src/proto/lookup/limits.rs @@ -1,15 +1,15 @@ -//! Discovery protocol rate limiting and backoff. +//! Mesh lookup protocol rate limiting and backoff. //! //! Two complementary mechanisms: //! -//! - **`DiscoveryBackoff`** (originator-side, optional): Exponential +//! - **`LookupBackoff`** (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 +//! - **`LookupForwardRateLimiter`** (transit-side): Per-target minimum //! interval for forwarded requests. Defense-in-depth against misbehaving //! 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 /// 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. @@ -35,11 +35,11 @@ const DEFAULT_BACKOFF_BASE_SECS: u64 = 0; /// Default maximum backoff cap. `0` = disabled. 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 /// re-initiation with increasing delays. Cleared on topology changes. -pub struct DiscoveryBackoff { +pub struct LookupBackoff { /// Maps target → (suppress_until, consecutive_failures). pub(crate) entries: BTreeMap, /// Base backoff in milliseconds (first failure). @@ -55,7 +55,7 @@ pub(crate) struct BackoffEntry { failures: u32, } -impl DiscoveryBackoff { +impl LookupBackoff { /// Create with default parameters (disabled — base/cap = 0). pub fn new() -> Self { 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 { Self::new() } } // ============================================================================ -// Transit-side: Discovery Forward Rate Limiter +// Transit-side: Lookup Forward Rate Limiter // ============================================================================ /// 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. 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 /// and enforces a minimum interval to prevent floods from misbehaving /// nodes generating fresh request_ids. -pub struct DiscoveryForwardRateLimiter(PerAddrRateLimiter); +pub struct LookupForwardRateLimiter(PerAddrRateLimiter); -impl DiscoveryForwardRateLimiter { +impl LookupForwardRateLimiter { /// Create with default parameters (2s interval). pub fn new() -> Self { Self(PerAddrRateLimiter::new( @@ -201,7 +201,7 @@ impl DiscoveryForwardRateLimiter { } } -impl Default for DiscoveryForwardRateLimiter { +impl Default for LookupForwardRateLimiter { fn default() -> Self { Self::new() } diff --git a/src/proto/discovery/mod.rs b/src/proto/lookup/mod.rs similarity index 61% rename from src/proto/discovery/mod.rs rename to src/proto/lookup/mod.rs index 883b85e..22a5a12 100644 --- a/src/proto/discovery/mod.rs +++ b/src/proto/lookup/mod.rs @@ -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 -//! `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-migrates-with-subsystem policy. //! //! The sans-IO decision core lives in `core.rs`: it defines the `RoutingView` //! 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 //! bytes, calls the planner, and drives the returned actions. @@ -21,15 +21,13 @@ mod wire; mod tests; pub(crate) use core::{ - DiscoveryAction, ForwardOutcome, InitiateDecision, RequestOutcome, ResponseRoute, + ForwardOutcome, InitiateDecision, LookupAction, RequestOutcome, ResponseRoute, ResponseRouteDecision, RoutingView, classify_request, classify_response, initiate_failed, initiate_gate, on_response_accepted, plan_forward, plan_initiate, plan_response_route, poll_pending, }; -pub(crate) use limits::{ - DiscoveryBackoff, DiscoveryForwardRateLimiter, MAX_RECENT_DISCOVERY_REQUESTS, -}; +pub(crate) use limits::{LookupBackoff, LookupForwardRateLimiter, MAX_RECENT_LOOKUP_REQUESTS}; #[cfg(test)] pub(crate) use state::RecentRequest; -pub(crate) use state::{Discovery, PendingLookup}; +pub(crate) use state::{Lookup, PendingLookup}; pub use wire::{LookupRequest, LookupResponse}; diff --git a/src/proto/discovery/state.rs b/src/proto/lookup/state.rs similarity index 73% rename from src/proto/discovery/state.rs rename to src/proto/lookup/state.rs index aa188f0..45f1c26 100644 --- a/src/proto/discovery/state.rs +++ b/src/proto/lookup/state.rs @@ -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 -//! 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 //! `Node`. use alloc::collections::BTreeMap; -use super::limits::{DiscoveryBackoff, DiscoveryForwardRateLimiter}; +use super::limits::{LookupBackoff, LookupForwardRateLimiter}; use crate::NodeAddr; /// 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 { /// When the lookup was first initiated. pub initiated_ms: u64, @@ -64,30 +64,27 @@ impl PendingLookup { } } -/// Discovery-subsystem state. -pub(crate) struct Discovery { - /// Recent discovery requests (dedup + reverse-path forwarding). +/// Mesh lookup subsystem state. +pub(crate) struct Lookup { + /// Recent lookup requests (dedup + reverse-path forwarding). /// Maps request_id → RecentRequest. pub(crate) recent_requests: BTreeMap, - /// 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. pub(crate) pending_lookups: BTreeMap, - /// Backoff for failed discovery lookups (originator-side). - pub(crate) backoff: DiscoveryBackoff, - /// Rate limiter for forwarded discovery requests (transit-side). - pub(crate) forward_limiter: DiscoveryForwardRateLimiter, + /// Backoff for failed lookups (originator-side). + pub(crate) backoff: LookupBackoff, + /// Rate limiter for forwarded lookup requests (transit-side). + pub(crate) forward_limiter: LookupForwardRateLimiter, } -impl Discovery { - /// Create discovery state with the given backoff and forward limiter. +impl Lookup { + /// Create mesh lookup state with the given backoff and forward limiter. /// /// The two limiters are constructed by the caller so each `Node` /// constructor can supply its own configured/default variant, matching /// the pre-refactor initialization exactly. - pub(crate) fn new( - backoff: DiscoveryBackoff, - forward_limiter: DiscoveryForwardRateLimiter, - ) -> Self { + pub(crate) fn new(backoff: LookupBackoff, forward_limiter: LookupForwardRateLimiter) -> Self { Self { recent_requests: 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 — /// observability stays out of the pure core. pub(crate) fn reset_backoff(&mut self) -> usize { diff --git a/src/proto/discovery/tests/core.rs b/src/proto/lookup/tests/core.rs similarity index 70% rename from src/proto/discovery/tests/core.rs rename to src/proto/lookup/tests/core.rs index 4859e2a..a783c3b 100644 --- a/src/proto/discovery/tests/core.rs +++ b/src/proto/lookup/tests/core.rs @@ -1,11 +1,10 @@ -//! Tests for the sans-IO discovery decision core. +//! Tests for the sans-IO lookup decision core. use super::util::{ - MockRoutingView, action_peers, empty_discovery, make_request, make_request_id, - suppressing_discovery, + MockRoutingView, action_peers, empty_lookup, make_request, make_request_id, suppressing_lookup, }; use crate::TreeCoordinate; -use crate::proto::discovery::*; +use crate::proto::lookup::*; use crate::testutil::make_node_addr; #[test] @@ -125,10 +124,12 @@ fn initiate_returns_empty_when_nothing_reaches_target() { #[test] fn response_route_uses_recorded_reverse_path() { - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); let from = make_node_addr(9); - disc.recent_requests.insert(42, RecentRequest::new(from, 0)); - match plan_response_route(&disc, 42) { + lookup + .recent_requests + .insert(42, RecentRequest::new(from, 0)); + match plan_response_route(&lookup, 42) { ResponseRouteDecision::ReversePath(peer) => assert_eq!(peer, from), ResponseRouteDecision::NeedsTreeRoute => panic!("expected ReversePath"), } @@ -136,9 +137,9 @@ fn response_route_uses_recorded_reverse_path() { #[test] fn response_route_needs_tree_route_when_no_record() { - let disc = empty_discovery(); + let lookup = empty_lookup(); assert!(matches!( - plan_response_route(&disc, 42), + plan_response_route(&lookup, 42), ResponseRouteDecision::NeedsTreeRoute )); } @@ -146,40 +147,42 @@ fn response_route_needs_tree_route_when_no_record() { #[test] fn classify_response_transit_on_fresh_forwarded_request() { let from_peer = make_node_addr(0x11); - let mut disc = empty_discovery(); - disc.recent_requests + let mut lookup = empty_lookup(); + lookup + .recent_requests .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), _ => panic!("expected 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] fn classify_response_already_forwarded_on_second_call() { let from_peer = make_node_addr(0x22); - let mut disc = empty_discovery(); - disc.recent_requests + let mut lookup = empty_lookup(); + lookup + .recent_requests .insert(7, RecentRequest::new(from_peer, 1000)); assert!(matches!( - classify_response(&mut disc, 7), + classify_response(&mut lookup, 7), ResponseRoute::Transit { .. } )); assert!(matches!( - classify_response(&mut disc, 7), + classify_response(&mut lookup, 7), ResponseRoute::AlreadyForwarded )); } #[test] fn classify_response_originator_when_request_absent() { - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); assert!(matches!( - classify_response(&mut disc, 999), + classify_response(&mut lookup, 999), ResponseRoute::Originator )); } @@ -187,34 +190,35 @@ fn classify_response_originator_when_request_absent() { #[test] fn on_response_accepted_clears_state_and_emits_effects() { 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. - disc.backoff.record_failure(&target, 1000); - assert!(!disc.backoff.is_empty(), "precondition: backoff seeded"); - disc.pending_lookups + lookup.backoff.record_failure(&target, 1000); + assert!(!lookup.backoff.is_empty(), "precondition: backoff seeded"); + lookup + .pending_lookups .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 now_ms = 12_345u64; 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. assert!( - disc.backoff.is_empty(), + lookup.backoff.is_empty(), "backoff entry must clear on success" ); assert!( - !disc.pending_lookups.contains_key(&target), + !lookup.pending_lookups.contains_key(&target), "pending lookup must be dropped" ); // Exactly the four effect actions, in order. assert_eq!(actions.len(), 4, "expected four effect actions"); match &actions[0] { - DiscoveryAction::CacheCoords { + LookupAction::CacheCoords { target: t, now_ms: n, path_mtu: p, @@ -227,7 +231,7 @@ fn on_response_accepted_clears_state_and_emits_effects() { _ => panic!("action[0] must be CacheCoords"), } match &actions[1] { - DiscoveryAction::WritePathMtu { + LookupAction::WritePathMtu { target: t, path_mtu: p, } => { @@ -237,11 +241,11 @@ fn on_response_accepted_clears_state_and_emits_effects() { _ => panic!("action[1] must be WritePathMtu"), } 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"), } 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"), } } @@ -249,17 +253,19 @@ fn on_response_accepted_clears_state_and_emits_effects() { #[test] fn poll_pending_no_action_before_first_deadline() { let target = make_node_addr(0x30); - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); 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. - 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.timeouts.is_empty(), "no timeout before deadline"); // 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.last_sent_ms, t0); } @@ -267,17 +273,19 @@ fn poll_pending_no_action_before_first_deadline() { #[test] fn poll_pending_retries_at_first_deadline() { let target = make_node_addr(0x31); - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); 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. - 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!(outcome.timeouts.is_empty()); // 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.last_sent_ms, t0 + 1000); } @@ -285,17 +293,17 @@ fn poll_pending_retries_at_first_deadline() { #[test] fn poll_pending_final_timeout_at_max_attempt() { 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. let tn = 50_000u64; let mut entry = PendingLookup::new(tn); entry.attempt = 4; 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. - 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_eq!( outcome.timeouts, @@ -305,89 +313,98 @@ fn poll_pending_final_timeout_at_max_attempt() { // Entry removed and a backoff failure recorded. assert!( - !disc.pending_lookups.contains_key(&target), + !lookup.pending_lookups.contains_key(&target), "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 --- #[test] 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 my_addr = make_node_addr(0x99); let target = make_node_addr(0xAA); 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)); // Recorded for reverse-path forwarding. - assert!(disc.recent_requests.contains_key(&1)); - assert_eq!(disc.recent_requests.get(&1).unwrap().from_peer, from); + assert!(lookup.recent_requests.contains_key(&1)); + assert_eq!(lookup.recent_requests.get(&1).unwrap().from_peer, from); } #[test] fn classify_request_duplicate_on_second_call() { - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); let from = make_node_addr(0x01); let my_addr = make_node_addr(0x99); let target = make_node_addr(0xAA); let request = make_request_id(1, target, 3); 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 )); 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 )); } #[test] fn classify_request_dedup_cache_full() { - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); let from = make_node_addr(0x01); let my_addr = make_node_addr(0x99); let target = make_node_addr(0xAA); // Fill the cache to max_recent with distinct request_ids. let max_recent = 3usize; for id in 100..(100 + max_recent as u64) { - disc.recent_requests + lookup + .recent_requests .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); - 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), _ => panic!("expected DedupCacheFull"), } // 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] fn classify_request_respond_as_target() { - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); let from = make_node_addr(0x01); let my_addr = make_node_addr(0xAA); // target == my_addr let request = make_request_id(1, my_addr, 3); 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 )); // Recorded before the target decision. - assert!(disc.recent_requests.contains_key(&1)); + assert!(lookup.recent_requests.contains_key(&1)); } #[test] 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 my_addr = make_node_addr(0x99); 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); 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 )); } #[test] fn classify_request_forward_rate_limited() { - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); let from = make_node_addr(0x01); let my_addr = make_node_addr(0x99); let target = make_node_addr(0xAA); // Pre-seed the forward limiter so should_forward(target) returns false // 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); 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 )); } #[test] fn classify_request_purges_expired_entries() { - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); let from = make_node_addr(0x01); let my_addr = make_node_addr(0x99); let target = make_node_addr(0xAA); // Seed an entry that is expired at now_ms with the given expiry window. // is_expired: now - timestamp > expiry_ms → expired. - disc.recent_requests + lookup + .recent_requests .insert(55, RecentRequest::new(from, 1000)); // now_ms = 10_000, expiry_ms = 5000 → 9000 > 5000 → expired. 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)); // 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. - assert!(disc.recent_requests.contains_key(&1)); + assert!(lookup.recent_requests.contains_key(&1)); } #[test] fn poll_pending_full_ladder_end_to_end() { let target = make_node_addr(0x33); - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); 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]; // 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)]); // 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)]); // 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)]); // attempt 4 is max → final timeout at deadline 8s after last send 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_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 --- @@ -467,12 +487,12 @@ fn poll_pending_full_ladder_end_to_end() { #[test] fn initiate_gate_deduplicated_when_pending() { let dest = make_node_addr(0x40); - let mut disc = empty_discovery(); - disc.pending_lookups.insert(dest, PendingLookup::new(500)); + let mut lookup = empty_lookup(); + lookup.pending_lookups.insert(dest, PendingLookup::new(500)); // reachable=true would otherwise Proceed, but the pending entry wins. assert!(matches!( - initiate_gate(&mut disc, &dest, 1000, true), + initiate_gate(&mut lookup, &dest, 1000, true), InitiateDecision::Deduplicated )); } @@ -480,48 +500,48 @@ fn initiate_gate_deduplicated_when_pending() { #[test] fn initiate_gate_suppressed_by_backoff() { 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). - disc.backoff.record_failure(&dest, 1000); + lookup.backoff.record_failure(&dest, 1000); assert!( - disc.backoff.is_suppressed(&dest, 1000), + lookup.backoff.is_suppressed(&dest, 1000), "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), _ => panic!("expected Suppressed"), } // No pending entry was inserted on the suppress path. - assert!(!disc.pending_lookups.contains_key(&dest)); + assert!(!lookup.pending_lookups.contains_key(&dest)); } #[test] fn initiate_gate_bloom_miss_records_failure() { let dest = make_node_addr(0x42); - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); assert!(matches!( - initiate_gate(&mut disc, &dest, 1000, false), + initiate_gate(&mut lookup, &dest, 1000, false), InitiateDecision::BloomMiss )); // A backoff failure was recorded, and no pending entry created. - assert_eq!(disc.backoff.failure_count(&dest), 1); - assert!(!disc.pending_lookups.contains_key(&dest)); + assert_eq!(lookup.backoff.failure_count(&dest), 1); + assert!(!lookup.pending_lookups.contains_key(&dest)); } #[test] fn initiate_gate_proceed_inserts_pending() { let dest = make_node_addr(0x43); - let mut disc = empty_discovery(); + let mut lookup = empty_lookup(); let now_ms = 7_777u64; assert!(matches!( - initiate_gate(&mut disc, &dest, now_ms, true), + initiate_gate(&mut lookup, &dest, now_ms, true), InitiateDecision::Proceed )); // The pending entry now exists, stamped with now_ms. - let entry = disc + let entry = lookup .pending_lookups .get(&dest) .expect("Proceed must insert a pending lookup"); @@ -532,13 +552,15 @@ fn initiate_gate_proceed_inserts_pending() { #[test] fn initiate_failed_drops_pending_and_records_failure() { let dest = make_node_addr(0x44); - let mut disc = empty_discovery(); - disc.pending_lookups.insert(dest, PendingLookup::new(1000)); + let mut lookup = empty_lookup(); + lookup + .pending_lookups + .insert(dest, PendingLookup::new(1000)); - initiate_failed(&mut disc, &dest, 1000); + initiate_failed(&mut lookup, &dest, 1000); assert!( - !disc.pending_lookups.contains_key(&dest), + !lookup.pending_lookups.contains_key(&dest), "pending entry must be dropped" ); - assert_eq!(disc.backoff.failure_count(&dest), 1); + assert_eq!(lookup.backoff.failure_count(&dest), 1); } diff --git a/src/proto/discovery/tests/limits.rs b/src/proto/lookup/tests/limits.rs similarity index 78% rename from src/proto/discovery/tests/limits.rs rename to src/proto/lookup/tests/limits.rs index 2422284..993df74 100644 --- a/src/proto/discovery/tests/limits.rs +++ b/src/proto/lookup/tests/limits.rs @@ -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; -// --- DiscoveryBackoff tests --- +// --- LookupBackoff tests --- #[test] fn test_backoff_not_suppressed_initially() { - let backoff = DiscoveryBackoff::new(); + let backoff = LookupBackoff::new(); assert!(!backoff.is_suppressed(&addr(1), 0)); } @@ -15,7 +15,7 @@ fn test_backoff_not_suppressed_initially() { fn test_backoff_suppressed_after_failure() { // Backoff is opt-in; exercise the suppression path with explicit params. 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); assert!(backoff.is_suppressed(&addr(1), now)); // Different target not affected @@ -25,7 +25,7 @@ fn test_backoff_suppressed_after_failure() { #[test] fn test_backoff_cleared_on_success() { 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); assert!(backoff.is_suppressed(&addr(1), now)); @@ -36,7 +36,7 @@ fn test_backoff_cleared_on_success() { #[test] fn test_backoff_reset_all() { let now = 1_000; - let mut backoff = DiscoveryBackoff::new(); + let mut backoff = LookupBackoff::new(); backoff.record_failure(&addr(1), now); backoff.record_failure(&addr(2), now); assert_eq!(backoff.len(), 2); @@ -49,7 +49,7 @@ fn test_backoff_reset_all() { #[test] fn test_backoff_exponential() { let now = 1_000; - let mut backoff = DiscoveryBackoff::with_params(1, 300); + let mut backoff = LookupBackoff::with_params(1, 300); // First failure: 1s backoff backoff.record_failure(&addr(1), now); @@ -67,7 +67,7 @@ fn test_backoff_exponential() { #[test] fn test_backoff_expires() { 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); // With 0s backoff, should not be suppressed assert!(!backoff.is_suppressed(&addr(1), now)); @@ -76,7 +76,7 @@ fn test_backoff_expires() { #[test] fn test_backoff_capped() { let now = 1_000; - let mut backoff = DiscoveryBackoff::with_params(1, 10); + let mut backoff = LookupBackoff::with_params(1, 10); // Record many failures for _ in 0..20 { @@ -89,18 +89,18 @@ fn test_backoff_capped() { assert!(remaining <= 11_000); } -// --- DiscoveryForwardRateLimiter tests --- +// --- LookupForwardRateLimiter tests --- #[test] fn test_forward_first_allowed() { - let mut limiter = DiscoveryForwardRateLimiter::new(); + let mut limiter = LookupForwardRateLimiter::new(); assert!(limiter.should_forward(&addr(1), 0)); } #[test] fn test_forward_rapid_rate_limited() { 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)); @@ -109,7 +109,7 @@ fn test_forward_rapid_rate_limited() { #[test] fn test_forward_different_targets_independent() { 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(2), now)); assert!(!limiter.should_forward(&addr(1), now)); @@ -119,7 +119,7 @@ fn test_forward_different_targets_independent() { #[test] fn test_forward_allowed_after_interval() { 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)); // Advance past the minimum interval. @@ -129,7 +129,7 @@ fn test_forward_allowed_after_interval() { #[test] fn test_forward_cleanup_removes_old() { 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(2), now)); assert_eq!(limiter.len(), 2); @@ -142,7 +142,7 @@ fn test_forward_cleanup_removes_old() { #[test] fn test_forward_cleanup_preserves_recent() { let now = 1_000; - let mut limiter = DiscoveryForwardRateLimiter::new(); + let mut limiter = LookupForwardRateLimiter::new(); assert!(limiter.should_forward(&addr(1), now)); assert_eq!(limiter.len(), 1); diff --git a/src/proto/discovery/tests/mod.rs b/src/proto/lookup/tests/mod.rs similarity index 58% rename from src/proto/discovery/tests/mod.rs rename to src/proto/lookup/tests/mod.rs index 9e7a06d..5850239 100644 --- a/src/proto/discovery/tests/mod.rs +++ b/src/proto/lookup/tests/mod.rs @@ -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`. mod core; diff --git a/src/proto/discovery/tests/util.rs b/src/proto/lookup/tests/util.rs similarity index 78% rename from src/proto/discovery/tests/util.rs rename to src/proto/lookup/tests/util.rs index f81c946..85305c7 100644 --- a/src/proto/discovery/tests/util.rs +++ b/src/proto/lookup/tests/util.rs @@ -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 crate::proto::discovery::{ - Discovery, DiscoveryAction, DiscoveryBackoff, DiscoveryForwardRateLimiter, LookupRequest, - LookupResponse, RoutingView, +use crate::proto::lookup::{ + Lookup, LookupAction, LookupBackoff, LookupForwardRateLimiter, LookupRequest, LookupResponse, + RoutingView, }; use crate::testutil::make_node_addr; 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() } -pub(super) fn action_peers(actions: &[DiscoveryAction]) -> Vec { +pub(super) fn action_peers(actions: &[LookupAction]) -> Vec { actions .iter() .map(|action| match action { - DiscoveryAction::SendLink { peer, .. } => *peer, + LookupAction::SendLink { peer, .. } => *peer, _ => panic!("expected SendLink, got a different action variant"), }) .collect() } -pub(super) fn empty_discovery() -> Discovery { - Discovery::new( - DiscoveryBackoff::default(), - DiscoveryForwardRateLimiter::default(), +pub(super) fn empty_lookup() -> Lookup { + Lookup::new( + LookupBackoff::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. -pub(super) fn suppressing_discovery() -> Discovery { - Discovery::new( - DiscoveryBackoff::with_params(30, 300), - DiscoveryForwardRateLimiter::default(), +pub(super) fn suppressing_lookup() -> Lookup { + Lookup::new( + LookupBackoff::with_params(30, 300), + LookupForwardRateLimiter::default(), ) } diff --git a/src/proto/discovery/tests/wire.rs b/src/proto/lookup/tests/wire.rs similarity index 97% rename from src/proto/discovery/tests/wire.rs rename to src/proto/lookup/tests/wire.rs index 29f5dc1..52b7137 100644 --- a/src/proto/discovery/tests/wire.rs +++ b/src/proto/lookup/tests/wire.rs @@ -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 crate::proto::discovery::{LookupRequest, LookupResponse}; +use crate::proto::lookup::{LookupRequest, LookupResponse}; use crate::testutil::make_node_addr; #[test] diff --git a/src/proto/discovery/wire.rs b/src/proto/lookup/wire.rs similarity index 98% rename from src/proto/discovery/wire.rs rename to src/proto/lookup/wire.rs index c31c539..33e5090 100644 --- a/src/proto/discovery/wire.rs +++ b/src/proto/lookup/wire.rs @@ -1,4 +1,4 @@ -//! Discovery messages: LookupRequest and LookupResponse. +//! Mesh lookup messages: LookupRequest and LookupResponse. use crate::NodeAddr; use crate::proto::Error; @@ -7,7 +7,7 @@ use crate::proto::stp::TreeCoordinate; use crate::proto::stp::{decode_coords, encode_coords}; 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. /// Each transit node forwards only to tree peers whose bloom filter diff --git a/src/proto/mod.rs b/src/proto/mod.rs index ac08f03..d466711 100644 --- a/src/proto/mod.rs +++ b/src/proto/mod.rs @@ -9,10 +9,10 @@ pub use error::Error; pub(crate) mod bloom; pub(crate) mod codec; pub(crate) mod coord; -pub(crate) mod discovery; pub(crate) mod fmp; pub(crate) mod fsp; pub(crate) mod link; +pub(crate) mod lookup; pub(crate) mod math; pub(crate) mod mmp; pub(crate) mod rate_limit;