Add type-aware REQUEST_SYNC sync rounds (#853)

Co-authored-by: jack <jackjackbits@users.noreply.github.com>
This commit is contained in:
jack
2025-10-22 12:24:29 +02:00
committed by GitHub
co-authored by jack
parent 98fa02cc16
commit b6ce4fae43
5 changed files with 447 additions and 84 deletions
+193 -83
View File
@@ -8,6 +8,55 @@ final class GossipSyncManager {
func signPacketForBroadcast(_ packet: BitchatPacket) -> BitchatPacket
}
private struct PacketStore {
private(set) var packets: [String: BitchatPacket] = [:]
private(set) var order: [String] = []
mutating func insert(idHex: String, packet: BitchatPacket, capacity: Int) {
guard capacity > 0 else { return }
if packets[idHex] != nil {
packets[idHex] = packet
return
}
packets[idHex] = packet
order.append(idHex)
while order.count > capacity {
let victim = order.removeFirst()
packets.removeValue(forKey: victim)
}
}
func allPackets(isFresh: (BitchatPacket) -> Bool) -> [BitchatPacket] {
order.compactMap { key in
guard let packet = packets[key], isFresh(packet) else { return nil }
return packet
}
}
mutating func remove(where shouldRemove: (BitchatPacket) -> Bool) {
var nextOrder: [String] = []
for key in order {
guard let packet = packets[key] else { continue }
if shouldRemove(packet) {
packets.removeValue(forKey: key)
} else {
nextOrder.append(key)
}
}
order = nextOrder
}
mutating func removeExpired(isFresh: (BitchatPacket) -> Bool) {
remove { !isFresh($0) }
}
}
private struct SyncSchedule {
let types: SyncTypeFlags
let interval: TimeInterval
var lastSent: Date
}
struct Config {
var seenCapacity: Int = 1000 // max packets per sync (cap across types)
var gcsMaxBytes: Int = 400 // filter size budget (128..1024)
@@ -16,25 +65,43 @@ final class GossipSyncManager {
var maintenanceIntervalSeconds: TimeInterval = 30.0
var stalePeerCleanupIntervalSeconds: TimeInterval = 60.0
var stalePeerTimeoutSeconds: TimeInterval = 60.0
var fragmentCapacity: Int = 600
var fileTransferCapacity: Int = 200
var fragmentSyncIntervalSeconds: TimeInterval = 30.0
var fileTransferSyncIntervalSeconds: TimeInterval = 60.0
var messageSyncIntervalSeconds: TimeInterval = 15.0
}
private let myPeerID: PeerID
private let config: Config
weak var delegate: Delegate?
// Storage: broadcast messages (ordered by insert), and latest announce per sender
private var messages: [String: BitchatPacket] = [:] // idHex -> packet
private var messageOrder: [String] = []
// Storage: broadcast packets by type, and latest announce per sender
private var messages = PacketStore()
private var fragments = PacketStore()
private var fileTransfers = PacketStore()
private var latestAnnouncementByPeer: [PeerID: (id: String, packet: BitchatPacket)] = [:]
// Timer
private var periodicTimer: DispatchSourceTimer?
private let queue = DispatchQueue(label: "mesh.sync", qos: .utility)
private var lastStalePeerCleanup: Date = .distantPast
private var syncSchedules: [SyncSchedule] = []
init(myPeerID: PeerID, config: Config = Config()) {
self.myPeerID = myPeerID
self.config = config
var schedules: [SyncSchedule] = []
if config.seenCapacity > 0 && config.messageSyncIntervalSeconds > 0 {
schedules.append(SyncSchedule(types: .publicMessages, interval: config.messageSyncIntervalSeconds, lastSent: .distantPast))
}
if config.fragmentCapacity > 0 && config.fragmentSyncIntervalSeconds > 0 {
schedules.append(SyncSchedule(types: .fragment, interval: config.fragmentSyncIntervalSeconds, lastSent: .distantPast))
}
if config.fileTransferCapacity > 0 && config.fileTransferSyncIntervalSeconds > 0 {
schedules.append(SyncSchedule(types: .fileTransfer, interval: config.fileTransferSyncIntervalSeconds, lastSent: .distantPast))
}
syncSchedules = schedules
}
func start() {
@@ -55,7 +122,18 @@ final class GossipSyncManager {
func scheduleInitialSyncToPeer(_ peerID: PeerID, delaySeconds: TimeInterval = 5.0) {
queue.asyncAfter(deadline: .now() + delaySeconds) { [weak self] in
self?.sendRequestSync(to: peerID)
guard let self = self else { return }
self.sendRequestSync(to: peerID, types: .publicMessages)
if self.config.fragmentCapacity > 0 && self.config.fragmentSyncIntervalSeconds > 0 {
self.queue.asyncAfter(deadline: .now() + 0.5) { [weak self] in
self?.sendRequestSync(to: peerID, types: .fragment)
}
}
if self.config.fileTransferCapacity > 0 && self.config.fileTransferSyncIntervalSeconds > 0 {
self.queue.asyncAfter(deadline: .now() + 1.0) { [weak self] in
self?.sendRequestSync(to: peerID, types: .fileTransfer)
}
}
}
}
@@ -87,47 +165,45 @@ final class GossipSyncManager {
}
private func _onPublicPacketSeen(_ packet: BitchatPacket) {
let mt = MessageType(rawValue: packet.type)
guard let messageType = MessageType(rawValue: packet.type) else { return }
let isBroadcastRecipient: Bool = {
guard let r = packet.recipientID else { return true }
return r.count == 8 && r.allSatisfy { $0 == 0xFF }
}()
let isBroadcastMessage = (mt == .message && isBroadcastRecipient)
let isAnnounce = (mt == .announce)
guard isBroadcastMessage || isAnnounce else { return }
// Reject expired packets to prevent ghost peers and old messages
guard isPacketFresh(packet) else { return }
if isAnnounce {
switch messageType {
case .announce:
guard isPacketFresh(packet) else { return }
guard isAnnouncementFresh(packet) else {
let sender = PeerID(hexData: packet.senderID)
removeState(for: sender)
return
}
}
let idHex = PacketIdUtil.computeId(packet).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 idHex = PacketIdUtil.computeId(packet).hexEncodedString()
let sender = PeerID(hexData: packet.senderID)
latestAnnouncementByPeer[sender] = (id: idHex, packet: packet)
case .message:
guard isBroadcastRecipient else { return }
guard isPacketFresh(packet) else { return }
let idHex = PacketIdUtil.computeId(packet).hexEncodedString()
messages.insert(idHex: idHex, packet: packet, capacity: max(1, config.seenCapacity))
case .fragment:
guard isBroadcastRecipient else { return }
guard isPacketFresh(packet) else { return }
let idHex = PacketIdUtil.computeId(packet).hexEncodedString()
fragments.insert(idHex: idHex, packet: packet, capacity: max(1, config.fragmentCapacity))
case .fileTransfer:
guard isBroadcastRecipient else { return }
guard isPacketFresh(packet) else { return }
let idHex = PacketIdUtil.computeId(packet).hexEncodedString()
fileTransfers.insert(idHex: idHex, packet: packet, capacity: max(1, config.fileTransferCapacity))
default:
break
}
}
private func sendRequestSync() {
let payload = buildGcsPayload()
private func sendRequestSync(for types: SyncTypeFlags) {
let payload = buildGcsPayload(for: types)
let pkt = BitchatPacket(
type: MessageType.requestSync.rawValue,
senderID: Data(hexString: myPeerID.id) ?? Data(),
@@ -141,8 +217,8 @@ final class GossipSyncManager {
delegate?.sendPacket(signed)
}
private func sendRequestSync(to peerID: PeerID) {
let payload = buildGcsPayload()
private func sendRequestSync(to peerID: PeerID, types: SyncTypeFlags) {
let payload = buildGcsPayload(for: types)
var recipient = Data()
var temp = peerID.id
while temp.count >= 2 && recipient.count < 8 {
@@ -170,6 +246,7 @@ final class GossipSyncManager {
}
private func _handleRequestSync(from peerID: PeerID, request: RequestSyncPacket) {
let requestedTypes = (request.types ?? .publicMessages)
// Decode GCS into sorted set and prepare membership checker
let sorted = GCSFilter.decodeToSortedSet(p: request.p, m: request.m, data: request.data)
func mightContain(_ id: Data) -> Bool {
@@ -177,60 +254,100 @@ final class GossipSyncManager {
return GCSFilter.contains(sortedValues: sorted, candidate: bucket)
}
// 1) Announcements: send latest per peer if requester lacks them (and not expired)
for (_, pair) in latestAnnouncementByPeer {
let (idHex, pkt) = pair
guard isPacketFresh(pkt) else { continue }
let idBytes = Data(hexString: idHex) ?? Data()
if !mightContain(idBytes) {
var toSend = pkt
toSend.ttl = 0
delegate?.sendPacket(to: peerID, packet: toSend)
if requestedTypes.contains(.announce) {
for (_, pair) in latestAnnouncementByPeer {
let (idHex, pkt) = pair
guard isPacketFresh(pkt) else { continue }
let idBytes = Data(hexString: idHex) ?? Data()
if !mightContain(idBytes) {
var toSend = pkt
toSend.ttl = 0
delegate?.sendPacket(to: peerID, packet: toSend)
}
}
}
// 2) Broadcast messages: send all missing (and not expired)
let toSendMsgs = messageOrder.compactMap { messages[$0] }
for pkt in toSendMsgs {
guard isPacketFresh(pkt) else { continue }
let idBytes = PacketIdUtil.computeId(pkt)
if !mightContain(idBytes) {
var toSend = pkt
toSend.ttl = 0
delegate?.sendPacket(to: peerID, packet: toSend)
if requestedTypes.contains(.message) {
let toSendMsgs = messages.allPackets(isFresh: isPacketFresh)
for pkt in toSendMsgs {
let idBytes = PacketIdUtil.computeId(pkt)
if !mightContain(idBytes) {
var toSend = pkt
toSend.ttl = 0
delegate?.sendPacket(to: peerID, packet: toSend)
}
}
}
if requestedTypes.contains(.fragment) {
let frags = fragments.allPackets(isFresh: isPacketFresh)
for pkt in frags {
let idBytes = PacketIdUtil.computeId(pkt)
if !mightContain(idBytes) {
var toSend = pkt
toSend.ttl = 0
delegate?.sendPacket(to: peerID, packet: toSend)
}
}
}
if requestedTypes.contains(.fileTransfer) {
let files = fileTransfers.allPackets(isFresh: isPacketFresh)
for pkt in files {
let idBytes = PacketIdUtil.computeId(pkt)
if !mightContain(idBytes) {
var toSend = pkt
toSend.ttl = 0
delegate?.sendPacket(to: peerID, packet: toSend)
}
}
}
}
// Build REQUEST_SYNC payload using current candidates and GCS params
private func buildGcsPayload() -> Data {
// Collect candidates: latest announce per peer + broadcast messages (only fresh)
private func buildGcsPayload(for types: SyncTypeFlags) -> Data {
var candidates: [BitchatPacket] = []
candidates.reserveCapacity(latestAnnouncementByPeer.count + messageOrder.count)
for (_, pair) in latestAnnouncementByPeer {
if isPacketFresh(pair.packet) {
if types.contains(.announce) {
for (_, pair) in latestAnnouncementByPeer where isPacketFresh(pair.packet) {
candidates.append(pair.packet)
}
}
for id in messageOrder {
if let p = messages[id], isPacketFresh(p) {
candidates.append(p)
}
if types.contains(.message) {
candidates.append(contentsOf: messages.allPackets(isFresh: isPacketFresh))
}
if types.contains(.fragment) {
candidates.append(contentsOf: fragments.allPackets(isFresh: isPacketFresh))
}
if types.contains(.fileTransfer) {
candidates.append(contentsOf: fileTransfers.allPackets(isFresh: isPacketFresh))
}
if candidates.isEmpty {
let p = GCSFilter.deriveP(targetFpr: config.gcsTargetFpr)
let req = RequestSyncPacket(p: p, m: 1, data: Data(), types: types)
return req.encode()
}
// Sort by timestamp desc
candidates.sort { $0.timestamp > $1.timestamp }
let p = GCSFilter.deriveP(targetFpr: config.gcsTargetFpr)
let nMax = GCSFilter.estimateMaxElements(sizeBytes: config.gcsMaxBytes, p: p)
let cap = max(1, config.seenCapacity)
let cap: Int
if types == .fragment {
cap = max(1, config.fragmentCapacity)
} else if types == .fileTransfer {
cap = max(1, config.fileTransferCapacity)
} else {
cap = max(1, config.seenCapacity)
}
let takeN = min(candidates.count, min(nMax, cap))
if takeN <= 0 {
let req = RequestSyncPacket(p: p, m: 1, data: Data())
let req = RequestSyncPacket(p: p, m: 1, data: Data(), types: types)
return req.encode()
}
let ids: [Data] = candidates.prefix(takeN).map { PacketIdUtil.computeId($0) }
let params = GCSFilter.buildFilter(ids: ids, maxBytes: config.gcsMaxBytes, targetFpr: config.gcsTargetFpr)
let req = RequestSyncPacket(p: params.p, m: params.m, data: params.data)
let req = RequestSyncPacket(p: params.p, m: params.m, data: params.data, types: types)
return req.encode()
}
@@ -241,20 +358,21 @@ final class GossipSyncManager {
isPacketFresh(pair.packet)
}
// Remove expired messages
let expiredMessageIds = messages.compactMap { id, pkt in
isPacketFresh(pkt) ? nil : id
}
for id in expiredMessageIds {
messages.removeValue(forKey: id)
messageOrder.removeAll { $0 == id }
}
messages.removeExpired(isFresh: isPacketFresh)
fragments.removeExpired(isFresh: isPacketFresh)
fileTransfers.removeExpired(isFresh: isPacketFresh)
}
private func performPeriodicMaintenance(now: Date = Date()) {
cleanupExpiredMessages()
cleanupStaleAnnouncementsIfNeeded(now: now)
sendRequestSync()
for index in syncSchedules.indices {
guard syncSchedules[index].interval > 0 else { continue }
if syncSchedules[index].lastSent == .distantPast || now.timeIntervalSince(syncSchedules[index].lastSent) >= syncSchedules[index].interval {
syncSchedules[index].lastSent = now
sendRequestSync(for: syncSchedules[index].types)
}
}
}
private func cleanupStaleAnnouncementsIfNeeded(now: Date) {
@@ -288,17 +406,9 @@ final class GossipSyncManager {
private func removeState(for peerID: PeerID) {
_ = latestAnnouncementByPeer.removeValue(forKey: peerID)
// Remove messages from this peer
// Collect IDs to remove first to avoid concurrent modification
let messageIdsToRemove = messages.compactMap { (id, message) -> String? in
PeerID(hexData: message.senderID) == peerID ? id : nil
}
// Remove messages and update messageOrder
for id in messageIdsToRemove {
messages.removeValue(forKey: id)
messageOrder.removeAll { $0 == id }
}
messages.remove { PeerID(hexData: $0.senderID) == peerID }
fragments.remove { PeerID(hexData: $0.senderID) == peerID }
fileTransfers.remove { PeerID(hexData: $0.senderID) == peerID }
}
}
@@ -318,7 +428,7 @@ extension GossipSyncManager {
func _messageCount(for peerID: PeerID) -> Int {
queue.sync {
messages.values.filter { PeerID(hexData: $0.senderID) == peerID }.count
messages.allPackets { _ in true }.filter { PeerID(hexData: $0.senderID) == peerID }.count
}
}
}