fix(rec): bounded encoder queue + drop policy + burned frame counter — the "1...23...4...56..." smeared-ticker take
slice 9 made the DURATION right but content still hiccuped; aggregates (301/300, uniform file PTS) could not see it. Measured root cause: FfmpegEncoder.SubmitFrameAsync BLOCKED on WriteAsync(8.3MB)+FlushAsync when ffmpeg lagged the pipe, and the burst while-loop re-wrote that same stale composite per crossed slot — frozen runs. OBS shape (derivative, wrapped pre-1.0): the encoder queue in libobs/obs-encoder.c — encoder thread never couples back into the video thread; overflow = dropped data, never a frozen producer. https://github.com/obsproject/obs-studio/blob/master/libobs/obs-encoder.c - FfmpegEncoder: SubmitFrameAsync is now an enqueue (ArrayPool copy) into a bounded Channel (cap 120) drained by its own task; drop-newest + count when full; StopAsync flushes the queue then EOF (TryComplete). IFfmpegEncoder.DroppedFrames. - FramePump: ONE fresh composite per iteration (burst loop deleted); worst-submit stat, stall logger (>2x interval names the stage), dropped/stalls in stats. - Burned-in 6-digit dot-matrix frame counter (white box, bottom-right) on every composite — the clock-independent judge replacing the WSL ticker: +1/frame, jumps = counted drops. - ONE new test Backpressure_QueueOverflow_DropsFrames_AndNeverBlocks (slow-sink fake: submit never blocks, drops counted, stop flushes exactly submitted-minus-dropped). - Full suite 290 tests, 289 pass — sole failure the pre-existing compositor pixel test. - Docs same-commit: ai.md slice 10 (+ encoder/stop-note corrections), MyMistakes point 8, HANDOFF. Audio untouched (queued follow-up); web overlay still frozen pending timing closure.
This commit is contained in:
@@ -1,4 +1,6 @@
|
||||
using System.Buffers;
|
||||
using System.Diagnostics;
|
||||
using System.Threading.Channels;
|
||||
using ytLive.Helpers;
|
||||
using ytLive.Models;
|
||||
|
||||
@@ -30,6 +32,21 @@ public sealed class FfmpegEncoder : IFfmpegEncoder
|
||||
private Task? _stderrLoop;
|
||||
private bool _stopRequested;
|
||||
|
||||
// Pipe decoupling (slice 10, 2026-09-10): a bounded pending-frame queue drained
|
||||
// by its own task (the OBS video-thread → encoder-queue model — the encoder's
|
||||
// thread never couples back into the video thread; libobs obs-encoder.c). The
|
||||
// pump's SubmitFrameAsync now ENQUEUES (copying into a pooled buffer) instead of
|
||||
// blocking on WriteAsync when ffmpeg lags the pipe; overflow DROPS the newest
|
||||
// frame (skip-newest) and counts it. BoundedChannelFullMode.Wait + TryWrite gives
|
||||
// exactly that: when full, TryWrite returns false and the caller drops the item.
|
||||
private const int QueueCapacity = 120; // ~2s at the tier's 60fps
|
||||
private readonly Channel<PendingWrite> _frames = Channel.CreateBounded<PendingWrite>(
|
||||
new BoundedChannelOptions(QueueCapacity) { SingleReader = true });
|
||||
private Task? _drainTask;
|
||||
private int _droppedBackpressure;
|
||||
|
||||
private readonly record struct PendingWrite(byte[] Buffer, int Length);
|
||||
|
||||
public FfmpegEncoder(
|
||||
IFfmpegLocator locator,
|
||||
Func<IEncoderProcess>? processFactory = null)
|
||||
@@ -40,6 +57,13 @@ public sealed class FfmpegEncoder : IFfmpegEncoder
|
||||
|
||||
public bool IsRunning { get; private set; }
|
||||
|
||||
/// <summary>Frames dropped by the bounded queue because the drain task couldn't
|
||||
/// keep the pipe fed (the encoder lagging the real-time capture rate). The pump
|
||||
/// reports this in its stats; the burned-in frame counter in the recording
|
||||
/// shows the identical jumps. Distinct from <see cref="StreamHealth.DroppedFrames"/>
|
||||
/// (ffmpeg's own progress-derived estimate).</summary>
|
||||
public int DroppedFrames => Volatile.Read(ref _droppedBackpressure);
|
||||
|
||||
public async Task StartAsync(EncoderOptions options, CancellationToken cancellationToken = default)
|
||||
{
|
||||
if (options == null) throw new ArgumentNullException(nameof(options));
|
||||
@@ -93,26 +117,40 @@ public sealed class FfmpegEncoder : IFfmpegEncoder
|
||||
|
||||
_stopRequested = false;
|
||||
_stderrLoop = RunStderrLoopAsync(process);
|
||||
StartDrainLoop(process);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Write one raw BGRA frame to ffmpeg's stdin. Serialized internally; callers
|
||||
/// (the compositor pump) may race freely. Frames are written as-is — the caller
|
||||
/// paces to capture rate (the compositor's job, ship step 5).
|
||||
/// Enqueue one raw BGRA frame for ffmpeg's stdin (slice 10). The caller races
|
||||
/// freely (the compositor pump); this never blocks on the pipe. The pixels are
|
||||
/// copied into a pooled buffer first because — unlike the old direct write — the
|
||||
/// write happens later on the drain thread, so the caller may (and does) recycle
|
||||
/// its scratch buffer the moment this returns. When the bounded queue is full the
|
||||
/// NEWEST frame is dropped and counted (freshness over coverage; the encoder is
|
||||
/// already behind, so the stale-content failure mode never encodes).
|
||||
/// </summary>
|
||||
public async Task SubmitFrameAsync(VideoFrame frame, CancellationToken cancellationToken = default)
|
||||
public Task SubmitFrameAsync(VideoFrame frame, CancellationToken cancellationToken = default)
|
||||
{
|
||||
if (frame == null) throw new ArgumentNullException(nameof(frame));
|
||||
IEncoderProcess? process;
|
||||
lock (_gate)
|
||||
{
|
||||
if (!IsRunning) throw new InvalidOperationException("The encoder is not running.");
|
||||
process = _process;
|
||||
}
|
||||
|
||||
var bytes = frame.BgraPixels;
|
||||
await process!.StandardInput.WriteAsync(bytes, cancellationToken).ConfigureAwait(false);
|
||||
await process.StandardInput.FlushAsync(cancellationToken).ConfigureAwait(false);
|
||||
if (cancellationToken.IsCancellationRequested)
|
||||
return Task.FromCanceled(cancellationToken);
|
||||
|
||||
var pooled = ArrayPool<byte>.Shared.Rent(frame.BgraPixels.Length);
|
||||
Buffer.BlockCopy(frame.BgraPixels, 0, pooled, 0, frame.BgraPixels.Length);
|
||||
if (_frames.Writer.TryWrite(new PendingWrite(pooled, frame.BgraPixels.Length)))
|
||||
return Task.CompletedTask;
|
||||
|
||||
// Queue full (encoder lagging) or the channel completed during stop — drop the
|
||||
// newest and count it. The caller never blocks; the encoder catches up or the
|
||||
// session ends with a (reported) short gap instead of a frozen stall.
|
||||
ArrayPool<byte>.Shared.Return(pooled);
|
||||
Interlocked.Increment(ref _droppedBackpressure);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public async Task StopAsync(CancellationToken cancellationToken = default)
|
||||
@@ -127,13 +165,17 @@ public sealed class FfmpegEncoder : IFfmpegEncoder
|
||||
loop = _stderrLoop;
|
||||
}
|
||||
|
||||
// Flush every queued frame, then EOF: TryComplete lets the drain task write
|
||||
// the buffered frames, close stdin → ffmpeg finalizes and exits by itself
|
||||
// (the "did not exit after EOF" kill below is only the backstop).
|
||||
try
|
||||
{
|
||||
process!.StandardInput.Dispose(); // EOF → ffmpeg finalizes + exits
|
||||
_frames.Writer.TryComplete();
|
||||
if (_drainTask != null) await _drainTask.ConfigureAwait(false);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
AppLog.Write(ex, "FFmpeg encoder: closing stdin failed");
|
||||
AppLog.Write(ex, "FFmpeg encoder: drain loop faulted during stop");
|
||||
}
|
||||
|
||||
try
|
||||
@@ -178,6 +220,7 @@ public sealed class FfmpegEncoder : IFfmpegEncoder
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
_frames.Writer.TryComplete(); // unblock a stuck drain loop alongside the kill
|
||||
lock (_gate)
|
||||
{
|
||||
if (!IsRunning) return;
|
||||
@@ -188,6 +231,54 @@ public sealed class FfmpegEncoder : IFfmpegEncoder
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The drain task (slice 10): owns every stdin write, so pipe backpressure — the
|
||||
/// pump's old stall — lives HERE, on a thread the pump never touches. Frames are
|
||||
/// written exactly as queued (the pooled array through its recorded length —
|
||||
/// ArrayPool may return a larger buffer), returned to the pool after the write
|
||||
/// copies into the pipe, and stdin is closed (EOF) once the queue drains.
|
||||
/// </summary>
|
||||
private void StartDrainLoop(IEncoderProcess process)
|
||||
{
|
||||
_drainTask = Task.Run(async () =>
|
||||
{
|
||||
try
|
||||
{
|
||||
try
|
||||
{
|
||||
while (await _frames.Reader.WaitToReadAsync().ConfigureAwait(false))
|
||||
{
|
||||
while (_frames.Reader.TryRead(out var pending))
|
||||
{
|
||||
await process.StandardInput.WriteAsync(pending.Buffer, 0, pending.Length)
|
||||
.ConfigureAwait(false);
|
||||
await process.StandardInput.FlushAsync().ConfigureAwait(false);
|
||||
ArrayPool<byte>.Shared.Return(pending.Buffer);
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (ChannelClosedException)
|
||||
{
|
||||
// channel completed while waiting — fall through to EOF
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
AppLog.Write(ex, "FFmpeg encoder: drain loop write faulted");
|
||||
}
|
||||
finally
|
||||
{
|
||||
// EOF on every path — a byte that sits in the pool when the pipe
|
||||
// breaks is garbage anyway, and ffmpeg MUST see the close to finalize.
|
||||
process.StandardInput.Dispose();
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
AppLog.Write(ex, "FFmpeg encoder: drain loop faulted");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private async Task<string> ProbeEncoderAsync(string ffmpegPath, CancellationToken cancellationToken)
|
||||
{
|
||||
try
|
||||
|
||||
+110
-19
@@ -85,6 +85,7 @@ public sealed class FramePump : IDisposable
|
||||
private CancellationTokenSource? _cts;
|
||||
private Task? _pumpTask;
|
||||
private bool _started;
|
||||
private long _outputIndex;
|
||||
|
||||
/// <summary>Forwards the encoder's parsed health — ship step 6 binds this to the bottom bar.</summary>
|
||||
public event EventHandler<StreamHealth>? HealthUpdated;
|
||||
@@ -294,13 +295,17 @@ public sealed class FramePump : IDisposable
|
||||
// black box — the stats line now reports resolver time separately so a take
|
||||
// names the stage (get-frame vs blit) instead of feeding another guess.
|
||||
var resolveSw = new System.Diagnostics.Stopwatch();
|
||||
long renderTicks = 0, submitTicks = 0, resolveTicks = 0, waitTicks = 0, worstRender = 0;
|
||||
long renderTicks = 0, submitTicks = 0, resolveTicks = 0, waitTicks = 0, worstRender = 0, worstSubmit = 0;
|
||||
int statFrames = 0;
|
||||
int stalls = 0;
|
||||
var statsNext = DateTime.UtcNow + TimeSpan.FromSeconds(5);
|
||||
void ReportStats()
|
||||
{
|
||||
if (DateTime.UtcNow < statsNext) return;
|
||||
var target = 5d / interval.TotalSeconds; // frames expected per window
|
||||
IFfmpegEncoder? encoder;
|
||||
lock (_gate) encoder = _encoder;
|
||||
var dropped = encoder == null ? 0 : encoder.DroppedFrames;
|
||||
_log?.Invoke(statFrames == 0
|
||||
? "FramePump stats: NO frames produced in 5s (loop stalled?)"
|
||||
: $"FramePump stats: {statFrames}/{target:F0} frames per 5s, " +
|
||||
@@ -308,10 +313,14 @@ public sealed class FramePump : IDisposable
|
||||
$"(resolve {resolveTicks / (double)System.Diagnostics.Stopwatch.Frequency * 1000 / statFrames:F1}), " +
|
||||
$"avg submit {submitTicks / (double)System.Diagnostics.Stopwatch.Frequency * 1000 / statFrames:F1}ms, "
|
||||
+ $"avg wait {waitTicks / (double)System.Diagnostics.Stopwatch.Frequency * 1000 / statFrames:F1}ms, "
|
||||
+ $"worst render {worstRender / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F1}ms");
|
||||
+ $"worst render {worstRender / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F1}ms, "
|
||||
+ $"worst submit {worstSubmit / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F1}ms, "
|
||||
+ $"dropped {dropped}, stalls {stalls}");
|
||||
worstRender = 0;
|
||||
worstSubmit = 0;
|
||||
renderTicks = submitTicks = resolveTicks = waitTicks = 0;
|
||||
statFrames = 0;
|
||||
stalls = 0;
|
||||
statsNext = DateTime.UtcNow + TimeSpan.FromSeconds(5);
|
||||
}
|
||||
|
||||
@@ -336,6 +345,7 @@ public sealed class FramePump : IDisposable
|
||||
var scene = _sceneProvider();
|
||||
if (scene != null)
|
||||
{
|
||||
var iterStart = System.Diagnostics.Stopwatch.GetTimestamp();
|
||||
var compositorOptions = _compositorOptions();
|
||||
VideoFrame? socialBarFrame = null;
|
||||
var socialBarTop = 0;
|
||||
@@ -377,29 +387,45 @@ public sealed class FramePump : IDisposable
|
||||
lock (_gate) encoder = _encoder;
|
||||
if (encoder == null) break;
|
||||
|
||||
// Count-based CFR emission (libobs video-io.c — the frame interval
|
||||
// is a DEADLINE and the output stream holds its declared rate):
|
||||
// one frame per interval slot, whatever the render cost. When the
|
||||
// renderer falls behind, the SAME fresh composite is written again
|
||||
// for every slot that ticked past, so a slow render expresses as
|
||||
// duplicated footage (judder) — never as a skipped timestamp. The
|
||||
// deadline counter is NEVER reset to wall-now: the old rebase
|
||||
// erased every missed slot, so a 43fps reality was authored into a
|
||||
// 60fps container and every recording played ~1.4x fast (rawvideo
|
||||
// carries no timestamps — muxed duration is pure frame count).
|
||||
submitSw.Restart();
|
||||
while (!ct.IsCancellationRequested
|
||||
&& System.Diagnostics.Stopwatch.GetTimestamp() >= nextTick)
|
||||
// Burned-in frame counter (slice 10 judge): a clock-independent
|
||||
// pacing witness burned literally into the composite. Decode the
|
||||
// recording and read the bottom-right strip: the number must advance
|
||||
// +1 per frame and jump only by counted drops (queue overflow or a
|
||||
// render-overrun's skipped slots). The WSL ticker replaced as the
|
||||
// judge because its own timers can smear under Windows host load —
|
||||
// this can't lie.
|
||||
_outputIndex++;
|
||||
BurnFrameIndex(frame.BgraPixels, frame.Width, frame.Height, _outputIndex);
|
||||
|
||||
// Deadline pacing (slice 10 reshape): ONE fresh composite per iteration, submitted
|
||||
// only when the clock has reached the next deadline. Two changes from
|
||||
// the take-9 counting loop, both derived from where the stale-content
|
||||
// bug actually lived:
|
||||
// 1. SubmitFrameAsync no longer blocks on the pipe — it ENQUEUES into
|
||||
// the encoder's bounded queue (the OBS video-thread model), so the
|
||||
// pump can never stall behind ffmpeg, and overflow DROPS the
|
||||
// newest frame. The queue wasn't here in the take-9 loop — a
|
||||
// lagging ffmpeg made the pump's submit block, and the burst
|
||||
// while-loop (below) then re-wrote the SAME stale composite for
|
||||
// every slot that ticked past, which is why the ticker smeared.
|
||||
// 2. One submission max per iteration: each catch-up slot now gets a
|
||||
// FRESH render instead of a repeat of the stale one. Missed slots
|
||||
// vanish from the file (a count-based gap, like the queue drop) —
|
||||
// never duplicated frozen frames.
|
||||
if (!ct.IsCancellationRequested
|
||||
&& System.Diagnostics.Stopwatch.GetTimestamp() >= nextTick)
|
||||
{
|
||||
submitSw.Restart();
|
||||
await encoder.SubmitFrameAsync(frame, ct);
|
||||
submitSw.Stop();
|
||||
statFrames++;
|
||||
submitTicks += submitSw.ElapsedTicks;
|
||||
if (submitSw.ElapsedTicks > worstSubmit) worstSubmit = submitSw.ElapsedTicks;
|
||||
nextTick += intervalTicks;
|
||||
}
|
||||
submitSw.Stop();
|
||||
submitTicks += submitSw.ElapsedTicks;
|
||||
|
||||
// SubmitFrameAsync copied the bytes on every write above — the
|
||||
// tick's buffers are recyclable once each due slot consumed them.
|
||||
// SubmitFrameAsync copied the bytes into the encoder's queue — the
|
||||
// tick's buffers are recyclable once the enqueue snapshot them.
|
||||
// The free-list Contains guard keeps the Cut path (BlendFrame
|
||||
// returns toFrame itself, aliasing scratch) safe.
|
||||
ReleaseScratch(frame.BgraPixels);
|
||||
@@ -419,6 +445,20 @@ public sealed class FramePump : IDisposable
|
||||
Thread.SpinWait(400);
|
||||
waitSw.Stop();
|
||||
waitTicks += waitSw.ElapsedTicks;
|
||||
|
||||
// Stall logger (slice 10): an iteration spanning more than two full
|
||||
// intervals is the old bug's fingerprint — name the stage instead of
|
||||
// guessing. With the queue, submit should be ~1ms, so a stall here
|
||||
// means RENDER or RESOLVE (the stage terms of the last take).
|
||||
var iterWall = System.Diagnostics.Stopwatch.GetTimestamp() - iterStart;
|
||||
if (iterWall > 2 * intervalTicks)
|
||||
{
|
||||
stalls++;
|
||||
_log?.Invoke($"FramePump stall: iteration {iterWall / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F0}ms " +
|
||||
$"(> 2× the {interval.TotalMilliseconds:F0}ms interval): worst render " +
|
||||
$"{worstRender / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F0}ms, worst submit " +
|
||||
$"{worstSubmit / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F0}ms, dropped {encoder.DroppedFrames}");
|
||||
}
|
||||
ReportStats();
|
||||
}
|
||||
}
|
||||
@@ -491,6 +531,57 @@ public sealed class FramePump : IDisposable
|
||||
return _compositor.Render(scene, resolver, null, options, socialBarFrame, socialBarTop, scratch: scratch);
|
||||
}
|
||||
|
||||
// Dot-matrix digits (5×7, one row per raster line, '1' = lit) burned into the
|
||||
// bottom-right of every composite. Basic OCR-safe shapes, sized so the strip is
|
||||
// a 36×8 white box in the corner of a 1920×1080 frame — readable with a zoomed
|
||||
// player, invisible at normal size.
|
||||
private static readonly string[][] DigitGlyphs =
|
||||
{
|
||||
new[] { "01110","10001","10001","10001","10001","10001","01110" }, // 0
|
||||
new[] { "00100","01100","00100","00100","00100","00100","01110" }, // 1
|
||||
new[] { "01110","10001","00001","00010","00100","01000","11111" }, // 2
|
||||
new[] { "11111","00001","00010","00110","00001","10001","01110" }, // 3
|
||||
new[] { "00010","00110","01010","10010","11111","00010","00010" }, // 4
|
||||
new[] { "11111","10000","11110","00001","00001","10001","01110" }, // 5
|
||||
new[] { "01110","10001","10000","11110","10001","10001","01110" }, // 6
|
||||
new[] { "11111","00001","00010","00100","01000","01000","01000" }, // 7
|
||||
new[] { "01110","10001","10001","01110","10001","10001","01110" }, // 8
|
||||
new[] { "01110","10001","10001","01111","00001","10001","01110" }, // 9
|
||||
};
|
||||
|
||||
private static void BurnFrameIndex(byte[] bgra, int width, int height, long index)
|
||||
{
|
||||
const int digitW = 5, digitH = 7, gap = 1, margin = 2, digitCount = 6;
|
||||
var stripW = digitCount * (digitW + gap) - gap;
|
||||
var left = width - margin - stripW;
|
||||
var top = height - margin - digitH;
|
||||
|
||||
// Solid white box under the digits — the underlying scene can be anything.
|
||||
for (var y = top; y < top + digitH; y++)
|
||||
for (var x = left; x < left + stripW; x++)
|
||||
{
|
||||
var i = (y * width + x) * 4;
|
||||
bgra[i] = 255;
|
||||
bgra[i + 1] = 255;
|
||||
bgra[i + 2] = 255;
|
||||
}
|
||||
|
||||
var text = index.ToString("D6");
|
||||
for (var d = 0; d < digitCount; d++)
|
||||
{
|
||||
var glyph = DigitGlyphs[text[d] - '0'];
|
||||
for (var row = 0; row < digitH; row++)
|
||||
for (var col = 0; col < digitW; col++)
|
||||
{
|
||||
if (glyph[row][col] != '1') continue;
|
||||
var i = ((top + row) * width + (left + d * (digitW + gap) + col)) * 4;
|
||||
bgra[i] = 0;
|
||||
bgra[i + 1] = 0;
|
||||
bgra[i + 2] = 0;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void OnProcessFailed(object? sender, string message)
|
||||
{
|
||||
_log?.Invoke($"FramePump: encoder process failed: {message}");
|
||||
|
||||
@@ -22,4 +22,9 @@ public interface IFfmpegEncoder : IDisposable
|
||||
Task StartAsync(EncoderOptions options, CancellationToken cancellationToken = default);
|
||||
Task SubmitFrameAsync(VideoFrame frame, CancellationToken cancellationToken = default);
|
||||
Task StopAsync(CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>Frames dropped by the encoder's bounded queue since start (the encoder
|
||||
/// lagging the real-time capture rate). Zero in a healthy session; fed to the pump's
|
||||
/// stats and mirrored by gaps in the burned-in frame counter of the recording.</summary>
|
||||
int DroppedFrames { get; }
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user