diff --git a/bitchat/Protocols/BinaryProtocol.swift b/bitchat/Protocols/BinaryProtocol.swift index b95e4410..c7024e56 100644 --- a/bitchat/Protocols/BinaryProtocol.swift +++ b/bitchat/Protocols/BinaryProtocol.swift @@ -19,11 +19,12 @@ extension Data { } // Binary Protocol Format: -// Header (Fixed 13 bytes): +// Header (Fixed 17 bytes): // - Version: 1 byte // - Type: 1 byte // - TTL: 1 byte // - Timestamp: 8 bytes (UInt64) +// - SequenceNumber: 4 bytes (UInt32) // - Flags: 1 byte (bit 0: hasRecipient, bit 1: hasSignature) // - PayloadLength: 2 bytes (UInt16) // @@ -34,7 +35,7 @@ extension Data { // - Signature: 64 bytes (if hasSignature flag set) struct BinaryProtocol { - static let headerSize = 13 + static let headerSize = 17 static let senderIDSize = 8 static let recipientIDSize = 8 static let signatureSize = 64 @@ -77,6 +78,11 @@ struct BinaryProtocol { data.append(UInt8((packet.timestamp >> (i * 8)) & 0xFF)) } + // Sequence number (4 bytes, big-endian) + for i in (0..<4).reversed() { + data.append(UInt8((packet.sequenceNumber >> (i * 8)) & 0xFF)) + } + // Flags var flags: UInt8 = 0 if packet.recipientID != nil { @@ -165,6 +171,13 @@ struct BinaryProtocol { } offset += 8 + // Sequence number + let sequenceData = unpaddedData[offset..(maxSize: 1000) // Bounded to prevent memory growth + + // Message state tracking for better duplicate detection + private struct MessageState { + let firstSeen: Date + var lastSeen: Date + var seenCount: Int + var relayed: Bool + var acknowledged: Bool + + mutating func updateSeen() { + lastSeen = Date() + seenCount += 1 + } + } + + private var processedMessages = [String: MessageState]() // Track full message state + private let processedMessagesLock = NSLock() + private let maxProcessedMessages = 10000 + private let maxTTL: UInt8 = 7 // Maximum hops for long-distance delivery private var announcedToPeers = Set() // Track which peers we've announced to private var announcedPeers = Set() // Track peers who have already been announced @@ -391,10 +409,22 @@ class BluetoothMeshService: NSObject { private var ackTimer: Timer? private var advertisingTimer: Timer? // Timer for interval-based advertising - // Timing randomization for privacy + // Timing randomization for privacy (now with exponential distribution) private let minMessageDelay: TimeInterval = 0.01 // 10ms minimum for faster sync private let maxMessageDelay: TimeInterval = 0.1 // 100ms maximum for faster sync + // Dynamic RSSI threshold based on network size (smart compromise) + private var dynamicRSSIThreshold: Int { + let peerCount = activePeers.count + if peerCount < 5 { + return -95 // Maximum range when network is small + } else if peerCount < 10 { + return -92 // Slightly reduced range + } else { + return -90 // Conservative for larger networks + } + } + // Fragment handling with security limits private var incomingFragments: [String: [Int: Data]] = [:] // fragmentID -> [index: data] private var fragmentMetadata: [String: (originalType: UInt8, totalFragments: Int, timestamp: Date)] = [:] @@ -434,6 +464,15 @@ class BluetoothMeshService: NSObject { return max(activePeers.count, connectedPeripherals.count) } + // Sequence number tracking for duplicate detection + private var currentSequenceNumber: UInt32 = 0 + private var peerSequenceNumbers: [String: UInt32] = [:] // Track last seen sequence per peer + + private func getNextSequenceNumber() -> UInt32 { + currentSequenceNumber += 1 + return currentSequenceNumber + } + // Adaptive parameters based on network size private var adaptiveTTL: UInt8 { // Keep TTL high enough for messages to travel far @@ -450,21 +489,51 @@ class BluetoothMeshService: NSObject { } private var adaptiveRelayProbability: Double { - // Keep relay probability high enough to ensure delivery + // Adjust relay probability based on network size + // Small networks don't need 100% relay - RSSI will handle edge nodes let networkSize = estimatedNetworkSize - if networkSize <= 10 { - return 1.0 // 100% for small networks + if networkSize <= 2 { + return 0.0 // 0% for 2 nodes - no relay needed with only 2 peers + } else if networkSize <= 5 { + return 0.5 // 50% for very small networks + } else if networkSize <= 10 { + return 0.6 // 60% for small networks } else if networkSize <= 30 { - return 0.85 // 85% - most nodes relay + return 0.7 // 70% for medium networks } else if networkSize <= 50 { - return 0.7 // 70% - still high probability + return 0.6 // 60% for larger networks } else if networkSize <= 100 { - return 0.55 // 55% - over half relay + return 0.5 // 50% for big networks } else { - return 0.4 // 40% minimum - never go below this + return 0.4 // 40% minimum for very large networks } } + // Relay cancellation mechanism to prevent duplicate relays + private var pendingRelays: [String: DispatchWorkItem] = [:] // messageID -> relay task + private let pendingRelaysLock = NSLock() + private let relayCancellationWindow: TimeInterval = 0.05 // 50ms window to detect other relays + + // Global memory limits to prevent unbounded growth + private let maxPendingPrivateMessages = 100 + private let maxCachedMessagesSentToPeer = 1000 + private let maxMemoryUsageBytes = 50 * 1024 * 1024 // 50MB limit + private var memoryCleanupTimer: Timer? + + // Track acknowledged packets to prevent unnecessary retries + private var acknowledgedPackets = Set() + private let acknowledgedPacketsLock = NSLock() + + // Rate limiting for flood control with separate limits for different message types + private var messageRateLimiter: [String: [Date]] = [:] // peerID -> recent message timestamps + private var protocolMessageRateLimiter: [String: [Date]] = [:] // peerID -> protocol msg timestamps + private let rateLimiterLock = NSLock() + private let rateLimitWindow: TimeInterval = 60.0 // 1 minute window + private let maxChatMessagesPerPeerPerMinute = 300 // Allow fast typing (5 msgs/sec) + private let maxProtocolMessagesPerPeerPerMinute = 100 // Protocol messages + private let maxTotalMessagesPerMinute = 2000 // Increased for legitimate use + private var totalMessageTimestamps: [Date] = [] + // BLE advertisement for lightweight presence private var advertisementData: [String: Any] = [:] private var isAdvertising = false @@ -825,7 +894,9 @@ class BluetoothMeshService: NSObject { self.messageBloomFilter = OptimizedBloomFilter.adaptive(for: networkSize) // Clear other duplicate detection sets + self.processedMessagesLock.lock() self.processedMessages.removeAll() + self.processedMessagesLock.unlock() } } @@ -855,6 +926,9 @@ class BluetoothMeshService: NSObject { self?.cleanExpiredIdentityCache() } + // Start memory cleanup timer (every minute) + startMemoryCleanupTimer() + // Log handshake states periodically for debugging and clean up stale states #if DEBUG Timer.scheduledTimer(withTimeInterval: 30.0, repeats: true) { [weak self] _ in @@ -924,6 +998,7 @@ class BluetoothMeshService: NSObject { aggregationTimer?.invalidate() cleanupTimer?.invalidate() rotationTimer?.invalidate() + memoryCleanupTimer?.invalidate() } @objc private func appWillTerminate() { @@ -1011,27 +1086,19 @@ class BluetoothMeshService: NSObject { type: MessageType.announce.rawValue, ttl: 3, // Increase TTL so announce reaches all peers senderID: myPeerID, - payload: Data(vm.nickname.utf8) + payload: Data(vm.nickname.utf8), + sequenceNumber: getNextSequenceNumber() ) - // Initial send with random delay - let initialDelay = self.randomDelay() + // Single send with smart collision avoidance + let initialDelay = self.smartCollisionAvoidanceDelay(baseDelay: self.randomDelay()) DispatchQueue.main.asyncAfter(deadline: .now() + initialDelay) { [weak self] in self?.broadcastPacket(announcePacket) // Also send Noise identity announcement self?.sendNoiseIdentityAnnounce() } - - // Send multiple times for reliability with jittered delays - for baseDelay in [0.2, 0.5, 1.0] { - let jitteredDelay = baseDelay + self.randomDelay() - DispatchQueue.main.asyncAfter(deadline: .now() + jitteredDelay) { [weak self] in - guard let self = self else { return } - self.broadcastPacket(announcePacket) - } - } } func startAdvertising() { @@ -1152,7 +1219,8 @@ class BluetoothMeshService: NSObject { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), // milliseconds payload: messageData, signature: nil, - ttl: self.adaptiveTTL + ttl: self.adaptiveTTL, + sequenceNumber: self.getNextSequenceNumber() ) // Track this message to prevent duplicate sends @@ -1170,18 +1238,11 @@ class BluetoothMeshService: NSObject { self.recentlySentMessages.remove(msgID) } - // Add random delay before initial send - let initialDelay = self.randomDelay() + // Single send with smart collision avoidance + let initialDelay = self.smartCollisionAvoidanceDelay(baseDelay: self.randomDelay()) DispatchQueue.main.asyncAfter(deadline: .now() + initialDelay) { [weak self] in self?.broadcastPacket(packet) } - - // Single retry for reliability - let retryDelay = 0.3 + self.randomDelay() - DispatchQueue.main.asyncAfter(deadline: .now() + retryDelay) { [weak self] in - self?.broadcastPacket(packet) - // Re-sending message - } } } } @@ -1287,7 +1348,8 @@ class BluetoothMeshService: NSObject { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: encryptedPayload, signature: nil, - ttl: 3 + ttl: 3, + sequenceNumber: self.getNextSequenceNumber() ) self.broadcastPacket(outerPacket) @@ -1312,7 +1374,8 @@ class BluetoothMeshService: NSObject { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: encryptedPayload, signature: nil, // ACKs don't need signatures - ttl: 3 // Limited TTL for ACKs + ttl: 3, // Limited TTL for ACKs + sequenceNumber: self.getNextSequenceNumber() ) // Send immediately without delay (ACKs should be fast) @@ -1371,7 +1434,8 @@ class BluetoothMeshService: NSObject { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: receiptData, signature: nil, - ttl: 3 + ttl: 3, + sequenceNumber: self.getNextSequenceNumber() ) // Encrypt the entire inner packet @@ -1386,7 +1450,8 @@ class BluetoothMeshService: NSObject { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: encryptedInnerData, signature: nil, - ttl: 3 + ttl: 3, + sequenceNumber: self.getNextSequenceNumber() ) SecureLogger.log("Sending encrypted read receipt for message \(receipt.originalMessageID) to \(recipientID)", category: SecureLogger.noise, level: .info) @@ -1440,7 +1505,8 @@ class BluetoothMeshService: NSObject { type: MessageType.announce.rawValue, ttl: 3, // Allow relay for better reach senderID: myPeerID, - payload: Data(vm.nickname.utf8) + payload: Data(vm.nickname.utf8), + sequenceNumber: getNextSequenceNumber() ) if let data = packet.toBinaryData() { @@ -1467,7 +1533,8 @@ class BluetoothMeshService: NSObject { type: MessageType.leave.rawValue, ttl: 1, // Don't relay leave messages senderID: myPeerID, - payload: Data(vm.nickname.utf8) + payload: Data(vm.nickname.utf8), + sequenceNumber: getNextSequenceNumber() ) broadcastPacket(packet) @@ -1521,7 +1588,9 @@ class BluetoothMeshService: NSObject { hasNotifiedNetworkAvailable = false networkBecameEmptyTime = nil lastNetworkNotificationTime = nil + processedMessagesLock.lock() processedMessages.removeAll() + processedMessagesLock.unlock() incomingFragments.removeAll() // Clear all encryption queues @@ -1530,6 +1599,14 @@ class BluetoothMeshService: NSObject { encryptionQueuesLock.unlock() fragmentMetadata.removeAll() + // Cancel all pending relays + pendingRelaysLock.lock() + for (_, relay) in pendingRelays { + relay.cancel() + } + pendingRelays.removeAll() + pendingRelaysLock.unlock() + // Clear peer tracking lastHeardFromPeer.removeAll() @@ -1900,8 +1977,12 @@ class BluetoothMeshService: NSObject { return } + // Track which peers we've sent to (to avoid duplicates) + var sentToPeers = Set() + // Send to connected peripherals (as central) var sentToPeripherals = 0 + // Broadcasting to connected peripherals for (peerID, peripheral) in connectedPeripherals { if let characteristic = peripheralCharacteristics[peripheral] { // Check if peripheral is connected before writing @@ -1909,8 +1990,10 @@ class BluetoothMeshService: NSObject { // Additional safety check for characteristic properties if characteristic.properties.contains(.write) || characteristic.properties.contains(.writeWithoutResponse) { + // Writing packet to peripheral writeToPeripheral(data, peripheral: peripheral, characteristic: characteristic, peerID: peerID) sentToPeripherals += 1 + sentToPeers.insert(peerID) } } else { if let peerID = connectedPeripherals.first(where: { $0.value == peripheral })?.key { @@ -1921,15 +2004,18 @@ class BluetoothMeshService: NSObject { } } - // Send to subscribed centrals (as peripheral) + // Send to subscribed centrals (as peripheral) - but only if we didn't already send via peripheral connections var sentToCentrals = 0 - if let char = characteristic, !subscribedCentrals.isEmpty { - // Send to all subscribed centrals - // Note: Large packets should already be fragmented by the check at the beginning of broadcastPacket + if let char = characteristic, !subscribedCentrals.isEmpty && sentToPeripherals == 0 { + // Only send to centrals if we haven't sent via peripheral connections + // This prevents duplicate sends in 2-peer networks where peers connect both ways + // Broadcasting to subscribed centrals let success = peripheralManager?.updateValue(data, for: char, onSubscribedCentrals: nil) ?? false if success { sentToCentrals = subscribedCentrals.count } + } else if sentToPeripherals > 0 && !subscribedCentrals.isEmpty { + // Skip central broadcast - already sent via peripherals } // If no peers received the message, add to retry queue ONLY if it's our own message @@ -2031,33 +2117,128 @@ class BluetoothMeshService: NSObject { SecureLogger.log("Dropped message with stale timestamp. Age: \(timeDiff/1000)s from \(senderID)", category: SecureLogger.security, level: .warning) return } + + // Log message type for debugging + let messageTypeName: String + switch MessageType(rawValue: packet.type) { + case .message: + messageTypeName = "MESSAGE" + case .protocolAck: + messageTypeName = "PROTOCOL_ACK" + case .protocolNack: + messageTypeName = "PROTOCOL_NACK" + case .noiseHandshakeInit: + messageTypeName = "NOISE_HANDSHAKE_INIT" + case .noiseHandshakeResp: + messageTypeName = "NOISE_HANDSHAKE_RESP" + case .noiseIdentityAnnounce: + messageTypeName = "NOISE_IDENTITY_ANNOUNCE" + case .noiseEncrypted: + messageTypeName = "NOISE_ENCRYPTED" + case .leave: + messageTypeName = "LEAVE" + case .readReceipt: + messageTypeName = "READ_RECEIPT" + case .versionHello: + messageTypeName = "VERSION_HELLO" + case .versionAck: + messageTypeName = "VERSION_ACK" + case .systemValidation: + messageTypeName = "SYSTEM_VALIDATION" + default: + messageTypeName = "UNKNOWN(\(packet.type))" + } + + // Processing packet + + // Rate limiting check with message type awareness + let isHighPriority = [MessageType.protocolAck.rawValue, + MessageType.protocolNack.rawValue, + MessageType.noiseHandshakeInit.rawValue, + MessageType.noiseHandshakeResp.rawValue, + MessageType.noiseIdentityAnnounce.rawValue, + MessageType.leave.rawValue].contains(packet.type) + + + if senderID != self.myPeerID && isRateLimited(peerID: senderID, messageType: packet.type) { + if !isHighPriority { + SecureLogger.log("RATE_LIMITED: Dropped \(messageTypeName) from \(senderID) (seq#\(packet.sequenceNumber))", + category: SecureLogger.security, level: .warning) + return + } else { + SecureLogger.log("RATE_LIMITED: Allowing high-priority \(messageTypeName) from \(senderID)", + category: SecureLogger.security, level: .info) + } + } + + // Record message for rate limiting based on type + if senderID != self.myPeerID && !isHighPriority { + recordMessage(from: senderID, messageType: packet.type) + } - // For fragments, include packet type in messageID to avoid dropping CONTINUE/END fragments - let messageID: String - if packet.type == MessageType.fragmentStart.rawValue || - packet.type == MessageType.fragmentContinue.rawValue || - packet.type == MessageType.fragmentEnd.rawValue { - // Include both type and payload hash for fragments to ensure uniqueness - messageID = "\(packet.timestamp)-\(packet.senderID.hexEncodedString())-\(packet.type)-\(packet.payload.hashValue)" - } else { - // Include payload hash for absolute uniqueness (handles same-second messages) - messageID = "\(packet.timestamp)-\(packet.senderID.hexEncodedString())-\(packet.payload.prefix(64).hashValue)" + // Sequence-based duplicate detection for more reliability + let messageID = "\(senderID)-\(packet.sequenceNumber)" + + // Check if we've seen this exact message before (sequence + sender) + processedMessagesLock.lock() + if let existingState = processedMessages[messageID] { + // Update the state + var updatedState = existingState + updatedState.updateSeen() + processedMessages[messageID] = updatedState + processedMessagesLock.unlock() + + SecureLogger.log("Dropped duplicate message seq#\(packet.sequenceNumber) from \(senderID) (seen \(updatedState.seenCount) times)", category: SecureLogger.security, level: .debug) + // Cancel any pending relay for this message + cancelPendingRelay(messageID: messageID) + return + } + processedMessagesLock.unlock() + + // Check if this is an old message (sequence number rollback) + if let lastSeenSeq = peerSequenceNumbers[senderID] { + // Allow some out-of-order tolerance (up to 100 messages) + if packet.sequenceNumber < lastSeenSeq && (lastSeenSeq - packet.sequenceNumber) > 100 { + SecureLogger.log("Dropped old message seq#\(packet.sequenceNumber) from \(senderID), last seen: \(lastSeenSeq)", category: SecureLogger.security, level: .debug) + return + } } // Use bloom filter for efficient duplicate detection if messageBloomFilter.contains(messageID) { - // Also check exact set for accuracy (bloom filter can have false positives) - if processedMessages.contains(messageID) { - SecureLogger.log("Dropped duplicate message: \(messageID.prefix(20))... from \(senderID)", category: SecureLogger.security, level: .debug) - return - } else { - // False positive from Bloom filter - SecureLogger.log("Bloom filter false positive for message: \(messageID.prefix(20))...", category: SecureLogger.security, level: .debug) + // Double check with exact set (bloom filter can have false positives) + processedMessagesLock.lock() + let isProcessed = processedMessages[messageID] != nil + processedMessagesLock.unlock() + if !isProcessed { + SecureLogger.log("Bloom filter false positive for message: \(messageID)", category: SecureLogger.security, level: .debug) } } + // Record this message as processed messageBloomFilter.insert(messageID) - processedMessages.insert(messageID) + processedMessagesLock.lock() + processedMessages[messageID] = MessageState( + firstSeen: Date(), + lastSeen: Date(), + seenCount: 1, + relayed: false, + acknowledged: false + ) + + // Prune old entries if needed + if processedMessages.count > maxProcessedMessages { + // Remove oldest entries + let sortedByFirstSeen = processedMessages.sorted { $0.value.firstSeen < $1.value.firstSeen } + let toRemove = sortedByFirstSeen.prefix(processedMessages.count - maxProcessedMessages + 1000) + for (key, _) in toRemove { + processedMessages.removeValue(forKey: key) + } + } + processedMessagesLock.unlock() + + // Update last seen sequence number for this peer + peerSequenceNumbers[senderID] = max(packet.sequenceNumber, peerSequenceNumbers[senderID] ?? 0) // Log statistics periodically if messageBloomFilter.insertCount % 100 == 0 { @@ -2130,24 +2311,25 @@ class BluetoothMeshService: NSObject { var relayPacket = packet relayPacket.ttl -= 1 if relayPacket.ttl > 0 { - // Probabilistic flooding with smart relay decisions - let relayProb = self.adaptiveRelayProbability + // RSSI-based relay probability + let rssiValue = peerRSSI[senderID]?.intValue ?? -70 + let relayProb = self.calculateRelayProbability(baseProb: self.adaptiveRelayProbability, rssi: rssiValue) - // Always relay if TTL is high (fresh messages need to spread) - // or if we have few peers (ensure coverage in sparse networks) - let shouldRelay = relayPacket.ttl >= 4 || - self.activePeers.count <= 3 || - Double.random(in: 0...1) < relayProb + // Relay based on probability only - no TTL boost if base probability is 0 + let effectiveProb = relayProb > 0 ? relayProb : 0.0 + let shouldRelay = effectiveProb > 0 && Double.random(in: 0...1) < effectiveProb if shouldRelay { - SecureLogger.log("Relaying broadcast from \(senderID), TTL: \(relayPacket.ttl), peers: \(self.activePeers.count)", category: SecureLogger.noise, level: .debug) - // Add random delay to prevent collision storms - let delay = Double.random(in: minMessageDelay...maxMessageDelay) - DispatchQueue.main.asyncAfter(deadline: .now() + delay) { [weak self] in - self?.broadcastPacket(relayPacket) + // Relaying broadcast + // High priority messages relay immediately, others use exponential delay + if self.isHighPriorityMessage(type: relayPacket.type) { + self.broadcastPacket(relayPacket) + } else { + let delay = self.exponentialRelayDelay() + self.scheduleRelay(relayPacket, messageID: messageID, delay: delay) } } else { - SecureLogger.log("Dropped broadcast relay from \(senderID), TTL: \(relayPacket.ttl), prob: \(relayProb)", category: SecureLogger.noise, level: .debug) + // Dropped broadcast relay } } @@ -2237,19 +2419,21 @@ class BluetoothMeshService: NSObject { } } - // Private messages are important - use higher relay probability - let relayProb = min(self.adaptiveRelayProbability + 0.15, 1.0) // Boost by 15% + // Private messages are important - use RSSI-based relay with boost + let rssiValue = peerRSSI[senderID]?.intValue ?? -70 + let baseProb = min(self.adaptiveRelayProbability + 0.15, 1.0) // Boost by 15% + let relayProb = self.calculateRelayProbability(baseProb: baseProb, rssi: rssiValue) - // Always relay if TTL is high or we have few peers - let shouldRelay = relayPacket.ttl >= 4 || - self.activePeers.count <= 3 || - Double.random(in: 0...1) < relayProb + // Relay based on probability only - no forced relay for small networks + let shouldRelay = Double.random(in: 0...1) < relayProb if shouldRelay { - // Add random delay to prevent collision storms - let delay = Double.random(in: minMessageDelay...maxMessageDelay) - DispatchQueue.main.asyncAfter(deadline: .now() + delay) { [weak self] in - self?.broadcastPacket(relayPacket) + // High priority messages relay immediately, others use exponential delay + if self.isHighPriorityMessage(type: relayPacket.type) { + self.broadcastPacket(relayPacket) + } else { + let delay = self.exponentialRelayDelay() + self.scheduleRelay(relayPacket, messageID: messageID, delay: delay) } } } else { @@ -2475,9 +2659,7 @@ class BluetoothMeshService: NSObject { // Add small delay to prevent collision let delay = Double.random(in: 0.1...0.3) - DispatchQueue.main.asyncAfter(deadline: .now() + delay) { [weak self] in - self?.broadcastPacket(relayPacket) - } + self.scheduleRelay(relayPacket, messageID: messageID, delay: delay) } } else { } @@ -2523,7 +2705,8 @@ class BluetoothMeshService: NSObject { var relayPacket = packet relayPacket.ttl -= 1 if relayPacket.ttl > 0 { - self.broadcastPacket(relayPacket) + let delay = self.exponentialRelayDelay() + self.scheduleRelay(relayPacket, messageID: messageID, delay: delay) } @@ -2588,7 +2771,8 @@ class BluetoothMeshService: NSObject { var relayPacket = packet relayPacket.ttl -= 1 - self.broadcastPacket(relayPacket) + let delay = self.exponentialRelayDelay() + self.scheduleRelay(relayPacket, messageID: messageID, delay: delay) } case .readReceipt: @@ -2597,17 +2781,17 @@ class BluetoothMeshService: NSObject { isPeerIDOurs(recipientIDData.hexEncodedString()) { // This read receipt is for us let senderID = packet.senderID.hexEncodedString() - SecureLogger.log("Received read receipt from \(senderID)", category: SecureLogger.session, level: .info) + // Received read receipt // Check if payload is already decrypted (came through Noise) if let receipt = ReadReceipt.fromBinaryData(packet.payload) { // Already decrypted - process directly - SecureLogger.log("Processing read receipt for message \(receipt.originalMessageID) from \(receipt.readerID)", category: SecureLogger.session, level: .info) + // Processing read receipt DispatchQueue.main.async { self.delegate?.didReceiveReadReceipt(receipt) } } else if let receipt = ReadReceipt.decode(from: packet.payload) { // Fallback to JSON for backward compatibility - SecureLogger.log("Processing read receipt (JSON) for message \(receipt.originalMessageID) from \(receipt.readerID)", category: SecureLogger.session, level: .info) + // Processing read receipt (JSON) DispatchQueue.main.async { self.delegate?.didReceiveReadReceipt(receipt) } @@ -2634,7 +2818,8 @@ class BluetoothMeshService: NSObject { // Relay the read receipt if not for us var relayPacket = packet relayPacket.ttl -= 1 - self.broadcastPacket(relayPacket) + let delay = self.exponentialRelayDelay() + self.scheduleRelay(relayPacket, messageID: messageID, delay: delay) } case .noiseIdentityAnnounce: @@ -2646,11 +2831,11 @@ class BluetoothMeshService: NSObject { !isPeerIDOurs(recipientID.hexEncodedString()) { // Not for us, relay if TTL > 0 if packet.ttl > 0 { - SecureLogger.log("Relaying identity announce packet to \(recipientID.hexEncodedString()), TTL: \(packet.ttl)", - category: SecureLogger.session, level: .debug) + // Relay identity announce var relayPacket = packet relayPacket.ttl -= 1 - broadcastPacket(relayPacket) + let delay = self.exponentialRelayDelay() + self.scheduleRelay(relayPacket, messageID: messageID, delay: delay) } return } @@ -2778,7 +2963,7 @@ class BluetoothMeshService: NSObject { !isPeerIDOurs(recipientID.hexEncodedString()) { // Not for us, relay if TTL > 0 if packet.ttl > 0 { - SecureLogger.log("Relaying handshake init packet, TTL: \(packet.ttl)", category: SecureLogger.session, level: .debug) + // Relay handshake init var relayPacket = packet relayPacket.ttl -= 1 broadcastPacket(relayPacket) @@ -2845,7 +3030,7 @@ class BluetoothMeshService: NSObject { if !isPeerIDOurs(recipientIDStr) { // Not for us, relay if TTL > 0 if packet.ttl > 0 { - SecureLogger.log("Relaying handshake response packet, TTL: \(packet.ttl)", category: SecureLogger.session, level: .debug) + // Relay handshake response var relayPacket = packet relayPacket.ttl -= 1 broadcastPacket(relayPacket) @@ -2979,7 +3164,8 @@ class BluetoothMeshService: NSObject { timestamp: packet.timestamp, // Use original timestamp payload: fragmentPayload, signature: nil, // Fragments don't need signatures - ttl: packet.ttl + ttl: packet.ttl, + sequenceNumber: getNextSequenceNumber() ) // Send fragments with linear delay @@ -3178,9 +3364,10 @@ extension BluetoothMeshService: CBCentralManagerDelegate { // Optimize for 300m range - only connect to strong enough signals let rssiValue = RSSI.intValue - // Filter out very weak signals (below -90 dBm) to save battery - guard rssiValue > -90 else { - // Ignoring peripheral due to very weak signal + // Dynamic RSSI threshold based on peer count (smart compromise) + let rssiThreshold = dynamicRSSIThreshold + guard rssiValue > rssiThreshold else { + // Ignoring peripheral due to weak signal (threshold: \(rssiThreshold) dBm) return } @@ -3229,12 +3416,17 @@ extension BluetoothMeshService: CBCentralManagerDelegate { if let pooledPeripheral = connectionPool[peripheralID] { // Reuse existing peripheral from pool if pooledPeripheral.state == CBPeripheralState.disconnected { - // Reconnect if disconnected - central.connect(pooledPeripheral, options: [ + // Reconnect if disconnected with optimized parameters + let connectionOptions: [String: Any] = [ CBConnectPeripheralOptionNotifyOnConnectionKey: true, CBConnectPeripheralOptionNotifyOnDisconnectionKey: true, CBConnectPeripheralOptionNotifyOnNotificationKey: true - ]) + ] + + // Smart compromise: would set low latency for small networks if API supported it + // iOS/macOS don't expose connection interval control in public API + + central.connect(pooledPeripheral, options: connectionOptions) } return } @@ -3265,13 +3457,16 @@ extension BluetoothMeshService: CBCentralManagerDelegate { // Only attempt if under max attempts if attempts < maxConnectionAttempts { - // Use optimized connection parameters for better range + // Use optimized connection parameters based on peer count let connectionOptions: [String: Any] = [ CBConnectPeripheralOptionNotifyOnConnectionKey: true, CBConnectPeripheralOptionNotifyOnDisconnectionKey: true, CBConnectPeripheralOptionNotifyOnNotificationKey: true ] + // Smart compromise: would set low latency for small networks if API supported it + // iOS/macOS don't expose connection interval control in public API + central.connect(peripheral, options: connectionOptions) } } @@ -3279,10 +3474,10 @@ extension BluetoothMeshService: CBCentralManagerDelegate { func centralManager(_ central: CBCentralManager, didConnect peripheral: CBPeripheral) { let tempID = peripheral.identifier.uuidString - SecureLogger.log("Peripheral connected: \(tempID)", category: SecureLogger.session, level: .info) + // Peripheral connected - // Log current mapping state for debugging - SecureLogger.log("Current peerIDByPeripheralID mappings: \(peerIDByPeripheralID.keys.joined(separator: ", "))", + // Current mapping state for debugging + SecureLogger.log("Current peerIDByPeripheralID mappings:", category: SecureLogger.session, level: .debug) peripheral.delegate = self @@ -3532,20 +3727,18 @@ extension BluetoothMeshService: CBPeripheralDelegate { self.sendVersionHello(to: peripheral) // Send announce packet after version negotiation completes - // Send multiple times for reliability if let vm = self.delegate as? ChatViewModel { - // Send announces multiple times with delays - for delay in [0.3, 0.8, 1.5] { - DispatchQueue.main.asyncAfter(deadline: .now() + delay) { [weak self] in - guard let self = self else { return } - let announcePacket = BitchatPacket( - type: MessageType.announce.rawValue, - ttl: 3, - senderID: self.myPeerID, - payload: Data(vm.nickname.utf8) - ) - self.broadcastPacket(announcePacket) - } + // Send single announce with slight delay + DispatchQueue.main.asyncAfter(deadline: .now() + 0.3) { [weak self] in + guard let self = self else { return } + let announcePacket = BitchatPacket( + type: MessageType.announce.rawValue, + ttl: 3, + senderID: self.myPeerID, + payload: Data(vm.nickname.utf8), + sequenceNumber: self.getNextSequenceNumber() + ) + self.broadcastPacket(announcePacket) } // Also send targeted announce to this specific peripheral @@ -3559,7 +3752,8 @@ extension BluetoothMeshService: CBPeripheralDelegate { type: MessageType.announce.rawValue, ttl: 3, senderID: self.myPeerID, - payload: Data(vm.nickname.utf8) + payload: Data(vm.nickname.utf8), + sequenceNumber: self.getNextSequenceNumber() ) if let data = announcePacket.toBinaryData() { self.writeToPeripheral(data, peripheral: peripheral, characteristic: characteristic, peerID: nil) @@ -3896,6 +4090,359 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { return TimeInterval.random(in: minMessageDelay...maxMessageDelay) } + // MARK: - Range Optimization Methods + + // Calculate relay probability based on RSSI (edge nodes relay more) + private func calculateRelayProbability(baseProb: Double, rssi: Int) -> Double { + // RSSI typically ranges from -30 (very close) to -90 (far) + // We want edge nodes (weak RSSI) to relay more aggressively + let rssiFactor = Double(100 + rssi) / 100.0 // -90 = 0.1, -70 = 0.3, -50 = 0.5 + + // Invert the factor so weak signals relay more + let edgeFactor = max(0.5, 2.0 - rssiFactor) + + return min(1.0, baseProb * edgeFactor) + } + + // Exponential delay distribution to prevent synchronized collision storms + private func exponentialRelayDelay() -> TimeInterval { + // Use exponential distribution with mean of (min + max) / 2 + let meanDelay = (minMessageDelay + maxMessageDelay) / 2.0 + let lambda = 1.0 / meanDelay + + // Generate exponential random value + let u = Double.random(in: 0..<1) + let exponentialDelay = -log(1.0 - u) / lambda + + // Clamp to our bounds + return min(maxMessageDelay, max(minMessageDelay, exponentialDelay)) + } + + // Check if message type is high priority (should bypass aggregation) + private func isHighPriorityMessage(type: UInt8) -> Bool { + switch MessageType(rawValue: type) { + case .noiseHandshakeInit, .noiseHandshakeResp, .protocolAck, + .versionHello, .versionAck, .deliveryAck, .systemValidation: + return true + case .message, .announce, .leave, .readReceipt, .deliveryStatusRequest, + .fragmentStart, .fragmentContinue, .fragmentEnd, + .noiseIdentityAnnounce, .noiseEncrypted, .protocolNack, .none: + return false + } + } + + // Calculate exponential backoff for retries + private func calculateExponentialBackoff(retry: Int) -> TimeInterval { + // Start with 1s, double each time: 1s, 2s, 4s, 8s, 16s, max 30s + let baseDelay = 1.0 + let maxDelay = 30.0 + let delay = min(maxDelay, baseDelay * pow(2.0, Double(retry - 1))) + + // Add 10% jitter to prevent synchronized retries + let jitter = delay * 0.1 * (Double.random(in: -1...1)) + return delay + jitter + } + + // MARK: - Collision Avoidance + + // Calculate jitter based on node ID to spread transmissions + private func calculateNodeIDJitter() -> TimeInterval { + // Use hash of peer ID to generate consistent jitter for this node + let hashValue = myPeerID.hash + let normalizedHash = Double(abs(hashValue % 1000)) / 1000.0 // 0.0 to 0.999 + + // Jitter range: 0-20ms based on node ID + let jitterRange: TimeInterval = 0.02 // 20ms max jitter + return normalizedHash * jitterRange + } + + // Add smart delay to avoid collisions + private func smartCollisionAvoidanceDelay(baseDelay: TimeInterval) -> TimeInterval { + // Add node-specific jitter to base delay + let nodeJitter = calculateNodeIDJitter() + + // Add small random component to avoid perfect synchronization + let randomJitter = TimeInterval.random(in: 0...0.005) // 0-5ms additional random + + return baseDelay + nodeJitter + randomJitter + } + + // MARK: - Relay Cancellation + + private func scheduleRelay(_ packet: BitchatPacket, messageID: String, delay: TimeInterval) { + pendingRelaysLock.lock() + defer { pendingRelaysLock.unlock() } + + // Cancel any existing relay for this message + if let existingRelay = pendingRelays[messageID] { + existingRelay.cancel() + } + + // Apply smart collision avoidance to delay + let adjustedDelay = smartCollisionAvoidanceDelay(baseDelay: delay) + + // Create new relay task + let relayTask = DispatchWorkItem { [weak self] in + guard let self = self else { return } + + // Remove from pending when executed + self.pendingRelaysLock.lock() + self.pendingRelays.removeValue(forKey: messageID) + self.pendingRelaysLock.unlock() + + // Mark this message as relayed + self.processedMessagesLock.lock() + if var state = self.processedMessages[messageID] { + state.relayed = true + self.processedMessages[messageID] = state + } + self.processedMessagesLock.unlock() + + // Actually relay the packet + self.broadcastPacket(packet) + } + + // Store the task + pendingRelays[messageID] = relayTask + + // Schedule it with adjusted delay + DispatchQueue.global(qos: .default).asyncAfter(deadline: .now() + adjustedDelay, execute: relayTask) + } + + private func cancelPendingRelay(messageID: String) { + pendingRelaysLock.lock() + defer { pendingRelaysLock.unlock() } + + if let pendingRelay = pendingRelays[messageID] { + pendingRelay.cancel() + pendingRelays.removeValue(forKey: messageID) + // Cancelled pending relay - another node handled it + } + } + + // MARK: - Memory Management + + private func startMemoryCleanupTimer() { + memoryCleanupTimer?.invalidate() + memoryCleanupTimer = Timer.scheduledTimer(withTimeInterval: 60.0, repeats: true) { [weak self] _ in + self?.performMemoryCleanup() + } + } + + private func performMemoryCleanup() { + messageQueue.async(flags: .barrier) { [weak self] in + guard let self = self else { return } + + let startTime = Date() + + // Clean up old fragments + let fragmentTimeout = Date().addingTimeInterval(-self.fragmentTimeout) + var expiredFragments = [String]() + for (fragmentID, metadata) in self.fragmentMetadata { + if metadata.timestamp < fragmentTimeout { + expiredFragments.append(fragmentID) + } + } + for fragmentID in expiredFragments { + self.incomingFragments.removeValue(forKey: fragmentID) + self.fragmentMetadata.removeValue(forKey: fragmentID) + } + + // Limit pending private messages + if self.pendingPrivateMessages.count > self.maxPendingPrivateMessages { + let excess = self.pendingPrivateMessages.count - self.maxPendingPrivateMessages + let keysToRemove = Array(self.pendingPrivateMessages.keys.prefix(excess)) + for key in keysToRemove { + self.pendingPrivateMessages.removeValue(forKey: key) + } + SecureLogger.log("Removed \(excess) oldest pending private message queues", + category: SecureLogger.session, level: .info) + } + + // Limit cached messages sent to peer + if self.cachedMessagesSentToPeer.count > self.maxCachedMessagesSentToPeer { + let excess = self.cachedMessagesSentToPeer.count - self.maxCachedMessagesSentToPeer + let toRemove = self.cachedMessagesSentToPeer.prefix(excess) + self.cachedMessagesSentToPeer.subtract(toRemove) + SecureLogger.log("Removed \(excess) oldest cached message tracking entries", + category: SecureLogger.session, level: .debug) + } + + // Clean up pending relays + self.pendingRelaysLock.lock() + _ = self.pendingRelays.count + self.pendingRelaysLock.unlock() + + // Clean up rate limiters + self.cleanupRateLimiters() + + // Log memory status + _ = Date().timeIntervalSince(startTime) + // Memory cleanup completed + + // Estimate current memory usage and log if high + let estimatedMemory = self.estimateMemoryUsage() + if estimatedMemory > self.maxMemoryUsageBytes { + SecureLogger.log("Warning: Estimated memory usage \(estimatedMemory / 1024 / 1024)MB exceeds limit", + category: SecureLogger.session, level: .warning) + } + } + } + + private func estimateMemoryUsage() -> Int { + // Rough estimates based on typical sizes + let messageSize = 512 // Average message size + let fragmentSize = 512 // Average fragment size + + var totalBytes = 0 + + // Processed messages (string storage + MessageState) + processedMessagesLock.lock() + totalBytes += processedMessages.count * 150 // messageID strings + MessageState struct + processedMessagesLock.unlock() + + // Fragments + for (_, fragments) in incomingFragments { + totalBytes += fragments.count * fragmentSize + } + + // Pending private messages + for (_, messages) in pendingPrivateMessages { + totalBytes += messages.count * messageSize + } + + // Other caches + totalBytes += cachedMessagesSentToPeer.count * 50 // peerID strings + totalBytes += deliveredMessages.count * 100 // messageID strings + totalBytes += recentlySentMessages.count * 100 // messageID strings + + return totalBytes + } + + // MARK: - Rate Limiting + + private func isRateLimited(peerID: String, messageType: UInt8) -> Bool { + rateLimiterLock.lock() + defer { rateLimiterLock.unlock() } + + let now = Date() + let cutoff = now.addingTimeInterval(-rateLimitWindow) + + // Clean old timestamps from total counter + totalMessageTimestamps.removeAll { $0 < cutoff } + + // Check global rate limit with progressive throttling + if totalMessageTimestamps.count >= maxTotalMessagesPerMinute { + // Apply progressive throttling for global limit + let overageRatio = Double(totalMessageTimestamps.count - maxTotalMessagesPerMinute) / Double(maxTotalMessagesPerMinute) + let dropProbability = min(0.95, 0.5 + overageRatio * 0.5) // 50% to 95% drop rate + + if Double.random(in: 0...1) < dropProbability { + SecureLogger.log("Global rate limit throttling: \(totalMessageTimestamps.count) messages, drop prob: \(Int(dropProbability * 100))%", + category: SecureLogger.security, level: .warning) + return true + } + } + + // Determine which rate limiter to use based on message type + let isChatMessage = messageType == MessageType.message.rawValue + + let limiter = isChatMessage ? messageRateLimiter : protocolMessageRateLimiter + let maxPerMinute = isChatMessage ? maxChatMessagesPerPeerPerMinute : maxProtocolMessagesPerPeerPerMinute + + // Clean old timestamps for this peer + if var timestamps = limiter[peerID] { + timestamps.removeAll { $0 < cutoff } + + // Update the appropriate limiter + if isChatMessage { + messageRateLimiter[peerID] = timestamps + } else { + protocolMessageRateLimiter[peerID] = timestamps + } + + // Progressive throttling for per-peer limit + if timestamps.count >= maxPerMinute { + // Calculate how much over the limit we are + let overageRatio = Double(timestamps.count - maxPerMinute) / Double(maxPerMinute) + + // Progressive drop probability: starts at 30% when just over limit, increases to 90% + let dropProbability = min(0.9, 0.3 + overageRatio * 0.6) + + if Double.random(in: 0...1) < dropProbability { + let messageTypeStr = isChatMessage ? "chat" : "protocol" + SecureLogger.log("Peer \(peerID) \(messageTypeStr) rate throttling: \(timestamps.count) msgs/min, drop prob: \(Int(dropProbability * 100))%", + category: SecureLogger.security, level: .info) + return true + } else { + SecureLogger.log("Peer \(peerID) rate throttling: allowing message (\(timestamps.count) msgs/min)", + category: SecureLogger.security, level: .debug) + } + } + } + + return false + } + + private func recordMessage(from peerID: String, messageType: UInt8) { + rateLimiterLock.lock() + defer { rateLimiterLock.unlock() } + + let now = Date() + + // Record in global counter + totalMessageTimestamps.append(now) + + // Determine which rate limiter to use based on message type + let isChatMessage = messageType == MessageType.message.rawValue + + // Record for specific peer in the appropriate limiter + if isChatMessage { + if messageRateLimiter[peerID] == nil { + messageRateLimiter[peerID] = [] + } + messageRateLimiter[peerID]?.append(now) + } else { + if protocolMessageRateLimiter[peerID] == nil { + protocolMessageRateLimiter[peerID] = [] + } + protocolMessageRateLimiter[peerID]?.append(now) + } + } + + private func cleanupRateLimiters() { + rateLimiterLock.lock() + defer { rateLimiterLock.unlock() } + + let cutoff = Date().addingTimeInterval(-rateLimitWindow) + + // Clean global timestamps + totalMessageTimestamps.removeAll { $0 < cutoff } + + // Clean per-peer chat message timestamps + for (peerID, timestamps) in messageRateLimiter { + let filtered = timestamps.filter { $0 >= cutoff } + if filtered.isEmpty { + messageRateLimiter.removeValue(forKey: peerID) + } else { + messageRateLimiter[peerID] = filtered + } + } + + // Clean per-peer protocol message timestamps + for (peerID, timestamps) in protocolMessageRateLimiter { + let filtered = timestamps.filter { $0 >= cutoff } + if filtered.isEmpty { + protocolMessageRateLimiter.removeValue(forKey: peerID) + } else { + protocolMessageRateLimiter[peerID] = filtered + } + } + + SecureLogger.log("Rate limiter cleanup: tracking \(messageRateLimiter.count) chat peers, \(protocolMessageRateLimiter.count) protocol peers, \(totalMessageTimestamps.count) total messages", + category: SecureLogger.session, level: .debug) + } + // MARK: - Cover Traffic private func startCoverTraffic() { @@ -4066,7 +4613,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: encrypted, signature: nil, - ttl: 1 + ttl: 1, + sequenceNumber: self.getNextSequenceNumber() ) self.broadcastPacket(packet) @@ -4186,7 +4734,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: handshakeData, signature: nil, - ttl: 6 // Increased TTL for better delivery on startup + ttl: 6, // Increased TTL for better delivery on startup + sequenceNumber: getNextSequenceNumber() ) // Track packet for ACK @@ -4249,7 +4798,7 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { do { // Process handshake message if let response = try noiseService.processHandshakeMessage(from: peerID, message: message) { - SecureLogger.log("Handshake processing returned response of size \(response.count), sending back to \(peerID)", category: SecureLogger.noise, level: .info) + // Handshake response ready to send // Always send responses as handshake response type let packet = BitchatPacket( @@ -4259,7 +4808,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: response, signature: nil, - ttl: 6 // Increased TTL for better delivery on startup + ttl: 6, // Increased TTL for better delivery on startup + sequenceNumber: getNextSequenceNumber() ) // Track packet for ACK @@ -4273,8 +4823,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { // Check if handshake is complete let sessionEstablished = noiseService.hasEstablishedSession(with: peerID) - let newState = handshakeCoordinator.getHandshakeState(for: peerID) - SecureLogger.log("After processing handshake message - sessionEstablished: \(sessionEstablished), newState: \(newState)", category: SecureLogger.noise, level: .info) + _ = handshakeCoordinator.getHandshakeState(for: peerID) + // Handshake state updated if sessionEstablished { SecureLogger.logSecurityEvent(.handshakeCompleted(peerID: peerID)) @@ -4356,14 +4906,14 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { do { // Decrypt the message - SecureLogger.log("Attempting to decrypt Noise message from \(peerID), encrypted size: \(encryptedData.count)", category: SecureLogger.encryption, level: .debug) + // Attempting to decrypt let decryptedData = try noiseService.decrypt(encryptedData, from: peerID) - SecureLogger.log("Successfully decrypted message from \(peerID), decrypted size: \(decryptedData.count)", category: SecureLogger.encryption, level: .debug) + // Successfully decrypted message // Update last successful message time lastSuccessfulMessageTime[peerID] = Date() - // Send protocol ACK after successful decryption + // Send protocol ACK after successful decryption (only once per encrypted packet) sendProtocolAck(for: originalPacket, to: peerID) // If we can decrypt messages from this peer, they should be in activePeers @@ -4610,8 +5160,7 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { self?.delegate?.didDisconnectFromPeer(peerID) } } else { - SecureLogger.log("Version negotiation successful with \(peerID): agreed on v\(ack.agreedVersion) (server: \(ack.serverVersion), platform: \(ack.platform))", - category: SecureLogger.session, level: .info) + // Version negotiation successful negotiatedVersions[peerID] = ack.agreedVersion versionNegotiationState[peerID] = .ackReceived(version: ack.agreedVersion) @@ -4633,7 +5182,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { type: MessageType.versionHello.rawValue, ttl: 1, // Version negotiation is direct, no relay senderID: myPeerID, - payload: helloData + payload: helloData, + sequenceNumber: getNextSequenceNumber() ) // Mark that we initiated version negotiation @@ -4661,7 +5211,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: ackData, signature: nil, - ttl: 1 // Direct response, no relay + ttl: 1, // Direct response, no relay + sequenceNumber: getNextSequenceNumber() ) broadcastPacket(packet) @@ -4685,23 +5236,34 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { return } - SecureLogger.log("Received protocol ACK from \(peerID) for packet \(ack.originalPacketID), type: \(ack.packetType)", + SecureLogger.log("Received protocol ACK from \(peerID) for packet \(ack.originalPacketID), type: \(ack.packetType), ackID: \(ack.ackID)", category: SecureLogger.session, level: .debug) - // Remove from pending ACKs + // Remove from pending ACKs and mark as acknowledged + // Note: readUUID returns uppercase, but we track with lowercase + let normalizedPacketID = ack.originalPacketID.lowercased() _ = collectionsQueue.sync(flags: .barrier) { - pendingAcks.removeValue(forKey: ack.originalPacketID) + pendingAcks.removeValue(forKey: normalizedPacketID) } + // Track this packet as acknowledged to prevent future retries + acknowledgedPacketsLock.lock() + acknowledgedPackets.insert(ack.originalPacketID) + // Keep only recent acknowledged packets (last 1000) + if acknowledgedPackets.count > 1000 { + // Remove oldest entries (this is approximate since Set doesn't maintain order) + acknowledgedPackets = Set(Array(acknowledgedPackets).suffix(1000)) + } + acknowledgedPacketsLock.unlock() + // Handle specific packet types that need ACK confirmation if let messageType = MessageType(rawValue: ack.packetType) { switch messageType { case .noiseHandshakeInit, .noiseHandshakeResp: - SecureLogger.log("Handshake packet \(ack.originalPacketID) confirmed by \(peerID)", - category: SecureLogger.handshake, level: .info) + SecureLogger.log("Handshake confirmed by \(peerID)", category: SecureLogger.handshake, level: .info) case .noiseEncrypted: - SecureLogger.log("Encrypted message \(ack.originalPacketID) confirmed by \(peerID)", - category: SecureLogger.encryption, level: .debug) + // Encrypted message confirmed + break default: break } @@ -4778,6 +5340,11 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { // Generate packet ID from packet content hash let packetID = generatePacketID(for: packet) + // Debug: log packet details + _ = packet.senderID.prefix(4).hexEncodedString() + _ = packet.recipientID?.prefix(4).hexEncodedString() ?? "nil" + // Send protocol ACK + let ack = ProtocolAck( originalPacketID: packetID, senderID: packet.senderID.hexEncodedString(), @@ -4793,7 +5360,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: ack.toBinaryData(), signature: nil, - ttl: 3 // ACKs don't need to travel far + ttl: 3, // ACKs don't need to travel far + sequenceNumber: getNextSequenceNumber() ) broadcastPacket(ackPacket) @@ -4819,42 +5387,100 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: nack.toBinaryData(), signature: nil, - ttl: 3 // NACKs don't need to travel far + ttl: 3, // NACKs don't need to travel far + sequenceNumber: getNextSequenceNumber() ) broadcastPacket(nackPacket) } - // Generate unique packet ID from packet content + // Generate unique packet ID from immutable packet fields private func generatePacketID(for packet: BitchatPacket) -> String { - // Use hash of packet data for unique ID - if let data = packet.toBinaryData() { - let hash = SHA256.hash(data: data) - return hash.map { String(format: "%02x", $0) }.prefix(16).joined() - } - return UUID().uuidString + // Use only immutable fields for ID generation to ensure consistency + // across network hops (TTL changes, so can't use full packet data) + + // Create a deterministic ID using SHA256 of immutable fields + var data = Data() + data.append(packet.senderID) + data.append(contentsOf: withUnsafeBytes(of: packet.sequenceNumber) { Array($0) }) + data.append(contentsOf: withUnsafeBytes(of: packet.timestamp) { Array($0) }) + data.append(packet.type) + + let hash = SHA256.hash(data: data) + let hashData = Data(hash) + + // Take first 16 bytes for UUID format + let bytes = Array(hashData.prefix(16)) + + // Format as UUID + let p1 = String(format: "%02x%02x%02x%02x", bytes[0], bytes[1], bytes[2], bytes[3]) + let p2 = String(format: "%02x%02x", bytes[4], bytes[5]) + let p3 = String(format: "%02x%02x", bytes[6], bytes[7]) + let p4 = String(format: "%02x%02x", bytes[8], bytes[9]) + let p5 = String(format: "%02x%02x%02x%02x%02x%02x", bytes[10], bytes[11], bytes[12], bytes[13], bytes[14], bytes[15]) + + let result = "\(p1)-\(p2)-\(p3)-\(p4)-\(p5)" + + // Generated packet ID for tracking + + return result } // Track packets that need ACKs private func trackPacketForAck(_ packet: BitchatPacket) { let packetID = generatePacketID(for: packet) + // Debug: log packet details + _ = packet.senderID.prefix(4).hexEncodedString() + _ = packet.recipientID?.prefix(4).hexEncodedString() ?? "nil" + // Track packet for ACK + collectionsQueue.sync(flags: .barrier) { pendingAcks[packetID] = (packet: packet, timestamp: Date(), retries: 0) } - // Schedule timeout check - DispatchQueue.main.asyncAfter(deadline: .now() + ackTimeout) { [weak self] in + // Schedule timeout check with initial delay (using exponential backoff starting at 1s) + let initialDelay = calculateExponentialBackoff(retry: 1) // 1 second initial delay + DispatchQueue.main.asyncAfter(deadline: .now() + initialDelay) { [weak self] in self?.checkAckTimeout(for: packetID) } } // Check for ACK timeout and retry if needed private func checkAckTimeout(for packetID: String) { + // Check if already acknowledged + acknowledgedPacketsLock.lock() + let isAcknowledged = acknowledgedPackets.contains(packetID) + acknowledgedPacketsLock.unlock() + + if isAcknowledged { + // Already acknowledged, remove from pending and don't retry + _ = collectionsQueue.sync(flags: .barrier) { + pendingAcks.removeValue(forKey: packetID) + } + SecureLogger.log("Packet \(packetID) already acknowledged, cancelling retries", + category: SecureLogger.session, level: .debug) + return + } + collectionsQueue.sync(flags: .barrier) { [weak self] in guard let self = self, let pending = self.pendingAcks[packetID] else { return } + // Check if this is a handshake packet and we already have an established session + if pending.packet.type == MessageType.noiseHandshakeInit.rawValue || + pending.packet.type == MessageType.noiseHandshakeResp.rawValue { + // Extract peer ID from packet + let peerID = pending.packet.recipientID?.hexEncodedString() ?? "" + if !peerID.isEmpty && self.noiseService.hasEstablishedSession(with: peerID) { + // We have an established session, don't retry handshake packets + SecureLogger.log("Not retrying handshake packet \(packetID) - session already established with \(peerID)", + category: SecureLogger.handshake, level: .info) + self.pendingAcks.removeValue(forKey: packetID) + return + } + } + if pending.retries < self.maxAckRetries { // Retry sending the packet SecureLogger.log("ACK timeout for packet \(packetID), retrying (attempt \(pending.retries + 1))", @@ -4865,12 +5491,15 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { retries: pending.retries + 1) // Resend the packet + SecureLogger.log("Resending packet seq#\(pending.packet.sequenceNumber) due to ACK timeout (retry \(pending.retries + 1))", + category: SecureLogger.session, level: .debug) DispatchQueue.main.async { self.broadcastPacket(pending.packet) } - // Schedule next timeout check - DispatchQueue.main.asyncAfter(deadline: .now() + self.ackTimeout) { + // Schedule next timeout check with exponential backoff + let backoffDelay = self.calculateExponentialBackoff(retry: pending.retries + 1) + DispatchQueue.main.asyncAfter(deadline: .now() + backoffDelay) { self.checkAckTimeout(for: packetID) } } else { @@ -5032,7 +5661,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: announcementData, signature: nil, - ttl: adaptiveTTL + ttl: adaptiveTTL, + sequenceNumber: getNextSequenceNumber() ) if let targetPeer = specificPeerID { @@ -5155,7 +5785,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: messageData, signature: nil, - ttl: self.adaptiveTTL // Inner packet needs valid TTL for processing after decryption + ttl: self.adaptiveTTL, // Inner packet needs valid TTL for processing after decryption + sequenceNumber: getNextSequenceNumber() ) guard let innerData = innerPacket.toBinaryData() else { return } @@ -5177,7 +5808,8 @@ extension BluetoothMeshService: CBPeripheralManagerDelegate { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: encryptedData, signature: nil, - ttl: adaptiveTTL + ttl: adaptiveTTL, + sequenceNumber: getNextSequenceNumber() ) SecureLogger.log("Broadcasting encrypted private message \(msgID) to \(recipientPeerID)", category: SecureLogger.session, level: .info) diff --git a/bitchatTests/EndToEnd/PrivateChatE2ETests.swift b/bitchatTests/EndToEnd/PrivateChatE2ETests.swift index 171d47b7..bd7b2c93 100644 --- a/bitchatTests/EndToEnd/PrivateChatE2ETests.swift +++ b/bitchatTests/EndToEnd/PrivateChatE2ETests.swift @@ -365,7 +365,8 @@ final class PrivateChatE2ETests: XCTestCase { timestamp: packet.timestamp, payload: encrypted, signature: packet.signature, - ttl: packet.ttl + ttl: packet.ttl, + sequenceNumber: 1 ) self.bob.simulateIncomingPacket(encryptedPacket) } catch { diff --git a/bitchatTests/EndToEnd/PublicChatE2ETests.swift b/bitchatTests/EndToEnd/PublicChatE2ETests.swift index c1994d50..7f8e9887 100644 --- a/bitchatTests/EndToEnd/PublicChatE2ETests.swift +++ b/bitchatTests/EndToEnd/PublicChatE2ETests.swift @@ -122,7 +122,8 @@ final class PublicChatE2ETests: XCTestCase { timestamp: packet.timestamp, payload: relayPayload, signature: packet.signature, - ttl: packet.ttl - 1 + ttl: packet.ttl - 1, + sequenceNumber: 1 ) // Simulate relay to Charlie @@ -450,7 +451,8 @@ final class PublicChatE2ETests: XCTestCase { timestamp: packet.timestamp, payload: relayPayload, signature: packet.signature, - ttl: packet.ttl - 1 + ttl: packet.ttl - 1, + sequenceNumber: 1 ) // Relay to next hops diff --git a/bitchatTests/Integration/IntegrationTests.swift b/bitchatTests/Integration/IntegrationTests.swift index 078c8c95..bcdd04bd 100644 --- a/bitchatTests/Integration/IntegrationTests.swift +++ b/bitchatTests/Integration/IntegrationTests.swift @@ -197,7 +197,8 @@ final class IntegrationTests: XCTestCase { timestamp: packet.timestamp, payload: encrypted, signature: packet.signature, - ttl: packet.ttl + ttl: packet.ttl, + sequenceNumber: 1 ) self.nodes["Bob"]!.simulateIncomingPacket(encPacket) } @@ -698,7 +699,8 @@ final class IntegrationTests: XCTestCase { timestamp: packet.timestamp, payload: encrypted, signature: packet.signature, - ttl: packet.ttl + ttl: packet.ttl, + sequenceNumber: 1 ) self.nodes["Bob"]!.simulateIncomingPacket(encPacket) } @@ -805,7 +807,8 @@ final class IntegrationTests: XCTestCase { timestamp: packet.timestamp, payload: relayPayload, signature: packet.signature, - ttl: packet.ttl - 1 + ttl: packet.ttl - 1, + sequenceNumber: 1 ) for hop in nextHops { diff --git a/bitchatTests/Mocks/MockBluetoothMeshService.swift b/bitchatTests/Mocks/MockBluetoothMeshService.swift index b2ec7bbf..6c4a2071 100644 --- a/bitchatTests/Mocks/MockBluetoothMeshService.swift +++ b/bitchatTests/Mocks/MockBluetoothMeshService.swift @@ -73,7 +73,8 @@ class MockBluetoothMeshService: BluetoothMeshService { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: payload, signature: nil, - ttl: 3 + ttl: 3, + sequenceNumber: 1 ) sentMessages.append((message, packet)) @@ -111,7 +112,8 @@ class MockBluetoothMeshService: BluetoothMeshService { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: payload, signature: nil, - ttl: 3 + ttl: 3, + sequenceNumber: 1 ) sentMessages.append((message, packet)) diff --git a/bitchatTests/TestUtilities/TestHelpers.swift b/bitchatTests/TestUtilities/TestHelpers.swift index 41aa4e2d..be1bc588 100644 --- a/bitchatTests/TestUtilities/TestHelpers.swift +++ b/bitchatTests/TestUtilities/TestHelpers.swift @@ -55,7 +55,8 @@ class TestHelpers { recipientID: String? = nil, payload: Data = "test payload".data(using: .utf8)!, signature: Data? = nil, - ttl: UInt8 = 3 + ttl: UInt8 = 3, + sequenceNumber: UInt32 = 1 ) -> BitchatPacket { return BitchatPacket( type: type, @@ -64,7 +65,8 @@ class TestHelpers { timestamp: UInt64(Date().timeIntervalSince1970 * 1000), payload: payload, signature: signature, - ttl: ttl + ttl: ttl, + sequenceNumber: sequenceNumber ) }