feat(sync): recover direct-member descendant history

Repository collaboration events may reference only their immediate parent. Core Layer 3 filters follow repository root IDs, so a reply to an already discovered reply can remain absent forever when it omits repository and root tags.

Derive a stable, non-recursive frontier from locally stored direct root-thread members and rotate EOSE-closing filters through the existing historic queue, pagination, shared ledger, and request pacing. Run a complete initial pass, then rolling 24-hour passes; reconnects, daily reconciliation, and failed auxiliary batches restart completeness. Batch purpose keeps auxiliary failures out of core health state.

This deliberately provides the safe fallback baseline only: it does not retain descendant live subscriptions, recursively traverse threads, add configuration, persist cursors, or shard connections. Correctness assumes ordinary Layer 3 sync eventually stores direct root-thread members.

Validated with cargo check --lib, the descendant rotation unit test, and historic_sync_recovers_one_generation_of_parent_only_descendants. The scenario proves parent-only recovery while asserting a further recursive grandchild remains out of scope.
This commit is contained in:
DanConwayDev
2026-08-08 20:46:05 +00:00
parent d203f9de56
commit 5f84b8b55d
7 changed files with 390 additions and 9 deletions
+10
View File
@@ -7,6 +7,16 @@ 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. One low-priority, EOSE-closing
historic filter starts globally every five seconds while rotating across
relays: a complete initial pass followed by rolling 24-hour passes. The
frontier is deliberately non-recursive and uses the existing subscription
ledger, pagination, and request pacing rather than adding retained
subscriptions or connection sharding.
## [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,26 @@ live in
- **Function**: `build_root_event_tag_filters(root_events, since)`
- **Only for `SyncLevel::Full` repos** — purgatory announcements (`StateOnly`) skip this layer
### Scheduled 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. These filters are not retained live:
- one filter starts globally on the existing five-second maintenance cadence;
- every request uses the existing historic REQ+EOSE path, including the
per-connection subscription ledger, proactive pacing, and pagination;
- the first stable frontier receives a complete historic pass, then subsequent
rotations use a rolling 24-hour `since` window; and
- recovered descendants never enter the frontier, so this is deliberately one
additional generation rather than recursive thread traversal.
The rotation restarts with a complete pass after a reconnect or daily
reconciliation. This makes an interrupted filter recoverable without adding a
durable cursor, another scheduler, retained subscription pressure, or
multi-connection sharding.
### Combined Layer 2+3 (SyncLevel-Aware)
The `build_sync_level_aware_filters()` function combines both layers, partitioning repos by `SyncLevel`:
+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
+269 -8
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_RECENT_WINDOW_SECS: u64 = 24 * 60 * 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,49 @@ 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<Filter>,
next_filter: usize,
historic: bool,
}
impl Default for DescendantSyncRotation {
fn default() -> Self {
Self {
frontier: HashSet::new(),
filters: Vec::new(),
next_filter: 0,
historic: true,
}
}
}
impl DescendantSyncRotation {
fn refresh(&mut self, members: HashSet<EventId>, now: Timestamp) {
if self.historic
&& !self.filters.is_empty()
&& self.next_filter >= self.filters.len()
&& self.frontier == members
{
self.historic = false;
}
self.frontier = members.clone();
let since = (!self.historic)
.then(|| Timestamp::from(now.as_secs().saturating_sub(DESCENDANT_RECENT_WINDOW_SECS)));
self.filters = filters::tagged_one_of_our_root_event_filters(&members, since);
self.next_filter = 0;
}
}
/// Items included in a pending batch
#[derive(Debug, Clone, Default)]
pub struct PendingItems {
@@ -1172,6 +1219,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 at a time receives one low-priority repository-descendant
/// history query. Reusing this five-second maintenance cadence bounds the new
/// global query-start rate without another scheduler or configuration surface.
async fn run_purgatory_announcement_sync(
sync_manager: Arc<Mutex<SyncManager>>,
mut shutdown_rx: broadcast::Receiver<()>,
@@ -1187,6 +1238,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 +1467,11 @@ 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>,
/// Round-robin cursor so the maintenance loop starts at most one extra
/// query globally per 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 +1572,8 @@ 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_relay_cursor: 0,
disconnect_tx: None,
eose_tx: None,
subscription_closed_tx: None,
@@ -2507,14 +2566,20 @@ 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.failed && batch.purpose == PendingBatchPurpose::Descendants {
// Do not let an interrupted page turn the initial complete pass
// into a recent-only rotation. Starting this tiny auxiliary state
// afresh is simpler and safer than persisting per-filter cursors.
self.descendant_sync_rotations.remove(relay_url);
}
let mut relay_index = self.relay_sync_index.write().await;
@@ -2567,7 +2632,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 +2650,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 +2801,7 @@ 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");
self.descendant_sync_rotations.remove(relay_url);
// Get connection
let connection = match self.connections.get(relay_url) {
@@ -3528,16 +3595,151 @@ 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
}
/// Start at most one low-priority descendant query globally per tick.
///
/// Each relay first receives a complete historic pass over a stable
/// snapshot of direct thread members. Later passes use a rolling 24-hour
/// window. Every request uses the ordinary historic REQ+EOSE pagination
/// path, so it draws from the existing ledger and closes at EOSE.
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 relay_url = relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone();
self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1);
let root_events = targets[&relay_url].root_events.clone();
let needs_snapshot = self
.descendant_sync_rotations
.get(&relay_url)
.is_none_or(|rotation| rotation.next_filter >= rotation.filters.len());
if needs_snapshot {
let members = self.direct_thread_members(&root_events).await;
if members.is_empty() {
return;
}
let rotation = self
.descendant_sync_rotations
.entry(relay_url.clone())
.or_default();
rotation.refresh(members, Timestamp::now());
}
let (filter, historic, filter_index, filter_count) = {
let rotation = self.descendant_sync_rotations.get(&relay_url).unwrap();
(
rotation.filters[rotation.next_filter].clone(),
rotation.historic,
rotation.next_filter,
rotation.filters.len(),
)
};
if self
.historic_sync_with_options(
&relay_url,
vec![filter],
PendingItems::default(),
None,
PendingBatchPurpose::Descendants,
true,
)
.await
.is_some()
{
self.descendant_sync_rotations
.get_mut(&relay_url)
.unwrap()
.next_filter += 1;
tracing::debug!(
relay = %relay_url,
historic,
filter_index,
filter_count,
"Started scheduled repository descendant query"
);
}
}
/// 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 +4740,10 @@ 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);
// Check if this was an intentional disconnect (Disconnecting status)
let was_intentional = {
@@ -5945,6 +6151,26 @@ impl SyncManager {
filters: Vec<Filter>,
items: PendingItems,
since: Option<Timestamp>,
) -> Option<u64> {
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!(
@@ -5985,8 +6211,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();
@@ -6008,6 +6235,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,
@@ -6241,6 +6469,7 @@ impl SyncManager {
// Create PendingBatch for REQ+EOSE
let batch = PendingBatch {
batch_id,
purpose,
items,
outstanding_subs: subscription_ids,
sync_method: SyncMethod::ReqEose,
@@ -6349,6 +6578,33 @@ mod tests {
assert!(!should_use_semantic_fallback(19, 0));
}
#[test]
fn descendant_rotation_becomes_recent_only_after_a_stable_historic_pass() {
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.clone(), now);
assert!(rotation.historic);
assert!(rotation
.filters
.iter()
.all(|filter| !serde_json::to_value(filter)
.unwrap()
.as_object()
.unwrap()
.contains_key("since")));
rotation.next_filter = rotation.filters.len();
rotation.refresh(members, now);
assert!(!rotation.historic);
assert!(rotation.filters.iter().all(|filter| {
serde_json::to_value(filter).unwrap()["since"]
== serde_json::json!(now.as_secs() - DESCENDANT_RECENT_WINDOW_SECS)
}));
}
#[test]
fn group_filters_for_req_respects_count_and_byte_budgets() {
// Many small filters group by the count cap.
@@ -6863,6 +7119,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,
@@ -6965,6 +7222,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,
@@ -7000,6 +7258,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,
@@ -7318,6 +7577,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,
@@ -7341,6 +7601,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,
+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;
+82
View File
@@ -0,0 +1,82 @@
//! Scheduled repository-descendant synchronization scenarios.
use std::time::Duration;
use nostr_sdk::prelude::*;
use crate::common::{sync_helpers::*, TestRelay};
/// 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().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_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;
}