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. 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 562c06e..a796af4 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 { @@ -1542,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 // ============================================================================= @@ -1812,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() => { @@ -1913,15 +1928,12 @@ 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. - deferred_consolidation_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) @@ -2025,7 +2037,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, @@ -3168,11 +3179,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 @@ -3384,25 +3398,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); - // 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::(); + mpsc::channel::(1); // 4b. Create shutdown broadcast channel for graceful shutdown let (shutdown_tx, _shutdown_rx) = broadcast::channel(1); @@ -3424,7 +3433,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()); @@ -3540,12 +3548,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; @@ -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"); @@ -4805,6 +4811,7 @@ impl SyncManager { token, outcome, }) + .await .is_err() { connection.disconnect().await; @@ -5014,11 +5021,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; } @@ -6643,6 +6652,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, @@ -6655,11 +6675,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) @@ -7885,7 +7903,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 +8260,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}; @@ -8563,28 +8604,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), @@ -8597,16 +8629,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}" 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());