fix(relay): a hang-up before the OK is not the relay's answer

A relay that drops the socket between our EVENT frame and its OK has told
us nothing: the event may be stored, or it may not. We recorded that as
the relay's verdict and stopped waiting — even though the pool's own
outbox still owed the relay the event and would have flushed it on
reconnect. Nobody was listening by then, so the publish came back failed
and the event landed on the relay a second later anyway.

publishAndCollectResults now holds a transport failure as provisional for
one retry: it drops the tentative verdict, ignores the echoes of the same
drop, clears the backoff and dials, and takes the OK when the pool's
flush earns it. Everything happens inside the caller's existing timeout,
so no publish waits longer than it used to, and a relay that keeps
hanging up is still reported as a transport failure rather than a
success. transportRetries = 0 restores the old behaviour exactly.

Found through the Marmot interop harness, which was losing a message
every few runs to a loopback relay that was healthy a second later. The
same race is every publish that meets a network change on mobile.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
Claude
2026-09-09 09:49:25 +00:00
parent 7cfa0758d9
commit 7e187e39df
3 changed files with 296 additions and 6 deletions
+10 -6
View File
@@ -653,6 +653,16 @@ test we have.
that object's three VALUES, so every poll matched nothing and reported "never
received invite" for welcomes that had arrived and been accepted.
8. **A hang-up ended the publish wait.** A relay that dropped the socket
between our EVENT frame and its OK gave no verdict — the event may be stored,
it may not — and we recorded that as the relay's answer and stopped waiting.
The pool's own outbox would have re-sent on reconnect; nobody was listening
by then. `publishAndCollectResults` now keeps a transport failure provisional
for one retry, dials past the backoff, and takes the OK when it arrives; it
stays inside the caller's existing timeout, so nothing waits longer than
before. This was one lost message per few harness runs, and on mobile it is
every publish that races a network change.
### Harness defects (not ours)
- **Runs inherited each other's state.** wnd wipes B's and C's data dirs on
@@ -674,12 +684,6 @@ test we have.
### What is NOT done
- **The local relay drops a publish occasionally.** One run in several,
`amy` gets `disconnected before OK` from nostr-rs-relay and reports the send
as unconfirmed even though the event is on the relay a moment later. It is a
harness-relay flake, not a protocol failure, and it costs whichever test is
running at the time. Worth making the CLI's publish confirmation tolerate a
reconnect rather than papering over it in the tests.
- **Agent-text-stream publishes records but has nowhere to send them.** We
decode the `0x8006` policy, derive per-stream record keys, open records and
fold the transcript, we advertise the `0xF2D1` receive capability, and
@@ -69,6 +69,23 @@ class PublishResult(
}
}
/**
* How many times a relay that answered with a transport failure rather than an
* OK is re-sent to before the failure is reported. One retry covers the common
* case — a socket that dropped between our EVENT frame and the relay's OK —
* without turning a genuinely unreachable relay into a long stall, because the
* retries share the caller's existing publish timeout.
*/
const val DEFAULT_TRANSPORT_RETRIES = 1
/**
* Internal channel marker for "this relay is back up", so the wait loop — the
* one coroutine that owns the retry bookkeeping — can re-issue a send that a
* disconnected relay would have dropped. The NUL prefix keeps it out of reach
* of any real relay message, and it never surfaces in a [PublishResult].
*/
private const val RECONNECTED = "\u0000publish-retry-reconnected"
@OptIn(DelicateCoroutinesApi::class)
suspend fun INostrClient.publishAndConfirm(
event: Event,
@@ -107,6 +124,7 @@ suspend fun INostrClient.publishAndCollectResults(
event: Event,
relayList: Set<NormalizedRelayUrl>,
timeoutInSeconds: Long = 15,
transportRetries: Int = DEFAULT_TRANSPORT_RETRIES,
): Map<NormalizedRelayUrl, PublishResult> {
val resultChannel = Channel<DetailedResult>(UNLIMITED)
val mark = TimeSource.Monotonic.markNow()
@@ -132,6 +150,22 @@ suspend fun INostrClient.publishAndCollectResults(
}
}
/**
* A relay is only sendable once it is back up: publishing to a
* disconnected relay dials and drops the command, so a retry has to
* be re-issued from here rather than at the moment we noticed the
* hang-up.
*/
override fun onConnected(
relay: IRelayClient,
pingMillis: Int,
compressed: Boolean,
) {
if (relay.url in relayList) {
resultChannel.trySend(DetailedResult(relay.url, false, RECONNECTED))
}
}
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
@@ -165,18 +199,72 @@ suspend fun INostrClient.publishAndCollectResults(
val result =
async {
val receivedResults = mutableMapOf<NormalizedRelayUrl, PublishResult>()
// A relay that hung up or never connected gave no verdict on the
// event — it may have stored it, it may not. Re-send to that relay
// once (a Nostr event is idempotent under its own id, so the worst
// case is a duplicate the relay collapses) and keep waiting for the
// OK we were owed, instead of reporting a failed publish for a relay
// that is healthy a moment later. The retries live inside the
// caller's existing timeout, so nothing waits longer than before.
val retriesLeft = relayList.associateWith { transportRetries }.toMutableMap()
// The withTimeout block will cancel the coroutine if the loop takes too long
withTimeoutOrNull(timeoutInSeconds * 1000) {
val awaitingReconnect = mutableSetOf<NormalizedRelayUrl>()
while (receivedResults.size < relayList.size) {
val result = resultChannel.receive()
if (result.message == RECONNECTED) {
// The pool flushes what it still owes a relay as part of
// coming back up, so there is nothing to re-send here —
// this only reopens the relay to a fresh verdict.
awaitingReconnect.remove(result.relay)
continue
}
// One dropped socket can report itself more than once
// (the pool's disconnect and the relay client's both land
// here). While a relay is waiting to come back those are
// echoes of the drop we already answered, not new verdicts.
if (result.relay in awaitingReconnect) continue
val currentResult = receivedResults[result.relay]
// do not override a successful result.
if (currentResult == null || !currentResult.accepted) {
receivedResults[result.relay] = PublishResult(result.success, result.message, result.elapsedMs)
}
val recorded = receivedResults[result.relay]
if (recorded != null && recorded.isTransportFailure && (retriesLeft[result.relay] ?: 0) > 0) {
retriesLeft[result.relay] = retriesLeft.getValue(result.relay) - 1
// Drop the provisional verdict so the loop keeps waiting
// for this relay rather than treating the hang-up as its
// answer. If the retry also fails we record it again and
// report the transport failure as before.
receivedResults.remove(result.relay)
awaitingReconnect.add(result.relay)
Log.d("publishAndConfirm") {
"Retrying ${event.id} on ${result.relay} after ${recorded.message}"
}
// The event is still in the pool's outbox for this relay,
// so the dial is the whole job: the pool flushes what it
// owes the relay once the socket is back. Ignore the
// accumulated backoff — this is a user-visible publish
// waiting on it, not a background refresh.
resetBackoff()
reconnect(onlyIfChanged = false, ignoreRetryDelays = true)
}
}
}
// A relay whose last word was a transport failure and whose retry
// never came back inside the timeout still has to be reported: the
// caller promised a verdict for every listed relay, and "we retried"
// is not one.
for (relay in relayList) {
if (relay !in receivedResults && retriesLeft.getValue(relay) < transportRetries) {
receivedResults[relay] = PublishResult(false, PublishResult.DISCONNECTED)
}
}
receivedResults
}
@@ -0,0 +1,198 @@
/*
* 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.quartz.nip01Core.relay
import com.vitorpamplona.geode.InProcessRelays
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.PublishResult
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndCollectResults
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
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.signers.NostrSignerInternal
import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import java.util.concurrent.atomic.AtomicInteger
import kotlin.test.AfterTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
/**
* A relay that hangs up between our EVENT frame and its OK told us nothing:
* the event may be stored, or it may not. We reported it as a failed publish
* and never tried again, which is how the Marmot interop harness kept losing
* a message a run to `disconnected before OK` on a loopback relay that was
* perfectly healthy a second later.
*
* A Nostr event is idempotent under its own id, so re-sending it after a
* transport failure costs a duplicate the relay collapses and buys the OK we
* were owed. The retry stays inside the caller's existing timeout, so nothing
* waits longer than it used to.
*/
class PublishRetriesTransportFailureTest {
private val hub = InProcessRelays()
private val scope = CoroutineScope(Dispatchers.Default + SupervisorJob())
@AfterTest
fun tearDown() {
scope.cancel()
hub.close()
}
/**
* Wraps the in-process hub and hangs up the first [dropFirstConnections]
* sockets the moment they carry an `EVENT` frame — the relay took our
* bytes and vanished before answering, which is the case that has no
* verdict in it.
*/
private class HangsUpOnFirstEvent(
private val delegate: WebsocketBuilder,
private val dropFirstConnections: Int,
) : WebsocketBuilder {
val eventFramesSeen = AtomicInteger(0)
private val socketsBuilt = AtomicInteger(0)
override fun build(
url: NormalizedRelayUrl,
out: WebSocketListener,
): WebSocket {
val index = socketsBuilt.getAndIncrement()
val inner = delegate.build(url, out)
val hangUp = index < dropFirstConnections
return object : WebSocket by inner {
override fun send(msg: String): Boolean {
if (!msg.startsWith("[\"EVENT\"")) return inner.send(msg)
eventFramesSeen.incrementAndGet()
if (!hangUp) return inner.send(msg)
// Take the bytes, answer nothing, drop the socket.
inner.disconnect()
out.onClosed(1006, "abnormal closure")
return true
}
}
}
}
@Test
fun aDisconnectBeforeTheOkIsRetriedAndSucceeds() =
runBlocking {
val builder = HangsUpOnFirstEvent(hub, dropFirstConnections = 1)
val client = NostrClient(builder, scope)
val event = NostrSignerInternal(KeyPair()).sign(TextNoteEvent.build("survives a hang-up"))
val results =
client.publishAndCollectResults(
event = event,
relayList = setOf(InProcessRelays.DEFAULT_URL),
timeoutInSeconds = 20,
)
assertEquals(1, results.size)
val result = results.values.single()
assertTrue(
result.accepted,
"the retry must land the event: the relay was healthy, it just hung up before the OK " +
"(got \"${result.message}\")",
)
assertTrue(
builder.eventFramesSeen.get() >= 2,
"the event has to actually go out a second time, not just be re-reported",
)
client.disconnect()
}
@Test
fun aRelayThatKeepsHangingUpStillReportsTheTransportFailure() =
runBlocking {
// Every socket dies the same way, so no retry can help. The result
// must still name the transport failure rather than claim success
// or hide the relay.
val builder = HangsUpOnFirstEvent(hub, dropFirstConnections = Int.MAX_VALUE)
val client = NostrClient(builder, scope)
val event = NostrSignerInternal(KeyPair()).sign(TextNoteEvent.build("never lands"))
val results =
client.publishAndCollectResults(
event = event,
relayList = setOf(InProcessRelays.DEFAULT_URL),
timeoutInSeconds = 8,
)
val result = results.getValue(InProcessRelays.DEFAULT_URL)
assertFalse(result.accepted)
assertTrue(
result.isTransportFailure,
"a hang-up is never a verdict from the relay (got \"${result.message}\")",
)
client.disconnect()
}
@Test
fun aHealthyPublishStillTakesOneRoundTrip() =
runBlocking {
val builder = HangsUpOnFirstEvent(hub, dropFirstConnections = 0)
val client = NostrClient(builder, scope)
val event = NostrSignerInternal(KeyPair()).sign(TextNoteEvent.build("no retry needed"))
val results =
client.publishAndCollectResults(
event = event,
relayList = setOf(InProcessRelays.DEFAULT_URL),
timeoutInSeconds = 20,
)
assertTrue(results.values.single().accepted)
assertEquals(
1,
builder.eventFramesSeen.get(),
"an OK on the first try must not be followed by a speculative resend",
)
client.disconnect()
}
@Test
fun retriesCanBeTurnedOff() =
runBlocking {
val builder = HangsUpOnFirstEvent(hub, dropFirstConnections = 1)
val client = NostrClient(builder, scope)
val event = NostrSignerInternal(KeyPair()).sign(TextNoteEvent.build("one shot only"))
val results =
client.publishAndCollectResults(
event = event,
relayList = setOf(InProcessRelays.DEFAULT_URL),
timeoutInSeconds = 8,
transportRetries = 0,
)
assertEquals(PublishResult.DISCONNECTED, results.values.single().message)
assertEquals(1, builder.eventFramesSeen.get())
client.disconnect()
}
}