diff --git a/bitchat/Models/RequestSyncPacket.swift b/bitchat/Models/RequestSyncPacket.swift index b11b8c1a..10eb0c68 100644 --- a/bitchat/Models/RequestSyncPacket.swift +++ b/bitchat/Models/RequestSyncPacket.swift @@ -8,6 +8,14 @@ struct RequestSyncPacket { let p: Int let m: UInt32 let data: Data + let types: SyncTypeFlags? + + init(p: Int, m: UInt32, data: Data, types: SyncTypeFlags? = nil) { + self.p = p + self.m = m + self.data = data + self.types = types + } func encode() -> Data { var out = Data() @@ -25,6 +33,9 @@ struct RequestSyncPacket { putTLV(0x02, withUnsafeBytes(of: &mBE) { Data($0) }) // data putTLV(0x03, data) + if let typesData = types?.toData() { + putTLV(0x04, typesData) + } return out } @@ -33,6 +44,7 @@ struct RequestSyncPacket { var p: Int? = nil var m: UInt32? = nil var payload: Data? = nil + var types: SyncTypeFlags? = nil while off + 3 <= data.count { let t = Int(data[off]); off += 1 @@ -52,12 +64,16 @@ struct RequestSyncPacket { case 0x03: if v.count > maxAcceptBytes { return nil } payload = v + case 0x04: + if let decoded = SyncTypeFlags.decode(v) { + types = decoded + } 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) + return RequestSyncPacket(p: pp, m: mm, data: dd, types: types) } } diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index 0f20fcb1..141d2fbe 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -2803,6 +2803,9 @@ extension BLEService { let isActive = self.collectionsQueue.sync { self.activeTransfers[transferId] != nil } guard isActive else { return } } + if fragmentRecipient == nil || fragmentRecipient?.allSatisfy({ $0 == 0xFF }) == true { + self.gossipSyncManager?.onPublicPacketSeen(fragmentPacket) + } self.broadcastPacket(fragmentPacket) if let transferId = transferIdentifier { self.markFragmentSent(transferId: transferId) @@ -2904,6 +2907,14 @@ extension BLEService { // Sanity checks - add reasonable upper bound on total to prevent DoS guard total > 0 && total <= 10000 && index >= 0 && index < total else { return } + let isBroadcastFragment: Bool = { + guard let recipient = packet.recipientID else { return true } + return recipient.count == 8 && recipient.allSatisfy { $0 == 0xFF } + }() + if isBroadcastFragment { + gossipSyncManager?.onPublicPacketSeen(packet) + } + // Compute fragment key for this assembly let key = FragmentKey(sender: senderU64, id: fragU64) diff --git a/bitchat/Sync/GossipSyncManager.swift b/bitchat/Sync/GossipSyncManager.swift index e8b953b7..bd626a04 100644 --- a/bitchat/Sync/GossipSyncManager.swift +++ b/bitchat/Sync/GossipSyncManager.swift @@ -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 } } } diff --git a/bitchat/Sync/SyncTypeFlags.swift b/bitchat/Sync/SyncTypeFlags.swift new file mode 100644 index 00000000..15e22d89 --- /dev/null +++ b/bitchat/Sync/SyncTypeFlags.swift @@ -0,0 +1,104 @@ +import Foundation + +/// Bitfield describing which message types are covered by a REQUEST_SYNC round. +/// Matches the Android mapping (bit index -> message type). +struct SyncTypeFlags: OptionSet { + let rawValue: UInt64 + + init(rawValue: UInt64) { + self.rawValue = rawValue & 0x00FF_FFFF_FFFF_FFFF // Trim to max 8 bytes + } + + private static func bitIndex(for type: MessageType) -> Int? { + switch type { + case .announce: return 0 + case .message: return 1 + case .leave: return 2 + case .noiseHandshake: return 3 + case .noiseEncrypted: return 4 + case .fragment: return 5 + case .requestSync: return 6 + case .fileTransfer: return 7 + } + } + + private static func type(forBit index: Int) -> MessageType? { + switch index { + case 0: return .announce + case 1: return .message + case 2: return .leave + case 3: return .noiseHandshake + case 4: return .noiseEncrypted + case 5: return .fragment + case 6: return .requestSync + case 7: return .fileTransfer + default: + return nil + } + } + + static let announce = SyncTypeFlags(messageTypes: [.announce]) + static let message = SyncTypeFlags(messageTypes: [.message]) + static let fragment = SyncTypeFlags(messageTypes: [.fragment]) + static let fileTransfer = SyncTypeFlags(messageTypes: [.fileTransfer]) + + static let publicMessages = SyncTypeFlags(messageTypes: [.announce, .message]) + + init(messageTypes: [MessageType]) { + var raw: UInt64 = 0 + for type in messageTypes { + guard let bit = SyncTypeFlags.bitIndex(for: type) else { continue } + raw |= (1 << UInt64(bit)) + } + self.init(rawValue: raw) + } + + func contains(_ type: MessageType) -> Bool { + guard let bit = SyncTypeFlags.bitIndex(for: type) else { return false } + return contains(SyncTypeFlags(rawValue: 1 << UInt64(bit))) + } + + func union(_ other: SyncTypeFlags) -> SyncTypeFlags { + SyncTypeFlags(rawValue: rawValue | other.rawValue) + } + + func intersection(_ other: SyncTypeFlags) -> SyncTypeFlags { + SyncTypeFlags(rawValue: rawValue & other.rawValue) + } + + func toMessageTypes() -> [MessageType] { + guard rawValue != 0 else { return [] } + var types: [MessageType] = [] + for bit in 0..<64 { + guard (rawValue & (1 << UInt64(bit))) != 0 else { continue } + if let type = SyncTypeFlags.type(forBit: bit) { + types.append(type) + } + } + return types + } + + func toData() -> Data? { + guard rawValue != 0 else { return nil } + var value = rawValue + var bytes: [UInt8] = [] + while value > 0 && bytes.count < 8 { + bytes.append(UInt8(value & 0xFF)) + value >>= 8 + } + while let last = bytes.last, last == 0 { + bytes.removeLast() + } + guard !bytes.isEmpty, bytes.count <= 8 else { return nil } + return Data(bytes) + } + + static func decode(_ data: Data) -> SyncTypeFlags? { + guard (1...8).contains(data.count) else { return nil } + var raw: UInt64 = 0 + for (index, byte) in data.enumerated() { + raw |= UInt64(byte) << UInt64(index * 8) + } + return SyncTypeFlags(rawValue: raw) + } +} diff --git a/bitchatTests/GossipSyncManagerTests.swift b/bitchatTests/GossipSyncManagerTests.swift index a9ca9cf2..4816b504 100644 --- a/bitchatTests/GossipSyncManagerTests.swift +++ b/bitchatTests/GossipSyncManagerTests.swift @@ -125,16 +125,138 @@ struct GossipSyncManagerTests { #expect(manager._hasAnnouncement(for: PeerID(str: peerHex)) == false) #expect(manager._messageCount(for: PeerID(str: peerHex)) == 0) } + + @Test func maintenanceEmitsTypedSyncRequests() throws { + var config = GossipSyncManager.Config() + config.seenCapacity = 10 + config.fragmentCapacity = 5 + config.fileTransferCapacity = 4 + config.messageSyncIntervalSeconds = 1 + config.fragmentSyncIntervalSeconds = 1 + config.fileTransferSyncIntervalSeconds = 1 + config.maintenanceIntervalSeconds = 0 + + let manager = GossipSyncManager(myPeerID: myPeerID, config: config) + let delegate = RecordingDelegate() + manager.delegate = delegate + + let sender = try #require(Data(hexString: "1122334455667788")) + let now = UInt64(Date().timeIntervalSince1970 * 1000) + + let announcePacket = BitchatPacket( + type: MessageType.announce.rawValue, + senderID: sender, + recipientID: nil, + timestamp: now, + payload: Data(), + signature: nil, + ttl: 1 + ) + let messagePacket = BitchatPacket( + type: MessageType.message.rawValue, + senderID: sender, + recipientID: nil, + timestamp: now, + payload: Data([0x01]), + signature: nil, + ttl: 1 + ) + let fragmentPacket = BitchatPacket( + type: MessageType.fragment.rawValue, + senderID: sender, + recipientID: nil, + timestamp: now, + payload: Data([0xAA]), + signature: nil, + ttl: 1 + ) + let filePacket = BitchatPacket( + type: MessageType.fileTransfer.rawValue, + senderID: sender, + recipientID: nil, + timestamp: now, + payload: Data([0xBB]), + signature: nil, + ttl: 1, + version: 2 + ) + + manager.onPublicPacketSeen(announcePacket) + manager.onPublicPacketSeen(messagePacket) + manager.onPublicPacketSeen(fragmentPacket) + manager.onPublicPacketSeen(filePacket) + + manager._performMaintenanceSynchronously(now: Date()) + + let sentPackets = delegate.packets + #expect(sentPackets.count == 3) + let decoded = sentPackets.compactMap { RequestSyncPacket.decode(from: $0.payload) } + #expect(decoded.count == 3) + #expect(decoded[0].types == .publicMessages) + #expect(decoded[1].types == .fragment) + #expect(decoded[2].types == .fileTransfer) + } + + @Test func handleRequestSyncHonorsTypeFilter() async throws { + var config = GossipSyncManager.Config() + config.seenCapacity = 5 + config.fragmentCapacity = 5 + config.fileTransferCapacity = 0 + config.messageSyncIntervalSeconds = 0 + config.fragmentSyncIntervalSeconds = 0 + config.fileTransferSyncIntervalSeconds = 0 + + let manager = GossipSyncManager(myPeerID: myPeerID, config: config) + let delegate = RecordingDelegate() + manager.delegate = delegate + + let sender = try #require(Data(hexString: "aabbccddeeff0011")) + let now = UInt64(Date().timeIntervalSince1970 * 1000) + + let messagePacket = BitchatPacket( + type: MessageType.message.rawValue, + senderID: sender, + recipientID: nil, + timestamp: now, + payload: Data([0x10]), + signature: nil, + ttl: 1 + ) + + let fragmentPacket = BitchatPacket( + type: MessageType.fragment.rawValue, + senderID: sender, + recipientID: nil, + timestamp: now, + payload: Data([0x20]), + signature: nil, + ttl: 1 + ) + + manager.onPublicPacketSeen(messagePacket) + manager.onPublicPacketSeen(fragmentPacket) + + let peer = PeerID(str: "FFFFFFFFFFFFFFFF") + let request = RequestSyncPacket(p: 4, m: 1, data: Data(), types: .fragment) + manager.handleRequestSync(from: peer, request: request) + + try await sleep(0.01) + let sentPackets = delegate.packets + #expect(sentPackets.count == 1) + #expect(sentPackets[0].type == MessageType.fragment.rawValue) + } } private final class RecordingDelegate: GossipSyncManager.Delegate { var onSend: (() -> Void)? private(set) var lastPacket: BitchatPacket? + private(set) var packets: [BitchatPacket] = [] private let lock = NSLock() func sendPacket(_ packet: BitchatPacket) { lock.lock() lastPacket = packet + packets.append(packet) lock.unlock() onSend?() }