From 9e84f5e822499f93877092fe0d48b888b285ceef Mon Sep 17 00:00:00 2001 From: jack <212554440+jackjackbits@users.noreply.github.com> Date: Sun, 31 May 2026 14:16:11 +0200 Subject: [PATCH] [codex] Refactor BLE outbound scheduling and Noise queues (#1306) * Refactor BLE transport event handling * Make image output paths unique * Keep queued Nostr read receipts alive * Refine BLE ingress fanout * Rediscover BLE service after invalidation * Extract BLE notification retry buffer * Extract BLE inbound write buffer * Extract BLE fragment assembly buffer * Tidy secure log handling from device run * Extract BLE outbound fragment scheduler * Harden app CI media tests * Redact BLE message content from logs * Extract BLE Noise session queues * Fix BLE read receipt UI updates * Allow self-authored RSR ingress replies * Harden read receipt queue test timing * Extract BLE outbound policy and incoming file storage * Avoid duplicate BLE link snapshots during send * Canonicalize Nostr relay URLs --------- Co-authored-by: jack --- bitchat/Features/media/ImageUtils.swift | 23 +- bitchat/Nostr/GeoRelayDirectory.swift | 7 +- bitchat/Nostr/NostrRelayManager.swift | 17 +- bitchat/Nostr/NostrRelayURL.swift | 51 ++ .../Services/BLE/BLEIncomingFileStore.swift | 168 ++++ .../Services/BLE/BLENoisePayloadFactory.swift | 25 + .../Services/BLE/BLENoiseSessionQueues.swift | 46 ++ ...BLEOutboundFragmentTransferScheduler.swift | 181 +++++ .../BLE/BLEOutboundPacketPolicy.swift | 43 ++ bitchat/Services/BLE/BLEService.swift | 729 +++++------------- .../ViewModels/ChatDeliveryCoordinator.swift | 20 +- .../ChatTransportEventCoordinator.swift | 77 +- bitchatTests/ChatViewModelTests.swift | 51 ++ bitchatTests/Features/ImageUtilsTests.swift | 22 +- .../Nostr/GeoRelayDirectoryTests.swift | 2 + .../BLENoisePayloadFactoryTests.swift | 34 + .../Services/BLENoiseSessionQueuesTests.swift | 67 ++ ...tboundFragmentTransferSchedulerTests.swift | 148 ++++ .../Services/NostrRelayManagerTests.swift | 2 +- 19 files changed, 1084 insertions(+), 629 deletions(-) create mode 100644 bitchat/Nostr/NostrRelayURL.swift create mode 100644 bitchat/Services/BLE/BLEIncomingFileStore.swift create mode 100644 bitchat/Services/BLE/BLENoisePayloadFactory.swift create mode 100644 bitchat/Services/BLE/BLENoiseSessionQueues.swift create mode 100644 bitchat/Services/BLE/BLEOutboundFragmentTransferScheduler.swift create mode 100644 bitchat/Services/BLE/BLEOutboundPacketPolicy.swift create mode 100644 bitchatTests/Services/BLENoisePayloadFactoryTests.swift create mode 100644 bitchatTests/Services/BLENoiseSessionQueuesTests.swift create mode 100644 bitchatTests/Services/BLEOutboundFragmentTransferSchedulerTests.swift diff --git a/bitchat/Features/media/ImageUtils.swift b/bitchat/Features/media/ImageUtils.swift index f362c709..a75cc6db 100644 --- a/bitchat/Features/media/ImageUtils.swift +++ b/bitchat/Features/media/ImageUtils.swift @@ -16,7 +16,7 @@ enum ImageUtils { private static let compressionQuality: CGFloat = 0.82 private static let targetImageBytes: Int = 45_000 - static func processImage(at url: URL, maxDimension: CGFloat = 448) throws -> URL { + static func processImage(at url: URL, maxDimension: CGFloat = 448, outputDirectory: URL? = nil) throws -> URL { // Security H1: Check file size BEFORE reading into memory let attrs = try FileManager.default.attributesOfItem(atPath: url.path) guard let fileSize = attrs[.size] as? Int else { @@ -30,15 +30,15 @@ enum ImageUtils { let data = try Data(contentsOf: url) #if os(iOS) guard let image = UIImage(data: data) else { throw ImageUtilsError.invalidImage } - return try processImage(image, maxDimension: maxDimension) + return try processImage(image, maxDimension: maxDimension, outputDirectory: outputDirectory) #else guard let image = NSImage(data: data) else { throw ImageUtilsError.invalidImage } - return try processImage(image, maxDimension: maxDimension) + return try processImage(image, maxDimension: maxDimension, outputDirectory: outputDirectory) #endif } #if os(iOS) - static func processImage(_ image: UIImage, maxDimension: CGFloat = 448) throws -> URL { + static func processImage(_ image: UIImage, maxDimension: CGFloat = 448, outputDirectory: URL? = nil) throws -> URL { return try autoreleasepool { // Scale the image first let scaled = scaledImage(image, maxDimension: maxDimension) @@ -64,7 +64,7 @@ enum ImageUtils { } } - let outputURL = try makeOutputURL() + let outputURL = try makeOutputURL(outputDirectory: outputDirectory) try jpegData.write(to: outputURL, options: .atomic) return outputURL } @@ -106,7 +106,7 @@ enum ImageUtils { return data as Data } #else - static func processImage(_ image: NSImage, maxDimension: CGFloat = 448) throws -> URL { + static func processImage(_ image: NSImage, maxDimension: CGFloat = 448, outputDirectory: URL? = nil) throws -> URL { return try autoreleasepool { let scaled = scaledImage(image, maxDimension: maxDimension) guard let inputCG = scaled.cgImage(forProposedRect: nil, context: nil, hints: nil) else { @@ -142,7 +142,7 @@ enum ImageUtils { } } } - let outputURL = try makeOutputURL() + let outputURL = try makeOutputURL(outputDirectory: outputDirectory) try jpegData.write(to: outputURL, options: .atomic) return outputURL } @@ -186,12 +186,17 @@ enum ImageUtils { } #endif - private static func makeOutputURL() throws -> URL { + private static func makeOutputURL(outputDirectory: URL? = nil) throws -> URL { let formatter = DateFormatter() formatter.dateFormat = "yyyyMMdd_HHmmss" let fileName = "img_\(formatter.string(from: Date()))_\(UUID().uuidString).jpg" - let directory = try applicationFilesDirectory().appendingPathComponent("images/outgoing", isDirectory: true) + let directory: URL + if let outputDirectory { + directory = outputDirectory + } else { + directory = try applicationFilesDirectory().appendingPathComponent("images/outgoing", isDirectory: true) + } try FileManager.default.createDirectory(at: directory, withIntermediateDirectories: true, attributes: nil) return directory.appendingPathComponent(fileName) } diff --git a/bitchat/Nostr/GeoRelayDirectory.swift b/bitchat/Nostr/GeoRelayDirectory.swift index 653fc285..a759a834 100644 --- a/bitchat/Nostr/GeoRelayDirectory.swift +++ b/bitchat/Nostr/GeoRelayDirectory.swift @@ -381,12 +381,7 @@ final class GeoRelayDirectory { if idx == 0 && line.lowercased().contains("relay url") { continue } let parts = line.split(separator: ",").map { $0.trimmed } guard parts.count >= 3 else { continue } - var host = parts[0] - host = host.replacingOccurrences(of: "https://", with: "") - host = host.replacingOccurrences(of: "http://", with: "") - host = host.replacingOccurrences(of: "wss://", with: "") - host = host.replacingOccurrences(of: "ws://", with: "") - host = host.trimmingCharacters(in: CharacterSet(charactersIn: "/")) + guard let host = NostrRelayURL.directoryAddress(parts[0]) else { continue } guard let lat = Double(parts[1]), let lon = Double(parts[2]) else { continue } result.insert(Entry(host: host, lat: lat, lon: lon)) } diff --git a/bitchat/Nostr/NostrRelayManager.swift b/bitchat/Nostr/NostrRelayManager.swift index f12d0e9e..7428fc55 100644 --- a/bitchat/Nostr/NostrRelayManager.swift +++ b/bitchat/Nostr/NostrRelayManager.swift @@ -129,7 +129,7 @@ final class NostrRelayManager: ObservableObject { "wss://offchain.pub" // For local testing, you can add: "ws://localhost:8080" ] - private static let defaultRelaySet = Set(defaultRelays) + private static let defaultRelaySet = Set(defaultRelays.compactMap { NostrRelayURL.normalized($0) }) @Published private(set) var relays: [Relay] = [] @Published private(set) var isConnected = false @@ -390,8 +390,7 @@ final class NostrRelayManager: ObservableObject { // Target specific relays if provided; else default. Filter permanently failed relays. let baseUrls = relayUrls ?? Self.defaultRelays - let candidateUrls = baseUrls.filter { !isPermanentlyFailed($0) } - let urls = allowedRelayList(from: candidateUrls) + let urls = allowedRelayList(from: baseUrls).filter { !isPermanentlyFailed($0) } let requestState = SubscriptionRequestState(messageString: messageString, relayURLs: Set(urls)) if subscriptionRequestState[id] == requestState, subscriptionStateExists(id: id, requestState: requestState) { return @@ -474,7 +473,8 @@ final class NostrRelayManager: ObservableObject { private func allowedRelayList(from urls: [String]) -> [String] { var seen = Set() var result: [String] = [] - for url in urls { + for rawURL in urls { + guard let url = NostrRelayURL.normalized(rawURL) else { continue } if !allowDefaultRelays && Self.defaultRelaySet.contains(url) { continue } if seen.insert(url).inserted { result.append(url) @@ -881,7 +881,8 @@ final class NostrRelayManager: ObservableObject { /// Manually retry connection to a specific relay func retryConnection(to relayUrl: String) { - guard let index = relays.firstIndex(where: { $0.url == relayUrl }) else { return } + let normalizedRelayUrl = NostrRelayURL.normalized(relayUrl) ?? relayUrl + guard let index = relays.firstIndex(where: { $0.url == normalizedRelayUrl }) else { return } // Reset reconnection attempts relays[index].reconnectAttempts = 0 @@ -889,13 +890,13 @@ final class NostrRelayManager: ObservableObject { relays[index].lastError = nil // Disconnect if connected - if let connection = connections[relayUrl] { + if let connection = connections[normalizedRelayUrl] { connection.cancel(with: .goingAway, reason: nil) - connections.removeValue(forKey: relayUrl) + connections.removeValue(forKey: normalizedRelayUrl) } // Attempt immediate reconnection - connectToRelay(relayUrl) + connectToRelay(normalizedRelayUrl) } /// Get detailed status for all relays diff --git a/bitchat/Nostr/NostrRelayURL.swift b/bitchat/Nostr/NostrRelayURL.swift new file mode 100644 index 00000000..be9a6fb6 --- /dev/null +++ b/bitchat/Nostr/NostrRelayURL.swift @@ -0,0 +1,51 @@ +import Foundation + +enum NostrRelayURL { + static func normalized(_ rawValue: String, defaultScheme: String? = nil) -> String? { + var value = rawValue.trimmingCharacters(in: .whitespacesAndNewlines) + guard !value.isEmpty else { return nil } + + if !value.contains("://"), let defaultScheme { + value = "\(defaultScheme)://\(value)" + } + + guard var components = URLComponents(string: value), + let rawScheme = components.scheme?.lowercased(), + let rawHost = components.host?.lowercased(), + !rawHost.isEmpty else { + return nil + } + + switch rawScheme { + case "wss", "https": + components.scheme = "wss" + if components.port == 443 { + components.port = nil + } + case "ws", "http": + components.scheme = "ws" + if components.port == 80 { + components.port = nil + } + default: + return nil + } + + components.host = rawHost + if components.path == "/" { + components.path = "" + } + components.fragment = nil + + return components.string + } + + static func directoryAddress(_ rawValue: String) -> String? { + guard var normalized = normalized(rawValue, defaultScheme: "wss") else { return nil } + for prefix in ["wss://", "ws://"] where normalized.hasPrefix(prefix) { + normalized.removeFirst(prefix.count) + break + } + return normalized + } +} diff --git a/bitchat/Services/BLE/BLEIncomingFileStore.swift b/bitchat/Services/BLE/BLEIncomingFileStore.swift new file mode 100644 index 00000000..3379cfe5 --- /dev/null +++ b/bitchat/Services/BLE/BLEIncomingFileStore.swift @@ -0,0 +1,168 @@ +import BitLogger +import BitFoundation +import Foundation + +struct BLEIncomingFileStore { + private static let quotaBytes: Int64 = 100 * 1024 * 1024 + + private let fileManager: FileManager + private let baseDirectory: URL? + private let dateProvider: () -> Date + + init(fileManager: FileManager = .default, baseDirectory: URL? = nil, dateProvider: @escaping () -> Date = Date.init) { + self.fileManager = fileManager + self.baseDirectory = baseDirectory + self.dateProvider = dateProvider + } + + func save( + data: Data, + preferredName: String?, + subdirectory: String, + fallbackExtension: String?, + defaultPrefix: String + ) -> URL? { + do { + let base = try filesDirectory().appendingPathComponent(subdirectory, isDirectory: true) + try fileManager.createDirectory(at: base, withIntermediateDirectories: true, attributes: nil) + let sanitized = sanitizedFileName( + preferredName, + defaultName: "\(defaultPrefix)_\(Self.timestampString(from: dateProvider()))", + fallbackExtension: fallbackExtension + ) + let destination = uniqueFileURL(in: base, fileName: sanitized) + try data.write(to: destination, options: .atomic) + return destination + } catch { + SecureLogger.error("❌ Failed to persist incoming media: \(error)", category: .session) + return nil + } + } + + func enforceQuota(reservingBytes: Int) { + do { + let base = try filesDirectory() + let incomingDirs = [ + base.appendingPathComponent("voicenotes/incoming", isDirectory: true), + base.appendingPathComponent("images/incoming", isDirectory: true), + base.appendingPathComponent("files/incoming", isDirectory: true) + ] + var allFiles: [(url: URL, size: Int64, modified: Date)] = [] + + for dir in incomingDirs where fileManager.fileExists(atPath: dir.path) { + guard let contents = try? fileManager.contentsOfDirectory( + at: dir, + includingPropertiesForKeys: [.fileSizeKey, .contentModificationDateKey], + options: [.skipsHiddenFiles] + ) else { continue } + + for fileURL in contents { + guard let attrs = try? fileURL.resourceValues(forKeys: [.fileSizeKey, .contentModificationDateKey]), + let size = attrs.fileSize, + let modified = attrs.contentModificationDate else { continue } + allFiles.append((url: fileURL, size: Int64(size), modified: modified)) + } + } + + let currentUsage = allFiles.reduce(0) { $0 + $1.size } + let targetUsage = Self.quotaBytes - Int64(reservingBytes) + guard currentUsage > targetUsage else { return } + + let needToFree = currentUsage - targetUsage + var freedSpace: Int64 = 0 + for file in allFiles.sorted(by: { $0.modified < $1.modified }) { + guard freedSpace < needToFree else { break } + do { + try fileManager.removeItem(at: file.url) + freedSpace += file.size + SecureLogger.debug("🗑️ BCH-01-002: Deleted old incoming file to free space: \(file.url.lastPathComponent)", category: .security) + } catch { + SecureLogger.warning("⚠️ Failed to delete old file for quota: \(error)", category: .security) + } + } + + if freedSpace > 0 { + SecureLogger.info("📊 BCH-01-002: Freed \(ByteCountFormatter.string(fromByteCount: freedSpace, countStyle: .file)) to stay within incoming files quota", category: .security) + } + } catch { + SecureLogger.warning("⚠️ Could not enforce storage quota: \(error)", category: .security) + } + } + + private func filesDirectory() throws -> URL { + let root = try baseDirectory ?? fileManager.url( + for: .applicationSupportDirectory, + in: .userDomainMask, + appropriateFor: nil, + create: true + ) + let filesDir = root.appendingPathComponent("files", isDirectory: true) + try fileManager.createDirectory(at: filesDir, withIntermediateDirectories: true, attributes: nil) + return filesDir + } + + private func sanitizedFileName(_ name: String?, defaultName: String, fallbackExtension: String?) -> String { + var candidate = (name ?? "") + .replacingOccurrences(of: "\0", with: "") + .precomposedStringWithCanonicalMapping + .replacingOccurrences(of: "/", with: "_") + .replacingOccurrences(of: "\\", with: "_") + + let invalid = CharacterSet(charactersIn: "<>:\"|?*\0").union(.controlCharacters) + candidate = candidate.components(separatedBy: invalid).joined(separator: "_").trimmed + if candidate.isEmpty { candidate = defaultName } + if candidate.hasPrefix(".") { candidate = "_" + candidate } + + if candidate.count > 120 { + let ext = (candidate as NSString).pathExtension + let base = (candidate as NSString).deletingPathExtension + candidate = ext.isEmpty + ? String(candidate.prefix(120)) + : String(base.prefix(max(10, 120 - ext.count - 1))) + "." + ext + } + + if let fallbackExtension, (candidate as NSString).pathExtension.isEmpty { + candidate += ".\(fallbackExtension)" + } + + return candidate.isEmpty ? defaultName : candidate + } + + private func uniqueFileURL(in directory: URL, fileName: String) -> URL { + let directoryPath = directory.standardizedFileURL.path + func isInsideDirectory(_ url: URL) -> Bool { + url.standardizedFileURL.path.hasPrefix(directoryPath + "/") + } + + var candidate = directory.appendingPathComponent(fileName) + guard isInsideDirectory(candidate) else { + SecureLogger.warning("⚠️ Path traversal blocked: \(fileName)", category: .security) + return directory.appendingPathComponent("blocked_\(UUID().uuidString)") + } + + if !fileManager.fileExists(atPath: candidate.path) { + return candidate + } + + let baseName = (fileName as NSString).deletingPathExtension + let ext = (fileName as NSString).pathExtension + for counter in 1..<100 { + let newName = ext.isEmpty ? "\(baseName) (\(counter))" : "\(baseName) (\(counter)).\(ext)" + candidate = directory.appendingPathComponent(newName) + guard isInsideDirectory(candidate) else { + return directory.appendingPathComponent("blocked_\(UUID().uuidString)") + } + if !fileManager.fileExists(atPath: candidate.path) { + return candidate + } + } + + return directory.appendingPathComponent("\(baseName)_\(UUID().uuidString).\(ext.isEmpty ? "dat" : ext)") + } + + private static func timestampString(from date: Date) -> String { + let formatter = DateFormatter() + formatter.dateFormat = "yyyyMMdd_HHmmss" + return formatter.string(from: date) + } +} diff --git a/bitchat/Services/BLE/BLENoisePayloadFactory.swift b/bitchat/Services/BLE/BLENoisePayloadFactory.swift new file mode 100644 index 00000000..0aac077f --- /dev/null +++ b/bitchat/Services/BLE/BLENoisePayloadFactory.swift @@ -0,0 +1,25 @@ +import Foundation + +enum BLENoisePayloadFactory { + static func privateMessage(content: String, messageID: String) -> Data? { + guard let payload = PrivateMessagePacket(messageID: messageID, content: content).encode() else { + return nil + } + + return typedPayload(.privateMessage, payload: payload) + } + + static func readReceipt(originalMessageID: String) -> Data { + typedPayload(.readReceipt, payload: Data(originalMessageID.utf8)) + } + + static func delivered(messageID: String) -> Data { + typedPayload(.delivered, payload: Data(messageID.utf8)) + } + + static func typedPayload(_ type: NoisePayloadType, payload: Data) -> Data { + var typed = Data([type.rawValue]) + typed.append(payload) + return typed + } +} diff --git a/bitchat/Services/BLE/BLENoiseSessionQueues.swift b/bitchat/Services/BLE/BLENoiseSessionQueues.swift new file mode 100644 index 00000000..84eeaba8 --- /dev/null +++ b/bitchat/Services/BLE/BLENoiseSessionQueues.swift @@ -0,0 +1,46 @@ +import BitFoundation +import Foundation + +struct BLEPendingPrivateMessage: Equatable { + let content: String + let messageID: String +} + +struct BLENoiseSessionQueues { + private var privateMessagesByPeerID: [PeerID: [BLEPendingPrivateMessage]] = [:] + private var typedPayloadsByPeerID: [PeerID: [Data]] = [:] + + var isEmpty: Bool { + privateMessagesByPeerID.isEmpty && typedPayloadsByPeerID.isEmpty + } + + mutating func removeAll() { + privateMessagesByPeerID.removeAll() + typedPayloadsByPeerID.removeAll() + } + + mutating func appendPrivateMessage(content: String, messageID: String, for peerID: PeerID) { + privateMessagesByPeerID[peerID, default: []].append(BLEPendingPrivateMessage(content: content, messageID: messageID)) + } + + mutating func takePrivateMessages(for peerID: PeerID) -> [BLEPendingPrivateMessage] { + let messages = privateMessagesByPeerID[peerID] ?? [] + privateMessagesByPeerID.removeValue(forKey: peerID) + return messages + } + + mutating func prependPrivateMessages(_ messages: [BLEPendingPrivateMessage], for peerID: PeerID) { + guard !messages.isEmpty else { return } + privateMessagesByPeerID[peerID, default: []].insert(contentsOf: messages, at: 0) + } + + mutating func appendTypedPayload(_ payload: Data, for peerID: PeerID) { + typedPayloadsByPeerID[peerID, default: []].append(payload) + } + + mutating func takeTypedPayloads(for peerID: PeerID) -> [Data] { + let payloads = typedPayloadsByPeerID[peerID] ?? [] + typedPayloadsByPeerID.removeValue(forKey: peerID) + return payloads + } +} diff --git a/bitchat/Services/BLE/BLEOutboundFragmentTransferScheduler.swift b/bitchat/Services/BLE/BLEOutboundFragmentTransferScheduler.swift new file mode 100644 index 00000000..3e9d7fe7 --- /dev/null +++ b/bitchat/Services/BLE/BLEOutboundFragmentTransferScheduler.swift @@ -0,0 +1,181 @@ +import BitFoundation +import Foundation + +struct BLEOutboundFragmentTransferRequest { + let packet: BitchatPacket + let pad: Bool + let maxChunk: Int? + let directedPeer: PeerID? + let transferId: String? + + var resolvedTransferId: String? { + guard packet.type == MessageType.fileTransfer.rawValue else { return nil } + return transferId ?? packet.payload.sha256Hex() + } +} + +struct BLEOutboundFragmentTransferScheduler { + enum QueuePosition { + case front + case back + } + + enum SubmitResult { + case start(request: BLEOutboundFragmentTransferRequest, reservedTransferId: String?) + case queued(request: BLEOutboundFragmentTransferRequest, transferId: String?, position: QueuePosition) + } + + enum CancelResult { + case active(transferId: String, workItems: [DispatchWorkItem]) + case pending(transferId: String) + case missing + } + + enum SentResult: Equatable { + case progress(sentFragments: Int, totalFragments: Int) + case complete(sentFragments: Int, totalFragments: Int) + case missing + } + + private struct ActiveTransferState { + let totalFragments: Int + var sentFragments: Int + var workItems: [DispatchWorkItem] + } + + private var activeTransfers: [String: ActiveTransferState] = [:] + private var pendingTransfers: [BLEOutboundFragmentTransferRequest] = [] + + var activeCount: Int { + activeTransfers.count + } + + var pendingCount: Int { + pendingTransfers.count + } + + mutating func removeAll() -> [(id: String, workItems: [DispatchWorkItem])] { + let active = activeTransfers.map { ($0.key, $0.value.workItems) } + activeTransfers.removeAll() + pendingTransfers.removeAll() + return active + } + + mutating func submit( + _ request: BLEOutboundFragmentTransferRequest, + maxConcurrentTransfers: Int + ) -> SubmitResult { + guard let transferId = request.resolvedTransferId else { + return .start(request: request, reservedTransferId: nil) + } + + guard activeTransfers.count < maxConcurrentTransfers else { + pendingTransfers.append(request) + return .queued(request: request, transferId: transferId, position: .back) + } + + guard activeTransfers[transferId] == nil else { + pendingTransfers.insert(request, at: 0) + return .queued(request: request, transferId: transferId, position: .front) + } + + activeTransfers[transferId] = ActiveTransferState(totalFragments: 0, sentFragments: 0, workItems: []) + return .start(request: request, reservedTransferId: transferId) + } + + mutating func activateReservedTransfer( + id transferId: String, + totalFragments: Int, + workItems: [DispatchWorkItem] + ) -> Bool { + guard activeTransfers[transferId] != nil else { return false } + activeTransfers[transferId] = ActiveTransferState( + totalFragments: totalFragments, + sentFragments: 0, + workItems: workItems + ) + return true + } + + mutating func updateWorkItems(_ workItems: [DispatchWorkItem], for transferId: String) -> Bool { + guard var state = activeTransfers[transferId] else { return false } + state.workItems = workItems + activeTransfers[transferId] = state + return true + } + + mutating func releaseReservation(_ transferId: String) -> [DispatchWorkItem]? { + activeTransfers.removeValue(forKey: transferId)?.workItems + } + + func isActive(_ transferId: String) -> Bool { + activeTransfers[transferId] != nil + } + + mutating func cancelTransfer(_ transferId: String) -> CancelResult { + if let active = activeTransfers.removeValue(forKey: transferId) { + return .active(transferId: transferId, workItems: active.workItems) + } + + if let pendingIndex = pendingTransfers.firstIndex(where: { $0.resolvedTransferId == transferId || $0.transferId == transferId }) { + pendingTransfers.remove(at: pendingIndex) + return .pending(transferId: transferId) + } + + return .missing + } + + mutating func markFragmentSent(transferId: String) -> SentResult { + guard var state = activeTransfers[transferId] else { return .missing } + + state.sentFragments = min(state.sentFragments + 1, state.totalFragments) + let isComplete = state.sentFragments >= state.totalFragments + + if isComplete { + activeTransfers.removeValue(forKey: transferId) + return .complete(sentFragments: state.sentFragments, totalFragments: state.totalFragments) + } + + activeTransfers[transferId] = state + return .progress(sentFragments: state.sentFragments, totalFragments: state.totalFragments) + } + + mutating func reservePendingStarts(maxConcurrentTransfers: Int) -> [SubmitResult] { + var availableSlots = max(0, maxConcurrentTransfers - activeTransfers.count) + guard availableSlots > 0, !pendingTransfers.isEmpty else { return [] } + + var results: [SubmitResult] = [] + var blockedFront: [BLEOutboundFragmentTransferRequest] = [] + + while availableSlots > 0, !pendingTransfers.isEmpty { + let request = pendingTransfers.removeFirst() + availableSlots -= 1 + + guard let transferId = request.resolvedTransferId else { + results.append(.start(request: request, reservedTransferId: nil)) + continue + } + + guard activeTransfers.count < maxConcurrentTransfers else { + pendingTransfers.insert(request, at: 0) + results.append(.queued(request: request, transferId: transferId, position: .front)) + break + } + + guard activeTransfers[transferId] == nil else { + blockedFront.append(request) + results.append(.queued(request: request, transferId: transferId, position: .front)) + continue + } + + activeTransfers[transferId] = ActiveTransferState(totalFragments: 0, sentFragments: 0, workItems: []) + results.append(.start(request: request, reservedTransferId: transferId)) + } + + if !blockedFront.isEmpty { + pendingTransfers.insert(contentsOf: blockedFront, at: 0) + } + + return results + } +} diff --git a/bitchat/Services/BLE/BLEOutboundPacketPolicy.swift b/bitchat/Services/BLE/BLEOutboundPacketPolicy.swift new file mode 100644 index 00000000..f7117353 --- /dev/null +++ b/bitchat/Services/BLE/BLEOutboundPacketPolicy.swift @@ -0,0 +1,43 @@ +import BitFoundation +import Foundation + +enum BLEOutboundPacketPolicy { + private static let fragmentFrameOverhead = 13 + 8 + 8 + 13 + + static func messageID(for packet: BitchatPacket) -> String { + BLEIngressLinkRegistry.messageID(for: packet) + } + + static func padsBLEFrame(for packetType: UInt8) -> Bool { + switch MessageType(rawValue: packetType) { + case .noiseEncrypted, .noiseHandshake: + return true + case .none, .announce, .message, .leave, .requestSync, .fragment, .fileTransfer: + return false + } + } + + static func priority(for packet: BitchatPacket, data: Data) -> BLEOutboundWritePriority { + guard let messageType = MessageType(rawValue: packet.type) else { return .low } + switch messageType { + case .fragment: + return .fragment(totalFragments: fragmentTotalCount(from: packet.payload)) + case .fileTransfer: + return .fileTransfer + default: + return .high + } + } + + static func fragmentChunkSize(forLinkLimit limit: Int) -> Int { + max(64, limit - fragmentFrameOverhead) + } + + private static func fragmentTotalCount(from payload: Data) -> Int { + guard payload.count >= 12 else { return Int(UInt16.max) } + let totalHigh = Int(payload[10]) + let totalLow = Int(payload[11]) + let total = (totalHigh << 8) | totalLow + return max(total, 1) + } +} diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index cf4b3e37..c10bcf02 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -79,21 +79,12 @@ final class BLEService: NSObject { // 4. Efficient Message Deduplication private let messageDeduplicator = MessageDeduplicator() private var selfBroadcastMessageIDs: [String: (id: String, timestamp: Date)] = [:] - private lazy var mediaDateFormatter: DateFormatter = { - let formatter = DateFormatter() - formatter.dateFormat = "yyyyMMdd_HHmmss" - return formatter - }() private let meshTopology = MeshTopologyTracker() // 5. Fragment Reassembly (necessary for messages > MTU) private var fragmentAssemblyBuffer = BLEFragmentAssemblyBuffer() - private struct ActiveTransferState { - let totalFragments: Int - var sentFragments: Int - var workItems: [DispatchWorkItem] - } - private var activeTransfers: [String: ActiveTransferState] = [:] + private var outboundFragmentTransfers = BLEOutboundFragmentTransferScheduler() + private let incomingFileStore = BLEIncomingFileStore() // Backoff for peripherals that recently timed out connecting private var recentConnectTimeouts: [String: Date] = [:] // Peripheral UUID -> last timeout @@ -131,10 +122,8 @@ final class BLEService: NSObject { private let bleQueue = DispatchQueue(label: "mesh.bluetooth", qos: .userInitiated) private let bleQueueKey = DispatchSpecificKey() - // Queue for messages pending handshake completion - private var pendingMessagesAfterHandshake: [PeerID: [(content: String, messageID: String)]] = [:] - // Noise typed payloads (ACKs, read receipts, etc.) pending handshake - private var pendingNoisePayloadsAfterHandshake: [PeerID: [Data]] = [:] + // Noise messages and typed payloads pending handshake completion. + private var pendingNoiseSessionQueues = BLENoiseSessionQueues() // Queue for notifications that failed due to full queue private var pendingNotifications = BLEOutboundNotificationBuffer() @@ -148,15 +137,7 @@ final class BLEService: NSObject { // Ingress link tracking for duplicate and last-hop suppression private var ingressLinks = BLEIngressLinkRegistry() - private struct PendingFragmentTransfer { - let packet: BitchatPacket - let pad: Bool - let maxChunk: Int? - let directedPeer: PeerID? - let transferId: String? - } private var pendingPeripheralWrites = BLEOutboundWriteBuffer() - private var pendingFragmentTransfers: [PendingFragmentTransfer] = [] // Debounce duplicate disconnect notifies private var recentDisconnectNotifies: [PeerID: Date] = [:] // Store-and-forward for directed messages when we have no links @@ -344,20 +325,25 @@ final class BLEService: NSObject { func resetIdentityForPanic(currentNickname: String) { messageQueue.sync(flags: .barrier) { - pendingMessagesAfterHandshake.removeAll() - pendingNoisePayloadsAfterHandshake.removeAll() + pendingNoiseSessionQueues.removeAll() } - collectionsQueue.sync(flags: .barrier) { + let cancelledTransfers = collectionsQueue.sync(flags: .barrier) { pendingPeripheralWrites.removeAll() - pendingFragmentTransfers.removeAll() pendingNotifications.removeAll() + let transfers = outboundFragmentTransfers.removeAll() fragmentAssemblyBuffer.removeAll() pendingDirectedRelays.removeAll() ingressLinks.removeAll() recentPacketTimestamps.removeAll() scheduledRelays.values.forEach { $0.cancel() } scheduledRelays.removeAll() + return transfers + } + + for entry in cancelledTransfers { + entry.workItems.forEach { $0.cancel() } + TransferProgressManager.shared.cancel(id: entry.id) } bleQueue.sync { @@ -516,7 +502,7 @@ final class BLEService: NSObject { // Send immediately to all connected peers (synchronized access to BLE state) if let data = leavePacket.toBinaryData(padding: false) { - let leavePriority = priority(for: leavePacket, data: data) + let leavePriority = BLEOutboundPacketPolicy.priority(for: leavePacket, data: data) // Snapshot BLE state under bleQueue to avoid races with delegate callbacks let (peripheralStates, centralsCount, char) = bleQueue.sync { @@ -568,13 +554,11 @@ final class BLEService: NSObject { // Clear all sessions and peers let cancelledTransfers: [(id: String, items: [DispatchWorkItem])] = collectionsQueue.sync(flags: .barrier) { - let entries = activeTransfers.map { ($0.key, $0.value.workItems) } + let entries = outboundFragmentTransfers.removeAll().map { ($0.id, $0.workItems) } peers.removeAll() fragmentAssemblyBuffer.removeAll() - activeTransfers.removeAll() // Also clear pending message queues to avoid stale state across sessions - pendingMessagesAfterHandshake.removeAll() - pendingNoisePayloadsAfterHandshake.removeAll() + pendingNoiseSessionQueues.removeAll() pendingDirectedRelays.removeAll() return entries } @@ -672,17 +656,22 @@ final class BLEService: NSObject { func cancelTransfer(_ transferId: String) { collectionsQueue.async(flags: .barrier) { [weak self] in guard let self = self else { return } - if let state = self.activeTransfers.removeValue(forKey: transferId) { - state.workItems.forEach { $0.cancel() } + + switch self.outboundFragmentTransfers.cancelTransfer(transferId) { + case let .active(id, workItems): + workItems.forEach { $0.cancel() } TransferProgressManager.shared.cancel(id: transferId) - SecureLogger.debug("🛑 Cancelled transfer \(transferId.prefix(8))…", category: .session) + SecureLogger.debug("🛑 Cancelled transfer \(id.prefix(8))…", category: .session) self.messageQueue.async { [weak self] in self?.startNextPendingTransferIfNeeded() } - } else if let pendingIndex = self.pendingFragmentTransfers.firstIndex(where: { $0.transferId == transferId }) { - self.pendingFragmentTransfers.remove(at: pendingIndex) + + case let .pending(id): TransferProgressManager.shared.cancel(id: transferId) - SecureLogger.debug("🛑 Removed pending transfer \(transferId.prefix(8))… before start", category: .session) + SecureLogger.debug("🛑 Removed pending transfer \(id.prefix(8))… before start", category: .session) + + case .missing: + break } } } @@ -748,7 +737,7 @@ final class BLEService: NSObject { // This ensures 64-hex Noise keys are converted to the canonical routing format let targetID = peerID.toShort() guard let recipientData = Data(hexString: targetID.id) else { - SecureLogger.error("❌ Invalid recipient peer ID for file transfer: \(peerID)", category: .session) + SecureLogger.error("❌ Invalid recipient peer ID for file transfer: \(peerID.id.prefix(8))…", category: .session) return } @@ -774,35 +763,22 @@ final class BLEService: NSObject { func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) { - // Create typed payload: [type byte] + [message ID] - var payload = Data([NoisePayloadType.readReceipt.rawValue]) - payload.append(contentsOf: receipt.originalMessageID.utf8) + let payload = BLENoisePayloadFactory.readReceipt(originalMessageID: receipt.originalMessageID) if noiseService.hasEstablishedSession(with: peerID) { - SecureLogger.debug("📤 Sending READ receipt for message \(receipt.originalMessageID) to \(peerID)", category: .session) + SecureLogger.debug("📤 Sending READ receipt id=\(receipt.originalMessageID.prefix(8))… to \(peerID.id.prefix(8))…", category: .session) do { - let encrypted = try noiseService.encrypt(payload, for: peerID) - let packet = BitchatPacket( - type: MessageType.noiseEncrypted.rawValue, - senderID: myPeerIDData, - recipientID: Data(hexString: peerID.id), - timestamp: UInt64(Date().timeIntervalSince1970 * 1000), - payload: encrypted, - signature: nil, - ttl: messageTTL - ) - broadcastPacket(packet) + broadcastPacket(try makeEncryptedNoisePacket(payload, to: peerID)) } catch { SecureLogger.error("Failed to send read receipt: \(error)") } } else { // Queue for after handshake and initiate if needed - collectionsQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - self.pendingNoisePayloadsAfterHandshake[peerID, default: []].append(payload) + collectionsQueue.sync(flags: .barrier) { + pendingNoiseSessionQueues.appendTypedPayload(payload, for: peerID) } if !noiseService.hasSession(with: peerID) { initiateNoiseHandshake(with: peerID) } - SecureLogger.debug("🕒 Queued READ receipt for \(peerID) until handshake completes", category: .session) + SecureLogger.debug("🕒 Queued READ receipt for \(peerID.id.prefix(8))… until handshake completes", category: .session) } } @@ -889,7 +865,7 @@ final class BLEService: NSObject { } // Encode once using a small per-type padding policy, then delegate by type - let padForBLE = padPolicy(for: packetToSend.type) + let padForBLE = BLEOutboundPacketPolicy.padsBLEFrame(for: packetToSend.type) if packetToSend.type == MessageType.fileTransfer.rawValue { sendFragmentedPacket(packetToSend, pad: padForBLE, maxChunk: nil, directedOnlyPeer: nil, transferId: transferId) return @@ -905,52 +881,39 @@ final class BLEService: NSObject { sendGenericBroadcast(packetToSend, data: data, pad: padForBLE) } - // MARK: - Broadcast helpers (single responsibility) - private func padPolicy(for type: UInt8) -> Bool { - switch MessageType(rawValue: type) { - case .noiseEncrypted, .noiseHandshake: - return true - case .none, .announce, .message, .leave, .requestSync, .fragment, .fileTransfer: - return false - } - } - private func sendEncrypted(_ packet: BitchatPacket, data: Data, pad: Bool) { guard let recipientPeerID = PeerID(hexData: packet.recipientID) else { return } var sentEncrypted = false - let outboundPriority = priority(for: packet, data: data) + let outboundPriority = BLEOutboundPacketPolicy.priority(for: packet, data: data) // Per-link limits for the specific peer - var peripheralMaxLen: Int? - if let perUUID = (DispatchQueue.getSpecific(key: bleQueueKey) != nil) ? peerToPeripheralUUID[recipientPeerID] : bleQueue.sync(execute: { peerToPeripheralUUID[recipientPeerID] }) { - if let state = (DispatchQueue.getSpecific(key: bleQueueKey) != nil) ? peripherals[perUUID] : bleQueue.sync(execute: { peripherals[perUUID] }) { - peripheralMaxLen = state.peripheral.maximumWriteValueLength(for: .withoutResponse) + let directPeripheralState: PeripheralState? = { + if DispatchQueue.getSpecific(key: bleQueueKey) != nil { + return peerToPeripheralUUID[recipientPeerID].flatMap { peripherals[$0] } } - } - var centralMaxLen: Int? - do { - let (centrals, mapping) = snapshotSubscribedCentrals() - if let central = centrals.first(where: { mapping[$0.identifier.uuidString] == recipientPeerID }) { - centralMaxLen = central.maximumUpdateValueLength + return bleQueue.sync { + peerToPeripheralUUID[recipientPeerID].flatMap { peripherals[$0] } } - } - if let pm = peripheralMaxLen, data.count > pm { - let overhead = 13 + 8 + 8 + 13 - let chunk = max(64, pm - overhead) + }() + let (centrals, mapping) = snapshotSubscribedCentrals() + let recipientCentral = centrals.first { mapping[$0.identifier.uuidString] == recipientPeerID } + + if let peripheralMaxLen = directPeripheralState?.peripheral.maximumWriteValueLength(for: .withoutResponse), + data.count > peripheralMaxLen { + let chunk = BLEOutboundPacketPolicy.fragmentChunkSize(forLinkLimit: peripheralMaxLen) sendFragmentedPacket(packet, pad: pad, maxChunk: chunk, directedOnlyPeer: recipientPeerID) return } - if let cm = centralMaxLen, data.count > cm { - let overhead = 13 + 8 + 8 + 13 - let chunk = max(64, cm - overhead) + if let centralMaxLen = recipientCentral?.maximumUpdateValueLength, + data.count > centralMaxLen { + let chunk = BLEOutboundPacketPolicy.fragmentChunkSize(forLinkLimit: centralMaxLen) sendFragmentedPacket(packet, pad: pad, maxChunk: chunk, directedOnlyPeer: recipientPeerID) return } // Direct write via peripheral link - if let peripheralUUID = (DispatchQueue.getSpecific(key: bleQueueKey) != nil) ? peerToPeripheralUUID[recipientPeerID] : bleQueue.sync(execute: { peerToPeripheralUUID[recipientPeerID] }), - let state = (DispatchQueue.getSpecific(key: bleQueueKey) != nil) ? peripherals[peripheralUUID] : bleQueue.sync(execute: { peripherals[peripheralUUID] }), + if let state = directPeripheralState, state.isConnected, let characteristic = state.characteristic { writeOrEnqueue(data, to: state.peripheral, characteristic: characteristic, priority: outboundPriority) @@ -958,12 +921,12 @@ final class BLEService: NSObject { } // Notify via central link (dual-role) - if let characteristic = characteristic, !sentEncrypted { - let (centrals, mapping) = snapshotSubscribedCentrals() - for central in centrals where mapping[central.identifier.uuidString] == recipientPeerID { - let success = peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: [central]) ?? false - if success { sentEncrypted = true; break } - enqueuePendingNotification(data: data, centrals: [central], context: "encrypted") + if let characteristic = characteristic, !sentEncrypted, let recipientCentral { + let success = peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: [recipientCentral]) ?? false + if success { + sentEncrypted = true + } else { + enqueuePendingNotification(data: data, centrals: [recipientCentral], context: "encrypted") } } @@ -1006,7 +969,7 @@ final class BLEService: NSObject { private func sendOnAllLinks(packet: BitchatPacket, data: Data, pad: Bool, directedOnlyPeer: PeerID?) { // Determine the last-hop peer/link for this message to avoid echoing back over either BLE role. - let messageID = makeMessageID(for: packet) + let messageID = BLEOutboundPacketPolicy.messageID(for: packet) let ingressRecord = collectionsQueue.sync { ingressLinks.record(for: packet) } let excludedPeerLinks = links(to: ingressRecord?.peerID) let directedPeerHint: PeerID? = { @@ -1016,7 +979,7 @@ final class BLEService: NSObject { } return nil }() - let outboundPriority = priority(for: packet, data: data) + let outboundPriority = BLEOutboundPacketPolicy.priority(for: packet, data: data) let states = snapshotPeripheralStates() var minCentralWriteLen: Int? @@ -1024,35 +987,20 @@ final class BLEService: NSObject { let m = s.peripheral.maximumWriteValueLength(for: .withoutResponse) minCentralWriteLen = minCentralWriteLen.map { min($0, m) } ?? m } - var snapshotCentrals: [CBCentral] = [] - if let _ = characteristic { - let (centrals, _) = snapshotSubscribedCentrals() - snapshotCentrals = centrals - } - var minNotifyLen: Int? - if !snapshotCentrals.isEmpty { - minNotifyLen = snapshotCentrals.map { $0.maximumUpdateValueLength }.min() - } + let subscribedCentrals = characteristic == nil ? [] : snapshotSubscribedCentrals().0 + let minNotifyLen = subscribedCentrals.map { $0.maximumUpdateValueLength }.min() + // Avoid re-fragmenting fragment packets if packet.type != MessageType.fragment.rawValue, let minLen = [minCentralWriteLen, minNotifyLen].compactMap({ $0 }).min(), data.count > minLen { - let overhead = 13 + 8 + 8 + 13 - let chunk = max(64, minLen - overhead) + let chunk = BLEOutboundPacketPolicy.fragmentChunkSize(forLinkLimit: minLen) sendFragmentedPacket(packet, pad: pad, maxChunk: chunk, directedOnlyPeer: directedOnlyPeer) return } // Build link lists and apply K-of-N fanout for broadcasts; always exclude ingress link let connectedPeripheralIDs: [String] = states.filter { $0.isConnected }.map { $0.peripheral.identifier.uuidString } - let subscribedCentrals: [CBCentral] - var centralIDs: [String] = [] - if let _ = characteristic { - let (centrals, _) = snapshotSubscribedCentrals() - subscribedCentrals = centrals - centralIDs = centrals.map { $0.identifier.uuidString } - } else { - subscribedCentrals = [] - } + let centralIDs = subscribedCentrals.map { $0.identifier.uuidString } let selectedLinks = BLEFanoutSelector.selectLinks( peripheralIDs: connectedPeripheralIDs, @@ -1102,7 +1050,7 @@ final class BLEService: NSObject { // MARK: - Directed store-and-forward private func spoolDirectedPacket(_ packet: BitchatPacket, recipientPeerID: PeerID) { - let msgID = makeMessageID(for: packet) + let msgID = BLEOutboundPacketPolicy.messageID(for: packet) collectionsQueue.async(flags: .barrier) { [weak self] in guard let self = self else { return } var byMsg = self.pendingDirectedRelays[recipientPeerID] ?? [:] @@ -1219,9 +1167,9 @@ final class BLEService: NSObject { } // BCH-01-002: Enforce storage quota before saving - enforceIncomingFilesQuota(reservingBytes: filePacket.content.count) + incomingFileStore.enforceQuota(reservingBytes: filePacket.content.count) - guard let destination = saveIncomingFile( + guard let destination = incomingFileStore.save( data: filePacket.content, preferredName: filePacket.fileName, subdirectory: "\(mime.category.mediaDir)/incoming", @@ -1255,18 +1203,20 @@ final class BLEService: NSObject { } func sendFavoriteNotification(to peerID: PeerID, isFavorite: Bool) { - SecureLogger.debug("🔔 sendFavoriteNotification called - peerID: \(peerID), isFavorite: \(isFavorite)", category: .session) + SecureLogger.debug("🔔 sendFavoriteNotification peer=\(peerID.id.prefix(8))… isFavorite=\(isFavorite)", category: .session) // Include Nostr public key in the notification var content = isFavorite ? "[FAVORITED]" : "[UNFAVORITED]" + var includesNostrIdentity = false // Add our Nostr public key if available if let myNostrIdentity = try? idBridge.getCurrentNostrIdentity() { content += ":" + myNostrIdentity.npub - SecureLogger.debug("📝 Sending favorite notification with Nostr npub: \(myNostrIdentity.npub)", category: .session) + includesNostrIdentity = true + SecureLogger.debug("📝 Favorite notification includes Nostr npub=\(myNostrIdentity.npub.prefix(16))…", category: .session) } - SecureLogger.debug("📤 Sending favorite notification to \(peerID): \(content)", category: .session) + SecureLogger.debug("📤 Sending favorite notification to \(peerID.id.prefix(8))… isFavorite=\(isFavorite) includesNostrIdentity=\(includesNostrIdentity)", category: .session) sendPrivateMessage(content, to: peerID, messageID: UUID().uuidString) } @@ -1275,34 +1225,21 @@ final class BLEService: NSObject { } func sendDeliveryAck(for messageID: String, to peerID: PeerID) { - // Create typed payload: [type byte] + [message ID] - var payload = Data([NoisePayloadType.delivered.rawValue]) - payload.append(contentsOf: messageID.utf8) + let payload = BLENoisePayloadFactory.delivered(messageID: messageID) if noiseService.hasEstablishedSession(with: peerID) { do { - let encrypted = try noiseService.encrypt(payload, for: peerID) - let packet = BitchatPacket( - type: MessageType.noiseEncrypted.rawValue, - senderID: myPeerIDData, - recipientID: Data(hexString: peerID.id), - timestamp: UInt64(Date().timeIntervalSince1970 * 1000), - payload: encrypted, - signature: nil, - ttl: messageTTL - ) - broadcastPacket(packet) + broadcastPacket(try makeEncryptedNoisePacket(payload, to: peerID)) } catch { SecureLogger.error("Failed to send delivery ACK: \(error)") } } else { // Queue for after handshake and initiate if needed - collectionsQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - self.pendingNoisePayloadsAfterHandshake[peerID, default: []].append(payload) + collectionsQueue.sync(flags: .barrier) { + pendingNoiseSessionQueues.appendTypedPayload(payload, for: peerID) } if !noiseService.hasSession(with: peerID) { initiateNoiseHandshake(with: peerID) } - SecureLogger.debug("🕒 Queued DELIVERED ack for \(peerID) until handshake completes", category: .session) + SecureLogger.debug("🕒 Queued DELIVERED ack for \(peerID.id.prefix(8))… until handshake completes", category: .session) } } @@ -1324,180 +1261,6 @@ final class BLEService: NSObject { self.deliverTransportEvent(.peerListUpdated(currentPeerIDs)) } } - - // MARK: - Helper Functions - - private func applicationFilesDirectory() throws -> URL { - let base = try FileManager.default.url(for: .applicationSupportDirectory, in: .userDomainMask, appropriateFor: nil, create: true) - let filesDir = base.appendingPathComponent("files", isDirectory: true) - try FileManager.default.createDirectory(at: filesDir, withIntermediateDirectories: true, attributes: nil) - return filesDir - } - - private func sanitizeFileName(_ name: String?, defaultName: String, fallbackExtension: String?) -> String { - var candidate = name ?? "" - - // Security: Remove null bytes (path traversal vector) - candidate = candidate.replacingOccurrences(of: "\0", with: "") - - // Security: Unicode normalization prevents fullwidth character bypass - candidate = candidate.precomposedStringWithCanonicalMapping - - // Security: Remove ALL path separators (not just strip last component) - candidate = candidate.replacingOccurrences(of: "/", with: "_") - candidate = candidate.replacingOccurrences(of: "\\", with: "_") - - // Security: Remove control characters and dangerous filesystem chars - let invalid = CharacterSet(charactersIn: "<>:\"|?*\0").union(.controlCharacters) - candidate = candidate.components(separatedBy: invalid).joined(separator: "_") - - candidate = candidate.trimmed - if candidate.isEmpty { candidate = defaultName } - - // Security: Reject dotfiles (hidden file attacks) - if candidate.hasPrefix(".") { - candidate = "_" + candidate - } - - // Truncate while preserving extension - if candidate.count > 120 { - let ext = (candidate as NSString).pathExtension - let base = (candidate as NSString).deletingPathExtension - if ext.isEmpty { - candidate = String(candidate.prefix(120)) - } else { - let maxBase = max(10, 120 - ext.count - 1) - candidate = String(base.prefix(maxBase)) + "." + ext - } - } - - if let fallbackExtension = fallbackExtension, (candidate as NSString).pathExtension.isEmpty { - candidate += ".\(fallbackExtension)" - } - - if candidate.isEmpty { candidate = defaultName } - return candidate - } - - private func uniqueFileURL(in directory: URL, fileName: String) -> URL { - var candidate = directory.appendingPathComponent(fileName) - - // Security: Validate path doesn't escape directory - if !candidate.path.hasPrefix(directory.path) { - SecureLogger.warning("⚠️ Path traversal blocked: \(fileName)", category: .security) - return directory.appendingPathComponent("blocked_\(UUID().uuidString)") - } - - if !FileManager.default.fileExists(atPath: candidate.path) { - return candidate - } - - let baseName = (fileName as NSString).deletingPathExtension - let ext = (fileName as NSString).pathExtension - var counter = 1 - - // Limit iterations to prevent DoS - while counter < 100 { - let newName = ext.isEmpty ? "\(baseName) (\(counter))" : "\(baseName) (\(counter)).\(ext)" - candidate = directory.appendingPathComponent(newName) - - // Validate each iteration - guard candidate.path.hasPrefix(directory.path) else { - return directory.appendingPathComponent("blocked_\(UUID().uuidString)") - } - - if !FileManager.default.fileExists(atPath: candidate.path) { - return candidate - } - counter += 1 - } - - // Fallback: UUID to guarantee uniqueness - return directory.appendingPathComponent("\(baseName)_\(UUID().uuidString).\(ext.isEmpty ? "dat" : ext)") - } - - private func saveIncomingFile(data: Data, preferredName: String?, subdirectory: String, fallbackExtension: String?, defaultPrefix: String) -> URL? { - do { - let base = try applicationFilesDirectory().appendingPathComponent(subdirectory, isDirectory: true) - try FileManager.default.createDirectory(at: base, withIntermediateDirectories: true, attributes: nil) - let timestamp = mediaDateFormatter.string(from: Date()) - let defaultName = "\(defaultPrefix)_\(timestamp)" - let sanitized = sanitizeFileName(preferredName, defaultName: defaultName, fallbackExtension: fallbackExtension) - let destination = uniqueFileURL(in: base, fileName: sanitized) - try data.write(to: destination, options: .atomic) - return destination - } catch { - SecureLogger.error("❌ Failed to persist incoming media: \(error)", category: .session) - return nil - } - } - - // MARK: - Storage Quota Management (BCH-01-002) - - /// Maximum total storage for incoming files (100 MB) - private static let incomingFilesQuota: Int64 = 100 * 1024 * 1024 - - /// Enforces storage quota for incoming files by deleting oldest files when quota is exceeded. - /// Call before saving a new incoming file. - private func enforceIncomingFilesQuota(reservingBytes: Int) { - do { - let base = try applicationFilesDirectory() - let incomingDirs = [ - base.appendingPathComponent("voicenotes/incoming", isDirectory: true), - base.appendingPathComponent("images/incoming", isDirectory: true), - base.appendingPathComponent("files/incoming", isDirectory: true) - ] - - // Gather all incoming files with their sizes and modification dates - var allFiles: [(url: URL, size: Int64, modified: Date)] = [] - let fileManager = FileManager.default - - for dir in incomingDirs { - guard fileManager.fileExists(atPath: dir.path) else { continue } - guard let contents = try? fileManager.contentsOfDirectory( - at: dir, - includingPropertiesForKeys: [.fileSizeKey, .contentModificationDateKey], - options: [.skipsHiddenFiles] - ) else { continue } - - for fileURL in contents { - guard let attrs = try? fileURL.resourceValues(forKeys: [.fileSizeKey, .contentModificationDateKey]), - let size = attrs.fileSize, - let modified = attrs.contentModificationDate else { continue } - allFiles.append((url: fileURL, size: Int64(size), modified: modified)) - } - } - - // Calculate current usage - let currentUsage = allFiles.reduce(0) { $0 + $1.size } - let targetUsage = Self.incomingFilesQuota - Int64(reservingBytes) - - guard currentUsage > targetUsage else { return } - - // Sort by modification date (oldest first) and delete until under quota - let sortedFiles = allFiles.sorted { $0.modified < $1.modified } - var freedSpace: Int64 = 0 - let needToFree = currentUsage - targetUsage - - for file in sortedFiles { - guard freedSpace < needToFree else { break } - do { - try fileManager.removeItem(at: file.url) - freedSpace += file.size - SecureLogger.debug("🗑️ BCH-01-002: Deleted old incoming file to free space: \(file.url.lastPathComponent)", category: .security) - } catch { - SecureLogger.warning("⚠️ Failed to delete old file for quota: \(error)", category: .security) - } - } - - if freedSpace > 0 { - SecureLogger.info("📊 BCH-01-002: Freed \(ByteCountFormatter.string(fromByteCount: freedSpace, countStyle: .file)) to stay within incoming files quota", category: .security) - } - } catch { - SecureLogger.warning("⚠️ Could not enforce storage quota: \(error)", category: .security) - } - } - private func sendAnnounce(forceSend: Bool = false) { // Throttle announces to prevent flooding let now = Date() @@ -2257,7 +2020,7 @@ extension BLEService: CBPeripheralDelegate { let senderID = PeerID(hexData: packet.senderID) if packet.type != MessageType.announce.rawValue { - SecureLogger.debug("📦 Decoded notification packet type: \(packet.type) from sender: \(senderID)", category: .session) + SecureLogger.debug("📦 Decoded notification packet type: \(packet.type) from sender: \(senderID.id.prefix(8))…", category: .session) } if packet.type == MessageType.announce.rawValue { @@ -2440,7 +2203,7 @@ extension BLEService: CBPeripheralManagerDelegate { func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didSubscribeTo characteristic: CBCharacteristic) { let centralUUID = central.identifier.uuidString - SecureLogger.debug("📥 Central subscribed: \(centralUUID)", category: .session) + SecureLogger.debug("📥 Central subscribed: \(centralUUID.prefix(8))…", category: .session) subscribedCentrals.append(central) // BCH-01-004: Rate-limit subscription-triggered announces to prevent enumeration attacks @@ -2519,7 +2282,7 @@ extension BLEService: CBPeripheralManagerDelegate { } func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didUnsubscribeFrom characteristic: CBCharacteristic) { - SecureLogger.debug("📤 Central unsubscribed: \(central.identifier.uuidString)", category: .session) + SecureLogger.debug("📤 Central unsubscribed: \(central.identifier.uuidString.prefix(8))…", category: .session) subscribedCentrals.removeAll { $0.identifier == central.identifier } // Ensure we're still advertising for other devices to find us @@ -2648,7 +2411,7 @@ extension BLEService: CBPeripheralManagerDelegate { case let .oversized(metadata): logAccumulatedCentralWrite(metadata, centralUUID: centralUUID) - SecureLogger.warning("⚠️ Dropping oversized pending write buffer (\(metadata.accumulatedBytes) bytes) for central \(centralUUID)", category: .session) + SecureLogger.warning("⚠️ Dropping oversized pending write buffer (\(metadata.accumulatedBytes) bytes) for central \(centralUUID.prefix(8))…", category: .session) logFailedSingleWriteIfNeeded(hasMultiple: hasMultiple, sortedRequests: sorted) } } @@ -2659,7 +2422,7 @@ extension BLEService: CBPeripheralManagerDelegate { packetType != MessageType.announce.rawValue else { return } SecureLogger.debug( - "📥 Accumulated write from central \(centralUUID): size=\(metadata.accumulatedBytes) (+\(metadata.appendedBytes)) bytes (type=\(packetType)), offsets=\(metadata.offsets)", + "📥 Accumulated write from central \(centralUUID.prefix(8))…: size=\(metadata.accumulatedBytes) (+\(metadata.appendedBytes)) bytes (type=\(packetType)), offsets=\(metadata.offsets)", category: .session ) } @@ -2686,7 +2449,7 @@ extension BLEService: CBPeripheralManagerDelegate { } if packet.type != MessageType.announce.rawValue { - SecureLogger.debug("📦 Decoded (combined) packet type: \(packet.type) from sender: \(claimedSenderID)", category: .session) + SecureLogger.debug("📦 Decoded (combined) packet type: \(packet.type) from sender: \(claimedSenderID.id.prefix(8))…", category: .session) } if !subscribedCentrals.contains(central) { @@ -2936,7 +2699,7 @@ extension BLEService { private func configureNoiseServiceCallbacks(for service: NoiseEncryptionService) { service.onPeerAuthenticated = { [weak self] peerID, fingerprint in - SecureLogger.debug("🔐 Noise session authenticated with \(peerID), fingerprint: \(fingerprint.prefix(16))...") + SecureLogger.debug("🔐 Noise session authenticated with \(peerID.id.prefix(8))…, fingerprint: \(fingerprint.prefix(16))…") self?.messageQueue.async { [weak self] in self?.sendPendingMessagesAfterHandshake(for: peerID) self?.sendPendingNoisePayloadsAfterHandshake(for: peerID) @@ -2961,31 +2724,31 @@ extension BLEService { // No session yet - queue the payload SYNCHRONOUSLY before initiating handshake // to prevent race where fast handshake completion drains empty queue collectionsQueue.sync(flags: .barrier) { - if self.pendingNoisePayloadsAfterHandshake[peerID] == nil { - self.pendingNoisePayloadsAfterHandshake[peerID] = [] - } - self.pendingNoisePayloadsAfterHandshake[peerID]?.append(typedPayload) - SecureLogger.debug("📥 Queued noise payload for \(peerID) pending handshake", category: .session) + self.pendingNoiseSessionQueues.appendTypedPayload(typedPayload, for: peerID) + SecureLogger.debug("📥 Queued noise payload for \(peerID.id.prefix(8))… pending handshake", category: .session) } initiateNoiseHandshake(with: peerID) return } do { - let encrypted = try noiseService.encrypt(typedPayload, for: peerID) - let packet = BitchatPacket( - type: MessageType.noiseEncrypted.rawValue, - senderID: myPeerIDData, - recipientID: Data(hexString: peerID.id), - timestamp: UInt64(Date().timeIntervalSince1970 * 1000), - payload: encrypted, - signature: nil, - ttl: messageTTL - ) - broadcastPacket(packet) + broadcastPacket(try makeEncryptedNoisePacket(typedPayload, to: peerID)) } catch { SecureLogger.error("Failed to send verification payload: \(error)") } } + + private func makeEncryptedNoisePacket(_ typedPayload: Data, to peerID: PeerID) throws -> BitchatPacket { + let encrypted = try noiseService.encrypt(typedPayload, for: peerID) + return BitchatPacket( + type: MessageType.noiseEncrypted.rawValue, + senderID: myPeerIDData, + recipientID: Data(hexString: peerID.id), + timestamp: UInt64(Date().timeIntervalSince1970 * 1000), + payload: encrypted, + signature: nil, + ttl: messageTTL + ) + } // MARK: Link capability snapshots (thread-safe via bleQueue) @@ -3006,31 +2769,6 @@ extension BLEService { // MARK: Helpers: IDs, selection, and write backpressure - private func makeMessageID(for packet: BitchatPacket) -> String { - BLEIngressLinkRegistry.messageID(for: packet) - } - - 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 BLEOutboundWritePriority.fragment(totalFragments: total) - case .fileTransfer: - return .fileTransfer - default: - return .high - } - } - - private func fragmentTotalCount(from payload: Data) -> Int { - guard payload.count >= 12 else { return Int(UInt16.max) } - let totalHigh = Int(payload[10]) - let totalLow = Int(payload[11]) - let total = (totalHigh << 8) | totalLow - return max(total, 1) - } - 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 @@ -3139,52 +2877,18 @@ extension BLEService { // MARK: Private Message Handling private func sendPrivateMessage(_ content: String, to recipientID: PeerID, messageID: String) { - SecureLogger.debug("📨 Sending PM to \(recipientID): \(content.prefix(30))...", category: .session) + SecureLogger.debug("📨 Sending PM to \(recipientID.id.prefix(8))… id=\(messageID.prefix(8))… chars=\(content.count) bytes=\(content.utf8.count)", category: .session) // Check if we have an established Noise session if noiseService.hasEstablishedSession(with: recipientID) { // Encrypt and send do { - // Create TLV-encoded private message - let privateMessage = PrivateMessagePacket(messageID: messageID, content: content) - guard let tlvData = privateMessage.encode() else { + guard let messagePayload = BLENoisePayloadFactory.privateMessage(content: content, messageID: messageID) else { SecureLogger.error("Failed to encode private message with TLV") return } - // Create message payload with TLV: [type byte] + [TLV data] - var messagePayload = Data([NoisePayloadType.privateMessage.rawValue]) - messagePayload.append(tlvData) - - let encrypted = try noiseService.encrypt(messagePayload, for: recipientID) - - // Convert recipientID to Data (assuming it's a hex string) - var recipientData = Data() - var tempID = recipientID.id - while tempID.count >= 2 { - let hexByte = String(tempID.prefix(2)) - if let byte = UInt8(hexByte, radix: 16) { - recipientData.append(byte) - } - tempID = String(tempID.dropFirst(2)) - } - if tempID.count == 1 { - if let byte = UInt8(tempID, radix: 16) { - recipientData.append(byte) - } - } - - let packet = BitchatPacket( - type: MessageType.noiseEncrypted.rawValue, - senderID: myPeerIDData, - recipientID: recipientData, - timestamp: UInt64(Date().timeIntervalSince1970 * 1000), - payload: encrypted, - signature: nil, - ttl: messageTTL - ) - - broadcastPacket(packet) + broadcastPacket(try makeEncryptedNoisePacket(messagePayload, to: recipientID)) // Notify delegate that message was sent notifyUI { [weak self] in @@ -3195,14 +2899,11 @@ extension BLEService { } } else { // Queue message for sending after handshake completes - SecureLogger.debug("🤝 No session with \(recipientID), initiating handshake and queueing message", category: .session) + SecureLogger.debug("🤝 No session with \(recipientID.id.prefix(8))…, initiating handshake and queueing message", category: .session) // Queue the message (especially important for favorite notifications) collectionsQueue.sync(flags: .barrier) { - if pendingMessagesAfterHandshake[recipientID] == nil { - pendingMessagesAfterHandshake[recipientID] = [] - } - pendingMessagesAfterHandshake[recipientID]?.append((content, messageID)) + pendingNoiseSessionQueues.appendPrivateMessage(content: content, messageID: messageID, for: recipientID) } initiateNoiseHandshake(with: recipientID) @@ -3239,61 +2940,43 @@ extension BLEService { private func sendPendingMessagesAfterHandshake(for peerID: PeerID) { // Atomically take all pending messages to process (prevents concurrent modification) - let pendingMessages = collectionsQueue.sync(flags: .barrier) { () -> [(content: String, messageID: String)]? in - let messages = pendingMessagesAfterHandshake[peerID] - pendingMessagesAfterHandshake.removeValue(forKey: peerID) - return messages + let pendingMessages = collectionsQueue.sync(flags: .barrier) { () -> [BLEPendingPrivateMessage] in + pendingNoiseSessionQueues.takePrivateMessages(for: peerID) } - guard let messages = pendingMessages, !messages.isEmpty else { return } + guard !pendingMessages.isEmpty else { return } - SecureLogger.debug("📤 Sending \(messages.count) pending messages after handshake to \(peerID)", category: .session) + SecureLogger.debug("📤 Sending \(pendingMessages.count) pending messages after handshake to \(peerID.id.prefix(8))…", category: .session) // Track failed messages for re-queuing - var failedMessages: [(content: String, messageID: String)] = [] + var failedMessages: [BLEPendingPrivateMessage] = [] // Send each pending message directly (we know session is established) - for (content, messageID) in messages { + for message in pendingMessages { do { // Use the same TLV format as normal sends to keep receiver decoding consistent - let privateMessage = PrivateMessagePacket(messageID: messageID, content: content) - guard let tlvData = privateMessage.encode() else { + guard let messagePayload = BLENoisePayloadFactory.privateMessage(content: message.content, messageID: message.messageID) else { SecureLogger.error("Failed to encode pending private message TLV") - failedMessages.append((content, messageID)) + failedMessages.append(message) continue } - var messagePayload = Data([NoisePayloadType.privateMessage.rawValue]) - messagePayload.append(tlvData) - - let encrypted = try noiseService.encrypt(messagePayload, for: peerID) - - let packet = BitchatPacket( - type: MessageType.noiseEncrypted.rawValue, - senderID: myPeerIDData, - recipientID: Data(hexString: peerID.id), - timestamp: UInt64(Date().timeIntervalSince1970 * 1000), - payload: encrypted, - signature: nil, - ttl: messageTTL - ) - // We're already on messageQueue from the callback - broadcastPacket(packet) + broadcastPacket(try makeEncryptedNoisePacket(messagePayload, to: peerID)) // Notify delegate that message was sent notifyUI { [weak self] in - self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: messageID, status: .sent)) + self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: message.messageID, status: .sent)) } - SecureLogger.debug("✅ Sent pending message \(messageID) to \(peerID) after handshake", category: .session) + SecureLogger.debug("✅ Sent pending message id=\(message.messageID.prefix(8))… to \(peerID.id.prefix(8))… after handshake", category: .session) } catch { SecureLogger.error("Failed to send pending message after handshake: \(error)") - failedMessages.append((content, messageID)) + failedMessages.append(message) // Notify delegate of failure notifyUI { [weak self] in - self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: messageID, status: .failed(reason: "Encryption failed"))) + self?.deliverTransportEvent(.messageDeliveryStatusUpdated(messageID: message.messageID, status: .failed(reason: "Encryption failed"))) } } } @@ -3302,12 +2985,9 @@ extension BLEService { if !failedMessages.isEmpty { collectionsQueue.async(flags: .barrier) { [weak self] in guard let self = self else { return } - if self.pendingMessagesAfterHandshake[peerID] == nil { - self.pendingMessagesAfterHandshake[peerID] = [] - } // Prepend failed messages to maintain order - self.pendingMessagesAfterHandshake[peerID]?.insert(contentsOf: failedMessages, at: 0) - SecureLogger.warning("⚠️ Re-queued \(failedMessages.count) failed messages for \(peerID)", category: .session) + self.pendingNoiseSessionQueues.prependPrivateMessages(failedMessages, for: peerID) + SecureLogger.warning("⚠️ Re-queued \(failedMessages.count) failed messages for \(peerID.id.prefix(8))…", category: .session) } } } @@ -3315,68 +2995,49 @@ extension BLEService { // MARK: Fragmentation (Required for messages > BLE MTU) private func sendFragmentedPacket(_ packet: BitchatPacket, pad: Bool, maxChunk: Int? = nil, directedOnlyPeer: PeerID? = nil, transferId: String? = nil) { - let context = PendingFragmentTransfer(packet: packet, pad: pad, maxChunk: maxChunk, directedPeer: directedOnlyPeer, transferId: transferId) - if packet.type == MessageType.fileTransfer.rawValue { - let shouldQueue = collectionsQueue.sync { - self.activeTransfers.count >= TransportConfig.bleMaxConcurrentTransfers - } - if shouldQueue { - queueFragmentTransfer(context, prioritizeFront: false) - return - } + let request = BLEOutboundFragmentTransferRequest( + packet: packet, + pad: pad, + maxChunk: maxChunk, + directedPeer: directedOnlyPeer, + transferId: transferId + ) + + let result = collectionsQueue.sync(flags: .barrier) { + outboundFragmentTransfers.submit(request, maxConcurrentTransfers: TransportConfig.bleMaxConcurrentTransfers) } - startFragmentedPacket(context) + handleFragmentTransferSubmitResult(result) } - private func queueFragmentTransfer(_ context: PendingFragmentTransfer, prioritizeFront: Bool) { - collectionsQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - if prioritizeFront { - self.pendingFragmentTransfers.insert(context, at: 0) + private func handleFragmentTransferSubmitResult(_ result: BLEOutboundFragmentTransferScheduler.SubmitResult) { + switch result { + case let .start(request, reservedTransferId): + startFragmentedPacket(request, reservedTransferId: reservedTransferId) + + case let .queued(_, transferId, _): + if let transferId { + SecureLogger.debug("🚦 Queued media transfer \(transferId.prefix(8))… waiting for slot", category: .session) } else { - self.pendingFragmentTransfers.append(context) + SecureLogger.debug("🚦 Queued fragment transfer waiting for slot", category: .session) } } - if let transferId = context.transferId { - SecureLogger.debug("🚦 Queued media transfer \(transferId.prefix(8))… waiting for slot", category: .session) - } else { - SecureLogger.debug("🚦 Queued fragment transfer waiting for slot", category: .session) - } } - private func startFragmentedPacket(_ context: PendingFragmentTransfer) { - let packet = context.packet - let isFileTransfer = packet.type == MessageType.fileTransfer.rawValue - var reservedTransferId: String? + private func startFragmentedPacket(_ request: BLEOutboundFragmentTransferRequest, reservedTransferId: String?) { + let packet = request.packet - let releaseReservedSlot: (String) -> Void = { id in + let releaseReservedSlot: (String) -> Void = { [weak self] id in + guard let self = self else { return } TransferProgressManager.shared.cancel(id: id) self.collectionsQueue.async(flags: .barrier) { [weak self] in - self?.activeTransfers.removeValue(forKey: id) + _ = self?.outboundFragmentTransfers.releaseReservation(id) } self.messageQueue.async { [weak self] in self?.startNextPendingTransferIfNeeded() } } - if isFileTransfer { - let candidateId = context.transferId ?? packet.payload.sha256Hex() - var didReserve = false - collectionsQueue.sync(flags: .barrier) { - if self.activeTransfers.count < TransportConfig.bleMaxConcurrentTransfers, - self.activeTransfers[candidateId] == nil { - self.activeTransfers[candidateId] = ActiveTransferState(totalFragments: 0, sentFragments: 0, workItems: []) - didReserve = true - } - } - guard didReserve else { - queueFragmentTransfer(context, prioritizeFront: true) - return - } - reservedTransferId = candidateId - } - - guard let fullData = packet.toBinaryData(padding: context.pad) else { + guard let fullData = packet.toBinaryData(padding: request.pad) else { if let id = reservedTransferId { releaseReservedSlot(id) } @@ -3398,7 +3059,7 @@ extension BLEService { calculatedChunk = max(64, bleMaxMTU - overhead) } - let chunk = context.maxChunk ?? calculatedChunk + let chunk = request.maxChunk ?? calculatedChunk let safeChunk = max(64, chunk) let fragments = stride(from: 0, to: fullData.count, by: safeChunk).map { offset in Data(fullData[offset..= state.totalFragments - if isComplete { - self.activeTransfers.removeValue(forKey: transferId) - } else { - self.activeTransfers[transferId] = state + guard let self = self else { return } + + switch self.outboundFragmentTransfers.markFragmentSent(transferId: transferId) { + case .progress, .complete: + TransferProgressManager.shared.recordFragmentSent(id: transferId) + + case .missing: + return } - TransferProgressManager.shared.recordFragmentSent(id: transferId) - if isComplete { + + if !self.outboundFragmentTransfers.isActive(transferId) { self.messageQueue.async { [weak self] in self?.startNextPendingTransferIfNeeded() } @@ -3517,20 +3177,13 @@ extension BLEService { } private func startNextPendingTransferIfNeeded() { - collectionsQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - let limit = TransportConfig.bleMaxConcurrentTransfers - var availableSlots = max(0, limit - self.activeTransfers.count) - guard availableSlots > 0, !self.pendingFragmentTransfers.isEmpty else { return } - var toStart: [PendingFragmentTransfer] = [] - while availableSlots > 0, !self.pendingFragmentTransfers.isEmpty { - toStart.append(self.pendingFragmentTransfers.removeFirst()) - availableSlots -= 1 - } - for context in toStart { - self.messageQueue.async { [weak self] in - self?.startFragmentedPacket(context) - } + let results = collectionsQueue.sync(flags: .barrier) { + outboundFragmentTransfers.reservePendingStarts(maxConcurrentTransfers: TransportConfig.bleMaxConcurrentTransfers) + } + + for result in results { + messageQueue.async { [weak self] in + self?.handleFragmentTransferSubmitResult(result) } } } @@ -3627,7 +3280,7 @@ extension BLEService { // Only log non-announce packets to reduce noise if packet.type != MessageType.announce.rawValue { // Log packet details for debugging - SecureLogger.debug("📦 Handling packet type \(packet.type) from \(senderID), messageID: \(messageID)", category: .session) + SecureLogger.debug("📦 Handling packet type \(packet.type) from \(senderID.id.prefix(8))…, messageID: \(messageID.prefix(24))…", category: .session) } // Efficient deduplication @@ -3639,7 +3292,7 @@ extension BLEService { // Announce packets (type 1) are sent every 10 seconds for peer discovery // It's normal to see these as duplicates - don't log them to reduce noise if packet.type != MessageType.announce.rawValue { - SecureLogger.debug("⚠️ Duplicate packet ignored: \(messageID)", category: .session) + SecureLogger.debug("⚠️ Duplicate packet ignored: \(messageID.prefix(24))…", category: .session) } // In sparse graphs (<=2 neighbors), keep the pending relay to ensure bridging. // In denser graphs, cancel the pending relay to reduce redundant floods. @@ -3744,7 +3397,7 @@ extension BLEService { private func handleAnnounce(_ packet: BitchatPacket, from peerID: PeerID) { guard let announcement = AnnouncementPacket.decode(from: packet.payload) else { - SecureLogger.error("❌ Failed to decode announce packet from \(peerID)", category: .session) + SecureLogger.error("❌ Failed to decode announce packet from \(peerID.id.prefix(8))…", category: .session) return } @@ -3865,7 +3518,7 @@ extension BLEService { lastReconnectLogAt[peerID] = now } } else if existingPeer?.nickname != announcement.nickname { - SecureLogger.debug("🔄 Peer \(peerID) changed nickname: \(existingPeer?.nickname ?? "Unknown") -> \(announcement.nickname)", category: .session) + SecureLogger.debug("🔄 Peer \(peerID.id.prefix(8))… changed nickname: \(existingPeer?.nickname ?? "Unknown") -> \(announcement.nickname)", category: .session) } } } @@ -3929,7 +3582,7 @@ extension BLEService { // Handle REQUEST_SYNC: decode payload and respond with missing packets via sync manager private func handleRequestSync(_ packet: BitchatPacket, from peerID: PeerID) { guard let req = RequestSyncPacket.decode(from: packet.payload) else { - SecureLogger.warning("⚠️ Malformed REQUEST_SYNC from \(peerID)", category: .session) + SecureLogger.warning("⚠️ Malformed REQUEST_SYNC from \(peerID.id.prefix(8))…", category: .session) return } gossipSyncManager?.handleRequestSync(from: peerID, request: req) @@ -4025,7 +3678,7 @@ extension BLEService { let hasDirectLink = directLink.hasPeripheral || directLink.hasCentral let pathTag = hasDirectLink ? "direct" : "mesh" - SecureLogger.debug("💬 [\(senderNickname)] TTL:\(packet.ttl) (\(pathTag)): \(String(content.prefix(50)))\(content.count > 50 ? "..." : "")", category: .session) + SecureLogger.debug("💬 [\(senderNickname)] TTL:\(packet.ttl) (\(pathTag)) chars=\(content.count) bytes=\(packet.payload.count)", category: .session) let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000) var resolvedSelfMessageID: String? = nil @@ -4086,7 +3739,7 @@ extension BLEService { } if recipientID != myPeerID { - SecureLogger.debug("🔐 Encrypted message not for me (for \(recipientID), I am \(myPeerID))", category: .session) + SecureLogger.debug("🔐 Encrypted message not for me (for \(recipientID.id.prefix(8))…, I am \(myPeerID.id.prefix(8))…)", category: .session) return } @@ -4106,7 +3759,7 @@ extension BLEService { return } - SecureLogger.debug("🔐 Decrypted noise payload type \(noisePayloadType.description) from \(peerID)", category: .session) + SecureLogger.debug("🔐 Decrypted noise payload type \(noisePayloadType.description) from \(peerID.id.prefix(8))…", category: .session) switch noisePayloadType { case .privateMessage: @@ -4138,14 +3791,14 @@ extension BLEService { } catch NoiseEncryptionError.sessionNotEstablished { // We received an encrypted message before establishing a session with this peer. // Trigger a handshake so future messages can be decrypted. - SecureLogger.debug("🔑 Encrypted message from \(peerID) without session; initiating handshake") + SecureLogger.debug("🔑 Encrypted message from \(peerID.id.prefix(8))… without session; initiating handshake") if !noiseService.hasSession(with: peerID) { initiateNoiseHandshake(with: peerID) } } catch { // Decryption failed - clear the corrupted session and re-initiate handshake // This handles cases where session state got out of sync (nonce mismatch, etc.) - SecureLogger.error("❌ Failed to decrypt message from \(peerID): \(error) - clearing session and re-initiating handshake") + SecureLogger.error("❌ Failed to decrypt message from \(peerID.id.prefix(8))…: \(error) - clearing session and re-initiating handshake") noiseService.clearSession(for: peerID) initiateNoiseHandshake(with: peerID) } @@ -4155,27 +3808,15 @@ extension BLEService { private func sendPendingNoisePayloadsAfterHandshake(for peerID: PeerID) { let payloads = collectionsQueue.sync(flags: .barrier) { () -> [Data] in - let list = pendingNoisePayloadsAfterHandshake[peerID] ?? [] - pendingNoisePayloadsAfterHandshake.removeValue(forKey: peerID) - return list + pendingNoiseSessionQueues.takeTypedPayloads(for: peerID) } guard !payloads.isEmpty else { return } - SecureLogger.debug("📤 Sending \(payloads.count) pending noise payloads to \(peerID) after handshake", category: .session) + SecureLogger.debug("📤 Sending \(payloads.count) pending noise payloads to \(peerID.id.prefix(8))… after handshake", category: .session) for payload in payloads { do { - let encrypted = try noiseService.encrypt(payload, for: peerID) - let packet = BitchatPacket( - type: MessageType.noiseEncrypted.rawValue, - senderID: myPeerIDData, - recipientID: Data(hexString: peerID.id), - timestamp: UInt64(Date().timeIntervalSince1970 * 1000), - payload: encrypted, - signature: nil, - ttl: messageTTL - ) - broadcastPacket(packet) + broadcastPacket(try makeEncryptedNoisePacket(payload, to: peerID)) } catch { - SecureLogger.error("❌ Failed to send pending noise payload to \(peerID): \(error)") + SecureLogger.error("❌ Failed to send pending noise payload to \(peerID.id.prefix(8))…: \(error)") } } } @@ -4334,7 +3975,7 @@ extension BLEService { // Cleanup: remove peers that are not connected and past reachability retention if !peer.isConnected { if age > retention { - SecureLogger.debug("🗑️ Removing stale peer after reachability window: \(peerID) (\(peer.nickname))", category: .session) + SecureLogger.debug("🗑️ Removing stale peer after reachability window: \(peerID.id.prefix(8))… (\(peer.nickname))", category: .session) // Also remove any stored announcement from sync candidates gossipSyncManager?.removeAnnouncementForPeer(peerID) peers.removeValue(forKey: peerID) diff --git a/bitchat/ViewModels/ChatDeliveryCoordinator.swift b/bitchat/ViewModels/ChatDeliveryCoordinator.swift index 59f3ba74..19f3a333 100644 --- a/bitchat/ViewModels/ChatDeliveryCoordinator.swift +++ b/bitchat/ViewModels/ChatDeliveryCoordinator.swift @@ -44,7 +44,23 @@ final class ChatDeliveryCoordinator { } @MainActor - func updateMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) { + func deliveryStatus(for messageID: String) -> DeliveryStatus? { + if let message = viewModel.messages.first(where: { $0.id == messageID }) { + return message.deliveryStatus + } + + for messages in viewModel.privateChats.values { + if let message = messages.first(where: { $0.id == messageID }) { + return message.deliveryStatus + } + } + + return nil + } + + @MainActor + @discardableResult + func updateMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool { var didUpdateStatus = false if let index = viewModel.messages.firstIndex(where: { $0.id == messageID }) { @@ -72,6 +88,8 @@ final class ChatDeliveryCoordinator { viewModel.privateChats = privateChats viewModel.objectWillChange.send() } + + return didUpdateStatus } } diff --git a/bitchat/ViewModels/ChatTransportEventCoordinator.swift b/bitchat/ViewModels/ChatTransportEventCoordinator.swift index 028942b2..a671f7b3 100644 --- a/bitchat/ViewModels/ChatTransportEventCoordinator.swift +++ b/bitchat/ViewModels/ChatTransportEventCoordinator.swift @@ -209,34 +209,34 @@ private extension ChatTransportEventCoordinator { viewModel.meshService.sendDeliveryAck(for: packet.messageID, to: peerID) case .delivered: - guard let messageID = String(data: payload, encoding: .utf8), - let name = viewModel.unifiedPeerService.getPeer(by: peerID)?.nickname, - let (foundPeerID, index) = findMessageIndex( - for: messageID, - peerID: peerID, - in: viewModel - ) else { return } + guard let messageID = String(data: payload, encoding: .utf8) else { return } - if case .read = viewModel.privateChats[foundPeerID]?[index].deliveryStatus { return } + let name = deliveryStatusName(for: peerID, in: viewModel) + let didUpdate = viewModel.deliveryCoordinator.updateMessageDeliveryStatus( + messageID, + status: .delivered(to: name, at: Date()) + ) - viewModel.privateChats[foundPeerID]?[index].deliveryStatus = .delivered(to: name, at: Date()) - viewModel.objectWillChange.send() + if !didUpdate { + if case .read? = viewModel.deliveryCoordinator.deliveryStatus(for: messageID) { + SecureLogger.debug("📬 Ignored stale delivered ACK for already-read message id=\(messageID.prefix(8))… from \(peerID.id.prefix(8))…", category: .session) + } else { + SecureLogger.debug("📬 Delivered ACK for unknown message id=\(messageID.prefix(8))… from \(peerID.id.prefix(8))…", category: .session) + } + } case .readReceipt: - guard let messageID = String(data: payload, encoding: .utf8), - let name = viewModel.unifiedPeerService.getPeer(by: peerID)?.nickname, - let (foundPeerID, index) = findMessageIndex( - for: messageID, - peerID: peerID, - in: viewModel - ), - let messages = viewModel.privateChats[foundPeerID], - index < messages.count else { return } + guard let messageID = String(data: payload, encoding: .utf8) else { return } - messages[index].deliveryStatus = .read(by: name, at: Date()) - viewModel.privateChats[foundPeerID] = messages - viewModel.privateChatManager.objectWillChange.send() - viewModel.objectWillChange.send() + let name = deliveryStatusName(for: peerID, in: viewModel) + let didUpdate = viewModel.deliveryCoordinator.updateMessageDeliveryStatus( + messageID, + status: .read(by: name, at: Date()) + ) + + if !didUpdate { + SecureLogger.debug("📖 Read receipt for unknown message id=\(messageID.prefix(8))… from \(peerID.id.prefix(8))…", category: .session) + } case .verifyChallenge: viewModel.verificationCoordinator.handleVerifyChallengePayload(from: peerID, payload: payload) @@ -247,34 +247,7 @@ private extension ChatTransportEventCoordinator { } @MainActor - func findMessageIndex( - for messageID: String, - peerID: PeerID, - in viewModel: ChatViewModel - ) -> (peerID: PeerID, index: Int)? { - if let messages = viewModel.privateChats[peerID], - let index = messages.firstIndex(where: { $0.id == messageID }) { - return (peerID, index) - } - - if peerID.bare.count == 16, - let peer = viewModel.unifiedPeerService.getPeer(by: peerID), - !peer.noisePublicKey.isEmpty { - let longID = PeerID(hexData: peer.noisePublicKey) - if let messages = viewModel.privateChats[longID], - let index = messages.firstIndex(where: { $0.id == messageID }) { - return (longID, index) - } - } - - if peerID.bare.count == 64 { - let shortID = peerID.toShort() - if let messages = viewModel.privateChats[shortID], - let index = messages.firstIndex(where: { $0.id == messageID }) { - return (shortID, index) - } - } - - return nil + func deliveryStatusName(for peerID: PeerID, in viewModel: ChatViewModel) -> String { + viewModel.unifiedPeerService.getPeer(by: peerID)?.nickname ?? viewModel.resolveNickname(for: peerID) } } diff --git a/bitchatTests/ChatViewModelTests.swift b/bitchatTests/ChatViewModelTests.swift index 19737944..304b96d5 100644 --- a/bitchatTests/ChatViewModelTests.swift +++ b/bitchatTests/ChatViewModelTests.swift @@ -516,6 +516,57 @@ struct ChatViewModelNoisePayloadTests { #expect(delivered) } + + @Test @MainActor + func didReceiveNoisePayload_readReceiptUpdatesBeforePeerNicknameIsKnown() async { + let (viewModel, _) = makeTestableViewModel() + let peerID = PeerID(str: "0000000000000005") + + let message = BitchatMessage( + id: "pm-read-before-name", + sender: viewModel.nickname, + content: "Waiting on read receipt", + timestamp: Date(), + isRelay: false, + originalSender: nil, + isPrivate: true, + recipientNickname: "Peer", + senderPeerID: viewModel.meshService.myPeerID, + mentions: nil, + deliveryStatus: .sent + ) + viewModel.privateChats[peerID] = [message] + + viewModel.didReceiveNoisePayload( + from: peerID, + type: .readReceipt, + payload: Data("pm-read-before-name".utf8), + timestamp: Date() + ) + + let privateChatUpdated = await TestHelpers.waitUntil({ + guard let status = viewModel.privateChats[peerID]?.first?.deliveryStatus else { return false } + if case .read = status { + return true + } + return false + }, timeout: TestConstants.defaultTimeout) + + let conversationStoreUpdated = await TestHelpers.waitUntil({ + let messages = viewModel.conversationStore.directMessages( + for: peerID, + identityResolver: viewModel.identityResolver + ) + guard let status = messages.first?.deliveryStatus else { return false } + if case .read = status { + return true + } + return false + }, timeout: TestConstants.defaultTimeout) + + #expect(privateChatUpdated) + #expect(conversationStoreUpdated) + } } // MARK: - Formatting Tests diff --git a/bitchatTests/Features/ImageUtilsTests.swift b/bitchatTests/Features/ImageUtilsTests.swift index 08273244..a130b4ea 100644 --- a/bitchatTests/Features/ImageUtilsTests.swift +++ b/bitchatTests/Features/ImageUtilsTests.swift @@ -11,6 +11,10 @@ private func makeTemporaryFileURL(_ name: String) -> URL { FileManager.default.temporaryDirectory.appendingPathComponent(name) } +private func makeTemporaryDirectoryURL(_ name: String) -> URL { + FileManager.default.temporaryDirectory.appendingPathComponent(name, isDirectory: true) +} + #if os(iOS) private func makePlatformImage(size: CGSize) -> UIImage { UIGraphicsImageRenderer(size: size).image { context in @@ -55,11 +59,14 @@ struct ImageUtilsTests { @Test func processImage_writesCompressedJpeg() throws { let image = makePlatformImage(size: CGSize(width: 1024, height: 768)) - let outputURL = try ImageUtils.processImage(image, maxDimension: 256) - defer { try? FileManager.default.removeItem(at: outputURL) } + let outputDirectory = makeTemporaryDirectoryURL("image-output-\(UUID().uuidString)") + defer { try? FileManager.default.removeItem(at: outputDirectory) } + + let outputURL = try ImageUtils.processImage(image, maxDimension: 256, outputDirectory: outputDirectory) let data = try Data(contentsOf: outputURL) + #expect(outputURL.deletingLastPathComponent() == outputDirectory) #expect(outputURL.pathExtension.lowercased() == "jpg") #expect(data.starts(with: Data([0xFF, 0xD8]))) #expect(data.count > 0) @@ -68,12 +75,11 @@ struct ImageUtilsTests { @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) - } + let outputDirectory = makeTemporaryDirectoryURL("image-output-\(UUID().uuidString)") + defer { try? FileManager.default.removeItem(at: outputDirectory) } + + let firstURL = try ImageUtils.processImage(image, maxDimension: 64, outputDirectory: outputDirectory) + let secondURL = try ImageUtils.processImage(image, maxDimension: 64, outputDirectory: outputDirectory) #expect(firstURL != secondURL) #expect(FileManager.default.fileExists(atPath: firstURL.path)) diff --git a/bitchatTests/Nostr/GeoRelayDirectoryTests.swift b/bitchatTests/Nostr/GeoRelayDirectoryTests.swift index b516f708..ddd3b971 100644 --- a/bitchatTests/Nostr/GeoRelayDirectoryTests.swift +++ b/bitchatTests/Nostr/GeoRelayDirectoryTests.swift @@ -10,7 +10,9 @@ final class GeoRelayDirectoryTests: XCTestCase { relay url,lat,lon wss://one.example/,10,20 https://one.example,10,20 + wss://one.example:443/,10,20 http://two.example/,11,21 + wss://two.example:443,11,21 invalid row ws://three.example,not-a-lat,22 """ diff --git a/bitchatTests/Services/BLENoisePayloadFactoryTests.swift b/bitchatTests/Services/BLENoisePayloadFactoryTests.swift new file mode 100644 index 00000000..ab743c9a --- /dev/null +++ b/bitchatTests/Services/BLENoisePayloadFactoryTests.swift @@ -0,0 +1,34 @@ +import Foundation +import Testing +@testable import bitchat + +struct BLENoisePayloadFactoryTests { + @Test + func privateMessagePayloadPrefixesTLVWithNoiseType() throws { + let payload = try #require(BLENoisePayloadFactory.privateMessage(content: "secret", messageID: "pm-1")) + + #expect(payload.first == NoisePayloadType.privateMessage.rawValue) + + let packet = try #require(PrivateMessagePacket.decode(from: Data(payload.dropFirst()))) + #expect(packet.messageID == "pm-1") + #expect(packet.content == "secret") + } + + @Test + func receiptPayloadsUseMessageIDBytes() { + let read = BLENoisePayloadFactory.readReceipt(originalMessageID: "read-id") + let delivered = BLENoisePayloadFactory.delivered(messageID: "delivered-id") + + #expect(read.first == NoisePayloadType.readReceipt.rawValue) + #expect(String(data: Data(read.dropFirst()), encoding: .utf8) == "read-id") + #expect(delivered.first == NoisePayloadType.delivered.rawValue) + #expect(String(data: Data(delivered.dropFirst()), encoding: .utf8) == "delivered-id") + } + + @Test + func typedPayloadKeepsOpaqueDataUnchanged() { + let payload = BLENoisePayloadFactory.typedPayload(.verifyChallenge, payload: Data([0xCA, 0xFE])) + + #expect(payload == Data([NoisePayloadType.verifyChallenge.rawValue, 0xCA, 0xFE])) + } +} diff --git a/bitchatTests/Services/BLENoiseSessionQueuesTests.swift b/bitchatTests/Services/BLENoiseSessionQueuesTests.swift new file mode 100644 index 00000000..33cb0fe7 --- /dev/null +++ b/bitchatTests/Services/BLENoiseSessionQueuesTests.swift @@ -0,0 +1,67 @@ +import BitFoundation +import Foundation +import Testing +@testable import bitchat + +struct BLENoiseSessionQueuesTests { + @Test + func privateMessagesDrainInPeerOrderAndClearOnlyThatPeer() { + let firstPeer = PeerID(str: "aaaaaaaaaaaaaaaa") + let secondPeer = PeerID(str: "bbbbbbbbbbbbbbbb") + var queues = BLENoiseSessionQueues() + + queues.appendPrivateMessage(content: "first", messageID: "m1", for: firstPeer) + queues.appendPrivateMessage(content: "second", messageID: "m2", for: firstPeer) + queues.appendPrivateMessage(content: "other", messageID: "m3", for: secondPeer) + + let drained = queues.takePrivateMessages(for: firstPeer) + + #expect(drained == [ + BLEPendingPrivateMessage(content: "first", messageID: "m1"), + BLEPendingPrivateMessage(content: "second", messageID: "m2") + ]) + #expect(queues.takePrivateMessages(for: firstPeer).isEmpty) + #expect(queues.takePrivateMessages(for: secondPeer) == [ + BLEPendingPrivateMessage(content: "other", messageID: "m3") + ]) + } + + @Test + func prependPrivateMessagesRestoresFailedMessagesAheadOfNewerOnes() { + let peerID = PeerID(str: "aaaaaaaaaaaaaaaa") + var queues = BLENoiseSessionQueues() + + queues.appendPrivateMessage(content: "new", messageID: "m2", for: peerID) + queues.prependPrivateMessages([ + BLEPendingPrivateMessage(content: "retry", messageID: "m1") + ], for: peerID) + + #expect(queues.takePrivateMessages(for: peerID).map(\.messageID) == ["m1", "m2"]) + } + + @Test + func typedPayloadsDrainIndependentlyFromPrivateMessages() { + let peerID = PeerID(str: "aaaaaaaaaaaaaaaa") + var queues = BLENoiseSessionQueues() + + queues.appendPrivateMessage(content: "queued", messageID: "m1", for: peerID) + queues.appendTypedPayload(Data([0x01]), for: peerID) + queues.appendTypedPayload(Data([0x02]), for: peerID) + + #expect(queues.takeTypedPayloads(for: peerID) == [Data([0x01]), Data([0x02])]) + #expect(queues.takeTypedPayloads(for: peerID).isEmpty) + #expect(queues.takePrivateMessages(for: peerID).map(\.messageID) == ["m1"]) + } + + @Test + func removeAllClearsBothQueueTypes() { + let peerID = PeerID(str: "aaaaaaaaaaaaaaaa") + var queues = BLENoiseSessionQueues() + + queues.appendPrivateMessage(content: "queued", messageID: "m1", for: peerID) + queues.appendTypedPayload(Data([0x01]), for: peerID) + queues.removeAll() + + #expect(queues.isEmpty) + } +} diff --git a/bitchatTests/Services/BLEOutboundFragmentTransferSchedulerTests.swift b/bitchatTests/Services/BLEOutboundFragmentTransferSchedulerTests.swift new file mode 100644 index 00000000..91b54399 --- /dev/null +++ b/bitchatTests/Services/BLEOutboundFragmentTransferSchedulerTests.swift @@ -0,0 +1,148 @@ +import BitFoundation +import Foundation +import Testing +@testable import bitchat + +struct BLEOutboundFragmentTransferSchedulerTests { + @Test + func submitStartsPublicMessageWithoutTransferReservation() { + var scheduler = BLEOutboundFragmentTransferScheduler() + let request = makeRequest(type: MessageType.message.rawValue, transferId: nil) + + let result = scheduler.submit(request, maxConcurrentTransfers: 1) + + if case let .start(_, reservedTransferId) = result { + #expect(reservedTransferId == nil) + #expect(scheduler.activeCount == 0) + #expect(scheduler.pendingCount == 0) + } else { + Issue.record("Expected non-file fragments to start without reserving a transfer slot") + } + } + + @Test + func submitQueuesFileTransferWhenSlotsAreFull() { + var scheduler = BLEOutboundFragmentTransferScheduler() + let first = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "first") + let second = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "second") + + guard case let .start(_, firstReservation?) = scheduler.submit(first, maxConcurrentTransfers: 1) else { + Issue.record("Expected first file transfer to reserve a slot") + return + } + #expect(firstReservation == "first") + + let result = scheduler.submit(second, maxConcurrentTransfers: 1) + + if case let .queued(_, transferId, position) = result { + #expect(transferId == "second") + #expect(position == .back) + #expect(scheduler.activeCount == 1) + #expect(scheduler.pendingCount == 1) + } else { + Issue.record("Expected second file transfer to queue while slots are full") + } + } + + @Test + func submitQueuesDuplicateActiveTransferAtFront() { + var scheduler = BLEOutboundFragmentTransferScheduler() + let request = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "same") + + _ = scheduler.submit(request, maxConcurrentTransfers: 2) + let result = scheduler.submit(request, maxConcurrentTransfers: 2) + + if case let .queued(_, transferId, position) = result { + #expect(transferId == "same") + #expect(position == .front) + #expect(scheduler.activeCount == 1) + #expect(scheduler.pendingCount == 1) + } else { + Issue.record("Expected duplicate active transfer to queue at the front") + } + } + + @Test + func cancelActiveTransferReturnsScheduledWorkItems() { + var scheduler = BLEOutboundFragmentTransferScheduler() + let request = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "active") + _ = scheduler.submit(request, maxConcurrentTransfers: 1) + let workItem = DispatchWorkItem {} + + let didActivate = scheduler.activateReservedTransfer(id: "active", totalFragments: 2, workItems: [workItem]) + #expect(didActivate) + + if case let .active(transferId, workItems) = scheduler.cancelTransfer("active") { + #expect(transferId == "active") + #expect(workItems.count == 1) + #expect(scheduler.activeCount == 0) + } else { + Issue.record("Expected active transfer cancellation to return its work items") + } + } + + @Test + func completedTransferFreesSlotForPendingTransfer() { + var scheduler = BLEOutboundFragmentTransferScheduler() + let first = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "first") + let second = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "second") + + _ = scheduler.submit(first, maxConcurrentTransfers: 1) + let didActivate = scheduler.activateReservedTransfer(id: "first", totalFragments: 2, workItems: []) + #expect(didActivate) + _ = scheduler.submit(second, maxConcurrentTransfers: 1) + + #expect(scheduler.markFragmentSent(transferId: "first") == .progress(sentFragments: 1, totalFragments: 2)) + #expect(scheduler.markFragmentSent(transferId: "first") == .complete(sentFragments: 2, totalFragments: 2)) + + let starts = scheduler.reservePendingStarts(maxConcurrentTransfers: 1) + #expect(starts.count == 1) + + if case let .start(_, reservedTransferId?) = starts.first { + #expect(reservedTransferId == "second") + #expect(scheduler.activeCount == 1) + #expect(scheduler.pendingCount == 0) + } else { + Issue.record("Expected pending transfer to reserve the freed slot") + } + } + + @Test + func removeAllReturnsActiveWorkItemsAndDropsPendingTransfers() { + var scheduler = BLEOutboundFragmentTransferScheduler() + let active = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "active") + let pending = makeRequest(type: MessageType.fileTransfer.rawValue, transferId: "pending") + let workItem = DispatchWorkItem {} + + _ = scheduler.submit(active, maxConcurrentTransfers: 1) + let didActivate = scheduler.activateReservedTransfer(id: "active", totalFragments: 1, workItems: [workItem]) + #expect(didActivate) + _ = scheduler.submit(pending, maxConcurrentTransfers: 1) + + let removed = scheduler.removeAll() + + #expect(removed.count == 1) + #expect(removed.first?.id == "active") + #expect(removed.first?.workItems.count == 1) + #expect(scheduler.activeCount == 0) + #expect(scheduler.pendingCount == 0) + } + + private func makeRequest(type: UInt8, transferId: String?) -> BLEOutboundFragmentTransferRequest { + BLEOutboundFragmentTransferRequest( + packet: BitchatPacket( + type: type, + senderID: Data([0x00, 0x11, 0x22, 0x33, 0x44, 0x55, 0x66, 0x77]), + recipientID: nil, + timestamp: 0x0102030405, + payload: Data((transferId ?? "payload").utf8), + signature: nil, + ttl: 3 + ), + pad: false, + maxChunk: nil, + directedPeer: nil, + transferId: transferId + ) + } +} diff --git a/bitchatTests/Services/NostrRelayManagerTests.swift b/bitchatTests/Services/NostrRelayManagerTests.swift index ee9fb3d2..a7447980 100644 --- a/bitchatTests/Services/NostrRelayManagerTests.swift +++ b/bitchatTests/Services/NostrRelayManagerTests.swift @@ -208,7 +208,7 @@ final class NostrRelayManagerTests: XCTestCase { let relayTwo = "wss://relay-two.example" let context = makeContext(permission: .denied) - context.manager.ensureConnections(to: [relayOne, relayOne, relayTwo]) + context.manager.ensureConnections(to: [relayOne, "wss://relay-one.example:443/", "WSS://RELAY-TWO.EXAMPLE:443"]) let connected = await waitUntil { Set(context.manager.getRelayStatuses().map(\.url)) == Set([relayOne, relayTwo]) &&