mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
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.
This commit is contained in:
@@ -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.
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user