From 5ae7c9a3a755cc86879185f3ce87ba79ae615859 Mon Sep 17 00:00:00 2001 From: Olivier Giroux Date: Fri, 2 Oct 2026 22:12:00 -0400 Subject: [PATCH 1/6] fix: repair lost video packets with retransmissions before giving up A lost video packet turned into a keyframe recovery and a visible freeze, while the official client repairs the same losses unseen: in a session on the same route it requested 133 packets and received all of them. OpenNOW requested packets only once their gap had already been given up, and the video replay window was RFC 3711's 64 packets, 7.5 ms at ~8,500 packets/s, less than one round trip, so a resent packet was rejected as a replay anyway. - Missing packets are requested while their gap is still open, on the official client's logged schedule: 1 ms after the loss, then up to 3 retries 4 ms apart (NvstNackTracker). - A gap whose first packet was requested waits up to the official 52 ms for it instead of becoming loss once the reorder window passes it. - The video replay window holds 2,048 packets, the official NACK queue length. Audio keeps 64. - The counters line adds replayed, late and duplicate drops and the packets a retransmission repaired. --- .../BifrostFree/NvstMjolnirReceiver.swift | 17 +++-- GFN/NVST/BifrostFree/NvstNackTracker.swift | 67 +++++++++++++++++++ GFN/NVST/BifrostFree/NvstVideoReceiver.swift | 50 +++++++++++++- GFN/NVST/BifrostFree/SrtpCryptography.swift | 50 ++++++++++---- OPN/Stream/NvstBifrostFreeTransport.swift | 2 +- Tests/GFN/NVST/NvstMjolnirReceiverTests.swift | 56 ++++++++++++++++ Tests/GFN/NVST/NvstNackTrackerTests.swift | 42 ++++++++++++ Tests/GFN/NVST/SrtpCryptographyTests.swift | 28 ++++++++ 8 files changed, 292 insertions(+), 20 deletions(-) create mode 100644 GFN/NVST/BifrostFree/NvstNackTracker.swift create mode 100644 Tests/GFN/NVST/NvstNackTrackerTests.swift diff --git a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift index 59a1be43..d233fde5 100644 --- a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift @@ -252,11 +252,18 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { /// packet by packet. @discardableResult public func requestRetransmission(firstMissing: UInt64, lastMissing: UInt64) -> Int { - guard let ssrc = receiver.feedbackCounters.boundSSRC, lastMissing >= firstMissing else { return 0 } - let span = Int(lastMissing - firstMissing) + 1 + guard lastMissing >= firstMissing else { return 0 } // Beyond a burst this size a keyframe arrives sooner than the retransmissions would. - guard span <= NvstRtcp.maximumNackEntries * 16 else { return 0 } - let missing = (firstMissing...lastMissing).map { UInt16(truncatingIfNeeded: $0) } + guard lastMissing - firstMissing < UInt64(NvstRtcp.maximumNackEntries * 16) else { return 0 } + return requestRetransmission(of: Array(firstMissing...lastMissing)) + } + + /// Asks the seat to retransmit these packets, named in extended RTP sequence space. + @discardableResult + public func requestRetransmission(of indices: [UInt64]) -> Int { + guard let ssrc = receiver.feedbackCounters.boundSSRC, !indices.isEmpty else { return 0 } + let span = indices.count + let missing = indices.map { UInt16(truncatingIfNeeded: $0) } guard let nack = NvstRtcp.genericNack(senderSSRC: NvstVideoReceiver.clientSSRC, mediaSSRC: ssrc, missing: missing) else { return 0 } @@ -404,6 +411,8 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { let handler = onDrop callbackLock.unlock() handler?(reason) + case .retransmissionWanted(let indices): + requestRetransmission(of: indices) } } } diff --git a/GFN/NVST/BifrostFree/NvstNackTracker.swift b/GFN/NVST/BifrostFree/NvstNackTracker.swift new file mode 100644 index 00000000..6ce82c4b --- /dev/null +++ b/GFN/NVST/BifrostFree/NvstNackTracker.swift @@ -0,0 +1,67 @@ +// When to ask the seat to resend a missing video packet. +// + +import Foundation + +/// Schedules retransmission requests for missing video packets on the official client's timing, +/// as its log prints it: the first request 1 ms after a packet goes missing, then up to three +/// retries 4 ms apart, within the 52 ms it holds a frame for them. +struct NvstNackTracker { + static let initialDelayNanoseconds: UInt64 = 1_000_000 + static let retryIntervalNanoseconds: UInt64 = 4_000_000 + static let maximumRetries = 3 + static let maximumWaitNanoseconds: UInt64 = 52_000_000 + + private struct Request { + let firstMissedAt: UInt64 + var lastSentAt: UInt64? + var sendCount = 0 + } + + private var requests: [UInt64: Request] = [:] + + var isEmpty: Bool { requests.isEmpty } + + /// The missing indices to request now, recording that they were requested. + mutating func due(missing: [UInt64], now: UInt64) -> [UInt64] { + var due: [UInt64] = [] + for index in missing { + var request = requests[index] ?? Request(firstMissedAt: now) + let isDue: Bool + if let lastSentAt = request.lastSentAt { + isDue = request.sendCount <= Self.maximumRetries + && now &- lastSentAt >= Self.retryIntervalNanoseconds + } else { + isDue = now &- request.firstMissedAt >= Self.initialDelayNanoseconds + } + if isDue { + request.lastSentAt = now + request.sendCount += 1 + due.append(index) + } + requests[index] = request + } + return due + } + + /// Whether `index` was requested and is still inside the time its retransmission takes. + func isAwaitingRetransmission(of index: UInt64, now: UInt64) -> Bool { + guard let request = requests[index], request.sendCount > 0 else { return false } + return now &- request.firstMissedAt < Self.maximumWaitNanoseconds + } + + /// Forgets an index that arrived. True when it had been requested, so the arrival is a repair. + mutating func arrived(_ index: UInt64) -> Bool { + (requests.removeValue(forKey: index)?.sendCount ?? 0) > 0 + } + + /// Forgets every index below `index`: delivered, or given up on. + mutating func forget(below index: UInt64) { + guard !requests.isEmpty else { return } + requests = requests.filter { $0.key >= index } + } + + mutating func reset() { + requests.removeAll() + } +} diff --git a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift index ac5ac610..df179b5a 100644 --- a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift @@ -40,6 +40,9 @@ public enum NvstReceiveEvent: Sendable { /// telling the seat not to reference `frameIndex` (the frame the hole was in) — the latter /// without the keyframe's cost. case chainBroken(frameIndex: UInt32) + /// Packets missing from a gap that is still open, due a retransmission request now. The gap + /// keeps waiting for them; `recoveryNeeded` follows only if they never arrive. + case retransmissionWanted([UInt64]) } public struct NvstReceiverStats: Equatable, Sendable { @@ -53,6 +56,10 @@ public struct NvstReceiverStats: Equatable, Sendable { public var latePackets: UInt64 = 0 /// Packets received twice; the second copy is dropped. public var duplicatePackets: UInt64 = 0 + /// Packets the SRTP replay window rejected: already seen, or older than the window. + public var replayedPackets: UInt64 = 0 + /// Requested packets that arrived while their gap was still open. + public var retransmissionRepairedPackets: UInt64 = 0 /// Packets rebuilt by FEC recovery and injected back into the reorder path. public var recoveredPackets: UInt64 = 0 /// The largest single finalized-loss range, in packets. @@ -142,6 +149,10 @@ public final class NvstVideoReceiver: @unchecked Sendable { /// The armed bound. The packet repair window is ~140 ms at 5K's ~8,500 packets/s, so the flat /// bound cut a repairable gap off early — and the packet window alone freezes a calm scene. public static let fecRepairMaximumWaitCeilingNanoseconds: UInt64 = 250_000_000 + /// How far behind the newest packet a retransmission may still be accepted: the official + /// client's NACK queue length. RFC 3711's 64 is 7.5 ms at ~8,500 packets/s, less than one + /// round trip, so every resent packet was rejected as a replay. + public static let replayWindowPackets = 2048 public enum ReceiverError: LocalizedError, Equatable, Sendable { case unsupportedProfile(String) @@ -162,7 +173,10 @@ public final class NvstVideoReceiver: @unchecked Sendable { private let srtcpMasterSalt: Data private let reorderWindow: Int private let reassembler: NvstFrameReassembler - private var replay = SrtpReplayWindow() + private var replay = SrtpReplayWindow(size: NvstVideoReceiver.replayWindowPackets) + private var nackTracker = NvstNackTracker() + /// When missing packets were last scanned for; nil makes the next packet scan at once. + private var lastNackScanAt: UInt64? private var reorder: [UInt64: NvstRtpVideoPacket] = [:] private var nextIndex: UInt64? private var openGap: (index: UInt64, since: UInt64)? @@ -265,6 +279,8 @@ public final class NvstVideoReceiver: @unchecked Sendable { defer { lock.unlock() } reorder.removeAll() nextIndex = nil + nackTracker.reset() + lastNackScanAt = nil reassembler.reset() } @@ -330,6 +346,7 @@ public final class NvstVideoReceiver: @unchecked Sendable { stageNanoseconds.unprotect += DispatchTime.now().uptimeNanoseconds - stageStart } catch NvstRtpParseError.replayed { stats.droppedPackets += 1 + stats.replayedPackets += 1 return [.dropped(.replayed)] } catch { stats.droppedPackets += 1 @@ -650,8 +667,12 @@ extension NvstVideoReceiver { } if index > expected { stats.outOfOrderPackets += 1 - if openGap?.index != expected { openGap = (expected, uptimeNanoseconds()) } + if openGap?.index != expected { + openGap = (expected, uptimeNanoseconds()) + lastNackScanAt = nil + } } + if !nackTracker.isEmpty, nackTracker.arrived(index) { stats.retransmissionRepairedPackets += 1 } // Armed FEC holds a gap past the plain reorder window: a block's parity trails an early-frame // hole by up to `fecRepairReorderWindow` packets, so only the wall-clock bound ends it. let isPastFecWait = index > expected && isPastFecRepairWait() @@ -673,9 +694,25 @@ extension NvstVideoReceiver { nextIndex = cursor + 1 } if reorder.isEmpty { openGap = nil } + if let nextIndex { nackTracker.forget(below: nextIndex) } + requestMissingPackets(events: &events) return ready } + /// Asks for the packets still missing below the newest buffered one, at most once a + /// millisecond, so a lost packet can be resent before its gap gives up. + private func requestMissingPackets(events: inout [NvstReceiveEvent]) { + guard let expected = nextIndex, !reorder.isEmpty else { return } + let now = uptimeNanoseconds() + if let lastNackScanAt, now &- lastNackScanAt < NvstNackTracker.initialDelayNanoseconds { return } + guard let newest = reorder.keys.max(), newest > expected else { return } + lastNackScanAt = now + let limit = NvstRtcp.maximumNackEntries * 17 + let missing = (expected.. Bool { guard !isPastFecWait else { return true } guard depth >= UInt64(reorderWindow) else { return false } - return depth >= UInt64(Self.fecRepairReorderWindow) || !fecRecovery.snapshot.isArmed + if depth >= UInt64(Self.fecRepairReorderWindow) { return true } + return !fecRecovery.snapshot.isArmed && !isWithinRetransmissionWait() + } + + /// Whether the open gap's first packet was requested and may still be resent in time. + private func isWithinRetransmissionWait() -> Bool { + guard let expected = nextIndex else { return false } + return nackTracker.isAwaitingRetransmission(of: expected, now: uptimeNanoseconds()) } /// Whether the open gap has outlived its wall-clock chance of repair. diff --git a/GFN/NVST/BifrostFree/SrtpCryptography.swift b/GFN/NVST/BifrostFree/SrtpCryptography.swift index b12b0cb1..3a9684b2 100644 --- a/GFN/NVST/BifrostFree/SrtpCryptography.swift +++ b/GFN/NVST/BifrostFree/SrtpCryptography.swift @@ -293,10 +293,20 @@ public enum SrtpKeyDerivation { /// Replay window for the video stream (64 packets, matching the observed vendor behavior). public struct SrtpReplayWindow { + /// RFC 3711's minimum window, and what every stream used before video needed retransmissions. + public static let minimumSize = 64 + + /// How far behind the highest accepted index a packet may still arrive, rounded up to 64. + public let size: UInt64 private var highestIndex: UInt64? - private var seen: UInt64 = 0 + /// One bit per index, slotted by `index % size`. + private var seen: [UInt64] - public init() {} + public init(size: Int = SrtpReplayWindow.minimumSize) { + let words = (max(size, Self.minimumSize) + 63) / 64 + self.size = UInt64(words * 64) + seen = Array(repeating: 0, count: words) + } /// The extended index a sequence number most plausibly belongs to (RFC 3711 §3.3.1). /// @@ -327,27 +337,43 @@ public struct SrtpReplayWindow { public func wouldAccept(_ index: UInt64) -> Bool { guard let highest = highestIndex else { return true } if index > highest { return true } - let age = highest - index - return age < 64 && (seen & (1 << age)) == 0 + return highest - index < size && !isSeen(index) } public mutating func accept(_ index: UInt64) -> Bool { guard let highest = highestIndex else { highestIndex = index - seen = 1 // the highest itself is already seen (age 0) + setSeen(index, true) return true } if index > highest { - let delta = index - highest - // Every known index ages by delta; the previous highest becomes age delta. - seen = (delta >= 64) ? 0 : (seen << delta) | (1 << delta) + // The slots the window slides over held indices a full window older; they start unseen. + if index - highest >= size { + seen = Array(repeating: 0, count: seen.count) + } else { + for slid in (highest + 1)...index { setSeen(slid, false) } + } highestIndex = index - seen |= 1 + setSeen(index, true) return true } - let age = highest - index - if age >= 64 || (seen & (1 << age)) != 0 { return false } - seen |= (1 << age) + if highest - index >= size || isSeen(index) { return false } + setSeen(index, true) return true } + + private func isSeen(_ index: UInt64) -> Bool { + let slot = index % size + return seen[Int(slot / 64)] & (1 << (slot % 64)) != 0 + } + + private mutating func setSeen(_ index: UInt64, _ isSeen: Bool) { + let slot = index % size + let bit: UInt64 = 1 << (slot % 64) + if isSeen { + seen[Int(slot / 64)] |= bit + } else { + seen[Int(slot / 64)] &= ~bit + } + } } diff --git a/OPN/Stream/NvstBifrostFreeTransport.swift b/OPN/Stream/NvstBifrostFreeTransport.swift index 0964b41d..209088ab 100644 --- a/OPN/Stream/NvstBifrostFreeTransport.swift +++ b/OPN/Stream/NvstBifrostFreeTransport.swift @@ -475,7 +475,7 @@ public actor NvstBifrostFreeTransport: NativeNVSTTransport { logger?("NVST audio tracks=\(bundle?.remoteAudioTrackCount ?? 0) pktIn=\(audio?.packets ?? 0) bytesIn=\(audio?.bytes ?? 0)" + " samples=\(audio?.samples ?? 0) concealed=\(audio?.concealed ?? 0) discarded=\(audio?.discarded ?? 0) ssrc=\(audio?.ssrc.map(String.init) ?? "-")") let video = videoPipeline?.snapshot ?? NvstVideoPipeline.Counters() - logger?("NVST counters auth=\(stats.authenticatedPackets) fec=\(stats.fecPackets) dropped=\(stats.droppedPackets) rtpLoss=\(stats.finalizedLossPackets) parityLoss=\(stats.parityOnlyLossPackets) frames=\(stats.framesEmitted) keyframes=\(stats.keyframesEmitted) recoveries=\(stats.recoveries) sofFlagged=\(stats.startOfFrameFlagged) sofOk=\(stats.startOfFrameAccepted) abandoned=\(stats.abandonedFrames) rrFail=\(stats.receiverReportFailures)\(stats.lastReceiverReportFailure.map { " rrErr=\($0)" } ?? "") multiBlock=\(stats.multiBlockPackets) maxBlock=\(stats.highestFecLastBlock) decoded=\(decoder?.decodedFrameCount ?? 0) decodeFailed=\(decoder?.failedFrameCount ?? 0) decodeErr=\(decoder?.failureStatusSummary ?? "-") noParamSets=\(video.missingParameterSetFrames) idrOut=\(idrRequestsSent) invalidOut=\(invalidationsSent) inputOut=\(inputEventsSent) padOut=\(gamepadPacketsSent) padFail=\(gamepadSendFailures) padDropped=\(gamepadPacketsDroppedForUnannouncedPad) padReg=\(didRegisterGamepad) textTyped=\(textCharactersTyped) textDroppedBytes=\(textBytesDropped) inputReady=\(bundle?.isInputReady == true) rrOut=\(stats.receiverReportsSent) frac=\(stats.lastFractionLost) lost=\(stats.lastCumulativeLost) jitter=\(stats.lastJitter) seqSpan=\(stats.sequenceSpan) negFps=\(negotiatedFps.map(String.init) ?? "nil") mediaSeconds=\(String(format: "%.2f", Double(stats.lastRtpTimestamp &- (stats.firstRtpTimestamp ?? 0)) / Double(NvstVideoToolboxDecoder.clockRate))) fidxChanges=\(stats.frameIndexChanges) \(perSecondSeries(stats: stats, included: includingPerSecondSeries)) paceOut=\(video.pacingReportsSent) paceFail=\(video.pacingReportFailures) ackOut=\(video.frameAcksSent) ackFail=\(video.frameAckFailures) qosOut=\(qosReportsSent) qosFail=\(qosReportFailures) rtpStatsOut=\(rtpStatsReportsSent) ccStatsOut=\(controlStatsReportsSent) ssrc=\(stats.boundSSRC.map { String(format: "0x%08x", $0) } ?? "-")") + logger?("NVST counters auth=\(stats.authenticatedPackets) fec=\(stats.fecPackets) dropped=\(stats.droppedPackets) replayed=\(stats.replayedPackets) late=\(stats.latePackets) dup=\(stats.duplicatePackets) nackRepaired=\(stats.retransmissionRepairedPackets) rtpLoss=\(stats.finalizedLossPackets) parityLoss=\(stats.parityOnlyLossPackets) frames=\(stats.framesEmitted) keyframes=\(stats.keyframesEmitted) recoveries=\(stats.recoveries) sofFlagged=\(stats.startOfFrameFlagged) sofOk=\(stats.startOfFrameAccepted) abandoned=\(stats.abandonedFrames) rrFail=\(stats.receiverReportFailures)\(stats.lastReceiverReportFailure.map { " rrErr=\($0)" } ?? "") multiBlock=\(stats.multiBlockPackets) maxBlock=\(stats.highestFecLastBlock) decoded=\(decoder?.decodedFrameCount ?? 0) decodeFailed=\(decoder?.failedFrameCount ?? 0) decodeErr=\(decoder?.failureStatusSummary ?? "-") noParamSets=\(video.missingParameterSetFrames) idrOut=\(idrRequestsSent) invalidOut=\(invalidationsSent) inputOut=\(inputEventsSent) padOut=\(gamepadPacketsSent) padFail=\(gamepadSendFailures) padDropped=\(gamepadPacketsDroppedForUnannouncedPad) padReg=\(didRegisterGamepad) textTyped=\(textCharactersTyped) textDroppedBytes=\(textBytesDropped) inputReady=\(bundle?.isInputReady == true) rrOut=\(stats.receiverReportsSent) frac=\(stats.lastFractionLost) lost=\(stats.lastCumulativeLost) jitter=\(stats.lastJitter) seqSpan=\(stats.sequenceSpan) negFps=\(negotiatedFps.map(String.init) ?? "nil") mediaSeconds=\(String(format: "%.2f", Double(stats.lastRtpTimestamp &- (stats.firstRtpTimestamp ?? 0)) / Double(NvstVideoToolboxDecoder.clockRate))) fidxChanges=\(stats.frameIndexChanges) \(perSecondSeries(stats: stats, included: includingPerSecondSeries)) paceOut=\(video.pacingReportsSent) paceFail=\(video.pacingReportFailures) ackOut=\(video.frameAcksSent) ackFail=\(video.frameAckFailures) qosOut=\(qosReportsSent) qosFail=\(qosReportFailures) rtpStatsOut=\(rtpStatsReportsSent) ccStatsOut=\(controlStatsReportsSent) ssrc=\(stats.boundSSRC.map { String(format: "0x%08x", $0) } ?? "-")") // Which pipeline stage a latency spike lives in. `peak*` are per-stage session maxima, so a // single 500 ms stall is still visible after the average has recovered. logger?(String(format: "NVST frame stages slow=%d frames=%llu resyncs=%d skipped=%d abandoned=%d lastLatency=%.1fms inputSendTotal=%.0fms inputSendPeak=%.1fms", diff --git a/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift b/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift index ffccea50..c4c34f9d 100644 --- a/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift +++ b/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift @@ -128,6 +128,62 @@ struct NvstMjolnirReceiverTests { #expect(stats.droppedPackets == 0) } + /// A lost packet is requested while its gap is open, and the gap waits for the resend instead + /// of becoming loss once the reorder window passes it, so the frames behind it are not lost. + @Test func aRequestedPacketThatIsResentInTimeRepairsTheGap() throws { + let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) + let clock = OSAllocatedUnfairLock(initialState: UInt64(0)) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { clock.withLock { $0 } }) + let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] + func feed(_ sequence: UInt16) throws -> [NvstReceiveEvent] { + receiver.process(datagram: try NvstReceiverFixtures.seal( + NvstReceiverFixtures.packet(sequence: sequence, frameIndex: UInt32(sequence), flags: 0x07, media: media), + sequence: sequence, handoff: handoff)) + } + func requested(_ events: [NvstReceiveEvent]) -> [UInt64] { + events.flatMap { event -> [UInt64] in + if case .retransmissionWanted(let indices) = event { return indices } else { return [] } + } + } + _ = try feed(1) + #expect(requested(try feed(3)).isEmpty) + clock.withLock { $0 = NvstNackTracker.initialDelayNanoseconds } + #expect(requested(try feed(4)) == [2]) + // Well past the four-packet window and the 64-packet RFC 3711 replay window. + for sequence in UInt16(5)...UInt16(80) { + #expect(NvstReceiverFixtures.recoveries(try feed(sequence)) == 0) + } + let repaired = try feed(2) + #expect(NvstReceiverFixtures.frames(repaired).count == 79) + #expect(NvstReceiverFixtures.recoveries(repaired) == 0) + let stats = receiver.snapshot + #expect(stats.retransmissionRepairedPackets == 1) + #expect(stats.finalizedLossPackets == 0) + #expect(stats.replayedPackets == 0) + } + + /// A resend that never comes ends the wait: the gap is loss once its request has expired. + @Test func aRequestedPacketThatNeverArrivesBecomesLossAfterTheWait() throws { + let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) + let clock = OSAllocatedUnfairLock(initialState: UInt64(0)) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { clock.withLock { $0 } }) + let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] + func feed(_ sequence: UInt16) throws -> [NvstReceiveEvent] { + receiver.process(datagram: try NvstReceiverFixtures.seal( + NvstReceiverFixtures.packet(sequence: sequence, frameIndex: UInt32(sequence), flags: 0x07, media: media), + sequence: sequence, handoff: handoff)) + } + _ = try feed(1) + _ = try feed(3) + clock.withLock { $0 = NvstNackTracker.initialDelayNanoseconds } + for sequence in UInt16(4)...UInt16(8) { + #expect(NvstReceiverFixtures.recoveries(try feed(sequence)) == 0) + } + clock.withLock { $0 = NvstNackTracker.maximumWaitNanoseconds } + #expect(NvstReceiverFixtures.recoveries(try feed(9)) == 1) + #expect(receiver.snapshot.finalizedLossPackets == 1) + } + /// A second gap among the survivors finalizes on its own later packet rather than being /// silently absorbed by the first flush. @Test func interleavedGapsFinalizeIndependently() throws { diff --git a/Tests/GFN/NVST/NvstNackTrackerTests.swift b/Tests/GFN/NVST/NvstNackTrackerTests.swift new file mode 100644 index 00000000..94d150c8 --- /dev/null +++ b/Tests/GFN/NVST/NvstNackTrackerTests.swift @@ -0,0 +1,42 @@ +import Foundation +import Testing +@testable import OpenNOW + +struct NvstNackTrackerTests { + @Test func aMissingPacketIsRequestedAfterTheInitialDelay() { + var tracker = NvstNackTracker() + #expect(tracker.due(missing: [7], now: 0).isEmpty) + #expect(tracker.due(missing: [7], now: NvstNackTracker.initialDelayNanoseconds - 1).isEmpty) + #expect(tracker.due(missing: [7], now: NvstNackTracker.initialDelayNanoseconds) == [7]) + } + + @Test func aRequestIsRetriedThreeTimesAtTheRetryInterval() { + var tracker = NvstNackTracker() + _ = tracker.due(missing: [7], now: 0) + var sendTimes: [UInt64] = [] + var now = NvstNackTracker.initialDelayNanoseconds + while now <= NvstNackTracker.maximumWaitNanoseconds { + if !tracker.due(missing: [7], now: now).isEmpty { sendTimes.append(now) } + now += 1_000_000 + } + #expect(sendTimes.count == 1 + NvstNackTracker.maximumRetries) + #expect(zip(sendTimes, sendTimes.dropFirst()).allSatisfy { $1 - $0 == NvstNackTracker.retryIntervalNanoseconds }) + } + + @Test func aRequestedArrivalCountsAsARepair() { + var tracker = NvstNackTracker() + _ = tracker.due(missing: [7, 8], now: 0) + _ = tracker.due(missing: [7, 8], now: NvstNackTracker.initialDelayNanoseconds) + let notRequested = NvstNackTracker() + var unrequested = notRequested + _ = unrequested.due(missing: [9], now: 0) + let repaired = tracker.arrived(7) + let notRepaired = unrequested.arrived(9) + #expect(repaired) + #expect(!notRepaired) + #expect(tracker.isAwaitingRetransmission(of: 8, now: NvstNackTracker.initialDelayNanoseconds)) + #expect(!tracker.isAwaitingRetransmission(of: 8, now: NvstNackTracker.maximumWaitNanoseconds)) + tracker.forget(below: 9) + #expect(tracker.isEmpty) + } +} diff --git a/Tests/GFN/NVST/SrtpCryptographyTests.swift b/Tests/GFN/NVST/SrtpCryptographyTests.swift index fe127c5b..3568d35c 100644 --- a/Tests/GFN/NVST/SrtpCryptographyTests.swift +++ b/Tests/GFN/NVST/SrtpCryptographyTests.swift @@ -136,4 +136,32 @@ struct SrtpCryptographyTests { #expect(forwardAccepted) #expect(forward.estimatedIndex(for: 0) == 0x0001_0000) } + + /// A retransmission returns a round trip after the loss, thousands of packets behind the + /// newest at video rates; a window that small rejected every one of them as a replay. + @Test func aWideReplayWindowAcceptsALateRetransmissionOnce() { + var window = SrtpReplayWindow(size: 2048) + #expect(window.size == 2048) + let newest = window.accept(5_000) + let resent = window.accept(4_000) + let again = window.accept(4_000) + let tooOld = window.accept(5_000 - 2048) + #expect(newest) + #expect(resent) + #expect(!again) + #expect(!tooOld) + let slid = window.accept(6_000) + let behindSlid = window.accept(4_100) + #expect(slid) + #expect(behindSlid) + } + + @Test func theDefaultReplayWindowKeepsRfc3711sSixtyFour() { + var window = SrtpReplayWindow() + #expect(window.size == 64) + let newest = window.accept(200) + #expect(newest) + #expect(window.wouldAccept(137)) + #expect(!window.wouldAccept(136)) + } } From c43f827fefce91a2c5638506e0f583351dae7c2d Mon Sep 17 00:00:00 2001 From: Olivier Giroux Date: Fri, 2 Oct 2026 22:26:47 -0400 Subject: [PATCH 2/6] fix: send retransmission requests as the version 2 control command The seat ignored the RTCP NACKs: 173 requests for 564 packets repaired 4. The official client announces rtpNackVersion 2 and, for that version, sends control command 0x317 (NvscClientPipeline:: createAndSendNackRequest -> ServerControl::sendRtpNackRequest) rather than an RTCP NACK. NvstRtpNackRequest builds its payload as RtpSourceQueueExtV2::createNackRequest does: version, stream index and entry count, then per entry a little-endian u16 sequence number and a u64 mask of the following 64 packets. Requests go out on the partially reliable control stream through the video pipeline's bundle; the RTCP NACK stays as the fallback before the bundle is up. Receiver loss tests run on a fixed clock, so a slow test machine no longer turns their gaps into retransmission waits. --- .../BifrostFree/NvstMjolnirReceiver.swift | 17 +++++- GFN/NVST/BifrostFree/NvstRtpNackRequest.swift | 59 +++++++++++++++++++ OPN/Stream/NvstBifrostFreeVideo.swift | 3 + OPN/Stream/NvstVideoPipeline.swift | 11 ++++ Tests/GFN/NVST/NvstMjolnirFeedbackTests.swift | 12 ++-- Tests/GFN/NVST/NvstMjolnirReceiverTests.swift | 26 ++++---- Tests/GFN/NVST/NvstRtpNackRequestTests.swift | 32 ++++++++++ 7 files changed, 140 insertions(+), 20 deletions(-) create mode 100644 GFN/NVST/BifrostFree/NvstRtpNackRequest.swift create mode 100644 Tests/GFN/NVST/NvstRtpNackRequestTests.swift diff --git a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift index d233fde5..4c4bfe43 100644 --- a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift @@ -65,6 +65,9 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { /// knows it (a hole inside a frame), nil for a finalized RTP loss whose frame is unknown. public var onRecoveryNeeded: (@Sendable (UInt32?) -> Void)? public var onDrop: (@Sendable (NvstReceiveDrop) -> Void)? + /// Sends a retransmission request on the control channel and returns whether it went out. + /// Unset, requests fall back to an RTCP NACK on this socket. + public var onRetransmissionWanted: (@Sendable ([UInt16]) -> Bool)? /// Diagnostics that would otherwise be invisible, such as a feedback timer producing nothing. public var onDiagnostic: (@Sendable (String) -> Void)? @@ -412,7 +415,19 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { callbackLock.unlock() handler?(reason) case .retransmissionWanted(let indices): - requestRetransmission(of: indices) + callbackLock.lock() + let handler = onRetransmissionWanted + callbackLock.unlock() + guard let handler else { + requestRetransmission(of: indices) + continue + } + let sequenceNumbers = indices.prefix(NvstRtpNackRequest.maximumSequenceNumbers).map { UInt16(truncatingIfNeeded: $0) } + guard handler(sequenceNumbers) else { continue } + counterLock.lock() + nacksSent += 1 + nackedPackets += sequenceNumbers.count + counterLock.unlock() } } } diff --git a/GFN/NVST/BifrostFree/NvstRtpNackRequest.swift b/GFN/NVST/BifrostFree/NvstRtpNackRequest.swift new file mode 100644 index 00000000..81e3269b --- /dev/null +++ b/GFN/NVST/BifrostFree/NvstRtpNackRequest.swift @@ -0,0 +1,59 @@ +// The version 2 retransmission request the seat answers. +// + +import Foundation + +/// Control command `0x317`, which the official client sends for `rtpNackVersion` 2 +/// (`NvscClientPipeline::createAndSendNackRequest` → `ServerControl::sendRtpNackRequest`) instead +/// of an RTCP NACK. Layout from `RtpSourceQueueExtV2::createNackRequest`: version, stream index +/// and entry count, one byte each, then per entry a little-endian u16 sequence number and a +/// little-endian u64 whose bit n names sequence + n + 1. At most 64 sequence numbers per request. +public struct NvstRtpNackRequest: Equatable, Sendable { + public static let version: UInt8 = 2 + public static let maximumSequenceNumbers = 64 + + public struct Entry: Equatable, Sendable { + public let sequenceNumber: UInt16 + public let followingMask: UInt64 + } + + public let streamIndex: UInt8 + public let entries: [Entry] + + /// `sequenceNumbers` in RTP order; only the first `maximumSequenceNumbers` are named. + public init(streamIndex: UInt8 = 0, sequenceNumbers: [UInt16]) { + self.streamIndex = streamIndex + let named = sequenceNumbers.prefix(Self.maximumSequenceNumbers) + var entries: [Entry] = [] + var index = named.startIndex + while index < named.endIndex { + let base = named[index] + var mask: UInt64 = 0 + index += 1 + while index < named.endIndex { + let offset = named[index] &- base + guard offset >= 1, offset <= 64 else { break } + mask |= 1 << UInt64(offset - 1) + index += 1 + } + entries.append(Entry(sequenceNumber: base, followingMask: mask)) + } + self.entries = entries + } + + public var payload: Data { + var writer = NvstByteWriter(capacity: 3 + entries.count * 10) + writer.u8(Self.version) + writer.u8(streamIndex) + writer.u8(UInt8(entries.count)) + for entry in entries { + writer.u16LE(entry.sequenceNumber) + writer.u64LE(entry.followingMask) + } + return writer.data + } + + public var command: NvstControlCommand { + NvstControlCommand(code: .rtpNackRequest, payload: payload) + } +} diff --git a/OPN/Stream/NvstBifrostFreeVideo.swift b/OPN/Stream/NvstBifrostFreeVideo.swift index f9c6b4fc..31b82ec7 100644 --- a/OPN/Stream/NvstBifrostFreeVideo.swift +++ b/OPN/Stream/NvstBifrostFreeVideo.swift @@ -102,6 +102,9 @@ extension NvstBifrostFreeTransport { let pipeline = makeVideoPipeline(handoff: handoff, decoder: decoder, receiver: receiver, mediaContinuation: mediaContinuation) videoPipeline = pipeline receiver.onAccessUnit = { [weak pipeline] unit in pipeline?.submit(unit) } + receiver.onRetransmissionWanted = { [weak pipeline] sequenceNumbers in + pipeline?.requestRetransmission(of: sequenceNumbers) ?? false + } receiver.onRecoveryNeeded = { [weak self, weak receiver] brokenFrameIndex in // Reached only for gaps too wide to repair by retransmission; the receiver NACKs the // rest itself. A broken reference chain then only recovers with a fresh keyframe. diff --git a/OPN/Stream/NvstVideoPipeline.swift b/OPN/Stream/NvstVideoPipeline.swift index 5387a6bb..dbeaf817 100644 --- a/OPN/Stream/NvstVideoPipeline.swift +++ b/OPN/Stream/NvstVideoPipeline.swift @@ -675,3 +675,14 @@ extension NvstVideoPipeline { lock.unlock() } } + +extension NvstVideoPipeline { + /// Asks the seat to resend these video packets. False until the bundle is up. + public func requestRetransmission(of sequenceNumbers: [UInt16]) -> Bool { + lock.lock() + let channel = bundle + lock.unlock() + guard let channel else { return false } + return channel.sendPartiallyReliableControl(NvstRtpNackRequest(sequenceNumbers: sequenceNumbers).command) + } +} diff --git a/Tests/GFN/NVST/NvstMjolnirFeedbackTests.swift b/Tests/GFN/NVST/NvstMjolnirFeedbackTests.swift index f8f460fa..6192ba2c 100644 --- a/Tests/GFN/NVST/NvstMjolnirFeedbackTests.swift +++ b/Tests/GFN/NVST/NvstMjolnirFeedbackTests.swift @@ -20,7 +20,7 @@ struct NvstMjolnirFeedbackTests { @Test func aGapBeyondTheReorderWindowRequestsRecoveryAndNeverEmitsAPartialFrame() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let sof = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: 1, frameIndex: 1, flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65]), sequence: 1, handoff: handoff) #expect(NvstReceiverFixtures.frames(receiver.process(datagram: sof)).isEmpty) // Jump far past the window: the reference chain is broken. @@ -60,7 +60,7 @@ struct NvstMjolnirFeedbackTests { /// producing nothing looks like a slow stream rather than a dead feedback plane. @Test func aReceiverReportIsProducedOnceAnSsrcIsBound() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) // No SSRC bound yet, so nothing to report about. #expect(receiver.pollReceiverReport() == nil) @@ -81,7 +81,7 @@ struct NvstMjolnirFeedbackTests { /// The sender's rate control reads these, so they have to be real. @Test func receiverReportsCarryRealReceptionStatistics() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 8) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) for sequence in UInt16(1)...UInt16(6) { let datagram = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: sequence, frameIndex: UInt32(sequence), flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65]), @@ -105,7 +105,7 @@ struct NvstMjolnirFeedbackTests { @Test func outOfOrderPacketsAreDeliveredInSequence() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 8) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let sof = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: 1, frameIndex: 3, flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65]), sequence: 1, handoff: handoff) let eof = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: 3, frameIndex: 3, flags: 0x03, media: [0xcc]), sequence: 3, handoff: handoff) let middle = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: 2, frameIndex: 3, flags: 0x01, media: [0xbb]), sequence: 2, handoff: handoff) @@ -120,7 +120,7 @@ struct NvstMjolnirFeedbackTests { @Test func fecRepairPacketsAreCountedNotDecoded() throws { let handoff = NvstReceiverFixtures.makeHandoff() - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let repair = NvstVideoPacketTests.buildPacket(sequence: 1, frameIndex: 1, flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65], fecWord: 0x00c0_3420) let events = receiver.process(datagram: try NvstReceiverFixtures.seal(repair, sequence: 1, handoff: handoff)) #expect(NvstReceiverFixtures.frames(events).isEmpty) @@ -129,7 +129,7 @@ struct NvstMjolnirFeedbackTests { @Test func srtcpReceiverReportsWaitForTheMediaSsrcThenRespectTheInterval() throws { let handoff = NvstReceiverFixtures.makeHandoff() - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let origin = Date(timeIntervalSince1970: 1_000_000) // No authenticated packet yet: nothing to report about. #expect(receiver.pollReceiverReport(now: origin) == nil) diff --git a/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift b/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift index c4c34f9d..4f21d541 100644 --- a/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift +++ b/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift @@ -7,7 +7,7 @@ import Testing struct NvstMjolnirReceiverTests { @Test func authenticatedPacketsReassembleIntoAnAnnexBAccessUnit() throws { let handoff = NvstReceiverFixtures.makeHandoff() - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let sof = NvstReceiverFixtures.packet(sequence: 1, frameIndex: 42, flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65]) let eof = NvstReceiverFixtures.packet(sequence: 2, frameIndex: 42, flags: 0x03, media: [0x00, 0x00, 0x01, 0x41]) @@ -27,7 +27,7 @@ struct NvstMjolnirReceiverTests { // The seat advertises AEAD_AES_256_GCM explicitly on some builds; the tag is 16 bytes // there, not NVIDIA's default 8. let handoff = NvstReceiverFixtures.makeHandoff(profile: .aeadAes256Gcm) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let unit = NvstReceiverFixtures.packet(sequence: 7, frameIndex: 1, flags: 0x07, media: [0x00, 0x00, 0x00, 0x01, 0x65]) let emitted = NvstReceiverFixtures.frames(receiver.process(datagram: try NvstReceiverFixtures.seal(unit, sequence: 7, handoff: handoff))) #expect(emitted.count == 1) @@ -41,7 +41,7 @@ struct NvstMjolnirReceiverTests { @Test func tamperedAndReplayedPacketsAreDropped() throws { let handoff = NvstReceiverFixtures.makeHandoff() - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let sof = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: 5, frameIndex: 9, flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65]), sequence: 5, handoff: handoff) #expect(NvstReceiverFixtures.frames(receiver.process(datagram: sof)).isEmpty) // Same packet again: the replay window must reject it. @@ -54,7 +54,7 @@ struct NvstMjolnirReceiverTests { @Test func aPacketFromAForeignSsrcNeverJoinsTheStream() throws { let handoff = NvstReceiverFixtures.makeHandoff() - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) var foreign = NvstReceiverFixtures.packet(sequence: 1, frameIndex: 1, flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65]) foreign.replaceSubrange(8..<12, with: Data([0xde, 0xad, 0xbe, 0xef])) // Sealed under its own SSRC so authentication passes; only the handoff disagrees. @@ -72,7 +72,7 @@ struct NvstMjolnirReceiverTests { /// socket terminated the app. The attacker needs the port, not the key. @Test func anUnauthenticatedPacketFarAheadOfTheWindowIsDroppedRatherThanFatal() throws { let handoff = NvstReceiverFixtures.makeHandoff() - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] let first = try NvstReceiverFixtures.seal( NvstReceiverFixtures.packet(sequence: 100, frameIndex: 1, flags: 0x07, media: media), @@ -105,7 +105,7 @@ struct NvstMjolnirReceiverTests { /// frames they carried — uncounted, because the destroyed packets hit no drop counter. @Test func aSingleLostPacketDeliversTheBufferedPacketsBehindIt() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] // A complete single-packet frame per sequence number. Sequence 2 never arrives. let first = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: 1, frameIndex: 1, flags: 0x07, media: media), sequence: 1, handoff: handoff) @@ -188,7 +188,7 @@ struct NvstMjolnirReceiverTests { /// silently absorbed by the first flush. @Test func interleavedGapsFinalizeIndependently() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] func feed(_ sequence: UInt16) throws -> [NvstReceiveEvent] { receiver.process(datagram: try NvstReceiverFixtures.seal( @@ -213,7 +213,7 @@ struct NvstMjolnirReceiverTests { /// it must ask for a keyframe, not a retransmission of an arbitrary sequence number. @Test func aStreamSequenceHoleAsksForAKeyframeNotARetransmission() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let sof = try NvstReceiverFixtures.seal(NvstReceiverFixtures.packet(sequence: 1, frameIndex: 7, flags: 0x05, media: [0x00, 0x00, 0x00, 0x01, 0x65]), sequence: 1, handoff: handoff) #expect(NvstReceiverFixtures.frames(receiver.process(datagram: sof)).isEmpty) @@ -236,7 +236,7 @@ struct NvstMjolnirReceiverTests { /// destroyed wholesale. @Test func oneLostPacketInsideAFrameCostsExactlyThatFrame() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 32) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let packetsPerFrame: UInt16 = 50 let droppedSequence: UInt16 = 75 // mid-frame 2 var emitted = 0 @@ -269,7 +269,7 @@ struct NvstMjolnirReceiverTests { /// no compounding, no stuck reassembler, no double-counted recovery. @Test func isolatedLossesAcrossAStreamDoNotCompound() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 32) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) let packetsPerFrame: UInt16 = 30 let frameCount: UInt16 = 20 let dropped: Set = [100, 250, 400] // frames 4, 9, 14 — all far from the stream tail @@ -298,7 +298,7 @@ struct NvstMjolnirReceiverTests { /// the seat never learns anything was lost. @Test func aLostPacketIsRepairedFromParityAndTheFrameStillEmits() throws { let handoff = NvstReceiverFixtures.makeHandoff() - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) func fecWord(index: UInt32) -> UInt32 { (50 << 4) | (index << 12) | (2 << 22) } // One frame per FEC block: two sources (SOF, EOF) and one parity shard (50% of 2). @@ -353,7 +353,7 @@ struct NvstMjolnirReceiverTests { /// race and turn a repairable hole into a lost frame plus a keyframe round trip. @Test func anOpenGapWaitsForFecRepairBeforeBecomingLoss() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) func fecWord(index: UInt32) -> UInt32 { (50 << 4) | (index << 12) | (2 << 22) } func blockPackets(frameIndex: UInt32, baseSequence: UInt16, seed: UInt8) -> [(UInt16, Data)] { @@ -420,7 +420,7 @@ struct NvstMjolnirReceiverTests { /// frame is delivered at once instead of waiting on the gap and then asking for a keyframe. @Test func aLostParityPacketOfACompleteBlockIsSteppedOver() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) - let receiver = try NvstVideoReceiver(handoff: handoff) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { 0 }) func fecWord(index: UInt32) -> UInt32 { (50 << 4) | (index << 12) | (2 << 22) } func blockPackets(frameIndex: UInt32, baseSequence: UInt16, seed: UInt8) -> [(UInt16, Data)] { diff --git a/Tests/GFN/NVST/NvstRtpNackRequestTests.swift b/Tests/GFN/NVST/NvstRtpNackRequestTests.swift new file mode 100644 index 00000000..cc59e065 --- /dev/null +++ b/Tests/GFN/NVST/NvstRtpNackRequestTests.swift @@ -0,0 +1,32 @@ +import Foundation +import Testing +@testable import OpenNOW + +/// Layout from `RtpSourceQueueExtV2::createNackRequest` in the official client. +struct NvstRtpNackRequestTests { + @Test func packetsWithin64OfAnEntryRideInItsMask() { + let request = NvstRtpNackRequest(sequenceNumbers: [100, 101, 103, 200]) + #expect(request.entries == [ + NvstRtpNackRequest.Entry(sequenceNumber: 100, followingMask: 0b101), + NvstRtpNackRequest.Entry(sequenceNumber: 200, followingMask: 0), + ]) + #expect([UInt8](request.payload) == [0x02, 0x00, 0x02, + 0x64, 0x00, 0x05, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, + 0xc8, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00]) + #expect(request.command.code == .rtpNackRequest) + #expect(NvstControlCommandCode.rtpNackRequest.rawValue == 0x0317) + } + + @Test func theMaskSpansTheSequenceWrapAndEndsAtSixtyFour() { + let request = NvstRtpNackRequest(sequenceNumbers: [0xffff, 0x0000, 0xffff &+ 64, 0xffff &+ 65]) + #expect(request.entries == [ + NvstRtpNackRequest.Entry(sequenceNumber: 0xffff, followingMask: 1 | (1 << 63)), + NvstRtpNackRequest.Entry(sequenceNumber: 0xffff &+ 65, followingMask: 0), + ]) + } + + @Test func atMostSixtyFourSequenceNumbersAreNamed() { + let request = NvstRtpNackRequest(sequenceNumbers: (0..<100).map { UInt16($0 * 100) }) + #expect(request.entries.count == NvstRtpNackRequest.maximumSequenceNumbers) + } +} From 24832df59f0db44dfd14352a7f951e6ee98a937b Mon Sep 17 00:00:00 2001 From: Olivier Giroux Date: Sat, 3 Oct 2026 05:20:30 -0400 Subject: [PATCH 3/6] fix: wait a round trip before retrying a retransmission request A retry waited 4 ms against a ~15 ms round trip, so a lost packet was requested up to four times before the first answer could arrive; in a two-second loss burst ~720 resent copies were dropped as duplicates. With useRtdForRtpNackToggle the official client waits the round trip plus 4 ms, and the receiver now does the same, from the control connection's measured round trip. The counters line adds nackRetries. --- GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift | 5 +++++ GFN/NVST/BifrostFree/NvstNackTracker.swift | 17 ++++++++++++++--- GFN/NVST/BifrostFree/NvstVideoReceiver.swift | 11 +++++++++++ OPN/Stream/NvstBifrostFreeTransport.swift | 5 ++++- Tests/GFN/NVST/NvstNackTrackerTests.swift | 15 ++++++++++++++- 5 files changed, 48 insertions(+), 5 deletions(-) diff --git a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift index 4c4bfe43..12a467d8 100644 --- a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift @@ -214,6 +214,11 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { } /// Asks the peer for a fresh keyframe on this socket's SRTCP path. + /// The measured round trip retransmission retries wait for. + public func useRetransmissionRoundTrip(milliseconds: Double) { + receiver.useRetransmissionRoundTrip(milliseconds: milliseconds) + } + public func requestKeyframe() { guard let ssrc = receiver.feedbackCounters.boundSSRC else { return } let pli = NvstRtcp.pictureLossIndication(senderSSRC: NvstVideoReceiver.clientSSRC, mediaSSRC: ssrc) diff --git a/GFN/NVST/BifrostFree/NvstNackTracker.swift b/GFN/NVST/BifrostFree/NvstNackTracker.swift index 6ce82c4b..6318e563 100644 --- a/GFN/NVST/BifrostFree/NvstNackTracker.swift +++ b/GFN/NVST/BifrostFree/NvstNackTracker.swift @@ -5,13 +5,19 @@ import Foundation /// Schedules retransmission requests for missing video packets on the official client's timing, /// as its log prints it: the first request 1 ms after a packet goes missing, then up to three -/// retries 4 ms apart, within the 52 ms it holds a frame for them. +/// retries, each one round trip plus 4 ms after the last (`useRtdForRtpNackToggle`), within the +/// 52 ms it holds a frame for them. A retry sooner than the round trip asks for a packet that is +/// already on its way, and the seat sends it again. struct NvstNackTracker { static let initialDelayNanoseconds: UInt64 = 1_000_000 - static let retryIntervalNanoseconds: UInt64 = 4_000_000 + static let extraRetryWaitNanoseconds: UInt64 = 4_000_000 static let maximumRetries = 3 static let maximumWaitNanoseconds: UInt64 = 52_000_000 + /// The wait before a retry: the extra wait alone until a round trip has been measured. + private(set) var retryIntervalNanoseconds = NvstNackTracker.extraRetryWaitNanoseconds + private(set) var retryCount = 0 + private struct Request { let firstMissedAt: UInt64 var lastSentAt: UInt64? @@ -22,6 +28,10 @@ struct NvstNackTracker { var isEmpty: Bool { requests.isEmpty } + mutating func useRoundTrip(nanoseconds: UInt64) { + retryIntervalNanoseconds = nanoseconds + Self.extraRetryWaitNanoseconds + } + /// The missing indices to request now, recording that they were requested. mutating func due(missing: [UInt64], now: UInt64) -> [UInt64] { var due: [UInt64] = [] @@ -30,11 +40,12 @@ struct NvstNackTracker { let isDue: Bool if let lastSentAt = request.lastSentAt { isDue = request.sendCount <= Self.maximumRetries - && now &- lastSentAt >= Self.retryIntervalNanoseconds + && now &- lastSentAt >= retryIntervalNanoseconds } else { isDue = now &- request.firstMissedAt >= Self.initialDelayNanoseconds } if isDue { + if request.sendCount > 0 { retryCount += 1 } request.lastSentAt = now request.sendCount += 1 due.append(index) diff --git a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift index df179b5a..aac16d7e 100644 --- a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift @@ -60,6 +60,8 @@ public struct NvstReceiverStats: Equatable, Sendable { public var replayedPackets: UInt64 = 0 /// Requested packets that arrived while their gap was still open. public var retransmissionRepairedPackets: UInt64 = 0 + /// Requests repeated for a packet already asked for. + public var retransmissionRetries: UInt64 = 0 /// Packets rebuilt by FEC recovery and injected back into the reorder path. public var recoveredPackets: UInt64 = 0 /// The largest single finalized-loss range, in packets. @@ -253,6 +255,14 @@ public final class NvstVideoReceiver: @unchecked Sendable { public var snapshot: NvstReceiverStats { lock.lock(); defer { lock.unlock() }; return stats } + /// The measured round trip retransmission retries wait for. + public func useRetransmissionRoundTrip(milliseconds: Double) { + guard milliseconds > 0, milliseconds.isFinite else { return } + lock.lock() + nackTracker.useRoundTrip(nanoseconds: UInt64(milliseconds * 1_000_000)) + lock.unlock() + } + /// Just the counters the feedback reports need. /// /// `stats` carries per-second arrays that are appended to on every frame, so snapshotting the @@ -710,6 +720,7 @@ extension NvstVideoReceiver { let limit = NvstRtcp.maximumNackEntries * 17 let missing = (expected.. Date: Sun, 4 Oct 2026 21:28:15 -0400 Subject: [PATCH 4/6] fix: request only the packets one retransmission request can name A scan could collect up to 136 missing packets and the tracker counted all of them as requested, but a 0x317 request names at most 64. The rest were never asked for, and their gap still waited up to 52 ms for them. The scan now stops at 64. A request the control channel could not take, before the bundle is up, goes out as the RTCP NACK instead of being counted as sent. --- .../BifrostFree/NvstMjolnirReceiver.swift | 9 ++++++-- GFN/NVST/BifrostFree/NvstVideoReceiver.swift | 3 ++- Tests/GFN/NVST/NvstMjolnirReceiverTests.swift | 21 +++++++++++++++++++ 3 files changed, 30 insertions(+), 3 deletions(-) diff --git a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift index 12a467d8..7f05516a 100644 --- a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift @@ -427,8 +427,13 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { requestRetransmission(of: indices) continue } - let sequenceNumbers = indices.prefix(NvstRtpNackRequest.maximumSequenceNumbers).map { UInt16(truncatingIfNeeded: $0) } - guard handler(sequenceNumbers) else { continue } + let sequenceNumbers = indices.map { UInt16(truncatingIfNeeded: $0) } + // The tracker already counts these as requested, so a request the control channel + // could not take still goes out, as the RTCP NACK. + guard handler(sequenceNumbers) else { + requestRetransmission(of: indices) + continue + } counterLock.lock() nacksSent += 1 nackedPackets += sequenceNumbers.count diff --git a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift index aac16d7e..f3b53e2f 100644 --- a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift @@ -717,7 +717,8 @@ extension NvstVideoReceiver { if let lastNackScanAt, now &- lastNackScanAt < NvstNackTracker.initialDelayNanoseconds { return } guard let newest = reorder.keys.max(), newest > expected else { return } lastNackScanAt = now - let limit = NvstRtcp.maximumNackEntries * 17 + // No more than one request can name, or the rest would count as requested without being sent. + let limit = NvstRtpNackRequest.maximumSequenceNumbers let missing = (expected.. [NvstReceiveEvent] { + receiver.process(datagram: try NvstReceiverFixtures.seal( + NvstReceiverFixtures.packet(sequence: sequence, frameIndex: UInt32(sequence), flags: 0x07, media: media), + sequence: sequence, handoff: handoff)) + } + _ = try feed(1) + _ = try feed(102) + clock.withLock { $0 = NvstNackTracker.initialDelayNanoseconds } + let requested = try feed(103).flatMap { event -> [UInt64] in + if case .retransmissionWanted(let indices) = event { return indices } else { return [] } + } + #expect(requested == Array(2...65)) + } + /// A resend that never comes ends the wait: the gap is loss once its request has expired. @Test func aRequestedPacketThatNeverArrivesBecomesLossAfterTheWait() throws { let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) From 561e04a45bd9a6093c6849e7c5e4ff60a688da4e Mon Sep 17 00:00:00 2001 From: Anderson Shindy Oki Date: Mon, 5 Oct 2026 12:28:47 +0900 Subject: [PATCH 5/6] fix: hold a gap for the measured round trip and stop when resends never arrive --- .../BifrostFree/NvstMjolnirReceiver.swift | 4 +- GFN/NVST/BifrostFree/NvstNackTracker.swift | 39 +++++- GFN/NVST/BifrostFree/NvstVideoReceiver.swift | 64 +++++++-- OPN/Stream/NvstBifrostFreeTransport.swift | 2 +- Tests/GFN/NVST/NvstMjolnirReceiverTests.swift | 2 +- Tests/GFN/NVST/NvstNackTrackerTests.swift | 51 ++++++- .../NVST/NvstRetransmissionWaitTests.swift | 125 ++++++++++++++++++ 7 files changed, 259 insertions(+), 28 deletions(-) create mode 100644 Tests/GFN/NVST/NvstRetransmissionWaitTests.swift diff --git a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift index 7f05516a..25f85574 100644 --- a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift @@ -213,12 +213,12 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { readSource = nil } - /// Asks the peer for a fresh keyframe on this socket's SRTCP path. - /// The measured round trip retransmission retries wait for. + /// The measured round trip the retransmission retries and the gap hold follow. public func useRetransmissionRoundTrip(milliseconds: Double) { receiver.useRetransmissionRoundTrip(milliseconds: milliseconds) } + /// Asks the peer for a fresh keyframe on this socket's SRTCP path. public func requestKeyframe() { guard let ssrc = receiver.feedbackCounters.boundSSRC else { return } let pli = NvstRtcp.pictureLossIndication(senderSSRC: NvstVideoReceiver.clientSSRC, mediaSSRC: ssrc) diff --git a/GFN/NVST/BifrostFree/NvstNackTracker.swift b/GFN/NVST/BifrostFree/NvstNackTracker.swift index 6318e563..a5e8c3d3 100644 --- a/GFN/NVST/BifrostFree/NvstNackTracker.swift +++ b/GFN/NVST/BifrostFree/NvstNackTracker.swift @@ -5,17 +5,35 @@ import Foundation /// Schedules retransmission requests for missing video packets on the official client's timing, /// as its log prints it: the first request 1 ms after a packet goes missing, then up to three -/// retries, each one round trip plus 4 ms after the last (`useRtdForRtpNackToggle`), within the -/// 52 ms it holds a frame for them. A retry sooner than the round trip asks for a packet that is -/// already on its way, and the seat sends it again. +/// retries, each one round trip plus 4 ms after the last (`useRtdForRtpNackToggle`). +/// +/// A retry sooner than the round trip asks for a packet that is already on its way, and the seat +/// sends it again, so retries stay off until a round trip has been measured. The time a requested +/// packet is held for follows the same measurement: the first request's answer arrives one round +/// trip after it, and one retry is allowed, so the hold covers two. The official client's flat +/// 52 ms only covers that at a ~20 ms round trip — below it the gap waited longer than the answer +/// could take, and above it the answer arrived after the gap had already fallen back. struct NvstNackTracker { static let initialDelayNanoseconds: UInt64 = 1_000_000 static let extraRetryWaitNanoseconds: UInt64 = 4_000_000 static let maximumRetries = 3 - static let maximumWaitNanoseconds: UInt64 = 52_000_000 + /// The wait the official client holds a frame for, and what is used until a round trip is known. + static let defaultWaitNanoseconds: UInt64 = 52_000_000 + /// Never give up on a requested packet sooner than this, whatever the round trip reads: the + /// answer's arrival has to absorb scheduling jitter on both ends. + static let minimumWaitNanoseconds: UInt64 = 16_000_000 + /// The receiver's own wall-clock bound on an unrepaired gap (`fecRepairMaximumWaitNanoseconds`). + /// A hold longer than this can never end, because that bound finalizes the gap first. + static let maximumWaitCeilingNanoseconds: UInt64 = 100_000_000 + /// A round trip above this is not a measurement worth deriving a wait from. + static let maximumRoundTripNanoseconds: UInt64 = 2_000_000_000 /// The wait before a retry: the extra wait alone until a round trip has been measured. private(set) var retryIntervalNanoseconds = NvstNackTracker.extraRetryWaitNanoseconds + /// Whether a second request may go out at all. Off until a round trip is known. + private(set) var retriesEnabled = false + /// How long a requested packet is held for, before the gap falls back to a keyframe. + private(set) var maximumWaitNanoseconds = NvstNackTracker.defaultWaitNanoseconds private(set) var retryCount = 0 private struct Request { @@ -29,7 +47,12 @@ struct NvstNackTracker { var isEmpty: Bool { requests.isEmpty } mutating func useRoundTrip(nanoseconds: UInt64) { - retryIntervalNanoseconds = nanoseconds + Self.extraRetryWaitNanoseconds + let roundTrip = min(nanoseconds, Self.maximumRoundTripNanoseconds) + retryIntervalNanoseconds = roundTrip + Self.extraRetryWaitNanoseconds + retriesEnabled = true + let derived = Self.initialDelayNanoseconds + 2 * retryIntervalNanoseconds + maximumWaitNanoseconds = min(Self.maximumWaitCeilingNanoseconds, + max(Self.minimumWaitNanoseconds, derived)) } /// The missing indices to request now, recording that they were requested. @@ -39,7 +62,8 @@ struct NvstNackTracker { var request = requests[index] ?? Request(firstMissedAt: now) let isDue: Bool if let lastSentAt = request.lastSentAt { - isDue = request.sendCount <= Self.maximumRetries + isDue = retriesEnabled + && request.sendCount <= Self.maximumRetries && now &- lastSentAt >= retryIntervalNanoseconds } else { isDue = now &- request.firstMissedAt >= Self.initialDelayNanoseconds @@ -58,7 +82,7 @@ struct NvstNackTracker { /// Whether `index` was requested and is still inside the time its retransmission takes. func isAwaitingRetransmission(of index: UInt64, now: UInt64) -> Bool { guard let request = requests[index], request.sendCount > 0 else { return false } - return now &- request.firstMissedAt < Self.maximumWaitNanoseconds + return now &- request.firstMissedAt < maximumWaitNanoseconds } /// Forgets an index that arrived. True when it had been requested, so the arrival is a repair. @@ -69,6 +93,7 @@ struct NvstNackTracker { /// Forgets every index below `index`: delivered, or given up on. mutating func forget(below index: UInt64) { guard !requests.isEmpty else { return } + guard requests.keys.contains(where: { $0 < index }) else { return } requests = requests.filter { $0.key >= index } } diff --git a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift index f3b53e2f..18db3aed 100644 --- a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift @@ -58,10 +58,15 @@ public struct NvstReceiverStats: Equatable, Sendable { public var duplicatePackets: UInt64 = 0 /// Packets the SRTP replay window rejected: already seen, or older than the window. public var replayedPackets: UInt64 = 0 - /// Requested packets that arrived while their gap was still open. + /// Requested packets that arrived over SRTP while their gap was still open. public var retransmissionRepairedPackets: UInt64 = 0 /// Requests repeated for a packet already asked for. public var retransmissionRetries: UInt64 = 0 + /// Retransmission requests emitted for packets still missing from an open gap. + public var retransmissionRequestsSent: UInt64 = 0 + /// The retransmission wait was switched off: requests went out and none was ever repaired, + /// so the seat is not answering them and a gap must not be held for an answer that never comes. + public var retransmissionWaitDisabled = false /// Packets rebuilt by FEC recovery and injected back into the reorder path. public var recoveredPackets: UInt64 = 0 /// The largest single finalized-loss range, in packets. @@ -155,6 +160,13 @@ public final class NvstVideoReceiver: @unchecked Sendable { /// client's NACK queue length. RFC 3711's 64 is 7.5 ms at ~8,500 packets/s, less than one /// round trip, so every resent packet was rejected as a replay. public static let replayWindowPackets = 2048 + /// Requests after which a seat that has repaired nothing is treated as one that will not. + public static let retransmissionWaitDisableRequestCount: UInt64 = 8 + /// ...and how long after the first request that verdict is allowed to be reached, so a resend + /// still in flight is never mistaken for one that is not coming. + public static let retransmissionWaitDisableDelayNanoseconds: UInt64 = 200_000_000 + /// A round trip measurement beyond this is not one to derive a retry or a hold from. + private static let maximumRoundTripMilliseconds = 10_000.0 public enum ReceiverError: LocalizedError, Equatable, Sendable { case unsupportedProfile(String) @@ -179,6 +191,10 @@ public final class NvstVideoReceiver: @unchecked Sendable { private var nackTracker = NvstNackTracker() /// When missing packets were last scanned for; nil makes the next packet scan at once. private var lastNackScanAt: UInt64? + /// When the first retransmission request went out, for the disable rule's settling delay. + private var firstRetransmissionRequestAt: UInt64? + /// Set once requests have gone unanswered long enough to stop holding gaps for them. + private var retransmissionWaitDisabled = false private var reorder: [UInt64: NvstRtpVideoPacket] = [:] private var nextIndex: UInt64? private var openGap: (index: UInt64, since: UInt64)? @@ -255,11 +271,12 @@ public final class NvstVideoReceiver: @unchecked Sendable { public var snapshot: NvstReceiverStats { lock.lock(); defer { lock.unlock() }; return stats } - /// The measured round trip retransmission retries wait for. + /// The measured round trip the retransmission retries and the gap hold follow. public func useRetransmissionRoundTrip(milliseconds: Double) { guard milliseconds > 0, milliseconds.isFinite else { return } + let clamped = min(milliseconds, Self.maximumRoundTripMilliseconds) lock.lock() - nackTracker.useRoundTrip(nanoseconds: UInt64(milliseconds * 1_000_000)) + nackTracker.useRoundTrip(nanoseconds: UInt64(clamped * 1_000_000)) lock.unlock() } @@ -291,6 +308,9 @@ public final class NvstVideoReceiver: @unchecked Sendable { nextIndex = nil nackTracker.reset() lastNackScanAt = nil + firstRetransmissionRequestAt = nil + retransmissionWaitDisabled = false + stats.retransmissionWaitDisabled = false reassembler.reset() } @@ -385,7 +405,9 @@ public final class NvstVideoReceiver: @unchecked Sendable { var events: [NvstReceiveEvent] = [] let assembleStart = DispatchTime.now().uptimeNanoseconds defer { stageNanoseconds.assemble += DispatchTime.now().uptimeNanoseconds - assembleStart } - let ordered = arrivals.flatMap { pushReorder(index: $0.index, packet: $0.packet, events: &events) } + let ordered = arrivals.flatMap { + pushReorder(index: $0.index, packet: $0.packet, isFecRecovered: $0.isFecRecovered, events: &events) + } stageNanoseconds.reorder += DispatchTime.now().uptimeNanoseconds - assembleStart for candidate in ordered { reassembleLocked(candidate, events: &events) @@ -422,13 +444,13 @@ public final class NvstVideoReceiver: @unchecked Sendable { /// ordinary duplicate. private func arrivalsLocked(plaintext: Data, packet: NvstRtpVideoPacket, - extendedIndex: UInt64) -> [(index: UInt64, packet: NvstRtpVideoPacket)] { - var arrivals: [(index: UInt64, packet: NvstRtpVideoPacket)] = [(extendedIndex, packet)] + extendedIndex: UInt64) -> [(index: UInt64, packet: NvstRtpVideoPacket, isFecRecovered: Bool)] { + var arrivals: [(index: UInt64, packet: NvstRtpVideoPacket, isFecRecovered: Bool)] = [(extendedIndex, packet, false)] for recoveredPlaintext in fecRecovery.observe(plaintext: plaintext, packet: packet) { guard let repaired = try? NvstVideoPacketParser.parse(recoveredPlaintext), repaired.ssrc == boundSSRC else { continue } stats.recoveredPackets += 1 - arrivals.append((replay.estimatedIndex(for: repaired.sequenceNumber), repaired)) + arrivals.append((replay.estimatedIndex(for: repaired.sequenceNumber), repaired, true)) } return arrivals } @@ -660,7 +682,10 @@ public final class NvstVideoReceiver: @unchecked Sendable { extension NvstVideoReceiver { // MARK: - Reorder - private func pushReorder(index: UInt64, packet: NvstRtpVideoPacket, events: inout [NvstReceiveEvent]) -> [NvstRtpVideoPacket] { + private func pushReorder(index: UInt64, + packet: NvstRtpVideoPacket, + isFecRecovered: Bool, + events: inout [NvstReceiveEvent]) -> [NvstRtpVideoPacket] { let expected = nextIndex ?? index if nextIndex == nil { nextIndex = index } if index < expected { @@ -682,7 +707,12 @@ extension NvstVideoReceiver { lastNackScanAt = nil } } - if !nackTracker.isEmpty, nackTracker.arrived(index) { stats.retransmissionRepairedPackets += 1 } + // Only a packet that came in over SRTP can be credited to the request. An FEC rebuild of a + // requested packet is not the retransmission arriving, and counting it as one would keep + // the hold enabled for a seat that answers nothing. + if !nackTracker.isEmpty, nackTracker.arrived(index), !isFecRecovered { + stats.retransmissionRepairedPackets += 1 + } // Armed FEC holds a gap past the plain reorder window: a block's parity trails an early-frame // hole by up to `fecRepairReorderWindow` packets, so only the wall-clock bound ends it. let isPastFecWait = index > expected && isPastFecRepairWait() @@ -722,7 +752,20 @@ extension NvstVideoReceiver { let missing = (expected..= Self.retransmissionWaitDisableRequestCount, + stats.retransmissionRepairedPackets == 0, + now &- firstRequestAt >= Self.retransmissionWaitDisableDelayNanoseconds { + retransmissionWaitDisabled = true + stats.retransmissionWaitDisabled = true + } } private func recordRecovery(first: UInt64, last: UInt64, events: inout [NvstReceiveEvent]) { @@ -771,6 +814,7 @@ extension NvstVideoReceiver { /// Whether the open gap's first packet was requested and may still be resent in time. private func isWithinRetransmissionWait() -> Bool { + guard !retransmissionWaitDisabled else { return false } guard let expected = nextIndex else { return false } return nackTracker.isAwaitingRetransmission(of: expected, now: uptimeNanoseconds()) } diff --git a/OPN/Stream/NvstBifrostFreeTransport.swift b/OPN/Stream/NvstBifrostFreeTransport.swift index f54a4a8d..da9602a9 100644 --- a/OPN/Stream/NvstBifrostFreeTransport.swift +++ b/OPN/Stream/NvstBifrostFreeTransport.swift @@ -475,7 +475,7 @@ public actor NvstBifrostFreeTransport: NativeNVSTTransport { logger?("NVST audio tracks=\(bundle?.remoteAudioTrackCount ?? 0) pktIn=\(audio?.packets ?? 0) bytesIn=\(audio?.bytes ?? 0)" + " samples=\(audio?.samples ?? 0) concealed=\(audio?.concealed ?? 0) discarded=\(audio?.discarded ?? 0) ssrc=\(audio?.ssrc.map(String.init) ?? "-")") let video = videoPipeline?.snapshot ?? NvstVideoPipeline.Counters() - logger?("NVST counters auth=\(stats.authenticatedPackets) fec=\(stats.fecPackets) dropped=\(stats.droppedPackets) replayed=\(stats.replayedPackets) late=\(stats.latePackets) dup=\(stats.duplicatePackets) nackRepaired=\(stats.retransmissionRepairedPackets) nackRetries=\(stats.retransmissionRetries) rtpLoss=\(stats.finalizedLossPackets) parityLoss=\(stats.parityOnlyLossPackets) frames=\(stats.framesEmitted) keyframes=\(stats.keyframesEmitted) recoveries=\(stats.recoveries) sofFlagged=\(stats.startOfFrameFlagged) sofOk=\(stats.startOfFrameAccepted) abandoned=\(stats.abandonedFrames) rrFail=\(stats.receiverReportFailures)\(stats.lastReceiverReportFailure.map { " rrErr=\($0)" } ?? "") multiBlock=\(stats.multiBlockPackets) maxBlock=\(stats.highestFecLastBlock) decoded=\(decoder?.decodedFrameCount ?? 0) decodeFailed=\(decoder?.failedFrameCount ?? 0) decodeErr=\(decoder?.failureStatusSummary ?? "-") noParamSets=\(video.missingParameterSetFrames) idrOut=\(idrRequestsSent) invalidOut=\(invalidationsSent) inputOut=\(inputEventsSent) padOut=\(gamepadPacketsSent) padFail=\(gamepadSendFailures) padDropped=\(gamepadPacketsDroppedForUnannouncedPad) padReg=\(didRegisterGamepad) textTyped=\(textCharactersTyped) textDroppedBytes=\(textBytesDropped) inputReady=\(bundle?.isInputReady == true) rrOut=\(stats.receiverReportsSent) frac=\(stats.lastFractionLost) lost=\(stats.lastCumulativeLost) jitter=\(stats.lastJitter) seqSpan=\(stats.sequenceSpan) negFps=\(negotiatedFps.map(String.init) ?? "nil") mediaSeconds=\(String(format: "%.2f", Double(stats.lastRtpTimestamp &- (stats.firstRtpTimestamp ?? 0)) / Double(NvstVideoToolboxDecoder.clockRate))) fidxChanges=\(stats.frameIndexChanges) \(perSecondSeries(stats: stats, included: includingPerSecondSeries)) paceOut=\(video.pacingReportsSent) paceFail=\(video.pacingReportFailures) ackOut=\(video.frameAcksSent) ackFail=\(video.frameAckFailures) qosOut=\(qosReportsSent) qosFail=\(qosReportFailures) rtpStatsOut=\(rtpStatsReportsSent) ccStatsOut=\(controlStatsReportsSent) ssrc=\(stats.boundSSRC.map { String(format: "0x%08x", $0) } ?? "-")") + logger?("NVST counters auth=\(stats.authenticatedPackets) fec=\(stats.fecPackets) dropped=\(stats.droppedPackets) replayed=\(stats.replayedPackets) late=\(stats.latePackets) dup=\(stats.duplicatePackets) nackRepaired=\(stats.retransmissionRepairedPackets) nackRetries=\(stats.retransmissionRetries) nackOut=\(stats.retransmissionRequestsSent) nackWaitOff=\(stats.retransmissionWaitDisabled) rtpLoss=\(stats.finalizedLossPackets) parityLoss=\(stats.parityOnlyLossPackets) frames=\(stats.framesEmitted) keyframes=\(stats.keyframesEmitted) recoveries=\(stats.recoveries) sofFlagged=\(stats.startOfFrameFlagged) sofOk=\(stats.startOfFrameAccepted) abandoned=\(stats.abandonedFrames) rrFail=\(stats.receiverReportFailures)\(stats.lastReceiverReportFailure.map { " rrErr=\($0)" } ?? "") multiBlock=\(stats.multiBlockPackets) maxBlock=\(stats.highestFecLastBlock) decoded=\(decoder?.decodedFrameCount ?? 0) decodeFailed=\(decoder?.failedFrameCount ?? 0) decodeErr=\(decoder?.failureStatusSummary ?? "-") noParamSets=\(video.missingParameterSetFrames) idrOut=\(idrRequestsSent) invalidOut=\(invalidationsSent) inputOut=\(inputEventsSent) padOut=\(gamepadPacketsSent) padFail=\(gamepadSendFailures) padDropped=\(gamepadPacketsDroppedForUnannouncedPad) padReg=\(didRegisterGamepad) textTyped=\(textCharactersTyped) textDroppedBytes=\(textBytesDropped) inputReady=\(bundle?.isInputReady == true) rrOut=\(stats.receiverReportsSent) frac=\(stats.lastFractionLost) lost=\(stats.lastCumulativeLost) jitter=\(stats.lastJitter) seqSpan=\(stats.sequenceSpan) negFps=\(negotiatedFps.map(String.init) ?? "nil") mediaSeconds=\(String(format: "%.2f", Double(stats.lastRtpTimestamp &- (stats.firstRtpTimestamp ?? 0)) / Double(NvstVideoToolboxDecoder.clockRate))) fidxChanges=\(stats.frameIndexChanges) \(perSecondSeries(stats: stats, included: includingPerSecondSeries)) paceOut=\(video.pacingReportsSent) paceFail=\(video.pacingReportFailures) ackOut=\(video.frameAcksSent) ackFail=\(video.frameAckFailures) qosOut=\(qosReportsSent) qosFail=\(qosReportFailures) rtpStatsOut=\(rtpStatsReportsSent) ccStatsOut=\(controlStatsReportsSent) ssrc=\(stats.boundSSRC.map { String(format: "0x%08x", $0) } ?? "-")") // Which pipeline stage a latency spike lives in. `peak*` are per-stage session maxima, so a // single 500 ms stall is still visible after the average has recovered. logger?(String(format: "NVST frame stages slow=%d frames=%llu resyncs=%d skipped=%d abandoned=%d lastLatency=%.1fms inputSendTotal=%.0fms inputSendPeak=%.1fms", diff --git a/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift b/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift index 4fefe074..0497046c 100644 --- a/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift +++ b/Tests/GFN/NVST/NvstMjolnirReceiverTests.swift @@ -200,7 +200,7 @@ struct NvstMjolnirReceiverTests { for sequence in UInt16(4)...UInt16(8) { #expect(NvstReceiverFixtures.recoveries(try feed(sequence)) == 0) } - clock.withLock { $0 = NvstNackTracker.maximumWaitNanoseconds } + clock.withLock { $0 = NvstNackTracker.defaultWaitNanoseconds } #expect(NvstReceiverFixtures.recoveries(try feed(9)) == 1) #expect(receiver.snapshot.finalizedLossPackets == 1) } diff --git a/Tests/GFN/NVST/NvstNackTrackerTests.swift b/Tests/GFN/NVST/NvstNackTrackerTests.swift index 5fb9398e..31a1b000 100644 --- a/Tests/GFN/NVST/NvstNackTrackerTests.swift +++ b/Tests/GFN/NVST/NvstNackTrackerTests.swift @@ -10,18 +10,55 @@ struct NvstNackTrackerTests { #expect(tracker.due(missing: [7], now: NvstNackTracker.initialDelayNanoseconds) == [7]) } - @Test func aRequestIsRetriedThreeTimesAtTheRetryInterval() { + @Test func aRequestIsRetriedAtTheRetryIntervalWithinTheRetryBudget() { var tracker = NvstNackTracker() + tracker.useRoundTrip(nanoseconds: 16_000_000) _ = tracker.due(missing: [7], now: 0) var sendTimes: [UInt64] = [] - var now = NvstNackTracker.initialDelayNanoseconds - while now <= NvstNackTracker.maximumWaitNanoseconds { + var now: UInt64 = 0 + while now <= tracker.maximumWaitNanoseconds { if !tracker.due(missing: [7], now: now).isEmpty { sendTimes.append(now) } now += 1_000_000 } - #expect(sendTimes.count == 1 + NvstNackTracker.maximumRetries) - #expect(zip(sendTimes, sendTimes.dropFirst()).allSatisfy { $1 - $0 == NvstNackTracker.extraRetryWaitNanoseconds }) - #expect(tracker.retryCount == NvstNackTracker.maximumRetries) + #expect(sendTimes.first == NvstNackTracker.initialDelayNanoseconds) + #expect(zip(sendTimes, sendTimes.dropFirst()).allSatisfy { $1 - $0 == tracker.retryIntervalNanoseconds }) + #expect(sendTimes.count <= 1 + NvstNackTracker.maximumRetries) + #expect(tracker.retryCount == sendTimes.count - 1) + } + + /// A fixed short retry re-asks for a packet that is already on its way, and the seat sends it + /// again. Without a measured round trip there is nothing to wait for, so there is one request. + @Test func aRequestIsNotRetriedBeforeTheRoundTripIsKnown() { + var tracker = NvstNackTracker() + _ = tracker.due(missing: [7], now: 0) + #expect(tracker.due(missing: [7], now: NvstNackTracker.initialDelayNanoseconds) == [7]) + var now = NvstNackTracker.initialDelayNanoseconds + 1_000_000 + while now <= NvstNackTracker.defaultWaitNanoseconds { + #expect(tracker.due(missing: [7], now: now).isEmpty) + now += 1_000_000 + } + #expect(tracker.retryCount == 0) + } + + /// The hold has to cover the answer to the first request and to one retry; below the official + /// client's 52 ms it fell back before the answer could arrive, and above it the answer arrived + /// after the fallback. It stays inside the receiver's own wall-clock bound on an unrepaired gap. + @Test func theHoldCoversTwoRoundTripsAndIsClamped() { + var tracker = NvstNackTracker() + #expect(tracker.maximumWaitNanoseconds == NvstNackTracker.defaultWaitNanoseconds) + + tracker.useRoundTrip(nanoseconds: 1_000_000) + #expect(tracker.maximumWaitNanoseconds == NvstNackTracker.minimumWaitNanoseconds) + + tracker.useRoundTrip(nanoseconds: 20_000_000) + #expect(tracker.maximumWaitNanoseconds == 49_000_000) + + tracker.useRoundTrip(nanoseconds: 60_000_000) + #expect(tracker.maximumWaitNanoseconds == NvstNackTracker.maximumWaitCeilingNanoseconds) + + // A nonsensical measurement cannot push the wait past that bound, nor overflow the derivation. + tracker.useRoundTrip(nanoseconds: .max) + #expect(tracker.maximumWaitNanoseconds == NvstNackTracker.maximumWaitCeilingNanoseconds) } /// With the round trip known, a retry waits for the answer the first request could bring. @@ -48,7 +85,7 @@ struct NvstNackTrackerTests { #expect(repaired) #expect(!notRepaired) #expect(tracker.isAwaitingRetransmission(of: 8, now: NvstNackTracker.initialDelayNanoseconds)) - #expect(!tracker.isAwaitingRetransmission(of: 8, now: NvstNackTracker.maximumWaitNanoseconds)) + #expect(!tracker.isAwaitingRetransmission(of: 8, now: NvstNackTracker.defaultWaitNanoseconds)) tracker.forget(below: 9) #expect(tracker.isEmpty) } diff --git a/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift b/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift new file mode 100644 index 00000000..295d9317 --- /dev/null +++ b/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift @@ -0,0 +1,125 @@ +import Foundation +import os +import Testing +@testable import OpenNOW + +/// The retransmission wait: how long a gap is held for a resend, and when it stops being held. +@Suite(.serialized) +struct NvstRetransmissionWaitTests { + /// The hold follows the measured round trip: a short one ends the wait well inside the official + /// client's flat 52 ms, so the fallback to a keyframe starts as soon as the answer could have + /// arrived instead of after a wait that outlasts it. + @Test func aShorterRoundTripEndsTheWaitSooner() throws { + let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) + let clock = OSAllocatedUnfairLock(initialState: UInt64(0)) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { clock.withLock { $0 } }) + // A 4 ms round trip derives a 17 ms hold: two round trips plus the initial delay, floored. + receiver.useRetransmissionRoundTrip(milliseconds: 4) + let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] + func feed(_ sequence: UInt16) throws -> [NvstReceiveEvent] { + receiver.process(datagram: try NvstReceiverFixtures.seal( + NvstReceiverFixtures.packet(sequence: sequence, frameIndex: UInt32(sequence), flags: 0x07, media: media), + sequence: sequence, handoff: handoff)) + } + _ = try feed(1) + _ = try feed(3) + clock.withLock { $0 = NvstNackTracker.initialDelayNanoseconds } + _ = try feed(4) + // Past the 17 ms hold but inside the flat 52 ms one, so only the derived hold ends the wait. + clock.withLock { $0 = 20_000_000 } + #expect(NvstReceiverFixtures.recoveries(try feed(5)) == 0) + #expect(NvstReceiverFixtures.recoveries(try feed(6)) == 1) + } + + /// A seat that has repaired none of the requests it was sent is not going to. After enough of + /// them the wait is switched off, so a gap falls back as it did before the requests existed + /// rather than holding a frame for an answer that never comes. + @Test func aSeatThatNeverRepairsStopsTheWaitFromHoldingGaps() throws { + let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) + let clock = OSAllocatedUnfairLock(initialState: UInt64(0)) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { clock.withLock { $0 } }) + receiver.useRetransmissionRoundTrip(milliseconds: 4) + let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] + func feed(_ sequence: UInt16) throws -> [NvstReceiveEvent] { + receiver.process(datagram: try NvstReceiverFixtures.seal( + NvstReceiverFixtures.packet(sequence: sequence, frameIndex: UInt32(sequence), flags: 0x07, media: media), + sequence: sequence, handoff: handoff)) + } + // One gap per round, each requested once and never answered. + var sequence: UInt16 = 1 + for _ in 0..<8 { + let base = sequence + sequence += 8 + clock.withLock { $0 += 30_000_000 } + _ = try feed(base) + _ = try feed(base &+ 2) + clock.withLock { $0 += NvstNackTracker.initialDelayNanoseconds } + _ = try feed(base &+ 3) + for offset in UInt16(4)...6 { _ = try feed(base &+ offset) } + } + let stats = receiver.snapshot + #expect(stats.retransmissionRequestsSent >= NvstVideoReceiver.retransmissionWaitDisableRequestCount) + #expect(stats.retransmissionRepairedPackets == 0) + #expect(stats.retransmissionWaitDisabled) + + // The next gap is no longer held: it is loss as soon as the reorder window passes it. + let base = sequence + _ = try feed(base) + _ = try feed(base &+ 2) + clock.withLock { $0 += NvstNackTracker.initialDelayNanoseconds } + _ = try feed(base &+ 3) + var recoveries = 0 + for offset in UInt16(4)...6 { recoveries += NvstReceiverFixtures.recoveries(try feed(base &+ offset)) } + #expect(recoveries == 1) + } + + /// A packet the seat never sent but FEC rebuilt is not a retransmission arriving: crediting it + /// would keep the wait enabled for a seat that answers nothing, and overstate the repair count. + @Test func anFecRebuildOfARequestedPacketIsNotCreditedToTheRequest() throws { + let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) + let clock = OSAllocatedUnfairLock(initialState: UInt64(0)) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { clock.withLock { $0 } }) + + func fecWord(index: UInt32) -> UInt32 { (50 << 4) | (index << 12) | (2 << 22) } + func blockPackets(frameIndex: UInt32, baseSequence: UInt16, seed: UInt8) -> [(UInt16, Data)] { + let sof = NvstReceiverFixtures.packet(sequence: baseSequence, frameIndex: frameIndex, flags: 0x05, + media: [0x00, 0x00, 0x00, 0x01, 0x65, seed], fecWord: fecWord(index: 0)) + let eof = NvstReceiverFixtures.packet(sequence: baseSequence + 1, frameIndex: frameIndex, flags: 0x03, + media: [0xbb, seed, 0xcc], fecWord: fecWord(index: 1)) + let size = max(sof.count, eof.count) + let shards = [sof, eof].map { source -> [UInt8] in + let bytes = [UInt8](source) + return bytes.count == size ? bytes : bytes + [UInt8](repeating: 0, count: size - bytes.count) + } + let parity = NvstReedSolomon.encode(data: shards, parityCount: 1, size: size)! + let parityHeader = NvstReceiverFixtures.packet(sequence: baseSequence + 2, frameIndex: frameIndex, flags: 0x00, + media: [], fecWord: fecWord(index: 2)) + return [(baseSequence, sof), (baseSequence + 1, eof), + (baseSequence + 2, parityHeader + Data(parity[0][NvstFecRecovery.headerLength...]))] + } + // Clean blocks first: recovery arms only after verifying the scheme against this stream. + for round in 0.. Date: Mon, 5 Oct 2026 13:02:22 +0900 Subject: [PATCH 6/6] fix: reach the retransmission wait verdict on the next scan, not the next request --- GFN/NVST/BifrostFree/NvstVideoReceiver.swift | 27 ++++++----- .../NVST/NvstRetransmissionWaitTests.swift | 46 +++++++++++++++++++ 2 files changed, 62 insertions(+), 11 deletions(-) diff --git a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift index 18db3aed..b7a33c82 100644 --- a/GFN/NVST/BifrostFree/NvstVideoReceiver.swift +++ b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift @@ -747,6 +747,7 @@ extension NvstVideoReceiver { if let lastNackScanAt, now &- lastNackScanAt < NvstNackTracker.initialDelayNanoseconds { return } guard let newest = reorder.keys.max(), newest > expected else { return } lastNackScanAt = now + disableRetransmissionWaitIfUnanswered(now: now) // No more than one request can name, or the rest would count as requested without being sent. let limit = NvstRtpNackRequest.maximumSequenceNumbers let missing = (expected..= Self.retransmissionWaitDisableRequestCount, - stats.retransmissionRepairedPackets == 0, - now &- firstRequestAt >= Self.retransmissionWaitDisableDelayNanoseconds { - retransmissionWaitDisabled = true - stats.retransmissionWaitDisabled = true - } + } + + /// A seat that has repaired none of the requests it was sent is not going to: stop holding gaps + /// for an answer that never comes, which is what the wait costs when it is wrong. Every scan + /// asks, not only the ones that send a request, so a session whose requests have all gone out + /// stops holding the gaps behind them instead of waiting for a further request to judge it by. + private func disableRetransmissionWaitIfUnanswered(now: UInt64) { + guard !retransmissionWaitDisabled, + stats.retransmissionRequestsSent >= Self.retransmissionWaitDisableRequestCount, + stats.retransmissionRepairedPackets == 0, + let firstRequestAt = firstRetransmissionRequestAt, + now &- firstRequestAt >= Self.retransmissionWaitDisableDelayNanoseconds else { return } + retransmissionWaitDisabled = true + stats.retransmissionWaitDisabled = true } private func recordRecovery(first: UInt64, last: UInt64, events: inout [NvstReceiveEvent]) { diff --git a/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift b/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift index 295d9317..3a7c8905 100644 --- a/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift +++ b/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift @@ -73,6 +73,52 @@ struct NvstRetransmissionWaitTests { #expect(recoveries == 1) } + /// The verdict is reached on time, not only when another request happens to go out: eight + /// requests, no repair, 200 ms, and the next scan ends the wait even though it sends nothing. + @Test func theWaitIsDisabledWithoutAnotherRequestGoingOut() throws { + let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 4) + let clock = OSAllocatedUnfairLock(initialState: UInt64(0)) + let receiver = try NvstVideoReceiver(handoff: handoff, uptimeNanoseconds: { clock.withLock { $0 } }) + // A measured round trip, so a retry brings the request count up inside a few short gaps. + receiver.useRetransmissionRoundTrip(milliseconds: 4) + let retryInterval = NvstNackTracker.extraRetryWaitNanoseconds + 4_000_000 + let media: [UInt8] = [0x00, 0x00, 0x00, 0x01, 0x65] + func feed(_ sequence: UInt16) throws -> [NvstReceiveEvent] { + receiver.process(datagram: try NvstReceiverFixtures.seal( + NvstReceiverFixtures.packet(sequence: sequence, frameIndex: UInt32(sequence), flags: 0x07, media: media), + sequence: sequence, handoff: handoff)) + } + var sequence: UInt16 = 1 + _ = try feed(sequence) + // Four one-packet holes, each requested and retried once before its hold expires. Every + // request goes out inside the first 70 ms and none of them is ever repaired. + for _ in 0..<4 { + sequence += 2 + _ = try feed(sequence) + clock.withLock { $0 += NvstNackTracker.initialDelayNanoseconds } + sequence += 1 + _ = try feed(sequence) + clock.withLock { $0 += retryInterval } + sequence += 1 + _ = try feed(sequence) + clock.withLock { $0 += retryInterval } + sequence += 1 + _ = try feed(sequence) + } + let asked = receiver.snapshot + #expect(asked.retransmissionRequestsSent == NvstVideoReceiver.retransmissionWaitDisableRequestCount) + #expect(asked.retransmissionRepairedPackets == 0) + #expect(!asked.retransmissionWaitDisabled) + + // Past the delay, the next scan sends nothing and still ends the wait. + clock.withLock { $0 = 202_000_000 } + sequence += 2 + _ = try feed(sequence) + let stats = receiver.snapshot + #expect(stats.retransmissionRequestsSent == asked.retransmissionRequestsSent) + #expect(stats.retransmissionWaitDisabled) + } + /// A packet the seat never sent but FEC rebuilt is not a retransmission arriving: crediting it /// would keep the wait enabled for a seat that answers nothing, and overstate the repair count. @Test func anFecRebuildOfARequestedPacketIsNotCreditedToTheRequest() throws {