mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix(sync): retire obsolete relay ownership passively
Repository addressable announcements replace their relay lists, but the sync index previously unioned every URL ever observed. Superseded relays therefore remained live indefinitely, and fully expired purgatory ownership survived after its revival window.
Replace announcement-owned relay sets in SelfSubscriber, reconcile StateOnly purgatory ownership from each current snapshot, and treat queued actions as dirty-relay signals recomputed against current state. Ownership removal deliberately leaves healthy connections and subscriptions alone. A natural live CLOSED rebuilds coverage from the current index; a natural connection end reconciles confirmed state, reconnecting shared relays with current items and retiring wholly unowned relays.
Correctness assumes addressable kind 30617 replacement semantics and that root events add sync work without changing announcement ownership. Soft-expired purgatory entries remain until the snapshot drops them, and promoted Full entries are never downgraded. Multi-relay fan-out policy and active session compaction are excluded.
Validation: cargo test --lib passed (665 tests). The passive-retirement integration scenario passed and proves replacement keeps the old session live until the source naturally disconnects, after which it is not reconnected. The broader live_sync selection passed 4/5; test_live_sync_layer3_events failed identically standalone on the unchanged 35894d7 base, confirming an existing unrelated regression.
This commit is contained in:
@@ -15,6 +15,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
REQ, serializes transient work against that reserve, and covers overflow with
|
||||
paced five-minute history plus overlap instead of retrying the same
|
||||
impossible live set.
|
||||
- Replace relay ownership when an addressable repository announcement changes,
|
||||
rather than retaining every relay URL ever listed. Fully expired StateOnly
|
||||
purgatory targets are pruned while soft-expired and promoted repositories are
|
||||
preserved. Healthy connections and live requests are left undisturbed;
|
||||
naturally closed subscriptions rebuild from current ownership, and naturally
|
||||
ended connections omit obsolete items when deciding whether to reconnect.
|
||||
- Raise retained subscription state per connection from 1 MiB to 5 MiB. A
|
||||
production 34-filter repository-sync live set reached roughly 1.2 MiB, so
|
||||
rust-nostr's newly introduced default repeatedly closed part of persistent
|
||||
|
||||
@@ -82,9 +82,18 @@ pub struct RepoSyncNeeds {
|
||||
|
||||
**Two sources populate `RepoSyncIndex`:**
|
||||
|
||||
1. **`SelfSubscriber`** — monitors the relay's own event stream for accepted announcements (kinds 30617, 1617, 1618, 1621). Adds entries with `SyncLevel::Full`. When an announcement is promoted from purgatory to the database, the SelfSubscriber sees it and upgrades the entry to `Full`.
|
||||
1. **`SelfSubscriber`** — monitors the relay's own event stream for accepted announcements (kinds 30617, 1617, 1618, 1621). Adds entries with `SyncLevel::Full`. When an announcement is promoted from purgatory to the database, the SelfSubscriber sees it and upgrades the entry to `Full`. Because kind 30617 is addressable, a newer announcement replaces that repository's relay set; root events add work without changing relay ownership. Both the old and new relay sets are marked dirty.
|
||||
|
||||
2. **Purgatory announcement sync timer** (`run_purgatory_announcement_sync`, every 5 seconds) — iterates `purgatory.announcements_for_sync()` and ensures each purgatory announcement has a `SyncLevel::StateOnly` entry in `RepoSyncIndex`. This is the only registration path for purgatory announcements because they are not saved to the database and therefore never seen by the SelfSubscriber.
|
||||
2. **Purgatory announcement sync timer** (`run_purgatory_announcement_sync`, every 5 seconds) — reconciles `SyncLevel::StateOnly` entries against `purgatory.announcements_for_sync()`. Relay sets are replaced from the current snapshot. Soft-expired announcements remain present through their 24-hour revival window; only fully expired StateOnly entries are removed. Promoted `Full` entries are never downgraded or pruned by this path. This is the only registration path for purgatory announcements because they are not saved to the database and therefore never seen by the SelfSubscriber.
|
||||
|
||||
Ownership removal is deliberately passive. It does not close a healthy
|
||||
connection or replace a working live request. If a live request receives
|
||||
`CLOSED`, its replacement filters are derived from the current index and omit
|
||||
obsolete items. When the connection itself ends naturally, confirmed state is
|
||||
reconciled before reconnect: shared relays reconnect with current items only,
|
||||
while an unreferenced relay is retired instead of reconnected. This preserves
|
||||
continuous live coverage while allowing stale ownership to drain at ordinary
|
||||
lifecycle boundaries.
|
||||
|
||||
### RelaySyncIndex (Confirmed State + Connection)
|
||||
|
||||
@@ -711,7 +720,7 @@ fn compute_actions(
|
||||
- **`register_relay()`**: Creates RelayConnection object, stores in HashMap, returns immediately
|
||||
- **`try_connect_relay()`**: Attempts connection using `connection.connect()` with timeout
|
||||
- **`handle_connect_or_reconnect()`**: Spawns event loop, updates state, decides sync strategy (fresh_start/quick_reconnect)
|
||||
- **`handle_disconnect()`**: Updates state to Disconnected, clears pending batches, keeps RelayConnection object
|
||||
- **`handle_disconnect()`**: Reconciles reusable state with current ownership after a natural disconnect; unreferenced relays are retired, while shared relays reconnect with current items only
|
||||
- **`retry_disconnected_relays()`**: Called every 2s, retries relays that pass health tracker checks
|
||||
|
||||
### Sync Entry Points
|
||||
|
||||
@@ -633,6 +633,11 @@ impl RelayHealthTracker {
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Forget lifecycle state for a relay intentionally retired by ownership reconciliation.
|
||||
pub fn forget_relay(&self, relay_url: &str) {
|
||||
self.health.remove(relay_url);
|
||||
}
|
||||
|
||||
/// Get a clone of the health info for a relay
|
||||
pub fn get_health(&self, relay_url: &str) -> Option<RelayHealth> {
|
||||
self.health
|
||||
|
||||
@@ -394,6 +394,14 @@ impl SyncMetrics {
|
||||
self.relays_tracked_total.inc();
|
||||
}
|
||||
|
||||
/// Remove current-state metric series for a relay that is no longer tracked.
|
||||
pub fn forget_relay(&self, relay: &str) {
|
||||
let _ = self.relay_connected.remove_label_values(&[relay]);
|
||||
let _ = self.relay_status.remove_label_values(&[relay]);
|
||||
let _ = self.relay_failures.remove_label_values(&[relay]);
|
||||
self.relays_tracked_total.dec();
|
||||
}
|
||||
|
||||
/// Get current tracked relay count.
|
||||
pub fn get_tracked_count(&self) -> i64 {
|
||||
self.relays_tracked_total.get()
|
||||
@@ -627,6 +635,14 @@ mod tests {
|
||||
metrics.inc_tracked_count();
|
||||
assert_eq!(metrics.get_tracked_count(), 6);
|
||||
|
||||
metrics.record_connection_status(
|
||||
"wss://obsolete.example",
|
||||
crate::sync::ConnectionStatus::Connected,
|
||||
);
|
||||
metrics.record_health_state("wss://obsolete.example", HealthState::Healthy);
|
||||
metrics.forget_relay("wss://obsolete.example");
|
||||
assert_eq!(metrics.get_tracked_count(), 5);
|
||||
|
||||
// Test connected count
|
||||
metrics.update_connected_count(3);
|
||||
assert_eq!(metrics.get_connected_count(), 3);
|
||||
|
||||
+221
-87
@@ -368,6 +368,45 @@ impl RelayState {
|
||||
}
|
||||
}
|
||||
|
||||
fn reconcile_purgatory_relay_ownership(
|
||||
index: &mut HashMap<String, RepoSyncNeeds>,
|
||||
announcements: &[(String, HashSet<String>)],
|
||||
) -> HashSet<String> {
|
||||
let active: HashMap<&str, &HashSet<String>> = announcements
|
||||
.iter()
|
||||
.map(|(repo_id, relays)| (repo_id.as_str(), relays))
|
||||
.collect();
|
||||
let mut dirty_relays = HashSet::new();
|
||||
|
||||
index.retain(|repo_id, needs| {
|
||||
if needs.sync_level == SyncLevel::Full {
|
||||
return true;
|
||||
}
|
||||
if active.contains_key(repo_id.as_str()) {
|
||||
return true;
|
||||
}
|
||||
dirty_relays.extend(needs.relays.iter().cloned());
|
||||
false
|
||||
});
|
||||
|
||||
for (repo_id, relays) in announcements {
|
||||
let entry = index
|
||||
.entry(repo_id.clone())
|
||||
.or_insert_with(|| RepoSyncNeeds {
|
||||
relays: HashSet::new(),
|
||||
root_events: HashSet::new(),
|
||||
sync_level: SyncLevel::StateOnly,
|
||||
});
|
||||
if entry.sync_level == SyncLevel::StateOnly && entry.relays != *relays {
|
||||
dirty_relays.extend(entry.relays.iter().cloned());
|
||||
dirty_relays.extend(relays.iter().cloned());
|
||||
entry.relays = relays.clone();
|
||||
}
|
||||
}
|
||||
|
||||
dirty_relays
|
||||
}
|
||||
|
||||
/// Method used for synchronization
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum SyncMethod {
|
||||
@@ -2826,21 +2865,10 @@ impl SyncManager {
|
||||
loop {
|
||||
// Wait for an event without holding the lock
|
||||
tokio::select! {
|
||||
action = action_rx.recv() => {
|
||||
match action {
|
||||
Some(add_filters) => {
|
||||
// SelfSubscriber actions are dirty-relay signals.
|
||||
// Recompute against current pending and confirmed
|
||||
// state instead of trusting a queued full-index
|
||||
// snapshot that may already be stale.
|
||||
let mut manager = sync_manager.lock().await;
|
||||
manager
|
||||
.recompute_new_sync_filters_for_relay(&add_filters.relay_url)
|
||||
.await;
|
||||
}
|
||||
None => break,
|
||||
}
|
||||
}
|
||||
// Lifecycle completions release resources and unblock later
|
||||
// work. Drain them ahead of a continuously-ready discovery
|
||||
// queue so retirement cannot remain half-finished under load.
|
||||
biased;
|
||||
disconnect = disconnect_rx.recv() => {
|
||||
match disconnect {
|
||||
Some(notification) => {
|
||||
@@ -2896,6 +2924,21 @@ impl SyncManager {
|
||||
manager.process_deferred_consolidation(&relay_url).await;
|
||||
}
|
||||
}
|
||||
action = action_rx.recv() => {
|
||||
match action {
|
||||
Some(add_filters) => {
|
||||
// SelfSubscriber actions are dirty-relay signals.
|
||||
// Recompute against current pending and confirmed
|
||||
// state instead of trusting a queued full-index
|
||||
// snapshot that may already be stale.
|
||||
let mut manager = sync_manager.lock().await;
|
||||
manager
|
||||
.recompute_new_sync_filters_for_relay(&add_filters.relay_url)
|
||||
.await;
|
||||
}
|
||||
None => break,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -3967,31 +4010,20 @@ impl SyncManager {
|
||||
let announcements = self.purgatory.announcements_for_sync();
|
||||
let announcement_events = self.purgatory.announcement_events_for_sync();
|
||||
|
||||
if announcements.is_empty() {
|
||||
return;
|
||||
// Replace StateOnly relay ownership from the current purgatory
|
||||
// snapshot and remove entries only after full expiry. Soft-expired
|
||||
// announcements remain in the snapshot for their 24-hour revival
|
||||
// window; promoted Full entries are never downgraded or removed here.
|
||||
let dirty_relays = {
|
||||
let mut index = self.repo_sync_index.write().await;
|
||||
reconcile_purgatory_relay_ownership(&mut index, &announcements)
|
||||
};
|
||||
for relay_url in dirty_relays {
|
||||
self.recompute_new_sync_filters_for_relay(&relay_url).await;
|
||||
}
|
||||
|
||||
// Register any new entries in repo_sync_index as StateOnly.
|
||||
{
|
||||
let mut index = self.repo_sync_index.write().await;
|
||||
for (repo_id, relays) in &announcements {
|
||||
let entry = index.entry(repo_id.clone()).or_insert_with(|| {
|
||||
tracing::debug!(
|
||||
repo_id = %repo_id,
|
||||
"Registering purgatory announcement in repo_sync_index as StateOnly"
|
||||
);
|
||||
RepoSyncNeeds {
|
||||
relays: std::collections::HashSet::new(),
|
||||
root_events: std::collections::HashSet::new(),
|
||||
sync_level: SyncLevel::StateOnly,
|
||||
}
|
||||
});
|
||||
// Don't downgrade an already-Full entry
|
||||
// Add any new relay URLs
|
||||
for relay in relays {
|
||||
entry.relays.insert(relay.clone());
|
||||
}
|
||||
}
|
||||
if announcements.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
// A maintainer announcement or state may have been fetched and rejected
|
||||
@@ -4463,57 +4495,26 @@ impl SyncManager {
|
||||
self.cancel_deferred_consolidation(relay_url, "relay disconnect");
|
||||
|
||||
if was_intentional {
|
||||
// Intentional disconnect - complete cleanup by removing state
|
||||
tracing::info!(relay = %relay_url, "Event loop terminated for intentional disconnect, completing cleanup");
|
||||
|
||||
// Update metrics to Disconnected before cleanup
|
||||
if let Some(ref metrics) = self.metrics {
|
||||
metrics.record_connection_status(relay_url, ConnectionStatus::Disconnected);
|
||||
}
|
||||
|
||||
// Remove from relay_sync_index
|
||||
{
|
||||
let mut index = self.relay_sync_index.write().await;
|
||||
if index.remove(relay_url).is_some() {
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
"Removed relay from relay_sync_index"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Remove from pending_sync_index
|
||||
{
|
||||
let mut pending = self.pending_sync_index.write().await;
|
||||
if pending.remove(relay_url).is_some() {
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
"Removed relay from pending_sync_index"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Remove the connection object (won't reconnect)
|
||||
if self.connections.remove(relay_url).is_some() {
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
"Removed connection from connections map"
|
||||
);
|
||||
}
|
||||
|
||||
// Drop pending missing-event recovery along with the relay.
|
||||
self.missing_event_recovery
|
||||
.lock()
|
||||
.unwrap()
|
||||
.clear_relay(relay_url);
|
||||
|
||||
// Update metrics - decrement connected count
|
||||
if let Some(ref metrics) = self.metrics {
|
||||
metrics.dec_connected_count();
|
||||
}
|
||||
|
||||
tracing::info!(relay = %relay_url, "Intentional disconnect cleanup complete");
|
||||
self.complete_ended_session(relay_url, true).await;
|
||||
} else {
|
||||
// Ownership changes are passive: do not disturb a healthy session.
|
||||
// Once the connection ends naturally, however, reconcile its
|
||||
// confirmed state with the latest index before deciding whether it
|
||||
// should ever reconnect.
|
||||
let desired = {
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
algorithms::derive_relay_targets(&repo_index).remove(relay_url)
|
||||
};
|
||||
let Some(desired) = desired else {
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
"Naturally ended relay session is no longer owned; retiring without reconnect"
|
||||
);
|
||||
self.complete_ended_session(relay_url, true).await;
|
||||
return;
|
||||
};
|
||||
|
||||
// Unexpected disconnect - update state but keep for reconnection
|
||||
tracing::warn!(relay = %relay_url, "Unexpected relay disconnect detected");
|
||||
|
||||
@@ -4521,6 +4522,9 @@ impl SyncManager {
|
||||
{
|
||||
let mut index = self.relay_sync_index.write().await;
|
||||
if let Some(state) = index.get_mut(relay_url) {
|
||||
state.repos = desired.repos;
|
||||
state.state_only_repos = desired.state_only_repos;
|
||||
state.root_events = desired.root_events;
|
||||
state.connection_status = ConnectionStatus::Disconnected;
|
||||
state.disconnected_at = Some(Timestamp::now());
|
||||
tracing::info!(
|
||||
@@ -4573,6 +4577,48 @@ impl SyncManager {
|
||||
}
|
||||
}
|
||||
|
||||
/// Remove all state owned by an ended session that must not reconnect.
|
||||
///
|
||||
/// Connected sessions reach this after their event loop terminates. An
|
||||
/// already-disconnected session has no event loop left to notify us, so
|
||||
/// `disconnect_relay` calls it directly without decrementing the connected
|
||||
/// gauge a second time.
|
||||
async fn complete_ended_session(&mut self, relay_url: &str, was_connected: bool) {
|
||||
if let Some(ref metrics) = self.metrics {
|
||||
metrics.record_connection_status(relay_url, ConnectionStatus::Disconnected);
|
||||
}
|
||||
|
||||
self.relay_sync_index.write().await.remove(relay_url);
|
||||
self.pending_sync_index.write().await.remove(relay_url);
|
||||
self.connections.remove(relay_url);
|
||||
self.missing_event_recovery
|
||||
.lock()
|
||||
.unwrap()
|
||||
.clear_relay(relay_url);
|
||||
self.health_tracker.forget_relay(relay_url);
|
||||
|
||||
if let Some(ref metrics) = self.metrics {
|
||||
if was_connected {
|
||||
metrics.dec_connected_count();
|
||||
}
|
||||
metrics.forget_relay(relay_url);
|
||||
}
|
||||
|
||||
tracing::info!(relay = %relay_url, "Ended relay session cleanup complete");
|
||||
|
||||
let still_desired = {
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
algorithms::derive_relay_targets(&repo_index).contains_key(relay_url)
|
||||
};
|
||||
if still_desired && self.register_relay(relay_url.to_string(), false).await {
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
"Relay remains shared by current repository announcements; starting a clean session"
|
||||
);
|
||||
self.schedule_connect_relay(relay_url).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Re-process events from hot cache after their dependencies become available
|
||||
///
|
||||
/// This helper consolidates the common pattern of re-processing rejected events
|
||||
@@ -5415,6 +5461,18 @@ impl SyncManager {
|
||||
///
|
||||
/// Used by check_disconnects for cleanup of empty relays.
|
||||
async fn disconnect_relay(&mut self, relay_url: &str) {
|
||||
let prior_status = {
|
||||
let index = self.relay_sync_index.read().await;
|
||||
index.get(relay_url).map(|state| state.connection_status)
|
||||
};
|
||||
if prior_status == Some(ConnectionStatus::Disconnecting) {
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
"Relay disconnect already in progress"
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
tracing::info!(relay = %relay_url, "Initiating disconnect for empty relay");
|
||||
|
||||
// Mark relay as Disconnecting (keep state for event loop to drain)
|
||||
@@ -5444,6 +5502,13 @@ impl SyncManager {
|
||||
metrics.record_connection_status(relay_url, ConnectionStatus::Disconnecting);
|
||||
}
|
||||
|
||||
// An unexpected disconnect has already ended the event loop. There is
|
||||
// no future notification to finish an intentional retirement.
|
||||
if prior_status == Some(ConnectionStatus::Disconnected) {
|
||||
self.complete_ended_session(relay_url, false).await;
|
||||
return;
|
||||
}
|
||||
|
||||
tracing::info!(relay = %relay_url, "Disconnect initiated, waiting for event loop termination");
|
||||
}
|
||||
|
||||
@@ -6830,6 +6895,75 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn purgatory_snapshot_replaces_state_only_relays_and_prunes_full_expiry() {
|
||||
let active = "30617:owner:active".to_string();
|
||||
let expired = "30617:owner:expired".to_string();
|
||||
let promoted = "30617:owner:promoted".to_string();
|
||||
let old = "wss://old.example".to_string();
|
||||
let new = "wss://new.example".to_string();
|
||||
let expired_relay = "wss://expired.example".to_string();
|
||||
let promoted_relay = "wss://promoted.example".to_string();
|
||||
let mut index = HashMap::from([
|
||||
(
|
||||
active.clone(),
|
||||
RepoSyncNeeds {
|
||||
relays: HashSet::from([old.clone()]),
|
||||
sync_level: SyncLevel::StateOnly,
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
expired.clone(),
|
||||
RepoSyncNeeds {
|
||||
relays: HashSet::from([expired_relay.clone()]),
|
||||
sync_level: SyncLevel::StateOnly,
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
promoted.clone(),
|
||||
RepoSyncNeeds {
|
||||
relays: HashSet::from([promoted_relay.clone()]),
|
||||
sync_level: SyncLevel::Full,
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
]);
|
||||
|
||||
let dirty = reconcile_purgatory_relay_ownership(
|
||||
&mut index,
|
||||
&[(active.clone(), HashSet::from([new.clone()]))],
|
||||
);
|
||||
|
||||
assert_eq!(index[&active].relays, HashSet::from([new.clone()]));
|
||||
assert!(!index.contains_key(&expired));
|
||||
assert_eq!(index[&promoted].relays, HashSet::from([promoted_relay]));
|
||||
assert_eq!(dirty, HashSet::from([old, new, expired_relay]));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn soft_expired_purgatory_entry_remains_while_present_in_snapshot() {
|
||||
let repo = "30617:owner:revivable".to_string();
|
||||
let relay = "wss://owner.example".to_string();
|
||||
let mut index = HashMap::from([(
|
||||
repo.clone(),
|
||||
RepoSyncNeeds {
|
||||
relays: HashSet::from([relay.clone()]),
|
||||
sync_level: SyncLevel::StateOnly,
|
||||
..Default::default()
|
||||
},
|
||||
)]);
|
||||
|
||||
let dirty = reconcile_purgatory_relay_ownership(
|
||||
&mut index,
|
||||
&[(repo.clone(), HashSet::from([relay]))],
|
||||
);
|
||||
|
||||
assert!(index.contains_key(&repo));
|
||||
assert!(dirty.is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_rejected_events_index_tracks_announcements() {
|
||||
// Create a rejected events index with 2 minute hot cache, 7 day cold index
|
||||
|
||||
+91
-24
@@ -43,6 +43,12 @@ enum LoopControl {
|
||||
struct PendingUpdates {
|
||||
/// Repos discovered since last batch, keyed by repo addressable ref
|
||||
repos: HashMap<String, RepoSyncNeeds>,
|
||||
/// Latest announcement relay set for each addressable repository.
|
||||
///
|
||||
/// Root events add work without changing ownership. Announcements replace
|
||||
/// relay ownership, so keeping that distinction prevents obsolete URLs
|
||||
/// from becoming permanent members of the sync index.
|
||||
relay_replacements: HashMap<String, HashSet<String>>,
|
||||
}
|
||||
|
||||
impl PendingUpdates {
|
||||
@@ -50,23 +56,38 @@ impl PendingUpdates {
|
||||
fn new() -> Self {
|
||||
Self {
|
||||
repos: HashMap::new(),
|
||||
relay_replacements: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Add or update a repo with its relays and root events
|
||||
fn add_repo(
|
||||
&mut self,
|
||||
repo_id: String,
|
||||
relays: HashSet<String>,
|
||||
root_events: HashSet<EventId>,
|
||||
) {
|
||||
let entry = self.repos.entry(repo_id).or_insert_with(|| RepoSyncNeeds {
|
||||
relays: HashSet::new(),
|
||||
root_events: HashSet::new(),
|
||||
sync_level: SyncLevel::Full,
|
||||
});
|
||||
entry.relays.extend(relays);
|
||||
entry.root_events.extend(root_events);
|
||||
/// Record the latest addressable announcement for a repository.
|
||||
fn replace_announcement_relays(&mut self, repo_id: String, relays: HashSet<String>) {
|
||||
let entry = self
|
||||
.repos
|
||||
.entry(repo_id.clone())
|
||||
.or_insert_with(|| RepoSyncNeeds {
|
||||
relays: HashSet::new(),
|
||||
root_events: HashSet::new(),
|
||||
sync_level: SyncLevel::Full,
|
||||
});
|
||||
entry.relays = relays.clone();
|
||||
self.relay_replacements.insert(repo_id, relays);
|
||||
}
|
||||
|
||||
/// Add root-event work without changing announcement relay ownership.
|
||||
fn add_root_event(&mut self, repo_id: String, relays: HashSet<String>, event_id: EventId) {
|
||||
let entry = self
|
||||
.repos
|
||||
.entry(repo_id.clone())
|
||||
.or_insert_with(|| RepoSyncNeeds {
|
||||
relays: HashSet::new(),
|
||||
root_events: HashSet::new(),
|
||||
sync_level: SyncLevel::Full,
|
||||
});
|
||||
if !self.relay_replacements.contains_key(&repo_id) {
|
||||
entry.relays.extend(relays);
|
||||
}
|
||||
entry.root_events.insert(event_id);
|
||||
}
|
||||
|
||||
/// Check if there are any pending updates
|
||||
@@ -75,8 +96,16 @@ impl PendingUpdates {
|
||||
}
|
||||
|
||||
/// Take all pending updates, leaving empty
|
||||
fn take(&mut self) -> HashMap<String, RepoSyncNeeds> {
|
||||
std::mem::take(&mut self.repos)
|
||||
fn take(
|
||||
&mut self,
|
||||
) -> (
|
||||
HashMap<String, RepoSyncNeeds>,
|
||||
HashMap<String, HashSet<String>>,
|
||||
) {
|
||||
(
|
||||
std::mem::take(&mut self.repos),
|
||||
std::mem::take(&mut self.relay_replacements),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -194,7 +223,7 @@ impl SelfSubscriber {
|
||||
for event in announcements.iter() {
|
||||
if let Some(repo_id) = Self::extract_repo_id(event) {
|
||||
let relays = Self::extract_relay_urls(event);
|
||||
pending.add_repo(repo_id, relays, HashSet::new());
|
||||
pending.replace_announcement_relays(repo_id, relays);
|
||||
announcements_loaded += 1;
|
||||
}
|
||||
}
|
||||
@@ -273,7 +302,7 @@ impl SelfSubscriber {
|
||||
// 30617 announcements don't contribute to root_events - those are
|
||||
// the 1617/1618/1621 event IDs that get added when we receive
|
||||
// root events via handle_root_event. See mod.rs:71 for details.
|
||||
pending.add_repo(repo_id.clone(), relays.clone(), HashSet::new());
|
||||
pending.replace_announcement_relays(repo_id.clone(), relays.clone());
|
||||
tracing::info!(
|
||||
event_id = %event.id,
|
||||
repo_id = %repo_id,
|
||||
@@ -538,9 +567,7 @@ impl SelfSubscriber {
|
||||
// Also add root event to pending - this ensures batch processing runs
|
||||
// and creates Layer 3 filters for events referencing this root event.
|
||||
// CRITICAL: Include relays so derive_relay_targets knows where to send filters!
|
||||
let mut root_events = HashSet::new();
|
||||
root_events.insert(event.id);
|
||||
pending.add_repo(repo_ref.clone(), relays.clone(), root_events);
|
||||
pending.add_root_event(repo_ref.clone(), relays.clone(), event.id);
|
||||
|
||||
tracing::debug!(
|
||||
event_id = %event.id,
|
||||
@@ -562,7 +589,7 @@ impl SelfSubscriber {
|
||||
/// Updates the RepoSyncIndex with discovered repos, then marks only the
|
||||
/// relays changed by this batch for SyncManager recomputation.
|
||||
async fn process_batch(&self, pending: &mut PendingUpdates) {
|
||||
let updates = pending.take();
|
||||
let (updates, relay_replacements) = pending.take();
|
||||
|
||||
if updates.is_empty() {
|
||||
return;
|
||||
@@ -583,7 +610,7 @@ impl SelfSubscriber {
|
||||
);
|
||||
}
|
||||
|
||||
let dirty_relays = dirty_relays(&updates);
|
||||
let mut dirty_relays = dirty_relays(&updates);
|
||||
|
||||
// Update RepoSyncIndex
|
||||
let mut index = self.repo_sync_index.write().await;
|
||||
@@ -601,7 +628,12 @@ impl SelfSubscriber {
|
||||
// already exists as StateOnly (purgatory announcement) and is now being
|
||||
// promoted (git data arrived and the event was broadcast via notify_event).
|
||||
entry.sync_level = SyncLevel::Full;
|
||||
entry.relays.extend(needs.relays);
|
||||
if let Some(replacement) = relay_replacements.get(&repo_id) {
|
||||
dirty_relays.extend(entry.relays.iter().cloned());
|
||||
entry.relays = replacement.clone();
|
||||
} else {
|
||||
entry.relays.extend(needs.relays);
|
||||
}
|
||||
entry.root_events.extend(needs.root_events);
|
||||
|
||||
tracing::debug!(
|
||||
@@ -730,4 +762,39 @@ mod tests {
|
||||
"a fresh historic batch must not rebuild work for unrelated indexed relays"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn announcement_replacement_does_not_union_obsolete_relays() {
|
||||
let repo = "30617:owner:repo".to_string();
|
||||
let old = "wss://old.example".to_string();
|
||||
let new = "wss://new.example".to_string();
|
||||
let mut pending = PendingUpdates::new();
|
||||
|
||||
pending.replace_announcement_relays(repo.clone(), HashSet::from([old]));
|
||||
pending.replace_announcement_relays(repo.clone(), HashSet::from([new.clone()]));
|
||||
|
||||
let (updates, replacements) = pending.take();
|
||||
assert_eq!(updates[&repo].relays, HashSet::from([new.clone()]));
|
||||
assert_eq!(replacements[&repo], HashSet::from([new]));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn root_event_does_not_override_latest_announcement_relays() {
|
||||
let repo = "30617:owner:repo".to_string();
|
||||
let relay = "wss://owner.example".to_string();
|
||||
let root = EventId::from_byte_array([7; 32]);
|
||||
let mut pending = PendingUpdates::new();
|
||||
|
||||
pending.replace_announcement_relays(repo.clone(), HashSet::from([relay.clone()]));
|
||||
pending.add_root_event(
|
||||
repo.clone(),
|
||||
HashSet::from(["wss://stale.example".to_string()]),
|
||||
root,
|
||||
);
|
||||
|
||||
let (updates, replacements) = pending.take();
|
||||
assert_eq!(updates[&repo].relays, HashSet::from([relay.clone()]));
|
||||
assert_eq!(updates[&repo].root_events, HashSet::from([root]));
|
||||
assert_eq!(replacements[&repo], HashSet::from([relay]));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -234,6 +234,11 @@ impl MockRelay {
|
||||
|
||||
/// Stop the mock relay.
|
||||
pub async fn stop(mut self) {
|
||||
// End active WebSocket sessions as well as refusing new accepts. Tests
|
||||
// that exercise natural remote disconnects need the same lifecycle a
|
||||
// real relay shutdown produces.
|
||||
self.relay.shutdown();
|
||||
|
||||
// Send shutdown signal
|
||||
if let Some(tx) = self.shutdown_tx.take() {
|
||||
let _ = tx.send(());
|
||||
|
||||
@@ -24,6 +24,102 @@ use nostr_sdk::prelude::*;
|
||||
|
||||
use crate::common::{sync_helpers::*, MockRelay, TestRelay};
|
||||
|
||||
#[tokio::test]
|
||||
async fn replacement_announcement_retires_old_relay_only_after_natural_disconnect() {
|
||||
let old_source = MockRelay::start().await;
|
||||
let new_source = MockRelay::start().await;
|
||||
let syncing = TestRelay::start_with_sync(None).await;
|
||||
let keys = Keys::generate();
|
||||
let repo_id = "replacement-retires-old-relay";
|
||||
let initial_domains = [old_source.domain(), syncing.domain()];
|
||||
let initial_refs = initial_domains
|
||||
.iter()
|
||||
.map(String::as_str)
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let (_initial, _git_dir) =
|
||||
setup_announcement_on_relay(&syncing, &keys, &initial_refs, repo_id).await;
|
||||
|
||||
let old_label = format!(
|
||||
r#"ngit_sync_relay_connected{{relay="{}"}}"#,
|
||||
old_source.url()
|
||||
);
|
||||
let new_label = format!(
|
||||
r#"ngit_sync_relay_connected{{relay="{}"}}"#,
|
||||
new_source.url()
|
||||
);
|
||||
let wait_for_metrics = |required: String, absent: Option<String>| {
|
||||
let relay_url = syncing.url();
|
||||
async move {
|
||||
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
|
||||
loop {
|
||||
let metrics = fetch_metrics(&relay_url).await.unwrap_or_default();
|
||||
if metrics.contains(&required)
|
||||
&& absent.as_ref().is_none_or(|label| !metrics.contains(label))
|
||||
{
|
||||
return true;
|
||||
}
|
||||
if tokio::time::Instant::now() >= deadline {
|
||||
return false;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
}
|
||||
}
|
||||
};
|
||||
assert!(
|
||||
wait_for_metrics(old_label.clone(), None).await,
|
||||
"initial announcement should establish the old relay session"
|
||||
);
|
||||
|
||||
let npub = keys.public_key().to_bech32().unwrap();
|
||||
let replacement_domains = [new_source.domain(), syncing.domain()];
|
||||
let replacement = EventBuilder::new(Kind::GitRepoAnnouncement, "replacement")
|
||||
.tags(vec![
|
||||
Tag::identifier(repo_id),
|
||||
Tag::custom(
|
||||
"clone",
|
||||
replacement_domains
|
||||
.iter()
|
||||
.map(|domain| format!("http://{domain}/{npub}/{repo_id}.git"))
|
||||
.collect::<Vec<_>>(),
|
||||
),
|
||||
Tag::custom(
|
||||
"relays",
|
||||
replacement_domains
|
||||
.iter()
|
||||
.map(|domain| format!("ws://{domain}"))
|
||||
.collect::<Vec<_>>(),
|
||||
),
|
||||
])
|
||||
.custom_created_at(Timestamp::from(Timestamp::now().as_secs() + 1))
|
||||
.finalize(&keys)
|
||||
.unwrap();
|
||||
let client = TestClient::new(syncing.url(), keys.clone()).await.unwrap();
|
||||
client.send_event(&replacement).await.unwrap();
|
||||
|
||||
assert!(
|
||||
wait_for_metrics(new_label, None).await,
|
||||
"replacement should establish the new relay"
|
||||
);
|
||||
assert!(
|
||||
fetch_metrics(&syncing.url())
|
||||
.await
|
||||
.unwrap_or_default()
|
||||
.contains(&old_label),
|
||||
"ownership replacement must not tear down the healthy old relay session"
|
||||
);
|
||||
|
||||
old_source.stop().await;
|
||||
assert!(
|
||||
wait_for_metrics(String::new(), Some(old_label)).await,
|
||||
"after the old relay ends naturally it should retire instead of reconnecting"
|
||||
);
|
||||
|
||||
client.disconnect().await;
|
||||
syncing.stop().await;
|
||||
new_source.stop().await;
|
||||
}
|
||||
|
||||
/// A source relay's active-REQ cap must not silently remove one of the
|
||||
/// repository filter variants.
|
||||
///
|
||||
|
||||
Reference in New Issue
Block a user