diff --git a/bitchat/Sync/GossipSyncManager.swift b/bitchat/Sync/GossipSyncManager.swift index 85ded306..7f3cebbb 100644 --- a/bitchat/Sync/GossipSyncManager.swift +++ b/bitchat/Sync/GossipSyncManager.swift @@ -13,6 +13,9 @@ final class GossipSyncManager { var gcsMaxBytes: Int = 400 // filter size budget (128..1024) var gcsTargetFpr: Double = 0.01 // 1% var maxMessageAgeSeconds: TimeInterval = 900 // 15 min - discard older messages + var maintenanceIntervalSeconds: TimeInterval = 30.0 + var stalePeerCleanupIntervalSeconds: TimeInterval = 60.0 + var stalePeerTimeoutSeconds: TimeInterval = 60.0 } private let myPeerID: PeerID @@ -27,6 +30,7 @@ final class GossipSyncManager { // Timer private var periodicTimer: DispatchSourceTimer? private let queue = DispatchQueue(label: "mesh.sync", qos: .utility) + private var lastStalePeerCleanup: Date = .distantPast init(myPeerID: PeerID, config: Config = Config()) { self.myPeerID = myPeerID @@ -36,10 +40,10 @@ final class GossipSyncManager { func start() { stop() let timer = DispatchSource.makeTimerSource(queue: queue) - timer.schedule(deadline: .now() + 30.0, repeating: 30.0, leeway: .seconds(1)) + let interval = max(0.1, config.maintenanceIntervalSeconds) + timer.schedule(deadline: .now() + interval, repeating: interval, leeway: .seconds(1)) timer.setEventHandler { [weak self] in - self?.cleanupExpiredMessages() - self?.sendRequestSync() + self?.performPeriodicMaintenance() } timer.resume() periodicTimer = timer @@ -73,6 +77,15 @@ final class GossipSyncManager { return packet.timestamp >= cutoffMs } + private func isAnnouncementFresh(_ packet: BitchatPacket) -> Bool { + guard config.stalePeerTimeoutSeconds > 0 else { return true } + let nowMs = UInt64(Date().timeIntervalSince1970 * 1000) + let timeoutMs = UInt64(config.stalePeerTimeoutSeconds * 1000) + guard nowMs >= timeoutMs else { return true } + let cutoffMs = nowMs - timeoutMs + return packet.timestamp >= cutoffMs + } + private func _onPublicPacketSeen(_ packet: BitchatPacket) { let mt = MessageType(rawValue: packet.type) let isBroadcastRecipient: Bool = { @@ -86,6 +99,14 @@ final class GossipSyncManager { // Reject expired packets to prevent ghost peers and old messages guard isPacketFresh(packet) else { return } + if isAnnounce { + guard isAnnouncementFresh(packet) else { + let sender = packet.senderID.hexEncodedString().lowercased() + removeState(forNormalizedPeerID: sender) + return + } + } + let idHex = PacketIdUtil.computeId(packet).hexEncodedString() if isBroadcastMessage { @@ -100,7 +121,7 @@ final class GossipSyncManager { } } } else if isAnnounce { - let sender = packet.senderID.hexEncodedString() + let sender = packet.senderID.hexEncodedString().lowercased() latestAnnouncementByPeer[sender] = (id: idHex, packet: packet) } } @@ -230,6 +251,34 @@ final class GossipSyncManager { } } + private func performPeriodicMaintenance(now: Date = Date()) { + cleanupExpiredMessages() + cleanupStaleAnnouncementsIfNeeded(now: now) + sendRequestSync() + } + + private func cleanupStaleAnnouncementsIfNeeded(now: Date) { + guard now.timeIntervalSince(lastStalePeerCleanup) >= config.stalePeerCleanupIntervalSeconds else { + return + } + lastStalePeerCleanup = now + cleanupStaleAnnouncements(now: now) + } + + private func cleanupStaleAnnouncements(now: Date) { + let timeoutMs = UInt64(config.stalePeerTimeoutSeconds * 1000) + let nowMs = UInt64(now.timeIntervalSince1970 * 1000) + guard nowMs >= timeoutMs else { return } + let cutoff = nowMs - timeoutMs + let stalePeerIDs = latestAnnouncementByPeer.compactMap { (peerHex, pair) -> String? in + pair.packet.timestamp < cutoff ? peerHex.lowercased() : nil + } + guard !stalePeerIDs.isEmpty else { return } + for peerKey in stalePeerIDs { + removeState(forNormalizedPeerID: peerKey) + } + } + // Explicit removal hook for LEAVE/stale peer func removeAnnouncementForPeer(_ peerID: PeerID) { queue.async { [weak self] in @@ -239,8 +288,11 @@ final class GossipSyncManager { private func _removeAnnouncementForPeer(_ peerID: PeerID) { let normalizedPeerID = peerID.id.lowercased() - _ = latestAnnouncementByPeer.removeValue(forKey: normalizedPeerID) + removeState(forNormalizedPeerID: normalizedPeerID) + } + private func removeState(forNormalizedPeerID normalizedPeerID: String) { + _ = latestAnnouncementByPeer.removeValue(forKey: normalizedPeerID) // Remove messages from this peer // Collect IDs to remove first to avoid concurrent modification let messageIdsToRemove = messages.compactMap { (id, message) -> String? in @@ -254,3 +306,25 @@ final class GossipSyncManager { } } } + +#if DEBUG +extension GossipSyncManager { + func _performMaintenanceSynchronously(now: Date = Date()) { + queue.sync { + performPeriodicMaintenance(now: now) + } + } + + func _hasAnnouncement(for peerID: PeerID) -> Bool { + queue.sync { + latestAnnouncementByPeer[peerID.id.lowercased()] != nil + } + } + + func _messageCount(for peerID: PeerID) -> Int { + queue.sync { + messages.values.filter { $0.senderID.hexEncodedString().lowercased() == peerID.id.lowercased() }.count + } + } +} +#endif diff --git a/bitchatTests/GossipSyncManagerTests.swift b/bitchatTests/GossipSyncManagerTests.swift index 6455e326..c37610e7 100644 --- a/bitchatTests/GossipSyncManagerTests.swift +++ b/bitchatTests/GossipSyncManagerTests.swift @@ -46,6 +46,90 @@ final class GossipSyncManagerTests: XCTestCase { XCTAssertEqual(lastPacket.type, MessageType.requestSync.rawValue) XCTAssertNotNil(RequestSyncPacket.decode(from: lastPacket.payload)) } + + func testStaleAnnouncementsArePurgedWithMessages() { + var config = GossipSyncManager.Config() + config.stalePeerCleanupIntervalSeconds = 0 + config.stalePeerTimeoutSeconds = 5 + + let manager = GossipSyncManager(myPeerID: "0102030405060708", config: config) + let peerHex = "0011223344556677" + let senderData = Data(hexString: peerHex) ?? Data() + let initialTimestampMs = UInt64(Date().timeIntervalSince1970 * 1000) + + let announcePacket = BitchatPacket( + type: MessageType.announce.rawValue, + senderID: senderData, + recipientID: nil, + timestamp: initialTimestampMs, + payload: Data(), + signature: nil, + ttl: 1 + ) + + let messagePacket = BitchatPacket( + type: MessageType.message.rawValue, + senderID: senderData, + recipientID: nil, + timestamp: initialTimestampMs, + payload: Data([0x01]), + signature: nil, + ttl: 1 + ) + + manager.onPublicPacketSeen(announcePacket) + manager.onPublicPacketSeen(messagePacket) + + // Flush queue without triggering stale cleanup yet + manager._performMaintenanceSynchronously(now: Date()) + XCTAssertTrue(manager._hasAnnouncement(for: PeerID(str: peerHex))) + XCTAssertEqual(manager._messageCount(for: PeerID(str: peerHex)), 1) + + // Run cleanup past the timeout + let future = Date().addingTimeInterval(config.stalePeerTimeoutSeconds + 1) + manager._performMaintenanceSynchronously(now: future) + XCTAssertFalse(manager._hasAnnouncement(for: PeerID(str: peerHex))) + XCTAssertEqual(manager._messageCount(for: PeerID(str: peerHex)), 0) + } + + func testIgnoresAnnounceOlderThanStaleTimeout() { + var config = GossipSyncManager.Config() + config.stalePeerTimeoutSeconds = 5 + config.maxMessageAgeSeconds = 100 + + let manager = GossipSyncManager(myPeerID: "0102030405060708", config: config) + let peerHex = "8899aabbccddeeff" + let senderData = Data(hexString: peerHex) ?? Data() + let staleTimestampMs = UInt64(Date().addingTimeInterval(-(config.stalePeerTimeoutSeconds + 1)).timeIntervalSince1970 * 1000) + + let freshMessage = BitchatPacket( + type: MessageType.message.rawValue, + senderID: senderData, + recipientID: nil, + timestamp: UInt64(Date().timeIntervalSince1970 * 1000), + payload: Data([0xAA]), + signature: nil, + ttl: 1 + ) + manager.onPublicPacketSeen(freshMessage) + + let announcePacket = BitchatPacket( + type: MessageType.announce.rawValue, + senderID: senderData, + recipientID: nil, + timestamp: staleTimestampMs, + payload: Data(), + signature: nil, + ttl: 1 + ) + + manager.onPublicPacketSeen(announcePacket) + + manager._performMaintenanceSynchronously() + + XCTAssertFalse(manager._hasAnnouncement(for: PeerID(str: peerHex))) + XCTAssertEqual(manager._messageCount(for: PeerID(str: peerHex)), 0) + } } private final class RecordingDelegate: GossipSyncManager.Delegate {