Merge #1a76cd3d: Fix relay retry churn and avoidable response stalls

nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqsp5akd8h6qc7k40glf0a8d9wuw7qrw5uljn2wa7k8s5w2v04ykewg6jm5uw

PR-Author: DanConwayDev's Agent
nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0

CoverNote:

Small HTTP/WebSocket responses could wait for delayed acknowledgements, while checkpoint writes and metrics rendering performed blocking work on async workers. Background discovery also lost retry history during cleanup and could accept incomplete fetches as successful history.

Review the five commits independently:

| Commit | Scope | Change |
| --- | --- | --- |
| `ddcfa3c` | User responses | Enable TCP_NODELAY on accepted sockets; include delayed-ACK and concurrent LMDB read benchmarks. |
| `fe36d6f` | Shared runtime | Move periodic checkpoints to a blocking worker, release snapshot locks before I/O, and join active writes before the final shutdown snapshot. |
| `f55ce48` | Shared runtime | Render metrics on a blocking worker; retain a shared permit through completion so canceled scrapes cannot start overlapping scans. |
| `5d62037` | Background sync | Preserve discovery ownership and failure history through cleanup and reconnection without adding persistent subscriptions. |
| `df9280e` | Background sync | Require the exact subscription’s EOSE and a drained event stream before accepting discovered history. |

Each commit includes its tests, architecture documentation and changelog entry. Dependency versions and inbound relay query behavior match the base; the transport change applies to accepted connections. The two sync fixes affect outbound discovery.

Validation on the rewritten tree:

- Full workspace tests: 3,034 passed, 0 failed, 16 ignored.
- Formatting and Clippy with warnings denied: passed.
- Nix package build and its library checks: passed.
- Both opt-in response benchmarks passed. Twenty-five two-event reads with delayed ACK took 8.65 ms total. Concurrent LMDB reads reached EOSE at 1, 4, 16 and 32 readers; the maximum per-reader time for three 32-event batches at 32 readers was 220.87 ms.

The five retained fixes match the previously reviewed implementation. The withdrawn SDK workaround and pin are excluded. Investigation reports are absent from the final tree; the original history is preserved locally on `archive/relay-timeout-performance-2026-09-14`.

