diff --git a/bitchat/Services/MessageDeduplicationService.swift b/bitchat/Services/MessageDeduplicationService.swift index 88f249ed..b513bb87 100644 --- a/bitchat/Services/MessageDeduplicationService.swift +++ b/bitchat/Services/MessageDeduplicationService.swift @@ -12,6 +12,8 @@ import Foundation /// Generic LRU (Least Recently Used) cache for deduplication. /// Uses an efficient O(1) lookup with periodic compaction. +/// Thread-safe via @MainActor - all callers are already on main actor. +@MainActor final class LRUDeduplicationCache { private var map: [String: Value] = [:] private var order: [String] = [] @@ -157,6 +159,8 @@ enum ContentNormalizer { /// Service that manages message deduplication using LRU caches. /// Provides separate caches for content-based dedup and Nostr event ID dedup. +/// Thread-safe via @MainActor - all callers are already on main actor. +@MainActor final class MessageDeduplicationService { /// Cache for content-based near-duplicate detection 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/MessageDeduplicationServiceTests.swift b/bitchatTests/MessageDeduplicationServiceTests.swift index 8f265afb..251fe808 100644 --- a/bitchatTests/MessageDeduplicationServiceTests.swift +++ b/bitchatTests/MessageDeduplicationServiceTests.swift @@ -12,6 +12,8 @@ import Foundation // MARK: - LRU Deduplication Cache Tests +@Suite("LRU Deduplication Cache") +@MainActor struct LRUDeduplicationCacheTests { // MARK: - Basic Operations @@ -265,6 +267,8 @@ struct ContentNormalizerTests { // MARK: - Message Deduplication Service Tests +@Suite("Message Deduplication Service") +@MainActor struct MessageDeduplicationServiceTests { // MARK: - Content Deduplication @@ -467,4 +471,134 @@ struct MessageDeduplicationServiceTests { #expect(service.contentTimestamp(for: "hello world") == now) #expect(service.contentTimestamp(for: "Hello World") == now) } + + // MARK: - Thread Safety Tests (via @MainActor enforcement) + + @Test("Concurrent content recording is safe via MainActor") + func concurrentContentRecording() async { + let service = MessageDeduplicationService(contentCapacity: 1000, nostrEventCapacity: 1000) + let iterations = 100 + + // All operations run on MainActor due to @MainActor annotation + // This test verifies the pattern works correctly + await withTaskGroup(of: Void.self) { group in + for i in 0..(capacity: 500) + let iterations = 100 + + await withTaskGroup(of: Void.self) { group in + // Write tasks + for i in 0..(capacity: 10) + let iterations = 100 + + await withTaskGroup(of: Void.self) { group in + for i in 0.. +// + +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..