From 186d9a8ccf71a73ec3c963068ff3e335c6c3e8a5 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 21:53:46 +0000 Subject: [PATCH 1/2] refactor(sync): make repository filter chunks deterministic Minimum-churn live consolidation can preserve an existing subscription only when equivalent desired coverage is rendered identically. Layer-2 repository and state filters currently inherit randomized HashSet iteration order, so the same repository set may be split and serialized differently on each derivation. Sort repository references and extracted identifiers before byte-budget chunking, matching the existing deterministic root-event path. This changes neither the selected events nor the byte and filter limits; it only stabilizes group identity and makes future lifecycle decisions reproducible. This commit deliberately does not alter subscription replacement, packing policy, or descendant ownership. Those behavior changes remain isolated in the following commit. Validation: nix develop -c cargo test --lib repo_and_state_filter_serialization_is_insertion_order_independent passed. The focused test constructs equivalent sets through opposite insertion orders and compares serialized repository and state filters. --- src/sync/filters.rs | 37 +++++++++++++++++++++++++++++++++++-- 1 file changed, 35 insertions(+), 2 deletions(-) diff --git a/src/sync/filters.rs b/src/sync/filters.rs index f24cacb..d6de410 100644 --- a/src/sync/filters.rs +++ b/src/sync/filters.rs @@ -112,7 +112,8 @@ pub fn state_event_filters_for_our_repos( } let mut filters = Vec::new(); - let identifier_vec: Vec<_> = identifiers.iter().collect(); + let mut identifier_vec: Vec<_> = identifiers.iter().collect(); + identifier_vec.sort_unstable(); // Batch identifiers per filter within the serialized byte budget for chunk in chunk_values_by_bytes(&identifier_vec) { @@ -153,7 +154,8 @@ pub fn tagged_one_of_our_repo_event_filters( } let mut filters = Vec::new(); - let repo_refs: Vec<_> = repos.iter().collect(); + let mut repo_refs: Vec<_> = repos.iter().collect(); + repo_refs.sort_unstable(); for chunk in chunk_values_by_bytes(&repo_refs) { // Lowercase 'a' tag - standard addressable reference @@ -395,6 +397,37 @@ mod tests { assert_eq!(filters.len(), 3); } + #[test] + fn repo_and_state_filter_serialization_is_insertion_order_independent() { + let ascending: HashSet = (0..700) + .map(|index| format!("30617:pubkey:{index:04}-{}", "x".repeat(48))) + .collect(); + let descending: HashSet = (0..700) + .rev() + .map(|index| format!("30617:pubkey:{index:04}-{}", "x".repeat(48))) + .collect(); + + let repo_ascending: Vec = tagged_one_of_our_repo_event_filters(&ascending, None) + .into_iter() + .map(|filter| filter.as_json()) + .collect(); + let repo_descending: Vec = tagged_one_of_our_repo_event_filters(&descending, None) + .into_iter() + .map(|filter| filter.as_json()) + .collect(); + assert_eq!(repo_ascending, repo_descending); + + let state_ascending: Vec = state_event_filters_for_our_repos(&ascending, None) + .into_iter() + .map(|filter| filter.as_json()) + .collect(); + let state_descending: Vec = state_event_filters_for_our_repos(&descending, None) + .into_iter() + .map(|filter| filter.as_json()) + .collect(); + assert_eq!(state_ascending, state_descending); + } + #[test] fn test_chunk_values_by_bytes_edge_cases() { // Empty input -> no chunks. From 9a1fa8de01298e296dd771c2e4d178dac2f169bd Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 22:09:02 +0000 Subject: [PATCH 2/2] feat(sync): consolidate useful live tail capacity Repository discovery arrives in five-second batches, but the previous incremental path accumulated subscriptions until a coarse threshold rebuilt the complete live set. Reopening healthy core groups and separately managed descendant coverage creates avoidable relay load and coverage gaps. Use the actual filter-count and serialized-byte grouping function as the capacity authority. Preserve protected descendant and count-full groups. Repack the complete mutable core tail when that releases slots; when byte limits mean the complete tail cannot shrink, retire only the smallest groups whose participation reduces the incremental slot cost. Failed CLOSE or replacement open restores the exact retired tail, while capacity pressure retains the deferred full-regroup backstop. Correctness assumes the held-permit map describes the filters open in the current connection session and that adding one existing group can increase grouped output by at most its prior one slot. Reconnect, exceptional full restoration, and multi-connection sharding are deliberately unchanged. Validated with cargo check --lib and cargo test --lib (680 passed). Tests cover deterministic grouping, full/descendant preservation, thirteen small partial groups becoming two replacements and releasing eleven slots, seventeen byte-bound groups replacing only one to avoid a new slot, byte-full preservation, replacement-open rollback, and partial-CLOSE rollback. --- CHANGELOG.md | 9 + docs/explanation/sync-scaling-constraints.md | 31 +- src/sync/mod.rs | 155 +---- src/sync/relay_connection.rs | 602 +++++++++++++++++++ 4 files changed, 658 insertions(+), 139 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ef99923..203d080 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 history provides eventual coverage. The frontier is deliberately non-recursive and connection sharding remains unnecessary. +### Changed + +- Minimise live-subscription churn as repository coverage grows: stable full + core groups and descendant groups remain open. The mutable core tail is + repacked in full when doing so releases slots; otherwise only the smallest + useful subset is replaced to absorb new filters. This accounts for both + filter-count and serialized-byte limits. Repository filter chunks are + deterministic, and failed tail replacement restores the previous tail. + ## [2.1.2] - 2026-08-08 ngit-grasp 2.1.2 is a patch release improving repository-event sync under diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 1f5d4c6..bd8f973 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -323,16 +323,27 @@ NIP-11 `max_subscriptions` sets B for each new connection session; when it is absent B falls back to 20. Advertised values below that floor are honoured (notably nostream's default 10). Two slots remain reserved. Core live filter groups are packed first and admitted atomically against the advertised -subscription-count budget: if the complete core live set cannot fit, -the existing subscriptions are consolidated into the byte- and filter-count -bounded REQ groups first. If the consolidated core live set still cannot fit, -partial coverage is not opened, historic work is deferred, and a warning -surfaces the condition. Reconnect and consolidation build Layer 1, 2, and 3 -coverage together and reserve the whole grouped set before opening it. A -failure while opening group N sends CLOSE for every earlier group and returns -their slots, so callers never inherit hidden partial coverage. Replacement -also remembers the exact previous grouping and restores it if the new complete -set fails at runtime. Multi-connection sharding remains the later lever. +subscription-count budget. Incremental five-second batches preserve full core +REQs and separately owned descendant REQs. The planner uses the same +filter-count and serialized-byte grouping rules as wire submission. It repacks +the complete mutable core tail with the new filters when that releases at least +one slot; otherwise it retires only the smallest useful subset which reduces +the incremental slot cost. Thus byte-bound groups are not rebuilt merely +because they contain fewer than the maximum filter count. Repository and +identifier inputs are sorted before byte chunking so equivalent coverage has +stable group identity. +If the changed tail cannot fit the count or learned byte budget, no extension +is opened and historic recovery remains available. Capacity pressure is the +backstop which may schedule a complete regroup after outstanding historic +batches drain; an ordinary tail update never rebuilds stable full groups. + +Reconnect and exceptional full restoration still build Layer 1, 2, and 3 +coverage together and reserve the whole grouped set before opening it. Both +full replacement and tail replacement remember the exact previous grouping: +a failure while opening a replacement closes every newly opened group and +restores the retired groups. A partial CLOSE failure likewise reopens any tail +groups which were already closed before reporting the failure. Multi-connection +sharding remains the later lever. A transient slot is released only after EOSE has caused CLOSE to be enqueued, after relay CLOSED, or after connection teardown. Because NIP-01 provides no diff --git a/src/sync/mod.rs b/src/sync/mod.rs index b16a595..d3cc367 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -974,10 +974,6 @@ struct ConnectAttemptResult { /// Quick reconnect window in seconds (15 minutes) const QUICK_RECONNECT_WINDOW_SECS: u64 = 15 * 60; -/// Maximum incremental filter fragmentation above the desired live baseline -/// before triggering consolidation. -const CONSOLIDATION_THRESHOLD: usize = 70; - /// Bound concurrent DNS and websocket handshakes so a large relay list cannot /// exhaust network resources while keeping the sync actor responsive. const MAX_CONCURRENT_CONNECT_ATTEMPTS: usize = 8; @@ -1134,21 +1130,6 @@ async fn begin_connect_attempt( Some(permit) } -fn consolidation_fragmentation( - current_count: usize, - new_count: usize, - desired_baseline: usize, -) -> usize { - current_count - .saturating_add(new_count) - .saturating_sub(desired_baseline) -} - -fn should_consolidate(current_count: usize, new_count: usize, desired_baseline: usize) -> bool { - consolidation_fragmentation(current_count, new_count, desired_baseline) - > CONSOLIDATION_THRESHOLD -} - fn grouped_subscription_count(filters: &[Filter]) -> usize { group_filters_for_req(filters).len() } @@ -3230,17 +3211,6 @@ impl SyncManager { return; } - // Step 3: Check if consolidation is needed BEFORE adding new filters - if !self - .maybe_consolidate( - &action.relay_url, - grouped_subscription_count(&action.filters), - ) - .await - { - return; - } - // Subscribe to each filter and collect subscription IDs tracing::info!( relay = %action.relay_url, @@ -3251,15 +3221,28 @@ impl SyncManager { "handle_add_filters: calling sync_live and historic_sync" ); - if let Err(error) = self - .sync_live(&action.relay_url, &action.filters) - .await - { + if let Err(error) = self.sync_live(&action.relay_url, &action.filters).await { tracing::warn!( relay = %action.relay_url, %error, "Live coverage could not be extended; continuing bounded historic sync" ); + if error.starts_with("Minimum-churn live extension needs") + || error.starts_with("Minimum-churn live extension exceeds") + { + let has_pending_batches = self.has_pending_batches(&action.relay_url).await; + if self + .deferred_consolidations + .request(&action.relay_url, has_pending_batches) + { + let _ = self.consolidate(&action.relay_url).await; + } else { + tracing::info!( + relay = %action.relay_url, + "Minimum-churn extension reached capacity; full consolidation deferred until pending batches drain" + ); + } + } } self.historic_sync(&action.relay_url, action.filters, action.items, None) .await; @@ -5542,77 +5525,6 @@ impl SyncManager { .is_some_and(|batches| !batches.is_empty()) } - async fn desired_live_filter_count(&self, relay_url: &str) -> usize { - let target = { - let repo_index = self.repo_sync_index.read().await; - algorithms::derive_relay_targets(&repo_index).remove(relay_url) - }; - let desired_repo_filters = target.map_or_else(Vec::new, |target| { - filters::build_sync_level_aware_filters( - &target.repos, - &target.state_only_repos, - &target.root_events, - None, - ) - }); - - // Every connected relay carries one consolidated generic announcement - // subscription in addition to its repository-specific desired filters. - 1 + grouped_subscription_count(&desired_repo_filters) - } - - /// Check if incremental fragmentation exceeds the consolidation threshold. - /// - /// The desired live set is an irreducible baseline, so a repository whose - /// consolidated filters already exceed 70 must remain stable. - async fn maybe_consolidate(&mut self, relay_url: &str, new_count: usize) -> bool { - let current_count = if let Some(connection) = self.connections.get(relay_url) { - connection.subscription_count().await - } else { - 0 - }; - let desired_baseline = self.desired_live_filter_count(relay_url).await; - let fragmentation = consolidation_fragmentation(current_count, new_count, desired_baseline); - let budget_pressure = self - .connections - .get(relay_url) - .is_some_and(|connection| connection.needs_live_consolidation(new_count)); - - if budget_pressure || should_consolidate(current_count, new_count, desired_baseline) { - let has_pending_batches = self.has_pending_batches(relay_url).await; - if !self - .deferred_consolidations - .request(relay_url, has_pending_batches) - { - tracing::info!( - relay = %relay_url, - current_count, - new_count, - desired_baseline, - fragmentation, - threshold = CONSOLIDATION_THRESHOLD, - budget_pressure, - "Live subscription pressure requires consolidation; deferring until pending batches drain" - ); - return false; - } - - tracing::info!( - relay = %relay_url, - current_count = current_count, - new_count = new_count, - desired_baseline, - fragmentation, - threshold = CONSOLIDATION_THRESHOLD, - budget_pressure, - "Live subscription pressure requires consolidation" - ); - - return self.consolidate(relay_url).await; - } - true - } - async fn process_deferred_consolidation(&mut self, relay_url: &str) { let has_pending_batches = self.has_pending_batches(relay_url).await; if !self @@ -6317,8 +6229,16 @@ impl SyncManager { if filter_groups.is_empty() { return Ok(Vec::new()); } + let protected_subscription_ids = self + .descendant_live_coverage + .get(relay_url) + .map(|coverage| coverage.subscription_ids.clone()) + .unwrap_or_default(); connection - .subscribe_live_filter_groups(filter_groups) + .extend_live_filter_groups_minimally( + filter_groups, + &protected_subscription_ids, + ) .await .inspect_err(|error| { tracing::error!(relay = %relay_url, error = %error, "Failed to create complete live subscription set"); @@ -7128,29 +7048,6 @@ mod tests { ); } - #[test] - fn desired_baseline_above_threshold_does_not_reconsolidate() { - let desired_baseline = 178; - - assert!( - !should_consolidate(desired_baseline, 1, desired_baseline), - "the irreducible desired set must be a stable post-rebuild baseline" - ); - assert!( - !should_consolidate( - desired_baseline + CONSOLIDATION_THRESHOLD - 1, - 1, - desired_baseline, - ), - "the configured fragmentation headroom is allowed above the baseline" - ); - assert!(should_consolidate( - desired_baseline + CONSOLIDATION_THRESHOLD, - 1, - desired_baseline, - )); - } - #[test] fn grouped_pagination_advances_only_filters_that_fill_a_page() { let keys = Keys::generate(); diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 4993ec2..d93bd86 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -332,6 +332,95 @@ struct HeldLiveSubscription { filters: Vec, } +#[derive(Debug)] +struct LiveTailReplacement { + retired: Vec<(SubscriptionId, Vec)>, + replacement_groups: Vec>, +} + +fn plan_live_tail_extension( + current: Vec<(SubscriptionId, Vec)>, + protected: &std::collections::HashSet, + new_groups: Vec>, + max_filters: usize, +) -> LiveTailReplacement { + let mut new_filters: Vec = new_groups.into_iter().flatten().collect(); + new_filters.sort_unstable_by_key(|filter| filter.as_json()); + + let max_filters = max_filters.max(1); + let mut candidates: Vec<_> = current + .into_iter() + .filter(|(subscription_id, filters)| { + !protected.contains(subscription_id) && filters.len() < max_filters + }) + .collect(); + candidates.sort_unstable_by_key(|(subscription_id, filters)| { + ( + filters + .iter() + .map(|filter| filter.as_json()) + .collect::>(), + subscription_id.to_string(), + ) + }); + + let mut all_filters = new_filters.clone(); + all_filters.extend( + candidates + .iter() + .flat_map(|(_, filters)| filters.iter().cloned()), + ); + all_filters.sort_unstable_by_key(|filter| filter.as_json()); + let all_groups = super::group_filters_for_req_with_max(&all_filters, max_filters); + + // Repacking the complete tail maximises released slots whenever it can + // actually release one. If byte limits make the complete tail just as + // large, retire only groups which reduce the incremental slot cost. + if all_groups.len() < candidates.len() { + return LiveTailReplacement { + retired: candidates, + replacement_groups: all_groups, + }; + } + + let mut retired = Vec::new(); + let mut replacement_filters = new_filters; + let mut replacement_groups = + super::group_filters_for_req_with_max(&replacement_filters, max_filters); + let mut net_new_slots = replacement_groups.len() as isize; + while net_new_slots > 0 { + let mut best = None; + for (index, (_, filters)) in candidates.iter().enumerate() { + 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; + if trial_net_new_slots < net_new_slots + && 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)); + } + } + let Some((index, filters, groups, net_slots)) = best else { + break; + }; + retired.push(candidates.remove(index)); + replacement_filters = filters; + replacement_groups = groups; + net_new_slots = net_slots; + } + + LiveTailReplacement { + retired, + replacement_groups, + } +} + #[derive(Clone, Copy)] struct ReleasedLiveSubscription { generation: u64, @@ -1117,6 +1206,190 @@ impl RelayConnection { .await } + /// Extend core live coverage while preserving full groups and separately + /// owned auxiliary subscriptions. Only core groups which can absorb at + /// least one new filter are closed and repacked with the new filters. + /// + /// This is a best-effort transaction: any failure after CLOSE restores the + /// exact retired filter groups before it returns. Admission remains subject + /// to the connection ledger, so a caller can fall back to a full regroup + /// when the minimally changed set cannot fit. + pub async fn extend_live_filter_groups_minimally( + &self, + filter_groups: Vec>, + protected_subscription_ids: &[SubscriptionId], + ) -> Result, String> { + self.extend_live_filter_groups_minimally_with( + filter_groups, + protected_subscription_ids, + |subscription_id| async move { + self.client + .unsubscribe(&subscription_id) + .await + .map(|_| ()) + .map_err(|error| format!("{subscription_id}: {error}")) + }, + |filters, permit| self.subscribe_filters_with_live_permit(filters, None, Some(permit)), + ) + .await + } + + async fn extend_live_filter_groups_minimally_with( + &self, + filter_groups: Vec>, + protected_subscription_ids: &[SubscriptionId], + mut close_subscription: C, + mut subscribe_group: F, + ) -> Result, String> + where + C: FnMut(SubscriptionId) -> CFut, + CFut: Future>, + F: FnMut(Vec, SessionPermit) -> Fut, + Fut: Future>, + { + if filter_groups.is_empty() { + return Ok(Vec::new()); + } + let new_filter_count: usize = filter_groups.iter().map(Vec::len).sum(); + + let current: Vec<_> = self + .live_req_permits_held + .lock() + .expect("live permit map poisoned") + .iter() + .map(|(id, held)| (id.clone(), held.filters.clone())) + .collect(); + let protected: std::collections::HashSet<_> = + protected_subscription_ids.iter().cloned().collect(); + let plan = plan_live_tail_extension( + current.clone(), + &protected, + filter_groups, + self.max_filters_per_req(), + ); + + let preserved_count = current.len().saturating_sub(plan.retired.len()); + let retired_count = plan.retired.len(); + let replacement_count = plan.replacement_groups.len(); + let released_slot_count = retired_count.saturating_sub(replacement_count); + let additional_slot_count = replacement_count.saturating_sub(retired_count); + let target_count = preserved_count.saturating_add(plan.replacement_groups.len()); + let usable = self + .subscription_usable_slots + .load(std::sync::atomic::Ordering::Relaxed); + if target_count > usable { + return Err(format!( + "Minimum-churn live extension needs {target_count} of {usable} usable slots for {}", + self.url + )); + } + + if let Some(remote_limit) = self.remote_subscription_byte_limit() { + let retired_ids: std::collections::HashSet<_> = + plan.retired.iter().map(|(id, _)| id).collect(); + let preserved_bytes: usize = current + .iter() + .filter(|(id, _)| !retired_ids.contains(id)) + .map(|(_, filters)| { + ClientMessage::req(SubscriptionId::generate(), filters.clone()) + .as_json() + .len() + }) + .sum(); + let replacement_bytes: usize = plan + .replacement_groups + .iter() + .map(|filters| { + ClientMessage::req(SubscriptionId::generate(), filters.clone()) + .as_json() + .len() + }) + .sum(); + if preserved_bytes + .checked_add(replacement_bytes) + .and_then(|used| used.checked_add(super::SUBSCRIPTION_BYTE_RESERVED_MARGIN)) + .is_none_or(|used| used > remote_limit) + { + return Err(format!( + "Minimum-churn live extension exceeds the learned subscription-state limit for {}", + self.url + )); + } + } + + let mut close_error = None; + for (subscription_id, _) in &plan.retired { + if let Err(error) = close_subscription(subscription_id.clone()).await { + close_error = Some(error); + break; + } + self.release_live_req_permit(subscription_id); + } + if let Some(close_error) = close_error { + let still_held: std::collections::HashSet<_> = self + .live_req_permits_held + .lock() + .expect("live permit map poisoned") + .keys() + .cloned() + .collect(); + let closed_groups: Vec<_> = plan + .retired + .iter() + .filter(|(id, _)| !still_held.contains(id)) + .map(|(_, filters)| filters.clone()) + .collect(); + let restoration = self + .subscribe_live_filter_groups_with(closed_groups, &mut subscribe_group) + .await; + return match restoration { + Ok(_) => Err(format!( + "Minimum-churn CLOSE failed; prior tail restored: {close_error}" + )), + Err(restoration_error) => Err(format!( + "Minimum-churn CLOSE failed ({close_error}); prior tail restoration failed ({restoration_error})" + )), + }; + } + + match self + .subscribe_live_filter_groups_with(plan.replacement_groups, &mut subscribe_group) + .await + { + Ok(ids) => { + tracing::info!( + relay = %self.url, + new_filter_count, + preserved_group_count = preserved_count, + retired_group_count = retired_count, + replacement_group_count = replacement_count, + released_slot_count, + additional_slot_count, + "Extended core live coverage with minimum churn" + ); + Ok(ids) + } + Err(replacement_error) => { + let previous_groups: Vec<_> = plan + .retired + .into_iter() + .map(|(_, filters)| filters) + .collect(); + match self + .subscribe_live_filter_groups_with(previous_groups, &mut subscribe_group) + .await + { + Ok(_) => Err(format!( + "Minimum-churn live extension failed; prior tail restored: {replacement_error}" + )), + Err(restoration_error) => Err(format!( + "Minimum-churn live extension failed ({replacement_error}); prior tail restoration failed ({restoration_error})" + )), + } + } + } + } + async fn replace_live_filter_groups_with( &self, filter_groups: Vec>, @@ -3119,6 +3392,335 @@ mod tests { assert_eq!(connection.subscription_budget().available_permits(), 7); } + #[test] + fn minimum_churn_extension_preserves_full_and_auxiliary_groups() { + let full_id = SubscriptionId::new("full-core"); + let tail_id = SubscriptionId::new("partial-core"); + let auxiliary_id = SubscriptionId::new("descendants"); + let full = vec![ + Filter::new().kind(Kind::Custom(23_000)), + Filter::new().kind(Kind::Custom(23_001)), + Filter::new().kind(Kind::Custom(23_002)), + ]; + let tail = vec![Filter::new().kind(Kind::Custom(23_100))]; + let auxiliary = vec![Filter::new().kind(Kind::Custom(23_200))]; + let new = vec![vec![ + Filter::new().kind(Kind::Custom(23_300)), + Filter::new().kind(Kind::Custom(23_301)), + ]]; + let protected = std::collections::HashSet::from([auxiliary_id.clone()]); + + let plan = plan_live_tail_extension( + vec![ + (full_id, full), + (tail_id.clone(), tail), + (auxiliary_id, auxiliary), + ], + &protected, + new, + 3, + ); + + assert_eq!(plan.retired.len(), 1); + assert_eq!(plan.retired[0].0, tail_id); + assert_eq!(plan.replacement_groups.len(), 1); + assert_eq!(plan.replacement_groups[0].len(), 3); + } + + #[test] + 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 new_filter = Filter::new().kind(Kind::Custom(23_400)); + + let plan = plan_live_tail_extension( + vec![(byte_full_id.clone(), byte_full.clone())], + &std::collections::HashSet::new(), + vec![vec![new_filter.clone()]], + 10, + ); + + assert!(plan.retired.is_empty()); + assert_eq!(plan.replacement_groups, vec![vec![new_filter]]); + } + + #[test] + fn minimum_churn_extension_rebuilds_every_partial_tail_and_releases_slots() { + let tails: Vec<_> = (0..13) + .map(|index| { + ( + SubscriptionId::new(format!("tail-{index}")), + vec![Filter::new().kind(Kind::Custom(23_450 + index))], + ) + }) + .collect(); + + let plan = plan_live_tail_extension( + tails, + &std::collections::HashSet::new(), + vec![vec![Filter::new().kind(Kind::Custom(23_499))]], + 10, + ); + + 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.retired.len() - plan.replacement_groups.len(), 11); + } + + #[test] + fn minimum_churn_extension_retires_one_byte_bound_group_to_avoid_a_new_slot() { + let tails: Vec<_> = (0..17) + .map(|index| { + ( + SubscriptionId::new(format!("byte-tail-{index}")), + vec![Filter::new().custom_tag( + SingleLetterTag::LOWERCASE_A, + format!("{index:02}{}", "x".repeat(60_000)), + )], + ) + }) + .collect(); + let new_filter = Filter::new().custom_tag( + SingleLetterTag::LOWERCASE_A, + format!("00{}", "y".repeat(1_000)), + ); + + let plan = plan_live_tail_extension( + tails, + &std::collections::HashSet::new(), + vec![vec![new_filter]], + 10, + ); + + assert_eq!(plan.retired.len(), 1); + assert_eq!(plan.replacement_groups.len(), 1); + } + + #[tokio::test] + async fn minimum_churn_extension_keeps_full_and_auxiliary_subscriptions_open() { + let relay = LocalRelayBuilder::default().build(); + relay.run().await.expect("start local relay"); + let connection = RelayConnection::new( + relay.url().await.to_string(), + Keys::generate(), + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(3).await.expect("connect local relay"); + + let full: Vec<_> = (0..connection.max_filters_per_req()) + .map(|index| Filter::new().kind(Kind::Custom(23_500 + index as u16))) + .collect(); + let partial = vec![Filter::new().kind(Kind::Custom(23_600))]; + let core_ids = connection + .subscribe_live_filter_groups(vec![full, partial]) + .await + .expect("open initial core groups"); + let auxiliary_ids = connection + .subscribe_auxiliary_live_filter_groups(vec![vec![ + Filter::new().kind(Kind::Custom(23_700)), + ]]) + .await + .expect("open auxiliary group"); + + let replacement_ids = connection + .extend_live_filter_groups_minimally( + vec![vec![ + Filter::new().kind(Kind::Custom(23_800)), + Filter::new().kind(Kind::Custom(23_801)), + ]], + &auxiliary_ids, + ) + .await + .expect("extend only the mutable core tail"); + + let held = connection + .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[1]), + "partial core tail was not replaced" + ); + assert!( + held.contains_key(&auxiliary_ids[0]), + "auxiliary group was replaced" + ); + assert_eq!(replacement_ids.len(), 1); + assert_eq!(held[&replacement_ids[0]].filters.len(), 3); + drop(held); + + connection.disconnect().await; + relay.shutdown(); + } + + #[tokio::test] + async fn failed_minimum_churn_extension_restores_only_the_retired_tail() { + let relay = LocalRelayBuilder::default().build(); + relay.run().await.expect("start local relay"); + let connection = RelayConnection::new( + relay.url().await.to_string(), + Keys::generate(), + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(3).await.expect("connect local relay"); + + let full: Vec<_> = (0..connection.max_filters_per_req()) + .map(|index| Filter::new().kind(Kind::Custom(24_000 + index as u16))) + .collect(); + let tail = vec![Filter::new().kind(Kind::Custom(24_100))]; + let core_ids = connection + .subscribe_live_filter_groups(vec![full, tail.clone()]) + .await + .expect("open initial core groups"); + let auxiliary_ids = connection + .subscribe_auxiliary_live_filter_groups(vec![vec![ + Filter::new().kind(Kind::Custom(24_200)), + ]]) + .await + .expect("open auxiliary group"); + + let attempts = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let held = std::sync::Arc::clone(&connection.live_req_permits_held); + let error = connection + .extend_live_filter_groups_minimally_with( + vec![vec![Filter::new().kind(Kind::Custom(24_300))]], + &auxiliary_ids, + |subscription_id| { + let client = connection.client.clone(); + async move { + client + .unsubscribe(&subscription_id) + .await + .map(|_| ()) + .map_err(|error| error.to_string()) + } + }, + move |filters, permit| { + let attempt = attempts.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let held = std::sync::Arc::clone(&held); + async move { + if attempt == 0 { + return Err("controlled tail replacement failure".to_string()); + } + let sub_id = SubscriptionId::new(format!("restored-tail-{attempt}")); + held.lock().expect("live permit map poisoned").insert( + sub_id.clone(), + HeldLiveSubscription { + generation: permit.generation, + _ledger_slot: permit.permit, + filters, + }, + ); + Ok(sub_id) + } + }, + ) + .await + .expect_err("controlled replacement must fail"); + + assert!(error.contains("prior tail restored"), "{error}"); + let held = connection + .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(&auxiliary_ids[0]), + "auxiliary group was churned" + ); + assert!( + held.values().any(|subscription| subscription.filters == tail), + "the exact retired tail was not restored" + ); + assert_eq!(held.len(), 3); + drop(held); + + connection.disconnect().await; + relay.shutdown(); + } + + #[tokio::test] + async fn partial_close_failure_restores_the_already_closed_tail_group() { + let relay = LocalRelayBuilder::default().build(); + relay.run().await.expect("start local relay"); + let connection = RelayConnection::new( + relay.url().await.to_string(), + Keys::generate(), + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(3).await.expect("connect local relay"); + + let first_tail = vec![Filter::new().custom_tag( + SingleLetterTag::LOWERCASE_A, + format!("a{}", "x".repeat(30_000)), + )]; + let second_tail = vec![Filter::new().custom_tag( + SingleLetterTag::LOWERCASE_A, + format!("c{}", "x".repeat(30_000)), + )]; + connection + .subscribe_live_filter_groups(vec![first_tail.clone(), second_tail.clone()]) + .await + .expect("open two partial core groups"); + + let close_attempts = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let close_client = connection.client.clone(); + let close_attempts_for_call = std::sync::Arc::clone(&close_attempts); + let error = connection + .extend_live_filter_groups_minimally_with( + vec![vec![ + Filter::new().custom_tag( + SingleLetterTag::LOWERCASE_A, + format!("b{}", "x".repeat(60_000)), + ), + Filter::new().custom_tag( + SingleLetterTag::LOWERCASE_A, + format!("d{}", "x".repeat(60_000)), + ), + ]], + &[], + move |subscription_id| { + let attempt = close_attempts_for_call + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let client = close_client.clone(); + async move { + if attempt == 1 { + return Err("controlled second CLOSE failure".to_string()); + } + client + .unsubscribe(&subscription_id) + .await + .map(|_| ()) + .map_err(|error| error.to_string()) + } + }, + |filters, permit| { + connection.subscribe_filters_with_live_permit(filters, None, Some(permit)) + }, + ) + .await + .expect_err("the controlled second CLOSE must fail"); + + assert!(error.contains("prior tail restored"), "{error}"); + let held = connection.live_filter_groups(); + assert!(held.contains(&first_tail)); + assert!(held.contains(&second_tail)); + assert_eq!(held.len(), 2); + assert_eq!(close_attempts.load(std::sync::atomic::Ordering::Relaxed), 2); + + connection.disconnect().await; + relay.shutdown(); + } + #[tokio::test] async fn complete_live_set_can_fill_usable_budget_but_history_defers() { let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());