Improve mesh media throughput (#845)

* Improve mesh media throughput

* Reserve media slots atomically

* Prioritize small fragment trains

---------

Co-authored-by: jack <jackjackbits@users.noreply.github.com>
This commit is contained in:
jack
2025-10-21 12:53:54 +02:00
committed by GitHub
co-authored by jack
parent 450955525c
commit 0812cafd55
6 changed files with 227 additions and 50 deletions
+5 -5
View File
@@ -13,10 +13,10 @@ enum ImageUtilsError: Error {
}
enum ImageUtils {
private static let compressionQuality: CGFloat = 0.85
private static let targetImageBytes: Int = 60_000
private static let compressionQuality: CGFloat = 0.82
private static let targetImageBytes: Int = 45_000
static func processImage(at url: URL, maxDimension: CGFloat = 512) throws -> URL {
static func processImage(at url: URL, maxDimension: CGFloat = 448) 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 {
@@ -38,7 +38,7 @@ enum ImageUtils {
}
#if os(iOS)
static func processImage(_ image: UIImage, maxDimension: CGFloat = 512) throws -> URL {
static func processImage(_ image: UIImage, maxDimension: CGFloat = 448) throws -> URL {
return try autoreleasepool {
// Scale the image first
let scaled = scaledImage(image, maxDimension: maxDimension)
@@ -106,7 +106,7 @@ enum ImageUtils {
return data as Data
}
#else
static func processImage(_ image: NSImage, maxDimension: CGFloat = 512) throws -> URL {
static func processImage(_ image: NSImage, maxDimension: CGFloat = 448) throws -> URL {
return try autoreleasepool {
let scaled = scaledImage(image, maxDimension: maxDimension)
guard let inputCG = scaled.cgImage(forProposedRect: nil, context: nil, hints: nil) else {
+3 -2
View File
@@ -14,6 +14,7 @@ final class VoiceRecorder: NSObject, AVAudioRecorderDelegate {
private let queue = DispatchQueue(label: "com.bitchat.voice-recorder")
private let paddingInterval: TimeInterval = 0.5
private let maxRecordingDuration: TimeInterval = 120
private var recorder: AVAudioRecorder?
private var currentURL: URL?
@@ -75,14 +76,14 @@ final class VoiceRecorder: NSObject, AVAudioRecorderDelegate {
AVFormatIDKey: kAudioFormatMPEG4AAC,
AVSampleRateKey: 16_000,
AVNumberOfChannelsKey: 1,
AVEncoderBitRateKey: 20_000
AVEncoderBitRateKey: 16_000
]
let audioRecorder = try AVAudioRecorder(url: outputURL, settings: settings)
audioRecorder.delegate = self
audioRecorder.isMeteringEnabled = true
audioRecorder.prepareToRecord()
audioRecorder.record()
audioRecorder.record(forDuration: maxRecordingDuration)
recorder = audioRecorder
currentURL = outputURL
+202 -41
View File
@@ -147,7 +147,35 @@ final class BLEService: NSObject {
private var ingressByMessageID: [String: (link: LinkID, timestamp: Date)] = [:]
// Backpressure-aware write queue per peripheral
private var pendingPeripheralWrites: [String: [Data]] = [:]
private struct OutboundPriority: Comparable {
let level: Int
let suborder: Int
static let high = OutboundPriority(level: 0, suborder: 0)
static func fragment(totalFragments: Int) -> OutboundPriority {
OutboundPriority(level: 1, suborder: max(1, min(totalFragments, Int(UInt16.max))))
}
static let fileTransfer = OutboundPriority(level: 2, suborder: Int.max - 1)
static let low = OutboundPriority(level: 2, suborder: Int.max)
static func < (lhs: OutboundPriority, rhs: OutboundPriority) -> Bool {
if lhs.level != rhs.level { return lhs.level < rhs.level }
return lhs.suborder < rhs.suborder
}
}
private struct PendingWrite {
let priority: OutboundPriority
let data: Data
}
private struct PendingFragmentTransfer {
let packet: BitchatPacket
let pad: Bool
let maxChunk: Int?
let directedPeer: PeerID?
let transferId: String?
}
private var pendingPeripheralWrites: [String: [PendingWrite]] = [:]
private var pendingFragmentTransfers: [PendingFragmentTransfer] = []
// Debounce duplicate disconnect notifies
private var recentDisconnectNotifies: [PeerID: Date] = [:]
// Store-and-forward for directed messages when we have no links
@@ -303,6 +331,7 @@ final class BLEService: NSObject {
recentAnnounceBySender.removeAll()
recentAnnounceOrder.removeAll()
pendingPeripheralWrites.removeAll()
pendingFragmentTransfers.removeAll()
pendingNotifications.removeAll()
pendingDirectedRelays.removeAll()
ingressByMessageID.removeAll()
@@ -454,10 +483,11 @@ final class BLEService: NSObject {
// Send immediately to all connected peers
if let data = leavePacket.toBinaryData(padding: false) {
let leavePriority = priority(for: leavePacket, data: data)
// Send to peripherals we're connected to as central
for state in peripherals.values where state.isConnected {
if let characteristic = state.characteristic {
writeOrEnqueue(data, to: state.peripheral, characteristic: characteristic)
writeOrEnqueue(data, to: state.peripheral, characteristic: characteristic, priority: leavePriority)
}
}
@@ -591,10 +621,19 @@ final class BLEService: NSObject {
func cancelTransfer(_ transferId: String) {
collectionsQueue.async(flags: .barrier) { [weak self] in
guard let self = self, let state = self.activeTransfers.removeValue(forKey: transferId) else { return }
state.workItems.forEach { $0.cancel() }
TransferProgressManager.shared.cancel(id: transferId)
SecureLogger.debug("🛑 Cancelled transfer \(transferId.prefix(8))", category: .session)
guard let self = self else { return }
if let state = self.activeTransfers.removeValue(forKey: transferId) {
state.workItems.forEach { $0.cancel() }
TransferProgressManager.shared.cancel(id: transferId)
SecureLogger.debug("🛑 Cancelled transfer \(transferId.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)
TransferProgressManager.shared.cancel(id: transferId)
SecureLogger.debug("🛑 Removed pending transfer \(transferId.prefix(8))… before start", category: .session)
}
}
}
@@ -737,6 +776,8 @@ final class BLEService: NSObject {
guard let recipientPeerID = PeerID(hexData: packet.recipientID) else { return }
var sentEncrypted = false
let outboundPriority = 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] }) {
@@ -769,7 +810,7 @@ final class BLEService: NSObject {
let state = (DispatchQueue.getSpecific(key: bleQueueKey) != nil) ? peripherals[peripheralUUID] : bleQueue.sync(execute: { peripherals[peripheralUUID] }),
state.isConnected,
let characteristic = state.characteristic {
writeOrEnqueue(data, to: state.peripheral, characteristic: characteristic)
writeOrEnqueue(data, to: state.peripheral, characteristic: characteristic, priority: outboundPriority)
sentEncrypted = true
}
@@ -826,6 +867,7 @@ final class BLEService: NSObject {
}
return nil
}()
let outboundPriority = priority(for: packet, data: data)
let states = snapshotPeripheralStates()
var minCentralWriteLen: Int?
@@ -901,7 +943,7 @@ final class BLEService: NSObject {
let pid = s.peripheral.identifier.uuidString
guard selectedPeripheralIDs.contains(pid) else { continue }
if let ch = s.characteristic {
writeOrEnqueue(data, to: s.peripheral, characteristic: ch)
writeOrEnqueue(data, to: s.peripheral, characteristic: ch, priority: outboundPriority)
}
}
// Notify selected subscribed centrals
@@ -980,7 +1022,7 @@ final class BLEService: NSObject {
// Fire-and-forget principle: always use .withoutResponse for speed
// CoreBluetooth will handle fragmentation at L2CAP layer
writeOrEnqueue(data, to: peripheral, characteristic: characteristic)
writeOrEnqueue(data, to: peripheral, characteristic: characteristic, priority: .high)
}
private func handleFileTransfer(_ packet: BitchatPacket, from peerID: PeerID) {
@@ -2347,7 +2389,28 @@ extension BLEService {
return Set(scored.prefix(k).map { $0.id })
}
private func writeOrEnqueue(_ data: Data, to peripheral: CBPeripheral, characteristic: CBCharacteristic) {
private func priority(for packet: BitchatPacket, data: Data) -> OutboundPriority {
guard let messageType = MessageType(rawValue: packet.type) else { return .low }
switch messageType {
case .fragment:
let total = fragmentTotalCount(from: packet.payload)
return OutboundPriority.fragment(totalFragments: total)
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: OutboundPriority) {
// BLE operations run on bleQueue; keep queue affinity
bleQueue.async { [weak self] in
guard let self = self else { return }
@@ -2363,18 +2426,20 @@ extension BLEService {
if newSize > capBytes {
SecureLogger.warning("⚠️ Dropping oversized write chunk (\(newSize)B) for peripheral \(uuid)", category: .session)
} else {
// Append and trim from the front to respect cap
var total = queue.reduce(0) { $0 + $1.count }
queue.append(data)
total += newSize
let item = PendingWrite(priority: priority, data: data)
var total = queue.reduce(0) { $0 + $1.data.count } + newSize
let insertIndex = queue.firstIndex { item.priority < $0.priority } ?? queue.count
queue.insert(item, at: insertIndex)
if total > capBytes {
var removedBytes = 0
while total > capBytes && !queue.isEmpty {
let removed = queue.removeFirst()
removedBytes += removed.count
total -= removed.count
let removed = queue.removeLast()
removedBytes += removed.data.count
total -= removed.data.count
}
if removedBytes > 0 {
SecureLogger.warning("📉 Trimmed pending write buffer for \(uuid) by \(removedBytes)B to \(total)B", category: .session)
}
SecureLogger.warning("📉 Trimmed pending write buffer for \(uuid) by \(removedBytes)B to \(total)B", category: .session)
}
self.pendingPeripheralWrites[uuid] = queue.isEmpty ? nil : queue
}
@@ -2388,7 +2453,7 @@ extension BLEService {
bleQueue.async { [weak self] in
guard let self = self else { return }
guard let state = self.peripherals[uuid], let ch = state.characteristic else { return }
var queueCopy: [Data] = []
var queueCopy: [PendingWrite] = []
self.collectionsQueue.sync {
queueCopy = self.pendingPeripheralWrites[uuid] ?? []
}
@@ -2396,7 +2461,7 @@ extension BLEService {
var sent = 0
for item in queueCopy {
if peripheral.canSendWriteWithoutResponse {
peripheral.writeValue(item, for: ch, type: .withoutResponse)
peripheral.writeValue(item.data, for: ch, type: .withoutResponse)
sent += 1
} else {
break
@@ -2405,12 +2470,13 @@ extension BLEService {
if sent > 0 {
self.collectionsQueue.async(flags: .barrier) {
var q = self.pendingPeripheralWrites[uuid] ?? []
if sent <= q.count {
q.removeFirst(sent)
} else {
q.removeAll()
if sent > 0 {
let toRemove = min(sent, q.count)
if toRemove > 0 {
q.removeFirst(toRemove)
}
self.pendingPeripheralWrites[uuid] = q.isEmpty ? nil : q
}
self.pendingPeripheralWrites[uuid] = q.isEmpty ? nil : q
}
}
}
@@ -2601,16 +2667,86 @@ 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) {
guard let fullData = packet.toBinaryData(padding: pad) else { return }
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
}
}
startFragmentedPacket(context)
}
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)
} else {
self.pendingFragmentTransfers.append(context)
}
}
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?
let releaseReservedSlot: (String) -> Void = { id in
TransferProgressManager.shared.cancel(id: id)
self.collectionsQueue.async(flags: .barrier) { [weak self] in
self?.activeTransfers.removeValue(forKey: 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 {
if let id = reservedTransferId {
releaseReservedSlot(id)
}
return
}
// Fragment the unpadded frame; each fragment will be encoded independently
let fragmentID = Data((0..<8).map { _ in UInt8.random(in: 0...255) })
let chunk = maxChunk ?? defaultFragmentSize
let chunk = context.maxChunk ?? defaultFragmentSize
let safeChunk = max(64, chunk)
let fragments = stride(from: 0, to: fullData.count, by: safeChunk).map { offset in
Data(fullData[offset..<min(offset + safeChunk, fullData.count)])
}
guard !fragments.isEmpty else { return }
guard !fragments.isEmpty else {
if let id = reservedTransferId {
releaseReservedSlot(id)
}
return
}
// Lightweight pacing to reduce floods and allow BLE buffers to drain
// Also briefly pause scanning during long fragment trains to save battery
@@ -2619,18 +2755,16 @@ extension BLEService {
bleQueue.async { [weak self] in
guard let self = self, let c = self.centralManager, c.state == .poweredOn else { return }
if c.isScanning { c.stopScan() }
// Resume scanning after we expect last fragment to be sent
let expectedMs = min(TransportConfig.bleExpectedWriteMaxMs, totalFragments * TransportConfig.bleExpectedWritePerFragmentMs) // ~8ms per fragment
let expectedMs = min(TransportConfig.bleExpectedWriteMaxMs, totalFragments * TransportConfig.bleExpectedWritePerFragmentMs)
self.bleQueue.asyncAfter(deadline: .now() + .milliseconds(expectedMs)) { [weak self] in
self?.startScanning()
}
}
}
let perFragMs = (directedOnlyPeer != nil || packet.recipientID != nil) ? TransportConfig.bleFragmentSpacingDirectedMs : TransportConfig.bleFragmentSpacingMs
let perFragMs = (context.directedPeer != nil || packet.recipientID != nil) ? TransportConfig.bleFragmentSpacingDirectedMs : TransportConfig.bleFragmentSpacingMs
let transferIdentifier: String? = {
guard packet.type == MessageType.fileTransfer.rawValue else { return nil }
let id = transferId ?? packet.payload.sha256Hex()
guard let id = reservedTransferId else { return nil }
collectionsQueue.sync(flags: .barrier) {
self.activeTransfers[id] = ActiveTransferState(totalFragments: totalFragments, sentFragments: 0, workItems: [])
}
@@ -2647,10 +2781,9 @@ extension BLEService {
payload.append(contentsOf: withUnsafeBytes(of: UInt16(fragments.count).bigEndian) { Data($0) })
payload.append(packet.type)
payload.append(fragment)
// Choose recipient for the fragment: directed override if provided
let fragmentRecipient: Data? = {
if let only = directedOnlyPeer { return Data(hexString: only.id) }
if let only = context.directedPeer { return Data(hexString: only.id) }
return packet.recipientID
}()
@@ -2700,10 +2833,36 @@ extension BLEService {
collectionsQueue.async(flags: .barrier) { [weak self] in
guard let self = self, var state = self.activeTransfers[transferId] else { return }
state.sentFragments = min(state.sentFragments + 1, state.totalFragments)
self.activeTransfers[transferId] = state
TransferProgressManager.shared.recordFragmentSent(id: transferId)
if state.sentFragments >= state.totalFragments {
let isComplete = state.sentFragments >= state.totalFragments
if isComplete {
self.activeTransfers.removeValue(forKey: transferId)
} else {
self.activeTransfers[transferId] = state
}
TransferProgressManager.shared.recordFragmentSent(id: transferId)
if isComplete {
self.messageQueue.async { [weak self] in
self?.startNextPendingTransferIfNeeded()
}
}
}
}
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)
}
}
}
}
@@ -2814,8 +2973,9 @@ extension BLEService {
}
// Decode the original packet bytes we reassembled, so flags/compression are preserved
if let originalPacket = BinaryProtocol.decode(reassembled) {
if var originalPacket = BinaryProtocol.decode(reassembled) {
SecureLogger.debug("✅ Reassembled packet id=\(String(format: "%016llx", fragU64)) type=\(originalPacket.type) bytes=\(reassembled.count)", category: .session)
originalPacket.ttl = 0
handleReceivedPacket(originalPacket, from: peerID)
} else {
SecureLogger.error("❌ Failed to decode reassembled packet (type=\(originalType), total=\(total))", category: .session)
@@ -2931,6 +3091,7 @@ extension BLEService {
senderIsSelf: senderID == myPeerID,
isEncrypted: packet.type == MessageType.noiseEncrypted.rawValue,
isDirectedEncrypted: (packet.type == MessageType.noiseEncrypted.rawValue) && (packet.recipientID != nil),
isFragment: packet.type == MessageType.fragment.rawValue,
isDirectedFragment: packet.type == MessageType.fragment.rawValue && packet.recipientID != nil,
isHandshake: packet.type == MessageType.noiseHandshake.rawValue,
isAnnounce: packet.type == MessageType.announce.rawValue,
+11
View File
@@ -13,6 +13,7 @@ struct RelayController {
senderIsSelf: Bool,
isEncrypted: Bool,
isDirectedEncrypted: Bool,
isFragment: Bool,
isDirectedFragment: Bool,
isHandshake: Bool,
isAnnounce: Bool,
@@ -36,6 +37,16 @@ struct RelayController {
return RelayDecision(shouldRelay: true, newTTL: newTTL, delayMs: delayMs)
}
if isFragment {
let ttlLimit = min(ttlCap, TransportConfig.bleFragmentRelayTtlCap)
guard ttlLimit > 1 else {
return RelayDecision(shouldRelay: false, newTTL: ttlLimit, delayMs: 0)
}
let newTTL = ttlLimit &- 1
let delayMs = Int.random(in: TransportConfig.bleFragmentRelayMinDelayMs...TransportConfig.bleFragmentRelayMaxDelayMs)
return RelayDecision(shouldRelay: true, newTTL: newTTL, delayMs: delayMs)
}
// TTL clamping for broadcast
// - Dense graphs: keep lower but still allow multi-hop bridging
// - Announces get a bit more headroom
+4
View File
@@ -8,6 +8,10 @@ enum TransportConfig {
static let messageTTLDefault: UInt8 = 7 // Default TTL for mesh flooding
static let bleMaxInFlightAssemblies: Int = 128 // Cap concurrent fragment assemblies
static let bleHighDegreeThreshold: Int = 6 // For adaptive TTL/probabilistic relays
static let bleMaxConcurrentTransfers: Int = 2 // Limit simultaneous large media sends
static let bleFragmentRelayMinDelayMs: Int = 8 // Faster forwarding for media fragments
static let bleFragmentRelayMaxDelayMs: Int = 25 // Upper jitter bound for fragment relays
static let bleFragmentRelayTtlCap: UInt8 = 5 // Clamp fragment TTL to contain floods
// UI / Storage Caps
static let privateChatCap: Int = 1337
+2 -2
View File
@@ -5,9 +5,9 @@ enum FileTransferLimits {
/// Absolute ceiling enforced for any file payload (voice, image, other).
static let maxPayloadBytes: Int = 1 * 1024 * 1024 // 1 MiB
/// Voice notes stay small for low-latency relays.
static let maxVoiceNoteBytes: Int = 1 * 1024 * 1024 // 1 MiB
static let maxVoiceNoteBytes: Int = 512 * 1024 // 512 KiB
/// Compressed images after downscaling should comfortably fit under this budget.
static let maxImageBytes: Int = 1 * 1024 * 1024 // 1 MiB
static let maxImageBytes: Int = 512 * 1024 // 512 KiB
/// Worst-case size once TLV metadata and binary packet framing are included for the largest payloads.
static let maxFramedFileBytes: Int = {
let maxMetadataBytes = Int(UInt16.max) * 2 // fileName + mimeType TLVs