From c3bcd132bc740b8eb1600b741a5008b341373514 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Fri, 7 Aug 2026 16:49:30 +0000 Subject: [PATCH] fix(sync): reject stale connection setup success Production soak of bccbb54 showed nostr-pub.wellorder.net disconnect at 16:30:57 while its connection worker was still fetching NIP-11. The delayed result was accepted at 16:31:19, marked Healthy, and started fresh sync against an already dead SDK session. Because the event loop was spawned after the disconnect notification, that lifecycle could also remain falsely connected. Revalidate the SDK relay status when a successful setup result reaches the sync actor. A peer that vanished after the handshake but before setup completed is handled through the existing failed-attempt path, which restores Disconnected state and records backoff without starting subscriptions or an event loop. A deterministic dual-protocol fixture coordinates WebSocket teardown with the NIP-11 request and proves the stale result is rejected. The earlier flapping fixture now includes a short established dwell so it continues to model post-setup failure rather than this new setup race. Correctness assumes RelayStatus::Connected is the authoritative final setup gate. A disconnect immediately after that gate remains handled by the normal event-loop notification path. NIP-11 timeout policy, reconnect pacing, and broader lifecycle serialization are excluded. Validated with cargo test --lib (650 passed), cargo test --test sync (89 passed, 1 ignored), both focused reconnect lifecycle scenarios, and nix build .#ngit-grasp. --- docs/explanation/grasp-02-proactive-sync.md | 6 + src/sync/mod.rs | 24 +++- src/sync/relay_connection.rs | 14 +++ tests/common/flapping_relay.rs | 6 +- tests/common/mod.rs | 1 + tests/common/setup_drop_relay.rs | 133 ++++++++++++++++++++ tests/sync.rs | 1 + tests/sync/reconnect_backoff.rs | 2 +- tests/sync/stale_connect_result.rs | 76 +++++++++++ 9 files changed, 260 insertions(+), 3 deletions(-) create mode 100644 tests/common/setup_drop_relay.rs create mode 100644 tests/sync/stale_connect_result.rs diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 0c2dad5..c8fec8b 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -307,6 +307,12 @@ This allows operators to monitor sync progress and distinguish between "connecte **Critical**: Event loops die on disconnect and cannot be reused. +A successful WebSocket handshake is followed by a bounded NIP-11 fetch for +session limits. The connection worker revalidates the SDK relay status when +that setup result reaches the sync actor; if the peer disconnected meanwhile, +the stale success is recorded as a failed attempt and no event loop or sync +work is started for the dead session. + ```mermaid flowchart LR CONN[Connection Success] --> SPAWN[handle_connect_or_reconnect
spawns event loop] diff --git a/src/sync/mod.rs b/src/sync/mod.rs index b56464e..59a4050 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -3663,7 +3663,29 @@ impl SyncManager { return; } - match result.outcome { + let outcome = if matches!(result.outcome, ConnectAttemptOutcome::Connected { .. }) { + let still_connected = if let Some(connection) = self.connections.get(&result.relay_url) + { + connection.is_connected().await + } else { + false + }; + if still_connected { + result.outcome + } else { + tracing::warn!( + relay = %result.relay_url, + "Rejecting stale connection success after relay disconnected during setup" + ); + ConnectAttemptOutcome::Failed( + "Relay disconnected before connection setup completed".to_string(), + ) + } + } else { + result.outcome + }; + + match outcome { ConnectAttemptOutcome::Connected { advertised_default_limit, advertised_max_subscriptions, diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 271948e..33ac32a 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -620,6 +620,20 @@ impl RelayConnection { } } + /// Whether the SDK still considers this relay's WebSocket established. + /// + /// Connection setup performs bounded HTTP work after the handshake. The + /// peer can disappear during that window, so callers must revalidate the + /// session before committing its successful lifecycle transition. + pub async fn is_connected(&self) -> bool { + self.client + .relay(&self.url) + .await + .ok() + .flatten() + .is_some_and(|relay| relay.status() == RelayStatus::Connected) + } + /// Configure the one per-session ledger before any subscriptions open. pub fn reset_subscription_budget(&self, advertised: Option) { self.clear_subscription_permits(); diff --git a/tests/common/flapping_relay.rs b/tests/common/flapping_relay.rs index 857fbdd..cd08627 100644 --- a/tests/common/flapping_relay.rs +++ b/tests/common/flapping_relay.rs @@ -51,7 +51,11 @@ impl FlappingRelay { let Ok(message) = message else { break }; if matches!(message, Message::Text(ref text) if text.contains("\"REQ\"")) { // Dropping the socket after observable normal - // setup models the production peer reset. + // setup and a short established dwell models + // the production peer reset. The dwell also + // distinguishes this fixture from a peer that + // vanishes during connection setup. + tokio::time::sleep(Duration::from_millis(100)).await; return; } } diff --git a/tests/common/mod.rs b/tests/common/mod.rs index aeaca7b..f3461e4 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -13,6 +13,7 @@ pub mod upload_pack_counting_proxy; pub mod port; pub mod purgatory_helpers; pub mod relay; +pub mod setup_drop_relay; pub mod sync_helpers; pub use git_server::{SimpleGitServer, SmartGitServer}; diff --git a/tests/common/setup_drop_relay.rs b/tests/common/setup_drop_relay.rs new file mode 100644 index 0000000..1c23d4f --- /dev/null +++ b/tests/common/setup_drop_relay.rs @@ -0,0 +1,133 @@ +//! Relay fixture that disconnects the WebSocket during its NIP-11 setup fetch. + +use std::sync::Arc; +use std::time::Duration; + +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::{TcpListener, TcpStream}; +use tokio::sync::{oneshot, watch}; + +/// Endpoint for reproducing a stale successful connection result. +pub struct SetupDropRelay { + url: String, + shutdown_tx: Option>, + handle: Option>, +} + +impl SetupDropRelay { + pub async fn start() -> Self { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("SetupDropRelay failed to bind"); + let port = listener.local_addr().expect("setup-drop local_addr").port(); + let (nip11_tx, nip11_rx) = watch::channel(false); + let (dropped_tx, dropped_rx) = watch::channel(false); + let nip11_tx = Arc::new(nip11_tx); + let dropped_tx = Arc::new(dropped_tx); + let (shutdown_tx, mut shutdown_rx) = oneshot::channel(); + + let handle = tokio::spawn(async move { + loop { + tokio::select! { + accepted = listener.accept() => { + let Ok((stream, _)) = accepted else { break }; + let nip11_tx = nip11_tx.clone(); + let nip11_rx = nip11_rx.clone(); + let dropped_tx = dropped_tx.clone(); + let dropped_rx = dropped_rx.clone(); + tokio::spawn(async move { + handle_connection( + stream, + nip11_tx, + nip11_rx, + dropped_tx, + dropped_rx, + ) + .await; + }); + } + _ = &mut shutdown_rx => break, + } + } + }); + + Self { + url: format!("ws://127.0.0.1:{port}"), + shutdown_tx: Some(shutdown_tx), + handle: Some(handle), + } + } + + pub fn url(&self) -> &str { + &self.url + } + + pub async fn stop(mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + if let Some(handle) = self.handle.take() { + let _ = handle.await; + } + } +} + +async fn handle_connection( + mut stream: TcpStream, + nip11_tx: Arc>, + mut nip11_rx: watch::Receiver, + dropped_tx: Arc>, + mut dropped_rx: watch::Receiver, +) { + let mut header = [0_u8; 4096]; + let header_len = loop { + let Ok(length) = stream.peek(&mut header).await else { + return; + }; + if length == 0 { + return; + } + if header[..length] + .windows(4) + .any(|bytes| bytes == b"\r\n\r\n") + { + break length; + } + tokio::task::yield_now().await; + }; + let request = String::from_utf8_lossy(&header[..header_len]).to_ascii_lowercase(); + + if request.contains("upgrade: websocket") { + let Ok(_websocket) = tokio_tungstenite::accept_async(stream).await else { + return; + }; + while !*nip11_rx.borrow() && nip11_rx.changed().await.is_ok() {} + // Dropping the live socket is the behavior under test. Give the SDK a + // bounded propagation window before allowing the HTTP setup to finish. + drop(_websocket); + let _ = dropped_tx.send(true); + } else { + let mut request_bytes = vec![0_u8; header_len]; + if stream.read_exact(&mut request_bytes).await.is_err() { + return; + } + let _ = nip11_tx.send(true); + while !*dropped_rx.borrow() && dropped_rx.changed().await.is_ok() {} + tokio::time::sleep(Duration::from_millis(100)).await; + let body = r#"{"limitation":{"max_subscriptions":20}}"#; + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/nostr+json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ); + let _ = stream.write_all(response.as_bytes()).await; + } +} + +impl Drop for SetupDropRelay { + fn drop(&mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + } +} diff --git a/tests/sync.rs b/tests/sync.rs index 625f628..7949928 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -44,5 +44,6 @@ mod sync { pub mod purgatory_fetch; pub mod reconnect_backoff; pub mod req_concurrency; + pub mod stale_connect_result; pub mod tag_variations; } diff --git a/tests/sync/reconnect_backoff.rs b/tests/sync/reconnect_backoff.rs index 6b0aacf..ed496fc 100644 --- a/tests/sync/reconnect_backoff.rs +++ b/tests/sync/reconnect_backoff.rs @@ -67,7 +67,7 @@ async fn flapping_relay_handshakes_do_not_reset_exponential_backoff() { let logs = wait_for_log( &syncing.log_path(), "consecutive_failures=1", - Duration::from_secs(5), + Duration::from_secs(20), ) .await; assert!(logs.contains("preserving failure streak until stable")); diff --git a/tests/sync/stale_connect_result.rs b/tests/sync/stale_connect_result.rs new file mode 100644 index 0000000..04b10b3 --- /dev/null +++ b/tests/sync/stale_connect_result.rs @@ -0,0 +1,76 @@ +//! A relay that vanishes during NIP-11 setup must not be reported connected. + +use std::time::Duration; + +use nostr_sdk::prelude::*; + +use crate::common::setup_drop_relay::SetupDropRelay; +use crate::common::{TestClient, TestRelay}; + +async fn wait_for_log(log_path: &std::path::Path, needle: &str, timeout: Duration) -> String { + tokio::time::timeout(timeout, async { + loop { + let contents = std::fs::read_to_string(log_path).unwrap_or_default(); + if contents.contains(needle) { + return contents; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap_or_else(|_| panic!("relay log never contained {needle:?}")) +} + +#[tokio::test] +async fn disconnect_during_nip11_fetch_rejects_stale_connection_success() { + let source = SetupDropRelay::start().await; + let syncing = TestRelay::start_with_sync(None).await; + let keys = Keys::generate(); + let identifier = "stale-connect-result"; + let npub = keys.public_key().to_bech32().expect("npub"); + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom( + "clone", + vec![format!( + "http://{}/{npub}/{identifier}.git", + syncing.domain() + )], + ), + Tag::custom( + "relays", + vec![ + format!("ws://{}", syncing.domain()), + source.url().to_string(), + ], + ), + ]) + .finalize(&keys) + .expect("sign announcement"); + let client = TestClient::new(syncing.url(), keys) + .await + .expect("connect publishing client"); + client + .send_event(&announcement) + .await + .expect("publish setup-drop relay announcement"); + + wait_for_log( + &syncing.log_path(), + "Rejecting stale connection success after relay disconnected during setup", + Duration::from_secs(15), + ) + .await; + let logs = wait_for_log( + &syncing.log_path(), + "degraded, backoff 1s", + Duration::from_secs(5), + ) + .await; + assert!(!logs.contains("First connection - initiating fresh_start")); + + client.disconnect().await; + syncing.stop().await; + source.stop().await; +}