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.
390 lines
15 KiB
C#
390 lines
15 KiB
C#
using System.Buffers;
|
|
using System.Diagnostics;
|
|
using System.Threading.Channels;
|
|
using ytLive.Helpers;
|
|
using ytLive.Models;
|
|
|
|
namespace ytLive.Services.Encoder;
|
|
|
|
/// <summary>
|
|
/// The default <see cref="IFfmpegEncoder"/>: spawns <c>ffmpeg.exe</c> (resolved via
|
|
/// <see cref="IFfmpegLocator"/>), feeds raw BGRA frames into stdin, and parses the
|
|
/// <c>-stats</c> progress lines into <see cref="StreamHealth"/>. Encoder choice is
|
|
/// probed from the binary's <c>-encoders</c> listing (hardware NVENC/QSV/AMF first,
|
|
/// OpenH264 software fallback — never libx264, see the license posture) unless
|
|
/// <see cref="EncoderOptions.VideoEncoder"/> forces one.
|
|
///
|
|
/// Graceful stop = close stdin (EOF) → ffmpeg finalizes the FLV and exits by itself;
|
|
/// a watchdogs kill fires only if it hasn't exited shortly after EOF.
|
|
/// </summary>
|
|
public sealed class FfmpegEncoder : IFfmpegEncoder
|
|
{
|
|
public event EventHandler<StreamHealth>? HealthUpdated;
|
|
public event EventHandler<string>? ProcessFailed;
|
|
|
|
private readonly IFfmpegLocator _locator;
|
|
private readonly Func<IEncoderProcess> _processFactory;
|
|
|
|
private readonly object _gate = new();
|
|
private IEncoderProcess? _process;
|
|
private EncoderOptions? _options;
|
|
private StreamHealth _health = new() { Status = StreamStatus.Offline };
|
|
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)
|
|
{
|
|
_locator = locator;
|
|
_processFactory = processFactory ?? (() => new FfmpegEncoderProcess());
|
|
}
|
|
|
|
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));
|
|
var hasStream = options.StreamEnabled && !string.IsNullOrWhiteSpace(options.RtmpUrl);
|
|
var hasRecord = options.RecordEnabled && !string.IsNullOrWhiteSpace(options.RecordPath);
|
|
if (!hasStream && !hasRecord)
|
|
throw new ArgumentException(
|
|
"At least one output is required — a stream needs an RTMP URL, a recording needs an output path.",
|
|
nameof(options));
|
|
|
|
lock (_gate)
|
|
{
|
|
if (IsRunning) throw new InvalidOperationException("The encoder is already running.");
|
|
_options = options;
|
|
}
|
|
|
|
var ffmpegPath = await _locator.LocateAsync(cancellationToken).ConfigureAwait(false);
|
|
var encoder = options.VideoEncoder ?? await ProbeEncoderAsync(ffmpegPath, cancellationToken).ConfigureAwait(false);
|
|
var args = FfmpegArgs.Build(options, encoder);
|
|
|
|
var startInfo = new ProcessStartInfo
|
|
{
|
|
FileName = ffmpegPath,
|
|
UseShellExecute = false,
|
|
RedirectStandardInput = true,
|
|
RedirectStandardOutput = true,
|
|
RedirectStandardError = true,
|
|
CreateNoWindow = true,
|
|
};
|
|
foreach (var arg in args) startInfo.ArgumentList.Add(arg);
|
|
|
|
IEncoderProcess process;
|
|
try
|
|
{
|
|
process = _processFactory();
|
|
process.Start(startInfo);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
AppLog.Write(ex, "FFmpeg encoder: failed to start subprocess");
|
|
lock (_gate) _options = null;
|
|
throw;
|
|
}
|
|
|
|
lock (_gate)
|
|
{
|
|
_process = process;
|
|
IsRunning = true;
|
|
_health = new StreamHealth { Status = StreamStatus.Streaming };
|
|
}
|
|
|
|
_stopRequested = false;
|
|
_stderrLoop = RunStderrLoopAsync(process);
|
|
StartDrainLoop(process);
|
|
}
|
|
|
|
/// <summary>
|
|
/// 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 Task SubmitFrameAsync(VideoFrame frame, CancellationToken cancellationToken = default)
|
|
{
|
|
if (frame == null) throw new ArgumentNullException(nameof(frame));
|
|
lock (_gate)
|
|
{
|
|
if (!IsRunning) throw new InvalidOperationException("The encoder is not running.");
|
|
}
|
|
|
|
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)
|
|
{
|
|
IEncoderProcess? process;
|
|
Task? loop;
|
|
lock (_gate)
|
|
{
|
|
if (!IsRunning) return;
|
|
_stopRequested = true;
|
|
process = _process;
|
|
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
|
|
{
|
|
_frames.Writer.TryComplete();
|
|
if (_drainTask != null) await _drainTask.ConfigureAwait(false);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
AppLog.Write(ex, "FFmpeg encoder: drain loop faulted during stop");
|
|
}
|
|
|
|
try
|
|
{
|
|
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
|
|
timeout.CancelAfter(TimeSpan.FromSeconds(10));
|
|
await process!.WaitForExitAsync(timeout.Token).ConfigureAwait(false);
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
AppLog.Write("FFmpeg encoder: did not exit after stdin EOF — killing");
|
|
process!.Kill();
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
AppLog.Write(ex, "FFmpeg encoder: waiting for exit failed");
|
|
process!.Kill();
|
|
}
|
|
|
|
try
|
|
{
|
|
if (loop != null) await loop.ConfigureAwait(false);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
AppLog.Write(ex, "FFmpeg encoder: stderr loop faulted during stop");
|
|
}
|
|
|
|
process.Dispose();
|
|
|
|
lock (_gate)
|
|
{
|
|
IsRunning = false;
|
|
_health.Status = StreamStatus.Offline;
|
|
_health.LastError = null;
|
|
_process = null;
|
|
_options = null;
|
|
}
|
|
|
|
AppLog.Write($"FFmpeg encoder stopped (exit {process.ExitCode})");
|
|
}
|
|
|
|
public void Dispose()
|
|
{
|
|
_frames.Writer.TryComplete(); // unblock a stuck drain loop alongside the kill
|
|
lock (_gate)
|
|
{
|
|
if (!IsRunning) return;
|
|
_process?.Kill();
|
|
_process?.Dispose();
|
|
_process = null;
|
|
IsRunning = false;
|
|
}
|
|
}
|
|
|
|
/// <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
|
|
{
|
|
var startInfo = new ProcessStartInfo
|
|
{
|
|
FileName = ffmpegPath,
|
|
UseShellExecute = false,
|
|
RedirectStandardOutput = true,
|
|
RedirectStandardError = true,
|
|
CreateNoWindow = true,
|
|
};
|
|
startInfo.ArgumentList.Add("-hide_banner");
|
|
startInfo.ArgumentList.Add("-encoders");
|
|
|
|
using var probe = _processFactory();
|
|
probe.Start(startInfo);
|
|
using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
|
|
cts.CancelAfter(TimeSpan.FromSeconds(15));
|
|
var output = await probe.StandardOutput.ReadToEndAsync(cts.Token).ConfigureAwait(false);
|
|
return FfmpegEncoderPicker.Pick(output);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
AppLog.Write(ex, "FFmpeg encoder: encoder probe failed — falling back to software");
|
|
return FfmpegEncoderPicker.Preference[^1];
|
|
}
|
|
}
|
|
|
|
private Task RunStderrLoopAsync(IEncoderProcess process)
|
|
{
|
|
return Task.Run(async () =>
|
|
{
|
|
try
|
|
{
|
|
while (true)
|
|
{
|
|
var line = await process.StandardError.ReadLineAsync().ConfigureAwait(false);
|
|
if (line == null) break;
|
|
OnStderrLine(process, line);
|
|
}
|
|
|
|
// Read the exit code defensively: the stop path disposes the process
|
|
// while this loop may still be waking (2026-09-01 "No process is
|
|
// associated with this object" — a cosmetic race that poisoned the log).
|
|
int exitCode;
|
|
try
|
|
{
|
|
exitCode = process.ExitCode;
|
|
}
|
|
catch (Exception)
|
|
{
|
|
exitCode = -1;
|
|
}
|
|
|
|
var stillRunning = false;
|
|
var stopping = false;
|
|
lock (_gate)
|
|
{
|
|
stillRunning = IsRunning;
|
|
stopping = _stopRequested;
|
|
}
|
|
if (stillRunning && !stopping)
|
|
{
|
|
// ANY exit we did not ask for is a failure — including code 0.
|
|
// The old `code != 0` gate let ffmpeg's clean early exit vanish
|
|
// without a word while the pump kept feeding a corpse
|
|
// (2026-09-01 first real recording: died at ~0.9s, zero logs).
|
|
AppLog.Write($"FFmpeg encoder: subprocess exited unexpectedly (code {exitCode})");
|
|
_health.LastError = $"FFmpeg exited with code {exitCode}";
|
|
ProcessFailed?.Invoke(this, $"FFmpeg exited with code {exitCode}");
|
|
}
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
AppLog.Write(ex, "FFmpeg encoder: stderr loop faulted");
|
|
}
|
|
});
|
|
}
|
|
|
|
private void OnStderrLine(IEncoderProcess process, string line)
|
|
{
|
|
var progress = FfmpegProgressParser.TryParse(line);
|
|
if (progress == null)
|
|
{
|
|
// THE EYES: everything ffmpeg says that isn't a progress line — startup
|
|
// config, warnings, the real reason it quit. Was silently discarded
|
|
// (2026-09-01: encoder died at ~0.9s, log silent). Truncated so a
|
|
// \r-choked stats blob can't flood; noise now beats blindness.
|
|
var text = line.Trim('\r', ' ', '\t');
|
|
if (text.Length > 0)
|
|
AppLog.Write($"ffmpeg: {(text.Length > 400 ? text[..400] + "…" : text)}");
|
|
return;
|
|
}
|
|
|
|
var dropped = Math.Max(0, (long)Math.Round(progress.Value.Fps * progress.Value.Duration.TotalSeconds) - progress.Value.Frame);
|
|
|
|
lock (_gate)
|
|
{
|
|
_health.CurrentBitrate = progress.Value.BitrateKbps;
|
|
_health.FPS = progress.Value.Fps;
|
|
_health.DroppedFrames = (int)dropped;
|
|
_health.StreamDuration = progress.Value.Duration;
|
|
}
|
|
|
|
HealthUpdated?.Invoke(this, _health);
|
|
}
|
|
}
|