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..