From 0dd999af6bf4265063ccfe97febd67210701d8fc Mon Sep 17 00:00:00 2001 From: Islam <2553451+qalandarov@users.noreply.github.com> Date: Sun, 19 Oct 2025 19:35:41 +0100 Subject: [PATCH] PeerID 29/n: BLEService + remove dupe funcs from #823 (#838) * 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> --- bitchat/Services/BLEService.swift | 757 ++++-------------------------- 1 file changed, 96 insertions(+), 661 deletions(-) diff --git a/bitchat/Services/BLEService.swift b/bitchat/Services/BLEService.swift index e9693306..88c31965 100644 --- a/bitchat/Services/BLEService.swift +++ b/bitchat/Services/BLEService.swift @@ -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.. 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.. 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.. 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