mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-08-11 00:37:41 +00:00
fix(quartz): update NIP-66 relay records instead of rebuilding them
A kind:30166 is addressable, so RelayReachabilityStore keeps exactly one record per (monitor, relay) — but it is not necessarily the only thing writing per-relay knowledge under that identity. Both write paths built the record from their own tags and inserted it, so every update deleted whatever else was in that slot. Observed while adding a "this url is an alias of that one" tag alongside the monitor: `[d, n, rtt-open]` became `[d, redirect]` on our write, and the monitor's next observation turned it back into `[d, n, rtt-open]`. Nothing looks wrong at any point — the event still signs, still parses, still reads as a valid NIP-66 record. It just says less than it did, and the reader downstream cannot tell. Writing is now an edit: read this monitor's current record, carry across every tag the writer does not own — including tags this version of quartz has never heard of — and replace only what it measured. `n` and the three `rtt-*` types are owned by both paths, so a dead update still clears a stale rtt and liveness keeps meaning what it meant. `R` is owned only by the observation path, which is the one that learns whether a relay challenged us; writeOne leaves it alone rather than deleting what it cannot re-measure. Only OUR records are merged. Folding another monitor's tags into a document signed with this key would republish their claims as ours. The timestamp is now `max(now, current + 1)` rather than `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 ahead — are ordinary. That is not theoretical: it silently swallowed a repair pass in the caller that found this bug, which reported success having written nothing. The reads are batched per call rather than per relay, so a flush over N relays costs one extra query, not N. Test plan: ./gradlew :quartz:jvmTest — 4,071 tests, all passing, including four new cases in RelayReachabilityStoreTest covering a foreign tag surviving an update, an update against a record stamped an hour ahead, a dead update clearing its rtt, and another monitor's record not being merged. ./gradlew :quartz:spotlessApply clean.
This commit is contained in:
+95
-9
@@ -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,6 +31,8 @@ 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
|
||||
|
||||
@@ -139,8 +142,9 @@ 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)
|
||||
for (relay in reachable) writeOne(relay, up = true, now, rttOpenMs, current[relay])
|
||||
for (relay in dead) if (relay !in reachable) writeOne(relay, up = false, now, rttOpenMs, current[relay])
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -154,8 +158,9 @@ 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)
|
||||
for ((relay, rtt) in reachableRttMs) writeOne(relay, up = true, now, rtt.coerceAtLeast(0), current[relay])
|
||||
for (relay in dead) if (relay !in reachableRttMs) writeOne(relay, up = false, now, 0, current[relay])
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -171,10 +176,11 @@ 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 })
|
||||
var written = 0
|
||||
for (o in observations) {
|
||||
if (!o.reachable && o.error == null) continue
|
||||
writeObserved(o, now)
|
||||
for (o in reported) {
|
||||
writeObserved(o, now, current[o.url])
|
||||
written++
|
||||
}
|
||||
return written
|
||||
@@ -183,9 +189,14 @@ class RelayReachabilityStore(
|
||||
private suspend fun writeObserved(
|
||||
o: RelayObserver.Observation,
|
||||
now: Long,
|
||||
current: RelayDiscoveryEvent?,
|
||||
) {
|
||||
// `R` is included in the owned set here and NOT in [writeOne]: an
|
||||
// observation knows whether this relay challenged us, so it may clear a
|
||||
// requirement that no longer holds. writeOne never learns that, so it
|
||||
// leaves the tag alone rather than deleting what it cannot re-measure.
|
||||
val template =
|
||||
RelayDiscoveryEvent.build(o.url, createdAt = now) {
|
||||
edit(o.url, now, current, OWNED_LIVENESS + RequirementTag.TAG_NAME) {
|
||||
networkType(networkTypeOf(o.url))
|
||||
if (o.reachable) {
|
||||
// Liveness is the presence of rtt-open, per NIP-66. A relay we
|
||||
@@ -210,16 +221,91 @@ class RelayReachabilityStore(
|
||||
up: Boolean,
|
||||
now: Long,
|
||||
rttOpenMs: Long,
|
||||
current: RelayDiscoveryEvent?,
|
||||
) {
|
||||
val template =
|
||||
RelayDiscoveryEvent.build(relay, createdAt = now) {
|
||||
edit(relay, now, current, OWNED_LIVENESS) {
|
||||
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.
|
||||
*
|
||||
* [owned] is what this writer measured and may therefore replace.
|
||||
* Everything else 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 ahead of ours — are ordinary. An update lost that way is
|
||||
* indistinguishable from one that had nothing to say.
|
||||
*/
|
||||
private fun edit(
|
||||
relay: NormalizedRelayUrl,
|
||||
now: Long,
|
||||
current: RelayDiscoveryEvent?,
|
||||
owned: Set<String>,
|
||||
measured: TagArrayBuilder<RelayDiscoveryEvent>.() -> Unit,
|
||||
) = RelayDiscoveryEvent.build(
|
||||
relay,
|
||||
current?.content ?: "",
|
||||
createdAt = maxOf(now, (current?.createdAt ?: 0L) + 1),
|
||||
) {
|
||||
current?.tags?.forEach { tag ->
|
||||
if (tag.firstOrNull() != "d" && tag.firstOrNull() !in owned) add(tag)
|
||||
}
|
||||
measured()
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 held =
|
||||
store.query<RelayDiscoveryEvent>(
|
||||
Filter(
|
||||
kinds = listOf(RelayDiscoveryEvent.KIND),
|
||||
authors = listOf(signer.pubKey),
|
||||
tags = mapOf("d" to relays.map { it.url }.distinct()),
|
||||
),
|
||||
)
|
||||
val out = HashMap<NormalizedRelayUrl, RelayDiscoveryEvent>(held.size)
|
||||
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 {
|
||||
/**
|
||||
* The tags this class measures on every write, and may therefore
|
||||
* replace. A dead record must be able to CLEAR a stale rtt — liveness
|
||||
* is the presence of `rtt-open` — so all three rtt types are owned even
|
||||
* though only `rtt-open` is written by every path.
|
||||
*/
|
||||
private val OWNED_LIVENESS =
|
||||
setOf(NetworkTypeTag.TAG_NAME, RttType.OPEN.tagName, RttType.READ.tagName, RttType.WRITE.tagName)
|
||||
|
||||
/** Default freshness window: a relay's status is trusted for a day, then re-probed. */
|
||||
const val DEFAULT_TTL_SECONDS = 24L * 60 * 60
|
||||
|
||||
|
||||
+106
@@ -21,11 +21,15 @@
|
||||
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.RttType
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
@@ -110,4 +114,106 @@ 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)
|
||||
|
||||
cache.recordProbed(mapOf(live1 to 120L), emptySet(), now = 1_700_003_600)
|
||||
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 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)))
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user