diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index 15fa363..b58785f 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -611,9 +611,8 @@ async fn run_observed_git_command_with_policy( let first_log = tokio::time::Instant::now() + GIT_LONG_RUNNING_AFTER; let mut long_running = tokio::time::interval_at(first_log, GIT_LONG_RUNNING_REPEAT); - let mut inactivity_check = tokio::time::interval( - policy.inactivity_limit.min(Duration::from_secs(1)), - ); + let mut inactivity_check = + tokio::time::interval(policy.inactivity_limit.min(Duration::from_secs(1))); inactivity_check.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); let mut stalled = false; let status = loop { @@ -1395,7 +1394,10 @@ mod fetch_helper_tests { if !alive { break; } - assert!(tokio::time::Instant::now() < deadline, "descendant survived"); + assert!( + tokio::time::Instant::now() < deadline, + "descendant survived" + ); tokio::task::yield_now().await; } } @@ -1457,9 +1459,14 @@ mod fetch_helper_tests { .stdout(std::process::Stdio::piped()) .stderr(std::process::Stdio::piped()) .kill_on_drop(true); - let output = run_observed_git_command(follow_up, "test.example", "fetch_batch", GitFetchRole::Primary) - .await - .unwrap(); + let output = run_observed_git_command( + follow_up, + "test.example", + "fetch_batch", + GitFetchRole::Primary, + ) + .await + .unwrap(); assert!(output.status.success()); assert_eq!(output.stdout, b"ready"); } diff --git a/src/sync/discovery.rs b/src/sync/discovery.rs index 7735e31..fac4262 100644 --- a/src/sync/discovery.rs +++ b/src/sync/discovery.rs @@ -162,9 +162,11 @@ pub fn accepted_repository_authors<'a>( continue; } authors.insert(event.pubkey); - for maintainer in event.tags.iter().filter(|tag| { - tag.as_slice().first().map(String::as_str) == Some("maintainers") - }) { + for maintainer in event + .tags + .iter() + .filter(|tag| tag.as_slice().first().map(String::as_str) == Some("maintainers")) + { authors.extend( maintainer.as_slice()[1..] .iter() diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 913598d..e3b990b 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -854,7 +854,11 @@ impl Default for DescendantSyncRotation { } fn filter_group_fingerprint(filters: &[Filter]) -> String { - filters.iter().map(Filter::as_json).collect::>().join("\n") + filters + .iter() + .map(Filter::as_json) + .collect::>() + .join("\n") } fn rotation_fingerprint(filters: &[Filter], max_filters_per_req: usize) -> Vec { @@ -869,7 +873,12 @@ impl DescendantSyncRotation { let previous: HashMap> = self .filters .drain(..) - .map(|cursor| (filter_group_fingerprint(&cursor.filters), cursor.last_successful_until)) + .map(|cursor| { + ( + filter_group_fingerprint(&cursor.filters), + cursor.last_successful_until, + ) + }) .collect(); self.filters = group_filters_for_req_with_max(&filters, max_filters_per_req) .into_iter() @@ -981,11 +990,7 @@ fn packed_historic_filters(items: &PendingItems, frontier: &DescendantFrontier) .collect(); let mut filters = filters::state_event_filters_for_our_repos(&all_repos, None); - let coordinate_values: HashSet<_> = items - .repos - .union(&frontier.coordinates) - .cloned() - .collect(); + let coordinate_values: HashSet<_> = items.repos.union(&frontier.coordinates).cloned().collect(); filters.extend(filters::tagged_one_of_our_repo_event_filters( &coordinate_values, None, @@ -1247,8 +1252,7 @@ fn policy_refusal(message: &str) -> Option { return None; } let message = message.to_ascii_lowercase(); - if message.contains("unsupported filter") || message.contains("filter validation failed") - { + if message.contains("unsupported filter") || message.contains("filter validation failed") { Some(PolicyRefusal::FilterIncompatible) } else if message.contains("not a member") || message.contains("membership") { Some(PolicyRefusal::MembershipRequired) @@ -1924,8 +1928,7 @@ pub struct SyncManager { connect_attempt_result_tx: Option>, /// Wakes the actor when a batch completion may unblock consolidation. deferred_consolidation_tx: Option>, - nip65_discovery_result_tx: - Option>, + nip65_discovery_result_tx: Option>, /// Channel for broadcasting shutdown signal to all background tasks shutdown_tx: Option>, /// Prometheus metrics for sync operations (None if metrics disabled) @@ -2532,9 +2535,8 @@ impl SyncManager { .collect(); { let mut pending = self.pending_sync_index.write().await; - if let Some(batch) = pending - .get_mut(&relay_url_for_retry) - .and_then(|batches| { + if let Some(batch) = + pending.get_mut(&relay_url_for_retry).and_then(|batches| { batches.iter_mut().find(|batch| batch.batch_id == batch_id) }) { @@ -2572,9 +2574,8 @@ impl SyncManager { "Failed to create retry subscription for missing events" ); let mut pending = self.pending_sync_index.write().await; - if let Some(batch) = pending - .get_mut(&relay_url_for_retry) - .and_then(|batches| { + if let Some(batch) = + pending.get_mut(&relay_url_for_retry).and_then(|batches| { batches.iter_mut().find(|batch| batch.batch_id == batch_id) }) { @@ -3438,10 +3439,7 @@ impl SyncManager { if let Some(ref bootstrap_url) = self.bootstrap_relay_url.clone() { match canonical_relay_key(bootstrap_url) { Ok(relay_url) => { - if self - .register_relay(relay_url.clone(), true, false) - .await - { + if self.register_relay(relay_url.clone(), true, false).await { self.schedule_connect_relay(&relay_url).await; } } @@ -3626,9 +3624,7 @@ impl SyncManager { // A relay first reached through NIP-65 may later become an ordinary // repository target. Promote that existing connection before reading // its state so subsequent reconnects and sync work use the full role. - let promoted_from_discovery = self - .nip65_discovery_only_relays - .remove(&action.relay_url); + let promoted_from_discovery = self.nip65_discovery_only_relays.remove(&action.relay_url); if promoted_from_discovery { tracing::info!( relay = %action.relay_url, @@ -3946,11 +3942,10 @@ impl SyncManager { sub_id = %sub_id, "EOSE received, notifying SyncManager" ); - let _ = eose_tx - .send(EoseNotification { - relay_url: relay_url_clone.clone(), - sub_id, - }); + let _ = eose_tx.send(EoseNotification { + relay_url: relay_url_clone.clone(), + sub_id, + }); } RelayEvent::Notice(notice) => { if is_rate_limit_message(¬ice) { @@ -4018,14 +4013,13 @@ impl SyncManager { ); } } - let _ = subscription_closed_tx - .send(SubscriptionClosedNotification { - relay_url: relay_url_clone.clone(), - subscription_id, - reason, - generation: live_generation, - live_filter_count, - }); + let _ = subscription_closed_tx.send(SubscriptionClosedNotification { + relay_url: relay_url_clone.clone(), + subscription_id, + reason, + generation: live_generation, + live_filter_count, + }); } RelayEvent::Shutdown => { tracing::info!(relay = %relay_url_clone, "Relay shutdown detected"); @@ -4474,14 +4468,10 @@ impl SyncManager { return; } - let reconcile_relay = - relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone(); + let reconcile_relay = relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone(); self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1); - self.reconcile_descendant_mode( - &reconcile_relay, - &targets[&reconcile_relay], - ) - .await; + self.reconcile_descendant_mode(&reconcile_relay, &targets[&reconcile_relay]) + .await; let constrained: Vec = relay_urls .into_iter() @@ -4890,10 +4880,8 @@ impl SyncManager { retained_identity.iter().cloned(), ¤t_authors, ); - self.nip65_discovery.relay_lists = discovery::latest_relay_lists( - latest_identity, - ¤t_authors, - ); + self.nip65_discovery.relay_lists = + discovery::latest_relay_lists(latest_identity, ¤t_authors); self.nip65_discovery.author_inboxes = self .nip65_discovery .relay_lists @@ -4946,8 +4934,8 @@ impl SyncManager { .next_attempt_at .retain(|(relay, author), _| { current_sources - .get(author) - .is_some_and(|sources| sources.contains(relay)) + .get(author) + .is_some_and(|sources| sources.contains(relay)) }); self.nip65_discovery.next_inventory_at = Some(now + nip65_inventory_interval()); let overlay = discovery::build_inbox_root_overlay( @@ -4997,11 +4985,9 @@ impl SyncManager { if authors.is_empty() { continue; } - self.nip65_discovery.in_flight.extend( - authors - .iter() - .map(|author| (source.clone(), *author)), - ); + self.nip65_discovery + .in_flight + .extend(authors.iter().map(|author| (source.clone(), *author))); let filter = Filter::new() .kinds([Kind::Metadata, Kind::RelayList]) .authors(authors.iter().copied()) @@ -5011,7 +4997,9 @@ impl SyncManager { }; let source_relay = source.clone(); tokio::spawn(async move { - let outcome = connection.fetch_events(filter, Duration::from_secs(30)).await; + let outcome = connection + .fetch_events(filter, Duration::from_secs(30)) + .await; let _ = result_tx.send(Nip65DiscoveryResult { source_relay, authors, @@ -5039,10 +5027,7 @@ impl SyncManager { } } - async fn install_nip65_overlay( - &mut self, - new_overlay: HashMap>, - ) { + async fn install_nip65_overlay(&mut self, new_overlay: HashMap>) { let old_overlay = std::mem::replace(&mut self.nip65_discovery.inbox_roots, new_overlay); if old_overlay == self.nip65_discovery.inbox_roots { return; @@ -5086,7 +5071,10 @@ impl SyncManager { ) .await; if event.kind == Kind::RelayList - && matches!(process_result, ProcessResult::Saved | ProcessResult::Duplicate) + && matches!( + process_result, + ProcessResult::Saved | ProcessResult::Duplicate + ) { represented_by_source.insert(event.pubkey); } @@ -5109,10 +5097,9 @@ impl SyncManager { for author in &result.authors { let retry_after = nip65_author_retry_after(query_succeeded, &represented_by_source, author); - self.nip65_discovery.next_attempt_at.insert( - (result.source_relay.clone(), *author), - now + retry_after, - ); + self.nip65_discovery + .next_attempt_at + .insert((result.source_relay.clone(), *author), now + retry_after); } let mut changed = false; @@ -5126,8 +5113,7 @@ impl SyncManager { .get(&author) .is_none_or(|current| { candidate.created_at > current.created_at - || (candidate.created_at == current.created_at - && candidate.id < current.id) + || (candidate.created_at == current.created_at && candidate.id < current.id) }); if replace { self.nip65_discovery @@ -5193,8 +5179,11 @@ impl SyncManager { return; } let now = Instant::now(); - let has_due_author = self.nip65_discovery.author_sources.iter().any( - |(author, sources)| { + let has_due_author = self + .nip65_discovery + .author_sources + .iter() + .any(|(author, sources)| { sources.contains(source) && !self .nip65_discovery @@ -5205,8 +5194,7 @@ impl SyncManager { .next_attempt_at .get(&(source.to_string(), *author)) .is_none_or(|due| *due <= now) - }, - ); + }); if !has_due_author { tracing::debug!( relay = %source, @@ -5666,10 +5654,7 @@ impl SyncManager { Instant::now() + dependency_relay_retention(), ); if !self.connections.contains_key(&relay_url) { - if !self - .register_relay(relay_url.clone(), false, false) - .await - { + if !self.register_relay(relay_url.clone(), false, false).await { continue; } self.schedule_connect_relay(&relay_url).await; @@ -6655,8 +6640,8 @@ impl SyncManager { } else { policy_refusal(reason) }; - if let Some(category) = policy_category - .filter(|category| *category != PolicyRefusal::FilterIncompatible) + if let Some(category) = + policy_category.filter(|category| *category != PolicyRefusal::FilterIncompatible) { let removed_batch = { let mut pending = self.pending_sync_index.write().await; @@ -6951,8 +6936,7 @@ impl SyncManager { self.dependency_relay_deadlines .retain(|_, deadline| *deadline > now); - let mut desired_relays: HashSet = - self.derive_targets().await.into_keys().collect(); + let mut desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); desired_relays.extend(self.dependency_relay_deadlines.keys().cloned()); // Collect relays to disconnect @@ -7065,8 +7049,7 @@ impl SyncManager { /// /// For each eligible relay, a reconnection is queued via schedule_connect_relay. async fn retry_disconnected_relays(&mut self) { - let desired_relays: HashSet = - self.derive_targets().await.into_keys().collect(); + let desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); // Collect relays to reconnect let to_reconnect: Vec = { @@ -7626,9 +7609,7 @@ impl SyncManager { } } } - for (idx, (subscription_id, filter)) in - planned_subscriptions.iter().enumerate() - { + for (idx, (subscription_id, filter)) in planned_subscriptions.iter().enumerate() { let result = if let Some(conn) = self.connections.get(relay_url) { conn.subscribe_filter_with_id( filter.clone(), @@ -7649,12 +7630,9 @@ impl SyncManager { ); let completed_batch = { let mut pending = self.pending_sync_index.write().await; - if let Some(batch) = pending - .get_mut(relay_url) - .and_then(|batches| { - batches.iter_mut().find(|batch| batch.batch_id == batch_id) - }) - { + if let Some(batch) = pending.get_mut(relay_url).and_then(|batches| { + batches.iter_mut().find(|batch| batch.batch_id == batch_id) + }) { batch.outstanding_subs.remove(subscription_id); } take_drained_batch_as_failed(&mut pending, relay_url, batch_id) @@ -7843,20 +7821,12 @@ mod tests { let represented_by_source = HashSet::new(); assert_eq!( - nip65_author_retry_after( - true, - &represented_by_source, - &author_with_retained_list, - ), + nip65_author_retry_after(true, &represented_by_source, &author_with_retained_list,), nip65_discovery_retry_interval() ); let represented_by_source = HashSet::from([author_with_retained_list]); assert_eq!( - nip65_author_retry_after( - true, - &represented_by_source, - &author_with_retained_list, - ), + nip65_author_retry_after(true, &represented_by_source, &author_with_retained_list,), nip65_discovery_refresh_interval() ); } @@ -7879,14 +7849,8 @@ mod tests { assert_eq!(window.saved, 1); assert_eq!(window.duplicate, 1); assert_eq!(window.queue_delay, std::time::Duration::from_millis(20)); - assert_eq!( - window.max_queue_delay, - std::time::Duration::from_millis(12) - ); - assert_eq!( - window.processing_time, - std::time::Duration::from_millis(10) - ); + assert_eq!(window.max_queue_delay, std::time::Duration::from_millis(12)); + assert_eq!(window.processing_time, std::time::Duration::from_millis(10)); assert_eq!( window.max_processing_time, std::time::Duration::from_millis(7) @@ -8013,18 +7977,18 @@ mod tests { MAX_FILTERS_PER_REQ, ); let (filter_index, first, until) = rotation.next_request(now).unwrap(); - assert!(first.iter().all(|filter| { - serde_json::to_value(filter).unwrap().get("since").is_none() - })); + assert!(first + .iter() + .all(|filter| { serde_json::to_value(filter).unwrap().get("since").is_none() })); rotation.mark_started(41, filter_index, until); assert!(rotation.next_request(now).is_none()); assert!(rotation.mark_completed(41, false)); let (retry_index, retry, retry_until) = rotation.next_request(now).unwrap(); assert_eq!(retry_index, filter_index); - assert!(retry.iter().all(|filter| { - serde_json::to_value(filter).unwrap().get("since").is_none() - })); + assert!(retry + .iter() + .all(|filter| { serde_json::to_value(filter).unwrap().get("since").is_none() })); rotation.mark_started(42, retry_index, retry_until); assert!(rotation.mark_completed(42, true)); @@ -8102,7 +8066,11 @@ mod tests { // One state filter plus one a/A/q and one e/E/q family. Appending // descendants separately would produce thirteen filters here. assert_eq!(packed.len(), 7); - let json = packed.iter().map(Filter::as_json).collect::>().join("\n"); + let json = packed + .iter() + .map(Filter::as_json) + .collect::>() + .join("\n"); assert!(json.contains(&root.to_hex())); assert!(json.contains(&descendant.to_hex())); assert!(json.contains("30617:owner:repo")); @@ -8715,7 +8683,10 @@ mod tests { relay, &subscription_id )); - assert!(attempts.is_empty(), "a refused retry must retire its marker"); + assert!( + attempts.is_empty(), + "a refused retry must retire its marker" + ); } #[test] @@ -9144,7 +9115,10 @@ mod tests { let subscription_ids: Vec<_> = (0..169) .map(|index| SubscriptionId::new(format!("hydration-{index}"))) .collect(); - let event_ids = [EventId::from_byte_array([1; 32]), EventId::from_byte_array([2; 32])]; + let event_ids = [ + EventId::from_byte_array([1; 32]), + EventId::from_byte_array([2; 32]), + ]; register_negentropy_hydration_attempt( &mut batch, @@ -9157,7 +9131,12 @@ mod tests { assert_eq!(batch.requested_event_ids.as_ref().unwrap().len(), 2); assert_eq!(batch.received_event_ids, Some(HashSet::new())); for event_id in event_ids { - if batch.requested_event_ids.as_ref().unwrap().contains(&event_id) { + if batch + .requested_event_ids + .as_ref() + .unwrap() + .contains(&event_id) + { batch.received_event_ids.as_mut().unwrap().insert(event_id); } } diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 06c7f80..3e38ebe 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -399,14 +399,12 @@ fn plan_live_tail_extension( let mut trial_filters = replacement_filters.clone(); trial_filters.extend(filters.iter().cloned()); trial_filters.sort_unstable_by_key(|filter| filter.as_json()); - let trial_groups = - super::group_filters_for_req_with_max(&trial_filters, max_filters); - let trial_net_new_slots = - trial_groups.len() as isize - (retired.len() + 1) as isize; + let trial_groups = super::group_filters_for_req_with_max(&trial_filters, max_filters); + let trial_net_new_slots = trial_groups.len() as isize - (retired.len() + 1) as isize; if trial_net_new_slots < net_new_slots - && best.as_ref().is_none_or( - |(_, _, _, best_net)| trial_net_new_slots < *best_net, - ) + && best + .as_ref() + .is_none_or(|(_, _, _, best_net)| trial_net_new_slots < *best_net) { best = Some((index, trial_filters, trial_groups, trial_net_new_slots)); } @@ -2048,9 +2046,8 @@ impl RelayConnection { // The relay can answer an empty or cached query before `subscribe` // returns. Register transient ownership against a caller-chosen ID // first so an immediate EOSE/CLOSED cannot race past local accounting. - let transient_sub_id = transient_class.map(|_| { - requested_subscription_id.unwrap_or_else(SubscriptionId::generate) - }); + let transient_sub_id = transient_class + .map(|_| requested_subscription_id.unwrap_or_else(SubscriptionId::generate)); if let (Some(sub_id), Some(permit)) = (&transient_sub_id, transient_permit) { self.hold_transient_req_permit(sub_id.clone(), permit); } @@ -2209,10 +2206,7 @@ impl RelayConnection { } } - pub(super) async fn retire_auth_refused_subscription( - &self, - subscription_id: &SubscriptionId, - ) { + pub(super) async fn retire_auth_refused_subscription(&self, subscription_id: &SubscriptionId) { self.release_transient_req_permit(subscription_id); self.release_live_req_permit(subscription_id); if let Err(error) = self.client.unsubscribe(subscription_id).await { @@ -2379,10 +2373,7 @@ impl RelayConnection { } /// Release the live ledger slot held for `sub_id`, if any. - fn release_live_req_permit( - &self, - sub_id: &SubscriptionId, - ) -> Option { + fn release_live_req_permit(&self, sub_id: &SubscriptionId) -> Option { self.live_req_permits_held .lock() .expect("live permit map poisoned") @@ -3323,7 +3314,10 @@ mod tests { fn filter_count_learning_halves_and_resets_per_session() { let connection = permissive_connection("wss://strict.example", Keys::generate()); - assert_eq!(connection.max_filters_per_req(), super::super::MAX_FILTERS_PER_REQ); + assert_eq!( + connection.max_filters_per_req(), + super::super::MAX_FILTERS_PER_REQ + ); assert_eq!(connection.reduce_max_filters_per_req(8), Some(4)); assert_eq!(connection.reduce_max_filters_per_req(4), Some(2)); assert_eq!(connection.reduce_max_filters_per_req(2), Some(1)); @@ -3331,7 +3325,10 @@ mod tests { assert_eq!(connection.reduce_max_filters_per_req(8), None); connection.reset_subscription_budget(None); - assert_eq!(connection.max_filters_per_req(), super::super::MAX_FILTERS_PER_REQ); + assert_eq!( + connection.max_filters_per_req(), + super::super::MAX_FILTERS_PER_REQ + ); } #[tokio::test] @@ -3581,10 +3578,7 @@ mod tests { fn minimum_churn_extension_preserves_a_byte_full_partial_group() { let byte_full_id = SubscriptionId::new("byte-full-core"); let large_value = "x".repeat(crate::sync::REQ_MESSAGE_BYTE_BUDGET); - let byte_full = vec![Filter::new().custom_tag( - SingleLetterTag::LOWERCASE_A, - large_value, - )]; + let byte_full = vec![Filter::new().custom_tag(SingleLetterTag::LOWERCASE_A, large_value)]; let new_filter = Filter::new().kind(Kind::Custom(23_400)); let plan = plan_live_tail_extension( @@ -3618,7 +3612,10 @@ mod tests { assert_eq!(plan.retired.len(), 13); assert_eq!(plan.replacement_groups.len(), 2); - assert_eq!(plan.replacement_groups.iter().map(Vec::len).sum::(), 14); + assert_eq!( + plan.replacement_groups.iter().map(Vec::len).sum::(), + 14 + ); assert_eq!(plan.retired.len() - plan.replacement_groups.len(), 11); } @@ -3673,7 +3670,7 @@ mod tests { .expect("open initial core groups"); let auxiliary_ids = connection .subscribe_auxiliary_live_filter_groups(vec![vec![ - Filter::new().kind(Kind::Custom(23_700)), + Filter::new().kind(Kind::Custom(23_700)) ]]) .await .expect("open auxiliary group"); @@ -3693,7 +3690,10 @@ mod tests { .live_req_permits_held .lock() .expect("live permit map poisoned"); - assert!(held.contains_key(&core_ids[0]), "full core group was replaced"); + assert!( + held.contains_key(&core_ids[0]), + "full core group was replaced" + ); assert!( !held.contains_key(&core_ids[1]), "partial core tail was not replaced" @@ -3732,7 +3732,7 @@ mod tests { .expect("open initial core groups"); let auxiliary_ids = connection .subscribe_auxiliary_live_filter_groups(vec![vec![ - Filter::new().kind(Kind::Custom(24_200)), + Filter::new().kind(Kind::Custom(24_200)) ]]) .await .expect("open auxiliary group"); @@ -3781,13 +3781,17 @@ mod tests { .live_req_permits_held .lock() .expect("live permit map poisoned"); - assert!(held.contains_key(&core_ids[0]), "full core group was churned"); + assert!( + held.contains_key(&core_ids[0]), + "full core group was churned" + ); assert!( held.contains_key(&auxiliary_ids[0]), "auxiliary group was churned" ); assert!( - held.values().any(|subscription| subscription.filters == tail), + held.values() + .any(|subscription| subscription.filters == tail), "the exact retired tail was not restored" ); assert_eq!(held.len(), 3); @@ -3839,8 +3843,8 @@ mod tests { ]], &[], move |subscription_id| { - let attempt = close_attempts_for_call - .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let attempt = + close_attempts_for_call.fetch_add(1, std::sync::atomic::Ordering::Relaxed); let client = close_client.clone(); async move { if attempt == 1 { diff --git a/src/sync/self_subscriber.rs b/src/sync/self_subscriber.rs index b971a75..b9a5327 100644 --- a/src/sync/self_subscriber.rs +++ b/src/sync/self_subscriber.rs @@ -734,18 +734,14 @@ mod tests { let repo = format!("30617:{}:{identifier}", keys.public_key()); let relay = "wss://source.example"; let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") - .tags([ - Tag::identifier(identifier), - Tag::custom("relays", [relay]), - ]) + .tags([Tag::identifier(identifier), Tag::custom("relays", [relay])]) .finalize(&keys) .expect("build announcement"); 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 database: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded()); database .save_event(&announcement) .await diff --git a/tests/common/mod.rs b/tests/common/mod.rs index f3461e4..e336a53 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -8,13 +8,13 @@ pub mod git_server; pub mod mock_relay; pub mod neg_limiting_proxy; pub mod nip09_helpers; -pub mod req_limiting_proxy; -pub mod upload_pack_counting_proxy; pub mod port; pub mod purgatory_helpers; pub mod relay; +pub mod req_limiting_proxy; pub mod setup_drop_relay; pub mod sync_helpers; +pub mod upload_pack_counting_proxy; pub use git_server::{SimpleGitServer, SmartGitServer}; pub use mock_relay::MockRelay; diff --git a/tests/common/neg_limiting_proxy.rs b/tests/common/neg_limiting_proxy.rs index fe55088..aa4e8b6 100644 --- a/tests/common/neg_limiting_proxy.rs +++ b/tests/common/neg_limiting_proxy.rs @@ -265,12 +265,7 @@ fn parse_neg_frame(text: &str) -> Option { let value = serde_json::from_str::(text).ok()?; let array = value.as_array()?; let kind = array.first()?.as_str()?; - let subid = || { - array - .get(1) - .and_then(|v| v.as_str()) - .map(|s| s.to_string()) - }; + let subid = || array.get(1).and_then(|v| v.as_str()).map(|s| s.to_string()); match kind { "NEG-OPEN" => Some(NegFrame::Open(subid()?)), "NEG-CLOSE" | "NEG-ERR" => Some(NegFrame::Close(subid()?)), diff --git a/tests/sync.rs b/tests/sync.rs index dbbacbe..0c1f56d 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -42,8 +42,8 @@ mod sync { pub mod metrics; pub mod naughty_list_scheduling; pub mod neg_concurrency; - pub mod purgatory_fetch; pub mod proactive_sync_plus; + pub mod purgatory_fetch; pub mod reconnect_backoff; pub mod req_concurrency; pub mod stale_connect_result; diff --git a/tests/sync/live_sync.rs b/tests/sync/live_sync.rs index 82632af..ca6ae27 100644 --- a/tests/sync/live_sync.rs +++ b/tests/sync/live_sync.rs @@ -223,9 +223,7 @@ async fn live_sync_regroups_after_filter_count_refusal() { let refusal_deadline = tokio::time::Instant::now() + Duration::from_secs(10); loop { let metrics = fetch_metrics(&syncing.url()).await.unwrap_or_default(); - if metrics.contains("ngit_sync_policy_refusals_total") - && metrics.contains("filter_count") - { + if metrics.contains("ngit_sync_policy_refusals_total") && metrics.contains("filter_count") { break; } assert!( diff --git a/tests/sync/naughty_list_scheduling.rs b/tests/sync/naughty_list_scheduling.rs index c5e9637..add65f8 100644 --- a/tests/sync/naughty_list_scheduling.rs +++ b/tests/sync/naughty_list_scheduling.rs @@ -88,9 +88,7 @@ async fn naughty_relay_is_not_scheduled_for_reconnection() { assert_eq!(accepted.load(Ordering::SeqCst), 1, "first dial is required"); let reconnect_deadline = tokio::time::Instant::now() + RECONNECT_OBSERVATION; - while tokio::time::Instant::now() < reconnect_deadline - && accepted.load(Ordering::SeqCst) == 1 - { + while tokio::time::Instant::now() < reconnect_deadline && accepted.load(Ordering::SeqCst) == 1 { tokio::time::sleep(Duration::from_millis(100)).await; } assert_eq!( diff --git a/tests/sync/neg_concurrency.rs b/tests/sync/neg_concurrency.rs index 0686934..5bf4c15 100644 --- a/tests/sync/neg_concurrency.rs +++ b/tests/sync/neg_concurrency.rs @@ -176,7 +176,11 @@ async fn startup_historic_sync_stays_within_relay_neg_concurrency_limit() { // 6. The issues must arrive via the bounded historic sync (sampled ends // of the range cover both byte-budgeted chunks). - for issue_id in [issue_ids[0], issue_ids[ISSUE_COUNT / 2], issue_ids[ISSUE_COUNT - 1]] { + for issue_id in [ + issue_ids[0], + issue_ids[ISSUE_COUNT / 2], + issue_ids[ISSUE_COUNT - 1], + ] { assert!( wait_for_event_on_relay( syncing.url(), @@ -191,18 +195,15 @@ async fn startup_historic_sync_stays_within_relay_neg_concurrency_limit() { // 7. Wait for negentropy activity to include the root-event batch and // settle. Layer-1 plus the repository batch plus the six root-event // rounds put the floor at eight. - let (opened, rejected) = wait_for_neg_quiescence( - &proxy, - 8, - Duration::from_secs(3), - Duration::from_secs(60), - ) - .await - .expect("negentropy rounds should reach the root-event batch and settle"); + let (opened, rejected) = + wait_for_neg_quiescence(&proxy, 8, Duration::from_secs(3), Duration::from_secs(60)) + .await + .expect("negentropy rounds should reach the root-event batch and settle"); // 8. The regression assertions. assert_eq!( - rejected, 0, + rejected, + 0, "relay exceeded the per-connection NEG concurrency limit \ (opened: {opened}, peak: {})", proxy.peak_concurrent() diff --git a/tests/sync/proactive_sync_plus.rs b/tests/sync/proactive_sync_plus.rs index b7f65a0..64ab3ff 100644 --- a/tests/sync/proactive_sync_plus.rs +++ b/tests/sync/proactive_sync_plus.rs @@ -5,8 +5,8 @@ use std::time::Duration; use nostr_sdk::prelude::*; use crate::common::{ - build_layer2_issue_event, repo_coord, send_to_relay_url, setup_announcement_on_relay, - wait_for_event_on_relay, MockRelay, TestClient, TestRelay, reserve_port, + build_layer2_issue_event, repo_coord, reserve_port, send_to_relay_url, + setup_announcement_on_relay, wait_for_event_on_relay, MockRelay, TestClient, TestRelay, }; #[tokio::test] @@ -157,9 +157,9 @@ async fn root_author_inbox_reuses_existing_root_sync_pipeline() { "the advertised outbox connection should be identified as discovery-only" ); assert!( - !sync_log.lines().any(|line| { - line.contains("Starting fresh_start") && line.contains(outbox.url()) - }), + !sync_log + .lines() + .any(|line| { line.contains("Starting fresh_start") && line.contains(outbox.url()) }), "querying an advertised outbox must not start ordinary repository sync against it" ); let draining_reply = EventBuilder::new(Kind::TextNote, "reply on naturally draining inbox") diff --git a/tests/sync/req_concurrency.rs b/tests/sync/req_concurrency.rs index 022a969..ad19f9a 100644 --- a/tests/sync/req_concurrency.rs +++ b/tests/sync/req_concurrency.rs @@ -226,7 +226,11 @@ async fn startup_historic_sync_stays_within_relay_req_concurrency_limit() { // 6. The issues must arrive via historic sync (sampled ends of the // range cover the byte-budgeted chunks). - for issue_id in [issue_ids[0], issue_ids[ISSUE_COUNT / 2], issue_ids[ISSUE_COUNT - 1]] { + for issue_id in [ + issue_ids[0], + issue_ids[ISSUE_COUNT / 2], + issue_ids[ISSUE_COUNT - 1], + ] { assert!( wait_for_event_on_relay( syncing.url(), @@ -263,7 +267,8 @@ async fn startup_historic_sync_stays_within_relay_req_concurrency_limit() { // 9. The regression assertions (cumulative across both phases; the // trickle-shaped phase 1 must stay within the bound too). assert_eq!( - rejected, 0, + rejected, + 0, "relay exceeded the per-connection REQ concurrency limit \ (opened: {opened}, peak: {})", proxy.peak_concurrent()