Merge pull request #3882 from vitorpamplona/fix/quartz-nip66-record-merge

fix(quartz): update NIP-66 relay records instead of rebuilding them
This commit is contained in:
Vitor Pamplona
2026-08-08 14:50:34 -04:00
committed by GitHub
2 changed files with 455 additions and 11 deletions
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.quartz.nip66RelayMonitor.reachability
import com.vitorpamplona.quartz.nip01Core.core.TagArrayBuilder
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
@@ -30,8 +31,11 @@ import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.networkType
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.requirement
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.rtt
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.NetworkType
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.NetworkTypeTag
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.RequirementTag
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.RttType
import com.vitorpamplona.quartz.utils.TimeUtils
import kotlinx.coroutines.CancellationException
/**
* A durable, shareable relay-reachability cache backed by an [IEventStore] as
@@ -139,8 +143,11 @@ class RelayReachabilityStore(
now: Long = TimeUtils.now(),
rttOpenMs: Long = 0,
) {
for (relay in reachable) writeOne(relay, up = true, now, rttOpenMs)
for (relay in dead) if (relay !in reachable) writeOne(relay, up = false, now, rttOpenMs)
val current = currentRecords(reachable + dead)
val up = reachable.toList()
val down = dead.filterNot { it in reachable }
eachRelay(up.size) { i -> writeOne(up[i], up = true, now, rttOpenMs, current[up[i]]) }
eachRelay(down.size) { i -> writeOne(down[i], up = false, now, rttOpenMs, current[down[i]]) }
}
/**
@@ -154,8 +161,11 @@ class RelayReachabilityStore(
dead: Set<NormalizedRelayUrl>,
now: Long = TimeUtils.now(),
) {
for ((relay, rtt) in reachableRttMs) writeOne(relay, up = true, now, rtt.coerceAtLeast(0))
for (relay in dead) if (relay !in reachableRttMs) writeOne(relay, up = false, now, 0)
val current = currentRecords(reachableRttMs.keys + dead)
val up = reachableRttMs.toList()
val down = dead.filterNot { it in reachableRttMs }
eachRelay(up.size) { i -> writeOne(up[i].first, up = true, now, up[i].second.coerceAtLeast(0), current[up[i].first]) }
eachRelay(down.size) { i -> writeOne(down[i], up = false, now, 0, current[down[i]]) }
}
/**
@@ -171,21 +181,55 @@ class RelayReachabilityStore(
observations: Collection<RelayObserver.Observation>,
now: Long = TimeUtils.now(),
): Int {
val reported = observations.filter { it.reachable || it.error != null }
val current = currentRecords(reported.map { it.url })
return eachRelay(reported.size) { i -> writeObserved(reported[i], now, current[reported[i].url]) }
}
/**
* Run every relay's write, then fail if any of them did.
*
* The read-modify-write spans a store round trip and [IEventStore] offers
* no read inside a transaction, so a concurrent writer to the same address
* can win the race and our now-stale insert is REJECTED. One such throw
* must not end the loop and drop every relay after it — but it must not
* vanish either: [RelayObserver.collectUnreported] has already cleared the
* flags by the time this runs, so a swallowed failure loses the
* measurement for good and the caller cannot tell an empty run from a
* failed one. So: attempt all, remember the first failure, rethrow it.
*
* Cancellation is never caught. A shutdown flush wrapped in `withTimeout`
* — the pattern [RelayMonitor.close] prescribes — would otherwise be
* unabortable, grinding through every remaining relay with each write
* throwing and being swallowed.
*/
private suspend inline fun eachRelay(
count: Int,
write: (Int) -> Unit,
): Int {
var first: Exception? = null
var written = 0
for (o in observations) {
if (!o.reachable && o.error == null) continue
writeObserved(o, now)
written++
for (i in 0 until count) {
try {
write(i)
written++
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
if (first == null) first = e
}
}
first?.let { throw it }
return written
}
private suspend fun writeObserved(
o: RelayObserver.Observation,
now: Long,
current: RelayDiscoveryEvent?,
) {
val template =
RelayDiscoveryEvent.build(o.url, createdAt = now) {
edit(o.url, now, current, ::ownsLiveness) {
networkType(networkTypeOf(o.url))
if (o.reachable) {
// Liveness is the presence of rtt-open, per NIP-66. A relay we
@@ -200,7 +244,7 @@ class RelayReachabilityStore(
// Observed, not read off NIP-11: this relay actually challenged
// us. A relay advertising open reads and then demanding AUTH is
// exactly what a monitor exists to catch.
if (o.authRequired) requirement("auth")
if (o.authRequired) requirement(AUTH_REQUIREMENT)
}
store.insert(signer.sign(template))
}
@@ -210,16 +254,145 @@ class RelayReachabilityStore(
up: Boolean,
now: Long,
rttOpenMs: Long,
current: RelayDiscoveryEvent?,
) {
val template =
RelayDiscoveryEvent.build(relay, createdAt = now) {
edit(relay, now, current, ::ownsLiveness) {
networkType(networkTypeOf(relay))
if (up) rtt(RttType.OPEN, rttOpenMs)
}
store.insert(signer.sign(template))
}
/**
* Build this monitor's next record for [relay] as an EDIT of [current]
* rather than a fresh document.
*
* A 30166 is addressable, so a relay has exactly one record per monitor —
* and this class is not necessarily its only writer. Anything else keeping
* per-relay knowledge under the same identity (an operator marking a relay
* as a mirror of another, a crawler recording which kinds it served) writes
* into this same slot, and a build-from-scratch silently deletes it. The
* result still signs, still parses, and still reads as a valid NIP-66
* record — it just says less than it did, and the reader downstream has no
* way to know something was lost.
*
* [owns] decides what this writer measured and may therefore replace — a
* predicate rather than a set of names, because ownership is sometimes per
* VALUE: an observation may clear `R auth` without touching the `R pow`
* another writer measured. Everything it does not claim is carried across
* untouched, including tags this version of quartz has never heard of.
*
* The timestamp is `max(now, current + 1)`, not `now`: a store enforcing
* replaceable semantics REJECTS a record that is not strictly newer than
* the one it replaces, and two writers inside the same second — or a peer
* whose clock runs slightly ahead — are ordinary. An update lost that way
* is indistinguishable from one that had nothing to say.
*
* The bump is deliberately NOT capped to some window past `now`. Capping
* it looks prudent and is worse: a record already further ahead than the
* cap can then never be replaced at all, because every stamp we are willing
* to write is older than what is stored, so the relay's live/dead verdict
* freezes until the wall clock catches up. It does not even buy the thing
* it appears to — [snapshot] selects on `since` alone, so a future-stamped
* record sits inside the freshness window either way. A record stamped
* ahead of the clock is a defect in whatever produced it; this class's job
* is to keep updating it, not to freeze it.
*/
private fun edit(
relay: NormalizedRelayUrl,
now: Long,
current: RelayDiscoveryEvent?,
owns: (Array<String>) -> Boolean,
measured: TagArrayBuilder<RelayDiscoveryEvent>.() -> Unit,
) = RelayDiscoveryEvent.build(
relay,
current?.content ?: "",
createdAt = maxOf(now, (current?.createdAt ?: 0L) + 1),
) {
current?.tags?.forEach { tag ->
if (tag.firstOrNull() != "d" && !owns(tag)) add(tag)
}
measured()
}
/**
* The tags this class measures, and may therefore replace.
*
* Everything here expires together with the record: a 30166 carries ONE
* `created_at` for the whole document, so a tag carried across is re-dated
* as a current measurement. Keeping a `rtt-read` from an earlier
* observation beside a fresh `rtt-open` would republish a stale latency as
* today's — and [RelayObserver] documents exactly how wrong a queued rtt
* can be. So this class's own liveness facts are rewritten wholesale on
* every write, including clearing `R auth` when nothing re-asserts it: a
* permanent auth flag is worse than a missing one, because it discourages
* the very connection that could clear it.
*
* Both polarities of the auth requirement are owned. Owning only the
* positive form let `R !auth` survive while `requirement("auth")` appended
* the opposite, publishing a record that asserted both at once.
*
* Everything NOT matched here — `R pow`, `R payment`, annotations another
* writer keeps on this address, tags this version has never heard of — is
* somebody else's measurement and is carried across untouched.
*/
private fun ownsLiveness(tag: Array<String>): Boolean =
when (tag.firstOrNull()) {
NetworkTypeTag.TAG_NAME -> true
in ALL_RTT -> true
RequirementTag.TAG_NAME -> tag.getOrNull(1) in AUTH_REQUIREMENT_FORMS
else -> false
}
/**
* This monitor's own current record for each relay, in one query.
*
* Only OUR records: merging another monitor's tags into a document signed
* with this key would republish their claims as ours.
*/
private suspend fun currentRecords(relays: Collection<NormalizedRelayUrl>): Map<NormalizedRelayUrl, RelayDiscoveryEvent> {
if (relays.isEmpty()) return emptyMap()
val out = HashMap<NormalizedRelayUrl, RelayDiscoveryEvent>()
// CHUNKED: a `d` filter binds one host parameter per url, and callers
// pass the whole relay universe — the fan-out this module measures
// itself against is 16,507 relays (see RelayObserver). Measured on
// BundledSQLiteDriver, 32,765 `d` values pass and 32,766 fails with
// "too many SQL variables"; the throw lands BEFORE anything is
// written, so an entire probe run's records are lost rather than one
// relay's. The headroom here is deliberate — the ceiling is a property
// of the driver, not of this query.
for (chunk in relays.map { it.url }.distinct().chunked(RELAYS_PER_QUERY)) {
val held =
store.query<RelayDiscoveryEvent>(
Filter(kinds = listOf(RelayDiscoveryEvent.KIND), authors = listOf(signer.pubKey), tags = mapOf("d" to chunk)),
)
for (ev in held) {
val relay = ev.relay() ?: continue
val seen = out[relay]
if (seen == null || ev.createdAt > seen.createdAt) out[relay] = ev
}
}
return out
}
companion object {
/** Every rtt tag name. A DEAD record must clear all of them: liveness is the presence of `rtt-open`. */
private val ALL_RTT = setOf(RttType.OPEN.tagName, RttType.READ.tagName, RttType.WRITE.tagName)
/** The one NIP-66 requirement an observation can prove: the relay challenged us. */
const val AUTH_REQUIREMENT = "auth"
/** Both polarities, so an update cannot leave the record asserting `auth` and `!auth` at once. */
private val AUTH_REQUIREMENT_FORMS = setOf(AUTH_REQUIREMENT, "!$AUTH_REQUIREMENT")
/**
* Urls per `d` lookup. Well under a bundled SQLite's 32,766-variable
* ceiling, and in the same range as the author chunking elsewhere in
* this codebase.
*/
const val RELAYS_PER_QUERY = 500
/** Default freshness window: a relay's status is trusted for a day, then re-probed. */
const val DEFAULT_TTL_SECONDS = 24L * 60 * 60
@@ -21,11 +21,16 @@
package com.vitorpamplona.quartz.nip66RelayMonitor.reachability
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
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.signers.NostrSignerInternal
import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.RelayDiscoveryEvent
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.NetworkType
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.RequirementTag
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.RttType
import kotlinx.coroutines.runBlocking
import kotlin.test.Test
import kotlin.test.assertEquals
@@ -110,4 +115,270 @@ class RelayReachabilityStoreTest {
// A host that merely contains ".onion" as a substring is clearnet, not Tor.
assertEquals(NetworkType.CLEARNET, RelayReachabilityStore.networkTypeOf(fakeOnion))
}
// ---- one address, more than one writer ---------------------------------
private suspend fun tagsOf(
store: EventStore,
signer: NostrSignerInternal,
relay: NormalizedRelayUrl,
): List<Array<String>> =
store
.query<RelayDiscoveryEvent>(
Filter(kinds = listOf(RelayDiscoveryEvent.KIND), authors = listOf(signer.pubKey), tags = mapOf("d" to listOf(relay.url))),
).maxByOrNull { it.createdAt }
?.tags
?.toList()
.orEmpty()
private fun names(tags: List<Array<String>>) = tags.mapNotNull { it.firstOrNull() }.toSet()
/**
* A 30166 is addressable, so this monitor has one record per relay — and it
* is not necessarily the only thing writing per-relay knowledge under that
* identity. A record rebuilt from this writer's own tags deletes the rest,
* and the loss is invisible: the event still signs and still parses.
*/
@Test
fun `an update keeps tags this writer does not own`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_000)
// Something else records what it knows about the same relay,
// keeping what the monitor already put there.
val existing = tagsOf(store, signer, live1)
val withExtra =
RelayDiscoveryEvent.build(live1, "", createdAt = 1_700_000_001) {
existing.forEach { if (it.firstOrNull() != "d") add(it) }
add(arrayOf("redirect", "wss://canonical.example.com/"))
}
store.insert(signer.sign(withExtra))
cache.recordProbed(mapOf(live1 to 131L), emptySet(), now = 1_700_000_002)
val after = tagsOf(store, signer, live1)
assertTrue("redirect" in names(after), "the update erased another writer's tag: ${names(after)}")
// ...and still replaced what it does own.
assertEquals("131", after.first { it[0] == RttType.OPEN.tagName }[1])
}
/**
* A store enforcing replaceable semantics rejects a record that is not
* strictly newer than the one it replaces. Two writers inside one second,
* or a peer whose clock runs ahead, are ordinary — and an update lost that
* way looks exactly like one that had nothing to say.
*/
@Test
fun `an update lands even when the record it replaces is newer than the clock`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
// Ordinary skew: the record ahead of our clock by seconds, which is
// what two writers or a slightly fast peer produce.
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_010)
cache.recordProbed(mapOf(live1 to 131L), emptySet(), now = 1_700_000_000)
assertEquals("131", tagsOf(store, signer, live1).first { it[0] == RttType.OPEN.tagName }[1])
}
/**
* A record stamped ahead of the clock is a defect in whatever produced it.
* Capping our stamp to some window past `now` looks prudent and is worse:
* every stamp we would write is then older than what is stored, so the
* insert is rejected and the relay's live/dead verdict freezes until the
* wall clock catches up. It does not even buy freshness — snapshot()
* selects on `since` alone, so the future record is inside the window
* either way.
*/
@Test
fun `a record stamped ahead of the clock can still be updated`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
val now = 1_700_000_000L
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = now + 86_400)
cache.recordProbed(mapOf(live1 to 131L), emptySet(), now = now)
assertEquals("131", tagsOf(store, signer, live1).first { it[0] == RttType.OPEN.tagName }[1])
}
/**
* A 30166 carries ONE created_at, so a carried tag is re-dated as a current
* measurement. This class's own liveness facts must therefore be rewritten
* wholesale, or a stale rtt-read is republished as today's number — which
* aggregators rank on.
*/
@Test
fun `a later observation does not re-date an older latency`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
val observed = RelayObserver()
observed.record(live1, true, 100L, null)
observed.observationOf(live1)?.rttReadMs = 55L
cache.record(observed.collectUnreported(), now = 1_700_000_000)
assertTrue(RttType.READ.tagName in names(tagsOf(store, signer, live1)))
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_100)
val after = tagsOf(store, signer, live1)
assertTrue(RttType.READ.tagName !in names(after), "a stale rtt-read was carried onto a fresh record")
assertEquals("120", after.first { it[0] == RttType.OPEN.tagName }[1])
}
/**
* Only [writeObserved] can clear `auth`, and it needs a connection the flag
* discourages — so carrying it forward would make it permanent. It expires
* with the rest of this class's liveness facts.
*/
@Test
fun `an auth requirement does not outlive the observation that set it`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
val walled = RelayObserver()
walled.record(live1, true, 100L, null)
walled.observationOf(live1)?.authRequired = true
cache.record(walled.collectUnreported(), now = 1_700_000_000)
assertTrue(RequirementTag.TAG_NAME in names(tagsOf(store, signer, live1)))
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_100)
assertTrue(RequirementTag.TAG_NAME !in names(tagsOf(store, signer, live1)), "the auth wall became permanent")
}
/** Owning only the positive form left `!auth` in place while appending `auth`. */
@Test
fun `an update never leaves the record asserting both auth polarities`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_000)
val existing = tagsOf(store, signer, live1)
val negated =
RelayDiscoveryEvent.build(live1, "", createdAt = 1_700_000_001) {
existing.forEach { if (it.firstOrNull() != "d") add(it) }
add(arrayOf(RequirementTag.TAG_NAME, "!auth"))
}
store.insert(signer.sign(negated))
val walled = RelayObserver()
walled.record(live1, true, 100L, null)
walled.observationOf(live1)?.authRequired = true
cache.record(walled.collectUnreported(), now = 1_700_000_002)
val reqs = tagsOf(store, signer, live1).filter { it[0] == RequirementTag.TAG_NAME }.map { it[1] }
assertEquals(listOf(RelayReachabilityStore.AUTH_REQUIREMENT), reqs, "record asserts contradictory requirements: " + reqs)
}
/**
* collectUnreported() has already cleared the flags by the time a write
* runs, so a swallowed failure loses the measurement for good and an empty
* run is indistinguishable from a failed one.
*/
@Test
fun `a failing write is reported, not swallowed`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
store.close()
var threw = false
try {
cache.recordProbed(mapOf(live1 to 120L, live2 to 130L), emptySet(), now = 1_700_000_000)
} catch (e: Exception) {
threw = true
}
assertTrue(threw, "every write failed and the run reported success")
}
/** An observation proves `R auth` and nothing else; other requirements belong to whoever measured them. */
@Test
fun `an observation clears only the auth requirement`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_000)
val existing = tagsOf(store, signer, live1)
val withReqs =
RelayDiscoveryEvent.build(live1, "", createdAt = 1_700_000_001) {
existing.forEach { if (it.firstOrNull() != "d") add(it) }
add(arrayOf(RequirementTag.TAG_NAME, "pow"))
add(arrayOf(RequirementTag.TAG_NAME, RelayReachabilityStore.AUTH_REQUIREMENT))
}
store.insert(signer.sign(withReqs))
// Reached without a challenge: auth no longer holds, pow was never ours.
val observed = RelayObserver()
observed.record(live1, true, 100L, null)
cache.record(observed.collectUnreported(), now = 1_700_000_002)
val after = tagsOf(store, signer, live1).filter { it[0] == RequirementTag.TAG_NAME }.map { it[1] }
assertTrue("pow" in after, "another writer's requirement was erased: " + after)
assertTrue(RelayReachabilityStore.AUTH_REQUIREMENT !in after, "auth should have been cleared: " + after)
}
/** One `d` filter binds one host parameter per url, and callers pass the whole relay universe. */
@Test
fun `a flush wider than one query chunk still writes every relay`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
val many = (0 until RelayReachabilityStore.RELAYS_PER_QUERY * 2 + 7).map { RelayUrlNormalizer.normalize("wss://r" + it + ".example.com") }
cache.record(reachable = many.toSet(), dead = emptySet(), now = 1_700_000_000)
assertEquals(many.size, cache.snapshot(now = 1_700_000_100).live.size)
}
/** A relay that went down must lose its rtt, or it still reads as live. */
@Test
fun `a dead update clears the rtt it replaces`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_000)
cache.record(reachable = emptySet(), dead = setOf(live1), now = 1_700_000_100)
assertTrue(RttType.OPEN.tagName !in names(tagsOf(store, signer, live1)))
assertTrue(cache.snapshot(now = 1_700_000_200).isKnownDead(live1))
}
/** Only OUR records merge: republishing another monitor's tags under this key would launder their claims. */
@Test
fun `another monitor's record is not merged into ours`() =
runBlocking {
val store = store()
val signer = NostrSignerInternal(KeyPair())
val other = NostrSignerInternal(KeyPair())
val cache = RelayReachabilityStore(store, signer, ttlSeconds = 3600)
val theirs =
RelayDiscoveryEvent.build(live1, "", createdAt = 1_700_000_000) {
add(arrayOf("redirect", "wss://not-ours.example.com/"))
}
store.insert(other.sign(theirs))
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_000_100)
assertTrue("redirect" !in names(tagsOf(store, signer, live1)))
}
}