From e887e04f40cbfdb4c44c1f48cba60ffda9dae6d0 Mon Sep 17 00:00:00 2001 From: jack Date: Sun, 4 Jan 2026 13:23:32 -1000 Subject: [PATCH] fix: add thread safety to NostrTransport read receipt queue MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Synchronized access to readQueue and isSendingReadAcks using the existing concurrent DispatchQueue with barrier flags: - sendReadReceipt(): wrap enqueue in barrier async - processReadQueueIfNeeded(): extract item within barrier context - scheduleNextReadAck(): wrap callback in barrier async This fixes race conditions where concurrent calls could corrupt the read queue or cause check-then-act bugs on isSendingReadAcks. Also adds thread safety tests: - concurrentReadReceiptEnqueue: 100 concurrent enqueue operations - readQueueProcessingUnderLoad: concurrent enqueue during processing - isPeerReachableThreadSafety: concurrent read access test 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 --- bitchat/Services/NostrTransport.swift | 23 ++-- .../Services/NostrTransportTests.swift | 108 ++++++++++++++++++ 2 files changed, 122 insertions(+), 9 deletions(-) create mode 100644 bitchatTests/Services/NostrTransportTests.swift diff --git a/bitchat/Services/NostrTransport.swift b/bitchat/Services/NostrTransport.swift index 294274bc..3e2dccdc 100644 --- a/bitchat/Services/NostrTransport.swift +++ b/bitchat/Services/NostrTransport.swift @@ -149,8 +149,11 @@ final class NostrTransport: Transport, @unchecked Sendable { func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) { // Enqueue and process with throttling to avoid relay rate limits - readQueue.append(QueuedRead(receipt: receipt, peerID: peerID)) - processReadQueueIfNeeded() + // Use barrier to synchronize access to readQueue + queue.async(flags: .barrier) { [weak self] in + self?.readQueue.append(QueuedRead(receipt: receipt, peerID: peerID)) + self?.processReadQueueIfNeeded() + } } func sendFavoriteNotification(to peerID: PeerID, isFavorite: Bool) { @@ -260,16 +263,17 @@ extension NostrTransport { // MARK: - Private Helpers extension NostrTransport { + /// Must be called within a barrier on `queue` private func processReadQueueIfNeeded() { guard !isSendingReadAcks else { return } guard !readQueue.isEmpty else { return } isSendingReadAcks = true - sendNextReadAck() + let item = readQueue.removeFirst() + sendReadAckItem(item) } - private func sendNextReadAck() { - guard !readQueue.isEmpty else { isSendingReadAcks = false; return } - let item = readQueue.removeFirst() + /// Sends a single read ack item (called after extraction from queue within barrier) + private func sendReadAckItem(_ item: QueuedRead) { Task { @MainActor in guard let recipientNpub = resolveRecipientNpub(for: item.peerID) else { scheduleNextReadAck(); return } guard let senderIdentity = try? idBridge.getCurrentNostrIdentity() else { scheduleNextReadAck(); return } @@ -301,9 +305,10 @@ extension NostrTransport { private func scheduleNextReadAck() { DispatchQueue.main.asyncAfter(deadline: .now() + readAckInterval) { [weak self] in - guard let self = self else { return } - self.isSendingReadAcks = false - self.processReadQueueIfNeeded() + self?.queue.async(flags: .barrier) { [weak self] in + self?.isSendingReadAcks = false + self?.processReadQueueIfNeeded() + } } } diff --git a/bitchatTests/Services/NostrTransportTests.swift b/bitchatTests/Services/NostrTransportTests.swift new file mode 100644 index 00000000..ce92228d --- /dev/null +++ b/bitchatTests/Services/NostrTransportTests.swift @@ -0,0 +1,108 @@ +// +// NostrTransportTests.swift +// bitchat +// +// This is free and unencumbered software released into the public domain. +// For more information, see +// + +import Foundation +import Testing +@testable import bitchat + +@Suite("NostrTransport Thread Safety Tests") +struct NostrTransportTests { + + @Test("Concurrent read receipt enqueue does not crash") + @MainActor + func concurrentReadReceiptEnqueue() async throws { + let keychain = MockKeychain() + let idBridge = NostrIdentityBridge(keychain: keychain) + let transport = NostrTransport(keychain: keychain, idBridge: idBridge) + + // Create 100 concurrent read receipt submissions + let iterations = 100 + + await withTaskGroup(of: Void.self) { group in + for i in 0..