using System.Buffers; using System.Diagnostics; using System.Threading.Channels; using ytLive.Helpers; using ytLive.Models; namespace ytLive.Services.Encoder; /// /// The default : spawns ffmpeg.exe (resolved via /// ), feeds raw BGRA frames into stdin, and parses the /// -stats progress lines into . Encoder choice is /// probed from the binary's -encoders listing (hardware NVENC/QSV/AMF first, /// OpenH264 software fallback — never libx264, see the license posture) unless /// 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. /// public sealed class FfmpegEncoder : IFfmpegEncoder { public event EventHandler? HealthUpdated; public event EventHandler? ProcessFailed; private readonly IFfmpegLocator _locator; private readonly Func _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 _frames = Channel.CreateBounded( 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? processFactory = null) { _locator = locator; _processFactory = processFactory ?? (() => new FfmpegEncoderProcess()); } public bool IsRunning { get; private set; } /// 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 /// (ffmpeg's own progress-derived estimate). 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); } /// /// 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). /// 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.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.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; } } /// /// 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. /// 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.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 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); } }