diff --git a/bitchat/Services/BLE/BLEFragmentAssemblyBuffer.swift b/bitchat/Services/BLE/BLEFragmentAssemblyBuffer.swift index ad6e007c..83311584 100644 --- a/bitchat/Services/BLE/BLEFragmentAssemblyBuffer.swift +++ b/bitchat/Services/BLE/BLEFragmentAssemblyBuffer.swift @@ -108,8 +108,15 @@ struct BLEFragmentAssemblyBuffer { return .oversized(header: header, projectedSize: projectedSize, limit: limit, started: started) } + // Only actual progress resets the stall clock: fragment packets + // bypass the packet deduplicator, so relayed duplicates of an + // already-held index must not keep suppressing the targeted + // REQUEST_SYNC for a stalled stream. + let isNewIndex = fragmentsByKey[header.key]?[header.index] == nil fragmentsByKey[header.key]?[header.index] = header.fragmentData - metadataByKey[header.key]?.lastFragmentAt = now + if isNewIndex { + metadataByKey[header.key]?.lastFragmentAt = now + } guard let fragments = fragmentsByKey[header.key], fragments.count == header.total else { @@ -156,14 +163,18 @@ struct BLEFragmentAssemblyBuffer { /// reassemblies that have not seen a new fragment for `stalledAfter` /// seconds — candidates for a targeted REQUEST_SYNC. Each returned /// stream is marked so it is not re-requested within `retryAfter`. - /// Directed reassemblies are excluded: peers only archive broadcast - /// fragments for gossip sync, so a targeted request cannot recover them. + /// At most `RequestSyncPacket.maxFragmentIdFilterCount` streams are + /// returned per pass — the wire filter cannot carry more — selected + /// oldest-stall first; overflow streams stay unmarked and eligible for + /// the next pass. Directed reassemblies are excluded: peers only archive + /// broadcast fragments for gossip sync, so a targeted request cannot + /// recover them. mutating func stalledBroadcastFragmentIDs( stalledAfter: TimeInterval, retryAfter: TimeInterval, now: Date = Date() ) -> [Data] { - var stalled: [Data] = [] + var candidates: [(key: BLEFragmentKey, lastFragmentAt: Date)] = [] for (key, metadata) in metadataByKey { guard metadata.isBroadcast, let fragments = fragmentsByKey[key], @@ -171,10 +182,24 @@ struct BLEFragmentAssemblyBuffer { now.timeIntervalSince(metadata.lastFragmentAt) >= stalledAfter else { continue } if let lastRequest = metadata.lastResyncRequestAt, now.timeIntervalSince(lastRequest) < retryAfter { continue } - metadataByKey[key]?.lastResyncRequestAt = now - stalled.append(withUnsafeBytes(of: key.id.bigEndian) { Data($0) }) + candidates.append((key: key, lastFragmentAt: metadata.lastFragmentAt)) + } + + // Mark only the streams that will actually go on the wire, so the + // overflow is not silently suppressed for `retryAfter`. + let selected = candidates + .sorted { + if $0.lastFragmentAt != $1.lastFragmentAt { + return $0.lastFragmentAt < $1.lastFragmentAt + } + return ($0.key.sender, $0.key.id) < ($1.key.sender, $1.key.id) + } + .prefix(RequestSyncPacket.maxFragmentIdFilterCount) + + return selected.map { candidate in + metadataByKey[candidate.key]?.lastResyncRequestAt = now + return withUnsafeBytes(of: candidate.key.id.bigEndian) { Data($0) } } - return stalled } private static func assemblyLimit(for originalType: UInt8) -> Int { diff --git a/bitchatTests/Services/BLEFragmentAssemblyBufferTests.swift b/bitchatTests/Services/BLEFragmentAssemblyBufferTests.swift index ace2cfbe..a4423ec9 100644 --- a/bitchatTests/Services/BLEFragmentAssemblyBufferTests.swift +++ b/bitchatTests/Services/BLEFragmentAssemblyBufferTests.swift @@ -192,6 +192,64 @@ struct BLEFragmentAssemblyBufferTests { #expect(afterCompletion.isEmpty) } + @Test + func duplicateFragmentsDoNotResetStallClock() throws { + var buffer = BLEFragmentAssemblyBuffer() + let fragmentID = Data((20...27).map { UInt8($0) }) + let packet = makePacket(payload: makePayload(count: 256)) + let fragments = try makeFragments(for: packet, chunkSize: 128, fragmentID: fragmentID) + let first = try #require(BLEFragmentHeader(packet: fragments[0])) + + let t0 = Date(timeIntervalSince1970: 100) + _ = buffer.append(first, maxInFlightAssemblies: 8, now: t0) + + // Relay duplicates of the same index arrive every few seconds; they + // bring no new data, so they must not keep the stream "fresh". + _ = buffer.append(first, maxInFlightAssemblies: 8, now: t0.addingTimeInterval(3)) + _ = buffer.append(first, maxInFlightAssemblies: 8, now: t0.addingTimeInterval(5)) + + let stalled = buffer.stalledBroadcastFragmentIDs(stalledAfter: 5, retryAfter: 10, now: t0.addingTimeInterval(6)) + #expect(stalled == [fragmentID]) + } + + @Test + func overflowStalledStreamsRotateAcrossPasses() throws { + var buffer = BLEFragmentAssemblyBuffer() + let cap = RequestSyncPacket.maxFragmentIdFilterCount + let streamCount = cap + 10 + let t0 = Date(timeIntervalSince1970: 100) + + // Incomplete broadcast assemblies with staggered last-fragment times + // (stream 0 is the oldest stall). + var ids: [Data] = [] + for i in 0..> 8), UInt8(i & 0xFF)]) + ids.append(fragmentID) + let header = try #require(BLEFragmentHeader(packet: makeFragmentPacket( + fragmentID: fragmentID, + index: 0, + total: 2, + originalType: MessageType.message.rawValue, + fragmentData: Data([0x01]) + ))) + _ = buffer.append(header, maxInFlightAssemblies: streamCount, now: t0.addingTimeInterval(Double(i))) + } + + // All streams are stalled; only the cap's worth (oldest first) is + // requested and rate-limited, the overflow stays eligible. + let firstPassAt = t0.addingTimeInterval(Double(streamCount) + 5) + let firstPass = buffer.stalledBroadcastFragmentIDs(stalledAfter: 5, retryAfter: 60, now: firstPassAt) + #expect(firstPass == Array(ids.prefix(cap))) + + // Next pass picks up exactly the overflow streams. + let secondPass = buffer.stalledBroadcastFragmentIDs(stalledAfter: 5, retryAfter: 60, now: firstPassAt.addingTimeInterval(1)) + #expect(secondPass == Array(ids.suffix(streamCount - cap))) + + // Nothing left until a retry window lapses. + let thirdPass = buffer.stalledBroadcastFragmentIDs(stalledAfter: 5, retryAfter: 60, now: firstPassAt.addingTimeInterval(2)) + #expect(thirdPass.isEmpty) + } + @Test func directedAssembliesAreNeverReportedAsStalled() throws { var buffer = BLEFragmentAssemblyBuffer()