mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
Merge #be6422ae: security: bound peer-controlled transient state
nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqstuepz4mknkrlxznjptkwkeu4c2up5ya23tun8ke7rlr7qdptzdsc2ae0fd PR-Author: DanConwayDev's Agent nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0 CoverNote: Issue nostr:nevent1qqsv09nf4dy084sk8vyl5yyg5v6e3ujfrxmuc93zmrc68nqeayup62cpz3mhxue69uhhyetvv9ujumn8d96zuer9wca4ndyx follows the auth-required OOM incident by auditing every transient collection whose lifetime can be influenced by a peer. This stack fixes four concrete retention surfaces: - EOSE/CLOSED actor notifications now have fixed backpressure while the independent terminal lane still releases permits; - auth retry markers accept only subscription IDs owned by the current session ledger; - deferred consolidation processes its deduplicated actor-owned state directly instead of accumulating self-wakeups; - connection and NIP-65 worker result channels encode their existing producer bounds. It also exports fixed-cardinality aggregate gauges for important long-lived sync state and adds the architecture inventory recording producer, cleanup owner, cardinality class, terminal behavior, and operational signal for each subsystem. Normal retained event/database growth remains out of scope. Validation: - `cargo test --locked --lib`: 749 passed; - serialized `cargo test --locked --test sync -- --test-threads=1`: 96 passed, one deliberately ignored; - formatting and warnings-as-errors clippy passed throughout the atomic commits; - gitnostr.com is running exact revised tip `70dec5ac295dbedd340a9be5a8db873f4b4040fb`; - production startup exercised 173 queued canonical relay attempts while execution remained under the existing eight-worker semaphore; the corrected gauge distinguishes this externally bounded queue from simultaneous dials; - retained state remained small under startup traffic: four pending batches/subscriptions, zero auth retries and deferred consolidations, ten temporary dependency relays, 86 purgatory dependency attempts, and eight descendant rotations; - zero service restarts; memory settled around 1.21 GiB with a 1.22 GiB peak; - no lifecycle-channel saturation, forged-auth growth, panic, or proposal-related application error was observed. SDK error logs were malformed third-party events, unsupported NEG responses, an ordinary local reset, and an expected private-relay auth rejection. Recommendation: ready to merge.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
@@ -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::<i64>(),
|
||||
15
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_rejected_events_metrics() {
|
||||
let registry = create_test_registry();
|
||||
|
||||
+113
-88
@@ -1173,11 +1173,14 @@ struct SubscriptionClosedNotification {
|
||||
live_filter_count: Option<usize>,
|
||||
}
|
||||
|
||||
fn lifecycle_notification_channel<T>() -> (
|
||||
tokio::sync::mpsc::UnboundedSender<T>,
|
||||
tokio::sync::mpsc::UnboundedReceiver<T>,
|
||||
) {
|
||||
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<T>(
|
||||
) -> (tokio::sync::mpsc::Sender<T>, tokio::sync::mpsc::Receiver<T>) {
|
||||
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<String>>,
|
||||
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<tokio::sync::mpsc::Sender<DisconnectNotification>>,
|
||||
/// Channel for EOSE notifications (set during run)
|
||||
eose_tx: Option<tokio::sync::mpsc::UnboundedSender<EoseNotification>>,
|
||||
eose_tx: Option<tokio::sync::mpsc::Sender<EoseNotification>>,
|
||||
/// Serializes CLOSED recovery and pending-batch cleanup through the actor.
|
||||
subscription_closed_tx:
|
||||
Option<tokio::sync::mpsc::UnboundedSender<SubscriptionClosedNotification>>,
|
||||
subscription_closed_tx: Option<tokio::sync::mpsc::Sender<SubscriptionClosedNotification>>,
|
||||
/// Returns connection outcomes to the sync actor for serialized state changes.
|
||||
connect_attempt_result_tx: Option<tokio::sync::mpsc::UnboundedSender<ConnectAttemptResult>>,
|
||||
/// Wakes the actor when a batch completion may unblock consolidation.
|
||||
deferred_consolidation_tx: Option<tokio::sync::mpsc::UnboundedSender<String>>,
|
||||
nip65_discovery_result_tx: Option<tokio::sync::mpsc::UnboundedSender<Nip65DiscoveryResult>>,
|
||||
connect_attempt_result_tx: Option<tokio::sync::mpsc::Sender<ConnectAttemptResult>>,
|
||||
nip65_discovery_result_tx: Option<tokio::sync::mpsc::Sender<Nip65DiscoveryResult>>,
|
||||
/// Channel for broadcasting shutdown signal to all background tasks
|
||||
shutdown_tx: Option<broadcast::Sender<()>>,
|
||||
/// 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::<DisconnectNotification>(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::<EoseNotification>();
|
||||
let (subscription_closed_tx, mut subscription_closed_rx) =
|
||||
lifecycle_notification_channel::<SubscriptionClosedNotification>();
|
||||
|
||||
// 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::<ConnectAttemptResult>();
|
||||
mpsc::channel::<ConnectAttemptResult>(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::<String>();
|
||||
let (nip65_discovery_result_tx, mut nip65_discovery_result_rx) =
|
||||
mpsc::unbounded_channel::<Nip65DiscoveryResult>();
|
||||
mpsc::channel::<Nip65DiscoveryResult>(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::<EoseNotification>();
|
||||
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}"
|
||||
|
||||
@@ -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());
|
||||
|
||||
Reference in New Issue
Block a user