diff --git a/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift b/GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift index 59a1be43..25f85574 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)? @@ -210,6 +213,11 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { readSource = nil } + /// 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 } @@ -252,11 +260,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 +419,25 @@ public final class NvstMjolnirReceiver: @unchecked Sendable { let handler = onDrop callbackLock.unlock() handler?(reason) + case .retransmissionWanted(let indices): + callbackLock.lock() + let handler = onRetransmissionWanted + callbackLock.unlock() + guard let handler else { + requestRetransmission(of: indices) + 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 + counterLock.unlock() } } } diff --git a/GFN/NVST/BifrostFree/NvstNackTracker.swift b/GFN/NVST/BifrostFree/NvstNackTracker.swift new file mode 100644 index 00000000..a5e8c3d3 --- /dev/null +++ b/GFN/NVST/BifrostFree/NvstNackTracker.swift @@ -0,0 +1,103 @@ +// 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, 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 + /// 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 { + let firstMissedAt: UInt64 + var lastSentAt: UInt64? + var sendCount = 0 + } + + private var requests: [UInt64: Request] = [:] + + var isEmpty: Bool { requests.isEmpty } + + mutating func useRoundTrip(nanoseconds: UInt64) { + 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. + 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 = retriesEnabled + && request.sendCount <= Self.maximumRetries + && 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) + } + 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 < 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 } + guard requests.keys.contains(where: { $0 < index }) else { return } + requests = requests.filter { $0.key >= index } + } + + mutating func reset() { + requests.removeAll() + } +} 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/GFN/NVST/BifrostFree/NvstVideoReceiver.swift b/GFN/NVST/BifrostFree/NvstVideoReceiver.swift index ac5ac610..b7a33c82 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,17 @@ 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 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. @@ -142,6 +156,17 @@ 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 + /// 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) @@ -162,7 +187,14 @@ 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? + /// 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)? @@ -239,6 +271,15 @@ public final class NvstVideoReceiver: @unchecked Sendable { public var snapshot: NvstReceiverStats { lock.lock(); defer { lock.unlock() }; return stats } + /// 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(clamped * 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 @@ -265,6 +306,11 @@ public final class NvstVideoReceiver: @unchecked Sendable { defer { lock.unlock() } reorder.removeAll() nextIndex = nil + nackTracker.reset() + lastNackScanAt = nil + firstRetransmissionRequestAt = nil + retransmissionWaitDisabled = false + stats.retransmissionWaitDisabled = false reassembler.reset() } @@ -330,6 +376,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 @@ -358,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) @@ -395,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 } @@ -633,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 { @@ -650,7 +702,16 @@ extension NvstVideoReceiver { } if index > expected { stats.outOfOrderPackets += 1 - if openGap?.index != expected { openGap = (expected, uptimeNanoseconds()) } + if openGap?.index != expected { + openGap = (expected, uptimeNanoseconds()) + lastNackScanAt = nil + } + } + // 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. @@ -673,9 +734,45 @@ 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 + 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, + let firstRequestAt = firstRetransmissionRequestAt, + now &- firstRequestAt >= Self.retransmissionWaitDisableDelayNanoseconds else { return } + retransmissionWaitDisabled = true + stats.retransmissionWaitDisabled = true + } + private func recordRecovery(first: UInt64, last: UInt64, events: inout [NvstReceiveEvent]) { reassembler.reset() stats.recoveries += 1 @@ -716,7 +813,15 @@ extension NvstVideoReceiver { private func isGapFinalized(depth: UInt64, isPastFecWait: Bool) -> 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 !retransmissionWaitDisabled else { return false } + 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..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) 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", @@ -487,6 +487,9 @@ public actor NvstBifrostFreeTransport: NativeNVSTTransport { + " decoderSessions=\(decoder?.sessionCreationCount ?? 0)" + " hwDecode=\(decoder?.isHardwareAccelerated == true)" + " decode\(decoder?.stageTimingSummary ?? "-")") + if let controlRoundTrip = await session?.controlRoundTripMilliseconds() { + receiver.useRetransmissionRoundTrip(milliseconds: controlRoundTrip) + } await logHudCounters(receiver: receiver, stats: stats) } 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 ffccea50..0497046c 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) @@ -128,11 +128,88 @@ 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) + } + + /// One request names at most 64 packets, so a wider gap asks for its first 64 and no more: a + /// packet counted as requested but never named would hold its gap open for nothing. + @Test func aWideGapRequestsOnlyWhatOneRequestCanName() throws { + let handoff = NvstReceiverFixtures.makeHandoff(reorderWindow: 128) + 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(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) + 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.defaultWaitNanoseconds } + #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 { 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( @@ -157,7 +234,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) @@ -180,7 +257,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 @@ -213,7 +290,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 @@ -242,7 +319,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). @@ -297,7 +374,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)] { @@ -364,7 +441,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/NvstNackTrackerTests.swift b/Tests/GFN/NVST/NvstNackTrackerTests.swift new file mode 100644 index 00000000..31a1b000 --- /dev/null +++ b/Tests/GFN/NVST/NvstNackTrackerTests.swift @@ -0,0 +1,92 @@ +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 aRequestIsRetriedAtTheRetryIntervalWithinTheRetryBudget() { + var tracker = NvstNackTracker() + tracker.useRoundTrip(nanoseconds: 16_000_000) + _ = tracker.due(missing: [7], now: 0) + var sendTimes: [UInt64] = [] + 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.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. + @Test func aRetryWaitsOneRoundTripPlusTheExtraWait() { + var tracker = NvstNackTracker() + tracker.useRoundTrip(nanoseconds: 16_000_000) + _ = tracker.due(missing: [7], now: 0) + let first = NvstNackTracker.initialDelayNanoseconds + #expect(tracker.due(missing: [7], now: first) == [7]) + #expect(tracker.due(missing: [7], now: first + 19_999_999).isEmpty) + #expect(tracker.due(missing: [7], now: first + 20_000_000) == [7]) + #expect(tracker.retryCount == 1) + } + + @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.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..3a7c8905 --- /dev/null +++ b/Tests/GFN/NVST/NvstRetransmissionWaitTests.swift @@ -0,0 +1,171 @@ +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) + } + + /// 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 { + 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..