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());