From 87cbf5a03c2bbb28b26d33793e07ab10d7b9ad70 Mon Sep 17 00:00:00 2001 From: jack Date: Mon, 27 Jul 2026 00:49:27 +0100 Subject: [PATCH] Unify the message and collections queues into one serial engine queue MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The old model ran a concurrent message queue over a second concurrent collections queue whose barrier flags served as the real mutual exclusion — every field carried an ownership comment, and correctness lived in per-site discipline. The message queue is now a single serial engine queue that owns all mesh protocol state; the collections queue, its 98 sync/async hops, and every barrier flag are gone. Cross-thread callers go through onEngine, which documents and (in debug) enforces the transport's sync-edge order: main and test threads may block on the engine, the engine may block on bleQueue and the crypto/identity queues, and nothing may block the other way. The debug trap caught two latent inversions the leaf-lock structure had been masking: the verified-announce rebind path re-resolved the ingress link through the engine from inside its bleQueue critical section (it now receives the already-resolved link), and the noise session-generation closures sync-re-entered the engine from the noise manager's queue while their own engine slot was blocked on it (they now touch engine state directly, which the held slot makes exclusive). BLE throughput is orders of magnitude below what one serial queue sustains; the full suite runs at identical speed. Co-Authored-By: Claude Fable 5 --- bitchat/Services/BLE/BLEService.swift | 407 ++++++++++++++------------ 1 file changed, 219 insertions(+), 188 deletions(-) diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index e2df4097..3cc9f236 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -255,8 +255,8 @@ final class BLEService: NSObject { private var subscriptionAnnounceLimiter = BLESubscriptionAnnounceLimiter() // 3. Peer Information (single source of truth). Lock-backed so the main - // actor reads it directly instead of blocking on collectionsQueue; - // mutations still only run on the transport's queues. + // actor reads it directly instead of blocking on the engine queue; + // mutations still only run on the engine. private let peerRegistry = BLEPeerRegistryStore() // 4. Efficient Message Deduplication @@ -280,14 +280,14 @@ final class BLEService: NSObject { // Verified one-time prekey bundles gossiped by other peers, used to seal // courier mail forward-secretly. Injectable for tests. var prekeyBundleStore: PrekeyBundleStore = .shared - // Throttle for re-broadcasting our own (unchanged) bundle; guarded by - // collectionsQueue barriers. + // Throttle for re-broadcasting our own (unchanged) bundle + // (engine-confined). private var lastPrekeyBundleSentAt: Date? // Prekey bundles that arrived before their owner's verified announce bound - // a signing key. The receive queue is concurrent, so a bundle can race - // ahead of the announce it depends on; we retain the latest such bundle per - // owner (bounded) and re-attempt attribution when the announce lands. - // Guarded by collectionsQueue barriers. + // a signing key. Over the air a bundle can still arrive before the + // announce it depends on; we retain the latest such bundle per owner + // (bounded) and re-attempt attribution when the announce lands. + // Engine-confined. private var pendingPrekeyBundles: [PeerID: BitchatPacket] = [:] private static let pendingPrekeyBundleCap = 64 // Gateway mode: sink for received nostrCarrier packets (set by app @@ -322,15 +322,15 @@ final class BLEService: NSObject { #endif private var selfBroadcastTracker = BLESelfBroadcastTracker() private let meshTopology = MeshTopologyTracker() - // Route health for originated source routes; guarded by collectionsQueue. + // Route health for originated source routes (engine-confined). private var sourceRouteFailures = BLESourceRouteFailureCache() // Mesh diagnostics: outstanding /ping probes keyed by nonce, plus the // inbound ping budget — keyed by the ingress link (the directly connected // peer that delivered the packet), since the unsigned claimed sender is // spoofable — so a directed unencrypted probe cannot be turned into an - // amplification primitive. Both are owned by collectionsQueue barriers - // like the other mutable collections. + // amplification primitive. Both are engine-confined like the other + // mutable collections. private struct PendingMeshPing { let peerID: PeerID let sentAt: Date @@ -350,7 +350,7 @@ final class BLEService: NSObject { private lazy var privateMediaTransferAdmissions = BLEPrivateMediaTransferAdmissionRegistry { [weak self] transferId in self?.handlePrivateMediaAdmissionExpiry(transferId) } - // All six maps below are protected by `collectionsQueue`. A fresh Noise + // All six maps below are engine-confined. A fresh Noise // authentication rotates the generation UUID, so stale proof timers and // proof packets cannot classify a replacement session. private var privateMediaSessionGenerations: [PeerID: UUID] = [:] @@ -404,11 +404,35 @@ final class BLEService: NSObject { // MARK: - Queues - private let messageQueue = DispatchQueue(label: "mesh.message", attributes: .concurrent) - private let collectionsQueue = DispatchQueue(label: "mesh.collections", attributes: .concurrent) + /// The engine queue: one serial domain that owns every piece of mesh + /// protocol state (the former concurrent message queue and the separate + /// collections queue it guarded state with). BLE throughput is far below + /// what one queue serializes comfortably, and a single writer makes the + /// old per-field ownership comments and barrier discipline structural. + private let messageQueue = DispatchQueue(label: "mesh.message") private let messageQueueKey = DispatchSpecificKey() private let bleQueue = DispatchQueue(label: "mesh.bluetooth", qos: .userInitiated) private let bleQueueKey = DispatchSpecificKey() + + /// Runs `body` exclusively with respect to all engine-owned state. + /// Executes inline when already on the engine queue; otherwise blocks + /// until the engine drains the work ahead of it. + /// + /// Sync-edge order (deadlock freedom): main and test threads may + /// sync-wait on the engine; the engine sync-waits on bleQueue + /// (`readLinkState`) and on the crypto/identity services' internal + /// queues. None of those may ever sync-wait back on the engine — + /// bleQueue callers hop with `messageQueue.async` instead, and debug + /// builds trap any violation here. + private func onEngine(_ body: () -> T) -> T { + #if DEBUG + dispatchPrecondition(condition: .notOnQueue(bleQueue)) + #endif + if DispatchQueue.getSpecific(key: messageQueueKey) != nil { + return body() + } + return messageQueue.sync(execute: body) + } // Noise messages and typed payloads pending handshake completion. private var pendingNoiseSessionQueues = BLENoiseSessionQueues() @@ -436,7 +460,7 @@ final class BLEService: NSObject { // ephemeral key, and bridge drops multiply across relays/couriers), so // envelope-level dedup can't catch them; dedup on the inner ID before // delivery so a duplicate costs one decrypt instead of a delivery + ack - // + handshake each. Owned by collectionsQueue barriers. + // + handshake each. Engine-confined. private var openedCourierMessageIDs = BoundedIDSet(capacity: TransportConfig.courierOpenedMessageIDCap) private let logRateLimiter = BLELogRateLimiter(defaultMinimumInterval: 5) @@ -735,7 +759,7 @@ final class BLEService: NSObject { // generation-bound handoffs that raced this barrier reject themselves. // Clear the old identity's bounded early-ciphertext queue again after // those callbacks drain so none can repopulate it after the first wipe. - messageQueue.sync(flags: .barrier) { + messageQueue.sync { noisePacketHandler.resetForPanic() } clearEmergencySessionState() @@ -763,16 +787,14 @@ final class BLEService: NSObject { gossipSyncManager = nil // Discard deferred pre-panic ciphertext behind any in-flight receive // handlers so none can repopulate the handler's bounded queue. - messageQueue.sync(flags: .barrier) { + messageQueue.sync { noisePacketHandler.resetForPanic() } - // pendingNoiseSessionQueues is owned by collectionsQueue everywhere - // else, so clear it there too rather than on messageQueue. - collectionsQueue.sync(flags: .barrier) { + onEngine { pendingNoiseSessionQueues.removeAll() } - let panicReset = collectionsQueue.sync(flags: .barrier) { + let panicReset = onEngine { let transfers = outboundFragmentTransfers.removeAll() fragmentAssemblyBuffer.removeAll() pendingDirectedRelays.removeAll() @@ -812,7 +834,7 @@ final class BLEService: NSObject { // must never observe the new Noise service alongside the old peer ID // (it would sign with the new identity while carrying the old sender). // refreshPeerIdentity() executes inline here via its re-entrancy check. - messageQueue.sync(flags: .barrier) { + messageQueue.sync { noiseService.clearEphemeralStateForPanic() noiseService.clearPersistentIdentity() @@ -830,7 +852,7 @@ final class BLEService: NSObject { // would force-send an announce and break that silence). localIdentityState.setNickname(currentNickname) messageDeduplicator.reset() - messageQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in self?.selfBroadcastTracker.removeAll() } requestPeerDataPublish() @@ -1069,7 +1091,7 @@ final class BLEService: NSObject { private func clearEmergencySessionState() { // Clear all sessions and peers - let cancelled = collectionsQueue.sync(flags: .barrier) { + let cancelled = onEngine { let entries = outboundFragmentTransfers.removeAll().map { (id: $0.id, items: $0.workItems) } @@ -1149,7 +1171,7 @@ final class BLEService: NSObject { let normalizedPeerID = peerID.toShort() let currentNoiseGeneration = noiseService.sessionGeneration(for: normalizedPeerID) - return collectionsQueue.sync { + return onEngine { guard let generation = privateMediaSessionGenerations[normalizedPeerID], generation == currentNoiseGeneration, @@ -1194,7 +1216,7 @@ final class BLEService: NSObject { sessionGeneration: UUID?, authenticatedState: BLEAuthenticatedPeerStateObservation?, timedOut: BLEPrivateMediaProofTimeoutMarker? - ) = collectionsQueue.sync { + ) = onEngine { let info = peerRegistry.info(for: normalizedPeerID) return ( info?.capabilities ?? [], @@ -1270,7 +1292,7 @@ final class BLEService: NSObject { return } - let generation = self.collectionsQueue.sync { + let generation = onEngine { self.privateMediaSessionGenerations[normalizedPeerID] } let fingerprint = self.privateMediaPolicyFingerprint( @@ -1283,7 +1305,7 @@ final class BLEService: NSObject { } let requestID = UUID() - let registration = self.collectionsQueue.sync(flags: .barrier) { + let registration = onEngine { () -> (registered: Bool, shouldSchedule: Bool, nonce: UUID, generation: UUID?) in let generation = self.privateMediaSessionGenerations[normalizedPeerID] if var pending = self.pendingPrivateMediaPolicyResolutions[normalizedPeerID] { @@ -1374,7 +1396,7 @@ final class BLEService: NSObject { sessionGeneration: UUID?, nonce: UUID ) { - let expiration = collectionsQueue.sync(flags: .barrier) { + let expiration = onEngine { () -> (expired: Bool, completions: [@MainActor (PrivateMediaSendPolicy) -> Void]) in let pending = pendingPrivateMediaPolicyResolutions[peerID] let pendingMatches = pending?.timeoutNonce == nonce @@ -1512,10 +1534,10 @@ final class BLEService: NSObject { // MARK: Messaging private func handlePrivateMediaAdmissionExpiry(_ transferId: String) { - // Expiry can be discovered from the BLE maintenance queue or while a - // caller already owns collectionsQueue. Cleanup is therefore - // fire-and-forget; never synchronously re-enter the collections lock. - collectionsQueue.async(flags: .barrier) { [weak self] in + // Expiry can be discovered from the BLE maintenance queue or from an + // engine slot. Cleanup is therefore fire-and-forget; never + // synchronously re-enter the engine. + messageQueue.async { [weak self] in _ = self?.pendingNoiseSessionQueues.removeTypedPayload(transferId: transferId) } TransferProgressManager.shared.rejectBeforeStart( @@ -1533,7 +1555,7 @@ final class BLEService: NSObject { // Noise cleanup remains asynchronous, but deferred private-media work // cannot pass another admission boundary after this returns. privateMediaTransferAdmissions.cancel(transferId) - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } switch self.outboundFragmentTransfers.cancelTransfer(transferId) { @@ -1808,7 +1830,7 @@ final class BLEService: NSObject { self.privateMediaTransferAdmissions.finish(transferId) return } - let queued = self.collectionsQueue.sync(flags: .barrier) { + let queued = onEngine { self.privateMediaTransferAdmissions.withActive(transferId) { self.pendingNoiseSessionQueues.appendTypedPayload( typedPayload, @@ -1824,7 +1846,7 @@ final class BLEService: NSObject { } SecureLogger.debug("📥 Queued private file for \(targetID.id.prefix(8))… pending handshake", category: .session) guard self.privateMediaTransferAdmissions.isActive(transferId) else { - self.collectionsQueue.sync(flags: .barrier) { + onEngine { _ = self.pendingNoiseSessionQueues.removeTypedPayload(transferId: transferId) } self.privateMediaTransferAdmissions.finish(transferId) @@ -1951,7 +1973,7 @@ final class BLEService: NSObject { // Queue for after handshake; initiate only while the peer is // around to answer (see sendDeliveryAck — absent senders must // not turn queued acks into handshake floods). - collectionsQueue.sync(flags: .barrier) { + onEngine { pendingNoiseSessionQueues.appendTypedPayload(payload, for: peerID) } if !noiseService.hasSession(with: peerID), isPeerReachable(peerID) { @@ -2035,7 +2057,7 @@ final class BLEService: NSObject { } private func recordIngressIfNew(_ packet: BitchatPacket, link: BLEIngressLinkID, peerID: PeerID) -> Bool { - return collectionsQueue.sync(flags: .barrier) { + return onEngine { ingressLinks.recordIfNew( packet, link: link, @@ -2350,7 +2372,7 @@ final class BLEService: NSObject { requireNoiseAuthenticatedPeerLink: Bool = false ) -> Bool { guard !isPanicSuspended else { return false } - let ingressRecord = collectionsQueue.sync { ingressLinks.record(for: packet) } + let ingressRecord = onEngine { ingressLinks.record(for: packet) } var excludedPeerLinks = links(to: ingressRecord?.peerID) if requireNoiseAuthenticatedPeerLink { guard let directedOnlyPeer else { return false } @@ -2486,7 +2508,7 @@ final class BLEService: NSObject { // MARK: - Directed store-and-forward private func spoolDirectedPacket(_ packet: BitchatPacket, recipientPeerID: PeerID) { let msgID = BLEOutboundPacketPolicy.messageID(for: packet) - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } if self.pendingDirectedRelays.enqueue( packet: packet, @@ -2502,7 +2524,7 @@ final class BLEService: NSObject { private func flushDirectedSpool() { guard !isPanicSuspended else { return } // Move items out and attempt broadcast; if still no links, they'll be re-spooled - let toSend = collectionsQueue.sync(flags: .barrier) { + let toSend = onEngine { pendingDirectedRelays.drainUnexpired( now: Date(), window: TransportConfig.bleDirectedSpoolWindowSeconds @@ -2678,7 +2700,7 @@ final class BLEService: NSObject { // queued barrier must still observe the path as pending. If // insertion wins first, the next MainActor snapshot sees the // new bubble and protects the path explicitly. - self?.messageQueue.async(flags: .barrier) { + self?.messageQueue.async { self?.incomingFileStore.finishIncomingFileDelivery( at: storedURL ) @@ -2687,7 +2709,7 @@ final class BLEService: NSObject { isPrivateMediaSenderBlocked: { [weak self] peerID in guard let self else { return false } let senderStaticKey = self.noiseService.getPeerPublicKeyData(peerID) - ?? self.collectionsQueue.sync { + ?? onEngine { self.peerRegistry.info(for: peerID)?.noisePublicKey } guard let senderStaticKey else { return false } @@ -2766,7 +2788,7 @@ final class BLEService: NSObject { // initiating a handshake broadcast turns one undeliverable ack // into a mesh-wide flood. The queued ack flushes whenever a // session eventually establishes. - collectionsQueue.sync(flags: .barrier) { + onEngine { pendingNoiseSessionQueues.appendTypedPayload(payload, for: peerID) } if !noiseService.hasSession(with: peerID), isPeerReachable(peerID) { @@ -2781,7 +2803,7 @@ final class BLEService: NSObject { /// keeps delayed/relayed leaves verifiable after the live registry entry /// has aged out. private func handleLeave(_ packet: BitchatPacket, from peerID: PeerID) -> Bool { - let registrySigningKey = collectionsQueue.sync { + let registrySigningKey = onEngine { peerRegistry.info(for: peerID)?.signingPublicKey } let verifiedViaRegistry = registrySigningKey.map { @@ -2837,7 +2859,7 @@ final class BLEService: NSObject { // related state snapshots. Serialize the whole operation with identity // rotation instead of letting CoreBluetooth and maintenance callbacks // execute it directly on their own queues. - messageQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in self?.sendAnnounceNow(forceSend: forceSend) } } @@ -3032,7 +3054,7 @@ extension BLEService: GossipSyncManager.Delegate { } func getConnectedPeers() -> [PeerID] { - return collectionsQueue.sync { + return onEngine { peerRegistry.connectedPeerIDs } } @@ -3548,7 +3570,7 @@ extension BLEService { /// avoiding wall-clock sleeps that become flaky under a parallel suite. func _test_drainFragmentPipeline() async { await withCheckedContinuation { continuation in - messageQueue.async(flags: .barrier) { + messageQueue.async { // Reassembled packets are reinjected synchronously on // `messageQueue`; their UI delivery task is therefore already // enqueued before this later MainActor marker. @@ -3638,7 +3660,7 @@ extension BLEService { messageID: String, for peerID: PeerID ) { - collectionsQueue.sync(flags: .barrier) { + onEngine { pendingNoiseSessionQueues.appendPrivateMessage( content: content, messageID: messageID, @@ -3653,7 +3675,7 @@ extension BLEService { for peerID: PeerID ) { guard privateMediaTransferAdmissions.begin(transferId) == .admitted else { return } - collectionsQueue.sync(flags: .barrier) { + onEngine { pendingNoiseSessionQueues.appendTypedPayload( payload, transferId: transferId, @@ -3667,14 +3689,14 @@ extension BLEService { } func _test_hasPendingPrivateMediaPolicyResolution(for peerID: PeerID) -> Bool { - collectionsQueue.sync { + onEngine { pendingPrivateMediaPolicyResolutions[peerID.toShort()] != nil } } func _test_forcePrivateMediaProofTimeout(for peerID: PeerID) { let normalizedPeerID = peerID.toShort() - let target = collectionsQueue.sync { + let target = onEngine { () -> (fingerprint: String, generation: UUID?, nonce: UUID)? in if let watchdog = privateMediaProofWatchdogs[normalizedPeerID] { return ( @@ -3704,7 +3726,7 @@ extension BLEService { func _test_privateMediaTransferState( transferId: String ) -> (admissionActive: Bool, pendingNoise: Bool, activeScheduler: Int, pendingScheduler: Int) { - let scheduler = collectionsQueue.sync { + let scheduler = onEngine { ( pendingNoiseSessionQueues.containsTypedPayload(transferId: transferId), outboundFragmentTransfers.activeCount, @@ -3737,10 +3759,9 @@ extension BLEService { } func _test_drainPrivateMediaSendPipeline() async { - let collectionsQueue = self.collectionsQueue await withCheckedContinuation { continuation in - messageQueue.async { - collectionsQueue.async(flags: .barrier) { + self.messageQueue.async { [weak self] in + self?.messageQueue.async { continuation.resume() } } @@ -3759,10 +3780,9 @@ extension BLEService { } func _test_drainNoiseMessagePipeline() async { - let collectionsQueue = self.collectionsQueue await withCheckedContinuation { continuation in - messageQueue.async(flags: .barrier) { - collectionsQueue.async(flags: .barrier) { + self.messageQueue.async { + self.messageQueue.async { continuation.resume() } } @@ -3773,7 +3793,7 @@ extension BLEService { /// this to prove same-generation reconciliation is idempotent. func _test_reconcileCurrentNoiseSession(for peerID: PeerID) { let normalizedPeerID = peerID.toShort() - messageQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self, let generation = self.noiseService.sessionGeneration( for: normalizedPeerID @@ -4439,7 +4459,7 @@ extension BLEService: PrivateMediaDeletionPersisting { completion: @escaping @MainActor (Bool) -> Void ) { let fileStore = incomingFileStore - messageQueue.async(flags: .barrier) { + messageQueue.async { guard let reservation = fileStore .reservePrivateMediaDeletion( messageIDs: messageIDs, @@ -4467,7 +4487,7 @@ extension BLEService: PrivateMediaDeletionPersisting { @MainActor func removeLegacyPrivateMediaPayload(relativePath: String) { let fileStore = incomingFileStore - messageQueue.async(flags: .barrier) { + messageQueue.async { fileStore.removeLegacyIncomingFile(relativePath: relativePath) } } @@ -4722,7 +4742,7 @@ extension BLEService { localPeerIDData: myPeerIDData, isRecipientConnected: { self.isPeerConnected($0) }, shouldAttemptRoute: { peer in - self.collectionsQueue.sync(flags: .barrier) { + onEngine { self.sourceRouteFailures.shouldAttemptRoute(to: peer, now: now) } }, @@ -4746,7 +4766,7 @@ extension BLEService { SecureLogger.error("❌ Failed to re-sign packet with route", category: .security) return packet // Return original packet if signing fails } - collectionsQueue.sync(flags: .barrier) { + onEngine { sourceRouteFailures.noteRoutedSend(to: recipient, now: now) } return signedPacket @@ -4794,7 +4814,7 @@ extension BLEService { ) let timeout = DispatchWorkItem { [weak self] in guard let self else { return } - let expired = self.collectionsQueue.sync(flags: .barrier) { + let expired = onEngine { self.pendingMeshPings.removeValue(forKey: nonce) } guard let expired else { return } @@ -4808,7 +4828,7 @@ extension BLEService { expired.completion(nil) } } - self.collectionsQueue.sync(flags: .barrier) { + onEngine { self.pendingMeshPings[nonce] = PendingMeshPing( peerID: PeerID(hexData: recipientData), sentAt: Date(), @@ -4840,7 +4860,7 @@ extension BLEService { SecureLogger.debug("⚠️ Malformed ping via \(linkPeerID.id.prefix(8))…", category: .session) return } - let allowed = collectionsQueue.sync(flags: .barrier) { + let allowed = onEngine { meshPingResponseLimiter.shouldRespond(to: linkPeerID, now: Date()) } guard allowed else { @@ -4868,7 +4888,7 @@ extension BLEService { private func handleMeshPong(_ packet: BitchatPacket, from peerID: PeerID) { guard packet.recipientID == myPeerIDData else { return } guard let pong = MeshPingPayload.decode(packet.payload) else { return } - let pending = collectionsQueue.sync(flags: .barrier) { () -> PendingMeshPing? in + let pending = onEngine { () -> PendingMeshPing? in guard pendingMeshPings[pong.nonce]?.peerID == peerID else { return nil } return pendingMeshPings.removeValue(forKey: pong.nonce) } @@ -4967,7 +4987,7 @@ extension BLEService { /// handshake. An old session keyed only by peer ID is insufficient: a /// replayed announce can rebind an attacker's link to that ID. private func markNoiseAuthenticatedIngressLink(for packet: BitchatPacket, peerID: PeerID) { - guard let link = collectionsQueue.sync(execute: { ingressLinks.link(for: packet) }) else { return } + guard let link = onEngine({ ingressLinks.link(for: packet) }) else { return } readLinkState { store in guard boundPeerID(for: link, in: store) == peerID else { return } noiseAuthenticatedLinkOwners[link] = peerID @@ -4975,7 +4995,7 @@ extension BLEService { } private func isNoiseAuthenticatedIngressLink(for packet: BitchatPacket, peerID: PeerID) -> Bool { - guard let link = collectionsQueue.sync(execute: { ingressLinks.link(for: packet) }) else { return false } + guard let link = onEngine({ ingressLinks.link(for: packet) }) else { return false } return readLinkState { store in noiseAuthenticatedLinkOwners[link] == peerID && boundPeerID(for: link, in: store) == peerID } @@ -4996,14 +5016,14 @@ extension BLEService { /// A peer-level session can outlive the physical link that established it. /// Revalidate a fresh direct link with an ordinary XX exchange, retiring /// cached sending keys atomically before message 1 can leave. + /// + /// Takes the already-resolved ingress link: both callers run inside the + /// rebind's bleQueue critical section, which must never sync-wait on the + /// engine (the engine sync-waits on bleQueue via `readLinkState`). private func refreshNoiseSessionForVerifiedDirectLink( - _ packet: BitchatPacket, + link: BLEIngressLinkID, peerID: PeerID ) { - guard let link = collectionsQueue.sync(execute: { ingressLinks.link(for: packet) }) else { - return - } - let hasEstablishedSession = noiseService.hasEstablishedSession(with: peerID) let authenticatedPeerLinks = currentNoiseAuthenticatedLinks(to: peerID) let shouldRevalidate = readLinkState { store in @@ -5033,7 +5053,7 @@ extension BLEService { // Authentication can be reported while an initiator is still // returning XX message 3. Serialize generation-bound state and // every post-handshake drain behind the handshake packet handler. - self?.messageQueue.async(flags: .barrier) { [weak self] in + self?.messageQueue.async { [weak self] in self?.handleNoisePeerAuthenticated( peerID: peerID, fingerprint: fingerprint, @@ -5043,7 +5063,7 @@ extension BLEService { } service.onRekeyHandshakeReady = { [weak self, weak service] peerID, initiation in - self?.messageQueue.async(flags: .barrier) { + self?.messageQueue.async { [weak self, weak service] in guard let self, let service, @@ -5066,7 +5086,7 @@ extension BLEService { #if DEBUG self._test_beforeHandshakeRecoveryEnqueued?(request.peerID) #endif - self.messageQueue.async(flags: .barrier) { + self.messageQueue.async { [weak self, weak service] in guard let self, let service, @@ -5113,7 +5133,7 @@ extension BLEService { guard let self, let service else { return } // The manager makes restored keys visible atomically. Reconcile // transport state and queued sends as the next serialized phase. - self.messageQueue.async(flags: .barrier) { [weak self, weak service] in + self.messageQueue.async { [weak self, weak service] in guard let self, let service, self.noiseService === service, @@ -5149,43 +5169,51 @@ extension BLEService { sessionGeneration generation: UUID, deferOutboundUntilConvergence: Bool = false ) { + // Every caller enqueues this on the engine queue, and that matters + // below: the closure handed to `withCurrentSessionGeneration` runs + // on the noise manager's queue while this engine slot stays blocked, + // so it may touch engine-owned state directly — the engine cannot + // run anything else concurrently. It must NOT go through `onEngine` + // (a sync back onto the engine from inside a queue the engine is + // sync-waiting on is a self-deadlock). + #if DEBUG + dispatchPrecondition(condition: .onQueue(messageQueue)) + #endif let normalizedPeerID = peerID.toShort() guard let transition = noiseService.withCurrentSessionGeneration( for: normalizedPeerID, expected: generation, { - collectionsQueue.sync(flags: .barrier) { - () -> ( - watchdog: (fingerprint: String, nonce: UUID)?, - rejected: [@MainActor (PrivateMediaSendPolicy) -> Void] - ) in - guard privateMediaSessionGenerations[normalizedPeerID] != generation else { - return (nil, []) - } - let watchdogNonce = UUID() - privateMediaSessionGenerations[normalizedPeerID] = generation - authenticatedPeerStates.removeValue(forKey: normalizedPeerID) - privateMediaProofTimeoutMarkers.removeValue(forKey: normalizedPeerID) - privateMediaProofWatchdogs[normalizedPeerID] = BLEPrivateMediaProofWatchdog( - fingerprint: fingerprint, - sessionGeneration: generation, - timeoutNonce: watchdogNonce - ) - authenticatedPeerStateSendProgress[normalizedPeerID] = - BLEAuthenticatedPeerStateSendProgress(sessionGeneration: generation) - - guard var pending = pendingPrivateMediaPolicyResolutions[normalizedPeerID] else { - return ((fingerprint, watchdogNonce), []) - } - guard pending.fingerprint.caseInsensitiveCompare(fingerprint) == .orderedSame else { - pendingPrivateMediaPolicyResolutions.removeValue(forKey: normalizedPeerID) - return ((fingerprint, watchdogNonce), Array(pending.completions.values)) - } - pending.sessionGeneration = generation - pending.timeoutNonce = watchdogNonce - pendingPrivateMediaPolicyResolutions[normalizedPeerID] = pending - return ((pending.fingerprint, watchdogNonce), []) + () -> ( + watchdog: (fingerprint: String, nonce: UUID)?, + rejected: [@MainActor (PrivateMediaSendPolicy) -> Void] + ) in + guard privateMediaSessionGenerations[normalizedPeerID] != generation else { + return (nil, []) } + let watchdogNonce = UUID() + privateMediaSessionGenerations[normalizedPeerID] = generation + authenticatedPeerStates.removeValue(forKey: normalizedPeerID) + privateMediaProofTimeoutMarkers.removeValue(forKey: normalizedPeerID) + privateMediaProofWatchdogs[normalizedPeerID] = BLEPrivateMediaProofWatchdog( + fingerprint: fingerprint, + sessionGeneration: generation, + timeoutNonce: watchdogNonce + ) + authenticatedPeerStateSendProgress[normalizedPeerID] = + BLEAuthenticatedPeerStateSendProgress(sessionGeneration: generation) + + guard var pending = pendingPrivateMediaPolicyResolutions[normalizedPeerID] else { + return ((fingerprint, watchdogNonce), []) + } + guard pending.fingerprint.caseInsensitiveCompare(fingerprint) == .orderedSame else { + pendingPrivateMediaPolicyResolutions.removeValue(forKey: normalizedPeerID) + return ((fingerprint, watchdogNonce), Array(pending.completions.values)) + } + pending.sessionGeneration = generation + pending.timeoutNonce = watchdogNonce + pendingPrivateMediaPolicyResolutions[normalizedPeerID] = pending + return ((pending.fingerprint, watchdogNonce), []) } ) else { return } @@ -5257,7 +5285,7 @@ extension BLEService { private func sendAuthenticatedPeerState(to peerID: PeerID, echo: Bool) { let normalizedPeerID = peerID.toShort() - let shouldSend = collectionsQueue.sync(flags: .barrier) { + let shouldSend = onEngine { guard let generation = privateMediaSessionGenerations[normalizedPeerID], var progress = authenticatedPeerStateSendProgress[normalizedPeerID], progress.sessionGeneration == generation else { return false } @@ -5290,6 +5318,13 @@ extension BLEService { from peerID: PeerID, sessionGeneration generation: UUID ) { + // Engine-only, like handleNoisePeerAuthenticated: the closure below + // runs on the noise manager's queue while this engine slot stays + // blocked, so it accesses engine-owned state directly instead of + // sync-re-entering the engine (self-deadlock). + #if DEBUG + dispatchPrecondition(condition: .onQueue(messageQueue)) + #endif let normalizedPeerID = peerID.toShort() guard let state = AuthenticatedPeerStatePacket.decode(from: payload) else { SecureLogger.warning( @@ -5312,9 +5347,7 @@ extension BLEService { expected: generation, { () -> (accepted: Bool, completions: [@MainActor (PrivateMediaSendPolicy) -> Void]) in - guard collectionsQueue.sync(execute: { - privateMediaSessionGenerations[normalizedPeerID] == generation - }) else { + guard privateMediaSessionGenerations[normalizedPeerID] == generation else { return (false, []) } @@ -5334,33 +5367,26 @@ extension BLEService { identityManager.markPrivateMediaCapable(fingerprint: fingerprint) } - let completions = collectionsQueue.sync(flags: .barrier) { - () -> [@MainActor (PrivateMediaSendPolicy) -> Void] in - guard privateMediaSessionGenerations[normalizedPeerID] == generation else { - return [] - } - peerRegistry.mutate { - $0.bindAuthenticatedSigningPublicKey( - state.signingPublicKey, - for: normalizedPeerID - ) - } - authenticatedPeerStates[normalizedPeerID] = BLEAuthenticatedPeerStateObservation( - fingerprint: fingerprint, - sessionGeneration: generation, - capabilities: state.capabilities + peerRegistry.mutate { + $0.bindAuthenticatedSigningPublicKey( + state.signingPublicKey, + for: normalizedPeerID ) - privateMediaProofTimeoutMarkers.removeValue(forKey: normalizedPeerID) - privateMediaProofWatchdogs.removeValue(forKey: normalizedPeerID) - guard let pending = pendingPrivateMediaPolicyResolutions.removeValue( - forKey: normalizedPeerID - ), pending.fingerprint.caseInsensitiveCompare(fingerprint) == .orderedSame, - pending.sessionGeneration == generation else { - return [] - } - return Array(pending.completions.values) } - return (true, completions) + authenticatedPeerStates[normalizedPeerID] = BLEAuthenticatedPeerStateObservation( + fingerprint: fingerprint, + sessionGeneration: generation, + capabilities: state.capabilities + ) + privateMediaProofTimeoutMarkers.removeValue(forKey: normalizedPeerID) + privateMediaProofWatchdogs.removeValue(forKey: normalizedPeerID) + guard let pending = pendingPrivateMediaPolicyResolutions.removeValue( + forKey: normalizedPeerID + ), pending.fingerprint.caseInsensitiveCompare(fingerprint) == .orderedSame, + pending.sessionGeneration == generation else { + return (true, []) + } + return (true, Array(pending.completions.values)) } ), application.accepted else { return } @@ -5374,7 +5400,7 @@ extension BLEService { private func noteNoiseSessionCleared(for peerID: PeerID) { let normalizedPeerID = peerID.toShort() - let reset = collectionsQueue.sync(flags: .barrier) { + let reset = onEngine { () -> (fingerprint: String, nonce: UUID)? in privateMediaSessionGenerations.removeValue(forKey: normalizedPeerID) authenticatedPeerStates.removeValue(forKey: normalizedPeerID) @@ -5423,7 +5449,7 @@ extension BLEService { if DispatchQueue.getSpecific(key: messageQueueKey) != nil { swap() } else { - messageQueue.sync(flags: .barrier, execute: swap) + messageQueue.sync(execute: swap) } } @@ -5443,7 +5469,7 @@ extension BLEService { // No established session yet - queue the payload synchronously // before initiating a handshake // to prevent race where fast handshake completion drains empty queue - collectionsQueue.sync(flags: .barrier) { + onEngine { self.pendingNoiseSessionQueues.appendTypedPayload(typedPayload, for: peerID) SecureLogger.debug("📥 Queued noise payload for \(peerID.id.prefix(8))… pending handshake", category: .session) } @@ -5465,7 +5491,7 @@ extension BLEService { let encrypted: Data let isPrivateFile = NoisePayloadType.isPrivateFile(rawValue: typedPayload.first) if isPrivateFile { - let provenGeneration: UUID? = collectionsQueue.sync { + let provenGeneration: UUID? = onEngine { () -> UUID? in guard let generation = privateMediaSessionGenerations[peerID], let authenticated = authenticatedPeerStates[peerID], @@ -5645,7 +5671,7 @@ extension BLEService { /// gateway watches courier drops for. func verifiedPeersWithNoiseKeys() -> [(peerID: PeerID, noiseKey: Data)] { let now = Date() - return collectionsQueue.sync { + return onEngine { peerRegistry.snapshotByID.values.compactMap { info in guard info.isVerifiedNickname, let key = info.noisePublicKey, @@ -5742,7 +5768,7 @@ extension BLEService { // so dedup here on the inner message ID — before delivery, ack, // and handshake work. A duplicate costs only the decrypt above // and at most one ack ever goes out per message ID. - let firstOpen = collectionsQueue.sync(flags: .barrier) { + let firstOpen = onEngine { openedCourierMessageIDs.insert(innerMessageID) } guard firstOpen else { @@ -5902,7 +5928,7 @@ extension BLEService { /// Forced sends (bundle changed after consumption) go immediately. private func sendPrekeyBundle(force: Bool = false) { let now = Date() - let shouldSend: Bool = collectionsQueue.sync(flags: .barrier) { + let shouldSend: Bool = onEngine { if !force, let last = lastPrekeyBundleSentAt, now.timeIntervalSince(last) < TransportConfig.prekeyBundleRebroadcastSeconds { @@ -5970,7 +5996,7 @@ extension BLEService { // ahead of the announce that binds the key. Reading the live registry // and stashing atomically closes the check-then-act gap against // handleAnnounce's drain (see drainPendingPrekeyBundles). - let signingKey: Data? = collectionsQueue.sync(flags: .barrier) { + let signingKey: Data? = onEngine { if let info = peerRegistry.info(for: owner), info.noisePublicKey == bundle.noiseStaticPublicKey, let key = info.signingPublicKey { @@ -6015,7 +6041,7 @@ extension BLEService { /// announce, in a barrier ordered after the registry write, so a bundle /// stashed before the write is always observed here. private func drainPendingPrekeyBundles(for owner: PeerID) { - let pending: BitchatPacket? = collectionsQueue.sync(flags: .barrier) { + let pending: BitchatPacket? = onEngine { pendingPrekeyBundles.removeValue(forKey: owner) } guard let packet = pending, @@ -6422,7 +6448,7 @@ extension BLEService { SecureLogger.debug("🤝 No session with \(recipientID.id.prefix(8))…, initiating handshake and queueing message", category: .session) // Queue the message (especially important for favorite notifications) - collectionsQueue.sync(flags: .barrier) { + onEngine { pendingNoiseSessionQueues.appendPrivateMessage(content: content, messageID: messageID, for: recipientID) } @@ -6444,7 +6470,7 @@ extension BLEService { ) else { return } - messageQueue.async(flags: .barrier) { + messageQueue.async { [weak self, weak service] in guard let self, let service, @@ -6486,7 +6512,7 @@ extension BLEService { with: peerID, retryOnTimeout: true ) - messageQueue.async(flags: .barrier) { [weak self, weak service] in + messageQueue.async { [weak self, weak service] in guard let self, let service, self.noiseService === service else { @@ -6513,7 +6539,7 @@ extension BLEService { private func sendPendingMessagesAfterHandshake(for peerID: PeerID) { // Atomically take all pending messages to process (prevents concurrent modification) - let pendingMessages = collectionsQueue.sync(flags: .barrier) { () -> [BLEPendingPrivateMessage] in + let pendingMessages = onEngine { () -> [BLEPendingPrivateMessage] in pendingNoiseSessionQueues.takePrivateMessages(for: peerID) } @@ -6556,7 +6582,7 @@ extension BLEService { // Re-queue any failed messages for retry on next handshake if !failedMessages.isEmpty { - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } // Prepend failed messages to maintain order self.pendingNoiseSessionQueues.prependPrivateMessages(failedMessages, for: peerID) @@ -6588,12 +6614,12 @@ extension BLEService { requireNoiseAuthenticatedPeerLink: requireNoiseAuthenticatedPeerLink ) - let result: BLEOutboundFragmentTransferScheduler.SubmitResult? = collectionsQueue.sync(flags: .barrier) { + let result: BLEOutboundFragmentTransferScheduler.SubmitResult? = onEngine { if requiresPrivateMediaAdmission { guard let transferId else { return nil } - // This lock is taken while the scheduler is already protected - // by collectionsQueue. Cancellation takes the admission lock - // synchronously but never waits on collectionsQueue, avoiding + // This lock is taken while the scheduler is already + // engine-confined. Cancellation takes the admission lock + // synchronously but never waits on the engine, avoiding // lock inversion while giving submit/cancel one linear order. return privateMediaTransferAdmissions.withActive(transferId) { outboundFragmentTransfers.submit( @@ -6658,7 +6684,7 @@ extension BLEService { let releaseReservedSlot: (String) -> Void = { [weak self] id in guard let self = self else { return } TransferProgressManager.shared.cancel(id: id) - self.collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in _ = self?.outboundFragmentTransfers.releaseReservation(id) } self.messageQueue.async { [weak self] in @@ -6693,7 +6719,7 @@ extension BLEService { let transferIdentifier: String? if let id = reservedTransferId { - let activated = collectionsQueue.sync(flags: .barrier) { + let activated = onEngine { self.outboundFragmentTransfers.activateReservedTransfer( id: id, totalFragments: plan.totalFragments, @@ -6752,7 +6778,7 @@ extension BLEService { let workItem = DispatchWorkItem { [weak self] in guard let self = self else { return } if let transferId = transferIdentifier { - let isActive = self.collectionsQueue.sync { self.outboundFragmentTransfers.isActive(transferId) } + let isActive = onEngine { self.outboundFragmentTransfers.isActive(transferId) } guard isActive else { return } } if fragmentPacket.recipientID == nil || fragmentPacket.recipientID?.allSatisfy({ $0 == 0xFF }) == true { @@ -6769,7 +6795,7 @@ extension BLEService { if let transferId = transferIdentifier { let workItems = scheduledItems.map { $0.item } - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in _ = self?.outboundFragmentTransfers.updateWorkItems(workItems, for: transferId) } } @@ -6784,7 +6810,7 @@ extension BLEService { // MARK: - Fragmentation (Required for messages > BLE MTU) private func markFragmentSent(transferId: String) { - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } switch self.outboundFragmentTransfers.markFragmentSent(transferId: transferId) { @@ -6804,7 +6830,7 @@ extension BLEService { } private func startNextPendingTransferIfNeeded() { - let results = collectionsQueue.sync(flags: .barrier) { + let results = onEngine { outboundFragmentTransfers.reservePendingStarts(maxConcurrentTransfers: TransportConfig.bleMaxConcurrentTransfers) } @@ -6819,7 +6845,7 @@ extension BLEService { if DispatchQueue.getSpecific(key: messageQueueKey) != nil { fragmentHandler.handle(packet, from: peerID) } else { - messageQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in self?.fragmentHandler.handle(packet, from: peerID) } } @@ -6839,7 +6865,7 @@ extension BLEService { guard let self = self else { return .stored(header: header, started: false) } - return self.collectionsQueue.sync(flags: .barrier) { + return onEngine { self.fragmentAssemblyBuffer.append(header, maxInFlightAssemblies: self.maxInFlightAssemblies) } }, @@ -6890,7 +6916,7 @@ extension BLEService { capturePanicLifecycleGeneration() else { return } - messageQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self, self.isCurrentPanicLifecycleGeneration( lifecycleGeneration @@ -6925,7 +6951,7 @@ extension BLEService { // Track recent traffic timestamps for adaptive behavior; the same // barrier hop confirms route health for the packet's originator. - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } self.recentTrafficTracker.recordPacket(at: Date()) self.sourceRouteFailures.noteInboundActivity(from: senderID) @@ -7032,7 +7058,7 @@ extension BLEService { let connectedCount = peerRegistry.connectedCount if BLEReceivePipeline.shouldCancelScheduledRelayForDuplicate(connectedPeerCount: connectedCount) { - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in self?.scheduledRelays.cancel(messageID: messageID) } } @@ -7053,7 +7079,7 @@ extension BLEService { let work = DispatchWorkItem { [weak self] in guard let self = self else { return } - self.collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in self?.scheduledRelays.remove(messageID: messageID) } var relayPacket = packet @@ -7061,7 +7087,7 @@ extension BLEService { self.broadcastPacket(relayPacket) } - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in self?.scheduledRelays.schedule(work, messageID: messageID) } messageQueue.asyncAfter(deadline: .now() + .milliseconds(decision.delayMs), execute: work) @@ -7148,7 +7174,7 @@ extension BLEService { /// sender owns the link it arrived on, so rebind the link to the new ID /// and retire the old identity. private func rebindLinkAfterVerifiedDirectAnnounce(_ packet: BitchatPacket, to peerID: PeerID) { - guard let link = (collectionsQueue.sync { ingressLinks.link(for: packet) }) else { return } + guard let link = (onEngine { ingressLinks.link(for: packet) }) else { return } bleQueue.async { [weak self] in guard let self else { return } let linkUUID: String @@ -7164,7 +7190,7 @@ extension BLEService { guard let previousPeerID else { return } guard previousPeerID != peerID else { self.refreshNoiseSessionForVerifiedDirectLink( - packet, + link: link, peerID: peerID ) return @@ -7206,7 +7232,7 @@ extension BLEService { // section. No observer may see the new binding while a cached // peer-level sender is still considered established. self.refreshNoiseSessionForVerifiedDirectLink( - packet, + link: link, peerID: peerID ) SecureLogger.debug("🔄 Rebinding link after peer-ID rotation: \(previousPeerID.id.prefix(8))… → \(peerID.id.prefix(8))…", category: .session) @@ -7254,7 +7280,7 @@ extension BLEService { /// retirement per peer per cooldown window, and the peer keeps a live /// link either way. private func retireRedundantPeripheralLinks(_ packet: BitchatPacket, to peerID: PeerID) { - let ingressLink = collectionsQueue.sync { ingressLinks.link(for: packet) } + let ingressLink = onEngine { ingressLinks.link(for: packet) } bleQueue.async { [weak self] in guard let self else { return } let now = Date() @@ -7365,7 +7391,7 @@ extension BLEService { now: { Date() }, existingPeerKeys: { [weak self] peerID in guard let self = self else { return (nil, nil) } - return self.collectionsQueue.sync { + return onEngine { let info = self.peerRegistry.info(for: peerID) return (info?.noisePublicKey, info?.signingPublicKey) } @@ -7397,7 +7423,7 @@ extension BLEService { // connected. See the caller in BLEAnnounceHandler for why the // residual forged-presence window this leaves is accepted. guard let self else { return false } - guard let link = (self.collectionsQueue.sync { self.ingressLinks.link(for: packet) }) else { return false } + guard let link = (onEngine { self.ingressLinks.link(for: packet) }) else { return false } let boundPeerID: PeerID? = self.readLinkState { store in switch link { case .peripheral(let peripheralUUID): @@ -7410,7 +7436,7 @@ extension BLEService { return boundPeerID != peerID }, withRegistryBarrier: { [weak self] body in - self?.collectionsQueue.sync(flags: .barrier) { body() } + self?.onEngine { body() } }, upsertVerifiedAnnounce: { [weak self] peerID, announcement, isConnected, now in // Called from inside withRegistryBarrier; access registry directly. @@ -7762,15 +7788,20 @@ extension BLEService { }, decrypt: { [weak self] payload, peerID in guard let self = self else { throw NoiseEncryptionError.sessionNotEstablished } + // Decrypt runs on the engine queue; the readiness callback + // fires on the noise manager's queue while this engine slot + // stays blocked, so it reads engine-owned state directly + // (sync-re-entering the engine there would self-deadlock). + #if DEBUG + dispatchPrecondition(condition: .onQueue(self.messageQueue)) + #endif let result = try self.noiseService.decryptWithSessionGeneration( payload, from: peerID, establishedGenerationIsReady: { generation in - self.collectionsQueue.sync { - self.privateMediaSessionGenerations[ - peerID.toShort() - ] == generation - } + self.privateMediaSessionGenerations[ + peerID.toShort() + ] == generation } ) return BLENoiseDecryptionResult( @@ -7813,7 +7844,7 @@ extension BLEService { // MARK: Helper Functions private func sendPendingNoisePayloadsAfterHandshake(for peerID: PeerID) { - let payloads = collectionsQueue.sync(flags: .barrier) { () -> [BLEPendingTypedPayload] in + let payloads = onEngine { () -> [BLEPendingTypedPayload] in pendingNoiseSessionQueues.takeTypedPayloads(for: peerID) } guard !payloads.isEmpty else { return } @@ -7831,7 +7862,7 @@ extension BLEService { // Handshake completion alone is insufficient. Put the // exact payload back until authenticated 0x21 state // arrives; that handler calls this drain again. - collectionsQueue.sync(flags: .barrier) { + onEngine { pendingNoiseSessionQueues.appendTypedPayload( pending.payload, transferId: pending.transferId, @@ -8055,7 +8086,7 @@ extension BLEService { // Clean old fragments (> configured seconds old), then ask peers for // the specific fragment streams whose reassembly has stalled instead // of waiting for the next periodic GCS fragment round. - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self else { return } let cutoff = now.addingTimeInterval(-TransportConfig.bleFragmentLifetimeSeconds) self.fragmentAssemblyBuffer.removeExpired(before: cutoff) @@ -8076,14 +8107,14 @@ extension BLEService { connectionScheduler.pruneConnectionTimeouts(before: timeoutCutoff) // Clean up stale scheduled relays that somehow persisted (> 2s) - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } // Nothing to compare times to; just cap the size defensively self.scheduledRelays.removeAllIfOverCapacity(512) } // Clean ingress link records older than configured seconds - collectionsQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } let cutoff = now.addingTimeInterval(-TransportConfig.bleIngressRecordLifetimeSeconds) if !self.ingressLinks.isEmpty { @@ -8096,7 +8127,7 @@ extension BLEService { ) } - messageQueue.async(flags: .barrier) { [weak self] in + messageQueue.async { [weak self] in guard let self = self else { return } guard !self.selfBroadcastTracker.isEmpty else { return } let cutoff = now.addingTimeInterval(-TransportConfig.messageDedupMaxAgeSeconds)