diff --git a/src/sync/self_subscriber.rs b/src/sync/self_subscriber.rs index 5c43754..08b41e5 100644 --- a/src/sync/self_subscriber.rs +++ b/src/sync/self_subscriber.rs @@ -558,49 +558,42 @@ impl SelfSubscriber { return; }; - // Look up repo in repo_sync_index - add root event directly and also to pending + // A live root can follow its announcement before the batch timer has + // published that repository. Resolve against the pending announcement + // as well as the committed index so that arrival order cannot lose it. let mut index = self.repo_sync_index.write().await; - if let Some(repo_sync) = index.get_mut(&repo_ref) { - // Add event.id to root_events set in the index (immediate availability) + let relays = if let Some(repo_sync) = index.get_mut(&repo_ref) { repo_sync.root_events.insert(event.id); - - // Clone the relays before releasing the lock - Layer 3 filters need to be - // sent to the same relays as Layer 2 filters for this repo - let relays = repo_sync.relays.clone(); - - // Release lock before modifying pending - drop(index); - - // Keep compact metadata even while the repository is StateOnly. - // Discovery consults the current RepoSyncIndex level, so promotion - // activates an already-observed root without a database rescan. - self.root_candidate_index.write().await.insert( - event.id, - super::discovery::AcceptedRoot { - id: event.id, - author: event.pubkey, - repository: repo_ref.clone(), - }, - ); - - // Also add root event to pending - this ensures batch processing runs - // and creates Layer 3 filters for events referencing this root event. - // CRITICAL: Include relays so derive_relay_targets knows where to send filters! - pending.add_root_event(repo_ref.clone(), relays.clone(), event.id); - - tracing::debug!( - event_id = %event.id, - repo_ref = %repo_ref, - relay_count = relays.len(), - "Added root event to index and pending for Layer 3 filter creation" - ); + repo_sync.relays.clone() + } else if let Some(needs) = pending.repos.get(&repo_ref) { + needs.relays.clone() } else { tracing::debug!( event_id = %event.id, repo_ref = %repo_ref, "Root event references unknown repo" ); - } + return; + }; + drop(index); + + // Discovery consults the published repository index, so staging this + // candidate does not activate it before the announcement batch commits. + self.root_candidate_index.write().await.insert( + event.id, + super::discovery::AcceptedRoot { + id: event.id, + author: event.pubkey, + repository: repo_ref.clone(), + }, + ); + pending.add_root_event(repo_ref.clone(), relays.clone(), event.id); + tracing::debug!( + event_id = %event.id, + repo_ref = %repo_ref, + relay_count = relays.len(), + "Added root event to index and pending for Layer 3 filter creation" + ); } /// Process accumulated batch @@ -763,6 +756,48 @@ mod tests { assert_eq!(root_event_repo_ref(&event), Some(coordinate)); } + #[tokio::test] + async fn live_root_survives_an_uncommitted_announcement_batch() { + let keys = Keys::generate(); + let repo = format!("30617:{}:live-staging", keys.public_key()); + let relay = "wss://source.example".to_string(); + let root = EventBuilder::new(Kind::GitIssue, "") + .tag(Tag::custom("a", [repo.clone()])) + .finalize(&keys) + .expect("build root"); + let database: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded()); + let repo_sync_index = Arc::new(RwLock::new(HashMap::new())); + let root_candidate_index = Arc::new(RwLock::new(HashMap::new())); + let (action_tx, _action_rx) = mpsc::channel(1); + let subscriber = SelfSubscriber::new( + "ws://127.0.0.1:1".to_string(), + "127.0.0.1:1".to_string(), + Arc::clone(&repo_sync_index), + Arc::clone(&root_candidate_index), + action_tx, + database, + LocalRelay::new(), + ); + let mut pending = PendingUpdates::new(); + subscriber.handle_root_event(&root, &mut pending).await; + assert!( + root_candidate_index.read().await.is_empty(), + "unknown repositories stay excluded" + ); + pending.replace_announcement_relays(repo.clone(), HashSet::from([relay])); + subscriber.handle_root_event(&root, &mut pending).await; + assert!( + repo_sync_index.read().await.is_empty(), + "the batch remains private" + ); + assert!(root_candidate_index.read().await.contains_key(&root.id)); + subscriber.process_batch(&mut pending).await; + assert_eq!( + repo_sync_index.read().await[&repo].root_events, + HashSet::from([root.id]) + ); + } + #[tokio::test] async fn startup_reconstruction_promotes_candidates_only_after_batch_completion() { let keys = Keys::generate();