mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A
This commit is contained in:
@@ -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<DownSub>()
|
||||
|
||||
/**
|
||||
* 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<NormalizedRelayUrl>()
|
||||
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<Filter>?,
|
||||
) {
|
||||
// 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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<Event>(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<Event>(Filter()).map { it.id },
|
||||
)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user