From ddcfa3c5f38816bd9ebeaf52714149204443db0c Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 14 Sep 2026 10:05:18 +0000 Subject: [PATCH 1/5] fix(http): avoid delayed-ACK stalls in small relay responses Accepted sockets retained Nagle while the relay flushes individual EVENT and EOSE frames. A delayed-ACK peer could hold small response batches for about 40 ms, including when the immediate peer is a loopback proxy. Enable TCP_NODELAY before handing each accepted socket to Hyper. Reject only that connection if setting the option fails. Protocol framing, proxy configuration and client deadlines remain unchanged. This addresses a measured transport delay, not every reported production timeout. Retain opt-in delayed-ACK and concurrent LMDB read benchmarks for the HTTP/WebSocket path. The controlled delayed-ACK reproduction originally fell from 1.050 s to 15.17 ms for 25 two-event reads. Benchmarks have bounded completion deadlines and validate complete EVENT/EOSE responses; latency assertions remain opt-in because busy CI workers can distort timings. Validation: the original before/after benchmark passed; the rebuilt PR is validated again after restoring its released SDK dependency. --- CHANGELOG.md | 5 ++ docs/explanation/architecture.md | 3 + src/http/mod.rs | 7 ++ tests/relay_response_latency.rs | 146 +++++++++++++++++++++++++++++++ 4 files changed, 161 insertions(+) create mode 100644 tests/relay_response_latency.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index dcbedff..8c42beb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- Avoid delayed-acknowledgement stalls in small WebSocket and Git responses by + enabling `TCP_NODELAY` on accepted connections. + ## [3.0.2] - 2026-09-11 This release contains no production runtime changes. It improves release and diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 575900e..f7f359a 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -75,6 +75,9 @@ runtime: - Bind the TCP listener (supports `:0` for kernel-assigned ports; the resolved address back-fills `bind_address` and, if empty, `domain`) +- Disable Nagle on accepted sockets so small EVENT/EOSE and streaming Git + writes do not wait for the peer's delayed acknowledgements, including + when the immediate peer is a local reverse proxy - Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51) - Set up shared storage (LMDB or Memory), purgatory, sync manager, and background maintenance tasks diff --git a/src/http/mod.rs b/src/http/mod.rs index bedd3dc..1c2c2c6 100644 --- a/src/http/mod.rs +++ b/src/http/mod.rs @@ -1073,6 +1073,13 @@ pub async fn run_server_on_listener( loop { let (socket, addr) = listener.accept().await?; + // WebSocket EVENT/EOSE and streaming Git responses flush small writes. + // Nagle can hold the tail until the peer's delayed ACK (about 40 ms + // even through a loopback proxy), multiplying multi-query latency. + if let Err(error) = socket.set_nodelay(true) { + tracing::warn!(peer = %addr, %error, "Could not disable Nagle for client socket"); + continue; + } let io = TokioIo::new(socket); let service = HttpService::new( relay.clone(), diff --git a/tests/relay_response_latency.rs b/tests/relay_response_latency.rs new file mode 100644 index 0000000..52ff71a --- /dev/null +++ b/tests/relay_response_latency.rs @@ -0,0 +1,146 @@ +//! Opt-in Linux delayed-ACK benchmark for the accepted HTTP/WebSocket socket. + +mod common; + +#[tokio::test] +#[ignore = "local load benchmark; run explicitly with --nocapture"] +async fn concurrent_lmdb_read_batches_reach_eose() { + use common::{TestClient, TestRelay}; + use futures_util::{SinkExt, StreamExt}; + use nostr_sdk::prelude::*; + use std::time::{Duration, Instant}; + use tokio_tungstenite::tungstenite::Message; + + let relay = TestRelay::start_with_lmdb().await; + let keys = relay.owner_keys().clone(); + let client = TestClient::new(relay.url(), keys.clone()).await.unwrap(); + for index in 0..32 { + let event = EventBuilder::new(Kind::TextNote, format!("{index}:{}", "x".repeat(8192))) + .finalize(&keys) + .unwrap(); + client.send_event(&event).await.unwrap(); + } + for concurrency in [1, 4, 16, 32] { + let url = relay.url().to_string(); + let reads = (0..concurrency).map(|_| { + let url = url.clone(); + async move { + let (mut stream, _) = tokio_tungstenite::connect_async(url).await.unwrap(); + let started = Instant::now(); + for _ in 0..3 { + stream + .send(Message::Text(r#"["REQ","load",{"kinds":[1]}]"#.into())) + .await + .unwrap(); + let mut events = 0; + loop { + let frame = stream.next().await.unwrap().unwrap(); + if !frame.is_text() { + continue; + } + let message: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + match message[0].as_str() { + Some("EVENT") => events += 1, + Some("EOSE") => break, + other => panic!("unexpected reply: {other:?}"), + } + } + assert_eq!(events, 32); + } + let elapsed = started.elapsed(); + stream.close(None).await.unwrap(); + elapsed + } + }); + let times = tokio::time::timeout( + Duration::from_secs(30), + futures_util::future::join_all(reads), + ) + .await + .expect("all readers must reach EOSE under local load"); + eprintln!( + "{concurrency} concurrent LMDB readers, three 32-event batches each: max {:?}", + times.iter().max().unwrap() + ); + } + client.disconnect().await; + relay.stop().await; +} + +#[cfg(target_os = "linux")] +#[tokio::test] +#[ignore = "socket latency benchmark; run explicitly on an otherwise idle worker"] +async fn small_event_batches_do_not_wait_for_delayed_ack() { + use std::os::fd::AsRawFd; + use std::time::{Duration, Instant}; + + use common::{TestClient, TestRelay}; + use futures_util::{SinkExt, StreamExt}; + use nostr_sdk::prelude::*; + use tokio_tungstenite::tungstenite::Message; + + let relay = TestRelay::start().await; + let keys = relay.owner_keys().clone(); + let client = TestClient::new(relay.url(), keys.clone()).await.unwrap(); + for index in 0..2 { + let event = EventBuilder::new(Kind::TextNote, format!("event {index}")) + .finalize(&keys) + .unwrap(); + client.send_event(&event).await.unwrap(); + } + let socket = tokio::net::TcpStream::connect(relay.domain()) + .await + .unwrap(); + socket.set_nodelay(true).unwrap(); + let (mut stream, _) = tokio_tungstenite::client_async(relay.url(), socket) + .await + .unwrap(); + let started = Instant::now(); + tokio::time::timeout(Duration::from_secs(10), async { + for _ in 0..25 { + let disabled: libc::c_int = 0; + // Explicitly model a peer that uses delayed acknowledgements. + // The descriptor and option value remain alive for setsockopt. + let result = unsafe { + libc::setsockopt( + stream.get_ref().as_raw_fd(), + libc::IPPROTO_TCP, + libc::TCP_QUICKACK, + &disabled as *const _ as *const libc::c_void, + std::mem::size_of_val(&disabled) as libc::socklen_t, + ) + }; + assert_eq!(result, 0); + stream + .send(Message::Text(r#"["REQ","latency",{"kinds":[1]}]"#.into())) + .await + .unwrap(); + let mut events = 0; + loop { + let frame = stream.next().await.unwrap().unwrap(); + if !frame.is_text() { + continue; + } + let message: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + match message[0].as_str() { + Some("EVENT") => events += 1, + Some("EOSE") => break, + other => panic!("unexpected reply: {other:?}"), + } + } + assert_eq!(events, 2); + } + }) + .await + .expect("batch reads must complete"); + let elapsed = started.elapsed(); + eprintln!("25 two-event reads with delayed ACK: {elapsed:?}"); + assert!( + elapsed < Duration::from_millis(500), + "small batches were delayed: {elapsed:?}" + ); + client.disconnect().await; + relay.stop().await; +} From fe36d6f7e4a05e71a6e85f62392ad11589afcfcc Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 14 Sep 2026 10:05:48 +0000 Subject: [PATCH 2/5] fix(storage): isolate periodic checkpoints from async workers Recovery snapshots serialized JSON and performed durable writes inline on a Tokio worker. Rejected-event snapshots also retained cache read locks through the write, making ingestion wait for filesystem I/O. Run one checkpoint at a time on the blocking pool, release cache guards once the owned snapshot exists, and measure the next interval from write completion. Join an active checkpoint before the final shutdown snapshot: started blocking work cannot be canceled safely, and an older snapshot must not replace the final one. Keep the on-disk format and atomic replacement unchanged. Startup and final shutdown writes remain synchronous, outside the serving interval. This isolates periodic I/O; it does not speed up storage or establish that checkpoint stalls caused the historical production timeouts. Validation: current-thread runtime tests cover progress during a blocked checkpoint, shutdown joining that write, and stopping before the first run. The existing snapshot persistence tests cover the unchanged file format. --- CHANGELOG.md | 3 + docs/explanation/architecture.md | 6 +- src/checkpoint.rs | 110 +++++++++++++++++++++++++++++++ src/lib.rs | 1 + src/server.rs | 13 ++-- src/sync/rejected_index.rs | 9 +++ 6 files changed, 134 insertions(+), 8 deletions(-) create mode 100644 src/checkpoint.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 8c42beb..1a7652e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Keep periodic recovery checkpoints off async workers, release snapshot locks + before disk I/O, and finish active checkpoints before the shutdown snapshot. + - Avoid delayed-acknowledgement stalls in small WebSocket and Git responses by enabling `TCP_NODELAY` on accepted connections. diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index f7f359a..3058ca3 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -92,7 +92,11 @@ runtime: Transient failures retry with a capped backoff, and in private mode the identity stays local so the relay's existence is never advertised - Atomically checkpoint purgatory and rejected-event recovery state every 60 - seconds without consuming the checkpoint during restore + seconds without consuming the checkpoint during restore. Serialization and + durable filesystem writes run on a blocking worker, with at most one + checkpoint active and the next interval measured from completion. Shutdown + joins any active checkpoint before writing the final snapshot, preventing + an older background write from replacing it - Start non-blocking storage-integrity and authorization-integrity passes after database initialization. The former checks family objects and thin-view wiring; the latter reconciles each served ref against accepted State, PR, diff --git a/src/checkpoint.rs b/src/checkpoint.rs new file mode 100644 index 0000000..1d7867a --- /dev/null +++ b/src/checkpoint.rs @@ -0,0 +1,110 @@ +//! Serialized blocking checkpoints with an explicit shutdown barrier. + +use std::sync::Arc; +use std::time::Duration; + +use tokio::sync::watch; +use tokio::task::JoinHandle; + +pub(crate) struct CheckpointTask { + stop: watch::Sender, + task: JoinHandle<()>, +} + +impl CheckpointTask { + pub(crate) fn start(interval: Duration, save: impl Fn() + Send + Sync + 'static) -> Self { + let (stop, mut stopped) = watch::channel(false); + let save = Arc::new(save); + let task = tokio::spawn(async move { + loop { + // Schedule from completion: slow storage must not accumulate + // overdue checkpoints or start overlapping snapshots. + tokio::select! { + biased; + _ = stopped.changed() => break, + _ = tokio::time::sleep(interval) => {} + } + let save = save.clone(); + // Do not race shutdown against this join. Blocking tasks + // cannot be aborted after they start, and an old snapshot + // must never overwrite the final shutdown snapshot. + if let Err(error) = tokio::task::spawn_blocking(move || save()).await { + tracing::error!(%error, "Checkpoint worker failed"); + } + } + }); + Self { stop, task } + } + + pub(crate) async fn shutdown(self) { + let _ = self.stop.send(true); + let _ = self.task.await; + } +} + +#[cfg(test)] +mod tests { + use std::sync::{mpsc, Mutex}; + + use futures_util::FutureExt; + use tokio::sync::oneshot; + + use super::*; + + #[tokio::test] + async fn blocking_checkpoint_yields_runtime_and_shutdown_joins_it() { + let (started_tx, started_rx) = oneshot::channel(); + let started_tx = Mutex::new(Some(started_tx)); + let (release_tx, release_rx) = mpsc::channel(); + let release_rx = Mutex::new(release_rx); + let (finished_tx, finished_rx) = oneshot::channel(); + let finished_tx = Mutex::new(Some(finished_tx)); + let checkpoint = CheckpointTask::start(Duration::from_millis(1), move || { + started_tx.lock().unwrap().take().unwrap().send(()).unwrap(); + // A bounded blocking I/O surrogate. The sole runtime thread must + // remain free to send release; an inline checkpoint fails here. + release_rx + .lock() + .unwrap() + .recv_timeout(Duration::from_secs(5)) + .unwrap(); + finished_tx + .lock() + .unwrap() + .take() + .unwrap() + .send(()) + .unwrap(); + }); + + tokio::time::timeout(Duration::from_secs(5), started_rx) + .await + .unwrap() + .unwrap(); + let shutdown = checkpoint.shutdown(); + tokio::pin!(shutdown); + assert!( + shutdown.as_mut().now_or_never().is_none(), + "shutdown must join the active snapshot" + ); + release_tx.send(()).unwrap(); + tokio::time::timeout(Duration::from_secs(5), shutdown) + .await + .unwrap(); + assert!(finished_rx.now_or_never().unwrap().is_ok()); + } + + #[tokio::test(start_paused = true)] + async fn shutdown_cancels_a_checkpoint_that_has_not_started() { + use std::sync::atomic::{AtomicUsize, Ordering}; + + let saves = Arc::new(AtomicUsize::new(0)); + let worker_saves = saves.clone(); + let checkpoint = CheckpointTask::start(Duration::from_secs(60), move || { + worker_saves.fetch_add(1, Ordering::SeqCst); + }); + checkpoint.shutdown().await; + tokio::time::advance(Duration::from_secs(120)).await; + assert_eq!(saves.load(Ordering::SeqCst), 0); + } +} diff --git a/src/lib.rs b/src/lib.rs index 69a927b..c3962a0 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,5 +1,6 @@ mod atomic_file; pub mod audit_cleanup; +mod checkpoint; pub mod cleanup_empty_repos; pub mod config; pub mod git; diff --git a/src/server.rs b/src/server.rs index d9ae464..5cef2cd 100644 --- a/src/server.rs +++ b/src/server.rs @@ -66,6 +66,7 @@ pub struct RelayServer { /// cleanup, purgatory sync loop). Aborted on shutdown so the host /// process does not leak tasks per relay instance. background_tasks: Vec>, + checkpoint: crate::checkpoint::CheckpointTask, git_data_path: String, } @@ -363,11 +364,8 @@ impl RelayServer { let checkpoint_purgatory = purgatory.clone(); let checkpoint_rejected = rejected_events_index.clone(); let checkpoint_root = PathBuf::from(config.effective_git_data_path()); - background_tasks.push(tokio::spawn(async move { - let first = tokio::time::Instant::now() + SYNC_STATE_CHECKPOINT_INTERVAL; - let mut interval = tokio::time::interval_at(first, SYNC_STATE_CHECKPOINT_INTERVAL); - loop { - interval.tick().await; + let checkpoint = + crate::checkpoint::CheckpointTask::start(SYNC_STATE_CHECKPOINT_INTERVAL, move || { let purgatory_path = checkpoint_root.join("purgatory-state.json"); if let Err(error) = checkpoint_purgatory.save_to_disk(&purgatory_path) { warn!(%error, "Failed to checkpoint purgatory state"); @@ -376,8 +374,7 @@ impl RelayServer { if let Err(error) = checkpoint_rejected.save_to_disk(&rejected_path) { warn!(%error, "Failed to checkpoint rejected-events cache"); } - } - })); + }); info!( interval_secs = SYNC_STATE_CHECKPOINT_INTERVAL.as_secs(), "Crash-safe sync-state checkpoint task started" @@ -492,6 +489,7 @@ impl RelayServer { private_access, deletion_cleanup, background_tasks, + checkpoint, git_data_path, }) } @@ -547,6 +545,7 @@ impl RelayServer { task.abort(); let _ = task.await; } + self.checkpoint.shutdown().await; self.deletion_cleanup.shutdown().await; diff --git a/src/sync/rejected_index.rs b/src/sync/rejected_index.rs index 718fca3..93b27f4 100644 --- a/src/sync/rejected_index.rs +++ b/src/sync/rejected_index.rs @@ -1357,6 +1357,15 @@ impl RejectedEventsIndex { related_dependencies: serializable_related_entries, }; + // The snapshot owns its data. Ingestion must not wait on JSON + // encoding, fsync or rename while holding these synchronous locks. + drop(( + hot_entries, + cold_entries, + unrecoverable_entries, + related_entries, + )); + // Replace the previous checkpoint only after the new snapshot is // complete and durable. An abrupt stop must leave one valid version. let json = serde_json::to_string_pretty(&state)?; From f55ce4885ba726bfadcd35b64161431676076002 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 14 Sep 2026 10:05:48 +0000 Subject: [PATCH 3/5] fix(metrics): isolate scrape rendering from async workers Each metrics scrape traverses repository directories and encodes the response synchronously. Slow filesystem metadata operations could occupy a worker that also serves WebSocket and Git requests. Render on a blocking worker. Acquire a shared permit before spawning it and keep that permit inside the worker even if the HTTP request is canceled, so concurrent scrapes cannot start overlapping directory scans. Report a worker failure as an HTTP 500 response rather than losing the request task. Preserve the synchronous rendering API, metric contents and monitoring configuration. This does not cache metrics or claim a measured production speedup; it prevents scrape work from running on the async executor. Validation: a current-thread test checks that a queued scrape yields, capacity is released and async rendering preserves repository counts. --- CHANGELOG.md | 3 +++ docs/explanation/architecture.md | 2 ++ src/http/mod.rs | 13 ++++++++- src/metrics/mod.rs | 45 ++++++++++++++++++++++++++++++-- 4 files changed, 60 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1a7652e..311a2a2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Render metrics on a blocking worker and serialize scrapes so repository + counting does not occupy async workers serving relay and Git requests. + - Keep periodic recovery checkpoints off async workers, release snapshot locks before disk I/O, and finish active checkpoints before the shutdown snapshot. diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 3058ca3..3186132 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -78,6 +78,8 @@ runtime: - Disable Nagle on accepted sockets so small EVENT/EOSE and streaming Git writes do not wait for the peer's delayed acknowledgements, including when the immediate peer is a local reverse proxy +- Render metrics on a blocking worker, serializing concurrent scrapes before + spawning filesystem scans so monitoring cannot occupy the network executor - Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51) - Set up shared storage (LMDB or Memory), purgatory, sync manager, and background maintenance tasks diff --git a/src/http/mod.rs b/src/http/mod.rs index 1c2c2c6..7523073 100644 --- a/src/http/mod.rs +++ b/src/http/mod.rs @@ -923,7 +923,18 @@ impl Service> for HttpService { if let Some(ref metrics) = self.metrics { let metrics = metrics.clone(); return Box::pin(async move { - let output = metrics.render(); + let output = match metrics.render_async().await { + Ok(output) => output, + Err(error) => { + tracing::error!(%error, "Metrics render worker failed"); + return Ok(add_cors_headers( + Response::builder().header("server", "ngit-grasp"), + ) + .status(500) + .body(full_body("Metrics temporarily unavailable")) + .unwrap()); + } + }; Ok( add_cors_headers(Response::builder().header("server", "ngit-grasp")) .status(200) diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index bdd79f9..b9644a6 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -564,6 +564,7 @@ pub fn record_manual_ejection( #[derive(Clone)] pub struct Metrics { inner: Arc, + render_permit: Arc, } struct MetricsInner { @@ -628,6 +629,7 @@ impl Metrics { let inner = MetricsInner::new(abuse_threshold, git_data_path); Self { inner: Arc::new(inner), + render_permit: Arc::new(tokio::sync::Semaphore::new(1)), } } @@ -794,6 +796,24 @@ impl Metrics { // === Rendering === + /// Render a scrape without blocking the async executor on filesystem + /// traversal or encoding. Concurrent scrapes wait before spawning work; + /// cancelling an HTTP request cannot release a running worker's permit. + pub async fn render_async(&self) -> Result { + let permit = self + .render_permit + .clone() + .acquire_owned() + .await + .expect("metrics render semaphore is never closed"); + let metrics = self.clone(); + tokio::task::spawn_blocking(move || { + let _permit = permit; + metrics.render() + }) + .await + } + /// Render all metrics in Prometheus text format. /// /// This method: @@ -1104,6 +1124,21 @@ impl Drop for GitOperationTimer { mod tests { use super::*; + async fn concurrent_scrapes_wait_without_blocking_the_runtime(metrics: &Metrics) { + use futures_util::FutureExt; + let permit = metrics.render_permit.clone().acquire_owned().await.unwrap(); + let render = metrics.render_async(); + tokio::pin!(render); + assert!(render.as_mut().now_or_never().is_none()); + drop(permit); + let output = tokio::time::timeout(std::time::Duration::from_secs(5), render) + .await + .unwrap() + .unwrap(); + assert!(output.contains("ngit_uptime_seconds")); + assert_eq!(metrics.render_permit.available_permits(), 1); + } + #[test] fn test_count_repositories_on_disk() { use std::fs; @@ -1153,8 +1188,8 @@ mod tests { /// /// If additional Metrics tests are needed, they should either be added to this /// test or use a separate test-specific Prometheus registry. - #[test] - fn test_metrics_with_repository_counting() { + #[tokio::test] + async fn test_metrics_with_repository_counting() { use std::fs; use tempfile::TempDir; @@ -1211,5 +1246,11 @@ mod tests { // Render should count 3 repos let output = metrics.render(); assert!(output.contains("ngit_repositories_total 3")); + concurrent_scrapes_wait_without_blocking_the_runtime(&metrics).await; + assert!(metrics + .render_async() + .await + .unwrap() + .contains("ngit_repositories_total 3")); } } From 5d62037d94ce0049e2c21eeeb577ecd83bd1aa76 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 14 Sep 2026 10:06:15 +0000 Subject: [PATCH 4/5] fix(sync): preserve discovery relay retry ownership Historical mailbox and profile discovery owned connections without owning persistent subscription targets. Idle cleanup could therefore retire a failed discovery connection, erase its health history and let the next probe bypass the intended backoff. Include discovery sources, mailbox scopes and in-flight fetches in connection ownership for cleanup and reconnect scheduling. Preserve failed or policy-limited sessions until recovery or scope removal, including when a disconnect notification has not yet recorded the failure. This changes background synchronization only. Discovery ownership does not add live subscriptions, and healthy idle sessions can still retire. Existing backoff durations and persistent repository targets are unchanged. Validation: unit coverage checks ownership across deferred/in-flight work and scope removal. A real-socket mailbox-only reconnect regression checks that cleanup retains the failure streak through reconnection; it failed before the ownership correction. --- CHANGELOG.md | 3 ++ docs/explanation/architecture.md | 8 ++++ src/sync/mod.rs | 72 +++++++++++++++++++++++++++++++- tests/sync/reconnect_backoff.rs | 41 ++++++++++++++++++ 4 files changed, 122 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 311a2a2..4ac0d60 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Preserve background discovery relay backoff across idle-connection cleanup + so unavailable mailbox and profile sources do not restart their retry history. + - Render metrics on a blocking worker and serialize scrapes so repository counting does not occupy async workers serving relay and Git requests. diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 3186132..d49fd77 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -744,6 +744,14 @@ is deliberately historical-only: owning or maintaining a repository never adds that author's inbox relays to ordinary persistent live repository targets. Private instances omit repository-coordinate mailbox expansion. +Discovery source and mailbox scope also own connection retry state. The empty +relay checker and reconnect scheduler include that scope, including in-flight +fetches and future probe deadlines, without adding persistent subscriptions. +Healthy idle discovery sessions can retire after their fetch; failed or +policy-limited sessions retain their backoff until recovery or scope removal. +This prevents cleanup from erasing a failure every two seconds and allowing +discovery to dial the same unavailable endpoint from a fresh retry history. + ### Rejected Events Index The rejected events index solves two critical problems during sync: diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 41cc3d7..fb4972f 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1960,6 +1960,20 @@ struct Nip65DiscoveryState { } impl Nip65DiscoveryState { + /// Connection ownership includes historical discovery work, even while + /// its next query is deferred. It is not persistent subscription scope. + fn connection_targets(&self) -> HashSet { + self.author_sources + .values() + .flatten() + .chain(self.mailbox_roots.keys()) + .chain(self.mailbox_repositories.keys()) + .chain(self.mailbox_probes_in_flight.iter()) + .chain(self.in_flight.iter().map(|(relay, _)| relay)) + .cloned() + .collect() + } + fn has_mailbox_scope(&self, relay: &str) -> bool { self.mailbox_roots.contains_key(relay) || self.mailbox_repositories.contains_key(relay) } @@ -6422,6 +6436,22 @@ impl SyncManager { if !self.nip65_discovery_only_relays.contains(source) { return; } + // Retirement forgets connection health. A failed or policy-limited + // source must retain its session state until recovery or removal from + // discovery ownership, otherwise the next probe bypasses its backoff. + if self.health_tracker.get_failure_count(source) > 0 + || self.health_tracker.is_subscription_paused(source) + { + return; + } + if let Some(connection) = self.connections.get(source) { + if !connection.is_connected().await { + // The disconnect notification may still be queued behind a + // fetch result. Let it record the unexpected failure instead + // of relabeling the session as an intentional retirement. + return; + } + } let now = Instant::now(); let has_author_work = self.nip65_discovery @@ -7221,7 +7251,12 @@ impl SyncManager { // Once the connection ends naturally, however, reconcile its // confirmed state with the latest index before deciding whether it // should ever reconnect. - let desired = self.derive_targets().await.remove(relay_url); + let desired = self.derive_targets().await.remove(relay_url).or_else(|| { + self.nip65_discovery + .connection_targets() + .contains(relay_url) + .then(RelaySyncNeeds::default) + }); let Some(desired) = desired else { tracing::info!( relay = %relay_url, @@ -8420,6 +8455,7 @@ impl SyncManager { let mut desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); desired_relays.extend(self.dependency_relay_deadlines.keys().cloned()); + desired_relays.extend(self.nip65_discovery.connection_targets()); // Collect relays to disconnect let to_disconnect: Vec = { @@ -8531,7 +8567,8 @@ impl SyncManager { /// /// For each eligible relay, a reconnection is queued via schedule_connect_relay. async fn retry_disconnected_relays(&mut self) { - let desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); + let mut desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); + desired_relays.extend(self.nip65_discovery.connection_targets()); // Collect relays to reconnect let to_reconnect: Vec = { @@ -10971,6 +11008,37 @@ mod tests { ); } + #[test] + fn discovery_connection_ownership_survives_retry_deadlines_and_in_flight_removal() { + let author = Keys::generate().public_key(); + let profile = "wss://profile.example".to_string(); + let mailbox = "wss://mailbox.example".to_string(); + let retired = "wss://retired.example".to_string(); + let mut discovery = Nip65DiscoveryState::default(); + discovery + .author_sources + .insert(author, HashSet::from([profile.clone()])); + discovery + .mailbox_repositories + .insert(mailbox.clone(), HashSet::from(["repo".into()])); + discovery + .mailbox_probe_next_at + .insert(mailbox.clone(), Instant::now() + Duration::from_secs(3600)); + discovery.in_flight.insert((retired.clone(), author)); + + let disconnected = RelayState::default(); + for relay in [&profile, &mailbox, &retired] { + assert!(!disconnected + .is_disconnect_candidate(false, discovery.connection_targets().contains(relay))); + } + discovery.in_flight.clear(); + assert!(!discovery.connection_targets().contains(&retired)); + discovery.author_sources.clear(); + discovery.install_mailbox_overlay(HashMap::new(), HashMap::new(), Instant::now()); + assert!(discovery.connection_targets().is_empty()); + assert!(disconnected.is_disconnect_candidate(false, false)); + } + #[test] fn disconnected_empty_relay_can_still_be_cleaned_up() { let mut source = RelayState::default(); diff --git a/tests/sync/reconnect_backoff.rs b/tests/sync/reconnect_backoff.rs index ed496fc..3730563 100644 --- a/tests/sync/reconnect_backoff.rs +++ b/tests/sync/reconnect_backoff.rs @@ -7,6 +7,47 @@ use nostr_sdk::prelude::*; use crate::common::flapping_relay::FlappingRelay; use crate::common::{TestClient, TestRelay}; +#[tokio::test] +async fn discovery_only_mailbox_keeps_failure_history_through_cleanup() { + use crate::common::{send_to_relay_url, setup_announcement_on_relay, MockRelay}; + + let index = MockRelay::start().await; + let flapping = FlappingRelay::start().await; + let owner = Keys::generate(); + let relay_list = EventBuilder::new(Kind::RelayList, "") + .tags([Tag::custom("r", vec![flapping.url(), "read"])]) + .finalize(&owner) + .unwrap(); + send_to_relay_url(index.url(), &relay_list).await.unwrap(); + let syncing = TestRelay::start_with_sync(Some(index.url().to_string())).await; + let domain = syncing.domain(); + let (_announcement, _git) = + setup_announcement_on_relay(&syncing, &owner, &[&domain], "discovery-mailbox-backoff") + .await; + + // Owner inbox discovery has no ordinary repository/root live target. + // The two-second cleanup pass previously erased its first failure before + // the next handshake, so it could never reach this recovery state. + flapping + .wait_for_connections(2, Duration::from_secs(60)) + .await; + let logs = wait_for_log( + &syncing.log_path(), + "consecutive_failures=1", + Duration::from_secs(15), + ) + .await; + assert!(logs.lines().any(|line| { + line.contains(flapping.url()) + && line.contains("consecutive_failures=1") + && line.contains("preserving failure streak until stable") + })); + + syncing.stop().await; + flapping.stop().await; + index.stop().await; +} + async fn wait_for_log(log_path: &std::path::Path, needle: &str, timeout: Duration) -> String { tokio::time::timeout(timeout, async { loop { From df9280e596f0829c8d064d1eb983545928670514 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 14 Sep 2026 10:06:15 +0000 Subject: [PATCH 5/5] fix(sync): require EOSE before accepting discovery history The SDK convenience fetch can return partial events as success when a peer disconnects or the timeout expires. Mailbox and profile discovery then mistake an incomplete response for completed history or absence. Keep the SDK's validated, auto-closing event stream and AUTH handling, but independently observe the exact subscription's EOSE. Bound the whole fetch, retain the existing unique-event buffer cap and reject premature termination. Subscription permits and cancellation remain fetch-scoped. This changes outbound background discovery, not inbound user REQ handling. It does not alter historical pagination, discovery coverage, timeout budgets or the SDK's intentional completed-ID CLOSED behavior. Validation: the partial-delivery regression fails against the original implementation. Exact-relay, unregistered-relay, disconnect and silent-peer deadline cases exercise the completed-fetch contract over real sockets. --- CHANGELOG.md | 3 + docs/explanation/grasp-02-proactive-sync.md | 6 + src/sync/relay_connection.rs | 133 +++++++++++++++++++- 3 files changed, 138 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4ac0d60..ef4adda 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Require EOSE before accepting background discovery history, instead of treating + a disconnected or timed-out partial response as a completed fetch. + - Preserve background discovery relay backoff across idle-connection cleanup so unavailable mailbox and profile sources do not restart their retry history. diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 691e4c9..5ed61a8 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -40,6 +40,12 @@ Key Architectural Points: stable after rebuilding - **Quick Reconnect** (< 15mins) - doesn't do a full reconciliation vs fresh start (longer disconnect or relaunch binary) - **Background timers** handle relay connection health and metrics, handling reconnects after backoff and recovery after rate-limiting +- **Completed discovery reads** require both the exact REQ's EOSE and a + drained, validated event stream. A disconnect or local deadline before EOSE + is a failed fetch, even if some events arrived. This avoids treating partial + pages as successful mailbox history or an empty response as proof that no + NIP-65 relay list exists. The SDK still handles event validation and AUTH + retries; ngit-grasp additionally verifies the terminal protocol evidence. Sections: diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 429fa94..967ac43 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -2231,11 +2231,64 @@ impl RelayConnection { })?; self.ensure_current_session(&ledger_slot)?; - relay - .fetch_events(filter) - .timeout(timeout) + // The SDK event stream can end normally on timeout or disconnect, + // returning partial history as Ok. Discovery must not mistake that + // for a completed page (or proof that an author has no relay list). + // Observe this exact subscription's EOSE independently, retaining SDK + // event validation, deduplication and NIP-42 retries in the data lane. + let subscription_id = SubscriptionId::generate(); + let mut notifications = relay.notifications(); + let fetch = async { + let mut stream = relay + .stream_events(filter) + .with_id(subscription_id.clone()) + .await + .map_err(|error| error.to_string())?; + let mut events = std::collections::BTreeSet::new(); + let mut received_eose = false; + let mut drained = false; + while !received_eose || !drained { + tokio::select! { + item = stream.next(), if !drained => { + match item { + Some(Ok(event)) => { + // Match the SDK fetch API's existing buffer cap. + if events.len() >= 10_000 && !events.contains(&event) { + return Err("too many fetched events".to_string()); + } + events.insert(event); + } + Some(Err(error)) => return Err(error.to_string()), + None => drained = true, + } + } + notification = notifications.next(), if !received_eose => { + match notification { + Some(RelayNotification::Message { message }) => { + if matches!(*message, RelayMessage::EndOfStoredEvents(ref id) if id.as_ref() == &subscription_id) { + received_eose = true; + } + } + Some(RelayNotification::RelayStatus { + status: RelayStatus::Disconnected | RelayStatus::Terminated | RelayStatus::Banned, + }) | None => { + return Err("relay disconnected before EOSE".to_string()); + } + _ => {} + } + } + } + } + Ok(events.into_iter().collect()) + }; + tokio::time::timeout(timeout, fetch) .await - .map(|events| events.into_iter().collect()) + .map_err(|_| { + format!( + "Failed to fetch events from {}: timed out before complete EOSE", + self.url + ) + })? .map_err(|error| format!("Failed to fetch events from {}: {}", self.url, error)) } @@ -3109,6 +3162,78 @@ mod tests { assert!(!connection.supports_negentropy().await); } + #[tokio::test] + async fn fetch_events_requires_eose_after_partial_delivery() { + use futures_util::SinkExt; + use tokio_tungstenite::tungstenite::Message; + + for disconnect in [true, false] { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + let event = EventBuilder::new(Kind::TextNote, "incomplete history") + .finalize(&Keys::generate()) + .unwrap(); + let server = tokio::spawn(async move { + let (socket, _) = listener.accept().await.unwrap(); + let mut ws = tokio_tungstenite::accept_async(socket).await.unwrap(); + while let Some(frame) = ws.next().await { + let frame = frame.unwrap(); + if !frame.is_text() { + continue; + } + let request: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + if request[0] != "REQ" { + continue; + } + ws.send(Message::Text( + serde_json::json!(["EVENT", request[1], event]) + .to_string() + .into(), + )) + .await + .unwrap(); + if disconnect { + ws.close(None).await.unwrap(); + } else { + // Stay connected without EOSE until the client's + // bounded fetch cancellation sends CLOSE. + tokio::time::timeout(Duration::from_secs(5), async { + while let Some(frame) = ws.next().await { + let Ok(frame) = frame else { break }; + if frame.is_close() + || frame.to_text().is_ok_and(|text| text.contains("CLOSE")) + { + break; + } + } + }) + .await + .unwrap(); + } + break; + } + }); + let connection = permissive_connection(&url, Keys::generate()); + connection.connect(3).await.unwrap(); + let result = connection + .fetch_events( + Filter::new().kind(Kind::TextNote), + Duration::from_millis(250), + ) + .await; + connection.disconnect().await; + tokio::time::timeout(Duration::from_secs(5), server) + .await + .unwrap() + .unwrap(); + assert!( + result.is_err(), + "partial fetch without EOSE was accepted (disconnect={disconnect}): {result:?}" + ); + } + } + #[tokio::test] async fn fetch_events_targets_the_connections_exact_relay() { let configured = LocalRelayBuilder::default().build();