From 03ced618ce4ed8201043246a96acab725440475b Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Mon, 13 Jul 2026 05:34:28 +0000 Subject: [PATCH] node/peering: drive Nostr open-discovery through the reconciler overlay layer Cut the overlay (Nostr open-discovery) enqueue over to the sans-IO reconciler. A gate-checked reconcile_overlay wrapper runs the overlay layer alone at the discovery tick slot; the monolithic reconcile would re-fire the always-on retry-dial and double the per-tick dial cap, so the overlay slot must call only its layer. The driver builds the candidate pool from the overlay advert cache (self excluded), the cooldown and configured-npub sets, and the startup/steady max-age, then feeds the reconciler; the emitted enqueue set drives the per-enqueue identity-cache and alias pre-seed exactly as before. Enqueued entries are dialed at the retry slot, preserving the two-phase cadence. Delete the now-dead enqueue-budget and expiry helpers. Behavior-neutral: same candidates enqueued in the same order with the same budget, cap, and expiry. Opportunistic transport-neighbor growth remains imperative and is cut over next. Unit test count unchanged (1607). --- src/node/lifecycle/mod.rs | 317 +++++++++------------------------- src/node/peering/reconcile.rs | 33 ++++ src/node/tests/discovery.rs | 9 +- 3 files changed, 123 insertions(+), 236 deletions(-) diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 5209473..7405eac 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -2152,25 +2152,36 @@ impl Node { } async fn queue_open_discovery_retries(&mut self, bootstrap: &std::sync::Arc) { - self.run_open_discovery_sweep(bootstrap, None, "per-tick") - .await; + self.run_open_discovery_sweep(bootstrap, None).await; } - /// Open-discovery cache sweep. Iterates the cached overlay adverts and - /// queues retries for non-configured, not-yet-connected peers. + /// Open-discovery cache sweep — drains the cached overlay adverts and + /// enqueues retries for eligible peers via the sans-IO reconciler's overlay + /// layer. /// - /// `max_age_secs`, if set, filters out adverts whose `created_at` is - /// older than `now - max_age_secs`. The per-tick sweep passes `None` - /// (relies on the cache's own `valid_until_ms` filter); the one-shot - /// startup sweep passes `Some(startup_sweep_max_age_secs)`. + /// The driver builds the [`DiscoveryPools`] overlay input from + /// `bootstrap.cached_open_discovery_candidates(64)` (the I/O), excluding the + /// node's own advert (obligation O6 — the sans-IO core has no self-identity + /// input), and supplies the configured-npub set, the per-npub cooldown set, + /// and the startup-sweep max-age. `max_age_secs` is `None` for the per-tick + /// sweep and `Some(startup_sweep_max_age_secs)` for the one-shot startup + /// sweep. /// - /// `caller` is a short label included in log lines so per-tick and - /// startup sweeps are distinguishable in operator-facing logs. + /// The core reproduces the old sweep's full skip order, configured-advert + /// expedite, and enqueue budget internally: it inserts due-now entries into + /// the relocated `retry_pending`, and the retry-slot `process_pending_retries` + /// dials them later in the same tick (two-phase, design §2.5). Per design §10 + /// this calls the reconciler's `reconcile_overlay` (the overlay layer only) — + /// NOT the monolithic `reconcile()` — so the always-on retry-dial phase does + /// not re-fire at this slot (which would double-dial the due entries and + /// apply the per-tick 16-cap twice). The returned `ScheduleRetry` actions + /// name exactly the newly enqueued peers; the driver consumes them to + /// pre-seed the alias + identity caches for that set (the old sweep's + /// per-enqueue `peer_aliases` / `register_identity` side effects). pub(in crate::node) async fn run_open_discovery_sweep( &mut self, bootstrap: &std::sync::Arc, max_age_secs: Option, - caller: &'static str, ) { if !self.config().node.rendezvous.nostr.enabled || self.config().node.rendezvous.nostr.policy @@ -2179,192 +2190,82 @@ impl Node { return; } + let now_ms = Self::now_ms(); + let self_node_addr = *self.identity().node_addr(); + let configured_npubs = self .config() .peers() .iter() .map(|peer| peer.npub.clone()) .collect::>(); - let now_ms = Self::now_ms(); - let now_secs = now_ms / 1000; - let mut enqueue_budget = self.open_discovery_enqueue_budget(&configured_npubs); - if enqueue_budget == 0 { - debug!( - caller = %caller, - "open-discovery sweep: enqueue budget is 0, skipping" - ); - return; - } + // Drain the cached overlay adverts (the I/O). Exclude our own advert + // (O6): the sans-IO core has no self-identity input, so the driver + // filters self here, reproducing the old sweep's self-skip. The cooldown + // set mirrors the old per-candidate `bootstrap.cooldown_until` skip. + // + // `candidate_identities` keeps each forwarded candidate's `PeerIdentity` + // keyed by NodeAddr so that, once the core decides which candidates to + // enqueue, the driver can pre-seed the alias + identity caches for + // exactly that set (below) — reproducing the old sweep's per-enqueue + // `peer_aliases` / `register_identity` side effects byte-for-byte. let candidates = bootstrap.cached_open_discovery_candidates(64).await; - let cached_count = candidates.len(); - let mut enqueued = 0usize; - let mut skipped_age = 0usize; - let mut skipped_configured = 0usize; - let mut skipped_self = 0usize; - let mut skipped_connected = 0usize; - let mut skipped_retry_pending = 0usize; - let mut skipped_connecting = 0usize; - let mut skipped_no_endpoints = 0usize; - let mut skipped_invalid_npub = 0usize; - let mut skipped_cooldown = 0usize; - + let mut overlay = Vec::with_capacity(candidates.len()); + let mut overlay_cooldown = HashSet::new(); + let mut candidate_identities: HashMap = HashMap::new(); for (npub, endpoints, created_at_secs) in candidates { - if enqueue_budget == 0 { - break; - } - - if let Some(max_age) = max_age_secs - && now_secs.saturating_sub(created_at_secs) > max_age - { - skipped_age = skipped_age.saturating_add(1); - continue; - } - - if configured_npubs.contains(&npub) { - if let Ok(peer_identity) = PeerIdentity::from_npub(&npub) { - let node_addr = *peer_identity.node_addr(); - if !self.peers.contains_key(&node_addr) - && !self.is_connecting_to_peer(&node_addr) - && let Some(state) = - self.peering.reconciler.retry_pending.get_mut(&node_addr) - && state.retry_after_ms > now_ms - { - state.retry_after_ms = now_ms; - debug!( - peer = %peer_identity.short_npub(), - caller = %caller, - "open-discovery sweep: fresh configured-peer advert expedited retry" - ); - } - } - skipped_configured = skipped_configured.saturating_add(1); - continue; - } - - let peer_identity = match PeerIdentity::from_npub(&npub) { - Ok(identity) => identity, - Err(_) => { - skipped_invalid_npub = skipped_invalid_npub.saturating_add(1); + if let Ok(identity) = PeerIdentity::from_npub(&npub) { + let node_addr = *identity.node_addr(); + if node_addr == self_node_addr { continue; } - }; - let node_addr = *peer_identity.node_addr(); - if node_addr == *self.identity().node_addr() { - skipped_self = skipped_self.saturating_add(1); - continue; - } - if self.peers.contains_key(&node_addr) { - skipped_connected = skipped_connected.saturating_add(1); - continue; - } - if self - .peering - .reconciler - .retry_pending - .contains_key(&node_addr) - { - skipped_retry_pending = skipped_retry_pending.saturating_add(1); - continue; + candidate_identities.insert(node_addr, identity); } if bootstrap.cooldown_until(&npub, now_ms).is_some() { - skipped_cooldown = skipped_cooldown.saturating_add(1); - continue; + overlay_cooldown.insert(npub.clone()); } - let connecting = self.connections.values().any(|conn| { - conn.expected_identity() - .map(|id| id.node_addr() == &node_addr) - .unwrap_or(false) - }); - if connecting { - skipped_connecting = skipped_connecting.saturating_add(1); - continue; - } - - let mut addresses = Vec::new(); - let mut priority = 120u8; - let seen_at_ms = Self::now_ms(); - for endpoint in endpoints { - let Some(candidate) = - crate::nostr::RendezvousDriver::overlay_endpoint_to_peer_address( - &endpoint, priority, seen_at_ms, - ) - else { - continue; - }; - if addresses.iter().any(|existing: &PeerAddress| { - existing.transport == candidate.transport && existing.addr == candidate.addr - }) { - continue; - } - addresses.push(candidate); - priority = priority.saturating_add(1); - } - if addresses.is_empty() { - skipped_no_endpoints = skipped_no_endpoints.saturating_add(1); - continue; - } - - self.peer_aliases - .entry(node_addr) - .or_insert_with(|| peer_identity.short_npub()); - self.register_identity(node_addr, peer_identity.pubkey_full()); - - let mut state = super::peering::retry::RetryState::new(PeerConfig { - npub: npub.clone(), - alias: None, - addresses, - connect_policy: ConnectPolicy::AutoConnect, - auto_reconnect: true, - via_nostr: false, - }); - state.reconnect = false; - state.retry_after_ms = now_ms; - state.expires_at_ms = Some(self.open_discovery_retry_expires_at_ms(now_ms)); - self.peering - .reconciler - .retry_pending - .insert(node_addr, state); - info!( - caller = %caller, - peer = %peer_identity.short_npub(), - advert_age_secs = now_secs.saturating_sub(created_at_secs), - "open-discovery sweep: queued retry for cached advert" - ); - enqueue_budget = enqueue_budget.saturating_sub(1); - enqueued = enqueued.saturating_add(1); + overlay.push((npub, endpoints, created_at_secs)); } - // Always log a one-line summary on the startup sweep so operators - // can verify it ran. Per-tick sweeps are noisier; only summarize - // when something happened. - let total_skipped = skipped_age - + skipped_configured - + skipped_self - + skipped_connected - + skipped_retry_pending - + skipped_connecting - + skipped_no_endpoints - + skipped_invalid_npub - + skipped_cooldown; - let should_summarize = caller == "startup" || enqueued > 0; - if should_summarize { - info!( - caller = %caller, - cached = cached_count, - queued = enqueued, - skipped_age = skipped_age, - skipped_configured = skipped_configured, - skipped_self = skipped_self, - skipped_connected = skipped_connected, - skipped_retry_pending = skipped_retry_pending, - skipped_connecting = skipped_connecting, - skipped_no_endpoints = skipped_no_endpoints, - skipped_invalid_npub = skipped_invalid_npub, - skipped_cooldown = skipped_cooldown, - skipped_total = total_skipped, - "open-discovery sweep complete" - ); + let pools = DiscoveryPools { + overlay, + configured_npubs, + overlay_cooldown, + startup_sweep_max_age_secs: max_age_secs, + ..DiscoveryPools::default() + }; + + let policy = self.build_peering_policy(Vec::new()); + let observed = self.observe_peering(); + let budget = self.build_peering_budget(); + let gate = Gate::from_state(self.supervisor.state); + + // Two-phase enqueue (design §2.5): reconcile_overlay inserts due-now + // entries into retry_pending; the retry slot dials them. The core emits + // a `ScheduleRetry` for exactly the NEW enqueues (the configured-advert + // expedite bumps `retry_after_ms` without emitting), which is precisely + // the set the old sweep pre-seeded — its `register_identity` / + // `peer_aliases` sat after every skip `continue`, so only enqueued + // candidates reached it. Consuming the core's own output here to drive + // that I/O bookkeeping is legitimate sans-IO: the core decides WHICH to + // enqueue, the driver performs the side effects for exactly that set. + let actions = self + .peering + .reconciler + .reconcile_overlay(&policy, &observed, &budget, &pools, now_ms, gate); + + for action in actions { + let PeeringAction::ScheduleRetry { peer, .. } = action else { + continue; + }; + let Some(identity) = candidate_identities.get(&peer).copied() else { + continue; + }; + self.peer_aliases + .entry(peer) + .or_insert_with(|| identity.short_npub()); + self.register_identity(peer, identity.pubkey_full()); } } @@ -2411,7 +2312,7 @@ impl Node { .rendezvous .nostr .startup_sweep_max_age_secs; - self.run_open_discovery_sweep(bootstrap, Some(max_age_secs), "startup") + self.run_open_discovery_sweep(bootstrap, Some(max_age_secs)) .await; self.supervisor.nostr_rendezvous.set_startup_sweep_done(); } @@ -2489,26 +2390,6 @@ impl Node { } } - fn available_outbound_slots(&self) -> usize { - let connection_used = self - .connections - .len() - .saturating_add(self.peering.pending_connects.len()); - let connection_slots = if self.max_connections() == 0 { - usize::MAX - } else { - self.max_connections().saturating_sub(connection_used) - }; - - let peer_slots = if self.max_peers() == 0 { - usize::MAX - } else { - self.max_peers().saturating_sub(self.peers.len()) - }; - - connection_slots.min(peer_slots) - } - fn outbound_handshake_slots(&self) -> usize { let used = self .connections @@ -2565,38 +2446,6 @@ impl Node { .min(MAX_DISCOVERY_CONNECTS_PER_TICK) } - fn open_discovery_enqueue_budget(&self, configured_npubs: &HashSet) -> usize { - let current_open_discovery_pending = self - .peering - .reconciler - .retry_pending - .values() - .filter(|state| !configured_npubs.contains(&state.peer_config.npub)) - .count(); - - let cap_remaining = self - .config() - .node - .rendezvous - .nostr - .open_discovery_max_pending - .saturating_sub(current_open_discovery_pending); - - cap_remaining.min(self.available_outbound_slots()) - } - - fn open_discovery_retry_expires_at_ms(&self, now_ms: u64) -> u64 { - now_ms.saturating_add( - self.config() - .node - .rendezvous - .nostr - .advert_ttl_secs - .saturating_mul(1000) - .saturating_mul(OPEN_DISCOVERY_RETRY_LIFETIME_MULTIPLIER), - ) - } - /// Capture the advertisable-endpoint inputs of every operational /// transport into a snapshot the rendezvous driver can turn into an /// `OverlayAdvert` without borrowing the transport table across the diff --git a/src/node/peering/reconcile.rs b/src/node/peering/reconcile.rs index 1037f98..284bd7c 100644 --- a/src/node/peering/reconcile.rs +++ b/src/node/peering/reconcile.rs @@ -339,6 +339,39 @@ impl PeeringReconciler { actions } + /// Per-layer wrapper: run **only** the overlay-enqueue layer (design §10). + /// + /// The monolithic [`reconcile`] runs the always-on retry-dial phase on every + /// call, so the driver must NOT call it at the overlay (nostr-poll) cadence + /// slot — that would re-fire the retry-dial there, dialing the due entries a + /// second time in the tick and applying the per-tick 16-cap more than once + /// (today the cap applies exactly once, at the retry slot). This wrapper is + /// gate-checked and then calls [`Self::layer_overlay_enqueue`] only. + /// + /// On `NotRunning` / `Suspended` it returns no actions. It does **not** + /// clear `retry_pending` on `Suspended` — that clear is owned by the drain + /// gate (`enter_drain`) and the retry-slot [`reconcile`], not the overlay + /// slot (design §10). The returned [`PeeringAction::ScheduleRetry`] items are + /// observability only; the durable mutation is the insert the layer performed + /// into `retry_pending`, dialed at the later retry slot (two-phase, §2.5). + pub(in crate::node) fn reconcile_overlay( + &mut self, + policy: &Policy, + observed: &Observed, + budget: &Budget, + pools: &DiscoveryPools, + now: u64, + gate: Gate, + ) -> Vec { + match gate { + Gate::NotRunning | Gate::Suspended => return Vec::new(), + Gate::Reconciling => {} + } + let mut actions = Vec::new(); + self.layer_overlay_enqueue(policy, observed, budget, pools, now, &mut actions); + actions + } + /// Layer 1a — config bring-up floor (subsumes `initiate_peer_connections`). /// /// Emits a `Connect` for every auto-connect peer not already connected or diff --git a/src/node/tests/discovery.rs b/src/node/tests/discovery.rs index a809d17..fb14493 100644 --- a/src/node/tests/discovery.rs +++ b/src/node/tests/discovery.rs @@ -1020,9 +1020,14 @@ async fn test_open_discovery_sweep_queues_eligible_skips_filtered() { bootstrap.insert_advert_for_test(npub.clone(), advert).await; } + // The sweep now runs through the gate-checked reconciler overlay layer, + // which is inert unless the node is Running/Degraded. In production the + // sweep fires only from the rx_loop tick (which spins after `start()` + // returns Running), so drive the node into `Running` to reflect that. + node.supervisor.state = crate::node::NodeState::Running; + // Run the sweep. - node.run_open_discovery_sweep(&bootstrap, Some(3_600), "test") - .await; + node.run_open_discovery_sweep(&bootstrap, Some(3_600)).await; // Eligible peer was queued. assert!(