feat(quartz): read+write relay checks — observed facts only, no NIP-11 claims

RelayProber.readWriteCheck(relays, signer) is the deeper check pair for
relays already proven live (warm sockets from a probe that just ran):

- READ: a real limit-1 REQ the relay must query its store for, timed
  REQ→first answer (honest rtt-read on an open socket).
- WRITE: one ephemeral RelayProbeWriteTest event signed by the monitor
  key, timed publish→OK (honest rtt-write). An OK false is a measured
  policy answer, kept with its NIP-01 machine-readable reason; only
  silence leaves the write side unobserved (writeAccepted = null).

publishAndCollectResults now stamps each OK with its elapsedMs (a
rejection is still a round trip; -1 when the relay never answered), so
any caller gets write latency for free.

toDiscoveryEventTemplate(readWrite = ...) folds the pair into the 30166
template: rtt-read/rtt-write when measured, R auth / R pow when the
write was refused with auth-required:/pow:. NIP-11-derived tags (N
supported NIPs, k kinds, T type) are deliberately NOT emitted — those
are relay self-claims, and publishing them under a monitor signature
without per-NIP compliance tests would launder claims into
measurements. Per-NIP/per-kind compliance suites can come later as
opt-in checks; open/read/write is the default surface.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01A9sSyh1QLJD3PZ18tVPPVK
This commit is contained in:
Claude
2026-08-04 14:41:15 +00:00
parent a7eec1d605
commit fd5bd994a1
3 changed files with 264 additions and 13 deletions
@@ -34,6 +34,7 @@ import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.withTimeoutOrNull
import kotlin.time.TimeSource
/**
* One relay's verdict on a published event: [accepted] plus the reason the
@@ -45,6 +46,12 @@ import kotlinx.coroutines.withTimeoutOrNull
class PublishResult(
val accepted: Boolean,
val message: String,
/**
* Milliseconds from the publish to this relay's OK (true or false — a rejection
* is still a measured round trip), or -1 when the relay never answered with an
* OK. On an already-open socket this is an honest NIP-66 `rtt-write`.
*/
val elapsedMs: Long = -1,
) {
/**
* True when this failure came from the transport (never connected,
@@ -102,6 +109,7 @@ suspend fun INostrClient.publishAndCollectResults(
timeoutInSeconds: Long = 15,
): Map<NormalizedRelayUrl, PublishResult> {
val resultChannel = Channel<DetailedResult>(UNLIMITED)
val mark = TimeSource.Monotonic.markNow()
Log.d("publishAndConfirm") { "Waiting for ${relayList.size} responses" }
@@ -134,7 +142,7 @@ suspend fun INostrClient.publishAndCollectResults(
when (msg) {
is OkMessage -> {
if (msg.eventId == event.id) {
resultChannel.trySend(DetailedResult(relay.url, msg.success, msg.message))
resultChannel.trySend(DetailedResult(relay.url, msg.success, msg.message, mark.elapsedNow().inWholeMilliseconds))
Log.d("publishAndConfirm") { "onSendResponse Received response for ${msg.eventId} from relay ${relay.url} message ${msg.message} success ${msg.success}" }
}
}
@@ -160,7 +168,7 @@ suspend fun INostrClient.publishAndCollectResults(
val currentResult = receivedResults[result.relay]
// do not override a successful result.
if (currentResult == null || !currentResult.accepted) {
receivedResults[result.relay] = PublishResult(result.success, result.message)
receivedResults[result.relay] = PublishResult(result.success, result.message, result.elapsedMs)
}
}
}
@@ -191,4 +199,5 @@ private class DetailedResult(
val relay: NormalizedRelayUrl,
val success: Boolean,
val message: String,
val elapsedMs: Long = -1,
)
@@ -22,6 +22,8 @@ package com.vitorpamplona.quartz.nip66RelayMonitor.reachability
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.PublishResult
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndCollectResults
import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient
@@ -30,6 +32,7 @@ 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.EventTemplate
import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.RelayDiscoveryEvent
@@ -83,6 +86,19 @@ class RelayProber(
val error: String?,
)
/**
* One relay's read+write check outcome (see [readWriteCheck]). Latencies are
* -1 when unobserved; [writeAccepted] is null when the relay never answered
* the write with an OK (transport failure or silence).
*/
class ReadWriteVerdict(
val relay: NormalizedRelayUrl,
val rttReadMs: Long,
val rttWriteMs: Long,
val writeAccepted: Boolean?,
val writeMessage: String?,
)
class Result(
val verdicts: List<Verdict>,
val elapsedMs: Long,
@@ -153,6 +169,55 @@ class RelayProber(
}
}
/**
* The deeper, still-honest check pair: READ (a real limit-[readLimit] REQ the
* relay must query its store for) and WRITE (one ephemeral [RelayProbeWriteTest]
* event signed by [signer], the monitor key, timed to its OK). Everything is a
* direct observation — nothing is copied from the relay's NIP-11 self-claims.
*
* Run it against relays ALREADY PROVEN LIVE — typically [Result.reachable] of a
* [probe] that just ran, while the pool's sockets are still open. On a warm
* socket both numbers are honest NIP-66 rtts (`rtt-read`, `rtt-write`); against
* a cold relay they silently include the dial, so don't.
*
* A write REJECTION is still a measurement: `OK false` proves the write path
* works and documents policy ([ReadWriteVerdict.writeMessage] keeps the NIP-01
* machine-readable reason; [toDiscoveryEventTemplate] maps `auth-required:` and
* `pow:` to `R` tags). Only silence leaves [ReadWriteVerdict.writeAccepted] null.
*/
suspend fun readWriteCheck(
relays: Collection<NormalizedRelayUrl>,
signer: NostrSigner,
timeoutMs: Long = 15_000,
waveSize: Int = 1000,
readLimit: Int = 1,
): Map<NormalizedRelayUrl, ReadWriteVerdict> {
val out = HashMap<NormalizedRelayUrl, ReadWriteVerdict>()
val distinct = relays.toSet()
for (wave in distinct.chunked(waveSize.coerceAtLeast(1))) {
val reads = HashMap<NormalizedRelayUrl, Long>()
probeWave(wave, timeoutMs, readTestFilter(readLimit)) { reads[it.relay] = it.rttEoseMs }
val event = signer.sign(RelayProbeWriteTest.build())
val writes = client.publishAndCollectResults(event, wave.toSet(), (timeoutMs / 1000).coerceAtLeast(1))
for (relay in wave) {
// Only a real OK (true or false) counts as an answer; transport
// failures and silence leave the write side unobserved.
val answered = writes[relay]?.takeUnless { it.isTransportFailure || it.message == PublishResult.NO_RESPONSE }
out[relay] =
ReadWriteVerdict(
relay = relay,
rttReadMs = reads[relay] ?: -1,
rttWriteMs = answered?.elapsedMs ?: -1,
writeAccepted = answered?.accepted,
writeMessage = answered?.message,
)
}
}
return out
}
private suspend fun probeWave(
wave: List<NormalizedRelayUrl>,
timeoutMs: Long,
@@ -329,22 +394,37 @@ class RelayProber(
* is the normalized relay url; sign it with the consumer's own monitor key (per
* NIP-66 a monitor is its own identity, so the prober never signs on its own).
*
* Only facts this probe actually observed are tagged:
* Only facts a probe actually observed are tagged:
* - `n` network type inferred from the url (clearnet/tor/i2p);
* - `rtt-open` when the relay was reachable — the measured WS-upgrade round trip,
* or 0 for "reachable, latency not observed" (liveness is the tag's PRESENCE);
* - `R auth` when the relay answered the probe REQ with a NIP-42 `auth-required`
* CLOSED — an observed auth wall, not a NIP-11 claim.
* - `rtt-read`/`rtt-write` when a [RelayProber.readWriteCheck] result is passed
* as [readWrite] and actually measured that side;
* - `R auth` when the relay answered a probe with a NIP-42 `auth-required`
* (CLOSED on the REQ, or OK-false on the write) and `R pow` when the write
* was refused with a `pow:` reason — observed walls, not NIP-11 claims.
*
* [Verdict.rttEoseMs] is deliberately NOT written as `rtt-read`: it is measured from
* the wave start, so it bundles dial, TLS and any handshake queueing with the read —
* publishing it as a read round trip would hand aggregators an inflated latency.
* The [RelayObserver]/[RelayMonitor] path supplies honest `rtt-read`/`rtt-write`
* from real traffic instead.
* NIP-11-derived tags (`N` supported NIPs, `k` kinds, `T` type) are deliberately
* absent: those are the relay's self-claims, and asserting them under a monitor
* signature without a per-NIP compliance test would launder claims into
* measurements. [Verdict.rttEoseMs] is likewise never written as `rtt-read` — it
* is wave-relative (dial + TLS + queueing), so the honest read number only comes
* from [readWrite] (or the [RelayObserver]/[RelayMonitor] real-traffic path).
*/
fun RelayProber.Verdict.toDiscoveryEventTemplate(createdAt: Long = TimeUtils.now()): EventTemplate<RelayDiscoveryEvent> =
fun RelayProber.Verdict.toDiscoveryEventTemplate(
createdAt: Long = TimeUtils.now(),
readWrite: RelayProber.ReadWriteVerdict? = null,
): EventTemplate<RelayDiscoveryEvent> =
RelayDiscoveryEvent.build(relay, createdAt = createdAt) {
networkType(RelayReachabilityStore.networkTypeOf(relay))
if (reachable) rtt(RttType.OPEN, rttOpenMs.coerceAtLeast(0))
if (error?.startsWith("closed:auth-required") == true) requirement("auth")
if (readWrite != null) {
if (readWrite.rttReadMs >= 0) rtt(RttType.READ, readWrite.rttReadMs)
if (readWrite.rttWriteMs >= 0) rtt(RttType.WRITE, readWrite.rttWriteMs)
}
val authWalled =
error?.startsWith("closed:auth-required") == true ||
readWrite?.writeMessage?.startsWith("auth-required") == true
if (authWalled) requirement("auth")
if (readWrite?.writeMessage?.startsWith("pow:") == true) requirement("pow")
}
@@ -20,13 +20,20 @@
*/
package com.vitorpamplona.quartz.nip66RelayMonitor.reachability
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
import com.vitorpamplona.quartz.nip01Core.relay.client.EmptyNostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.OkMessage
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.Command
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.EventTemplate
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
@@ -46,10 +53,12 @@ import kotlin.test.assertTrue
*/
@OptIn(ExperimentalCoroutinesApi::class)
class RelayProberFlowTest {
/** Captures the probe subscription so the test can play the relays. */
/** Captures the probe subscription and publish so the test can play the relays. */
private class ScriptedClient : INostrClient by EmptyNostrClient() {
var listener: SubscriptionListener? = null
var sentFilters: Map<NormalizedRelayUrl, List<Filter>>? = null
var published: Event? = null
val connListeners = mutableListOf<RelayConnectionListener>()
override fun subscribe(
subId: String,
@@ -59,6 +68,49 @@ class RelayProberFlowTest {
this.listener = listener
this.sentFilters = filters
}
override fun publish(
event: Event,
relayList: Set<NormalizedRelayUrl>,
) {
published = event
}
override fun addConnectionListener(listener: RelayConnectionListener) {
connListeners += listener
}
override fun removeConnectionListener(listener: RelayConnectionListener) {
connListeners -= listener
}
/** Plays a relay's OK answer for the published event to every armed listener. */
fun answerOk(
relay: NormalizedRelayUrl,
success: Boolean,
message: String,
) {
val ok = OkMessage(published!!.id, success, message)
connListeners.toList().forEach { it.onIncomingMessage(FakeRelayClient(relay), "", ok) }
}
}
private class FakeRelayClient(
override val url: NormalizedRelayUrl,
) : IRelayClient {
override fun connect() = Unit
override fun needsToReconnect() = false
override fun connectAndSyncFiltersIfDisconnected(ignoreRetryDelays: Boolean) = Unit
override fun isConnected() = true
override fun sendOrConnectAndSync(cmd: Command) = Unit
override fun sendIfConnected(cmd: Command) = Unit
override fun disconnect() = Unit
}
private val fast = RelayUrlNormalizer.normalize("wss://fast.example.com")
@@ -190,6 +242,82 @@ class RelayProberFlowTest {
assertTrue(listOf("expiration", "5060") in template.tags.map { it.toList() })
}
// ------------------------------------------------------------------
// readWriteCheck — honest read + write measurements, nothing claimed
// ------------------------------------------------------------------
private suspend fun ScriptedClient.playReadThenWrite(
relay: NormalizedRelayUrl,
ok: Boolean?,
okMessage: String = "",
) {
delay(50)
listener!!.onEose(relay, null) // read phase answers
while (published == null) delay(10) // write phase begins
if (ok != null) answerOk(relay, ok, okMessage)
}
@Test
fun readWriteCheckMeasuresBothSides() =
runTest {
val client = ScriptedClient()
val signer = NostrSignerInternal(KeyPair())
var result: Map<NormalizedRelayUrl, RelayProber.ReadWriteVerdict>? = null
val check =
launch {
result = RelayProber(client).readWriteCheck(listOf(fast), signer, timeoutMs = 5_000)
}
launch { client.playReadThenWrite(fast, ok = true) }
check.join()
val verdict = result!![fast]!!
assertTrue(verdict.rttReadMs >= 0, "an answered read must be measured")
assertTrue(verdict.rttWriteMs >= 0, "an answered write must be measured")
assertEquals(true, verdict.writeAccepted)
assertEquals(20166, client.published!!.kind, "the write test must use the ephemeral probe event")
}
@Test
fun writeRejectionIsAnAnswerNotAFailure() =
runTest {
val client = ScriptedClient()
val signer = NostrSignerInternal(KeyPair())
var result: Map<NormalizedRelayUrl, RelayProber.ReadWriteVerdict>? = null
val check =
launch {
result = RelayProber(client).readWriteCheck(listOf(walled), signer, timeoutMs = 5_000)
}
launch { client.playReadThenWrite(walled, ok = false, okMessage = "pow: 28 bits needed") }
check.join()
val verdict = result!![walled]!!
assertEquals(false, verdict.writeAccepted, "OK false is a measured policy answer")
assertEquals("pow: 28 bits needed", verdict.writeMessage)
assertTrue(verdict.rttWriteMs >= 0, "a rejection is still a round trip")
}
@Test
fun silentWriteLeavesTheWriteSideUnobserved() =
runTest {
val client = ScriptedClient()
val signer = NostrSignerInternal(KeyPair())
var result: Map<NormalizedRelayUrl, RelayProber.ReadWriteVerdict>? = null
val check =
launch {
result = RelayProber(client).readWriteCheck(listOf(fast), signer, timeoutMs = 2_000)
}
launch { client.playReadThenWrite(fast, ok = null) }
check.join()
val verdict = result!![fast]!!
assertTrue(verdict.rttReadMs >= 0)
assertNull(verdict.writeAccepted, "silence is not evidence about the write path")
assertEquals(-1, verdict.rttWriteMs)
}
// ------------------------------------------------------------------
// toDiscoveryEventTemplate — only observed facts become tags
// ------------------------------------------------------------------
@@ -257,6 +385,40 @@ class RelayProberFlowTest {
assertNull(tagsOf(template).firstOrNull { it[0] == "R" })
}
@Test
fun readWriteResultsBecomeRttTags() {
val verdict = RelayProber.Verdict(fast, reachable = true, rttOpenMs = 100, rttEoseMs = 300, error = null)
val readWrite = RelayProber.ReadWriteVerdict(fast, rttReadMs = 40, rttWriteMs = 55, writeAccepted = true, writeMessage = "")
val tags = tagsOf(verdict.toDiscoveryEventTemplate(readWrite = readWrite))
assertTrue(listOf("rtt-read", "40") in tags)
assertTrue(listOf("rtt-write", "55") in tags)
}
@Test
fun unobservedReadWriteSidesStayUntagged() {
val verdict = RelayProber.Verdict(fast, reachable = true, rttOpenMs = 100, rttEoseMs = -1, error = null)
val readWrite = RelayProber.ReadWriteVerdict(fast, rttReadMs = -1, rttWriteMs = -1, writeAccepted = null, writeMessage = null)
val tags = tagsOf(verdict.toDiscoveryEventTemplate(readWrite = readWrite))
assertNull(tags.firstOrNull { it[0] == "rtt-read" })
assertNull(tags.firstOrNull { it[0] == "rtt-write" })
}
@Test
fun writeRejectionReasonsBecomeRequirementTags() {
val verdict = RelayProber.Verdict(walled, reachable = true, rttOpenMs = 100, rttEoseMs = -1, error = null)
val pow = RelayProber.ReadWriteVerdict(walled, -1, 30, writeAccepted = false, writeMessage = "pow: 28 bits needed")
assertTrue(listOf("R", "pow") in tagsOf(verdict.toDiscoveryEventTemplate(readWrite = pow)))
val auth = RelayProber.ReadWriteVerdict(walled, -1, 30, writeAccepted = false, writeMessage = "auth-required: sign in")
assertTrue(listOf("R", "auth") in tagsOf(verdict.toDiscoveryEventTemplate(readWrite = auth)))
val blocked = RelayProber.ReadWriteVerdict(walled, -1, 30, writeAccepted = false, writeMessage = "blocked: not welcome")
assertNull(tagsOf(verdict.toDiscoveryEventTemplate(readWrite = blocked)).firstOrNull { it[0] == "R" })
}
@Test
fun onionRelayIsTaggedTor() {
val onion = RelayUrlNormalizer.normalize("ws://someonionaddressabcdefghijklmnop.onion")