From c2cabf3c4710490d08358380dfa4b0bb96a78d74 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 3 Jul 2026 22:26:15 +0000 Subject: [PATCH] =?UTF-8?q?feat(geode):=20mirror=20survives=20upstream=20r?= =?UTF-8?q?estarts=20=E2=80=94=20retry=20pump,=20ping=20keepalive,=20e2e?= =?UTF-8?q?=20proof?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Measured over the real OkHttp transport (upstream killed mid-mirror and restarted on the same port): NostrClient re-dials once on disconnect and then falls back to its 60s keep-alive, which showed up as a 61s mirror blackout when the immediate re-dial raced the port rebind. - Retry pump: the worker nudges the pool every 5s. Cheap and safe — reconnectIfNeedsTo skips connected relays and each relay's exponential backoff (1s doubling, 5min cap) still gates real dial attempts, so dead upstreams aren't hammered. Restart recovery drops from 61s to ~6s in the new reconnect test. - WebSocket pings (120s, same as the Android relay pool): without them a half-open connection (network drop, no FIN) never fires onDisconnected and the mirror would stall silently forever. - MirrorWorkerReconnectTest: end-to-end over a real Ktor port — ride through the drop, re-subscribe, drop the duplicate replay, pull the event that only exists on the new instance. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A --- .../geode/mirror/MirrorWorker.kt | 44 ++++- .../geode/mirror/MirrorWorkerReconnectTest.kt | 157 ++++++++++++++++++ 2 files changed, 200 insertions(+), 1 deletion(-) create mode 100644 geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt index 2350bf72c5..d12fffc496 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt @@ -37,8 +37,10 @@ import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.delay import kotlinx.coroutines.launch import okhttp3.OkHttpClient +import java.time.Duration import java.util.concurrent.atomic.AtomicLong /** @@ -102,7 +104,24 @@ class MirrorWorker( ) : AutoCloseable { private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) - private val okhttp: OkHttpClient? = if (websocketBuilder == null) OkHttpClient.Builder().build() else null + /** + * The ping interval matters on a long-running daemon: a half-open + * upstream connection (network drop with no FIN) never fires + * onDisconnected on its own, so without pings the client would + * believe it is connected forever and stop mirroring silently. A + * missed pong fails the socket, which routes into the normal + * disconnect → backoff → re-dial path. Same value the Android app + * uses for its relay pool. + */ + private val okhttp: OkHttpClient? = + if (websocketBuilder == null) { + OkHttpClient + .Builder() + .pingInterval(Duration.ofSeconds(PING_INTERVAL_SECS)) + .build() + } else { + null + } private val client = NostrClient( @@ -204,6 +223,21 @@ class MirrorWorker( ) } client.connect() + + // Retry pump. NostrClient re-dials once on disconnect and then + // relies on its 60s keep-alive — measured as a 61s mirror blackout + // when an upstream restarts and the immediate re-dial races the + // port rebind. Poking more often costs nothing: reconnectIfNeedsTo + // skips connected relays, and each relay's exponential backoff + // (1s doubling, 5min cap) still gates actual dial attempts, so a + // long-dead upstream is not hammered — only a briefly-restarting + // one is picked back up in seconds instead of a minute. + scope.launch { + while (true) { + delay(RECONNECT_POKE_MS) + client.reconnect(onlyIfChanged = true, ignoreRetryDelays = false) + } + } } /** @@ -222,4 +256,12 @@ class MirrorWorker( okhttp?.dispatcher?.executorService?.shutdown() okhttp?.connectionPool?.evictAll() } + + private companion object { + /** Matches the Android app's relay-pool WebSocket ping interval. */ + const val PING_INTERVAL_SECS = 120L + + /** How often the retry pump nudges disconnected upstreams. */ + const val RECONNECT_POKE_MS = 5_000L + } } diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt new file mode 100644 index 0000000000..9268ca0384 --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt @@ -0,0 +1,157 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.mirror + +import com.vitorpamplona.geode.KtorRelay +import com.vitorpamplona.geode.RelayEngine +import com.vitorpamplona.geode.testing.preload +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl +import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore +import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.coroutines.delay +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import org.junit.After +import kotlin.test.Test +import kotlin.test.assertEquals + +/** + * The mirror is a long-running daemon, so surviving its upstream is not + * optional. This drives the REAL production transport (OkHttp WebSocket + * against a real Ktor port — no in-process shortcut): the upstream is + * stopped mid-mirror and a new instance is brought up on the same port. + * The worker must ride through the disconnect (NostrClient's backoff + + * keep-alive own the re-dial), re-send its REQ, drop the replayed + * duplicate, and pull the event that only exists on the new instance. + * + * Timing note: after a stable connection drops, the client's backoff is + * reset to 1s and re-dial attempts double from there. The worker's 5s + * retry pump nudges the pool well ahead of NostrClient's own 60s + * keep-alive (without it, this test measured a 61s blackout when the + * immediate re-dial raced the port rebind; with it, ~6s). The generous + * timeout only covers a worst-case scheduling stall on a loaded CI + * runner. + */ +class MirrorWorkerReconnectTest { + private val downstreamStore = EventStore(null) + private val downstream = + RelayEngine( + url = "ws://127.0.0.1:7797/".normalizeRelayUrl(), + store = downstreamStore, + parallelVerify = true, + ) + + private var worker: MirrorWorker? = null + private var upstreamServer: KtorRelay? = null + private var upstreamEngine: RelayEngine? = null + + @After + fun tearDown() { + worker?.close() + upstreamServer?.stop(gracePeriodMillis = 0, timeoutMillis = 1_000) + upstreamEngine?.close() + downstream.close() + } + + private fun startUpstream(port: Int): KtorRelay { + val engine = RelayEngine(url = "ws://127.0.0.1:7796/".normalizeRelayUrl()) + val server = KtorRelay(engine, host = "127.0.0.1", port = port).start() + upstreamEngine = engine + upstreamServer = server + return server + } + + private fun forgedEvent(idSeed: Int): Event = + Event( + id = idSeed.toString().padStart(64, '0'), + pubKey = "1".repeat(64), + createdAt = TimeUtils.now() - idSeed, + kind = 1, + tags = emptyArray(), + content = "forged $idSeed", + sig = "f".repeat(128), + ) + + private suspend fun awaitDownstreamCount(expected: Int) = + withTimeout(120_000) { + while (downstreamStore.count(Filter()) < expected) delay(50) + } + + @Test + fun mirrorSurvivesUpstreamRestart() = + runBlocking { + val beforeRestart = forgedEvent(1) + val afterRestart = forgedEvent(2) + + // First upstream instance on an OS-assigned port. + val first = startUpstream(port = 0) + val port = + first.url + .normalizeRelayUrl() + .url + .substringAfterLast(':') + .trimEnd('/') + .toInt() + upstreamEngine!!.preload(beforeRestart) + + // Real OkHttp transport: websocketBuilder deliberately omitted. + val mirror = + MirrorWorker( + upstreams = + listOf( + MirrorUpstream( + url = first.url.normalizeRelayUrl(), + trusted = true, + backfillSeconds = 3600, + ), + ), + server = downstream.server, + ).also { worker = it } + mirror.start() + awaitDownstreamCount(1) + + // Kill the upstream: every socket drops, the port closes. + first.stop(gracePeriodMillis = 0, timeoutMillis = 1_000) + upstreamEngine!!.close() + + // Bring up a NEW instance on the SAME port. Its store has the + // old event (so the re-sent REQ replays a duplicate) plus one + // that only exists post-restart. + val second = startUpstream(port = port) + upstreamEngine!!.preload(beforeRestart, afterRestart) + + // The worker must reconnect on its own — no external poke — + // re-subscribe, and pull the new event. + awaitDownstreamCount(2) + + assertEquals( + setOf(beforeRestart.id, afterRestart.id), + downstreamStore.query(Filter()).map { it.id }.toSet(), + ) + // The duplicate replay of the pre-restart event was dropped by + // the store, not double-inserted. + assertEquals(2, downstreamStore.count(Filter())) + + second.stop(gracePeriodMillis = 0, timeoutMillis = 1_000) + } +}