From ffa0d7aa4f96f4273b99369fc7ae3851258debee Mon Sep 17 00:00:00 2001 From: jack <212554440+jackjackbits@users.noreply.github.com> Date: Sun, 31 May 2026 14:24:09 +0200 Subject: [PATCH] [codex] Extract BLE link state store (#1310) * Refactor BLE transport event handling * Make image output paths unique * Keep queued Nostr read receipts alive * Refine BLE ingress fanout * Rediscover BLE service after invalidation * Extract BLE notification retry buffer * Extract BLE inbound write buffer * Extract BLE fragment assembly buffer * Tidy secure log handling from device run * Extract BLE outbound fragment scheduler * Harden app CI media tests * Redact BLE message content from logs * Extract BLE Noise session queues * Fix BLE read receipt UI updates * Allow self-authored RSR ingress replies * Harden read receipt queue test timing * Extract BLE outbound policy and incoming file storage * Avoid duplicate BLE link snapshots during send * Canonicalize Nostr relay URLs * Extract BLE link state store * Extract BLE connection scheduler --------- Co-authored-by: jack --- .../Services/BLE/BLEConnectionScheduler.swift | 271 +++++++++ bitchat/Services/BLE/BLELinkStateStore.swift | 196 +++++++ bitchat/Services/BLE/BLEService.swift | 534 ++++++------------ .../BLEConnectionSchedulerTests.swift | 156 +++++ .../Services/BLELinkStateStoreTests.swift | 46 ++ 5 files changed, 835 insertions(+), 368 deletions(-) create mode 100644 bitchat/Services/BLE/BLEConnectionScheduler.swift create mode 100644 bitchat/Services/BLE/BLELinkStateStore.swift create mode 100644 bitchatTests/Services/BLEConnectionSchedulerTests.swift create mode 100644 bitchatTests/Services/BLELinkStateStoreTests.swift diff --git a/bitchat/Services/BLE/BLEConnectionScheduler.swift b/bitchat/Services/BLE/BLEConnectionScheduler.swift new file mode 100644 index 00000000..a799fc8e --- /dev/null +++ b/bitchat/Services/BLE/BLEConnectionScheduler.swift @@ -0,0 +1,271 @@ +import Foundation + +struct BLEConnectionCandidate { + let peripheral: Peripheral + let peripheralID: String + let rssi: Int + let name: String + let isConnectable: Bool + let discoveredAt: Date +} + +struct BLEExistingConnectionState { + let isConnecting: Bool + let isConnected: Bool + let lastConnectionAttempt: Date? +} + +enum BLEPeripheralConnectionState { + case disconnected + case connecting + case connected +} + +enum BLEDiscoveryDecision: Equatable { + case ignore + case queued + case scheduleRetry(after: TimeInterval) + case cancelStaleConnection + case connectNow +} + +enum BLEConnectionQueueDecision { + case none + case retryAfter(TimeInterval) + case connect(BLEConnectionCandidate) +} + +final class BLEConnectionScheduler { + private let maxCentralLinks: Int + private let connectRateLimitInterval: TimeInterval + private let candidateCap: Int + private let weakLinkCooldownSeconds: TimeInterval + private let weakLinkRSSICutoff: Int + private let recentTimeoutWindowSeconds: TimeInterval + private let recentTimeoutCountThreshold: Int + + private var lastGlobalConnectAttempt: Date = .distantPast + private var candidates: [BLEConnectionCandidate] = [] + private var failureCounts: [String: Int] = [:] + private var recentConnectTimeouts: [String: Date] = [:] + private var lastIsolatedAt: Date? + + private let initialDynamicRSSIThreshold: Int + private(set) var dynamicRSSIThreshold: Int + + var candidateCount: Int { + candidates.count + } + + init( + maxCentralLinks: Int = TransportConfig.bleMaxCentralLinks, + connectRateLimitInterval: TimeInterval = TransportConfig.bleConnectRateLimitInterval, + candidateCap: Int = TransportConfig.bleConnectionCandidatesMax, + weakLinkCooldownSeconds: TimeInterval = TransportConfig.bleWeakLinkCooldownSeconds, + weakLinkRSSICutoff: Int = TransportConfig.bleWeakLinkRSSICutoff, + recentTimeoutWindowSeconds: TimeInterval = TransportConfig.bleRecentTimeoutWindowSeconds, + recentTimeoutCountThreshold: Int = TransportConfig.bleRecentTimeoutCountThreshold, + dynamicRSSIThreshold: Int = TransportConfig.bleDynamicRSSIThresholdDefault + ) { + self.maxCentralLinks = maxCentralLinks + self.connectRateLimitInterval = connectRateLimitInterval + self.candidateCap = candidateCap + self.weakLinkCooldownSeconds = weakLinkCooldownSeconds + self.weakLinkRSSICutoff = weakLinkRSSICutoff + self.recentTimeoutWindowSeconds = recentTimeoutWindowSeconds + self.recentTimeoutCountThreshold = recentTimeoutCountThreshold + self.initialDynamicRSSIThreshold = dynamicRSSIThreshold + self.dynamicRSSIThreshold = dynamicRSSIThreshold + } + + func handleDiscovery( + _ candidate: BLEConnectionCandidate, + connectedOrConnectingCount: Int, + existingState: BLEExistingConnectionState?, + peripheralState: BLEPeripheralConnectionState, + now: Date + ) -> BLEDiscoveryDecision { + guard candidate.isConnectable else { return .ignore } + + if candidate.rssi <= dynamicRSSIThreshold { + enqueue(candidate) + return .queued + } + + if connectedOrConnectingCount >= maxCentralLinks { + enqueue(candidate) + return .queued + } + + if let retryDelay = rateLimitRetryDelay(now: now) { + enqueue(candidate) + return .scheduleRetry(after: retryDelay) + } + + if let existingState { + if existingState.isConnected || existingState.isConnecting { + return .ignore + } + + if let lastAttempt = existingState.lastConnectionAttempt, + now.timeIntervalSince(lastAttempt) < 2.0 { + return .ignore + } + } + + if let lastTimeout = recentConnectTimeouts[candidate.peripheralID], + now.timeIntervalSince(lastTimeout) < 15 { + return .ignore + } + + switch peripheralState { + case .disconnected: + return .connectNow + case .connecting, .connected: + return .cancelStaleConnection + } + } + + func enqueue(_ candidate: BLEConnectionCandidate) { + if let existingIndex = candidates.firstIndex(where: { $0.peripheralID == candidate.peripheralID }) { + candidates[existingIndex] = candidate + } else { + candidates.append(candidate) + } + + candidates.sort { + if $0.rssi != $1.rssi { return $0.rssi > $1.rssi } + return $0.discoveredAt < $1.discoveredAt + } + if candidates.count > candidateCap { + candidates.removeLast(candidates.count - candidateCap) + } + } + + func nextCandidate( + connectedOrConnectingCount: Int, + isAlreadyConnectingOrConnected: (String) -> Bool, + now: Date + ) -> BLEConnectionQueueDecision { + guard connectedOrConnectingCount < maxCentralLinks else { return .none } + + if let retryDelay = rateLimitRetryDelay(now: now) { + return .retryAfter(retryDelay) + } + + while !candidates.isEmpty { + candidates.sort { score($0, now: now) > score($1, now: now) } + let candidate = candidates.removeFirst() + guard candidate.isConnectable else { continue } + + if let delay = weakLinkRetryDelay(for: candidate, now: now) { + enqueue(candidate) + return .retryAfter(delay) + } + + if isAlreadyConnectingOrConnected(candidate.peripheralID) { + continue + } + + return .connect(candidate) + } + + return .none + } + + func recordConnectionAttempt(at now: Date) { + lastGlobalConnectAttempt = now + } + + func recordConnectionSuccess(peripheralID: String) { + failureCounts[peripheralID] = 0 + recentConnectTimeouts.removeValue(forKey: peripheralID) + } + + func recordConnectionFailure(peripheralID: String) { + failureCounts[peripheralID, default: 0] += 1 + } + + func recordDisconnectError(peripheralID: String, at now: Date) { + recentConnectTimeouts[peripheralID] = now + } + + func recordConnectionTimeout(peripheralID: String, at now: Date) { + recentConnectTimeouts[peripheralID] = now + recordConnectionFailure(peripheralID: peripheralID) + } + + func pruneConnectionTimeouts(before cutoff: Date) { + recentConnectTimeouts = recentConnectTimeouts.filter { $0.value >= cutoff } + } + + func reset() { + lastGlobalConnectAttempt = .distantPast + candidates.removeAll() + failureCounts.removeAll() + recentConnectTimeouts.removeAll() + lastIsolatedAt = nil + dynamicRSSIThreshold = initialDynamicRSSIThreshold + } + + @discardableResult + func updateRSSIThreshold( + connectedCount: Int, + connectedOrConnectingLinkCount: Int, + now: Date + ) -> Int { + if connectedCount == 0 { + if lastIsolatedAt == nil { lastIsolatedAt = now } + let isolatedAt = lastIsolatedAt ?? now + let elapsed = now.timeIntervalSince(isolatedAt) + dynamicRSSIThreshold = elapsed > TransportConfig.bleIsolationRelaxThresholdSeconds + ? TransportConfig.bleRSSIIsolatedRelaxed + : TransportConfig.bleRSSIIsolatedBase + return dynamicRSSIThreshold + } + + lastIsolatedAt = nil + var threshold = TransportConfig.bleDynamicRSSIThresholdDefault + if connectedOrConnectingLinkCount >= maxCentralLinks || candidates.count >= candidateCap { + threshold = TransportConfig.bleRSSIConnectedThreshold + } + + let recentTimeouts = recentConnectTimeouts.filter { + now.timeIntervalSince($0.value) < recentTimeoutWindowSeconds + }.count + if recentTimeouts >= recentTimeoutCountThreshold { + threshold = max(threshold, TransportConfig.bleRSSIHighTimeoutThreshold) + } + + dynamicRSSIThreshold = threshold + return threshold + } + + private func rateLimitRetryDelay(now: Date) -> TimeInterval? { + let elapsed = now.timeIntervalSince(lastGlobalConnectAttempt) + guard elapsed < connectRateLimitInterval else { return nil } + return connectRateLimitInterval - elapsed + 0.05 + } + + private func weakLinkRetryDelay( + for candidate: BLEConnectionCandidate, + now: Date + ) -> TimeInterval? { + guard let lastTimeout = recentConnectTimeouts[candidate.peripheralID] else { return nil } + let elapsed = now.timeIntervalSince(lastTimeout) + guard elapsed < weakLinkCooldownSeconds && candidate.rssi <= weakLinkRSSICutoff else { return nil } + let remaining = weakLinkCooldownSeconds - elapsed + return min(max(2.0, remaining), 15.0) + } + + private func score(_ candidate: BLEConnectionCandidate, now: Date) -> Int { + let failures = failureCounts[candidate.peripheralID] ?? 0 + let penalty = min(20, 1 << min(4, failures)) + let timeoutBias = recentConnectTimeouts[candidate.peripheralID].map { + now.timeIntervalSince($0) < 60 ? 10 : 0 + } ?? 0 + let base = (candidate.isConnectable ? 1000 : 0) + (candidate.rssi + 100) * 2 + let recency = -Int(now.timeIntervalSince(candidate.discoveredAt) * 10) + return base + recency - penalty - timeoutBias + } +} diff --git a/bitchat/Services/BLE/BLELinkStateStore.swift b/bitchat/Services/BLE/BLELinkStateStore.swift new file mode 100644 index 00000000..b46edaec --- /dev/null +++ b/bitchat/Services/BLE/BLELinkStateStore.swift @@ -0,0 +1,196 @@ +import BitFoundation +import CoreBluetooth +import Foundation + +struct BLEPeripheralLinkState { + let peripheral: CBPeripheral + var characteristic: CBCharacteristic? + var peerID: PeerID? + var isConnecting: Bool + var isConnected: Bool + var lastConnectionAttempt: Date? + var assembler: NotificationStreamAssembler +} + +struct BLEDirectLinkState: Equatable { + let hasPeripheral: Bool + let hasCentral: Bool +} + +struct BLESubscribedCentralSnapshot { + let centrals: [CBCentral] + let peerIDsByCentralUUID: [String: PeerID] + + func central(for peerID: PeerID) -> CBCentral? { + centrals.first { peerIDsByCentralUUID[$0.identifier.uuidString] == peerID } + } +} + +final class BLELinkStateStore { + private(set) var peripherals: [String: BLEPeripheralLinkState] = [:] + private(set) var peerToPeripheralUUID: [PeerID: String] = [:] + private(set) var subscribedCentrals: [CBCentral] = [] + private(set) var centralToPeerID: [String: PeerID] = [:] + + var peripheralStates: [BLEPeripheralLinkState] { + Array(peripherals.values) + } + + var subscribedCentralSnapshot: BLESubscribedCentralSnapshot { + BLESubscribedCentralSnapshot( + centrals: subscribedCentrals, + peerIDsByCentralUUID: centralToPeerID + ) + } + + var subscribedCentralCount: Int { + subscribedCentrals.count + } + + var connectedOrConnectingPeripheralCount: Int { + peripherals.values.filter { $0.isConnected || $0.isConnecting }.count + } + + func state(forPeripheralID peripheralID: String) -> BLEPeripheralLinkState? { + peripherals[peripheralID] + } + + func setPeripheralState(_ state: BLEPeripheralLinkState, for peripheralID: String) { + peripherals[peripheralID] = state + } + + @discardableResult + func updatePeripheral( + _ peripheralID: String, + _ update: (inout BLEPeripheralLinkState) -> Void + ) -> BLEPeripheralLinkState? { + guard var state = peripherals[peripheralID] else { return nil } + update(&state) + peripherals[peripheralID] = state + return state + } + + func beginConnecting(to peripheral: CBPeripheral, at date: Date) { + setPeripheralState( + BLEPeripheralLinkState( + peripheral: peripheral, + characteristic: nil, + peerID: nil, + isConnecting: true, + isConnected: false, + lastConnectionAttempt: date, + assembler: NotificationStreamAssembler() + ), + for: peripheral.identifier.uuidString + ) + } + + func markConnected(_ peripheral: CBPeripheral) { + let peripheralID = peripheral.identifier.uuidString + if updatePeripheral(peripheralID, { + $0.isConnecting = false + $0.isConnected = true + }) == nil { + setPeripheralState( + BLEPeripheralLinkState( + peripheral: peripheral, + characteristic: nil, + peerID: nil, + isConnecting: false, + isConnected: true, + lastConnectionAttempt: nil, + assembler: NotificationStreamAssembler() + ), + for: peripheralID + ) + } + } + + func updateCharacteristic(_ characteristic: CBCharacteristic, forPeripheralID peripheralID: String) { + updatePeripheral(peripheralID) { + $0.characteristic = characteristic + } + } + + func directPeripheralState(for peerID: PeerID) -> BLEPeripheralLinkState? { + peerToPeripheralUUID[peerID].flatMap { peripherals[$0] } + } + + func directLinkState(for peerID: PeerID) -> BLEDirectLinkState { + let peripheralUUID = peerToPeripheralUUID[peerID] + let hasPeripheral = peripheralUUID.flatMap { peripherals[$0]?.isConnected } ?? false + let hasCentral = centralToPeerID.values.contains(peerID) + return BLEDirectLinkState(hasPeripheral: hasPeripheral, hasCentral: hasCentral) + } + + func links(to peerID: PeerID?) -> Set { + guard let peerID else { return [] } + + var links: Set = [] + if let peripheralUUID = peerToPeripheralUUID[peerID] { + links.insert(.peripheral(peripheralUUID)) + } + for (centralUUID, mappedPeerID) in centralToPeerID where mappedPeerID == peerID { + links.insert(.central(centralUUID)) + } + return links + } + + func peerID(forPeripheralID peripheralID: String) -> PeerID? { + peripherals[peripheralID]?.peerID + } + + func peerID(forCentralUUID centralUUID: String) -> PeerID? { + centralToPeerID[centralUUID] + } + + func addSubscribedCentral(_ central: CBCentral) { + guard !subscribedCentrals.contains(central) else { return } + subscribedCentrals.append(central) + } + + func removeSubscribedCentral(_ central: CBCentral) -> PeerID? { + let centralUUID = central.identifier.uuidString + subscribedCentrals.removeAll { $0.identifier == central.identifier } + return centralToPeerID.removeValue(forKey: centralUUID) + } + + func bindCentral(_ centralUUID: String, to peerID: PeerID) { + centralToPeerID[centralUUID] = peerID + } + + func bindPeripheral(_ peripheralUUID: String, to peerID: PeerID) { + if updatePeripheral(peripheralUUID, { $0.peerID = peerID }) != nil { + peerToPeripheralUUID[peerID] = peripheralUUID + } + } + + func removePeripheral(_ peripheralID: String) -> PeerID? { + let peerID = peripherals.removeValue(forKey: peripheralID)?.peerID + if let peerID { + peerToPeripheralUUID.removeValue(forKey: peerID) + } + return peerID + } + + func clearPeripherals() -> [PeerID] { + let peerIDs = peripherals.compactMap { $0.value.peerID } + peripherals.removeAll() + peerToPeripheralUUID.removeAll() + return peerIDs + } + + func clearCentrals() -> [PeerID] { + let peerIDs = Array(centralToPeerID.values) + subscribedCentrals.removeAll() + centralToPeerID.removeAll() + return peerIDs + } + + func clearAll() { + peripherals.removeAll() + peerToPeripheralUUID.removeAll() + subscribedCentrals.removeAll() + centralToPeerID.removeAll() + } +} diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index c10bcf02..5460a63d 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -34,23 +34,9 @@ final class BLEService: NSObject { private let highDegreeThreshold = TransportConfig.bleHighDegreeThreshold // for adaptive TTL/probabilistic relays // MARK: - Core State (5 Essential Collections) - - // 1. Consolidated Peripheral Tracking - private struct PeripheralState { - let peripheral: CBPeripheral - var characteristic: CBCharacteristic? - var peerID: PeerID? - var isConnecting: Bool = false - var isConnected: Bool = false - var lastConnectionAttempt: Date? = nil - var assembler = NotificationStreamAssembler() - } - private var peripherals: [String: PeripheralState] = [:] // UUID -> PeripheralState - private var peerToPeripheralUUID: [PeerID: String] = [:] // PeerID -> Peripheral UUID - - // 2. BLE Centrals (when acting as peripheral) - private var subscribedCentrals: [CBCentral] = [] - private var centralToPeerID: [String: PeerID] = [:] // Central UUID -> Peer ID mapping + + // 1. Consolidated BLE link tracking for both central and peripheral roles. + private var linkStateStore = BLELinkStateStore() // BCH-01-004: Rate-limiting for subscription-triggered announces // Tracks subscription attempts per central to prevent enumeration attacks @@ -85,8 +71,6 @@ final class BLEService: NSObject { private var fragmentAssemblyBuffer = BLEFragmentAssemblyBuffer() private var outboundFragmentTransfers = BLEOutboundFragmentTransferScheduler() private let incomingFileStore = BLEIncomingFileStore() - // Backoff for peripherals that recently timed out connecting - private var recentConnectTimeouts: [String: Date] = [:] // Peripheral UUID -> last timeout // Simple announce throttling private var lastAnnounceSent = Date.distantPast @@ -156,20 +140,7 @@ final class BLEService: NSObject { private var maintenanceCounter = 0 // Track maintenance cycles // MARK: - Connection budget & scheduling (central role) - private let maxCentralLinks = TransportConfig.bleMaxCentralLinks - private let connectRateLimitInterval: TimeInterval = TransportConfig.bleConnectRateLimitInterval - private var lastGlobalConnectAttempt: Date = .distantPast - private struct ConnectionCandidate { - let peripheral: CBPeripheral - let rssi: Int - let name: String - let isConnectable: Bool - let discoveredAt: Date - } - private var connectionCandidates: [ConnectionCandidate] = [] - private var failureCounts: [String: Int] = [:] // Peripheral UUID -> failures - private var lastIsolatedAt: Date? = nil - private var dynamicRSSIThreshold: Int = TransportConfig.bleDynamicRSSIThresholdDefault + private var connectionScheduler = BLEConnectionScheduler() // MARK: - Adaptive scanning duty-cycle private var scanDutyTimer: DispatchSourceTimer? @@ -348,7 +319,7 @@ final class BLEService: NSObject { bleQueue.sync { pendingWriteBuffers.removeAll() - recentConnectTimeouts.removeAll() + connectionScheduler.reset() } recentDisconnectNotifies.removeAll() @@ -506,7 +477,7 @@ final class BLEService: NSObject { // Snapshot BLE state under bleQueue to avoid races with delegate callbacks let (peripheralStates, centralsCount, char) = bleQueue.sync { - (Array(peripherals.values), subscribedCentrals.count, characteristic) + (linkStateStore.peripheralStates, linkStateStore.subscribedCentralCount, characteristic) } // Send to peripherals we're connected to as central @@ -543,7 +514,7 @@ final class BLEService: NSObject { peripheralManager?.stopAdvertising() // Disconnect all peripherals (synchronized access) - let peripheralsToDisconnect = bleQueue.sync { Array(peripherals.values) } + let peripheralsToDisconnect = bleQueue.sync { linkStateStore.peripheralStates } for state in peripheralsToDisconnect { centralManager?.cancelPeripheralConnection(state.peripheral) } @@ -573,10 +544,8 @@ final class BLEService: NSObject { // Clear peripheral references (synchronized access to avoid races with BLE callbacks) bleQueue.sync { - peripherals.removeAll() - peerToPeripheralUUID.removeAll() - subscribedCentrals.removeAll() - centralToPeerID.removeAll() + linkStateStore.clearAll() + connectionScheduler.reset() centralSubscriptionRateLimits.removeAll() } meshTopology.reset() @@ -888,16 +857,8 @@ final class BLEService: NSObject { let outboundPriority = BLEOutboundPacketPolicy.priority(for: packet, data: data) // Per-link limits for the specific peer - let directPeripheralState: PeripheralState? = { - if DispatchQueue.getSpecific(key: bleQueueKey) != nil { - return peerToPeripheralUUID[recipientPeerID].flatMap { peripherals[$0] } - } - return bleQueue.sync { - peerToPeripheralUUID[recipientPeerID].flatMap { peripherals[$0] } - } - }() - let (centrals, mapping) = snapshotSubscribedCentrals() - let recipientCentral = centrals.first { mapping[$0.identifier.uuidString] == recipientPeerID } + let directPeripheralState = snapshotDirectPeripheralState(for: recipientPeerID) + let recipientCentral = snapshotSubscribedCentrals().central(for: recipientPeerID) if let peripheralMaxLen = directPeripheralState?.peripheral.maximumWriteValueLength(for: .withoutResponse), data.count > peripheralMaxLen { @@ -987,7 +948,7 @@ final class BLEService: NSObject { let m = s.peripheral.maximumWriteValueLength(for: .withoutResponse) minCentralWriteLen = minCentralWriteLen.map { min($0, m) } ?? m } - let subscribedCentrals = characteristic == nil ? [] : snapshotSubscribedCentrals().0 + let subscribedCentrals = characteristic == nil ? [] : snapshotSubscribedCentrals().centrals let minNotifyLen = subscribedCentrals.map { $0.maximumUpdateValueLength }.min() // Avoid re-fragmenting fragment packets @@ -1378,14 +1339,14 @@ extension BLEService: CBCentralManagerDelegate { for peripheral in restoredPeripherals { let identifier = peripheral.identifier.uuidString peripheral.delegate = self - let existing = peripherals[identifier] + let existing = linkStateStore.state(forPeripheralID: identifier) let assembler = existing?.assembler ?? NotificationStreamAssembler() let characteristic = existing?.characteristic let peerID = existing?.peerID let wasConnecting = existing?.isConnecting ?? false let wasConnected = existing?.isConnected ?? false - let restoredState = PeripheralState( + let restoredState = BLEPeripheralLinkState( peripheral: peripheral, characteristic: characteristic, peerID: peerID, @@ -1394,7 +1355,7 @@ extension BLEService: CBCentralManagerDelegate { lastConnectionAttempt: existing?.lastConnectionAttempt, assembler: assembler ) - peripherals[identifier] = restoredState + linkStateStore.setPeripheralState(restoredState, for: identifier) } captureBluetoothStatus(context: "central-restore") @@ -1418,12 +1379,12 @@ extension BLEService: CBCentralManagerDelegate { SecureLogger.info("📴 Bluetooth powered off - cleaning up central state", category: .session) central.stopScan() // Mark all peripheral connections as disconnected (they are now invalid) - let peerIDs: [PeerID] = peripherals.compactMap { $0.value.peerID } - for state in peripherals.values { + let peripheralStates = linkStateStore.peripheralStates + let peerIDs: [PeerID] = peripheralStates.compactMap(\.peerID) + for state in peripheralStates { central.cancelPeripheralConnection(state.peripheral) } - peripherals.removeAll() - peerToPeripheralUUID.removeAll() + _ = linkStateStore.clearPeripherals() // Notify UI of disconnections for peerID in peerIDs { notifyUI { [weak self] in @@ -1435,8 +1396,7 @@ extension BLEService: CBCentralManagerDelegate { // User denied Bluetooth permission SecureLogger.warning("🚫 Bluetooth unauthorized - user denied permission", category: .session) central.stopScan() - peripherals.removeAll() - peerToPeripheralUUID.removeAll() + _ = linkStateStore.clearPeripherals() case .unsupported: // Device doesn't support BLE @@ -1481,160 +1441,47 @@ extension BLEService: CBCentralManagerDelegate { let advertisedName = advertisementData[CBAdvertisementDataLocalNameKey] as? String ?? (peripheralID.prefix(6) + "…") let isConnectable = (advertisementData[CBAdvertisementDataIsConnectable] as? NSNumber)?.boolValue ?? true let rssiValue = RSSI.intValue - - // Skip if peripheral is not connectable (per advertisement data) - guard isConnectable else { return } - // Skip immediate connect if signal too weak for current conditions; enqueue instead - if rssiValue <= dynamicRSSIThreshold { - connectionCandidates.append(ConnectionCandidate(peripheral: peripheral, rssi: rssiValue, name: String(advertisedName), isConnectable: isConnectable, discoveredAt: Date())) - // Keep list tidy - connectionCandidates.sort { (a, b) in - if a.rssi != b.rssi { return a.rssi > b.rssi } - return a.discoveredAt < b.discoveredAt - } - if connectionCandidates.count > TransportConfig.bleConnectionCandidatesMax { - connectionCandidates.removeLast(connectionCandidates.count - TransportConfig.bleConnectionCandidatesMax) - } - return - } - - // Budget: limit simultaneous central links (connected + connecting) - let currentCentralLinks = peripherals.values.filter { $0.isConnected || $0.isConnecting }.count - if currentCentralLinks >= maxCentralLinks { - // Enqueue as candidate; we'll attempt later as slots open - connectionCandidates.append(ConnectionCandidate(peripheral: peripheral, rssi: rssiValue, name: String(advertisedName), isConnectable: isConnectable, discoveredAt: Date())) - // Keep candidate list tidy: prefer stronger RSSI, then recency; cap list - connectionCandidates.sort { (a, b) in - if a.rssi != b.rssi { return a.rssi > b.rssi } - return a.discoveredAt < b.discoveredAt - } - if connectionCandidates.count > TransportConfig.bleConnectionCandidatesMax { - connectionCandidates.removeLast(connectionCandidates.count - TransportConfig.bleConnectionCandidatesMax) - } - return - } + let candidate = BLEConnectionCandidate( + peripheral: peripheral, + peripheralID: peripheralID, + rssi: rssiValue, + name: String(advertisedName), + isConnectable: isConnectable, + discoveredAt: Date() + ) + let existingState = linkStateStore.state(forPeripheralID: peripheralID).map(BLEExistingConnectionState.init) - // Rate limit global connect attempts - let sinceLast = Date().timeIntervalSince(lastGlobalConnectAttempt) - if sinceLast < connectRateLimitInterval { - connectionCandidates.append(ConnectionCandidate(peripheral: peripheral, rssi: rssiValue, name: String(advertisedName), isConnectable: isConnectable, discoveredAt: Date())) - connectionCandidates.sort { (a, b) in - if a.rssi != b.rssi { return a.rssi > b.rssi } - return a.discoveredAt < b.discoveredAt - } - // Schedule a deferred attempt after rate-limit interval - let delay = connectRateLimitInterval - sinceLast + 0.05 + switch connectionScheduler.handleDiscovery( + candidate, + connectedOrConnectingCount: linkStateStore.connectedOrConnectingPeripheralCount, + existingState: existingState, + peripheralState: peripheral.state.connectionSchedulerState, + now: candidate.discoveredAt + ) { + case .ignore, .queued: + return + case .scheduleRetry(let delay): bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in self?.tryConnectFromQueue() } return - } - - // Check if we already have this peripheral - if let state = peripherals[peripheralID] { - if state.isConnected || state.isConnecting { - return // Already connected or connecting - } - - // Add backoff for reconnection attempts - if let lastAttempt = state.lastConnectionAttempt { - let timeSinceLastAttempt = Date().timeIntervalSince(lastAttempt) - if timeSinceLastAttempt < 2.0 { - return // Wait at least 2 seconds between connection attempts - } - } - } - - // Backoff if this peripheral recently timed out connection within the last 15 seconds - if let lastTimeout = recentConnectTimeouts[peripheralID], Date().timeIntervalSince(lastTimeout) < 15 { - return - } - - // Check peripheral state - but cancel if stale - if peripheral.state == .connecting || peripheral.state == .connected { - // iOS might have stale state - force disconnect and retry + case .cancelStaleConnection: central.cancelPeripheralConnection(peripheral) - // Will retry on next discovery return - } - - // Only log when we're actually attempting connection - // Discovered BLE peripheral - - // Store the peripheral and mark as connecting - peripherals[peripheralID] = PeripheralState( - peripheral: peripheral, - characteristic: nil, - peerID: nil, - isConnecting: true, - isConnected: false, - lastConnectionAttempt: Date(), - assembler: NotificationStreamAssembler() - ) - peripheral.delegate = self - - // Connect to the peripheral with options for faster connection - SecureLogger.debug("📱 Connect: \(advertisedName) [RSSI:\(rssiValue)]", category: .session) - - // Use connection options for faster reconnection - let options: [String: Any] = [ - CBConnectPeripheralOptionNotifyOnConnectionKey: true, - CBConnectPeripheralOptionNotifyOnDisconnectionKey: true, - CBConnectPeripheralOptionNotifyOnNotificationKey: true - ] - central.connect(peripheral, options: options) - lastGlobalConnectAttempt = Date() - - // Set a timeout for the connection attempt (slightly longer for reliability) - // Use BLE queue to mutate BLE-related state consistently - bleQueue.asyncAfter(deadline: .now() + TransportConfig.bleConnectTimeoutSeconds) { [weak self] in - guard let self = self, - let state = self.peripherals[peripheralID], - state.isConnecting && !state.isConnected else { return } - - // Double-check actual CBPeripheral state to avoid canceling a just-connected peripheral - // This prevents a race where connection completes just as timeout fires - guard peripheral.state != .connected else { - SecureLogger.debug("⏱️ Timeout fired but peripheral already connected: \(advertisedName)", category: .session) - return - } - - // Connection timed out - cancel it - SecureLogger.debug("⏱️ Timeout: \(advertisedName)", category: .session) - central.cancelPeripheralConnection(peripheral) - self.peripherals[peripheralID] = nil - self.recentConnectTimeouts[peripheralID] = Date() - self.failureCounts[peripheralID, default: 0] += 1 - // Try next candidate if any - self.tryConnectFromQueue() + case .connectNow: + beginCentralConnection(candidate, using: central, logPrefix: "📱 Connect") } } -func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeripheral) { + func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeripheral) { let peripheralID = peripheral.identifier.uuidString // Update state to connected - if var state = peripherals[peripheralID] { - state.isConnecting = false - state.isConnected = true - peripherals[peripheralID] = state - } else { - // Create new state if not found - peripherals[peripheralID] = PeripheralState( - peripheral: peripheral, - characteristic: nil, - peerID: nil, - isConnecting: false, - isConnected: true, - lastConnectionAttempt: nil, - assembler: NotificationStreamAssembler() - ) - } + linkStateStore.markConnected(peripheral) // Reset backoff state on success - failureCounts[peripheralID] = 0 - recentConnectTimeouts.removeValue(forKey: peripheralID) + connectionScheduler.recordConnectionSuccess(peripheralID: peripheralID) SecureLogger.debug("✅ Connected: \(peripheral.name ?? "Unknown") [\(peripheralID)]", category: .session) @@ -1646,22 +1493,18 @@ func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeriph let peripheralID = peripheral.identifier.uuidString // Find the peer ID if we have it - let peerID = peripherals[peripheralID]?.peerID + let peerID = linkStateStore.peerID(forPeripheralID: peripheralID) SecureLogger.debug("📱 Disconnect: \(peerID?.id ?? peripheralID)\(error != nil ? " (\(error!.localizedDescription))" : "")", category: .session) // If disconnect carried an error (often timeout), apply short backoff to avoid thrash if error != nil { - recentConnectTimeouts[peripheralID] = Date() + connectionScheduler.recordDisconnectError(peripheralID: peripheralID, at: Date()) } - // Clean up references - peripherals.removeValue(forKey: peripheralID) - - // Clean up peer mappings + // Clean up references and peer mappings + _ = linkStateStore.removePeripheral(peripheralID) if let peerID { - peerToPeripheralUUID.removeValue(forKey: peerID) - // Do not remove peer; mark as not connected but retain for reachability collectionsQueue.sync(flags: .barrier) { if var info = peers[peerID] { @@ -1703,74 +1546,72 @@ func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeriph let peripheralID = peripheral.identifier.uuidString // Clean up the references - peripherals.removeValue(forKey: peripheralID) + _ = linkStateStore.removePeripheral(peripheralID) SecureLogger.error("❌ Failed to connect to peripheral: \(peripheral.name ?? "Unknown") [\(peripheralID)] - Error: \(error?.localizedDescription ?? "Unknown")", category: .session) - failureCounts[peripheralID, default: 0] += 1 + connectionScheduler.recordConnectionFailure(peripheralID: peripheralID) // Try next candidate bleQueue.async { [weak self] in self?.tryConnectFromQueue() } } } // MARK: - Connection scheduling helpers +private extension BLEExistingConnectionState { + init(_ state: BLEPeripheralLinkState) { + self.init( + isConnecting: state.isConnecting, + isConnected: state.isConnected, + lastConnectionAttempt: state.lastConnectionAttempt + ) + } +} + +private extension CBPeripheralState { + var connectionSchedulerState: BLEPeripheralConnectionState { + switch self { + case .connected: + return .connected + case .connecting: + return .connecting + case .disconnected, .disconnecting: + return .disconnected + @unknown default: + return .disconnected + } + } +} + extension BLEService { private func tryConnectFromQueue() { guard let central = centralManager, central.state == .poweredOn else { return } - // Check budget and rate limit - let current = peripherals.values.filter { $0.isConnected || $0.isConnecting }.count - guard current < maxCentralLinks else { return } - let delta = Date().timeIntervalSince(lastGlobalConnectAttempt) - guard delta >= connectRateLimitInterval else { - let delay = connectRateLimitInterval - delta + 0.05 - bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in self?.tryConnectFromQueue() } - return - } - // Pull best candidate by composite score - guard !connectionCandidates.isEmpty else { return } - // compute score: connectable> RSSI > recency, with backoff penalty - func score(_ c: ConnectionCandidate) -> Int { - let uuid = c.peripheral.identifier.uuidString - // Penalty if recently timed out (exponential) - let fails = failureCounts[uuid] ?? 0 - let penalty = min(20, (1 << min(4, fails))) // 1,2,4,8,16 cap 16-20 - let timeoutRecent = recentConnectTimeouts[uuid] - let timeoutBias = (timeoutRecent != nil && Date().timeIntervalSince(timeoutRecent!) < 60) ? 10 : 0 - let base = (c.isConnectable ? 1000 : 0) + (c.rssi + 100) * 2 - let rec = -Int(Date().timeIntervalSince(c.discoveredAt) * 10) - return base + rec - penalty - timeoutBias - } - connectionCandidates.sort { score($0) > score($1) } - let candidate = connectionCandidates.removeFirst() - guard candidate.isConnectable else { return } - let peripheral = candidate.peripheral - let peripheralID = peripheral.identifier.uuidString - // Weak-link cooldown: if we recently timed out and RSSI is very weak, delay retries - if let lastTO = recentConnectTimeouts[peripheralID] { - let elapsed = Date().timeIntervalSince(lastTO) - if elapsed < TransportConfig.bleWeakLinkCooldownSeconds && candidate.rssi <= TransportConfig.bleWeakLinkRSSICutoff { - // Requeue the candidate and try again later - connectionCandidates.append(candidate) - let remaining = TransportConfig.bleWeakLinkCooldownSeconds - elapsed - let delay = min(max(2.0, remaining), 15.0) - bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in self?.tryConnectFromQueue() } - return - } - } - if peripherals[peripheralID]?.isConnected == true || peripherals[peripheralID]?.isConnecting == true { - // Already in progress; skip - bleQueue.async { [weak self] in self?.tryConnectFromQueue() } - return - } - // Initiate connection - peripherals[peripheralID] = PeripheralState( - peripheral: peripheral, - characteristic: nil, - peerID: nil, - isConnecting: true, - isConnected: false, - lastConnectionAttempt: Date(), - assembler: NotificationStreamAssembler() + + let decision = connectionScheduler.nextCandidate( + connectedOrConnectingCount: linkStateStore.connectedOrConnectingPeripheralCount, + isAlreadyConnectingOrConnected: { [linkStateStore] peripheralID in + let state = linkStateStore.state(forPeripheralID: peripheralID) + return state?.isConnected == true || state?.isConnecting == true + }, + now: Date() ) + + switch decision { + case .none: + return + case .retryAfter(let delay): + bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in self?.tryConnectFromQueue() } + case .connect(let candidate): + beginCentralConnection(candidate, using: central, logPrefix: "⏩ Queue connect") + } + } + + private func beginCentralConnection( + _ candidate: BLEConnectionCandidate, + using central: CBCentralManager, + logPrefix: String + ) { + let peripheral = candidate.peripheral + let peripheralID = candidate.peripheralID + linkStateStore.beginConnecting(to: peripheral, at: Date()) peripheral.delegate = self let options: [String: Any] = [ CBConnectPeripheralOptionNotifyOnConnectionKey: true, @@ -1778,8 +1619,25 @@ extension BLEService { CBConnectPeripheralOptionNotifyOnNotificationKey: true ] central.connect(peripheral, options: options) - lastGlobalConnectAttempt = Date() - SecureLogger.debug("⏩ Queue connect: \(candidate.name) [RSSI:\(candidate.rssi)]", category: .session) + connectionScheduler.recordConnectionAttempt(at: Date()) + SecureLogger.debug("\(logPrefix): \(candidate.name) [RSSI:\(candidate.rssi)]", category: .session) + + bleQueue.asyncAfter(deadline: .now() + TransportConfig.bleConnectTimeoutSeconds) { [weak self] in + guard let self = self, + let state = self.linkStateStore.state(forPeripheralID: peripheralID), + state.isConnecting && !state.isConnected else { return } + + guard peripheral.state != .connected else { + SecureLogger.debug("⏱️ Timeout fired but peripheral already connected: \(candidate.name)", category: .session) + return + } + + SecureLogger.debug("⏱️ Timeout: \(candidate.name)", category: .session) + central.cancelPeripheralConnection(peripheral) + _ = self.linkStateStore.removePeripheral(peripheralID) + self.connectionScheduler.recordConnectionTimeout(peripheralID: peripheralID, at: Date()) + self.tryConnectFromQueue() + } } } @@ -1906,10 +1764,7 @@ extension BLEService: CBPeripheralDelegate { // Store characteristic in our consolidated structure let peripheralID = peripheral.identifier.uuidString - if var state = peripherals[peripheralID] { - state.characteristic = characteristic - peripherals[peripheralID] = state - } + linkStateStore.updateCharacteristic(characteristic, forPeripheralID: peripheralID) // Subscribe for notifications if characteristic.properties.contains(.notify) { @@ -1944,7 +1799,7 @@ extension BLEService: CBPeripheralDelegate { private func bufferNotificationChunk(_ chunk: Data, from peripheral: CBPeripheral) { let peripheralUUID = peripheral.identifier.uuidString - var state = peripherals[peripheralUUID] ?? PeripheralState( + var state = linkStateStore.state(forPeripheralID: peripheralUUID) ?? BLEPeripheralLinkState( peripheral: peripheral, characteristic: nil, peerID: nil, @@ -1957,7 +1812,7 @@ extension BLEService: CBPeripheralDelegate { var assembler = state.assembler let result = assembler.append(chunk) state.assembler = assembler - peripherals[peripheralUUID] = state + linkStateStore.setPeripheralState(state, for: peripheralUUID) for byte in result.droppedPrefixes { SecureLogger.warning("⚠️ Dropping byte from BLE stream (unexpected prefix \(String(format: "%02x", byte)))", category: .session) @@ -1969,7 +1824,7 @@ extension BLEService: CBPeripheralDelegate { // Codex review identified TOCTOU in this patch. // Enforce per-link sender binding immediately within the same notification batch. - // NOTE: `processNotificationPacket` may bind `peripherals[peripheralUUID].peerID` when an announce + // 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? = state.peerID @@ -2001,7 +1856,7 @@ extension BLEService: CBPeripheralDelegate { packet.ttl == messageTTL { boundPeerID = claimedSenderID state.peerID = claimedSenderID - peripherals[peripheralUUID] = state + linkStateStore.bindPeripheral(peripheralUUID, to: claimedSenderID) } if !recordIngressIfNew(packet, link: .peripheral(peripheralUUID), peerID: context.receivedFromPeerID) { @@ -2025,11 +1880,7 @@ extension BLEService: CBPeripheralDelegate { if packet.type == MessageType.announce.rawValue { if packet.ttl == messageTTL { - if var state = peripherals[peripheralUUID] { - state.peerID = senderID - peripherals[peripheralUUID] = state - } - peerToPeripheralUUID[senderID] = peripheralUUID + linkStateStore.bindPeripheral(peripheralUUID, to: senderID) refreshLocalTopology() } @@ -2065,10 +1916,9 @@ extension BLEService: CBPeripheralDelegate { guard shouldRediscover else { return } let peripheralID = peripheral.identifier.uuidString - if var state = peripherals[peripheralID] { - state.characteristic = nil - state.assembler = NotificationStreamAssembler() - peripherals[peripheralID] = state + linkStateStore.updatePeripheral(peripheralID) { + $0.characteristic = nil + $0.assembler = NotificationStreamAssembler() } SecureLogger.debug("🔄 BitChat service changed for \(peripheral.name ?? peripheral.identifier.uuidString), rediscovering", category: .session) @@ -2123,9 +1973,7 @@ extension BLEService: CBPeripheralManagerDelegate { SecureLogger.info("📴 Bluetooth powered off - cleaning up peripheral state", category: .session) peripheral.stopAdvertising() // Clear subscribed centrals (they are now invalid) - let centralPeerIDs = centralToPeerID.values.map { $0 } - subscribedCentrals.removeAll() - centralToPeerID.removeAll() + let centralPeerIDs = linkStateStore.clearCentrals() centralSubscriptionRateLimits.removeAll() characteristic = nil // Notify UI of disconnections @@ -2139,8 +1987,7 @@ extension BLEService: CBPeripheralManagerDelegate { // User denied Bluetooth permission SecureLogger.warning("🚫 Bluetooth unauthorized for peripheral role", category: .session) peripheral.stopAdvertising() - subscribedCentrals.removeAll() - centralToPeerID.removeAll() + _ = linkStateStore.clearCentrals() centralSubscriptionRateLimits.removeAll() characteristic = nil @@ -2204,7 +2051,7 @@ extension BLEService: CBPeripheralManagerDelegate { func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didSubscribeTo characteristic: CBCharacteristic) { let centralUUID = central.identifier.uuidString SecureLogger.debug("📥 Central subscribed: \(centralUUID.prefix(8))…", category: .session) - subscribedCentrals.append(central) + linkStateStore.addSubscribedCentral(central) // BCH-01-004: Rate-limit subscription-triggered announces to prevent enumeration attacks let now = Date() @@ -2283,7 +2130,7 @@ extension BLEService: CBPeripheralManagerDelegate { func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didUnsubscribeFrom characteristic: CBCharacteristic) { SecureLogger.debug("📤 Central unsubscribed: \(central.identifier.uuidString.prefix(8))…", category: .session) - subscribedCentrals.removeAll { $0.identifier == central.identifier } + let removedPeerID = linkStateStore.removeSubscribedCentral(central) // Ensure we're still advertising for other devices to find us if peripheral.isAdvertising == false { @@ -2292,8 +2139,7 @@ extension BLEService: CBPeripheralManagerDelegate { } // Find and disconnect the peer associated with this central - let centralUUID = central.identifier.uuidString - if let peerID = centralToPeerID[centralUUID] { + if let peerID = removedPeerID { // Mark peer as not connected; retain for reachability collectionsQueue.sync(flags: .barrier) { if var info = peers[peerID] { @@ -2302,8 +2148,6 @@ extension BLEService: CBPeripheralManagerDelegate { } } - // Clean up mappings - centralToPeerID.removeValue(forKey: centralUUID) refreshLocalTopology() // Update UI immediately @@ -2439,7 +2283,7 @@ extension BLEService: CBPeripheralManagerDelegate { let context = makeIngressPacketContext( for: packet, claimedSenderID: claimedSenderID, - boundPeerID: centralToPeerID[centralUUID], + boundPeerID: linkStateStore.peerID(forCentralUUID: centralUUID), linkDescription: "Central \(centralUUID.prefix(8))…" ) guard let context else { return } @@ -2452,13 +2296,11 @@ extension BLEService: CBPeripheralManagerDelegate { SecureLogger.debug("📦 Decoded (combined) packet type: \(packet.type) from sender: \(claimedSenderID.id.prefix(8))…", category: .session) } - if !subscribedCentrals.contains(central) { - subscribedCentrals.append(central) - } + linkStateStore.addSubscribedCentral(central) if packet.type == MessageType.announce.rawValue, packet.ttl == messageTTL { - centralToPeerID[centralUUID] = claimedSenderID + linkStateStore.bindCentral(centralUUID, to: claimedSenderID) refreshLocalTopology() } @@ -2537,7 +2379,7 @@ extension BLEService { ( connected: peers.values.filter { $0.isConnected }.count, known: peers.count, - candidates: connectionCandidates.count + candidates: connectionScheduler.candidateCount ) } @@ -2662,39 +2504,12 @@ extension BLEService { /// Safely fetch the current direct-link state for a peer using the BLE queue. private func linkState(for peerID: PeerID) -> (hasPeripheral: Bool, hasCentral: Bool) { - let computeState = { () -> (Bool, Bool) in - let peripheralUUID = self.peerToPeripheralUUID[peerID] - let hasPeripheral = peripheralUUID.flatMap { self.peripherals[$0]?.isConnected } ?? false - let hasCentral = self.centralToPeerID.values.contains(peerID) - return (hasPeripheral, hasCentral) - } - - if DispatchQueue.getSpecific(key: bleQueueKey) != nil { - return computeState() - } else { - return bleQueue.sync { computeState() } - } + let state = readLinkState { $0.directLinkState(for: peerID) } + return (state.hasPeripheral, state.hasCentral) } private func links(to peerID: PeerID?) -> Set { - guard let peerID else { return [] } - - let computeLinks = { () -> Set in - var links: Set = [] - if let peripheralUUID = self.peerToPeripheralUUID[peerID] { - links.insert(.peripheral(peripheralUUID)) - } - for (centralUUID, mappedPeerID) in self.centralToPeerID where mappedPeerID == peerID { - links.insert(.central(centralUUID)) - } - return links - } - - if DispatchQueue.getSpecific(key: bleQueueKey) != nil { - return computeLinks() - } else { - return bleQueue.sync { computeLinks() } - } + readLinkState { $0.links(to: peerID) } } private func configureNoiseServiceCallbacks(for service: NoiseEncryptionService) { @@ -2751,20 +2566,25 @@ extension BLEService { } // MARK: Link capability snapshots (thread-safe via bleQueue) - - private func snapshotPeripheralStates() -> [PeripheralState] { + + private func readLinkState(_ body: (BLELinkStateStore) -> T) -> T { if DispatchQueue.getSpecific(key: bleQueueKey) != nil { - return Array(peripherals.values) + return body(linkStateStore) } else { - return bleQueue.sync { Array(peripherals.values) } + return bleQueue.sync { body(linkStateStore) } } } - private func snapshotSubscribedCentrals() -> ([CBCentral], [String: PeerID]) { - if DispatchQueue.getSpecific(key: bleQueueKey) != nil { - return (self.subscribedCentrals, self.centralToPeerID) - } else { - return bleQueue.sync { (self.subscribedCentrals, self.centralToPeerID) } - } + + private func snapshotDirectPeripheralState(for peerID: PeerID) -> BLEPeripheralLinkState? { + readLinkState { $0.directPeripheralState(for: peerID) } + } + + private func snapshotPeripheralStates() -> [BLEPeripheralLinkState] { + readLinkState(\.peripheralStates) + } + + private func snapshotSubscribedCentrals() -> BLESubscribedCentralSnapshot { + readLinkState(\.subscribedCentralSnapshot) } // MARK: Helpers: IDs, selection, and write backpressure @@ -2802,7 +2622,7 @@ extension BLEService { let uuid = peripheral.identifier.uuidString bleQueue.async { [weak self] in guard let self = self else { return } - guard let state = self.peripherals[uuid], let ch = state.characteristic else { return } + guard let state = self.linkStateStore.state(forPeripheralID: uuid), let ch = state.characteristic else { return } // Atomically take all pending items from the queue to avoid race conditions // where new items could be enqueued between read and update @@ -2841,7 +2661,7 @@ extension BLEService { private func drainAllPendingWrites() { let uuids = collectionsQueue.sync { pendingPeripheralWrites.peripheralIDs } for uuid in uuids { - guard let state = peripherals[uuid], state.isConnected else { continue } + guard let state = linkStateStore.state(forPeripheralID: uuid), state.isConnected else { continue } drainPendingWrites(for: state.peripheral) } } @@ -4022,7 +3842,7 @@ extension BLEService { // Clean old connection timeout backoff entries (> window) let timeoutCutoff = now.addingTimeInterval(-TransportConfig.bleConnectTimeoutBackoffWindowSeconds) - recentConnectTimeouts = recentConnectTimeouts.filter { $0.value >= timeoutCutoff } + connectionScheduler.pruneConnectionTimeouts(before: timeoutCutoff) // Clean up stale scheduled relays that somehow persisted (> 2s) collectionsQueue.async(flags: .barrier) { [weak self] in @@ -4117,32 +3937,10 @@ extension BLEService { } private func updateRSSIThreshold(connectedCount: Int) { - // Adjust RSSI threshold based on connectivity, candidate pressure, and failures - if connectedCount == 0 { - // Isolated: relax floor slowly to hunt for distant nodes - if lastIsolatedAt == nil { lastIsolatedAt = Date() } - let iso = lastIsolatedAt ?? Date() - let elapsed = Date().timeIntervalSince(iso) - if elapsed > TransportConfig.bleIsolationRelaxThresholdSeconds { - dynamicRSSIThreshold = TransportConfig.bleRSSIIsolatedRelaxed - } else { - dynamicRSSIThreshold = TransportConfig.bleRSSIIsolatedBase - } - return - } - lastIsolatedAt = nil - // Base threshold when connected - var threshold = TransportConfig.bleDynamicRSSIThresholdDefault - // If we're at budget or queue is large, prefer closer peers - let linkCount = peripherals.values.filter { $0.isConnected || $0.isConnecting }.count - if linkCount >= maxCentralLinks || connectionCandidates.count > TransportConfig.bleConnectionCandidatesMax { - threshold = TransportConfig.bleRSSIConnectedThreshold - } - // If we have many recent timeouts, raise further - let recentTimeouts = recentConnectTimeouts.filter { Date().timeIntervalSince($0.value) < TransportConfig.bleRecentTimeoutWindowSeconds }.count - if recentTimeouts >= TransportConfig.bleRecentTimeoutCountThreshold { - threshold = max(threshold, TransportConfig.bleRSSIHighTimeoutThreshold) - } - dynamicRSSIThreshold = threshold + connectionScheduler.updateRSSIThreshold( + connectedCount: connectedCount, + connectedOrConnectingLinkCount: linkStateStore.connectedOrConnectingPeripheralCount, + now: Date() + ) } } diff --git a/bitchatTests/Services/BLEConnectionSchedulerTests.swift b/bitchatTests/Services/BLEConnectionSchedulerTests.swift new file mode 100644 index 00000000..f88588f5 --- /dev/null +++ b/bitchatTests/Services/BLEConnectionSchedulerTests.swift @@ -0,0 +1,156 @@ +import Foundation +import Testing +@testable import bitchat + +struct BLEConnectionSchedulerTests { + @Test + func discoveryQueuesWeakSignalCandidates() { + let scheduler = BLEConnectionScheduler(dynamicRSSIThreshold: -80) + let now = Date() + let candidate = makeCandidate(id: "p1", rssi: -85, now: now) + + let decision = scheduler.handleDiscovery( + candidate, + connectedOrConnectingCount: 0, + existingState: nil, + peripheralState: .disconnected, + now: now + ) + + #expect(decision == .queued) + #expect(scheduler.candidateCount == 1) + } + + @Test + func discoveryQueuesAndSchedulesRetryWhenRateLimited() { + let scheduler = BLEConnectionScheduler(connectRateLimitInterval: 1.0) + let now = Date() + scheduler.recordConnectionAttempt(at: now) + + let decision = scheduler.handleDiscovery( + makeCandidate(id: "p1", rssi: -50, now: now), + connectedOrConnectingCount: 0, + existingState: nil, + peripheralState: .disconnected, + now: now.addingTimeInterval(0.25) + ) + + guard case .scheduleRetry(let delay) = decision else { + Issue.record("Expected scheduleRetry, got \(decision)") + return + } + #expect(delay > 0) + #expect(scheduler.candidateCount == 1) + } + + @Test + func nextCandidateSelectsBestScoredCandidate() { + let scheduler = BLEConnectionScheduler() + let now = Date() + scheduler.enqueue(makeCandidate(id: "weak", rssi: -88, now: now)) + scheduler.enqueue(makeCandidate(id: "strong", rssi: -60, now: now.addingTimeInterval(1))) + + let decision = scheduler.nextCandidate( + connectedOrConnectingCount: 0, + isAlreadyConnectingOrConnected: { _ in false }, + now: now.addingTimeInterval(2) + ) + + guard case .connect(let candidate) = decision else { + Issue.record("Expected connect decision") + return + } + #expect(candidate.peripheralID == "strong") + } + + @Test + func enqueueReplacesExistingPeripheralCandidate() { + let scheduler = BLEConnectionScheduler() + let now = Date() + scheduler.enqueue(makeCandidate(id: "same", rssi: -90, now: now)) + scheduler.enqueue(makeCandidate(id: "same", rssi: -55, now: now.addingTimeInterval(1))) + + let decision = scheduler.nextCandidate( + connectedOrConnectingCount: 0, + isAlreadyConnectingOrConnected: { _ in false }, + now: now.addingTimeInterval(2) + ) + + guard case .connect(let candidate) = decision else { + Issue.record("Expected connect decision") + return + } + #expect(scheduler.candidateCount == 0) + #expect(candidate.peripheralID == "same") + #expect(candidate.rssi == -55) + } + + @Test + func weakTimedOutCandidateIsRequeuedWithRetryDelay() { + let scheduler = BLEConnectionScheduler( + weakLinkCooldownSeconds: 30, + weakLinkRSSICutoff: -90 + ) + let now = Date() + scheduler.recordConnectionTimeout(peripheralID: "weak", at: now) + scheduler.enqueue(makeCandidate(id: "weak", rssi: -95, now: now)) + + let decision = scheduler.nextCandidate( + connectedOrConnectingCount: 0, + isAlreadyConnectingOrConnected: { _ in false }, + now: now.addingTimeInterval(5) + ) + + guard case .retryAfter(let delay) = decision else { + Issue.record("Expected retryAfter decision") + return + } + #expect(delay == 15) + #expect(scheduler.candidateCount == 1) + } + + @Test + func rssiThresholdTightensAfterRepeatedRecentTimeouts() { + let scheduler = BLEConnectionScheduler() + let now = Date() + scheduler.recordConnectionTimeout(peripheralID: "p1", at: now) + scheduler.recordConnectionTimeout(peripheralID: "p2", at: now) + scheduler.recordConnectionTimeout(peripheralID: "p3", at: now) + + let threshold = scheduler.updateRSSIThreshold( + connectedCount: 1, + connectedOrConnectingLinkCount: 1, + now: now.addingTimeInterval(1) + ) + + #expect(threshold == TransportConfig.bleRSSIHighTimeoutThreshold) + #expect(scheduler.dynamicRSSIThreshold == TransportConfig.bleRSSIHighTimeoutThreshold) + } + + @Test + func rssiThresholdTightensWhenCandidateQueueIsFull() { + let scheduler = BLEConnectionScheduler(candidateCap: 2) + let now = Date() + scheduler.enqueue(makeCandidate(id: "p1", rssi: -65, now: now)) + scheduler.enqueue(makeCandidate(id: "p2", rssi: -66, now: now)) + + let threshold = scheduler.updateRSSIThreshold( + connectedCount: 1, + connectedOrConnectingLinkCount: 1, + now: now + ) + + #expect(threshold == TransportConfig.bleRSSIConnectedThreshold) + } +} + +private func makeCandidate(id: String, rssi: Int, now: Date) -> BLEConnectionCandidate { + BLEConnectionCandidate( + peripheral: id, + peripheralID: id, + rssi: rssi, + name: id, + isConnectable: true, + discoveredAt: now + ) +} diff --git a/bitchatTests/Services/BLELinkStateStoreTests.swift b/bitchatTests/Services/BLELinkStateStoreTests.swift new file mode 100644 index 00000000..9a5179df --- /dev/null +++ b/bitchatTests/Services/BLELinkStateStoreTests.swift @@ -0,0 +1,46 @@ +import BitFoundation +import Testing +@testable import bitchat + +struct BLELinkStateStoreTests { + @Test + func centralBindingExposesDirectLinkStateAndLinks() { + let store = BLELinkStateStore() + let peerID = PeerID(str: "1122334455667788") + + store.bindCentral("central-a", to: peerID) + + #expect(store.peerID(forCentralUUID: "central-a") == peerID) + #expect(store.directLinkState(for: peerID) == BLEDirectLinkState(hasPeripheral: false, hasCentral: true)) + #expect(store.links(to: peerID) == [.central("central-a")]) + } + + @Test + func linksReturnsAllCentralBindingsForPeer() { + let store = BLELinkStateStore() + let peerID = PeerID(str: "1122334455667788") + let otherPeerID = PeerID(str: "8899aabbccddeeff") + + store.bindCentral("central-a", to: peerID) + store.bindCentral("central-b", to: peerID) + store.bindCentral("central-c", to: otherPeerID) + + #expect(store.links(to: peerID) == [.central("central-a"), .central("central-b")]) + } + + @Test + func clearCentralsReturnsPreviouslyBoundPeerIDsAndClearsLookups() { + let store = BLELinkStateStore() + let firstPeerID = PeerID(str: "1122334455667788") + let secondPeerID = PeerID(str: "8899aabbccddeeff") + + store.bindCentral("central-a", to: firstPeerID) + store.bindCentral("central-b", to: secondPeerID) + + let removedPeerIDs = Set(store.clearCentrals()) + + #expect(removedPeerIDs == Set([firstPeerID, secondPeerID])) + #expect(store.peerID(forCentralUUID: "central-a") == nil) + #expect(store.links(to: firstPeerID).isEmpty) + } +}