From 401f36ee82761d9ced67ed243aed472f1ee99495 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 4 Jul 2026 02:32:21 +0000 Subject: [PATCH] fix(geode): mirror trust bound to origin relay; reconnect advances since watermark Two mirror hardening fixes from the audit: - Trust leak: the skip-verify decision was keyed on the subscription id, but the client pool dispatches EVENTs by subscription id alone and every [[mirror]] upstream shares one client. A hostile untrusted upstream could answer with the trusted upstream's subscription id and ride its skip-verify into the store. startDown now drops any event whose delivering relay isn't the one that subscription dialed (relay != up.url). New MirrorWorkerTrustOriginTest injects a foreign subId frame from a hostile relay and asserts it never lands. - Reconnect replay: `since` was frozen at boot, so a long-lived daemon re-streamed the whole backfill window on every upstream flap. The down path now tracks the newest ingested created_at and, on the connected->disconnected edge, advances the REQ's since to watermark - overlap before the reconnect re-sends it. Advancing only on disconnect means a stable link never re-queries. The reconnect test now asserts the watermark advanced. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A --- .../geode/mirror/MirrorWorker.kt | 123 +++++++++++-- .../geode/mirror/MirrorWorkerReconnectTest.kt | 7 + .../mirror/MirrorWorkerTrustOriginTest.kt | 165 ++++++++++++++++++ 3 files changed, 280 insertions(+), 15 deletions(-) create mode 100644 geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerTrustOriginTest.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 60643a24a7..eae9d20b0c 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt @@ -180,15 +180,63 @@ class MirrorWorker( val rejected = AtomicLong(0) /** - * Deliveries dropped by the [MirrorUpstream.filter] re-check before - * ever reaching the store — an upstream sending these is answering - * outside the REQ it was given. + * Deliveries dropped before ever reaching the store: events outside + * the [MirrorUpstream.filter] scope (an upstream answering outside + * the REQ it was given) and events delivered by a different relay + * than the one the subscription dialed (a peer answering with a + * subscription id that isn't its own). */ val filtered = AtomicLong(0) /** Local events handed to the client's outbox for an up-direction upstream. */ val sentUp = AtomicLong(0) + /** + * How many times a down subscription advanced its `since` watermark + * on a reconnect (test/observability hook). Each advance is one + * avoided full-window replay. + */ + val sinceAdvances = AtomicLong(0) + + /** + * One down subscription's re-subscribe state. The upstream replays + * everything at or after the REQ's `since` on every (re)connect; left + * at the boot-time value, a month-old daemon re-streams a month of + * events on every flap. So we track the newest `created_at` ingested + * from this upstream and, on a disconnect, advance the REQ's `since` + * to `watermark - overlap` before the reconnect re-sends it. The + * overlap re-requests a small tail (dup-safe against the store's + * unique-id constraint) to cover out-of-order streaming and clock + * skew. Advancing only on disconnect keeps a healthy connection from + * ever re-querying — no steady-state cost. + */ + private inner class DownSub( + val subId: String, + val up: MirrorUpstream, + val scopedBase: Filter, + val listener: SubscriptionListener, + initialSince: Long, + val watermark: AtomicLong, + ) { + @Volatile + var issuedSince: Long = initialSince + + fun advanceSinceOnReconnect() { + val candidate = watermark.get() - WATERMARK_OVERLAP_SECS + if (candidate > issuedSince) { + issuedSince = candidate + client.subscribe( + subId = subId, + filters = mapOf(up.url to listOf(scopedBase.copy(since = candidate))), + listener = listener, + ) + sinceAdvances.incrementAndGet() + } + } + } + + private val downSubs = mutableListOf() + /** * Recently exchanged event ids, one set per up-capable upstream — * the echo suppressor for [MirrorDirection.BOTH]. An event pulled @@ -254,22 +302,37 @@ class MirrorWorker( // The operator's filter scopes the subscriptions in both // directions; the mirror owns the time window (since) and - // never bounds the result (limit). - val scopedFilter = - (up.filter ?: Filter()).copy( - since = since - up.backfillSeconds, - limit = null, - ) + // never bounds the result (limit). scopedBase carries no + // since — the down path applies (and later advances) it via + // the watermark; the up path uses the fixed initial value. + val scopedBase = (up.filter ?: Filter()).copy(since = null, limit = null) + val initialSince = since - up.backfillSeconds if (up.direction != MirrorDirection.UP) { - startDown(i, up, scopedFilter, exchanged) + downSubs += startDown(i, up, scopedBase, initialSince, exchanged) } if (up.direction != MirrorDirection.DOWN) { - startUp(up, scopedFilter, exchanged) + startUp(up, scopedBase.copy(since = initialSince), exchanged) } } client.connect() + // Watermark advance: when an upstream drops, bump its REQ's since + // to the newest event we ingested from it (minus an overlap) so + // the reconnect doesn't replay the whole window since boot. Only + // fires on the connected→disconnected edge, so a stable link + // never re-queries. + scope.launch { + var prev = emptySet() + client.connectedRelaysFlow().collect { current -> + val dropped = prev - current + if (dropped.isNotEmpty()) { + downSubs.forEach { if (it.up.url in dropped) it.advanceSinceOnReconnect() } + } + prev = current + } + } + // 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 @@ -290,9 +353,14 @@ class MirrorWorker( private fun startDown( index: Int, up: MirrorUpstream, - scopedFilter: Filter, + scopedBase: Filter, + initialSince: Long, exchanged: RecentIds?, - ) { + ): DownSub { + // watermark tracks the newest created_at ingested from this + // upstream; seeded at initialSince so a still-catching-up + // backfill can never advance the since backwards. + val watermark = AtomicLong(initialSince) val listener = object : SubscriptionListener { override fun onEvent( @@ -301,6 +369,18 @@ class MirrorWorker( relay: NormalizedRelayUrl, forFilters: List?, ) { + // The trusted identity is the relay THIS subscription + // dialed. The pool dispatches EVENTs by subscription + // id alone and every upstream shares this one client, + // so a hostile co-configured upstream could answer + // with another subscription's id and ride its trust — + // bind the decision to the delivering relay, not the + // sub id. + if (relay != up.url) { + filtered.incrementAndGet() + Log.w("MirrorWorker") { "dropped event delivered by ${relay.url} on ${up.url.url}'s subscription: ${event.id}" } + return + } // strfry-router parity: never take the upstream's // word for what matched. Re-checking the configured // scope here means even a trusted (skip-verify) @@ -316,6 +396,8 @@ class MirrorWorker( // this subscription — it already exists locally. if (exchanged?.contains(event.id) == true) return exchanged?.add(event.id) + // Advance the reconnect watermark past this event. + watermark.updateAndGet { if (event.createdAt > it) event.createdAt else it } inbound.trySend(Inbound(event, up.trusted)) } @@ -327,11 +409,13 @@ class MirrorWorker( Log.w("MirrorWorker") { "cannot reach upstream ${relay.url}: $message" } } } + val subId = "geode-mirror-$index" client.subscribe( - subId = "geode-mirror-$index", - filters = mapOf(up.url to listOf(scopedFilter)), + subId = subId, + filters = mapOf(up.url to listOf(scopedBase.copy(since = initialSince))), listener = listener, ) + return DownSub(subId, up, scopedBase, listener, initialSince, watermark) } /** @@ -392,6 +476,15 @@ class MirrorWorker( /** How often the retry pump nudges disconnected upstreams. */ const val RECONNECT_POKE_MS = 5_000L + /** + * Overlap (seconds) subtracted from the reconnect `since` + * watermark. Re-requests a small tail on reconnect to cover + * out-of-order streaming and clock skew; the store's unique-id + * constraint drops the duplicates. Bounds worst-case replay after + * a flap to this window instead of everything since boot. + */ + const val WATERMARK_OVERLAP_SECS = 300L + /** * Per-upstream echo-suppression LRU size (BOTH direction only). * Covers the burst window between pulling an event down and the diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt index 9268ca0384..bb6b007287 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerReconnectTest.kt @@ -34,6 +34,7 @@ import kotlinx.coroutines.withTimeout import org.junit.After import kotlin.test.Test import kotlin.test.assertEquals +import kotlin.test.assertTrue /** * The mirror is a long-running daemon, so surviving its upstream is not @@ -152,6 +153,12 @@ class MirrorWorkerReconnectTest { // the store, not double-inserted. assertEquals(2, downstreamStore.count(Filter())) + // The disconnect advanced the REQ's since watermark: the + // reconnect subscribed from ~(newest event - overlap), not + // from the boot-time `now - backfill`. Without this the + // reconnect would replay the whole backfill window every flap. + assertTrue(mirror.sinceAdvances.get() >= 1) + second.stop(gracePeriodMillis = 0, timeoutMillis = 1_000) } } diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerTrustOriginTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerTrustOriginTest.kt new file mode 100644 index 0000000000..28d7701ebc --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorWorkerTrustOriginTest.kt @@ -0,0 +1,165 @@ +/* + * 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.InProcessRelays +import com.vitorpamplona.geode.RelayEngine +import com.vitorpamplona.geode.testing.publish +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder +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 +import kotlin.test.assertTrue + +/** + * The `trusted = true` skip-verify decision must be bound to the relay + * that DELIVERED the event, not to the subscription id it arrived on. + * The client pool dispatches EVENT messages by subscription id alone and + * every `[[mirror]]` upstream shares one client — so a hostile untrusted + * upstream that answers with the *trusted* upstream's subscription id + * would otherwise ride its skip-verify straight into the store. + */ +class MirrorWorkerTrustOriginTest { + private val trustedUrl = RelayUrlNormalizer.normalize("ws://trusted.relay/") + private val hostileUrl = RelayUrlNormalizer.normalize("ws://hostile.relay/") + private val downstreamUrl = RelayUrlNormalizer.normalize("ws://downstream.relay/") + + private val hub = InProcessRelays() + + private val downstreamStore = EventStore(null) + private val downstream = + RelayEngine( + url = downstreamUrl, + store = downstreamStore, + parallelVerify = true, + ) + + private var worker: MirrorWorker? = null + + @After + fun tearDown() { + worker?.close() + downstream.close() + hub.close() + } + + 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), + ) + + /** + * A relay that never answers the REQ it was given; instead, on + * connect it injects one EVENT frame carrying a subscription id that + * belongs to a DIFFERENT upstream's subscription. + */ + private class HijackingWebSocket( + private val out: WebSocketListener, + private val frame: String, + ) : WebSocket { + private var connected = false + + override fun needsReconnect(): Boolean = !connected + + override fun connect() { + connected = true + out.onOpen(0, false) + out.onMessage(frame) + } + + override fun disconnect() { + connected = false + } + + override fun send(msg: String): Boolean = true + } + + @Test + fun hostileUpstreamCannotRideAnotherSubscriptionsTrust() = + runBlocking { + val forged = forgedEvent(1) + // The trusted upstream is configured FIRST, so its down + // subscription id is "geode-mirror-0" — which the hostile + // relay claims in its injected frame. + val hijackFrame = """["EVENT","geode-mirror-0",${forged.toJson()}]""" + + val builder = + object : WebsocketBuilder { + override fun build( + url: NormalizedRelayUrl, + out: WebSocketListener, + ): WebSocket = + if (url == hostileUrl) { + HijackingWebSocket(out, hijackFrame) + } else { + hub.build(url, out) + } + } + + val mirror = + MirrorWorker( + upstreams = + listOf( + MirrorUpstream(trustedUrl, trusted = true, backfillSeconds = 3600), + MirrorUpstream(hostileUrl, trusted = false, backfillSeconds = 3600), + ), + server = downstream.server, + websocketBuilder = builder, + ).also { worker = it } + mirror.start() + + // The injected frame must be dropped at the origin check — + // observable as a `filtered` tick, never as a stored row. + withTimeout(15_000) { + while (mirror.filtered.get() < 1) delay(25) + } + assertTrue(downstreamStore.query(Filter(ids = listOf(forged.id))).isEmpty()) + + // The genuinely trusted upstream still works on the same + // subscription id the hostile relay tried to claim. + val real = forgedEvent(2) + hub.getOrCreate(trustedUrl).publish(real) + withTimeout(15_000) { + while (downstreamStore.count(Filter()) < 1) delay(25) + } + assertEquals( + listOf(real.id), + downstreamStore.query(Filter()).map { it.id }, + ) + } +}