fix(sync): preserve discovery relay retry ownership

Historical mailbox and profile discovery owned connections without owning
persistent subscription targets. Idle cleanup could therefore retire a
failed discovery connection, erase its health history and let the next
probe bypass the intended backoff.

Include discovery sources, mailbox scopes and in-flight fetches in
connection ownership for cleanup and reconnect scheduling. Preserve failed
or policy-limited sessions until recovery or scope removal, including
when a disconnect notification has not yet recorded the failure.

This changes background synchronization only. Discovery ownership does
not add live subscriptions, and healthy idle sessions can still retire.
Existing backoff durations and persistent repository targets are unchanged.

Validation: unit coverage checks ownership across deferred/in-flight work
and scope removal. A real-socket mailbox-only reconnect regression checks
that cleanup retains the failure streak through reconnection; it failed
before the ownership correction.
This commit is contained in:
DanConwayDev
2026-09-14 10:06:15 +00:00
parent f55ce4885b
commit 5d62037d94
4 changed files with 122 additions and 2 deletions
+3
View File
@@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed
- Preserve background discovery relay backoff across idle-connection cleanup
so unavailable mailbox and profile sources do not restart their retry history.
- Render metrics on a blocking worker and serialize scrapes so repository
counting does not occupy async workers serving relay and Git requests.
+8
View File
@@ -744,6 +744,14 @@ is deliberately historical-only: owning or maintaining a repository never adds
that author's inbox relays to ordinary persistent live repository targets.
Private instances omit repository-coordinate mailbox expansion.
Discovery source and mailbox scope also own connection retry state. The empty
relay checker and reconnect scheduler include that scope, including in-flight
fetches and future probe deadlines, without adding persistent subscriptions.
Healthy idle discovery sessions can retire after their fetch; failed or
policy-limited sessions retain their backoff until recovery or scope removal.
This prevents cleanup from erasing a failure every two seconds and allowing
discovery to dial the same unavailable endpoint from a fresh retry history.
### Rejected Events Index
The rejected events index solves two critical problems during sync:
+70 -2
View File
@@ -1960,6 +1960,20 @@ struct Nip65DiscoveryState {
}
impl Nip65DiscoveryState {
/// Connection ownership includes historical discovery work, even while
/// its next query is deferred. It is not persistent subscription scope.
fn connection_targets(&self) -> HashSet<String> {
self.author_sources
.values()
.flatten()
.chain(self.mailbox_roots.keys())
.chain(self.mailbox_repositories.keys())
.chain(self.mailbox_probes_in_flight.iter())
.chain(self.in_flight.iter().map(|(relay, _)| relay))
.cloned()
.collect()
}
fn has_mailbox_scope(&self, relay: &str) -> bool {
self.mailbox_roots.contains_key(relay) || self.mailbox_repositories.contains_key(relay)
}
@@ -6422,6 +6436,22 @@ impl SyncManager {
if !self.nip65_discovery_only_relays.contains(source) {
return;
}
// Retirement forgets connection health. A failed or policy-limited
// source must retain its session state until recovery or removal from
// discovery ownership, otherwise the next probe bypasses its backoff.
if self.health_tracker.get_failure_count(source) > 0
|| self.health_tracker.is_subscription_paused(source)
{
return;
}
if let Some(connection) = self.connections.get(source) {
if !connection.is_connected().await {
// The disconnect notification may still be queued behind a
// fetch result. Let it record the unexpected failure instead
// of relabeling the session as an intentional retirement.
return;
}
}
let now = Instant::now();
let has_author_work =
self.nip65_discovery
@@ -7221,7 +7251,12 @@ impl SyncManager {
// Once the connection ends naturally, however, reconcile its
// confirmed state with the latest index before deciding whether it
// should ever reconnect.
let desired = self.derive_targets().await.remove(relay_url);
let desired = self.derive_targets().await.remove(relay_url).or_else(|| {
self.nip65_discovery
.connection_targets()
.contains(relay_url)
.then(RelaySyncNeeds::default)
});
let Some(desired) = desired else {
tracing::info!(
relay = %relay_url,
@@ -8420,6 +8455,7 @@ impl SyncManager {
let mut desired_relays: HashSet<String> = self.derive_targets().await.into_keys().collect();
desired_relays.extend(self.dependency_relay_deadlines.keys().cloned());
desired_relays.extend(self.nip65_discovery.connection_targets());
// Collect relays to disconnect
let to_disconnect: Vec<String> = {
@@ -8531,7 +8567,8 @@ impl SyncManager {
///
/// For each eligible relay, a reconnection is queued via schedule_connect_relay.
async fn retry_disconnected_relays(&mut self) {
let desired_relays: HashSet<String> = self.derive_targets().await.into_keys().collect();
let mut desired_relays: HashSet<String> = self.derive_targets().await.into_keys().collect();
desired_relays.extend(self.nip65_discovery.connection_targets());
// Collect relays to reconnect
let to_reconnect: Vec<String> = {
@@ -10971,6 +11008,37 @@ mod tests {
);
}
#[test]
fn discovery_connection_ownership_survives_retry_deadlines_and_in_flight_removal() {
let author = Keys::generate().public_key();
let profile = "wss://profile.example".to_string();
let mailbox = "wss://mailbox.example".to_string();
let retired = "wss://retired.example".to_string();
let mut discovery = Nip65DiscoveryState::default();
discovery
.author_sources
.insert(author, HashSet::from([profile.clone()]));
discovery
.mailbox_repositories
.insert(mailbox.clone(), HashSet::from(["repo".into()]));
discovery
.mailbox_probe_next_at
.insert(mailbox.clone(), Instant::now() + Duration::from_secs(3600));
discovery.in_flight.insert((retired.clone(), author));
let disconnected = RelayState::default();
for relay in [&profile, &mailbox, &retired] {
assert!(!disconnected
.is_disconnect_candidate(false, discovery.connection_targets().contains(relay)));
}
discovery.in_flight.clear();
assert!(!discovery.connection_targets().contains(&retired));
discovery.author_sources.clear();
discovery.install_mailbox_overlay(HashMap::new(), HashMap::new(), Instant::now());
assert!(discovery.connection_targets().is_empty());
assert!(disconnected.is_disconnect_candidate(false, false));
}
#[test]
fn disconnected_empty_relay_can_still_be_cleaned_up() {
let mut source = RelayState::default();
+41
View File
@@ -7,6 +7,47 @@ use nostr_sdk::prelude::*;
use crate::common::flapping_relay::FlappingRelay;
use crate::common::{TestClient, TestRelay};
#[tokio::test]
async fn discovery_only_mailbox_keeps_failure_history_through_cleanup() {
use crate::common::{send_to_relay_url, setup_announcement_on_relay, MockRelay};
let index = MockRelay::start().await;
let flapping = FlappingRelay::start().await;
let owner = Keys::generate();
let relay_list = EventBuilder::new(Kind::RelayList, "")
.tags([Tag::custom("r", vec![flapping.url(), "read"])])
.finalize(&owner)
.unwrap();
send_to_relay_url(index.url(), &relay_list).await.unwrap();
let syncing = TestRelay::start_with_sync(Some(index.url().to_string())).await;
let domain = syncing.domain();
let (_announcement, _git) =
setup_announcement_on_relay(&syncing, &owner, &[&domain], "discovery-mailbox-backoff")
.await;
// Owner inbox discovery has no ordinary repository/root live target.
// The two-second cleanup pass previously erased its first failure before
// the next handshake, so it could never reach this recovery state.
flapping
.wait_for_connections(2, Duration::from_secs(60))
.await;
let logs = wait_for_log(
&syncing.log_path(),
"consecutive_failures=1",
Duration::from_secs(15),
)
.await;
assert!(logs.lines().any(|line| {
line.contains(flapping.url())
&& line.contains("consecutive_failures=1")
&& line.contains("preserving failure streak until stable")
}));
syncing.stop().await;
flapping.stop().await;
index.stop().await;
}
async fn wait_for_log(log_path: &std::path::Path, needle: &str, timeout: Duration) -> String {
tokio::time::timeout(timeout, async {
loop {