[codex] Refine BLE ingress fanout (#1280)

* 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

* Allow self-authored RSR ingress replies

* Harden read receipt queue test timing

---------

Co-authored-by: jack <jackjackbits@users.noreply.github.com>
This commit is contained in:
jack
2026-05-31 14:08:30 +02:00
committed by GitHub
co-authored by jack
parent ab0da61533
commit df36b19afe
18 changed files with 1503 additions and 383 deletions
+263 -377
View File
@@ -3,7 +3,6 @@ import BitFoundation
import Foundation
import CoreBluetooth
import Combine
import CryptoKit
#if os(iOS)
import UIKit
#endif
@@ -88,9 +87,7 @@ final class BLEService: NSObject {
private let meshTopology = MeshTopologyTracker()
// 5. Fragment Reassembly (necessary for messages > MTU)
private struct FragmentKey: Hashable { let sender: UInt64; let id: UInt64 }
private var incomingFragments: [FragmentKey: [Int: Data]] = [:]
private var fragmentMetadata: [FragmentKey: (type: UInt8, total: Int, timestamp: Date)] = [:]
private var fragmentAssemblyBuffer = BLEFragmentAssemblyBuffer()
private struct ActiveTransferState {
let totalFragments: Int
var sentFragments: Int
@@ -139,21 +136,17 @@ final class BLEService: NSObject {
// Noise typed payloads (ACKs, read receipts, etc.) pending handshake
private var pendingNoisePayloadsAfterHandshake: [PeerID: [Data]] = [:]
// Queue for notifications that failed due to full queue
private var pendingNotifications: [(data: Data, centrals: [CBCentral]?)] = []
private var pendingNotifications = BLEOutboundNotificationBuffer<CBCentral>()
// Accumulate long write chunks per central until a full frame decodes
private var pendingWriteBuffers: [String: Data] = [:]
private var pendingWriteBuffers = BLEInboundWriteBuffer()
// Relay jitter scheduling to reduce redundant floods
private var scheduledRelays: [String: DispatchWorkItem] = [:]
// Track short-lived traffic bursts to adapt announces/scanning under load
private var recentPacketTimestamps: [Date] = []
// Ingress link tracking for last-hop suppression
private enum LinkID: Hashable {
case peripheral(String)
case central(String)
}
private var ingressByMessageID: [String: (link: LinkID, timestamp: Date)] = [:]
// Ingress link tracking for duplicate and last-hop suppression
private var ingressLinks = BLEIngressLinkRegistry()
private struct PendingFragmentTransfer {
let packet: BitchatPacket
@@ -359,8 +352,9 @@ final class BLEService: NSObject {
pendingPeripheralWrites.removeAll()
pendingFragmentTransfers.removeAll()
pendingNotifications.removeAll()
fragmentAssemblyBuffer.removeAll()
pendingDirectedRelays.removeAll()
ingressByMessageID.removeAll()
ingressLinks.removeAll()
recentPacketTimestamps.removeAll()
scheduledRelays.values.forEach { $0.cancel() }
scheduledRelays.removeAll()
@@ -576,8 +570,7 @@ final class BLEService: NSObject {
let cancelledTransfers: [(id: String, items: [DispatchWorkItem])] = collectionsQueue.sync(flags: .barrier) {
let entries = activeTransfers.map { ($0.key, $0.value.workItems) }
peers.removeAll()
incomingFragments.removeAll()
fragmentMetadata.removeAll()
fragmentAssemblyBuffer.removeAll()
activeTransfers.removeAll()
// Also clear pending message queues to avoid stale state across sessions
pendingMessagesAfterHandshake.removeAll()
@@ -819,42 +812,28 @@ final class BLEService: NSObject {
case unknown
}
private struct IngressPacketContext {
let receivedFromPeerID: PeerID
let validationPeerID: PeerID
}
private func requiresDirectSenderBinding(_ packet: BitchatPacket) -> Bool {
packet.type == MessageType.announce.rawValue && packet.ttl == messageTTL
}
private func isSelfAuthoredSyncResponse(_ packet: BitchatPacket) -> Bool {
packet.isRSR && packet.ttl == 0
}
private func makeIngressPacketContext(
for packet: BitchatPacket,
claimedSenderID: PeerID,
boundPeerID: PeerID?,
linkDescription: String
) -> IngressPacketContext? {
if claimedSenderID == myPeerID,
!isSelfAuthoredSyncResponse(packet) {
) -> BLEIngressPacketContext? {
switch BLEIngressLinkRegistry.packetContext(
for: packet,
claimedSenderID: claimedSenderID,
boundPeerID: boundPeerID,
localPeerID: myPeerID,
directAnnounceTTL: messageTTL
) {
case .success(let context):
return context
case .failure(.selfLoopback):
SecureLogger.debug("↩️ Dropping BLE self-loopback packet type \(packet.type) from \(linkDescription)", category: .session)
return nil
}
if let boundPeerID, boundPeerID != claimedSenderID, requiresDirectSenderBinding(packet) {
case .failure(.directSenderMismatch(let boundPeerID, let claimedSenderID)):
SecureLogger.warning("🚫 SECURITY: Sender ID spoofing attempt detected! \(linkDescription) claimed to be \(claimedSenderID.id.prefix(8))… but is bound to \(boundPeerID.id.prefix(8))", category: .security)
return nil
}
let receivedFromPeerID = boundPeerID ?? claimedSenderID
let validationPeerID = packet.isRSR ? receivedFromPeerID : claimedSenderID
return IngressPacketContext(
receivedFromPeerID: receivedFromPeerID,
validationPeerID: validationPeerID
)
}
private func validatePacket(_ packet: BitchatPacket, from peerID: PeerID, connectionSource: ConnectionSource = .unknown) -> Bool {
@@ -887,18 +866,14 @@ final class BLEService: NSObject {
return true
}
private func recordIngressIfNew(_ packet: BitchatPacket, link: LinkID) -> Bool {
let messageID = makeMessageID(for: packet)
let now = Date()
private func recordIngressIfNew(_ packet: BitchatPacket, link: BLEIngressLinkID, peerID: PeerID) -> Bool {
return collectionsQueue.sync(flags: .barrier) {
if let existing = ingressByMessageID[messageID],
now.timeIntervalSince(existing.timestamp) <= TransportConfig.bleIngressRecordLifetimeSeconds {
return false
}
ingressByMessageID[messageID] = (link, now)
return true
ingressLinks.recordIfNew(
packet,
link: link,
peerID: peerID,
lifetime: TransportConfig.bleIngressRecordLifetimeSeconds
)
}
}
@@ -1005,9 +980,14 @@ final class BLEService: NSObject {
private func enqueuePendingNotification(data: Data, centrals: [CBCentral]?, context: String, attempt: Int = 0) {
collectionsQueue.async(flags: .barrier) { [weak self] in
guard let self = self else { return }
if self.pendingNotifications.count < TransportConfig.blePendingNotificationsCapCount {
self.pendingNotifications.append((data: data, centrals: centrals))
SecureLogger.debug("📋 Queued \(context) packet for retry (pending=\(self.pendingNotifications.count))", category: .session)
let result = self.pendingNotifications.enqueue(
data: data,
targets: centrals,
capCount: TransportConfig.blePendingNotificationsCapCount
)
if case let .enqueued(count) = result {
SecureLogger.debug("📋 Queued \(context) packet for retry (pending=\(count))", category: .session)
return
}
@@ -1025,9 +1005,10 @@ final class BLEService: NSObject {
}
private func sendOnAllLinks(packet: BitchatPacket, data: Data, pad: Bool, directedOnlyPeer: PeerID?) {
// Determine last-hop link for this message to avoid echoing back
// Determine the last-hop peer/link for this message to avoid echoing back over either BLE role.
let messageID = makeMessageID(for: packet)
let ingressLink: LinkID? = collectionsQueue.sync { ingressByMessageID[messageID]?.link }
let ingressRecord = collectionsQueue.sync { ingressLinks.record(for: packet) }
let excludedPeerLinks = links(to: ingressRecord?.peerID)
let directedPeerHint: PeerID? = {
if let explicit = directedOnlyPeer { return explicit }
if let recipient = PeerID(str: packet.recipientID?.hexEncodedString()), !recipient.isEmpty {
@@ -1073,35 +1054,19 @@ final class BLEService: NSObject {
subscribedCentrals = []
}
// Exclude ingress link
var allowedPeripheralIDs = connectedPeripheralIDs
var allowedCentralIDs = centralIDs
if let ingress = ingressLink {
switch ingress {
case .peripheral(let id):
allowedPeripheralIDs.removeAll { $0 == id }
case .central(let id):
allowedCentralIDs.removeAll { $0 == id }
}
}
// For broadcast (no directed peer) and non-fragment, choose a subset deterministically
// Special-case control/presence messages: do NOT subset to maximize immediate coverage
var selectedPeripheralIDs = Set(allowedPeripheralIDs)
var selectedCentralIDs = Set(allowedCentralIDs)
if directedPeerHint == nil
&& packet.type != MessageType.fragment.rawValue
&& packet.type != MessageType.announce.rawValue
&& packet.type != MessageType.requestSync.rawValue {
let kp = subsetSizeForFanout(allowedPeripheralIDs.count)
let kc = subsetSizeForFanout(allowedCentralIDs.count)
selectedPeripheralIDs = selectDeterministicSubset(ids: allowedPeripheralIDs, k: kp, seed: messageID)
selectedCentralIDs = selectDeterministicSubset(ids: allowedCentralIDs, k: kc, seed: messageID)
}
let selectedLinks = BLEFanoutSelector.selectLinks(
peripheralIDs: connectedPeripheralIDs,
centralIDs: centralIDs,
ingressLink: ingressRecord?.link,
excludedLinks: excludedPeerLinks,
directedPeerHint: directedPeerHint,
packetType: packet.type,
messageID: messageID
)
// If directed and we currently have no links to forward on, spool for a short window
if let only = directedPeerHint,
selectedPeripheralIDs.isEmpty && selectedCentralIDs.isEmpty,
selectedLinks.peripheralIDs.isEmpty && selectedLinks.centralIDs.isEmpty,
(packet.type == MessageType.noiseEncrypted.rawValue || packet.type == MessageType.noiseHandshake.rawValue) {
spoolDirectedPacket(packet, recipientPeerID: only)
}
@@ -1109,14 +1074,14 @@ final class BLEService: NSObject {
// Writes to selected connected peripherals
for s in states where s.isConnected {
let pid = s.peripheral.identifier.uuidString
guard selectedPeripheralIDs.contains(pid) else { continue }
guard selectedLinks.peripheralIDs.contains(pid) else { continue }
if let ch = s.characteristic {
writeOrEnqueue(data, to: s.peripheral, characteristic: ch, priority: outboundPriority)
}
}
// Notify selected subscribed centrals
if let ch = characteristic {
let targets = subscribedCentrals.filter { selectedCentralIDs.contains($0.identifier.uuidString) }
let targets = subscribedCentrals.filter { selectedLinks.centralIDs.contains($0.identifier.uuidString) }
if !targets.isEmpty {
let success = peripheralManager?.updateValue(data, for: ch, onSubscribedCentrals: targets) ?? false
if !success {
@@ -2055,6 +2020,15 @@ extension BLEService {
}
}
private extension BLEService {
static func shouldRediscoverBitChatService(
invalidatedServiceUUIDs: [CBUUID],
cachedServiceUUIDs: [CBUUID]?
) -> Bool {
invalidatedServiceUUIDs.contains(serviceUUID) || cachedServiceUUIDs?.contains(serviceUUID) != true
}
}
#if DEBUG
// Test-only helper to inject packets into the receive pipeline
extension BLEService {
@@ -2096,7 +2070,17 @@ extension BLEService {
}
func _test_recordIngressIfNew(packet: BitchatPacket, linkID: String) -> Bool {
recordIngressIfNew(packet, link: .central(linkID))
recordIngressIfNew(packet, link: .central(linkID), peerID: PeerID(hexData: packet.senderID))
}
static func _test_shouldRediscoverBitChatService(
invalidatedServiceUUIDs: [CBUUID],
cachedServiceUUIDs: [CBUUID]?
) -> Bool {
shouldRediscoverBitChatService(
invalidatedServiceUUIDs: invalidatedServiceUUIDs,
cachedServiceUUIDs: cachedServiceUUIDs
)
}
}
#endif
@@ -2257,7 +2241,7 @@ extension BLEService: CBPeripheralDelegate {
peripherals[peripheralUUID] = state
}
if !recordIngressIfNew(packet, link: .peripheral(peripheralUUID)) {
if !recordIngressIfNew(packet, link: .peripheral(peripheralUUID), peerID: context.receivedFromPeerID) {
continue
}
processNotificationPacket(
@@ -2309,18 +2293,23 @@ extension BLEService: CBPeripheralDelegate {
func peripheral(_ peripheral: CBPeripheral, didModifyServices invalidatedServices: [CBService]) {
SecureLogger.warning("⚠️ Services modified for \(peripheral.name ?? peripheral.identifier.uuidString)", category: .session)
// Check if our service was invalidated (peer app quit)
let hasOurService = peripheral.services?.contains { $0.uuid == BLEService.serviceUUID } ?? false
if !hasOurService {
// Service is gone - disconnect
SecureLogger.warning("❌ BitChat service removed - disconnecting from \(peripheral.name ?? peripheral.identifier.uuidString)", category: .session)
centralManager?.cancelPeripheralConnection(peripheral)
} else {
// Try to rediscover
peripheral.discoverServices([BLEService.serviceUUID])
let shouldRediscover = BLEService.shouldRediscoverBitChatService(
invalidatedServiceUUIDs: invalidatedServices.map(\.uuid),
cachedServiceUUIDs: peripheral.services?.map(\.uuid)
)
guard shouldRediscover else { return }
let peripheralID = peripheral.identifier.uuidString
if var state = peripherals[peripheralID] {
state.characteristic = nil
state.assembler = NotificationStreamAssembler()
peripherals[peripheralID] = state
}
SecureLogger.debug("🔄 BitChat service changed for \(peripheral.name ?? peripheral.identifier.uuidString), rediscovering", category: .session)
peripheral.discoverServices([BLEService.serviceUUID])
}
func peripheral(_ peripheral: CBPeripheral, didUpdateNotificationStateFor characteristic: CBCharacteristic, error: Error?) {
@@ -2571,55 +2560,51 @@ extension BLEService: CBPeripheralManagerDelegate {
func peripheralManagerIsReady(toUpdateSubscribers peripheral: CBPeripheralManager) {
SecureLogger.debug("📤 Peripheral manager ready to send more notifications", category: .session)
// Retry pending notifications now that queue has space
drainPendingNotifications(logPrefix: "✅ Sent")
}
private func drainPendingNotifications(logPrefix: String) {
collectionsQueue.async(flags: .barrier) { [weak self] in
guard let self = self,
let characteristic = self.characteristic,
!self.pendingNotifications.isEmpty else { return }
let pending = self.pendingNotifications
self.pendingNotifications.removeAll()
// Try to send pending notifications
var sentCount = 0
for (index, (data, centrals)) in pending.enumerated() {
if let centrals = centrals {
// Send to specific centrals
let success = self.peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: centrals) ?? false
if !success {
// Still full, re-queue this and all remaining items
let remaining = pending.dropFirst(index)
self.pendingNotifications.append(contentsOf: remaining)
SecureLogger.debug("⚠️ Notification queue still full after \(sentCount) sent, re-queuing \(remaining.count) items", category: .session)
break // Stop trying, wait for next ready callback
} else {
sentCount += 1
}
} else {
// Broadcast to all
let success = self.peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: nil) ?? false
if !success {
// Still full, re-queue this and all remaining items
let remaining = pending.dropFirst(index)
self.pendingNotifications.append(contentsOf: remaining)
SecureLogger.debug("⚠️ Notification queue still full after \(sentCount) sent, re-queuing \(remaining.count) items", category: .session)
break
} else {
sentCount += 1
}
}
}
let pending = self.pendingNotifications.takeAll()
let sentCount = self.sendPendingNotifications(pending, characteristic: characteristic)
if sentCount > 0 {
SecureLogger.debug("✅ Sent \(sentCount) pending notifications from retry queue", category: .session)
SecureLogger.debug("\(logPrefix) \(sentCount) pending notifications from retry queue", category: .session)
}
if !self.pendingNotifications.isEmpty {
SecureLogger.debug("📋 Still have \(self.pendingNotifications.count) pending notifications", category: .session)
}
}
}
private func sendPendingNotifications(_ pending: [BLEPendingNotification<CBCentral>], characteristic: CBMutableCharacteristic) -> Int {
var sentCount = 0
for (index, notification) in pending.enumerated() {
let success = peripheralManager?.updateValue(
notification.data,
for: characteristic,
onSubscribedCentrals: notification.targets
) ?? false
guard success else {
let remaining = Array(pending.dropFirst(index))
pendingNotifications.prepend(remaining)
SecureLogger.debug("⚠️ Notification queue still full after \(sentCount) sent, re-queuing \(remaining.count) items", category: .session)
break
}
sentCount += 1
}
return sentCount
}
func peripheralManager(_ peripheral: CBPeripheralManager, didReceiveWrite requests: [CBATTRequest]) {
// Suppress logs for single write requests to reduce noise
@@ -2641,85 +2626,85 @@ extension BLEService: CBPeripheralManagerDelegate {
// Sort by offset ascending
let sorted = group.sorted { $0.offset < $1.offset }
let hasMultiple = sorted.count > 1 || (sorted.first?.offset ?? 0) > 0
// Always merge into a persistent per-central buffer to handle multi-callback long writes
var combined = pendingWriteBuffers[centralUUID] ?? Data()
var appendedBytes = 0
var offsets: [Int] = []
for r in sorted {
guard let chunk = r.value, !chunk.isEmpty else { continue }
offsets.append(r.offset)
let end = r.offset + chunk.count
if combined.count < end {
combined.append(Data(repeating: 0, count: end - combined.count))
}
// Write chunk into the correct position (supports out-of-order and overlapping writes)
combined.replaceSubrange(r.offset..<end, with: chunk)
appendedBytes += chunk.count
}
pendingWriteBuffers[centralUUID] = combined
// Peek type byte for debug: version is at 0, type at 1 when well-formed
if combined.count >= 2 {
let peekType = combined[1]
if peekType != MessageType.announce.rawValue {
SecureLogger.debug("📥 Accumulated write from central \(centralUUID): size=\(combined.count) (+\(appendedBytes)) bytes (type=\(peekType)), offsets=\(offsets)", category: .session)
}
let chunks = sorted.compactMap { request -> BLEInboundWriteChunk? in
guard let data = request.value, !data.isEmpty else { return nil }
return BLEInboundWriteChunk(offset: request.offset, data: data)
}
// Try decode the accumulated buffer
if let packet = BinaryProtocol.decode(combined) {
// Clear buffer on success
pendingWriteBuffers.removeValue(forKey: centralUUID)
let result = pendingWriteBuffers.append(
chunks: chunks,
for: centralUUID,
capBytes: TransportConfig.blePendingWriteBufferCapBytes
)
let claimedSenderID = PeerID(hexData: packet.senderID)
let context = makeIngressPacketContext(
for: packet,
claimedSenderID: claimedSenderID,
boundPeerID: centralToPeerID[centralUUID],
linkDescription: "Central \(centralUUID.prefix(8))"
)
guard let context else { continue }
switch result {
case let .decoded(packet, metadata):
logAccumulatedCentralWrite(metadata, centralUUID: centralUUID)
processDecodedCentralWrite(packet, centralUUID: centralUUID, central: sorted[0].central)
if !validatePacket(packet, from: context.validationPeerID, connectionSource: .central(centralUUID)) {
continue
}
case let .waiting(metadata):
logAccumulatedCentralWrite(metadata, centralUUID: centralUUID)
logFailedSingleWriteIfNeeded(hasMultiple: hasMultiple, sortedRequests: sorted)
if packet.type != MessageType.announce.rawValue {
SecureLogger.debug("📦 Decoded (combined) packet type: \(packet.type) from sender: \(claimedSenderID)", category: .session)
}
if !subscribedCentrals.contains(sorted[0].central) {
subscribedCentrals.append(sorted[0].central)
}
if packet.type == MessageType.announce.rawValue {
if packet.ttl == messageTTL {
centralToPeerID[centralUUID] = claimedSenderID
refreshLocalTopology()
}
if !recordIngressIfNew(packet, link: .central(centralUUID)) {
continue
}
handleReceivedPacket(packet, from: context.receivedFromPeerID)
} else {
if !recordIngressIfNew(packet, link: .central(centralUUID)) {
continue
}
handleReceivedPacket(packet, from: context.receivedFromPeerID)
}
} else {
// If buffer grows suspiciously large, reset to avoid memory leak
if combined.count > TransportConfig.blePendingWriteBufferCapBytes { // cap for safety
pendingWriteBuffers.removeValue(forKey: centralUUID)
SecureLogger.warning("⚠️ Dropping oversized pending write buffer (\(combined.count) bytes) for central \(centralUUID)", category: .session)
}
// If this was a single short write and still failed, log the raw chunk for debugging
if !hasMultiple, let only = sorted.first, let raw = only.value {
let prefix = raw.prefix(16).map { String(format: "%02x", $0) }.joined(separator: " ")
SecureLogger.error("❌ Failed to decode packet from central (len=\(raw.count), prefix=\(prefix))", category: .session)
}
case let .oversized(metadata):
logAccumulatedCentralWrite(metadata, centralUUID: centralUUID)
SecureLogger.warning("⚠️ Dropping oversized pending write buffer (\(metadata.accumulatedBytes) bytes) for central \(centralUUID)", category: .session)
logFailedSingleWriteIfNeeded(hasMultiple: hasMultiple, sortedRequests: sorted)
}
}
}
}
private func logAccumulatedCentralWrite(_ metadata: BLEInboundWriteAppendMetadata, centralUUID: String) {
guard let packetType = metadata.packetType,
packetType != MessageType.announce.rawValue else { return }
SecureLogger.debug(
"📥 Accumulated write from central \(centralUUID): size=\(metadata.accumulatedBytes) (+\(metadata.appendedBytes)) bytes (type=\(packetType)), offsets=\(metadata.offsets)",
category: .session
)
}
private func logFailedSingleWriteIfNeeded(hasMultiple: Bool, sortedRequests: [CBATTRequest]) {
guard !hasMultiple, let raw = sortedRequests.first?.value else { return }
let prefix = raw.prefix(16).map { String(format: "%02x", $0) }.joined(separator: " ")
SecureLogger.error("❌ Failed to decode packet from central (len=\(raw.count), prefix=\(prefix))", category: .session)
}
private func processDecodedCentralWrite(_ packet: BitchatPacket, centralUUID: String, central: CBCentral) {
let claimedSenderID = PeerID(hexData: packet.senderID)
let context = makeIngressPacketContext(
for: packet,
claimedSenderID: claimedSenderID,
boundPeerID: centralToPeerID[centralUUID],
linkDescription: "Central \(centralUUID.prefix(8))"
)
guard let context else { return }
guard validatePacket(packet, from: context.validationPeerID, connectionSource: .central(centralUUID)) else {
return
}
if packet.type != MessageType.announce.rawValue {
SecureLogger.debug("📦 Decoded (combined) packet type: \(packet.type) from sender: \(claimedSenderID)", category: .session)
}
if !subscribedCentrals.contains(central) {
subscribedCentrals.append(central)
}
if packet.type == MessageType.announce.rawValue,
packet.ttl == messageTTL {
centralToPeerID[centralUUID] = claimedSenderID
refreshLocalTopology()
}
guard recordIngressIfNew(packet, link: .central(centralUUID), peerID: context.receivedFromPeerID) else {
return
}
handleReceivedPacket(packet, from: context.receivedFromPeerID)
}
}
// MARK: - Advertising Builders & Alias Rotation
@@ -2860,6 +2845,10 @@ extension BLEService {
}
private func forwardAlongRouteIfNeeded(_ packet: BitchatPacket) -> Bool {
if PeerID(hexData: packet.recipientID) == myPeerID {
return true
}
guard let route = packet.route, !route.isEmpty else { return false }
let myRoutingData = routingData(for: myPeerID) ?? (myPeerIDData.isEmpty ? nil : myPeerIDData)
guard let selfData = myRoutingData else { return false }
@@ -2923,6 +2912,27 @@ extension BLEService {
return bleQueue.sync { computeState() }
}
}
private func links(to peerID: PeerID?) -> Set<BLEIngressLinkID> {
guard let peerID else { return [] }
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) {
service.onPeerAuthenticated = { [weak self] peerID, fingerprint in
@@ -2997,37 +3007,7 @@ extension BLEService {
// MARK: Helpers: IDs, selection, and write backpressure
private func makeMessageID(for packet: BitchatPacket) -> String {
let senderID = packet.senderID.hexEncodedString()
let digestPrefix = packet.payload.sha256Hash().prefix(4).hexEncodedString()
return "\(senderID)-\(packet.timestamp)-\(packet.type)-\(digestPrefix)"
}
private func subsetSizeForFanout(_ n: Int) -> Int {
guard n > 0 else { return 0 }
if n <= 2 { return n }
// approx ceil(log2(n)) + 1 without floating point
var v = n - 1
var bits = 0
while v > 0 { v >>= 1; bits += 1 }
return min(n, max(1, bits + 1))
}
private func selectDeterministicSubset(ids: [String], k: Int, seed: String) -> Set<String> {
guard k > 0 && ids.count > k else { return Set(ids) }
// Stable order by SHA256(seed || "::" || id)
var scored: [(score: [UInt8], id: String)] = []
for id in ids {
let msg = (seed + "::" + id).data(using: .utf8) ?? Data()
let digest = Array(SHA256.hash(data: msg))
scored.append((digest, id))
}
scored.sort { a, b in
for i in 0..<min(a.score.count, b.score.count) {
if a.score[i] != b.score[i] { return a.score[i] < b.score[i] }
}
return a.id < b.id
}
return Set(scored.prefix(k).map { $0.id })
BLEIngressLinkRegistry.messageID(for: packet)
}
private func priority(for packet: BitchatPacket, data: Data) -> BLEOutboundWritePriority {
@@ -3116,37 +3096,7 @@ extension BLEService {
/// Periodically try to drain pending notifications as a backup mechanism
private func drainPendingNotificationsIfPossible() {
collectionsQueue.async(flags: .barrier) { [weak self] in
guard let self = self,
let characteristic = self.characteristic,
!self.pendingNotifications.isEmpty else { return }
let pending = self.pendingNotifications
self.pendingNotifications.removeAll()
var sentCount = 0
for (index, (data, centrals)) in pending.enumerated() {
let success: Bool
if let centrals = centrals {
success = self.peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: centrals) ?? false
} else {
success = self.peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: nil) ?? false
}
if !success {
// Re-queue this and all remaining items
let remaining = pending.dropFirst(index)
self.pendingNotifications.append(contentsOf: remaining)
break
} else {
sentCount += 1
}
}
if sentCount > 0 {
SecureLogger.debug("🔄 Periodic drain: sent \(sentCount) pending notifications", category: .session)
}
}
drainPendingNotifications(logPrefix: "🔄 Periodic drain: sent")
}
/// Periodically try to drain pending writes for all connected peripherals
@@ -3601,102 +3551,19 @@ extension BLEService {
return
}
// Minimum header: 8 bytes ID + 2 index + 2 total + 1 type
guard packet.payload.count >= 13 else { return }
guard let header = BLEFragmentHeader(packet: packet) else { return }
// Compute compact fragment key (sender: 8 bytes, id: 8 bytes), big-endian
var senderU64: UInt64 = 0
for b in packet.senderID.prefix(8) { senderU64 = (senderU64 << 8) | UInt64(b) }
var fragU64: UInt64 = 0
for b in packet.payload.prefix(8) { fragU64 = (fragU64 << 8) | UInt64(b) }
// Parse big-endian UInt16 safely without alignment assumptions
let idxHi = UInt16(packet.payload[8])
let idxLo = UInt16(packet.payload[9])
let index = Int((idxHi << 8) | idxLo)
let totHi = UInt16(packet.payload[10])
let totLo = UInt16(packet.payload[11])
let total = Int((totHi << 8) | totLo)
let originalType = packet.payload[12]
let fragmentData = packet.payload.suffix(from: 13)
// Sanity checks - add reasonable upper bound on total to prevent DoS
guard total > 0 && total <= 10000 && index >= 0 && index < total else { return }
let isBroadcastFragment: Bool = {
guard let recipient = packet.recipientID else { return true }
return recipient.count == 8 && recipient.allSatisfy { $0 == 0xFF }
}()
if isBroadcastFragment {
if header.isBroadcastFragment {
gossipSyncManager?.onPublicPacketSeen(packet)
}
// Compute fragment key for this assembly
let key = FragmentKey(sender: senderU64, id: fragU64)
// Critical section: Store fragment and check completion status
var shouldReassemble: Bool = false
var fragmentsToReassemble: [Int: Data]? = nil
collectionsQueue.sync(flags: .barrier) {
if incomingFragments[key] == nil {
// Cap in-flight assemblies to prevent memory/battery blowups
if incomingFragments.count >= maxInFlightAssemblies {
// Evict the oldest assembly by timestamp
if let oldest = fragmentMetadata.min(by: { $0.value.timestamp < $1.value.timestamp })?.key {
incomingFragments.removeValue(forKey: oldest)
fragmentMetadata.removeValue(forKey: oldest)
}
}
incomingFragments[key] = [:]
fragmentMetadata[key] = (originalType, total, Date())
SecureLogger.debug("📦 Started fragment assembly id=\(String(format: "%016llx", fragU64)) total=\(total)", category: .session)
}
// Check cumulative size before storing this fragment
let currentSize = incomingFragments[key]?.values.reduce(0) { $0 + $1.count } ?? 0
let assemblyLimit: Int = {
if originalType == MessageType.fileTransfer.rawValue {
// Allow headroom for TLV metadata and binary framing overhead.
return FileTransferLimits.maxFramedFileBytes
}
return FileTransferLimits.maxPayloadBytes
}()
let projectedSize = currentSize + fragmentData.count
guard projectedSize <= assemblyLimit else {
// Exceeds size limit - evict this assembly
SecureLogger.warning(
"🚫 Fragment assembly exceeds size limit (\(projectedSize) bytes > \(assemblyLimit)), evicting. Type=\(originalType) Index=\(index)/\(total)",
category: .security
)
incomingFragments.removeValue(forKey: key)
fragmentMetadata.removeValue(forKey: key)
shouldReassemble = false
fragmentsToReassemble = nil
return
}
incomingFragments[key]?[index] = Data(fragmentData)
SecureLogger.debug("📦 Fragment \(index + 1)/\(total) (len=\(fragmentData.count)) for id=\(String(format: "%016llx", fragU64))", category: .session)
// Check if complete
if let fragments = incomingFragments[key], fragments.count == total {
shouldReassemble = true
fragmentsToReassemble = fragments
} else {
shouldReassemble = false
fragmentsToReassemble = nil
}
let assemblyResult = collectionsQueue.sync(flags: .barrier) {
fragmentAssemblyBuffer.append(header, maxInFlightAssemblies: maxInFlightAssemblies)
}
// Heavy work outside lock: reassemble and decode
guard shouldReassemble, let fragments = fragmentsToReassemble else { return }
logFragmentAssemblyResult(assemblyResult)
var reassembled = Data()
for i in 0..<total {
if let fragment = fragments[i] {
reassembled.append(fragment)
}
}
guard case let .complete(completedHeader, reassembled, _) = assemblyResult else { return }
// Decode the original packet bytes we reassembled, so flags/compression are preserved
if var originalPacket = BinaryProtocol.decode(reassembled) {
@@ -3706,18 +3573,37 @@ extension BLEService {
if !validatePacket(originalPacket, from: innerSender) {
// Cleanup below
} else {
SecureLogger.debug("✅ Reassembled packet id=\(String(format: "%016llx", fragU64)) type=\(originalPacket.type) bytes=\(reassembled.count)", category: .session)
SecureLogger.debug("✅ Reassembled packet id=\(completedHeader.idLogString) type=\(originalPacket.type) bytes=\(reassembled.count)", category: .session)
originalPacket.ttl = 0
handleReceivedPacket(originalPacket, from: peerID)
}
} else {
SecureLogger.error("❌ Failed to decode reassembled packet (type=\(originalType), total=\(total))", category: .session)
SecureLogger.error("❌ Failed to decode reassembled packet (type=\(completedHeader.originalType), total=\(completedHeader.total))", category: .session)
}
}
private func logFragmentAssemblyResult(_ result: BLEFragmentAssemblyBuffer.AppendResult) {
func logStartedIfNeeded(header: BLEFragmentHeader, started: Bool) {
if started {
SecureLogger.debug("📦 Started fragment assembly id=\(header.idLogString) total=\(header.total)", category: .session)
}
}
// Critical section: Cleanup completed assembly
collectionsQueue.sync(flags: .barrier) {
incomingFragments.removeValue(forKey: key)
fragmentMetadata.removeValue(forKey: key)
switch result {
case let .stored(header, started):
logStartedIfNeeded(header: header, started: started)
SecureLogger.debug("📦 Fragment \(header.index + 1)/\(header.total) (len=\(header.fragmentData.count)) for id=\(header.idLogString)", category: .session)
case let .complete(header, _, started):
logStartedIfNeeded(header: header, started: started)
SecureLogger.debug("📦 Fragment \(header.index + 1)/\(header.total) (len=\(header.fragmentData.count)) for id=\(header.idLogString)", category: .session)
case let .oversized(header, projectedSize, limit, started):
logStartedIfNeeded(header: header, started: started)
SecureLogger.warning(
"🚫 Fragment assembly exceeds size limit (\(projectedSize) bytes > \(limit)), evicting. Type=\(header.originalType) Index=\(header.index)/\(header.total)",
category: .security
)
}
}
@@ -3827,6 +3713,7 @@ extension BLEService {
let decision = RelayController.decide(
ttl: packet.ttl,
senderIsSelf: senderID == myPeerID,
recipientIsSelf: PeerID(hexData: packet.recipientID) == myPeerID,
isEncrypted: packet.type == MessageType.noiseEncrypted.rawValue,
isDirectedEncrypted: (packet.type == MessageType.noiseEncrypted.rawValue) && (packet.recipientID != nil),
isFragment: packet.type == MessageType.fragment.rawValue,
@@ -4193,8 +4080,6 @@ extension BLEService {
}
private func handleNoiseEncrypted(_ packet: BitchatPacket, from peerID: PeerID) {
SecureLogger.debug("🔐 handleNoiseEncrypted called for packet from \(peerID)")
guard let recipientID = PeerID(hexData: packet.recipientID) else {
SecureLogger.warning("⚠️ Encrypted message has no recipient ID", category: .session)
return
@@ -4215,8 +4100,15 @@ extension BLEService {
// First byte indicates the payload type
let payloadType = decrypted[0]
let payloadData = decrypted.dropFirst()
switch NoisePayloadType(rawValue: payloadType) {
guard let noisePayloadType = NoisePayloadType(rawValue: payloadType) else {
SecureLogger.warning("⚠️ Unknown noise payload type: \(payloadType)")
return
}
SecureLogger.debug("🔐 Decrypted noise payload type \(noisePayloadType.description) from \(peerID)", category: .session)
switch noisePayloadType {
case .privateMessage:
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
notifyUI { [weak self] in
@@ -4242,8 +4134,6 @@ extension BLEService {
notifyUI { [weak self] in
self?.deliverTransportEvent(.noisePayloadReceived(peerID: peerID, type: .verifyResponse, payload: Data(payloadData), timestamp: ts))
}
case .none:
SecureLogger.warning("⚠️ Unknown noise payload type: \(payloadType)")
}
} catch NoiseEncryptionError.sessionNotEstablished {
// We received an encrypted message before establishing a session with this peer.
@@ -4486,11 +4376,7 @@ extension BLEService {
// Clean old fragments (> configured seconds old)
collectionsQueue.sync(flags: .barrier) {
let cutoff = now.addingTimeInterval(-TransportConfig.bleFragmentLifetimeSeconds)
let oldFragments = fragmentMetadata.filter { $0.value.timestamp < cutoff }.map { $0.key }
for fragmentID in oldFragments {
incomingFragments.removeValue(forKey: fragmentID)
fragmentMetadata.removeValue(forKey: fragmentID)
}
fragmentAssemblyBuffer.removeExpired(before: cutoff)
}
// Clean old connection timeout backoff entries (> window)
@@ -4512,8 +4398,8 @@ extension BLEService {
collectionsQueue.async(flags: .barrier) { [weak self] in
guard let self = self else { return }
let cutoff = now.addingTimeInterval(-TransportConfig.bleIngressRecordLifetimeSeconds)
if !self.ingressByMessageID.isEmpty {
self.ingressByMessageID = self.ingressByMessageID.filter { $0.value.timestamp >= cutoff }
if !self.ingressLinks.isEmpty {
self.ingressLinks.prune(before: cutoff)
}
// Clean expired directed spooled items
if !self.pendingDirectedRelays.isEmpty {