diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientPublishExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientPublishExt.kt index 9cbec380b4..112345d123 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientPublishExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientPublishExt.kt @@ -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 { val resultChannel = Channel(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, ) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProber.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProber.kt index 58663c2728..86e57c7e58 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProber.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProber.kt @@ -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, 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, + signer: NostrSigner, + timeoutMs: Long = 15_000, + waveSize: Int = 1000, + readLimit: Int = 1, + ): Map { + val out = HashMap() + val distinct = relays.toSet() + for (wave in distinct.chunked(waveSize.coerceAtLeast(1))) { + val reads = HashMap() + 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, 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 = +fun RelayProber.Verdict.toDiscoveryEventTemplate( + createdAt: Long = TimeUtils.now(), + readWrite: RelayProber.ReadWriteVerdict? = null, +): EventTemplate = 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") } diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProberFlowTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProberFlowTest.kt index 2efc801324..0195e52b53 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProberFlowTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayProberFlowTest.kt @@ -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>? = null + var published: Event? = null + val connListeners = mutableListOf() override fun subscribe( subId: String, @@ -59,6 +68,49 @@ class RelayProberFlowTest { this.listener = listener this.sentFilters = filters } + + override fun publish( + event: Event, + relayList: Set, + ) { + 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? = 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? = 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? = 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")