diff --git a/src/sync/mod.rs b/src/sync/mod.rs index e7325b8..cf3b8f7 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -60,6 +60,7 @@ use nostr_sdk::prelude::LocalRelay; const MAX_PURGATORY_DEPENDENCY_EVENTS_PER_TICK: usize = 32; const MAX_PURGATORY_FILTER_ACTIONS_PER_TICK: usize = 1; +const MAX_PURGATORY_DEPENDENCY_IDS_PER_QUERY: usize = 100; const SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS: usize = 20; const SEMANTIC_FALLBACK_MAX_DELIVERED_PERCENT: usize = 10; @@ -121,6 +122,12 @@ fn select_purgatory_dependency_events( selected } +#[derive(Default)] +struct DependencyRelayBatch { + event_ids: HashSet, + identifiers: HashSet, +} + /// Return one stable identity for a relay URL throughout all sync indexes. /// /// URL parsers represent a root path as `/`, so announcements that alternate @@ -149,6 +156,7 @@ fn is_own_sync_target(relay_url: &str, service_domain: &str) -> bool { url_matches_service_domain(relay_url, service_domain) } +#[cfg(test)] fn connections_for_relay_urls( connections: &HashMap, relay_urls: &[String], @@ -3621,10 +3629,15 @@ impl SyncManager { purgatory_dependency_retry_after(), MAX_PURGATORY_DEPENDENCY_EVENTS_PER_TICK, ); + let mut dependency_refetch_batches = HashMap::new(); for event in dependency_events { - self.reprocess_purgatory_announcement_dependencies(&event) - .await; + self.reprocess_purgatory_announcement_dependencies( + &event, + &mut dependency_refetch_batches, + ) + .await; } + self.spawn_batched_purgatory_dependency_refetch(dependency_refetch_batches); // Recompute all outstanding actions on every pass. A bounded pass may // intentionally defer some relays, so restricting this to newly seen @@ -3672,7 +3685,11 @@ impl SyncManager { /// `process_event_static`, but runs while the owner announcement is still in /// purgatory. Announcement policy already treats purgatory announcements as /// maintainer authority; doing the retry here makes arrival order irrelevant. - async fn reprocess_purgatory_announcement_dependencies(&mut self, event: &Event) { + async fn reprocess_purgatory_announcement_dependencies( + &mut self, + event: &Event, + dependency_refetch_batches: &mut HashMap, + ) { let announcement = match crate::nostr::events::RepositoryAnnouncement::from_event(event.clone()) { Ok(announcement) => announcement, @@ -3807,7 +3824,6 @@ impl SyncManager { if !hot_dependency_events.is_empty() { Self::process_purgatory_dependency_events( hot_dependency_events, - &announcement.identifier, &self.database, &self.write_policy, &self.local_relay, @@ -3855,46 +3871,48 @@ impl SyncManager { }) .collect() }; - self.spawn_purgatory_dependency_refetch( - &connected_dependency_relays, - dependency_refetch_ids, - &announcement.identifier, - ); + for relay_url in connected_dependency_relays { + let batch = dependency_refetch_batches.entry(relay_url).or_default(); + batch + .event_ids + .extend(dependency_refetch_ids.iter().copied()); + batch.identifiers.insert(announcement.identifier.clone()); + } } } /// Fetch cold-cache dependency IDs without blocking the sync manager. /// /// Candidate IDs stay in the rejected index until policy processing - /// succeeds. Failed or empty requests are retried on a bounded cadence, and - /// all relay requests for one repository run in parallel. - fn spawn_purgatory_dependency_refetch( + /// succeeds. All due repositories are combined by relay before requests + /// are spawned, so the maintenance cadence consumes one query per bounded + /// ID chunk rather than one query per repository and relay. + fn spawn_batched_purgatory_dependency_refetch( &self, - relay_urls: &[String], - event_ids: HashSet, - identifier: &str, + mut batches: HashMap, ) { - let connections = connections_for_relay_urls(&self.connections, relay_urls); - - if connections.is_empty() { - tracing::debug!( - identifier = %identifier, - event_count = event_ids.len(), - "Cannot refetch purgatory dependencies before a relay connection is available" - ); + if batches.is_empty() { return; } - let event_ids: Vec = self - .reserve_dependency_refetch_attempts(event_ids) - .into_iter() - .collect(); - - if event_ids.is_empty() { + let due_event_ids = self.reserve_dependency_refetch_attempts( + batches + .values() + .flat_map(|batch| batch.event_ids.iter().copied()), + ); + if due_event_ids.is_empty() { return; } + for batch in batches.values_mut() { + batch + .event_ids + .retain(|event_id| due_event_ids.contains(event_id)); + } + batches.retain(|relay_url, batch| { + !batch.event_ids.is_empty() && self.connections.contains_key(relay_url) + }); - let identifier = identifier.to_string(); + let connections = self.connections.clone(); let database = self.database.clone(); let write_policy = self.write_policy.clone(); let local_relay = self.local_relay.clone(); @@ -3902,24 +3920,37 @@ impl SyncManager { let dependency_refetch_attempts = self.dependency_refetch_attempts.clone(); tokio::spawn(async move { - let fetches = connections.into_iter().map(|(relay_url, connection)| { - let filter = Filter::new().ids(event_ids.clone()); - async move { - let result = connection - .fetch_events(filter, Duration::from_secs(5)) - .await; - (relay_url, result) + let mut fetches = Vec::new(); + for (relay_url, batch) in batches { + let Some(connection) = connections.get(&relay_url).cloned() else { + continue; + }; + let event_ids: Vec = batch.event_ids.into_iter().collect(); + let repository_count = batch.identifiers.len(); + for chunk in event_ids.chunks(MAX_PURGATORY_DEPENDENCY_IDS_PER_QUERY) { + let chunk = chunk.to_vec(); + let connection = connection.clone(); + let relay_url = relay_url.clone(); + fetches.push(async move { + let result = connection + .fetch_events( + Filter::new().ids(chunk.iter().copied()).limit(chunk.len()), + Duration::from_secs(5), + ) + .await; + (relay_url, repository_count, chunk.len(), result) + }); } - }); + } let mut fetched_events = Vec::new(); - for (relay_url, result) in join_all(fetches).await { + for (relay_url, repository_count, requested_count, result) in join_all(fetches).await { match result { Ok(events) => { tracing::info!( relay = %relay_url, - identifier = %identifier, - requested_count = event_ids.len(), + repository_count, + requested_count, fetched_count = events.len(), "Fetched purgatory dependencies by exact event ID" ); @@ -3932,10 +3963,10 @@ impl SyncManager { Err(error) => { tracing::warn!( relay = %relay_url, - identifier = %identifier, - event_count = event_ids.len(), + repository_count, + event_count = requested_count, error = %error, - "Failed to refetch purgatory dependencies" + "Failed to refetch batched purgatory dependencies" ); } } @@ -3943,7 +3974,6 @@ impl SyncManager { Self::process_purgatory_dependency_events( fetched_events, - &identifier, &database, &write_policy, &local_relay, @@ -3988,7 +4018,6 @@ impl SyncManager { #[allow(clippy::too_many_arguments)] async fn process_purgatory_dependency_events( mut events: Vec<(String, Event)>, - identifier: &str, database: &SharedDatabase, write_policy: &Nip34WritePolicy, local_relay: &LocalRelay, @@ -4025,7 +4054,14 @@ impl SyncManager { } if result == ProcessResult::Purgatory && event.kind == Kind::RepoState { - write_policy.purgatory().enqueue_sync_immediate(identifier); + if let Some(identifier) = event + .tags + .iter() + .find(|tag| tag.kind() == "d") + .and_then(|tag| tag.content()) + { + write_policy.purgatory().enqueue_sync_immediate(identifier); + } } } } diff --git a/tests/sync/maintainer_reprocessing.rs b/tests/sync/maintainer_reprocessing.rs index d0f3643..75438b3 100644 --- a/tests/sync/maintainer_reprocessing.rs +++ b/tests/sync/maintainer_reprocessing.rs @@ -36,7 +36,9 @@ use std::time::Duration; use nostr_sdk::prelude::*; -use crate::common::{sync_helpers::*, TestRelay}; +use crate::common::{ + censoring_proxy::CensoringProxy, mock_relay::MockRelay, sync_helpers::*, TestRelay, +}; async fn wait_for_log(path: &Path, needle: &str, timeout: Duration) -> bool { let deadline = tokio::time::Instant::now() + timeout; @@ -588,6 +590,108 @@ async fn test_invitee_only_acceptance_recovers_cold_owner_invitation() { owner_relay.stop().await; } +/// Cold dependency polling for separate repositories sharing a relay is +/// combined into one exact-ID query on each maintenance round. +#[tokio::test] +async fn unresolved_repositories_share_one_dependency_poll_per_relay() { + let source = MockRelay::start().await; + let proxy = CensoringProxy::start(source.url()).await; + let invitee_relay = + TestRelay::start_with_sync_and_rejected_hot_cache(Some(proxy.url().to_string()), 1).await; + let invitee_keys = Keys::generate(); + let owners = [Keys::generate(), Keys::generate()]; + let identifiers = ["batched-dependency-one", "batched-dependency-two"]; + + let mut invitations = Vec::new(); + for (owner, identifier) in owners.iter().zip(identifiers) { + let invitation = EventBuilder::new(Kind::GitRepoAnnouncement, "Owner invitation") + .tags([ + Tag::identifier(identifier), + Tag::custom( + "clone", + [format!( + "https://example.invalid/{}/{}.git", + owner.public_key().to_hex(), + identifier + )], + ), + Tag::custom("relays", [proxy.url().to_string()]), + Tag::custom("maintainers", [invitee_keys.public_key().to_hex()]), + ]) + .finalize(owner) + .expect("Failed to create owner invitation"); + send_to_relay_url(source.url(), &invitation) + .await + .expect("Failed to seed owner invitation"); + invitations.push(invitation); + } + + for invitation in &invitations { + let invitation_note = invitation + .id + .to_bech32() + .expect("Failed to encode invitation event ID"); + assert!( + wait_for_log( + &invitee_relay.log_path(), + &invitation_note, + Duration::from_secs(10), + ) + .await, + "Invitee relay should reject and index both owner invitations" + ); + proxy.withhold(invitation.id); + } + + // Passage of the configured one-second hot-cache TTL is required so both + // repositories exercise cold exact-ID recovery on the same relay. + tokio::time::sleep(Duration::from_secs(2)).await; + + let acceptances: Vec = owners + .iter() + .zip(identifiers) + .zip(&invitations) + .map(|((owner, identifier), invitation)| { + repository_announcement( + &invitee_keys, + &[&invitee_relay], + &[owner.public_key()], + identifier, + ) + .custom_created_at(Timestamp::from_secs(invitation.created_at.as_secs() + 1)) + .finalize(&invitee_keys) + .expect("Failed to create invitee acceptance") + }) + .collect(); + let client = Client::default(); + client + .add_relay(invitee_relay.url()) + .await + .expect("Failed to add invitee relay"); + client.connect().await; + for acceptance in &acceptances { + client + .send_event(acceptance) + .await + .expect("Failed to publish invitee acceptance"); + } + client.disconnect().await; + + assert!( + wait_for_log( + &invitee_relay.log_path(), + "repository_count=2 requested_count=2 fetched_count=0", + Duration::from_secs(20), + ) + .await, + "Separate repositories should share one dependency query to their common relay" + ); + + invitee_relay.stop().await; + proxy.stop().await; + source.stop().await; +} + /// Listing a maintainer immediately authorizes that maintainer's state for the /// inviting owner's repository; reciprocal acceptance is not required. ///