From 6cd7a43b4bf15f98947540a1478b269a177f52c3 Mon Sep 17 00:00:00 2001 From: jack Date: Sat, 25 Jul 2026 20:19:41 +0200 Subject: [PATCH] Invalidate queued BLE ingress during panic --- bitchat/Services/BLE/BLEService.swift | 154 +++++++++++++++++++++---- bitchatTests/BLEServiceCoreTests.swift | 154 +++++++++++++++++++++++++ 2 files changed, 286 insertions(+), 22 deletions(-) diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index d8d7a6c6..f38e090a 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -104,6 +104,10 @@ final class BLEService: NSObject { // Test-only tap on the outbound pipeline so multi-node tests can ferry // packets between in-process service instances. var _test_onOutboundPacket: ((BitchatPacket) -> Void)? + /// May block a synthetic CoreBluetooth receive callback immediately + /// before it hands a packet to `messageQueue`. + var _test_beforeReceivePacketHandoff: (() -> Void)? + var _test_onReceivePacketHandoff: (() -> Void)? #endif private var selfBroadcastTracker = BLESelfBroadcastTracker() private let meshTopology = MeshTopologyTracker() @@ -119,6 +123,7 @@ final class BLEService: NSObject { private struct PendingMeshPing { let peerID: PeerID let sentAt: Date + let lifecycleGeneration: UInt64 let completion: @MainActor (MeshPingResult?) -> Void let timeout: DispatchWorkItem } @@ -483,11 +488,15 @@ final class BLEService: NSObject { setPanicSuspended(true) gossipSyncManager?.stop() gossipSyncManager = nil - // Wait out sends that passed admission before the gate closed. Work - // queued behind this barrier observes `isPanicSuspended` and exits, - // so the radio stop below is a true synchronous boundary. - messageQueue.sync(flags: .barrier) {} + // Stop the radio and drain CoreBluetooth's delegate queue first. A + // callback may already have passed its initial suspension check; the + // bleQueue drain forces its final messageQueue handoff to happen + // before the receive barrier below. stopServicesImmediatelyForPanic() + // Drain every receive/send submitted by callbacks that finished ahead + // of the radio stop. Later callbacks observe the closed lifecycle, and + // generation-bound handoffs that raced this barrier reject themselves. + messageQueue.sync(flags: .barrier) {} clearEmergencySessionState() } @@ -795,21 +804,30 @@ final class BLEService: NSObject { private func clearEmergencySessionState() { // Clear all sessions and peers - let cancelledTransfers: [(id: String, items: [DispatchWorkItem])] = collectionsQueue.sync(flags: .barrier) { - let entries = outboundFragmentTransfers.removeAll().map { ($0.id, $0.workItems) } + let cancelled = collectionsQueue.sync(flags: .barrier) { + let entries = outboundFragmentTransfers.removeAll().map { + (id: $0.id, items: $0.workItems) + } + let pingTimeouts = pendingMeshPings.values.map(\.timeout) + pendingMeshPings.removeAll() + meshPingResponseLimiter = SyncResponseRateLimiter( + maxResponses: TransportConfig.meshPingInboundMaxPerLink, + window: TransportConfig.meshPingInboundWindowSeconds + ) peerRegistry.removeAll() fragmentAssemblyBuffer.removeAll() sourceRouteFailures = BLESourceRouteFailureCache() // Also clear pending message queues to avoid stale state across sessions pendingNoiseSessionQueues.removeAll() pendingDirectedRelays.removeAll() - return entries + return (transfers: entries, pingTimeouts: pingTimeouts) } - for entry in cancelledTransfers { + for entry in cancelled.transfers { entry.items.forEach { $0.cancel() } TransferProgressManager.shared.cancel(id: entry.id) } + cancelled.pingTimeouts.forEach { $0.cancel() } // Clear processed messages messageDeduplicator.reset() @@ -1602,22 +1620,40 @@ final class BLEService: NSObject { } func collectArchivedPublicMessages(completion: @escaping @MainActor ([ArchivedPublicMessage]) -> Void) { + guard let generation = capturePanicLifecycleGeneration() else { + return + } guard let sync = gossipSyncManager else { - Task { @MainActor in completion([]) } + notifyUI { [weak self] in + guard let self, + self.isCurrentPanicLifecycleGeneration(generation) else { + return + } + completion([]) + } return } sync.collectPublicMessagePackets { [weak self] packets in - guard let self = self else { - Task { @MainActor in completion([]) } + guard let self, + self.isCurrentPanicLifecycleGeneration(generation) else { return } // Signature verification and registry lookups run on messageQueue // like the live receive path. self.messageQueue.async { + guard self.isCurrentPanicLifecycleGeneration(generation) else { + return + } let decoded = packets .compactMap { self.decodeArchivedPublicMessage($0) } .sorted { $0.timestamp < $1.timestamp } - Task { @MainActor in completion(decoded) } + self.notifyUI { [weak self] in + guard let self, + self.isCurrentPanicLifecycleGeneration(generation) else { + return + } + completion(decoded) + } } } } @@ -2414,6 +2450,28 @@ private extension BLEService { #if DEBUG // Test-only helper to inject packets into the receive pipeline extension BLEService { + /// Queues an event through the same MainActor hop as production receive + /// handlers so panic-boundary tests can deterministically invalidate it. + func _test_emitTransportEvent(_ event: TransportEvent) { + emitTransportEvent(event) + } + + var _test_isPanicIngressOpen: Bool { + capturePanicLifecycleGeneration() != nil + } + + /// Models a CoreBluetooth delegate callback without requiring a physical + /// peripheral. The callback itself runs on `bleQueue`, exactly where the + /// panic radio-stop barrier must linearize it. + func _test_handlePacketFromBLEQueue( + _ packet: BitchatPacket, + fromPeerID: PeerID + ) { + bleQueue.async { [weak self] in + self?.handleReceivedPacket(packet, from: fromPeerID) + } + } + func _test_handlePacket(_ packet: BitchatPacket, fromPeerID: PeerID, preseedPeer: Bool = true, signingPublicKey: Data? = nil) { if preseedPeer { // Ensure the synthetic peer is known and marked verified for public-message tests @@ -3149,8 +3207,18 @@ extension BLEService { /// Notify UI on the MainActor to satisfy Swift concurrency isolation private func notifyUI(_ block: @escaping @MainActor () -> Void) { - // Always hop onto the MainActor so calls to @MainActor delegates are safe - Task { @MainActor in + // Capture the panic lifecycle before queueing the MainActor hop. A + // receive callback can enqueue UI delivery immediately before panic + // clears application state; rechecking here prevents that stale work + // from repopulating the wiped conversation store afterward. + guard let generation = capturePanicLifecycleGeneration() else { + return + } + Task { @MainActor [weak self] in + guard let self, + self.isCurrentPanicLifecycleGeneration(generation) else { + return + } block() } } @@ -3315,14 +3383,24 @@ extension BLEService { /// The completion fires exactly once on the main actor: with RTT/hops /// when the matching pong returns, or nil after the timeout window. func sendMeshPing(to peerID: PeerID, completion: @escaping @MainActor (MeshPingResult?) -> Void) { + guard let generation = capturePanicLifecycleGeneration() else { + return + } messageQueue.async { [weak self] in guard let self, + self.isCurrentPanicLifecycleGeneration(generation), let recipientData = peerID.toShort().routingData, let payload = MeshPingPayload( nonce: Data((0.. Bool { + let deadline = DispatchTime.now().uptimeNanoseconds + + UInt64(timeout * 1_000_000_000) + while service._test_isPanicIngressOpen, + DispatchTime.now().uptimeNanoseconds < deadline { + Thread.sleep(forTimeInterval: 0.001) + } + return !service._test_isPanicIngressOpen + } +} + private func makeService() -> BLEService { let keychain = MockKeychain() let identityManager = MockIdentityManager(keychain) @@ -744,3 +888,13 @@ private final class PublicCaptureDelegate: BitchatDelegate { return publicMessages } } + +@MainActor +private final class TransportEventCaptureDelegate: TransportEventDelegate { + private(set) var messageIDs: [String] = [] + + func didReceiveTransportEvent(_ event: TransportEvent) { + guard case .messageReceived(let message) = event else { return } + messageIDs.append(message.id) + } +}