c45cbc93b7
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.
421 lines
16 KiB
C#
421 lines
16 KiB
C#
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);
|
||
}
|
||
}
|