diff --git a/bitchat/Models/RequestSyncPacket.swift b/bitchat/Models/RequestSyncPacket.swift new file mode 100644 index 00000000..b11b8c1a --- /dev/null +++ b/bitchat/Models/RequestSyncPacket.swift @@ -0,0 +1,63 @@ +import Foundation + +// REQUEST_SYNC payload TLV (type, length16, value) +// - 0x01: P (uint8) — Golomb-Rice parameter +// - 0x02: M (uint32, big-endian) — hash range (N * 2^P) +// - 0x03: data (opaque) — GR bitstream bytes (MSB-first) +struct RequestSyncPacket { + let p: Int + let m: UInt32 + let data: Data + + func encode() -> Data { + var out = Data() + func putTLV(_ t: UInt8, _ v: Data) { + out.append(t) + let len = UInt16(v.count) + out.append(UInt8((len >> 8) & 0xFF)) + out.append(UInt8(len & 0xFF)) + out.append(v) + } + // P + putTLV(0x01, Data([UInt8(p & 0xFF)])) + // M (uint32) + var mBE = m.bigEndian + putTLV(0x02, withUnsafeBytes(of: &mBE) { Data($0) }) + // data + putTLV(0x03, data) + return out + } + + static func decode(from data: Data, maxAcceptBytes: Int = 1024) -> RequestSyncPacket? { + var off = 0 + var p: Int? = nil + var m: UInt32? = nil + var payload: Data? = nil + + while off + 3 <= data.count { + let t = Int(data[off]); off += 1 + guard off + 2 <= data.count else { return nil } + let len = (Int(data[off]) << 8) | Int(data[off+1]); off += 2 + guard off + len <= data.count else { return nil } + let v = data.subdata(in: off..<(off+len)); off += len + switch t { + case 0x01: + if v.count == 1 { p = Int(v[0]) } + case 0x02: + if v.count == 4 { + var mm: UInt32 = 0 + for b in v { mm = (mm << 8) | UInt32(b) } + m = mm + } + case 0x03: + if v.count > maxAcceptBytes { return nil } + payload = v + default: + break // forward compatible; ignore unknown TLVs + } + } + + guard let pp = p, let mm = m, let dd = payload, pp >= 1, mm > 0 else { return nil } + return RequestSyncPacket(p: pp, m: mm, data: dd) + } +} diff --git a/bitchat/Protocols/BitchatProtocol.swift b/bitchat/Protocols/BitchatProtocol.swift index 0bc55536..6045cb21 100644 --- a/bitchat/Protocols/BitchatProtocol.swift +++ b/bitchat/Protocols/BitchatProtocol.swift @@ -125,6 +125,7 @@ enum MessageType: UInt8 { case announce = 0x01 // "I'm here" with nickname case message = 0x02 // Public chat message case leave = 0x03 // "I'm leaving" + case requestSync = 0x21 // GCS filter-based sync request (local-only) // Noise encryption case noiseHandshake = 0x10 // Handshake (init or response determined by payload) @@ -138,6 +139,7 @@ enum MessageType: UInt8 { case .announce: return "announce" case .message: return "message" case .leave: return "leave" + case .requestSync: return "requestSync" case .noiseHandshake: return "noiseHandshake" case .noiseEncrypted: return "noiseEncrypted" case .fragment: return "fragment" diff --git a/bitchat/Services/BLEService.swift b/bitchat/Services/BLEService.swift index 509dbd6f..7f37a841 100644 --- a/bitchat/Services/BLEService.swift +++ b/bitchat/Services/BLEService.swift @@ -138,6 +138,9 @@ final class BLEService: NSObject { private var pendingDirectedRelays: [String: [String: (packet: BitchatPacket, enqueuedAt: Date)]] = [:] // Debounce for 'reconnected' logs private var lastReconnectLogAt: [String: Date] = [:] + + // MARK: - Gossip Sync + private var gossipSyncManager: GossipSyncManager? // MARK: - Maintenance Timer @@ -398,9 +401,15 @@ final class BLEService: NSObject { } timer.resume() maintenanceTimer = timer - + // Publish initial empty state requestPeerDataPublish() + + // Initialize gossip sync manager + let sync = GossipSyncManager(myPeerID: myPeerID) + sync.delegate = self + sync.start() + self.gossipSyncManager = sync } func setNickname(_ nickname: String) { @@ -775,6 +784,8 @@ final class BLEService: NSObject { self.messageDeduplicator.markProcessed(dedupID) // Call synchronously since we're already on background queue self.broadcastPacket(signedPacket) + // Track our own broadcast for sync + self.gossipSyncManager?.onPublicPacketSeen(signedPacket) } } } @@ -1099,10 +1110,13 @@ final class BLEService: NSObject { } // For broadcast (no directed peer) and non-fragment, choose a subset deterministically - // Special-case announces: do NOT subset to maximize reach for presence + // Special-case control/presence messages: do NOT subset to maximize immediate coverage var selectedPeripheralIDs = Set(allowedPeripheralIDs) var selectedCentralIDs = Set(allowedCentralIDs) - if directedOnlyPeer == nil && packet.type != MessageType.fragment.rawValue && packet.type != MessageType.announce.rawValue { + if directedOnlyPeer == nil + && packet.type != MessageType.fragment.rawValue + && packet.type != MessageType.announce.rawValue + && packet.type != MessageType.requestSync.rawValue { let kp = subsetSizeForFanout(allowedPeripheralIDs.count) let kc = subsetSizeForFanout(allowedCentralIDs.count) selectedPeripheralIDs = selectDeterministicSubset(ids: allowedPeripheralIDs, k: kp, seed: messageID) @@ -1133,6 +1147,12 @@ final class BLEService: NSObject { } } + // Directed send helper (unicast to a specific peerID) without altering packet contents + private func sendPacketDirected(_ packet: BitchatPacket, to peerID: String) { + 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) { let msgID = makeMessageID(for: packet) @@ -1339,7 +1359,10 @@ final class BLEService: NSObject { // Efficient deduplication // Important: do not dedup fragment packets globally (each piece must pass) - if packet.type != MessageType.fragment.rawValue && messageDeduplicator.isDuplicate(messageID) { + // Special case: allow our own packets recovered via sync (TTL==0) to pass + // through even if we've marked them as seen at send time. + let allowSelfSyncReplay = (packet.ttl == 0) && (senderID == myPeerID) + if packet.type != MessageType.fragment.rawValue && !allowSelfSyncReplay && messageDeduplicator.isDuplicate(messageID) { // Announce packets (type 1) are sent every 10 seconds for peer discovery // It's normal to see these as duplicates - don't log them to reduce noise if packet.type != MessageType.announce.rawValue { @@ -1383,6 +1406,9 @@ final class BLEService: NSObject { case .message: handleMessage(packet, from: senderID) + case .requestSync: + handleRequestSync(packet, from: senderID) + case .noiseHandshake: handleNoiseHandshake(packet, from: senderID) @@ -1583,12 +1609,17 @@ final class BLEService: NSObject { // 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) + // Schedule initial unicast sync to this peer + self.gossipSyncManager?.scheduleInitialSyncToPeer(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) @@ -1610,17 +1641,33 @@ final class BLEService: NSObject { } } } + + // 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(fromPeerID: peerID, request: req) + } // Mention parsing moved to ChatViewModel private func handleMessage(_ packet: BitchatPacket, from peerID: String) { - // Ignore self-origin public messages that may be seen again via relay - if peerID == myPeerID { return } + // 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 = "" - if let info = peers[peerID], info.isVerifiedNickname { + // 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 = peers[peerID], info.isVerifiedNickname { // Known verified peer path accepted = true senderNickname = info.nickname @@ -1648,6 +1695,22 @@ final class BLEService: NSObject { } } } + // 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 { @@ -1780,6 +1843,8 @@ final class BLEService: NSObject { // Remove the peer when they leave peers.removeValue(forKey: peerID) } + // Remove any stored announcement for sync purposes + gossipSyncManager?.removeAnnouncementForPeer(peerID) // Send on main thread notifyUI { [weak self] in guard let self = self else { return } @@ -1861,6 +1926,8 @@ final class BLEService: NSObject { self?.broadcastPacket(signedPacket) } } + // Ensure our own announce is included in sync state + gossipSyncManager?.onPublicPacketSeen(signedPacket) } func sendDeliveryAck(for messageID: String, to peerID: String) { @@ -2088,6 +2155,8 @@ final class BLEService: NSObject { if !peer.isConnected { if age > retention { SecureLogger.debug("🗑️ Removing stale peer after reachability window: \(peerID) (\(peer.nickname))", category: .session) + // Also remove any stored announcement from sync candidates + gossipSyncManager?.removeAnnouncementForPeer(peerID) peers.removeValue(forKey: peerID) removedOfflineCount += 1 } @@ -2249,6 +2318,29 @@ final class BLEService: NSObject { } } +// MARK: - GossipSyncManager Delegate +extension BLEService: GossipSyncManager.Delegate { + func sendPacket(_ packet: BitchatPacket) { + if DispatchQueue.getSpecific(key: messageQueueKey) != nil { + broadcastPacket(packet) + } else { + messageQueue.async { [weak self] in self?.broadcastPacket(packet) } + } + } + + func sendPacket(to peerID: String, packet: BitchatPacket) { + if DispatchQueue.getSpecific(key: messageQueueKey) != nil { + sendPacketDirected(packet, to: peerID) + } else { + messageQueue.async { [weak self] in self?.sendPacketDirected(packet, to: peerID) } + } + } + + func signPacketForBroadcast(_ packet: BitchatPacket) -> BitchatPacket { + return noiseService.signPacket(packet) ?? packet + } +} + // MARK: - CBCentralManagerDelegate extension BLEService: CBCentralManagerDelegate { diff --git a/bitchat/Sync/GCSFilter.swift b/bitchat/Sync/GCSFilter.swift new file mode 100644 index 00000000..d05856fc --- /dev/null +++ b/bitchat/Sync/GCSFilter.swift @@ -0,0 +1,183 @@ +import Foundation +import CryptoKit + +// Golomb-Coded Set (GCS) filter utilities for sync. +// Hashing: +// - Packet ID is 16 bytes (see PacketIdUtil). For GCS mapping, use h64 = first 8 bytes of SHA-256 over the 16-byte ID. +// - Map to [0, M) via (h64 % M). +// Encoding (v1): +// - Sort mapped values ascending; encode deltas (first is v0, then vi - v{i-1}) as positive integers x >= 1. +// - Golomb-Rice with parameter P: q = (x - 1) >> P encoded as unary (q ones then a zero), then write P-bit remainder r = (x - 1) & ((1<
Int {
+ let f = max(0.000001, min(0.25, targetFpr))
+ // ceil(log2(1/f))
+ let p = Int(ceil(log2(1.0 / f)))
+ return max(1, p)
+ }
+
+ // Estimate max elements that fit in size bytes: bits per element ~= P + 2 (approx)
+ static func estimateMaxElements(sizeBytes: Int, p: Int) -> Int {
+ let bits = max(8, sizeBytes * 8)
+ let per = max(3, p + 2)
+ return max(1, bits / per)
+ }
+
+ static func buildFilter(ids: [Data], maxBytes: Int, targetFpr: Double) -> Params {
+ let p = deriveP(targetFpr: targetFpr)
+ let cap = estimateMaxElements(sizeBytes: maxBytes, p: p)
+ let n = min(ids.count, cap)
+ let selected = Array(ids.prefix(n))
+ // Map to [0, M)
+ let mInit = UInt32(n << p)
+ var mapped = selected.map { id16 -> UInt64 in
+ let h = h64(id16)
+ return UInt64(h % UInt64(max(1, mInit)))
+ }.sorted()
+ var encoded = encode(sorted: mapped, p: p)
+ var trimmedN = n
+ // Trim if over budget
+ while encoded.count > maxBytes && trimmedN > 0 {
+ trimmedN = (trimmedN * 9) / 10 // drop ~10%
+ mapped = Array(mapped.prefix(trimmedN))
+ encoded = encode(sorted: mapped, p: p)
+ }
+ let finalM = UInt32(max(1, trimmedN << p))
+ return Params(p: p, m: finalM, data: encoded)
+ }
+
+ static func decodeToSortedSet(p: Int, m: UInt32, data: Data) -> [UInt64] {
+ var values: [UInt64] = []
+ let reader = BitReader(data)
+ var acc: UInt64 = 0
+ while true {
+ guard let q = reader.readUnary() else { break }
+ guard let r = reader.readBits(count: p) else { break }
+ let x = (UInt64(q) << UInt64(p)) + UInt64(r) + 1
+ acc &+= x
+ if acc >= UInt64(m) { break }
+ values.append(acc)
+ }
+ return values
+ }
+
+ static func contains(sortedValues: [UInt64], candidate: UInt64) -> Bool {
+ var lo = 0
+ var hi = sortedValues.count - 1
+ while lo <= hi {
+ let mid = (lo + hi) >> 1
+ let v = sortedValues[mid]
+ if v == candidate { return true }
+ if v < candidate { lo = mid + 1 } else { hi = mid - 1 }
+ }
+ return false
+ }
+
+ private static func h64(_ id16: Data) -> UInt64 {
+ var hasher = SHA256()
+ hasher.update(data: id16)
+ let d = hasher.finalize()
+ let db = Data(d)
+ var x: UInt64 = 0
+ let take = min(8, db.count)
+ for i in 0..