Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 33 additions & 4 deletions GFN/NVST/BifrostFree/NvstMjolnirReceiver.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)?

Expand Down Expand Up @@ -211,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)
Expand Down Expand Up @@ -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 }
Expand Down Expand Up @@ -404,6 +419,20 @@ 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.prefix(NvstRtpNackRequest.maximumSequenceNumbers).map { UInt16(truncatingIfNeeded: $0) }
guard handler(sequenceNumbers) else { continue }
counterLock.lock()
nacksSent += 1
nackedPackets += sequenceNumbers.count
counterLock.unlock()
}
}
}
Expand Down
78 changes: 78 additions & 0 deletions GFN/NVST/BifrostFree/NvstNackTracker.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
// 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`), 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 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?
var sendCount = 0
}

private var requests: [UInt64: Request] = [:]

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] = []
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 >= 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 < 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()
}
}
59 changes: 59 additions & 0 deletions GFN/NVST/BifrostFree/NvstRtpNackRequest.swift
Original file line number Diff line number Diff line change
@@ -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)
}
}
61 changes: 58 additions & 3 deletions GFN/NVST/BifrostFree/NvstVideoReceiver.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -53,6 +56,12 @@ 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
/// 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.
Expand Down Expand Up @@ -142,6 +151,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)
Expand All @@ -162,7 +175,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)?
Expand Down Expand Up @@ -239,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
Expand All @@ -265,6 +289,8 @@ public final class NvstVideoReceiver: @unchecked Sendable {
defer { lock.unlock() }
reorder.removeAll()
nextIndex = nil
nackTracker.reset()
lastNackScanAt = nil
reassembler.reset()
}

Expand Down Expand Up @@ -330,6 +356,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
Expand Down Expand Up @@ -650,8 +677,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()
Expand All @@ -673,9 +704,26 @@ 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..<newest).lazy.filter { self.reorder[$0] == nil }.prefix(limit)
let due = nackTracker.due(missing: Array(missing), now: now)
stats.retransmissionRetries = UInt64(nackTracker.retryCount)
if !due.isEmpty { events.append(.retransmissionWanted(due)) }
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

private func recordRecovery(first: UInt64, last: UInt64, events: inout [NvstReceiveEvent]) {
reassembler.reset()
stats.recoveries += 1
Expand Down Expand Up @@ -716,7 +764,14 @@ 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 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.
Expand Down
Loading