mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-25 08:45:19 +00:00
* PeerID 28/n: `ChatViewModel.getShortIDForNoiseKey` * PeerID 29/n: `BLEService` + remove dupe funcs from #823 * `handleFileTransfer` to use PeerID * `sendMessage` and `sendPrivateMessage` --------- Co-authored-by: jack <212554440+jackjackbits@users.noreply.github.com>
This commit is contained in:
@@ -334,7 +334,7 @@ final class BLEService: NSObject {
|
||||
}
|
||||
|
||||
// Ensure this runs on message queue to avoid main thread blocking
|
||||
func sendMessage(_ content: String, mentions: [String] = [], to recipientID: String? = nil, messageID: String? = nil, timestamp: Date? = nil) {
|
||||
func sendMessage(_ content: String, mentions: [String] = [], to recipientID: PeerID? = nil, messageID: String? = nil, timestamp: Date? = nil) {
|
||||
// Call directly if already on messageQueue, otherwise dispatch
|
||||
if DispatchQueue.getSpecific(key: messageQueueKey) == nil {
|
||||
messageQueue.async { [weak self] in
|
||||
@@ -605,7 +605,7 @@ final class BLEService: NSObject {
|
||||
}
|
||||
|
||||
func sendPrivateMessage(_ content: String, to peerID: PeerID, recipientNickname: String, messageID: String) {
|
||||
sendPrivateMessage(content, to: peerID.id, messageID: messageID)
|
||||
sendPrivateMessage(content, to: peerID, messageID: messageID)
|
||||
}
|
||||
|
||||
func sendFileBroadcast(_ filePacket: BitchatFilePacket, transferId: String) {
|
||||
@@ -735,12 +735,12 @@ final class BLEService: NSObject {
|
||||
|
||||
private func sendEncrypted(_ packet: BitchatPacket, data: Data, pad: Bool) {
|
||||
guard let recipientID = packet.recipientID else { return }
|
||||
let recipientPeerID = recipientID.hexEncodedString()
|
||||
let recipientPeerID = PeerID(hexData: recipientID)
|
||||
var sentEncrypted = false
|
||||
|
||||
// Per-link limits for the specific peer
|
||||
var peripheralMaxLen: Int?
|
||||
if let perUUID = (DispatchQueue.getSpecific(key: bleQueueKey) != nil) ? peerToPeripheralUUID[PeerID(str: recipientPeerID)] : bleQueue.sync(execute: { peerToPeripheralUUID[PeerID(str: recipientPeerID)] }) {
|
||||
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)
|
||||
}
|
||||
@@ -748,7 +748,7 @@ final class BLEService: NSObject {
|
||||
var centralMaxLen: Int?
|
||||
do {
|
||||
let (centrals, mapping) = snapshotSubscribedCentrals()
|
||||
if let central = centrals.first(where: { mapping[$0.identifier.uuidString] == PeerID(str: recipientPeerID) }) {
|
||||
if let central = centrals.first(where: { mapping[$0.identifier.uuidString] == recipientPeerID }) {
|
||||
centralMaxLen = central.maximumUpdateValueLength
|
||||
}
|
||||
}
|
||||
@@ -766,7 +766,7 @@ final class BLEService: NSObject {
|
||||
}
|
||||
|
||||
// Direct write via peripheral link
|
||||
if let peripheralUUID = (DispatchQueue.getSpecific(key: bleQueueKey) != nil) ? peerToPeripheralUUID[PeerID(str: recipientPeerID)] : bleQueue.sync(execute: { peerToPeripheralUUID[PeerID(str: recipientPeerID)] }),
|
||||
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] }),
|
||||
state.isConnected,
|
||||
let characteristic = state.characteristic {
|
||||
@@ -777,7 +777,7 @@ 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] == PeerID(str: recipientPeerID) {
|
||||
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")
|
||||
@@ -816,13 +816,13 @@ final class BLEService: NSObject {
|
||||
}
|
||||
}
|
||||
|
||||
private func sendOnAllLinks(packet: BitchatPacket, data: Data, pad: Bool, directedOnlyPeer: String?) {
|
||||
private func sendOnAllLinks(packet: BitchatPacket, data: Data, pad: Bool, directedOnlyPeer: PeerID?) {
|
||||
// Determine last-hop link for this message to avoid echoing back
|
||||
let messageID = makeMessageID(for: packet)
|
||||
let ingressLink: LinkID? = collectionsQueue.sync { ingressByMessageID[messageID]?.link }
|
||||
let directedPeerHint: String? = {
|
||||
let directedPeerHint: PeerID? = {
|
||||
if let explicit = directedOnlyPeer { return explicit }
|
||||
if let recipient = packet.recipientID?.hexEncodedString(), !recipient.isEmpty {
|
||||
if let recipient = PeerID(str: packet.recipientID?.hexEncodedString()), !recipient.isEmpty {
|
||||
return recipient
|
||||
}
|
||||
return nil
|
||||
@@ -915,20 +915,20 @@ final class BLEService: NSObject {
|
||||
}
|
||||
|
||||
// Directed send helper (unicast to a specific peerID) without altering packet contents
|
||||
private func sendPacketDirected(_ packet: BitchatPacket, to peerID: String) {
|
||||
private func sendPacketDirected(_ packet: BitchatPacket, to peerID: PeerID) {
|
||||
guard let data = packet.toBinaryData(padding: false) else { return }
|
||||
sendOnAllLinks(packet: packet, data: data, pad: false, directedOnlyPeer: peerID)
|
||||
}
|
||||
|
||||
// MARK: - Directed store-and-forward
|
||||
private func spoolDirectedPacket(_ packet: BitchatPacket, recipientPeerID: String) {
|
||||
private func spoolDirectedPacket(_ packet: BitchatPacket, recipientPeerID: PeerID) {
|
||||
let msgID = makeMessageID(for: packet)
|
||||
collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
guard let self = self else { return }
|
||||
var byMsg = self.pendingDirectedRelays[PeerID(str: recipientPeerID)] ?? [:]
|
||||
var byMsg = self.pendingDirectedRelays[recipientPeerID] ?? [:]
|
||||
if byMsg[msgID] == nil {
|
||||
byMsg[msgID] = (packet: packet, enqueuedAt: Date())
|
||||
self.pendingDirectedRelays[PeerID(str: recipientPeerID)] = byMsg
|
||||
self.pendingDirectedRelays[recipientPeerID] = byMsg
|
||||
SecureLogger.debug("🧳 Spooling directed packet for \(recipientPeerID) mid=\(msgID.prefix(8))…", category: .session)
|
||||
}
|
||||
}
|
||||
@@ -983,615 +983,8 @@ final class BLEService: NSObject {
|
||||
// CoreBluetooth will handle fragmentation at L2CAP layer
|
||||
writeOrEnqueue(data, to: peripheral, characteristic: characteristic)
|
||||
}
|
||||
|
||||
// MARK: - Fragmentation (Required for messages > BLE MTU)
|
||||
|
||||
private func sendFragmentedPacket(_ packet: BitchatPacket, pad: Bool, maxChunk: Int? = nil, directedOnlyPeer: String? = nil, transferId: String? = nil) {
|
||||
guard let fullData = packet.toBinaryData(padding: pad) else { 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 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 }
|
||||
|
||||
// Lightweight pacing to reduce floods and allow BLE buffers to drain
|
||||
// Also briefly pause scanning during long fragment trains to save battery
|
||||
let totalFragments = fragments.count
|
||||
if totalFragments > 4 {
|
||||
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
|
||||
self.bleQueue.asyncAfter(deadline: .now() + .milliseconds(expectedMs)) { [weak self] in
|
||||
self?.startScanning()
|
||||
}
|
||||
}
|
||||
}
|
||||
let perFragMs = (directedOnlyPeer != 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()
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
self.activeTransfers[id] = ActiveTransferState(totalFragments: totalFragments, sentFragments: 0, workItems: [])
|
||||
}
|
||||
TransferProgressManager.shared.start(id: id, totalFragments: totalFragments)
|
||||
return id
|
||||
}()
|
||||
|
||||
var scheduledItems: [(item: DispatchWorkItem, index: Int)] = []
|
||||
|
||||
for (index, fragment) in fragments.enumerated() {
|
||||
var payload = Data()
|
||||
payload.append(fragmentID)
|
||||
payload.append(contentsOf: withUnsafeBytes(of: UInt16(index).bigEndian) { Data($0) })
|
||||
payload.append(contentsOf: withUnsafeBytes(of: UInt16(fragments.count).bigEndian) { Data($0) })
|
||||
payload.append(packet.type)
|
||||
payload.append(fragment)
|
||||
|
||||
let fragmentRecipient: Data? = {
|
||||
if let only = directedOnlyPeer { return Data(hexString: only) }
|
||||
return packet.recipientID
|
||||
}()
|
||||
|
||||
let fragmentPacket = BitchatPacket(
|
||||
type: MessageType.fragment.rawValue,
|
||||
senderID: packet.senderID,
|
||||
recipientID: fragmentRecipient,
|
||||
timestamp: packet.timestamp,
|
||||
payload: payload,
|
||||
signature: nil,
|
||||
ttl: packet.ttl
|
||||
)
|
||||
|
||||
let workItem = DispatchWorkItem { [weak self] in
|
||||
guard let self = self else { return }
|
||||
if let transferId = transferIdentifier {
|
||||
let isActive = self.collectionsQueue.sync { self.activeTransfers[transferId] != nil }
|
||||
guard isActive else { return }
|
||||
}
|
||||
self.broadcastPacket(fragmentPacket)
|
||||
if let transferId = transferIdentifier {
|
||||
self.markFragmentSent(transferId: transferId)
|
||||
}
|
||||
}
|
||||
|
||||
scheduledItems.append((item: workItem, index: index))
|
||||
}
|
||||
|
||||
if let transferId = transferIdentifier {
|
||||
let workItems = scheduledItems.map { $0.item }
|
||||
collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
guard let self = self, var state = self.activeTransfers[transferId] else { return }
|
||||
state.workItems = workItems
|
||||
self.activeTransfers[transferId] = state
|
||||
}
|
||||
}
|
||||
|
||||
for (workItem, index) in scheduledItems {
|
||||
let delayMs = index * perFragMs
|
||||
messageQueue.asyncAfter(deadline: .now() + .milliseconds(delayMs), execute: workItem)
|
||||
}
|
||||
}
|
||||
|
||||
private func markFragmentSent(transferId: String) {
|
||||
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 {
|
||||
self.activeTransfers.removeValue(forKey: transferId)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private func handleFragment(_ packet: BitchatPacket, from peerID: String) {
|
||||
// Don't process our own fragments
|
||||
if peerID == myPeerID {
|
||||
return
|
||||
}
|
||||
|
||||
// Minimum header: 8 bytes ID + 2 index + 2 total + 1 type
|
||||
guard packet.payload.count >= 13 else { return }
|
||||
|
||||
// Compute compact fragment key (sender: 8 bytes, id: 8 bytes), big-endian
|
||||
var senderU64: UInt64 = 0
|
||||
for b in packet.senderID.prefix(8) { senderU64 = (senderU64 << 8) | UInt64(b) }
|
||||
var fragU64: UInt64 = 0
|
||||
for b in packet.payload.prefix(8) { fragU64 = (fragU64 << 8) | UInt64(b) }
|
||||
// Parse big-endian UInt16 safely without alignment assumptions
|
||||
let idxHi = UInt16(packet.payload[8])
|
||||
let idxLo = UInt16(packet.payload[9])
|
||||
let index = Int((idxHi << 8) | idxLo)
|
||||
let totHi = UInt16(packet.payload[10])
|
||||
let totLo = UInt16(packet.payload[11])
|
||||
let total = Int((totHi << 8) | totLo)
|
||||
let originalType = packet.payload[12]
|
||||
let fragmentData = packet.payload.suffix(from: 13)
|
||||
|
||||
// Sanity checks - add reasonable upper bound on total to prevent DoS
|
||||
guard total > 0 && total <= 10000 && index >= 0 && index < total else { return }
|
||||
|
||||
// Compute fragment key for this assembly
|
||||
let key = FragmentKey(sender: senderU64, id: fragU64)
|
||||
|
||||
// Critical section: Store fragment and check completion status
|
||||
var shouldReassemble: Bool = false
|
||||
var fragmentsToReassemble: [Int: Data]? = nil
|
||||
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
if incomingFragments[key] == nil {
|
||||
// Cap in-flight assemblies to prevent memory/battery blowups
|
||||
if incomingFragments.count >= maxInFlightAssemblies {
|
||||
// Evict the oldest assembly by timestamp
|
||||
if let oldest = fragmentMetadata.min(by: { $0.value.timestamp < $1.value.timestamp })?.key {
|
||||
incomingFragments.removeValue(forKey: oldest)
|
||||
fragmentMetadata.removeValue(forKey: oldest)
|
||||
}
|
||||
}
|
||||
incomingFragments[key] = [:]
|
||||
fragmentMetadata[key] = (originalType, total, Date())
|
||||
SecureLogger.debug("📦 Started fragment assembly id=\(String(format: "%016llx", fragU64)) total=\(total)", category: .session)
|
||||
}
|
||||
|
||||
// Check cumulative size before storing this fragment
|
||||
let currentSize = incomingFragments[key]?.values.reduce(0) { $0 + $1.count } ?? 0
|
||||
let assemblyLimit: Int = {
|
||||
if originalType == MessageType.fileTransfer.rawValue {
|
||||
// Allow headroom for TLV metadata and binary framing overhead.
|
||||
return FileTransferLimits.maxFramedFileBytes
|
||||
}
|
||||
return FileTransferLimits.maxPayloadBytes
|
||||
}()
|
||||
guard currentSize + fragmentData.count <= assemblyLimit else {
|
||||
// Exceeds size limit - evict this assembly
|
||||
SecureLogger.warning(
|
||||
"🚫 Fragment assembly exceeds size limit (\(currentSize + fragmentData.count) bytes > \(assemblyLimit)), evicting",
|
||||
category: .security
|
||||
)
|
||||
incomingFragments.removeValue(forKey: key)
|
||||
fragmentMetadata.removeValue(forKey: key)
|
||||
shouldReassemble = false
|
||||
fragmentsToReassemble = nil
|
||||
return
|
||||
}
|
||||
|
||||
incomingFragments[key]?[index] = Data(fragmentData)
|
||||
SecureLogger.debug("📦 Fragment \(index + 1)/\(total) (len=\(fragmentData.count)) for id=\(String(format: "%016llx", fragU64))", category: .session)
|
||||
|
||||
// Check if complete
|
||||
if let fragments = incomingFragments[key], fragments.count == total {
|
||||
shouldReassemble = true
|
||||
fragmentsToReassemble = fragments
|
||||
} else {
|
||||
shouldReassemble = false
|
||||
fragmentsToReassemble = nil
|
||||
}
|
||||
}
|
||||
|
||||
// Heavy work outside lock: reassemble and decode
|
||||
guard shouldReassemble, let fragments = fragmentsToReassemble else { return }
|
||||
|
||||
var reassembled = Data()
|
||||
for i in 0..<total {
|
||||
if let fragment = fragments[i] {
|
||||
reassembled.append(fragment)
|
||||
}
|
||||
}
|
||||
|
||||
// Decode the original packet bytes we reassembled, so flags/compression are preserved
|
||||
if let originalPacket = BinaryProtocol.decode(reassembled) {
|
||||
SecureLogger.debug("✅ Reassembled packet id=\(String(format: "%016llx", fragU64)) type=\(originalPacket.type) bytes=\(reassembled.count)", category: .session)
|
||||
handleReceivedPacket(originalPacket, from: peerID)
|
||||
} else {
|
||||
SecureLogger.error("❌ Failed to decode reassembled packet (type=\(originalType), total=\(total))", category: .session)
|
||||
}
|
||||
|
||||
// Critical section: Cleanup completed assembly
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
incomingFragments.removeValue(forKey: key)
|
||||
fragmentMetadata.removeValue(forKey: key)
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Packet Reception
|
||||
|
||||
private func handleReceivedPacket(_ packet: BitchatPacket, from peerID: String) {
|
||||
// Deduplication (thread-safe)
|
||||
let senderID = packet.senderID.hexEncodedString()
|
||||
// Include packet type in message ID to prevent collisions between different packet types
|
||||
let messageID = "\(senderID)-\(packet.timestamp)-\(packet.type)"
|
||||
|
||||
// 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)
|
||||
}
|
||||
|
||||
// Efficient deduplication
|
||||
// Important: do not dedup fragment packets globally (each piece must pass)
|
||||
// Special case: allow our own packets recovered via sync (TTL==0) to pass
|
||||
// through even if we've marked them as seen at send time.
|
||||
let allowSelfSyncReplay = (packet.ttl == 0) && (senderID == myPeerID)
|
||||
if packet.type != MessageType.fragment.rawValue && !allowSelfSyncReplay && messageDeduplicator.isDuplicate(messageID) {
|
||||
// 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)
|
||||
}
|
||||
// In sparse graphs (<=2 neighbors), keep the pending relay to ensure bridging.
|
||||
// In denser graphs, cancel the pending relay to reduce redundant floods.
|
||||
let connectedCount = collectionsQueue.sync { peers.values.filter { $0.isConnected }.count }
|
||||
if connectedCount > 2 {
|
||||
collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
if let task = self?.scheduledRelays.removeValue(forKey: messageID) {
|
||||
task.cancel()
|
||||
}
|
||||
}
|
||||
}
|
||||
return // Duplicate ignored
|
||||
}
|
||||
|
||||
// Update peer info without verbose logging - update the peer we received from, not the original sender
|
||||
updatePeerLastSeen(PeerID(str: peerID))
|
||||
|
||||
// Track recent traffic timestamps for adaptive behavior
|
||||
collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
guard let self = self else { return }
|
||||
let now = Date()
|
||||
self.recentPacketTimestamps.append(now)
|
||||
// keep last N timestamps within window
|
||||
let cutoff = now.addingTimeInterval(-TransportConfig.bleRecentPacketWindowSeconds)
|
||||
if self.recentPacketTimestamps.count > TransportConfig.bleRecentPacketWindowMaxCount {
|
||||
self.recentPacketTimestamps.removeFirst(self.recentPacketTimestamps.count - TransportConfig.bleRecentPacketWindowMaxCount)
|
||||
}
|
||||
self.recentPacketTimestamps.removeAll { $0 < cutoff }
|
||||
}
|
||||
|
||||
|
||||
// Process by type
|
||||
switch MessageType(rawValue: packet.type) {
|
||||
case .announce:
|
||||
handleAnnounce(packet, from: senderID)
|
||||
|
||||
case .message:
|
||||
handleMessage(packet, from: senderID)
|
||||
|
||||
case .requestSync:
|
||||
handleRequestSync(packet, from: senderID)
|
||||
|
||||
case .noiseHandshake:
|
||||
handleNoiseHandshake(packet, from: PeerID(str: senderID))
|
||||
|
||||
case .noiseEncrypted:
|
||||
handleNoiseEncrypted(packet, from: PeerID(str: senderID))
|
||||
|
||||
case .fragment:
|
||||
handleFragment(packet, from: senderID)
|
||||
|
||||
case .fileTransfer:
|
||||
handleFileTransfer(packet, from: senderID)
|
||||
|
||||
case .leave:
|
||||
handleLeave(packet, from: PeerID(str: senderID))
|
||||
|
||||
case .none:
|
||||
SecureLogger.warning("⚠️ Unknown message type: \(packet.type)", category: .session)
|
||||
break
|
||||
}
|
||||
|
||||
// Relay if TTL > 1 and we're not the original sender
|
||||
// Relay decision and scheduling (extracted via RelayController)
|
||||
do {
|
||||
let degree = collectionsQueue.sync { peers.values.filter { $0.isConnected }.count }
|
||||
let decision = RelayController.decide(
|
||||
ttl: packet.ttl,
|
||||
senderIsSelf: senderID == myPeerID,
|
||||
isEncrypted: packet.type == MessageType.noiseEncrypted.rawValue,
|
||||
isDirectedEncrypted: (packet.type == MessageType.noiseEncrypted.rawValue) && (packet.recipientID != nil),
|
||||
isDirectedFragment: packet.type == MessageType.fragment.rawValue && packet.recipientID != nil,
|
||||
isHandshake: packet.type == MessageType.noiseHandshake.rawValue,
|
||||
isAnnounce: packet.type == MessageType.announce.rawValue,
|
||||
degree: degree,
|
||||
highDegreeThreshold: highDegreeThreshold
|
||||
)
|
||||
guard decision.shouldRelay else { return }
|
||||
let work = DispatchWorkItem { [weak self] in
|
||||
guard let self = self else { return }
|
||||
// Remove scheduled task before executing
|
||||
self.collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
_ = self?.scheduledRelays.removeValue(forKey: messageID)
|
||||
}
|
||||
var relayPacket = packet
|
||||
relayPacket.ttl = decision.newTTL
|
||||
self.broadcastPacket(relayPacket)
|
||||
}
|
||||
// Track the scheduled relay so duplicates can cancel it
|
||||
collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
self?.scheduledRelays[messageID] = work
|
||||
}
|
||||
messageQueue.asyncAfter(deadline: .now() + .milliseconds(decision.delayMs), execute: work)
|
||||
}
|
||||
}
|
||||
|
||||
private func handleAnnounce(_ packet: BitchatPacket, from peerID: String) {
|
||||
guard let announcement = AnnouncementPacket.decode(from: packet.payload) else {
|
||||
SecureLogger.error("❌ Failed to decode announce packet from \(peerID)", category: .session)
|
||||
return
|
||||
}
|
||||
|
||||
// Verify that the sender's derived ID from the announced noise public key matches the packet senderID
|
||||
// This helps detect relayed or spoofed announces. Only warn in release; assert in debug.
|
||||
let derivedFromKey = PeerID(publicKey: announcement.noisePublicKey).id
|
||||
if derivedFromKey != peerID {
|
||||
SecureLogger.warning("⚠️ Announce sender mismatch: derived \(derivedFromKey.prefix(8))… vs packet \(peerID.prefix(8))…", category: .security)
|
||||
|
||||
}
|
||||
|
||||
// Don't add ourselves as a peer
|
||||
if peerID == myPeerID {
|
||||
return
|
||||
}
|
||||
|
||||
// Suppress announce logs to reduce noise
|
||||
|
||||
// Precompute signature verification outside barrier to reduce contention
|
||||
let existingPeerForVerify = collectionsQueue.sync { peers[PeerID(str: peerID)] }
|
||||
var verifiedAnnounce = false
|
||||
if packet.signature != nil {
|
||||
verifiedAnnounce = noiseService.verifyPacketSignature(packet, publicKey: announcement.signingPublicKey)
|
||||
if !verifiedAnnounce {
|
||||
SecureLogger.warning("⚠️ Signature verification for announce failed \(peerID.prefix(8))", category: .security)
|
||||
}
|
||||
}
|
||||
if let existingKey = existingPeerForVerify?.noisePublicKey, existingKey != announcement.noisePublicKey {
|
||||
SecureLogger.warning("⚠️ Announce key mismatch for \(peerID.prefix(8))… — keeping unverified", category: .security)
|
||||
verifiedAnnounce = false
|
||||
}
|
||||
|
||||
// Track if this is a new or reconnected peer
|
||||
var isNewPeer = false
|
||||
var isReconnectedPeer = false
|
||||
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
// Check if we have an actual BLE connection to this peer
|
||||
let peripheralUUID = peerToPeripheralUUID[PeerID(str: peerID)]
|
||||
let hasPeripheralConnection = peripheralUUID != nil && peripherals[peripheralUUID!]?.isConnected == true
|
||||
|
||||
// Check if this peer is subscribed to us as a central
|
||||
// Note: We can't identify which specific central is which peer without additional mapping
|
||||
let hasCentralSubscription = centralToPeerID.values.contains(PeerID(str: peerID))
|
||||
|
||||
// Direct announces arrive with full TTL (no prior hop)
|
||||
let isDirectAnnounce = (packet.ttl == messageTTL)
|
||||
|
||||
// Check if we already have this peer (might be reconnecting)
|
||||
let existingPeer = peers[PeerID(str: peerID)]
|
||||
let wasDisconnected = existingPeer?.isConnected == false
|
||||
|
||||
// Set flags for use outside the sync block
|
||||
isNewPeer = (existingPeer == nil)
|
||||
isReconnectedPeer = wasDisconnected
|
||||
|
||||
// Use precomputed verification result
|
||||
let verified = verifiedAnnounce
|
||||
|
||||
// Require verified announce; ignore otherwise (no backward compatibility)
|
||||
if !verified {
|
||||
SecureLogger.warning("❌ Ignoring unverified announce from \(peerID.prefix(8))…", category: .security)
|
||||
return
|
||||
}
|
||||
|
||||
// Update or create peer info
|
||||
if let existing = existingPeer, existing.isConnected {
|
||||
// Update lastSeen and identity info
|
||||
peers[PeerID(str: peerID)] = PeerInfo(
|
||||
peerID: existing.peerID,
|
||||
nickname: announcement.nickname,
|
||||
isConnected: isDirectAnnounce || hasPeripheralConnection || hasCentralSubscription,
|
||||
noisePublicKey: announcement.noisePublicKey,
|
||||
signingPublicKey: announcement.signingPublicKey,
|
||||
isVerifiedNickname: true,
|
||||
lastSeen: Date()
|
||||
)
|
||||
} else {
|
||||
// New peer or reconnecting peer
|
||||
peers[PeerID(str: peerID)] = PeerInfo(
|
||||
peerID: PeerID(str: peerID),
|
||||
nickname: announcement.nickname,
|
||||
isConnected: isDirectAnnounce || hasPeripheralConnection || hasCentralSubscription,
|
||||
noisePublicKey: announcement.noisePublicKey,
|
||||
signingPublicKey: announcement.signingPublicKey,
|
||||
isVerifiedNickname: true,
|
||||
lastSeen: Date()
|
||||
)
|
||||
}
|
||||
|
||||
// Log connection status only for direct connectivity changes; debounce to reduce spam
|
||||
if isDirectAnnounce || hasPeripheralConnection || hasCentralSubscription {
|
||||
let now = Date()
|
||||
if existingPeer == nil {
|
||||
SecureLogger.debug("🆕 New peer: \(announcement.nickname)", category: .session)
|
||||
} else if wasDisconnected {
|
||||
// Debounce 'reconnected' logs within short window
|
||||
if let last = lastReconnectLogAt[PeerID(str: peerID)], now.timeIntervalSince(last) < TransportConfig.bleReconnectLogDebounceSeconds {
|
||||
// Skip duplicate log
|
||||
} else {
|
||||
SecureLogger.debug("🔄 Peer \(announcement.nickname) reconnected", category: .session)
|
||||
lastReconnectLogAt[PeerID(str: peerID)] = now
|
||||
}
|
||||
} else if existingPeer?.nickname != announcement.nickname {
|
||||
SecureLogger.debug("🔄 Peer \(peerID) changed nickname: \(existingPeer?.nickname ?? "Unknown") -> \(announcement.nickname)", category: .session)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Persist cryptographic identity and signing key for robust offline verification
|
||||
do {
|
||||
// Derive fingerprint from Noise public key
|
||||
let hash = SHA256.hash(data: announcement.noisePublicKey)
|
||||
let fingerprint = hash.map { String(format: "%02x", $0) }.joined()
|
||||
identityManager.upsertCryptographicIdentity(
|
||||
fingerprint: fingerprint,
|
||||
noisePublicKey: announcement.noisePublicKey,
|
||||
signingPublicKey: announcement.signingPublicKey,
|
||||
claimedNickname: announcement.nickname
|
||||
)
|
||||
}
|
||||
|
||||
// Record this announce for lightweight rebroadcast buffer (exclude self)
|
||||
// (recentAnnounces has been removed in the refactor)
|
||||
|
||||
// Notify UI on main thread
|
||||
notifyUI { [weak self] in
|
||||
guard let self = self else { return }
|
||||
|
||||
// Get current peer list (after addition)
|
||||
let currentPeerIDs = self.collectionsQueue.sync { Array(self.peers.keys) }
|
||||
|
||||
// Only notify of connection for new or reconnected peers when it is a direct announce
|
||||
if (packet.ttl == self.messageTTL) && (isNewPeer || isReconnectedPeer) {
|
||||
self.delegate?.didConnectToPeer(PeerID(str: peerID))
|
||||
// Schedule initial unicast sync to this peer
|
||||
self.gossipSyncManager?.scheduleInitialSyncToPeer(PeerID(str: peerID), delaySeconds: 1.0)
|
||||
}
|
||||
|
||||
self.requestPeerDataPublish()
|
||||
self.delegate?.didUpdatePeerList(currentPeerIDs)
|
||||
}
|
||||
|
||||
// Track for sync (include our own and others' announces)
|
||||
gossipSyncManager?.onPublicPacketSeen(packet)
|
||||
|
||||
// Send announce back for bidirectional discovery (only once per peer)
|
||||
let announceBackID = "announce-back-\(peerID)"
|
||||
let shouldSendBack = !messageDeduplicator.contains(announceBackID)
|
||||
if shouldSendBack {
|
||||
messageDeduplicator.markProcessed(announceBackID)
|
||||
}
|
||||
|
||||
if shouldSendBack {
|
||||
// Reciprocate announce for bidirectional discovery
|
||||
// Force send to ensure the peer receives our announce
|
||||
sendAnnounce(forceSend: true)
|
||||
}
|
||||
|
||||
// Afterglow: on first-seen peers, schedule a short re-announce to push presence one more hop
|
||||
if isNewPeer {
|
||||
let delay = Double.random(in: 0.3...0.6)
|
||||
messageQueue.asyncAfter(deadline: .now() + delay) { [weak self] in
|
||||
self?.sendAnnounce(forceSend: true)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Handle REQUEST_SYNC: decode payload and respond with missing packets via sync manager
|
||||
private func handleRequestSync(_ packet: BitchatPacket, from peerID: String) {
|
||||
guard let req = RequestSyncPacket.decode(from: packet.payload) else {
|
||||
SecureLogger.warning("⚠️ Malformed REQUEST_SYNC from \(peerID)", category: .session)
|
||||
return
|
||||
}
|
||||
gossipSyncManager?.handleRequestSync(from: PeerID(str: peerID), request: req)
|
||||
}
|
||||
|
||||
// Mention parsing moved to ChatViewModel
|
||||
|
||||
private func handleMessage(_ packet: BitchatPacket, from peerID: String) {
|
||||
// Ignore self-origin public messages except when returned via sync (TTL==0).
|
||||
// This allows our own messages to be surfaced when they come back via
|
||||
// the sync path without re-processing regular relayed copies.
|
||||
if peerID == myPeerID && packet.ttl != 0 { return }
|
||||
|
||||
var accepted = false
|
||||
var senderNickname: String = ""
|
||||
|
||||
// Snapshot peers dictionary to avoid mutating-while-iterating crashes when checking collisions.
|
||||
let peersSnapshot = collectionsQueue.sync { peers }
|
||||
|
||||
// If the packet is from ourselves (e.g., recovered via sync TTL==0), accept immediately
|
||||
if peerID == myPeerID {
|
||||
accepted = true
|
||||
senderNickname = myNickname
|
||||
}
|
||||
else if let info = peersSnapshot[PeerID(str: peerID)], info.isVerifiedNickname {
|
||||
// Known verified peer path
|
||||
accepted = true
|
||||
senderNickname = info.nickname
|
||||
// Handle nickname collisions
|
||||
let hasCollision = peersSnapshot.values.contains { $0.isConnected && $0.nickname == info.nickname && $0.peerID.id != peerID } || (myNickname == info.nickname)
|
||||
if hasCollision {
|
||||
senderNickname += "#" + String(peerID.prefix(4))
|
||||
}
|
||||
} else {
|
||||
// Fallback: verify signature using persisted signing key for this peerID's fingerprint prefix
|
||||
if let signature = packet.signature, let packetData = packet.toBinaryDataForSigning() {
|
||||
// Find candidate identities by peerID prefix (16 hex)
|
||||
let candidates = identityManager.getCryptoIdentitiesByPeerIDPrefix(PeerID(str: peerID))
|
||||
for candidate in candidates {
|
||||
if let signingKey = candidate.signingPublicKey,
|
||||
noiseService.verifySignature(signature, for: packetData, publicKey: signingKey) {
|
||||
accepted = true
|
||||
// Prefer persisted social petname or claimed nickname
|
||||
if let social = identityManager.getSocialIdentity(for: candidate.fingerprint) {
|
||||
senderNickname = social.localPetname ?? social.claimedNickname
|
||||
} else {
|
||||
senderNickname = "anon" + String(peerID.prefix(4))
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
// If still not accepted and this is a sync-returned packet (TTL==0),
|
||||
// accept with a generic nickname so history can be restored even for
|
||||
// peers we haven't verified yet.
|
||||
if !accepted && packet.ttl == 0 {
|
||||
accepted = true
|
||||
senderNickname = "anon" + String(peerID.prefix(4))
|
||||
}
|
||||
}
|
||||
|
||||
// Track broadcast messages for sync (treat nil or 0xFF..0xFF as broadcast)
|
||||
let isBroadcastRecipient: Bool = {
|
||||
guard let r = packet.recipientID else { return true }
|
||||
return r.count == 8 && r.allSatisfy { $0 == 0xFF }
|
||||
}()
|
||||
if isBroadcastRecipient && packet.type == MessageType.message.rawValue {
|
||||
gossipSyncManager?.onPublicPacketSeen(packet)
|
||||
}
|
||||
|
||||
guard accepted else {
|
||||
SecureLogger.warning("🚫 Dropping public message from unverified or unknown peer \(peerID.prefix(8))…", category: .security)
|
||||
return
|
||||
}
|
||||
|
||||
guard let content = String(data: packet.payload, encoding: .utf8) else {
|
||||
SecureLogger.error("❌ Failed to decode message payload as UTF-8", category: .session)
|
||||
return
|
||||
}
|
||||
// Determine if we have a direct link to the sender
|
||||
let hasDirectLink: Bool = collectionsQueue.sync {
|
||||
let perUUID = peerToPeripheralUUID[PeerID(str: peerID)]
|
||||
let perConnected = perUUID != nil && peripherals[perUUID!]?.isConnected == true
|
||||
let hasCentral = centralToPeerID.values.contains(PeerID(str: peerID))
|
||||
return perConnected || hasCentral
|
||||
}
|
||||
|
||||
let pathTag = hasDirectLink ? "direct" : "mesh"
|
||||
SecureLogger.debug("💬 [\(senderNickname)] TTL:\(packet.ttl) (\(pathTag)): \(String(content.prefix(50)))\(content.count > 50 ? "..." : "")", category: .session)
|
||||
|
||||
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
|
||||
notifyUI { [weak self] in
|
||||
self?.delegate?.didReceivePublicMessage(from: PeerID(str: peerID), nickname: senderNickname, content: content, timestamp: ts)
|
||||
}
|
||||
}
|
||||
|
||||
private func handleFileTransfer(_ packet: BitchatPacket, from peerID: String) {
|
||||
private func handleFileTransfer(_ packet: BitchatPacket, from peerID: PeerID) {
|
||||
if peerID == myPeerID && packet.ttl != 0 { return }
|
||||
|
||||
var accepted = false
|
||||
@@ -1602,22 +995,22 @@ final class BLEService: NSObject {
|
||||
if peerID == myPeerID {
|
||||
accepted = true
|
||||
senderNickname = myNickname
|
||||
} else if let info = peersSnapshot[PeerID(str: peerID)], info.isVerifiedNickname {
|
||||
} else if let info = peersSnapshot[peerID], info.isVerifiedNickname {
|
||||
accepted = true
|
||||
senderNickname = info.nickname
|
||||
let hasCollision = peersSnapshot.values.contains { $0.isConnected && $0.nickname == info.nickname && $0.peerID.id != peerID } || (myNickname == info.nickname)
|
||||
let hasCollision = peersSnapshot.values.contains { $0.isConnected && $0.nickname == info.nickname && $0.peerID != peerID } || (myNickname == info.nickname)
|
||||
if hasCollision {
|
||||
senderNickname += "#" + String(peerID.prefix(4))
|
||||
senderNickname += "#" + String(peerID.id.prefix(4))
|
||||
}
|
||||
} else if let info = peersSnapshot[PeerID(str: peerID)], info.isConnected {
|
||||
} else if let info = peersSnapshot[peerID], info.isConnected {
|
||||
accepted = true
|
||||
senderNickname = info.nickname.isEmpty ? "anon" + String(peerID.prefix(4)) : info.nickname
|
||||
let hasCollision = peersSnapshot.values.contains { $0.isConnected && $0.nickname == info.nickname && $0.peerID.id != peerID } || (myNickname == info.nickname)
|
||||
senderNickname = info.nickname.isEmpty ? "anon" + String(peerID.id.prefix(4)) : info.nickname
|
||||
let hasCollision = peersSnapshot.values.contains { $0.isConnected && $0.nickname == info.nickname && $0.peerID != peerID } || (myNickname == info.nickname)
|
||||
if hasCollision {
|
||||
senderNickname += "#" + String(peerID.prefix(4))
|
||||
senderNickname += "#" + String(peerID.id.prefix(4))
|
||||
}
|
||||
} else if let signature = packet.signature, let packetData = packet.toBinaryDataForSigning() {
|
||||
let candidates = identityManager.getCryptoIdentitiesByPeerIDPrefix(PeerID(str: peerID))
|
||||
let candidates = identityManager.getCryptoIdentitiesByPeerIDPrefix(peerID)
|
||||
for candidate in candidates {
|
||||
if let signingKey = candidate.signingPublicKey,
|
||||
noiseService.verifySignature(signature, for: packetData, publicKey: signingKey) {
|
||||
@@ -1625,22 +1018,22 @@ final class BLEService: NSObject {
|
||||
if let social = identityManager.getSocialIdentity(for: candidate.fingerprint) {
|
||||
senderNickname = social.localPetname ?? social.claimedNickname
|
||||
} else {
|
||||
senderNickname = "anon" + String(peerID.prefix(4))
|
||||
senderNickname = "anon" + String(peerID.id.prefix(4))
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
if !accepted && packet.ttl == 0 {
|
||||
accepted = true
|
||||
senderNickname = "anon" + String(peerID.prefix(4))
|
||||
senderNickname = "anon" + String(peerID.id.prefix(4))
|
||||
}
|
||||
} else if packet.ttl == 0 {
|
||||
accepted = true
|
||||
senderNickname = "anon" + String(peerID.prefix(4))
|
||||
senderNickname = "anon" + String(peerID.id.prefix(4))
|
||||
}
|
||||
|
||||
guard accepted else {
|
||||
SecureLogger.warning("🚫 Dropping file transfer from unverified or unknown peer \(peerID.prefix(8))…", category: .security)
|
||||
SecureLogger.warning("🚫 Dropping file transfer from unverified or unknown peer \(peerID.id.prefix(8))…", category: .security)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1673,14 +1066,14 @@ final class BLEService: NSObject {
|
||||
|
||||
// Validate MIME type against whitelist
|
||||
guard isAllowedMimeType(mime) else {
|
||||
SecureLogger.warning("🚫 MIME REJECT: '\(mime)' not in whitelist. Size=\(filePacket.content.count)b from \(peerID.prefix(8))...", category: .security)
|
||||
SecureLogger.warning("🚫 MIME REJECT: '\(mime)' not in whitelist. Size=\(filePacket.content.count)b from \(peerID.id.prefix(8))...", category: .security)
|
||||
return
|
||||
}
|
||||
|
||||
// Validate content matches declared MIME type (magic byte check)
|
||||
guard validateContentMatchesMime(data: filePacket.content, declaredMime: mime) else {
|
||||
let prefix = filePacket.content.prefix(20).map { String(format: "%02x", $0) }.joined(separator: " ")
|
||||
SecureLogger.warning("🚫 MAGIC REJECT: MIME='\(mime)' size=\(filePacket.content.count)b prefix=[\(prefix)] from \(peerID.prefix(8))...", category: .security)
|
||||
SecureLogger.warning("🚫 MAGIC REJECT: MIME='\(mime)' size=\(filePacket.content.count)b prefix=[\(prefix)] from \(peerID.id.prefix(8))...", category: .security)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1735,7 +1128,7 @@ final class BLEService: NSObject {
|
||||
}()
|
||||
|
||||
if isPrivateMessage {
|
||||
updatePeerLastSeen(PeerID(str: peerID))
|
||||
updatePeerLastSeen(peerID)
|
||||
}
|
||||
|
||||
let ts = Date(timeIntervalSince1970: Double(packet.timestamp) / 1000)
|
||||
@@ -1747,10 +1140,10 @@ final class BLEService: NSObject {
|
||||
originalSender: nil,
|
||||
isPrivate: isPrivateMessage,
|
||||
recipientNickname: nil,
|
||||
senderPeerID: PeerID(str: peerID)
|
||||
senderPeerID: peerID
|
||||
)
|
||||
|
||||
SecureLogger.debug("📁 Stored incoming media from \(peerID.prefix(8))… -> \(destination.lastPathComponent)", category: .session)
|
||||
SecureLogger.debug("📁 Stored incoming media from \(peerID.id.prefix(8))… -> \(destination.lastPathComponent)", category: .session)
|
||||
|
||||
notifyUI { [weak self] in
|
||||
self?.delegate?.didReceiveMessage(message)
|
||||
@@ -1770,7 +1163,7 @@ final class BLEService: NSObject {
|
||||
}
|
||||
|
||||
SecureLogger.debug("📤 Sending favorite notification to \(peerID): \(content)", category: .session)
|
||||
sendPrivateMessage(content, to: peerID.id, messageID: UUID().uuidString)
|
||||
sendPrivateMessage(content, to: peerID, messageID: UUID().uuidString)
|
||||
}
|
||||
|
||||
func sendBroadcastAnnounce() {
|
||||
@@ -2119,7 +1512,7 @@ extension BLEService: GossipSyncManager.Delegate {
|
||||
}
|
||||
|
||||
func sendPacket(to peerID: PeerID, packet: BitchatPacket) {
|
||||
sendPacketDirected(packet, to: peerID.id)
|
||||
sendPacketDirected(packet, to: peerID)
|
||||
}
|
||||
|
||||
func signPacketForBroadcast(_ packet: BitchatPacket) -> BitchatPacket {
|
||||
@@ -3164,8 +2557,7 @@ extension BLEService {
|
||||
|
||||
// MARK: Private Message Handling
|
||||
|
||||
private func sendPrivateMessage(_ content: String, to recipientID: String, messageID: String) {
|
||||
let recipientID = PeerID(str: recipientID)
|
||||
private func sendPrivateMessage(_ content: String, to recipientID: PeerID, messageID: String) {
|
||||
SecureLogger.debug("📨 Sending PM to \(recipientID): \(content.prefix(30))...", category: .session)
|
||||
|
||||
// Check if we have an established Noise session
|
||||
@@ -3323,7 +2715,7 @@ extension BLEService {
|
||||
|
||||
// MARK: Fragmentation (Required for messages > BLE MTU)
|
||||
|
||||
private func sendFragmentedPacket(_ packet: BitchatPacket, pad: Bool, maxChunk: Int? = nil, directedOnlyPeer: PeerID? = nil) {
|
||||
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 }
|
||||
// Fragment the unpadded frame; each fragment will be encoded independently
|
||||
|
||||
@@ -3333,6 +2725,8 @@ extension BLEService {
|
||||
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 }
|
||||
|
||||
// Lightweight pacing to reduce floods and allow BLE buffers to drain
|
||||
// Also briefly pause scanning during long fragment trains to save battery
|
||||
let totalFragments = fragments.count
|
||||
@@ -3347,6 +2741,19 @@ extension BLEService {
|
||||
}
|
||||
}
|
||||
}
|
||||
let perFragMs = (directedOnlyPeer != 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()
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
self.activeTransfers[id] = ActiveTransferState(totalFragments: totalFragments, sentFragments: 0, workItems: [])
|
||||
}
|
||||
TransferProgressManager.shared.start(id: id, totalFragments: totalFragments)
|
||||
return id
|
||||
}()
|
||||
|
||||
var scheduledItems: [(item: DispatchWorkItem, index: Int)] = []
|
||||
|
||||
for (index, fragment) in fragments.enumerated() {
|
||||
var payload = Data()
|
||||
@@ -3371,11 +2778,47 @@ extension BLEService {
|
||||
signature: nil,
|
||||
ttl: packet.ttl
|
||||
)
|
||||
// Pace fragments with small jitter to avoid bursts
|
||||
let perFragMs = (directedOnlyPeer != nil || packet.recipientID != nil) ? TransportConfig.bleFragmentSpacingDirectedMs : TransportConfig.bleFragmentSpacingMs
|
||||
|
||||
let workItem = DispatchWorkItem { [weak self] in
|
||||
guard let self = self else { return }
|
||||
if let transferId = transferIdentifier {
|
||||
let isActive = self.collectionsQueue.sync { self.activeTransfers[transferId] != nil }
|
||||
guard isActive else { return }
|
||||
}
|
||||
self.broadcastPacket(fragmentPacket)
|
||||
if let transferId = transferIdentifier {
|
||||
self.markFragmentSent(transferId: transferId)
|
||||
}
|
||||
}
|
||||
|
||||
scheduledItems.append((item: workItem, index: index))
|
||||
}
|
||||
|
||||
if let transferId = transferIdentifier {
|
||||
let workItems = scheduledItems.map { $0.item }
|
||||
collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
guard let self = self, var state = self.activeTransfers[transferId] else { return }
|
||||
state.workItems = workItems
|
||||
self.activeTransfers[transferId] = state
|
||||
}
|
||||
}
|
||||
|
||||
for (workItem, index) in scheduledItems {
|
||||
let delayMs = index * perFragMs
|
||||
messageQueue.asyncAfter(deadline: .now() + .milliseconds(delayMs)) { [weak self] in
|
||||
self?.broadcastPacket(fragmentPacket)
|
||||
messageQueue.asyncAfter(deadline: .now() + .milliseconds(delayMs), execute: workItem)
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Fragmentation (Required for messages > BLE MTU)
|
||||
|
||||
private func markFragmentSent(transferId: String) {
|
||||
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 {
|
||||
self.activeTransfers.removeValue(forKey: transferId)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -3436,6 +2879,7 @@ extension BLEService {
|
||||
}
|
||||
incomingFragments[key] = [:]
|
||||
fragmentMetadata[key] = (originalType, total, Date())
|
||||
SecureLogger.debug("📦 Started fragment assembly id=\(String(format: "%016llx", fragU64)) total=\(total)", category: .session)
|
||||
}
|
||||
|
||||
// Check cumulative size before storing this fragment
|
||||
@@ -3462,6 +2906,7 @@ extension BLEService {
|
||||
}
|
||||
|
||||
incomingFragments[key]?[index] = Data(fragmentData)
|
||||
SecureLogger.debug("📦 Fragment \(index + 1)/\(total) (len=\(fragmentData.count)) for id=\(String(format: "%016llx", fragU64))", category: .session)
|
||||
|
||||
// Check if complete
|
||||
if let fragments = incomingFragments[key], fragments.count == total {
|
||||
@@ -3485,6 +2930,7 @@ extension BLEService {
|
||||
|
||||
// Decode the original packet bytes we reassembled, so flags/compression are preserved
|
||||
if let originalPacket = BinaryProtocol.decode(reassembled) {
|
||||
SecureLogger.debug("✅ Reassembled packet id=\(String(format: "%016llx", fragU64)) type=\(originalPacket.type) bytes=\(reassembled.count)", category: .session)
|
||||
handleReceivedPacket(originalPacket, from: peerID)
|
||||
} else {
|
||||
SecureLogger.error("❌ Failed to decode reassembled packet (type=\(originalType), total=\(total))", category: .session)
|
||||
@@ -3581,7 +3027,7 @@ extension BLEService {
|
||||
handleFragment(packet, from: senderID)
|
||||
|
||||
case .fileTransfer:
|
||||
handleFileTransfer(packet, from: senderID.id)
|
||||
handleFileTransfer(packet, from: senderID)
|
||||
|
||||
case .leave:
|
||||
handleLeave(packet, from: senderID)
|
||||
@@ -3759,18 +3205,7 @@ extension BLEService {
|
||||
)
|
||||
|
||||
// Record this announce for lightweight rebroadcast buffer (exclude self)
|
||||
if peerID != myPeerID {
|
||||
collectionsQueue.async(flags: .barrier) { [weak self] in
|
||||
guard let self = self else { return }
|
||||
self.recentAnnounceBySender[peerID] = packet
|
||||
if !self.recentAnnounceOrder.contains(peerID) { self.recentAnnounceOrder.append(peerID) }
|
||||
// Trim to cap, oldest first
|
||||
while self.recentAnnounceOrder.count > self.recentAnnounceBufferCap {
|
||||
let victim = self.recentAnnounceOrder.removeFirst()
|
||||
self.recentAnnounceBySender.removeValue(forKey: victim)
|
||||
}
|
||||
}
|
||||
}
|
||||
// (recentAnnounces has been removed in the refactor)
|
||||
|
||||
// Notify UI on main thread
|
||||
notifyUI { [weak self] in
|
||||
|
||||
Reference in New Issue
Block a user