mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
Merge #599cc928: sync: minimise live subscription churn
nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqs9n8xf9rnczwhv6u6yagtp95u8y45fd2pahtthkfensw5nnaekx0sls8puq PR-Author: DanConwayDev's Agent nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0 CoverNote: Built on the descendant-sync commits merged by #2b377e20 and implements nostr:nevent1qqswhdur5f7z8zwy726za854qhahnt9k5ytru2h6zqyysrzdskwzk9cpz3mhxue69uhhyetvv9ujumn8d96zuer9wcykq3rj. The proposal and current master share exact descendant tip `9f3c943`, so it can merge without rebasing or duplicating commits. Repository discovery is batched every five seconds. Previously each batch added live subscriptions until a coarse threshold rebuilt the complete live set, needlessly closing healthy core and descendant subscriptions. This proposal is split for review: 1. `186d9a8` makes repository filter chunking deterministic, giving stable grouping independent of HashMap insertion order. 2. `9a1fa8d` uses the real filter-count and serialized-byte grouping result to select useful consolidation work. Full core and descendant groups remain open. If repacking the complete mutable tail releases slots, it does so; if byte limits mean the tail cannot shrink, it replaces only the smallest useful subset needed to reduce additional slot usage. Failed CLOSE or replacement-open operations restore the exact retired tail. Capacity pressure retains the existing deferred full-regroup backstop while historic recovery continues. Validation: - `cargo check --lib` - `cargo test --lib`: 680 passed - tests cover stable full/auxiliary IDs against a real embedded relay, deterministic grouping, count and byte constraints, replacement-open rollback, and partial-CLOSE rollback - 13 small partial groups plus one new filter become 2 replacements, releasing 11 slots - 17 byte-bound groups plus one small filter replace only 1 group, avoiding a new slot without rebuilding the other 16 - a byte-full group which cannot help is left untouched Disposable archive production refinement: - the count-only version exposed the exact byte-bound case: a one-filter update on relay.ngit.dev retired 17 groups and reopened 17, releasing no slots - exact final tip `9a1fa8d`, under the equivalent roughly 50,000-ID / 167-chunk startup workload, preserved 17 groups and replaced exactly 1 for the same one-filter update, with zero additional slots - the final-tip service remained active with zero restarts, zero candidate-specific rate-limit, ledger-overrun, rollback, or panic errors, and about 349 MB peak during the comparison window - the public gitnostr.com service remained on released v2.1.2 throughout Earlier archive evidence also established that descendant permanent coverage, multi-subscription splitting, constrained fallback rotation, and historic completion continue correctly beneath this change. Recommendation: ready to merge.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
+35
-2
@@ -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<String> = (0..700)
|
||||
.map(|index| format!("30617:pubkey:{index:04}-{}", "x".repeat(48)))
|
||||
.collect();
|
||||
let descending: HashSet<String> = (0..700)
|
||||
.rev()
|
||||
.map(|index| format!("30617:pubkey:{index:04}-{}", "x".repeat(48)))
|
||||
.collect();
|
||||
|
||||
let repo_ascending: Vec<String> = tagged_one_of_our_repo_event_filters(&ascending, None)
|
||||
.into_iter()
|
||||
.map(|filter| filter.as_json())
|
||||
.collect();
|
||||
let repo_descending: Vec<String> = 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<String> = state_event_filters_for_our_repos(&ascending, None)
|
||||
.into_iter()
|
||||
.map(|filter| filter.as_json())
|
||||
.collect();
|
||||
let state_descending: Vec<String> = 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.
|
||||
|
||||
+26
-129
@@ -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();
|
||||
|
||||
@@ -332,6 +332,95 @@ struct HeldLiveSubscription {
|
||||
filters: Vec<Filter>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct LiveTailReplacement {
|
||||
retired: Vec<(SubscriptionId, Vec<Filter>)>,
|
||||
replacement_groups: Vec<Vec<Filter>>,
|
||||
}
|
||||
|
||||
fn plan_live_tail_extension(
|
||||
current: Vec<(SubscriptionId, Vec<Filter>)>,
|
||||
protected: &std::collections::HashSet<SubscriptionId>,
|
||||
new_groups: Vec<Vec<Filter>>,
|
||||
max_filters: usize,
|
||||
) -> LiveTailReplacement {
|
||||
let mut new_filters: Vec<Filter> = 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::<Vec<_>>(),
|
||||
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<Vec<Filter>>,
|
||||
protected_subscription_ids: &[SubscriptionId],
|
||||
) -> Result<Vec<SubscriptionId>, 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<C, CFut, F, Fut>(
|
||||
&self,
|
||||
filter_groups: Vec<Vec<Filter>>,
|
||||
protected_subscription_ids: &[SubscriptionId],
|
||||
mut close_subscription: C,
|
||||
mut subscribe_group: F,
|
||||
) -> Result<Vec<SubscriptionId>, String>
|
||||
where
|
||||
C: FnMut(SubscriptionId) -> CFut,
|
||||
CFut: Future<Output = Result<(), String>>,
|
||||
F: FnMut(Vec<Filter>, SessionPermit) -> Fut,
|
||||
Fut: Future<Output = Result<SubscriptionId, String>>,
|
||||
{
|
||||
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<F, Fut>(
|
||||
&self,
|
||||
filter_groups: Vec<Vec<Filter>>,
|
||||
@@ -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::<usize>(), 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());
|
||||
|
||||
Reference in New Issue
Block a user