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();