mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
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.
This commit is contained in:
@@ -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.
|
**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
|
```mermaid
|
||||||
flowchart LR
|
flowchart LR
|
||||||
CONN[Connection Success] --> SPAWN[handle_connect_or_reconnect<br/>spawns event loop]
|
CONN[Connection Success] --> SPAWN[handle_connect_or_reconnect<br/>spawns event loop]
|
||||||
|
|||||||
+23
-1
@@ -3663,7 +3663,29 @@ impl SyncManager {
|
|||||||
return;
|
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 {
|
ConnectAttemptOutcome::Connected {
|
||||||
advertised_default_limit,
|
advertised_default_limit,
|
||||||
advertised_max_subscriptions,
|
advertised_max_subscriptions,
|
||||||
|
|||||||
@@ -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.
|
/// Configure the one per-session ledger before any subscriptions open.
|
||||||
pub fn reset_subscription_budget(&self, advertised: Option<usize>) {
|
pub fn reset_subscription_budget(&self, advertised: Option<usize>) {
|
||||||
self.clear_subscription_permits();
|
self.clear_subscription_permits();
|
||||||
|
|||||||
@@ -51,7 +51,11 @@ impl FlappingRelay {
|
|||||||
let Ok(message) = message else { break };
|
let Ok(message) = message else { break };
|
||||||
if matches!(message, Message::Text(ref text) if text.contains("\"REQ\"")) {
|
if matches!(message, Message::Text(ref text) if text.contains("\"REQ\"")) {
|
||||||
// Dropping the socket after observable normal
|
// 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;
|
return;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ pub mod upload_pack_counting_proxy;
|
|||||||
pub mod port;
|
pub mod port;
|
||||||
pub mod purgatory_helpers;
|
pub mod purgatory_helpers;
|
||||||
pub mod relay;
|
pub mod relay;
|
||||||
|
pub mod setup_drop_relay;
|
||||||
pub mod sync_helpers;
|
pub mod sync_helpers;
|
||||||
|
|
||||||
pub use git_server::{SimpleGitServer, SmartGitServer};
|
pub use git_server::{SimpleGitServer, SmartGitServer};
|
||||||
|
|||||||
@@ -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<oneshot::Sender<()>>,
|
||||||
|
handle: Option<tokio::task::JoinHandle<()>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<watch::Sender<bool>>,
|
||||||
|
mut nip11_rx: watch::Receiver<bool>,
|
||||||
|
dropped_tx: Arc<watch::Sender<bool>>,
|
||||||
|
mut dropped_rx: watch::Receiver<bool>,
|
||||||
|
) {
|
||||||
|
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(());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -44,5 +44,6 @@ mod sync {
|
|||||||
pub mod purgatory_fetch;
|
pub mod purgatory_fetch;
|
||||||
pub mod reconnect_backoff;
|
pub mod reconnect_backoff;
|
||||||
pub mod req_concurrency;
|
pub mod req_concurrency;
|
||||||
|
pub mod stale_connect_result;
|
||||||
pub mod tag_variations;
|
pub mod tag_variations;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -67,7 +67,7 @@ async fn flapping_relay_handshakes_do_not_reset_exponential_backoff() {
|
|||||||
let logs = wait_for_log(
|
let logs = wait_for_log(
|
||||||
&syncing.log_path(),
|
&syncing.log_path(),
|
||||||
"consecutive_failures=1",
|
"consecutive_failures=1",
|
||||||
Duration::from_secs(5),
|
Duration::from_secs(20),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
assert!(logs.contains("preserving failure streak until stable"));
|
assert!(logs.contains("preserving failure streak until stable"));
|
||||||
|
|||||||
@@ -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;
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user