diff --git a/bitchat/Services/BLE/BLELinkAuthState.swift b/bitchat/Services/BLE/BLELinkAuthState.swift index ce31e350..3989c3c4 100644 --- a/bitchat/Services/BLE/BLELinkAuthState.swift +++ b/bitchat/Services/BLE/BLELinkAuthState.swift @@ -11,9 +11,9 @@ import Foundation /// stop a replayed announce from flip-flopping bindings or survivor /// selection. /// -/// bleQueue-confined today, alongside the link bindings it qualifies; -/// both move to the engine together in the option-B boundary flip -/// (docs/BLE-ARCHITECTURE-V3.md). +/// Engine-owned (option-B boundary, docs/BLE-ARCHITECTURE-V3.md), +/// alongside the link bindings it qualifies: BLEService debug-traps any +/// access off the engine queue. struct BLELinkAuthState { private var authenticatedOwners: [BLEIngressLinkID: PeerID] = [:] private var reconnectPolicy = BLENoiseReconnectPolicy() diff --git a/bitchat/Services/BLE/BLELinkBindings.swift b/bitchat/Services/BLE/BLELinkBindings.swift index d4c5342c..840acdfa 100644 --- a/bitchat/Services/BLE/BLELinkBindings.swift +++ b/bitchat/Services/BLE/BLELinkBindings.swift @@ -5,15 +5,18 @@ import Foundation /// belongs to, in both roles, plus each peer's preferred peripheral link /// for directed sends and fanout collapse. /// -/// Split from the physical link-state store so the option-B boundary flip -/// (docs/BLE-ARCHITECTURE-V3.md) can move ownership of *who owns a link* -/// to the engine without touching *what links exist*. bleQueue-confined -/// today, alongside the physical store and `BLELinkAuthState`. +/// Engine-owned (option-B boundary, docs/BLE-ARCHITECTURE-V3.md), +/// alongside `BLELinkAuthState`: *who owns a link* lives on the engine, +/// *what links exist* stays on bleQueue in the physical store. BLEService +/// debug-traps any access off the engine queue. /// /// Lifecycle contract: bindings are only created for live physical links -/// (callers guard existence) and are retired through -/// `peripheralRemoved`/`centralRemoved`/`clear*` when the physical link -/// goes, so binding queries never see departed links. +/// (callers check liveness through `readLinkState`) and are retired +/// through `peripheralRemoved`/`centralRemoved`/`clear*` on an engine hop +/// queued by the physical teardown. A binding can therefore briefly +/// outlive its departed link; queries that need liveness join against the +/// physical store, and everything converges once the queued retirement +/// runs. struct BLELinkBindings { private var peripheralPeers: [String: PeerID] = [:] private var centralPeers: [String: PeerID] = [:] diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index 5ad02cf9..cd366608 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -206,14 +206,35 @@ final class BLEService: NSObject { // 1. Consolidated BLE link tracking for both central and peripheral roles. private var linkStateStore = BLELinkStateStore() - // Per-link Noise authentication and rebind containment (bleQueue-owned, - // like the link store — courier handover needs the stronger fact that a + // The engine-owned identity domain: per-link Noise authentication + + // rebind containment (courier handover needs the stronger fact that a // session was established *on this current ingress link*, not merely - // that some session exists for the claimed ID). - private var linkAuth = BLELinkAuthState() - // Identity↔link bindings, split from the physical link store so the - // option-B flip can move ownership to the engine (bleQueue-owned). - private var linkBindings = BLELinkBindings() + // that some session exists for the claimed ID), and the identity↔link + // bindings that qualify every attribution decision. + // + // Owned by the engine queue since the option-B flip: bleQueue hands + // decoded packets up as (packet, linkID) and the engine attributes + // them; bleQueue never touches these. A binding can therefore briefly + // outlive its physical link (the delegate's retirement hop is async) — + // every query that needs liveness joins against the physical store, + // which the engine may sync-read via `readLinkState`. + private var _linkAuth = BLELinkAuthState() + private var _linkBindings = BLELinkBindings() + private var linkAuth: BLELinkAuthState { + get { assertLinkIdentityEngineOwned(); return _linkAuth } + set { assertLinkIdentityEngineOwned(); _linkAuth = newValue } + } + private var linkBindings: BLELinkBindings { + get { assertLinkIdentityEngineOwned(); return _linkBindings } + set { assertLinkIdentityEngineOwned(); _linkBindings = newValue } + } + /// Debug-traps any identity-domain access off the engine queue — the + /// mechanical form of the option-B ownership contract. + private func assertLinkIdentityEngineOwned() { + #if DEBUG + dispatchPrecondition(condition: .onQueue(messageQueue)) + #endif + } // BCH-01-004: Rate-limiting for subscription-triggered announces. private var subscriptionAnnounceLimiter = BLESubscriptionAnnounceLimiter() @@ -279,9 +300,6 @@ final class BLEService: NSObject { /// May block in tests to hold the serial message queue immediately before /// the deferred private-media admission check. var _test_beforePrivateMediaDeferredSend: ((String) -> Void)? - /// May block announce handling after verified-link rebind work is queued. - /// Tests use this boundary to prove rebind and reconnect are serialized. - var _test_afterVerifiedDirectRebindEnqueued: (() -> Void)? /// May block the convergence-recovery callback on its global-queue thread /// before it enqueues onto `messageQueue`. Tests use this boundary to /// force the quarantine-restore handler to win the dispatch race. @@ -367,12 +385,14 @@ final class BLEService: NSObject { /// 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. + /// Sync-edge order (deadlock freedom): main, test threads, and the + /// gossip manager's mesh.sync queue 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. (The engine only ever async-dispatches into mesh.sync; its + /// queue.sync helpers are DEBUG test entry points on test threads.) private func onEngine(_ body: () -> T) -> T { #if DEBUG dispatchPrecondition(condition: .notOnQueue(bleQueue)) @@ -758,6 +778,10 @@ final class BLEService: NSObject { ingressLinks.removeAll() recentTrafficTracker.removeAll() scheduledRelays.cancelAll() + // Proofs and revalidation epochs die with the identity; the + // rebind/retirement cooldowns deliberately survive (see + // BLELinkAuthState.removeAll). + linkAuth.removeAll() // These callbacks belong to pre-panic transfer state. Invoking // them would let queued UI work recreate or resend wiped media. privateMediaSessions.panicReset() @@ -775,7 +799,6 @@ final class BLEService: NSObject { pendingPeripheralWrites.removeAll() pendingNotifications.removeAll() pendingWriteBuffers.removeAll() - linkAuth.removeAll() radio.reset() } disconnectNotifyDebouncer.removeAll() @@ -1051,6 +1074,10 @@ final class BLEService: NSObject { // Also clear pending message queues to avoid stale state across sessions pendingNoiseSessionQueues.removeAll() pendingDirectedRelays.removeAll() + // Identity domain is engine-owned: bindings and link proofs + // clear here, physical link state clears on bleQueue below. + linkBindings.removeAll() + linkAuth.removeAll() return (transfers: entries, pingTimeouts: pingTimeouts) } @@ -1066,8 +1093,6 @@ final class BLEService: NSObject { // Clear peripheral references (synchronized access to avoid races with BLE callbacks) bleQueue.sync { linkStateStore.clearAll() - linkBindings.removeAll() - linkAuth.removeAll() radio.reset() subscriptionAnnounceLimiter.removeAll() } @@ -2185,9 +2210,12 @@ final class BLEService: NSObject { } } - /// Serializes the final authenticated-link check with CoreBluetooth's - /// notification admission on `bleQueue`, closing the rebind/disconnect - /// race between fanout planning and the actual handoff. + /// The authenticated-link eligibility check runs here on the engine — + /// the queue that owns bindings and rebinds — so fanout planning and + /// the final check are serialized against identity changes by + /// construction. Only the physical admission (updateValue and the + /// backpressure queue) hops to `bleQueue`; a central that physically + /// departs in between is a harmless no-op delivery. private func notifyOrEnqueueIfAccepted( data: Data, centrals: [CBCentral], @@ -2195,18 +2223,19 @@ final class BLEService: NSObject { context: String, requiredAuthenticatedPeer: PeerID? ) -> Bool { - let accept = { [self] in - let eligible: [CBCentral] - if let peerID = requiredAuthenticatedPeer { - eligible = centrals.filter { central in - let link = BLEIngressLinkID.central(central.identifier.uuidString) - return linkAuth.isAuthenticated(link, for: peerID) - && linkBindings.peer(forCentralUUID: central.identifier.uuidString) == peerID - } - } else { - eligible = centrals + let eligible: [CBCentral] + if let peerID = requiredAuthenticatedPeer { + eligible = centrals.filter { central in + let link = BLEIngressLinkID.central(central.identifier.uuidString) + return linkAuth.isAuthenticated(link, for: peerID) + && linkBindings.peer(forCentralUUID: central.identifier.uuidString) == peerID } - guard !eligible.isEmpty else { return false } + } else { + eligible = centrals + } + guard !eligible.isEmpty else { return false } + + let accept = { [self] in if peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: eligible) == true { return true } @@ -2216,10 +2245,7 @@ final class BLEService: NSObject { context: context ) } - - if DispatchQueue.getSpecific(key: bleQueueKey) != nil { - return accept() - } + // queue-contract-ok: engine → bleQueue is the sanctioned sync direction. return bleQueue.sync(execute: accept) } @@ -2260,7 +2286,7 @@ final class BLEService: NSObject { let centralIDs = subscribedCentrals.map { $0.identifier.uuidString } let peripheralPeerBindings = Dictionary(uniqueKeysWithValues: connectedStates.compactMap { state -> (String, PeerID)? in let uuid = state.peripheral.identifier.uuidString - return readLinkState { _ in linkBindings.peer(forPeripheralID: uuid) }.map { (uuid, $0) } + return linkBindings.peer(forPeripheralID: uuid).map { (uuid, $0) } }) let plan = BLEOutboundLinkPlanner.plan( packet: packet, @@ -2273,10 +2299,7 @@ final class BLEService: NSObject { excludedLinks: excludedPeerLinks, peripheralPeerBindings: peripheralPeerBindings, centralPeerBindings: centralSnapshot.peerIDsByCentralUUID, - // Perf note: this is a third bleQueue hop per send; if send-path - // profiling ever flags it, fold it into snapshotPeripheralStates - // as a combined snapshot. - preferredPeripheralPerPeer: readLinkState { _ in linkBindings.preferredPeripheralBindings }, + preferredPeripheralPerPeer: linkBindings.preferredPeripheralBindings, directAnnounceTTL: messageTTL, directedOnlyPeer: directedOnlyPeer, requireDirectPeerLink: requireDirectPeerLink || requireNoiseAuthenticatedPeerLink @@ -2696,9 +2719,7 @@ final class BLEService: NSObject { // A valid departure retires transport state too; otherwise // canDeliverSecurely could remain true for a peer we just removed. clearNoiseSession(for: peerID) - readLinkState { _ in - _ = linkAuth.retireLinks(ownedBy: peerID) - } + _ = linkAuth.retireLinks(ownedBy: peerID) // Remove the peer when they leave peerRegistry.mutate { _ = $0.remove(peerID) } // Remove any stored announcement for sync purposes @@ -2903,12 +2924,22 @@ final class BLEService: NSObject { // MARK: - GossipSyncManager Delegate extension BLEService: GossipSyncManager.Delegate { + // Gossip calls arrive on the manager's own serial queue; sends read + // the engine-owned bindings, so they enter an engine slot. The sync + // hop is safe: mesh.sync sits above the engine in the sync order — + // production engine code only ever queue.async's into the manager + // (the queue.sync helpers are DEBUG test entry points that run on + // test threads), so no reverse edge exists. func sendPacket(_ packet: BitchatPacket) { - broadcastPacket(packet) + onEngine { + broadcastPacket(packet) + } } func sendPacket(to peerID: PeerID, packet: BitchatPacket) { - sendPacketDirected(packet, to: peerID) + onEngine { + sendPacketDirected(packet, to: peerID) + } } func signPacketForBroadcast(_ packet: BitchatPacket) -> BitchatPacket { @@ -3011,18 +3042,22 @@ extension BLEService: CBCentralManagerDelegate { // not issue stop/cancel commands now; they are rejected as API // misuse. Retire our link state locally instead. SecureLogger.info("📴 Bluetooth powered off - cleaning up central state", category: .session) - let peripheralStates = linkStateStore.peripheralStates - for state in peripheralStates { - let peripheralID = state.peripheral.identifier.uuidString + let peripheralIDs = linkStateStore.peripheralStates.map { $0.peripheral.identifier.uuidString } + for peripheralID in peripheralIDs { pendingPeripheralWrites.discardAll(for: peripheralID) - linkAuth.retireLink(.peripheral(peripheralID)) } linkStateStore.clearPeripherals() - let peerIDs = linkBindings.clearPeripherals() - // Notify UI of disconnections - for peerID in peerIDs { - notifyUI { [weak self] in - self?.notifyPeerDisconnectedDebounced(peerID) + messageQueue.async { [weak self] in + guard let self else { return } + for peripheralID in peripheralIDs { + self.linkAuth.retireLink(.peripheral(peripheralID)) + } + let peerIDs = self.linkBindings.clearPeripherals() + // Notify UI of disconnections + for peerID in peerIDs { + self.notifyUI { [weak self] in + self?.notifyPeerDisconnectedDebounced(peerID) + } } } @@ -3030,7 +3065,9 @@ extension BLEService: CBCentralManagerDelegate { // User denied Bluetooth permission SecureLogger.warning("🚫 Bluetooth unauthorized - user denied permission", category: .session) linkStateStore.clearPeripherals() - _ = linkBindings.clearPeripherals() + messageQueue.async { [weak self] in + _ = self?.linkBindings.clearPeripherals() + } case .unsupported: // Device doesn't support BLE @@ -3083,11 +3120,8 @@ extension BLEService: CBCentralManagerDelegate { func centralManager(_ central: CBCentralManager, didDisconnectPeripheral peripheral: CBPeripheral, error: Error?) { let peripheralID = peripheral.identifier.uuidString - - // Find the peer ID if we have it - let peerID = linkBindings.peer(forPeripheralID: peripheralID) - - SecureLogger.debug("📱 Disconnect: \(peerID?.id ?? peripheralID)\(error != nil ? " (\(error!.localizedDescription))" : "")", category: .session) + + SecureLogger.debug("📱 Disconnect: \(peripheralID)\(error != nil ? " (\(error!.localizedDescription))" : "")", category: .session) // If disconnect carried an error (often timeout), apply short backoff to avoid thrash if error != nil { @@ -3113,24 +3147,46 @@ extension BLEService: CBCentralManagerDelegate { } #endif - // Clean up references and peer mappings - tearDownPeripheralLink(peripheralID) - // A duplicate link can drop while the peer stays live on another - // (the dual-role central link, or a second bound link after a - // restore): peer-disconnect bookkeeping only runs once the peer's - // last live link is gone. removePeripheral just repaired the reverse - // map onto a connected survivor, so directLinkState is accurate - // here. The scan restart and connect-slot refill below stay - // unguarded — they respond to the physical drop regardless of - // remaining logical links. - let remainingLinks = peerID.map { directLinkState(for: $0) } - let peerStillLinked = (remainingLinks?.hasPeripheral ?? false) || (remainingLinks?.hasCentral ?? false) - if let peerID, !peerStillLinked { - // Do not remove peer; mark as not connected but retain for reachability - peerRegistry.mutate { $0.markDisconnected(peerID) } - refreshLocalTopology() - } + // Physical teardown now; identity retirement and peer-disconnect + // bookkeeping on the engine, which owns the bindings. The scan + // restart and connect-slot refill below stay on bleQueue — they + // respond to the physical drop regardless of remaining logical + // links. + discardPeripheralLinkPhysical(peripheralID) + messageQueue.async { [weak self] in + guard let self else { return } + // A duplicate link can drop while the peer stays live on + // another (the dual-role central link, or a second bound link + // after a restore): peer-disconnect bookkeeping only runs once + // the peer's last live link is gone. The retirement just + // repaired the reverse map onto a connected survivor, so + // directLinkState is accurate here. + let peerID = self.retirePeripheralLinkIdentity(peripheralID) + if let peerID { + SecureLogger.debug("📱 Disconnected link was bound to \(peerID.id.prefix(8))…", category: .session) + } + let remainingLinks = peerID.map { self.directLinkState(for: $0) } + let peerStillLinked = (remainingLinks?.hasPeripheral ?? false) || (remainingLinks?.hasCentral ?? false) + if let peerID, !peerStillLinked { + // Do not remove peer; mark as not connected but retain for reachability + self.peerRegistry.mutate { $0.markDisconnected(peerID) } + self.refreshLocalTopology() + } + // Notify delegate about disconnection on main thread (direct link dropped) + self.notifyUI { [weak self] in + guard let self = self else { return } + + // Get current peer list (after removal) + let currentPeerIDs = self.peerRegistry.peerIDs + + if let peerID, !peerStillLinked { + self.notifyPeerDisconnectedDebounced(peerID) + } + self.requestPeerDataPublish() + self.deliverTransportEvent(.peerListUpdated(currentPeerIDs)) + } + } // Restart scanning with allow duplicates for faster rediscovery if centralManager?.state == .poweredOn { @@ -3142,27 +3198,16 @@ extension BLEService: CBCentralManagerDelegate { } // Attempt to fill freed slot from queue bleQueue.async { [weak self] in self?.radio.tryConnectFromQueue() } - - // Notify delegate about disconnection on main thread (direct link dropped) - notifyUI { [weak self] in - guard let self = self else { return } - - // Get current peer list (after removal) - let currentPeerIDs = self.peerRegistry.peerIDs - - if let peerID, !peerStillLinked { - self.notifyPeerDisconnectedDebounced(peerID) - } - self.requestPeerDataPublish() - self.deliverTransportEvent(.peerListUpdated(currentPeerIDs)) - } } func centralManager(_ central: CBCentralManager, didFailToConnect peripheral: CBPeripheral, error: Error?) { let peripheralID = peripheral.identifier.uuidString - - // Clean up the references - tearDownPeripheralLink(peripheralID) + + // Clean up the references: physical now, identity on the engine. + discardPeripheralLinkPhysical(peripheralID) + messageQueue.async { [weak self] in + self?.retirePeripheralLinkIdentity(peripheralID) + } SecureLogger.error("❌ Failed to connect to peripheral: \(peripheral.name ?? "Unknown") [\(peripheralID)] - Error: \(error?.localizedDescription ?? "Unknown")", category: .session) radio.recordConnectionFailure(peripheralID: peripheralID) @@ -3190,48 +3235,59 @@ extension BLEService: BLERadioControllerDelegate { } func radioTearDownPeripheralLink(_ peripheralID: String) { - tearDownPeripheralLink(peripheralID) + // bleQueue (the controller's queue): physical discard now, identity + // retirement on the engine. + discardPeripheralLinkPhysical(peripheralID) + messageQueue.async { [weak self] in + self?.retirePeripheralLinkIdentity(peripheralID) + } } - /// Retires one peripheral link's transport bookkeeping: its write - /// backpressure, its Noise link proof and reconnect epoch, and the - /// link-state entry plus binding (which repairs the peer's reverse - /// mapping onto a surviving duplicate link). bleQueue-confined. - func tearDownPeripheralLink(_ peripheralID: String) { + /// bleQueue half of a peripheral-link teardown: the link's write + /// backpressure and its physical link-state entry. Identity retirement + /// (proof, epoch, binding repair) rides a separate engine hop — + /// `retirePeripheralLinkIdentity`. bleQueue-confined. + func discardPeripheralLinkPhysical(_ peripheralID: String) { pendingPeripheralWrites.discardAll(for: peripheralID) - linkAuth.retireLink(.peripheral(peripheralID)) - removePeripheralLink(peripheralID) + linkStateStore.removePeripheral(peripheralID) } - /// Physical removal plus binding retirement, one unit: the preferred - /// link repairs onto a connected survivor, preferring a writable one + /// Engine half of a peripheral-link teardown: retires the link's Noise + /// proof and revalidation epoch, and its binding — repairing the peer's + /// preferred link onto a connected survivor, preferring a writable one /// (a link mid-service-rediscovery would strand directed sends until - /// its characteristic comes back). bleQueue-confined. + /// its characteristic comes back). Returns the peer that owned the + /// binding. Engine-confined. @discardableResult - func removePeripheralLink(_ peripheralID: String) -> PeerID? { - linkStateStore.removePeripheral(peripheralID) + func retirePeripheralLinkIdentity(_ peripheralID: String) -> PeerID? { + linkAuth.retireLink(.peripheral(peripheralID)) return linkBindings.peripheralRemoved(peripheralID) { remaining in - let alive = remaining.compactMap { uuid -> (uuid: String, writable: Bool)? in - guard let state = linkStateStore.state(forPeripheralID: uuid), - state.isConnected else { return nil } - return (uuid, state.characteristic != nil) + let alive = readLinkState { store in + remaining.compactMap { uuid -> (uuid: String, writable: Bool)? in + guard let state = store.state(forPeripheralID: uuid), + state.isConnected else { return nil } + return (uuid, state.characteristic != nil) + } } return (alive.first(where: \.writable) ?? alive.first)?.uuid } } /// Binds only live physical links, preserving the store-era guard that - /// a binding can never outlive (or precede) its link. bleQueue-confined. + /// a binding can never precede its link (a lost race against a + /// concurrent physical removal is healed by that removal's queued + /// identity retirement). Engine-confined. func bindPeripheralLink(_ peripheralUUID: String, to peerID: PeerID) { - guard linkStateStore.state(forPeripheralID: peripheralUUID) != nil else { return } + guard readLinkState({ $0.state(forPeripheralID: peripheralUUID) }) != nil else { return } linkBindings.bindPeripheral(peripheralUUID, to: peerID) } - /// Whether the peer holds a live direct link in either role. - /// bleQueue-confined (physical liveness + bindings in one view). + /// Whether the peer holds a live direct link in either role: bindings + /// (engine) joined against physical liveness (readLinkState). + /// Engine-confined. func directLinkState(for peerID: PeerID) -> BLEDirectLinkState { let hasPeripheral = linkBindings.preferredPeripheralUUID(for: peerID) - .flatMap { linkStateStore.state(forPeripheralID: $0)?.isConnected } ?? false + .flatMap { uuid in readLinkState { $0.state(forPeripheralID: uuid)?.isConnected } } ?? false return BLEDirectLinkState( hasPeripheral: hasPeripheral, hasCentral: linkBindings.hasCentral(boundTo: peerID) @@ -3239,16 +3295,16 @@ extension BLEService: BLERadioControllerDelegate { } /// The peer's preferred peripheral link state, when physically present. - /// bleQueue-confined. + /// Engine-confined. func directPeripheralState(for peerID: PeerID) -> BLEPeripheralLinkState? { linkBindings.preferredPeripheralUUID(for: peerID) - .flatMap { linkStateStore.state(forPeripheralID: $0) } + .flatMap { uuid in readLinkState { $0.state(forPeripheralID: uuid) } } } - /// Subscribed centrals with their bindings, one view. bleQueue-confined. + /// Subscribed centrals with their bindings, one view. Engine-confined. func subscribedCentralSnapshot() -> BLESubscribedCentralSnapshot { BLESubscribedCentralSnapshot( - centrals: linkStateStore.subscribedCentrals, + centrals: readLinkState(\.subscribedCentrals), peerIDsByCentralUUID: linkBindings.centralPeersByUUID ) } @@ -3376,22 +3432,22 @@ extension BLEService { } func _test_bindCentral(_ centralUUID: String, to peerID: PeerID) { - bleQueue.sync { linkBindings.bindCentral(centralUUID, to: peerID) } + onEngine { linkBindings.bindCentral(centralUUID, to: peerID) } } func _test_centralBinding(_ centralUUID: String) -> PeerID? { - bleQueue.sync { linkBindings.peer(forCentralUUID: centralUUID) } + onEngine { linkBindings.peer(forCentralUUID: centralUUID) } } func _test_markNoiseAuthenticatedCentral(_ centralUUID: String, to peerID: PeerID) { - bleQueue.sync { + onEngine { guard linkBindings.peer(forCentralUUID: centralUUID) == peerID else { return } linkAuth.markAuthenticated(.central(centralUUID), owner: peerID) } } func _test_isNoiseAuthenticatedCentral(_ centralUUID: String, for peerID: PeerID) -> Bool { - bleQueue.sync { + onEngine { linkAuth.isAuthenticated(.central(centralUUID), for: peerID) } } @@ -3702,72 +3758,25 @@ extension BLEService: CBPeripheralDelegate { SecureLogger.error("❌ Invalid BLE frame length; reset notification stream", category: .session) } - // Codex review identified TOCTOU in this patch. - // Enforce per-link sender binding immediately within the same notification batch. - // NOTE: `processNotificationPacket` may bind the stored peer ID when an announce - // is processed, but `state` above is a snapshot. Track a local binding that we update as soon as - // we see a binding-eligible announce so subsequent frames can't spoof a different sender. - var boundPeerID: PeerID? = linkBindings.peer(forPeripheralID: peripheralUUID) - + // Attribution — spoof rejection, announce binding, ingress + // recording — is engine work now (the engine owns the bindings). + // Frames hop up in decode order; the engine's serial slot ordering + // gives the same same-batch spoof protection the old bleQueue-side + // batch-local binding enforced: an announce that binds this link is + // attributed before every frame that rode behind it. for frame in result.frames { guard let packet = BinaryProtocol.decode(frame) else { let prefix = frame.prefix(16).map { String(format: "%02x", $0) }.joined(separator: " ") SecureLogger.error("❌ Failed to decode assembled notification frame (len=\(frame.count), prefix=\(prefix))", category: .session) continue } - - let claimedSenderID = PeerID(hexData: packet.senderID) - let context = acceptedIngressContext( - for: packet, - claimedSenderID: claimedSenderID, - boundPeerID: boundPeerID, + ingestDecodedPacket( + packet, + link: .peripheral(peripheralUUID), linkDescription: "Peripheral \(peripheralUUID.prefix(8))…" ) - - guard let context else { continue } - - // If this is a direct-link announce, bind immediately for the remainder of this batch. - if boundPeerID == nil, - packet.type == MessageType.announce.rawValue, - packet.ttl == messageTTL { - boundPeerID = claimedSenderID - bindPeripheralLink(peripheralUUID, to: claimedSenderID) - } - - if !recordIngressIfNew(packet, link: .peripheral(peripheralUUID), peerID: context.receivedFromPeerID) { - continue - } - processNotificationPacket( - packet, - from: peripheral, - peripheralUUID: peripheralUUID, - receivedFrom: context.receivedFromPeerID - ) } } - - private func processNotificationPacket(_ packet: BitchatPacket, from _: CBPeripheral, peripheralUUID: String, receivedFrom peerID: PeerID) { - let senderID = PeerID(hexData: packet.senderID) - - if packet.type != MessageType.announce.rawValue { - SecureLogger.debug("📦 Decoded notification packet type: \(packet.type) from sender: \(senderID.id.prefix(8))…", category: .session) - } - - if packet.type == MessageType.announce.rawValue, - packet.ttl == messageTTL { - // Only bind an unbound link here: this runs before signature - // verification, so a bound link must not be re-bound by a raw - // announce (spoofable). Rotation rebinds happen after the announce - // verifies (rebindLinkAfterVerifiedDirectAnnounce). - let boundPeerID = linkBindings.peer(forPeripheralID: peripheralUUID) - if boundPeerID == nil || boundPeerID == senderID { - bindPeripheralLink(peripheralUUID, to: senderID) - refreshLocalTopology() - } - } - - handleReceivedPacket(packet, from: peerID) - } func peripheral(_ peripheral: CBPeripheral, didWriteValueFor characteristic: CBCharacteristic, error: Error?) { if let error = error { @@ -3862,21 +3871,23 @@ extension BLEService: CBPeripheralManagerDelegate { // Bluetooth was turned off - clean up peripheral state SecureLogger.info("📴 Bluetooth powered off - cleaning up peripheral state", category: .session) // Clear subscribed centrals (they are now invalid) - let centralSnapshot = subscribedCentralSnapshot() - for central in centralSnapshot.centrals { - let centralID = central.identifier.uuidString - linkAuth.retireLink(.central(centralID)) - } + let centralIDs = linkStateStore.subscribedCentrals.map { $0.identifier.uuidString } pendingNotifications.removeAll() pendingWriteBuffers.removeAll() linkStateStore.clearCentrals() - let centralPeerIDs = linkBindings.clearCentrals() subscriptionAnnounceLimiter.removeAll() characteristic = nil - // Notify UI of disconnections - for peerID in centralPeerIDs { - notifyUI { [weak self] in - self?.notifyPeerDisconnectedDebounced(peerID) + messageQueue.async { [weak self] in + guard let self else { return } + for centralID in centralIDs { + self.linkAuth.retireLink(.central(centralID)) + } + let centralPeerIDs = self.linkBindings.clearCentrals() + // Notify UI of disconnections + for peerID in centralPeerIDs { + self.notifyUI { [weak self] in + self?.notifyPeerDisconnectedDebounced(peerID) + } } } @@ -3884,9 +3895,11 @@ extension BLEService: CBPeripheralManagerDelegate { // User denied Bluetooth permission SecureLogger.warning("🚫 Bluetooth unauthorized for peripheral role", category: .session) linkStateStore.clearCentrals() - _ = linkBindings.clearCentrals() subscriptionAnnounceLimiter.removeAll() characteristic = nil + messageQueue.async { [weak self] in + _ = self?.linkBindings.clearCentrals() + } case .unsupported: // Device doesn't support BLE peripheral role @@ -3992,38 +4005,41 @@ extension BLEService: CBPeripheralManagerDelegate { func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didUnsubscribeFrom characteristic: CBCharacteristic) { let centralID = central.identifier.uuidString SecureLogger.debug("📤 Central unsubscribed: \(centralID.prefix(8))…", category: .session) + // bleQueue: physical retirement now. pendingNotifications.removeTarget { $0.identifier.uuidString == centralID } - linkAuth.retireLink(.central(centralID)) linkStateStore.removeSubscribedCentral(central) - let removedPeerID = linkBindings.centralRemoved(centralID) - + // Ensure we're still advertising for other devices to find us if !isPanicSuspended, peripheral.isAdvertising == false { SecureLogger.debug("📡 Restarting advertising after central unsubscribed", category: .session) peripheral.startAdvertising(BLERadioController.advertisementData()) } - - // Find and disconnect the peer associated with this central - if let peerID = removedPeerID { + + // Identity retirement and peer-disconnect bookkeeping on the + // engine, which owns the bindings. + messageQueue.async { [weak self] in + guard let self else { return } + self.linkAuth.retireLink(.central(centralID)) + guard let peerID = self.linkBindings.centralRemoved(centralID) else { return } // The remote side retiring a redundant duplicate connection // arrives here as an unsubscribe while the peer stays live on // its other links; only the peer's last link disconnecting // counts. If every link truly dropped, the surviving-link // callbacks (didDisconnectPeripheral, or this one again) run // the bookkeeping. - guard linkBindings.links(to: peerID).isEmpty else { return } + guard self.linkBindings.links(to: peerID).isEmpty else { return } // Mark peer as not connected; retain for reachability - peerRegistry.mutate { $0.markDisconnected(peerID) } - - refreshLocalTopology() - + self.peerRegistry.mutate { $0.markDisconnected(peerID) } + + self.refreshLocalTopology() + // Update UI immediately - notifyUI { [weak self] in + self.notifyUI { [weak self] in guard let self = self else { return } - + // Get current peer list (after removal) let currentPeerIDs = self.peerRegistry.peerIDs - + self.notifyPeerDisconnectedDebounced(peerID) // Publish snapshots so UnifiedPeerService can refresh icons promptly self.requestPeerDataPublish() @@ -4150,37 +4166,16 @@ extension BLEService: CBPeripheralManagerDelegate { } private func processDecodedCentralWrite(_ packet: BitchatPacket, centralUUID: String, central: CBCentral) { - let claimedSenderID = PeerID(hexData: packet.senderID) - let context = acceptedIngressContext( - for: packet, - claimedSenderID: claimedSenderID, - boundPeerID: linkBindings.peer(forCentralUUID: centralUUID), + // bleQueue: physical bookkeeping only. A writer is a live central + // whether or not it subscribed; track it so directed replies and + // the fanout planner can reach it. + linkStateStore.addSubscribedCentral(central) + // Attribution is engine work (the engine owns the bindings). + ingestDecodedPacket( + packet, + link: .central(centralUUID), linkDescription: "Central \(centralUUID.prefix(8))…" ) - guard let context else { return } - - if packet.type != MessageType.announce.rawValue { - SecureLogger.debug("📦 Decoded (combined) packet type: \(packet.type) from sender: \(claimedSenderID.id.prefix(8))…", category: .session) - } - - linkStateStore.addSubscribedCentral(central) - - if packet.type == MessageType.announce.rawValue, - packet.ttl == messageTTL { - // Same rule as the peripheral path: raw announces only bind - // unbound links; rotation rebinds require a verified announce. - let boundPeerID = linkBindings.peer(forCentralUUID: centralUUID) - if boundPeerID == nil || boundPeerID == claimedSenderID { - linkBindings.bindCentral(centralUUID, to: claimedSenderID) - refreshLocalTopology() - } - } - - guard recordIngressIfNew(packet, link: .central(centralUUID), peerID: context.receivedFromPeerID) else { - return - } - - handleReceivedPacket(packet, from: context.receivedFromPeerID) } } @@ -4711,14 +4706,15 @@ extension BLEService { return plan.shouldSuppressFloodRelay } - /// Safely fetch the current direct-link state for a peer using the BLE queue. + /// The current direct-link state for a peer. Engine-confined (bindings + /// joined against physical liveness inside directLinkState). private func linkState(for peerID: PeerID) -> (hasPeripheral: Bool, hasCentral: Bool) { - let state = readLinkState { _ in directLinkState(for: peerID) } + let state = directLinkState(for: peerID) return (state.hasPeripheral, state.hasCentral) } private func links(to peerID: PeerID?) -> Set { - readLinkState { _ in linkBindings.links(to: peerID) } + linkBindings.links(to: peerID) } @@ -4726,19 +4722,16 @@ extension BLEService { /// Marks the exact physical ingress link that completed a fresh Noise /// handshake. An old session keyed only by peer ID is insufficient: a /// replayed announce can rebind an attacker's link to that ID. + /// Engine-confined. private func markNoiseAuthenticatedIngressLink(for packet: BitchatPacket, peerID: PeerID) { guard let link = ingressLinks.link(for: packet) else { return } - readLinkState { store in - guard linkBindings.boundPeer(for: link) == peerID else { return } - linkAuth.markAuthenticated(link, owner: peerID) - } + guard linkBindings.boundPeer(for: link) == peerID else { return } + linkAuth.markAuthenticated(link, owner: peerID) } private func isNoiseAuthenticatedIngressLink(for packet: BitchatPacket, peerID: PeerID) -> Bool { guard let link = ingressLinks.link(for: packet) else { return false } - return readLinkState { store in - linkAuth.isAuthenticated(link, for: peerID) && linkBindings.boundPeer(for: link) == peerID - } + return linkAuth.isAuthenticated(link, for: peerID) && linkBindings.boundPeer(for: link) == peerID } private func hasCurrentNoiseAuthenticatedLink(to peerID: PeerID) -> Bool { @@ -4746,37 +4739,36 @@ extension BLEService { } private func currentNoiseAuthenticatedLinks(to peerID: PeerID) -> Set { - readLinkState { store in - Set(linkAuth.links(ownedBy: peerID).filter { link in - linkBindings.boundPeer(for: link) == peerID - }) - } + Set(linkAuth.links(ownedBy: peerID).filter { link in + linkBindings.boundPeer(for: link) == peerID + }) } /// 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`). + /// Takes the already-resolved ingress link. Engine-confined: it runs + /// inside the rebind's engine slot, so no observer can see the new + /// binding while a cached peer-level sender is still considered + /// established. private func refreshNoiseSessionForVerifiedDirectLink( link: BLEIngressLinkID, peerID: PeerID ) { let hasEstablishedSession = noiseService.hasEstablishedSession(with: peerID) let authenticatedPeerLinks = currentNoiseAuthenticatedLinks(to: peerID) - let shouldRevalidate = readLinkState { store in - guard linkBindings.boundPeer(for: link) == peerID else { - return false - } - return linkAuth.shouldRevalidate( + let shouldRevalidate: Bool + if linkBindings.boundPeer(for: link) == peerID { + shouldRevalidate = linkAuth.shouldRevalidate( on: link, for: peerID, hasEstablishedSession: hasEstablishedSession, hasAuthenticatedPeerLink: !authenticatedPeerLinks.isEmpty, now: Date() ) + } else { + shouldRevalidate = false } guard shouldRevalidate else { return } @@ -5318,11 +5310,13 @@ extension BLEService { /// replay-rebound link, or process-local spool is not delivery. @discardableResult func deliverBridgedEnvelope(_ envelope: CourierEnvelope, to peerID: PeerID) -> Bool { - guard hasCurrentNoiseAuthenticatedLink(to: peerID) else { return false } guard let payload = envelope.encode() else { return false } let packet = makeCourierPacket(payload, to: peerID) return onEngine { - sendPacketDirected( + // Engine slot: the auth-link check and the directed send see one + // consistent view of the identity domain. + guard hasCurrentNoiseAuthenticatedLink(to: peerID) else { return false } + return sendPacketDirected( packet, to: peerID, requireDirectPeerLink: true, @@ -5813,7 +5807,10 @@ extension BLEService { } } - // MARK: Link capability snapshots (thread-safe via bleQueue) + // MARK: Link capability snapshots + // Physical link state is bleQueue-owned; the engine (and main) may + // sync-read it here. The bindings half of a combined view comes from + // the engine-owned identity domain directly. private func readLinkState(_ body: (BLELinkStateStore) -> T) -> T { if DispatchQueue.getSpecific(key: bleQueueKey) != nil { @@ -5824,7 +5821,7 @@ extension BLEService { } private func snapshotDirectPeripheralState(for peerID: PeerID) -> BLEPeripheralLinkState? { - readLinkState { _ in directPeripheralState(for: peerID) } + directPeripheralState(for: peerID) } private func snapshotPeripheralStates() -> [BLEPeripheralLinkState] { @@ -5832,7 +5829,7 @@ extension BLEService { } private func snapshotSubscribedCentrals() -> BLESubscribedCentralSnapshot { - readLinkState { _ in subscribedCentralSnapshot() } + subscribedCentralSnapshot() } // MARK: Helpers: IDs, selection, and write backpressure @@ -5868,6 +5865,11 @@ extension BLEService { /// peripheral's bounded retry queue. Unlike `writeOrEnqueue`, the return /// value distinguishes a retained queue item from one rejected or trimmed /// immediately, which lets durable courier state commit truthfully. + /// + /// The authenticated-link eligibility check runs on the engine (which + /// owns bindings and rebinds, so it is serialized against identity + /// changes by construction); only the physical admission hops to + /// `bleQueue`. private func writeOrEnqueueIfAccepted( _ data: Data, to peripheral: CBPeripheral, @@ -5875,20 +5877,20 @@ extension BLEService { priority: BLEOutboundWritePriority, requiredAuthenticatedPeer: PeerID? ) -> Bool { + let uuid = peripheral.identifier.uuidString + if let peerID = requiredAuthenticatedPeer { + let link = BLEIngressLinkID.peripheral(uuid) + guard linkBindings.peer(forPeripheralID: uuid) == peerID, + linkAuth.isAuthenticated(link, for: peerID) else { + return false + } + } let accept = { [self] in - let uuid = peripheral.identifier.uuidString guard let state = linkStateStore.state(forPeripheralID: uuid), state.isConnected, state.characteristic?.uuid == characteristic.uuid else { return false } - if let peerID = requiredAuthenticatedPeer { - let link = BLEIngressLinkID.peripheral(uuid) - guard linkBindings.peer(forPeripheralID: uuid) == peerID, - linkAuth.isAuthenticated(link, for: peerID) else { - return false - } - } if peripheral.canSendWriteWithoutResponse { peripheral.writeValue(data, for: characteristic, type: .withoutResponse) @@ -6467,7 +6469,81 @@ extension BLEService { } // MARK: Packet Reception - + + /// The bleQueue → engine handoff for every frame the link layer + /// decodes: the radio side hands up (packet, linkID) and all + /// attribution — binding lookup, spoof rejection, raw-announce + /// binding, ingress recording — happens on the engine, the queue that + /// owns the identity domain. Captures the panic lifecycle at the + /// handoff, like `handleReceivedPacket`. + /// + /// Per-link frame order is preserved end to end (bleQueue and the + /// engine are both serial), so an announce that binds a link is + /// attributed before the directed frames that ride behind it — the + /// same-batch spoof protection the old bleQueue-side attribution + /// enforced with a batch-local binding. + private func ingestDecodedPacket( + _ packet: BitchatPacket, + link: BLEIngressLinkID, + linkDescription: String + ) { + guard let lifecycleGeneration = capturePanicLifecycleGeneration() else { return } + messageQueue.async { [weak self] in + guard let self, + self.isCurrentPanicLifecycleGeneration(lifecycleGeneration) else { + return + } + self.attributeAndHandlePacket(packet, link: link, linkDescription: linkDescription) + } + } + + /// Engine-confined attribution: resolves the link's bound owner, + /// admits or rejects the claimed sender, lets a direct raw announce + /// bind an unbound link (rotation rebinds still require a verified + /// announce — `rebindLinkAfterVerifiedDirectAnnounce`), records + /// ingress, and hands the packet to the handler pipeline. + private func attributeAndHandlePacket( + _ packet: BitchatPacket, + link: BLEIngressLinkID, + linkDescription: String + ) { + let claimedSenderID = PeerID(hexData: packet.senderID) + let context = acceptedIngressContext( + for: packet, + claimedSenderID: claimedSenderID, + boundPeerID: linkBindings.boundPeer(for: link), + linkDescription: linkDescription + ) + guard let context else { return } + + if packet.type != MessageType.announce.rawValue { + SecureLogger.debug("📦 Decoded packet type: \(packet.type) from sender: \(claimedSenderID.id.prefix(8))… (\(linkDescription))", category: .session) + } + + if packet.type == MessageType.announce.rawValue, + packet.ttl == messageTTL { + // Raw announces only bind unbound links: this runs before + // signature verification, so a bound link must not be re-bound + // by a raw announce (spoofable). + let boundPeerID = linkBindings.boundPeer(for: link) + if boundPeerID == nil || boundPeerID == claimedSenderID { + switch link { + case .peripheral(let peripheralUUID): + bindPeripheralLink(peripheralUUID, to: claimedSenderID) + case .central(let centralUUID): + linkBindings.bindCentral(centralUUID, to: claimedSenderID) + } + refreshLocalTopology() + } + } + + guard recordIngressIfNew(packet, link: link, peerID: context.receivedFromPeerID) else { + return + } + + handleReceivedPacket(packet, from: context.receivedFromPeerID) + } + private func handleReceivedPacket(_ packet: BitchatPacket, from peerID: PeerID) { let isNoisePacket = packet.type == MessageType.noiseHandshake.rawValue || packet.type == MessageType.noiseEncrypted.rawValue @@ -6708,9 +6784,6 @@ extension BLEService { // consolidate duplicate same-role connections onto that link. if let result, result.isVerified, result.isDirectAnnounce { rebindLinkAfterVerifiedDirectAnnounce(packet, to: result.peerID) - #if DEBUG - _test_afterVerifiedDirectRebindEnqueued?() - #endif retireRedundantPeripheralLinks(packet, to: result.peerID) } @@ -6761,92 +6834,90 @@ extension BLEService { /// spoofed. A signature-verified direct announce proves the claimed /// sender owns the link it arrived on, so rebind the link to the new ID /// and retire the old identity. + /// Engine-confined: the whole rebind — containment checks, proof + /// retirement, binding flip, reconnect decision, and rotated-identity + /// retirement — is one engine slot, so no observer can see a + /// half-applied rotation. Only the physical connection cancels hop to + /// bleQueue. private func rebindLinkAfterVerifiedDirectAnnounce(_ packet: BitchatPacket, to peerID: PeerID) { guard let link = ingressLinks.link(for: packet) else { return } - bleQueue.async { [weak self] in - guard let self else { return } - let linkUUID: String - let previousPeerID: PeerID? - switch link { - case .peripheral(let peripheralUUID): - linkUUID = peripheralUUID - previousPeerID = self.linkBindings.peer(forPeripheralID: peripheralUUID) - case .central(let centralUUID): - linkUUID = centralUUID - previousPeerID = self.linkBindings.peer(forCentralUUID: centralUUID) - } - guard let previousPeerID else { return } - guard previousPeerID != peerID else { - self.refreshNoiseSessionForVerifiedDirectLink( - link: link, - peerID: peerID - ) - return - } - - // The signature does not authenticate directness (TTL is excluded - // from signing because relays mutate it), so a "verified direct" - // announce can be a replay of another peer's fresh announce with - // its TTL restored. Contain what a forged rebind could do: - // never steal an identity another live link already owns, and - // allow at most one rebind per link per cooldown window so two - // identities can't fight over a link in a replay flip-flop. - guard self.linkBindings.links(to: peerID).isEmpty else { - SecureLogger.warning("🚫 Refusing link rebind to \(peerID.id.prefix(8))…: identity already owns another live link", category: .security) - return - } - let now = Date() - guard self.linkAuth.permitRebind( - linkUUID: linkUUID, - now: now, - cooldown: TransportConfig.bleLinkRebindCooldownSeconds - ) else { - SecureLogger.warning("🚫 Refusing link rebind to \(peerID.id.prefix(8))…: rebind cooldown active for this link", category: .security) - return - } - - // A Noise proof belongs to the old physical binding. Never carry - // it across an announce-driven rebind, whose direct TTL is - // replayable; the new owner must complete a fresh handshake. - self.linkAuth.retireLink(link) - switch link { - case .peripheral(let peripheralUUID): - self.bindPeripheralLink(peripheralUUID, to: peerID) - case .central(let centralUUID): - self.linkBindings.bindCentral(centralUUID, to: peerID) - } - // Keep the rebind and reconnect decision in one bleQueue critical - // section. No observer may see the new binding while a cached - // peer-level sender is still considered established. - self.refreshNoiseSessionForVerifiedDirectLink( + let linkUUID: String + let previousPeerID: PeerID? + switch link { + case .peripheral(let peripheralUUID): + linkUUID = peripheralUUID + previousPeerID = linkBindings.peer(forPeripheralID: peripheralUUID) + case .central(let centralUUID): + linkUUID = centralUUID + previousPeerID = linkBindings.peer(forCentralUUID: centralUUID) + } + guard let previousPeerID else { return } + guard previousPeerID != peerID else { + refreshNoiseSessionForVerifiedDirectLink( link: link, peerID: peerID ) - SecureLogger.debug("🔄 Rebinding link after peer-ID rotation: \(previousPeerID.id.prefix(8))… → \(peerID.id.prefix(8))…", category: .session) - self.refreshLocalTopology() - // The announce that triggered this rebind was upserted as - // disconnected: the registry ran while the link still belonged - // to the previous ID (the ambiguous state BLEAnnounceHandler - // denies the connected shortcut). The rebind has now - // containment-checked the claim and the identity owns a live - // link, so promote it — otherwise a healed rotation leaves a - // live link that reads as disconnected until the next announce. - self.messageQueue.async { [weak self] in - self?.promoteReboundPeerToConnected(peerID) - } - // Any other peripheral links still bound to the rotated-away ID - // are stale duplicates of the same physical device (its restored - // connections outlived the relaunch that rotated the ID): cancel - // them now instead of leaving ghost links that spray duplicate - // traffic until the inactivity timeout. - self.cancelBoundPeripheralLinks(to: previousPeerID, keeping: linkUUID) - // Retire the rotated-away ID only once its last link is gone; a - // remaining stale link heals the same way or ages out. - guard self.linkBindings.links(to: previousPeerID).isEmpty else { return } - self.messageQueue.async { [weak self] in - self?.retireRotatedPeer(previousPeerID) - } + return } + + // The signature does not authenticate directness (TTL is excluded + // from signing because relays mutate it), so a "verified direct" + // announce can be a replay of another peer's fresh announce with + // its TTL restored. Contain what a forged rebind could do: + // never steal an identity another live link already owns, and + // allow at most one rebind per link per cooldown window so two + // identities can't fight over a link in a replay flip-flop. + guard linkBindings.links(to: peerID).isEmpty else { + SecureLogger.warning("🚫 Refusing link rebind to \(peerID.id.prefix(8))…: identity already owns another live link", category: .security) + return + } + let now = Date() + guard linkAuth.permitRebind( + linkUUID: linkUUID, + now: now, + cooldown: TransportConfig.bleLinkRebindCooldownSeconds + ) else { + SecureLogger.warning("🚫 Refusing link rebind to \(peerID.id.prefix(8))…: rebind cooldown active for this link", category: .security) + return + } + + // A Noise proof belongs to the old physical binding. Never carry + // it across an announce-driven rebind, whose direct TTL is + // replayable; the new owner must complete a fresh handshake. + linkAuth.retireLink(link) + switch link { + case .peripheral(let peripheralUUID): + bindPeripheralLink(peripheralUUID, to: peerID) + case .central(let centralUUID): + linkBindings.bindCentral(centralUUID, to: peerID) + } + // Same engine slot as the rebind: no observer may see the new + // binding while a cached peer-level sender is still considered + // established. + refreshNoiseSessionForVerifiedDirectLink( + link: link, + peerID: peerID + ) + SecureLogger.debug("🔄 Rebinding link after peer-ID rotation: \(previousPeerID.id.prefix(8))… → \(peerID.id.prefix(8))…", category: .session) + refreshLocalTopology() + // The announce that triggered this rebind was upserted as + // disconnected: the registry ran while the link still belonged + // to the previous ID (the ambiguous state BLEAnnounceHandler + // denies the connected shortcut). The rebind has now + // containment-checked the claim and the identity owns a live + // link, so promote it — otherwise a healed rotation leaves a + // live link that reads as disconnected until the next announce. + promoteReboundPeerToConnected(peerID) + // Any other peripheral links still bound to the rotated-away ID + // are stale duplicates of the same physical device (its restored + // connections outlived the relaunch that rotated the ID): cancel + // them now instead of leaving ghost links that spray duplicate + // traffic until the inactivity timeout. + cancelBoundPeripheralLinks(to: previousPeerID, keeping: linkUUID) + // Retire the rotated-away ID only once its last link is gone; a + // remaining stale link heals the same way or ages out. + guard linkBindings.links(to: previousPeerID).isEmpty else { return } + retireRotatedPeer(previousPeerID) } /// After a restore relaunch the same phone can reappear under a fresh @@ -6868,38 +6939,37 @@ extension BLEService { /// link either way. private func retireRedundantPeripheralLinks(_ packet: BitchatPacket, to peerID: PeerID) { let ingressLink = ingressLinks.link(for: packet) - bleQueue.async { [weak self] in - guard let self else { return } - let now = Date() - var ingressPeripheralUUID: String? - if case .peripheral(let uuid) = ingressLink { - ingressPeripheralUUID = uuid - } - guard let keptUUID = BLERedundantLinkPolicy.keptPeripheralUUID( - ingressPeripheralUUID: ingressPeripheralUUID, - mostRecentlyBoundUUID: self.linkBindings.preferredPeripheralUUID(for: peerID), - links: self.peripheralLinkPolicySnapshot(), - peerID: peerID - ) else { return } - - guard self.linkAuth.permitRedundantRetirement( - peerID: peerID, - now: now, - cooldown: TransportConfig.bleLinkRebindCooldownSeconds - ) else { return } - // The survivor becomes the peer's reverse-mapped link so directed - // sends follow the consolidation. - self.bindPeripheralLink(keptUUID, to: peerID) - self.cancelBoundPeripheralLinks(to: peerID, keeping: keptUUID) - self.refreshLocalTopology() + let now = Date() + var ingressPeripheralUUID: String? + if case .peripheral(let uuid) = ingressLink { + ingressPeripheralUUID = uuid } + guard let keptUUID = BLERedundantLinkPolicy.keptPeripheralUUID( + ingressPeripheralUUID: ingressPeripheralUUID, + mostRecentlyBoundUUID: linkBindings.preferredPeripheralUUID(for: peerID), + links: peripheralLinkPolicySnapshot(), + peerID: peerID + ) else { return } + + guard linkAuth.permitRedundantRetirement( + peerID: peerID, + now: now, + cooldown: TransportConfig.bleLinkRebindCooldownSeconds + ) else { return } + // The survivor becomes the peer's reverse-mapped link so directed + // sends follow the consolidation. + bindPeripheralLink(keptUUID, to: peerID) + cancelBoundPeripheralLinks(to: peerID, keeping: keptUUID) + refreshLocalTopology() } /// Cancels our central-role connections whose link is bound to `peerID`, - /// except `keptUUID`. bleQueue only. Each entry is removed from the link - /// store BEFORE cancelling so didDisconnectPeripheral sees no peer - /// binding and skips its peer-disconnect bookkeeping — the peer is still - /// live (on the kept link, or under its rotated identity). + /// except `keptUUID`. Engine-confined: each binding is retired BEFORE + /// the cancel is issued, so didDisconnectPeripheral's identity hop sees + /// no peer binding and skips its peer-disconnect bookkeeping — the peer + /// is still live (on the kept link, or under its rotated identity). + /// Only the physical discard and the CoreBluetooth cancel hop to + /// bleQueue. private func cancelBoundPeripheralLinks(to peerID: PeerID, keeping keptUUID: String?) { let retiring = BLERedundantLinkPolicy.peripheralUUIDsToRetire( links: peripheralLinkPolicySnapshot(), @@ -6907,25 +6977,36 @@ extension BLEService { keeping: keptUUID ?? "" ) for uuid in retiring { - guard let state = linkStateStore.state(forPeripheralID: uuid) else { continue } - tearDownPeripheralLink(uuid) + retirePeripheralLinkIdentity(uuid) SecureLogger.info( "🔗 Retiring redundant link \(uuid.prefix(8))… bound to \(peerID.id.prefix(8))…\(keptUUID.map { " (keeping \($0.prefix(8))…)" } ?? "")", category: .session ) - centralManager?.cancelPeripheralConnection(state.peripheral) + bleQueue.async { [weak self] in + guard let self, + let state = self.linkStateStore.state(forPeripheralID: uuid) else { return } + self.discardPeripheralLinkPhysical(uuid) + self.centralManager?.cancelPeripheralConnection(state.peripheral) + } } } - /// bleQueue only (reads the link store). + /// Engine-confined: physical link rows joined with their engine-owned + /// bindings. private func peripheralLinkPolicySnapshot() -> [BLERedundantLinkPolicy.PeripheralLink] { - linkStateStore.peripheralStates.map { - let uuid = $0.peripheral.identifier.uuidString - return BLERedundantLinkPolicy.PeripheralLink( - uuid: uuid, - peerID: linkBindings.peer(forPeripheralID: uuid), + let physical = readLinkState { store in + store.peripheralStates.map { + (uuid: $0.peripheral.identifier.uuidString, + isConnected: $0.isConnected, + hasCharacteristic: $0.characteristic != nil) + } + } + return physical.map { + BLERedundantLinkPolicy.PeripheralLink( + uuid: $0.uuid, + peerID: linkBindings.peer(forPeripheralID: $0.uuid), isConnected: $0.isConnected, - hasCharacteristic: $0.characteristic != nil + hasCharacteristic: $0.hasCharacteristic ) } } @@ -7008,10 +7089,7 @@ extension BLEService { // residual forged-presence window this leaves is accepted. guard let self else { return false } guard let link = self.ingressLinks.link(for: packet) else { return false } - let boundPeerID: PeerID? = self.readLinkState { _ in - self.linkBindings.boundPeer(for: link) - } - guard let boundPeerID else { return false } + guard let boundPeerID = self.linkBindings.boundPeer(for: link) else { return false } return boundPeerID != peerID }, withRegistryBarrier: { [weak self] body in @@ -7605,6 +7683,14 @@ extension BLEService { #endif private func checkPeerConnectivity() { + // Maintenance ticks on bleQueue; connectivity reconciliation reads + // the engine-owned bindings, so it rides an engine slot. + messageQueue.async { [weak self] in + self?.checkPeerConnectivityOnEngine() + } + } + + private func checkPeerConnectivityOnEngine() { let now = Date() let peerIDsForLinkState: [PeerID] = peerRegistry.peerIDs var cachedLinkStates: [PeerID: BLEPeerLinkPresence] = [:] diff --git a/bitchatTests/BLEServiceCoreTests.swift b/bitchatTests/BLEServiceCoreTests.swift index f0a573cd..13f2bb25 100644 --- a/bitchatTests/BLEServiceCoreTests.swift +++ b/bitchatTests/BLEServiceCoreTests.swift @@ -554,26 +554,19 @@ struct BLEServiceCoreTests { ) let replay = try #require(victim.signPacket(unsigned), "Failed to sign replayed announce") #expect(ble._test_recordIngressIfNew(packet: replay, linkID: attackerLink)) - let rebindGate = VerifiedDirectRebindGate() - ble._test_afterVerifiedDirectRebindEnqueued = rebindGate.pause - defer { - rebindGate.release() - ble._test_afterVerifiedDirectRebindEnqueued = nil - } ble._test_handlePacket(replay, fromPeerID: victimPeerID, preseedPeer: false) - let announcePaused = await TestHelpers.waitUntil( - { rebindGate.hasPaused }, + // The rebind, its Noise-proof retirement, and the ordinary + // reconnect preparation are one engine slot: no observer can see + // the new binding while the victim's stale sending keys are still + // available. Once the binding is visible, the keys must already be + // gone. + let rebound = await TestHelpers.waitUntil( + { ble._test_centralBinding(attackerLink) == victimPeerID }, timeout: TestConstants.longTimeout ) - try #require(announcePaused) - - // Rebind and ordinary reconnect preparation are one bleQueue - // critical section. Once the binding is visible, stale sending keys - // must already be unavailable. - #expect(ble._test_centralBinding(attackerLink) == victimPeerID) + try #require(rebound) #expect(!ble.canDeliverSecurely(to: victimPeerID)) - rebindGate.release() let outbound = OutboundPacketTap() ble._test_onOutboundPacket = { outbound.record($0) } @@ -1469,35 +1462,6 @@ private final class SessionReconcileCounter: @unchecked Sendable { } } -private final class VerifiedDirectRebindGate: @unchecked Sendable { - private let condition = NSCondition() - private var paused = false - private var released = false - - var hasPaused: Bool { - condition.lock() - defer { condition.unlock() } - return paused - } - - func pause() { - condition.lock() - paused = true - condition.broadcast() - while !released { - condition.wait() - } - condition.unlock() - } - - func release() { - condition.lock() - released = true - condition.broadcast() - condition.unlock() - } -} - private final class ReceivePacketHandoffGate: @unchecked Sendable { private let condition = NSCondition() private var paused = false diff --git a/docs/BLE-ARCHITECTURE-V3.md b/docs/BLE-ARCHITECTURE-V3.md index 60400ceb..d72df1bf 100644 --- a/docs/BLE-ARCHITECTURE-V3.md +++ b/docs/BLE-ARCHITECTURE-V3.md @@ -158,6 +158,39 @@ throughput is nowhere near what one serial queue sustains. (it makes no peer decisions); (b) bindings + link-auth migrate to the engine, converting `readLinkState` callers; (c) the delegates shrink to event emission and move behind the port. + + **(a) and (b) are done.** (a) landed as `BLERadioController` + (#1539). (b) landed in two steps: #1540 cohered the loose maps into + `BLELinkAuthState` + `BLELinkBindings` (still bleQueue-owned, + behavior-identical), and the option-B flip then moved ownership to + the engine. Since the flip: + + - `linkAuth`/`linkBindings` are engine-owned behind a DEBUG + `dispatchPrecondition` trap; bleQueue code cannot touch them. + - The receive path is in its sans-I/O shape: bleQueue decodes + frames and hands `(packet, linkID)` up through + `ingestDecodedPacket` (which captures the panic lifecycle at the + handoff); `attributeAndHandlePacket` resolves the sender binding, + admits or rejects the claimed sender, applies raw-announce + binding, and records ingress — all on the engine. Per-link frame + order is preserved end to end (both queues are serial), which + supersedes the old batch-local TOCTOU binding. + - The rotation rebind is one engine slot + (`rebindLinkAfterVerifiedDirectAnnounce`): containment checks, + proof retirement, binding flip, reconnect decision, and + rotated-identity retirement, with only CoreBluetooth cancels + hopping to bleQueue. + - Authenticated-send eligibility (`notifyOrEnqueueIfAccepted`, + `writeOrEnqueueIfAccepted`) is checked on the engine — serialized + against rebinds by construction — and only the physical admission + (updateValue / write / backpressure queues) runs on bleQueue. + - Teardown splits: bleQueue delegates do physical work inline + (`discardPeripheralLinkPhysical`) and queue the identity half + (`retirePeripheralLinkIdentity`, binding survivor repair) to the + engine. A binding can briefly outlive its physical link; queries + that need liveness join against the physical store via + `readLinkState` (the engine→bleQueue sync direction), and the + queued retirement converges the two. 2. **Sans-I/O engine core + simulator.** Make the engine formally `handle(event) -> [Effect]`, feed it from a `SimulatedLinkLayer`, and move the multi-node E2E suite onto deterministic simulation (no