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
This commit is contained in:
Claude
2026-03-30 17:37:23 +00:00
parent c44f59b1ce
commit ef8b9be0e8
2 changed files with 53 additions and 44 deletions
@@ -31,6 +31,10 @@ import com.vitorpamplona.quartz.nipBEBle.protocol.BleChunkAssembler
import com.vitorpamplona.quartz.nipBEBle.protocol.BleMessageChunker import com.vitorpamplona.quartz.nipBEBle.protocol.BleMessageChunker
import com.vitorpamplona.quartz.nipBEBle.transport.BleTransport import com.vitorpamplona.quartz.nipBEBle.transport.BleTransport
import com.vitorpamplona.quartz.utils.Log 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. * 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)) * client.sendIfConnected(EventCmd(myEvent))
* ``` * ```
*/ */
@OptIn(ExperimentalAtomicApi::class, ExperimentalCoroutinesApi::class)
class BleNostrClient( class BleNostrClient(
val peer: BlePeer, val peer: BlePeer,
val transport: BleTransport, val transport: BleTransport,
@@ -64,8 +69,8 @@ class BleNostrClient(
private var connected = false private var connected = false
private val assembler = BleChunkAssembler() private val assembler = BleChunkAssembler()
private val sendQueue = ArrayDeque<Array<ByteArray>>() private val sendQueue = Channel<Array<ByteArray>>(Channel.UNLIMITED)
private var isSending = false private val isSending = AtomicBoolean(false)
override fun isConnected(): Boolean = connected override fun isConnected(): Boolean = connected
@@ -96,12 +101,9 @@ class BleNostrClient(
val json = OptimizedJsonMapper.toJson(cmd) val json = OptimizedJsonMapper.toJson(cmd)
val chunks = BleMessageChunker.splitIntoChunks(json, chunkSize) val chunks = BleMessageChunker.splitIntoChunks(json, chunkSize)
synchronized(sendQueue) { sendQueue.trySend(chunks)
sendQueue.addLast(chunks) if (!isSending.exchange(true)) {
if (!isSending) { sendNextMessage()
isSending = true
sendNextMessage()
}
} }
listener.onSent(this, json, cmd, true) listener.onSent(this, json, cmd, true)
@@ -110,10 +112,7 @@ class BleNostrClient(
override fun disconnect() { override fun disconnect() {
connected = false connected = false
assembler.reset() assembler.reset()
synchronized(sendQueue) { drainQueue()
sendQueue.clear()
isSending = false
}
transport.disconnectFromPeer(peer) transport.disconnectFromPeer(peer)
listener.onDisconnected(this) listener.onDisconnected(this)
} }
@@ -132,10 +131,7 @@ class BleNostrClient(
fun onDisconnected() { fun onDisconnected() {
connected = false connected = false
assembler.reset() assembler.reset()
synchronized(sendQueue) { drainQueue()
sendQueue.clear()
isSending = false
}
listener.onDisconnected(this) listener.onDisconnected(this)
} }
@@ -165,21 +161,28 @@ class BleNostrClient(
private var currentChunkIndex = 0 private var currentChunkIndex = 0
private fun sendNextMessage() { private fun sendNextMessage() {
val chunks = val result = sendQueue.tryReceive()
synchronized(sendQueue) { if (result.isFailure) {
sendQueue.removeFirstOrNull() isSending.store(false)
} // Double-check: an item may have been added after tryReceive
if (chunks == null) { // but before we cleared the flag.
synchronized(sendQueue) { if (sendQueue.isEmpty) return
isSending = false if (!isSending.exchange(true)) {
sendNextMessage()
} }
return return
} }
currentChunks = chunks currentChunks = result.getOrNull()
currentChunkIndex = 0 currentChunkIndex = 0
sendNextChunk() sendNextChunk()
} }
private fun drainQueue() {
while (sendQueue.tryReceive().isSuccess) {}
currentChunks = null
isSending.store(false)
}
private fun sendNextChunk() { private fun sendNextChunk() {
val chunks = currentChunks ?: return val chunks = currentChunks ?: return
if (currentChunkIndex >= chunks.size) { if (currentChunkIndex >= chunks.size) {
@@ -29,6 +29,10 @@ import com.vitorpamplona.quartz.nipBEBle.protocol.BleChunkAssembler
import com.vitorpamplona.quartz.nipBEBle.protocol.BleMessageChunker import com.vitorpamplona.quartz.nipBEBle.protocol.BleMessageChunker
import com.vitorpamplona.quartz.nipBEBle.transport.BleTransport import com.vitorpamplona.quartz.nipBEBle.transport.BleTransport
import com.vitorpamplona.quartz.utils.Log 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. * 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)) * server.sendMessage(EventMessage(subId, event))
* ``` * ```
*/ */
@OptIn(ExperimentalAtomicApi::class, ExperimentalCoroutinesApi::class)
class BleNostrServer( class BleNostrServer(
val peer: BlePeer, val peer: BlePeer,
val transport: BleTransport, val transport: BleTransport,
@@ -60,8 +65,8 @@ class BleNostrServer(
) { ) {
private var connected = false private var connected = false
private val assembler = BleChunkAssembler() private val assembler = BleChunkAssembler()
private val sendQueue = ArrayDeque<Array<ByteArray>>() private val sendQueue = Channel<Array<ByteArray>>(Channel.UNLIMITED)
private var isSending = false private val isSending = AtomicBoolean(false)
fun isConnected(): Boolean = connected fun isConnected(): Boolean = connected
@@ -79,10 +84,7 @@ class BleNostrServer(
fun onClientDisconnected() { fun onClientDisconnected() {
connected = false connected = false
assembler.reset() assembler.reset()
synchronized(sendQueue) { drainQueue()
sendQueue.clear()
isSending = false
}
listener.onClientDisconnected(this) listener.onClientDisconnected(this)
} }
@@ -110,12 +112,9 @@ class BleNostrServer(
val json = OptimizedJsonMapper.toJson(msg) val json = OptimizedJsonMapper.toJson(msg)
val chunks = BleMessageChunker.splitIntoChunks(json, chunkSize) val chunks = BleMessageChunker.splitIntoChunks(json, chunkSize)
synchronized(sendQueue) { sendQueue.trySend(chunks)
sendQueue.addLast(chunks) if (!isSending.exchange(true)) {
if (!isSending) { sendNextMessage()
isSending = true
sendNextMessage()
}
} }
} }
@@ -139,21 +138,28 @@ class BleNostrServer(
private var currentChunkIndex = 0 private var currentChunkIndex = 0
private fun sendNextMessage() { private fun sendNextMessage() {
val chunks = val result = sendQueue.tryReceive()
synchronized(sendQueue) { if (result.isFailure) {
sendQueue.removeFirstOrNull() isSending.store(false)
} // Double-check: an item may have been added after tryReceive
if (chunks == null) { // but before we cleared the flag.
synchronized(sendQueue) { if (sendQueue.isEmpty) return
isSending = false if (!isSending.exchange(true)) {
sendNextMessage()
} }
return return
} }
currentChunks = chunks currentChunks = result.getOrNull()
currentChunkIndex = 0 currentChunkIndex = 0
sendNextChunk() sendNextChunk()
} }
private fun drainQueue() {
while (sendQueue.tryReceive().isSuccess) {}
currentChunks = null
isSending.store(false)
}
private fun sendNextChunk() { private fun sendNextChunk() {
val chunks = currentChunks ?: return val chunks = currentChunks ?: return
if (currentChunkIndex >= chunks.size) { if (currentChunkIndex >= chunks.size) {