Merge #2b377e20: feat(sync): schedule direct-member descendant history

nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqszkdm7yq5lkte4ycw72r3esglyp5z2g70cdqm0gv9zkfgyrttk2kq0su8wr

PR-Author: DanConwayDev's Agent
nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0

CoverNote:

Closes #727b0992 with the capacity-aware hybrid design agreed in review.

Repository collaboration events may reference only their immediate parent. Existing repository/root filters can therefore miss replies to a discovered reply when those events omit repository and root tags.

The proposal is intentionally split into four reviewable commits:

- `d203f9d` refactors the existing historic REQ path around one shared queued submission helper, without changing scheduling.
- `5f84b8b` includes known direct-member descendants in ordinary historic coverage.
- `4dced40` retains separately consolidated descendant live coverage when core live coverage, a two-slot control reserve, and at least one transient historic slot all fit the relay per-connection budget.
- `9f3c943` otherwise rotates descendant filters through the same transient ledger, starting at most one queued or active descendant request per relay, on a five-second cadence.

Fallback requests use REQ+EOSE, advance their cursor only after successful EOSE, and overlap the prior successful upper bound by 15 minutes. They do not skip a filter merely because the transient slot is busy: the request waits in the existing queue. Permanent coverage is paired with the complete historic baseline so events created before installation are not missed.

The design deliberately does not add request-class priority, recursive descendant traversal, durable cursors, multi-connection sharding, new configuration, or minimum-churn core consolidation. The latter is tracked separately at nostr:nevent1qqswhdur5f7z8zwy726za854qhahnt9k5ytru2h6zqyysrzdskwzk9cpz3mhxue69uhhyetvv9ujumn8d96zuer9wcykq3rj.

Local validation:

- `cargo check --lib`: passed
- `cargo test --lib`: 673 passed
- forced-fallback integration at advertised capacity 5: passed
- permanent-live integration: passed
- 600-member unit scenario confirms descendant filters split across multiple live subscriptions
- full `cargo test --test sync`: 83 passed, 10 failed, 1 ignored; failures retain established load-sensitive signatures while both new descendant scenarios pass under whole-file load

Disposable archive production evidence:

- exact tip `9f3c943` and the behavior-identical pre-amend executable both started with zero restarts and no candidate-specific rate-limit, panic, or ledger-overrun errors
- dependency sync queued 2,293 IDs in 8 chunks and 34,868 IDs in 117 chunks without starving later work
- permanent descendant coverage was installed across real relays; examples include 24 filters split over 12 subscriptions on relay.damus.io, 9 over 3 on grasp.budabit.club, and 6 over 2 on git.nostrhub.io
- the naturally constrained git.shakespeare.diy session had a 20-slot fallback budget; 27 descendant filters could not fit alongside 1,279 repositories and 961 roots, so filters 0, 1, and 2 rotated sequentially and completed successfully without skipping the busy transient queue
- the refined exact tip remained stable with bounded memory (about 517 MB peak during its observed window)
- the stacked follow-up executable retains the same descendant behavior and is now exercising a 49,253-ID dependency batch on the same archive dataset without restart or rate-limit failures

