From ef8b9be0e88efe93e30a39575dc09dd965d8afca Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 30 Mar 2026 17:37:23 +0000 Subject: [PATCH] refactor: replace synchronized with Channel + AtomicBoolean in BLE relay classes Use Channel(UNLIMITED) as a lock-free send queue and AtomicBoolean with atomic exchange for the isSending flag, eliminating all synchronized blocks while maintaining thread safety via a double-check pattern. https://claude.ai/code/session_01KLRz7FJgWRtgCcZe3qD1VX --- .../quartz/nipBEBle/relay/BleNostrClient.kt | 51 ++++++++++--------- .../quartz/nipBEBle/relay/BleNostrServer.kt | 46 +++++++++-------- 2 files changed, 53 insertions(+), 44 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrClient.kt index 371b701d7..2c75d9576 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrClient.kt @@ -31,6 +31,10 @@ import com.vitorpamplona.quartz.nipBEBle.protocol.BleChunkAssembler import com.vitorpamplona.quartz.nipBEBle.protocol.BleMessageChunker import com.vitorpamplona.quartz.nipBEBle.transport.BleTransport import com.vitorpamplona.quartz.utils.Log +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.channels.Channel +import kotlin.concurrent.atomics.AtomicBoolean +import kotlin.concurrent.atomics.ExperimentalAtomicApi /** * A Nostr relay client that communicates over BLE instead of WebSockets. @@ -54,6 +58,7 @@ import com.vitorpamplona.quartz.utils.Log * client.sendIfConnected(EventCmd(myEvent)) * ``` */ +@OptIn(ExperimentalAtomicApi::class, ExperimentalCoroutinesApi::class) class BleNostrClient( val peer: BlePeer, val transport: BleTransport, @@ -64,8 +69,8 @@ class BleNostrClient( private var connected = false private val assembler = BleChunkAssembler() - private val sendQueue = ArrayDeque>() - private var isSending = false + private val sendQueue = Channel>(Channel.UNLIMITED) + private val isSending = AtomicBoolean(false) override fun isConnected(): Boolean = connected @@ -96,12 +101,9 @@ class BleNostrClient( val json = OptimizedJsonMapper.toJson(cmd) val chunks = BleMessageChunker.splitIntoChunks(json, chunkSize) - synchronized(sendQueue) { - sendQueue.addLast(chunks) - if (!isSending) { - isSending = true - sendNextMessage() - } + sendQueue.trySend(chunks) + if (!isSending.exchange(true)) { + sendNextMessage() } listener.onSent(this, json, cmd, true) @@ -110,10 +112,7 @@ class BleNostrClient( override fun disconnect() { connected = false assembler.reset() - synchronized(sendQueue) { - sendQueue.clear() - isSending = false - } + drainQueue() transport.disconnectFromPeer(peer) listener.onDisconnected(this) } @@ -132,10 +131,7 @@ class BleNostrClient( fun onDisconnected() { connected = false assembler.reset() - synchronized(sendQueue) { - sendQueue.clear() - isSending = false - } + drainQueue() listener.onDisconnected(this) } @@ -165,21 +161,28 @@ class BleNostrClient( private var currentChunkIndex = 0 private fun sendNextMessage() { - val chunks = - synchronized(sendQueue) { - sendQueue.removeFirstOrNull() - } - if (chunks == null) { - synchronized(sendQueue) { - isSending = false + val result = sendQueue.tryReceive() + if (result.isFailure) { + isSending.store(false) + // Double-check: an item may have been added after tryReceive + // but before we cleared the flag. + if (sendQueue.isEmpty) return + if (!isSending.exchange(true)) { + sendNextMessage() } return } - currentChunks = chunks + currentChunks = result.getOrNull() currentChunkIndex = 0 sendNextChunk() } + private fun drainQueue() { + while (sendQueue.tryReceive().isSuccess) {} + currentChunks = null + isSending.store(false) + } + private fun sendNextChunk() { val chunks = currentChunks ?: return if (currentChunkIndex >= chunks.size) { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrServer.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrServer.kt index 2a0b493e6..5a8d6d39b 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrServer.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipBEBle/relay/BleNostrServer.kt @@ -29,6 +29,10 @@ import com.vitorpamplona.quartz.nipBEBle.protocol.BleChunkAssembler import com.vitorpamplona.quartz.nipBEBle.protocol.BleMessageChunker import com.vitorpamplona.quartz.nipBEBle.transport.BleTransport import com.vitorpamplona.quartz.utils.Log +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.channels.Channel +import kotlin.concurrent.atomics.AtomicBoolean +import kotlin.concurrent.atomics.ExperimentalAtomicApi /** * Handles the server (relay) side of a BLE Nostr connection. @@ -52,6 +56,7 @@ import com.vitorpamplona.quartz.utils.Log * server.sendMessage(EventMessage(subId, event)) * ``` */ +@OptIn(ExperimentalAtomicApi::class, ExperimentalCoroutinesApi::class) class BleNostrServer( val peer: BlePeer, val transport: BleTransport, @@ -60,8 +65,8 @@ class BleNostrServer( ) { private var connected = false private val assembler = BleChunkAssembler() - private val sendQueue = ArrayDeque>() - private var isSending = false + private val sendQueue = Channel>(Channel.UNLIMITED) + private val isSending = AtomicBoolean(false) fun isConnected(): Boolean = connected @@ -79,10 +84,7 @@ class BleNostrServer( fun onClientDisconnected() { connected = false assembler.reset() - synchronized(sendQueue) { - sendQueue.clear() - isSending = false - } + drainQueue() listener.onClientDisconnected(this) } @@ -110,12 +112,9 @@ class BleNostrServer( val json = OptimizedJsonMapper.toJson(msg) val chunks = BleMessageChunker.splitIntoChunks(json, chunkSize) - synchronized(sendQueue) { - sendQueue.addLast(chunks) - if (!isSending) { - isSending = true - sendNextMessage() - } + sendQueue.trySend(chunks) + if (!isSending.exchange(true)) { + sendNextMessage() } } @@ -139,21 +138,28 @@ class BleNostrServer( private var currentChunkIndex = 0 private fun sendNextMessage() { - val chunks = - synchronized(sendQueue) { - sendQueue.removeFirstOrNull() - } - if (chunks == null) { - synchronized(sendQueue) { - isSending = false + val result = sendQueue.tryReceive() + if (result.isFailure) { + isSending.store(false) + // Double-check: an item may have been added after tryReceive + // but before we cleared the flag. + if (sendQueue.isEmpty) return + if (!isSending.exchange(true)) { + sendNextMessage() } return } - currentChunks = chunks + currentChunks = result.getOrNull() currentChunkIndex = 0 sendNextChunk() } + private fun drainQueue() { + while (sendQueue.tryReceive().isSuccess) {} + currentChunks = null + isSending.store(false) + } + private fun sendNextChunk() { val chunks = currentChunks ?: return if (currentChunkIndex >= chunks.size) {