From 1093b3ce8aa0a1e2f7be6bd06522e0a2baaaa7f9 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 12 Sep 2026 20:17:29 +0000 Subject: [PATCH 1/7] fix(notifications): count connected relays from the pool's live socket state The always-on notification read its "Connected to N relays" from RelayPool.connectedRelays, a set that only moves on the onConnected / onDisconnected callbacks. That set over-reports: a relay that sent a WebSocket CLOSE frame never produces a callback (the app does not answer onClosing, so OkHttp fires neither onClosed nor onFailure, and the later cancel() is silent too), so the URL lingers until the 120s ping path finally fails. After the feeds tore down in the background this left hundreds of already-dropped relays in the count for minutes, with no subscription in the "show details" breakdown to justify any of them. Expose the pool's ground truth instead: RelayPool.connectedRelayUrls() reads each member's isConnected(), surfaced as INostrClient.connectedRelays() (defaulting to the flow's value for pool-less clients). The notification keeps the flows only as a trigger, merging in availableRelaysFlow because that is the flow that moves when the pool drops such a relay, and re-reads the live count on each sample. The "show details" breakdown and the Active Subscriptions screen read the same source so all three agree. A pool test pins the drift: with a socket layer that never confirms the close, removing a relay leaves it in the flow but out of the snapshot. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_014cq6vrfQkASwxgqXXY4py8 --- .../notifications/NotificationRelayService.kt | 34 ++++++--- .../notifications/RelayPurposeSummary.kt | 5 +- .../ActiveSubscriptionsViewModel.kt | 3 +- .../nip01Core/relay/client/INostrClient.kt | 9 +++ .../nip01Core/relay/client/NostrClient.kt | 2 + .../nip01Core/relay/client/pool/RelayPool.kt | 20 ++++++ .../pool/RelayPoolConnectedSnapshotTest.kt | 72 +++++++++++++++++++ 7 files changed, 133 insertions(+), 12 deletions(-) create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt index 74870a5645..78817e7574 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt @@ -50,6 +50,7 @@ import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.flow.collectLatest +import kotlinx.coroutines.flow.merge import kotlinx.coroutines.flow.sample import kotlinx.coroutines.launch @@ -299,7 +300,8 @@ class NotificationRelayService : Service() { * Tor, network changes). Without this, the client disconnects 30s after * the UI stops collecting. * - * 2. connectedRelaysFlow: Updates the persistent notification with relay count. + * 2. connectedRelaysFlow + availableRelaysFlow: re-read the pool's live relay count + * (client.connectedRelays()) and refresh the persistent notification. * * The service does NOT create its own relay subscriptions. Instead, it relies on * the AccountFilterAssembler subscription that lives in the Compose tree (LoggedInPage). @@ -320,17 +322,29 @@ class NotificationRelayService : Service() { } launch { + // The two flows are only the *trigger*; the number comes from + // client.connectedRelays(), which reads each pool member's live socket + // state. connectedRelaysFlow() alone over-reported: it is maintained by + // OkHttp callbacks, and a relay that sent a WebSocket CLOSE frame never + // produces one (the app doesn't answer onClosing, so neither onClosed nor + // onFailure fires, and the later cancel() is silent too). After the feeds + // tore down in the background that left hundreds of already-closed relays + // in the flow for minutes, with no subscription to justify a single one of + // them. availableRelaysFlow() is merged in because that is the flow that + // moves when the pool drops such a relay -- the connected flow, by + // definition, doesn't. + // // sample() caps how often we touch the notification. During feed - // load/teardown connectedRelaysFlow churns dozens of times per second; - // posting on every delta blows past Android's notification rate limit - // (~10/s), which silently drops updates and leaves the visible count - // stuck on a stale intermediate value. One refresh per second stays - // well under the limit and always lands the settled count. - Amethyst.instance.client - .connectedRelaysFlow() + // load/teardown these flows churn dozens of times per second; posting on + // every delta blows past Android's notification rate limit (~10/s), which + // silently drops updates and leaves the visible count stuck on a stale + // intermediate value. One refresh per second stays well under the limit + // and always lands the settled count. + val client = Amethyst.instance.client + merge(client.connectedRelaysFlow(), client.availableRelaysFlow()) .sample(NOTIFICATION_REFRESH_MS) - .collectLatest { relays -> - val count = relays.size + .collectLatest { + val count = client.connectedRelays().size if (count != connectedRelayCount) { connectedRelayCount = count updateNotification(count) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt index 4039025be3..c205734df0 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt @@ -56,7 +56,10 @@ object RelayPurposeSummary { val named = mutableMapOf>() val browsing = mutableSetOf() - client.connectedRelaysFlow().value.forEach { relay -> + // Same source as the count above it (see NotificationRelayService): the pool's live socket + // state, not the callback-maintained flow, so the breakdown never has to explain relays the + // pool has already dropped. + client.connectedRelays().forEach { relay -> client .activeRequests(relay) .values diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt index a5db379f0b..352ca674d6 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt @@ -156,7 +156,8 @@ class ActiveSubscriptionsViewModel : ViewModel() { withContext(Dispatchers.Default) { val client = Amethyst.instance.client aggregateSubscriptions( - client.connectedRelaysFlow().value.associateWith { relay -> + // The pool's live socket state, same as the always-on notification this screen explains. + client.connectedRelays().associateWith { relay -> client.activeRequests(relay).values.flatten() }, ) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt index b77d033c7d..72d9e164f9 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt @@ -37,6 +37,15 @@ interface INostrClient : AutoCloseable { fun availableRelaysFlow(): StateFlow> + /** + * The relays whose socket is up at the moment of the call, read from the pool itself rather than + * from the callback-maintained [connectedRelaysFlow]. The flow is for reacting to changes; this is + * for reporting a count, because the flow keeps a relay whose socket died without a callback + * (see `RelayPool.connectedRelayUrls`) until the ping timeout catches up. Defaults to the flow's + * current value for clients that have no pool behind them. + */ + fun connectedRelays(): Set = connectedRelaysFlow().value + fun connect() fun disconnect() diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt index 21c16d9b64..63f3d0d3ce 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt @@ -493,6 +493,8 @@ class NostrClient( override fun connectedRelaysFlow() = relayPool.connectedRelays + override fun connectedRelays() = relayPool.connectedRelayUrls() + override fun availableRelaysFlow() = relayPool.availableRelays override fun close() { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt index 62fe531d86..99ad6aeb6f 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt @@ -267,4 +267,24 @@ class RelayPool( ) = listener.onSent(relay, cmdStr, cmd, success) fun connectedRelaysCount(): Int = relays.count { url, relay -> relay.isConnected() } + + /** + * The relays whose socket is up *right now*, read from each pool member's [IRelayClient.isConnected]. + * + * This is the ground truth; [connectedRelays] is a projection of it that only moves on the + * [onConnected] / [onDisconnected] callbacks. The two drift whenever a socket dies without a + * callback: a relay that sent a WebSocket CLOSE frame leaves OkHttp waiting for a reply that never + * comes, so neither `onClosed` nor `onFailure` fires, and a later `cancel()` is silent too. The URL + * then sits in [connectedRelays] until the 120s ping path finally fails, while this snapshot has + * already dropped it (a removed relay is no longer a member; a still-desired one reads + * `isConnected() == false` as soon as its socket closes). Use this for anything a person reads as + * "how many relays am I connected to"; keep the flow for change notification. + */ + fun connectedRelayUrls(): Set { + val urls = mutableSetOf() + relays.forEach { url, relay -> + if (relay.isConnected()) urls.add(url) + } + return urls + } } diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt new file mode 100644 index 0000000000..b1bb89789d --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt @@ -0,0 +1,72 @@ +/* + * 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.client.pool + +import com.vitorpamplona.quartz.nip01Core.relay.client.single.basic.FakeWebsocketBuilder +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlin.test.Test +import kotlin.test.assertEquals + +/** + * [RelayPool.connectedRelayUrls] must report what the pool actually holds open, even when the + * socket layer never confirms a close. + * + * [com.vitorpamplona.quartz.nip01Core.relay.client.single.basic.FakeWebSocket.disconnect] is a + * no-op that never calls back, which is exactly what OkHttp does after a relay sent a WebSocket + * CLOSE frame the app never answered: `cancel()` then fires neither `onClosed` nor `onFailure`. + * The callback-driven [RelayPool.connectedRelays] flow keeps such a relay for minutes; the + * snapshot drops it the moment the pool lets go of it. + */ +class RelayPoolConnectedSnapshotTest { + private val url = NormalizedRelayUrl("wss://relay.example.com/") + + @Test + fun snapshotMatchesTheFlowWhileTheSocketIsOpen() { + val sockets = FakeWebsocketBuilder() + val pool = RelayPool(sockets) + + pool.getOrCreateRelay(url).connect() + sockets.lastListener.onOpen(pingMillis = 10, compression = false) + + assertEquals(setOf(url), pool.connectedRelays.value) + assertEquals(setOf(url), pool.connectedRelayUrls()) + assertEquals(1, pool.connectedRelaysCount()) + } + + @Test + fun removedRelayLeavesTheSnapshotEvenWhenTheCloseIsSilent() { + val sockets = FakeWebsocketBuilder() + val pool = RelayPool(sockets) + + pool.getOrCreateRelay(url).connect() + sockets.lastListener.onOpen(pingMillis = 10, compression = false) + + // The pool drops the relay (no subscription wants it anymore) and cancels its socket, + // but the socket layer stays silent. + pool.removeRelay(url) + + // Documents the drift this test exists for: the flow still carries the relay... + assertEquals(setOf(url), pool.connectedRelays.value) + // ...while the snapshot already reflects that nothing is held open. + assertEquals(emptySet(), pool.connectedRelayUrls()) + assertEquals(0, pool.connectedRelaysCount()) + } +} From b8135e76c4ab8cd3de467b68636f3c31c27f5a61 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 12 Sep 2026 22:52:37 +0000 Subject: [PATCH 2/7] fix(relay): answer a relay's WebSocket CLOSE frame so OkHttp can finish the handshake OkHttp fires onClosed only once both peers have sent a CLOSE frame, and sending ours is the application's job (WebSocketListener KDoc; RealWebSocket emits onClosed solely from the writer once our Close is dequeued with the peer's code already set; the bundled WebSocketEcho recipe answers onClosing with close(1000, null)). BasicOkHttpWebSocket never implemented onClosing, so a relay-initiated close left the socket half-closed: no onClosed, no onFailure, send() still accepted and silently discarded, and a later cancel() silent as well because no reader was left to fail. The relay client kept believing it was connected, with its REQs live, until the 120s ping path finally failed up to two intervals later. Answer onClosing with close(1000, null). Verified against OkHttp 5.5.0: onClosed then fires at once whether the relay still holds the TCP session or has already dropped it, and the existing onClosed path in BasicRelayClient marks the connection closed and lets the pool reconnect under its normal backoff. Always 1000 rather than echoing the relay's code, since close() validates the code it writes and relays may send reserved ones. The test drives the wrapper against a minimal RFC 6455 server on a loopback ServerSocket (no new dependency): the relay sends CLOSE, the client must answer with its own CLOSE frame and report onClosed with the relay's code. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_014cq6vrfQkASwxgqXXY4py8 --- .../sockets/okhttp/BasicOkHttpWebSocket.kt | 22 ++ .../BasicOkHttpWebSocketCloseHandshakeTest.kt | 216 ++++++++++++++++++ 2 files changed, 238 insertions(+) create mode 100644 quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt index f0a9c61cec..80b2679e73 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt @@ -96,6 +96,28 @@ class BasicOkHttpWebSocket( incomingMessages.trySendBlocking(text) } + override fun onClosing( + webSocket: OkHttpWebSocket, + code: Int, + reason: String, + ) { + // The relay sent a CLOSE frame. OkHttp's contract (WebSocketListener KDoc, + // RealWebSocket, and its own WebSocketEcho recipe) is that onClosed fires + // only once BOTH peers have sent a close, and sending ours is the + // application's job. Left unanswered, the socket sits half-closed: no + // onClosed, no onFailure, send() still accepted and silently discarded, and + // a later cancel() is silent too -- so the relay client kept believing it + // was connected, with its REQs live, until OkHttp's 120s ping path finally + // failed up to two intervals later. Answering completes the handshake and + // OkHttp reports onClosed at once, whether or not the relay still holds the + // TCP session open. + // + // Always 1000 rather than echoing `code`: close() validates the code it is + // asked to write and throws on the reserved ones (1005, 1006, 1015), and a + // relay may send anything. + webSocket.close(1000, null) + } + override fun onClosed( webSocket: OkHttpWebSocket, code: Int, diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt new file mode 100644 index 0000000000..56c525b946 --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt @@ -0,0 +1,216 @@ +/* + * 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.sockets.okhttp + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener +import okhttp3.OkHttpClient +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Test +import java.io.InputStream +import java.io.OutputStream +import java.net.ServerSocket +import java.net.Socket +import java.security.MessageDigest +import java.util.Base64 +import java.util.concurrent.CountDownLatch +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicInteger +import java.util.concurrent.atomic.AtomicReference +import kotlin.concurrent.thread + +/** + * A relay-initiated close must be answered, or OkHttp never finishes the handshake. + * + * OkHttp fires `onClosed` only once BOTH peers have sent a CLOSE frame, and sending ours is the + * application's job (its `WebSocketEcho` recipe answers `onClosing` with `close(1000, null)`). + * Before [BasicOkHttpWebSocket] did that, a relay's CLOSE frame left the socket half-closed: no + * `onClosed`, no `onFailure`, `send()` still accepted and discarded, and a later `cancel()` silent + * too. The relay client kept believing it was connected, with its REQs live, until the ping path + * failed up to two ping intervals later. + * + * Driven against a minimal RFC 6455 server on a loopback [ServerSocket] rather than a mock + * server library, so the test needs no new dependency and controls the exact frames on the wire. + */ +class BasicOkHttpWebSocketCloseHandshakeTest { + /** Handshakes one client, sends it a CLOSE frame on demand, and records the frames it sends back. */ + private class TinyRelay : AutoCloseable { + private val server = ServerSocket(0) + val url = NormalizedRelayUrl("ws://127.0.0.1:${server.localPort}/") + + private val handshaken = CountDownLatch(1) + val clientCloseFrame = CountDownLatch(1) + val clientCloseCode = AtomicInteger(-1) + + private var socket: Socket? = null + private var out: OutputStream? = null + + private val thread = + thread(isDaemon = true, name = "tiny-relay") { + runCatching { + val s = server.accept() + socket = s + val input = s.getInputStream() + val output = s.getOutputStream() + out = output + handshake(input, output) + handshaken.countDown() + readFrames(input) + } + } + + private fun handshake( + input: InputStream, + output: OutputStream, + ) { + var key: String? = null + val line = StringBuilder() + while (true) { + val c = input.read() + check(c != -1) { "EOF during handshake" } + if (c == '\n'.code) { + val l = line.toString().trim() + if (l.isEmpty()) break + if (l.lowercase().startsWith("sec-websocket-key:")) key = l.substring(18).trim() + line.setLength(0) + } else if (c != '\r'.code) { + line.append(c.toChar()) + } + } + val accept = + Base64.getEncoder().encodeToString( + MessageDigest.getInstance("SHA-1").digest((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").toByteArray()), + ) + output.write( + ( + "HTTP/1.1 101 Switching Protocols\r\n" + + "Upgrade: websocket\r\n" + + "Connection: Upgrade\r\n" + + "Sec-WebSocket-Accept: $accept\r\n\r\n" + ).toByteArray(Charsets.ISO_8859_1), + ) + output.flush() + } + + /** Client frames are masked; decode enough to spot a CLOSE and read its status code. */ + private fun readFrames(input: InputStream) { + while (true) { + val b0 = input.read() + if (b0 == -1) return + val b1 = input.read() + if (b1 == -1) return + val opcode = b0 and 0x0F + var len = b1 and 0x7F + if (len == 126) { + len = (input.read() shl 8) or input.read() + } else if (len == 127) { + len = 0 + repeat(8) { len = (len shl 8) or input.read() } + } + val masked = (b1 and 0x80) != 0 + val mask = if (masked) ByteArray(4) { input.read().toByte() } else ByteArray(4) + val payload = ByteArray(len) { i -> (input.read() xor mask[i % 4].toInt()).toByte() } + if (opcode == 0x8) { + if (len >= 2) { + clientCloseCode.set(((payload[0].toInt() and 0xFF) shl 8) or (payload[1].toInt() and 0xFF)) + } + clientCloseFrame.countDown() + } + } + } + + fun awaitClient() = handshaken.await(5, TimeUnit.SECONDS) + + /** Server-initiated CLOSE, status 1000, unmasked as servers send it. The TCP session stays open. */ + fun sendClose() { + val output = checkNotNull(out) { "no client yet" } + output.write(byteArrayOf(0x88.toByte(), 0x02, 0x03, 0xE8.toByte())) + output.flush() + } + + override fun close() { + runCatching { socket?.close() } + runCatching { server.close() } + thread.join(2_000) + } + } + + private class Recorder : WebSocketListener { + val opened = CountDownLatch(1) + val closed = CountDownLatch(1) + val closedCode = AtomicInteger(-1) + val failure = AtomicReference(null) + + override fun onOpen( + pingMillis: Int, + compression: Boolean, + ) = opened.countDown() + + override suspend fun onMessage(text: String) {} + + override fun onClosed( + code: Int, + reason: String, + ) { + closedCode.set(code) + closed.countDown() + } + + override fun onFailure( + t: Throwable, + code: Int?, + response: String?, + ) { + failure.set(t) + } + } + + @Test + fun `a relay initiated close is answered and reported as closed`() { + TinyRelay().use { relay -> + val recorder = Recorder() + val client = OkHttpClient() + val socket = BasicOkHttpWebSocket(relay.url, { client }, recorder) + + socket.connect() + assertTrue("relay never saw the client", relay.awaitClient()) + assertTrue("no onOpen", recorder.opened.await(5, TimeUnit.SECONDS)) + + relay.sendClose() + + // The half of the handshake that is ours to send. + assertTrue("client never answered the relay's CLOSE frame", relay.clientCloseFrame.await(5, TimeUnit.SECONDS)) + assertEquals(1000, relay.clientCloseCode.get()) + + // And the terminal callback the relay client's bookkeeping depends on. + assertTrue("onClosed never fired", recorder.closed.await(5, TimeUnit.SECONDS)) + assertEquals("the relay's status code is what gets reported", 1000, recorder.closedCode.get()) + assertNull("a clean handshake is not a failure", recorder.failure.get()) + + // The usual teardown afterwards must stay a harmless no-op. + socket.disconnect() + + client.dispatcher.executorService.shutdown() + } + } +} From d108ba0cc87bfafead9686c63de9dfbd3074d22f Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 12 Sep 2026 23:28:47 +0000 Subject: [PATCH 3/7] fix(relay): answer CLOSE in the Android socket too, and stop trusting socket callbacks for pool bookkeeping Follow-ups from an audit of the two previous commits on this branch. The Android app does not use quartz's BasicOkHttpWebSocket: AppModules wires its own OkHttpWebSocket, a near-twin that decides needsReconnect() from the OkHttpClient in use. The onClosing answer therefore only reached desktop, CLI and geode. Mirror it here, with the same loopback RFC 6455 test. The unit-test android.util.Log stub gets a primitive-signature isLoggable (the boxed one it had was a different method to the JVM, which is why OkHttp could never be built in this module's tests) plus println, so the socket can be exercised for real. RelayPool now clears its connected flow itself on every path where it lets a relay go -- removeRelay, removeAllRelays and disconnect -- instead of waiting for a callback the socket layer may never deliver. The previous commit added a members-read snapshot and migrated three callers; the other consumers of connectedRelaysFlow (drawer status, connection-time accounting, "wait until connected" loops) were still reading the stale set. BasicRelayClient retires its listener before tearing a socket down and ignores anything a retired socket reports afterwards. OkHttp delivers onClosed from its writer thread and a cancelled socket's failure later still; with the pool rebuilding sessions via disconnect()+connect(), a late callback from the old socket could null the new one, orphaning a live connection and dialing a third. disconnect() also reports onDisconnected itself now, so a client-initiated teardown never depends on the socket layer confirming it. KDoc and comments that described the unanswered-close behaviour in the present tense, or claimed a still-desired relay reads isConnected()=false on a silent close, are corrected. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_014cq6vrfQkASwxgqXXY4py8 --- .../notifications/NotificationRelayService.kt | 18 +- .../service/okhttp/OkHttpWebSocket.kt | 18 ++ amethyst/src/test/java/android/util/Log.java | 13 +- .../OkHttpWebSocketCloseHandshakeTest.kt | 211 ++++++++++++++++++ .../nip01Core/relay/client/INostrClient.kt | 11 +- .../nip01Core/relay/client/pool/RelayPool.kt | 21 +- .../client/single/basic/BasicRelayClient.kt | 37 ++- .../pool/RelayPoolConnectedSnapshotTest.kt | 51 +++-- .../basic/BasicRelayClientStaleSocketTest.kt | 123 ++++++++++ 9 files changed, 458 insertions(+), 45 deletions(-) create mode 100644 amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt index 78817e7574..d7ee24f75a 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt @@ -324,15 +324,15 @@ class NotificationRelayService : Service() { launch { // The two flows are only the *trigger*; the number comes from // client.connectedRelays(), which reads each pool member's live socket - // state. connectedRelaysFlow() alone over-reported: it is maintained by - // OkHttp callbacks, and a relay that sent a WebSocket CLOSE frame never - // produces one (the app doesn't answer onClosing, so neither onClosed nor - // onFailure fires, and the later cancel() is silent too). After the feeds - // tore down in the background that left hundreds of already-closed relays - // in the flow for minutes, with no subscription to justify a single one of - // them. availableRelaysFlow() is merged in because that is the flow that - // moves when the pool drops such a relay -- the connected flow, by - // definition, doesn't. + // state. connectedRelaysFlow() alone used to over-report: it is fed by + // socket callbacks, and until the OkHttp sockets answered a relay's CLOSE + // frame a relay-initiated close produced none (no onClosed, no onFailure, + // and a silent cancel() afterwards). After the feeds tore down in the + // background that left hundreds of already-dropped relays in the flow for + // minutes, with no subscription to justify a single one of them. Reading + // the members directly cannot be fooled that way, and availableRelaysFlow() + // is merged in because that is the flow that moves when the pool drops a + // relay. // // sample() caps how often we touch the notification. During feed // load/teardown these flows churn dozens of times per second; posting on diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt index 4c70c12498..55b5994060 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt @@ -109,6 +109,24 @@ class OkHttpWebSocket( incomingMessages.trySendBlocking(text) } + override fun onClosing( + webSocket: okhttp3.WebSocket, + code: Int, + reason: String, + ) { + // The relay sent a CLOSE frame. OkHttp fires onClosed only once BOTH peers have sent + // one, and sending ours is the application's job (WebSocketListener KDoc; its own + // WebSocketEcho recipe does exactly this). Unanswered, the socket sat half-closed: + // no onClosed, no onFailure, send() still accepted and discarded, a later cancel() + // silent too -- so the relay client believed it was connected until the 120s ping + // path failed up to two intervals later. Mirrors BasicOkHttpWebSocket in quartz; + // the two classes differ only in how needsReconnect() is decided. + // + // Always 1000 rather than echoing `code`: close() validates the code it writes and + // throws on the reserved ones (1005, 1006, 1015), and a relay may send anything. + webSocket.close(1000, null) + } + override fun onClosed( webSocket: okhttp3.WebSocket, code: Int, diff --git a/amethyst/src/test/java/android/util/Log.java b/amethyst/src/test/java/android/util/Log.java index 245a440a12..4272c9bdd5 100644 --- a/amethyst/src/test/java/android/util/Log.java +++ b/amethyst/src/test/java/android/util/Log.java @@ -1,8 +1,17 @@ package android.util; public class Log { - public static Boolean isLoggable(String tag, Integer msg) { - return true; + // Primitive signature on purpose: OkHttp's Android platform probe (AndroidLog.enableLogging) + // links against `boolean isLoggable(String, int)`, and a boxed variant is a different method. + // Answering false keeps OkHttp from installing its Android log handler, which would route + // every internal task-runner trace through println() below on the dispatcher threads. + public static boolean isLoggable(String tag, int level) { + return false; + } + + public static int println(int priority, String tag, String msg) { + System.out.println(tag + ": " + msg); + return 0; } public static int d(String tag, String msg) { diff --git a/amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt b/amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt new file mode 100644 index 0000000000..cddd2f9f6f --- /dev/null +++ b/amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt @@ -0,0 +1,211 @@ +/* + * 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.amethyst.service.okhttp + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener +import okhttp3.OkHttpClient +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Test +import java.io.InputStream +import java.io.OutputStream +import java.net.ServerSocket +import java.net.Socket +import java.security.MessageDigest +import java.util.Base64 +import java.util.concurrent.CountDownLatch +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicInteger +import java.util.concurrent.atomic.AtomicReference +import kotlin.concurrent.thread + +/** + * The Android app's socket must answer a relay-initiated close, like quartz's + * `BasicOkHttpWebSocket` (see `BasicOkHttpWebSocketCloseHandshakeTest` there for the full story). + * OkHttp fires `onClosed` only once BOTH peers have sent a CLOSE frame; until [OkHttpWebSocket] + * sent ours, a relay's close left the socket half-closed and the relay client believed it was + * connected until the ping path failed minutes later. + * + * Same minimal RFC 6455 loopback server as the quartz test: the two socket classes live in + * different modules with no shared test fixtures, and the harness is small enough to carry twice. + */ +class OkHttpWebSocketCloseHandshakeTest { + private class TinyRelay : AutoCloseable { + private val server = ServerSocket(0) + val url = NormalizedRelayUrl("ws://127.0.0.1:${server.localPort}/") + + private val handshaken = CountDownLatch(1) + val clientCloseFrame = CountDownLatch(1) + val clientCloseCode = AtomicInteger(-1) + + private var socket: Socket? = null + private var out: OutputStream? = null + + private val thread = + thread(isDaemon = true, name = "tiny-relay") { + runCatching { + val s = server.accept() + socket = s + val input = s.getInputStream() + val output = s.getOutputStream() + out = output + handshake(input, output) + handshaken.countDown() + readFrames(input) + } + } + + private fun handshake( + input: InputStream, + output: OutputStream, + ) { + var key: String? = null + val line = StringBuilder() + while (true) { + val c = input.read() + check(c != -1) { "EOF during handshake" } + if (c == '\n'.code) { + val l = line.toString().trim() + if (l.isEmpty()) break + if (l.lowercase().startsWith("sec-websocket-key:")) key = l.substring(18).trim() + line.setLength(0) + } else if (c != '\r'.code) { + line.append(c.toChar()) + } + } + val accept = + Base64.getEncoder().encodeToString( + MessageDigest.getInstance("SHA-1").digest((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").toByteArray()), + ) + output.write( + ( + "HTTP/1.1 101 Switching Protocols\r\n" + + "Upgrade: websocket\r\n" + + "Connection: Upgrade\r\n" + + "Sec-WebSocket-Accept: $accept\r\n\r\n" + ).toByteArray(Charsets.ISO_8859_1), + ) + output.flush() + } + + /** Client frames are masked; decode enough to spot a CLOSE and read its status code. */ + private fun readFrames(input: InputStream) { + while (true) { + val b0 = input.read() + if (b0 == -1) return + val b1 = input.read() + if (b1 == -1) return + val opcode = b0 and 0x0F + var len = b1 and 0x7F + if (len == 126) { + len = (input.read() shl 8) or input.read() + } else if (len == 127) { + len = 0 + repeat(8) { len = (len shl 8) or input.read() } + } + val masked = (b1 and 0x80) != 0 + val mask = if (masked) ByteArray(4) { input.read().toByte() } else ByteArray(4) + val payload = ByteArray(len) { i -> (input.read() xor mask[i % 4].toInt()).toByte() } + if (opcode == 0x8) { + if (len >= 2) { + clientCloseCode.set(((payload[0].toInt() and 0xFF) shl 8) or (payload[1].toInt() and 0xFF)) + } + clientCloseFrame.countDown() + } + } + } + + fun awaitClient() = handshaken.await(5, TimeUnit.SECONDS) + + /** Server-initiated CLOSE, status 1000, unmasked as servers send it. The TCP session stays open. */ + fun sendClose() { + val output = checkNotNull(out) { "no client yet" } + output.write(byteArrayOf(0x88.toByte(), 0x02, 0x03, 0xE8.toByte())) + output.flush() + } + + override fun close() { + runCatching { socket?.close() } + runCatching { server.close() } + thread.join(2_000) + } + } + + private class Recorder : WebSocketListener { + val opened = CountDownLatch(1) + val closed = CountDownLatch(1) + val closedCode = AtomicInteger(-1) + val failure = AtomicReference(null) + + override fun onOpen( + pingMillis: Int, + compression: Boolean, + ) = opened.countDown() + + override suspend fun onMessage(text: String) {} + + override fun onClosed( + code: Int, + reason: String, + ) { + closedCode.set(code) + closed.countDown() + } + + override fun onFailure( + t: Throwable, + code: Int?, + response: String?, + ) { + failure.set(t) + } + } + + @Test + fun `a relay initiated close is answered and reported as closed`() { + TinyRelay().use { relay -> + val recorder = Recorder() + val client = OkHttpClient() + val socket = OkHttpWebSocket(relay.url, { client }, recorder) + + socket.connect() + assertTrue("relay never saw the client", relay.awaitClient()) + assertTrue("no onOpen", recorder.opened.await(5, TimeUnit.SECONDS)) + + relay.sendClose() + + assertTrue("client never answered the relay's CLOSE frame", relay.clientCloseFrame.await(5, TimeUnit.SECONDS)) + assertEquals(1000, relay.clientCloseCode.get()) + + assertTrue("onClosed never fired", recorder.closed.await(5, TimeUnit.SECONDS)) + assertEquals("the relay's status code is what gets reported", 1000, recorder.closedCode.get()) + assertNull("a clean handshake is not a failure", recorder.failure.get()) + + // After onClosed the wrapper has let go of its socket, so this must be a no-op. + socket.disconnect() + assertTrue("a closed socket needs a fresh dial", socket.needsReconnect()) + + client.dispatcher.executorService.shutdown() + } + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt index 72d9e164f9..d29629cfdf 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt @@ -38,11 +38,12 @@ interface INostrClient : AutoCloseable { fun availableRelaysFlow(): StateFlow> /** - * The relays whose socket is up at the moment of the call, read from the pool itself rather than - * from the callback-maintained [connectedRelaysFlow]. The flow is for reacting to changes; this is - * for reporting a count, because the flow keeps a relay whose socket died without a callback - * (see `RelayPool.connectedRelayUrls`) until the ping timeout catches up. Defaults to the flow's - * current value for clients that have no pool behind them. + * The relays whose socket is up at the moment of the call, read from the pool's members rather + * than from the callback-maintained [connectedRelaysFlow]. The flow is for reacting to changes; + * this is for reporting a count: it is computed from what the pool actually holds, so a terminal + * callback the socket layer never delivered cannot leave a relay in it that the pool has already + * dropped (see `RelayPool.connectedRelayUrls`). Defaults to the flow's current value for clients + * that have no pool behind them. */ fun connectedRelays(): Set = connectedRelaysFlow().value diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt index 99ad6aeb6f..04bb6c42c0 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt @@ -119,6 +119,9 @@ class RelayPool( relays.forEach { url, relay -> relay.disconnect() } + // We just tore every socket down; don't leave the answer to the socket layer's + // callbacks (see [connectedRelayUrls] for how those can go missing). + _connectedRelays.update { emptySet() } } fun sendOrConnectAndSync( @@ -210,6 +213,9 @@ class RelayPool( val relayInPool = relays.remove(relay) if (relayInPool != null) { relayInPool.disconnect() + // A relay that is no longer a member cannot be connected, whatever its socket + // layer reports (or fails to report) later. + _connectedRelays.update { it - relay } return true } return false @@ -226,6 +232,7 @@ class RelayPool( disconnect() relays.clear() _availableRelays.update { emptySet() } + _connectedRelays.update { emptySet() } } } @@ -271,13 +278,13 @@ class RelayPool( /** * The relays whose socket is up *right now*, read from each pool member's [IRelayClient.isConnected]. * - * This is the ground truth; [connectedRelays] is a projection of it that only moves on the - * [onConnected] / [onDisconnected] callbacks. The two drift whenever a socket dies without a - * callback: a relay that sent a WebSocket CLOSE frame leaves OkHttp waiting for a reply that never - * comes, so neither `onClosed` nor `onFailure` fires, and a later `cancel()` is silent too. The URL - * then sits in [connectedRelays] until the 120s ping path finally fails, while this snapshot has - * already dropped it (a removed relay is no longer a member; a still-desired one reads - * `isConnected() == false` as soon as its socket closes). Use this for anything a person reads as + * [connectedRelays] is a projection of this that moves on the [onConnected] / [onDisconnected] + * callbacks plus the pool's own removals and [disconnect]. The two can still drift for a relay that + * is *still a member* whose socket layer lost a terminal callback: before the OkHttp sockets + * answered a relay's CLOSE frame, that was every relay-initiated close (no `onClosed`, no + * `onFailure`, and a silent `cancel()` afterwards), and the flow carried such relays for minutes + * after the pool had let go of them. Reading the members directly cannot be fooled by a callback + * that never came for a relay the pool no longer holds. Use this for anything a person reads as * "how many relays am I connected to"; keep the flow for change notification. */ fun connectedRelayUrls(): Set { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt index eec65748e9..2b1ef0beff 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt @@ -83,6 +83,17 @@ open class BasicRelayClient( private var socket: WebSocket? = null + /** + * The listener wired to the socket this client currently owns. Every callback checks it is still + * the current one before touching any state: a socket this client has already replaced (see + * [disconnect] + [connect], which is how the pool rebuilds a session) or torn down can still + * report on its own thread afterwards -- OkHttp delivers `onClosed` from its writer thread once + * the close handshake completes, and a cancelled socket's failure lands later still. Without + * the check such a late callback would null the NEW socket out from under the client, leaving a + * live connection orphaned and dialing a third one on the next pass. + */ + @Volatile private var currentListener: MyWebsocketListener? = null + // True if it has received the onOpen call from the socket. // @Volatile: written on the serialized socket-callback thread, read from the // relay-pool/timer thread (see RelayLoadingCursors for the same pattern). @@ -132,7 +143,9 @@ open class BasicRelayClient( lastConnectTentativeInSeconds = nowInSeconds() - socket = socketBuilder.build(url, MyWebsocketListener()) + val newListener = MyWebsocketListener() + currentListener = newListener + socket = socketBuilder.build(url, newListener) socket?.connect() } catch (e: Exception) { if (e is CancellationException) throw e @@ -148,15 +161,20 @@ open class BasicRelayClient( } inner class MyWebsocketListener : WebSocketListener { + /** True once this client moved on to another socket, or tore this one down itself. */ + private fun isStale() = currentListener !== this + override fun onOpen( pingMillis: Int, compression: Boolean, ) { + if (isStale()) return markConnectionAsReady(compression) listener.onConnected(this@BasicRelayClient, pingMillis, compression) } override suspend fun onMessage(text: String) { + if (isStale()) return try { val msg = decoder.decode(text) listener.onIncomingMessage(this@BasicRelayClient, text, msg) @@ -171,6 +189,7 @@ open class BasicRelayClient( code: Int, reason: String, ) { + if (isStale()) return markConnectionAsClosed() listener.onDisconnected(this@BasicRelayClient) } @@ -180,6 +199,9 @@ open class BasicRelayClient( code: Int?, response: String?, ) { + // A session this client already retired: disconnect() reported it when it happened. + if (isStale()) return + // socket is already closed // socket?.disconnect() @@ -279,10 +301,21 @@ open class BasicRelayClient( lastConnectTentativeInSeconds = 0L // this is not an error, so prepare to reconnect as soon as requested. delayToConnectInSeconds = DELAY_TO_RECONNECT_IN_SECS connectedAtInSeconds = 0L - socket?.disconnect() + val closing = socket + // Retire the session before touching the socket: whatever its layer reports from here + // on is about a socket this client no longer owns, and is ignored (see currentListener). + currentListener = null socket = null isReady = false usingCompression = false + if (closing != null) { + closing.disconnect() + // Report the teardown ourselves instead of waiting for the socket layer to confirm + // it. It might not: OkHttp's cancel() raises no callback when no reader is left to + // fail, which is exactly the state a relay-initiated close leaves behind, so the + // pool's bookkeeping used to keep such a relay "connected" until the ping timeout. + listener.onDisconnected(this) + } } override fun connectAndSyncFiltersIfDisconnected(ignoreRetryDelays: Boolean) { diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt index b1bb89789d..1cd21c217b 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt @@ -26,25 +26,30 @@ import kotlin.test.Test import kotlin.test.assertEquals /** - * [RelayPool.connectedRelayUrls] must report what the pool actually holds open, even when the - * socket layer never confirms a close. + * Both views of "connected" -- the callback-fed [RelayPool.connectedRelays] flow and the + * members-read [RelayPool.connectedRelayUrls] snapshot -- must agree the moment the pool lets a + * relay go, even when the socket layer never confirms the close. * * [com.vitorpamplona.quartz.nip01Core.relay.client.single.basic.FakeWebSocket.disconnect] is a - * no-op that never calls back, which is exactly what OkHttp does after a relay sent a WebSocket - * CLOSE frame the app never answered: `cancel()` then fires neither `onClosed` nor `onFailure`. - * The callback-driven [RelayPool.connectedRelays] flow keeps such a relay for minutes; the - * snapshot drops it the moment the pool lets go of it. + * no-op that never calls back, which is exactly what OkHttp did after a relay sent a CLOSE frame + * the app never answered: `cancel()` then fired neither `onClosed` nor `onFailure`. The flow used + * to keep such a relay for minutes; now the pool clears it on removal and on disconnect, and the + * client reports its own teardown, so neither view depends on that callback. */ class RelayPoolConnectedSnapshotTest { private val url = NormalizedRelayUrl("wss://relay.example.com/") - @Test - fun snapshotMatchesTheFlowWhileTheSocketIsOpen() { + private fun openedPool(): Pair { val sockets = FakeWebsocketBuilder() val pool = RelayPool(sockets) - pool.getOrCreateRelay(url).connect() sockets.lastListener.onOpen(pingMillis = 10, compression = false) + return sockets to pool + } + + @Test + fun snapshotMatchesTheFlowWhileTheSocketIsOpen() { + val (_, pool) = openedPool() assertEquals(setOf(url), pool.connectedRelays.value) assertEquals(setOf(url), pool.connectedRelayUrls()) @@ -52,21 +57,27 @@ class RelayPoolConnectedSnapshotTest { } @Test - fun removedRelayLeavesTheSnapshotEvenWhenTheCloseIsSilent() { - val sockets = FakeWebsocketBuilder() - val pool = RelayPool(sockets) + fun removingARelayClearsBothViewsEvenWhenTheCloseIsSilent() { + val (_, pool) = openedPool() - pool.getOrCreateRelay(url).connect() - sockets.lastListener.onOpen(pingMillis = 10, compression = false) - - // The pool drops the relay (no subscription wants it anymore) and cancels its socket, - // but the socket layer stays silent. + // No subscription wants it anymore: the pool drops it and cancels its socket, and the + // socket layer stays silent. pool.removeRelay(url) - // Documents the drift this test exists for: the flow still carries the relay... - assertEquals(setOf(url), pool.connectedRelays.value) - // ...while the snapshot already reflects that nothing is held open. + assertEquals(emptySet(), pool.connectedRelays.value) assertEquals(emptySet(), pool.connectedRelayUrls()) assertEquals(0, pool.connectedRelaysCount()) } + + @Test + fun disconnectingThePoolClearsBothViewsEvenWhenTheCloseIsSilent() { + val (_, pool) = openedPool() + + // The host is putting the client down (app backgrounded, connectivity lost). + pool.disconnect() + + assertEquals(emptySet(), pool.connectedRelays.value) + assertEquals(emptySet(), pool.connectedRelayUrls()) + assertEquals(setOf(url), pool.availableRelays.value, "still a member, just not connected") + } } diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt new file mode 100644 index 0000000000..0b15ff4be9 --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt @@ -0,0 +1,123 @@ +/* + * 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.client.single.basic + +import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener +import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +/** + * A socket this client has retired must not be able to reach back into its state. + * + * Socket layers report on their own threads, after the fact: OkHttp delivers `onClosed` from its + * writer thread once a close handshake completes, and a cancelled socket's failure lands later + * still. The pool rebuilds a session with `disconnect()` + `connect()`, so a late callback from + * the old socket used to null the new one out from under the client -- orphaning a live + * connection and dialing a third one on the next pass. + */ +class BasicRelayClientStaleSocketTest { + private class Counting : RelayConnectionListener { + var connected = 0 + var disconnected = 0 + + override fun onConnected( + relay: IRelayClient, + pingMillis: Int, + compressed: Boolean, + ) { + connected++ + } + + override fun onDisconnected(relay: IRelayClient) { + disconnected++ + } + } + + private val url = NormalizedRelayUrl("wss://relay.example.com/") + + @Test + fun `disconnect reports the teardown itself instead of waiting for the socket layer`() { + val sockets = FakeWebsocketBuilder() + val events = Counting() + val client = BasicRelayClient(url, sockets, events) + + client.connect() + sockets.lastListener.onOpen(50, false) + assertTrue(client.isConnected()) + + // FakeWebSocket.disconnect() never calls back -- the case OkHttp's cancel() leaves after + // a relay-initiated close. + client.disconnect() + + assertFalse(client.isConnected()) + assertEquals(1, events.disconnected, "the client itself must say the session ended") + } + + @Test + fun `a late callback from a replaced socket cannot touch the new session`() { + val sockets = FakeWebsocketBuilder() + val events = Counting() + val client = BasicRelayClient(url, sockets, events) + + client.connect() + val first = sockets.lastListener + first.onOpen(50, false) + + // What RelayPool.reconnectIfNeedsTo does when the transport changed under a live socket. + client.disconnect() + client.connect() + val second = sockets.lastListener + second.onOpen(50, false) + assertTrue(client.isConnected()) + assertEquals(2, events.connected) + assertEquals(1, events.disconnected) + + // The old socket finally reports, in every way OkHttp can. + first.onClosed(1000, "late close handshake") + first.onFailure(RuntimeException("Socket closed"), null, null) + first.onOpen(50, false) + + assertTrue(client.isConnected(), "the live session must survive its predecessor's callbacks") + assertEquals(2, events.connected, "no phantom connect") + assertEquals(1, events.disconnected, "no phantom disconnect") + } + + @Test + fun `a callback after disconnect is not a second disconnect`() { + val sockets = FakeWebsocketBuilder() + val events = Counting() + val client = BasicRelayClient(url, sockets, events) + + client.connect() + val socket = sockets.lastListener + socket.onOpen(50, false) + client.disconnect() + + // The cancel's own failure arriving afterwards, as it does for a healthy socket. + socket.onFailure(RuntimeException("Canceled"), null, null) + + assertEquals(1, events.disconnected, "disconnect() already reported this session") + } +} From 08529c3a748e1be0a9f885b22e703f953e2f6edb Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 13 Sep 2026 01:36:28 +0000 Subject: [PATCH 4/7] refactor(relay): make the socket adapters own the session, and take the guard out of BasicRelayClient The previous commit taught BasicRelayClient to remember which listener it had wired and to ignore callbacks from any other. That put the knowledge "this report is about a socket I already threw away" in the wrong layer: the transport adapter is the one that owns the socket, and OkHttp names the socket in every callback, so the adapter can tell for free. BasicRelayClient goes back to exactly its previous code. The contract it relies on is now written on the WebSocket interface and kept by every transport: a session ends with exactly one terminal callback, and disconnect() reports onClosed synchronously and forwards nothing from that socket afterwards -- what InProcessWebSocket has always done. Both OkHttp adapters (quartz's BasicOkHttpWebSocket and the Android app's OkHttpWebSocket) now check ownership on every callback, claim the terminal report under a lock so a session cannot be reported twice, and answer disconnect() themselves instead of waiting for OkHttp, which raises nothing for a cancel when no reader is left to fail and otherwise raises it later on its own thread. The reconnect race this closes is the same one as before: with disconnect()+connect() back to back, the old socket's late failure used to land on the new connection and cancel it. The client-level stale-socket test is replaced by adapter-level tests on both modules: disconnect() reports once and synchronously, OkHttp's own reaction to the cancel never surfaces, and a relay-initiated close followed by disconnect() is still one report. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_014cq6vrfQkASwxgqXXY4py8 --- .../service/okhttp/OkHttpWebSocket.kt | 78 +++++++---- .../OkHttpWebSocketCloseHandshakeTest.kt | 34 ++++- .../client/single/basic/BasicRelayClient.kt | 37 +----- .../nip01Core/relay/sockets/WebSocket.kt | 15 +++ .../pool/RelayPoolConnectedSnapshotTest.kt | 5 +- .../basic/BasicRelayClientStaleSocketTest.kt | 123 ------------------ .../sockets/okhttp/BasicOkHttpWebSocket.kt | 70 +++++++--- .../BasicOkHttpWebSocketCloseHandshakeTest.kt | 35 ++++- 8 files changed, 189 insertions(+), 208 deletions(-) delete mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt index 55b5994060..33880cb1cb 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt @@ -40,8 +40,17 @@ class OkHttpWebSocket( val httpClient: (url: NormalizedRelayUrl) -> OkHttpClient, val out: WebSocketListener, ) : WebSocket { + private val lock = Any() private var usingOkHttp: OkHttpClient? = null - private var socket: okhttp3.WebSocket? = null + + /** + * The OkHttp socket this adapter currently owns, or null once the session has ended -- by the + * relay closing it, by a network failure, or by [disconnect]. Only the owned socket may reach + * [out], and the terminal callbacks claim the slot under [lock], so a session ends with exactly + * one report however it ends. See quartz's `BasicOkHttpWebSocket` for the full reasoning; the + * two adapters differ only in how [needsReconnect] is decided. + */ + @Volatile private var socket: okhttp3.WebSocket? = null fun buildRequest() = Request.Builder().url(url.url).build() @@ -68,8 +77,13 @@ class OkHttpWebSocket( } override fun connect() { - usingOkHttp = httpClient(url) - socket = usingOkHttp?.newWebSocket(buildRequest(), OkHttpWebsocketListener(out)) + val client = httpClient(url) + // Under the lock so a callback racing this dial waits until the socket is owned rather + // than being dropped as foreign. + synchronized(lock) { + usingOkHttp = client + socket = client.newWebSocket(buildRequest(), OkHttpWebsocketListener(out)) + } } inner class OkHttpWebsocketListener( @@ -91,21 +105,38 @@ class OkHttpWebSocket( } } + /** Only the socket this adapter still owns may reach [out]. */ + private fun isOwned(webSocket: okhttp3.WebSocket) = synchronized(lock) { socket === webSocket } + + /** Claims the session's single terminal report. False if it already ended. */ + private fun endSession(webSocket: okhttp3.WebSocket): Boolean { + val ended = synchronized(lock) { (socket === webSocket).also { if (it) socket = null } } + if (ended) { + incomingMessages.close() + job.cancel() + scope.cancel() + } + return ended + } + override fun onOpen( webSocket: okhttp3.WebSocket, response: Response, - ) = out.onOpen( - (response.receivedResponseAtMillis - response.sentRequestAtMillis).toInt(), - response.headers["Sec-WebSocket-Extensions"]?.contains("permessage-deflate") ?: false, - ) + ) { + if (!isOwned(webSocket)) return + out.onOpen( + (response.receivedResponseAtMillis - response.sentRequestAtMillis).toInt(), + response.headers["Sec-WebSocket-Extensions"]?.contains("permessage-deflate") ?: false, + ) + } override fun onMessage( webSocket: okhttp3.WebSocket, text: String, ) { - // Asynchronously send the received message to the channel. - // `trySendBlocking` is used here for simplicity within the callback, - // but it's important to understand potential thread blocking if the buffer is full. + if (!isOwned(webSocket)) return + // Never blocks (unlimited channel): the OkHttp reader thread must + // stay free to keep draining the socket. incomingMessages.trySendBlocking(text) } @@ -119,8 +150,7 @@ class OkHttpWebSocket( // WebSocketEcho recipe does exactly this). Unanswered, the socket sat half-closed: // no onClosed, no onFailure, send() still accepted and discarded, a later cancel() // silent too -- so the relay client believed it was connected until the 120s ping - // path failed up to two intervals later. Mirrors BasicOkHttpWebSocket in quartz; - // the two classes differ only in how needsReconnect() is decided. + // path failed up to two intervals later. // // Always 1000 rather than echoing `code`: close() validates the code it writes and // throws on the reserved ones (1005, 1006, 1015), and a relay may send anything. @@ -132,12 +162,7 @@ class OkHttpWebSocket( code: Int, reason: String, ) { - // Close the channel on failure, and propagate the error. - incomingMessages.close() - job.cancel() - scope.cancel() - - socket = null + if (!endSession(webSocket)) return out.onClosed(code, reason) } @@ -146,12 +171,7 @@ class OkHttpWebSocket( t: Throwable, response: Response?, ) { - // Close the channel on failure, and propagate the error. - incomingMessages.close() - job.cancel() - scope.cancel() - - socket = null + if (!endSession(webSocket)) return out.onFailure(t, response?.code, response?.message) } } @@ -171,9 +191,13 @@ class OkHttpWebSocket( } override fun disconnect() { - // uses cancel to kill the SEND stack that might be waiting - socket?.cancel() - socket = null + // Claim the session ourselves and cancel (which also kills a SEND stack that might be + // waiting): OkHttp's cancel() raises no callback when no reader is left to fail, and when + // it does the failure arrives later on its own thread. The relay client needs the answer + // now, and must not hear from this socket again. + val closing = synchronized(lock) { socket?.also { socket = null } } ?: return + closing.cancel() + out.onClosed(1000, "client disconnect") } override fun send(msg: String): Boolean = socket?.send(msg) ?: false diff --git a/amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt b/amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt index cddd2f9f6f..907a0b1ec4 100644 --- a/amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt +++ b/amethyst/src/test/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocketCloseHandshakeTest.kt @@ -154,6 +154,7 @@ class OkHttpWebSocketCloseHandshakeTest { private class Recorder : WebSocketListener { val opened = CountDownLatch(1) val closed = CountDownLatch(1) + val closedCount = AtomicInteger(0) val closedCode = AtomicInteger(-1) val failure = AtomicReference(null) @@ -169,6 +170,7 @@ class OkHttpWebSocketCloseHandshakeTest { reason: String, ) { closedCode.set(code) + closedCount.incrementAndGet() closed.countDown() } @@ -201,11 +203,41 @@ class OkHttpWebSocketCloseHandshakeTest { assertEquals("the relay's status code is what gets reported", 1000, recorder.closedCode.get()) assertNull("a clean handshake is not a failure", recorder.failure.get()) - // After onClosed the wrapper has let go of its socket, so this must be a no-op. + // The session already ended; the usual teardown afterwards must not report it twice. socket.disconnect() + assertEquals("one terminal report per session", 1, recorder.closedCount.get()) assertTrue("a closed socket needs a fresh dial", socket.needsReconnect()) client.dispatcher.executorService.shutdown() } } + + @Test + fun `disconnect reports the session end once, synchronously, and drops what OkHttp says afterwards`() { + TinyRelay().use { relay -> + val recorder = Recorder() + val client = OkHttpClient() + val socket = OkHttpWebSocket(relay.url, { client }, recorder) + + socket.connect() + assertTrue("relay never saw the client", relay.awaitClient()) + assertTrue("no onOpen", recorder.opened.await(5, TimeUnit.SECONDS)) + + socket.disconnect() + + // Reported before disconnect() returned: the relay client dials the replacement + // right after this call and must not hear from the old socket later. + assertEquals("disconnect() must report synchronously", 1, recorder.closedCount.get()) + assertEquals(1000, recorder.closedCode.get()) + assertTrue(socket.needsReconnect()) + + // OkHttp's own reaction to cancel() -- a failure on its reader thread -- and the relay's + // reaction to the dropped TCP session must both be swallowed. + Thread.sleep(500) + assertEquals("no second report", 1, recorder.closedCount.get()) + assertNull("the cancel's failure must not surface", recorder.failure.get()) + + client.dispatcher.executorService.shutdown() + } + } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt index 2b1ef0beff..eec65748e9 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClient.kt @@ -83,17 +83,6 @@ open class BasicRelayClient( private var socket: WebSocket? = null - /** - * The listener wired to the socket this client currently owns. Every callback checks it is still - * the current one before touching any state: a socket this client has already replaced (see - * [disconnect] + [connect], which is how the pool rebuilds a session) or torn down can still - * report on its own thread afterwards -- OkHttp delivers `onClosed` from its writer thread once - * the close handshake completes, and a cancelled socket's failure lands later still. Without - * the check such a late callback would null the NEW socket out from under the client, leaving a - * live connection orphaned and dialing a third one on the next pass. - */ - @Volatile private var currentListener: MyWebsocketListener? = null - // True if it has received the onOpen call from the socket. // @Volatile: written on the serialized socket-callback thread, read from the // relay-pool/timer thread (see RelayLoadingCursors for the same pattern). @@ -143,9 +132,7 @@ open class BasicRelayClient( lastConnectTentativeInSeconds = nowInSeconds() - val newListener = MyWebsocketListener() - currentListener = newListener - socket = socketBuilder.build(url, newListener) + socket = socketBuilder.build(url, MyWebsocketListener()) socket?.connect() } catch (e: Exception) { if (e is CancellationException) throw e @@ -161,20 +148,15 @@ open class BasicRelayClient( } inner class MyWebsocketListener : WebSocketListener { - /** True once this client moved on to another socket, or tore this one down itself. */ - private fun isStale() = currentListener !== this - override fun onOpen( pingMillis: Int, compression: Boolean, ) { - if (isStale()) return markConnectionAsReady(compression) listener.onConnected(this@BasicRelayClient, pingMillis, compression) } override suspend fun onMessage(text: String) { - if (isStale()) return try { val msg = decoder.decode(text) listener.onIncomingMessage(this@BasicRelayClient, text, msg) @@ -189,7 +171,6 @@ open class BasicRelayClient( code: Int, reason: String, ) { - if (isStale()) return markConnectionAsClosed() listener.onDisconnected(this@BasicRelayClient) } @@ -199,9 +180,6 @@ open class BasicRelayClient( code: Int?, response: String?, ) { - // A session this client already retired: disconnect() reported it when it happened. - if (isStale()) return - // socket is already closed // socket?.disconnect() @@ -301,21 +279,10 @@ open class BasicRelayClient( lastConnectTentativeInSeconds = 0L // this is not an error, so prepare to reconnect as soon as requested. delayToConnectInSeconds = DELAY_TO_RECONNECT_IN_SECS connectedAtInSeconds = 0L - val closing = socket - // Retire the session before touching the socket: whatever its layer reports from here - // on is about a socket this client no longer owns, and is ignored (see currentListener). - currentListener = null + socket?.disconnect() socket = null isReady = false usingCompression = false - if (closing != null) { - closing.disconnect() - // Report the teardown ourselves instead of waiting for the socket layer to confirm - // it. It might not: OkHttp's cancel() raises no callback when no reader is left to - // fail, which is exactly the state a relay-initiated close leaves behind, so the - // pool's bookkeeping used to keep such a relay "connected" until the ping timeout. - listener.onDisconnected(this) - } } override fun connectAndSyncFiltersIfDisconnected(ignoreRetryDelays: Boolean) { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/WebSocket.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/WebSocket.kt index 525d5b992b..6c1172230a 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/WebSocket.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/WebSocket.kt @@ -20,11 +20,26 @@ */ package com.vitorpamplona.quartz.nip01Core.relay.sockets +/** + * One socket session towards a relay, as the relay client sees it. + * + * The contract every implementation keeps, and that [com.vitorpamplona.quartz.nip01Core.relay.client.single.basic.BasicRelayClient] + * relies on for its bookkeeping: + * + * - A session ends with **exactly one** terminal callback on its [WebSocketListener], `onClosed` + * or `onFailure`, however it ends. + * - [disconnect] ends the session itself: it reports `onClosed` **synchronously**, before + * returning, and nothing from that socket reaches the listener afterwards. The relay client + * may dial a new socket immediately, so a late report from the old one -- which OkHttp + * delivers on its own threads for a cancel, and never delivers at all for a relay-initiated + * close it was not allowed to finish -- must be swallowed by the adapter, not forwarded. + */ interface WebSocket { fun needsReconnect(): Boolean fun connect() + /** Ends the session now. Reports `onClosed` synchronously if one was open; a no-op otherwise. */ fun disconnect() fun send(msg: String): Boolean diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt index 1cd21c217b..40ebdc94a7 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt @@ -33,8 +33,9 @@ import kotlin.test.assertEquals * [com.vitorpamplona.quartz.nip01Core.relay.client.single.basic.FakeWebSocket.disconnect] is a * no-op that never calls back, which is exactly what OkHttp did after a relay sent a CLOSE frame * the app never answered: `cancel()` then fired neither `onClosed` nor `onFailure`. The flow used - * to keep such a relay for minutes; now the pool clears it on removal and on disconnect, and the - * client reports its own teardown, so neither view depends on that callback. + * to keep such a relay for minutes; now the pool clears it on removal and on disconnect itself + * (and the real transports report a disconnect synchronously, see [com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket]), + * so neither view depends on a callback that may never come. */ class RelayPoolConnectedSnapshotTest { private val url = NormalizedRelayUrl("wss://relay.example.com/") diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt deleted file mode 100644 index 0b15ff4be9..0000000000 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/single/basic/BasicRelayClientStaleSocketTest.kt +++ /dev/null @@ -1,123 +0,0 @@ -/* - * 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.client.single.basic - -import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener -import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient -import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl -import kotlin.test.Test -import kotlin.test.assertEquals -import kotlin.test.assertFalse -import kotlin.test.assertTrue - -/** - * A socket this client has retired must not be able to reach back into its state. - * - * Socket layers report on their own threads, after the fact: OkHttp delivers `onClosed` from its - * writer thread once a close handshake completes, and a cancelled socket's failure lands later - * still. The pool rebuilds a session with `disconnect()` + `connect()`, so a late callback from - * the old socket used to null the new one out from under the client -- orphaning a live - * connection and dialing a third one on the next pass. - */ -class BasicRelayClientStaleSocketTest { - private class Counting : RelayConnectionListener { - var connected = 0 - var disconnected = 0 - - override fun onConnected( - relay: IRelayClient, - pingMillis: Int, - compressed: Boolean, - ) { - connected++ - } - - override fun onDisconnected(relay: IRelayClient) { - disconnected++ - } - } - - private val url = NormalizedRelayUrl("wss://relay.example.com/") - - @Test - fun `disconnect reports the teardown itself instead of waiting for the socket layer`() { - val sockets = FakeWebsocketBuilder() - val events = Counting() - val client = BasicRelayClient(url, sockets, events) - - client.connect() - sockets.lastListener.onOpen(50, false) - assertTrue(client.isConnected()) - - // FakeWebSocket.disconnect() never calls back -- the case OkHttp's cancel() leaves after - // a relay-initiated close. - client.disconnect() - - assertFalse(client.isConnected()) - assertEquals(1, events.disconnected, "the client itself must say the session ended") - } - - @Test - fun `a late callback from a replaced socket cannot touch the new session`() { - val sockets = FakeWebsocketBuilder() - val events = Counting() - val client = BasicRelayClient(url, sockets, events) - - client.connect() - val first = sockets.lastListener - first.onOpen(50, false) - - // What RelayPool.reconnectIfNeedsTo does when the transport changed under a live socket. - client.disconnect() - client.connect() - val second = sockets.lastListener - second.onOpen(50, false) - assertTrue(client.isConnected()) - assertEquals(2, events.connected) - assertEquals(1, events.disconnected) - - // The old socket finally reports, in every way OkHttp can. - first.onClosed(1000, "late close handshake") - first.onFailure(RuntimeException("Socket closed"), null, null) - first.onOpen(50, false) - - assertTrue(client.isConnected(), "the live session must survive its predecessor's callbacks") - assertEquals(2, events.connected, "no phantom connect") - assertEquals(1, events.disconnected, "no phantom disconnect") - } - - @Test - fun `a callback after disconnect is not a second disconnect`() { - val sockets = FakeWebsocketBuilder() - val events = Counting() - val client = BasicRelayClient(url, sockets, events) - - client.connect() - val socket = sockets.lastListener - socket.onOpen(50, false) - client.disconnect() - - // The cancel's own failure arriving afterwards, as it does for a healthy socket. - socket.onFailure(RuntimeException("Canceled"), null, null) - - assertEquals(1, events.disconnected, "disconnect() already reported this session") - } -} diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt index 80b2679e73..dc03b99661 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt @@ -51,7 +51,21 @@ class BasicOkHttpWebSocket( } } - private var socket: OkHttpWebSocket? = null + private val lock = Any() + + /** + * The OkHttp socket this adapter currently owns, or null once the session has ended -- by the + * relay closing it, by a network failure, or by [disconnect]. + * + * OkHttp names the socket in every callback, and only the owned one may reach [out]. That is + * what makes this adapter honour the [WebSocket.disconnect] contract: after [disconnect] the + * slot is empty, so the failure OkHttp raises for its own `cancel()` on the reader thread, + * or the `onClosed` its writer thread delivers once a close handshake completes, is dropped + * instead of reaching a relay client that has already moved on to a new socket. The terminal + * callbacks claim the slot under [lock], so a session ends with exactly one report however it + * ends. + */ + @Volatile private var socket: OkHttpWebSocket? = null override fun needsReconnect() = socket == null @@ -79,18 +93,36 @@ class BasicOkHttpWebSocket( } } + /** Only the socket this adapter still owns may reach [out]. */ + private fun isOwned(webSocket: OkHttpWebSocket) = synchronized(lock) { socket === webSocket } + + /** Claims the session's single terminal report. False if it already ended. */ + private fun endSession(webSocket: OkHttpWebSocket): Boolean { + val ended = synchronized(lock) { (socket === webSocket).also { if (it) socket = null } } + if (ended) { + incomingMessages.close() + job.cancel() + scope.cancel() + } + return ended + } + override fun onOpen( webSocket: OkHttpWebSocket, response: Response, - ) = out.onOpen( - (response.receivedResponseAtMillis - response.sentRequestAtMillis).toInt(), - response.headers["Sec-WebSocket-Extensions"]?.contains("permessage-deflate") ?: false, - ) + ) { + if (!isOwned(webSocket)) return + out.onOpen( + (response.receivedResponseAtMillis - response.sentRequestAtMillis).toInt(), + response.headers["Sec-WebSocket-Extensions"]?.contains("permessage-deflate") ?: false, + ) + } override fun onMessage( webSocket: OkHttpWebSocket, text: String, ) { + if (!isOwned(webSocket)) return // Never blocks (unlimited channel): the OkHttp reader // thread must stay free to keep draining the socket. incomingMessages.trySendBlocking(text) @@ -123,11 +155,7 @@ class BasicOkHttpWebSocket( code: Int, reason: String, ) { - // Close the channel when the WebSocket connection is closed. - incomingMessages.close() - job.cancel() - scope.cancel() - + if (!endSession(webSocket)) return out.onClosed(code, reason) } @@ -136,21 +164,25 @@ class BasicOkHttpWebSocket( t: Throwable, response: Response?, ) { - // Close the channel on failure, and propagate the error. - incomingMessages.close() - job.cancel() - scope.cancel() - + if (!endSession(webSocket)) return out.onFailure(t, response?.code, response?.message) } } - socket = httpClient(url).newWebSocket(request, listener) + // Under the lock so a callback racing this dial (an instant failure lands on another + // thread) waits until the socket is owned, rather than being dropped as foreign. + synchronized(lock) { + socket = httpClient(url).newWebSocket(request, listener) + } } override fun disconnect() { - socket?.cancel() - socket = null + // Claim the session ourselves: OkHttp's cancel() raises no callback when no reader is + // left to fail (the state a relay-initiated close leaves behind), and when it does the + // failure arrives later on its own thread. The relay client needs the answer now. + val closing = synchronized(lock) { socket?.also { socket = null } } ?: return + closing.cancel() + out.onClosed(1000, "client disconnect") } override fun send(msg: String): Boolean = socket?.send(msg) ?: false @@ -162,6 +194,6 @@ class BasicOkHttpWebSocket( override fun build( url: NormalizedRelayUrl, out: WebSocketListener, - ) = BasicOkHttpWebSocket(url, httpClient, out) + ): WebSocket = BasicOkHttpWebSocket(url, httpClient, out) } } diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt index 56c525b946..27b43775e6 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocketCloseHandshakeTest.kt @@ -158,6 +158,7 @@ class BasicOkHttpWebSocketCloseHandshakeTest { private class Recorder : WebSocketListener { val opened = CountDownLatch(1) val closed = CountDownLatch(1) + val closedCount = AtomicInteger(0) val closedCode = AtomicInteger(-1) val failure = AtomicReference(null) @@ -173,6 +174,7 @@ class BasicOkHttpWebSocketCloseHandshakeTest { reason: String, ) { closedCode.set(code) + closedCount.incrementAndGet() closed.countDown() } @@ -207,8 +209,39 @@ class BasicOkHttpWebSocketCloseHandshakeTest { assertEquals("the relay's status code is what gets reported", 1000, recorder.closedCode.get()) assertNull("a clean handshake is not a failure", recorder.failure.get()) - // The usual teardown afterwards must stay a harmless no-op. + // The session already ended; the usual teardown afterwards must not report it twice. socket.disconnect() + assertEquals("one terminal report per session", 1, recorder.closedCount.get()) + assertTrue("a closed socket needs a fresh dial", socket.needsReconnect()) + + client.dispatcher.executorService.shutdown() + } + } + + @Test + fun `disconnect reports the session end once, synchronously, and drops what OkHttp says afterwards`() { + TinyRelay().use { relay -> + val recorder = Recorder() + val client = OkHttpClient() + val socket = BasicOkHttpWebSocket(relay.url, { client }, recorder) + + socket.connect() + assertTrue("relay never saw the client", relay.awaitClient()) + assertTrue("no onOpen", recorder.opened.await(5, TimeUnit.SECONDS)) + + socket.disconnect() + + // Reported before disconnect() returned: the relay client dials the replacement + // right after this call and must not hear from the old socket later. + assertEquals("disconnect() must report synchronously", 1, recorder.closedCount.get()) + assertEquals(1000, recorder.closedCode.get()) + assertTrue(socket.needsReconnect()) + + // OkHttp's own reaction to cancel() -- a failure on its reader thread -- and the relay's + // reaction to the dropped TCP session must both be swallowed. + Thread.sleep(500) + assertEquals("no second report", 1, recorder.closedCount.get()) + assertNull("the cancel's failure must not surface", recorder.failure.get()) client.dispatcher.executorService.shutdown() } From 0067db3608715a0b9299b522063573e0afcf01c9 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 13 Sep 2026 01:43:11 +0000 Subject: [PATCH 5/7] refactor(notifications): trigger the count refresh on the connected flow alone availableRelaysFlow was merged in when connectedRelaysFlow could not be trusted to move on a removal: a relay the pool had dropped could sit in it for minutes, so the one flow that did move on membership changes was used as an extra wake-up and the count re-read from the pool members. The pool now clears the connected flow itself whenever it lets a relay go, and every transport reports its session end exactly once, so that flow emits on every change the count can reflect and the extra trigger only added wake-ups on membership churn. The count still comes from client.connectedRelays(), the members' live socket state, which is what it is meant to show. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_014cq6vrfQkASwxgqXXY4py8 --- .../notifications/NotificationRelayService.kt | 29 ++++++++++--------- 1 file changed, 15 insertions(+), 14 deletions(-) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt index d7ee24f75a..18d0894f9d 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt @@ -50,7 +50,6 @@ import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.flow.collectLatest -import kotlinx.coroutines.flow.merge import kotlinx.coroutines.flow.sample import kotlinx.coroutines.launch @@ -300,8 +299,8 @@ class NotificationRelayService : Service() { * Tor, network changes). Without this, the client disconnects 30s after * the UI stops collecting. * - * 2. connectedRelaysFlow + availableRelaysFlow: re-read the pool's live relay count - * (client.connectedRelays()) and refresh the persistent notification. + * 2. connectedRelaysFlow: re-read the pool's live relay count (client.connectedRelays()) + * and refresh the persistent notification. * * The service does NOT create its own relay subscriptions. Instead, it relies on * the AccountFilterAssembler subscription that lives in the Compose tree (LoggedInPage). @@ -322,17 +321,18 @@ class NotificationRelayService : Service() { } launch { - // The two flows are only the *trigger*; the number comes from + // The flow is only the *trigger*; the number comes from // client.connectedRelays(), which reads each pool member's live socket - // state. connectedRelaysFlow() alone used to over-report: it is fed by - // socket callbacks, and until the OkHttp sockets answered a relay's CLOSE - // frame a relay-initiated close produced none (no onClosed, no onFailure, - // and a silent cancel() afterwards). After the feeds tore down in the - // background that left hundreds of already-dropped relays in the flow for - // minutes, with no subscription to justify a single one of them. Reading - // the members directly cannot be fooled that way, and availableRelaysFlow() - // is merged in because that is the flow that moves when the pool drops a - // relay. + // state. The flow's own value used to over-report: it is fed by socket + // callbacks, and until the OkHttp adapters answered a relay's CLOSE frame + // a relay-initiated close produced none (no onClosed, no onFailure, and a + // silent cancel() afterwards). After the feeds tore down in the background + // that left hundreds of already-dropped relays in it for minutes, with no + // subscription to justify a single one of them. The pool now clears the + // flow itself when it lets a relay go and every transport reports its + // session end exactly once, so the flow moves on every change that matters + // and is a sufficient trigger; the members are still read directly because + // that is the ground truth the count is meant to show. // // sample() caps how often we touch the notification. During feed // load/teardown these flows churn dozens of times per second; posting on @@ -341,7 +341,8 @@ class NotificationRelayService : Service() { // intermediate value. One refresh per second stays well under the limit // and always lands the settled count. val client = Amethyst.instance.client - merge(client.connectedRelaysFlow(), client.availableRelaysFlow()) + client + .connectedRelaysFlow() .sample(NOTIFICATION_REFRESH_MS) .collectLatest { val count = client.connectedRelays().size From b2431840261155872cd8e03b21fbf1733855cbac Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 13 Sep 2026 12:07:46 +0000 Subject: [PATCH 6/7] refactor(relay): count from the connected flow itself; drop the members snapshot The first commit on this branch added INostrClient.connectedRelays(), a snapshot read from the pool members, because connectedRelaysFlow could not be trusted: a relay the pool had dropped could sit in it for minutes. That has since been fixed at the source -- the pool clears the flow itself when it lets a relay go, and every transport reports its session end exactly once -- so the flow's value and the members' isConnected() move on the same transitions, and re-reading the members on every emission was a redundant pool walk. The notification takes its count from the emitted set, the breakdown and the Active Subscriptions screen go back to the flow, and the snapshot API is removed. The pool test keeps pinning what the app actually reads: the flow drops a relay on removal and on disconnect even when the socket layer never confirms the close. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_014cq6vrfQkASwxgqXXY4py8 --- .../notifications/NotificationRelayService.kt | 39 ++++++++----------- .../notifications/RelayPurposeSummary.kt | 5 +-- .../ActiveSubscriptionsViewModel.kt | 3 +- .../nip01Core/relay/client/INostrClient.kt | 10 ----- .../nip01Core/relay/client/NostrClient.kt | 2 - .../nip01Core/relay/client/pool/RelayPool.kt | 25 ++---------- ...tTest.kt => RelayPoolConnectedFlowTest.kt} | 35 ++++++++--------- 7 files changed, 37 insertions(+), 82 deletions(-) rename quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/{RelayPoolConnectedSnapshotTest.kt => RelayPoolConnectedFlowTest.kt} (69%) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt index 18d0894f9d..b94e006dd6 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/NotificationRelayService.kt @@ -299,8 +299,7 @@ class NotificationRelayService : Service() { * Tor, network changes). Without this, the client disconnects 30s after * the UI stops collecting. * - * 2. connectedRelaysFlow: re-read the pool's live relay count (client.connectedRelays()) - * and refresh the persistent notification. + * 2. connectedRelaysFlow: Updates the persistent notification with relay count. * * The service does NOT create its own relay subscriptions. Instead, it relies on * the AccountFilterAssembler subscription that lives in the Compose tree (LoggedInPage). @@ -321,31 +320,25 @@ class NotificationRelayService : Service() { } launch { - // The flow is only the *trigger*; the number comes from - // client.connectedRelays(), which reads each pool member's live socket - // state. The flow's own value used to over-report: it is fed by socket - // callbacks, and until the OkHttp adapters answered a relay's CLOSE frame - // a relay-initiated close produced none (no onClosed, no onFailure, and a - // silent cancel() afterwards). After the feeds tore down in the background - // that left hundreds of already-dropped relays in it for minutes, with no - // subscription to justify a single one of them. The pool now clears the - // flow itself when it lets a relay go and every transport reports its - // session end exactly once, so the flow moves on every change that matters - // and is a sufficient trigger; the members are still read directly because - // that is the ground truth the count is meant to show. + // This flow used to over-report: it is fed by socket callbacks, and until + // the OkHttp adapters answered a relay's CLOSE frame a relay-initiated close + // produced none (no onClosed, no onFailure, and a silent cancel() + // afterwards), so after the feeds tore down in the background it carried + // hundreds of already-dropped relays for minutes. The pool now clears it + // itself whenever it lets a relay go, and every transport reports its + // session end exactly once (see WebSocket), so what it emits is the count. // // sample() caps how often we touch the notification. During feed - // load/teardown these flows churn dozens of times per second; posting on - // every delta blows past Android's notification rate limit (~10/s), which - // silently drops updates and leaves the visible count stuck on a stale - // intermediate value. One refresh per second stays well under the limit - // and always lands the settled count. - val client = Amethyst.instance.client - client + // load/teardown connectedRelaysFlow churns dozens of times per second; + // posting on every delta blows past Android's notification rate limit + // (~10/s), which silently drops updates and leaves the visible count + // stuck on a stale intermediate value. One refresh per second stays + // well under the limit and always lands the settled count. + Amethyst.instance.client .connectedRelaysFlow() .sample(NOTIFICATION_REFRESH_MS) - .collectLatest { - val count = client.connectedRelays().size + .collectLatest { relays -> + val count = relays.size if (count != connectedRelayCount) { connectedRelayCount = count updateNotification(count) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt index c205734df0..4039025be3 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/notifications/RelayPurposeSummary.kt @@ -56,10 +56,7 @@ object RelayPurposeSummary { val named = mutableMapOf>() val browsing = mutableSetOf() - // Same source as the count above it (see NotificationRelayService): the pool's live socket - // state, not the callback-maintained flow, so the breakdown never has to explain relays the - // pool has already dropped. - client.connectedRelays().forEach { relay -> + client.connectedRelaysFlow().value.forEach { relay -> client .activeRequests(relay) .values diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt index 352ca674d6..a5db379f0b 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/subscriptions/ActiveSubscriptionsViewModel.kt @@ -156,8 +156,7 @@ class ActiveSubscriptionsViewModel : ViewModel() { withContext(Dispatchers.Default) { val client = Amethyst.instance.client aggregateSubscriptions( - // The pool's live socket state, same as the always-on notification this screen explains. - client.connectedRelays().associateWith { relay -> + client.connectedRelaysFlow().value.associateWith { relay -> client.activeRequests(relay).values.flatten() }, ) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt index d29629cfdf..b77d033c7d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt @@ -37,16 +37,6 @@ interface INostrClient : AutoCloseable { fun availableRelaysFlow(): StateFlow> - /** - * The relays whose socket is up at the moment of the call, read from the pool's members rather - * than from the callback-maintained [connectedRelaysFlow]. The flow is for reacting to changes; - * this is for reporting a count: it is computed from what the pool actually holds, so a terminal - * callback the socket layer never delivered cannot leave a relay in it that the pool has already - * dropped (see `RelayPool.connectedRelayUrls`). Defaults to the flow's current value for clients - * that have no pool behind them. - */ - fun connectedRelays(): Set = connectedRelaysFlow().value - fun connect() fun disconnect() diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt index 63f3d0d3ce..21c16d9b64 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt @@ -493,8 +493,6 @@ class NostrClient( override fun connectedRelaysFlow() = relayPool.connectedRelays - override fun connectedRelays() = relayPool.connectedRelayUrls() - override fun availableRelaysFlow() = relayPool.availableRelays override fun close() { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt index 04bb6c42c0..61cc4e7627 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPool.kt @@ -119,8 +119,9 @@ class RelayPool( relays.forEach { url, relay -> relay.disconnect() } - // We just tore every socket down; don't leave the answer to the socket layer's - // callbacks (see [connectedRelayUrls] for how those can go missing). + // We just tore every socket down; say so here rather than leaving it to each socket's + // own report. The transports do report a disconnect synchronously (see [WebSocket]), + // but this flow is what the app reads as "connected", and it must not depend on that. _connectedRelays.update { emptySet() } } @@ -274,24 +275,4 @@ class RelayPool( ) = listener.onSent(relay, cmdStr, cmd, success) fun connectedRelaysCount(): Int = relays.count { url, relay -> relay.isConnected() } - - /** - * The relays whose socket is up *right now*, read from each pool member's [IRelayClient.isConnected]. - * - * [connectedRelays] is a projection of this that moves on the [onConnected] / [onDisconnected] - * callbacks plus the pool's own removals and [disconnect]. The two can still drift for a relay that - * is *still a member* whose socket layer lost a terminal callback: before the OkHttp sockets - * answered a relay's CLOSE frame, that was every relay-initiated close (no `onClosed`, no - * `onFailure`, and a silent `cancel()` afterwards), and the flow carried such relays for minutes - * after the pool had let go of them. Reading the members directly cannot be fooled by a callback - * that never came for a relay the pool no longer holds. Use this for anything a person reads as - * "how many relays am I connected to"; keep the flow for change notification. - */ - fun connectedRelayUrls(): Set { - val urls = mutableSetOf() - relays.forEach { url, relay -> - if (relay.isConnected()) urls.add(url) - } - return urls - } } diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedFlowTest.kt similarity index 69% rename from quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt rename to quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedFlowTest.kt index 40ebdc94a7..e815fa167e 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedSnapshotTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayPoolConnectedFlowTest.kt @@ -26,59 +26,56 @@ import kotlin.test.Test import kotlin.test.assertEquals /** - * Both views of "connected" -- the callback-fed [RelayPool.connectedRelays] flow and the - * members-read [RelayPool.connectedRelayUrls] snapshot -- must agree the moment the pool lets a - * relay go, even when the socket layer never confirms the close. + * [RelayPool.connectedRelays] is what the app reads as "how many relays am I connected to", so it + * must drop a relay the moment the pool lets go of it, even when the socket layer never confirms + * the close. * * [com.vitorpamplona.quartz.nip01Core.relay.client.single.basic.FakeWebSocket.disconnect] is a * no-op that never calls back, which is exactly what OkHttp did after a relay sent a CLOSE frame - * the app never answered: `cancel()` then fired neither `onClosed` nor `onFailure`. The flow used - * to keep such a relay for minutes; now the pool clears it on removal and on disconnect itself - * (and the real transports report a disconnect synchronously, see [com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket]), - * so neither view depends on a callback that may never come. + * the app never answered: `cancel()` then fired neither `onClosed` nor `onFailure`, and the flow + * carried such relays for minutes. The real transports now report a disconnect synchronously (see + * [com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket]), but this flow must not depend on + * that: the pool clears it itself on removal and on disconnect. */ -class RelayPoolConnectedSnapshotTest { +class RelayPoolConnectedFlowTest { private val url = NormalizedRelayUrl("wss://relay.example.com/") - private fun openedPool(): Pair { + private fun openedPool(): RelayPool { val sockets = FakeWebsocketBuilder() val pool = RelayPool(sockets) pool.getOrCreateRelay(url).connect() sockets.lastListener.onOpen(pingMillis = 10, compression = false) - return sockets to pool + return pool } @Test - fun snapshotMatchesTheFlowWhileTheSocketIsOpen() { - val (_, pool) = openedPool() + fun anOpenSocketIsConnected() { + val pool = openedPool() assertEquals(setOf(url), pool.connectedRelays.value) - assertEquals(setOf(url), pool.connectedRelayUrls()) assertEquals(1, pool.connectedRelaysCount()) } @Test - fun removingARelayClearsBothViewsEvenWhenTheCloseIsSilent() { - val (_, pool) = openedPool() + fun removingARelayDropsItEvenWhenTheCloseIsSilent() { + val pool = openedPool() // No subscription wants it anymore: the pool drops it and cancels its socket, and the // socket layer stays silent. pool.removeRelay(url) assertEquals(emptySet(), pool.connectedRelays.value) - assertEquals(emptySet(), pool.connectedRelayUrls()) assertEquals(0, pool.connectedRelaysCount()) } @Test - fun disconnectingThePoolClearsBothViewsEvenWhenTheCloseIsSilent() { - val (_, pool) = openedPool() + fun disconnectingThePoolDropsEveryRelayEvenWhenTheCloseIsSilent() { + val pool = openedPool() // The host is putting the client down (app backgrounded, connectivity lost). pool.disconnect() assertEquals(emptySet(), pool.connectedRelays.value) - assertEquals(emptySet(), pool.connectedRelayUrls()) assertEquals(setOf(url), pool.availableRelays.value, "still a member, just not connected") } } From 3e9b629b8da66e46f7f9f0c97798dab98582b39d Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 13 Sep 2026 13:02:47 +0000 Subject: [PATCH 7/7] refactor(relay): one adapter is one session; drop the lock and the socket identity checks The adapters compared the socket OkHttp named in each callback against the one they held, and took a monitor to make the compare-and-null atomic and to close a window in connect() where a callback could arrive before the field was assigned. Neither is needed: the relay client builds a fresh adapter per dial and OkHttp binds exactly one socket to the listener created in connect(), so anything that reaches that listener is from this session by construction. The only question a callback has to ask is whether the session already ended, which is one AtomicBoolean claimed by whichever of onClosed, onFailure or disconnect() gets there first. A compare-and-set keeps the "exactly one terminal report" guarantee without a monitor, and there is no assignment window left to guard. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_014cq6vrfQkASwxgqXXY4py8 --- .../service/okhttp/OkHttpWebSocket.kt | 57 ++++++++--------- .../sockets/okhttp/BasicOkHttpWebSocket.kt | 61 +++++++++---------- 2 files changed, 56 insertions(+), 62 deletions(-) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt index 33880cb1cb..e90277072f 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpWebSocket.kt @@ -34,24 +34,26 @@ import kotlinx.coroutines.launch import okhttp3.OkHttpClient import okhttp3.Request import okhttp3.Response +import java.util.concurrent.atomic.AtomicBoolean class OkHttpWebSocket( val url: NormalizedRelayUrl, val httpClient: (url: NormalizedRelayUrl) -> OkHttpClient, val out: WebSocketListener, ) : WebSocket { - private val lock = Any() private var usingOkHttp: OkHttpClient? = null - /** - * The OkHttp socket this adapter currently owns, or null once the session has ended -- by the - * relay closing it, by a network failure, or by [disconnect]. Only the owned socket may reach - * [out], and the terminal callbacks claim the slot under [lock], so a session ends with exactly - * one report however it ends. See quartz's `BasicOkHttpWebSocket` for the full reasoning; the - * two adapters differ only in how [needsReconnect] is decided. - */ @Volatile private var socket: okhttp3.WebSocket? = null + /** + * Set once, by whichever of `onClosed`, `onFailure` or [disconnect] ends the session first. + * One adapter is one session (the relay client builds a fresh one per dial, and OkHttp binds + * exactly one socket to the listener), so a callback only has to ask whether the session + * already ended. See quartz's `BasicOkHttpWebSocket` for the full reasoning; the two adapters + * differ only in how [needsReconnect] is decided. + */ + private val ended = AtomicBoolean(false) + fun buildRequest() = Request.Builder().url(url.url).build() override fun needsReconnect(): Boolean { @@ -77,13 +79,10 @@ class OkHttpWebSocket( } override fun connect() { + if (socket != null || ended.get()) return val client = httpClient(url) - // Under the lock so a callback racing this dial waits until the socket is owned rather - // than being dropped as foreign. - synchronized(lock) { - usingOkHttp = client - socket = client.newWebSocket(buildRequest(), OkHttpWebsocketListener(out)) - } + usingOkHttp = client + socket = client.newWebSocket(buildRequest(), OkHttpWebsocketListener(out)) } inner class OkHttpWebsocketListener( @@ -105,25 +104,21 @@ class OkHttpWebSocket( } } - /** Only the socket this adapter still owns may reach [out]. */ - private fun isOwned(webSocket: okhttp3.WebSocket) = synchronized(lock) { socket === webSocket } - /** Claims the session's single terminal report. False if it already ended. */ - private fun endSession(webSocket: okhttp3.WebSocket): Boolean { - val ended = synchronized(lock) { (socket === webSocket).also { if (it) socket = null } } - if (ended) { - incomingMessages.close() - job.cancel() - scope.cancel() - } - return ended + private fun endSession(): Boolean { + if (!ended.compareAndSet(false, true)) return false + socket = null + incomingMessages.close() + job.cancel() + scope.cancel() + return true } override fun onOpen( webSocket: okhttp3.WebSocket, response: Response, ) { - if (!isOwned(webSocket)) return + if (ended.get()) return out.onOpen( (response.receivedResponseAtMillis - response.sentRequestAtMillis).toInt(), response.headers["Sec-WebSocket-Extensions"]?.contains("permessage-deflate") ?: false, @@ -134,7 +129,7 @@ class OkHttpWebSocket( webSocket: okhttp3.WebSocket, text: String, ) { - if (!isOwned(webSocket)) return + if (ended.get()) return // Never blocks (unlimited channel): the OkHttp reader thread must // stay free to keep draining the socket. incomingMessages.trySendBlocking(text) @@ -162,7 +157,7 @@ class OkHttpWebSocket( code: Int, reason: String, ) { - if (!endSession(webSocket)) return + if (!endSession()) return out.onClosed(code, reason) } @@ -171,7 +166,7 @@ class OkHttpWebSocket( t: Throwable, response: Response?, ) { - if (!endSession(webSocket)) return + if (!endSession()) return out.onFailure(t, response?.code, response?.message) } } @@ -195,7 +190,9 @@ class OkHttpWebSocket( // waiting): OkHttp's cancel() raises no callback when no reader is left to fail, and when // it does the failure arrives later on its own thread. The relay client needs the answer // now, and must not hear from this socket again. - val closing = synchronized(lock) { socket?.also { socket = null } } ?: return + val closing = socket ?: return + if (!ended.compareAndSet(false, true)) return + socket = null closing.cancel() out.onClosed(1000, "client disconnect") } diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt index dc03b99661..20a07cd991 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/BasicOkHttpWebSocket.kt @@ -35,6 +35,7 @@ import kotlinx.coroutines.launch import okhttp3.OkHttpClient import okhttp3.Request import okhttp3.Response +import java.util.concurrent.atomic.AtomicBoolean import okhttp3.WebSocket as OkHttpWebSocket import okhttp3.WebSocketListener as OkHttpWebSocketListener @@ -51,25 +52,27 @@ class BasicOkHttpWebSocket( } } - private val lock = Any() + @Volatile private var socket: OkHttpWebSocket? = null /** - * The OkHttp socket this adapter currently owns, or null once the session has ended -- by the - * relay closing it, by a network failure, or by [disconnect]. + * Set once, by whichever of `onClosed`, `onFailure` or [disconnect] ends the session first. * - * OkHttp names the socket in every callback, and only the owned one may reach [out]. That is - * what makes this adapter honour the [WebSocket.disconnect] contract: after [disconnect] the - * slot is empty, so the failure OkHttp raises for its own `cancel()` on the reader thread, - * or the `onClosed` its writer thread delivers once a close handshake completes, is dropped - * instead of reaching a relay client that has already moved on to a new socket. The terminal - * callbacks claim the slot under [lock], so a session ends with exactly one report however it - * ends. + * One adapter is one session: the relay client builds a fresh one per dial, and OkHttp binds + * exactly one socket to the listener created in [connect], so anything that reaches that + * listener is from this session by construction. The only question a callback has to ask is + * whether the session already ended -- which is what keeps the [WebSocket.disconnect] contract: + * after [disconnect] the failure OkHttp raises for its own `cancel()` on the reader thread, or + * the `onClosed` its writer thread delivers once a close handshake completes, is dropped rather + * than reaching a relay client that has already moved on. Claimed with a compare-and-set so a + * [disconnect] racing a terminal callback still yields exactly one report. */ - @Volatile private var socket: OkHttpWebSocket? = null + private val ended = AtomicBoolean(false) override fun needsReconnect() = socket == null override fun connect() { + if (socket != null || ended.get()) return + val request = Request.Builder().url(url.url).build() val listener = @@ -93,25 +96,21 @@ class BasicOkHttpWebSocket( } } - /** Only the socket this adapter still owns may reach [out]. */ - private fun isOwned(webSocket: OkHttpWebSocket) = synchronized(lock) { socket === webSocket } - /** Claims the session's single terminal report. False if it already ended. */ - private fun endSession(webSocket: OkHttpWebSocket): Boolean { - val ended = synchronized(lock) { (socket === webSocket).also { if (it) socket = null } } - if (ended) { - incomingMessages.close() - job.cancel() - scope.cancel() - } - return ended + private fun endSession(): Boolean { + if (!ended.compareAndSet(false, true)) return false + socket = null + incomingMessages.close() + job.cancel() + scope.cancel() + return true } override fun onOpen( webSocket: OkHttpWebSocket, response: Response, ) { - if (!isOwned(webSocket)) return + if (ended.get()) return out.onOpen( (response.receivedResponseAtMillis - response.sentRequestAtMillis).toInt(), response.headers["Sec-WebSocket-Extensions"]?.contains("permessage-deflate") ?: false, @@ -122,7 +121,7 @@ class BasicOkHttpWebSocket( webSocket: OkHttpWebSocket, text: String, ) { - if (!isOwned(webSocket)) return + if (ended.get()) return // Never blocks (unlimited channel): the OkHttp reader // thread must stay free to keep draining the socket. incomingMessages.trySendBlocking(text) @@ -155,7 +154,7 @@ class BasicOkHttpWebSocket( code: Int, reason: String, ) { - if (!endSession(webSocket)) return + if (!endSession()) return out.onClosed(code, reason) } @@ -164,23 +163,21 @@ class BasicOkHttpWebSocket( t: Throwable, response: Response?, ) { - if (!endSession(webSocket)) return + if (!endSession()) return out.onFailure(t, response?.code, response?.message) } } - // Under the lock so a callback racing this dial (an instant failure lands on another - // thread) waits until the socket is owned, rather than being dropped as foreign. - synchronized(lock) { - socket = httpClient(url).newWebSocket(request, listener) - } + socket = httpClient(url).newWebSocket(request, listener) } override fun disconnect() { // Claim the session ourselves: OkHttp's cancel() raises no callback when no reader is // left to fail (the state a relay-initiated close leaves behind), and when it does the // failure arrives later on its own thread. The relay client needs the answer now. - val closing = synchronized(lock) { socket?.also { socket = null } } ?: return + val closing = socket ?: return + if (!ended.compareAndSet(false, true)) return + socket = null closing.cancel() out.onClosed(1000, "client disconnect") }