Recommendation: ready to merge.
This commit is contained in:
DanConwayDev
2026-08-10 08:27:15 +01:00
10 changed files with 1015 additions and 68 deletions
+9
View File
@@ -7,6 +7,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased]
### Added
- Recover repository-event descendants which reference a direct thread member
but omit the repository and root-event tags. When the per-connection ledger
can retain complete descendant coverage while preserving control and
historic capacity, those filters stay live; otherwise bounded EOSE-closing
history provides eventual coverage. The frontier is deliberately
non-recursive and connection sharding remains unnecessary.
## [2.1.2] - 2026-08-08
ngit-grasp 2.1.2 is a patch release improving repository-event sync under
@@ -11,6 +11,9 @@ Features:
- Fetches all repository announcements from connected relays to discover new repos listing our service
- Discovers and dynamically connects to new relays listed by repository announcements we have accepted (with optional bootstrap relay to get started)
- Fetches events tagging repositories we are interested in, as well as events tagging Issues, Patches and PRs of these repositories
- Recovers one additional generation of events that tag those direct thread
members but omit repository and root-event tags, using scheduled history
queries rather than retained subscriptions
- Supports live sync and historic sync (tries NIP-77 negentropy but falls back to REQ+EOSE with 'until' based pagination)
- Plays nicely with other relays - connection backoff and rate-limiting detection with cooldown
- Does a full reconciliation daily
@@ -785,6 +788,37 @@ live in
- **Function**: `build_root_event_tag_filters(root_events, since)`
- **Only for `SyncLevel::Full` repos** — purgatory announcements (`StateOnly`) skip this layer
### Direct-Member Descendants
Some collaboration events reference only their immediate parent. Once the
ordinary Layer 3 filters have discovered an event which directly tags a
repository root, each source relay is also queried for events whose `e`, `E`,
or `q` tags reference that direct member.
- complete descendant filters are retained live when they fit after core live
coverage while preserving the two control-plane slots and at least one
transient historic slot;
- auxiliary subscriptions are tracked separately and retired before core
consolidation or restoration, so they never enter the core rollback set;
- when the complete live set does not fit, one EOSE-closing filter starts on
each constrained relay on the existing five-second maintenance cadence,
through the ordinary historic queue, pagination, shared ledger, and request
pacing;
- ordinary historic batches include the currently known direct-member filters;
- fallback filters keep an in-memory cursor, advance it only after successful
EOSE, and query from the preceding successful upper bound with 15 minutes of
overlap; and
- recovered descendants never enter the frontier, so this is deliberately one
additional generation rather than recursive thread traversal.
An unexpected auxiliary CLOSED retires the remaining descendant subscriptions
without rebuilding core coverage and falls back to history. A filter already
queued or active blocks another fallback filter for that relay; failure leaves
the same filter and cursor at the head. Reconnect and daily reconciliation
reconstruct the mode from current session capacity. This keeps the feature
complete without durable cursors, another capacity ledger, request-class
priority, or multi-connection sharding.
### Combined Layer 2+3 (SyncLevel-Aware)
The `build_sync_level_aware_filters()` function combines both layers, partitioning repos by `SyncLevel`:
+21 -8
View File
@@ -301,24 +301,31 @@ Derived from the tightest commonly observed values; all sizing below assumes:
## Our Approach: A Per-Connection Budget Ledger
Each relay connection owns one implemented budget ledger of B subscription slots. Three
consumers share it, in priority order:
Each relay connection owns one implemented budget ledger of B subscription
slots. Four consumers share it, in priority order:
1. **Live subscriptions** (persistent, `limit: 0`) — the product; sized first.
1. **Core live subscriptions** (persistent, `limit: 0`) — the product; sized
first.
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) — negentropy rounds,
3. **Historic sync and dependency recovery** (transient) — at least one usable
slot remains after live admission; negentropy rounds,
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.
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. Live filter groups
are packed first and admitted atomically against the advertised
subscription-count budget: if the complete live set cannot fit,
(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 live set still cannot fit,
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
@@ -346,6 +353,12 @@ live coverage. Each reconnect closes the retired ledger and creates a new
generation; queued or late borrowers therefore fail before sending on the new
SDK session and cannot inflate or bypass its capacity.
Descendant live subscriptions are kept outside the core rollback set. Core
consolidation or restoration first closes them, and aborts if CLOSE cannot be
sent, so auxiliary coverage cannot silently consume capacity needed by newly
required core filters. An auxiliary CLOSED retires its remaining group and
falls back to EOSE-closing history without rebuilding healthy core coverage.
Some relays additionally cap the cumulative serialized REQ state retained by
one connection. NIP-11 has no field for this limit, so it cannot be negotiated
before the first refusal. A CLOSED reason of the rust-nostr form `active
+3
View File
@@ -485,6 +485,7 @@ mod tests {
"wss://relay1.com".to_string(),
vec![super::super::PendingBatch {
batch_id: 1,
purpose: super::super::PendingBatchPurpose::Core,
items: super::super::PendingItems {
repos: vec!["repo1".to_string()].into_iter().collect(),
state_only_repos: HashSet::new(),
@@ -592,6 +593,7 @@ mod tests {
relay_url.to_string(),
vec![super::super::PendingBatch {
batch_id: 1,
purpose: super::super::PendingBatchPurpose::Core,
items: PendingItems {
repos: HashSet::new(),
state_only_repos: HashSet::from(["repo1".to_string()]),
@@ -681,6 +683,7 @@ mod tests {
"wss://relay1.com".to_string(),
vec![super::super::PendingBatch {
batch_id: 1,
purpose: super::super::PendingBatchPurpose::Core,
items: super::super::PendingItems {
repos: vec!["repo1".to_string()].into_iter().collect(),
state_only_repos: HashSet::new(),
+2 -1
View File
@@ -216,7 +216,8 @@ pub fn tagged_one_of_our_root_event_filters(
);
let mut filters = Vec::new();
let event_ids: Vec<String> = root_events.iter().map(|id| id.to_hex()).collect();
let mut event_ids: Vec<String> = root_events.iter().map(|id| id.to_hex()).collect();
event_ids.sort_unstable();
for (chunk_idx, chunk) in chunk_values_by_bytes(&event_ids).into_iter().enumerate() {
// Lowercase 'e' tag - standard event reference
+636 -59
View File
@@ -65,6 +65,7 @@ const MAX_PURGATORY_FILTER_ACTIONS_PER_TICK: usize = 1;
const MAX_PURGATORY_DEPENDENCY_IDS_PER_QUERY: usize = 100;
const SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS: usize = 20;
const SEMANTIC_FALLBACK_MAX_DELIVERED_PERCENT: usize = 10;
const DESCENDANT_FALLBACK_OVERLAP_SECS: u64 = 15 * 60;
fn should_use_semantic_fallback(requested_count: usize, received_count: usize) -> bool {
requested_count >= SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS
@@ -662,6 +663,9 @@ impl RelayPaginationSession {
pub struct PendingBatch {
/// Unique ID for this batch - for debugging/logging
pub batch_id: u64,
/// Why this batch exists. Auxiliary descendant discovery must not mutate
/// core historic-sync completion or health state.
pub purpose: PendingBatchPurpose,
/// The items this batch is syncing
pub items: PendingItems,
/// Subscription IDs that must ALL receive EOSE before confirming (for ReqEose)
@@ -688,6 +692,108 @@ pub struct PendingBatch {
pub failed: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PendingBatchPurpose {
Core,
Announcements,
Descendants,
}
#[derive(Debug)]
struct DescendantSyncRotation {
frontier: HashSet<EventId>,
filters: Vec<DescendantFilterCursor>,
next_filter: usize,
in_flight: Option<DescendantFilterInFlight>,
}
#[derive(Debug)]
struct DescendantFilterCursor {
filter: Filter,
last_successful_until: Option<Timestamp>,
}
#[derive(Debug, Clone, Copy)]
struct DescendantFilterInFlight {
batch_id: u64,
filter_index: usize,
until: Timestamp,
}
#[derive(Debug)]
struct DescendantLiveCoverage {
frontier: HashSet<EventId>,
subscription_ids: Vec<SubscriptionId>,
}
impl Default for DescendantSyncRotation {
fn default() -> Self {
Self {
frontier: HashSet::new(),
filters: Vec::new(),
next_filter: 0,
in_flight: None,
}
}
}
impl DescendantSyncRotation {
fn refresh(&mut self, members: HashSet<EventId>) {
let previous: HashMap<String, Option<Timestamp>> = self
.filters
.drain(..)
.map(|cursor| (cursor.filter.as_json(), cursor.last_successful_until))
.collect();
self.frontier = members.clone();
self.filters = filters::tagged_one_of_our_root_event_filters(&members, None)
.into_iter()
.map(|filter| DescendantFilterCursor {
last_successful_until: previous.get(&filter.as_json()).copied().flatten(),
filter,
})
.collect();
self.next_filter = 0;
self.in_flight = None;
}
fn next_request(&self, now: Timestamp) -> Option<(usize, 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))
}
fn mark_started(&mut self, batch_id: u64, filter_index: usize, until: Timestamp) {
self.in_flight = Some(DescendantFilterInFlight {
batch_id,
filter_index,
until,
});
}
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 {
return false;
};
self.in_flight = None;
if succeeded {
self.filters[in_flight.filter_index].last_successful_until = Some(in_flight.until);
self.next_filter = (in_flight.filter_index + 1) % self.filters.len();
}
true
}
}
/// Items included in a pending batch
#[derive(Debug, Clone, Default)]
pub struct PendingItems {
@@ -1172,6 +1278,10 @@ async fn run_daily_timer(
/// during negentropy reconciliation but failed to deliver on exact-ID fetches
/// (see [`missing_events`]). Recovery attempts are backed off per relay, so
/// the tick itself stays cheap when nothing is due.
///
/// Finally, one relay's live descendant admission is reconciled and every
/// constrained relay may advance one queued descendant history query. Reusing
/// this five-second cadence avoids another scheduler or configuration surface.
async fn run_purgatory_announcement_sync(
sync_manager: Arc<Mutex<SyncManager>>,
mut shutdown_rx: broadcast::Receiver<()>,
@@ -1187,6 +1297,7 @@ async fn run_purgatory_announcement_sync(
let mut manager = sync_manager.lock().await;
manager.sync_purgatory_announcements_to_index().await;
manager.tick_missing_event_recovery().await;
manager.tick_descendant_sync().await;
}
_ = shutdown_rx.recv() => {
tracing::debug!("Purgatory announcement sync timer received shutdown signal");
@@ -1415,6 +1526,14 @@ pub struct SyncManager {
/// Relays whose complete persistent filter set exceeds a learned remote
/// byte cap, mapped to their next bounded catch-up deadline.
byte_limited_live_relays: HashMap<String, Instant>,
/// Per-relay snapshots for low-priority, EOSE-closing descendant queries.
descendant_sync_rotations: HashMap<String, DescendantSyncRotation>,
/// Auxiliary persistent descendant coverage, kept separate from the core
/// desired set so it can be retired without replacing healthy core REQs.
descendant_live_coverage: HashMap<String, DescendantLiveCoverage>,
/// Round-robin cursor so only one relay performs descendant live-mode
/// derivation and admission work per maintenance tick.
descendant_relay_cursor: usize,
/// Channel for disconnect notifications (set during run)
disconnect_tx: Option<tokio::sync::mpsc::Sender<DisconnectNotification>>,
/// Channel for EOSE notifications (set during run)
@@ -1515,6 +1634,9 @@ impl SyncManager {
connect_attempt_semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_CONNECT_ATTEMPTS)),
deferred_consolidations: DeferredConsolidations::default(),
byte_limited_live_relays: HashMap::new(),
descendant_sync_rotations: HashMap::new(),
descendant_live_coverage: HashMap::new(),
descendant_relay_cursor: 0,
disconnect_tx: None,
eose_tx: None,
subscription_closed_tx: None,
@@ -2507,14 +2629,26 @@ impl SyncManager {
/// # Arguments
/// * `relay_url` - The relay URL the batch belongs to
/// * `batch` - The completed batch to confirm
async fn confirm_batch(&self, relay_url: &str, batch: PendingBatch) {
async fn confirm_batch(&mut self, relay_url: &str, batch: PendingBatch) {
let batch_id = batch.batch_id;
let full_repos_count = batch.items.repos.len();
let state_only_repos_count = batch.items.state_only_repos.len();
let events_count = batch.items.root_events.len();
let sync_method = batch.sync_method;
let is_generic_filter =
full_repos_count == 0 && state_only_repos_count == 0 && events_count == 0;
let is_generic_filter = batch.purpose == PendingBatchPurpose::Announcements;
if batch.purpose == PendingBatchPurpose::Descendants {
let succeeded = !batch.failed;
if let Some(rotation) = self.descendant_sync_rotations.get_mut(relay_url) {
rotation.mark_completed(batch_id, succeeded);
}
tracing::info!(
relay = %relay_url,
batch_id,
succeeded,
"Descendant historic query reached terminal batch state"
);
}
let mut relay_index = self.relay_sync_index.write().await;
@@ -2567,7 +2701,7 @@ impl SyncManager {
}
// Track if this batch failed (for ConnectedDegraded transition)
if batch.failed {
if batch.failed && batch.purpose != PendingBatchPurpose::Descendants {
state.historic_sync_had_failures = true;
// Failures unrelated to a relay's pending missing-event
// recovery mean full recovery must not restore its health.
@@ -2585,6 +2719,7 @@ impl SyncManager {
tracing::info!(
relay = %relay_url,
batch_id = batch_id,
purpose = ?batch.purpose,
sync_method = ?sync_method,
full_repos_confirmed = full_repos_count,
state_only_repos_confirmed = state_only_repos_count,
@@ -2735,6 +2870,10 @@ impl SyncManager {
async fn daily_sync(&mut self, relay_url: &str) {
tracing::info!(relay = %relay_url, "Starting daily sync");
self.cancel_deferred_consolidation(relay_url, "daily sync reset");
let _ = self
.close_descendant_live_coverage(relay_url, "daily sync reset")
.await;
self.descendant_sync_rotations.remove(relay_url);
// Get connection
let connection = match self.connections.get(relay_url) {
@@ -3528,16 +3667,276 @@ impl SyncManager {
// Use historic_sync with empty PendingItems for generic filters
// Generic filters (announcements) don't have associated repos or root_events
let items = PendingItems::default();
let _batch_id = self.historic_sync(relay_url, filters, items, since).await;
let _batch_id = self
.historic_sync_with_options(
relay_url,
filters,
items,
since,
PendingBatchPurpose::Announcements,
false,
)
.await;
}
async fn sync_generic_history(&mut self, relay_url: &str, since: Option<Timestamp>) {
let filters = vec![filters::build_announcement_filter(None)];
let _batch_id = self
.historic_sync(relay_url, filters, PendingItems::default(), since)
.historic_sync_with_options(
relay_url,
filters,
PendingItems::default(),
since,
PendingBatchPurpose::Announcements,
false,
)
.await;
}
/// Find the existing first-generation members of repository root threads.
///
/// Only events that directly reference a root are admitted to this
/// frontier. Events fetched by the descendant rotation are deliberately
/// not fed back into it, keeping this feature non-recursive.
async fn direct_thread_members(&self, root_events: &HashSet<EventId>) -> HashSet<EventId> {
let mut members = HashSet::new();
for filter in filters::tagged_one_of_our_root_event_filters(root_events, None) {
match self.database.query(filter).await {
Ok(events) => {
members.extend(
events
.iter()
.map(|event| event.id)
.filter(|event_id| !root_events.contains(event_id)),
);
}
Err(error) => {
tracing::warn!(
error = %error,
root_event_count = root_events.len(),
"Failed to derive direct repository thread members"
);
return HashSet::new();
}
}
}
members
}
async fn close_descendant_live_coverage(
&mut self,
relay_url: &str,
reason: &'static str,
) -> bool {
let Some(coverage) = self.descendant_live_coverage.remove(relay_url) else {
return true;
};
if let Some(connection) = self.connections.get(relay_url) {
if let Err(error) = connection
.close_live_subscriptions(&coverage.subscription_ids)
.await
{
tracing::warn!(
relay = %relay_url,
%error,
reason,
"Could not retire auxiliary descendant live coverage"
);
self.descendant_live_coverage
.insert(relay_url.to_string(), coverage);
return false;
}
}
tracing::info!(
relay = %relay_url,
subscription_count = coverage.subscription_ids.len(),
reason,
"Retired auxiliary descendant live coverage"
);
true
}
async fn reconcile_descendant_mode(
&mut self,
relay_url: &str,
root_events: &HashSet<EventId>,
) {
let members = self.direct_thread_members(root_events).await;
if members.is_empty() {
let _ = self
.close_descendant_live_coverage(relay_url, "frontier became empty")
.await;
self.descendant_sync_rotations.remove(relay_url);
return;
}
if self
.descendant_live_coverage
.get(relay_url)
.is_some_and(|coverage| coverage.frontier == members)
{
return;
}
if self.descendant_live_coverage.contains_key(relay_url)
&& !self
.close_descendant_live_coverage(relay_url, "frontier changed")
.await
{
return;
}
let live_since = Timestamp::from(
Timestamp::now()
.as_secs()
.saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS),
);
let historic_filters = filters::tagged_one_of_our_root_event_filters(&members, None);
let live_filters =
filters::tagged_one_of_our_root_event_filters(&members, Some(live_since));
if let Some(connection) = self.connections.get(relay_url).cloned() {
let groups = live_filter_groups(&live_filters, connection.max_filters_per_req());
if connection.can_admit_auxiliary_live_groups(&groups) {
match connection
.subscribe_auxiliary_live_filter_groups(groups)
.await
{
Ok(subscription_ids) => {
let filter_count = live_filters.len();
self.descendant_live_coverage.insert(
relay_url.to_string(),
DescendantLiveCoverage {
frontier: members,
subscription_ids: subscription_ids.clone(),
},
);
self.descendant_sync_rotations.remove(relay_url);
tracing::info!(
relay = %relay_url,
filter_count,
subscription_count = subscription_ids.len(),
"Installed auxiliary descendant live coverage"
);
// `limit:0` protects the future only. Pair every new
// live frontier with a complete EOSE-closing baseline
// so events stored before admission are not skipped.
let _ = self
.historic_sync_with_options(
relay_url,
historic_filters,
PendingItems::default(),
None,
PendingBatchPurpose::Descendants,
true,
)
.await;
return;
}
Err(error) => {
tracing::warn!(
relay = %relay_url,
%error,
"Descendant live admission failed; using historic rotation"
);
}
}
} else {
tracing::info!(
relay = %relay_url,
filter_count = live_filters.len(),
"Descendant live coverage does not fit while preserving historic capacity"
);
}
}
let rotation = self
.descendant_sync_rotations
.entry(relay_url.to_string())
.or_default();
if rotation.frontier != members {
rotation.refresh(members);
}
}
async fn start_descendant_fallback(&mut self, relay_url: &str) {
let Some((filter_index, filter, until)) = self
.descendant_sync_rotations
.get(relay_url)
.and_then(|rotation| rotation.next_request(Timestamp::now()))
else {
return;
};
let filter_count = self.descendant_sync_rotations[relay_url].filters.len();
if let Some(batch_id) = self
.historic_sync_with_options(
relay_url,
vec![filter],
PendingItems::default(),
None,
PendingBatchPurpose::Descendants,
true,
)
.await
{
if let Some(rotation) = self.descendant_sync_rotations.get_mut(relay_url) {
rotation.mark_started(batch_id, filter_index, until);
}
tracing::info!(
relay = %relay_url,
batch_id,
filter_index,
filter_count,
until = until.as_secs(),
"Started queued descendant fallback query"
);
}
}
/// Reconcile one live-mode decision and advance every constrained relay by
/// at most one EOSE-closing request on each five-second maintenance tick.
async fn tick_descendant_sync(&mut self) {
let targets = {
let index = self.repo_sync_index.read().await;
algorithms::derive_relay_targets(&index)
};
let states = self.relay_sync_index.read().await;
let mut relay_urls: Vec<String> = targets
.iter()
.filter(|(relay_url, needs)| {
!needs.root_events.is_empty()
&& states.get(*relay_url).is_some_and(|state| {
matches!(
state.connection_status,
ConnectionStatus::Connected
| ConnectionStatus::ConnectedHistoricSyncFailures
)
})
&& !self.health_tracker.is_subscription_paused(relay_url)
})
.map(|(relay_url, _)| relay_url.clone())
.collect();
drop(states);
relay_urls.sort_unstable();
if relay_urls.is_empty() {
return;
}
let reconcile_relay =
relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone();
self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1);
self.reconcile_descendant_mode(
&reconcile_relay,
&targets[&reconcile_relay].root_events,
)
.await;
let constrained: Vec<String> = relay_urls
.into_iter()
.filter(|relay_url| self.descendant_sync_rotations.contains_key(relay_url))
.collect();
for relay_url in constrained {
self.start_descendant_fallback(&relay_url).await;
}
}
/// Build the complete persistent coverage for one connection. L1 must be
/// admitted in the same transaction as rebuilt L2/L3 so a tight relay
/// budget cannot leave a successful generic REQ hiding partial repo
@@ -4538,6 +4937,11 @@ impl SyncManager {
async fn handle_disconnect(&mut self, relay_url: &str) {
// Learned page sizes and NIP-11 hints belong to the ended WebSocket session.
self.pagination_sessions.remove(relay_url);
// A connection can end after a descendant REQ was sent but before its
// EOSE. Restart the bounded full rotation so that page is not treated
// as complete on the next session.
self.descendant_sync_rotations.remove(relay_url);
self.descendant_live_coverage.remove(relay_url);
// Check if this was an intentional disconnect (Disconnecting status)
let was_intentional = {
@@ -5250,6 +5654,16 @@ impl SyncManager {
"Starting consolidation"
);
// Descendant coverage is auxiliary and independently reconstructible.
// Retire it before capturing the core rollback set so it can neither
// displace newly required core coverage nor become part of that set.
if !self
.close_descendant_live_coverage(relay_url, "core consolidation")
.await
{
return false;
}
let now = Timestamp::now();
let since = Timestamp::from(now.as_secs().saturating_sub(QUICK_RECONNECT_WINDOW_SECS));
let complete_live = self.complete_live_filters(relay_url, Some(since)).await;
@@ -5315,6 +5729,14 @@ impl SyncManager {
true
}
fn release_removed_descendant_batch(&mut self, relay_url: &str, batch: &PendingBatch) {
if batch.purpose == PendingBatchPurpose::Descendants {
if let Some(rotation) = self.descendant_sync_rotations.get_mut(relay_url) {
rotation.mark_completed(batch.batch_id, false);
}
}
}
async fn handle_subscription_closed(
&mut self,
relay_url: &str,
@@ -5323,6 +5745,33 @@ impl SyncManager {
live_generation: Option<u64>,
live_filter_count: Option<usize>,
) {
let is_descendant_live = self
.descendant_live_coverage
.get(relay_url)
.is_some_and(|coverage| coverage.subscription_ids.contains(&subscription_id));
if is_descendant_live {
let coverage = self
.descendant_live_coverage
.remove(relay_url)
.expect("descendant subscription belonged to tracked coverage");
if let Some(connection) = self.connections.get(relay_url) {
let _ = connection
.close_live_subscriptions(&coverage.subscription_ids)
.await;
}
self.descendant_sync_rotations
.entry(relay_url.to_string())
.or_default()
.refresh(coverage.frontier);
tracing::warn!(
relay = %relay_url,
sub_id = %subscription_id,
reason,
"Descendant live subscription closed; using historic fallback"
);
return;
}
if let Some(limit) = subscription_state_byte_limit(reason) {
tracing::warn!(
relay = %relay_url,
@@ -5331,9 +5780,14 @@ impl SyncManager {
);
self.byte_limited_live_relays
.insert(relay_url.to_string(), Instant::now());
if live_generation.is_none() {
let removed_batch = if live_generation.is_none() {
let mut pending = self.pending_sync_index.write().await;
take_batch_containing_subscription(&mut pending, relay_url, &subscription_id);
take_batch_containing_subscription(&mut pending, relay_url, &subscription_id)
} else {
None
};
if let Some(batch) = &removed_batch {
self.release_removed_descendant_batch(relay_url, batch);
}
let has_pending = self.has_pending_batches(relay_url).await;
if self
@@ -5351,6 +5805,7 @@ impl SyncManager {
};
if let Some(batch) = removed_batch {
self.release_removed_descendant_batch(relay_url, &batch);
tracing::warn!(
relay = %relay_url,
sub_id = %subscription_id,
@@ -5405,7 +5860,8 @@ impl SyncManager {
&subscription_id,
)
};
if removed_batch.is_some() {
if let Some(batch) = &removed_batch {
self.release_removed_descendant_batch(relay_url, batch);
self.recompute_new_sync_filters_for_relay(relay_url).await;
} else if let Some(generation) = live_generation {
self.restore_live_coverage_after_closed(relay_url, generation)
@@ -5421,6 +5877,9 @@ impl SyncManager {
let mut pending = self.pending_sync_index.write().await;
take_batch_containing_subscription(&mut pending, relay_url, &subscription_id)
};
if let Some(batch) = &removed_batch {
self.release_removed_descendant_batch(relay_url, batch);
}
self.health_tracker.record_policy_refusal(relay_url);
if let Some(metrics) = &self.metrics {
metrics.record_policy_refusal(relay_url, category.label());
@@ -5448,13 +5907,23 @@ impl SyncManager {
}
async fn restore_live_coverage_after_closed(&mut self, relay_url: &str, generation: u64) {
let Some(connection) = self.connections.get(relay_url) else {
let Some(current_generation) = self
.connections
.get(relay_url)
.map(RelayConnection::current_subscription_generation)
else {
return;
};
if connection.current_subscription_generation() != generation {
if current_generation != generation {
tracing::debug!(relay = %relay_url, generation, "Ignoring stale live CLOSED from retired session");
return;
}
if !self
.close_descendant_live_coverage(relay_url, "core live restoration")
.await
{
return;
}
let now = Timestamp::now();
let since = Timestamp::from(now.as_secs().saturating_sub(QUICK_RECONNECT_WINDOW_SECS));
let filters = self.complete_live_filters(relay_url, Some(since)).await;
@@ -5856,6 +6325,69 @@ impl SyncManager {
})
}
/// Submit grouped REQ+EOSE filters through the connection's transient
/// queue and return the subscriptions that actually started.
///
/// Callers submit only the current group while the connection waits for a
/// shared-ledger permit, rather than pre-queuing an unbounded historic
/// batch. The connection owns each permit until EOSE/CLOSED (or watchdog
/// recovery).
async fn subscribe_historic_filter_groups(
&self,
relay_url: &str,
batch_id: u64,
filters: &[Filter],
) -> (
HashSet<SubscriptionId>,
HashMap<SubscriptionId, PaginationState>,
) {
let mut subscription_ids = HashSet::new();
let mut pagination_state = HashMap::new();
let max_filters = self
.connections
.get(relay_url)
.map(RelayConnection::max_filters_per_req)
.unwrap_or(MAX_FILTERS_PER_REQ);
for (group_idx, filter_group) in group_filters_for_req_with_max(filters, max_filters)
.into_iter()
.enumerate()
{
tracing::debug!(
relay = %relay_url,
batch_id,
group_idx,
filter_count = filter_group.len(),
filters = ?filter_group,
"Subscribing to grouped filters in REQ+EOSE path"
);
if let Some(connection) = self.connections.get(relay_url) {
match connection
.subscribe_filters(filter_group.clone(), TransientRequestClass::HistoricPage)
.await
{
Ok(subscription_id) => {
subscription_ids.insert(subscription_id.clone());
pagination_state
.insert(subscription_id, PaginationState::new(filter_group));
}
Err(error) => {
tracing::error!(
relay = %relay_url,
batch_id,
group_idx,
error = %error,
"Failed to subscribe to filter in historic_sync"
);
}
}
}
}
(subscription_ids, pagination_state)
}
/// Sync historical events and track in PendingSyncIndex
///
/// This method handles historical synchronization for a set of filters,
@@ -5877,11 +6409,37 @@ impl SyncManager {
/// * `Some(batch_id)` - Batch was created and sync initiated
/// * `None` - No connection or sync failed to start
async fn historic_sync(
&mut self,
relay_url: &str,
mut filters: Vec<Filter>,
items: PendingItems,
since: Option<Timestamp>,
) -> Option<u64> {
if !items.root_events.is_empty() {
let members = self.direct_thread_members(&items.root_events).await;
filters.extend(filters::tagged_one_of_our_root_event_filters(
&members, None,
));
}
self.historic_sync_with_options(
relay_url,
filters,
items,
since,
PendingBatchPurpose::Core,
false,
)
.await
}
async fn historic_sync_with_options(
&mut self,
relay_url: &str,
filters: Vec<Filter>,
items: PendingItems,
since: Option<Timestamp>,
purpose: PendingBatchPurpose,
force_req_eose: bool,
) -> Option<u64> {
// DEBUG TRACING: Log all filters being passed to historic_sync
tracing::debug!(
@@ -5922,8 +6480,9 @@ impl SyncManager {
};
// Check if we should use negentropy
let use_negentropy =
!self.config.sync_disable_negentropy && connection.supports_negentropy().await;
let use_negentropy = !force_req_eose
&& !self.config.sync_disable_negentropy
&& connection.supports_negentropy().await;
// Generate batch ID
let batch_id = self.next_batch_id();
@@ -5945,6 +6504,7 @@ impl SyncManager {
// Create PendingBatch for negentropy (empty outstanding_subs and pagination_state)
let batch = PendingBatch {
batch_id,
purpose,
items: items.clone(),
outstanding_subs: HashSet::new(),
sync_method: SyncMethod::Negentropy,
@@ -6162,53 +6722,10 @@ impl SyncManager {
"Starting historic_sync with REQ+EOSE"
);
// Subscribe to each filter and collect subscription IDs
let mut subscription_ids = HashSet::new();
let mut pagination_state = HashMap::new();
// Keep several OR filters under each relay-visible subscription.
let max_filters = self
.connections
.get(relay_url)
.map(RelayConnection::max_filters_per_req)
.unwrap_or(MAX_FILTERS_PER_REQ);
for (idx, filter_group) in
group_filters_for_req_with_max(&filters_with_since, max_filters)
.into_iter()
.enumerate()
{
tracing::debug!(
relay = %relay_url,
batch_id = batch_id,
group_idx = idx,
filter_count = filter_group.len(),
filters = ?filter_group,
"Subscribing to grouped filters in REQ+EOSE path"
);
if let Some(conn) = self.connections.get(relay_url) {
let grouped_filters = filter_group;
match conn
.subscribe_filters(
grouped_filters.clone(),
TransientRequestClass::HistoricPage,
)
.await
{
Ok(sub_id) => {
subscription_ids.insert(sub_id.clone());
pagination_state.insert(sub_id, PaginationState::new(grouped_filters));
}
Err(e) => {
tracing::error!(
relay = %relay_url,
error = %e,
"Failed to subscribe to filter in historic_sync"
);
}
}
}
}
let (subscription_ids, pagination_state) = self
.subscribe_historic_filter_groups(relay_url, batch_id, &filters_with_since)
.await;
if subscription_ids.is_empty() && !filters_with_since.is_empty() {
tracing::warn!(
@@ -6221,6 +6738,7 @@ impl SyncManager {
// Create PendingBatch for REQ+EOSE
let batch = PendingBatch {
batch_id,
purpose,
items,
outstanding_subs: subscription_ids,
sync_method: SyncMethod::ReqEose,
@@ -6329,6 +6847,60 @@ mod tests {
assert!(!should_use_semantic_fallback(19, 0));
}
#[test]
fn descendant_rotation_advances_cursor_only_after_successful_eose() {
let member = EventId::from_byte_array([7; 32]);
let members = HashSet::from([member]);
let now = Timestamp::from_secs(200_000);
let mut rotation = DescendantSyncRotation::default();
rotation.refresh(members);
let (filter_index, first, until) = rotation.next_request(now).unwrap();
assert!(serde_json::to_value(first).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());
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();
assert_eq!(
serde_json::to_value(recent).unwrap()["since"],
serde_json::json!(now.as_secs() - DESCENDANT_FALLBACK_OVERLAP_SECS)
);
}
#[test]
fn large_descendant_frontier_splits_across_live_subscriptions() {
let members: HashSet<EventId> = (0..600u32)
.map(|index| {
let mut bytes = [0u8; 32];
bytes[..4].copy_from_slice(&index.to_be_bytes());
EventId::from_byte_array(bytes)
})
.collect();
let filters = filters::tagged_one_of_our_root_event_filters(&members, None);
let groups = live_filter_groups(&filters, MAX_FILTERS_PER_REQ);
assert!(filters.len() > 3, "the frontier must span byte chunks");
assert!(
groups.len() > 1,
"large complete descendant coverage must reserve several subscriptions"
);
assert!(groups.iter().all(|group| group.len() <= MAX_FILTERS_PER_REQ));
}
#[test]
fn group_filters_for_req_respects_count_and_byte_budgets() {
// Many small filters group by the count cap.
@@ -6843,6 +7415,7 @@ mod tests {
let still_pending = SubscriptionId::new("still-pending");
let make_batch = |batch_id, outstanding_subs| PendingBatch {
batch_id,
purpose: PendingBatchPurpose::Core,
items: PendingItems::default(),
outstanding_subs,
sync_method: SyncMethod::ReqEose,
@@ -6945,6 +7518,7 @@ mod tests {
let accepted = SubscriptionId::new("accepted");
let make_batch = |batch_id, subscription_id| PendingBatch {
batch_id,
purpose: PendingBatchPurpose::Core,
items: PendingItems::default(),
outstanding_subs: HashSet::from([subscription_id]),
sync_method: SyncMethod::ReqEose,
@@ -6980,6 +7554,7 @@ mod tests {
let completed_sub = SubscriptionId::new("completed-page");
let mut batch = PendingBatch {
batch_id: 73,
purpose: PendingBatchPurpose::Core,
items: PendingItems::default(),
outstanding_subs: HashSet::new(),
sync_method: SyncMethod::ReqEose,
@@ -7298,6 +7873,7 @@ mod tests {
// Test that PendingBatch properly tracks negentropy-specific fields
let batch = PendingBatch {
batch_id: 1,
purpose: PendingBatchPurpose::Core,
items: PendingItems::default(),
outstanding_subs: HashSet::new(),
sync_method: SyncMethod::Negentropy,
@@ -7321,6 +7897,7 @@ mod tests {
// Test that REQ+EOSE batches don't use negentropy fields
let batch = PendingBatch {
batch_id: 1,
purpose: PendingBatchPurpose::Core,
items: PendingItems::default(),
outstanding_subs: HashSet::new(),
sync_method: SyncMethod::ReqEose,
+124
View File
@@ -1170,6 +1170,100 @@ impl RelayConnection {
<= usable
}
/// Whether auxiliary persistent coverage fits while preserving one usable
/// ledger slot for transient history. The two control-plane slots have
/// already been removed from `usable` when the session ledger was built.
pub fn can_admit_auxiliary_live_groups(&self, filter_groups: &[Vec<Filter>]) -> bool {
let usable = self
.subscription_usable_slots
.load(std::sync::atomic::Ordering::Relaxed);
let held = self
.live_req_permits_held
.lock()
.expect("live permit map poisoned");
if held
.len()
.saturating_add(filter_groups.len())
.saturating_add(1)
> usable
{
return false;
}
let remote_limit = self.remote_subscription_byte_limit();
if remote_limit.is_none() {
return true;
}
let existing_bytes: usize = held
.values()
.map(|held| {
ClientMessage::req(SubscriptionId::generate(), held.filters.clone())
.as_json()
.len()
})
.sum();
let auxiliary_bytes: usize = filter_groups
.iter()
.map(|filters| {
ClientMessage::req(SubscriptionId::generate(), filters.clone())
.as_json()
.len()
})
.sum();
existing_bytes
.checked_add(auxiliary_bytes)
.and_then(|used| used.checked_add(super::SUBSCRIPTION_BYTE_RESERVED_MARGIN))
.is_some_and(|used| used <= remote_limit.unwrap())
}
/// Open separately tracked auxiliary live groups. The caller must first
/// use [`Self::can_admit_auxiliary_live_groups`] so one transient slot is
/// preserved; the shared ledger remains the final admission authority.
pub async fn subscribe_auxiliary_live_filter_groups(
&self,
filter_groups: Vec<Vec<Filter>>,
) -> Result<Vec<SubscriptionId>, String> {
if !self.can_admit_auxiliary_live_groups(&filter_groups) {
return Err(format!(
"Auxiliary live coverage would consume historic capacity for {}",
self.url
));
}
self.subscribe_live_filter_groups(filter_groups).await
}
/// Close a caller-owned subset of persistent subscriptions and return
/// their ledger slots only after CLOSE has been accepted by the SDK.
pub async fn close_live_subscriptions(
&self,
subscription_ids: &[SubscriptionId],
) -> Result<(), String> {
let mut failures = Vec::new();
for subscription_id in subscription_ids {
if let Err(error) = self.client.unsubscribe(subscription_id).await {
tracing::debug!(
relay = %self.url,
sub_id = %subscription_id,
%error,
"Failed to close caller-owned live subscription"
);
failures.push(format!("{subscription_id}: {error}"));
continue;
}
self.release_live_req_permit(subscription_id);
}
if failures.is_empty() {
Ok(())
} else {
Err(format!(
"Failed to close {} live subscriptions on {}: {}",
failures.len(),
self.url,
failures.join("; ")
))
}
}
pub fn complete_live_set_fits(&self, slots: usize) -> bool {
slots
<= self
@@ -2902,6 +2996,36 @@ mod tests {
assert_eq!(connection.subscription_budget().available_permits(), 8);
}
#[tokio::test]
async fn auxiliary_live_admission_preserves_one_historic_slot() {
let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());
connection.reset_subscription_budget(Some(10));
for index in 0..6 {
let permit = connection.acquire_subscription_slots(1).await.unwrap();
connection.live_req_permits_held.lock().unwrap().insert(
SubscriptionId::new(format!("core-{index}")),
HeldLiveSubscription {
generation: permit.generation,
_ledger_slot: permit.permit,
filters: vec![Filter::new().kind(Kind::Custom(20_100 + index))],
},
);
}
let one_group = vec![vec![Filter::new().kind(Kind::Custom(20_200))]];
assert!(connection.can_admit_auxiliary_live_groups(&one_group));
let two_groups = vec![
vec![Filter::new().kind(Kind::Custom(20_201))],
vec![Filter::new().kind(Kind::Custom(20_202))],
];
assert!(
!connection.can_admit_auxiliary_live_groups(&two_groups),
"two auxiliary groups would consume the final historic slot"
);
}
#[tokio::test]
async fn failed_live_group_rolls_back_opened_groups_and_full_capacity() {
let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());
+19
View File
@@ -84,6 +84,7 @@ struct RelayOptions {
relay_data_path: Option<PathBuf>,
deletion_lifecycle: Option<DeletionLifecycleOptions>,
rejected_hot_cache_duration_secs: Option<u64>,
relay_max_subscriptions: Option<usize>,
/// Run with the production outbound target policy (reject non-global
/// event-directed sync targets). The fixture default is permissive
/// because the entire test infrastructure lives on loopback.
@@ -309,6 +310,20 @@ impl TestRelay {
.await
}
/// Start a source relay advertising and enforcing a small per-connection
/// subscription budget. Sync scenarios use this to exercise graceful
/// degradation paths selected from NIP-11.
pub async fn start_with_relay_max_subscriptions(limit: usize) -> Self {
Self::start_internal(
port::reserve_port(),
RelayOptions {
relay_max_subscriptions: Some(limit),
..RelayOptions::default()
},
)
.await
}
/// Start a relay with LMDB backend (persistent side DBs on a temp dir).
pub async fn start_with_lmdb() -> Self {
Self::start_internal(
@@ -581,6 +596,10 @@ impl TestRelay {
);
}
if let Some(limit) = options.relay_max_subscriptions {
cmd.env("NGIT_RELAY_MAX_SUBSCRIPTIONS", limit.to_string());
}
// Add negentropy disable flag if requested
if options.disable_negentropy {
cmd.env("NGIT_SYNC_DISABLE_NEGENTROPY", "true");
+1
View File
@@ -33,6 +33,7 @@ mod common;
mod sync {
pub mod adaptive_pagination;
pub mod catchup;
pub mod descendant_sync;
pub mod discovery;
pub mod historic_recovery;
pub mod historic_sync;
+166
View File
@@ -0,0 +1,166 @@
//! Scheduled repository-descendant synchronization scenarios.
use std::time::Duration;
use nostr_sdk::prelude::*;
use crate::common::{sync_helpers::*, TestRelay};
async fn wait_for_log(path: &std::path::Path, needle: &str, timeout: Duration) -> bool {
let deadline = tokio::time::Instant::now() + timeout;
loop {
if std::fs::read_to_string(path)
.unwrap_or_default()
.contains(needle)
{
return true;
}
if tokio::time::Instant::now() >= deadline {
return false;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
/// Events which reference a direct repository-thread member, but not the
/// repository or its root event, are recovered by the scheduled historic
/// pass. The recovered events do not recursively extend the frontier.
#[tokio::test]
async fn historic_sync_recovers_one_generation_of_parent_only_descendants() {
let source = TestRelay::start_with_relay_max_subscriptions(5).await;
let keys = Keys::generate();
let repo_id = "scheduled-descendant-history";
// Seed the source before the syncing relay exists so this exercises the
// complete historic rotation, not ordinary live coverage.
let source_domains = [source.domain()];
let source_refs = source_domains
.iter()
.map(String::as_str)
.collect::<Vec<_>>();
let (_announcement, _source_git) =
setup_announcement_on_relay(&source, &keys, &source_refs, repo_id).await;
let source_client = TestClient::new(source.url(), keys.clone())
.await
.expect("connect to source relay");
let issue =
build_layer2_issue_event(&keys, &repo_coord(&keys, repo_id), "Repository thread root")
.expect("build issue");
let direct_reply = build_layer3_reply_with_e_tag(&keys, &issue.id, "Direct reply")
.expect("build direct reply");
let parent_only = build_layer3_reply_with_e_tag(
&keys,
&direct_reply.id,
"Reply visible only through the scheduled descendant query",
)
.expect("build parent-only descendant");
let recursive = build_layer3_reply_with_e_tag(
&keys,
&parent_only.id,
"Deliberately out-of-scope recursive descendant",
)
.expect("build recursive descendant");
for event in [&issue, &direct_reply, &parent_only, &recursive] {
source_client
.send_event(event)
.await
.expect("seed source event");
}
let syncing = TestRelay::start_with_sync(None).await;
let domains = [source.domain(), syncing.domain()];
let domain_refs = domains.iter().map(String::as_str).collect::<Vec<_>>();
let (_target_announcement, _target_git) =
setup_announcement_on_relay(&syncing, &keys, &domain_refs, repo_id).await;
assert!(
wait_for_event_on_relay(
syncing.url(),
Filter::new().id(parent_only.id),
Duration::from_secs(20),
)
.await,
"scheduled historic rotation should recover a parent-only descendant"
);
assert!(
wait_for_log(
&syncing.log_path(),
"Started queued descendant fallback query",
Duration::from_secs(5),
)
.await,
"the constrained source should exercise EOSE-closing fallback"
);
assert!(
!wait_for_event_on_relay(
syncing.url(),
Filter::new().id(recursive.id),
Duration::from_secs(7),
)
.await,
"recovered descendants must not recursively extend the frontier"
);
source_client.disconnect().await;
syncing.stop().await;
source.stop().await;
}
/// A relay with spare NIP-11 subscription capacity receives permanent
/// descendant coverage, so a later parent-only event arrives without waiting
/// for the fallback rotation.
#[tokio::test]
async fn descendant_live_coverage_is_preferred_when_capacity_remains() {
let source = TestRelay::start().await;
let keys = Keys::generate();
let repo_id = "live-descendant-coverage";
let source_domains = [source.domain()];
let source_refs = source_domains
.iter()
.map(String::as_str)
.collect::<Vec<_>>();
let (_announcement, _source_git) =
setup_announcement_on_relay(&source, &keys, &source_refs, repo_id).await;
let source_client = TestClient::new(source.url(), keys.clone())
.await
.expect("connect to source relay");
let issue = build_layer2_issue_event(&keys, &repo_coord(&keys, repo_id), "Live root")
.expect("build issue");
let direct_reply = build_layer3_reply_with_e_tag(&keys, &issue.id, "Direct reply")
.expect("build direct reply");
source_client.send_event(&issue).await.unwrap();
source_client.send_event(&direct_reply).await.unwrap();
let syncing = TestRelay::start_with_sync(None).await;
let domains = [source.domain(), syncing.domain()];
let domain_refs = domains.iter().map(String::as_str).collect::<Vec<_>>();
let (_target_announcement, _target_git) =
setup_announcement_on_relay(&syncing, &keys, &domain_refs, repo_id).await;
assert!(
wait_for_log(
&syncing.log_path(),
"Installed auxiliary descendant live coverage",
Duration::from_secs(20),
)
.await,
"spare source capacity should retain descendant coverage"
);
let parent_only =
build_layer3_reply_with_e_tag(&keys, &direct_reply.id, "Arrives through live coverage")
.expect("build parent-only descendant");
source_client.send_event(&parent_only).await.unwrap();
assert!(
wait_for_event_on_relay(
syncing.url(),
Filter::new().id(parent_only.id),
Duration::from_secs(10),
)
.await,
"parent-only descendant should arrive through retained live coverage"
);
source_client.disconnect().await;
syncing.stop().await;
source.stop().await;
}