diff --git a/docs/explanation/grasp-03-proactive-sync-plus.md b/docs/explanation/grasp-03-proactive-sync-plus.md index cec5ea2..692e57a 100644 --- a/docs/explanation/grasp-03-proactive-sync-plus.md +++ b/docs/explanation/grasp-03-proactive-sync-plus.md @@ -33,7 +33,14 @@ only one discovery batch may be in flight globally. Admission samples immediate transient capacity; if historic work wins the small race before the permit is acquired, that single batch may wait but discovery can never build a waiter queue. Its SDK-owned REQ draws from the same pacer and subscription ledger as -historic work. Authors returned by a successful query refresh after 24 hours; +historic work. Discovery sources use the managed connection safety and +reconnection machinery, but are a distinct control-plane role: connecting to a +user index or outbox does not start ordinary announcement, repository, or +descendant sync against it. At most one new discovery source is dialled per +maintenance pass, and an exclusively discovery connection retires after its +currently due author batches drain. A relay that independently becomes a +repository source is promoted to the ordinary lifecycle without opening a +duplicate connection. Authors returned by a successful query refresh after 24 hours; missing authors and failed queries retry after five minutes. Configured user index relays discover authors with no accepted local relay list, requesting the profile alongside it. Accepted NIP-65 diff --git a/src/sync/mod.rs b/src/sync/mod.rs index fdcf459..913598d 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1858,6 +1858,12 @@ pub struct SyncManager { rejected_events_index: Arc, /// Active relay connections - keyed by relay URL connections: HashMap, + /// Connections opened only to discover accepted authors' kind 0/10002. + /// + /// These share the ordinary connection, safety, and reconnect machinery, + /// but must not turn a user-index or outbox relay into a repository sync + /// source merely because it answered a control-plane lookup. + nip65_discovery_only_relays: HashSet, /// Adaptive pagination learning for each relay's current connection session. pagination_sessions: HashMap, /// Event-directed relay targets rejected by the outbound target policy. @@ -1998,6 +2004,7 @@ impl SyncManager { pending_sync_index: Arc::new(RwLock::new(HashMap::new())), rejected_events_index, connections: HashMap::new(), + nip65_discovery_only_relays: HashSet::new(), pagination_sessions: HashMap::new(), rejected_relay_targets: HashSet::new(), dependency_refetch_attempts: Arc::new(std::sync::Mutex::new(HashMap::new())), @@ -3431,7 +3438,10 @@ impl SyncManager { if let Some(ref bootstrap_url) = self.bootstrap_relay_url.clone() { match canonical_relay_key(bootstrap_url) { Ok(relay_url) => { - if self.register_relay(relay_url.clone(), true).await { + if self + .register_relay(relay_url.clone(), true, false) + .await + { self.schedule_connect_relay(&relay_url).await; } } @@ -3613,6 +3623,19 @@ impl SyncManager { return; } + // A relay first reached through NIP-65 may later become an ordinary + // repository target. Promote that existing connection before reading + // its state so subsequent reconnects and sync work use the full role. + let promoted_from_discovery = self + .nip65_discovery_only_relays + .remove(&action.relay_url); + if promoted_from_discovery { + tracing::info!( + relay = %action.relay_url, + "Promoting NIP-65 discovery connection to repository sync" + ); + } + // Step 1: Check if relay exists in relay_sync_index let connection_status = { let index = self.relay_sync_index.read().await; @@ -3630,7 +3653,10 @@ impl SyncManager { ); // Register relay (creates RelayConnection, initializes RelayState, updates metrics) - if self.register_relay(action.relay_url.clone(), false).await { + if self + .register_relay(action.relay_url.clone(), false, false) + .await + { self.schedule_connect_relay(&action.relay_url).await; } // Connection will trigger handle_connect_or_reconnect which will process items @@ -4035,6 +4061,14 @@ impl SyncManager { "Event loop and processor spawned for connected relay" ); + if self.nip65_discovery_only_relays.contains(relay_url) { + tracing::info!( + relay = %relay_url, + "NIP-65 discovery connection ready without repository sync" + ); + return; + } + // 3. Decide reconnection strategy based on OLD last_connected time // Use the value captured BEFORE the update to correctly detect first connections if let Some(last) = old_last_connected { @@ -4561,7 +4595,12 @@ impl SyncManager { /// registered connection (e.g. the bootstrap relay itself) reuses that /// connection rather than being re-authorized, which keeps the configured /// exception exact: only URLs that canonicalize to the same key share it. - async fn register_relay(&mut self, relay_url: String, is_bootstrap: bool) -> bool { + async fn register_relay( + &mut self, + relay_url: String, + is_bootstrap: bool, + nip65_discovery_only: bool, + ) -> bool { let relay_url = match canonical_relay_key(&relay_url) { Ok(relay_url) => relay_url, Err(error) => { @@ -4574,6 +4613,13 @@ impl SyncManager { } }; + // An ordinary sync registration upgrades a connection that was first + // opened only for NIP-65 discovery. Discovery must never downgrade an + // existing repository source. + if !nip65_discovery_only { + self.nip65_discovery_only_relays.remove(&relay_url); + } + // Create RelayConnection if not exists if !self.connections.contains_key(&relay_url) { let policy = self.outbound_target_policy(); @@ -4615,6 +4661,9 @@ impl SyncManager { policy, ); self.connections.insert(relay_url.clone(), connection); + if nip65_discovery_only { + self.nip65_discovery_only_relays.insert(relay_url.clone()); + } tracing::debug!(relay = %relay_url, "Registered new relay connection"); } @@ -4916,21 +4965,6 @@ impl SyncManager { return; } - let discovery_relays: HashSet = self - .nip65_discovery - .author_sources - .values() - .flatten() - .cloned() - .collect(); - for relay in discovery_relays { - if !self.connections.contains_key(&relay) - && self.register_relay(relay.clone(), false).await - { - self.schedule_connect_relay(&relay).await; - } - } - let mut candidates: HashMap> = HashMap::new(); for (author, sources) in &self.nip65_discovery.author_sources { for source in sources { @@ -4986,6 +5020,23 @@ impl SyncManager { }); break; } + + // No connected source could accept the single discovery query. Dial + // at most one new source per maintenance pass; registering the entire + // NIP-65 graph at once would turn a bounded query lane into an + // unbounded connection fan-out. + if self.nip65_discovery.in_flight.is_empty() { + if let Some(source) = candidates + .keys() + .filter(|source| !self.connections.contains_key(*source)) + .min() + .cloned() + { + if self.register_relay(source.clone(), false, true).await { + self.schedule_connect_relay(&source).await; + } + } + } } async fn install_nip65_overlay( @@ -5111,6 +5162,9 @@ impl SyncManager { } } if !changed { + self.retire_idle_nip65_discovery_source(&result.source_relay) + .await; + self.schedule_nip65_discovery().await; return; } let overlay = discovery::build_inbox_root_overlay( @@ -5124,6 +5178,42 @@ impl SyncManager { inbox_relays = self.nip65_discovery.inbox_roots.len(), "Updated proactive inbox coverage from NIP-65" ); + + self.retire_idle_nip65_discovery_source(&result.source_relay) + .await; + // Continue the single-flight round immediately. The health timer is a + // safety net, but ordinary sync work may hold the actor long enough + // that waiting for its next tick needlessly delays a newly discovered + // outbox. + self.schedule_nip65_discovery().await; + } + + async fn retire_idle_nip65_discovery_source(&mut self, source: &str) { + if !self.nip65_discovery_only_relays.contains(source) { + return; + } + let now = Instant::now(); + let has_due_author = self.nip65_discovery.author_sources.iter().any( + |(author, sources)| { + sources.contains(source) + && !self + .nip65_discovery + .in_flight + .contains(&(source.to_string(), *author)) + && self + .nip65_discovery + .next_attempt_at + .get(&(source.to_string(), *author)) + .is_none_or(|due| *due <= now) + }, + ); + if !has_due_author { + tracing::debug!( + relay = %source, + "Retiring idle NIP-65 discovery connection" + ); + self.disconnect_relay(source).await; + } } async fn handle_connect_attempt_result(&mut self, result: ConnectAttemptResult) { @@ -5576,7 +5666,10 @@ impl SyncManager { Instant::now() + dependency_relay_retention(), ); if !self.connections.contains_key(&relay_url) { - if !self.register_relay(relay_url.clone(), false).await { + if !self + .register_relay(relay_url.clone(), false, false) + .await + { continue; } self.schedule_connect_relay(&relay_url).await; @@ -5911,6 +6004,7 @@ impl SyncManager { self.relay_sync_index.write().await.remove(relay_url); self.pending_sync_index.write().await.remove(relay_url); self.connections.remove(relay_url); + self.nip65_discovery_only_relays.remove(relay_url); self.missing_event_recovery .lock() .unwrap() @@ -5927,7 +6021,11 @@ impl SyncManager { tracing::info!(relay = %relay_url, "Ended relay session cleanup complete"); let still_desired = self.derive_targets().await.contains_key(relay_url); - if still_desired && self.register_relay(relay_url.to_string(), false).await { + if still_desired + && self + .register_relay(relay_url.to_string(), false, false) + .await + { tracing::info!( relay = %relay_url, "Relay remains shared by current repository announcements; starting a clean session" diff --git a/tests/sync/proactive_sync_plus.rs b/tests/sync/proactive_sync_plus.rs index fcf0bf4..b7f65a0 100644 --- a/tests/sync/proactive_sync_plus.rs +++ b/tests/sync/proactive_sync_plus.rs @@ -148,6 +148,20 @@ async fn root_author_inbox_reuses_existing_root_sync_pipeline() { .await, "replacement NIP-65 inbox should become the desired root source" ); + let sync_log = std::fs::read_to_string(syncing.log_path()).expect("read syncing relay log"); + assert!( + sync_log.lines().any(|line| { + line.contains("NIP-65 discovery connection ready without repository sync") + && line.contains(outbox.url()) + }), + "the advertised outbox connection should be identified as discovery-only" + ); + assert!( + !sync_log.lines().any(|line| { + line.contains("Starting fresh_start") && line.contains(outbox.url()) + }), + "querying an advertised outbox must not start ordinary repository sync against it" + ); let draining_reply = EventBuilder::new(Kind::TextNote, "reply on naturally draining inbox") .tag(Tag::custom("e", vec![issue.id.to_hex()])) .finalize(&root_author)