mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-25 05:25:19 +00:00
wip
This commit is contained in:
@@ -0,0 +1,63 @@
|
|||||||
|
import Foundation
|
||||||
|
|
||||||
|
// REQUEST_SYNC payload TLV (type, length16, value)
|
||||||
|
// - 0x01: mBytes (uint16)
|
||||||
|
// - 0x02: k (uint8)
|
||||||
|
// - 0x03: bloom filter bits (opaque byte array of length mBytes)
|
||||||
|
struct RequestSyncPacket {
|
||||||
|
let mBytes: Int
|
||||||
|
let k: Int
|
||||||
|
let bits: 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)
|
||||||
|
}
|
||||||
|
// mBytes
|
||||||
|
var mb = UInt16(mBytes)
|
||||||
|
let mbData = withUnsafeBytes(of: &mb.bigEndian) { Data($0) }
|
||||||
|
putTLV(0x01, mbData)
|
||||||
|
// k
|
||||||
|
putTLV(0x02, Data([UInt8(k & 0xFF)]))
|
||||||
|
// bits
|
||||||
|
putTLV(0x03, bits)
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
static func decode(from data: Data) -> RequestSyncPacket? {
|
||||||
|
var off = 0
|
||||||
|
var mBytes: Int? = nil
|
||||||
|
var k: Int? = nil
|
||||||
|
var bits: 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 == 2 {
|
||||||
|
let mb = (Int(v[0]) << 8) | Int(v[1])
|
||||||
|
mBytes = mb
|
||||||
|
}
|
||||||
|
case 0x02:
|
||||||
|
if v.count == 1 { k = Int(v[0]) }
|
||||||
|
case 0x03:
|
||||||
|
bits = v
|
||||||
|
default:
|
||||||
|
break // forward compatible; ignore unknown TLVs
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
guard let mb = mBytes, let kk = k, let bb = bits, mb == bb.count else { return nil }
|
||||||
|
return RequestSyncPacket(mBytes: mb, k: kk, bits: bb)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@@ -125,6 +125,7 @@ enum MessageType: UInt8 {
|
|||||||
case announce = 0x01 // "I'm here" with nickname
|
case announce = 0x01 // "I'm here" with nickname
|
||||||
case message = 0x02 // Public chat message
|
case message = 0x02 // Public chat message
|
||||||
case leave = 0x03 // "I'm leaving"
|
case leave = 0x03 // "I'm leaving"
|
||||||
|
case requestSync = 0x21 // Bloom filter-based sync request (local-only)
|
||||||
|
|
||||||
// Noise encryption
|
// Noise encryption
|
||||||
case noiseHandshake = 0x10 // Handshake (init or response determined by payload)
|
case noiseHandshake = 0x10 // Handshake (init or response determined by payload)
|
||||||
@@ -138,6 +139,7 @@ enum MessageType: UInt8 {
|
|||||||
case .announce: return "announce"
|
case .announce: return "announce"
|
||||||
case .message: return "message"
|
case .message: return "message"
|
||||||
case .leave: return "leave"
|
case .leave: return "leave"
|
||||||
|
case .requestSync: return "requestSync"
|
||||||
case .noiseHandshake: return "noiseHandshake"
|
case .noiseHandshake: return "noiseHandshake"
|
||||||
case .noiseEncrypted: return "noiseEncrypted"
|
case .noiseEncrypted: return "noiseEncrypted"
|
||||||
case .fragment: return "fragment"
|
case .fragment: return "fragment"
|
||||||
|
|||||||
@@ -137,6 +137,9 @@ final class BLEService: NSObject {
|
|||||||
private var pendingDirectedRelays: [String: [String: (packet: BitchatPacket, enqueuedAt: Date)]] = [:]
|
private var pendingDirectedRelays: [String: [String: (packet: BitchatPacket, enqueuedAt: Date)]] = [:]
|
||||||
// Debounce for 'reconnected' logs
|
// Debounce for 'reconnected' logs
|
||||||
private var lastReconnectLogAt: [String: Date] = [:]
|
private var lastReconnectLogAt: [String: Date] = [:]
|
||||||
|
|
||||||
|
// MARK: - Gossip Sync
|
||||||
|
private var gossipSyncManager: GossipSyncManager?
|
||||||
|
|
||||||
// MARK: - Maintenance Timer
|
// MARK: - Maintenance Timer
|
||||||
|
|
||||||
@@ -395,9 +398,15 @@ final class BLEService: NSObject {
|
|||||||
}
|
}
|
||||||
timer.resume()
|
timer.resume()
|
||||||
maintenanceTimer = timer
|
maintenanceTimer = timer
|
||||||
|
|
||||||
// Publish initial empty state
|
// Publish initial empty state
|
||||||
requestPeerDataPublish()
|
requestPeerDataPublish()
|
||||||
|
|
||||||
|
// Initialize gossip sync manager
|
||||||
|
let sync = GossipSyncManager(myPeerID: myPeerID)
|
||||||
|
sync.delegate = self
|
||||||
|
sync.start()
|
||||||
|
self.gossipSyncManager = sync
|
||||||
}
|
}
|
||||||
|
|
||||||
func setNickname(_ nickname: String) {
|
func setNickname(_ nickname: String) {
|
||||||
@@ -772,6 +781,8 @@ final class BLEService: NSObject {
|
|||||||
self.messageDeduplicator.markProcessed(dedupID)
|
self.messageDeduplicator.markProcessed(dedupID)
|
||||||
// Call synchronously since we're already on background queue
|
// Call synchronously since we're already on background queue
|
||||||
self.broadcastPacket(signedPacket)
|
self.broadcastPacket(signedPacket)
|
||||||
|
// Track our own broadcast for sync
|
||||||
|
self.gossipSyncManager?.onPublicPacketSeen(signedPacket)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1130,6 +1141,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
|
// MARK: - Directed store-and-forward
|
||||||
private func spoolDirectedPacket(_ packet: BitchatPacket, recipientPeerID: String) {
|
private func spoolDirectedPacket(_ packet: BitchatPacket, recipientPeerID: String) {
|
||||||
let msgID = makeMessageID(for: packet)
|
let msgID = makeMessageID(for: packet)
|
||||||
@@ -1380,6 +1397,9 @@ final class BLEService: NSObject {
|
|||||||
case .message:
|
case .message:
|
||||||
handleMessage(packet, from: senderID)
|
handleMessage(packet, from: senderID)
|
||||||
|
|
||||||
|
case .requestSync:
|
||||||
|
handleRequestSync(packet, from: senderID)
|
||||||
|
|
||||||
case .noiseHandshake:
|
case .noiseHandshake:
|
||||||
handleNoiseHandshake(packet, from: senderID)
|
handleNoiseHandshake(packet, from: senderID)
|
||||||
|
|
||||||
@@ -1580,12 +1600,17 @@ final class BLEService: NSObject {
|
|||||||
// Only notify of connection for new or reconnected peers when it is a direct announce
|
// Only notify of connection for new or reconnected peers when it is a direct announce
|
||||||
if (packet.ttl == self.messageTTL) && (isNewPeer || isReconnectedPeer) {
|
if (packet.ttl == self.messageTTL) && (isNewPeer || isReconnectedPeer) {
|
||||||
self.delegate?.didConnectToPeer(peerID)
|
self.delegate?.didConnectToPeer(peerID)
|
||||||
|
// Schedule initial unicast sync to this peer
|
||||||
|
self.gossipSyncManager?.scheduleInitialSyncToPeer(peerID, delaySeconds: 5.0)
|
||||||
}
|
}
|
||||||
|
|
||||||
self.requestPeerDataPublish()
|
self.requestPeerDataPublish()
|
||||||
self.delegate?.didUpdatePeerList(currentPeerIDs)
|
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)
|
// Send announce back for bidirectional discovery (only once per peer)
|
||||||
let announceBackID = "announce-back-\(peerID)"
|
let announceBackID = "announce-back-\(peerID)"
|
||||||
let shouldSendBack = !messageDeduplicator.contains(announceBackID)
|
let shouldSendBack = !messageDeduplicator.contains(announceBackID)
|
||||||
@@ -1607,6 +1632,15 @@ 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
|
// Mention parsing moved to ChatViewModel
|
||||||
|
|
||||||
@@ -1647,6 +1681,11 @@ final class BLEService: NSObject {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Track broadcast messages (recipientID == nil indicates broadcast)
|
||||||
|
if packet.recipientID == nil && packet.type == MessageType.message.rawValue {
|
||||||
|
gossipSyncManager?.onPublicPacketSeen(packet)
|
||||||
|
}
|
||||||
|
|
||||||
guard accepted else {
|
guard accepted else {
|
||||||
SecureLogger.warning("🚫 Dropping public message from unverified or unknown peer \(peerID.prefix(8))…", category: .security)
|
SecureLogger.warning("🚫 Dropping public message from unverified or unknown peer \(peerID.prefix(8))…", category: .security)
|
||||||
return
|
return
|
||||||
@@ -1858,6 +1897,8 @@ final class BLEService: NSObject {
|
|||||||
self?.broadcastPacket(signedPacket)
|
self?.broadcastPacket(signedPacket)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
// Ensure our own announce is included in sync state
|
||||||
|
gossipSyncManager?.onPublicPacketSeen(signedPacket)
|
||||||
}
|
}
|
||||||
|
|
||||||
func sendDeliveryAck(for messageID: String, to peerID: String) {
|
func sendDeliveryAck(for messageID: String, to peerID: String) {
|
||||||
@@ -2246,6 +2287,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
|
// MARK: - CBCentralManagerDelegate
|
||||||
|
|
||||||
extension BLEService: CBCentralManagerDelegate {
|
extension BLEService: CBCentralManagerDelegate {
|
||||||
|
|||||||
@@ -0,0 +1,168 @@
|
|||||||
|
import Foundation
|
||||||
|
|
||||||
|
// Gossip-based sync manager using rotating Bloom filters
|
||||||
|
final class GossipSyncManager {
|
||||||
|
protocol Delegate: AnyObject {
|
||||||
|
func sendPacket(_ packet: BitchatPacket)
|
||||||
|
func sendPacket(to peerID: String, packet: BitchatPacket)
|
||||||
|
func signPacketForBroadcast(_ packet: BitchatPacket) -> BitchatPacket
|
||||||
|
}
|
||||||
|
|
||||||
|
struct Config {
|
||||||
|
var seenCapacity: Int = 100 // recent broadcast messages kept
|
||||||
|
var bloomMaxBytes: Int = 256 // up to 256 bytes
|
||||||
|
var bloomTargetFpr: Double = 0.01 // 1%
|
||||||
|
}
|
||||||
|
|
||||||
|
private let myPeerID: String
|
||||||
|
private let config: Config
|
||||||
|
weak var delegate: Delegate?
|
||||||
|
|
||||||
|
// Bloom filter
|
||||||
|
private let bloom: SeenPacketsBloomFilter
|
||||||
|
|
||||||
|
// Storage: broadcast messages (ordered by insert), and latest announce per sender
|
||||||
|
private var messages: [String: BitchatPacket] = [:] // idHex -> packet
|
||||||
|
private var messageOrder: [String] = []
|
||||||
|
private var latestAnnouncementByPeer: [String: (id: String, packet: BitchatPacket)] = [:]
|
||||||
|
|
||||||
|
// Timer
|
||||||
|
private var periodicTimer: DispatchSourceTimer?
|
||||||
|
private let queue = DispatchQueue(label: "mesh.sync", qos: .utility)
|
||||||
|
|
||||||
|
init(myPeerID: String, config: Config = Config()) {
|
||||||
|
self.myPeerID = myPeerID
|
||||||
|
self.config = config
|
||||||
|
self.bloom = SeenPacketsBloomFilter(maxBytes: config.bloomMaxBytes, targetFpr: config.bloomTargetFpr)
|
||||||
|
}
|
||||||
|
|
||||||
|
func start() {
|
||||||
|
stop()
|
||||||
|
let timer = DispatchSource.makeTimerSource(queue: queue)
|
||||||
|
timer.schedule(deadline: .now() + 30.0, repeating: 30.0, leeway: .seconds(1))
|
||||||
|
timer.setEventHandler { [weak self] in self?.sendRequestSync() }
|
||||||
|
timer.resume()
|
||||||
|
periodicTimer = timer
|
||||||
|
}
|
||||||
|
|
||||||
|
func stop() {
|
||||||
|
periodicTimer?.cancel(); periodicTimer = nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func scheduleInitialSyncToPeer(_ peerID: String, delaySeconds: TimeInterval = 5.0) {
|
||||||
|
queue.asyncAfter(deadline: .now() + delaySeconds) { [weak self] in
|
||||||
|
self?.sendRequestSync(to: peerID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func onPublicPacketSeen(_ packet: BitchatPacket) {
|
||||||
|
let mt = MessageType(rawValue: packet.type)
|
||||||
|
let isBroadcastMessage = (mt == .message && packet.recipientID == nil)
|
||||||
|
let isAnnounce = (mt == .announce)
|
||||||
|
guard isBroadcastMessage || isAnnounce else { return }
|
||||||
|
|
||||||
|
let idBytes = PacketIdUtil.computeId(packet)
|
||||||
|
bloom.add(idBytes)
|
||||||
|
let idHex = idBytes.hexEncodedString()
|
||||||
|
|
||||||
|
if isBroadcastMessage {
|
||||||
|
if messages[idHex] == nil {
|
||||||
|
messages[idHex] = packet
|
||||||
|
messageOrder.append(idHex)
|
||||||
|
// Enforce capacity
|
||||||
|
let cap = max(1, config.seenCapacity)
|
||||||
|
while messageOrder.count > cap {
|
||||||
|
let victim = messageOrder.removeFirst()
|
||||||
|
messages.removeValue(forKey: victim)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} else if isAnnounce {
|
||||||
|
let sender = packet.senderID.hexEncodedString()
|
||||||
|
latestAnnouncementByPeer[sender] = (id: idHex, packet: packet)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private func sendRequestSync() {
|
||||||
|
let snap = bloom.snapshotActive()
|
||||||
|
let payload = RequestSyncPacket(mBytes: snap.mBytes, k: snap.k, bits: snap.bits).encode()
|
||||||
|
let pkt = BitchatPacket(
|
||||||
|
type: MessageType.requestSync.rawValue,
|
||||||
|
senderID: Data(hexString: myPeerID) ?? Data(),
|
||||||
|
recipientID: nil, // broadcast
|
||||||
|
timestamp: UInt64(Date().timeIntervalSince1970 * 1000),
|
||||||
|
payload: payload,
|
||||||
|
signature: nil,
|
||||||
|
ttl: 0 // local-only
|
||||||
|
)
|
||||||
|
let signed = delegate?.signPacketForBroadcast(pkt) ?? pkt
|
||||||
|
delegate?.sendPacket(signed)
|
||||||
|
}
|
||||||
|
|
||||||
|
private func sendRequestSync(to peerID: String) {
|
||||||
|
let snap = bloom.snapshotActive()
|
||||||
|
let payload = RequestSyncPacket(mBytes: snap.mBytes, k: snap.k, bits: snap.bits).encode()
|
||||||
|
var recipient = Data()
|
||||||
|
var temp = peerID
|
||||||
|
while temp.count >= 2 && recipient.count < 8 {
|
||||||
|
let hexByte = String(temp.prefix(2))
|
||||||
|
if let b = UInt8(hexByte, radix: 16) { recipient.append(b) }
|
||||||
|
temp = String(temp.dropFirst(2))
|
||||||
|
}
|
||||||
|
let pkt = BitchatPacket(
|
||||||
|
type: MessageType.requestSync.rawValue,
|
||||||
|
senderID: Data(hexString: myPeerID) ?? Data(),
|
||||||
|
recipientID: recipient,
|
||||||
|
timestamp: UInt64(Date().timeIntervalSince1970 * 1000),
|
||||||
|
payload: payload,
|
||||||
|
signature: nil,
|
||||||
|
ttl: 0 // local-only
|
||||||
|
)
|
||||||
|
let signed = delegate?.signPacketForBroadcast(pkt) ?? pkt
|
||||||
|
delegate?.sendPacket(to: peerID, packet: signed)
|
||||||
|
}
|
||||||
|
|
||||||
|
func handleRequestSync(fromPeerID: String, request: RequestSyncPacket) {
|
||||||
|
// Build membership checker from provided parameters
|
||||||
|
let mBits = request.mBytes * 8
|
||||||
|
let k = request.k
|
||||||
|
func mightContain(_ id: Data) -> Bool {
|
||||||
|
// Same hashing as local bloom; compute indices, check MSB-first bits in request.bits
|
||||||
|
var h1: UInt64 = 1469598103934665603
|
||||||
|
var h2: UInt64 = 0x27d4eb2f165667c5
|
||||||
|
for b in id { h1 = (h1 ^ UInt64(b)) &* 1099511628211; h2 = (h2 ^ UInt64(b)) &* 0x100000001B3 }
|
||||||
|
for i in 0..<k {
|
||||||
|
let combined = h1 &+ (UInt64(i) &* h2)
|
||||||
|
let idx = Int((combined & 0x7fff_ffff_ffff_ffff) % UInt64(mBits))
|
||||||
|
let byteIndex = idx / 8
|
||||||
|
let bitIndex = idx % 8
|
||||||
|
let byte = request.bits[byteIndex]
|
||||||
|
let bit = ((Int(byte) >> (7 - bitIndex)) & 1) == 1
|
||||||
|
if !bit { return false }
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// 1) Announcements: send latest per peer if requester lacks them
|
||||||
|
for (_, pair) in latestAnnouncementByPeer {
|
||||||
|
let (idHex, pkt) = pair
|
||||||
|
let idBytes = Data(hexString: idHex) ?? Data()
|
||||||
|
if !mightContain(idBytes) {
|
||||||
|
var toSend = pkt
|
||||||
|
toSend.ttl = 0
|
||||||
|
delegate?.sendPacket(to: fromPeerID, packet: toSend)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2) Broadcast messages: send all missing
|
||||||
|
let toSendMsgs = messageOrder.compactMap { messages[$0] }
|
||||||
|
for pkt in toSendMsgs {
|
||||||
|
let idBytes = PacketIdUtil.computeId(pkt)
|
||||||
|
if !mightContain(idBytes) {
|
||||||
|
var toSend = pkt
|
||||||
|
toSend.ttl = 0
|
||||||
|
delegate?.sendPacket(to: fromPeerID, packet: toSend)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@@ -0,0 +1,22 @@
|
|||||||
|
import Foundation
|
||||||
|
import CryptoKit
|
||||||
|
|
||||||
|
// Deterministic packet ID used for sync Bloom membership
|
||||||
|
// ID = first 16 bytes of SHA-256 over: [type | senderID | timestamp | payload]
|
||||||
|
enum PacketIdUtil {
|
||||||
|
static func computeId(_ packet: BitchatPacket) -> Data {
|
||||||
|
var hasher = SHA256()
|
||||||
|
hasher.update(data: Data([packet.type]))
|
||||||
|
hasher.update(data: packet.senderID)
|
||||||
|
var tsBE = packet.timestamp.bigEndian
|
||||||
|
withUnsafeBytes(of: &tsBE) { raw in hasher.update(data: Data(raw)) }
|
||||||
|
hasher.update(data: packet.payload)
|
||||||
|
let digest = hasher.finalize()
|
||||||
|
return Data(digest.prefix(16))
|
||||||
|
}
|
||||||
|
|
||||||
|
static func computeIdHex(_ packet: BitchatPacket) -> String {
|
||||||
|
return computeId(packet).hexEncodedString()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@@ -0,0 +1,116 @@
|
|||||||
|
import Foundation
|
||||||
|
|
||||||
|
// Rotating Bloom filter for recently seen packet IDs (16-byte IDs)
|
||||||
|
final class SeenPacketsBloomFilter {
|
||||||
|
struct Snapshot { let mBytes: Int; let k: Int; let bits: Data }
|
||||||
|
|
||||||
|
private struct Filter { var mBits: Int; var k: Int; var bits: [UInt8]; var count: Int }
|
||||||
|
|
||||||
|
private let maxBytes: Int
|
||||||
|
private let targetFpr: Double
|
||||||
|
private let seed: UInt64 = 0x27d4eb2f165667c5
|
||||||
|
|
||||||
|
private let mBits: Int
|
||||||
|
private let kOptimal: Int
|
||||||
|
private let capacityOptimal: Int
|
||||||
|
|
||||||
|
private var active: Filter
|
||||||
|
private var standby: Filter
|
||||||
|
private var usingStandby: Bool = false
|
||||||
|
private let lock = NSLock()
|
||||||
|
|
||||||
|
init(maxBytes: Int = 256, targetFpr: Double = 0.01) {
|
||||||
|
self.maxBytes = max(1, maxBytes)
|
||||||
|
self.targetFpr = targetFpr
|
||||||
|
self.mBits = max(8, self.maxBytes * 8)
|
||||||
|
let (k, n) = SeenPacketsBloomFilter.deriveParams(mBits: self.mBits, fpr: targetFpr)
|
||||||
|
self.kOptimal = max(1, k)
|
||||||
|
self.capacityOptimal = max(1, n)
|
||||||
|
self.active = Filter(mBits: self.mBits, k: self.kOptimal, bits: [UInt8](repeating: 0, count: self.maxBytes), count: 0)
|
||||||
|
self.standby = Filter(mBits: self.mBits, k: self.kOptimal, bits: [UInt8](repeating: 0, count: self.maxBytes), count: 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
private static func deriveParams(mBits: Int, fpr: Double) -> (Int, Int) {
|
||||||
|
// n ≈ -(m (ln 2)^2) / ln p ; k ≈ (m/n) ln 2
|
||||||
|
let ln2 = log(2.0)
|
||||||
|
let n = max(1, Int(Double(-mBits) * ln2 * ln2 / log(fpr)))
|
||||||
|
let k = max(1, Int(ceil((Double(mBits) / Double(n)) * ln2)))
|
||||||
|
return (k, n)
|
||||||
|
}
|
||||||
|
|
||||||
|
private func indicesFor(id: Data, mBits: Int, k: Int) -> [Int] {
|
||||||
|
var h1: UInt64 = 1469598103934665603 // FNV-1a 64-bit offset
|
||||||
|
var h2: UInt64 = seed
|
||||||
|
for b in id { // treat as unsigned bytes
|
||||||
|
h1 = (h1 ^ UInt64(b)) &* 1099511628211
|
||||||
|
h2 = (h2 ^ UInt64(b)) &* 0x100000001B3
|
||||||
|
}
|
||||||
|
var result = [Int]()
|
||||||
|
result.reserveCapacity(k)
|
||||||
|
for i in 0..<k {
|
||||||
|
let combined = h1 &+ (UInt64(i) &* h2)
|
||||||
|
let idx = Int((combined & 0x7fff_ffff_ffff_ffff) % UInt64(mBits))
|
||||||
|
result.append(idx)
|
||||||
|
}
|
||||||
|
return result
|
||||||
|
}
|
||||||
|
|
||||||
|
func add(_ id: Data) {
|
||||||
|
lock.lock(); defer { lock.unlock() }
|
||||||
|
let startStandbyAt = capacityOptimal / 2
|
||||||
|
if !usingStandby && active.count >= startStandbyAt {
|
||||||
|
standby = Filter(mBits: mBits, k: kOptimal, bits: [UInt8](repeating: 0, count: maxBytes), count: 0)
|
||||||
|
usingStandby = true
|
||||||
|
}
|
||||||
|
insert(into: &active, id: id)
|
||||||
|
if usingStandby { insert(into: &standby, id: id) }
|
||||||
|
if active.count >= capacityOptimal {
|
||||||
|
active = standby
|
||||||
|
standby = Filter(mBits: mBits, k: kOptimal, bits: [UInt8](repeating: 0, count: maxBytes), count: 0)
|
||||||
|
usingStandby = false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private func insert(into filter: inout Filter, id: Data) {
|
||||||
|
let idxs = indicesFor(id: id, mBits: filter.mBits, k: filter.k)
|
||||||
|
for i in idxs {
|
||||||
|
let byteIndex = i / 8
|
||||||
|
let bitIndex = i % 8
|
||||||
|
filter.bits[byteIndex] = UInt8(Int(filter.bits[byteIndex]) | (1 << (7 - bitIndex)))
|
||||||
|
}
|
||||||
|
filter.count &+= 1
|
||||||
|
}
|
||||||
|
|
||||||
|
func mightContain(_ id: Data) -> Bool {
|
||||||
|
lock.lock(); defer { lock.unlock() }
|
||||||
|
let a = active
|
||||||
|
let idx = indicesFor(id: id, mBits: a.mBits, k: a.k)
|
||||||
|
var inActive = true
|
||||||
|
for i in idx {
|
||||||
|
let byteIndex = i / 8
|
||||||
|
let bitIndex = i % 8
|
||||||
|
let set = ((Int(a.bits[byteIndex]) >> (7 - bitIndex)) & 1) == 1
|
||||||
|
if !set { inActive = false; break }
|
||||||
|
}
|
||||||
|
if inActive { return true }
|
||||||
|
if usingStandby {
|
||||||
|
let s = standby
|
||||||
|
let idx2 = indicesFor(id: id, mBits: s.mBits, k: s.k)
|
||||||
|
for i in idx2 {
|
||||||
|
let byteIndex = i / 8
|
||||||
|
let bitIndex = i % 8
|
||||||
|
let set = ((Int(s.bits[byteIndex]) >> (7 - bitIndex)) & 1) == 1
|
||||||
|
if !set { return false }
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
func snapshotActive() -> Snapshot {
|
||||||
|
lock.lock(); defer { lock.unlock() }
|
||||||
|
let a = active
|
||||||
|
return Snapshot(mBytes: a.bits.count, k: a.k, bits: Data(a.bits))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
Reference in New Issue
Block a user