mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
style: restore rustfmt compliance
Apply the workspace rustfmt output that accumulated across recent sync changes so the formatting gate can reach the test phase again. This commit is deliberately mechanical: it contains no behavioral changes and leaves any test failures for separately reviewable fixes. Validated with: cargo fmt --all -- --check
This commit is contained in:
@@ -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");
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
+97
-118
@@ -854,7 +854,11 @@ impl Default for DescendantSyncRotation {
|
||||
}
|
||||
|
||||
fn filter_group_fingerprint(filters: &[Filter]) -> String {
|
||||
filters.iter().map(Filter::as_json).collect::<Vec<_>>().join("\n")
|
||||
filters
|
||||
.iter()
|
||||
.map(Filter::as_json)
|
||||
.collect::<Vec<_>>()
|
||||
.join("\n")
|
||||
}
|
||||
|
||||
fn rotation_fingerprint(filters: &[Filter], max_filters_per_req: usize) -> Vec<String> {
|
||||
@@ -869,7 +873,12 @@ impl DescendantSyncRotation {
|
||||
let previous: HashMap<String, Option<Timestamp>> = 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<PolicyRefusal> {
|
||||
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<tokio::sync::mpsc::UnboundedSender<ConnectAttemptResult>>,
|
||||
/// Wakes the actor when a batch completion may unblock consolidation.
|
||||
deferred_consolidation_tx: Option<tokio::sync::mpsc::UnboundedSender<String>>,
|
||||
nip65_discovery_result_tx:
|
||||
Option<tokio::sync::mpsc::UnboundedSender<Nip65DiscoveryResult>>,
|
||||
nip65_discovery_result_tx: Option<tokio::sync::mpsc::UnboundedSender<Nip65DiscoveryResult>>,
|
||||
/// Channel for broadcasting shutdown signal to all background tasks
|
||||
shutdown_tx: Option<broadcast::Sender<()>>,
|
||||
/// 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<String> = 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<String, HashSet<EventId>>,
|
||||
) {
|
||||
async fn install_nip65_overlay(&mut self, new_overlay: HashMap<String, HashSet<EventId>>) {
|
||||
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<String> =
|
||||
self.derive_targets().await.into_keys().collect();
|
||||
let mut desired_relays: HashSet<String> = 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<String> =
|
||||
self.derive_targets().await.into_keys().collect();
|
||||
let desired_relays: HashSet<String> = self.derive_targets().await.into_keys().collect();
|
||||
|
||||
// Collect relays to reconnect
|
||||
let to_reconnect: Vec<String> = {
|
||||
@@ -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::<Vec<_>>().join("\n");
|
||||
let json = packed
|
||||
.iter()
|
||||
.map(Filter::as_json)
|
||||
.collect::<Vec<_>>()
|
||||
.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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<ReleasedLiveSubscription> {
|
||||
fn release_live_req_permit(&self, sub_id: &SubscriptionId) -> Option<ReleasedLiveSubscription> {
|
||||
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::<usize>(), 14);
|
||||
assert_eq!(
|
||||
plan.replacement_groups.iter().map(Vec::len).sum::<usize>(),
|
||||
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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
+2
-2
@@ -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;
|
||||
|
||||
@@ -265,12 +265,7 @@ fn parse_neg_frame(text: &str) -> Option<NegFrame> {
|
||||
let value = serde_json::from_str::<serde_json::Value>(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()?)),
|
||||
|
||||
+1
-1
@@ -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;
|
||||
|
||||
@@ -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!(
|
||||
|
||||
@@ -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!(
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user