Files
ngit-grasp/tests/relay_query_completion.rs
T
DanConwayDev 99adfcbf66 fix(relay): preserve completed ID query delivery for clients
The SDK sends CLOSED immediately after EOSE for completed ID filters.
That triggers a cancellation race in nak 0.20.6: 39 of 45 production
lookups returned no event although raw queries received every response.

Pin the upstream SDK fixes and disable its automatic completed-ID closure.
Clients retain ordinary CLOSE ownership. The SDK also fixes its aggregate
query cap, so configure that cap as 20 per-filter pages to preserve our
existing documented multi-filter limit instead of truncating history.
Keep shared Nostr types on one Git source and update all flake packages
and the NixOS module with the verified source hash.

This relies on clients closing finished subscriptions; existing connection
and subscription bounds still apply. It does not change the configured
per-filter limit or fix cancellation in nak when using other servers.

Validation: the real-server nak regression failed before and passes 20
consecutive lookups after the fix. Multi-filter overlap/EOSE regression,
32-reader LMDB benchmark and delayed-ACK benchmark pass. Module rendering
preserves the default filter limit; full workspace and Nix builds follow.
2026-09-12 15:56:05 +00:00

88 lines
3.2 KiB
Rust

//! End-to-end query completion and per-filter limit compatibility.
mod common;
use common::{TestClient, TestRelay};
use futures_util::{SinkExt, StreamExt};
use nostr_sdk::prelude::*;
use std::time::Duration;
use tokio_tungstenite::tungstenite::Message;
#[tokio::test]
async fn merged_history_preserves_per_filter_limits_and_id_queries_stay_open() {
let relay = TestRelay::start_with_relay_filter_limit(2).await;
let keys = relay.owner_keys().clone();
let publisher = TestClient::new(relay.url(), keys.clone()).await.unwrap();
let mut events = Vec::new();
for kind in [Kind::TextNote, Kind::Repost] {
for index in 0..2 {
let event = EventBuilder::new(kind, format!("history {index}"))
.finalize(&keys)
.unwrap();
publisher.send_event(&event).await.unwrap();
events.push(event);
}
}
let (mut ws, _) = tokio_tungstenite::connect_async(relay.url()).await.unwrap();
tokio::time::timeout(Duration::from_secs(5), async {
// The aggregate must permit both complete per-filter pages, while
// deduplicating the overlapping ID filter.
for (id, filters, expected) in [
(
"merged",
vec![
serde_json::json!({"kinds":[1]}),
serde_json::json!({"kinds":[6]}),
serde_json::json!({"ids":[events[0].id]}),
],
4,
),
("id", vec![serde_json::json!({"ids":[events[0].id]})], 1),
(
"empty",
vec![serde_json::json!({"ids":["0".repeat(64)]})],
0,
),
] {
let mut request = vec![serde_json::json!("REQ"), serde_json::json!(id)];
request.extend(filters);
ws.send(Message::Text(
serde_json::to_string(&request).unwrap().into(),
))
.await
.unwrap();
let mut received = std::collections::HashSet::new();
loop {
let frame = ws.next().await.unwrap().unwrap();
if !frame.is_text() {
continue;
}
let message: serde_json::Value =
serde_json::from_str(frame.to_text().unwrap()).unwrap();
assert_ne!(
message[0], "CLOSED",
"completed queries must remain open: {message}"
);
assert_eq!(message[1], id);
match message[0].as_str() {
Some("EVENT") => {
received.insert(message[2]["id"].clone().to_string());
}
Some("EOSE") => break,
other => panic!("unexpected response: {other:?}"),
}
}
assert_eq!(received.len(), expected);
ws.send(Message::Text(
serde_json::json!(["CLOSE", id]).to_string().into(),
))
.await
.unwrap();
}
})
.await
.expect("every history response must reach EOSE");
ws.close(None).await.unwrap();
publisher.disconnect().await;
relay.stop().await;
}