Files
gramps c45cbc93b7 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.
2026-09-10 09:32:55 -07:00

421 lines
16 KiB
C#
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
using System.Diagnostics;
using System.Threading.Channels;
using Xunit;
using ytLive.Models;
using ytLive.Services;
using ytLive.Services.Encoder;
namespace ytLive.Tests;
/// <summary>
/// The FFmpeg subprocess encoder (TASK 4 ship step 3). The integration test drives
/// the full lifecycle against fakes: encoder probe → subprocess start with the right
/// args → raw frames into stdin → stderr progress parsed into health → graceful stop.
/// The units pin down the pure pieces (args, progress parser, encoder picker) and
/// the failure edges. No real ffmpeg binary, no network.
/// </summary>
public class FfmpegEncoderTests
{
private sealed class QueuedReader : TextReader
{
private readonly Channel<string?> _channel = Channel.CreateUnbounded<string?>();
public void Enqueue(string s) => _channel.Writer.TryWrite(s);
public void Complete() => _channel.Writer.TryComplete();
public override string? ReadLine() => ReadLineAsync().GetAwaiter().GetResult();
public override async Task<string?> ReadLineAsync()
{
try
{
return await _channel.Reader.ReadAsync().AsTask().ConfigureAwait(false);
}
catch (ChannelClosedException)
{
return null; // stream EOF — the process exited and the pipe closed
}
}
}
private sealed class FakeEncoderProcess : IEncoderProcess
{
public ProcessStartInfo? StartInfo { get; private set; }
public Stream Stdin { get; }
public Stream StandardInput => Stdin;
public TextReader StandardOutput { get; }
public bool HasExited { get; private set; }
public int ExitCode { get; private set; }
public bool Killed { get; private set; }
public bool Started { get; private set; }
private readonly TaskCompletionSource _exit =
new(TaskCreationOptions.RunContinuationsAsynchronously);
private readonly QueuedReader _error = new();
public FakeEncoderProcess(string probeOutput = "", Stream? stdin = null)
{
StandardOutput = new StringReader(probeOutput);
Stdin = stdin ?? new MemoryStream();
}
public TextReader StandardError => _error;
public void EnqueueStderr(string line) => _error.Enqueue(line);
public void Start(ProcessStartInfo startInfo)
{
StartInfo = startInfo;
Started = true;
}
public void SignalExit(int code = 0)
{
ExitCode = code;
HasExited = true;
_error.Complete();
_exit.TrySetResult();
}
public void Kill()
{
Killed = true;
_error.Complete();
_exit.TrySetResult();
}
public Task WaitForExitAsync(CancellationToken cancellationToken = default) => _exit.Task;
public void Dispose()
{
_error.Complete();
_exit.TrySetResult();
}
}
private sealed class FakeEncoderProcessFactory
{
private readonly Queue<FakeEncoderProcess> _processes = new();
public void Return(FakeEncoderProcess p) => _processes.Enqueue(p);
public FakeEncoderProcess Create() => _processes.Dequeue();
}
private sealed class StubLocator : IFfmpegLocator
{
public Task<string> LocateAsync(CancellationToken cancellationToken = default)
=> Task.FromResult(@"C:\tools\ffmpeg.exe");
}
private static byte[] BgraFrame(int w, int h, byte b, byte g, byte r)
{
var bytes = new byte[w * h * 4];
for (var i = 0; i < bytes.Length; i += 4)
{
bytes[i] = b;
bytes[i + 1] = g;
bytes[i + 2] = r;
bytes[i + 3] = 255;
}
return bytes;
}
[Fact]
public async Task Start_ProbesEncoder_FeedsFrames_ParsesHealth_StopsGracefully()
{
const string encoders =
" V..... h264_nvenc NVIDIA NVENC H.264 encoder (codec h264)\n" +
" V..... libopenh264 OpenH264 H.264 (codec h264)\n";
var probe = new FakeEncoderProcess(encoders);
var encoderProc = new FakeEncoderProcess();
var factory = new FakeEncoderProcessFactory();
factory.Return(probe);
factory.Return(encoderProc);
using var encoder = new FfmpegEncoder(new StubLocator(), factory.Create);
var health = new TaskCompletionSource<StreamHealth>(TaskCreationOptions.RunContinuationsAsynchronously);
encoder.HealthUpdated += (_, h) => health.TrySetResult(h);
var options = new EncoderOptions
{
RtmpUrl = "rtmp://a.rtmp.youtube.com/live2/abc-xyz",
Width = 1920,
Height = 1080,
Fps = 60,
BitrateKbps = 8000,
};
await encoder.StartAsync(options);
Assert.True(probe.Started, "the -encoders probe must run");
Assert.True(encoderProc.Started, "the encoder subprocess must start");
Assert.Equal(@"C:\tools\ffmpeg.exe", encoderProc.StartInfo!.FileName);
var args = encoderProc.StartInfo.ArgumentList.ToArray();
var c = Array.IndexOf(args, "-c:v");
Assert.True(
args[c + 1] == "h264_nvenc",
"hardware NVENC must be preferred over the listed openh264");
await encoder.SubmitFrameAsync(new VideoFrame(1920, 1080, BgraFrame(1920, 1080, 255, 0, 0)));
await encoder.SubmitFrameAsync(new VideoFrame(1920, 1080, BgraFrame(1920, 1080, 0, 255, 0)));
// The drain task writes async — poll for both frames instead of racing it.
const long twoFrames = 2 * 1920L * 1080 * 4;
var deadline = DateTime.UtcNow.AddSeconds(5);
while (encoderProc.Stdin.Length < twoFrames && DateTime.UtcNow < deadline)
await Task.Delay(10);
Assert.Equal(twoFrames, encoderProc.Stdin.Length);
encoderProc.EnqueueStderr(
"frame= 120 fps= 59.9 q=28.0 size= 1024KiB time=00:00:02.00 bitrate= 4000.1kbits/s speed=1.00x");
var h = await health.Task.WaitAsync(TimeSpan.FromSeconds(5));
Assert.Equal(4000.1, h.CurrentBitrate, 1);
Assert.Equal(59.9, h.FPS, 1);
Assert.Equal(TimeSpan.FromSeconds(2), h.StreamDuration);
encoderProc.SignalExit();
await encoder.StopAsync();
Assert.False(encoder.IsRunning);
Assert.False(encoderProc.Killed, "graceful stop must not kill the process");
Assert.Equal(StreamStatus.Offline, h.Status);
}
[Fact]
public async Task Start_WithoutRtmpUrl_Throws()
{
using var encoder = new FfmpegEncoder(new StubLocator());
await Assert.ThrowsAsync<ArgumentException>(() => encoder.StartAsync(new EncoderOptions()));
}
[Fact]
public void Args_RecordOnly_SkipsRtmp_WritesMp4()
{
var args = FfmpegArgs.Build(
new EncoderOptions
{
StreamEnabled = false,
RecordEnabled = true,
RecordPath = @"C:\videos\ty-20260829-1405-0000.mp4",
Width = 1920,
Height = 1080,
Fps = 60,
BitrateKbps = 8000,
},
"libopenh264").ToArray();
Assert.DoesNotContain("flv", args);
Assert.DoesNotContain("rtmp", args);
var idx = Array.IndexOf(args, "mp4");
Assert.True(idx > -1, "record-only must emit a -f mp4 output");
Assert.Equal("-f", args[idx - 1]);
Assert.Equal(@"C:\videos\ty-20260829-1405-0000.mp4", args[idx + 1]);
}
[Fact]
public void Args_DualOutput_HasTwoMapBlocks_FlvAndMp4()
{
var args = FfmpegArgs.Build(
new EncoderOptions
{
RtmpUrl = "rtmp://a.rtmp.youtube.com/live2/key",
StreamEnabled = true,
RecordEnabled = true,
RecordPath = @"C:\videos\ty-20260829-1405-0000.mp4",
Width = 1280,
Height = 720,
Fps = 30,
BitrateKbps = 4000,
},
"h264_nvenc").ToArray();
// two video maps (one per output block)
Assert.Equal(2, args.Count(a => a == "0:v"));
Assert.Equal(2, args.Count(a => a == "-c:v"));
Assert.Contains("flv", args);
Assert.Contains("mp4", args);
Assert.Contains("rtmp://a.rtmp.youtube.com/live2/key", args);
Assert.Contains(@"C:\videos\ty-20260829-1405-0000.mp4", args);
}
[Fact]
public async Task SubmitFrame_WhenNotRunning_Throws()
{
using var encoder = new FfmpegEncoder(new StubLocator());
await Assert.ThrowsAsync<InvalidOperationException>(
() => encoder.SubmitFrameAsync(new VideoFrame(2, 2, new byte[16])));
}
[Fact]
public async Task Stop_WithoutStart_IsNoop()
{
using var encoder = new FfmpegEncoder(new StubLocator());
await encoder.StopAsync();
Assert.False(encoder.IsRunning);
}
[Fact]
public async Task ProcessDeath_RaisesFailed()
{
var encoderProc = new FakeEncoderProcess();
var factory = new FakeEncoderProcessFactory();
factory.Return(new FakeEncoderProcess());
factory.Return(encoderProc);
using var encoder = new FfmpegEncoder(new StubLocator(), factory.Create);
var failed = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
encoder.ProcessFailed += (_, msg) => failed.TrySetResult(msg);
await encoder.StartAsync(new EncoderOptions { RtmpUrl = "rtmp://x/y" });
encoderProc.SignalExit(1);
var msg = await failed.Task.WaitAsync(TimeSpan.FromSeconds(5));
Assert.Contains("1", msg);
}
[Fact]
public void Args_Gop_IsFpsTimesFour()
{
var options = new EncoderOptions { Fps = 60 };
Assert.Equal(240, options.GopSize);
var args = FfmpegArgs.Build(options, "libopenh264").ToArray();
var g = Array.IndexOf(args, "-g");
Assert.Equal("240", args[g + 1]);
var keyint = Array.IndexOf(args, "-keyint_min");
Assert.Equal("240", args[keyint + 1]);
}
[Fact]
public void Args_IncludeInputOutputAndCompliance()
{
var args = FfmpegArgs.Build(
new EncoderOptions
{
RtmpUrl = "rtmp://a.rtmp.youtube.com/live2/key",
Width = 1920,
Height = 1080,
Fps = 60,
BitrateKbps = 8000,
},
"libopenh264").ToArray();
Assert.Contains("-f", args);
Assert.Contains("rawvideo", args);
Assert.Contains("pipe:0", args);
Assert.Contains("1920x1080", args);
Assert.Contains(@"\\.\pipe\ytllive_audio", args);
Assert.Contains("f32le", args);
Assert.Contains("-map", args);
Assert.Contains("0:v", args);
Assert.Contains("1:a", args);
Assert.Contains("-sc_threshold", args);
Assert.Contains("-bf", args);
Assert.Contains("yuv420p", args);
Assert.Contains("aac", args);
Assert.Contains("flv", args);
Assert.Contains("rtmp://a.rtmp.youtube.com/live2/key", args);
}
[Fact]
public void ProgressParser_ParsesRealStatsLine()
{
var p = FfmpegProgressParser.TryParse(
"frame= 123 fps= 59.9 q=28.0 size= 1024KiB time=00:00:02.04 bitrate= 4000.1kbits/s speed=1.00x");
Assert.NotNull(p);
Assert.Equal(123, p.Value.Frame);
Assert.Equal(59.9, p.Value.Fps, 1);
Assert.Equal(4000.1, p.Value.BitrateKbps, 1);
Assert.Equal(2.04, p.Value.Duration.TotalSeconds, 2);
Assert.Equal(1024 * 1024, p.Value.SizeBytes);
}
[Fact]
public void ProgressParser_IgnoresBannerAndErrors()
{
Assert.Null(FfmpegProgressParser.TryParse("ffmpeg version 6.1 Copyright (c) 2000-2026 the FFmpeg developers"));
Assert.Null(FfmpegProgressParser.TryParse("Error while opening encoder for output stream #0:0"));
Assert.Null(FfmpegProgressParser.TryParse(""));
}
[Fact]
public void EncoderPicker_PrefersHardware_NeverLibx264()
{
var listing =
" V..... libx264 libx264 H.264 / AVC (codec h264)\n" +
" V..... libopenh264 OpenH264 H.264 (codec h264)\n" +
" V..... h264_nvenc NVIDIA NVENC H.264 encoder (codec h264)\n";
Assert.Equal("h264_nvenc", FfmpegEncoderPicker.Pick(listing));
var onlyGpl = " V..... libx264 libx264 H.264 / AVC (codec h264)\n";
Assert.Equal("libopenh264", FfmpegEncoderPicker.Pick(onlyGpl));
var software = " V..... libopenh264 OpenH264 H.264 (codec h264)\n";
Assert.Equal("libopenh264", FfmpegEncoderPicker.Pick(software));
Assert.Equal("libopenh264", FfmpegEncoderPicker.Pick("no encoders at all"));
}
/// <summary>A stdin slow enough that the encoder can't keep up — the pump's old
/// blocking-write failure mode made demonstrable.</summary>
private sealed class SlowSink : Stream
{
public long TotalBytes;
public override bool CanRead => false;
public override bool CanSeek => false;
public override bool CanWrite => true;
public override long Length => throw new NotSupportedException();
public override long Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); }
public override void Flush() { }
public override Task FlushAsync(CancellationToken ct) => Task.CompletedTask;
public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException();
public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();
public override void SetLength(long value) => throw new NotSupportedException();
public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException();
public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken ct)
{
await Task.Delay(40, ct).ConfigureAwait(false);
Interlocked.Add(ref TotalBytes, count);
}
}
/// <summary>The ONE integration test for slice 10 (the bounded-queue reshape): a
/// lagging ffmpeg must NOT stall the frame producer. With 300 fast submits against
/// a ~25fps-max sink, every queue slot fills and the NEWEST frames are dropped and
/// counted; SubmitFrameAsync must return instantly (it never touched a pipe wait —
/// the old code blocked in WriteAsync and the pump froze); StopAsync still flushes
/// every accepted frame; and the sink ends with exactly (submitted − dropped) bytes
/// — the drop counter and the recording agree (the burned-in-frame-counter premise).</summary>
[Fact]
public async Task Backpressure_QueueOverflow_DropsFrames_AndNeverBlocks()
{
var sink = new SlowSink();
var encoderProc = new FakeEncoderProcess(stdin: sink);
var factory = new FakeEncoderProcessFactory();
factory.Return(new FakeEncoderProcess());
factory.Return(encoderProc);
using var encoder = new FfmpegEncoder(new StubLocator(), factory.Create);
await encoder.StartAsync(new EncoderOptions
{
RtmpUrl = "rtmp://x/y",
Width = 64,
Height = 48,
Fps = 60,
});
var frame = new VideoFrame(64, 48, BgraFrame(64, 48, 255, 0, 0));
const int submits = 300;
var sw = System.Diagnostics.Stopwatch.StartNew();
for (var i = 0; i < submits; i++)
await encoder.SubmitFrameAsync(frame);
sw.Stop();
Assert.True(sw.Elapsed < TimeSpan.FromSeconds(2),
$"SubmitFrameAsync blocked on a saturated queue: {sw.Elapsed}");
var dropped = encoder.DroppedFrames;
Assert.True(dropped > 0, "the bounded queue must overflow-drop when the encoder lags");
Assert.True(dropped < submits, $"the whole lot must not vanish ({dropped}/{submits})");
encoderProc.SignalExit();
await encoder.StopAsync();
Assert.Equal((submits - dropped) * 64 * 48 * 4L, sink.TotalBytes);
}
}