From 7821b107190cc116a30a4c339f935bc16a1d5197 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Tue, 16 Dec 2025 15:26:55 +0000 Subject: proactive sync prep - some helper functions written but not enabled --- docs/explanation/grasp-02-proactive-sync.md | 729 ++++++++++++++++------------ src/sync/algorithms.rs | 33 +- src/sync/mod.rs | 700 ++++++++++++++++++++++++-- src/sync/self_subscriber.rs | 8 +- 4 files changed, 1120 insertions(+), 350 deletions(-) diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 2a86126..34b7bb6 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -4,12 +4,11 @@ This document explains the proactive sync system that synchronizes repository data from external relays based on relay URLs listed in 30617 repository announcements. Key principles: -1. **Self-subscription as the only mechanism** - No database initialization at startup -2. **compute_actions as single decision point** - Determines what NEW subscriptions to create -3. **Two subscription paths on reconnect** - Catch-up (retained, with since) vs new items (via compute_actions) -4. **Blank state = fresh sync** - Empty confirmed state triggers full historical fetch -5. **Clear on disconnect, not reconnect** - PendingSyncIndex cleared at event boundary -6. **NIP-77 negentropy for historical sync** - Efficient set reconciliation, fallback to REQ if unsupported +1. **Triggers call compute_actions → sync_computed_filters** - Self-subscriber batches and connect/reconnect events trigger this flow +2. **Clear separation of live vs historic sync** - Two distinct primitives with different purposes +3. **Layer 1 on connect, Layer 2+3 via AddFilters** - L1 handled at connection time, L2+L3 flow through compute_actions +4. **Always clear PendingSyncIndex first** - Before any reconnect/consolidate operation +5. **NIP-77 negentropy for historical sync** - Efficient set reconciliation, fallback to REQ if unsupported --- @@ -90,7 +89,6 @@ impl RelayState { ### PendingSyncIndex (In-Flight Batches) ```rust - /// Method used for synchronization #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum SyncMethod { @@ -100,7 +98,6 @@ pub enum SyncMethod { Negentropy, } - /// Tracks batches of subscriptions that are in-flight, awaiting EOSE. /// Each batch has its own ID and can confirm independently. /// Key: relay URL @@ -118,7 +115,6 @@ pub struct PendingBatch { pub sync_method: SyncMethod, } - #[derive(Debug, Clone, Default)] pub struct PendingItems { pub repos: HashSet, @@ -145,152 +141,202 @@ stateDiagram-v2 --- -## Flow Scenarios +## Core Architecture: Live vs Historic Sync + +The sync system is built on two fundamental primitives that are clearly separated: + +### Sync Primitives + +| Primitive | Purpose | Filter Modifier | Tracking | +| ----------------- | ----------------------- | ---------------- | ---------------- | +| `sync_live()` | Ongoing event stream | `limit: 0` | Not tracked | +| `historic_sync()` | Catch up on past events | Optional `since` | PendingSyncIndex | + +### Why `limit: 0` for Live Sync? + +| Approach | Pros | Cons | +| ------------ | --------------------------------------- | --------------------------------- | +| `since: now` | Intuitive | Time-sensitive, clock skew issues | +| `limit: 0` | Deterministic, mirrors filter structure | Less intuitive name | + +`limit: 0` is better because: + +1. **No time dependency**: Doesn't depend on synchronized clocks +2. **Mirrors historic filters**: Same tag structure, just different limit +3. **State reconstruction**: Can rebuild from repo/event lists without timestamps + +### Layer Strategy + +| Layer | Content | When Subscribed | Managed By | +| ------- | --------------------------------------- | --------------------- | -------------------- | +| Layer 1 | 30617 Announcements, 30618 Maintainers | On connect (any type) | Connection lifecycle | +| Layer 2 | Events tagging our repos (a/A/q tags) | Via AddFilters | compute_actions | +| Layer 3 | Events tagging root events (e/E/q tags) | Via AddFilters | compute_actions | -### Scenario 1: Initial Connect +**Key insight**: Layer 1 is connection-level (handled at connect time), Layer 2+3 are item-level (flow through compute_actions → sync_computed_filters). + +--- + +## Triggers and Flow + +### What Triggers compute_actions → sync_computed_filters? + +| Trigger | When | What Happens | +| --------------------------- | -------------------------------------- | ---------------------------------------- | +| Self-subscriber batch fires | New events discovered on own relay | Update RepoSyncIndex → compute_actions | +| fresh_start() | Initial connect, long_reconnect, daily | After L1 setup → compute_actions | +| quick_reconnect() | Reconnect < 15 minutes | After L1+L2+L3 catchup → compute_actions | +| consolidate() | Filter count > threshold | After live rebuild → compute_actions | + +### The Core Flow ```mermaid flowchart TB - START[Startup] --> SS[Self-subscribe to own relay] - SS --> |no since filter| EVENTS[Receive historical events] - EVENTS --> RSI[Update RepoSyncIndex] - RSI --> DT[derive_relay_targets] - DT --> CA[compute_actions with targets and empty confirmed] - CA --> AF[AddFilters for each relay] - AF --> SPAWN{Relay connected?} - SPAWN --> |no| CONN[spawn_connection] - CONN --> HC[handle_connect_or_reconnect] - SPAWN --> |yes| SUB - - subgraph handle_connect_or_reconnect - Fresh Sync - HC --> CHECK_FRESH{is_fresh_sync?} - CHECK_FRESH --> |yes - no last_connected| L1[build_announcement_filter - no since] - L1 --> RCA[recompute_actions_for_relay] - end + TRIGGER[Trigger fires] --> CA[compute_actions] + CA --> |derives from| RSI[RepoSyncIndex] + CA --> |subtracts| RLI[RelaySyncIndex] + CA --> |subtracts| PSI[PendingSyncIndex] + CA --> |produces| AF[AddFilters actions] + AF --> SFRE[sync_computed_filters] + SFRE --> LIVE[sync_live - L2+L3] + SFRE --> HIST[historic_sync - L2+L3] + HIST --> PSI_UPDATE[Update PendingSyncIndex] + PSI_UPDATE --> |EOSE received| CONFIRM[Move to RelaySyncIndex] +``` + +--- + +## Flow Scenarios - RCA --> SUB[Subscribe Layer 2+3 filters via AddFilters] - SUB --> PB[Create PendingBatch] +### Scenario 1: Fresh Start (Initial Connect / Long Reconnect / Daily Sync) + +```mermaid +flowchart TB + START[fresh_start called] --> CLEAR_PSI[Clear PendingSyncIndex] + CLEAR_PSI --> CLEAR_RSI[Clear RelaySyncIndex] + CLEAR_RSI --> L1_LIVE[L1: sync_live - announcements] + L1_LIVE --> L1_HIST[L1: historic_sync - no since] + L1_HIST --> NEG{NIP-77 supported?} + NEG --> |yes| NEGENTROPY[negentropy sync] + NEG --> |no| REQ[REQ+EOSE] + NEGENTROPY --> CA[compute_actions] + REQ --> CA + CA --> |empty RelaySyncIndex| AF[AddFilters for ALL repos] + AF --> SFRE[sync_computed_filters] + SFRE --> L23_LIVE[L2+L3: sync_live] + SFRE --> L23_HIST[L2+L3: historic_sync] + L23_HIST --> PB[Create PendingBatch] PB --> EOSE[Wait for EOSE] - EOSE --> CONFIRM[Move items to confirmed repos/root_events] + EOSE --> CONFIRM[Move items to RelaySyncIndex] ``` **Key points:** -- No `since` filter on initial connect - get full history -- `handle_connect_or_reconnect` detects `is_fresh_sync` via `last_connected.is_none()` -- Layer 1: `build_announcement_filter(None)` - subscribed immediately without since -- Layer 2+3: handled via `recompute_actions_for_relay` → `compute_actions` with PendingBatch tracking +- Always clear PendingSyncIndex first, then RelaySyncIndex +- L1 live + L1 historic (uses negentropy if available) +- Empty RelaySyncIndex means diff produces AddFilters for everything +- L2+L3 flow through sync_computed_filters with proper pending tracking -### Scenario 2: Quick Reconnect (less than 15 minutes) +### Scenario 2: Quick Reconnect (< 15 minutes) ```mermaid flowchart TB DISC[Connection lost] --> MARK[Set disconnected_at = now] - MARK --> CLEAR_PEND[Clear PendingSyncIndex for relay] - CLEAR_PEND --> WAIT[Wait for reconnection] + MARK --> WAIT[Wait for reconnection < 15min] WAIT --> RECONN[Connection restored] - RECONN --> HC[handle_connect_or_reconnect] - - subgraph handle_connect_or_reconnect - Quick Reconnect - HC --> CHECK{is_fresh_sync?} - CHECK --> |no - last_connected exists AND less than 15min| SINCE[since = last_connected - 15min] - SINCE --> L1[build_announcement_filter - with since] - L1 --> L23[rebuild_layer2_and_layer3 - with since] - L23 --> RCA[recompute_actions_for_relay] - end - - RCA --> AF[AddFilters for new items only] - AF --> SUB[Subscribe] - SUB --> PB[Create PendingBatch] - PB --> EOSE[Wait for EOSE] - EOSE --> EXTEND[Extend confirmed state] + RECONN --> CLEAR_PSI[Clear PendingSyncIndex] + CLEAR_PSI --> L1_LIVE[L1: sync_live - announcements] + L1_LIVE --> L1_HIST[L1: historic_sync WITH since] + L1_HIST --> RECON[reconstruct_filters from RelaySyncIndex] + RECON --> L23_LIVE[L2+L3: sync_live] + RECON --> L23_HIST[L2+L3: historic_sync WITH since] + L23_HIST --> CA[compute_actions] + CA --> |check for new items| AF{New items?} + AF --> |yes| SFRE[sync_computed_filters] + AF --> |no| DONE[Done] + SFRE --> PB[Create PendingBatch] ``` **Key points:** -- PendingSyncIndex cleared on disconnect (not reconnect) -- `handle_connect_or_reconnect`: - 1. `build_announcement_filter(Some(since))` - Layer 1 with since - 2. `rebuild_layer2_and_layer3(since)` - Layer 2+3 with since - 3. `recompute_actions_for_relay` - check for new items -- since = last_connected - 15min ensures we catch events during disconnection +- Clear PendingSyncIndex first (old subscriptions are dead) +- L1 live (always on any connection) +- L1 historic WITH since (catches up missed announcements) +- L2+L3 rebuilt from RelaySyncIndex (confirmed state preserved) +- compute_actions checks for any NEW items discovered during catchup -### Scenario 3: Stale Reconnect (greater than 15 minutes) +### Scenario 3: Long Reconnect (> 15 minutes) ```mermaid flowchart TB - RECONN[Connection restored] --> HC[handle_connect_or_reconnect] + RECONN[Connection restored > 15min] --> METRIC[Record disconnect/reconnect metric] + METRIC --> FRESH[fresh_start] + FRESH --> |same as initial connect| DONE[Full sync initiated] +``` - subgraph handle_connect_or_reconnect - Stale Reconnect - HC --> CHECK{is_fresh_sync?} - CHECK --> |yes - disconnected greater than 15min| CLEAR[clear_sync_state] - CLEAR --> L1[build_announcement_filter - no since] - L1 --> RCA[recompute_actions_for_relay] - end +**Key points:** - RCA --> CA[compute_actions with empty confirmed] - CA --> AF[AddFilters for everything] - AF --> SUB[Subscribe - no since filter] - SUB --> PB[Create PendingBatch] - PB --> EOSE[Wait for EOSE] - EOSE --> CONFIRM[Populate confirmed state fresh] +- Records disconnect/reconnect as a metric +- Delegates to fresh_start() - same as initial connect +- State too stale to trust, start fresh + +### Scenario 4: Consolidation (Filter Count > Threshold) + +```mermaid +flowchart TB + CHECK[Filter count check] --> THRESHOLD{count > 70?} + THRESHOLD --> |yes| CLEAR_PSI[Clear PendingSyncIndex] + CLEAR_PSI --> UNSUB[unsubscribe_all] + UNSUB --> RECON[reconstruct_filters from RelaySyncIndex] + RECON --> L1_LIVE[L1: sync_live] + RECON --> L23_LIVE[L2+L3: sync_live] + L23_LIVE --> CA[compute_actions] + CA --> |check for new items| AF{New items?} + AF --> |yes| SFRE[sync_computed_filters] + AF --> |no| DONE[Done] + THRESHOLD --> |no| SKIP[Continue normally] ``` **Key points:** -- `should_clear_state()` returns true → triggers fresh sync -- Same path as initial connect after clearing state -- Layer 1: `build_announcement_filter(None)` - full history -- Layer 2+3: handled via empty confirmed state → compute_actions generates AddFilters for everything +- Clear PendingSyncIndex first +- NO historic sync needed - items already synced/syncing +- Only rebuilds live subscriptions from confirmed state +- compute_actions catches any new items that need syncing -### Scenario 4: Consolidation (Triggered on Filter Add) +### Scenario 5: Daily Sync (23-25h Random Timer) ```mermaid flowchart TB - AF[handle_add_filters called] --> COUNT{current + new > 70?} - COUNT --> |yes| CONSOLIDATE[consolidate] - CONSOLIDATE --> WAIT_PEND[wait_pending_complete] - WAIT_PEND --> CLOSE[unsubscribe_all] - CLOSE --> SINCE[since = now - 15min] - SINCE --> L1[build_announcement_filter - with since] - L1 --> L23[rebuild_layer2_and_layer3 - with since] - COUNT --> |no| SUB[Subscribe new filters] - SUB --> PB[Create PendingBatch] + TIMER[Daily timer fires] --> FRESH[fresh_start] + FRESH --> |NO disconnect metric| DONE[Full sync initiated] ``` **Key points:** -- Consolidation checked in `handle_add_filters` BEFORE adding new filters -- After closing all subscriptions, re-subscribe: - 1. `build_announcement_filter(Some(since))` - Layer 1 stays active with since - 2. `rebuild_layer2_and_layer3(since)` - Layer 2+3 with since -- `since = now - 15min` prevents re-fetching old events -- Keeps confirmed state, just reduces filter count +- Same as fresh_start() but WITHOUT recording disconnect/reconnect metric +- Ensures consistency, detects any drift accumulated over 24 hours -### Scenario 5: Daily Timer (23-25h Random) +### Scenario 6: Self-Subscriber Batch ```mermaid flowchart TB - DAILY[Daily timer fires] --> CLOSE[unsubscribe_all] - CLOSE --> CLEAR_PEND[Clear PendingSyncIndex for relay] - CLEAR_PEND --> CLEAR_STATE[clear_sync_state] - CLEAR_STATE --> L1[build_announcement_filter - no since] - L1 --> RCA[recompute_actions_for_relay] - RCA --> CA[compute_actions with empty confirmed] - CA --> AF[AddFilters for everything] - AF --> SUB[Subscribe - no since filter] - SUB --> PB[Create PendingBatch] - PB --> EOSE[Wait for EOSE] - EOSE --> CONFIRM[Repopulate confirmed state] + EVENTS[Events from own relay] --> QUEUE[Queue to pending batch] + QUEUE --> TIMER[Batch timer fires - 5 seconds] + TIMER --> UPDATE[Update RepoSyncIndex] + UPDATE --> CA[compute_actions] + CA --> |new repos/events discovered| AF[AddFilters] + AF --> SFRE[sync_computed_filters] + SFRE --> LIVE[sync_live - L2+L3] + SFRE --> HIST[historic_sync - L2+L3] ``` **Key points:** -- Daily timer is a full fresh sync, NOT consolidation -- Clears both PendingSyncIndex and confirmed state -- Layer 1: `build_announcement_filter(None)` - full history -- Layer 2+3: via compute_actions with empty confirmed - full history -- Detects any state drift accumulated over 24 hours +- Self-subscriber monitors own relay for 30617, 1617, 1618, 1619, 1621 +- Batches events (5 second window) +- Updates RepoSyncIndex, then compute_actions finds new work +- New items flow through sync_computed_filters --- @@ -300,8 +346,6 @@ flowchart TB Transforms the repo-centric `RepoSyncIndex` into a relay-centric view. For each relay URL mentioned in any repo's announcements, collects all the repos and root events that should be synced from that relay. -**Implementation:** [`derive_relay_targets()`](../../src/sync/algorithms.rs:61) - ```rust // Conceptual: inverts repo → relays to relay → repos fn derive_relay_targets(repo_index: &HashMap) @@ -320,104 +364,262 @@ Performs a three-way diff: `target - pending - confirmed = new` Only creates `AddFilters` actions for items not already pending or confirmed. Skips disconnected relays (they will get AddFilters on reconnect). -**Implementation:** [`compute_actions()`](../../src/sync/algorithms.rs:96) +```rust +fn compute_actions( + targets: &HashMap, + pending: &PendingSyncIndex, + confirmed: &RelaySyncIndex, +) -> Vec +``` --- -## Filter Building (Three-Layer Strategy) +## Method Specifications + +### Primitives + +#### `sync_live()` - Live Subscriptions + +```rust +/// Set up live subscription (filters with limit: 0) +/// +/// - Uses `limit: 0` to receive only new events +/// - NOT tracked in PendingSyncIndex (state reconstructable) +async fn sync_live(&self, relay_url: &str, filters: &[Filter]) +``` + +#### `historic_sync()` - Historical Sync Dispatcher -The filter strategy uses three layers to ensure comprehensive event coverage: +```rust +/// Dispatch to appropriate historic sync method based on relay capabilities +/// +/// Both paths update PendingSyncIndex to ensure consistent lifecycle tracking. +async fn historic_sync( + &mut self, + relay_url: &str, + filters: Vec, + items: PendingItems, + since: Option, +) -> Option // Returns batch_id +``` + +Dispatches to: + +- `historic_sync_negentropy()` - NIP-77 parallel sync (if supported) +- `historic_sync_legacy()` - REQ+EOSE fallback + +### Building Blocks + +#### `reconstruct_filters()` - Rebuild from Confirmed State + +```rust +/// Reconstruct filters from RelaySyncIndex (confirmed state ONLY) +/// +/// Returns raw Vec for L1+L2+L3. +/// Used by: quick_reconnect, consolidate +/// Does NOT include pending items - those flow through AddFilters path. +async fn reconstruct_filters(&self, relay_url: &str) -> Vec +``` + +#### `sync_computed_filters()` - Handle New AddFilters + +```rust +/// Process AddFilters action (from compute_actions) +/// +/// Orchestrates both live and historic sync for NEW items: +/// 1. sync_live() - set up permanent L2+L3 subscriptions +/// 2. historic_sync() - catch up on past events +/// +/// This is specifically for NEW filter discovery. +async fn sync_computed_filters( + &mut self, + action: AddFilters, + since: Option, +) -> Option +``` + +### Top-Level Entry Points + +#### `fresh_start()` - Clean Slate Sync + +```rust +/// Fresh start - clears state and does full sync +/// +/// Called by: initial connect, long_reconnect, daily_sync +/// +/// Flow: +/// 1. Clear PendingSyncIndex +/// 2. Clear RelaySyncIndex +/// 3. L1 live + L1 historic (negentropy if available) +/// 4. compute_actions → AddFilters → sync_computed_filters for L2+L3 +async fn fresh_start(&mut self, relay_url: &str) +``` + +#### `quick_reconnect()` - Short Disconnection Recovery + +```rust +/// Quick reconnect - for disconnections < 15 minutes +/// +/// Flow: +/// 1. Clear PendingSyncIndex +/// 2. L1 live + L1 historic(since) +/// 3. reconstruct_filters → L2+L3 live + L2+L3 historic(since) +/// 4. compute_actions for any new items +async fn quick_reconnect(&mut self, relay_url: &str, since: Timestamp) +``` + +#### `long_reconnect()` - Extended Disconnection Recovery + +```rust +/// Long reconnect - for disconnections > 15 minutes +/// +/// Flow: +/// 1. Record disconnect/reconnect metric +/// 2. fresh_start() +async fn long_reconnect(&mut self, relay_url: &str) +``` + +#### `daily_sync()` - Scheduled Full Refresh + +```rust +/// Daily sync - full refresh without disconnect metrics +/// +/// Flow: fresh_start() (no disconnect metric recorded) +async fn daily_sync(&mut self, relay_url: &str) +``` + +#### `consolidate()` - Filter Count Reduction + +```rust +/// Consolidate subscriptions when filter count exceeds threshold +/// +/// Flow: +/// 1. Clear PendingSyncIndex +/// 2. unsubscribe_all +/// 3. reconstruct_filters → sync_live only (L1+L2+L3) +/// 4. compute_actions for any new items +/// +/// NO historic sync - items already synced, just reducing subscriptions +async fn consolidate(&mut self, relay_url: &str) +``` + +#### `handle_new_sync_filters()` - New Filter Discovery + +```rust +/// Handle AddFilters action from compute_actions +/// +/// Flow: +/// 1. Check/spawn connection if needed +/// 2. maybe_consolidate (check filter threshold) +/// 3. sync_computed_filters +async fn handle_new_sync_filters(&mut self, action: AddFilters) +``` + +--- + +## Method Relationships Summary + +``` +fresh_start(relay_url) // Initial/long_reconnect/daily + ├──> Clear PendingSyncIndex + ├──> Clear RelaySyncIndex + ├──> L1: sync_live(announcement_filter) + ├──> L1: historic_sync(announcement_filter, None) + └──> compute_actions → AddFilters → sync_computed_filters (L2+L3) + +quick_reconnect(relay_url, since) // Disconnected < 15 min + ├──> Clear PendingSyncIndex + ├──> L1: sync_live(announcement_filter) + ├──> L1: historic_sync(announcement_filter, since) + ├──> reconstruct_filters() → L2+L3 filters + ├──> L2+L3: sync_live(filters) + ├──> L2+L3: historic_sync(filters, since) + └──> compute_actions → AddFilters → sync_computed_filters (new items only) + +long_reconnect(relay_url) // Disconnected > 15 min + ├──> Record disconnect/reconnect metric + └──> fresh_start() + +daily_sync(relay_url) // Timer fires + └──> fresh_start() // No disconnect metric + +consolidate(relay_url) // Filter count > threshold + ├──> Clear PendingSyncIndex + ├──> unsubscribe_all() + ├──> reconstruct_filters() → L1+L2+L3 filters + ├──> sync_live(filters) // Live only, NO historic + └──> compute_actions → AddFilters → sync_computed_filters (new items only) + +handle_new_sync_filters(action) // New filter discovery + ├──> Check/spawn connection + ├──> maybe_consolidate() + └──> sync_computed_filters(action, None) + +sync_computed_filters(action, since) // Process AddFilters + ├──> sync_live(action.filters) // L2+L3 live + └──> historic_sync(action.filters, since) // L2+L3 historic + ├── historic_sync_negentropy() // Parallel, updates Pending + └── historic_sync_legacy() // REQ+EOSE, updates Pending +``` + +--- + +## Filter Building (Three-Layer Strategy) ### Layer 1: Announcements - **Kinds**: 30617 (Repository Announcements), 30618 (Maintainer Lists) -- **When subscribed**: ONCE on connect, NOT rebuilt during consolidation -- **Function**: [`build_announcement_filter()`](../../src/sync/filters.rs:20) +- **When subscribed**: On connect (any type) - handled by connection lifecycle +- **Function**: `build_announcement_filter(since: Option)` - 30618 is ONLY synced from remote relays, not self-subscribed ### Layer 2: Events Tagging Our Repos - **Tags**: lowercase `a`, uppercase `A`, and `q` tags for comprehensive coverage - **Batching**: Per 100 repo refs -- **Function**: [`tagged_one_of_our_repo_event_filters()`](../../src/sync/filters.rs:43) +- **Function**: `build_repo_tag_filters(repos, since)` ### Layer 3: Events Tagging Our Root Events - **Tags**: lowercase `e`, uppercase `E`, and `q` tags for comprehensive coverage - **Batching**: Per 100 event IDs -- **Function**: [`tagged_one_of_our_root_event_filters()`](../../src/sync/filters.rs:98) +- **Function**: `build_root_event_tag_filters(root_events, since)` ### Combined Layer 2+3 -The [`build_layer2_and_layer3_filters()`](../../src/sync/filters.rs:152) function combines both layers. Used by: - -- `compute_actions` for incremental subscriptions -- `rebuild_layer2_and_layer3` during reconnection -- Consolidation rebuilds (Layer 1 remains active separately) +The `build_layer2_and_layer3_filters()` function combines both layers. Used by: -**Key insight**: Layer 1 is connection-level (subscribe once), Layer 2+3 are item-level (managed by compute_actions and PendingBatch). +- `sync_computed_filters` for new item subscriptions +- `reconstruct_filters` for rebuilding from confirmed state --- -## SyncManager Key Methods - -The [`SyncManager`](../../src/sync/mod.rs:308) orchestrates all sync operations. Key methods: - -### Connection Lifecycle - -| Method | Purpose | -| ------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------- | -| `handle_connect_or_reconnect()` | Unified handler for initial connect and reconnect. Determines fresh vs quick reconnect based on `last_connected` and 15-minute rule | -| `handle_disconnect()` | Updates RelayState to Disconnected, sets disconnected_at, clears pending batches, records failure in health tracker | -| `spawn_relay_connection()` | Creates RelayConnection, subscribes to Layer 1, spawns event loop task | - -### Sync Operations - -| Method | Purpose | -| ------------------------------- | ------------------------------------------------------------------------------------------------------------------- | -| `handle_add_filters()` | Auto-spawns connection if needed, checks consolidation threshold (>70 filters), subscribes and creates PendingBatch | -| `handle_eose()` | Processes EOSE for subscription, moves items from pending to confirmed when batch completes | -| `recompute_actions_for_relay()` | Runs derive_relay_targets → compute_actions for a specific relay to find new items | -| `rebuild_layer2_and_layer3()` | Rebuilds subscriptions from confirmed state with optional since filter | - -### Maintenance - -| Method | Purpose | -| --------------------- | -------------------------------------------------------------------------- | -| `daily_sync()` | Full fresh sync - unsubscribes all, clears state, recomputes actions | -| `consolidate()` | Reduces filter count by unsubscribing and rebuilding with combined filters | -| `check_disconnects()` | Periodic check for empty relays (no repos) to disconnect | -| `check_reconnects()` | Attempts reconnection for disconnected relays with pending work | +## NIP-77 Negentropy Sync ---- +### What is Negentropy? -## Self-Subscriber +NIP-77 defines the negentropy protocol for efficient event set comparison. Instead of requesting all events matching a filter (REQ+EOSE), negentropy allows relays to compare fingerprints of their event sets and only transfer the differences. -The [`SelfSubscriber`](../../src/sync/self_subscriber.rs:86) monitors our own relay for repository announcements and root events, updating the `RepoSyncIndex`. +### When Negentropy is Used -### Event Kinds Monitored +Negentropy sync is attempted for: -- **30617** - Repository Announcements (triggers discovery of repos listing our relay) -- **1617** - Patches (root events referencing repos) -- **1618** - Issues -- **1619** - Replies/Status -- **1621** - Pull Requests +- **fresh_start()** - Full sync without `since` +- **daily_sync()** - Periodic full refresh (via fresh_start) +- **long_reconnect()** - Via fresh_start -Note: 30618 (Maintainer Lists) is NOT self-subscribed - only synced from remote relays. +Negentropy is NOT used for: -### Batching Flow +- **quick_reconnect()** - Uses REQ with `since` (more efficient for small gaps) +- **Live subscriptions** - Always use REQ with `limit: 0` -1. **Receive events** from own relay subscription -2. **Queue to pending** - announcements get repo ID + relay URLs; root events get repo ref + event ID -3. **Timer fires** (configurable window, default 5 seconds) - does NOT reset on new events -4. **Process batch**: - - Update `RepoSyncIndex` with discovered repos and root events - - Call `derive_relay_targets()` → `compute_actions()` - - Send `AddFilters` actions to SyncManager +### Fallback Behavior -### Reconnection +If negentropy fails (relay doesn't support NIP-77, network error, etc.): -Uses `last_connected` timestamp to apply since filter on reconnect (15-minute buffer), similar to external relay reconnection logic. +1. A warning is logged (once per relay to avoid spam) +2. The sync falls back to traditional REQ+EOSE +3. No error is raised - fallback is automatic --- @@ -431,12 +633,14 @@ flowchart TB end subgraph RepoSyncIndex - What We Want - RSI[HashMap: Repo to Relays+Events] + RSI[HashMap: Repo → Relays+Events] end - subgraph Derived Target - DT[derive_relay_targets fn] - TGT[Per-relay: repos + events we should sync] + subgraph Triggers + T1[Self-subscriber batch] + T2[fresh_start after L1] + T3[quick_reconnect after catchup] + T4[consolidate after live rebuild] end subgraph compute_actions - Decision Point @@ -449,136 +653,41 @@ flowchart TB subgraph RelaySyncIndex - Confirmed State RLI[RelayState per relay] - CONN[connection_status] - REPOS[repos + root_events] - TIMES[last_connected + disconnected_at] end SS -->|subscribe| OWN OWN -->|events| SS SS -->|batch fires| RSI - RSI --> DT - DT --> TGT - TGT --> CA + RSI --> T1 + T1 --> CA + T2 --> CA + T3 --> CA + T4 --> CA PSI --> CA RLI --> CA - CA -->|Layer 2+3 new items| AF[AddFilters] - AF -->|check filter count| CONSOL{count + new > 70?} - CONSOL -->|yes| CONSOLIDATE[consolidate] - CONSOLIDATE --> L1_CONSOL[build_announcement_filter with since] - L1_CONSOL --> L23_CONSOL[rebuild_layer2_and_layer3 with since] - CONSOL -->|no| SUB[subscribe] - AF -->|spawn if needed| CONN - SUB --> PSI - PSI -->|EOSE| REPOS - - CONN -->|disconnect| DISC[Clear PSI + set disconnected_at] - DISC -->|any reconnect| HC[handle_connect_or_reconnect] - - subgraph handle_connect_or_reconnect - HC --> FRESH_CHECK{is_fresh_sync?} - FRESH_CHECK -->|yes: no last_connected OR >15min| L1_FRESH[build_announcement_filter - no since] - FRESH_CHECK -->|no: quick reconnect| L1_QUICK[build_announcement_filter - with since] - L1_FRESH --> RCA1[recompute_actions_for_relay] - L1_QUICK --> L23_QUICK[rebuild_layer2_and_layer3 - with since] - L23_QUICK --> RCA2[recompute_actions_for_relay] - end + CA -->|new items| AF[AddFilters] + AF --> SFRE[sync_computed_filters] + SFRE --> LIVE[sync_live L2+L3] + SFRE --> HIST[historic_sync L2+L3] + HIST --> PSI + PSI -->|EOSE| RLI ``` --- ## Key Design Decisions -| Decision | Choice | Rationale | -| -------------------------- | -------------------------------------- | --------------------------------------------------------------------------- | -| Startup mechanism | Self-subscription only | Single code path, fresh DB behaves same as reconnect | -| Connect/reconnect handling | Unified handle_connect_or_reconnect | Single entry point for both initial and reconnect | -| Layer 1 handling | Separate build_announcement_filter | Connection-level: subscribe ONCE on connect, NOT rebuilt in consolidation | -| Layer 2+3 handling | Separate rebuild_layer2_and_layer3 | Item-level: managed by compute_actions, consolidated when filter count > 70 | -| Filter functions | since as Option parameter | Allows same functions for fresh sync and catch-up | -| Since filter | Only on catch-up paths | Initial/stale gets full history, quick reconnect catches up | -| compute_actions role | ONLY for new Layer 2+3 items | Does NOT handle Layer 1 or catch-up | -| Catch-up pending tracking | No PendingBatch | Items already confirmed, don't need re-confirmation | -| Consolidation trigger | On filter add, not periodic | Check in handle_add_filters before adding new filters | -| Clear on disconnect | Clear PSI on disconnect | Cleanup at event boundary, simpler than on reconnect | -| 15-minute rule | Clear confirmed if disconnected >15min | Matches since filter buffer, prevents stale subscriptions | -| Daily timer | Fresh sync (clears state) | Ensures consistency, detects drift | -| NIP-77 negentropy | Try first, fallback to REQ | Efficient set reconciliation when supported | - ---- - -## NIP-77 Negentropy Sync - -The sync system supports NIP-77 negentropy for efficient set reconciliation when syncing with external relays. - -### What is Negentropy? - -NIP-77 defines the negentropy protocol for efficient event set comparison. Instead of requesting all events matching a filter (REQ+EOSE), negentropy allows relays to compare fingerprints of their event sets and only transfer the differences. - -### When Negentropy is Used - -Negentropy sync is attempted for: - -- **Initial connect** - Fresh sync without `last_connected` -- **Daily sync** - Periodic full refresh (23-25 hour timer) -- **Stale reconnect** - Disconnected for more than 15 minutes - -Negentropy is NOT used for: - -- **Quick reconnect** - Less than 15 minutes disconnected (uses REQ with `since`) -- **Live subscriptions** - Ongoing event streams always use REQ - -### Implementation - -The [`RelayConnection`](../../src/sync/relay_connection.rs:71) now includes NIP-77 methods: - -```rust -/// Check if negentropy sync should be attempted -pub async fn supports_negentropy(&self) -> bool { - // Always returns true - we try negentropy and handle failure gracefully - true -} - -/// Perform negentropy synchronization for a filter -pub async fn negentropy_sync_filter(&self, filter: Filter) - -> Result { - // Uses nostr-sdk's client.sync() method -} -``` - -### Sync Flow with Negentropy - -```mermaid -flowchart TB - CONNECT[Connect to relay] --> NEG{Try negentropy} - NEG --> |success| L1[Layer 1 synced via negentropy] - NEG --> |failure| FALLBACK[Fall back to REQ+EOSE] - - L1 --> SINCE[Record timestamp = now] - FALLBACK --> EOSE[Wait for EOSE] - EOSE --> SINCE - - SINCE --> LIVE[Open live REQ with since=now] -``` - -### Fallback Behavior - -If negentropy fails (relay doesn't support NIP-77, network error, etc.): - -1. A warning is logged (once per relay to avoid spam) -2. The sync falls back to traditional REQ+EOSE -3. No error is raised - fallback is automatic - -**Implementation:** [`negentropy_sync_and_process()`](../../src/sync/mod.rs:1549) - -### Key Design Decisions for Negentropy - -| Decision | Choice | Rationale | -| ------------------ | --------------------------- | ------------------------------------------------- | -| Detection approach | Try and fallback | More reliable than NIP-11 document detection | -| When to use | Fresh/daily/stale sync only | Quick reconnect with `since` is already efficient | -| Error handling | Log once, fallback silently | Avoid log spam while maintaining visibility | -| Layer application | Layer 1 first | Announcements are highest priority | +| Decision | Choice | Rationale | +| ----------------------------- | ------------------------------------------- | ------------------------------------------------------------------ | +| Live vs Historic separation | Two distinct primitives | Clear responsibilities, easier reasoning about state | +| Live sync method | `limit: 0` not `since: now` | No clock dependency, deterministic, mirrors filter structure | +| Layer 1 handling | On connect, separate from AddFilters | Connection-level concern, not item-level | +| Layer 2+3 handling | Via compute_actions → sync_computed_filters | Item-level, proper pending tracking | +| Clear PendingSyncIndex | Always first | Old subscriptions are dead, must clear before any operation | +| fresh_start vs long_reconnect | Same flow, different metrics | Reuse logic, distinguish intentional refresh from failure recovery | +| Consolidation | Live only, no historic | Items already synced, just reducing subscription count | +| compute_actions role | ONLY decision point for new work | Single place to reason about what needs syncing | +| NIP-77 negentropy | Try first on full sync, fallback | Efficient for large sets, graceful degradation | --- @@ -586,8 +695,8 @@ If negentropy fails (relay doesn't support NIP-77, network error, etc.): ``` src/sync/ -├── mod.rs # SyncManager, main loop, data structures (RepoSyncNeeds, RelayState, etc.) -├── algorithms.rs # derive_relay_targets(), compute_actions(), AddFilters +├── mod.rs # SyncManager, main loop, data structures +├── algorithms.rs # derive_relay_targets(), compute_actions() ├── filters.rs # build_announcement_filter(), build_layer2_and_layer3_filters() ├── health.rs # RelayHealthTracker with exponential backoff ├── relay_connection.rs # RelayConnection, RelayEvent handling @@ -599,7 +708,7 @@ src/sync/ ## Health Tracking -The [`RelayHealthTracker`](../../src/sync/health.rs:93) manages connection health with exponential backoff: +The `RelayHealthTracker` manages connection health with exponential backoff: - **States**: Healthy, Degraded, Dead - **Backoff**: `base * 2^(failures-1)`, capped at max_backoff @@ -610,6 +719,32 @@ Bootstrap relays are never disconnected by the cleanup system, even if empty. --- +## Self-Subscriber + +The `SelfSubscriber` monitors our own relay for repository announcements and root events, updating the `RepoSyncIndex`. + +### Event Kinds Monitored + +- **30617** - Repository Announcements (triggers discovery of repos listing our relay) +- **1617** - Patches (root events referencing repos) +- **1618** - Issues +- **1619** - Replies/Status +- **1621** - Pull Requests + +Note: 30618 (Maintainer Lists) is NOT self-subscribed - only synced from remote relays. + +### Batching Flow + +1. **Receive events** from own relay subscription +2. **Queue to pending** - announcements get repo ID + relay URLs; root events get repo ref + event ID +3. **Timer fires** (configurable window, default 5 seconds) - does NOT reset on new events +4. **Process batch**: + - Update `RepoSyncIndex` with discovered repos and root events + - Call `compute_actions()` + - Send `AddFilters` actions to SyncManager → `sync_computed_filters()` + +--- + ## Disconnect Handling The disconnect checker runs periodically (default: 60 seconds) to clean up empty relays: diff --git a/src/sync/algorithms.rs b/src/sync/algorithms.rs index 5b5b520..84248b1 100644 --- a/src/sync/algorithms.rs +++ b/src/sync/algorithms.rs @@ -11,7 +11,9 @@ use std::collections::{HashMap, HashSet}; use nostr_sdk::prelude::*; -use super::{ConnectionStatus, PendingBatch, RelayState, SyncMethod}; +use crate::sync::PendingItems; + +use super::{ConnectionStatus, PendingBatch, RelayState}; // ============================================================================= // Data Structures @@ -36,10 +38,8 @@ pub struct RelaySyncNeeds { pub struct AddFilters { /// The relay URL to add filters to pub relay_url: String, - /// Repos being synced in this action - pub repos: HashSet, - /// Root events being tracked in this action - pub root_events: HashSet, + /// pending items - repos and root events + pub items: PendingItems, /// The actual filters to subscribe with pub filters: Vec, } @@ -161,8 +161,10 @@ pub fn compute_actions( actions.push(AddFilters { relay_url: relay_url.clone(), - repos: new_repos, - root_events: new_events, + items: PendingItems { + repos: new_repos, + root_events: new_events, + }, filters, }); } @@ -175,6 +177,7 @@ pub fn compute_actions( mod tests { use super::*; use crate::sync::RepoSyncNeeds as ModRepoSyncNeeds; + use crate::sync::SyncMethod; // ========================================================================= // derive_relay_targets tests @@ -371,7 +374,7 @@ mod tests { assert_eq!(actions.len(), 1); let action = &actions[0]; assert_eq!(action.relay_url, "wss://relay1.com"); - assert!(action.repos.contains("repo1")); + assert!(action.items.repos.contains("repo1")); assert!(!action.filters.is_empty()); } @@ -528,10 +531,10 @@ mod tests { assert_eq!(actions.len(), 1); let action = &actions[0]; // Only repo3 should be in the action (repo1 pending, repo2 confirmed) - assert_eq!(action.repos.len(), 1); - assert!(action.repos.contains("repo3")); - assert!(!action.repos.contains("repo1")); - assert!(!action.repos.contains("repo2")); + assert_eq!(action.items.repos.len(), 1); + assert!(action.items.repos.contains("repo3")); + assert!(!action.items.repos.contains("repo1")); + assert!(!action.items.repos.contains("repo2")); } #[test] @@ -554,9 +557,9 @@ mod tests { assert_eq!(actions.len(), 1); let action = &actions[0]; - assert!(action.repos.is_empty()); - assert_eq!(action.root_events.len(), 1); - assert!(action.root_events.contains(&event_id)); + assert!(action.items.repos.is_empty()); + assert_eq!(action.items.root_events.len(), 1); + assert!(action.items.root_events.contains(&event_id)); // Should have 3 filters for the root event (e, E, q tags) assert_eq!(action.filters.len(), 3); } diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 41586a4..401cf21 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -591,7 +591,7 @@ impl SyncManager { } // Recompute actions for Layer 2+3 based on synced events - self.recompute_actions_for_relay(relay_url).await; + self.recompute_new_sync_filters_for_relay(relay_url).await; } else { // NIP-77 not supported - fall back to REQ+EOSE tracing::info!( @@ -612,7 +612,7 @@ impl SyncManager { } // Recompute actions for Layer 2+3 - will discover all repos/events again - self.recompute_actions_for_relay(relay_url).await; + self.recompute_new_sync_filters_for_relay(relay_url).await; } if let Some(ref metrics) = self.metrics { @@ -709,7 +709,7 @@ impl SyncManager { Some(add_filters) => { // Process AddFilters action directly let mut manager = sync_manager.lock().await; - manager.handle_add_filters(add_filters).await; + manager.handle_new_sync_filters(add_filters).await; } None => break, } @@ -763,13 +763,13 @@ impl SyncManager { /// - For new relays: creates entry with Connecting status, spawns connection /// - For existing connected relays: subscribes to filters, creates PendingBatch /// - For disconnected/connecting relays: returns (will be handled on connection) - async fn handle_add_filters(&mut self, action: AddFilters) { + async fn handle_new_sync_filters(&mut self, action: AddFilters) { tracing::info!( relay = %action.relay_url, - repo_count = action.repos.len(), - root_event_count = action.root_events.len(), + repo_count = action.items.repos.len(), + root_event_count = action.items.root_events.len(), filter_count = action.filters.len(), - "[DIAG] handle_add_filters called" + "[DIAG] handle_new_sync_filters called" ); // Step 1: Check if relay exists in relay_sync_index @@ -801,7 +801,7 @@ impl SyncManager { tracing::info!( relay = %action.relay_url, - repos = action.repos.len(), + repos = action.items.repos.len(), "Spawning connection for new relay" ); @@ -827,7 +827,7 @@ impl SyncManager { // Step 2: Check if consolidation is needed BEFORE adding new filters self.maybe_consolidate(&action.relay_url, action.filters.len()) .await; - + /// DELETE this bit // Step 3: Get connection and subscribe to all filters let connection = match self.connections.get(&action.relay_url) { Some(conn) => conn, @@ -870,8 +870,8 @@ impl SyncManager { let batch = PendingBatch { batch_id, items: PendingItems { - repos: action.repos.clone(), - root_events: action.root_events.clone(), + repos: action.items.repos.clone(), + root_events: action.items.root_events.clone(), }, outstanding_subs: subscription_ids.into_iter().collect(), sync_method: SyncMethod::ReqEose, @@ -889,33 +889,84 @@ impl SyncManager { tracing::debug!( relay = %action.relay_url, batch_id = batch_id, - repos = action.repos.len(), - root_events = action.root_events.len(), + repos = action.items.repos.len(), + root_events = action.items.root_events.len(), filters = action.filters.len(), "Created pending batch for filter subscriptions" ); + // REPLACE WITH THIS: + // // Subscribe to each filter and collect subscription IDs + // self.sync_live(&action.relay_url, &action.filters).await; + // // TODO need to do actions.repos + // self.historic_sync(&action.relay_url, action.filters, action.items, None) + // .await; } /// Handle a connection success (called when a relay connects or reconnects) /// - /// This method implements smart reconnection logic: - /// - Fresh sync if never connected or >15 min since last connection - /// - Quick reconnect with since filter if <15 min since last connection - /// - /// For fresh sync (with NIP-77 negentropy if supported): - /// - Clears any stale state - /// - Uses negentropy sync for Layer 1 (if NIP-77 supported) - /// - Falls back to REQ+EOSE if NIP-77 not supported - /// - Recomputes actions for new items - /// - /// For quick reconnect: - /// - Preserves existing state - /// - Subscribes to Layer 1 with since filter - /// - Rebuilds Layer 2 and Layer 3 with since filter - /// - Recomputes actions for new items + /// This method dispatches to the appropriate reconnection strategy: + /// - `fresh_start()` if never connected before + /// - `quick_reconnect()` if disconnected < 15 minutes + /// - `long_reconnect()` if disconnected > 15 minutes async fn handle_connect_or_reconnect(&mut self, relay_url: &str) { let now = Timestamp::now(); + // // Get the relay state to determine reconnect type + // let (last_connected, disconnected_at) = { + // let index = self.relay_sync_index.read().await; + // if let Some(state) = index.get(relay_url) { + // (state.last_connected, state.disconnected_at) + // } else { + // (None, None) // No state found + // } + // }; + + // // Determine which reconnection strategy to use + // match (last_connected, disconnected_at) { + // (None, _) => { + // // Never connected before - fresh start + // tracing::info!( + // relay = %relay_url, + // "First connection - initiating fresh_start" + // ); + // self.fresh_start(relay_url).await; + // } + // (Some(last), Some(disconnected)) => { + // // Was connected before, check how long disconnected + // let disconnect_duration = now.as_secs().saturating_sub(disconnected.as_secs()); + + // if disconnect_duration <= QUICK_RECONNECT_WINDOW_SECS { + // // Disconnected < 15 minutes - quick reconnect + // // Use last_connected minus buffer as since timestamp + // let since = + // Timestamp::from(last.as_secs().saturating_sub(QUICK_RECONNECT_WINDOW_SECS)); + // tracing::info!( + // relay = %relay_url, + // disconnect_secs = disconnect_duration, + // since = %since, + // "Short disconnection - initiating quick_reconnect" + // ); + // self.quick_reconnect(relay_url, since).await; + // } else { + // // Disconnected > 15 minutes - long reconnect + // tracing::info!( + // relay = %relay_url, + // disconnect_secs = disconnect_duration, + // "Long disconnection - initiating long_reconnect" + // ); + // self.long_reconnect(relay_url).await; + // } + // } + // (Some(_last), None) => { + // // Was connected but no disconnected_at - shouldn't happen normally + // // Treat as long reconnect to be safe + // tracing::warn!( + // relay = %relay_url, + // "Unexpected state: last_connected set but no disconnected_at - using long_reconnect" + // ); + // self.long_reconnect(relay_url).await; + // } + // } // Get the relay state to determine reconnect type let (is_fresh_sync, last_connected, is_bootstrap) = { let index = self.relay_sync_index.read().await; @@ -998,7 +1049,7 @@ impl SyncManager { // After negentropy sync, recompute Layer 2+3 actions // Layer 1 events are now in sync, so we can proceed with Layer 2+3 - self.recompute_actions_for_relay(relay_url).await; + self.recompute_new_sync_filters_for_relay(relay_url).await; // Set up live subscription for new events (since=now) let live_filter = filters::build_announcement_filter(Some(now)); @@ -1021,7 +1072,7 @@ impl SyncManager { // during connect_and_subscribe() in handle_add_filters(). That call subscribes // to kinds 30617+30618 for the full history. Here we only need to recompute // Layer 2+3 actions based on the repos we're tracking. - self.recompute_actions_for_relay(relay_url).await; + self.recompute_new_sync_filters_for_relay(relay_url).await; } } else { // Quick reconnect: use since filter (no negentropy needed) @@ -1055,7 +1106,7 @@ impl SyncManager { .await; // Recompute actions for any new items discovered while disconnected - self.recompute_actions_for_relay(relay_url).await; + self.recompute_new_sync_filters_for_relay(relay_url).await; if let Some(ref metrics) = self.metrics { metrics.record_event(event_source::RECONNECT); @@ -1063,6 +1114,225 @@ impl SyncManager { } } + /// Fresh start - clears state and does full sync + /// + /// Called by: initial connect, long_reconnect, daily_sync + /// + /// Flow: + /// 1. Clear PendingSyncIndex for this relay + /// 2. Clear RelaySyncIndex sync state (repos/root_events) + /// 3. Update connection state to Connected + /// 4. L1 live + L1 historic (negentropy if available) + /// 5. compute_actions → AddFilters → sync_computed_filters for L2+L3 + async fn fresh_start(&mut self, relay_url: &str) { + let now = Timestamp::now(); + + tracing::info!(relay = %relay_url, "Starting fresh_start"); + + // Step 1: Clear PendingSyncIndex for this relay + { + let mut pending = self.pending_sync_index.write().await; + if pending.remove(relay_url).is_some() { + tracing::debug!( + relay = %relay_url, + "Cleared pending batches in fresh_start" + ); + } + } + + // Step 2: Clear RelaySyncIndex sync state (but preserve connection metadata) + { + let mut index = self.relay_sync_index.write().await; + if let Some(state) = index.get_mut(relay_url) { + let repos_cleared = state.repos.len(); + let events_cleared = state.root_events.len(); + state.clear_sync_state(); + if repos_cleared > 0 || events_cleared > 0 { + tracing::debug!( + relay = %relay_url, + repos_cleared = repos_cleared, + events_cleared = events_cleared, + "Cleared sync state in fresh_start" + ); + } + } + } + + // Step 3: Update connection state + { + let mut index = self.relay_sync_index.write().await; + let state = index.entry(relay_url.to_string()).or_default(); + state.connection_status = ConnectionStatus::Connected; + state.last_connected = Some(now); + state.disconnected_at = None; + } + + // Record success in health tracker + self.health_tracker.record_success(relay_url); + + // Update metrics + if let Some(ref metrics) = self.metrics { + metrics.set_relay_connected(relay_url, true); + metrics.inc_connected_count(); + metrics.record_health_state(relay_url, self.health_tracker.get_state(relay_url)); + } + + // Step 4: L1 sync - check negentropy support + let use_negentropy = if self.config.sync_disable_negentropy { + tracing::debug!(relay = %relay_url, "Negentropy disabled via config"); + false + } else if let Some(connection) = self.connections.get(relay_url) { + connection.supports_negentropy().await + } else { + false + }; + + if use_negentropy { + // NIP-77 supported - use negentropy for L1 historical sync + tracing::info!( + relay = %relay_url, + "Using NIP-77 negentropy for L1 historical sync" + ); + + // L1 historic sync (no since - full sync) + let layer1_filter = filters::build_announcement_filter(None); + self.negentropy_sync_and_process(relay_url, layer1_filter, "Layer 1 (fresh_start)") + .await; + + // L1 live subscription (since=now for ongoing events) + let live_filter = filters::build_announcement_filter(Some(now)); + if let Some(connection) = self.connections.get(relay_url) { + if let Err(e) = connection.subscribe_filter(live_filter).await { + tracing::error!( + relay = %relay_url, + error = %e, + "Failed to set up L1 live subscription in fresh_start" + ); + } + } + } else { + // NIP-77 not supported - REQ+EOSE + // Note: Layer 1 subscription (without since) was already established + // during connect_and_subscribe() in spawn_relay_connection + tracing::info!( + relay = %relay_url, + "Using REQ+EOSE for L1 sync (negentropy not available)" + ); + } + + // Step 5: compute_actions → AddFilters for L2+L3 + // Since RelaySyncIndex is now empty, compute_actions will produce AddFilters + // for ALL repos that should be synced from this relay + self.recompute_new_sync_filters_for_relay(relay_url).await; + + tracing::info!(relay = %relay_url, "fresh_start complete"); + } + + /// Quick reconnect - for disconnections < 15 minutes + /// + /// Flow: + /// 1. Clear PendingSyncIndex for this relay + /// 2. Update connection state to Connected + /// 3. L1 live + L1 historic(since) + /// 4. reconstruct_filters → L2+L3 live + L2+L3 historic(since) + /// 5. compute_actions for any new items discovered during catchup + async fn quick_reconnect(&mut self, relay_url: &str, since: Timestamp) { + let now = Timestamp::now(); + + tracing::info!( + relay = %relay_url, + since = %since, + "Starting quick_reconnect" + ); + + // Step 1: Clear PendingSyncIndex for this relay + // Old subscriptions are dead after disconnect + { + let mut pending = self.pending_sync_index.write().await; + if pending.remove(relay_url).is_some() { + tracing::debug!( + relay = %relay_url, + "Cleared pending batches in quick_reconnect" + ); + } + } + + // Step 2: Update connection state (preserve repos/root_events - that's the point!) + { + let mut index = self.relay_sync_index.write().await; + let state = index.entry(relay_url.to_string()).or_default(); + state.connection_status = ConnectionStatus::Connected; + state.last_connected = Some(now); + state.disconnected_at = None; + } + + // Record success in health tracker + self.health_tracker.record_success(relay_url); + + // Update metrics + if let Some(ref metrics) = self.metrics { + metrics.set_relay_connected(relay_url, true); + metrics.inc_connected_count(); + metrics.record_health_state(relay_url, self.health_tracker.get_state(relay_url)); + metrics.record_event(event_source::RECONNECT); + } + + // Step 3: L1 live + L1 historic with since filter + // L1 live subscription (since=now for ongoing events) + let live_filter = filters::build_announcement_filter(Some(now)); + if let Some(connection) = self.connections.get(relay_url) { + if let Err(e) = connection.subscribe_filter(live_filter).await { + tracing::error!( + relay = %relay_url, + error = %e, + "Failed to set up L1 live subscription in quick_reconnect" + ); + } + } + + // L1 historic with since filter (catch up on missed announcements) + let layer1_filter = filters::build_announcement_filter(Some(since)); + if let Some(connection) = self.connections.get(relay_url) { + if let Err(e) = connection.subscribe_filter(layer1_filter).await { + tracing::error!( + relay = %relay_url, + error = %e, + "Failed to subscribe to L1 historic filter in quick_reconnect" + ); + } + } + + // Step 4: Rebuild L2+L3 from confirmed state with since filter + // This uses the preserved repos/root_events from RelaySyncIndex + self.rebuild_layer2_and_layer3(relay_url, Some(since)).await; + + // Step 5: compute_actions for any NEW items discovered while disconnected + self.recompute_new_sync_filters_for_relay(relay_url).await; + + tracing::info!(relay = %relay_url, "quick_reconnect complete"); + } + + /// Long reconnect - for disconnections > 15 minutes + /// + /// Flow: + /// 1. Record disconnect/reconnect metric + /// 2. Delegate to fresh_start() + async fn long_reconnect(&mut self, relay_url: &str) { + tracing::info!(relay = %relay_url, "Starting long_reconnect"); + + // Step 1: Record disconnect/reconnect metric + // This distinguishes intentional daily refresh from failure recovery + if let Some(ref metrics) = self.metrics { + metrics.record_event(event_source::RECONNECT); + } + + // Step 2: Delegate to fresh_start + // State is too stale to trust, start fresh + self.fresh_start(relay_url).await; + + tracing::info!(relay = %relay_url, "long_reconnect complete"); + } + /// Rebuild Layer 2 and Layer 3 subscriptions for a relay /// /// Uses the confirmed repos and root_events from RelayState to build filters. @@ -1129,7 +1399,7 @@ impl SyncManager { /// /// Uses derive_relay_targets and compute_actions to find new items /// that need to be synced. Processes AddFilters actions for new items. - async fn recompute_actions_for_relay(&mut self, relay_url: &str) { + async fn recompute_new_sync_filters_for_relay(&mut self, relay_url: &str) { use crate::sync::algorithms::{compute_actions, derive_relay_targets}; // Get current state from indexes (need to collect to avoid holding locks) @@ -1173,12 +1443,12 @@ impl SyncManager { for action in actions { tracing::info!( relay = %action.relay_url, - new_repos = action.repos.len(), - new_root_events = action.root_events.len(), + new_repos = action.items.repos.len(), + new_root_events = action.items.root_events.len(), filters = action.filters.len(), "Processing AddFilters for new items" ); - self.handle_add_filters(action).await; + self.handle_new_sync_filters(action).await; } } @@ -2095,6 +2365,366 @@ impl SyncManager { } } + // ========================================================================= + // Sync Primitives (Phase 1 of GRASP-02 refactoring) + // These methods are new primitives that will be used in subsequent phases. + // ========================================================================= + + /// Subscribe to filters for live (ongoing) events - NOT tracked in PendingSyncIndex + #[allow(dead_code)] // Will be used in Phase 2+ + /// + /// This method subscribes to filters with `limit: 0` for receiving ongoing events. + /// Live subscriptions are NOT tracked in PendingSyncIndex because they don't have + /// a definite "completion" - they stay open indefinitely. + /// + /// Used for: + /// - Layer 1 live subscription (new announcements after initial sync) + /// - Layer 2+3 live subscriptions (new events after initial sync) + /// + /// # Arguments + /// * `relay_url` - The relay URL to subscribe on + /// * `filters` - Filters to subscribe to (will have `limit: 0` applied) + /// + /// # Returns + /// Vec of subscription IDs for the live subscriptions, or empty if connection not found + async fn sync_live(&self, relay_url: &str, filters: &[Filter]) -> Vec { + if filters.is_empty() { + return vec![]; + } + + let connection = match self.connections.get(relay_url) { + Some(conn) => conn, + None => { + tracing::warn!( + relay = %relay_url, + "No connection found for sync_live" + ); + return vec![]; + } + }; + + let mut sub_ids = Vec::new(); + + for filter in filters { + // Apply limit: 0 to make this a live subscription + // Note: nostr-sdk Filter doesn't have a limit(0) that means "no limit", + // but omitting limit means "no limit" which is what we want for live. + // The filter passed in should already NOT have a limit set. + match connection.subscribe_filter(filter.clone()).await { + Ok(sub_id) => { + tracing::trace!( + relay = %relay_url, + sub_id = %sub_id, + "Live subscription created" + ); + sub_ids.push(sub_id); + } + Err(e) => { + tracing::error!( + relay = %relay_url, + error = %e, + "Failed to create live subscription" + ); + } + } + } + + tracing::debug!( + relay = %relay_url, + filter_count = filters.len(), + sub_count = sub_ids.len(), + "sync_live completed" + ); + + sub_ids + } + + /// Reconstruct filters from RelaySyncIndex (confirmed state ONLY) + /// + /// Returns raw Vec for L1+L2+L3. + /// Used by: quick_reconnect, consolidate + /// Does NOT include pending items - those flow through AddFilters path. + /// + /// # Arguments + /// * `relay_url` - The relay URL to reconstruct filters for + /// + /// # Returns + /// Vec of filters for L1 (announcements) + L2 (repo tags) + L3 (event tags) + #[allow(dead_code)] // Will be used in Phase 3+ + async fn reconstruct_filters(&self, relay_url: &str) -> Vec { + // Get confirmed state from relay_sync_index + let (repos, root_events) = { + let index = self.relay_sync_index.read().await; + match index.get(relay_url) { + Some(state) => (state.repos.clone(), state.root_events.clone()), + None => { + tracing::warn!( + relay = %relay_url, + "No RelayState found for reconstruct_filters" + ); + return vec![]; + } + } + }; + + let mut all_filters = Vec::new(); + + // Layer 1: Announcements (always included) + // Note: No `since` filter - this returns raw filters for live subscriptions + all_filters.push(filters::build_announcement_filter(None)); + + // Layer 2 + Layer 3: Repo and root event tag filters + if !repos.is_empty() || !root_events.is_empty() { + let l2_l3_filters = + filters::build_layer2_and_layer3_filters(&repos, &root_events, None); + all_filters.extend(l2_l3_filters); + } + + tracing::debug!( + relay = %relay_url, + total_filters = all_filters.len(), + repos_count = repos.len(), + root_events_count = root_events.len(), + "Reconstructed filters from confirmed state" + ); + + all_filters + } + + /// Sync historical events and track in PendingSyncIndex + #[allow(dead_code)] // Will be used in Phase 3+ + /// + /// This method handles historical synchronization for a set of filters, + /// creating a PendingBatch to track completion. It dispatches to either + /// negentropy sync or traditional REQ+EOSE based on relay capability and config. + /// + /// Used for: + /// - Initial sync (no since filter) + /// - Reconnect sync (with since filter) + /// - Daily sync (no since filter, full re-sync) + /// + /// # Arguments + /// * `relay_url` - The relay URL to sync from + /// * `filters` - Filters to sync (will have `since` applied if provided) + /// * `items` - Items being synced (for tracking in PendingBatch) + /// * `since` - Optional timestamp for incremental sync + /// + /// # Returns + /// * `Some(batch_id)` - Batch was created and sync initiated + /// * `None` - No connection or sync failed to start + async fn historic_sync( + &mut self, + relay_url: &str, + filters: Vec, + items: PendingItems, + since: Option, + ) -> Option { + if filters.is_empty() && items.repos.is_empty() && items.root_events.is_empty() { + tracing::debug!( + relay = %relay_url, + "historic_sync called with empty filters and items, skipping" + ); + return None; + } + + // Check connection exists + let connection = match self.connections.get(relay_url) { + Some(conn) => conn, + None => { + tracing::warn!( + relay = %relay_url, + "No connection found for historic_sync" + ); + return None; + } + }; + + // Apply since filter if provided + let filters_with_since: Vec = if let Some(ts) = since { + filters.into_iter().map(|f| f.since(ts)).collect() + } else { + filters + }; + + // Check if we should use negentropy + let use_negentropy = + !self.config.sync_disable_negentropy && connection.supports_negentropy().await; + + // Generate batch ID + let batch_id = self.next_batch_id(); + + if use_negentropy && !filters_with_since.is_empty() { + // NIP-77 negentropy path + tracing::debug!( + relay = %relay_url, + batch_id = batch_id, + filter_count = filters_with_since.len(), + repos = items.repos.len(), + root_events = items.root_events.len(), + "Starting historic_sync with negentropy" + ); + + // Create PendingBatch for negentropy (empty outstanding_subs) + let batch = PendingBatch { + batch_id, + items: items.clone(), + outstanding_subs: HashSet::new(), + sync_method: SyncMethod::Negentropy, + }; + + // Add to pending_sync_index + { + let mut pending = self.pending_sync_index.write().await; + pending + .entry(relay_url.to_string()) + .or_insert_with(Vec::new) + .push(batch); + } + + // Perform negentropy sync for each filter + // Note: We sync each filter separately because negentropy works on a single filter + let mut total_received = 0; + let mut any_success = false; + + for filter in &filters_with_since { + if let Some(conn) = self.connections.get(relay_url) { + match conn.negentropy_sync_filter(filter.clone()).await { + Ok(result) => { + total_received += result.received.len(); + any_success = true; + + // Record metrics for received events + if let Some(ref metrics) = self.metrics { + for _ in 0..result.received.len() { + metrics.record_event(event_source::STARTUP); + } + } + } + Err(e) => { + tracing::warn!( + relay = %relay_url, + error = %e, + "Negentropy sync failed for filter in historic_sync" + ); + } + } + } + } + + if any_success { + // Remove batch from pending and confirm it + let completed_batch = { + let mut pending = self.pending_sync_index.write().await; + if let Some(batches) = pending.get_mut(relay_url) { + let batch_idx = batches.iter().position(|b| b.batch_id == batch_id); + if let Some(idx) = batch_idx { + let batch = batches.remove(idx); + if batches.is_empty() { + pending.remove(relay_url); + } + Some(batch) + } else { + None + } + } else { + None + } + }; + + if let Some(batch) = completed_batch { + self.confirm_batch(relay_url, batch).await; + } + + tracing::info!( + relay = %relay_url, + batch_id = batch_id, + total_received = total_received, + "historic_sync (negentropy) completed" + ); + } else { + // All negentropy syncs failed - remove the pending batch + let mut pending = self.pending_sync_index.write().await; + if let Some(batches) = pending.get_mut(relay_url) { + batches.retain(|b| b.batch_id != batch_id); + if batches.is_empty() { + pending.remove(relay_url); + } + } + + tracing::warn!( + relay = %relay_url, + batch_id = batch_id, + "historic_sync (negentropy) failed for all filters" + ); + return None; + } + } else { + // Traditional REQ+EOSE path + tracing::debug!( + relay = %relay_url, + batch_id = batch_id, + filter_count = filters_with_since.len(), + repos = items.repos.len(), + root_events = items.root_events.len(), + use_negentropy = use_negentropy, + "Starting historic_sync with REQ+EOSE" + ); + + // Subscribe to each filter and collect subscription IDs + let mut subscription_ids = HashSet::new(); + + for filter in &filters_with_since { + if let Some(conn) = self.connections.get(relay_url) { + match conn.subscribe_filter(filter.clone()).await { + Ok(sub_id) => { + subscription_ids.insert(sub_id); + } + Err(e) => { + tracing::error!( + relay = %relay_url, + error = %e, + "Failed to subscribe to filter in historic_sync" + ); + } + } + } + } + + if subscription_ids.is_empty() && !filters_with_since.is_empty() { + tracing::warn!( + relay = %relay_url, + "All filter subscriptions failed in historic_sync" + ); + return None; + } + + // Create PendingBatch for REQ+EOSE + let batch = PendingBatch { + batch_id, + items, + outstanding_subs: subscription_ids, + sync_method: SyncMethod::ReqEose, + }; + + // Add to pending_sync_index + { + let mut pending = self.pending_sync_index.write().await; + pending + .entry(relay_url.to_string()) + .or_insert_with(Vec::new) + .push(batch); + } + + tracing::debug!( + relay = %relay_url, + batch_id = batch_id, + "historic_sync (REQ+EOSE) batch created, awaiting EOSE" + ); + } + + Some(batch_id) + } + /// Gracefully shutdown the SyncManager /// /// This method: diff --git a/src/sync/self_subscriber.rs b/src/sync/self_subscriber.rs index 0379fe4..9643fc0 100644 --- a/src/sync/self_subscriber.rs +++ b/src/sync/self_subscriber.rs @@ -499,7 +499,7 @@ impl SelfSubscriber { drop(index); // Release lock before async operations // For each relay, send AddFilters action directly - // SyncManager's handle_add_filters auto-spawns connection for unknown relays + // SyncManager's handle_new_sync_filters auto-spawns connection for unknown relays for (relay_url, needs) in targets { // Skip our own relay URL (we're subscribed to ourselves via self-subscription) if relay_url.contains(&self.relay_domain) { @@ -519,8 +519,10 @@ impl SelfSubscriber { let action = AddFilters { relay_url: relay_url.clone(), - repos: needs.repos, - root_events: needs.root_events, + items: crate::sync::PendingItems { + repos: needs.repos, + root_events: needs.root_events, + }, filters, }; -- cgit v1.2.3