From ab0da61533a1091ee2d50168d587985ff9ef18ff Mon Sep 17 00:00:00 2001 From: jack <212554440+jackjackbits@users.noreply.github.com> Date: Sun, 31 May 2026 13:58:27 +0200 Subject: [PATCH] [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 --- bitchat/Features/media/ImageUtils.swift | 2 +- .../Services/BLE/BLEOutboundWriteBuffer.swift | 83 +++++ bitchat/Services/BLE/BLEService.swift | 306 ++++++++++-------- bitchat/Services/NostrTransport.swift | 7 +- .../NotificationStreamAssembler.swift | 21 ++ bitchat/Services/Transport.swift | 51 +++ bitchat/Sync/GossipSyncManager.swift | 19 +- bitchat/ViewModels/ChatViewModel.swift | 7 +- .../ChatViewModelBootstrapper.swift | 1 + bitchatTests/BLEServiceCoreTests.swift | 71 ++++ bitchatTests/ChatViewModelTests.swift | 2 + bitchatTests/Features/ImageUtilsTests.swift | 15 + bitchatTests/GossipSyncManagerTests.swift | 37 ++- bitchatTests/Mocks/MockTransport.swift | 1 + .../NotificationStreamAssemblerTests.swift | 21 ++ bitchatTests/ProtocolContractTests.swift | 1 + .../BLEOutboundWriteBufferTests.swift | 75 +++++ .../Services/NostrTransportTests.swift | 1 + docs/ARCHITECTURE_V2.md | 6 + 19 files changed, 575 insertions(+), 152 deletions(-) create mode 100644 bitchat/Services/BLE/BLEOutboundWriteBuffer.swift create mode 100644 bitchatTests/Services/BLEOutboundWriteBufferTests.swift diff --git a/bitchat/Features/media/ImageUtils.swift b/bitchat/Features/media/ImageUtils.swift index af6a37af..f362c709 100644 --- a/bitchat/Features/media/ImageUtils.swift +++ b/bitchat/Features/media/ImageUtils.swift @@ -189,7 +189,7 @@ enum ImageUtils { private static func makeOutputURL() throws -> URL { let formatter = DateFormatter() formatter.dateFormat = "yyyyMMdd_HHmmss" - let fileName = "img_\(formatter.string(from: Date())).jpg" + let fileName = "img_\(formatter.string(from: Date()))_\(UUID().uuidString).jpg" let directory = try applicationFilesDirectory().appendingPathComponent("images/outgoing", isDirectory: true) try FileManager.default.createDirectory(at: directory, withIntermediateDirectories: true, attributes: nil) diff --git a/bitchat/Services/BLE/BLEOutboundWriteBuffer.swift b/bitchat/Services/BLE/BLEOutboundWriteBuffer.swift new file mode 100644 index 00000000..f1bbbee7 --- /dev/null +++ b/bitchat/Services/BLE/BLEOutboundWriteBuffer.swift @@ -0,0 +1,83 @@ +import Foundation + +struct BLEOutboundWritePriority: Comparable { + let level: Int + let suborder: Int + + static let high = BLEOutboundWritePriority(level: 0, suborder: 0) + + static func fragment(totalFragments: Int) -> BLEOutboundWritePriority { + BLEOutboundWritePriority(level: 1, suborder: max(1, min(totalFragments, Int(UInt16.max)))) + } + + static let fileTransfer = BLEOutboundWritePriority(level: 2, suborder: Int.max - 1) + static let low = BLEOutboundWritePriority(level: 2, suborder: Int.max) + + static func < (lhs: BLEOutboundWritePriority, rhs: BLEOutboundWritePriority) -> Bool { + if lhs.level != rhs.level { return lhs.level < rhs.level } + return lhs.suborder < rhs.suborder + } +} + +struct BLEPendingWrite { + let priority: BLEOutboundWritePriority + let data: Data +} + +struct BLEOutboundWriteBuffer { + enum EnqueueResult { + case enqueued(trimmedBytes: Int, remainingBytes: Int) + case oversized(bytes: Int) + } + + private var writesByPeripheralID: [String: [BLEPendingWrite]] = [:] + + var peripheralIDs: [String] { + Array(writesByPeripheralID.keys) + } + + mutating func removeAll() { + writesByPeripheralID.removeAll() + } + + mutating func enqueue( + data: Data, + for peripheralID: String, + priority: BLEOutboundWritePriority, + capBytes: Int + ) -> EnqueueResult { + guard data.count <= capBytes else { + return .oversized(bytes: data.count) + } + + var queue = writesByPeripheralID[peripheralID] ?? [] + let item = BLEPendingWrite(priority: priority, data: data) + let insertIndex = queue.firstIndex { item.priority < $0.priority } ?? queue.count + queue.insert(item, at: insertIndex) + + var total = queue.reduce(0) { $0 + $1.data.count } + var trimmedBytes = 0 + + while total > capBytes && !queue.isEmpty { + let removed = queue.removeLast() + trimmedBytes += removed.data.count + total -= removed.data.count + } + + writesByPeripheralID[peripheralID] = queue.isEmpty ? nil : queue + return .enqueued(trimmedBytes: trimmedBytes, remainingBytes: total) + } + + mutating func takeAll(for peripheralID: String) -> [BLEPendingWrite] { + let items = writesByPeripheralID[peripheralID] ?? [] + writesByPeripheralID[peripheralID] = nil + return items + } + + mutating func prepend(_ items: [BLEPendingWrite], for peripheralID: String) { + guard !items.isEmpty else { return } + var existing = writesByPeripheralID[peripheralID] ?? [] + existing.insert(contentsOf: items, at: 0) + writesByPeripheralID[peripheralID] = existing + } +} diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index 055c254f..be764f9b 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -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)) } } diff --git a/bitchat/Services/NostrTransport.swift b/bitchat/Services/NostrTransport.swift index b822a48b..d38adfc4 100644 --- a/bitchat/Services/NostrTransport.swift +++ b/bitchat/Services/NostrTransport.swift @@ -106,6 +106,7 @@ final class NostrTransport: Transport, @unchecked Sendable { // MARK: - Transport Protocol Conformance weak var delegate: BitchatDelegate? + weak var eventDelegate: TransportEventDelegate? weak var peerEventsDelegate: TransportPeerEventsDelegate? var peerSnapshotPublisher: AnyPublisher<[TransportPeerSnapshot], Never> { @@ -173,9 +174,9 @@ final class NostrTransport: Transport, @unchecked Sendable { func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) { // Enqueue and process with throttling to avoid relay rate limits // Use barrier to synchronize access to readQueue - queue.async(flags: .barrier) { [weak self] in - self?.readQueue.append(QueuedRead(receipt: receipt, peerID: peerID)) - self?.processReadQueueIfNeeded() + queue.async(flags: .barrier) { + self.readQueue.append(QueuedRead(receipt: receipt, peerID: peerID)) + self.processReadQueueIfNeeded() } } diff --git a/bitchat/Services/NotificationStreamAssembler.swift b/bitchat/Services/NotificationStreamAssembler.swift index 0a46ffbd..c4ebb25d 100644 --- a/bitchat/Services/NotificationStreamAssembler.swift +++ b/bitchat/Services/NotificationStreamAssembler.swift @@ -21,6 +21,19 @@ struct NotificationStreamAssembler { pendingFrameExpectedLength = 0 } + private mutating func discardLeadingPaddingIfPresent() -> Bool { + guard let first = buffer.first else { return false } + guard first != 1 && first != 2 else { return false } + let paddingLength = Int(first) + guard paddingLength > 0, paddingLength <= buffer.count else { return false } + guard buffer.prefix(paddingLength).allSatisfy({ $0 == first }) else { return false } + + buffer.removeFirst(paddingLength) + pendingFrameStartedAt = nil + pendingFrameExpectedLength = 0 + return true + } + mutating func append(_ chunk: Data) -> (frames: [Data], droppedPrefixes: [UInt8], reset: Bool) { guard !chunk.isEmpty else { return ([], [], false) } @@ -42,6 +55,9 @@ struct NotificationStreamAssembler { while buffer.count >= minimumFramePrefix { guard let version = buffer.first else { break } guard version == 1 || version == 2 else { + if discardLeadingPaddingIfPresent() { + continue + } dropped.append(buffer.removeFirst()) pendingFrameStartedAt = nil pendingFrameExpectedLength = 0 @@ -133,6 +149,11 @@ struct NotificationStreamAssembler { let frame = Data(buffer.prefix(frameLength)) frames.append(frame) buffer.removeFirst(frameLength) + _ = discardLeadingPaddingIfPresent() + } + + if discardLeadingPaddingIfPresent() { + return (frames, dropped, didReset) } if !buffer.isEmpty, buffer.allSatisfy({ $0 == 0 }) { diff --git a/bitchat/Services/Transport.swift b/bitchat/Services/Transport.swift index 8f7fd03d..72c6d542 100644 --- a/bitchat/Services/Transport.swift +++ b/bitchat/Services/Transport.swift @@ -1,6 +1,7 @@ import BitFoundation import Foundation import Combine +import CoreBluetooth /// Abstract transport interface used by ChatViewModel and services. /// BLEService implements this protocol; a future Nostr transport can too. @@ -12,9 +13,27 @@ struct TransportPeerSnapshot: Equatable, Hashable { let lastSeen: Date } +enum TransportEvent: @unchecked Sendable { + case messageReceived(BitchatMessage) + case publicMessageReceived(peerID: PeerID, nickname: String, content: String, timestamp: Date, messageID: String?) + case noisePayloadReceived(peerID: PeerID, type: NoisePayloadType, payload: Data, timestamp: Date) + case peerConnected(PeerID) + case peerDisconnected(PeerID) + case peerListUpdated([PeerID]) + case peerSnapshotsUpdated([TransportPeerSnapshot]) + case messageDeliveryStatusUpdated(messageID: String, status: DeliveryStatus) + case bluetoothStateUpdated(CBManagerState) +} + +protocol TransportEventDelegate: AnyObject { + @MainActor func didReceiveTransportEvent(_ event: TransportEvent) +} + protocol Transport: AnyObject { // Event sink var delegate: BitchatDelegate? { get set } + // Typed event sink for transport-domain events. Prefer this over BitchatDelegate for new code. + var eventDelegate: TransportEventDelegate? { get set } // Peer events (preferred over publishers for UI) var peerEventsDelegate: TransportPeerEventsDelegate? { get set } @@ -84,4 +103,36 @@ protocol TransportPeerEventsDelegate: AnyObject { @MainActor func didUpdatePeerSnapshots(_ peers: [TransportPeerSnapshot]) } +extension BitchatDelegate { + @MainActor + func receiveTransportEvent(_ event: TransportEvent) { + switch event { + case .messageReceived(let message): + didReceiveMessage(message) + case let .publicMessageReceived(peerID, nickname, content, timestamp, messageID): + didReceivePublicMessage( + from: peerID, + nickname: nickname, + content: content, + timestamp: timestamp, + messageID: messageID + ) + case let .noisePayloadReceived(peerID, type, payload, timestamp): + didReceiveNoisePayload(from: peerID, type: type, payload: payload, timestamp: timestamp) + case .peerConnected(let peerID): + didConnectToPeer(peerID) + case .peerDisconnected(let peerID): + didDisconnectFromPeer(peerID) + case .peerListUpdated(let peers): + didUpdatePeerList(peers) + case .peerSnapshotsUpdated: + break + case let .messageDeliveryStatusUpdated(messageID, status): + didUpdateMessageDeliveryStatus(messageID, status: status) + case .bluetoothStateUpdated(let state): + didUpdateBluetoothState(state) + } + } +} + extension BLEService: Transport {} diff --git a/bitchat/Sync/GossipSyncManager.swift b/bitchat/Sync/GossipSyncManager.swift index 1b65f9f9..3c8472d2 100644 --- a/bitchat/Sync/GossipSyncManager.swift +++ b/bitchat/Sync/GossipSyncManager.swift @@ -128,17 +128,15 @@ final class GossipSyncManager { func scheduleInitialSyncToPeer(_ peerID: PeerID, delaySeconds: TimeInterval = 5.0) { queue.asyncAfter(deadline: .now() + delaySeconds) { [weak self] in guard let self = self else { return } - self.sendRequestSync(to: peerID, types: .publicMessages) + + var types: SyncTypeFlags = .publicMessages if self.config.fragmentCapacity > 0 && self.config.fragmentSyncIntervalSeconds > 0 { - self.queue.asyncAfter(deadline: .now() + 0.5) { [weak self] in - self?.sendRequestSync(to: peerID, types: .fragment) - } + types.formUnion(.fragment) } if self.config.fileTransferCapacity > 0 && self.config.fileTransferSyncIntervalSeconds > 0 { - self.queue.asyncAfter(deadline: .now() + 1.0) { [weak self] in - self?.sendRequestSync(to: peerID, types: .fileTransfer) - } + types.formUnion(.fileTransfer) } + self.sendRequestSync(to: peerID, types: types) } } @@ -393,13 +391,18 @@ final class GossipSyncManager { cleanupStaleAnnouncementsIfNeeded(now: now) requestSyncManager.cleanup() // Cleanup expired sync requests + var dueTypes: SyncTypeFlags = [] for index in syncSchedules.indices { guard syncSchedules[index].interval > 0 else { continue } if syncSchedules[index].lastSent == .distantPast || now.timeIntervalSince(syncSchedules[index].lastSent) >= syncSchedules[index].interval { syncSchedules[index].lastSent = now - sendPeriodicSync(for: syncSchedules[index].types) + dueTypes.formUnion(syncSchedules[index].types) } } + + if !dueTypes.isEmpty { + sendPeriodicSync(for: dueTypes) + } } private func cleanupStaleAnnouncementsIfNeeded(now: Date) { diff --git a/bitchat/ViewModels/ChatViewModel.swift b/bitchat/ViewModels/ChatViewModel.swift index 77e418f2..10662a19 100644 --- a/bitchat/ViewModels/ChatViewModel.swift +++ b/bitchat/ViewModels/ChatViewModel.swift @@ -93,7 +93,7 @@ import UniformTypeIdentifiers /// Manages the application state and business logic for BitChat. /// Acts as the primary coordinator between UI components and backend services, /// implementing the BitchatDelegate protocol to handle network events. -final class ChatViewModel: ObservableObject, BitchatDelegate, CommandContextProvider, GeohashParticipantContext, MessageFormattingContext { +final class ChatViewModel: ObservableObject, BitchatDelegate, TransportEventDelegate, CommandContextProvider, GeohashParticipantContext, MessageFormattingContext { // Use MessageFormattingEngine.Patterns for regex matching (shared, precompiled) typealias Patterns = MessageFormattingEngine.Patterns @@ -1298,6 +1298,11 @@ final class ChatViewModel: ObservableObject, BitchatDelegate, CommandContextProv } // MARK: - Message Reception + + @MainActor + func didReceiveTransportEvent(_ event: TransportEvent) { + receiveTransportEvent(event) + } func didReceiveMessage(_ message: BitchatMessage) { transportEventCoordinator.didReceiveMessage(message) diff --git a/bitchat/ViewModels/ChatViewModelBootstrapper.swift b/bitchat/ViewModels/ChatViewModelBootstrapper.swift index b876ba08..2d7a8a03 100644 --- a/bitchat/ViewModels/ChatViewModelBootstrapper.swift +++ b/bitchat/ViewModels/ChatViewModelBootstrapper.swift @@ -130,6 +130,7 @@ private extension ChatViewModelBootstrapper { func configureTransport() { viewModel.meshService.delegate = viewModel + viewModel.meshService.eventDelegate = viewModel DispatchQueue.main.asyncAfter(deadline: .now() + TransportConfig.uiStartupInitialDelaySeconds) { [weak viewModel] in guard let viewModel else { return } diff --git a/bitchatTests/BLEServiceCoreTests.swift b/bitchatTests/BLEServiceCoreTests.swift index 8849f089..69df87c1 100644 --- a/bitchatTests/BLEServiceCoreTests.swift +++ b/bitchatTests/BLEServiceCoreTests.swift @@ -92,6 +92,77 @@ struct BLEServiceCoreTests { _ = await TestHelpers.waitUntil({ !ble.currentPeerSnapshots().isEmpty }, timeout: 0.3) #expect(ble.currentPeerSnapshots().isEmpty) } + + @Test + func ingressAllowsRelayedSenderOnBoundLink() async throws { + let ble = makeService() + let boundPeer = PeerID(str: "1122334455667788") + let relayedSender = PeerID(str: "8899aabbccddeeff") + let packet = makePublicPacket( + content: "Relayed", + sender: relayedSender, + timestamp: UInt64(Date().timeIntervalSince1970 * 1000) + ) + + #expect(ble._test_acceptsIngress(packet: packet, boundPeerID: boundPeer)) + } + + @Test + func ingressRejectsDirectAnnounceThatConflictsWithBoundLink() async throws { + let ble = makeService() + let boundPeer = PeerID(str: "1122334455667788") + let claimedPeer = PeerID(str: "8899aabbccddeeff") + let packet = BitchatPacket( + type: MessageType.announce.rawValue, + senderID: Data(hexString: claimedPeer.id) ?? Data(), + recipientID: nil, + timestamp: UInt64(Date().timeIntervalSince1970 * 1000), + payload: Data(), + signature: nil, + ttl: 7 + ) + + #expect(!ble._test_acceptsIngress(packet: packet, boundPeerID: boundPeer)) + } + + @Test + func ingressRejectsSelfLoopbackBeforeSpoofChecks() async throws { + let ble = makeService() + let packet = makePublicPacket( + content: "Loopback", + sender: ble.myPeerID, + timestamp: UInt64(Date().timeIntervalSince1970 * 1000) + ) + + #expect(!ble._test_acceptsIngress(packet: packet, boundPeerID: PeerID(str: "1122334455667788"))) + } + + @Test + func ingressAllowsSelfAuthoredRSRWithTTLZeroFromBoundPeer() async throws { + let ble = makeService() + var packet = makePublicPacket( + content: "Recovered by sync", + sender: ble.myPeerID, + timestamp: UInt64(Date().timeIntervalSince1970 * 1000) + ) + packet.isRSR = true + packet.ttl = 0 + + #expect(ble._test_acceptsIngress(packet: packet, boundPeerID: PeerID(str: "1122334455667788"))) + } + + @Test + func ingressRecordSuppressesSecondLinkDuplicate() async throws { + let ble = makeService() + let packet = makePublicPacket( + content: "Duplicate link copy", + sender: PeerID(str: "1122334455667788"), + timestamp: UInt64(Date().timeIntervalSince1970 * 1000) + ) + + #expect(ble._test_recordIngressIfNew(packet: packet, linkID: "central-a")) + #expect(!ble._test_recordIngressIfNew(packet: packet, linkID: "central-b")) + } } private func makeService() -> BLEService { diff --git a/bitchatTests/ChatViewModelTests.swift b/bitchatTests/ChatViewModelTests.swift index d59c911a..19737944 100644 --- a/bitchatTests/ChatViewModelTests.swift +++ b/bitchatTests/ChatViewModelTests.swift @@ -42,6 +42,7 @@ struct ChatViewModelInitializationTests { // The viewModel should set itself as the transport delegate #expect(transport.delegate === viewModel) + #expect(transport.eventDelegate === viewModel) } @Test @MainActor @@ -421,6 +422,7 @@ struct ChatViewModelReceivingTests { // Message may or may not appear due to rate limiting/pipeline batching // The important thing is no crash and delegate was called #expect(transport.delegate === viewModel) + #expect(transport.eventDelegate === viewModel) } @Test @MainActor diff --git a/bitchatTests/Features/ImageUtilsTests.swift b/bitchatTests/Features/ImageUtilsTests.swift index 880403bb..08273244 100644 --- a/bitchatTests/Features/ImageUtilsTests.swift +++ b/bitchatTests/Features/ImageUtilsTests.swift @@ -64,4 +64,19 @@ struct ImageUtilsTests { #expect(data.starts(with: Data([0xFF, 0xD8]))) #expect(data.count > 0) } + + @Test + func processImage_usesUniqueOutputURLs() throws { + let image = makePlatformImage(size: CGSize(width: 64, height: 64)) + let firstURL = try ImageUtils.processImage(image, maxDimension: 64) + let secondURL = try ImageUtils.processImage(image, maxDimension: 64) + defer { + try? FileManager.default.removeItem(at: firstURL) + try? FileManager.default.removeItem(at: secondURL) + } + + #expect(firstURL != secondURL) + #expect(FileManager.default.fileExists(atPath: firstURL.path)) + #expect(FileManager.default.fileExists(atPath: secondURL.path)) + } } diff --git a/bitchatTests/GossipSyncManagerTests.swift b/bitchatTests/GossipSyncManagerTests.swift index 9d45da98..54e65afc 100644 --- a/bitchatTests/GossipSyncManagerTests.swift +++ b/bitchatTests/GossipSyncManagerTests.swift @@ -195,12 +195,39 @@ struct GossipSyncManagerTests { manager._performMaintenanceSynchronously(now: Date()) let sentPackets = delegate.packets - #expect(sentPackets.count == 3) + #expect(sentPackets.count == 1) let decoded = sentPackets.compactMap { RequestSyncPacket.decode(from: $0.payload) } - #expect(decoded.count == 3) - #expect(decoded[0].types == .publicMessages) - #expect(decoded[1].types == .fragment) - #expect(decoded[2].types == .fileTransfer) + #expect(decoded.count == 1) + let types = try #require(decoded.first?.types) + #expect(types.contains(.announce)) + #expect(types.contains(.message)) + #expect(types.contains(.fragment)) + #expect(types.contains(.fileTransfer)) + } + + @Test func initialSyncCoalescesEnabledTypes() async throws { + var config = GossipSyncManager.Config() + config.seenCapacity = 10 + config.fragmentCapacity = 5 + config.fileTransferCapacity = 4 + config.fragmentSyncIntervalSeconds = 1 + config.fileTransferSyncIntervalSeconds = 1 + + let requestSyncManager = RequestSyncManager() + let manager = GossipSyncManager(myPeerID: myPeerID, config: config, requestSyncManager: requestSyncManager) + let delegate = RecordingDelegate() + manager.delegate = delegate + + manager.scheduleInitialSyncToPeer(PeerID(str: "FFFFFFFFFFFFFFFF"), delaySeconds: 0.0) + + try await TestHelpers.waitFor({ delegate.packets.count == 1 }, timeout: TestConstants.shortTimeout) + let packet = try #require(delegate.packets.first) + let request = try #require(RequestSyncPacket.decode(from: packet.payload)) + let types = try #require(request.types) + #expect(types.contains(.announce)) + #expect(types.contains(.message)) + #expect(types.contains(.fragment)) + #expect(types.contains(.fileTransfer)) } @Test func handleRequestSyncHonorsTypeFilter() async throws { diff --git a/bitchatTests/Mocks/MockTransport.swift b/bitchatTests/Mocks/MockTransport.swift index 2f077ed2..5ab7d7dc 100644 --- a/bitchatTests/Mocks/MockTransport.swift +++ b/bitchatTests/Mocks/MockTransport.swift @@ -19,6 +19,7 @@ final class MockTransport: Transport { // MARK: - Protocol Properties weak var delegate: BitchatDelegate? + weak var eventDelegate: TransportEventDelegate? weak var peerEventsDelegate: TransportPeerEventsDelegate? var myPeerID: PeerID = PeerID(str: "TESTPEER") diff --git a/bitchatTests/NotificationStreamAssemblerTests.swift b/bitchatTests/NotificationStreamAssemblerTests.swift index 7d78fb0a..de76bda0 100644 --- a/bitchatTests/NotificationStreamAssemblerTests.swift +++ b/bitchatTests/NotificationStreamAssemblerTests.swift @@ -97,6 +97,27 @@ struct NotificationStreamAssemblerTests { #expect(decoded.timestamp == packet.timestamp) } + @Test func discardsPKCSPaddingBetweenNotificationFrames() throws { + var assembler = NotificationStreamAssembler() + let packet1 = makePacket(timestamp: 0x111) + let packet2 = makePacket(timestamp: 0x222) + let paddedFrame1 = try #require(packet1.toBinaryData(padding: true), "Failed to encode first padded packet") + let paddedFrame2 = try #require(packet2.toBinaryData(padding: true), "Failed to encode second padded packet") + + var result = assembler.append(paddedFrame1) + #expect(result.frames.count == 1) + #expect(result.droppedPrefixes.isEmpty) + #expect(result.reset == false) + + result = assembler.append(paddedFrame2) + #expect(result.frames.count == 1) + #expect(result.droppedPrefixes.isEmpty) + #expect(result.reset == false) + + let decoded = try #require(BinaryProtocol.decode(result.frames[0]), "Failed to decode second frame") + #expect(decoded.timestamp == packet2.timestamp) + } + func testAssemblesCompressedLargeFrame() throws { var assembler = NotificationStreamAssembler() diff --git a/bitchatTests/ProtocolContractTests.swift b/bitchatTests/ProtocolContractTests.swift index df8a4b46..27541f47 100644 --- a/bitchatTests/ProtocolContractTests.swift +++ b/bitchatTests/ProtocolContractTests.swift @@ -15,6 +15,7 @@ private final class DefaultDelegateProbe: BitchatDelegate { private final class DefaultTransportProbe: Transport { weak var delegate: BitchatDelegate? + weak var eventDelegate: TransportEventDelegate? weak var peerEventsDelegate: TransportPeerEventsDelegate? let subject = CurrentValueSubject<[TransportPeerSnapshot], Never>([]) diff --git a/bitchatTests/Services/BLEOutboundWriteBufferTests.swift b/bitchatTests/Services/BLEOutboundWriteBufferTests.swift new file mode 100644 index 00000000..f8e64080 --- /dev/null +++ b/bitchatTests/Services/BLEOutboundWriteBufferTests.swift @@ -0,0 +1,75 @@ +import Foundation +import Testing +@testable import bitchat + +struct BLEOutboundWriteBufferTests { + @Test + func enqueueOrdersWritesByPriority() { + var buffer = BLEOutboundWriteBuffer() + let peerID = "peer-1" + + _ = buffer.enqueue( + data: Data(repeating: 0x01, count: 4), + for: peerID, + priority: .fileTransfer, + capBytes: 64 + ) + _ = buffer.enqueue( + data: Data(repeating: 0x02, count: 4), + for: peerID, + priority: .high, + capBytes: 64 + ) + _ = buffer.enqueue( + data: Data(repeating: 0x03, count: 4), + for: peerID, + priority: .fragment(totalFragments: 2), + capBytes: 64 + ) + + let writes = buffer.takeAll(for: peerID) + + #expect(writes.map { Int($0.data.first ?? 0) } == [0x02, 0x03, 0x01]) + } + + @Test + func enqueueTrimsLowestPriorityItemsToCap() { + var buffer = BLEOutboundWriteBuffer() + let peerID = "peer-1" + + _ = buffer.enqueue(data: Data(repeating: 0x01, count: 8), for: peerID, priority: .low, capBytes: 16) + _ = buffer.enqueue(data: Data(repeating: 0x02, count: 8), for: peerID, priority: .fileTransfer, capBytes: 16) + let result = buffer.enqueue(data: Data(repeating: 0x03, count: 8), for: peerID, priority: .high, capBytes: 16) + + if case let .enqueued(trimmedBytes, remainingBytes) = result { + #expect(trimmedBytes == 8) + #expect(remainingBytes == 16) + } else { + Issue.record("Expected buffered write to trim, not drop as oversized") + } + + let writes = buffer.takeAll(for: peerID) + + #expect(writes.map { Int($0.data.first ?? 0) } == [0x03, 0x02]) + } + + @Test + func enqueueRejectsOversizedSingleChunk() { + var buffer = BLEOutboundWriteBuffer() + + let result = buffer.enqueue( + data: Data(repeating: 0x01, count: 32), + for: "peer-1", + priority: .high, + capBytes: 16 + ) + + if case let .oversized(bytes) = result { + #expect(bytes == 32) + } else { + Issue.record("Expected oversized write to be rejected") + } + + #expect(buffer.peripheralIDs.isEmpty) + } +} diff --git a/bitchatTests/Services/NostrTransportTests.swift b/bitchatTests/Services/NostrTransportTests.swift index cb10d7fc..9541f3b8 100644 --- a/bitchatTests/Services/NostrTransportTests.swift +++ b/bitchatTests/Services/NostrTransportTests.swift @@ -306,6 +306,7 @@ struct NostrTransportTests { let secondPayload = try decodeEmbeddedPayload(from: secondEvent, recipient: recipient).payload #expect(secondPayload.type == .readReceipt) #expect(String(data: secondPayload.data, encoding: .utf8) == "read-2") + withExtendedLifetime(transport) {} } @Test("Concurrent read receipt enqueue does not crash") diff --git a/docs/ARCHITECTURE_V2.md b/docs/ARCHITECTURE_V2.md index 6b2c1037..6a6de351 100644 --- a/docs/ARCHITECTURE_V2.md +++ b/docs/ARCHITECTURE_V2.md @@ -38,3 +38,9 @@ This branch starts a larger simplification effort focused on performance, reliab 3. Move private/public conversation mutation behind the new store instead of still mirroring legacy message writes from `ChatViewModel`. 4. Replace remaining singleton-heavy seams with injected runtime services where practical. 5. Revisit actor isolation for identity and conversation state once the remaining message/peer models are safe to move off the main actor. + +## Transport Follow-Up + +The Bluetooth architecture branch begins step 2 by adding a typed `TransportEvent` boundary while preserving the legacy `BitchatDelegate` bridge. New transport code should emit typed events first, with delegate forwarding used only as a compatibility adapter during migration. + +The branch also starts carving performance-sensitive BLE scheduling state out of `BLEService`: pending write backpressure now lives in `BLEOutboundWriteBuffer`, giving the outbound hot path a focused, unit-tested component before deeper fragmentation and link-scheduler work.