mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-25 16:05:19 +00:00
fix: add thread safety to NostrTransport read receipt queue
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 <noreply@anthropic.com>
This commit is contained in:
@@ -149,8 +149,11 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
|||||||
|
|
||||||
func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) {
|
func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) {
|
||||||
// Enqueue and process with throttling to avoid relay rate limits
|
// Enqueue and process with throttling to avoid relay rate limits
|
||||||
readQueue.append(QueuedRead(receipt: receipt, peerID: peerID))
|
// Use barrier to synchronize access to readQueue
|
||||||
processReadQueueIfNeeded()
|
queue.async(flags: .barrier) { [weak self] in
|
||||||
|
self?.readQueue.append(QueuedRead(receipt: receipt, peerID: peerID))
|
||||||
|
self?.processReadQueueIfNeeded()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func sendFavoriteNotification(to peerID: PeerID, isFavorite: Bool) {
|
func sendFavoriteNotification(to peerID: PeerID, isFavorite: Bool) {
|
||||||
@@ -260,16 +263,17 @@ extension NostrTransport {
|
|||||||
// MARK: - Private Helpers
|
// MARK: - Private Helpers
|
||||||
|
|
||||||
extension NostrTransport {
|
extension NostrTransport {
|
||||||
|
/// Must be called within a barrier on `queue`
|
||||||
private func processReadQueueIfNeeded() {
|
private func processReadQueueIfNeeded() {
|
||||||
guard !isSendingReadAcks else { return }
|
guard !isSendingReadAcks else { return }
|
||||||
guard !readQueue.isEmpty else { return }
|
guard !readQueue.isEmpty else { return }
|
||||||
isSendingReadAcks = true
|
isSendingReadAcks = true
|
||||||
sendNextReadAck()
|
let item = readQueue.removeFirst()
|
||||||
|
sendReadAckItem(item)
|
||||||
}
|
}
|
||||||
|
|
||||||
private func sendNextReadAck() {
|
/// Sends a single read ack item (called after extraction from queue within barrier)
|
||||||
guard !readQueue.isEmpty else { isSendingReadAcks = false; return }
|
private func sendReadAckItem(_ item: QueuedRead) {
|
||||||
let item = readQueue.removeFirst()
|
|
||||||
Task { @MainActor in
|
Task { @MainActor in
|
||||||
guard let recipientNpub = resolveRecipientNpub(for: item.peerID) else { scheduleNextReadAck(); return }
|
guard let recipientNpub = resolveRecipientNpub(for: item.peerID) else { scheduleNextReadAck(); return }
|
||||||
guard let senderIdentity = try? idBridge.getCurrentNostrIdentity() else { scheduleNextReadAck(); return }
|
guard let senderIdentity = try? idBridge.getCurrentNostrIdentity() else { scheduleNextReadAck(); return }
|
||||||
@@ -301,9 +305,10 @@ extension NostrTransport {
|
|||||||
|
|
||||||
private func scheduleNextReadAck() {
|
private func scheduleNextReadAck() {
|
||||||
DispatchQueue.main.asyncAfter(deadline: .now() + readAckInterval) { [weak self] in
|
DispatchQueue.main.asyncAfter(deadline: .now() + readAckInterval) { [weak self] in
|
||||||
guard let self = self else { return }
|
self?.queue.async(flags: .barrier) { [weak self] in
|
||||||
self.isSendingReadAcks = false
|
self?.isSendingReadAcks = false
|
||||||
self.processReadQueueIfNeeded()
|
self?.processReadQueueIfNeeded()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,108 @@
|
|||||||
|
//
|
||||||
|
// NostrTransportTests.swift
|
||||||
|
// bitchat
|
||||||
|
//
|
||||||
|
// This is free and unencumbered software released into the public domain.
|
||||||
|
// For more information, see <https://unlicense.org>
|
||||||
|
//
|
||||||
|
|
||||||
|
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..<iterations {
|
||||||
|
group.addTask {
|
||||||
|
let receipt = ReadReceipt(
|
||||||
|
originalMessageID: UUID().uuidString,
|
||||||
|
readerID: PeerID(str: String(format: "%016x", i)),
|
||||||
|
readerNickname: "Reader\(i)"
|
||||||
|
)
|
||||||
|
let peerID = PeerID(str: String(format: "%016x", i))
|
||||||
|
transport.sendReadReceipt(receipt, to: peerID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// If we reach here without crashing, the test passes
|
||||||
|
// The concurrent enqueue operations completed without data races
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test("Read queue processes under concurrent load")
|
||||||
|
@MainActor
|
||||||
|
func readQueueProcessingUnderLoad() async throws {
|
||||||
|
let keychain = MockKeychain()
|
||||||
|
let idBridge = NostrIdentityBridge(keychain: keychain)
|
||||||
|
let transport = NostrTransport(keychain: keychain, idBridge: idBridge)
|
||||||
|
|
||||||
|
// Rapidly enqueue many receipts from multiple concurrent sources
|
||||||
|
let iterations = 50
|
||||||
|
|
||||||
|
// First batch - rapid fire
|
||||||
|
for i in 0..<iterations {
|
||||||
|
let receipt = ReadReceipt(
|
||||||
|
originalMessageID: UUID().uuidString,
|
||||||
|
readerID: PeerID(str: String(format: "%016x", i)),
|
||||||
|
readerNickname: "Reader\(i)"
|
||||||
|
)
|
||||||
|
transport.sendReadReceipt(receipt, to: PeerID(str: String(format: "%016x", i)))
|
||||||
|
}
|
||||||
|
|
||||||
|
// Give some time for processing to start
|
||||||
|
try await Task.sleep(nanoseconds: 100_000_000) // 100ms
|
||||||
|
|
||||||
|
// Second batch - while first might be processing
|
||||||
|
await withTaskGroup(of: Void.self) { group in
|
||||||
|
for i in iterations..<(iterations * 2) {
|
||||||
|
group.addTask {
|
||||||
|
let receipt = ReadReceipt(
|
||||||
|
originalMessageID: UUID().uuidString,
|
||||||
|
readerID: PeerID(str: String(format: "%016x", i)),
|
||||||
|
readerNickname: "Reader\(i)"
|
||||||
|
)
|
||||||
|
transport.sendReadReceipt(receipt, to: PeerID(str: String(format: "%016x", i)))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// If we reach here without crashing or deadlocking, test passes
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test("isPeerReachable is thread safe")
|
||||||
|
@MainActor
|
||||||
|
func isPeerReachableThreadSafety() async throws {
|
||||||
|
let keychain = MockKeychain()
|
||||||
|
let idBridge = NostrIdentityBridge(keychain: keychain)
|
||||||
|
let transport = NostrTransport(keychain: keychain, idBridge: idBridge)
|
||||||
|
|
||||||
|
let iterations = 100
|
||||||
|
|
||||||
|
// Concurrent reads on isPeerReachable
|
||||||
|
await withTaskGroup(of: Bool.self) { group in
|
||||||
|
for i in 0..<iterations {
|
||||||
|
group.addTask {
|
||||||
|
let peerID = PeerID(str: String(format: "%016x", i))
|
||||||
|
return transport.isPeerReachable(peerID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Collect results (all should be false since no favorites configured)
|
||||||
|
for await result in group {
|
||||||
|
#expect(result == false)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user