From 55e945a5ed1322fa74fc5370ebd82ea55ff513a2 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 1 Jul 2026 14:02:57 +0000 Subject: [PATCH] feat(quartz): idle watchdog for negentropySync instead of a fixed per-round timeout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The fixed `timeoutMs` (default 30s) applied to every reconcile round, but the FIRST round on a large relay is a legitimate long silence while the relay builds its whole negentropy fingerprint — observed at ~63-68s for a 3.5M-event kind:0 set on a real relay. So the old default spuriously failed big syncs with UNAVAILABLE ("reconcile round timed out"), even though the relay was working fine and would have answered seconds later. Replace it with `idleTimeoutMs` (default 120s): the maximum time the relay may go COMPLETELY SILENT before giving up. It resets on every message the relay sends — each NIP-77 round and every download EOSE/event — and on connect, so a genuinely slow but progressing sync runs for as long as it needs; only true silence trips it. Because the watchdog is fed by a connection-level listener that sees all of the relay's traffic, download activity extends the reconcile deadline and vice versa. `idleTimeoutMs = 0` disables it entirely (run until the socket drops); the initial connect and each download batch keep their own finite bounds so an unreachable relay or a single stuck batch still can't hang the pipeline. Liveness of a dead/half-open socket does not depend on this: the WebSocket keep-alive (ping/pong) detects it and the disconnect is already turned into a clean NEG-ERR abort. Verified against wss://wot.grapevine.network: the first round took 68.4s (the old 30s default would have thrown UNAVAILABLE) and the sync sailed through it under the 120s idle window. Negentropy unit suite passes. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01JmSyzdmKyiz3pPxUZ8Mg8Z --- .../NostrClientNegentropySyncEventsExt.kt | 16 +- .../NostrClientNegentropySyncExt.kt | 162 ++++++++++++++---- 2 files changed, 135 insertions(+), 43 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncEventsExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncEventsExt.kt index 37945e1d9c..4c064003a6 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncEventsExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncEventsExt.kt @@ -52,7 +52,7 @@ fun INostrClient.negentropySyncEvents( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, ): Flow = callbackFlow { negentropySync( @@ -61,7 +61,7 @@ fun INostrClient.negentropySyncEvents( maxEvents = maxEvents, maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, ) { event -> trySend(event) } @@ -77,7 +77,7 @@ fun INostrClient.negentropySyncEvents( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, ): Flow = negentropySyncEvents( relay = RelayUrlNormalizer.normalize(relay), @@ -85,7 +85,7 @@ fun INostrClient.negentropySyncEvents( maxEvents = maxEvents, maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, ) /** @@ -104,7 +104,7 @@ fun INostrClient.negentropySyncOrFetchEvents( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, ): Flow = callbackFlow { negentropySyncOrFetch( @@ -113,7 +113,7 @@ fun INostrClient.negentropySyncOrFetchEvents( maxEvents = maxEvents, maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, ) { event -> trySend(event) } @@ -129,7 +129,7 @@ fun INostrClient.negentropySyncOrFetchEvents( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, ): Flow = negentropySyncOrFetchEvents( relay = RelayUrlNormalizer.normalize(relay), @@ -137,5 +137,5 @@ fun INostrClient.negentropySyncOrFetchEvents( maxEvents = maxEvents, maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, ) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt index 2d1b85b0cf..689eeafdb6 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt @@ -42,8 +42,11 @@ import kotlinx.coroutines.flow.first import kotlinx.coroutines.joinAll import kotlinx.coroutines.launch import kotlinx.coroutines.withTimeoutOrNull +import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.ExperimentalAtomicApi import kotlin.coroutines.coroutineContext import kotlin.math.min +import kotlin.time.TimeSource /** * Outcome of a successful [negentropySync] run. @@ -107,12 +110,18 @@ class NegentropySyncResult( * @param maxConcurrentReqs upper bound on simultaneously-open download `REQ`s. Keep * it at or below the relay's per-connection subscription cap. * @param fetchBatch ids per download `REQ`. - * @param timeoutMs max wait for a single reconcile round or download page's `EOSE`. - * The relay builds its whole negentropy snapshot before the FIRST round responds, - * which is O(matched set) — for a multi-million-event filter that first response can - * take a minute or more (a real strfry took ~73s for an unbounded kind:0 set), so - * raise this for very large syncs. It is only a ceiling — download REQs return on - * their `EOSE`, so a generous value costs nothing in the common case. + * @param idleTimeoutMs the idle watchdog: the maximum time the relay may go + * **completely silent** before the sync gives up. It is NOT a per-round deadline — + * it **resets on every message the relay sends** (each NIP-77 round, every download + * `EOSE`/event) and on connect. So a genuinely slow but progressing sync runs for as + * long as it needs: only true silence trips it. This matters because the relay + * builds its whole negentropy snapshot before the FIRST round responds — O(matched + * set), a minute or more for a multi-million-event filter — and that first wait is a + * real silence, so keep this comfortably above the largest expected first-round build. + * A dead/half-open socket does NOT depend on this: the WebSocket keep-alive detects + * it and the disconnect is turned into a clean abort. Pass `0` to disable the + * watchdog entirely and run until the socket drops (download batches keep a finite + * internal idle bound regardless, so a single stuck batch can't hang the pipeline). * @param onProgress optional `(needSoFar, downloaded)` ticks as work proceeds. * @param onEvent called once per distinct event, on the relay reader thread. */ @@ -122,7 +131,7 @@ suspend fun INostrClient.negentropySync( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropySyncResult { @@ -151,7 +160,7 @@ suspend fun INostrClient.negentropySync( syncWindow( relay = relay, filter = filter, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, fetchBatch = fetchBatch, maxConcurrentReqs = maxConcurrentReqs, onWindow = { windows++ }, @@ -195,7 +204,7 @@ suspend fun INostrClient.negentropySync( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropySyncResult = @@ -205,7 +214,7 @@ suspend fun INostrClient.negentropySync( maxEvents = maxEvents, maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, onProgress = onProgress, onEvent = onEvent, ) @@ -248,7 +257,7 @@ suspend fun INostrClient.negentropySyncOrFetch( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropyOrFetchResult { @@ -274,15 +283,17 @@ suspend fun INostrClient.negentropySyncOrFetch( maxEvents = maxEvents, maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, onProgress = onProgress, ) { accept(it) } NegentropyOrFetchResult(delivered, pagedFallback = false, negentropy = result, fallbackCause = null) } catch (e: NegentropySyncException) { // Negentropy couldn't enumerate the set — page the whole filter instead, - // skipping anything the negentropy attempt already delivered. + // skipping anything the negentropy attempt already delivered. fetchAllPages + // has no "no timeout" mode, so a disabled watchdog maps to a finite page bound. val pageFilter = if (maxEvents > 0) filter.copy(limit = maxEvents) else filter - fetchAllPages(relay, listOf(pageFilter), timeoutMs) { event -> + val pageTimeoutMs = if (idleTimeoutMs > 0) idleTimeoutMs else DEFAULT_DOWNLOAD_IDLE_MS + fetchAllPages(relay, listOf(pageFilter), pageTimeoutMs) { event -> if (accept(event)) onProgress?.invoke(delivered, delivered) } NegentropyOrFetchResult(delivered, pagedFallback = true, negentropy = null, fallbackCause = e) @@ -295,7 +306,7 @@ suspend fun INostrClient.negentropySyncOrFetch( maxEvents: Int = 0, maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, - timeoutMs: Long = 30_000L, + idleTimeoutMs: Long = 120_000L, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropyOrFetchResult = @@ -305,7 +316,7 @@ suspend fun INostrClient.negentropySyncOrFetch( maxEvents = maxEvents, maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, - timeoutMs = timeoutMs, + idleTimeoutMs = idleTimeoutMs, onProgress = onProgress, onEvent = onEvent, ) @@ -321,7 +332,7 @@ suspend fun INostrClient.negentropySyncOrFetch( private suspend fun INostrClient.syncWindow( relay: NormalizedRelayUrl, filter: Filter, - timeoutMs: Long, + idleTimeoutMs: Long, fetchBatch: Int, maxConcurrentReqs: Int, onWindow: () -> Unit, @@ -330,7 +341,7 @@ private suspend fun INostrClient.syncWindow( ) { coroutineContext.ensureActive() - when (val outcome = downloadWindow(relay, filter, timeoutMs, fetchBatch, maxConcurrentReqs, onNeed, deliver)) { + when (val outcome = downloadWindow(relay, filter, idleTimeoutMs, fetchBatch, maxConcurrentReqs, onNeed, deliver)) { is ReconcileOutcome.Complete -> onWindow() is ReconcileOutcome.Overflow -> { @@ -347,8 +358,8 @@ private suspend fun INostrClient.syncWindow( ) } else { val mid = lo + (hi - lo) / 2 - syncWindow(relay, filter.copy(since = lo, until = mid), timeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onNeed, deliver) - syncWindow(relay, filter.copy(since = mid + 1, until = hi), timeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onNeed, deliver) + syncWindow(relay, filter.copy(since = lo, until = mid), idleTimeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onNeed, deliver) + syncWindow(relay, filter.copy(since = mid + 1, until = hi), idleTimeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onNeed, deliver) } } @@ -388,7 +399,7 @@ private sealed interface ReconcileOutcome { private suspend fun INostrClient.reconcileStreaming( relay: NormalizedRelayUrl, filter: Filter, - timeoutMs: Long, + idleTimeoutMs: Long, fetchBatch: Int, onNeed: (Int) -> Unit, sendBatch: suspend (List) -> Unit, @@ -402,13 +413,28 @@ private suspend fun INostrClient.reconcileStreaming( // the next one once we ack, and we ack only after this round's ids are queued. val incoming = Channel(Channel.UNLIMITED) + // Idle watchdog. Bumped on connect and on EVERY message this relay sends — + // including the download REQs' events, since this is a connection-level listener + // that sees all of them — so any progress anywhere in the pipeline pushes the + // reconcile deadline out. Only true silence trips it. + val clock = IdleClock() + val listener = object : RelayConnectionListener { + override fun onConnected( + relay: IRelayClient, + pingMillis: Int, + compressed: Boolean, + ) { + if (relay.url == targetUrl) clock.bump() + } + override fun onIncomingMessage( relay: IRelayClient, msgStr: String, msg: Message, ) { + if (relay.url == targetUrl) clock.bump() when (msg) { is NegMsgMessage -> if (msg.subId == subId) incoming.trySend(NegFrame.Msg(msg.message)) is NegErrMessage -> if (msg.subId == subId) incoming.trySend(NegFrame.Err(msg.reason)) @@ -426,22 +452,32 @@ private suspend fun INostrClient.reconcileStreaming( // NEG-OPEN is a one-shot command. Unlike a REQ — which the client replays // from its active-request state every time a relay (re)connects — a dropped // NEG-OPEN is never resent, so we must connect and wait until the relay is - // ready before sending it. + // ready before sending it. The connect itself keeps a finite bound even when + // the watchdog is disabled, so an unreachable relay can't hang here forever. relayClient.connect() + val connectBound = if (idleTimeoutMs > 0) idleTimeoutMs else DEFAULT_CONNECT_TIMEOUT_MS val connected = - withTimeoutOrNull(timeoutMs) { + withTimeoutOrNull(connectBound) { connectedRelaysFlow().first { targetUrl in it } } - if (connected == null) return ReconcileOutcome.Failed("could not connect within ${timeoutMs}ms") + if (connected == null) return ReconcileOutcome.Failed("could not connect within ${connectBound}ms") relayClient.sendIfConnected(session.open()) while (true) { - // Time only the wait for the relay's next frame — never our own - // back-pressured streaming of the previous frame's ids. + // Wait for the relay's next frame, giving up only after idleTimeoutMs of + // total silence (the wait resets whenever the relay sends anything — + // another round, or an event on a download REQ). A disconnect arrives as + // an Err frame, so a dead socket ends this promptly regardless. val frame = - withTimeoutOrNull(timeoutMs) { incoming.receive() } - ?: return ReconcileOutcome.Failed("reconcile round timed out after ${timeoutMs}ms") + incoming.receiveWithinIdle(clock, idleTimeoutMs) + ?: return ReconcileOutcome.Failed( + if (idleTimeoutMs > 0) { + "relay went silent for ${idleTimeoutMs}ms mid-reconcile" + } else { + "connection closed before reconcile completed" + }, + ) when (frame) { is NegFrame.Err -> @@ -508,7 +544,7 @@ private fun isOverflow(reason: String): Boolean = private suspend fun INostrClient.downloadWindow( relay: NormalizedRelayUrl, filter: Filter, - timeoutMs: Long, + idleTimeoutMs: Long, fetchBatch: Int, maxConcurrentReqs: Int, onNeed: (Int) -> Unit, @@ -527,7 +563,7 @@ private suspend fun INostrClient.downloadWindow( launch { for (batch in idBatches) { coroutineContext.ensureActive() - for (event in fetchByIds(relay, batch, timeoutMs)) { + for (event in fetchByIds(relay, batch, idleTimeoutMs)) { deliver(event) } } @@ -535,7 +571,7 @@ private suspend fun INostrClient.downloadWindow( } val outcome = - reconcileStreaming(relay, filter, timeoutMs, fetchBatch, onNeed) { batch -> + reconcileStreaming(relay, filter, idleTimeoutMs, fetchBatch, onNeed) { batch -> idBatches.send(batch) } @@ -558,13 +594,20 @@ private suspend fun INostrClient.downloadWindow( private suspend fun INostrClient.fetchByIds( relay: NormalizedRelayUrl, batch: List, - timeoutMs: Long, + idleTimeoutMs: Long, ): List { val subId = newSubId() val done = Channel(Channel.CONFLATED) val collected = ArrayList(batch.size) val seen = HashSet(batch.size) + // Per-batch idle clock: each event resets it, so a batch that keeps streaming is + // never cut off, but a batch that stalls (relay stops mid-flight) unblocks after + // the idle bound instead of hanging a worker. A download batch always keeps a + // finite bound even when the caller disabled the whole-sync watchdog. + val clock = IdleClock() + val batchIdleMs = if (idleTimeoutMs > 0) idleTimeoutMs else DEFAULT_DOWNLOAD_IDLE_MS + val listener = object : SubscriptionListener { override fun onEvent( @@ -573,6 +616,7 @@ private suspend fun INostrClient.fetchByIds( relay: NormalizedRelayUrl, forFilters: List?, ) { + clock.bump() if (seen.add(event.id)) collected.add(event) } @@ -602,9 +646,7 @@ private suspend fun INostrClient.fetchByIds( try { subscribe(subId, mapOf(relay to listOf(Filter(ids = batch))), listener) - withTimeoutOrNull(timeoutMs) { - done.receive() - } + done.receiveWithinIdle(clock, batchIdleMs) } finally { unsubscribe(subId) done.close() @@ -625,3 +667,53 @@ private const val DELIVERY_BUFFER = 256 * collides with an actual event. */ private val KEEP_ALIVE_ID = "f".repeat(64) + +/** + * Finite fallback bounds (ms) for the two waits that must stay bounded even when the + * whole-sync idle watchdog is disabled (`idleTimeoutMs = 0`): the initial connect, + * and each individual download batch. Keeping these finite means an unreachable relay + * or a single stuck batch can never hang the pipeline, while the reconcile rounds + * still honor "run until the socket drops". + */ +private const val DEFAULT_CONNECT_TIMEOUT_MS = 30_000L +private const val DEFAULT_DOWNLOAD_IDLE_MS = 60_000L + +/** + * Monotonic "last activity" marker for the idle watchdog. [bump] on every sign of + * life from the relay; [elapsedMs] reports the silence since the last bump. Thread + * safe: bumped from relay reader threads, read from the driver coroutine. + */ +@OptIn(ExperimentalAtomicApi::class) +private class IdleClock { + private val last = AtomicReference(TimeSource.Monotonic.markNow()) + + fun bump() { + last.store(TimeSource.Monotonic.markNow()) + } + + fun elapsedMs(): Long = last.load().elapsedNow().inWholeMilliseconds +} + +/** + * Receives the next item, giving up (returning `null`) only after [idleMs] elapse with + * no activity on [clock]. Because [clock] is bumped by *any* relay message — not just + * items on this channel — unrelated progress (e.g. download events arriving during a + * reconcile wait) keeps pushing the deadline out. [idleMs] `<= 0` disables the + * watchdog: it waits until an item arrives (a disconnect is delivered as an item, so + * a dead socket still unblocks it). + */ +private suspend fun Channel.receiveWithinIdle( + clock: IdleClock, + idleMs: Long, +): T? { + if (idleMs <= 0) return receive() + while (true) { + val remaining = idleMs - clock.elapsedMs() + if (remaining <= 0) return null + val item = withTimeoutOrNull(remaining) { receive() } + if (item != null) return item + // Timed out with nothing on this channel. If other activity bumped the clock + // meanwhile, the next `remaining` is positive and we wait again; otherwise it + // is <= 0 on the next iteration and we give up. + } +}