Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
42 changes: 38 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 @@ -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 }
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,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()
}
}
}
Expand Down
103 changes: 103 additions & 0 deletions GFN/NVST/BifrostFree/NvstNackTracker.swift
Original file line number Diff line number Diff line change
@@ -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()
}
}
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)
}
}
Loading