diff --git a/CHANGELOG.md b/CHANGELOG.md index 44b65b10..29594328 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,7 +10,18 @@ the public-API contract. ## [Unreleased] -_Nothing yet._ +### Fixed + +- **`LiveTelemetry`'s two bitrate fields measure the media played, not the bytes transferred** (#514). + Both were metered from the reader's transfer counter, which parts from playback on every route that + reads ahead: a 20 Mbps VOD stream read about 35 Mbps (prefetch and seek re-fetches counted as they + arrived), and a paused live session kept draining its origin into the DVR window while the divisor + stood still, so its average climbed for as long as the pause ran. The pumps now record the played + video and audio packets by presentation time, and the sampler charges what the playhead crossed: + `instantBitrateMbps` over about the last 10 s of playback, `averageBitrateMbps` over the session. + Both stand still through a pause on every route, live included, and a seek charges nothing for the + span it jumps. The transfer remains `networkThroughputMbps`. The remote-HLS bypass is unchanged (it + reports the variant's declared rates). ## [7.22.0] - 2026-09-28 diff --git a/Sources/AetherEngine/AetherEngine+Diagnostics.swift b/Sources/AetherEngine/AetherEngine+Diagnostics.swift index ca97196f..875dde21 100644 --- a/Sources/AetherEngine/AetherEngine+Diagnostics.swift +++ b/Sources/AetherEngine/AetherEngine+Diagnostics.swift @@ -506,6 +506,12 @@ extension AetherEngine { native: nativeVideoSession?.demuxerBytesFetched) } + /// AE#514: the played-media ledger of whichever path owns the pump, same precedence as the byte + /// counter above. nil before a pipeline exists and on the paths with no pump to feed one. + var playedMediaLedger: PlayedMediaLedger? { + softwareHost?.playedMediaLedger ?? nativeVideoSession?.playedMediaLedger + } + /// #306: the precedence itself, as a function, so the ordering is assertable without a live /// session on either path. Software first: only one of the two exists per session, and a /// software session's counter is the one that used to be dropped. diff --git a/Sources/AetherEngine/Diagnostics/LiveTelemetry.swift b/Sources/AetherEngine/Diagnostics/LiveTelemetry.swift index 35616425..8c5b8955 100644 --- a/Sources/AetherEngine/Diagnostics/LiveTelemetry.swift +++ b/Sources/AetherEngine/Diagnostics/LiveTelemetry.swift @@ -9,16 +9,20 @@ import Foundation /// interpret a sub-second decoded queue as the size of the compressed packet cache. public struct LiveTelemetry: Equatable, Sendable { // Enthusiast section - /// Mean rate over the last 10 s of the bytes the session pulled from its source, from the demuxer's - /// count on the loopback and software paths. nil until the window spans two ticks. The remote-HLS - /// bypass has no demuxer, and what AVPlayer transferred there is buffer fill at link speed rather than - /// the stream's rate, so on that route this is the playing variant's declared BANDWIDTH. + /// Rate of the media the playhead crossed over about the last 10 s of playback: the bytes of the + /// played video and audio packets presented in that span, over the media seconds it covered + /// (AE#514). Not the transfer, which is `networkThroughputMbps`: a read-ahead, a re-fetch after a + /// seek and a live source draining into its DVR window while paused all move bytes that nobody + /// played. Stands still through a pause, on live too. nil until the playhead has crossed a span the + /// session has bytes for. The remote-HLS bypass has no demuxer, so on that route this is the playing + /// variant's declared BANDWIDTH. public let instantBitrateMbps: Double? - /// Lifetime mean rate of the session, over the seconds it spent consuming media rather than over - /// wall-clock seconds since it started (AE#514). Metered from the same counter as `instantBitrateMbps`; on the remote-HLS bypass it is the variant's declared AVERAGE-BANDWIDTH, or BANDWIDTH where the master omits it. A pause therefore leaves this value standing - /// still instead of dragging it toward zero for as long as the pause lasts, and so does the tail - /// after end-of-media. nil until the session has both accrued active time and fetched something: - /// like `networkThroughputMbps`, a value that cannot be measured yet is a gap, never a zero. + /// Lifetime mean of the same quotient: every played byte over every media second played this + /// session (AE#514). A seek charges nothing for the span it jumped, a pause and the tail after + /// end-of-media charge nothing at all, and the prefetch counts once it is played, not when it + /// arrives. On the remote-HLS bypass it is the variant's declared AVERAGE-BANDWIDTH, or BANDWIDTH + /// where the master omits it. nil until something was played: like `networkThroughputMbps`, a value + /// that cannot be measured yet is a gap, never a zero. public let averageBitrateMbps: Double? /// Live bitrate of the audio bridge's encoded output, or nil when no bridge is active (stream-copy / /// AVPlayer-native path) or before the first delta. Measured from the bridge's cumulative output-byte diff --git a/Sources/AetherEngine/Diagnostics/LiveTelemetrySampler.swift b/Sources/AetherEngine/Diagnostics/LiveTelemetrySampler.swift index 7bc50ad6..104f8e06 100644 --- a/Sources/AetherEngine/Diagnostics/LiveTelemetrySampler.swift +++ b/Sources/AetherEngine/Diagnostics/LiveTelemetrySampler.swift @@ -118,16 +118,16 @@ final class LiveTelemetrySampler { private var lastDemuxerBytes: Int64 = 0 private var lastBridgeBytes: Int64 = 0 private var lastFramesEnqueued: Int = 0 - private var sessionStartBytes: Int64 = 0 - /// AE#514: wall-clock seconds this session spent in a phase that consumes media, accumulated one - /// tick at a time. The divisor of the lifetime average, in place of the wall clock since start. - /// Readable so a test can pin that the tick charges it, not only that the fold over it is right. - private(set) var activeSeconds: Double = 0 + /// AE#514: both bitrate fields, metered over what the playhead crossed. Readable so a test can pin + /// that the tick feeds it, not only that the meter is right. + private(set) var playedMeter = PlayedBitrateMeter() - /// Timestamp of the previous tick, anchoring the delta charged above. nil until the first tick, - /// which therefore charges nothing: that is also the tick which seeds `sessionStartBytes`, so the - /// numerator and the divisor start counting at the same instant. + /// The ledger the meter was last advanced against. A different one is a different session, whose + /// playhead and bytes have nothing to do with the old meter's. + private var meterLedger: ObjectIdentifier? + + /// Timestamp of the previous tick: the wall time the playhead had to cover its step in. private var lastTickTime: Date? /// [LagDiag] tick-over-tick state (#93 post-recovery lag diagnosis). @@ -160,9 +160,9 @@ final class LiveTelemetrySampler { lastDemuxerBytes = engine?.demuxerBytesFetched ?? 0 lastBridgeBytes = engine?.audioBridgeOutputBytesLifetime ?? 0 lastFramesEnqueued = 0 - activeSeconds = 0 + playedMeter = PlayedBitrateMeter() + meterLedger = nil lastTickTime = nil - sessionStartBytes = 0 lagLastClock = nil lagLastDroppedSum = 0 eomParkFrozenTicks = 0 @@ -231,78 +231,48 @@ final class LiveTelemetrySampler { route != .remoteBypass } - /// AE#514: whether a tick's second belongs in the lifetime average's divisor. - /// - /// That average used to divide by wall-clock time since the session started, which made a pause - /// permanently wrong in one direction. `demuxerBytesFetched` stops advancing once the forward - /// buffer is full, the wall clock does not, so a 2.8 Mbps file left paused for three minutes - /// reported 0.4 Mbps and afterwards climbed back only asymptotically: the paused seconds never - /// left the divisor again. Same shape at the end of a source, where the sampler keeps ticking - /// until the host tears the session down. - /// - /// The line is drawn at "is the session consuming media", not at "is the picture moving". A seek - /// and a rebuffer are where the bytes arrive hardest, and a stalled reader is a real part of what - /// this session averaged, so charging their seconds is what keeps the quotient equal to the rate - /// the session pulled at. Only a pause and the three phases with no live session behind them - /// stand outside it. - static func chargesActiveTime(_ phase: PlaybackPhase) -> Bool { - switch phase { - case .paused, .idle, .ended, .error: - return false - case .loading, .playing, .seeking, .rebuffering, .stalled: - return true + /// AE#514: where the playhead stands on the axis the played-media ledger is keyed on, the source + /// axis the pumps read packets on. Native folds AVPlayer's clock back with the producer's shift + /// (`source_pts - playlistShiftSeconds`) rather than reading `sourceTime`, which publishes the item + /// axis on a sequential origin (#368). nil where there is no ledger-fed pipeline to read it off. + nonisolated static func ledgerPlayhead(backend: PlaybackBackend, nativeClock: Double?, playlistShift: Double, + softwareSourceClock: Double?) -> Double? { + switch backend { + case .native: + guard let nativeClock, nativeClock.isFinite else { return nil } + return nativeClock + playlistShift + case .software: + return softwareSourceClock + case .aether, .none, .audio: + return nil } } - /// Lifetime mean of the bytes the session fetched over the seconds it spent consuming them (AE#514). - /// - /// nil rather than zero until both halves are measurable, mirroring `observedTransferMbps`: a host - /// cannot tell a confident 0.00 Mbps from a session that has not fetched anything yet. The old - /// wall-clock form published exactly that on its first tick, where the lifetime delta is zero by - /// construction because the same tick seeds the baseline it then subtracts. - static func averageBitrateMbps(lifetimeBytes: Int64, activeSeconds: Double) -> Double? { - guard activeSeconds > 0, lifetimeBytes > 0 else { return nil } - return Double(lifetimeBytes) * 8.0 / activeSeconds / 1_000_000.0 - } - private func tick() async { guard let engine = engine else { return } - // AE#514: this tick's wall-clock second goes into the lifetime average's divisor only if the - // session was consuming media in it. See `chargesActiveTime` for where the line runs. let tickTime = Date() - if let previous = lastTickTime, Self.chargesActiveTime(engine.playbackPhase) { - activeSeconds += tickTime.timeIntervalSince(previous) - } + let wallSeconds = lastTickTime.map { tickTime.timeIntervalSince($0) } ?? 0 lastTickTime = tickTime let route = engine.videoRoute - // Instant + average bitrate from demuxer byte counters (loopback and SW; the bypass overrides below) + // What the reader pulled from the source. Transfer, not media rate: it feeds the network fields + // and [LagDiag], never the bitrate fields (AE#514, see `playedMeter`). let demuxerBytes = engine.demuxerBytesFetched let bytesThisTick = max(0, demuxerBytes - lastDemuxerBytes) lastDemuxerBytes = demuxerBytes - if sessionStartBytes == 0 { sessionStartBytes = demuxerBytes } byteWindow.push(bytesThisTick) - var instantBitrateMbps: Double? - if byteWindow.count >= 2 { - let totalBytes = byteWindow.sum - let seconds = Double(byteWindow.count) - instantBitrateMbps = Double(totalBytes) * 8.0 / seconds / 1_000_000.0 - } else { - instantBitrateMbps = nil - } + let transferMbps: Double? = byteWindow.count >= 2 + ? Double(byteWindow.sum) * 8.0 / Double(byteWindow.count) / 1_000_000.0 + : nil let observedTransferMbps = Self.observedTransferMbps( windowBytes: byteWindow.sum, activeSeconds: byteWindow.activeCount, samples: byteWindow.count) - var averageBitrateMbps = Self.averageBitrateMbps( - lifetimeBytes: max(0, demuxerBytes - sessionStartBytes), - activeSeconds: activeSeconds) - // Live audio-bridge output bitrate from the bridge's cumulative encoded-byte counter. 0 on the // stream-copy / AVPlayer-native / video-only paths (no bridge), which surfaces as nil. let bridgeBytes = engine.audioBridgeOutputBytesLifetime @@ -420,6 +390,27 @@ final class LiveTelemetrySampler { accumulatedFrameDelaySeconds = nil } + // AE#514: the bitrate fields are the media the playhead crossed, never the transfer. + var instantBitrateMbps: Double? + var averageBitrateMbps: Double? + if let ledger = engine.playedMediaLedger { + let identity = ObjectIdentifier(ledger) + if meterLedger != identity { + meterLedger = identity + playedMeter = PlayedBitrateMeter() + } + playedMeter.advance( + to: Self.ledgerPlayhead( + backend: engine.playbackBackend, + nativeClock: nativeReadings?.currentTimeSeconds, + playlistShift: engine.playlistShiftSeconds, + softwareSourceClock: engine.softwareHost?.sourceClockSeconds), + wallSeconds: wallSeconds, + ledger: ledger) + instantBitrateMbps = playedMeter.instantMbps + averageBitrateMbps = playedMeter.averageMbps + } + // Remote-HLS bypass: no demuxer, so both rates are what the playing variant declares. if Self.bitrateCounter(for: route) == .declaredVariant { let declared = Self.declaredBitrates( @@ -434,7 +425,7 @@ final class LiveTelemetrySampler { engine.extractorYieldState.setForwardBuffer(forwardBufferSeconds) if let readings = nativeReadings { - emitLagDiag(engine: engine, readings: readings, netMbps: instantBitrateMbps) + emitLagDiag(engine: engine, readings: readings, netMbps: transferMbps) if Self.readsLoopbackPipeline(route) { evaluateEndOfMediaPark(engine: engine, readings: readings) } diff --git a/Sources/AetherEngine/Diagnostics/PlayedMediaLedger.swift b/Sources/AetherEngine/Diagnostics/PlayedMediaLedger.swift new file mode 100644 index 00000000..c6585e05 --- /dev/null +++ b/Sources/AetherEngine/Diagnostics/PlayedMediaLedger.swift @@ -0,0 +1,171 @@ +import Foundation + +/// AE#514 round 2: the bytes of the media the session plays, keyed by presentation time on the source +/// axis, so the bitrate fields can be metered from what the playhead crossed rather than from what the +/// reader transferred. +/// +/// The transfer counter and playback part ways on every route that reads ahead: VOD prefetches minutes +/// in front of the playhead, a seek discards a buffer and fetches it again, and a paused live session +/// keeps draining the origin into its DVR window (AE#443). A rate metered off that counter reported the +/// prefetch (35 Mbps for a 20 Mbps stream) and climbed without bound through a live pause. +/// +/// The pump records every packet of the played video and audio streams as it hands them on; the +/// sampler consumes the span the playhead crossed since its last tick. Thread-safe: the pump records +/// on its read thread, the sampler consumes on the main actor. +final class PlayedMediaLedger: @unchecked Sendable { + + enum Track: Int, CaseIterable { + case video = 0 + case audio = 1 + } + + private struct Entry { + let pts: Double + var bytes: Int + } + + /// A packet this far behind the newest one on its track is a re-read (a seek replaying a range the + /// reader already delivered), not B-frame reorder, which stays well under a second. + static let rereadThresholdSeconds: Double = 2.0 + + /// Entries held across both tracks. A VOD read-ahead of several minutes stays far below it; past + /// it the farthest-ahead packets go unrecorded, so the span they cover later reads as unmeasured + /// rather than as a lower rate. + static let maxEntries = 262_144 + + private let lock = NSLock() + private var tracks: [[Entry]] = Array(repeating: [], count: Track.allCases.count) + private var count = 0 + + /// Records one played packet. `pts` is in seconds on the source axis, the axis the playhead the + /// sampler hands `consume` is on. + func record(_ track: Track, pts: Double, bytes: Int) { + guard pts.isFinite, bytes > 0 else { return } + lock.lock() + defer { lock.unlock() } + var entries = tracks[track.rawValue] + tracks[track.rawValue] = [] + defer { tracks[track.rawValue] = entries } + + if let newest = entries.last?.pts, pts < newest - Self.rereadThresholdSeconds { + let cut = Self.lowerBound(entries, pts) + count -= entries.count - cut + entries.removeSubrange(cut...) + } + var index = entries.count + while index > 0, entries[index - 1].pts > pts { index -= 1 } + if index > 0, entries[index - 1].pts == pts { + entries[index - 1].bytes = bytes + return + } + guard count < Self.maxEntries else { return } + entries.insert(Entry(pts: pts, bytes: bytes), at: index) + count += 1 + } + + /// Bytes of every packet presented in `[from, to)`, which leave the ledger with everything before `to`. + func consume(from: Double, to: Double) -> Int64 { + lock.lock() + defer { lock.unlock() } + var total: Int64 = 0 + for track in tracks.indices { + let start = Self.lowerBound(tracks[track], from) + let end = Self.lowerBound(tracks[track], to) + if start < end { + for entry in tracks[track][start.. 0 else { return } + tracks[track].removeFirst(end) + count -= end + } + + private static func lowerBound(_ entries: [Entry], _ pts: Double) -> Int { + var low = 0, high = entries.count + while low < high { + let mid = (low + high) / 2 + if entries[mid].pts < pts { low = mid + 1 } else { high = mid } + } + return low + } +} + +/// AE#514 round 2: both bitrate fields over the media the playhead crossed. Advanced once per sampler +/// tick with the playhead on the ledger's axis. +/// +/// A step forward no larger than playback could have covered in the tick is charged: its bytes from the +/// ledger, its media seconds as the divisor. Anything else is a seek or a jump and charges nothing, and a +/// paused playhead does not move, so both values stand still through a pause on every route, live +/// included. A step whose span the ledger holds no bytes for is left out too (the ledger was not fed +/// there, or its axis does not line up with this playhead): an unmeasured span is a gap, never a zero. +struct PlayedBitrateMeter { + + /// Fastest playback the charge accepts, as media seconds per wall second, plus a fixed slack for + /// tick jitter. Past that the playhead jumped. + static let maxRate: Double = 4.0 + static let stepSlackSeconds: Double = 0.5 + + /// Charged ticks the instant value averages over, about the last ten seconds of playback. + static let windowTicks = 10 + + private var lastPlayhead: Double? + private var window: [(bytes: Int64, seconds: Double)] = [] + private(set) var lifetimeBytes: Int64 = 0 + private(set) var lifetimeSeconds: Double = 0 + + mutating func advance(to playhead: Double?, wallSeconds: Double, ledger: PlayedMediaLedger) { + guard let playhead, playhead.isFinite else { return } + defer { lastPlayhead = playhead } + guard let last = lastPlayhead else { + ledger.discard(below: playhead) + return + } + let step = playhead - last + guard step != 0 else { return } + guard step > 0, step <= max(0, wallSeconds) * Self.maxRate + Self.stepSlackSeconds else { + ledger.discard(below: playhead) + return + } + let bytes = ledger.consume(from: last, to: playhead) + guard bytes > 0 else { return } + lifetimeBytes += bytes + lifetimeSeconds += step + window.append((bytes, step)) + if window.count > Self.windowTicks { window.removeFirst(window.count - Self.windowTicks) } + } + + /// Mean rate of the last `windowTicks` charged ticks. nil before the first one. + var instantMbps: Double? { + let seconds = window.reduce(0) { $0 + $1.seconds } + let bytes = window.reduce(Int64(0)) { $0 + $1.bytes } + return Self.mbps(bytes: bytes, seconds: seconds) + } + + /// Mean rate of everything played this session. nil before the first charged tick. + var averageMbps: Double? { Self.mbps(bytes: lifetimeBytes, seconds: lifetimeSeconds) } + + private static func mbps(bytes: Int64, seconds: Double) -> Double? { + guard seconds > 0, bytes > 0 else { return nil } + return Double(bytes) * 8.0 / seconds / 1_000_000.0 + } +} diff --git a/Sources/AetherEngine/Native/SoftwarePlaybackHost.swift b/Sources/AetherEngine/Native/SoftwarePlaybackHost.swift index 1a64784d..b16f70c8 100644 --- a/Sources/AetherEngine/Native/SoftwarePlaybackHost.swift +++ b/Sources/AetherEngine/Native/SoftwarePlaybackHost.swift @@ -54,6 +54,10 @@ final class SoftwarePlaybackHost { demuxer?.avioBytesFetched } + /// AE#514: bytes of the played streams by presentation time on the source axis (the axis of + /// `sourceClockSeconds`), fed by the read loops once the timeline fold has been applied. + nonisolated let playedMediaLedger = PlayedMediaLedger() + @Published private(set) var isReady: Bool = false @Published private(set) var currentTime: Double = 0 /// Raw synchronizer clock in the SOURCE axis (same axis as demuxed packet PTS and @@ -1675,6 +1679,7 @@ final class SoftwarePlaybackHost { let getSubtitleTapSink: @Sendable () -> ((@Sendable (Int32, UnsafeMutablePointer, AVRational, Bool) -> Void)?) = { [weak self] in self?.subtitleTapSink } + let playedMediaLedger = playedMediaLedger // AE#560: both read loops hand every source packet to this before any branching. let recordingTap: @Sendable (UnsafeMutablePointer) -> Void = { [weak self] pkt in self?.tapForRecording(pkt) @@ -1780,7 +1785,8 @@ final class SoftwarePlaybackHost { subtitleTimeBases: subTimeBases, splitDisplaySetSubtitleStreamIndices: subSplitSetIndices, subtitleTapSink: getSubtitleTapSink, - recordingTap: recordingTap + recordingTap: recordingTap, + playedMedia: playedMediaLedger ) } let lookahead = audioLookahead @@ -1855,7 +1861,8 @@ final class SoftwarePlaybackHost { subtitleTimeBases: subTimeBases, splitDisplaySetSubtitleStreamIndices: subSplitSetIndices, subtitleTapSink: getSubtitleTapSink, - recordingTap: recordingTap + recordingTap: recordingTap, + playedMedia: playedMediaLedger ) } } @@ -1879,7 +1886,8 @@ final class SoftwarePlaybackHost { subtitleTimeBases: [Int32: AVRational] = [:], splitDisplaySetSubtitleStreamIndices: Set = [], subtitleTapSink: @Sendable () -> ((@Sendable (Int32, UnsafeMutablePointer, AVRational, Bool) -> Void)?) = { nil }, - recordingTap: @escaping @Sendable (UnsafeMutablePointer) -> Void = { _ in } + recordingTap: @escaping @Sendable (UnsafeMutablePointer) -> Void = { _ in }, + playedMedia: PlayedMediaLedger? = nil ) { let discontinuityThresholdSeconds = 10.0 var prevRawVideoPtsSec = Double.nan @@ -1993,6 +2001,7 @@ final class SoftwarePlaybackHost { if rawPts != Int64.min, tbSec > 0 { let ptsSec = Double(rawPts) * tbSec noteEdge(ptsSec) + playedMedia?.record(isVideo ? .video : .audio, pts: ptsSec, bytes: Int(packet.pointee.size)) if let data = packet.pointee.data, packet.pointee.size > 0 { let bytes = Data(bytes: data, count: Int(packet.pointee.size)) let isKey = isVideo && (packet.pointee.flags & AV_PKT_FLAG_KEY) != 0 @@ -2363,7 +2372,8 @@ final class SoftwarePlaybackHost { subtitleTimeBases: [Int32: AVRational] = [:], splitDisplaySetSubtitleStreamIndices: Set = [], subtitleTapSink: @Sendable () -> ((@Sendable (Int32, UnsafeMutablePointer, AVRational, Bool) -> Void)?) = { nil }, - recordingTap: @escaping @Sendable (UnsafeMutablePointer) -> Void = { _ in } + recordingTap: @escaping @Sendable (UnsafeMutablePointer) -> Void = { _ in }, + playedMedia: PlayedMediaLedger? = nil ) { // Clock arming: one-shot latch (seekClock is not idempotent -- re-calling snaps clock back to initialClockTime). Shared with host so a seek before first audio isn't overridden by a late re-arm. @@ -2746,6 +2756,16 @@ final class SoftwarePlaybackHost { } } + // AE#514: on the folded axis the clock runs on, which is what the sampler reads the playhead off. + if let playedMedia, streamIdx == videoStreamIndex || streamIdx == audioStreamIndex { + let tbSec = streamIdx == videoStreamIndex ? videoTimeBaseSeconds : audioTimeBaseSeconds + let ticks = packet.pointee.pts != Int64.min ? packet.pointee.pts : packet.pointee.dts + if ticks != Int64.min, tbSec > 0 { + playedMedia.record(streamIdx == videoStreamIndex ? .video : .audio, + pts: Double(ticks) * tbSec, bytes: Int(packet.pointee.size)) + } + } + // Fill ring before decode so the ring holds every packet. Audio appended for sync; only video keyframes tagged for eviction alignment. if isLive { let isVideo = streamIdx == videoStreamIndex diff --git a/Sources/AetherEngine/Video/HLSSegmentProducer.swift b/Sources/AetherEngine/Video/HLSSegmentProducer.swift index 1b00d275..7eca8b7c 100644 --- a/Sources/AetherEngine/Video/HLSSegmentProducer.swift +++ b/Sources/AetherEngine/Video/HLSSegmentProducer.swift @@ -510,6 +510,11 @@ final class HLSSegmentProducer: @unchecked Sendable { /// before `start()`; nil for hosts that drive the engine without one (`aetherctl`, tests). var sideReaderLinkGate: SideReaderLinkGate? + /// AE#514: the session's ledger of played-stream packet bytes, which the bitrate fields read at the + /// playhead. Set once before `start()`, like the gate above; nil for hosts that drive the engine + /// without a session (`aetherctl` probes, tests). + var playedMediaLedger: PlayedMediaLedger? + /// Forward-only producer restart counter; surfaced in live telemetry. Written on pump thread, read under packetCounterLock. var restartCount: Int { packetCounterLock.lock() @@ -2719,10 +2724,36 @@ final class HLSSegmentProducer: @unchecked Sendable { /// while playback listens to the bridged rendition. private func readNextSourcePacketMergedTapped() throws -> (packet: UnsafeMutablePointer, origin: PacketOrigin)? { let read = try readNextSourcePacketMerged() - if let read { tapForRecording(read.packet) } + if let read { + tapForRecording(read.packet) + recordPlayedMedia(read.packet, origin: read.origin) + } return read } + /// AE#514: the video stream and the audio stream this session plays, on the source axis the engine + /// folds AVPlayer's clock back onto (`source_pts - playlistShiftSeconds`). + private func recordPlayedMedia(_ packet: UnsafeMutablePointer, origin: PacketOrigin) { + guard let ledger = playedMediaLedger, packet.pointee.size > 0 else { return } + let index = packet.pointee.stream_index + let track: PlayedMediaLedger.Track + let timeBase: AVRational + if origin == .main, index == videoStreamIndex { + track = .video + timeBase = sourceVideoTimeBase + } else if let audio = audioConfig, index == audio.sourceStreamIndex, + (origin == .side) == (sideAudioDemuxer != nil) { + track = .audio + timeBase = audio.sourceTimeBase + } else { + return + } + let ticks = packet.pointee.pts != Int64.min ? packet.pointee.pts : packet.pointee.dts + guard ticks != Int64.min, timeBase.num > 0, timeBase.den > 0 else { return } + ledger.record(track, pts: Double(ticks) * Double(timeBase.num) / Double(timeBase.den), + bytes: Int(packet.pointee.size)) + } + private func readNextSourcePacketMerged() throws -> (packet: UnsafeMutablePointer, origin: PacketOrigin)? { guard let side = sideAudioDemuxer else { guard let packet = try demuxer.readPacket() else { return nil } diff --git a/Sources/AetherEngine/Video/HLSVideoEngine.swift b/Sources/AetherEngine/Video/HLSVideoEngine.swift index 2fda1e33..96fae422 100644 --- a/Sources/AetherEngine/Video/HLSVideoEngine.swift +++ b/Sources/AetherEngine/Video/HLSVideoEngine.swift @@ -2268,6 +2268,10 @@ public final class HLSVideoEngine: @unchecked Sendable { return (producer, cache, server, demuxer, audioBridge) } + /// AE#514: bytes of the played streams by presentation time, fed by every producer of the session + /// (initial, seek restart, live reopen, revive), read by the telemetry sampler at the playhead. + let playedMediaLedger = PlayedMediaLedger() + /// Bytes this session pulled from the SOURCE, across every demuxer it has had (see /// `retiredDemuxerBytes`). Not the same link as `LiveTelemetry.networkTransferredBytes`, which on /// the native path counts what AVPlayer pulled from the loopback server. @@ -2629,6 +2633,7 @@ public final class HLSVideoEngine: @unchecked Sendable { // below. The side readers read one gate for the whole session, so a restart must not leave // a gap where nobody claims the link. prod.sideReaderLinkGate = sideReaderLinkGate + prod.playedMediaLedger = playedMediaLedger prod.onFirstHDR10PlusDetected = { [weak self] in self?.notifyHDR10PlusOnce() } diff --git a/Tests/AetherEngineTests/Issue514AverageBitrateTests.swift b/Tests/AetherEngineTests/Issue514AverageBitrateTests.swift index 340f2293..8d2ceb2b 100644 --- a/Tests/AetherEngineTests/Issue514AverageBitrateTests.swift +++ b/Tests/AetherEngineTests/Issue514AverageBitrateTests.swift @@ -1,140 +1,210 @@ import Foundation import Testing -import AVFoundation @testable import AetherEngine -/// AE#514: `LiveTelemetry.averageBitrateMbps` divided the session's lifetime bytes by pure wall-clock -/// time. Bytes stop arriving while the transport is paused (the forward buffer is already full), the -/// divisor does not, so the reported average decayed toward zero for as long as the pause lasted and -/// never came back: the paused seconds stayed in the divisor for the rest of the session. The fix -/// charges a tick's second only when the session is in a phase that consumes media. -@MainActor +/// AE#514: both bitrate fields in `LiveTelemetry` were metered from the reader's transfer counter. +/// Round one fixed the divisor (a pause dragged the average toward zero); round two, from the +/// reporter's retest, fixes the numerator: transfer and playback part ways on every route that reads +/// ahead. VOD prefetch put a 20 Mbps stream at ~35 Mbps, and a paused live session kept draining the +/// origin into its DVR window while the divisor stood still, so its average climbed for as long as the +/// pause ran. Both fields now meter the bytes of the packets the playhead crossed over the media +/// seconds it crossed. struct Issue514AverageBitrateTests { - // MARK: - Which seconds belong in the divisor - - /// A pause and the three terminal phases are the session standing still. Everything else is the - /// session working, and the bytes it did or did not get in those seconds are part of its average. - @Test("only the phases that consume media charge a second to the divisor") - func pausedAndTerminalPhasesDoNotCharge() { - #expect(LiveTelemetrySampler.chargesActiveTime(.paused) == false) - #expect(LiveTelemetrySampler.chargesActiveTime(.idle) == false) - #expect(LiveTelemetrySampler.chargesActiveTime(.ended) == false) - #expect(LiveTelemetrySampler.chargesActiveTime(.error("boom")) == false) - - #expect(LiveTelemetrySampler.chargesActiveTime(.playing)) - #expect(LiveTelemetrySampler.chargesActiveTime(.loading)) - #expect(LiveTelemetrySampler.chargesActiveTime(.rebuffering)) - #expect(LiveTelemetrySampler.chargesActiveTime(.stalled(reconnecting: true))) - #expect(LiveTelemetrySampler.chargesActiveTime(.stalled(reconnecting: false))) + /// 20 Mbps split over a 25 fps video track and a ~47 Hz audio track, recorded as a pump would. + private static let videoBytesPerSecond = 2_400_000 + private static let audioBytesPerSecond = 100_000 + private static let mbps = Double(videoBytesPerSecond + audioBytesPerSecond) * 8 / 1_000_000 + + private static func feed(_ ledger: PlayedMediaLedger, from start: Double, to end: Double) { + var pts = start + while pts < end - 1e-9 { + ledger.record(.video, pts: pts, bytes: videoBytesPerSecond / 25) + pts += 1.0 / 25 + } + pts = start + while pts < end - 1e-9 { + ledger.record(.audio, pts: pts, bytes: audioBytesPerSecond * 1024 / 48_000) + pts += 1024.0 / 48_000 + } } - /// Deliberately NOT excluded, against the report's own suggestion: a seek is where the session - /// fetches hardest (the buffer is discarded and a new range pulled). Dropping those seconds while - /// keeping their bytes would push the average above the media's real rate on every scrub. - @Test("a seek charges its seconds, because it is also where the bytes arrive") - func seekingCharges() { - #expect(LiveTelemetrySampler.chargesActiveTime(.seeking)) + private static func play(_ meter: inout PlayedBitrateMeter, _ ledger: PlayedMediaLedger, + from start: Double, seconds: Int) -> Double { + var playhead = start + for _ in 0.. Bool { + guard let value else { return false } + return abs(value - expected) <= tolerance * expected + } - @Test("the average is the lifetime bytes over the active seconds") - func averageIsBytesOverActiveSeconds() { - // 2 Mbps for ten active seconds: 2 500 000 bytes. - let rate = LiveTelemetrySampler.averageBitrateMbps(lifetimeBytes: 2_500_000, activeSeconds: 10) - #expect(abs((rate ?? 0) - 2.0) < 0.001) + // MARK: - The reporter's table + + /// VOD, playing: minutes of read-ahead sit in the ledger the moment they are fetched, and count + /// only once the playhead crosses them. + @Test("VOD read-ahead does not inflate either field") + func vodReadAheadIsNotCounted() { + let ledger = PlayedMediaLedger() + Self.feed(ledger, from: 0, to: 240) + var meter = PlayedBitrateMeter() + meter.advance(to: 0, wallSeconds: 0, ledger: ledger) + _ = Self.play(&meter, ledger, from: 0, seconds: 30) + #expect(Self.near(meter.averageMbps, Self.mbps)) + #expect(Self.near(meter.instantMbps, Self.mbps)) } - /// Same rule as `observedTransferMbps` one field over: "not measurable yet" is a gap, never a - /// confident zero. Before the fix the very first tick published 0.00 Mbps, because it seeds - /// `sessionStartBytes` from the same counter it then subtracts. - @Test("a session with no active time and no bytes publishes nil, not zero") - func unmeasurableSessionPublishesNil() { - #expect(LiveTelemetrySampler.averageBitrateMbps(lifetimeBytes: 0, activeSeconds: 0) == nil) - #expect(LiveTelemetrySampler.averageBitrateMbps(lifetimeBytes: 350_000, activeSeconds: 0) == nil) - #expect(LiveTelemetrySampler.averageBitrateMbps(lifetimeBytes: 0, activeSeconds: 30) == nil) + /// VOD, paused: the reader keeps topping up the buffer, the playhead does not move. + @Test("a VOD pause leaves both fields standing") + func vodPauseFreezes() { + let ledger = PlayedMediaLedger() + Self.feed(ledger, from: 0, to: 60) + var meter = PlayedBitrateMeter() + meter.advance(to: 0, wallSeconds: 0, ledger: ledger) + let playhead = Self.play(&meter, ledger, from: 0, seconds: 30) + let average = meter.averageMbps, instant = meter.instantMbps + Self.feed(ledger, from: 60, to: 120) + for _ in 0..<180 { meter.advance(to: playhead, wallSeconds: 1, ledger: ledger) } + #expect(meter.averageMbps == average) + #expect(meter.instantMbps == instant) } - // MARK: - The reported session shape - - /// The reporter's steps, folded tick by tick: 30 s of a 2.8 Mbps file, a three-minute pause during - /// which the demuxer counter stands still, then playback again. - @Test("a three-minute pause leaves the average flat instead of collapsing it") - func pauseDoesNotDragTheAverageDown() { - let bytesPerSecond: Int64 = 350_000 // 2.8 Mbps - var lifetimeBytes: Int64 = 0 - var activeSeconds: Double = 0 - var wallClockSeconds: Double = 0 - - func tick(_ phase: PlaybackPhase, bytes: Int64) { - lifetimeBytes += bytes - wallClockSeconds += 1 - if LiveTelemetrySampler.chargesActiveTime(phase) { activeSeconds += 1 } + /// Live, paused: the pump keeps draining the origin at the broadcast rate for the whole pause + /// (AE#443). Nothing is played, so nothing may move; on resume the backlog counts as it plays. + @Test("a live pause does not climb, and the backlog counts at its real rate once played") + func livePauseDoesNotClimb() { + let ledger = PlayedMediaLedger() + var meter = PlayedBitrateMeter() + var edge = 1_000.0 + Self.feed(ledger, from: edge, to: edge + 2) + meter.advance(to: edge, wallSeconds: 0, ledger: ledger) + var playhead = edge + for _ in 0..<30 { + edge += 1 + Self.feed(ledger, from: edge + 1, to: edge + 2) + playhead += 1 + meter.advance(to: playhead, wallSeconds: 1, ledger: ledger) } + let average = meter.averageMbps + #expect(Self.near(average, Self.mbps)) - for _ in 0..<30 { tick(.playing, bytes: bytesPerSecond) } - let beforePause = LiveTelemetrySampler.averageBitrateMbps( - lifetimeBytes: lifetimeBytes, activeSeconds: activeSeconds) - #expect(abs((beforePause ?? 0) - 2.8) < 0.001) - - // Paused: the forward buffer is full, so the demuxer counter does not move. - for _ in 0..<180 { tick(.paused, bytes: 0) } - let duringPause = LiveTelemetrySampler.averageBitrateMbps( - lifetimeBytes: lifetimeBytes, activeSeconds: activeSeconds) - #expect(duringPause == beforePause, "the pause must not move the average at all") - - // What used to ship, for the record: the same bytes over wall-clock time. - let wallClockAverage = Double(lifetimeBytes) * 8.0 / wallClockSeconds / 1_000_000.0 - #expect(wallClockAverage < 0.5, "the old divisor collapsed a 2.8 Mbps session to \(wallClockAverage)") - - // And on resume it is still the media's rate, not a value climbing back out of a hole. - for _ in 0..<30 { tick(.playing, bytes: bytesPerSecond) } - let afterResume = LiveTelemetrySampler.averageBitrateMbps( - lifetimeBytes: lifetimeBytes, activeSeconds: activeSeconds) - #expect(abs((afterResume ?? 0) - 2.8) < 0.001) - } + for _ in 0..<300 { + edge += 1 + Self.feed(ledger, from: edge + 1, to: edge + 2) + meter.advance(to: playhead, wallSeconds: 1, ledger: ledger) + } + #expect(meter.averageMbps == average, "five minutes of paused live must not move the average") - /// End of media is the other unbounded divisor: the sampler runs until the host tears the session - /// down, so a snapshot left on screen after the last frame used to decay exactly like a pause. - @Test("the average stops moving once the source has ended") - func endedSessionFreezesTheAverage() { - var activeSeconds: Double = 0 - for _ in 0..<30 where LiveTelemetrySampler.chargesActiveTime(.playing) { activeSeconds += 1 } - let atEnd = activeSeconds - for _ in 0..<120 where LiveTelemetrySampler.chargesActiveTime(.ended) { activeSeconds += 1 } - #expect(activeSeconds == atEnd) + _ = Self.play(&meter, ledger, from: playhead, seconds: 60) + #expect(Self.near(meter.averageMbps, Self.mbps)) + #expect(Self.near(meter.instantMbps, Self.mbps)) } - // MARK: - Wiring - - /// The two functions above are only right if the tick actually asks them. A paused session accrues - /// no active time however long its sampler runs; the same sampler starts accruing on play. - @Test("the running sampler charges no active time while the transport is paused") - func samplerAccruesNoActiveTimeWhilePaused() async throws { - let engine = try AetherEngine() - engine.playbackBackend = .native - let item = AVPlayerItem(url: URL(fileURLWithPath: "/nonexistent-514.mp4")) - engine.currentAVPlayer = AVPlayer(playerItem: item) - engine.state = .paused - #expect(engine.playbackPhase == .paused) - - let sampler = LiveTelemetrySampler(engine: engine, nativeRead: { _, _ in - NativeAVFReadings(forwardBufferSeconds: 12.0) - }) - sampler.start() - defer { sampler.stop() } - - // Long enough for several 1 Hz ticks to have run and charged nothing. - try await Task.sleep(for: .milliseconds(2_500)) - #expect(sampler.activeSeconds == 0) - - engine.state = .playing - let started = ContinuousClock().now - while sampler.activeSeconds == 0 { - if ContinuousClock().now - started > .seconds(30) { break } - try await Task.sleep(for: .milliseconds(50)) + /// Live, playing: this was right before by accident of the clock, and must stay right. + @Test("live at the edge reads the broadcast rate") + func livePlayingReadsTheRate() { + let ledger = PlayedMediaLedger() + var meter = PlayedBitrateMeter() + Self.feed(ledger, from: 50, to: 53) + meter.advance(to: 50, wallSeconds: 0, ledger: ledger) + var playhead = 50.0 + for _ in 0..<60 { + Self.feed(ledger, from: playhead + 3, to: playhead + 4) + playhead += 1 + meter.advance(to: playhead, wallSeconds: 1, ledger: ledger) } - #expect(sampler.activeSeconds > 0) + #expect(Self.near(meter.averageMbps, Self.mbps)) + #expect(Self.near(meter.instantMbps, Self.mbps)) + } + + // MARK: - Seeks + + /// The span a seek jumps over was never played, and the buffer the seek discarded is fetched again: + /// neither may land in the fields. + @Test("a seek charges nothing for the jump, and a re-fetched range counts once") + func seekAndRefetch() { + let ledger = PlayedMediaLedger() + Self.feed(ledger, from: 0, to: 120) + var meter = PlayedBitrateMeter() + meter.advance(to: 0, wallSeconds: 0, ledger: ledger) + _ = Self.play(&meter, ledger, from: 0, seconds: 20) + let secondsBefore = meter.lifetimeSeconds + + // Forward seek to 600: a new range is fetched from there. + Self.feed(ledger, from: 600, to: 700) + meter.advance(to: 600, wallSeconds: 1, ledger: ledger) + #expect(meter.lifetimeSeconds == secondsBefore) + _ = Self.play(&meter, ledger, from: 600, seconds: 10) + + // Back to 100: the reader delivers 100..200 again, over entries it already held. + Self.feed(ledger, from: 100, to: 200) + meter.advance(to: 100, wallSeconds: 1, ledger: ledger) + _ = Self.play(&meter, ledger, from: 100, seconds: 30) + + #expect(abs(meter.lifetimeSeconds - 60) < 1e-6) + #expect(Self.near(meter.averageMbps, Self.mbps)) + } + + /// A playhead the ledger holds nothing for (not fed there, or on another axis) is unmeasured. + @Test("a playhead off the ledger's span publishes nil, never zero") + func unmeasuredSpanIsNil() { + let ledger = PlayedMediaLedger() + Self.feed(ledger, from: 5_000, to: 5_100) + var meter = PlayedBitrateMeter() + meter.advance(to: 0, wallSeconds: 0, ledger: ledger) + _ = Self.play(&meter, ledger, from: 0, seconds: 30) + #expect(meter.averageMbps == nil) + #expect(meter.instantMbps == nil) + } + + // MARK: - The ledger + + @Test("the same packet delivered twice is held once") + func duplicatePacketOverwrites() { + let ledger = PlayedMediaLedger() + ledger.record(.video, pts: 1.0, bytes: 100) + ledger.record(.video, pts: 1.04, bytes: 100) + ledger.record(.video, pts: 1.0, bytes: 100) + #expect(ledger.entryCount == 2) + #expect(ledger.consume(from: 0, to: 2) == 200) + } + + @Test("B-frame reorder inserts in place, a re-read from far behind drops what lies ahead") + func reorderAndReread() { + let ledger = PlayedMediaLedger() + for pts in [0.0, 0.12, 0.04, 0.08, 0.24, 0.16, 0.20] { ledger.record(.video, pts: pts, bytes: 10) } + #expect(ledger.consume(from: 0.04, to: 0.20) == 40) + ledger.record(.video, pts: 10, bytes: 10) + ledger.record(.video, pts: 11, bytes: 10) + ledger.record(.video, pts: 5, bytes: 10) + #expect(ledger.consume(from: 5, to: 12) == 10, "10 and 11 lie ahead of the re-read and are gone") + } + + @Test("consuming a span forgets everything before it") + func consumePrunes() { + let ledger = PlayedMediaLedger() + Self.feed(ledger, from: 0, to: 10) + _ = ledger.consume(from: 4, to: 5) + #expect(ledger.consume(from: 0, to: 5) == 0) + #expect(ledger.consume(from: 5, to: 10) > 0) + } + + // MARK: - The playhead + + @Test("native folds AVPlayer's clock back with the producer's shift; routes without a pump read nil") + func ledgerPlayheadPerBackend() { + #expect(LiveTelemetrySampler.ledgerPlayhead( + backend: .native, nativeClock: 12, playlistShift: 3_600, softwareSourceClock: nil) == 3_612) + #expect(LiveTelemetrySampler.ledgerPlayhead( + backend: .native, nativeClock: nil, playlistShift: 3_600, softwareSourceClock: nil) == nil) + #expect(LiveTelemetrySampler.ledgerPlayhead( + backend: .software, nativeClock: 12, playlistShift: 3_600, softwareSourceClock: 90) == 90) + #expect(LiveTelemetrySampler.ledgerPlayhead( + backend: .audio, nativeClock: 12, playlistShift: 0, softwareSourceClock: 12) == nil) } } diff --git a/docs/api.md b/docs/api.md index 9760f6f3..78714e88 100644 --- a/docs/api.md +++ b/docs/api.md @@ -925,7 +925,7 @@ as well. | Symbol | Notes | | --- | --- | | `AetherEngine.version` | The engine release this source descends from, as the string a published tag carries. SwiftPM resolves a package to a revision rather than to a tag, so an About panel or the header of a diagnostic log has nothing else to name the engine with. Between releases, and under a pin on an unreleased commit, it names the last published version the checkout descends from. | -| `diagnostics.liveTelemetry` | 1 Hz `LiveTelemetry?` snapshot while playing or paused, nil while idle. On a separate `ObservableObject` so its ticks cannot re-render a host observing the engine. On `.remoteBypass` it is fed from AVPlayer's access log alone: `instantBitrateMbps` and `averageBitrateMbps` are what the playing variant declares (BANDWIDTH, and AVERAGE-BANDWIDTH or BANDWIDTH where the master omits it), because the bytes AVPlayer transferred are buffer fill at link speed rather than the stream's rate, `networkThroughputMbps` is the access log's `observedBitrate`, `networkTransferredBytes` and `droppedFrameCount` its session totals, `forwardBufferSeconds` the loaded range ahead of the playhead. `avSyncGapMs` and `observedFps` are nil, and the loopback counters (`producerRestartCount`, `muxedBytesLifetime`, `serverBytesSentLifetime`, `serverRequestCount`, `demuxerBytesFetched`) read 0 because there is no loopback. | +| `diagnostics.liveTelemetry` | 1 Hz `LiveTelemetry?` snapshot while playing or paused, nil while idle. On a separate `ObservableObject` so its ticks cannot re-render a host observing the engine. On the loopback and software routes `instantBitrateMbps` and `averageBitrateMbps` are the rate of the MEDIA played, not of the transfer (AE#514): the bytes of the played video and audio packets the playhead crossed, over the media seconds it crossed, about the last 10 s for the first and the whole session for the second. Read-ahead, a re-fetch after a seek and a paused live source draining into its DVR window do not move them, and both stand still through a pause; the transfer is `networkThroughputMbps`. On `.remoteBypass` it is fed from AVPlayer's access log alone: `instantBitrateMbps` and `averageBitrateMbps` are what the playing variant declares (BANDWIDTH, and AVERAGE-BANDWIDTH or BANDWIDTH where the master omits it), because the bytes AVPlayer transferred are buffer fill at link speed rather than the stream's rate, `networkThroughputMbps` is the access log's `observedBitrate`, `networkTransferredBytes` and `droppedFrameCount` its session totals, `forwardBufferSeconds` the loaded range ahead of the playhead. `avSyncGapMs` and `observedFps` are nil, and the loopback counters (`producerRestartCount`, `muxedBytesLifetime`, `serverBytesSentLifetime`, `serverRequestCount`, `demuxerBytesFetched`) read 0 because there is no loopback. | | `LiveTelemetry.softwareCacheSeekHits`, `softwareCacheSeekMisses`, `softwareCacheSourceEpoch` | Optional cumulative software-VOD packet-cache counters. A hit repositions the retained consumer cursor without changing the source epoch; a miss repositions the demuxer and advances it. `nil` on other paths. `cachedBytes` includes retained compressed packet records on software VOD, distinct from decoded `displayCushionSeconds` and the underlying byte-reader window. | | `EngineLog.handler` | Mirror every info-level line into a host capture path. Fires from whatever thread emitted it, so it must be thread-safe and non-blocking. | | `EngineLog.subsystem`, `EngineLog.Category` | `de.superuser404.AetherEngine`, one category per subsystem: `engine`, `ffmpeg`, `session`, `muxer`, `demux`, `hls.server`, `audio.bridge`, `sw.playback`, `scrub`. |