From da16ddf1fa0c8a35cf9da2b254a87aa47c69fbbb Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 12 Aug 2026 22:28:28 +0000 Subject: [PATCH 1/6] fix(sync): bound terminal lifecycle notifications Outbound peers can repeat EOSE or CLOSED messages faster than the sync actor processes them. The relay data lane was bounded, but it forwarded terminal state into two unbounded actor channels, so a hostile or broken peer could turn actor contention into retained-memory growth. Use one fixed-capacity queue for each lifecycle notification type and await capacity in the per-relay event processor. Transient permit release remains on the independent terminal-control listener, so backpressure here cannot strand subscription-ledger slots. A full queue now bounds memory and naturally propagates pressure into the already-bounded per-relay data lane. The 1,000-entry capacity preserves the observed 169-terminal production startup burst while remaining independent of advertised peer limits. Repeated and reordered terminals remain idempotent in the existing actor handlers; missing terminals remain covered by the existing CLOSE-before-release watchdog. Other logically bounded worker-result channels and long-lived index ownership are deliberately left for the audit inventory. Validated with the production-sized burst test, a new exact-capacity repeated-terminal backpressure test, cargo fmt, and cargo clippy --locked --lib -- -D warnings. --- src/sync/mod.rs | 69 +++++++++++++++++++++++++++++++++++-------------- 1 file changed, 49 insertions(+), 20 deletions(-) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 562c06e..b3138a7 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1173,11 +1173,14 @@ struct SubscriptionClosedNotification { live_filter_count: Option, } -fn lifecycle_notification_channel() -> ( - tokio::sync::mpsc::UnboundedSender, - tokio::sync::mpsc::UnboundedReceiver, -) { - tokio::sync::mpsc::unbounded_channel() +// One global actor consumes terminal state changes from every outbound relay. +// A hostile peer can repeat EOSE/CLOSED indefinitely, so this queue must bound +// retained memory independently of the per-connection subscription ledger. +const LIFECYCLE_NOTIFICATION_CAPACITY: usize = 1_000; + +fn lifecycle_notification_channel( +) -> (tokio::sync::mpsc::Sender, tokio::sync::mpsc::Receiver) { + tokio::sync::mpsc::channel(LIFECYCLE_NOTIFICATION_CAPACITY) } fn is_rate_limit_message(message: &str) -> bool { @@ -1913,10 +1916,9 @@ pub struct SyncManager { /// Channel for disconnect notifications (set during run) disconnect_tx: Option>, /// Channel for EOSE notifications (set during run) - eose_tx: Option>, + eose_tx: Option>, /// Serializes CLOSED recovery and pending-batch cleanup through the actor. - subscription_closed_tx: - Option>, + subscription_closed_tx: Option>, /// Returns connection outcomes to the sync actor for serialized state changes. connect_attempt_result_tx: Option>, /// Wakes the actor when a batch completion may unblock consolidation. @@ -3935,10 +3937,12 @@ impl SyncManager { sub_id = %sub_id, "EOSE received, notifying SyncManager" ); - let _ = eose_tx.send(EoseNotification { - relay_url: relay_url_clone.clone(), - sub_id, - }); + let _ = eose_tx + .send(EoseNotification { + relay_url: relay_url_clone.clone(), + sub_id, + }) + .await; } RelayEvent::Notice(notice) => { if is_rate_limit_message(¬ice) { @@ -4006,13 +4010,15 @@ impl SyncManager { ); } } - let _ = subscription_closed_tx.send(SubscriptionClosedNotification { - relay_url: relay_url_clone.clone(), - subscription_id, - reason, - generation: live_generation, - live_filter_count, - }); + let _ = subscription_closed_tx + .send(SubscriptionClosedNotification { + relay_url: relay_url_clone.clone(), + subscription_id, + reason, + generation: live_generation, + live_filter_count, + }) + .await; } RelayEvent::Shutdown => { tracing::info!(relay = %relay_url_clone, "Relay shutdown detected"); @@ -7885,7 +7891,7 @@ mod tests { fn lifecycle_inbox_accepts_production_sized_terminal_burst() { let (tx, mut rx) = lifecycle_notification_channel::(); for index in 0..169 { - tx.send(EoseNotification { + tx.try_send(EoseNotification { relay_url: "wss://relay.example".to_string(), sub_id: SubscriptionId::new(format!("hydration-{index}")), }) @@ -8242,6 +8248,29 @@ mod tests { assert_ne!(first, second, "attempt identities must not be reused"); } + #[test] + fn lifecycle_notifications_apply_fixed_backpressure_under_peer_repetition() { + let (sender, mut receiver) = lifecycle_notification_channel(); + + for sequence in 0..LIFECYCLE_NOTIFICATION_CAPACITY { + sender + .try_send(sequence) + .expect("the documented lifecycle capacity should be usable"); + } + assert!( + matches!( + sender.try_send(LIFECYCLE_NOTIFICATION_CAPACITY), + Err(tokio::sync::mpsc::error::TrySendError::Full(_)) + ), + "repeated peer terminals must backpressure instead of retaining unbounded memory" + ); + + assert_eq!(receiver.try_recv().unwrap(), 0); + sender + .try_send(LIFECYCLE_NOTIFICATION_CAPACITY) + .expect("draining one terminal must release exactly one queue slot"); + } + #[tokio::test] async fn connect_attempt_semaphore_caps_parallel_workers() { use std::sync::atomic::{AtomicUsize, Ordering}; From cd6e5612631b1c5ddcdd0e87e781876818fff807 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 12 Aug 2026 22:30:46 +0000 Subject: [PATCH 2/6] fix(sync): reject unowned auth retry identities The auth-required recovery marker was keyed by the CLOSED subscription id supplied by a relay. Any first-seen id was inserted before checking whether ngit-grasp had opened that subscription, so a peer could retain an unbounded set of forged identities without consuming a subscription-ledger slot. Expose a session-local ownership check over the transient and live permit maps. Auth-required handling now ignores CLOSED ids that do not own a current ledger permit; legitimate subscriptions still receive exactly one NIP-42 retry and a second refusal still retires the SDK subscription. This deliberately does not impose another numeric limit: marker cardinality is bounded by the existing per-connection subscription budget and is cleared on disconnect. Ordinary non-auth CLOSED handling remains idempotent and unchanged. Validated with a unit test distinguishing an owned live subscription from a peer-forged id, the existing one-retry test, cargo fmt, and cargo clippy --locked --lib -- -D warnings. --- src/sync/mod.rs | 19 +++++++++++++----- src/sync/relay_connection.rs | 39 ++++++++++++++++++++++++++++++++++++ 2 files changed, 53 insertions(+), 5 deletions(-) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index b3138a7..abb04f0 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -6649,6 +6649,17 @@ impl SyncManager { .get(relay_url) .is_some_and(|coverage| coverage.subscription_ids.contains(&subscription_id)); let policy_category = if is_auth_required_message(reason) { + let Some(connection) = self.connections.get(relay_url) else { + return; + }; + if !connection.holds_subscription_permit(&subscription_id) { + tracing::debug!( + relay = %relay_url, + sub_id = %subscription_id, + "Ignoring auth-required CLOSED for an unowned subscription" + ); + return; + } if reserve_authentication_retry( &mut self.auth_required_attempts, relay_url, @@ -6661,11 +6672,9 @@ impl SyncManager { ); return; } - if let Some(connection) = self.connections.get(relay_url) { - connection - .retire_auth_refused_subscription(&subscription_id) - .await; - } + connection + .retire_auth_refused_subscription(&subscription_id) + .await; Some(PolicyRefusal::AuthenticationRequired) } else { policy_refusal(reason) diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index bc72b57..f7d152d 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -2185,6 +2185,22 @@ impl RelayConnection { self.client.subscriptions().await.len() } + /// Whether this session still owns a ledger permit for a subscription. + /// + /// Peer-supplied CLOSED identifiers must not create application retry + /// state unless they name work admitted through our bounded ledger. + pub(super) fn holds_subscription_permit(&self, subscription_id: &SubscriptionId) -> bool { + self.transient_req_permits_held + .lock() + .expect("transient permit map poisoned") + .contains_key(subscription_id) + || self + .live_req_permits_held + .lock() + .expect("live permit map poisoned") + .contains_key(subscription_id) + } + async fn retire_peer_closed_subscription(&self, subscription_id: &SubscriptionId) { self.release_transient_req_permit(subscription_id); // A peer CLOSED is terminal for this subscription. In particular, @@ -3413,6 +3429,29 @@ mod tests { assert_eq!(connection.subscription_budget().available_permits(), 8); } + #[tokio::test] + async fn peer_subscription_identity_is_bounded_by_owned_ledger_permits() { + let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate()); + connection.reset_subscription_budget(Some(10)); + let owned = SubscriptionId::new("owned-live"); + let forged = SubscriptionId::new("peer-forged"); + let permit = connection.acquire_subscription_slots(1).await.unwrap(); + connection.live_req_permits_held.lock().unwrap().insert( + owned.clone(), + HeldLiveSubscription { + generation: permit.generation, + _ledger_slot: permit.permit, + filters: vec![Filter::new().kind(Kind::TextNote)], + }, + ); + + assert!(connection.holds_subscription_permit(&owned)); + assert!( + !connection.holds_subscription_permit(&forged), + "a peer-selected CLOSED id must not become application-owned state" + ); + } + #[tokio::test] async fn auxiliary_live_admission_preserves_one_historic_slot() { let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate()); From 4d5c13ca9fc2965184b6a73f99ca0b33f1221cbf Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 12 Aug 2026 22:33:12 +0000 Subject: [PATCH 3/6] refactor(sync): remove deferred self-wakeup queue Deferred live-filter consolidation sent one unbounded actor message after every completing batch while retaining the actual desired work in a deduplicated relay set. If the actor was busy, redundant self-addressed wakeups could accumulate even though only one consolidation per relay was meaningful. Batch confirmation already executes under exclusive actor ownership after releasing every index guard. It now checks the deduplicated deferral set directly and processes a ready consolidation in place. The finite synchronous-completion path is type-erased with Box::pin because consolidation can immediately confirm an empty historic batch; the deferral is removed before that path, so it cannot recurse indefinitely. This changes no subscription policy or consolidation timing: pending batches still prevent consolidation, reset/disconnect still cancel it, and the final batch completion still re-derives dropped work. It only removes the redundant notification transport and its retained-memory surface. Validated with the deferred-until-final-batch and cancellation/idempotence tests, cargo fmt, and cargo clippy --locked --lib -- -D warnings. --- src/sync/mod.rs | 57 ++++++++----------------------------------------- 1 file changed, 9 insertions(+), 48 deletions(-) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index abb04f0..9f3f996 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1545,18 +1545,6 @@ impl DeferredConsolidations { } } -fn notify_deferred_consolidation_after_batch_completion( - deferred: &DeferredConsolidations, - sender: Option<&tokio::sync::mpsc::UnboundedSender>, - relay_url: &str, -) { - if deferred.contains(relay_url) { - if let Some(sender) = sender { - let _ = sender.send(relay_url.to_string()); - } - } -} - // ============================================================================= // Daily Timer // ============================================================================= @@ -1921,8 +1909,6 @@ pub struct SyncManager { subscription_closed_tx: Option>, /// Returns connection outcomes to the sync actor for serialized state changes. connect_attempt_result_tx: Option>, - /// Wakes the actor when a batch completion may unblock consolidation. - deferred_consolidation_tx: Option>, nip65_discovery_result_tx: Option>, /// Channel for broadcasting shutdown signal to all background tasks shutdown_tx: Option>, @@ -2027,7 +2013,6 @@ impl SyncManager { eose_tx: None, subscription_closed_tx: None, connect_attempt_result_tx: None, - deferred_consolidation_tx: None, nip65_discovery_result_tx: None, shutdown_tx: None, metrics: sync_metrics, @@ -3170,11 +3155,14 @@ impl SyncManager { // Release lock before checking if historic sync is complete drop(relay_index); - notify_deferred_consolidation_after_batch_completion( - &self.deferred_consolidations, - self.deferred_consolidation_tx.as_ref(), - relay_url, - ); + // Batch completion already runs inside the sync actor and all index + // guards are released. Process a ready deferral directly instead of + // retaining self-addressed wakeups in another queue. + if self.deferred_consolidations.contains(relay_url) { + // Consolidation can synchronously complete an empty historic + // batch, so erase that finite recursive future from the type. + Box::pin(self.process_deferred_consolidation(relay_url)).await; + } // Spawn background task to check if historic sync is complete // This avoids blocking the confirm_batch flow for 6 seconds @@ -3399,10 +3387,6 @@ impl SyncManager { let (connect_attempt_result_tx, mut connect_attempt_result_rx) = mpsc::unbounded_channel::(); - // Batch completion can be synchronous, so use a non-blocking wakeup - // instead of waiting for pending work while holding the actor lock. - let (deferred_consolidation_tx, mut deferred_consolidation_rx) = - mpsc::unbounded_channel::(); let (nip65_discovery_result_tx, mut nip65_discovery_result_rx) = mpsc::unbounded_channel::(); @@ -3426,7 +3410,6 @@ impl SyncManager { self.eose_tx = Some(eose_tx.clone()); self.subscription_closed_tx = Some(subscription_closed_tx); self.connect_attempt_result_tx = Some(connect_attempt_result_tx); - self.deferred_consolidation_tx = Some(deferred_consolidation_tx); self.nip65_discovery_result_tx = Some(nip65_discovery_result_tx); self.shutdown_tx = Some(shutdown_tx.clone()); @@ -3542,12 +3525,6 @@ impl SyncManager { } } } - relay_url = deferred_consolidation_rx.recv() => { - if let Some(relay_url) = relay_url { - let mut manager = sync_manager.lock().await; - manager.process_deferred_consolidation(&relay_url).await; - } - } result = nip65_discovery_result_rx.recv() => { if let Some(result) = result { let mut manager = sync_manager.lock().await; @@ -8601,28 +8578,19 @@ mod tests { fn deferred_consolidation_runs_only_after_final_batch_completion() { let relay_url = "wss://relay.example"; let mut deferred = DeferredConsolidations::default(); - let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel(); assert!( !deferred.request(relay_url, true), "an in-flight batch must defer instead of waiting in the sync actor" ); - notify_deferred_consolidation_after_batch_completion(&deferred, Some(&sender), relay_url); - assert_eq!( - receiver.try_recv().unwrap(), - relay_url, - "batch completion must wake the actor without awaiting under its lock" - ); assert!( !deferred.take_ready(relay_url, true), "consolidation must remain queued while another batch is pending" ); - notify_deferred_consolidation_after_batch_completion(&deferred, Some(&sender), relay_url); - assert_eq!(receiver.try_recv().unwrap(), relay_url); assert!( deferred.take_ready(relay_url, false), - "the final batch completion must make consolidation runnable" + "the final actor-owned batch completion must make consolidation runnable" ); assert!( !deferred.take_ready(relay_url, false), @@ -8635,16 +8603,9 @@ mod tests { for reason in ["reset", "disconnect"] { let relay_url = format!("wss://{reason}.example"); let mut deferred = DeferredConsolidations::default(); - let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel(); assert!(!deferred.request(&relay_url, true)); - notify_deferred_consolidation_after_batch_completion( - &deferred, - Some(&sender), - &relay_url, - ); assert!(deferred.cancel(&relay_url)); - assert_eq!(receiver.try_recv().unwrap(), relay_url); assert!( !deferred.take_ready(&relay_url, false), "a queued wakeup must not resurrect consolidation after {reason}" From 826b0fb67f0fca71c7b4d68b2a7673a1ad207913 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 12 Aug 2026 22:35:14 +0000 Subject: [PATCH 4/6] fix(sync): bound worker result channels Connection and NIP-65 discovery workers returned completion through unbounded channels. Their producers were logically constrained, but encoding those assumptions only in scheduler state left retained-memory safety dependent on every future call site preserving that discipline. Use a result queue equal to MAX_CONCURRENT_CONNECT_ATTEMPTS for connection workers and a one-entry queue for the deliberately single-flight NIP-65 lookup. Workers await result capacity; they own no sync-actor lock, and connection shutdown still disconnects a result whose receiver has gone away. No work concurrency or retry cadence changes. The capacities mirror the existing producer bounds rather than introducing lower operational limits. The adjacent lifecycle comment is corrected to describe the independent permit-release lane and safe bounded backpressure. Validated with connection token/semaphore tests, partial NIP-65 retry coverage, cargo fmt, and cargo clippy --locked --lib -- -D warnings. --- src/sync/mod.rs | 32 +++++++++++++++++--------------- 1 file changed, 17 insertions(+), 15 deletions(-) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 9f3f996..c5cf41f 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1908,8 +1908,8 @@ pub struct SyncManager { /// Serializes CLOSED recovery and pending-batch cleanup through the actor. subscription_closed_tx: Option>, /// Returns connection outcomes to the sync actor for serialized state changes. - connect_attempt_result_tx: Option>, - nip65_discovery_result_tx: Option>, + connect_attempt_result_tx: Option>, + nip65_discovery_result_tx: Option>, /// Channel for broadcasting shutdown signal to all background tasks shutdown_tx: Option>, /// Prometheus metrics for sync operations (None if metrics disabled) @@ -3374,21 +3374,20 @@ impl SyncManager { let (disconnect_tx, mut disconnect_rx) = mpsc::channel::(100); // 3. Create EOSE channel for spawned tasks -> manager communication - // Lifecycle notifications must never backpressure the ordered EVENT - // processor. Their production is already bounded by the subscription - // ledger, while the actor may legitimately stay busy opening a large, - // paced historic batch for longer than a fixed inbox can absorb. + // The independent terminal-control listener has already released + // transient resources. These actor notifications may therefore apply + // bounded backpressure without stranding subscription-ledger slots. let (eose_tx, mut eose_rx) = lifecycle_notification_channel::(); let (subscription_closed_tx, mut subscription_closed_rx) = lifecycle_notification_channel::(); - // 4. Connection workers never mutate manager state. Their unbounded - // result channel cannot make a completed worker wait behind the actor. + // Connection workers never mutate manager state. At most the global + // connection-attempt cap can complete while the actor is busy. let (connect_attempt_result_tx, mut connect_attempt_result_rx) = - mpsc::unbounded_channel::(); + mpsc::channel::(MAX_CONCURRENT_CONNECT_ATTEMPTS); let (nip65_discovery_result_tx, mut nip65_discovery_result_rx) = - mpsc::unbounded_channel::(); + mpsc::channel::(1); // 4b. Create shutdown broadcast channel for graceful shutdown let (shutdown_tx, _shutdown_rx) = broadcast::channel(1); @@ -4788,6 +4787,7 @@ impl SyncManager { token, outcome, }) + .await .is_err() { connection.disconnect().await; @@ -4997,11 +4997,13 @@ impl SyncManager { let outcome = connection .fetch_events(filter, Duration::from_secs(30)) .await; - let _ = result_tx.send(Nip65DiscoveryResult { - source_relay, - authors, - outcome, - }); + let _ = result_tx + .send(Nip65DiscoveryResult { + source_relay, + authors, + outcome, + }) + .await; }); break; } From 18213db4cb000f096510141a614f7de765bf22de Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 12 Aug 2026 22:37:28 +0000 Subject: [PATCH 5/6] feat(metrics): expose retained sync state totals The peer-state audit found that important long-lived collections were visible only through incidental logs. That made slow cardinality growth difficult to distinguish from expected event throughput during production soaks. Export one fixed-label aggregate gauge and refresh eight operational classes from the existing two-second health pass: pending batches and subscriptions, auth retries, connection attempts, purgatory dependency attempts, temporary dependency relays, deferred consolidations, and descendant rotations. No relay, event, or repository identifier enters a label. The metric is observational only and introduces no admission limit. Its class vocabulary is deliberately fixed in the sync manager so monitoring cannot recreate the peer-controlled cardinality problem removed by the preceding metrics fix. Validated with a registry test for aggregate values and series count, cargo fmt, and cargo clippy --locked --lib -- -D warnings. --- src/sync/metrics.rs | 44 ++++++++++++++++++++++++++++++++++++++++++++ src/sync/mod.rs | 24 ++++++++++++++++++++++++ 2 files changed, 68 insertions(+) diff --git a/src/sync/metrics.rs b/src/sync/metrics.rs index ecda789..fe063b3 100644 --- a/src/sync/metrics.rs +++ b/src/sync/metrics.rs @@ -42,6 +42,8 @@ pub struct SyncMetrics { relays_connected_total: IntGauge, /// Relays marked as dead relays_dead_total: IntGauge, + /// Aggregate cardinality of long-lived sync-manager state by fixed class. + retained_state_current: IntGaugeVec, // === Rejected Events Index Metrics (unified with event_type label) === /// Current number of entries in hot cache (by event_type: announcement, state) @@ -144,6 +146,15 @@ impl SyncMetrics { ))?; registry.register(Box::new(relays_dead_total.clone()))?; + let retained_state_current = IntGaugeVec::new( + Opts::new( + "ngit_sync_retained_state_current", + "Current retained sync-manager entries by fixed state class", + ), + &["class"], + )?; + registry.register(Box::new(retained_state_current.clone()))?; + // Rejected events metrics (unified with event_type label) let rejected_hot_cache_current = IntGaugeVec::new( Opts::new( @@ -228,6 +239,7 @@ impl SyncMetrics { relays_tracked_total, relays_connected_total, relays_dead_total, + retained_state_current, rejected_hot_cache_current, rejected_hot_cache_hits_total, rejected_hot_cache_misses_total, @@ -425,6 +437,14 @@ impl SyncMetrics { self.relays_dead_total.get() } + /// Record one aggregate long-lived state cardinality. Callers must use the + /// fixed class vocabulary documented by the sync manager. + pub fn set_retained_state(&self, class: &str, count: usize) { + self.retained_state_current + .with_label_values(&[class]) + .set(count as i64); + } + // === Rejected Events Recording Methods (unified with event_type parameter) === /// Update hot cache current size gauge for a specific event type. @@ -663,6 +683,30 @@ mod tests { assert!(metrics2.is_err()); } + #[test] + fn retained_state_metric_uses_fixed_aggregate_classes() { + let registry = create_test_registry(); + let metrics = SyncMetrics::register(®istry).unwrap(); + + metrics.set_retained_state("pending_batches", 12); + metrics.set_retained_state("queued_connection_attempts", 3); + + let family = registry + .gather() + .into_iter() + .find(|family| family.name() == "ngit_sync_retained_state_current") + .expect("retained-state metric must be exported"); + assert_eq!(family.get_metric().len(), 2); + assert_eq!( + family + .get_metric() + .iter() + .map(|metric| metric.get_gauge().value() as i64) + .sum::(), + 15 + ); + } + #[test] fn test_rejected_events_metrics() { let registry = create_test_registry(); diff --git a/src/sync/mod.rs b/src/sync/mod.rs index c5cf41f..a796af4 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1803,6 +1803,30 @@ async fn run_health_and_metrics_checker( let entries = naughty_list.get_all(); metrics.update_naughty_list(entries); } + + let (pending_batches, pending_subscriptions) = { + let pending = manager.pending_sync_index.read().await; + ( + pending.values().map(Vec::len).sum(), + pending + .values() + .flatten() + .map(|batch| batch.outstanding_subs.len()) + .sum(), + ) + }; + for (class, count) in [ + ("pending_batches", pending_batches), + ("pending_subscriptions", pending_subscriptions), + ("auth_retries", manager.auth_required_attempts.len()), + ("queued_connection_attempts", manager.in_flight_connect_attempts.len()), + ("purgatory_dependencies", manager.purgatory_dependency_attempts.len()), + ("dependency_relays", manager.dependency_relay_deadlines.len()), + ("deferred_consolidations", manager.deferred_consolidations.relays.len()), + ("descendant_rotations", manager.descendant_sync_rotations.len()), + ] { + metrics.set_retained_state(class, count); + } } } _ = shutdown_rx.recv() => { From 70dec5ac295dbedd340a9be5a8db873f4b4040fb Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 12 Aug 2026 22:39:32 +0000 Subject: [PATCH 6/6] docs(architecture): inventory peer-controlled state The auth-required OOM incident showed that concurrency controls alone do not establish who owns retained state or how it ends. That reasoning was scattered across implementation comments, leaving future queue and retry changes vulnerable to recreating an unbounded lifecycle. Add a subsystem inventory classifying inbound, outbound sync, purgatory, Git, retry, channel, registry, and task state as static, time-bounded, session-bounded, or externally bounded. Each row records its producer, cleanup owner, missing/repeated terminal behavior, and operational signal. Link the invariants from the main architecture and summarize the completed security/observability work in the changelog. The inventory deliberately excludes normal database growth from valid admitted events. It records external bounds only where transient keys are deduplicated against that admitted state, and establishes that arbitrary wire identifiers may never become application-owned state. Validated against the current implementation after removing every production unbounded mpsc channel; git diff --check passes. --- CHANGELOG.md | 7 +++ docs/explanation/architecture.md | 4 ++ docs/explanation/peer-controlled-state.md | 69 +++++++++++++++++++++++ 3 files changed, 80 insertions(+) create mode 100644 docs/explanation/peer-controlled-state.md diff --git a/CHANGELOG.md b/CHANGELOG.md index 7e0a51c..aefb13e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,12 +14,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Security +- Bound every sync-actor notification channel and require `auth-required` + retry IDs to own a current subscription-ledger permit. Repeated terminals, + forged subscription IDs, and delayed worker results can no longer create + unbounded retained memory. - Added opt-in trusted-proxy CIDRs so WebSocket connection policy, per-IP metrics, abuse indicators, and logs can use the real client address without trusting spoofable forwarding headers from arbitrary direct peers. ### Added +- Add fixed-cardinality aggregate metrics for important long-lived sync state + and document the producer, cleanup owner, bound, and terminal behavior of + every peer-influenced transient subsystem. - Add bounded, operator-configurable inbox fallback coverage when a successful Sync+ user-index query finds no accepted NIP-65 relay list for a root author. - Add a default-on `NGIT_SYNC_PLUS_ENABLED` opt-out and advertise GRASP-03 in diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 63a3295..d0180b1 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -1,5 +1,9 @@ # ngit-grasp Architecture +Transient collections, queues, tasks, and peer-selected identifiers follow the +ownership and cleanup invariants in +[Peer-controlled state ownership](peer-controlled-state.md). + ## Executive Summary `ngit-grasp` implements the GRASP protocol in Rust with **inline authorization** rather than Git hooks. Git push operations are intercepted and validated at the HTTP handler level before reaching the Git repository, eliminating the need for pre-receive hooks. diff --git a/docs/explanation/peer-controlled-state.md b/docs/explanation/peer-controlled-state.md new file mode 100644 index 0000000..27209a2 --- /dev/null +++ b/docs/explanation/peer-controlled-state.md @@ -0,0 +1,69 @@ +# Peer-controlled state ownership + +This document records the memory-lifecycle audit prompted by the 2026-08-09 +`auth-required` incident. It covers state whose cardinality or lifetime can be +influenced by an inbound client, an outbound relay, or an advertised Git +source. Ordinary growth of admitted events and repositories is intentionally +out of scope: that is retained application data, not transient peer state. + +The classifications used below are: + +- **Static**: a local constant, configured ceiling, or semaphore bounds the + number of entries. +- **Time**: every entry has an expiry or a finite retry budget. +- **External**: entries are deduplicated against admitted repository/event + state and therefore cannot outgrow that retained state. +- **Session**: entries are bounded by a connection ledger and cleared when the + connection generation ends. + +No transient peer-controlled collection is intentionally unbounded. An +external bound is acceptable only where admission policy already owns the +larger retained set; it is not an excuse to copy arbitrary wire input. + +## Inventory + +| State | Producer and cleanup owner | Bound and terminal behaviour | Observation | +|---|---|---|---| +| Inbound WebSocket connections | rust-nostr accepts sockets; disconnect owns cleanup | External to the process by the configured/OS connection ceiling. Per-IP accounting uses only a trusted proxy chain. Repeated disconnect is idempotent. | `ngit_connections_active`, unique-IP and abuser gauges | +| Inbound subscriptions | rust-nostr REQ/CLOSE handling | Static per connection: configured subscription count, cumulative subscription-state bytes, request/filter sizes, and event rate. Disconnect clears the connection registry. | NIP-11 limits, connection logs | +| Git HTTP request/response streams | HTTP handlers and child-process pumps | Static queues of `STREAM_CHANNEL_DEPTH = 8`; body and WebSocket message limits bound buffered input. Receiver loss terminates pumps. | request status and process-resource metrics | +| Per-relay EVENT data lane | `RelayConnection::run_event_loop`; processor task owns drain | Static 1,000 entries per connected relay. Disconnect drops the receiver. Delayed or missing terminals cannot enlarge it. | pipeline queue-depth/delay log | +| EOSE/CLOSED actor lanes | per-relay event processors; sync actor owns drain | Static 1,000 entries globally per lane. Repeated terminals backpressure the data lane; permit release is independently handled by the terminal-control listener. | retained-state and delayed-EOSE logs | +| Transient and live subscription permits | subscription ledger; EOSE/CLOSED/CLOSE/disconnect own release | Session bound by NIP-11 `max_subscriptions` or the fallback budget. The watchdog sends CLOSE before releasing a missing-terminal slot. Repeated terminals are idempotent. | subscription-ledger and watchdog logs | +| Authentication retry markers | first `auth-required` CLOSED; second refusal/disconnect owns cleanup | Session bound: an ID is admitted only while it owns a live or transient ledger permit. Forged and stale peer IDs are ignored. | `ngit_sync_retained_state_current{class="auth_retries"}` | +| Connection attempts/workers/results | one token per canonical derived relay; actor owns result | External queue: at most one queued/running task per relay in admitted coverage. Static execution: eight semaphore holders and an eight-entry result channel. Stale/reordered results cannot consume a newer token. | connection-attempt counters and `queued_connection_attempts` gauge | +| Self-subscriber actions and disconnect notices | subscriber/connection tasks; actor owns drain | Static 100-entry channels. Producers await capacity; duplicate dirty-relay actions are re-derived from indexes. | actor and connection logs | +| Pending sync batches, pagination, and requested/received IDs | historic/negentropy admission; terminal handler owns completion | Session/external: open subscriptions require ledger permits; pagination state is nested in a pending batch. EOSE/CLOSED, watchdog CLOSE, disconnect, or daily reset removes it. | retained pending-batch/subscription gauges | +| Missing-event recovery | incomplete negentropy hydration; recovery task owns progress/expiry | Time and static per attempt: 300 IDs per pass, 16 source batches retained, 12 zero-progress attempts. Duplicate IDs merge. | recovery outcome logs | +| Rejected-event hot/cold indexes | write policy; invalidation and cleanup tasks | Time: full events expire after two minutes and metadata after seven days; duplicate IDs replace/merge. Startup restores remaining lifetime rather than resetting it. | hot/cold current and expiry metrics | +| Purgatory event maps and sync queue | write policy; promotion/cleanup owns removal | Time/external: entries are keyed by admitted event/repository identity, normally expire at 30 minutes, and soft-expired announcements at 24 hours. Queue entries deduplicate by identifier and disappear on completion or event expiry. | purgatory counts, queue and Git-process metrics | +| Per-domain Git throttle queues | incomplete purgatory fetch; throttle manager owns drain | External: one entry per purgatory identifier/domain, merged on repeat. Completion, URL exhaustion, or purgatory expiry removes useful work. Request history is time-windowed. | domain/fetch logs and Git-process metrics | +| Dependency retry attempts and temporary relays | rejected/purgatory dependency discovery; maintenance owns expiry | Time/external: event IDs are pruned against current purgatory input; temporary relays have explicit deadlines. Repeats overwrite timestamps. | retained dependency gauges | +| NIP-65 discovery state/results | accepted root authors; discovery scheduler owns completion | External plus static work: author/source maps derive from accepted roots, one discovery query is in flight, and its result channel has capacity one. Missing results time out and clear in-flight state through result handling/disconnect refresh. | discovery logs and relay gauges | +| Deferred consolidation | capacity refusal; final batch/reset/disconnect owns removal | External and deduplicated by relay. Final batch completion processes the set directly; no self-addressed notification queue remains. | retained deferred-consolidation gauge | +| Descendant rotations and auxiliary live coverage | accepted root coverage; EOSE/CLOSED/disconnect/daily reset own transition | Session/external: one state object per derived relay and at most one rotating request in flight per relay. Coverage IDs consume ledger slots. | retained rotation gauge and terminal logs | +| Health and naughty-list entries | connection failures; health checker owns recovery/expiry | External/time: one entry per canonical target, ordinary failures back off, persistent entries expire after 12 hours. Metrics expose only three fixed categories. | health gauges and aggregate naughty metrics | +| Outbound connection and desired-coverage indexes | admitted announcements/root events; disconnect/removal owns session state | External: canonical relay URLs and desired items derive from admitted or purgatory data. Reconnect replaces session-only state rather than duplicating it. | tracked/connected relay gauges | +| Spawned tasks | listeners, bounded workers, timers, and Git subprocesses | Static concurrency or one task per owned connection/request. Shutdown receivers terminate service tasks; subprocess watchdogs and receiver loss terminate request tasks. No task is spawned merely to retain an item that failed a full queue. | process task/resource metrics and lifecycle logs | + +## Invariants for future changes + +1. Wire-provided identifiers may index state only after they resolve to an + application-owned connection, subscription permit, admitted event, or + purgatory entry. +2. Backpressure must bound retained memory, not merely concurrent execution. + A semaphore in front of an unbounded waiting queue is not a memory bound. +3. Terminal handling is idempotent. Missing terminals require a finite + close-before-release path; repeated or reordered terminals must not create + replacement state. +4. Reconnect resets every session-class collection. A previous generation may + never release, retire, or restore a current generation's subscription. +5. Prometheus labels use fixed vocabularies for peer failures and retained + state. Peer URLs, event IDs, subscription IDs, and error text do not create + new diagnostic series unless their cardinality is already explicitly + bounded and retired. + +The aggregate `ngit_sync_retained_state_current` gauge is the first alerting +surface for unexpected transient growth. Process RSS and cgroup pressure remain +the final defence: a flat work-concurrency graph with a rising retained-state +class indicates a lifecycle bug rather than legitimate throughput.