[codex] Refactor BLE transport event handling (#1266)

* Refactor BLE transport event handling

* Make image output paths unique

* Keep queued Nostr read receipts alive

* Allow self-authored RSR ingress replies

---------

Co-authored-by: jack <jackjackbits@users.noreply.github.com>
This commit is contained in:
jack
2026-05-31 13:58:27 +02:00
committed by GitHub
co-authored by jack
parent 3be8fbf1c4
commit ab0da61533
19 changed files with 575 additions and 152 deletions
+172 -134
View File
@@ -155,27 +155,6 @@ final class BLEService: NSObject {
}
private var ingressByMessageID: [String: (link: LinkID, timestamp: Date)] = [:]
// Backpressure-aware write queue per peripheral
private struct OutboundPriority: Comparable {
let level: Int
let suborder: Int
static let high = OutboundPriority(level: 0, suborder: 0)
static func fragment(totalFragments: Int) -> OutboundPriority {
OutboundPriority(level: 1, suborder: max(1, min(totalFragments, Int(UInt16.max))))
}
static let fileTransfer = OutboundPriority(level: 2, suborder: Int.max - 1)
static let low = OutboundPriority(level: 2, suborder: Int.max)
static func < (lhs: OutboundPriority, rhs: OutboundPriority) -> Bool {
if lhs.level != rhs.level { return lhs.level < rhs.level }
return lhs.suborder < rhs.suborder
}
}
private struct PendingWrite {
let priority: OutboundPriority
let data: Data
}
private struct PendingFragmentTransfer {
let packet: BitchatPacket
let pad: Bool
@@ -183,7 +162,7 @@ final class BLEService: NSObject {
let directedPeer: PeerID?
let transferId: String?
}
private var pendingPeripheralWrites: [String: [PendingWrite]] = [:]
private var pendingPeripheralWrites = BLEOutboundWriteBuffer()
private var pendingFragmentTransfers: [PendingFragmentTransfer] = []
// Debounce duplicate disconnect notifies
private var recentDisconnectNotifies: [PeerID: Date] = [:]
@@ -465,8 +444,9 @@ final class BLEService: NSObject {
// MARK: - Transport Protocol Conformance
// MARK: Delegates
weak var delegate: BitchatDelegate?
weak var eventDelegate: TransportEventDelegate?
weak var peerEventsDelegate: TransportPeerEventsDelegate?
// MARK: Peer snapshots publisher (non-UI convenience)
@@ -839,6 +819,44 @@ 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) {
SecureLogger.debug("↩️ Dropping BLE self-loopback packet type \(packet.type) from \(linkDescription)", category: .session)
return nil
}
if let boundPeerID, boundPeerID != claimedSenderID, requiresDirectSenderBinding(packet) {
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 {
let currentTime = UInt64(Date().timeIntervalSince1970 * 1000)
@@ -869,6 +887,21 @@ final class BLEService: NSObject {
return true
}
private func recordIngressIfNew(_ packet: BitchatPacket, link: LinkID) -> Bool {
let messageID = makeMessageID(for: packet)
let now = Date()
return collectionsQueue.sync(flags: .barrier) {
if let existing = ingressByMessageID[messageID],
now.timeIntervalSince(existing.timestamp) <= TransportConfig.bleIngressRecordLifetimeSeconds {
return false
}
ingressByMessageID[messageID] = (link, now)
return true
}
}
// MARK: - Packet Broadcasting
private func broadcastPacket(_ packet: BitchatPacket, transferId: String? = nil) {
@@ -1253,9 +1286,7 @@ final class BLEService: NSObject {
SecureLogger.debug("📁 Stored incoming media from \(peerID.id.prefix(8))… -> \(destination.lastPathComponent)", category: .session)
notifyUI { [weak self] in
self?.delegate?.didReceiveMessage(message)
}
emitTransportEvent(.messageReceived(message))
}
func sendFavoriteNotification(to peerID: PeerID, isFavorite: Bool) {
@@ -1324,8 +1355,8 @@ final class BLEService: NSObject {
// Get current peer list (after removal)
let currentPeerIDs = self.collectionsQueue.sync { Array(self.peers.keys) }
self.delegate?.didDisconnectFromPeer(peerID)
self.delegate?.didUpdatePeerList(currentPeerIDs)
self.deliverTransportEvent(.peerDisconnected(peerID))
self.deliverTransportEvent(.peerListUpdated(currentPeerIDs))
}
}
@@ -1647,10 +1678,7 @@ extension BLEService: CBCentralManagerDelegate {
#endif
func centralManagerDidUpdateState(_ central: CBCentralManager) {
// Notify delegate about state change on main thread
Task { @MainActor in
self.delegate?.didUpdateBluetoothState(central.state)
}
emitTransportEvent(.bluetoothStateUpdated(central.state))
switch central.state {
case .poweredOn:
@@ -1939,7 +1967,7 @@ func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeriph
self.notifyPeerDisconnectedDebounced(peerID)
}
self.requestPeerDataPublish()
self.delegate?.didUpdatePeerList(currentPeerIDs)
self.deliverTransportEvent(.peerListUpdated(currentPeerIDs))
}
}
@@ -2056,6 +2084,20 @@ extension BLEService {
}
handleReceivedPacket(packet, from: fromPeerID)
}
func _test_acceptsIngress(packet: BitchatPacket, boundPeerID: PeerID?) -> Bool {
let claimedSenderID = PeerID(hexData: packet.senderID)
return makeIngressPacketContext(
for: packet,
claimedSenderID: claimedSenderID,
boundPeerID: boundPeerID,
linkDescription: "TestLink"
) != nil
}
func _test_recordIngressIfNew(packet: BitchatPacket, linkID: String) -> Bool {
recordIngressIfNew(packet, link: .central(linkID))
}
}
#endif
@@ -2193,19 +2235,16 @@ extension BLEService: CBPeripheralDelegate {
}
let claimedSenderID = PeerID(hexData: packet.senderID)
let context = makeIngressPacketContext(
for: packet,
claimedSenderID: claimedSenderID,
boundPeerID: boundPeerID,
linkDescription: "Peripheral \(peripheralUUID.prefix(8))"
)
let trustedSenderID: PeerID?
if let knownPeerID = boundPeerID {
if knownPeerID != claimedSenderID {
SecureLogger.warning("🚫 SECURITY: Sender ID spoofing attempt detected! Peripheral \(peripheralUUID.prefix(8))… claimed to be \(claimedSenderID.id.prefix(8))… but is bound to \(knownPeerID.id.prefix(8))", category: .security)
continue
}
trustedSenderID = knownPeerID
} else {
trustedSenderID = nil
}
guard let context else { continue }
if !validatePacket(packet, from: trustedSenderID ?? claimedSenderID, connectionSource: .peripheral(peripheralUUID)) {
if !validatePacket(packet, from: context.validationPeerID, connectionSource: .peripheral(peripheralUUID)) {
continue
}
@@ -2217,11 +2256,20 @@ extension BLEService: CBPeripheralDelegate {
state.peerID = claimedSenderID
peripherals[peripheralUUID] = state
}
processNotificationPacket(packet, from: peripheral, peripheralUUID: peripheralUUID)
if !recordIngressIfNew(packet, link: .peripheral(peripheralUUID)) {
continue
}
processNotificationPacket(
packet,
from: peripheral,
peripheralUUID: peripheralUUID,
receivedFrom: context.receivedFromPeerID
)
}
}
private func processNotificationPacket(_ packet: BitchatPacket, from peripheral: CBPeripheral, peripheralUUID: String) {
private func processNotificationPacket(_ packet: BitchatPacket, from peripheral: CBPeripheral, peripheralUUID: String, receivedFrom peerID: PeerID) {
let senderID = PeerID(hexData: packet.senderID)
if packet.type != MessageType.announce.rawValue {
@@ -2238,17 +2286,9 @@ extension BLEService: CBPeripheralDelegate {
refreshLocalTopology()
}
let msgID = makeMessageID(for: packet)
collectionsQueue.async(flags: .barrier) { [weak self] in
self?.ingressByMessageID[msgID] = (.peripheral(peripheralUUID), Date())
}
handleReceivedPacket(packet, from: senderID)
handleReceivedPacket(packet, from: peerID)
} else {
let msgID = makeMessageID(for: packet)
collectionsQueue.async(flags: .barrier) { [weak self] in
self?.ingressByMessageID[msgID] = (.peripheral(peripheralUUID), Date())
}
handleReceivedPacket(packet, from: senderID)
handleReceivedPacket(packet, from: peerID)
}
}
@@ -2524,7 +2564,7 @@ extension BLEService: CBPeripheralManagerDelegate {
self.notifyPeerDisconnectedDebounced(peerID)
// Publish snapshots so UnifiedPeerService can refresh icons promptly
self.requestPeerDataPublish()
self.delegate?.didUpdatePeerList(currentPeerIDs)
self.deliverTransportEvent(.peerListUpdated(currentPeerIDs))
}
}
}
@@ -2633,19 +2673,15 @@ extension BLEService: CBPeripheralManagerDelegate {
pendingWriteBuffers.removeValue(forKey: centralUUID)
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 }
let trustedSenderID: PeerID?
if let knownPeerID = centralToPeerID[centralUUID] {
if knownPeerID != claimedSenderID {
SecureLogger.warning("🚫 SECURITY: Sender ID spoofing attempt detected! Central \(centralUUID.prefix(8))… claimed to be \(claimedSenderID.id.prefix(8))… but is bound to \(knownPeerID.id.prefix(8))", category: .security)
continue
}
trustedSenderID = knownPeerID
} else {
trustedSenderID = nil
}
if !validatePacket(packet, from: trustedSenderID ?? claimedSenderID, connectionSource: .central(centralUUID)) {
if !validatePacket(packet, from: context.validationPeerID, connectionSource: .central(centralUUID)) {
continue
}
@@ -2660,19 +2696,15 @@ extension BLEService: CBPeripheralManagerDelegate {
centralToPeerID[centralUUID] = claimedSenderID
refreshLocalTopology()
}
// Record ingress link for last-hop suppression then process
let msgID = makeMessageID(for: packet)
collectionsQueue.async(flags: .barrier) { [weak self] in
self?.ingressByMessageID[msgID] = (.central(centralUUID), Date())
if !recordIngressIfNew(packet, link: .central(centralUUID)) {
continue
}
handleReceivedPacket(packet, from: claimedSenderID)
handleReceivedPacket(packet, from: context.receivedFromPeerID)
} else {
// Record ingress link for last-hop suppression then process
let msgID = makeMessageID(for: packet)
collectionsQueue.async(flags: .barrier) { [weak self] in
self?.ingressByMessageID[msgID] = (.central(centralUUID), Date())
if !recordIngressIfNew(packet, link: .central(centralUUID)) {
continue
}
handleReceivedPacket(packet, from: claimedSenderID)
handleReceivedPacket(packet, from: context.receivedFromPeerID)
}
} else {
// If buffer grows suspiciously large, reset to avoid memory leak
@@ -2709,13 +2741,28 @@ extension BLEService {
extension BLEService {
/// Notify UI on the MainActor to satisfy Swift concurrency isolation
private func notifyUI(_ block: @escaping () -> Void) {
private func notifyUI(_ block: @escaping @MainActor () -> Void) {
// Always hop onto the MainActor so calls to @MainActor delegates are safe
Task { @MainActor in
block()
}
}
private func emitTransportEvent(_ event: TransportEvent) {
notifyUI { [weak self] in
self?.deliverTransportEvent(event)
}
}
@MainActor
private func deliverTransportEvent(_ event: TransportEvent) {
if let eventDelegate {
eventDelegate.didReceiveTransportEvent(event)
} else {
delegate?.receiveTransportEvent(event)
}
}
private func logBluetoothStatus(_ context: String) {
bleQueue.async { [weak self] in
guard let self = self else { return }
@@ -2983,12 +3030,12 @@ extension BLEService {
return Set(scored.prefix(k).map { $0.id })
}
private func priority(for packet: BitchatPacket, data: Data) -> OutboundPriority {
private func priority(for packet: BitchatPacket, data: Data) -> BLEOutboundWritePriority {
guard let messageType = MessageType(rawValue: packet.type) else { return .low }
switch messageType {
case .fragment:
let total = fragmentTotalCount(from: packet.payload)
return OutboundPriority.fragment(totalFragments: total)
return BLEOutboundWritePriority.fragment(totalFragments: total)
case .fileTransfer:
return .fileTransfer
default:
@@ -3004,7 +3051,7 @@ extension BLEService {
return max(total, 1)
}
private func writeOrEnqueue(_ data: Data, to peripheral: CBPeripheral, characteristic: CBCharacteristic, priority: OutboundPriority) {
private func writeOrEnqueue(_ data: Data, to peripheral: CBPeripheral, characteristic: CBCharacteristic, priority: BLEOutboundWritePriority) {
// BLE operations run on bleQueue; keep queue affinity
bleQueue.async { [weak self] in
guard let self = self else { return }
@@ -3013,29 +3060,20 @@ extension BLEService {
peripheral.writeValue(data, for: characteristic, type: .withoutResponse)
} else {
self.collectionsQueue.async(flags: .barrier) {
var queue = self.pendingPeripheralWrites[uuid] ?? []
let capBytes = TransportConfig.blePendingWriteBufferCapBytes
let newSize = data.count
// If single chunk exceeds cap, drop it immediately
if newSize > capBytes {
SecureLogger.warning("⚠️ Dropping oversized write chunk (\(newSize)B) for peripheral \(uuid)", category: .session)
} else {
let item = PendingWrite(priority: priority, data: data)
var total = queue.reduce(0) { $0 + $1.data.count } + newSize
let insertIndex = queue.firstIndex { item.priority < $0.priority } ?? queue.count
queue.insert(item, at: insertIndex)
if total > capBytes {
var removedBytes = 0
while total > capBytes && !queue.isEmpty {
let removed = queue.removeLast()
removedBytes += removed.data.count
total -= removed.data.count
}
if removedBytes > 0 {
SecureLogger.warning("📉 Trimmed pending write buffer for \(uuid) by \(removedBytes)B to \(total)B", category: .session)
}
}
self.pendingPeripheralWrites[uuid] = queue.isEmpty ? nil : queue
let result = self.pendingPeripheralWrites.enqueue(
data: data,
for: uuid,
priority: priority,
capBytes: TransportConfig.blePendingWriteBufferCapBytes
)
switch result {
case .oversized(let bytes):
SecureLogger.warning("⚠️ Dropping oversized write chunk (\(bytes)B) for peripheral \(uuid)", category: .session)
case let .enqueued(trimmedBytes, remainingBytes) where trimmedBytes > 0:
SecureLogger.warning("📉 Trimmed pending write buffer for \(uuid) by \(trimmedBytes)B to \(remainingBytes)B", category: .session)
case .enqueued:
break
}
}
}
@@ -3050,10 +3088,8 @@ extension BLEService {
// Atomically take all pending items from the queue to avoid race conditions
// where new items could be enqueued between read and update
let itemsToSend: [PendingWrite] = self.collectionsQueue.sync(flags: .barrier) {
let items = self.pendingPeripheralWrites[uuid] ?? []
self.pendingPeripheralWrites[uuid] = nil
return items
let itemsToSend: [BLEPendingWrite] = self.collectionsQueue.sync(flags: .barrier) {
self.pendingPeripheralWrites.takeAll(for: uuid)
}
guard !itemsToSend.isEmpty else { return }
@@ -3072,10 +3108,7 @@ extension BLEService {
let unsent = Array(itemsToSend.dropFirst(sent))
if !unsent.isEmpty {
self.collectionsQueue.async(flags: .barrier) {
var existing = self.pendingPeripheralWrites[uuid] ?? []
// Prepend unsent items to maintain priority order
existing.insert(contentsOf: unsent, at: 0)
self.pendingPeripheralWrites[uuid] = existing
self.pendingPeripheralWrites.prepend(unsent, for: uuid)
}
}
}
@@ -3118,7 +3151,7 @@ extension BLEService {
/// Periodically try to drain pending writes for all connected peripherals
private func drainAllPendingWrites() {
let uuids = collectionsQueue.sync { Array(pendingPeripheralWrites.keys) }
let uuids = collectionsQueue.sync { pendingPeripheralWrites.peripheralIDs }
for uuid in uuids {
guard let state = peripherals[uuid], state.isConnected else { continue }
drainPendingWrites(for: state.peripheral)
@@ -3205,7 +3238,7 @@ extension BLEService {
// Notify delegate that message was sent
notifyUI { [weak self] in
self?.delegate?.didUpdateMessageDeliveryStatus(messageID, status: .sent)
self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: messageID, status: .sent))
}
} catch {
SecureLogger.error("Failed to encrypt message: \(error)")
@@ -3226,7 +3259,7 @@ extension BLEService {
// Notify delegate that message is pending
notifyUI { [weak self] in
self?.delegate?.didUpdateMessageDeliveryStatus(messageID, status: .sending)
self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: messageID, status: .sending))
}
}
}
@@ -3300,7 +3333,7 @@ extension BLEService {
// Notify delegate that message was sent
notifyUI { [weak self] in
self?.delegate?.didUpdateMessageDeliveryStatus(messageID, status: .sent)
self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: messageID, status: .sent))
}
SecureLogger.debug("✅ Sent pending message \(messageID) to \(peerID) after handshake", category: .session)
@@ -3310,7 +3343,7 @@ extension BLEService {
// Notify delegate of failure
notifyUI { [weak self] in
self?.delegate?.didUpdateMessageDeliveryStatus(messageID, status: .failed(reason: "Encryption failed"))
self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: messageID, status: .failed(reason: "Encryption failed")))
}
}
}
@@ -3972,13 +4005,13 @@ extension BLEService {
// Only notify of connection for new or reconnected peers when it is a direct announce
if (packet.ttl == self.messageTTL) && (isNewPeer || isReconnectedPeer) {
self.delegate?.didConnectToPeer(peerID)
self.deliverTransportEvent(.peerConnected(peerID))
// Schedule initial unicast sync to this peer
self.gossipSyncManager?.scheduleInitialSyncToPeer(peerID, delaySeconds: 1.0)
}
self.requestPeerDataPublish()
self.delegate?.didUpdatePeerList(currentPeerIDs)
self.deliverTransportEvent(.peerListUpdated(currentPeerIDs))
}
// Track for sync (include our own and others' announces)
@@ -4115,11 +4148,15 @@ extension BLEService {
resolvedSelfMessageID = selfBroadcastMessageIDs.removeValue(forKey: dedupID)?.id
}
notifyUI { [weak self] in
self?.delegate?.didReceivePublicMessage(from: peerID,
nickname: senderNickname,
content: content,
timestamp: ts,
messageID: resolvedSelfMessageID)
self?.deliverTransportEvent(
.publicMessageReceived(
peerID: peerID,
nickname: senderNickname,
content: content,
timestamp: ts,
messageID: resolvedSelfMessageID
)
)
}
}
@@ -4183,27 +4220,27 @@ extension BLEService {
case .privateMessage:
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
notifyUI { [weak self] in
self?.delegate?.didReceiveNoisePayload(from: peerID, type: .privateMessage, payload: Data(payloadData), timestamp: ts)
self?.deliverTransportEvent(.noisePayloadReceived(peerID: peerID, type: .privateMessage, payload: Data(payloadData), timestamp: ts))
}
case .delivered:
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
notifyUI { [weak self] in
self?.delegate?.didReceiveNoisePayload(from: peerID, type: .delivered, payload: Data(payloadData), timestamp: ts)
self?.deliverTransportEvent(.noisePayloadReceived(peerID: peerID, type: .delivered, payload: Data(payloadData), timestamp: ts))
}
case .readReceipt:
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
notifyUI { [weak self] in
self?.delegate?.didReceiveNoisePayload(from: peerID, type: .readReceipt, payload: Data(payloadData), timestamp: ts)
self?.deliverTransportEvent(.noisePayloadReceived(peerID: peerID, type: .readReceipt, payload: Data(payloadData), timestamp: ts))
}
case .verifyChallenge:
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
notifyUI { [weak self] in
self?.delegate?.didReceiveNoisePayload(from: peerID, type: .verifyChallenge, payload: Data(payloadData), timestamp: ts)
self?.deliverTransportEvent(.noisePayloadReceived(peerID: peerID, type: .verifyChallenge, payload: Data(payloadData), timestamp: ts))
}
case .verifyResponse:
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
notifyUI { [weak self] in
self?.delegate?.didReceiveNoisePayload(from: peerID, type: .verifyResponse, payload: Data(payloadData), timestamp: ts)
self?.deliverTransportEvent(.noisePayloadReceived(peerID: peerID, type: .verifyResponse, payload: Data(payloadData), timestamp: ts))
}
case .none:
SecureLogger.warning("⚠️ Unknown noise payload type: \(payloadType)")
@@ -4264,11 +4301,12 @@ extension BLEService {
}
// Debounced disconnect notifier to avoid duplicate disconnect callbacks within a short window
@MainActor
private func notifyPeerDisconnectedDebounced(_ peerID: PeerID) {
let now = Date()
let last = recentDisconnectNotifies[peerID]
if last == nil || now.timeIntervalSince(last!) >= TransportConfig.bleDisconnectNotifyDebounceSeconds {
delegate?.didDisconnectFromPeer(peerID)
deliverTransportEvent(.peerDisconnected(peerID))
recentDisconnectNotifies[peerID] = now
} else {
// Suppressed duplicate disconnect notification
@@ -4425,11 +4463,11 @@ extension BLEService {
let currentPeerIDs = self.collectionsQueue.sync { self.currentPeerIDs }
for peerID in disconnectedPeers {
self.delegate?.didDisconnectFromPeer(peerID)
self.deliverTransportEvent(.peerDisconnected(peerID))
}
// Publish snapshots so UnifiedPeerService updates connection/reachability icons
self.requestPeerDataPublish()
self.delegate?.didUpdatePeerList(currentPeerIDs)
self.deliverTransportEvent(.peerListUpdated(currentPeerIDs))
}
}