mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
feat(sync): demote reference coverage by priority tier
Motivation: the two-level core/descendant policy forced every core tag variant to stay live while treating every descendant variant alike. As repository scale grows, that either consumes the transient recovery slot or demotes useful canonical coverage together with low-value compatibility forms. Approach: keep announcements, state, repository a and root e filters essential. Reconcile root E, core q/A, descendant e/a and descendant q as an ordered complete-tier prefix against the existing NIP-11-aware ledger, filter-count grouping and byte budget. Pack adjacent tier filters together, reserve one historic slot, and feed all filters below the cutoff plus descendant E/A through the existing five-second cursor-overlapped REQ+EOSE queue. Correctness: historic sync still receives the full original filter set; only persistent admission changes. A tier is never partially admitted, fallback cursor state survives unchanged filters, CLOSED auxiliary subscriptions immediately return their demoted filters to rotation, and all consumers continue to share the existing ledger. Excluded scope: no connection sharding, new scheduler, request-class priority, configuration, or historic repository packing is included. Validation: nix develop -c cargo test --lib (685 passed); focused tier-cutoff, filter-tier and descendant-rotation tests passed.
This commit is contained in:
@@ -304,8 +304,9 @@ Derived from the tightest commonly observed values; all sizing below assumes:
|
||||
Each relay connection owns one implemented budget ledger of B subscription
|
||||
slots. Four consumers share it, in priority order:
|
||||
|
||||
1. **Core live subscriptions** (persistent, `limit: 0`) — the product; sized
|
||||
first.
|
||||
1. **Essential live subscriptions** (persistent, `limit: 0`) — announcements,
|
||||
repository states, canonical repository `a` references and canonical root
|
||||
`e` references are never demoted.
|
||||
2. **Reserved margin** (2 slots) — control-plane safety capacity kept beyond
|
||||
the live set (which includes Layer-1) for ad-hoc operations and recovery.
|
||||
3. **Historic sync and dependency recovery** (transient) — at least one usable
|
||||
@@ -313,21 +314,24 @@ slots. Four consumers share it, in priority order:
|
||||
REQ+EOSE pages/fallbacks/retries, and exact-ID purgatory polls draw from
|
||||
the remainder. NEG retains its four-round class cap and transient REQ its
|
||||
five-request class cap, but neither can exceed the shared residual.
|
||||
4. **Descendant live coverage** (auxiliary persistent) — admitted only when
|
||||
its complete separately grouped set fits after core coverage while still
|
||||
preserving the margin and a transient slot. Otherwise each constrained
|
||||
relay advances one cursor-overlapped REQ+EOSE filter per five-second tick
|
||||
through the same transient queue. Direct thread members contribute their
|
||||
event IDs; replaceable and addressable members also contribute their NIP-01
|
||||
coordinates. The resulting `e`/`E`/`q` and `a`/`A`/coordinate-`q` filters
|
||||
cover one descendant generation without recursively expanding the frontier.
|
||||
4. **Priority-tiered reference coverage** — remaining filters are considered
|
||||
in this order: root `E`; core compatibility `q`/`A`; descendant canonical
|
||||
`e`/`a`; descendant `q`. A complete tier remains persistent only when it
|
||||
fits after essential coverage while preserving the margin and a transient
|
||||
slot. Lower tiers advance one relay-compatible, cursor-overlapped REQ+EOSE
|
||||
filter group per five-second tick through the same transient queue. Descendant uppercase
|
||||
`E`/`A` references are historic-only. Direct thread members contribute
|
||||
event IDs and replaceable/addressable coordinates, covering one descendant
|
||||
generation without recursively expanding the frontier.
|
||||
|
||||
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. Incremental five-second batches preserve full core
|
||||
REQs and separately owned descendant REQs. The planner uses the same
|
||||
(notably nostream's default 10). Two slots remain reserved. Essential live
|
||||
filter groups are packed first and admitted atomically against the advertised
|
||||
subscription-count budget. Incremental five-second batches preserve full
|
||||
essential REQs and separately owned tiered reference REQs. Tier filters are
|
||||
packed across boundaries, while admission stops at a complete-tier boundary.
|
||||
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
|
||||
@@ -340,9 +344,10 @@ 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:
|
||||
Reconnect and exceptional full restoration rebuild essential coverage first;
|
||||
the five-second reconciler then admits the largest complete prefix of reference
|
||||
tiers which fits the refreshed session budget. 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
|
||||
|
||||
@@ -31,6 +31,46 @@ pub struct TieredFilter {
|
||||
pub filter: Filter,
|
||||
}
|
||||
|
||||
/// Build the live coverage that is never demoted under connection pressure.
|
||||
pub fn build_essential_live_filters(
|
||||
full_repos: &HashSet<String>,
|
||||
state_only_repos: &HashSet<String>,
|
||||
root_events: &HashSet<EventId>,
|
||||
since: Option<Timestamp>,
|
||||
) -> Vec<Filter> {
|
||||
let all_repos: HashSet<String> = full_repos.union(state_only_repos).cloned().collect();
|
||||
let mut filters = state_event_filters_for_our_repos(&all_repos, since);
|
||||
if !full_repos.is_empty() {
|
||||
filters.extend(
|
||||
tiered_repo_event_filters(full_repos, false, since)
|
||||
.into_iter()
|
||||
.filter(|entry| entry.tier == CoverageTier::EssentialCore)
|
||||
.map(|entry| entry.filter),
|
||||
);
|
||||
}
|
||||
filters.extend(
|
||||
tiered_root_event_filters(root_events, false, since)
|
||||
.into_iter()
|
||||
.filter(|entry| entry.tier == CoverageTier::EssentialCore)
|
||||
.map(|entry| entry.filter),
|
||||
);
|
||||
filters
|
||||
}
|
||||
|
||||
/// Build the priority-ordered core reference coverage that can move between a
|
||||
/// persistent subscription and paced REQ+EOSE rotation.
|
||||
pub fn tiered_auxiliary_core_filters(
|
||||
full_repos: &HashSet<String>,
|
||||
root_events: &HashSet<EventId>,
|
||||
since: Option<Timestamp>,
|
||||
) -> Vec<TieredFilter> {
|
||||
let mut filters = tiered_root_event_filters(root_events, false, since);
|
||||
filters.extend(tiered_repo_event_filters(full_repos, false, since));
|
||||
filters.retain(|entry| entry.tier != CoverageTier::EssentialCore);
|
||||
filters.sort_by_key(|entry| entry.tier);
|
||||
filters
|
||||
}
|
||||
|
||||
fn tagged_filters<T: AsRef<str>>(
|
||||
values: &[T],
|
||||
tag: SingleLetterTag,
|
||||
|
||||
+259
-98
@@ -707,7 +707,7 @@ struct DescendantFrontier {
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DescendantSyncRotation {
|
||||
frontier: DescendantFrontier,
|
||||
coverage_fingerprint: Vec<String>,
|
||||
filters: Vec<DescendantFilterCursor>,
|
||||
next_filter: usize,
|
||||
in_flight: Option<DescendantFilterInFlight>,
|
||||
@@ -715,7 +715,7 @@ struct DescendantSyncRotation {
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DescendantFilterCursor {
|
||||
filter: Filter,
|
||||
filters: Vec<Filter>,
|
||||
last_successful_until: Option<Timestamp>,
|
||||
}
|
||||
|
||||
@@ -728,14 +728,15 @@ struct DescendantFilterInFlight {
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DescendantLiveCoverage {
|
||||
frontier: DescendantFrontier,
|
||||
coverage_fingerprint: Vec<String>,
|
||||
fallback_filters: Vec<Filter>,
|
||||
subscription_ids: Vec<SubscriptionId>,
|
||||
}
|
||||
|
||||
impl Default for DescendantSyncRotation {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
frontier: DescendantFrontier::default(),
|
||||
coverage_fingerprint: Vec::new(),
|
||||
filters: Vec::new(),
|
||||
next_filter: 0,
|
||||
in_flight: None,
|
||||
@@ -743,40 +744,66 @@ impl Default for DescendantSyncRotation {
|
||||
}
|
||||
}
|
||||
|
||||
fn filter_group_fingerprint(filters: &[Filter]) -> String {
|
||||
filters.iter().map(Filter::as_json).collect::<Vec<_>>().join("\n")
|
||||
}
|
||||
|
||||
fn rotation_fingerprint(filters: &[Filter], max_filters_per_req: usize) -> Vec<String> {
|
||||
group_filters_for_req_with_max(filters, max_filters_per_req)
|
||||
.iter()
|
||||
.map(|group| filter_group_fingerprint(group))
|
||||
.collect()
|
||||
}
|
||||
|
||||
impl DescendantSyncRotation {
|
||||
fn refresh(&mut self, frontier: DescendantFrontier) {
|
||||
fn refresh(&mut self, filters: Vec<Filter>, max_filters_per_req: usize) {
|
||||
let previous: HashMap<String, Option<Timestamp>> = self
|
||||
.filters
|
||||
.drain(..)
|
||||
.map(|cursor| (cursor.filter.as_json(), cursor.last_successful_until))
|
||||
.map(|cursor| (filter_group_fingerprint(&cursor.filters), cursor.last_successful_until))
|
||||
.collect();
|
||||
self.filters = descendant_frontier_filters(&frontier, None)
|
||||
self.filters = group_filters_for_req_with_max(&filters, max_filters_per_req)
|
||||
.into_iter()
|
||||
.map(|filter| DescendantFilterCursor {
|
||||
last_successful_until: previous.get(&filter.as_json()).copied().flatten(),
|
||||
filter,
|
||||
.map(|filters| DescendantFilterCursor {
|
||||
last_successful_until: previous
|
||||
.get(&filter_group_fingerprint(&filters))
|
||||
.copied()
|
||||
.flatten(),
|
||||
filters,
|
||||
})
|
||||
.collect();
|
||||
self.frontier = frontier;
|
||||
self.coverage_fingerprint = self
|
||||
.filters
|
||||
.iter()
|
||||
.map(|cursor| filter_group_fingerprint(&cursor.filters))
|
||||
.collect();
|
||||
self.next_filter = 0;
|
||||
self.in_flight = None;
|
||||
}
|
||||
|
||||
fn next_request(&self, now: Timestamp) -> Option<(usize, Filter, Timestamp)> {
|
||||
fn next_request(&self, now: Timestamp) -> Option<(usize, Vec<Filter>, Timestamp)> {
|
||||
if self.in_flight.is_some() || self.filters.is_empty() {
|
||||
return None;
|
||||
}
|
||||
let filter_index = self.next_filter % self.filters.len();
|
||||
let cursor = &self.filters[filter_index];
|
||||
let mut filter = cursor.filter.clone().until(now);
|
||||
if let Some(last_until) = cursor.last_successful_until {
|
||||
filter = filter.since(Timestamp::from(
|
||||
last_until
|
||||
.as_secs()
|
||||
.saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS),
|
||||
));
|
||||
}
|
||||
Some((filter_index, filter, now))
|
||||
let filters = cursor
|
||||
.filters
|
||||
.iter()
|
||||
.cloned()
|
||||
.map(|filter| {
|
||||
let filter = filter.until(now);
|
||||
match cursor.last_successful_until {
|
||||
Some(last_until) => filter.since(Timestamp::from(
|
||||
last_until
|
||||
.as_secs()
|
||||
.saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS),
|
||||
)),
|
||||
None => filter,
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
Some((filter_index, filters, now))
|
||||
}
|
||||
|
||||
fn mark_started(&mut self, batch_id: u64, filter_index: usize, until: Timestamp) {
|
||||
@@ -788,7 +815,10 @@ impl DescendantSyncRotation {
|
||||
}
|
||||
|
||||
fn mark_completed(&mut self, batch_id: u64, succeeded: bool) -> bool {
|
||||
let Some(in_flight) = self.in_flight.filter(|request| request.batch_id == batch_id) else {
|
||||
let Some(in_flight) = self
|
||||
.in_flight
|
||||
.filter(|request| request.batch_id == batch_id)
|
||||
else {
|
||||
return false;
|
||||
};
|
||||
self.in_flight = None;
|
||||
@@ -825,8 +855,7 @@ fn descendant_frontier_filters(
|
||||
frontier: &DescendantFrontier,
|
||||
since: Option<Timestamp>,
|
||||
) -> Vec<Filter> {
|
||||
let mut filters =
|
||||
filters::tagged_one_of_our_root_event_filters(&frontier.event_ids, since);
|
||||
let mut filters = filters::tagged_one_of_our_root_event_filters(&frontier.event_ids, since);
|
||||
filters.extend(filters::tagged_one_of_our_repo_event_filters(
|
||||
&frontier.coordinates,
|
||||
since,
|
||||
@@ -834,6 +863,65 @@ fn descendant_frontier_filters(
|
||||
filters
|
||||
}
|
||||
|
||||
fn tiered_auxiliary_filters(
|
||||
repos: &HashSet<String>,
|
||||
root_events: &HashSet<EventId>,
|
||||
frontier: &DescendantFrontier,
|
||||
since: Option<Timestamp>,
|
||||
) -> Vec<filters::TieredFilter> {
|
||||
let mut entries = filters::tiered_auxiliary_core_filters(repos, root_events, since);
|
||||
entries.extend(filters::tiered_root_event_filters(
|
||||
&frontier.event_ids,
|
||||
true,
|
||||
since,
|
||||
));
|
||||
entries.extend(filters::tiered_repo_event_filters(
|
||||
&frontier.coordinates,
|
||||
true,
|
||||
since,
|
||||
));
|
||||
entries.sort_by_key(|entry| entry.tier);
|
||||
entries
|
||||
}
|
||||
|
||||
fn split_live_tier_prefix(
|
||||
entries: &[filters::TieredFilter],
|
||||
max_filters_per_req: usize,
|
||||
fits: impl Fn(&[Vec<Filter>]) -> bool,
|
||||
) -> (Vec<Filter>, Vec<Filter>) {
|
||||
use filters::CoverageTier;
|
||||
|
||||
let live_tiers = [
|
||||
CoverageTier::RootUppercase,
|
||||
CoverageTier::CoreCompatibility,
|
||||
CoverageTier::DescendantCanonical,
|
||||
CoverageTier::DescendantQuote,
|
||||
];
|
||||
let mut live = Vec::new();
|
||||
for tier in live_tiers {
|
||||
let mut candidate = live.clone();
|
||||
candidate.extend(
|
||||
entries
|
||||
.iter()
|
||||
.filter(|entry| entry.tier == tier)
|
||||
.map(|entry| entry.filter.clone()),
|
||||
);
|
||||
let groups = live_filter_groups(&candidate, max_filters_per_req);
|
||||
if fits(&groups) {
|
||||
live = candidate;
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
let live_json: HashSet<_> = live.iter().map(Filter::as_json).collect();
|
||||
let rotated = entries
|
||||
.iter()
|
||||
.filter(|entry| !live_json.contains(&entry.filter.as_json()))
|
||||
.map(|entry| entry.filter.clone())
|
||||
.collect();
|
||||
(live, rotated)
|
||||
}
|
||||
|
||||
/// Items included in a pending batch
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct PendingItems {
|
||||
@@ -940,8 +1028,7 @@ fn is_rate_limit_message(message: &str) -> bool {
|
||||
|
||||
fn is_filter_count_refusal(message: &str) -> bool {
|
||||
let message = message.to_ascii_lowercase();
|
||||
message.contains("invalid number of filters")
|
||||
|| message.contains("max filter count")
|
||||
message.contains("invalid number of filters") || message.contains("max filter count")
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
@@ -1085,7 +1172,10 @@ fn groups_within_subscription_byte_limit(
|
||||
let mut overflow = 0usize;
|
||||
for group in groups {
|
||||
let size = req_message_size(&group);
|
||||
if used.checked_add(size).is_some_and(|total| total <= live_budget) {
|
||||
if used
|
||||
.checked_add(size)
|
||||
.is_some_and(|total| total <= live_budget)
|
||||
{
|
||||
used += size;
|
||||
admitted.push(group);
|
||||
} else {
|
||||
@@ -3240,7 +3330,10 @@ impl SyncManager {
|
||||
}
|
||||
|
||||
// Step 2: Check if relay is rate-limited before creating new pending items
|
||||
if self.health_tracker.is_subscription_paused(&action.relay_url) {
|
||||
if self
|
||||
.health_tracker
|
||||
.is_subscription_paused(&action.relay_url)
|
||||
{
|
||||
tracing::debug!(
|
||||
relay = %action.relay_url,
|
||||
full_repos = action.items.repos.len(),
|
||||
@@ -3261,7 +3354,13 @@ impl SyncManager {
|
||||
"handle_add_filters: calling sync_live and historic_sync"
|
||||
);
|
||||
|
||||
if let Err(error) = self.sync_live(&action.relay_url, &action.filters).await {
|
||||
let essential_live = filters::build_essential_live_filters(
|
||||
&action.items.repos,
|
||||
&action.items.state_only_repos,
|
||||
&action.items.root_events,
|
||||
None,
|
||||
);
|
||||
if let Err(error) = self.sync_live(&action.relay_url, &essential_live).await {
|
||||
tracing::warn!(
|
||||
relay = %action.relay_url,
|
||||
%error,
|
||||
@@ -3525,8 +3624,7 @@ impl SyncManager {
|
||||
&& !is_filter_count_refusal(&reason)
|
||||
&& subscription_state_byte_limit(&reason).is_none()
|
||||
{
|
||||
let already_paused =
|
||||
health_tracker.is_rate_limited(&relay_url_clone);
|
||||
let already_paused = health_tracker.is_rate_limited(&relay_url_clone);
|
||||
if already_paused {
|
||||
tracing::debug!(
|
||||
relay = %relay_url_clone,
|
||||
@@ -3785,26 +3883,25 @@ impl SyncManager {
|
||||
async fn reconcile_descendant_mode(
|
||||
&mut self,
|
||||
relay_url: &str,
|
||||
root_events: &HashSet<EventId>,
|
||||
target: &algorithms::RelaySyncNeeds,
|
||||
) {
|
||||
let frontier = self.direct_thread_members(root_events).await;
|
||||
if frontier.event_ids.is_empty() {
|
||||
let _ = self
|
||||
.close_descendant_live_coverage(relay_url, "frontier became empty")
|
||||
.await;
|
||||
self.descendant_sync_rotations.remove(relay_url);
|
||||
return;
|
||||
}
|
||||
let frontier = self.direct_thread_members(&target.root_events).await;
|
||||
let historic_entries =
|
||||
tiered_auxiliary_filters(&target.repos, &target.root_events, &frontier, None);
|
||||
let desired_fingerprint: Vec<_> = historic_entries
|
||||
.iter()
|
||||
.map(|entry| entry.filter.as_json())
|
||||
.collect();
|
||||
if self
|
||||
.descendant_live_coverage
|
||||
.get(relay_url)
|
||||
.is_some_and(|coverage| coverage.frontier == frontier)
|
||||
.is_some_and(|coverage| coverage.coverage_fingerprint == desired_fingerprint)
|
||||
{
|
||||
return;
|
||||
}
|
||||
if self.descendant_live_coverage.contains_key(relay_url)
|
||||
&& !self
|
||||
.close_descendant_live_coverage(relay_url, "frontier changed")
|
||||
.close_descendant_live_coverage(relay_url, "tiered coverage changed")
|
||||
.await
|
||||
{
|
||||
return;
|
||||
@@ -3815,11 +3912,18 @@ impl SyncManager {
|
||||
.as_secs()
|
||||
.saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS),
|
||||
);
|
||||
let historic_filters = descendant_frontier_filters(&frontier, None);
|
||||
let live_filters = descendant_frontier_filters(&frontier, Some(live_since));
|
||||
if let Some(connection) = self.connections.get(relay_url).cloned() {
|
||||
let (live_filters, rotated_filters) = split_live_tier_prefix(
|
||||
&historic_entries,
|
||||
connection.max_filters_per_req(),
|
||||
|groups| connection.can_admit_auxiliary_live_groups(groups),
|
||||
);
|
||||
let live_filters: Vec<_> = live_filters
|
||||
.into_iter()
|
||||
.map(|filter| filter.since(live_since))
|
||||
.collect();
|
||||
let groups = live_filter_groups(&live_filters, connection.max_filters_per_req());
|
||||
if connection.can_admit_auxiliary_live_groups(&groups) {
|
||||
if !groups.is_empty() {
|
||||
match connection
|
||||
.subscribe_auxiliary_live_filter_groups(groups)
|
||||
.await
|
||||
@@ -3831,18 +3935,19 @@ impl SyncManager {
|
||||
self.descendant_live_coverage.insert(
|
||||
relay_url.to_string(),
|
||||
DescendantLiveCoverage {
|
||||
frontier,
|
||||
coverage_fingerprint: desired_fingerprint.clone(),
|
||||
fallback_filters: rotated_filters.clone(),
|
||||
subscription_ids: subscription_ids.clone(),
|
||||
},
|
||||
);
|
||||
self.descendant_sync_rotations.remove(relay_url);
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
event_id_count,
|
||||
coordinate_count,
|
||||
filter_count,
|
||||
subscription_count = subscription_ids.len(),
|
||||
"Installed auxiliary descendant live coverage"
|
||||
rotated_filter_count = rotated_filters.len(),
|
||||
"Installed priority-bounded auxiliary live coverage"
|
||||
);
|
||||
// `limit:0` protects the future only. Pair every new
|
||||
// live frontier with a complete EOSE-closing baseline
|
||||
@@ -3850,13 +3955,26 @@ impl SyncManager {
|
||||
let _ = self
|
||||
.historic_sync_with_options(
|
||||
relay_url,
|
||||
historic_filters,
|
||||
historic_entries
|
||||
.iter()
|
||||
.map(|entry| entry.filter.clone())
|
||||
.collect(),
|
||||
PendingItems::default(),
|
||||
None,
|
||||
PendingBatchPurpose::Descendants,
|
||||
true,
|
||||
)
|
||||
.await;
|
||||
let rotation = self
|
||||
.descendant_sync_rotations
|
||||
.entry(relay_url.to_string())
|
||||
.or_default();
|
||||
let max_filters = connection.max_filters_per_req();
|
||||
if rotation.coverage_fingerprint
|
||||
!= rotation_fingerprint(&rotated_filters, max_filters)
|
||||
{
|
||||
rotation.refresh(rotated_filters, max_filters);
|
||||
}
|
||||
return;
|
||||
}
|
||||
Err(error) => {
|
||||
@@ -3871,22 +3989,31 @@ impl SyncManager {
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
filter_count = live_filters.len(),
|
||||
"Descendant live coverage does not fit while preserving historic capacity"
|
||||
"No auxiliary live tier fits while preserving historic capacity"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
let rotated_filters: Vec<_> = historic_entries
|
||||
.into_iter()
|
||||
.map(|entry| entry.filter)
|
||||
.collect();
|
||||
let max_filters = self
|
||||
.connections
|
||||
.get(relay_url)
|
||||
.map(|connection| connection.max_filters_per_req())
|
||||
.unwrap_or(MAX_FILTERS_PER_REQ);
|
||||
let rotation = self
|
||||
.descendant_sync_rotations
|
||||
.entry(relay_url.to_string())
|
||||
.or_default();
|
||||
if rotation.frontier != frontier {
|
||||
rotation.refresh(frontier);
|
||||
if rotation.coverage_fingerprint != rotation_fingerprint(&rotated_filters, max_filters) {
|
||||
rotation.refresh(rotated_filters, max_filters);
|
||||
}
|
||||
}
|
||||
|
||||
async fn start_descendant_fallback(&mut self, relay_url: &str) {
|
||||
let Some((filter_index, filter, until)) = self
|
||||
let Some((filter_index, filters, until)) = self
|
||||
.descendant_sync_rotations
|
||||
.get(relay_url)
|
||||
.and_then(|rotation| rotation.next_request(Timestamp::now()))
|
||||
@@ -3897,7 +4024,7 @@ impl SyncManager {
|
||||
if let Some(batch_id) = self
|
||||
.historic_sync_with_options(
|
||||
relay_url,
|
||||
vec![filter],
|
||||
filters,
|
||||
PendingItems::default(),
|
||||
None,
|
||||
PendingBatchPurpose::Descendants,
|
||||
@@ -3953,7 +4080,7 @@ impl SyncManager {
|
||||
self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1);
|
||||
self.reconcile_descendant_mode(
|
||||
&reconcile_relay,
|
||||
&targets[&reconcile_relay].root_events,
|
||||
&targets[&reconcile_relay],
|
||||
)
|
||||
.await;
|
||||
|
||||
@@ -3978,7 +4105,7 @@ impl SyncManager {
|
||||
let mut filters = vec![filters::build_announcement_filter(None)];
|
||||
let index = self.relay_sync_index.read().await;
|
||||
if let Some(state) = index.get(relay_url) {
|
||||
filters.extend(filters::build_sync_level_aware_filters(
|
||||
filters.extend(filters::build_essential_live_filters(
|
||||
&state.repos,
|
||||
&state.state_only_repos,
|
||||
&state.root_events,
|
||||
@@ -5649,8 +5776,7 @@ impl SyncManager {
|
||||
return false;
|
||||
}
|
||||
|
||||
let unbounded_groups =
|
||||
live_filter_groups(&complete_live, connection.max_filters_per_req());
|
||||
let unbounded_groups = live_filter_groups(&complete_live, connection.max_filters_per_req());
|
||||
let remote_limit = connection.remote_subscription_byte_limit();
|
||||
let (complete_groups, overflow_groups) =
|
||||
groups_within_subscription_byte_limit(unbounded_groups, remote_limit, 0);
|
||||
@@ -5717,10 +5843,15 @@ impl SyncManager {
|
||||
.close_live_subscriptions(&coverage.subscription_ids)
|
||||
.await;
|
||||
}
|
||||
let max_filters = self
|
||||
.connections
|
||||
.get(relay_url)
|
||||
.map(|connection| connection.max_filters_per_req())
|
||||
.unwrap_or(MAX_FILTERS_PER_REQ);
|
||||
self.descendant_sync_rotations
|
||||
.entry(relay_url.to_string())
|
||||
.or_default()
|
||||
.refresh(coverage.frontier);
|
||||
.refresh(coverage.fallback_filters, max_filters);
|
||||
tracing::warn!(
|
||||
relay = %relay_url,
|
||||
sub_id = %subscription_id,
|
||||
@@ -5748,10 +5879,7 @@ impl SyncManager {
|
||||
self.release_removed_descendant_batch(relay_url, batch);
|
||||
}
|
||||
let has_pending = self.has_pending_batches(relay_url).await;
|
||||
if self
|
||||
.deferred_consolidations
|
||||
.request(relay_url, has_pending)
|
||||
{
|
||||
if self.deferred_consolidations.request(relay_url, has_pending) {
|
||||
let _ = self.consolidate(relay_url).await;
|
||||
}
|
||||
return;
|
||||
@@ -5841,10 +5969,7 @@ impl SyncManager {
|
||||
self.health_tracker.record_policy_refusal(relay_url);
|
||||
if let Some(metrics) = &self.metrics {
|
||||
metrics.record_policy_refusal(relay_url, category.label());
|
||||
metrics.record_health_state(
|
||||
relay_url,
|
||||
self.health_tracker.get_state(relay_url),
|
||||
);
|
||||
metrics.record_health_state(relay_url, self.health_tracker.get_state(relay_url));
|
||||
}
|
||||
tracing::warn!(
|
||||
relay = %relay_url,
|
||||
@@ -5915,10 +6040,8 @@ impl SyncManager {
|
||||
};
|
||||
|
||||
if self.has_pending_batches(&relay_url).await {
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url,
|
||||
now + Duration::from_secs(10),
|
||||
);
|
||||
self.byte_limited_live_relays
|
||||
.insert(relay_url, now + Duration::from_secs(10));
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -5929,17 +6052,13 @@ impl SyncManager {
|
||||
.get(&relay_url)
|
||||
.is_some_and(|state| state.connection_status.is_live_sync_active());
|
||||
if !connected {
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url,
|
||||
now + byte_limited_catchup_interval(),
|
||||
);
|
||||
self.byte_limited_live_relays
|
||||
.insert(relay_url, now + byte_limited_catchup_interval());
|
||||
return;
|
||||
}
|
||||
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url.clone(),
|
||||
now + byte_limited_catchup_interval(),
|
||||
);
|
||||
self.byte_limited_live_relays
|
||||
.insert(relay_url.clone(), now + byte_limited_catchup_interval());
|
||||
let overlap = byte_limited_catchup_interval() + Duration::from_secs(60);
|
||||
let since = Timestamp::from(Timestamp::now().as_secs().saturating_sub(overlap.as_secs()));
|
||||
let filters = self.complete_live_filters(&relay_url, Some(since)).await;
|
||||
@@ -6880,31 +6999,74 @@ mod tests {
|
||||
let now = Timestamp::from_secs(200_000);
|
||||
let mut rotation = DescendantSyncRotation::default();
|
||||
|
||||
rotation.refresh(frontier);
|
||||
rotation.refresh(
|
||||
descendant_frontier_filters(&frontier, None),
|
||||
MAX_FILTERS_PER_REQ,
|
||||
);
|
||||
let (filter_index, first, until) = rotation.next_request(now).unwrap();
|
||||
assert!(serde_json::to_value(first).unwrap().get("since").is_none());
|
||||
assert!(first.iter().all(|filter| {
|
||||
serde_json::to_value(filter).unwrap().get("since").is_none()
|
||||
}));
|
||||
rotation.mark_started(41, filter_index, until);
|
||||
assert!(rotation.next_request(now).is_none());
|
||||
|
||||
assert!(rotation.mark_completed(41, false));
|
||||
let (retry_index, retry, retry_until) = rotation.next_request(now).unwrap();
|
||||
assert_eq!(retry_index, filter_index);
|
||||
assert!(serde_json::to_value(retry).unwrap().get("since").is_none());
|
||||
assert!(retry.iter().all(|filter| {
|
||||
serde_json::to_value(filter).unwrap().get("since").is_none()
|
||||
}));
|
||||
rotation.mark_started(42, retry_index, retry_until);
|
||||
assert!(rotation.mark_completed(42, true));
|
||||
|
||||
// Complete the other two e/E/q variants so the rotation returns to
|
||||
// the first filter with its successful cursor and overlap.
|
||||
for batch_id in [43, 44] {
|
||||
let (index, _, upper) = rotation.next_request(now).unwrap();
|
||||
rotation.mark_started(batch_id, index, upper);
|
||||
assert!(rotation.mark_completed(batch_id, true));
|
||||
}
|
||||
let (_, recent, _) = rotation.next_request(Timestamp::from_secs(201_000)).unwrap();
|
||||
// All three e/E/q variants share one relay-compatible REQ, so the
|
||||
// next turn returns to that group with its successful overlap cursor.
|
||||
let (_, recent, _) = rotation
|
||||
.next_request(Timestamp::from_secs(201_000))
|
||||
.unwrap();
|
||||
assert!(recent.iter().all(|filter| {
|
||||
serde_json::to_value(filter).unwrap()["since"]
|
||||
== serde_json::json!(now.as_secs() - DESCENDANT_FALLBACK_OVERLAP_SECS)
|
||||
}));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_cutoff_keeps_only_complete_priority_tiers_live() {
|
||||
use filters::{CoverageTier, TieredFilter};
|
||||
|
||||
let entries = vec![
|
||||
TieredFilter {
|
||||
tier: CoverageTier::RootUppercase,
|
||||
filter: Filter::new().kind(Kind::Custom(31_001)),
|
||||
},
|
||||
TieredFilter {
|
||||
tier: CoverageTier::CoreCompatibility,
|
||||
filter: Filter::new().kind(Kind::Custom(31_002)),
|
||||
},
|
||||
TieredFilter {
|
||||
tier: CoverageTier::CoreCompatibility,
|
||||
filter: Filter::new().kind(Kind::Custom(31_003)),
|
||||
},
|
||||
TieredFilter {
|
||||
tier: CoverageTier::DescendantCanonical,
|
||||
filter: Filter::new().kind(Kind::Custom(31_004)),
|
||||
},
|
||||
TieredFilter {
|
||||
tier: CoverageTier::HistoricOnly,
|
||||
filter: Filter::new().kind(Kind::Custom(31_005)),
|
||||
},
|
||||
];
|
||||
|
||||
let (live, rotated) = split_live_tier_prefix(&entries, 2, |groups| groups.len() <= 1);
|
||||
assert_eq!(
|
||||
serde_json::to_value(recent).unwrap()["since"],
|
||||
serde_json::json!(now.as_secs() - DESCENDANT_FALLBACK_OVERLAP_SECS)
|
||||
live.len(),
|
||||
1,
|
||||
"the two-filter compatibility tier must not be split"
|
||||
);
|
||||
assert_eq!(rotated.len(), 4);
|
||||
assert!(rotated.iter().any(|filter| {
|
||||
filter.as_json() == Filter::new().kind(Kind::Custom(31_005)).as_json()
|
||||
}));
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -6930,7 +7092,9 @@ mod tests {
|
||||
groups.len() > 1,
|
||||
"large complete descendant coverage must reserve several subscriptions"
|
||||
);
|
||||
assert!(groups.iter().all(|group| group.len() <= MAX_FILTERS_PER_REQ));
|
||||
assert!(groups
|
||||
.iter()
|
||||
.all(|group| group.len() <= MAX_FILTERS_PER_REQ));
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -7510,11 +7674,8 @@ mod tests {
|
||||
let second = vec![Filter::new().kind(Kind::Metadata).limit(0)];
|
||||
let limit = SUBSCRIPTION_BYTE_RESERVED_MARGIN + req_message_size(&first);
|
||||
|
||||
let (admitted, overflow) = groups_within_subscription_byte_limit(
|
||||
vec![first.clone(), second],
|
||||
Some(limit),
|
||||
0,
|
||||
);
|
||||
let (admitted, overflow) =
|
||||
groups_within_subscription_byte_limit(vec![first.clone(), second], Some(limit), 0);
|
||||
|
||||
assert_eq!(admitted, vec![first]);
|
||||
assert_eq!(overflow, 1);
|
||||
|
||||
@@ -139,7 +139,7 @@ async fn descendant_live_coverage_is_preferred_when_capacity_remains() {
|
||||
assert!(
|
||||
wait_for_log(
|
||||
&syncing.log_path(),
|
||||
"Installed auxiliary descendant live coverage",
|
||||
"Installed priority-bounded auxiliary live coverage",
|
||||
Duration::from_secs(20),
|
||||
)
|
||||
.await,
|
||||
|
||||
Reference in New Issue
Block a user