fix(sync): batch purgatory dependency polls by relay

Production on gitnostr.com recorded four unresolved repositories producing 193 exact-ID relay fetch log entries and 482 requested event copies during the 2026-08-06 06:26:59-06:42:35 UTC window. The 30-second recovery cadence is intentional so authorization and state events are found promptly when they appear, but each purgatory repository currently spawns its own request to every hinted relay and competes with ordinary synchronization for per-connection query-rate budget.

Build one maintenance-round plan keyed by canonical connected relay, union all due cold dependency IDs from every selected purgatory announcement, and issue bounded 100-ID queries per relay. Preserve the existing 30-second cadence, parallel relay I/O, exact-ID retry accounting, authorization ordering, and early exit when dependencies resolve. Include repository_count and requested_count in the existing fetch log so production batching is directly observable.

The 100-ID chunk assumes an exact-ID filter can return at most one match per requested ID and stays within the lowest audited default result cap. This deliberately does not add live subscriptions, alter purgatory retention, add exponential backoff, or change general historic-sync requests.

Validated with cargo check --workspace --all-targets, all nine maintainer-reprocessing integration scenarios, and 47 sync unit tests. The new scenario indexes two cold invitations for separate repositories behind one relay and observes one recovery query with repository_count=2 and requested_count=2.
This commit is contained in:
DanConwayDev
2026-08-06 10:02:56 +00:00
parent dfb054a5d0
commit 6dfe0a7e25
2 changed files with 188 additions and 48 deletions
+83 -47
View File
@@ -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<EventId>,
identifiers: HashSet<String>,
}
/// 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<T: Clone>(
connections: &HashMap<String, T>,
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<String, DependencyRelayBatch>,
) {
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<EventId>,
identifier: &str,
mut batches: HashMap<String, DependencyRelayBatch>,
) {
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<EventId> = 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<EventId> = 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);
}
}
}
}
+105 -1
View File
@@ -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<Event> = 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.
///