mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
fix(sync): isolate NIP-65 discovery connections
Production comparison of the proactive Sync+ candidate showed that user-index relays entered the ordinary fresh-start lifecycle. A control-plane kind 0/10002 lookup therefore also started announcement, repository, and descendant sync, contaminating the experiment and allowing the discovered identity graph to fan out connection work. Track connections opened solely for NIP-65 discovery and skip ordinary fresh-start work for that role. Dial at most one new discovery source per maintenance pass, advance the single-flight discovery round immediately after completion, and retire a discovery-only connection after its currently due author batches drain. If repository declarations independently make the same relay a sync source, promote the existing connection rather than opening a duplicate. Correctness assumes the configured bootstrap remains an ordinary sync source even when it also serves as an identity index. Exclusively discovery-owned relays still share outbound authorization, connection scheduling, subscription pacing, and the per-connection ledger. Multi-connection sharding and changes to the accepted-author scope remain excluded. Validation: nix develop -c cargo check --lib; nix develop -c cargo test --test sync root_author_inbox_reuses_existing_root_sync_pipeline -- --nocapture (pass, 18.88s). The integration scenario now proves an advertised outbox is discovery-only and never receives fresh-start repository sync.
This commit is contained in:
@@ -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
|
||||
|
||||
+118
-20
@@ -1858,6 +1858,12 @@ pub struct SyncManager {
|
||||
rejected_events_index: Arc<RejectedEventsIndex>,
|
||||
/// Active relay connections - keyed by relay URL
|
||||
connections: HashMap<String, RelayConnection>,
|
||||
/// 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<String>,
|
||||
/// Adaptive pagination learning for each relay's current connection session.
|
||||
pagination_sessions: HashMap<String, RelayPaginationSession>,
|
||||
/// 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<String> = 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<String, Vec<PublicKey>> = 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"
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user