mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
improve sync design
This commit is contained in:
@@ -0,0 +1,871 @@
|
||||
# GRASP-02: Proactive Sync v3 - Event-Driven Design
|
||||
|
||||
## Overview
|
||||
|
||||
This document presents v3 of the proactive sync design. Key principles:
|
||||
|
||||
1. **Self-subscription as the only mechanism** - No database initialization at startup
|
||||
2. **Batch-based pending tracking** - Each batch confirms independently
|
||||
3. **Single action type** - AddFilters only, auto-spawn connections
|
||||
4. **Three-way state model** - RepoSyncIndex (want) → PendingSyncIndex (in-flight) → RelaySyncIndex (confirmed)
|
||||
|
||||
---
|
||||
|
||||
## Data Model
|
||||
|
||||
### RepoSyncIndex (Source of Truth)
|
||||
|
||||
```rust
|
||||
/// What we WANT to sync - derived from events received via self-subscription.
|
||||
/// Updated immediately when self-subscriber batch fires.
|
||||
/// Key: repo addressable ref ("30617:pubkey:identifier")
|
||||
pub type RepoSyncIndex = Arc<RwLock<HashMap<String, RepoSyncNeeds>>>;
|
||||
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct RepoSyncNeeds {
|
||||
/// Relay URLs listed in this repo's 30617 announcement
|
||||
pub relays: HashSet<String>,
|
||||
/// Root event IDs (1617/1618/1619/1621) that reference this repo
|
||||
pub root_events: HashSet<EventId>,
|
||||
}
|
||||
```
|
||||
|
||||
### RelaySyncIndex (Confirmed State + Connection)
|
||||
|
||||
```rust
|
||||
/// What we've CONFIRMED syncing - includes connection state for integrated lifecycle.
|
||||
/// Key: relay URL
|
||||
pub type RelaySyncIndex = Arc<RwLock<HashMap<String, RelayState>>>;
|
||||
|
||||
/// Connection status for a relay
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum ConnectionStatus {
|
||||
/// Not currently connected
|
||||
Disconnected,
|
||||
/// Connection attempt in progress
|
||||
Connecting,
|
||||
/// Successfully connected and subscribed
|
||||
Connected,
|
||||
}
|
||||
|
||||
/// Complete state for a single relay - combines sync needs with connection lifecycle
|
||||
#[derive(Debug)]
|
||||
pub struct RelayState {
|
||||
/// Repos we've confirmed syncing from this relay
|
||||
pub repos: HashSet<String>,
|
||||
/// Root events we've confirmed tracking
|
||||
pub root_events: HashSet<EventId>,
|
||||
/// If true, never disconnect this relay
|
||||
pub is_bootstrap: bool,
|
||||
/// Current connection status
|
||||
pub connection_status: ConnectionStatus,
|
||||
/// When we last successfully connected (for since filter on reconnect)
|
||||
pub last_connected: Option<Timestamp>,
|
||||
/// When we disconnected (for 15-minute state retention rule)
|
||||
pub disconnected_at: Option<Timestamp>,
|
||||
/// The active connection (None if disconnected)
|
||||
pub connection: Option<RelayConnection>,
|
||||
}
|
||||
|
||||
impl RelayState {
|
||||
/// Check if state should be cleared based on 15-minute rule
|
||||
pub fn should_clear_state(&self) -> bool {
|
||||
match self.disconnected_at {
|
||||
Some(disconnected) => {
|
||||
let now = Timestamp::now();
|
||||
now.as_u64().saturating_sub(disconnected.as_u64()) > 900 // 15 minutes
|
||||
}
|
||||
None => false, // Still connected or never connected
|
||||
}
|
||||
}
|
||||
|
||||
/// Clear repos and root_events (called when reconnect takes > 15 minutes)
|
||||
pub fn clear_sync_state(&mut self) {
|
||||
self.repos.clear();
|
||||
self.root_events.clear();
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### PendingSyncIndex (In-Flight Batches)
|
||||
|
||||
```rust
|
||||
/// Tracks batches of subscriptions that are in-flight, awaiting EOSE.
|
||||
/// Each batch has its own ID and can confirm independently.
|
||||
/// Key: relay URL
|
||||
pub type PendingSyncIndex = Arc<RwLock<HashMap<String, Vec<PendingBatch>>>>;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct PendingBatch {
|
||||
/// Unique ID for this batch (for debugging/logging)
|
||||
pub batch_id: u64,
|
||||
/// The items this batch is syncing
|
||||
pub items: PendingItems,
|
||||
/// Subscription IDs that must ALL receive EOSE before confirming
|
||||
pub outstanding_subs: HashSet<SubscriptionId>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct PendingItems {
|
||||
pub repos: HashSet<String>,
|
||||
pub root_events: HashSet<EventId>,
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## State Flow
|
||||
|
||||
```mermaid
|
||||
flowchart TB
|
||||
subgraph Input
|
||||
SS[SelfSubscriber]
|
||||
OWN[Own Relay]
|
||||
end
|
||||
|
||||
subgraph RepoSyncIndex - Want
|
||||
RSI[HashMap of Repo to Relays+Events]
|
||||
end
|
||||
|
||||
subgraph Derived Target
|
||||
DT[derive_relay_targets fn]
|
||||
TGT[Per-relay: repos + events we should sync]
|
||||
end
|
||||
|
||||
subgraph PendingSyncIndex - In Flight
|
||||
PSI[Vec of PendingBatch per relay]
|
||||
end
|
||||
|
||||
subgraph RelaySyncIndex - State + Connection
|
||||
RLI[RelayState per relay]
|
||||
CONN[connection: Option of RelayConnection]
|
||||
STATUS[connection_status: Connected/Disconnected/Connecting]
|
||||
REPOS[repos + root_events]
|
||||
end
|
||||
|
||||
SS -->|subscribe| OWN
|
||||
OWN -->|events| SS
|
||||
SS -->|batch fires| RSI
|
||||
RSI --> DT
|
||||
DT --> TGT
|
||||
TGT -->|diff: target - pending - confirmed| DIFF[Compute new items]
|
||||
PSI --> DIFF
|
||||
RLI --> DIFF
|
||||
DIFF -->|skip if disconnected| CHECK{Connected?}
|
||||
CHECK -->|yes| AF[AddFilters]
|
||||
CHECK -->|no| QUEUE[Queued in RelayState.repos]
|
||||
AF -->|subscribe| CONN
|
||||
AF -->|create batch| PSI
|
||||
CONN -->|EOSE| PSI
|
||||
PSI -->|batch complete| REPOS
|
||||
CONN -->|disconnect event| DISC[Mark Disconnected + set disconnected_at]
|
||||
DISC -->|reconnect| RECONN[On Reconnect]
|
||||
RECONN -->|check 15min rule| RULE{disconnected > 15min?}
|
||||
RULE -->|yes| CLEAR[Clear repos/root_events]
|
||||
RULE -->|no| RETAIN[Keep retained state]
|
||||
CLEAR --> REGEN[Regenerate AddFilters from RepoSyncIndex]
|
||||
RETAIN --> RESUB[Resubscribe with since filter]
|
||||
```
|
||||
|
||||
### Connection Lifecycle Integration
|
||||
|
||||
The `RelayState` struct now owns both the connection and sync state:
|
||||
|
||||
```rust
|
||||
// On disconnect (detected via RelayPoolNotification::Shutdown or handle_notifications returning)
|
||||
fn handle_disconnect(&mut self, relay_url: &str) {
|
||||
if let Some(state) = self.relay_sync_index.write().await.get_mut(relay_url) {
|
||||
state.connection_status = ConnectionStatus::Disconnected;
|
||||
state.disconnected_at = Some(Timestamp::now());
|
||||
state.connection = None;
|
||||
|
||||
// Clear any pending batches for this relay
|
||||
self.pending_sync_index.write().await.remove(relay_url);
|
||||
}
|
||||
}
|
||||
|
||||
// On reconnect
|
||||
async fn handle_reconnect(&mut self, relay_url: &str) -> Result<(), Error> {
|
||||
let mut index = self.relay_sync_index.write().await;
|
||||
let state = index.get_mut(relay_url).ok_or("Relay not in index")?;
|
||||
|
||||
// Apply 15-minute state retention rule
|
||||
if state.should_clear_state() {
|
||||
tracing::info!("Reconnect after >15min for {}, clearing state", relay_url);
|
||||
state.clear_sync_state();
|
||||
}
|
||||
|
||||
// Create new connection
|
||||
state.connection_status = ConnectionStatus::Connecting;
|
||||
let connection = RelayConnection::new(relay_url.to_string());
|
||||
|
||||
// Connect with since filter if we have last_connected
|
||||
let since = state.last_connected.map(|ts| {
|
||||
Timestamp::from(ts.as_u64().saturating_sub(900)) // -15 min buffer
|
||||
});
|
||||
|
||||
connection.connect_and_subscribe_with_since(since).await?;
|
||||
|
||||
state.connection = Some(connection);
|
||||
state.connection_status = ConnectionStatus::Connected;
|
||||
state.last_connected = Some(Timestamp::now());
|
||||
state.disconnected_at = None;
|
||||
|
||||
drop(index); // Release lock
|
||||
|
||||
// Regenerate AddFilters from current state (either retained or fresh from RepoSyncIndex)
|
||||
self.regenerate_filters_for_relay(relay_url).await;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Regenerate AddFilters for a relay after reconnection
|
||||
async fn regenerate_filters_for_relay(&mut self, relay_url: &str) {
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
let targets = derive_relay_targets(&repo_index);
|
||||
|
||||
if let Some(target) = targets.get(relay_url) {
|
||||
// Build filters for everything this relay should sync
|
||||
let filters = build_filters(&target.repos, &target.root_events);
|
||||
|
||||
// Create and process AddFilters action
|
||||
let action = AddFilters {
|
||||
relay_url: relay_url.to_string(),
|
||||
repos: target.repos.clone(),
|
||||
root_events: target.root_events.clone(),
|
||||
filters,
|
||||
};
|
||||
|
||||
self.handle_add_filters(action).await;
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Action Type
|
||||
|
||||
```rust
|
||||
/// Action sent from SelfSubscriber to SyncManager.
|
||||
/// SyncManager auto-spawns relay connections if they don't exist.
|
||||
pub struct AddFilters {
|
||||
pub relay_url: String,
|
||||
/// Items this action covers (for pending tracking)
|
||||
pub repos: HashSet<String>,
|
||||
pub root_events: HashSet<EventId>,
|
||||
/// Pre-batched filters (each with <= 100 tags)
|
||||
pub filters: Vec<Filter>,
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Core Algorithms
|
||||
|
||||
### 1. derive_relay_targets
|
||||
|
||||
Transform RepoSyncIndex into per-relay sync targets:
|
||||
|
||||
```rust
|
||||
fn derive_relay_targets(
|
||||
repo_index: &HashMap<String, RepoSyncNeeds>
|
||||
) -> HashMap<String, RelaySyncNeeds> {
|
||||
let mut targets: HashMap<String, RelaySyncNeeds> = HashMap::new();
|
||||
|
||||
for (repo_ref, needs) in repo_index {
|
||||
for relay_url in &needs.relays {
|
||||
let target = targets.entry(relay_url.clone()).or_default();
|
||||
target.repos.insert(repo_ref.clone());
|
||||
target.root_events.extend(needs.root_events.iter().cloned());
|
||||
}
|
||||
}
|
||||
|
||||
targets
|
||||
}
|
||||
```
|
||||
|
||||
### 2. compute_actions (Three-Way Diff)
|
||||
|
||||
```rust
|
||||
fn compute_actions(
|
||||
targets: &HashMap<String, RelaySyncNeeds>,
|
||||
pending: &HashMap<String, Vec<PendingBatch>>,
|
||||
confirmed: &HashMap<String, RelayState>,
|
||||
) -> Vec<AddFilters> {
|
||||
let mut actions = Vec::new();
|
||||
|
||||
for (relay_url, target) in targets {
|
||||
// Skip disconnected relays - they'll get AddFilters on reconnect
|
||||
if let Some(state) = confirmed.get(relay_url) {
|
||||
if state.connection_status != ConnectionStatus::Connected {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
// Collect all pending items for this relay
|
||||
let pending_repos: HashSet<_> = pending.get(relay_url)
|
||||
.map(|batches| batches.iter()
|
||||
.flat_map(|b| b.items.repos.iter().cloned())
|
||||
.collect())
|
||||
.unwrap_or_default();
|
||||
let pending_events: HashSet<_> = pending.get(relay_url)
|
||||
.map(|batches| batches.iter()
|
||||
.flat_map(|b| b.items.root_events.iter().cloned())
|
||||
.collect())
|
||||
.unwrap_or_default();
|
||||
|
||||
// Collect confirmed items for this relay
|
||||
let confirmed_repos = confirmed.get(relay_url)
|
||||
.map(|c| &c.repos)
|
||||
.unwrap_or(&HashSet::new());
|
||||
let confirmed_events = confirmed.get(relay_url)
|
||||
.map(|c| &c.root_events)
|
||||
.unwrap_or(&HashSet::new());
|
||||
|
||||
// New = target - pending - confirmed
|
||||
let new_repos: HashSet<_> = target.repos.iter()
|
||||
.filter(|r| !pending_repos.contains(*r) && !confirmed_repos.contains(*r))
|
||||
.cloned()
|
||||
.collect();
|
||||
let new_events: HashSet<_> = target.root_events.iter()
|
||||
.filter(|e| !pending_events.contains(*e) && !confirmed_events.contains(*e))
|
||||
.cloned()
|
||||
.collect();
|
||||
|
||||
if !new_repos.is_empty() || !new_events.is_empty() {
|
||||
let filters = build_filters(&new_repos, &new_events);
|
||||
actions.push(AddFilters {
|
||||
relay_url: relay_url.clone(),
|
||||
repos: new_repos,
|
||||
root_events: new_events,
|
||||
filters,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
actions
|
||||
}
|
||||
```
|
||||
|
||||
### 3. handle_add_filters (SyncManager)
|
||||
|
||||
```rust
|
||||
impl SyncManager {
|
||||
async fn handle_add_filters(&mut self, action: AddFilters) {
|
||||
let AddFilters { relay_url, repos, root_events, filters } = action;
|
||||
|
||||
// Auto-spawn connection if needed
|
||||
if !self.connections.contains_key(&relay_url) {
|
||||
self.spawn_connection(&relay_url).await;
|
||||
}
|
||||
|
||||
let conn = self.connections.get(&relay_url).unwrap();
|
||||
|
||||
// Subscribe and collect subscription IDs
|
||||
// nostr-sdk 0.44: subscribe returns Output<Vec<SubscriptionId>>
|
||||
// since we're only subscribed to one relay per connection
|
||||
let mut sub_ids = HashSet::new();
|
||||
for filter in filters {
|
||||
// cloned filter for each subscription call
|
||||
match conn.client.subscribe(filter, None).await {
|
||||
Ok(output) => {
|
||||
// Output contains subscription IDs for each relay
|
||||
for sub_id in output.val {
|
||||
sub_ids.insert(sub_id);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!("Failed to subscribe: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Create pending batch
|
||||
let batch = PendingBatch {
|
||||
batch_id: self.next_batch_id(),
|
||||
items: PendingItems { repos, root_events },
|
||||
outstanding_subs: sub_ids,
|
||||
};
|
||||
|
||||
// Add to pending index
|
||||
self.pending_sync_index.write().await
|
||||
.entry(relay_url)
|
||||
.or_default()
|
||||
.push(batch);
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### 4. handle_eose (Batch Completion)
|
||||
|
||||
```rust
|
||||
impl SyncManager {
|
||||
async fn handle_eose(&mut self, relay_url: &str, sub_id: SubscriptionId) {
|
||||
let mut pending = self.pending_sync_index.write().await;
|
||||
|
||||
if let Some(batches) = pending.get_mut(relay_url) {
|
||||
// Find which batch this subscription belongs to
|
||||
for batch in batches.iter_mut() {
|
||||
if batch.outstanding_subs.remove(&sub_id) {
|
||||
// Check if batch is now complete
|
||||
if batch.outstanding_subs.is_empty() {
|
||||
// Move items to confirmed
|
||||
let items = batch.items.clone();
|
||||
drop(pending); // Release lock before acquiring another
|
||||
|
||||
let mut confirmed = self.relay_sync_index.write().await;
|
||||
let relay_confirmed = confirmed
|
||||
.entry(relay_url.to_string())
|
||||
.or_default();
|
||||
relay_confirmed.repos.extend(items.repos);
|
||||
relay_confirmed.root_events.extend(items.root_events);
|
||||
|
||||
tracing::info!(
|
||||
"Batch {} complete for {} - confirmed {} repos, {} events",
|
||||
batch.batch_id, relay_url,
|
||||
items.repos.len(), items.root_events.len()
|
||||
);
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
// Clean up completed batches
|
||||
if let Some(batches) = pending.get_mut(relay_url) {
|
||||
batches.retain(|b| !b.outstanding_subs.is_empty());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Self-Subscriber Flow
|
||||
|
||||
### State Tracking
|
||||
|
||||
```rust
|
||||
pub struct SelfSubscriber {
|
||||
own_relay_url: String,
|
||||
relay_domain: String,
|
||||
repo_sync_index: RepoSyncIndex,
|
||||
pending_sync_index: PendingSyncIndex,
|
||||
relay_sync_index: RelaySyncIndex,
|
||||
action_tx: mpsc::Sender<AddFilters>,
|
||||
/// Timestamp of last successful connection - used for since filter on reconnection
|
||||
last_connected: Option<Timestamp>,
|
||||
/// Is this the first connection attempt since startup?
|
||||
is_initial_connect: bool,
|
||||
}
|
||||
```
|
||||
|
||||
### On Startup
|
||||
|
||||
```rust
|
||||
impl SelfSubscriber {
|
||||
async fn run(mut self) {
|
||||
// Connect to own relay
|
||||
let client = Client::new(Keys::generate());
|
||||
client.add_relay(&self.own_relay_url).await?;
|
||||
client.connect().await;
|
||||
|
||||
// Track connection time
|
||||
self.last_connected = Some(Timestamp::now());
|
||||
|
||||
// Subscribe WITHOUT since filter (get all historical) on first connect
|
||||
let filter = Filter::new().kinds([
|
||||
Kind::Custom(30617), // Repository announcements
|
||||
Kind::GitPatch, // 1617
|
||||
Kind::Custom(1618), // PRs
|
||||
Kind::Custom(1619), // PR updates
|
||||
Kind::GitIssue, // 1621
|
||||
]);
|
||||
|
||||
client.subscribe(filter, None).await?;
|
||||
self.is_initial_connect = false;
|
||||
|
||||
// Run event loop with batching
|
||||
self.event_loop(&client).await;
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### On Reconnection
|
||||
|
||||
```rust
|
||||
impl SelfSubscriber {
|
||||
async fn reconnect(&mut self, client: &Client) -> Result<(), Error> {
|
||||
// Reconnect to own relay
|
||||
client.connect().await;
|
||||
|
||||
// On reconnection ONLY, use since filter based on last_connected
|
||||
let since = match self.last_connected {
|
||||
Some(ts) => Timestamp::from(ts.as_u64().saturating_sub(900)), // -15 minutes buffer
|
||||
None => Timestamp::from(0), // Shouldn't happen, but fall back to full sync
|
||||
};
|
||||
|
||||
// Update last_connected AFTER computing since
|
||||
self.last_connected = Some(Timestamp::now());
|
||||
|
||||
let filter = Filter::new()
|
||||
.kinds([
|
||||
Kind::Custom(30617),
|
||||
Kind::GitPatch,
|
||||
Kind::Custom(1618),
|
||||
Kind::Custom(1619),
|
||||
Kind::GitIssue,
|
||||
])
|
||||
.since(since);
|
||||
|
||||
client.subscribe(filter, None).await?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### Batching Logic
|
||||
|
||||
```rust
|
||||
impl SelfSubscriber {
|
||||
async fn event_loop(&self, client: &Client) {
|
||||
let mut pending_events: Vec<Event> = Vec::new();
|
||||
let mut batch_timer: Option<Instant> = None;
|
||||
let batch_window = Duration::from_secs(5);
|
||||
|
||||
loop {
|
||||
let timeout = batch_timer
|
||||
.map(|t| batch_window.saturating_sub(t.elapsed()))
|
||||
.unwrap_or(Duration::from_secs(60));
|
||||
|
||||
tokio::select! {
|
||||
notification = client.notifications().recv() => {
|
||||
if let Ok(RelayPoolNotification::Event { event, .. }) = notification {
|
||||
pending_events.push(*event);
|
||||
|
||||
// Start timer on first event (does NOT reset)
|
||||
if batch_timer.is_none() {
|
||||
batch_timer = Some(Instant::now());
|
||||
}
|
||||
}
|
||||
}
|
||||
_ = tokio::time::sleep(timeout), if batch_timer.is_some() => {
|
||||
// Batch window elapsed
|
||||
self.process_batch(pending_events.drain(..).collect()).await;
|
||||
batch_timer = None;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn process_batch(&self, events: Vec<Event>) {
|
||||
// 1. Update RepoSyncIndex
|
||||
for event in events {
|
||||
match event.kind.as_u16() {
|
||||
30617 => self.handle_announcement(&event).await,
|
||||
1617 | 1618 | 1619 | 1621 => self.handle_root_event(&event).await,
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Derive targets and compute actions
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
let targets = derive_relay_targets(&repo_index);
|
||||
|
||||
let pending = self.pending_sync_index.read().await;
|
||||
let confirmed = self.relay_sync_index.read().await;
|
||||
|
||||
let actions = compute_actions(&targets, &pending, &confirmed);
|
||||
|
||||
drop(repo_index);
|
||||
drop(pending);
|
||||
drop(confirmed);
|
||||
|
||||
// 3. Send actions to SyncManager
|
||||
for action in actions {
|
||||
let _ = self.action_tx.send(action).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Bootstrap Relay
|
||||
|
||||
```rust
|
||||
impl SyncManager {
|
||||
async fn initialize_bootstrap(&mut self) {
|
||||
if let Some(url) = &self.config.bootstrap_relay_url {
|
||||
// Pre-mark as bootstrap (never removed)
|
||||
self.relay_sync_index.write().await.insert(
|
||||
url.clone(),
|
||||
RelaySyncNeeds {
|
||||
repos: HashSet::new(),
|
||||
root_events: HashSet::new(),
|
||||
is_bootstrap: true,
|
||||
}
|
||||
);
|
||||
|
||||
// Send Layer 1 filter
|
||||
let filters = vec![
|
||||
Filter::new().kinds([Kind::Custom(30617), Kind::Custom(30618)])
|
||||
];
|
||||
|
||||
self.handle_add_filters(AddFilters {
|
||||
relay_url: url.clone(),
|
||||
repos: HashSet::new(), // Layer 1 doesn't track specific repos
|
||||
root_events: HashSet::new(),
|
||||
filters,
|
||||
}).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Disconnect Handling
|
||||
|
||||
Direct in SyncManager (not via action):
|
||||
|
||||
```rust
|
||||
impl SyncManager {
|
||||
async fn check_disconnects(&mut self) {
|
||||
let confirmed = self.relay_sync_index.read().await;
|
||||
|
||||
for (relay_url, state) in confirmed.iter() {
|
||||
if state.is_bootstrap {
|
||||
continue; // Never disconnect bootstrap
|
||||
}
|
||||
|
||||
if state.repos.is_empty() && state.root_events.is_empty() {
|
||||
// No repos - disconnect
|
||||
self.disconnect_relay(relay_url).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn disconnect_relay(&mut self, relay_url: &str) {
|
||||
self.relay_sync_index.write().await.remove(relay_url);
|
||||
self.pending_sync_index.write().await.remove(relay_url);
|
||||
|
||||
if let Some(conn) = self.connections.remove(relay_url) {
|
||||
conn.disconnect().await;
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Relay Connection Lifecycle
|
||||
|
||||
### State Machine for External Relays
|
||||
|
||||
```mermaid
|
||||
stateDiagram-v2
|
||||
[*] --> Connecting: spawn_connection
|
||||
Connecting --> Connected: success
|
||||
Connecting --> Backoff: failure
|
||||
Connected --> Disconnected: connection lost
|
||||
Connected --> [*]: intentional disconnect
|
||||
Disconnected --> Backoff: record_failure
|
||||
Backoff --> Connecting: backoff elapsed
|
||||
Backoff --> Dead: 24h continuous failures
|
||||
Dead --> Connecting: daily retry
|
||||
```
|
||||
|
||||
### Health Integration
|
||||
|
||||
Uses `RelayHealthTracker` from [`src/sync/health.rs`](../../src/sync/health.rs):
|
||||
|
||||
```rust
|
||||
impl SyncManager {
|
||||
/// Spawn a connection with health tracking
|
||||
async fn spawn_connection(&mut self, relay_url: &str) {
|
||||
// Check if we should attempt connection
|
||||
if !self.health_tracker.should_attempt_connection(relay_url) {
|
||||
let remaining = self.health_tracker.get_remaining_backoff(relay_url);
|
||||
tracing::debug!(
|
||||
"Skipping connection to {} - backoff {:?}",
|
||||
relay_url,
|
||||
remaining
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
match self.try_connect(relay_url).await {
|
||||
Ok(conn) => {
|
||||
self.health_tracker.record_success(relay_url);
|
||||
self.connections.insert(relay_url.to_string(), conn);
|
||||
}
|
||||
Err(e) => {
|
||||
self.health_tracker.record_failure(relay_url);
|
||||
tracing::warn!("Connection to {} failed: {}", relay_url, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### Reconnection Loop
|
||||
|
||||
Each relay connection runs its own reconnection loop:
|
||||
|
||||
```rust
|
||||
impl RelayConnection {
|
||||
async fn run_with_reconnection(
|
||||
mut self,
|
||||
health_tracker: Arc<RelayHealthTracker>,
|
||||
event_tx: mpsc::Sender<RelayEvent>,
|
||||
) {
|
||||
loop {
|
||||
// Check backoff before attempting
|
||||
if !health_tracker.should_attempt_connection(&self.url) {
|
||||
if let Some(remaining) = health_tracker.get_remaining_backoff(&self.url) {
|
||||
tokio::time::sleep(remaining).await;
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
// Attempt connection
|
||||
match self.connect_and_subscribe().await {
|
||||
Ok(()) => {
|
||||
health_tracker.record_success(&self.url);
|
||||
|
||||
// Track when we connected for since filter on reconnect
|
||||
let connected_at = Timestamp::now();
|
||||
|
||||
// Run event loop until disconnection
|
||||
self.run_event_loop(&event_tx).await;
|
||||
|
||||
// Connection lost - will reconnect with since filter
|
||||
health_tracker.record_failure(&self.url);
|
||||
|
||||
// On reconnect, use since = connected_at - 15 minutes
|
||||
self.set_reconnect_since(connected_at);
|
||||
}
|
||||
Err(e) => {
|
||||
health_tracker.record_failure(&self.url);
|
||||
tracing::warn!("Connection to {} failed: {}", self.url, e);
|
||||
}
|
||||
}
|
||||
|
||||
// Get backoff duration and wait
|
||||
let state = health_tracker.get_state(&self.url);
|
||||
if state == HealthState::Dead {
|
||||
// Dead relays retry once per 24 hours
|
||||
tokio::time::sleep(Duration::from_secs(24 * 3600)).await;
|
||||
}
|
||||
// Otherwise, loop will check should_attempt_connection
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### Backoff Configuration
|
||||
|
||||
From existing [`RelayHealthTracker`](../../src/sync/health.rs:91):
|
||||
|
||||
| Parameter | Value | Notes |
|
||||
|-----------|-------|-------|
|
||||
| Base backoff | 5 seconds | First failure |
|
||||
| Backoff multiplier | 2x | Exponential increase |
|
||||
| Max backoff | 1 hour (configurable) | `sync_max_backoff_secs` |
|
||||
| Dead threshold | 24 hours | Continuous failures |
|
||||
| Dead retry interval | 24 hours | Once per day |
|
||||
|
||||
---
|
||||
|
||||
## Consolidation
|
||||
|
||||
### Threshold-Based (70 filters)
|
||||
|
||||
```rust
|
||||
impl SyncManager {
|
||||
async fn maybe_consolidate(&mut self, relay_url: &str) {
|
||||
let filter_count = self.get_filter_count(relay_url).await;
|
||||
|
||||
if filter_count > 70 {
|
||||
self.consolidate(relay_url).await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn consolidate(&mut self, relay_url: &str) {
|
||||
// 1. Wait for all pending batches to complete
|
||||
self.wait_pending_complete(relay_url).await;
|
||||
|
||||
// 2. Close all subscriptions
|
||||
self.close_all_subs(relay_url).await;
|
||||
|
||||
// 3. Rebuild filters from confirmed state
|
||||
let confirmed = self.relay_sync_index.read().await;
|
||||
let state = confirmed.get(relay_url)?;
|
||||
let filters = build_filters(&state.repos, &state.root_events);
|
||||
|
||||
// 4. Resubscribe with since = now - 15 minutes
|
||||
let since = Timestamp::now() - 900;
|
||||
for filter in filters {
|
||||
self.subscribe(relay_url, filter.since(since)).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### Daily Timer (23-25h Random)
|
||||
|
||||
```rust
|
||||
impl SyncManager {
|
||||
async fn run_daily_consolidation(&self) {
|
||||
loop {
|
||||
let hours = 23 + rand::random::<f64>() * 2.0;
|
||||
tokio::time::sleep(Duration::from_secs_f64(hours * 3600.0)).await;
|
||||
|
||||
for relay_url in self.connections.keys() {
|
||||
self.consolidate(relay_url).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Key Design Decisions
|
||||
|
||||
| Decision | Choice | Rationale |
|
||||
|----------|--------|-----------|
|
||||
| Startup mechanism | Self-subscription only | Single code path, fresh DB behaves same as reconnect |
|
||||
| Since filter | Only on reconnection | Initial subscribe gets full history |
|
||||
| Pending tracking | Per-batch with batch ID | Independent confirmation, no blocking |
|
||||
| EOSE requirement | All subs in batch must complete | Single repo may need multiple filter subs |
|
||||
| Action type | Struct not enum | Only one action type needed |
|
||||
| Relay spawning | Auto-spawn on AddFilters | Simplifies action logic |
|
||||
| Disconnect | Direct in SyncManager | Not worth an action type |
|
||||
| Consolidation | 70 filters + daily timer | Threshold for growth, timer for staleness |
|
||||
| Timestamps | In-memory only | Not critical for correctness |
|
||||
| Health tracking | Reuse existing RelayHealthTracker | Already implements exponential backoff, dead relay detection |
|
||||
| Reconnection backoff | Exponential to 1h max | Prevents hammering failed relays |
|
||||
| Dead relay policy | 24h threshold, daily retry | Balance between giving up and resource waste |
|
||||
| last_connected tracking | Per-connection in-memory | Enables 15-minute buffer on reconnect |
|
||||
| Connection ownership | Inside RelayState | Ties connection lifecycle to sync state, simpler than separate maps |
|
||||
| State retention rule | Clear if disconnected >15min | Matches since filter buffer, prevents stale subscriptions |
|
||||
| Skip disconnected | compute_actions skips disconnected | Prevents queuing AddFilters for offline relays |
|
||||
| Reconnect triggers | handle_notifications returns or Shutdown | nostr-sdk signals disconnect via event loop exit |
|
||||
| On-reconnect flow | Regenerate AddFilters from RepoSyncIndex | Fresh subscriptions for what we actually need |
|
||||
|
||||
---
|
||||
|
||||
## Module Structure
|
||||
|
||||
```
|
||||
src/sync/
|
||||
├── mod.rs # SyncManager, main loop
|
||||
├── state.rs # RepoSyncIndex, RelaySyncIndex, PendingSyncIndex types
|
||||
├── actions.rs # AddFilters struct, compute_actions
|
||||
├── self_subscriber.rs # SelfSubscriber, batching logic
|
||||
├── relay_connection.rs # Per-relay WebSocket connection
|
||||
├── consolidation.rs # Consolidation logic, daily timer
|
||||
├── health.rs # Health tracking (reuse from v2)
|
||||
└── metrics.rs # Prometheus metrics (reuse from v2)
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,373 @@
|
||||
# State Structure Redesign Proposal v2
|
||||
|
||||
## The Core Problem
|
||||
|
||||
We need to transform:
|
||||
- **Repo Announcements** (30617) that list relays
|
||||
- **Root Events** (1617/1618/1619/1621) that tag repos
|
||||
|
||||
Into:
|
||||
- **Per-relay subscriptions**: which repos and root events to sync from each relay
|
||||
|
||||
And generate **RelayActions** when this mapping changes.
|
||||
|
||||
---
|
||||
|
||||
## Proposed Data Model
|
||||
|
||||
### 1. RepoIndex (Primary Source of Truth)
|
||||
|
||||
```rust
|
||||
/// Everything we know about repos we're tracking
|
||||
/// Key: repo addressable ref ("30617:pubkey:identifier")
|
||||
pub type RepoIndex = Arc<RwLock<HashMap<String, RepoInfo>>>;
|
||||
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct RepoInfo {
|
||||
/// Relay URLs listed in the repo's announcement
|
||||
pub relays: HashSet<String>,
|
||||
/// Root event IDs that reference this repo
|
||||
pub root_events: HashSet<EventId>,
|
||||
}
|
||||
```
|
||||
|
||||
**Updated by:** Database init, batch processing of new announcements/root events
|
||||
|
||||
### 2. RelayIndex (Applied State)
|
||||
|
||||
```rust
|
||||
/// What we've told each relay to sync
|
||||
/// Key: relay URL
|
||||
pub type RelayIndex = Arc<RwLock<HashMap<String, SyncTarget>>>;
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq)]
|
||||
pub struct SyncTarget {
|
||||
/// Repos we're syncing for this relay
|
||||
pub repos: HashSet<String>,
|
||||
/// Root events we're tracking
|
||||
pub root_events: HashSet<EventId>,
|
||||
}
|
||||
```
|
||||
|
||||
**Updated by:** SyncManager after RelayActions are applied
|
||||
|
||||
---
|
||||
|
||||
## The Transformation
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
subgraph Input
|
||||
RA[Repo Announcements]
|
||||
RE[Root Events]
|
||||
end
|
||||
|
||||
subgraph RepoIndex
|
||||
R1[repo_a: relays=X,Y events=1,2]
|
||||
R2[repo_b: relays=Y,Z events=3]
|
||||
end
|
||||
|
||||
subgraph Derived Target
|
||||
T1[relay_X: repos=a events=1,2]
|
||||
T2[relay_Y: repos=a,b events=1,2,3]
|
||||
T3[relay_Z: repos=b events=3]
|
||||
end
|
||||
|
||||
subgraph RelayIndex Applied
|
||||
A1[relay_X: repos=a events=1,2]
|
||||
A2[relay_Y: repos=a events=1,2]
|
||||
end
|
||||
|
||||
RA --> R1
|
||||
RA --> R2
|
||||
RE --> R1
|
||||
RE --> R2
|
||||
|
||||
R1 --> T1
|
||||
R1 --> T2
|
||||
R2 --> T2
|
||||
R2 --> T3
|
||||
```
|
||||
|
||||
The **diff** between Derived Target and RelayIndex produces RelayActions:
|
||||
- relay_Y needs AddFilters for repo_b and event 3
|
||||
- relay_Z needs SpawnRelay
|
||||
|
||||
---
|
||||
|
||||
## Algorithm: derive_target_from_repo_index
|
||||
|
||||
```rust
|
||||
/// Derive what we SHOULD be syncing from the repo data
|
||||
fn derive_relay_targets(repo_index: &HashMap<String, RepoInfo>) -> HashMap<String, SyncTarget> {
|
||||
let mut targets: HashMap<String, SyncTarget> = HashMap::new();
|
||||
|
||||
for (repo_ref, info) in repo_index {
|
||||
// For each relay that lists this repo
|
||||
for relay_url in &info.relays {
|
||||
let target = targets.entry(relay_url.clone()).or_default();
|
||||
target.repos.insert(repo_ref.clone());
|
||||
target.root_events.extend(info.root_events.iter().cloned());
|
||||
}
|
||||
}
|
||||
|
||||
targets
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Algorithm: process_batch
|
||||
|
||||
```rust
|
||||
async fn process_batch(&self, pending: &mut PendingUpdates) {
|
||||
// ============================================
|
||||
// STEP 1: Update RepoIndex from batch
|
||||
// ============================================
|
||||
|
||||
let mut repo_index = self.repo_index.write().await;
|
||||
|
||||
// 1a. Process root events - add to repo's root_events set
|
||||
for event in pending.root_events.drain(..) {
|
||||
for repo_ref in extract_repo_refs(&event) {
|
||||
repo_index.entry(repo_ref)
|
||||
.or_default()
|
||||
.root_events
|
||||
.insert(event.id);
|
||||
}
|
||||
}
|
||||
|
||||
// 1b. Process announcements - update repo's relay set
|
||||
for event in pending.announcements.drain(..) {
|
||||
if !lists_our_service(&event) {
|
||||
continue;
|
||||
}
|
||||
let repo_ref = build_repo_ref(&event);
|
||||
let relay_urls: HashSet<String> = extract_relay_urls(&event)
|
||||
.into_iter()
|
||||
.filter(|url| !is_own_relay(url))
|
||||
.collect();
|
||||
|
||||
// Replace relay set (handles updates that change relays)
|
||||
repo_index.entry(repo_ref)
|
||||
.or_default()
|
||||
.relays = relay_urls;
|
||||
}
|
||||
|
||||
// ============================================
|
||||
// STEP 2: Derive target state from RepoIndex
|
||||
// ============================================
|
||||
|
||||
let target = derive_relay_targets(&repo_index);
|
||||
drop(repo_index); // Release write lock
|
||||
|
||||
// ============================================
|
||||
// STEP 3: Diff target vs applied (RelayIndex)
|
||||
// ============================================
|
||||
|
||||
let applied = self.relay_index.read().await;
|
||||
let actions = compute_relay_actions(&target, &applied);
|
||||
drop(applied); // Release read lock
|
||||
|
||||
// ============================================
|
||||
// STEP 4: Send actions & update RelayIndex
|
||||
// ============================================
|
||||
|
||||
for action in actions {
|
||||
match &action {
|
||||
RelayAction::SpawnRelay { relay_url, repos_and_root_events } => {
|
||||
// Update RelayIndex with new relay
|
||||
let mut applied = self.relay_index.write().await;
|
||||
applied.insert(relay_url.clone(), SyncTarget {
|
||||
repos: repos_and_root_events.keys().cloned().collect(),
|
||||
root_events: repos_and_root_events.values()
|
||||
.flat_map(|e| e.iter().cloned())
|
||||
.collect(),
|
||||
});
|
||||
}
|
||||
RelayAction::AddFilters { relay_url, repos_and_new_root_event } => {
|
||||
// Update RelayIndex with additions
|
||||
let mut applied = self.relay_index.write().await;
|
||||
if let Some(target) = applied.get_mut(relay_url) {
|
||||
for (repo, events) in repos_and_new_root_event {
|
||||
target.repos.insert(repo.clone());
|
||||
target.root_events.extend(events.iter().cloned());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Send action to SyncManager
|
||||
let _ = self.action_tx.send(action).await;
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Algorithm: compute_relay_actions
|
||||
|
||||
```rust
|
||||
fn compute_relay_actions(
|
||||
target: &HashMap<String, SyncTarget>,
|
||||
applied: &HashMap<String, SyncTarget>,
|
||||
) -> Vec<RelayAction> {
|
||||
let mut actions = Vec::new();
|
||||
|
||||
for (relay_url, target_state) in target {
|
||||
match applied.get(relay_url) {
|
||||
None => {
|
||||
// New relay - spawn it
|
||||
let mut repos_and_events = HashMap::new();
|
||||
for repo in &target_state.repos {
|
||||
// Get events for this specific repo
|
||||
let events = target_state.root_events.clone(); // simplified
|
||||
repos_and_events.insert(repo.clone(), events);
|
||||
}
|
||||
actions.push(RelayAction::SpawnRelay {
|
||||
relay_url: relay_url.clone(),
|
||||
repos_and_root_events: repos_and_events,
|
||||
});
|
||||
}
|
||||
Some(applied_state) => {
|
||||
// Existing relay - check for new repos/events
|
||||
let new_repos: HashSet<_> = target_state.repos
|
||||
.difference(&applied_state.repos)
|
||||
.cloned()
|
||||
.collect();
|
||||
let new_events: HashSet<_> = target_state.root_events
|
||||
.difference(&applied_state.root_events)
|
||||
.cloned()
|
||||
.collect();
|
||||
|
||||
if !new_repos.is_empty() || !new_events.is_empty() {
|
||||
let mut repos_and_events = HashMap::new();
|
||||
for repo in &new_repos {
|
||||
repos_and_events.insert(repo.clone(), new_events.clone());
|
||||
}
|
||||
// Also handle new events for existing repos
|
||||
if !new_events.is_empty() && new_repos.is_empty() {
|
||||
for repo in &applied_state.repos {
|
||||
repos_and_events.insert(repo.clone(), new_events.clone());
|
||||
}
|
||||
}
|
||||
|
||||
actions.push(RelayAction::AddFilters {
|
||||
relay_url: relay_url.clone(),
|
||||
repos_and_new_root_event: repos_and_events,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Future: detect relay removal (in applied but not in target)
|
||||
|
||||
actions
|
||||
}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Handling Announcement Updates
|
||||
|
||||
When an announcement is **updated** and changes its relay list:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A[repo_a announcement updated] --> B[Old: relays X,Y]
|
||||
B --> C[New: relays Y,Z]
|
||||
C --> D[RepoIndex updated: repo_a.relays = Y,Z]
|
||||
D --> E[derive_relay_targets]
|
||||
E --> F[Target: X=empty, Y=repo_a, Z=repo_a]
|
||||
F --> G[Diff with Applied: X=repo_a, Y=repo_a]
|
||||
G --> H1[X: repo_a removed - future RemoveFilters]
|
||||
G --> H2[Z: new relay - SpawnRelay]
|
||||
```
|
||||
|
||||
The current RelayAction types only support growth (SpawnRelay, AddFilters). Removal would need a new `RemoveFilters` action type - this is a future enhancement.
|
||||
|
||||
---
|
||||
|
||||
## Name Mappings
|
||||
|
||||
| Current | Proposed | Semantics |
|
||||
|---------|----------|-----------|
|
||||
| `FollowingRepoRootEvents` | `RepoIndex` | Per-repo: relays + root events |
|
||||
| `SyncRelays` | `RelayIndex` | Per-relay: what we're syncing (applied state) |
|
||||
| - | `SyncTarget` | Struct for repos + events |
|
||||
| - | `RepoInfo` | Struct for relay set + event set |
|
||||
|
||||
---
|
||||
|
||||
## Data Flow Summary
|
||||
|
||||
```mermaid
|
||||
flowchart TB
|
||||
subgraph Batch Input
|
||||
RA[30617 Announcements]
|
||||
RE[Root Events 1617-1621]
|
||||
end
|
||||
|
||||
subgraph Step 1: Update Source
|
||||
RI[RepoIndex]
|
||||
end
|
||||
|
||||
subgraph Step 2: Derive Target
|
||||
DT[derive_relay_targets]
|
||||
TGT[Target HashMap]
|
||||
end
|
||||
|
||||
subgraph Step 3: Diff
|
||||
RLI[RelayIndex - Applied]
|
||||
DIFF[compute_relay_actions]
|
||||
end
|
||||
|
||||
subgraph Step 4: Apply
|
||||
ACT[RelayActions]
|
||||
SM[SyncManager]
|
||||
end
|
||||
|
||||
RA --> RI
|
||||
RE --> RI
|
||||
RI --> DT
|
||||
DT --> TGT
|
||||
TGT --> DIFF
|
||||
RLI --> DIFF
|
||||
DIFF --> ACT
|
||||
ACT --> SM
|
||||
ACT --> |update| RLI
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Files to Modify
|
||||
|
||||
| File | Changes |
|
||||
|------|---------|
|
||||
| [`src/sync/mod.rs`](src/sync/mod.rs) | Replace type aliases with RepoIndex/RelayIndex + structs |
|
||||
| [`src/sync/self_subscriber.rs`](src/sync/self_subscriber.rs) | Rewrite process_batch with new algorithm |
|
||||
|
||||
---
|
||||
|
||||
## Questions for Approval
|
||||
|
||||
1. **Naming**: Are `RepoIndex`/`RelayIndex` and `RepoInfo`/`SyncTarget` clear enough?
|
||||
|
||||
2. **When to update RelayIndex**: Should we:
|
||||
- (a) Update immediately when generating action (optimistic) ← proposed above
|
||||
- (b) Update only after SyncManager confirms action succeeded
|
||||
|
||||
3. **Bootstrap relay**: Keep special-casing it in RelayIndex (always present)?
|
||||
|
||||
4. **Future work**: Add `RemoveFilters` action for relay removal, or defer?
|
||||
|
||||
---
|
||||
|
||||
## Benefits
|
||||
|
||||
1. **Logical flow**: Source → Derived → Diff → Actions
|
||||
2. **Single source of truth**: RepoIndex is the authoritative data
|
||||
3. **Clear transformation**: `derive_relay_targets()` is a pure function
|
||||
4. **Handles updates**: Replacing `repo.relays` naturally handles announcement changes
|
||||
5. **Testable**: Each step can be unit tested independently
|
||||
+304
-31
@@ -37,17 +37,26 @@
|
||||
//! for the complete design context.
|
||||
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
|
||||
use nostr_relay_builder::prelude::{Event, Filter, Kind, TagKind};
|
||||
use nostr_relay_builder::prelude::{
|
||||
DatabaseEventStatus, Event, Filter, Kind, PolicyResult, SaveEventStatus, TagKind, WritePolicy,
|
||||
};
|
||||
use nostr_sdk::prelude::*;
|
||||
use nostr_sdk::EventId;
|
||||
use tokio::sync::RwLock;
|
||||
use tokio::sync::{mpsc, RwLock};
|
||||
|
||||
use crate::config::Config;
|
||||
use crate::nostr::builder::Nip34WritePolicy;
|
||||
use crate::nostr::events::{KIND_PR, KIND_PR_UPDATE, KIND_REPOSITORY_ANNOUNCEMENT};
|
||||
use crate::nostr::SharedDatabase;
|
||||
|
||||
mod relay_connection;
|
||||
mod self_subscriber;
|
||||
pub use relay_connection::{RelayConnection, RelayEvent};
|
||||
pub use self_subscriber::{RelayAction, SelfSubscriber};
|
||||
|
||||
// =============================================================================
|
||||
// Type Aliases for Sync State
|
||||
// =============================================================================
|
||||
@@ -176,7 +185,7 @@ pub fn new_sync_relays() -> SyncRelays {
|
||||
/// The SyncManager is responsible for:
|
||||
/// - Discovering relays from stored repository announcements
|
||||
/// - Maintaining connections to sync relays
|
||||
/// - Subscribing to events at external relays
|
||||
/// - Subscribing to events at external relays
|
||||
/// - Applying the acceptance policy to synced events
|
||||
///
|
||||
/// ## Lifecycle
|
||||
@@ -186,14 +195,16 @@ pub fn new_sync_relays() -> SyncRelays {
|
||||
///
|
||||
/// ## Current Status
|
||||
///
|
||||
/// This is a stub implementation. The core data structures are:
|
||||
/// Phase 2 implementation supports:
|
||||
/// - Layer 1 sync: Bootstrap relay connection with 30617/30618 filter
|
||||
/// - Event processing through write policy
|
||||
/// - Storage of accepted events
|
||||
///
|
||||
/// Core data structures:
|
||||
/// - [`FollowingRepoRootEvents`]: Repository root events we're following
|
||||
/// - [`SyncRelays`]: Relays we sync from with their repos and events
|
||||
///
|
||||
/// Full implementation will come in later phases.
|
||||
pub struct SyncManager {
|
||||
/// Bootstrap relay URL if configured
|
||||
#[allow(dead_code)]
|
||||
bootstrap_relay_url: Option<String>,
|
||||
|
||||
/// Our service domain for filtering repo announcements
|
||||
@@ -201,11 +212,9 @@ pub struct SyncManager {
|
||||
service_domain: String,
|
||||
|
||||
/// Database for querying/storing events
|
||||
#[allow(dead_code)]
|
||||
database: SharedDatabase,
|
||||
|
||||
/// Write policy for applying acceptance rules
|
||||
#[allow(dead_code)]
|
||||
write_policy: Nip34WritePolicy,
|
||||
|
||||
/// Repository root events we're following (Phase 1 data structure)
|
||||
@@ -219,6 +228,9 @@ pub struct SyncManager {
|
||||
/// Max backoff duration for relay reconnection
|
||||
#[allow(dead_code)]
|
||||
max_backoff_secs: u64,
|
||||
|
||||
/// Socket address used for sync source (for write policy)
|
||||
sync_source_addr: SocketAddr,
|
||||
}
|
||||
|
||||
impl SyncManager {
|
||||
@@ -238,6 +250,10 @@ impl SyncManager {
|
||||
write_policy: Nip34WritePolicy,
|
||||
config: &Config,
|
||||
) -> Self {
|
||||
// Create a synthetic SocketAddr for sync source identification
|
||||
// This is used when calling write_policy.admit_event() for synced events
|
||||
let sync_source_addr: SocketAddr = "0.0.0.0:0".parse().unwrap();
|
||||
|
||||
Self {
|
||||
bootstrap_relay_url,
|
||||
service_domain,
|
||||
@@ -246,6 +262,7 @@ impl SyncManager {
|
||||
following_repo_root_events: new_following_repo_root_events(),
|
||||
sync_relays: new_sync_relays(),
|
||||
max_backoff_secs: config.sync_max_backoff_secs,
|
||||
sync_source_addr,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -460,14 +477,12 @@ impl SyncManager {
|
||||
/// });
|
||||
/// ```
|
||||
///
|
||||
/// ## Current Status
|
||||
/// ## Implementation Status
|
||||
///
|
||||
/// This is a stub that logs and then waits indefinitely.
|
||||
/// Full implementation includes:
|
||||
/// - Phase 2: Database initialization queries ✓
|
||||
/// - Phase 3: Self-subscription for incremental updates
|
||||
/// - Phase 4-6: Filter building, connection management
|
||||
/// - Phase 7: Full sync loop
|
||||
/// - Phase 2: Layer 1 sync from bootstrap relay ✓
|
||||
/// - Phase 3: Self-subscription and relay discovery ✓
|
||||
/// - Phase 4-6: Filter building, connection management (TODO)
|
||||
/// - Phase 7: Full sync loop (TODO)
|
||||
pub async fn run(self) {
|
||||
tracing::info!(
|
||||
"SyncManager starting (bootstrap_relay={:?}, domain={})",
|
||||
@@ -475,27 +490,285 @@ impl SyncManager {
|
||||
self.service_domain
|
||||
);
|
||||
|
||||
// Phase 2: Initialize from database
|
||||
// Phase 3: Initialize state from database BEFORE spawning connections
|
||||
if let Err(e) = self.initialize_from_database().await {
|
||||
tracing::error!("Failed to initialize sync state from database: {}", e);
|
||||
// Continue anyway - we can still receive events via self-subscription
|
||||
tracing::error!("Failed to initialize from database: {}", e);
|
||||
// Continue anyway - we can still sync from bootstrap
|
||||
}
|
||||
|
||||
// Log initialization results
|
||||
{
|
||||
let following_count = self.following_repo_root_events.read().await.len();
|
||||
let sync_relays_count = self.sync_relays.read().await.len();
|
||||
tracing::info!(
|
||||
"Sync state initialized: {} repos tracked, {} sync relays",
|
||||
following_count,
|
||||
sync_relays_count
|
||||
);
|
||||
// Create channel for relay actions from self-subscriber
|
||||
let (action_tx, mut action_rx) = mpsc::channel::<RelayAction>(100);
|
||||
|
||||
// Construct our own relay URL for self-subscription
|
||||
let own_relay_url = format!("ws://{}", self.service_domain);
|
||||
|
||||
// Spawn self-subscriber task
|
||||
let self_subscriber = SelfSubscriber::new(
|
||||
own_relay_url.clone(),
|
||||
self.service_domain.clone(),
|
||||
Arc::clone(&self.following_repo_root_events),
|
||||
Arc::clone(&self.sync_relays),
|
||||
action_tx,
|
||||
);
|
||||
|
||||
tokio::spawn(async move {
|
||||
self_subscriber.run().await;
|
||||
});
|
||||
|
||||
tracing::info!("SelfSubscriber spawned for {}", own_relay_url);
|
||||
|
||||
// Track active relay connections (relay_url -> event_sender)
|
||||
let mut active_relays: HashMap<String, mpsc::Sender<RelayEvent>> = HashMap::new();
|
||||
|
||||
// Phase 2: Connect to bootstrap relay if configured
|
||||
if let Some(ref bootstrap_url) = self.bootstrap_relay_url {
|
||||
if let Some(event_tx) = self
|
||||
.spawn_relay_connection(bootstrap_url.clone(), None)
|
||||
.await
|
||||
{
|
||||
active_relays.insert(bootstrap_url.clone(), event_tx);
|
||||
}
|
||||
}
|
||||
|
||||
// Stub: wait indefinitely until full implementation (Phases 3-7)
|
||||
// This prevents the spawned task from immediately completing
|
||||
// Main coordination loop
|
||||
loop {
|
||||
tokio::time::sleep(std::time::Duration::from_secs(3600)).await;
|
||||
tokio::select! {
|
||||
// Handle relay actions from self-subscriber
|
||||
action = action_rx.recv() => {
|
||||
match action {
|
||||
Some(RelayAction::SpawnRelay { relay_url, repos_and_root_events }) => {
|
||||
tracing::info!("Spawning new relay connection to {}", relay_url);
|
||||
if !active_relays.contains_key(&relay_url) {
|
||||
if let Some(event_tx) = self.spawn_relay_connection(relay_url.clone(), Some(repos)).await {
|
||||
active_relays.insert(relay_url, event_tx);
|
||||
}
|
||||
}
|
||||
}
|
||||
Some(RelayAction::AddFilters { relay_url, repos_and_new_root_event }) => {
|
||||
tracing::debug!("AddFilters for {} - {} repos (not yet implemented)", relay_url, repos.len());
|
||||
// TODO: Implement filter updates for existing connections
|
||||
}
|
||||
None => {
|
||||
tracing::info!("Action channel closed, continuing without self-subscriber");
|
||||
}
|
||||
}
|
||||
}
|
||||
// Sleep to prevent busy loop when no events
|
||||
_ = tokio::time::sleep(std::time::Duration::from_secs(60)) => {
|
||||
// Periodic maintenance could go here
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Spawn a relay connection with optional Layer 2 filters.
|
||||
///
|
||||
/// Returns the event sender channel if successfully spawned.
|
||||
async fn spawn_relay_connection(
|
||||
&self,
|
||||
relay_url: String,
|
||||
repos: Option<HashMap<String, HashSet<EventId>>>,
|
||||
) -> Option<mpsc::Sender<RelayEvent>> {
|
||||
// Create channel for receiving events
|
||||
let (event_tx, event_rx) = mpsc::channel::<RelayEvent>(100);
|
||||
|
||||
// Create connection
|
||||
let connection = RelayConnection::new(relay_url.clone());
|
||||
|
||||
// Determine if this is bootstrap (no repos) or discovered relay (with repos)
|
||||
let is_bootstrap = repos.is_none();
|
||||
|
||||
match connection.connect_and_subscribe().await {
|
||||
Ok(()) => {
|
||||
if is_bootstrap {
|
||||
tracing::info!("Bootstrap relay connection established: {}", relay_url);
|
||||
} else {
|
||||
tracing::info!(
|
||||
"Discovered relay connection established: {} (with Layer 2 filters)",
|
||||
relay_url
|
||||
);
|
||||
|
||||
// Add Layer 2 subscription for repo events
|
||||
if let Some(ref repos) = repos {
|
||||
if let Err(e) = self.add_layer2_subscription(&connection, repos).await {
|
||||
tracing::warn!("Failed to add Layer 2 subscription: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Clone refs needed for event processing task
|
||||
let database = Arc::clone(&self.database);
|
||||
let write_policy = self.write_policy.clone();
|
||||
let sync_source_addr = self.sync_source_addr;
|
||||
|
||||
// Clone event_tx for the spawned task
|
||||
let event_tx_clone = event_tx.clone();
|
||||
|
||||
// Spawn event loop task
|
||||
let conn_url = relay_url.clone();
|
||||
tokio::spawn(async move {
|
||||
connection.run_event_loop(event_tx_clone).await;
|
||||
});
|
||||
|
||||
// Spawn event processing task
|
||||
tokio::spawn(async move {
|
||||
Self::process_relay_events(
|
||||
event_rx,
|
||||
database,
|
||||
write_policy,
|
||||
sync_source_addr,
|
||||
conn_url,
|
||||
)
|
||||
.await;
|
||||
});
|
||||
|
||||
Some(event_tx)
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to connect to relay {}: {}", relay_url, e);
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Add Layer 2 subscription for repo-related events.
|
||||
///
|
||||
/// Layer 2 filters subscribe to events with 'a' tags referencing repos we track.
|
||||
async fn add_layer2_subscription(
|
||||
&self,
|
||||
connection: &RelayConnection,
|
||||
repos: &HashMap<String, HashSet<EventId>>,
|
||||
) -> Result<(), String> {
|
||||
if repos.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Build repo refs list for filter
|
||||
let repo_refs: Vec<String> = repos.keys().cloned().collect();
|
||||
|
||||
tracing::debug!(
|
||||
"Adding Layer 2 subscription for {} repos to {}",
|
||||
repo_refs.len(),
|
||||
connection.url()
|
||||
);
|
||||
|
||||
// Chunk repo_refs into groups of 100 (per plan)
|
||||
for chunk in repo_refs.chunks(100) {
|
||||
// Build filter with lowercase 'a' tag for each repo ref
|
||||
let mut filter = Filter::new().kinds([
|
||||
Kind::GitPatch, // 1617
|
||||
Kind::Custom(1618), // PR
|
||||
Kind::Custom(1619), // PR update
|
||||
Kind::GitIssue, // 1621
|
||||
]);
|
||||
|
||||
// Add each repo ref as a custom tag filter
|
||||
for repo_ref in chunk {
|
||||
filter =
|
||||
filter.custom_tag(SingleLetterTag::lowercase(Alphabet::A), repo_ref.clone());
|
||||
}
|
||||
|
||||
// Subscribe to this filter
|
||||
if let Err(e) = connection.subscribe_filter(filter).await {
|
||||
return Err(format!("Failed to subscribe with Layer 2 filter: {}", e));
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Process events from a single relay connection.
|
||||
///
|
||||
/// This is a static method that runs in its own task.
|
||||
async fn process_relay_events(
|
||||
mut event_rx: mpsc::Receiver<RelayEvent>,
|
||||
database: SharedDatabase,
|
||||
write_policy: Nip34WritePolicy,
|
||||
sync_source_addr: SocketAddr,
|
||||
relay_url: String,
|
||||
) {
|
||||
tracing::debug!("Starting event processing for relay: {}", relay_url);
|
||||
|
||||
while let Some(relay_event) = event_rx.recv().await {
|
||||
match relay_event {
|
||||
RelayEvent::Event(event) => {
|
||||
Self::process_single_event_static(
|
||||
&event,
|
||||
&database,
|
||||
&write_policy,
|
||||
&sync_source_addr,
|
||||
&relay_url,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
RelayEvent::EndOfStoredEvents => {
|
||||
tracing::debug!("EOSE received from {}", relay_url);
|
||||
}
|
||||
RelayEvent::Closed(reason) => {
|
||||
tracing::warn!("Connection to {} closed: {}", relay_url, reason);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
tracing::info!("Event processing ended for relay: {}", relay_url);
|
||||
}
|
||||
|
||||
/// Process a single event (static version for use in spawned tasks).
|
||||
async fn process_single_event_static(
|
||||
event: &Event,
|
||||
database: &SharedDatabase,
|
||||
write_policy: &Nip34WritePolicy,
|
||||
sync_source_addr: &SocketAddr,
|
||||
relay_url: &str,
|
||||
) {
|
||||
let event_id = event.id;
|
||||
let kind = event.kind.as_u16();
|
||||
|
||||
// Check if event already exists in database
|
||||
match database.check_id(&event_id).await {
|
||||
Ok(DatabaseEventStatus::Saved) | Ok(DatabaseEventStatus::Deleted) => {
|
||||
tracing::trace!("Event {} already exists, skipping", event_id);
|
||||
return;
|
||||
}
|
||||
Ok(DatabaseEventStatus::NotExistent) => {} // Continue processing
|
||||
Err(e) => {
|
||||
tracing::warn!("Failed to check if event {} exists: {}", event_id, e);
|
||||
}
|
||||
}
|
||||
|
||||
// Pass through write policy
|
||||
let policy_result = write_policy.admit_event(event, sync_source_addr).await;
|
||||
|
||||
match policy_result {
|
||||
PolicyResult::Accept => match database.save_event(event).await {
|
||||
Ok(SaveEventStatus::Success) => {
|
||||
tracing::info!(
|
||||
"Synced event {} (kind {}) from {}",
|
||||
event_id,
|
||||
kind,
|
||||
relay_url
|
||||
);
|
||||
}
|
||||
Ok(_) => {
|
||||
tracing::debug!(
|
||||
"Event {} (kind {}) already stored or rejected by database",
|
||||
event_id,
|
||||
kind
|
||||
);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to save synced event {}: {}", event_id, e);
|
||||
}
|
||||
},
|
||||
PolicyResult::Reject(reason) => {
|
||||
tracing::debug!(
|
||||
"Rejected synced event {} (kind {}): {}",
|
||||
event_id,
|
||||
kind,
|
||||
reason
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,185 @@
|
||||
//! Relay Connection for Proactive Sync
|
||||
//!
|
||||
//! This module handles connecting to external relays and receiving events
|
||||
//! for the proactive sync system.
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
use nostr_sdk::prelude::*;
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
use crate::nostr::events::{KIND_REPOSITORY_ANNOUNCEMENT, KIND_REPOSITORY_STATE};
|
||||
|
||||
/// Events received from a relay connection
|
||||
#[derive(Debug)]
|
||||
pub enum RelayEvent {
|
||||
/// A nostr event was received
|
||||
Event(Event),
|
||||
/// End of stored events (EOSE) received
|
||||
EndOfStoredEvents,
|
||||
/// Connection was closed
|
||||
Closed(String),
|
||||
}
|
||||
|
||||
/// Connection to an external relay for syncing events.
|
||||
///
|
||||
/// RelayConnection handles:
|
||||
/// - Connecting to the relay
|
||||
/// - Subscribing with appropriate filters (Layer 1 for bootstrap)
|
||||
/// - Receiving events and sending them through a channel
|
||||
pub struct RelayConnection {
|
||||
/// The relay URL
|
||||
url: String,
|
||||
/// The nostr-sdk client
|
||||
client: Client,
|
||||
}
|
||||
|
||||
impl RelayConnection {
|
||||
/// Create a new relay connection.
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `url` - The WebSocket URL of the relay to connect to
|
||||
pub fn new(url: String) -> Self {
|
||||
// Create a client with generated keys (we're just subscribing, not publishing)
|
||||
let keys = Keys::generate();
|
||||
let client = Client::new(keys);
|
||||
|
||||
Self { url, client }
|
||||
}
|
||||
|
||||
/// Connect to the relay and subscribe with Layer 1 filter.
|
||||
///
|
||||
/// Layer 1 filter syncs announcement events (30617, 30618) which are
|
||||
/// the foundation for discovering repository relationships.
|
||||
///
|
||||
/// Returns the notification stream for receiving events.
|
||||
pub async fn connect_and_subscribe(&self) -> Result<(), String> {
|
||||
// Add the relay
|
||||
self.client
|
||||
.add_relay(&self.url)
|
||||
.await
|
||||
.map_err(|e| format!("Failed to add relay {}: {}", self.url, e))?;
|
||||
|
||||
// Connect to relay
|
||||
self.client.connect().await;
|
||||
|
||||
// Wait for connection to establish
|
||||
let mut connected = false;
|
||||
for _ in 0..30 {
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let relays = self.client.relays().await;
|
||||
if relays.values().any(|r| r.is_connected()) {
|
||||
connected = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if !connected {
|
||||
return Err(format!(
|
||||
"Failed to connect to relay {} after 3 seconds",
|
||||
self.url
|
||||
));
|
||||
}
|
||||
|
||||
tracing::info!("Connected to bootstrap relay: {}", self.url);
|
||||
|
||||
// Layer 1 filter: Repository announcements and state events
|
||||
// These are addressable events that define repositories
|
||||
let filter = Filter::new().kinds([
|
||||
Kind::Custom(KIND_REPOSITORY_ANNOUNCEMENT), // 30617
|
||||
Kind::Custom(KIND_REPOSITORY_STATE), // 30618
|
||||
]);
|
||||
|
||||
// Subscribe to the filter
|
||||
self.client
|
||||
.subscribe(filter, None)
|
||||
.await
|
||||
.map_err(|e| format!("Failed to subscribe: {}", e))?;
|
||||
|
||||
tracing::debug!(
|
||||
"Subscribed to Layer 1 events (kinds 30617, 30618) from {}",
|
||||
self.url
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Run the event loop, sending received events through the channel.
|
||||
///
|
||||
/// This method runs until the connection is closed or an error occurs.
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `event_sender` - Channel to send received events
|
||||
pub async fn run_event_loop(self, event_sender: mpsc::Sender<RelayEvent>) {
|
||||
tracing::debug!("Starting event loop for relay: {}", self.url);
|
||||
|
||||
// Handle notifications
|
||||
self.client
|
||||
.handle_notifications(|notification| async {
|
||||
match notification {
|
||||
RelayPoolNotification::Event { event, .. } => {
|
||||
tracing::debug!(
|
||||
"Received event {} (kind {}) from {}",
|
||||
event.id,
|
||||
event.kind.as_u16(),
|
||||
self.url
|
||||
);
|
||||
if event_sender.send(RelayEvent::Event(*event)).await.is_err() {
|
||||
tracing::warn!("Event channel closed, stopping relay connection");
|
||||
return Ok(true); // Stop handling
|
||||
}
|
||||
}
|
||||
RelayPoolNotification::Message { message, .. } => {
|
||||
if let RelayMessage::EndOfStoredEvents(_) = message {
|
||||
tracing::debug!("EOSE received from {}", self.url);
|
||||
if event_sender
|
||||
.send(RelayEvent::EndOfStoredEvents)
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
return Ok(true); // Stop handling
|
||||
}
|
||||
}
|
||||
}
|
||||
RelayPoolNotification::Shutdown => {
|
||||
tracing::info!("Relay {} shutting down", self.url);
|
||||
let _ = event_sender
|
||||
.send(RelayEvent::Closed("Shutdown".to_string()))
|
||||
.await;
|
||||
return Ok(true); // Stop handling
|
||||
}
|
||||
}
|
||||
Ok(false) // Continue handling
|
||||
})
|
||||
.await
|
||||
.ok(); // Ignore errors on shutdown
|
||||
|
||||
// Disconnect when done
|
||||
self.client.disconnect().await;
|
||||
tracing::info!("Disconnected from relay: {}", self.url);
|
||||
}
|
||||
|
||||
/// Get the relay URL
|
||||
pub fn url(&self) -> &str {
|
||||
&self.url
|
||||
}
|
||||
|
||||
/// Subscribe to an additional filter.
|
||||
///
|
||||
/// This is used to add Layer 2 filters for repo-related events after
|
||||
/// the initial connection is established.
|
||||
pub async fn subscribe_filter(&self, filter: Filter) -> Result<(), String> {
|
||||
self.client
|
||||
.subscribe(filter, None)
|
||||
.await
|
||||
.map_err(|e| format!("Failed to subscribe with filter: {}", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Get a reference to the client for additional operations.
|
||||
pub fn client(&self) -> &Client {
|
||||
&self.client
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,497 @@
|
||||
//! Self-Subscriber for Proactive Sync
|
||||
//!
|
||||
//! This module handles subscribing to our own relay to detect new events
|
||||
//! and trigger relay discovery from announcements.
|
||||
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::time::Duration;
|
||||
|
||||
use nostr_sdk::prelude::*;
|
||||
use tokio::sync::mpsc;
|
||||
use tokio::time::Instant;
|
||||
|
||||
use crate::nostr::events::{KIND_PR, KIND_PR_UPDATE, KIND_REPOSITORY_ANNOUNCEMENT};
|
||||
|
||||
use super::{FollowingRepoRootEvents, SyncManager, SyncRelays};
|
||||
|
||||
// =============================================================================
|
||||
// Types
|
||||
// =============================================================================
|
||||
|
||||
/// Actions to be taken by the SyncManager based on self-subscription events.
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum RelayAction {
|
||||
/// Spawn a new relay connection to sync from.
|
||||
/// Contains: relay_url, map of repo_refs to their event IDs for Layer 2 filtering.
|
||||
SpawnRelay {
|
||||
relay_url: String,
|
||||
repos_and_root_events: HashMap<String, HashSet<EventId>>,
|
||||
},
|
||||
/// Add filters to an existing relay connection.
|
||||
/// Contains: relay_url, additional repos to add.
|
||||
AddFilters {
|
||||
relay_url: String,
|
||||
repos_and_new_root_event: HashMap<String, HashSet<EventId>>,
|
||||
},
|
||||
}
|
||||
|
||||
/// Pending updates collected during batch window.
|
||||
#[derive(Debug, Default)]
|
||||
struct PendingUpdates {
|
||||
/// New announcements (kind 30617) - triggers relay discovery
|
||||
announcements: Vec<Event>,
|
||||
/// New root events (kinds 1617, 1618, 1619, 1621) - updates following set
|
||||
root_events: Vec<Event>,
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// SelfSubscriber
|
||||
// =============================================================================
|
||||
|
||||
/// Subscribes to our own relay to detect new events.
|
||||
///
|
||||
/// The self-subscriber:
|
||||
/// 1. Connects to our own relay
|
||||
/// 2. Subscribes to kinds 30617, 1617, 1618, 1619, 1621 (NOT 30618)
|
||||
/// 3. When events arrive, batches them
|
||||
/// 4. On batch timer fire, processes updates and sends relay actions
|
||||
pub struct SelfSubscriber {
|
||||
/// URL of our own relay to subscribe to
|
||||
own_relay_url: String,
|
||||
/// Our relay domain for checking if announcements list us
|
||||
relay_domain: String,
|
||||
/// Reference to following repo root events (shared with SyncManager)
|
||||
following_repo_root_events: FollowingRepoRootEvents,
|
||||
/// Reference to sync relays (shared with SyncManager)
|
||||
sync_relays: SyncRelays,
|
||||
/// Channel to send relay actions back to manager
|
||||
action_tx: mpsc::Sender<RelayAction>,
|
||||
}
|
||||
|
||||
impl SelfSubscriber {
|
||||
/// Create a new self-subscriber.
|
||||
pub fn new(
|
||||
own_relay_url: String,
|
||||
relay_domain: String,
|
||||
following_repo_root_events: FollowingRepoRootEvents,
|
||||
sync_relays: SyncRelays,
|
||||
action_tx: mpsc::Sender<RelayAction>,
|
||||
) -> Self {
|
||||
Self {
|
||||
own_relay_url,
|
||||
relay_domain,
|
||||
following_repo_root_events,
|
||||
sync_relays,
|
||||
action_tx,
|
||||
}
|
||||
}
|
||||
|
||||
/// Get the batch window duration from environment variable.
|
||||
///
|
||||
/// Default is 5 seconds, but can be overridden via NGIT_SYNC_BATCH_WINDOW_MS
|
||||
/// for faster tests (typically 200ms).
|
||||
fn get_batch_window() -> Duration {
|
||||
std::env::var("NGIT_SYNC_BATCH_WINDOW_MS")
|
||||
.ok()
|
||||
.and_then(|s| s.parse().ok())
|
||||
.map(Duration::from_millis)
|
||||
.unwrap_or(Duration::from_secs(5))
|
||||
}
|
||||
|
||||
/// Run the self-subscriber event loop.
|
||||
///
|
||||
/// This method:
|
||||
/// 1. Connects to our own relay
|
||||
/// 2. Subscribes to relevant event kinds
|
||||
/// 3. Receives events and batches them
|
||||
/// 4. On batch timer fire, processes and sends relay actions
|
||||
pub async fn run(self) {
|
||||
tracing::info!("SelfSubscriber starting for {}", self.own_relay_url);
|
||||
|
||||
// Create nostr-sdk client
|
||||
let keys = Keys::generate();
|
||||
let client = Client::new(keys);
|
||||
|
||||
// Connect to our own relay
|
||||
if let Err(e) = client.add_relay(&self.own_relay_url).await {
|
||||
tracing::error!("Failed to add own relay {}: {}", self.own_relay_url, e);
|
||||
return;
|
||||
}
|
||||
|
||||
client.connect().await;
|
||||
|
||||
// Wait for connection
|
||||
let mut connected = false;
|
||||
for _ in 0..30 {
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let relays = client.relays().await;
|
||||
if relays.values().any(|r| r.is_connected()) {
|
||||
connected = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if !connected {
|
||||
tracing::error!(
|
||||
"Failed to connect to own relay {} after 3 seconds",
|
||||
self.own_relay_url
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
tracing::info!("SelfSubscriber connected to {}", self.own_relay_url);
|
||||
|
||||
// Subscribe to kinds 30617, 1617, 1618, 1619, 1621 (NOT 30618 per v2 design)
|
||||
let filter = Filter::new()
|
||||
.kinds([
|
||||
Kind::Custom(KIND_REPOSITORY_ANNOUNCEMENT), // 30617
|
||||
Kind::GitPatch, // 1617
|
||||
Kind::Custom(KIND_PR), // 1618
|
||||
Kind::Custom(KIND_PR_UPDATE), // 1619
|
||||
Kind::GitIssue, // 1621
|
||||
])
|
||||
.since(Timestamp::now());
|
||||
|
||||
if let Err(e) = client.subscribe(filter, None).await {
|
||||
tracing::error!("Failed to subscribe to own relay: {}", e);
|
||||
return;
|
||||
}
|
||||
|
||||
tracing::info!("SelfSubscriber subscribed to event kinds on own relay");
|
||||
|
||||
// Batch state
|
||||
let mut pending = PendingUpdates::default();
|
||||
let mut batch_timer_started: Option<Instant> = None;
|
||||
let batch_window = Self::get_batch_window();
|
||||
|
||||
// Main event loop using notifications stream
|
||||
loop {
|
||||
// Calculate timeout for batch processing
|
||||
let timeout = if let Some(started) = batch_timer_started {
|
||||
let elapsed = started.elapsed();
|
||||
if elapsed >= batch_window {
|
||||
Duration::ZERO
|
||||
} else {
|
||||
batch_window - elapsed
|
||||
}
|
||||
} else {
|
||||
Duration::from_secs(60) // Long timeout when no batch pending
|
||||
};
|
||||
|
||||
// Wait for notification with timeout
|
||||
let notification = tokio::time::timeout(timeout, client.notifications().recv()).await;
|
||||
|
||||
match notification {
|
||||
Ok(Ok(notification)) => {
|
||||
match notification {
|
||||
RelayPoolNotification::Event { event, .. } => {
|
||||
let kind = event.kind.as_u16();
|
||||
|
||||
// Start batch timer on first event (does NOT reset)
|
||||
if batch_timer_started.is_none() {
|
||||
batch_timer_started = Some(Instant::now());
|
||||
tracing::debug!("Batch timer started");
|
||||
}
|
||||
|
||||
// Classify and add to pending
|
||||
if kind == KIND_REPOSITORY_ANNOUNCEMENT {
|
||||
tracing::debug!(
|
||||
"SelfSubscriber received announcement {}",
|
||||
event.id
|
||||
);
|
||||
pending.announcements.push(*event);
|
||||
} else {
|
||||
tracing::debug!(
|
||||
"SelfSubscriber received root event {} (kind {})",
|
||||
event.id,
|
||||
kind
|
||||
);
|
||||
pending.root_events.push(*event);
|
||||
}
|
||||
}
|
||||
RelayPoolNotification::Message { message, .. } => {
|
||||
if let RelayMessage::EndOfStoredEvents(_) = message {
|
||||
tracing::debug!("SelfSubscriber EOSE received");
|
||||
// Process any pending events after EOSE
|
||||
if !pending.announcements.is_empty()
|
||||
|| !pending.root_events.is_empty()
|
||||
{
|
||||
self.process_batch(&mut pending).await;
|
||||
batch_timer_started = None;
|
||||
}
|
||||
}
|
||||
}
|
||||
RelayPoolNotification::Shutdown => {
|
||||
tracing::info!("SelfSubscriber shutting down");
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(Err(_)) => {
|
||||
// Channel closed
|
||||
tracing::warn!("SelfSubscriber notification channel closed");
|
||||
break;
|
||||
}
|
||||
Err(_) => {
|
||||
// Timeout - check if batch should be processed
|
||||
if let Some(started) = batch_timer_started {
|
||||
if started.elapsed() >= batch_window {
|
||||
if !pending.announcements.is_empty() || !pending.root_events.is_empty()
|
||||
{
|
||||
self.process_batch(&mut pending).await;
|
||||
}
|
||||
batch_timer_started = None;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
client.disconnect().await;
|
||||
tracing::info!("SelfSubscriber disconnected");
|
||||
}
|
||||
|
||||
/// Process a batch of pending updates.
|
||||
async fn process_batch(&self, pending: &mut PendingUpdates) {
|
||||
tracing::debug!(
|
||||
"Processing batch: {} announcements, {} root events",
|
||||
pending.announcements.len(),
|
||||
pending.root_events.len()
|
||||
);
|
||||
|
||||
// Process root events first (update following_repo_root_events)
|
||||
for event in pending.root_events.drain(..) {
|
||||
let repo_refs = SyncManager::extract_all_repo_refs(&event);
|
||||
if !repo_refs.is_empty() {
|
||||
let mut guard = self.following_repo_root_events.write().await;
|
||||
for repo_ref in repo_refs {
|
||||
guard.entry(repo_ref).or_default().insert(event.id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Process announcements (relay discovery)
|
||||
for event in pending.announcements.drain(..) {
|
||||
self.process_announcement(&event).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Process an announcement event for relay discovery.
|
||||
async fn process_announcement(&self, event: &Event) {
|
||||
let repo_ref = SyncManager::build_repo_ref(event);
|
||||
let relay_urls = Self::extract_relay_urls_from_announcement(event);
|
||||
|
||||
// Check if this announcement lists our relay
|
||||
if !self.lists_our_service(event) {
|
||||
tracing::debug!(
|
||||
"Announcement {} does not list our service, skipping relay discovery",
|
||||
event.id
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
tracing::info!(
|
||||
"Processing announcement {} for repo {}, found {} relay URLs",
|
||||
event.id,
|
||||
repo_ref,
|
||||
relay_urls.len()
|
||||
);
|
||||
|
||||
// Get current events for this repo from following_repo_root_events
|
||||
let events = self
|
||||
.following_repo_root_events
|
||||
.read()
|
||||
.await
|
||||
.get(&repo_ref)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
|
||||
// For each relay URL in the announcement, check if we need to spawn or update
|
||||
for relay_url in relay_urls {
|
||||
if self.is_own_relay(&relay_url) {
|
||||
continue; // Skip our own relay
|
||||
}
|
||||
|
||||
let sync_relays_guard = self.sync_relays.read().await;
|
||||
let exists = sync_relays_guard.contains_key(&relay_url);
|
||||
drop(sync_relays_guard);
|
||||
|
||||
if exists {
|
||||
// Relay already known - check if we need to add this repo
|
||||
let mut guard = self.sync_relays.write().await;
|
||||
let relay_repos = guard.entry(relay_url.clone()).or_default();
|
||||
let is_new_repo = !relay_repos.contains_key(&repo_ref);
|
||||
|
||||
if is_new_repo {
|
||||
relay_repos.insert(repo_ref.clone(), events.clone());
|
||||
drop(guard);
|
||||
|
||||
// Send action to add filters
|
||||
let mut repos_filters = HashMap::new();
|
||||
repos_filters.insert(repo_ref.clone(), events.clone());
|
||||
|
||||
if let Err(e) = self
|
||||
.action_tx
|
||||
.send(RelayAction::AddFilters {
|
||||
relay_url: relay_url.clone(),
|
||||
repos_and_new_root_event: repos_filters,
|
||||
})
|
||||
.await
|
||||
{
|
||||
tracing::warn!("Failed to send AddFilters action: {}", e);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// New relay - add to sync_relays and spawn
|
||||
let mut guard = self.sync_relays.write().await;
|
||||
let mut repos = HashMap::new();
|
||||
repos.insert(repo_ref.clone(), events.clone());
|
||||
guard.insert(relay_url.clone(), repos.clone());
|
||||
drop(guard);
|
||||
|
||||
tracing::info!("Discovered new relay to sync from: {}", relay_url);
|
||||
|
||||
// Send action to spawn relay
|
||||
if let Err(e) = self
|
||||
.action_tx
|
||||
.send(RelayAction::SpawnRelay {
|
||||
relay_url: relay_url.clone(),
|
||||
repos_and_root_events: repos,
|
||||
})
|
||||
.await
|
||||
{
|
||||
tracing::warn!("Failed to send SpawnRelay action: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Extract relay URLs from an announcement event.
|
||||
///
|
||||
/// Looks for both 'relays' and 'clone' tags.
|
||||
fn extract_relay_urls_from_announcement(event: &Event) -> Vec<String> {
|
||||
let mut urls = Vec::new();
|
||||
|
||||
// Extract from 'relays' tag
|
||||
for tag in event.tags.iter() {
|
||||
if matches!(tag.kind(), TagKind::Relays) {
|
||||
let vec = tag.clone().to_vec();
|
||||
urls.extend(vec.into_iter().skip(1)); // Skip tag name
|
||||
}
|
||||
}
|
||||
|
||||
// Extract from 'clone' tag - parse URLs to get relay hints
|
||||
// Clone URLs look like: http://domain/repo.git or git://domain/repo.git
|
||||
// We want to construct ws://domain from these
|
||||
for tag in event.tags.iter() {
|
||||
if matches!(tag.kind(), TagKind::Clone) {
|
||||
let vec = tag.clone().to_vec();
|
||||
for url in vec.into_iter().skip(1) {
|
||||
if let Some(relay_url) = Self::clone_url_to_relay_url(&url) {
|
||||
if !urls.contains(&relay_url) {
|
||||
urls.push(relay_url);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
urls
|
||||
}
|
||||
|
||||
/// Convert a clone URL to a potential relay URL.
|
||||
///
|
||||
/// E.g., "http://127.0.0.1:8080/repo.git" -> "ws://127.0.0.1:8080"
|
||||
fn clone_url_to_relay_url(clone_url: &str) -> Option<String> {
|
||||
// Parse the URL to extract host:port
|
||||
if let Ok(url) = url::Url::parse(clone_url) {
|
||||
let host = url.host_str()?;
|
||||
let port = url.port();
|
||||
let scheme = if url.scheme() == "https" { "wss" } else { "ws" };
|
||||
|
||||
if let Some(port) = port {
|
||||
Some(format!("{}://{}:{}", scheme, host, port))
|
||||
} else {
|
||||
Some(format!("{}://{}", scheme, host))
|
||||
}
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
/// Check if event lists our service in the relays or clone tags.
|
||||
fn lists_our_service(&self, event: &Event) -> bool {
|
||||
// Check relays tag
|
||||
for tag in event.tags.iter() {
|
||||
if matches!(tag.kind(), TagKind::Relays) {
|
||||
let vec = tag.clone().to_vec();
|
||||
for url in vec.into_iter().skip(1) {
|
||||
if self.is_own_relay(&url) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Check clone tag
|
||||
for tag in event.tags.iter() {
|
||||
if matches!(tag.kind(), TagKind::Clone) {
|
||||
let vec = tag.clone().to_vec();
|
||||
for url in vec.into_iter().skip(1) {
|
||||
if url.contains(&self.relay_domain) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
false
|
||||
}
|
||||
|
||||
/// Check if a relay URL matches our relay.
|
||||
fn is_own_relay(&self, relay_url: &str) -> bool {
|
||||
relay_url.contains(&self.relay_domain)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_clone_url_to_relay_url_http() {
|
||||
let url = "http://127.0.0.1:8080/repo.git";
|
||||
let relay = SelfSubscriber::clone_url_to_relay_url(url);
|
||||
assert_eq!(relay, Some("ws://127.0.0.1:8080".to_string()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_clone_url_to_relay_url_https() {
|
||||
let url = "https://example.com/repo.git";
|
||||
let relay = SelfSubscriber::clone_url_to_relay_url(url);
|
||||
assert_eq!(relay, Some("wss://example.com".to_string()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_clone_url_to_relay_url_invalid() {
|
||||
let url = "not-a-valid-url";
|
||||
let relay = SelfSubscriber::clone_url_to_relay_url(url);
|
||||
assert_eq!(relay, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_get_batch_window_default() {
|
||||
// Clear env var if set
|
||||
std::env::remove_var("NGIT_SYNC_BATCH_WINDOW_MS");
|
||||
let window = SelfSubscriber::get_batch_window();
|
||||
assert_eq!(window, Duration::from_secs(5));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_get_batch_window_from_env() {
|
||||
std::env::set_var("NGIT_SYNC_BATCH_WINDOW_MS", "200");
|
||||
let window = SelfSubscriber::get_batch_window();
|
||||
assert_eq!(window, Duration::from_millis(200));
|
||||
std::env::remove_var("NGIT_SYNC_BATCH_WINDOW_MS");
|
||||
}
|
||||
}
|
||||
@@ -93,7 +93,7 @@ impl TestRelay {
|
||||
.env("NGIT_GIT_DATA_PATH", git_data_dir.path())
|
||||
.env("NGIT_DATABASE_BACKEND", "memory") // Force in-memory database for isolation
|
||||
.env("NGIT_OWNER_NPUB", &test_npub)
|
||||
.env("NGIT_SYNC_STARTUP_JITTER_MS", "0") // Disable jitter for tests
|
||||
.env("NGIT_SYNC_BATCH_WINDOW_MS", "200") // Fast batch window for tests (200ms instead of 5s default)
|
||||
.env("RUST_LOG", "warn") // Less logging during tests
|
||||
.stdout(Stdio::null())
|
||||
.stderr(Stdio::null()); // Disable stderr for cleaner test output
|
||||
|
||||
Reference in New Issue
Block a user