[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 <jackjackbits@users.noreply.github.com>
This commit is contained in:
jack
2026-05-31 14:24:09 +02:00
committed by GitHub
co-authored by jack
parent 9e84f5e822
commit ffa0d7aa4f
5 changed files with 835 additions and 368 deletions
@@ -0,0 +1,271 @@
import Foundation
struct BLEConnectionCandidate<Peripheral> {
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<Peripheral> {
case none
case retryAfter(TimeInterval)
case connect(BLEConnectionCandidate<Peripheral>)
}
final class BLEConnectionScheduler<Peripheral> {
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<Peripheral>] = []
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<Peripheral>,
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<Peripheral>) {
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<Peripheral> {
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<Peripheral>,
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<Peripheral>, 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
}
}
@@ -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<BLEIngressLinkID> {
guard let peerID else { return [] }
var links: Set<BLEIngressLinkID> = []
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()
}
}
+163 -365
View File
@@ -35,22 +35,8 @@ final class BLEService: NSObject {
// MARK: - Core State (5 Essential Collections) // MARK: - Core State (5 Essential Collections)
// 1. Consolidated Peripheral Tracking // 1. Consolidated BLE link tracking for both central and peripheral roles.
private struct PeripheralState { private var linkStateStore = BLELinkStateStore()
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
// BCH-01-004: Rate-limiting for subscription-triggered announces // BCH-01-004: Rate-limiting for subscription-triggered announces
// Tracks subscription attempts per central to prevent enumeration attacks // Tracks subscription attempts per central to prevent enumeration attacks
@@ -85,8 +71,6 @@ final class BLEService: NSObject {
private var fragmentAssemblyBuffer = BLEFragmentAssemblyBuffer() private var fragmentAssemblyBuffer = BLEFragmentAssemblyBuffer()
private var outboundFragmentTransfers = BLEOutboundFragmentTransferScheduler() private var outboundFragmentTransfers = BLEOutboundFragmentTransferScheduler()
private let incomingFileStore = BLEIncomingFileStore() private let incomingFileStore = BLEIncomingFileStore()
// Backoff for peripherals that recently timed out connecting
private var recentConnectTimeouts: [String: Date] = [:] // Peripheral UUID -> last timeout
// Simple announce throttling // Simple announce throttling
private var lastAnnounceSent = Date.distantPast private var lastAnnounceSent = Date.distantPast
@@ -156,20 +140,7 @@ final class BLEService: NSObject {
private var maintenanceCounter = 0 // Track maintenance cycles private var maintenanceCounter = 0 // Track maintenance cycles
// MARK: - Connection budget & scheduling (central role) // MARK: - Connection budget & scheduling (central role)
private let maxCentralLinks = TransportConfig.bleMaxCentralLinks private var connectionScheduler = BLEConnectionScheduler<CBPeripheral>()
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
// MARK: - Adaptive scanning duty-cycle // MARK: - Adaptive scanning duty-cycle
private var scanDutyTimer: DispatchSourceTimer? private var scanDutyTimer: DispatchSourceTimer?
@@ -348,7 +319,7 @@ final class BLEService: NSObject {
bleQueue.sync { bleQueue.sync {
pendingWriteBuffers.removeAll() pendingWriteBuffers.removeAll()
recentConnectTimeouts.removeAll() connectionScheduler.reset()
} }
recentDisconnectNotifies.removeAll() recentDisconnectNotifies.removeAll()
@@ -506,7 +477,7 @@ final class BLEService: NSObject {
// Snapshot BLE state under bleQueue to avoid races with delegate callbacks // Snapshot BLE state under bleQueue to avoid races with delegate callbacks
let (peripheralStates, centralsCount, char) = bleQueue.sync { 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 // Send to peripherals we're connected to as central
@@ -543,7 +514,7 @@ final class BLEService: NSObject {
peripheralManager?.stopAdvertising() peripheralManager?.stopAdvertising()
// Disconnect all peripherals (synchronized access) // Disconnect all peripherals (synchronized access)
let peripheralsToDisconnect = bleQueue.sync { Array(peripherals.values) } let peripheralsToDisconnect = bleQueue.sync { linkStateStore.peripheralStates }
for state in peripheralsToDisconnect { for state in peripheralsToDisconnect {
centralManager?.cancelPeripheralConnection(state.peripheral) centralManager?.cancelPeripheralConnection(state.peripheral)
} }
@@ -573,10 +544,8 @@ final class BLEService: NSObject {
// Clear peripheral references (synchronized access to avoid races with BLE callbacks) // Clear peripheral references (synchronized access to avoid races with BLE callbacks)
bleQueue.sync { bleQueue.sync {
peripherals.removeAll() linkStateStore.clearAll()
peerToPeripheralUUID.removeAll() connectionScheduler.reset()
subscribedCentrals.removeAll()
centralToPeerID.removeAll()
centralSubscriptionRateLimits.removeAll() centralSubscriptionRateLimits.removeAll()
} }
meshTopology.reset() meshTopology.reset()
@@ -888,16 +857,8 @@ final class BLEService: NSObject {
let outboundPriority = BLEOutboundPacketPolicy.priority(for: packet, data: data) let outboundPriority = BLEOutboundPacketPolicy.priority(for: packet, data: data)
// Per-link limits for the specific peer // Per-link limits for the specific peer
let directPeripheralState: PeripheralState? = { let directPeripheralState = snapshotDirectPeripheralState(for: recipientPeerID)
if DispatchQueue.getSpecific(key: bleQueueKey) != nil { let recipientCentral = snapshotSubscribedCentrals().central(for: recipientPeerID)
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 }
if let peripheralMaxLen = directPeripheralState?.peripheral.maximumWriteValueLength(for: .withoutResponse), if let peripheralMaxLen = directPeripheralState?.peripheral.maximumWriteValueLength(for: .withoutResponse),
data.count > peripheralMaxLen { data.count > peripheralMaxLen {
@@ -987,7 +948,7 @@ final class BLEService: NSObject {
let m = s.peripheral.maximumWriteValueLength(for: .withoutResponse) let m = s.peripheral.maximumWriteValueLength(for: .withoutResponse)
minCentralWriteLen = minCentralWriteLen.map { min($0, m) } ?? m 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() let minNotifyLen = subscribedCentrals.map { $0.maximumUpdateValueLength }.min()
// Avoid re-fragmenting fragment packets // Avoid re-fragmenting fragment packets
@@ -1378,14 +1339,14 @@ extension BLEService: CBCentralManagerDelegate {
for peripheral in restoredPeripherals { for peripheral in restoredPeripherals {
let identifier = peripheral.identifier.uuidString let identifier = peripheral.identifier.uuidString
peripheral.delegate = self peripheral.delegate = self
let existing = peripherals[identifier] let existing = linkStateStore.state(forPeripheralID: identifier)
let assembler = existing?.assembler ?? NotificationStreamAssembler() let assembler = existing?.assembler ?? NotificationStreamAssembler()
let characteristic = existing?.characteristic let characteristic = existing?.characteristic
let peerID = existing?.peerID let peerID = existing?.peerID
let wasConnecting = existing?.isConnecting ?? false let wasConnecting = existing?.isConnecting ?? false
let wasConnected = existing?.isConnected ?? false let wasConnected = existing?.isConnected ?? false
let restoredState = PeripheralState( let restoredState = BLEPeripheralLinkState(
peripheral: peripheral, peripheral: peripheral,
characteristic: characteristic, characteristic: characteristic,
peerID: peerID, peerID: peerID,
@@ -1394,7 +1355,7 @@ extension BLEService: CBCentralManagerDelegate {
lastConnectionAttempt: existing?.lastConnectionAttempt, lastConnectionAttempt: existing?.lastConnectionAttempt,
assembler: assembler assembler: assembler
) )
peripherals[identifier] = restoredState linkStateStore.setPeripheralState(restoredState, for: identifier)
} }
captureBluetoothStatus(context: "central-restore") captureBluetoothStatus(context: "central-restore")
@@ -1418,12 +1379,12 @@ extension BLEService: CBCentralManagerDelegate {
SecureLogger.info("📴 Bluetooth powered off - cleaning up central state", category: .session) SecureLogger.info("📴 Bluetooth powered off - cleaning up central state", category: .session)
central.stopScan() central.stopScan()
// Mark all peripheral connections as disconnected (they are now invalid) // Mark all peripheral connections as disconnected (they are now invalid)
let peerIDs: [PeerID] = peripherals.compactMap { $0.value.peerID } let peripheralStates = linkStateStore.peripheralStates
for state in peripherals.values { let peerIDs: [PeerID] = peripheralStates.compactMap(\.peerID)
for state in peripheralStates {
central.cancelPeripheralConnection(state.peripheral) central.cancelPeripheralConnection(state.peripheral)
} }
peripherals.removeAll() _ = linkStateStore.clearPeripherals()
peerToPeripheralUUID.removeAll()
// Notify UI of disconnections // Notify UI of disconnections
for peerID in peerIDs { for peerID in peerIDs {
notifyUI { [weak self] in notifyUI { [weak self] in
@@ -1435,8 +1396,7 @@ extension BLEService: CBCentralManagerDelegate {
// User denied Bluetooth permission // User denied Bluetooth permission
SecureLogger.warning("🚫 Bluetooth unauthorized - user denied permission", category: .session) SecureLogger.warning("🚫 Bluetooth unauthorized - user denied permission", category: .session)
central.stopScan() central.stopScan()
peripherals.removeAll() _ = linkStateStore.clearPeripherals()
peerToPeripheralUUID.removeAll()
case .unsupported: case .unsupported:
// Device doesn't support BLE // Device doesn't support BLE
@@ -1482,159 +1442,46 @@ extension BLEService: CBCentralManagerDelegate {
let isConnectable = (advertisementData[CBAdvertisementDataIsConnectable] as? NSNumber)?.boolValue ?? true let isConnectable = (advertisementData[CBAdvertisementDataIsConnectable] as? NSNumber)?.boolValue ?? true
let rssiValue = RSSI.intValue let rssiValue = RSSI.intValue
// Skip if peripheral is not connectable (per advertisement data) let candidate = BLEConnectionCandidate(
guard isConnectable else { return } peripheral: peripheral,
peripheralID: peripheralID,
rssi: rssiValue,
name: String(advertisedName),
isConnectable: isConnectable,
discoveredAt: Date()
)
let existingState = linkStateStore.state(forPeripheralID: peripheralID).map(BLEExistingConnectionState.init)
// Skip immediate connect if signal too weak for current conditions; enqueue instead switch connectionScheduler.handleDiscovery(
if rssiValue <= dynamicRSSIThreshold { candidate,
connectionCandidates.append(ConnectionCandidate(peripheral: peripheral, rssi: rssiValue, name: String(advertisedName), isConnectable: isConnectable, discoveredAt: Date())) connectedOrConnectingCount: linkStateStore.connectedOrConnectingPeripheralCount,
// Keep list tidy existingState: existingState,
connectionCandidates.sort { (a, b) in peripheralState: peripheral.state.connectionSchedulerState,
if a.rssi != b.rssi { return a.rssi > b.rssi } now: candidate.discoveredAt
return a.discoveredAt < b.discoveredAt ) {
} case .ignore, .queued:
if connectionCandidates.count > TransportConfig.bleConnectionCandidatesMax {
connectionCandidates.removeLast(connectionCandidates.count - TransportConfig.bleConnectionCandidatesMax)
}
return return
} case .scheduleRetry(let delay):
// 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
}
// 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
bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in
self?.tryConnectFromQueue() self?.tryConnectFromQueue()
} }
return return
} case .cancelStaleConnection:
// 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
central.cancelPeripheralConnection(peripheral) central.cancelPeripheralConnection(peripheral)
// Will retry on next discovery
return return
} case .connectNow:
beginCentralConnection(candidate, using: central, logPrefix: "📱 Connect")
// 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()
} }
} }
func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeripheral) { func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeripheral) {
let peripheralID = peripheral.identifier.uuidString let peripheralID = peripheral.identifier.uuidString
// Update state to connected // Update state to connected
if var state = peripherals[peripheralID] { linkStateStore.markConnected(peripheral)
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()
)
}
// Reset backoff state on success // Reset backoff state on success
failureCounts[peripheralID] = 0 connectionScheduler.recordConnectionSuccess(peripheralID: peripheralID)
recentConnectTimeouts.removeValue(forKey: peripheralID)
SecureLogger.debug("✅ Connected: \(peripheral.name ?? "Unknown") [\(peripheralID)]", category: .session) 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 let peripheralID = peripheral.identifier.uuidString
// Find the peer ID if we have it // 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) 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 disconnect carried an error (often timeout), apply short backoff to avoid thrash
if error != nil { if error != nil {
recentConnectTimeouts[peripheralID] = Date() connectionScheduler.recordDisconnectError(peripheralID: peripheralID, at: Date())
} }
// Clean up references // Clean up references and peer mappings
peripherals.removeValue(forKey: peripheralID) _ = linkStateStore.removePeripheral(peripheralID)
// Clean up peer mappings
if let peerID { if let peerID {
peerToPeripheralUUID.removeValue(forKey: peerID)
// Do not remove peer; mark as not connected but retain for reachability // Do not remove peer; mark as not connected but retain for reachability
collectionsQueue.sync(flags: .barrier) { collectionsQueue.sync(flags: .barrier) {
if var info = peers[peerID] { if var info = peers[peerID] {
@@ -1703,74 +1546,72 @@ func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeriph
let peripheralID = peripheral.identifier.uuidString let peripheralID = peripheral.identifier.uuidString
// Clean up the references // 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) 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 // Try next candidate
bleQueue.async { [weak self] in self?.tryConnectFromQueue() } bleQueue.async { [weak self] in self?.tryConnectFromQueue() }
} }
} }
// MARK: - Connection scheduling helpers // 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 { extension BLEService {
private func tryConnectFromQueue() { private func tryConnectFromQueue() {
guard let central = centralManager, central.state == .poweredOn else { return } guard let central = centralManager, central.state == .poweredOn else { return }
// Check budget and rate limit
let current = peripherals.values.filter { $0.isConnected || $0.isConnecting }.count let decision = connectionScheduler.nextCandidate(
guard current < maxCentralLinks else { return } connectedOrConnectingCount: linkStateStore.connectedOrConnectingPeripheralCount,
let delta = Date().timeIntervalSince(lastGlobalConnectAttempt) isAlreadyConnectingOrConnected: { [linkStateStore] peripheralID in
guard delta >= connectRateLimitInterval else { let state = linkStateStore.state(forPeripheralID: peripheralID)
let delay = connectRateLimitInterval - delta + 0.05 return state?.isConnected == true || state?.isConnecting == true
bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in self?.tryConnectFromQueue() } },
return now: Date()
}
// 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()
) )
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<CBPeripheral>,
using central: CBCentralManager,
logPrefix: String
) {
let peripheral = candidate.peripheral
let peripheralID = candidate.peripheralID
linkStateStore.beginConnecting(to: peripheral, at: Date())
peripheral.delegate = self peripheral.delegate = self
let options: [String: Any] = [ let options: [String: Any] = [
CBConnectPeripheralOptionNotifyOnConnectionKey: true, CBConnectPeripheralOptionNotifyOnConnectionKey: true,
@@ -1778,8 +1619,25 @@ extension BLEService {
CBConnectPeripheralOptionNotifyOnNotificationKey: true CBConnectPeripheralOptionNotifyOnNotificationKey: true
] ]
central.connect(peripheral, options: options) central.connect(peripheral, options: options)
lastGlobalConnectAttempt = Date() connectionScheduler.recordConnectionAttempt(at: Date())
SecureLogger.debug("⏩ Queue connect: \(candidate.name) [RSSI:\(candidate.rssi)]", category: .session) 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 // Store characteristic in our consolidated structure
let peripheralID = peripheral.identifier.uuidString let peripheralID = peripheral.identifier.uuidString
if var state = peripherals[peripheralID] { linkStateStore.updateCharacteristic(characteristic, forPeripheralID: peripheralID)
state.characteristic = characteristic
peripherals[peripheralID] = state
}
// Subscribe for notifications // Subscribe for notifications
if characteristic.properties.contains(.notify) { if characteristic.properties.contains(.notify) {
@@ -1944,7 +1799,7 @@ extension BLEService: CBPeripheralDelegate {
private func bufferNotificationChunk(_ chunk: Data, from peripheral: CBPeripheral) { private func bufferNotificationChunk(_ chunk: Data, from peripheral: CBPeripheral) {
let peripheralUUID = peripheral.identifier.uuidString let peripheralUUID = peripheral.identifier.uuidString
var state = peripherals[peripheralUUID] ?? PeripheralState( var state = linkStateStore.state(forPeripheralID: peripheralUUID) ?? BLEPeripheralLinkState(
peripheral: peripheral, peripheral: peripheral,
characteristic: nil, characteristic: nil,
peerID: nil, peerID: nil,
@@ -1957,7 +1812,7 @@ extension BLEService: CBPeripheralDelegate {
var assembler = state.assembler var assembler = state.assembler
let result = assembler.append(chunk) let result = assembler.append(chunk)
state.assembler = assembler state.assembler = assembler
peripherals[peripheralUUID] = state linkStateStore.setPeripheralState(state, for: peripheralUUID)
for byte in result.droppedPrefixes { for byte in result.droppedPrefixes {
SecureLogger.warning("⚠️ Dropping byte from BLE stream (unexpected prefix \(String(format: "%02x", byte)))", category: .session) 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. // Codex review identified TOCTOU in this patch.
// Enforce per-link sender binding immediately within the same notification batch. // 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 // 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. // we see a binding-eligible announce so subsequent frames can't spoof a different sender.
var boundPeerID: PeerID? = state.peerID var boundPeerID: PeerID? = state.peerID
@@ -2001,7 +1856,7 @@ extension BLEService: CBPeripheralDelegate {
packet.ttl == messageTTL { packet.ttl == messageTTL {
boundPeerID = claimedSenderID boundPeerID = claimedSenderID
state.peerID = claimedSenderID state.peerID = claimedSenderID
peripherals[peripheralUUID] = state linkStateStore.bindPeripheral(peripheralUUID, to: claimedSenderID)
} }
if !recordIngressIfNew(packet, link: .peripheral(peripheralUUID), peerID: context.receivedFromPeerID) { if !recordIngressIfNew(packet, link: .peripheral(peripheralUUID), peerID: context.receivedFromPeerID) {
@@ -2025,11 +1880,7 @@ extension BLEService: CBPeripheralDelegate {
if packet.type == MessageType.announce.rawValue { if packet.type == MessageType.announce.rawValue {
if packet.ttl == messageTTL { if packet.ttl == messageTTL {
if var state = peripherals[peripheralUUID] { linkStateStore.bindPeripheral(peripheralUUID, to: senderID)
state.peerID = senderID
peripherals[peripheralUUID] = state
}
peerToPeripheralUUID[senderID] = peripheralUUID
refreshLocalTopology() refreshLocalTopology()
} }
@@ -2065,10 +1916,9 @@ extension BLEService: CBPeripheralDelegate {
guard shouldRediscover else { return } guard shouldRediscover else { return }
let peripheralID = peripheral.identifier.uuidString let peripheralID = peripheral.identifier.uuidString
if var state = peripherals[peripheralID] { linkStateStore.updatePeripheral(peripheralID) {
state.characteristic = nil $0.characteristic = nil
state.assembler = NotificationStreamAssembler() $0.assembler = NotificationStreamAssembler()
peripherals[peripheralID] = state
} }
SecureLogger.debug("🔄 BitChat service changed for \(peripheral.name ?? peripheral.identifier.uuidString), rediscovering", category: .session) 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) SecureLogger.info("📴 Bluetooth powered off - cleaning up peripheral state", category: .session)
peripheral.stopAdvertising() peripheral.stopAdvertising()
// Clear subscribed centrals (they are now invalid) // Clear subscribed centrals (they are now invalid)
let centralPeerIDs = centralToPeerID.values.map { $0 } let centralPeerIDs = linkStateStore.clearCentrals()
subscribedCentrals.removeAll()
centralToPeerID.removeAll()
centralSubscriptionRateLimits.removeAll() centralSubscriptionRateLimits.removeAll()
characteristic = nil characteristic = nil
// Notify UI of disconnections // Notify UI of disconnections
@@ -2139,8 +1987,7 @@ extension BLEService: CBPeripheralManagerDelegate {
// User denied Bluetooth permission // User denied Bluetooth permission
SecureLogger.warning("🚫 Bluetooth unauthorized for peripheral role", category: .session) SecureLogger.warning("🚫 Bluetooth unauthorized for peripheral role", category: .session)
peripheral.stopAdvertising() peripheral.stopAdvertising()
subscribedCentrals.removeAll() _ = linkStateStore.clearCentrals()
centralToPeerID.removeAll()
centralSubscriptionRateLimits.removeAll() centralSubscriptionRateLimits.removeAll()
characteristic = nil characteristic = nil
@@ -2204,7 +2051,7 @@ extension BLEService: CBPeripheralManagerDelegate {
func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didSubscribeTo characteristic: CBCharacteristic) { func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didSubscribeTo characteristic: CBCharacteristic) {
let centralUUID = central.identifier.uuidString let centralUUID = central.identifier.uuidString
SecureLogger.debug("📥 Central subscribed: \(centralUUID.prefix(8))", category: .session) 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 // BCH-01-004: Rate-limit subscription-triggered announces to prevent enumeration attacks
let now = Date() let now = Date()
@@ -2283,7 +2130,7 @@ extension BLEService: CBPeripheralManagerDelegate {
func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didUnsubscribeFrom characteristic: CBCharacteristic) { func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didUnsubscribeFrom characteristic: CBCharacteristic) {
SecureLogger.debug("📤 Central unsubscribed: \(central.identifier.uuidString.prefix(8))", category: .session) 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 // Ensure we're still advertising for other devices to find us
if peripheral.isAdvertising == false { if peripheral.isAdvertising == false {
@@ -2292,8 +2139,7 @@ extension BLEService: CBPeripheralManagerDelegate {
} }
// Find and disconnect the peer associated with this central // Find and disconnect the peer associated with this central
let centralUUID = central.identifier.uuidString if let peerID = removedPeerID {
if let peerID = centralToPeerID[centralUUID] {
// Mark peer as not connected; retain for reachability // Mark peer as not connected; retain for reachability
collectionsQueue.sync(flags: .barrier) { collectionsQueue.sync(flags: .barrier) {
if var info = peers[peerID] { if var info = peers[peerID] {
@@ -2302,8 +2148,6 @@ extension BLEService: CBPeripheralManagerDelegate {
} }
} }
// Clean up mappings
centralToPeerID.removeValue(forKey: centralUUID)
refreshLocalTopology() refreshLocalTopology()
// Update UI immediately // Update UI immediately
@@ -2439,7 +2283,7 @@ extension BLEService: CBPeripheralManagerDelegate {
let context = makeIngressPacketContext( let context = makeIngressPacketContext(
for: packet, for: packet,
claimedSenderID: claimedSenderID, claimedSenderID: claimedSenderID,
boundPeerID: centralToPeerID[centralUUID], boundPeerID: linkStateStore.peerID(forCentralUUID: centralUUID),
linkDescription: "Central \(centralUUID.prefix(8))" linkDescription: "Central \(centralUUID.prefix(8))"
) )
guard let context else { return } 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) SecureLogger.debug("📦 Decoded (combined) packet type: \(packet.type) from sender: \(claimedSenderID.id.prefix(8))", category: .session)
} }
if !subscribedCentrals.contains(central) { linkStateStore.addSubscribedCentral(central)
subscribedCentrals.append(central)
}
if packet.type == MessageType.announce.rawValue, if packet.type == MessageType.announce.rawValue,
packet.ttl == messageTTL { packet.ttl == messageTTL {
centralToPeerID[centralUUID] = claimedSenderID linkStateStore.bindCentral(centralUUID, to: claimedSenderID)
refreshLocalTopology() refreshLocalTopology()
} }
@@ -2537,7 +2379,7 @@ extension BLEService {
( (
connected: peers.values.filter { $0.isConnected }.count, connected: peers.values.filter { $0.isConnected }.count,
known: peers.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. /// Safely fetch the current direct-link state for a peer using the BLE queue.
private func linkState(for peerID: PeerID) -> (hasPeripheral: Bool, hasCentral: Bool) { private func linkState(for peerID: PeerID) -> (hasPeripheral: Bool, hasCentral: Bool) {
let computeState = { () -> (Bool, Bool) in let state = readLinkState { $0.directLinkState(for: peerID) }
let peripheralUUID = self.peerToPeripheralUUID[peerID] return (state.hasPeripheral, state.hasCentral)
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() }
}
} }
private func links(to peerID: PeerID?) -> Set<BLEIngressLinkID> { private func links(to peerID: PeerID?) -> Set<BLEIngressLinkID> {
guard let peerID else { return [] } readLinkState { $0.links(to: peerID) }
let computeLinks = { () -> Set<BLEIngressLinkID> in
var links: Set<BLEIngressLinkID> = []
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() }
}
} }
private func configureNoiseServiceCallbacks(for service: NoiseEncryptionService) { private func configureNoiseServiceCallbacks(for service: NoiseEncryptionService) {
@@ -2752,19 +2567,24 @@ extension BLEService {
// MARK: Link capability snapshots (thread-safe via bleQueue) // MARK: Link capability snapshots (thread-safe via bleQueue)
private func snapshotPeripheralStates() -> [PeripheralState] { private func readLinkState<T>(_ body: (BLELinkStateStore) -> T) -> T {
if DispatchQueue.getSpecific(key: bleQueueKey) != nil { if DispatchQueue.getSpecific(key: bleQueueKey) != nil {
return Array(peripherals.values) return body(linkStateStore)
} else { } else {
return bleQueue.sync { Array(peripherals.values) } return bleQueue.sync { body(linkStateStore) }
} }
} }
private func snapshotSubscribedCentrals() -> ([CBCentral], [String: PeerID]) {
if DispatchQueue.getSpecific(key: bleQueueKey) != nil { private func snapshotDirectPeripheralState(for peerID: PeerID) -> BLEPeripheralLinkState? {
return (self.subscribedCentrals, self.centralToPeerID) readLinkState { $0.directPeripheralState(for: peerID) }
} else { }
return bleQueue.sync { (self.subscribedCentrals, self.centralToPeerID) }
} private func snapshotPeripheralStates() -> [BLEPeripheralLinkState] {
readLinkState(\.peripheralStates)
}
private func snapshotSubscribedCentrals() -> BLESubscribedCentralSnapshot {
readLinkState(\.subscribedCentralSnapshot)
} }
// MARK: Helpers: IDs, selection, and write backpressure // MARK: Helpers: IDs, selection, and write backpressure
@@ -2802,7 +2622,7 @@ extension BLEService {
let uuid = peripheral.identifier.uuidString let uuid = peripheral.identifier.uuidString
bleQueue.async { [weak self] in bleQueue.async { [weak self] in
guard let self = self else { return } 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 // Atomically take all pending items from the queue to avoid race conditions
// where new items could be enqueued between read and update // where new items could be enqueued between read and update
@@ -2841,7 +2661,7 @@ extension BLEService {
private func drainAllPendingWrites() { private func drainAllPendingWrites() {
let uuids = collectionsQueue.sync { pendingPeripheralWrites.peripheralIDs } let uuids = collectionsQueue.sync { pendingPeripheralWrites.peripheralIDs }
for uuid in uuids { 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) drainPendingWrites(for: state.peripheral)
} }
} }
@@ -4022,7 +3842,7 @@ extension BLEService {
// Clean old connection timeout backoff entries (> window) // Clean old connection timeout backoff entries (> window)
let timeoutCutoff = now.addingTimeInterval(-TransportConfig.bleConnectTimeoutBackoffWindowSeconds) 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) // Clean up stale scheduled relays that somehow persisted (> 2s)
collectionsQueue.async(flags: .barrier) { [weak self] in collectionsQueue.async(flags: .barrier) { [weak self] in
@@ -4117,32 +3937,10 @@ extension BLEService {
} }
private func updateRSSIThreshold(connectedCount: Int) { private func updateRSSIThreshold(connectedCount: Int) {
// Adjust RSSI threshold based on connectivity, candidate pressure, and failures connectionScheduler.updateRSSIThreshold(
if connectedCount == 0 { connectedCount: connectedCount,
// Isolated: relax floor slowly to hunt for distant nodes connectedOrConnectingLinkCount: linkStateStore.connectedOrConnectingPeripheralCount,
if lastIsolatedAt == nil { lastIsolatedAt = Date() } now: 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
} }
} }
@@ -0,0 +1,156 @@
import Foundation
import Testing
@testable import bitchat
struct BLEConnectionSchedulerTests {
@Test
func discoveryQueuesWeakSignalCandidates() {
let scheduler = BLEConnectionScheduler<String>(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<String>(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<String>()
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<String>()
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<String>(
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<String>()
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<String>(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<String> {
BLEConnectionCandidate(
peripheral: id,
peripheralID: id,
rssi: rssi,
name: id,
isConnectable: true,
discoveredAt: now
)
}
@@ -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)
}
}