Local latency and load measurements are diagnostic samples, not production guarantees. These changes do not establish that every historical seven-second timeout had the same cause.
This commit is contained in:
DanConwayDev
2026-09-14 11:37:20 +01:00
13 changed files with 615 additions and 17 deletions
+17
View File
@@ -7,6 +7,23 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased]
### 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.
- Render metrics on a blocking worker and serialize scrapes so repository
counting does not occupy async workers serving relay and Git requests.
- Keep periodic recovery checkpoints off async workers, release snapshot locks
before disk I/O, and finish active checkpoints before the shutdown snapshot.
- Avoid delayed-acknowledgement stalls in small WebSocket and Git responses by
enabling `TCP_NODELAY` on accepted connections.
## [3.0.2] - 2026-09-11
This release contains no production runtime changes. It improves release and
+18 -1
View File
@@ -75,6 +75,11 @@ runtime:
- Bind the TCP listener (supports `:0` for kernel-assigned ports; the
resolved address back-fills `bind_address` and, if empty, `domain`)
- Disable Nagle on accepted sockets so small EVENT/EOSE and streaming Git
writes do not wait for the peer's delayed acknowledgements, including
when the immediate peer is a local reverse proxy
- Render metrics on a blocking worker, serializing concurrent scrapes before
spawning filesystem scans so monitoring cannot occupy the network executor
- Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51)
- Set up shared storage (LMDB or Memory), purgatory, sync manager, and
background maintenance tasks
@@ -89,7 +94,11 @@ runtime:
Transient failures retry with a capped backoff, and in private mode the
identity stays local so the relay's existence is never advertised
- Atomically checkpoint purgatory and rejected-event recovery state every 60
seconds without consuming the checkpoint during restore
seconds without consuming the checkpoint during restore. Serialization and
durable filesystem writes run on a blocking worker, with at most one
checkpoint active and the next interval measured from completion. Shutdown
joins any active checkpoint before writing the final snapshot, preventing
an older background write from replacing it
- Start non-blocking storage-integrity and authorization-integrity passes after
database initialization. The former checks family objects and thin-view
wiring; the latter reconciles each served ref against accepted State, PR,
@@ -735,6 +744,14 @@ is deliberately historical-only: owning or maintaining a repository never adds
that author's inbox relays to ordinary persistent live repository targets.
Private instances omit repository-coordinate mailbox expansion.
Discovery source and mailbox scope also own connection retry state. The empty
relay checker and reconnect scheduler include that scope, including in-flight
fetches and future probe deadlines, without adding persistent subscriptions.
Healthy idle discovery sessions can retire after their fetch; failed or
policy-limited sessions retain their backoff until recovery or scope removal.
This prevents cleanup from erasing a failure every two seconds and allowing
discovery to dial the same unavailable endpoint from a fresh retry history.
### Rejected Events Index
The rejected events index solves two critical problems during sync:
@@ -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:
+110
View File
@@ -0,0 +1,110 @@
//! Serialized blocking checkpoints with an explicit shutdown barrier.
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::watch;
use tokio::task::JoinHandle;
pub(crate) struct CheckpointTask {
stop: watch::Sender<bool>,
task: JoinHandle<()>,
}
impl CheckpointTask {
pub(crate) fn start(interval: Duration, save: impl Fn() + Send + Sync + 'static) -> Self {
let (stop, mut stopped) = watch::channel(false);
let save = Arc::new(save);
let task = tokio::spawn(async move {
loop {
// Schedule from completion: slow storage must not accumulate
// overdue checkpoints or start overlapping snapshots.
tokio::select! {
biased;
_ = stopped.changed() => break,
_ = tokio::time::sleep(interval) => {}
}
let save = save.clone();
// Do not race shutdown against this join. Blocking tasks
// cannot be aborted after they start, and an old snapshot
// must never overwrite the final shutdown snapshot.
if let Err(error) = tokio::task::spawn_blocking(move || save()).await {
tracing::error!(%error, "Checkpoint worker failed");
}
}
});
Self { stop, task }
}
pub(crate) async fn shutdown(self) {
let _ = self.stop.send(true);
let _ = self.task.await;
}
}
#[cfg(test)]
mod tests {
use std::sync::{mpsc, Mutex};
use futures_util::FutureExt;
use tokio::sync::oneshot;
use super::*;
#[tokio::test]
async fn blocking_checkpoint_yields_runtime_and_shutdown_joins_it() {
let (started_tx, started_rx) = oneshot::channel();
let started_tx = Mutex::new(Some(started_tx));
let (release_tx, release_rx) = mpsc::channel();
let release_rx = Mutex::new(release_rx);
let (finished_tx, finished_rx) = oneshot::channel();
let finished_tx = Mutex::new(Some(finished_tx));
let checkpoint = CheckpointTask::start(Duration::from_millis(1), move || {
started_tx.lock().unwrap().take().unwrap().send(()).unwrap();
// A bounded blocking I/O surrogate. The sole runtime thread must
// remain free to send release; an inline checkpoint fails here.
release_rx
.lock()
.unwrap()
.recv_timeout(Duration::from_secs(5))
.unwrap();
finished_tx
.lock()
.unwrap()
.take()
.unwrap()
.send(())
.unwrap();
});
tokio::time::timeout(Duration::from_secs(5), started_rx)
.await
.unwrap()
.unwrap();
let shutdown = checkpoint.shutdown();
tokio::pin!(shutdown);
assert!(
shutdown.as_mut().now_or_never().is_none(),
"shutdown must join the active snapshot"
);
release_tx.send(()).unwrap();
tokio::time::timeout(Duration::from_secs(5), shutdown)
.await
.unwrap();
assert!(finished_rx.now_or_never().unwrap().is_ok());
}
#[tokio::test(start_paused = true)]
async fn shutdown_cancels_a_checkpoint_that_has_not_started() {
use std::sync::atomic::{AtomicUsize, Ordering};
let saves = Arc::new(AtomicUsize::new(0));
let worker_saves = saves.clone();
let checkpoint = CheckpointTask::start(Duration::from_secs(60), move || {
worker_saves.fetch_add(1, Ordering::SeqCst);
});
checkpoint.shutdown().await;
tokio::time::advance(Duration::from_secs(120)).await;
assert_eq!(saves.load(Ordering::SeqCst), 0);
}
}
+19 -1
View File
@@ -923,7 +923,18 @@ impl Service<Request<Incoming>> for HttpService {
if let Some(ref metrics) = self.metrics {
let metrics = metrics.clone();
return Box::pin(async move {
let output = metrics.render();
let output = match metrics.render_async().await {
Ok(output) => output,
Err(error) => {
tracing::error!(%error, "Metrics render worker failed");
return Ok(add_cors_headers(
Response::builder().header("server", "ngit-grasp"),
)
.status(500)
.body(full_body("Metrics temporarily unavailable"))
.unwrap());
}
};
Ok(
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
.status(200)
@@ -1073,6 +1084,13 @@ pub async fn run_server_on_listener(
loop {
let (socket, addr) = listener.accept().await?;
// WebSocket EVENT/EOSE and streaming Git responses flush small writes.
// Nagle can hold the tail until the peer's delayed ACK (about 40 ms
// even through a loopback proxy), multiplying multi-query latency.
if let Err(error) = socket.set_nodelay(true) {
tracing::warn!(peer = %addr, %error, "Could not disable Nagle for client socket");
continue;
}
let io = TokioIo::new(socket);
let service = HttpService::new(
relay.clone(),
+1
View File
@@ -1,5 +1,6 @@
mod atomic_file;
pub mod audit_cleanup;
mod checkpoint;
pub mod cleanup_empty_repos;
pub mod config;
pub mod git;
+43 -2
View File
@@ -564,6 +564,7 @@ pub fn record_manual_ejection(
#[derive(Clone)]
pub struct Metrics {
inner: Arc<MetricsInner>,
render_permit: Arc<tokio::sync::Semaphore>,
}
struct MetricsInner {
@@ -628,6 +629,7 @@ impl Metrics {
let inner = MetricsInner::new(abuse_threshold, git_data_path);
Self {
inner: Arc::new(inner),
render_permit: Arc::new(tokio::sync::Semaphore::new(1)),
}
}
@@ -794,6 +796,24 @@ impl Metrics {
// === Rendering ===
/// Render a scrape without blocking the async executor on filesystem
/// traversal or encoding. Concurrent scrapes wait before spawning work;
/// cancelling an HTTP request cannot release a running worker's permit.
pub async fn render_async(&self) -> Result<String, tokio::task::JoinError> {
let permit = self
.render_permit
.clone()
.acquire_owned()
.await
.expect("metrics render semaphore is never closed");
let metrics = self.clone();
tokio::task::spawn_blocking(move || {
let _permit = permit;
metrics.render()
})
.await
}
/// Render all metrics in Prometheus text format.
///
/// This method:
@@ -1104,6 +1124,21 @@ impl Drop for GitOperationTimer {
mod tests {
use super::*;
async fn concurrent_scrapes_wait_without_blocking_the_runtime(metrics: &Metrics) {
use futures_util::FutureExt;
let permit = metrics.render_permit.clone().acquire_owned().await.unwrap();
let render = metrics.render_async();
tokio::pin!(render);
assert!(render.as_mut().now_or_never().is_none());
drop(permit);
let output = tokio::time::timeout(std::time::Duration::from_secs(5), render)
.await
.unwrap()
.unwrap();
assert!(output.contains("ngit_uptime_seconds"));
assert_eq!(metrics.render_permit.available_permits(), 1);
}
#[test]
fn test_count_repositories_on_disk() {
use std::fs;
@@ -1153,8 +1188,8 @@ mod tests {
///
/// If additional Metrics tests are needed, they should either be added to this
/// test or use a separate test-specific Prometheus registry.
#[test]
fn test_metrics_with_repository_counting() {
#[tokio::test]
async fn test_metrics_with_repository_counting() {
use std::fs;
use tempfile::TempDir;
@@ -1211,5 +1246,11 @@ mod tests {
// Render should count 3 repos
let output = metrics.render();
assert!(output.contains("ngit_repositories_total 3"));
concurrent_scrapes_wait_without_blocking_the_runtime(&metrics).await;
assert!(metrics
.render_async()
.await
.unwrap()
.contains("ngit_repositories_total 3"));
}
}
+6 -7
View File
@@ -66,6 +66,7 @@ pub struct RelayServer {
/// cleanup, purgatory sync loop). Aborted on shutdown so the host
/// process does not leak tasks per relay instance.
background_tasks: Vec<JoinHandle<()>>,
checkpoint: crate::checkpoint::CheckpointTask,
git_data_path: String,
}
@@ -369,11 +370,8 @@ impl RelayServer {
let checkpoint_purgatory = purgatory.clone();
let checkpoint_rejected = rejected_events_index.clone();
let checkpoint_root = PathBuf::from(config.effective_git_data_path());
background_tasks.push(tokio::spawn(async move {
let first = tokio::time::Instant::now() + SYNC_STATE_CHECKPOINT_INTERVAL;
let mut interval = tokio::time::interval_at(first, SYNC_STATE_CHECKPOINT_INTERVAL);
loop {
interval.tick().await;
let checkpoint =
crate::checkpoint::CheckpointTask::start(SYNC_STATE_CHECKPOINT_INTERVAL, move || {
let purgatory_path = checkpoint_root.join("purgatory-state.json");
if let Err(error) = checkpoint_purgatory.save_to_disk(&purgatory_path) {
warn!(%error, "Failed to checkpoint purgatory state");
@@ -382,8 +380,7 @@ impl RelayServer {
if let Err(error) = checkpoint_rejected.save_to_disk(&rejected_path) {
warn!(%error, "Failed to checkpoint rejected-events cache");
}
}
}));
});
info!(
interval_secs = SYNC_STATE_CHECKPOINT_INTERVAL.as_secs(),
"Crash-safe sync-state checkpoint task started"
@@ -498,6 +495,7 @@ impl RelayServer {
private_access,
deletion_cleanup,
background_tasks,
checkpoint,
git_data_path,
})
}
@@ -553,6 +551,7 @@ impl RelayServer {
task.abort();
let _ = task.await;
}
self.checkpoint.shutdown().await;
self.deletion_cleanup.shutdown().await;
+70 -2
View File
@@ -1960,6 +1960,20 @@ struct Nip65DiscoveryState {
}
impl Nip65DiscoveryState {
/// Connection ownership includes historical discovery work, even while
/// its next query is deferred. It is not persistent subscription scope.
fn connection_targets(&self) -> HashSet<String> {
self.author_sources
.values()
.flatten()
.chain(self.mailbox_roots.keys())
.chain(self.mailbox_repositories.keys())
.chain(self.mailbox_probes_in_flight.iter())
.chain(self.in_flight.iter().map(|(relay, _)| relay))
.cloned()
.collect()
}
fn has_mailbox_scope(&self, relay: &str) -> bool {
self.mailbox_roots.contains_key(relay) || self.mailbox_repositories.contains_key(relay)
}
@@ -6422,6 +6436,22 @@ impl SyncManager {
if !self.nip65_discovery_only_relays.contains(source) {
return;
}
// Retirement forgets connection health. A failed or policy-limited
// source must retain its session state until recovery or removal from
// discovery ownership, otherwise the next probe bypasses its backoff.
if self.health_tracker.get_failure_count(source) > 0
|| self.health_tracker.is_subscription_paused(source)
{
return;
}
if let Some(connection) = self.connections.get(source) {
if !connection.is_connected().await {
// The disconnect notification may still be queued behind a
// fetch result. Let it record the unexpected failure instead
// of relabeling the session as an intentional retirement.
return;
}
}
let now = Instant::now();
let has_author_work =
self.nip65_discovery
@@ -7221,7 +7251,12 @@ impl SyncManager {
// Once the connection ends naturally, however, reconcile its
// confirmed state with the latest index before deciding whether it
// should ever reconnect.
let desired = self.derive_targets().await.remove(relay_url);
let desired = self.derive_targets().await.remove(relay_url).or_else(|| {
self.nip65_discovery
.connection_targets()
.contains(relay_url)
.then(RelaySyncNeeds::default)
});
let Some(desired) = desired else {
tracing::info!(
relay = %relay_url,
@@ -8420,6 +8455,7 @@ impl SyncManager {
let mut desired_relays: HashSet<String> = self.derive_targets().await.into_keys().collect();
desired_relays.extend(self.dependency_relay_deadlines.keys().cloned());
desired_relays.extend(self.nip65_discovery.connection_targets());
// Collect relays to disconnect
let to_disconnect: Vec<String> = {
@@ -8531,7 +8567,8 @@ impl SyncManager {
///
/// For each eligible relay, a reconnection is queued via schedule_connect_relay.
async fn retry_disconnected_relays(&mut self) {
let desired_relays: HashSet<String> = self.derive_targets().await.into_keys().collect();
let mut desired_relays: HashSet<String> = self.derive_targets().await.into_keys().collect();
desired_relays.extend(self.nip65_discovery.connection_targets());
// Collect relays to reconnect
let to_reconnect: Vec<String> = {
@@ -10971,6 +11008,37 @@ mod tests {
);
}
#[test]
fn discovery_connection_ownership_survives_retry_deadlines_and_in_flight_removal() {
let author = Keys::generate().public_key();
let profile = "wss://profile.example".to_string();
let mailbox = "wss://mailbox.example".to_string();
let retired = "wss://retired.example".to_string();
let mut discovery = Nip65DiscoveryState::default();
discovery
.author_sources
.insert(author, HashSet::from([profile.clone()]));
discovery
.mailbox_repositories
.insert(mailbox.clone(), HashSet::from(["repo".into()]));
discovery
.mailbox_probe_next_at
.insert(mailbox.clone(), Instant::now() + Duration::from_secs(3600));
discovery.in_flight.insert((retired.clone(), author));
let disconnected = RelayState::default();
for relay in [&profile, &mailbox, &retired] {
assert!(!disconnected
.is_disconnect_candidate(false, discovery.connection_targets().contains(relay)));
}
discovery.in_flight.clear();
assert!(!discovery.connection_targets().contains(&retired));
discovery.author_sources.clear();
discovery.install_mailbox_overlay(HashMap::new(), HashMap::new(), Instant::now());
assert!(discovery.connection_targets().is_empty());
assert!(disconnected.is_disconnect_candidate(false, false));
}
#[test]
fn disconnected_empty_relay_can_still_be_cleaned_up() {
let mut source = RelayState::default();
+9
View File
@@ -1357,6 +1357,15 @@ impl RejectedEventsIndex {
related_dependencies: serializable_related_entries,
};
// The snapshot owns its data. Ingestion must not wait on JSON
// encoding, fsync or rename while holding these synchronous locks.
drop((
hot_entries,
cold_entries,
unrecoverable_entries,
related_entries,
));
// Replace the previous checkpoint only after the new snapshot is
// complete and durable. An abrupt stop must leave one valid version.
let json = serde_json::to_string_pretty(&state)?;
+129 -4
View File
@@ -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))
}
@@ -3113,6 +3166,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 = TestRelay::start(LocalRelayBuilder::default()).await;
+146
View File
@@ -0,0 +1,146 @@
//! Opt-in Linux delayed-ACK benchmark for the accepted HTTP/WebSocket socket.
mod common;
#[tokio::test]
#[ignore = "local load benchmark; run explicitly with --nocapture"]
async fn concurrent_lmdb_read_batches_reach_eose() {
use common::{TestClient, TestRelay};
use futures_util::{SinkExt, StreamExt};
use nostr_sdk::prelude::*;
use std::time::{Duration, Instant};
use tokio_tungstenite::tungstenite::Message;
let relay = TestRelay::start_with_lmdb().await;
let keys = relay.owner_keys().clone();
let client = TestClient::new(relay.url(), keys.clone()).await.unwrap();
for index in 0..32 {
let event = EventBuilder::new(Kind::TextNote, format!("{index}:{}", "x".repeat(8192)))
.finalize(&keys)
.unwrap();
client.send_event(&event).await.unwrap();
}
for concurrency in [1, 4, 16, 32] {
let url = relay.url().to_string();
let reads = (0..concurrency).map(|_| {
let url = url.clone();
async move {
let (mut stream, _) = tokio_tungstenite::connect_async(url).await.unwrap();
let started = Instant::now();
for _ in 0..3 {
stream
.send(Message::Text(r#"["REQ","load",{"kinds":[1]}]"#.into()))
.await
.unwrap();
let mut events = 0;
loop {
let frame = stream.next().await.unwrap().unwrap();
if !frame.is_text() {
continue;
}
let message: serde_json::Value =
serde_json::from_str(frame.to_text().unwrap()).unwrap();
match message[0].as_str() {
Some("EVENT") => events += 1,
Some("EOSE") => break,
other => panic!("unexpected reply: {other:?}"),
}
}
assert_eq!(events, 32);
}
let elapsed = started.elapsed();
stream.close(None).await.unwrap();
elapsed
}
});
let times = tokio::time::timeout(
Duration::from_secs(30),
futures_util::future::join_all(reads),
)
.await
.expect("all readers must reach EOSE under local load");
eprintln!(
"{concurrency} concurrent LMDB readers, three 32-event batches each: max {:?}",
times.iter().max().unwrap()
);
}
client.disconnect().await;
relay.stop().await;
}
#[cfg(target_os = "linux")]
#[tokio::test]
#[ignore = "socket latency benchmark; run explicitly on an otherwise idle worker"]
async fn small_event_batches_do_not_wait_for_delayed_ack() {
use std::os::fd::AsRawFd;
use std::time::{Duration, Instant};
use common::{TestClient, TestRelay};
use futures_util::{SinkExt, StreamExt};
use nostr_sdk::prelude::*;
use tokio_tungstenite::tungstenite::Message;
let relay = TestRelay::start().await;
let keys = relay.owner_keys().clone();
let client = TestClient::new(relay.url(), keys.clone()).await.unwrap();
for index in 0..2 {
let event = EventBuilder::new(Kind::TextNote, format!("event {index}"))
.finalize(&keys)
.unwrap();
client.send_event(&event).await.unwrap();
}
let socket = tokio::net::TcpStream::connect(relay.domain())
.await
.unwrap();
socket.set_nodelay(true).unwrap();
let (mut stream, _) = tokio_tungstenite::client_async(relay.url(), socket)
.await
.unwrap();
let started = Instant::now();
tokio::time::timeout(Duration::from_secs(10), async {
for _ in 0..25 {
let disabled: libc::c_int = 0;
// Explicitly model a peer that uses delayed acknowledgements.
// The descriptor and option value remain alive for setsockopt.
let result = unsafe {
libc::setsockopt(
stream.get_ref().as_raw_fd(),
libc::IPPROTO_TCP,
libc::TCP_QUICKACK,
&disabled as *const _ as *const libc::c_void,
std::mem::size_of_val(&disabled) as libc::socklen_t,
)
};
assert_eq!(result, 0);
stream
.send(Message::Text(r#"["REQ","latency",{"kinds":[1]}]"#.into()))
.await
.unwrap();
let mut events = 0;
loop {
let frame = stream.next().await.unwrap().unwrap();
if !frame.is_text() {
continue;
}
let message: serde_json::Value =
serde_json::from_str(frame.to_text().unwrap()).unwrap();
match message[0].as_str() {
Some("EVENT") => events += 1,
Some("EOSE") => break,
other => panic!("unexpected reply: {other:?}"),
}
}
assert_eq!(events, 2);
}
})
.await
.expect("batch reads must complete");
let elapsed = started.elapsed();
eprintln!("25 two-event reads with delayed ACK: {elapsed:?}");
assert!(
elapsed < Duration::from_millis(500),
"small batches were delayed: {elapsed:?}"
);
client.disconnect().await;
relay.stop().await;
}
+41
View File
@@ -7,6 +7,47 @@ use nostr_sdk::prelude::*;
use crate::common::flapping_relay::FlappingRelay;
use crate::common::{TestClient, TestRelay};
#[tokio::test]
async fn discovery_only_mailbox_keeps_failure_history_through_cleanup() {
use crate::common::{send_to_relay_url, setup_announcement_on_relay, MockRelay};
let index = MockRelay::start().await;
let flapping = FlappingRelay::start().await;
let owner = Keys::generate();
let relay_list = EventBuilder::new(Kind::RelayList, "")
.tags([Tag::custom("r", vec![flapping.url(), "read"])])
.finalize(&owner)
.unwrap();
send_to_relay_url(index.url(), &relay_list).await.unwrap();
let syncing = TestRelay::start_with_sync(Some(index.url().to_string())).await;
let domain = syncing.domain();
let (_announcement, _git) =
setup_announcement_on_relay(&syncing, &owner, &[&domain], "discovery-mailbox-backoff")
.await;
// Owner inbox discovery has no ordinary repository/root live target.
// The two-second cleanup pass previously erased its first failure before
// the next handshake, so it could never reach this recovery state.
flapping
.wait_for_connections(2, Duration::from_secs(60))
.await;
let logs = wait_for_log(
&syncing.log_path(),
"consecutive_failures=1",
Duration::from_secs(15),
)
.await;
assert!(logs.lines().any(|line| {
line.contains(flapping.url())
&& line.contains("consecutive_failures=1")
&& line.contains("preserving failure streak until stable")
}));
syncing.stop().await;
flapping.stop().await;
index.stop().await;
}
async fn wait_for_log(log_path: &std::path::Path, needle: &str, timeout: Duration) -> String {
tokio::time::timeout(timeout, async {
loop {