mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
feat(sync): broadcast synced events to WebSocket subscribers
Enable recursive relay discovery by broadcasting synced events to WebSocket subscribers via LocalRelay.notify_event(). This allows the SelfSubscriber to receive 30617 announcements synced from external relays and discover additional relay URLs to connect to. Changes: - Pass LocalRelay to SyncManager::new() from main.rs - Add local_relay field to SyncManager struct - Call notify_event() after saving synced events to database - Enable test_recursive_relay_discovery_syncs_announcement test The test verifies that when relay_a syncs announcement_x from bootstrap relay_b (which lists relay_c), relay_a discovers and connects to relay_c to sync announcement_y. Fixes recursive relay discovery from bootstrap sync.
This commit is contained in:
@@ -58,6 +58,7 @@ async fn main() -> Result<()> {
|
||||
config.domain.clone(),
|
||||
relay_with_db.database.clone(),
|
||||
relay_with_db.write_policy.clone(),
|
||||
relay_with_db.relay.clone(),
|
||||
&config,
|
||||
);
|
||||
|
||||
|
||||
+29
-8
@@ -40,6 +40,7 @@ use tokio::sync::{broadcast, Mutex, RwLock};
|
||||
|
||||
use crate::config::Config;
|
||||
use crate::nostr::builder::{Nip34WritePolicy, SharedDatabase};
|
||||
use nostr_relay_builder::prelude::LocalRelay;
|
||||
|
||||
// =============================================================================
|
||||
// Type Aliases for Index Structures
|
||||
@@ -467,6 +468,8 @@ pub struct SyncManager {
|
||||
database: SharedDatabase,
|
||||
/// Write policy for validating incoming events
|
||||
write_policy: Nip34WritePolicy,
|
||||
/// Local relay for submitting synced events (enables broadcast to WebSocket subscribers)
|
||||
local_relay: LocalRelay,
|
||||
/// Configuration reference for sync settings
|
||||
config: Config,
|
||||
/// What we want to sync (source of truth)
|
||||
@@ -499,12 +502,14 @@ impl SyncManager {
|
||||
/// * `service_domain` - The domain this relay serves (for filtering repos)
|
||||
/// * `database` - Shared database for event storage
|
||||
/// * `write_policy` - Policy for validating events before storage
|
||||
/// * `local_relay` - Local relay for submitting synced events (enables WebSocket broadcast)
|
||||
/// * `config` - Configuration for sync settings
|
||||
pub fn new(
|
||||
bootstrap_relay_url: Option<String>,
|
||||
service_domain: String,
|
||||
database: SharedDatabase,
|
||||
write_policy: Nip34WritePolicy,
|
||||
local_relay: LocalRelay,
|
||||
config: &Config,
|
||||
) -> Self {
|
||||
Self {
|
||||
@@ -512,6 +517,7 @@ impl SyncManager {
|
||||
service_domain,
|
||||
database,
|
||||
write_policy,
|
||||
local_relay,
|
||||
config: config.clone(),
|
||||
repo_sync_index: Arc::new(RwLock::new(HashMap::new())),
|
||||
relay_sync_index: Arc::new(RwLock::new(HashMap::new())),
|
||||
@@ -1269,6 +1275,7 @@ impl SyncManager {
|
||||
|
||||
let database = Arc::clone(&self.database);
|
||||
let write_policy = self.write_policy.clone();
|
||||
let local_relay = self.local_relay.clone();
|
||||
let relay_sync_index = Arc::clone(&self.relay_sync_index);
|
||||
|
||||
// Check if this is a bootstrap relay
|
||||
@@ -1340,6 +1347,7 @@ impl SyncManager {
|
||||
&relay_url_clone,
|
||||
&database,
|
||||
&write_policy,
|
||||
&local_relay,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
@@ -1393,11 +1401,18 @@ impl SyncManager {
|
||||
}
|
||||
|
||||
/// Process a single event from a relay (static version for spawned tasks)
|
||||
///
|
||||
/// Processes events with dedup, policy check, database save, and broadcast:
|
||||
/// - Deduplication (skips if event already exists)
|
||||
/// - Write policy validation
|
||||
/// - Database save
|
||||
/// - Broadcast to WebSocket subscribers via notify_event (enables recursive relay discovery)
|
||||
async fn process_event_static(
|
||||
event: &Event,
|
||||
relay_url: &str,
|
||||
database: &SharedDatabase,
|
||||
write_policy: &Nip34WritePolicy,
|
||||
local_relay: &LocalRelay,
|
||||
) {
|
||||
use nostr_relay_builder::prelude::{PolicyResult, WritePolicy};
|
||||
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
|
||||
@@ -1421,7 +1436,7 @@ impl SyncManager {
|
||||
|
||||
match result {
|
||||
PolicyResult::Accept => {
|
||||
// Save event
|
||||
// Save event to database
|
||||
if let Err(e) = database.save_event(event).await {
|
||||
tracing::error!(
|
||||
event_id = %event.id,
|
||||
@@ -1429,14 +1444,20 @@ impl SyncManager {
|
||||
error = %e,
|
||||
"Failed to save synced event"
|
||||
);
|
||||
} else {
|
||||
tracing::debug!(
|
||||
event_id = %event.id,
|
||||
relay = %relay_url,
|
||||
kind = %event.kind.as_u16(),
|
||||
"Saved synced event"
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
// Broadcast to WebSocket subscribers (enables recursive relay discovery)
|
||||
// This allows SelfSubscriber to receive synced 30617 announcements
|
||||
let broadcast_success = local_relay.notify_event(event.clone());
|
||||
|
||||
tracing::debug!(
|
||||
event_id = %event.id,
|
||||
relay = %relay_url,
|
||||
kind = %event.kind.as_u16(),
|
||||
broadcast = broadcast_success,
|
||||
"Synced event saved and broadcast"
|
||||
);
|
||||
}
|
||||
PolicyResult::Reject(reason) => {
|
||||
tracing::debug!(
|
||||
|
||||
@@ -293,12 +293,7 @@ async fn test_layer2_discovery_with_chain() {
|
||||
/// 3. Discovers and connects to relay_c
|
||||
/// 4. Syncs announcement_y from relay_c
|
||||
///
|
||||
/// NOTE: This test is ignored because recursive relay discovery from synced
|
||||
/// announcements is not yet implemented. Currently, discovery only triggers
|
||||
/// when an announcement is directly submitted to a relay, not when it's
|
||||
/// synced from a bootstrap relay.
|
||||
#[tokio::test]
|
||||
#[ignore = "Recursive relay discovery from bootstrap sync not yet implemented"]
|
||||
async fn test_recursive_relay_discovery_syncs_announcement() {
|
||||
// 1. Start all three relays
|
||||
|
||||
|
||||
Reference in New Issue
Block a user