mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-06 07:28:23 +00:00
feat(sync): retain descendants when connection capacity permits
Scheduled history is a safe completeness fallback but delays new parent-only collaboration events even on relays with ample subscription capacity. Descendant coverage should be immediate when it does not compromise the existing core product or recovery headroom. Admit the complete grouped descendant filter set only when the session ledger can still preserve its two control-plane slots and one usable transient slot, and when the learned retained-byte limit still leaves one maximum-sized transient request. Track the auxiliary subscription IDs separately so frontier changes, CLOSED, daily reset, and core consolidation can retire only descendant coverage. A failed CLOSE retains its ledger ownership and aborts core replacement rather than allowing local accounting to run ahead of the relay. Core filter grouping and rollback remain unchanged; auxiliary filters are intentionally not packed into partially full core groups. Constrained sessions continue using the independently complete historic rotation from the preceding commit. Recursive traversal, connection sharding, and minimum-churn tail compaction remain excluded. Validated with cargo check --lib, the descendant rotation unit test, and auxiliary_live_admission_preserves_one_historic_slot using a nostream-shaped advertised budget of 10.
This commit is contained in:
+5
-6
@@ -10,12 +10,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
### 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.
|
||||
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
|
||||
|
||||
|
||||
@@ -788,24 +788,30 @@ 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
|
||||
### 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:
|
||||
or `q` tags reference that direct member.
|
||||
|
||||
- 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
|
||||
- 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
|
||||
the existing maintenance cadence through the ordinary historic queue,
|
||||
pagination, shared ledger, and request pacing;
|
||||
- the fallback receives a complete initial pass, then rolling 24-hour passes;
|
||||
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
|
||||
An unexpected auxiliary CLOSED retires the remaining descendant subscriptions
|
||||
without rebuilding core coverage and falls back to history. Reconnect and daily
|
||||
reconciliation reconstruct the mode from current session capacity. This keeps
|
||||
the feature complete without durable cursors, another capacity ledger, or
|
||||
multi-connection sharding.
|
||||
|
||||
### Combined Layer 2+3 (SyncLevel-Aware)
|
||||
|
||||
+166
-3
@@ -707,6 +707,12 @@ struct DescendantSyncRotation {
|
||||
historic: bool,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DescendantLiveCoverage {
|
||||
frontier: HashSet<EventId>,
|
||||
subscription_ids: Vec<SubscriptionId>,
|
||||
}
|
||||
|
||||
impl Default for DescendantSyncRotation {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
@@ -1469,6 +1475,9 @@ pub struct SyncManager {
|
||||
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 the maintenance loop starts at most one extra
|
||||
/// query globally per tick.
|
||||
descendant_relay_cursor: usize,
|
||||
@@ -1573,6 +1582,7 @@ impl SyncManager {
|
||||
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,
|
||||
@@ -2801,6 +2811,9 @@ 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
|
||||
@@ -3651,6 +3664,39 @@ impl SyncManager {
|
||||
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
|
||||
}
|
||||
|
||||
/// Start at most one low-priority descendant query globally per tick.
|
||||
///
|
||||
/// Each relay first receives a complete historic pass over a stable
|
||||
@@ -3688,16 +3734,85 @@ impl SyncManager {
|
||||
self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1);
|
||||
let root_events = targets[&relay_url].root_events.clone();
|
||||
|
||||
let mut refreshed_members = None;
|
||||
if let Some(existing) = self.descendant_live_coverage.get(&relay_url) {
|
||||
let members = self.direct_thread_members(&root_events).await;
|
||||
if members == existing.frontier {
|
||||
return;
|
||||
}
|
||||
refreshed_members = Some(members);
|
||||
if !self
|
||||
.close_descendant_live_coverage(&relay_url, "frontier changed")
|
||||
.await
|
||||
{
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
let members = match refreshed_members {
|
||||
Some(members) => members,
|
||||
None => self.direct_thread_members(&root_events).await,
|
||||
};
|
||||
if members.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
let now = Timestamp::now();
|
||||
let live_since = Timestamp::from(
|
||||
now.as_secs()
|
||||
.saturating_sub(DESCENDANT_RECENT_WINDOW_SECS.min(15 * 60)),
|
||||
);
|
||||
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.clone(),
|
||||
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"
|
||||
);
|
||||
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.clone())
|
||||
@@ -4744,6 +4859,7 @@ impl SyncManager {
|
||||
// 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 = {
|
||||
@@ -5456,6 +5572,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;
|
||||
@@ -5529,6 +5655,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, Timestamp::now());
|
||||
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,
|
||||
@@ -5654,13 +5807,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;
|
||||
|
||||
@@ -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());
|
||||
|
||||
Reference in New Issue
Block a user