using System;
using System.Threading;
using System.Threading.Tasks;
using ytLive.Models;
using ytLive.Services.Compositor;
namespace ytLive.Services.Encoder;
///
/// The live frame producer (TASK 4 ship step 5): the bridge between the capture
/// managers + compositor and the encoder. While live it snapshots the active
/// scene each tick, resolves every element to its latest frame, composites it
/// into the tier's output frame, and paces frames into the encoder at the tier's
/// FPS. All collaborators are constructor-injected seams (scene, resolver,
/// options, encoder factory, pacing delay) so the pump stays free of WPF and of
/// the capture managers and is fully hermetic in tests.
///
/// The RTMP URL comes from the options provider: until the live-stream create
/// flow lands (TASK 5) it yields null, so go-live runs the existing visual flow
/// without actually pushing.
///
public sealed class FramePump : IDisposable
{
private readonly Func _sceneProvider;
private readonly Func _frameResolver;
private readonly Func _compositorOptions;
private readonly Func _encoderOptions;
private readonly Func _encoderFactory;
private readonly Action? _log;
private readonly Func _pacingDelay;
private readonly Func<(VideoFrame? Frame, SocialBarPosition Position)>? _socialBar;
private readonly TransitionService? _transition;
private readonly SceneGraph? _sceneGraph;
private readonly SceneCompositor _compositor = new();
// Master-buffer scratch pool (take-4 starvation fix, slice 2 of 2): a fresh
// 8.3MB byte[] every tick is ~500MB/s of LOH churn — GC stalls masquerading
// as render cost. The pump pools ONLY buffers it handed out (reference-equality
// set), so bake-cache / social-bar / static-cache frames are never touched;
// SubmitFrameAsync copies the bytes to the encoder's stdin before returning,
// so recycling after submit is safe (MyMistakes recipe).
private readonly HashSet _ownedScratch = new(ReferenceEqualityComparer.Instance);
private readonly List _freeScratch = new();
private const int MaxScratchPooled = 4;
private byte[] AcquireScratch(int size)
{
for (var i = _freeScratch.Count - 1; i >= 0; i--)
{
if (_freeScratch[i].Length != size) continue;
var buffer = _freeScratch[i];
_freeScratch.RemoveAt(i);
return buffer;
}
var fresh = new byte[size];
_ownedScratch.Add(fresh);
return fresh;
}
private void ReleaseScratch(byte[]? buffer)
{
if (buffer == null || !_ownedScratch.Contains(buffer)) return;
if (_freeScratch.Count >= MaxScratchPooled || _freeScratch.Contains(buffer)) return;
_freeScratch.Add(buffer);
}
// Windows sleep quantum (take-9 finding, 2026-09-04): Task.Delay rounds every
// request up to the system clock tick (~15.6ms default — learn.microsoft.com/en-us/
// dotnet/api/system.threading.tasks.task.delay: "approximately 15 milliseconds on
// Windows systems"), so a pacer requesting 3-15ms actually sleeps 15.6ms. Takes
// 5-9 measured work ~25ms but period ~37ms: one padded wait per frame hid every
// compositor improvement. Established media-app practice (game-loop/OBS canon —
// stackoverflow.com/questions/5441464; and raise the resolution for the session —
// learn.microsoft.com/en-us/windows/win32/api/timeapi/nf-timeapi-timebeginperiod):
// timeBeginPeriod(1) while the pump runs, sleep only the BULK of the remainder,
// and spin the last ~2ms across the deadline.
[System.Runtime.InteropServices.DllImport("winmm.dll")]
private static extern uint timeBeginPeriod(uint uMilliseconds);
[System.Runtime.InteropServices.DllImport("winmm.dll")]
private static extern uint timeEndPeriod(uint uMilliseconds);
private static readonly long SpinTailTicks = System.Diagnostics.Stopwatch.Frequency * 2 / 1000; // 2ms
private readonly object _gate = new();
private IFfmpegEncoder? _encoder;
private CancellationTokenSource? _cts;
private Task? _pumpTask;
private bool _started;
/// Forwards the encoder's parsed health — ship step 6 binds this to the bottom bar.
public event EventHandler? HealthUpdated;
/// Raised when the encoder cannot start or dies mid-stream. The pump stops itself.
public event EventHandler? Failed;
public FramePump(
Func sceneProvider,
Func frameResolver,
Func compositorOptions,
Func encoderOptions,
Func encoderFactory,
Action? log = null,
Func? pacingDelay = null,
Func<(VideoFrame? Frame, SocialBarPosition Position)>? socialBar = null,
TransitionService? transition = null,
SceneGraph? sceneGraph = null)
{
_sceneProvider = sceneProvider ?? throw new ArgumentNullException(nameof(sceneProvider));
_frameResolver = frameResolver ?? throw new ArgumentNullException(nameof(frameResolver));
_compositorOptions = compositorOptions ?? throw new ArgumentNullException(nameof(compositorOptions));
_encoderOptions = encoderOptions ?? throw new ArgumentNullException(nameof(encoderOptions));
_encoderFactory = encoderFactory ?? throw new ArgumentNullException(nameof(encoderFactory));
_log = log;
_pacingDelay = pacingDelay ?? ((delay, ct) => Task.Delay(delay, ct));
_socialBar = socialBar;
_transition = transition;
_sceneGraph = sceneGraph;
}
public bool IsRunning { get; private set; }
/// Never throws: failures are logged and surfaced via ,
/// so the VM can fire-and-forget it from a sync command handler.
public async Task StartAsync(CancellationToken cancellationToken = default)
{
lock (_gate)
{
if (_started) return;
_started = true;
}
IFfmpegEncoder? encoder = null;
try
{
var options = _encoderOptions();
if (options == null)
{
_log?.Invoke("FramePump: no output configured (neither streaming nor recording) — encoder skipped");
lock (_gate) _started = false;
return;
}
encoder = _encoderFactory();
encoder.HealthUpdated += OnHealthUpdated;
encoder.ProcessFailed += OnProcessFailed;
await encoder.StartAsync(options, cancellationToken);
lock (_gate)
{
_encoder = encoder;
}
// IsRunning must be true before the loop starts: the loop reads it on
// its first iteration, and with a completed-task delay it can run
// synchronously on this thread before PumpAsync even returns.
IsRunning = true;
_cts = new CancellationTokenSource();
// The loop runs on the thread pool ON PURPOSE (take-8 finding, 2026-09-04):
// Task.Run installs no SynchronizationContext, so every await continuation
// stays off the UI dispatcher. Before this, the pump inherited the UI
// thread's sync context (StartAsync is fired from a command handler), so
// "render 22ms, wait 10ms" was the producer sitting in the dispatcher
// queue behind the live preview it is meant to be independent of — the
// stats quantum fix made the wait VISIBLE; this removes its cause.
// OBS's video threads are dedicated for exactly this reason.
_pumpTask = Task.Run(() => PumpAsync(options, _cts.Token));
_log?.Invoke($"FramePump started ({options.Width}×{options.Height} @ {options.Fps} fps)");
}
catch (Exception ex)
{
_log?.Invoke($"FramePump: start failed: {ex.Message}");
if (encoder != null)
{
encoder.HealthUpdated -= OnHealthUpdated;
encoder.ProcessFailed -= OnProcessFailed;
try
{
encoder.Dispose();
}
catch (Exception disposeEx)
{
_log?.Invoke($"FramePump: disposing failed encoder: {disposeEx.Message}");
}
}
lock (_gate)
{
_started = false;
IsRunning = false;
}
Failed?.Invoke(this, ex.Message);
}
}
public async Task StopAsync(CancellationToken cancellationToken = default)
{
IFfmpegEncoder? encoder;
Task? pump;
lock (_gate)
{
if (!_started && _encoder == null) return;
_started = false;
IsRunning = false;
encoder = _encoder;
pump = _pumpTask;
_cts?.Cancel();
}
// Stop the encoder BEFORE awaiting the pump: closing its stdin unblocks a
// write stuck on pipe backpressure, otherwise the pump could await forever.
if (encoder != null)
{
try
{
await encoder.StopAsync(cancellationToken);
}
catch (Exception ex)
{
_log?.Invoke($"FramePump: encoder stop failed: {ex.Message}");
}
}
if (pump != null)
{
try
{
await pump;
}
catch (Exception ex)
{
_log?.Invoke($"FramePump: pump loop faulted during stop: {ex.Message}");
}
}
if (encoder != null)
{
encoder.HealthUpdated -= OnHealthUpdated;
encoder.ProcessFailed -= OnProcessFailed;
try
{
encoder.Dispose();
}
catch (Exception ex)
{
_log?.Invoke($"FramePump: encoder dispose failed: {ex.Message}");
}
}
lock (_gate)
{
_encoder = null;
_cts = null;
_pumpTask = null;
}
_log?.Invoke("FramePump stopped");
}
public void Dispose()
{
try
{
StopAsync().GetAwaiter().GetResult();
}
catch (Exception ex)
{
_log?.Invoke($"FramePump: dispose failed: {ex.Message}");
}
_freeScratch.Clear();
_ownedScratch.Clear();
}
private async Task PumpAsync(EncoderOptions options, CancellationToken ct)
{
var interval = TimeSpan.FromSeconds(1d / Math.Max(1, options.Fps));
var lastTick = System.Diagnostics.Stopwatch.StartNew();
// Deadline pacing (2026-09-03, take-3 fix): the frame interval is a DEADLINE,
// not an afterthought sleep — the OBS libobs video-io.c pattern (researched
// before coding; see https://github.com/obsproject/obs-studio/blob/master/
// libobs/media-io/video-io.c). The old loop slept the FULL interval after
// each render, so period = render + submit + interval: at take-3's 258ms
// render that was 3.6fps stamped into a 60fps container — rawvideo stamps by
// arrival, so 30 wall-seconds muxed as a 2.1s time-lapse, no error anywhere.
var intervalTicks = Math.Max(1, (long)Math.Round(interval.TotalSeconds * System.Diagnostics.Stopwatch.Frequency));
var nextTick = System.Diagnostics.Stopwatch.GetTimestamp();
// Stage timing (2026-09-01, take two): rawvideo carries no per-frame
// timestamps — ffmpeg stamps frames by ARRIVAL at the declared fps. A producer
// slower than the declared rate yields a time-lapsed, short file (observed:
// 39 frames in 27 wall-seconds ≈ 30x at 60fps) with no error anywhere.
// Log the render/submit split every 5s so the next take names the stage.
var renderSw = new System.Diagnostics.Stopwatch();
var submitSw = new System.Diagnostics.Stopwatch();
var waitSw = new System.Diagnostics.Stopwatch();
// Resolve-vs-composite split (2026-09-04, take-6 ambiguity): "render" was a
// 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;
int statFrames = 0;
var statsNext = DateTime.UtcNow + TimeSpan.FromSeconds(5);
void ReportStats()
{
if (DateTime.UtcNow < statsNext) return;
var target = 5d / interval.TotalSeconds; // frames expected per window
_log?.Invoke(statFrames == 0
? "FramePump stats: NO frames produced in 5s (loop stalled?)"
: $"FramePump stats: {statFrames}/{target:F0} frames per 5s, " +
$"avg render {renderTicks / (double)System.Diagnostics.Stopwatch.Frequency * 1000 / statFrames:F1}ms " +
$"(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");
worstRender = 0;
renderTicks = submitTicks = resolveTicks = waitTicks = 0;
statFrames = 0;
statsNext = DateTime.UtcNow + TimeSpan.FromSeconds(5);
}
// One wrapper shared by every render of the run — resolve time accumulates
// inside the render measurement, and the stats line reports the split.
timeBeginPeriod(1); // pairs with timeEndPeriod in the finally — see field note
var previousGcMode = System.Runtime.GCSettings.LatencyMode;
System.Runtime.GCSettings.LatencyMode = System.Runtime.GCLatencyMode.SustainedLowLatency;
VideoFrame? TimedResolver(SceneElement element)
{
resolveSw.Restart();
var frame = _frameResolver(element);
resolveSw.Stop();
resolveTicks += resolveSw.ElapsedTicks;
return frame;
}
try
{
while (!ct.IsCancellationRequested)
{
var scene = _sceneProvider();
if (scene != null)
{
var compositorOptions = _compositorOptions();
VideoFrame? socialBarFrame = null;
var socialBarTop = 0;
if (_socialBar != null)
{
var (barFrame, position) = _socialBar();
socialBarFrame = barFrame;
if (barFrame != null)
socialBarTop = position == SocialBarPosition.Top
? 0
: compositorOptions.SourceRectHeight - barFrame.Height;
}
VideoFrame frame;
var scratchSize = compositorOptions.SourceRectWidth
* compositorOptions.SourceRectHeight * 4;
// The transition "from" frame lives in TransitionService.FromFrame
// (captured at Start by the VM). The old per-tick fromScene render
// here was dead weight — a full extra scene composite every
// transition tick that BlendFrame never read; removed with the
// pooling change because its buffer's only consumer was its own release.
renderSw.Restart();
var scratch = AcquireScratch(scratchSize);
frame = RenderScene(scene, compositorOptions, socialBarFrame, socialBarTop, scratch, TimedResolver);
if (_transition is { Active: true } transition)
{
frame = transition.BlendFrame(frame);
transition.Tick(lastTick.Elapsed.TotalMilliseconds);
}
// Restarted EVERY frame (transition or not): a transition begun
// after idle must not inherit a giant ElapsedMs and complete
// instantly on its first Tick.
lastTick.Restart();
renderSw.Stop();
renderTicks += renderSw.ElapsedTicks;
if (renderSw.ElapsedTicks > worstRender) worstRender = renderSw.ElapsedTicks;
IFfmpegEncoder? encoder;
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)
{
await encoder.SubmitFrameAsync(frame, ct);
statFrames++;
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.
// The free-list Contains guard keeps the Cut path (BlendFrame
// returns toFrame itself, aliasing scratch) safe.
ReleaseScratch(frame.BgraPixels);
ReleaseScratch(scratch);
// Sleep the BULK of the remainder, SPIN the 2ms tail — never hand
// a sub-tick remainder to the sleep quantum (the timeBeginPeriod
// note). If the renderer already ate the budget there is nothing
// left to sleep and the loop renders the next frame straight away.
waitSw.Restart();
var ahead = nextTick - System.Diagnostics.Stopwatch.GetTimestamp();
if (ahead > SpinTailTicks)
await _pacingDelay(TimeSpan.FromSeconds(
(ahead - SpinTailTicks) / (double)System.Diagnostics.Stopwatch.Frequency), ct);
while (!ct.IsCancellationRequested
&& System.Diagnostics.Stopwatch.GetTimestamp() < nextTick)
Thread.SpinWait(400);
waitSw.Stop();
waitTicks += waitSw.ElapsedTicks;
ReportStats();
}
}
}
catch (OperationCanceledException)
{
// normal stop
}
catch (Exception ex)
{
// A failure while the pump is supposed to run (encoder died under us,
// scene provider faulted, ...) stops the pump and surfaces once.
if (ct.IsCancellationRequested)
{
_log?.Invoke($"FramePump: pump exited during stop: {ex.Message}");
}
else
{
_log?.Invoke($"FramePump: pump loop faulted: {ex.Message}");
Failed?.Invoke(this, ex.Message);
}
}
finally
{
System.Runtime.GCSettings.LatencyMode = previousGcMode;
timeEndPeriod(1);
lock (_gate) IsRunning = false;
}
}
private void OnHealthUpdated(object? sender, StreamHealth health) => HealthUpdated?.Invoke(this, health);
/// Render a scene, using the baked-crust optimization when a
/// is wired in: bake/cache the static layers below the split point, then composite the
/// dynamic/above-split layers per frame. Without a SceneGraph, falls back to a full render
/// (identical output — see SceneCompositorTests). is a pooled
/// master buffer when supplied; the fully-static path returns the bake-cache frame itself
/// (never pooled — the release side checks owned-by-reference).
private VideoFrame RenderScene(
Scene scene,
CompositorOptions options,
VideoFrame? socialBarFrame,
int socialBarTop,
byte[]? scratch = null,
Func? resolver = null)
{
// per-tick composites use the (timed) resolver; the rare bake uses the raw one
// so bake cost lands in "render" but not "resolve".
resolver ??= _frameResolver;
if (_sceneGraph == null)
return _compositor.Render(scene, resolver, null, options, socialBarFrame, socialBarTop, scratch: scratch);
var split = _sceneGraph.GetSplitPoint(scene);
if (split == scene.Elements.Count)
{
// Fully static scene: bake once, reuse.
var baked = _sceneGraph.GetBakedBase(scene, _frameResolver, _compositorOptions);
if (baked != null)
return StretchMath.BilinearScale(baked, options.OutputWidth, options.OutputHeight);
}
var baseFrame = _sceneGraph.GetBakedBase(scene, _frameResolver, _compositorOptions);
if (baseFrame != null)
{
return SceneCompositor.CompositeLayers(
baseFrame, scene, split, resolver, options, socialBarFrame, socialBarTop, scratch: scratch);
}
// No static base (first layer is dynamic or empty scene) — full render.
return _compositor.Render(scene, resolver, null, options, socialBarFrame, socialBarTop, scratch: scratch);
}
private void OnProcessFailed(object? sender, string message)
{
_log?.Invoke($"FramePump: encoder process failed: {message}");
Failed?.Invoke(this, message);
_ = StopAsync();
}
}