Files
LlamaCasty/Services/Encoder/FramePump.cs
T
gramps 432adfdaef perf(render): take-4 slice 2 — IsOpaque memcpy, integer bilinear, pump scratch pool
Take 4: pacing held (sync perfect) but avg render stayed 58.9ms — 2M
managed row-walk iterations + a fresh 8.3MB buffer every tick (LOH churn
into GC stalls inside the render measurement).

- VideoFrame.IsOpaque: producer-contract flag (screen capture + webcam —
  DWM/MF fill alpha 255; media/chat/web/static NOT flagged). Full-cover
  aligned opaque backdrop = ONE Buffer.BlockCopy; black pre-fill skipped
  when it covers.
- General BlitContent: integer 8.8 fixed-point bilinear + blend, row
  invariants hoisted, no per-pixel division/Math.Round. Within ±1 of the
  float reference (pixel tests allow ±2). Research per derivative-work
  rule: libyuv row/scale kernels (chromium.googlesource.com/libyuv/libyuv).
- FramePump scratch pool (max 4, length-keyed, owned-by-reference):
  release strictly AFTER SubmitFrameAsync returns (stdin write copies);
  Contains-guard makes the transition Cut alias safe.
- Removed the dead per-tick fromScene render + fromSceneProvider seam —
  BlendFrame consumes TransitionService.FromFrame captured at Start; the
  pump's render fed nothing. MainViewModel call site updated (signature).

Bugs caught by the pixel probes pre-ship (recorded MyMistakes): first
Bilinear double-shifted both stages (solid-255 sampled to ~1 -> general
path drew nothing); sentinel 0xAB collided with an x+y pixel. FakeEncoder
snapshots submitted frames (mirrors real copy semantics under recycling).

ONE integration test: Pump_Pools_ScratchBuffers_Across_Frames_Without_
Stale_Pixels (alternating backdrops + repeated backing identity). Direct
pin: Composite_OpaqueFullCover_Backdrop_CopiesEveryPixel_Into_Scratch.
Clean build 0 warnings; 59/59 per-class + RealApp boot-smoke. take 5
verdict: expect avg render <= ~10ms, ~300/300 frames. Docs same commit:
ai.md pipeline section, TASKS.md TASK 18, HANDOFF rewritten (Unit B spec
+ settled decisions queued).
2026-09-04 09:35:59 -07:00

429 lines
18 KiB
C#
Raw 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;
using System.Threading;
using System.Threading.Tasks;
using ytLive.Models;
using ytLive.Services.Compositor;
namespace ytLive.Services.Encoder;
/// <summary>
/// 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.
/// </summary>
public sealed class FramePump : IDisposable
{
private readonly Func<Scene?> _sceneProvider;
private readonly Func<SceneElement, VideoFrame?> _frameResolver;
private readonly Func<CompositorOptions> _compositorOptions;
private readonly Func<EncoderOptions?> _encoderOptions;
private readonly Func<IFfmpegEncoder> _encoderFactory;
private readonly Action<string>? _log;
private readonly Func<TimeSpan, CancellationToken, Task> _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<byte[]> _ownedScratch = new(ReferenceEqualityComparer.Instance);
private readonly List<byte[]> _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);
}
private readonly object _gate = new();
private IFfmpegEncoder? _encoder;
private CancellationTokenSource? _cts;
private Task? _pumpTask;
private bool _started;
/// <summary>Forwards the encoder's parsed health — ship step 6 binds this to the bottom bar.</summary>
public event EventHandler<StreamHealth>? HealthUpdated;
/// <summary>Raised when the encoder cannot start or dies mid-stream. The pump stops itself.</summary>
public event EventHandler<string>? Failed;
public FramePump(
Func<Scene?> sceneProvider,
Func<SceneElement, VideoFrame?> frameResolver,
Func<CompositorOptions> compositorOptions,
Func<EncoderOptions?> encoderOptions,
Func<IFfmpegEncoder> encoderFactory,
Action<string>? log = null,
Func<TimeSpan, CancellationToken, Task>? 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; }
/// <summary>Never throws: failures are logged and surfaced via <see cref="Failed"/>,
/// so the VM can fire-and-forget it from a sync command handler.</summary>
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();
_pumpTask = 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();
long renderTicks = 0, submitTicks = 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, " +
$"avg submit {submitTicks / (double)System.Diagnostics.Stopwatch.Frequency * 1000 / statFrames:F1}ms");
renderTicks = submitTicks = 0;
statFrames = 0;
statsNext = DateTime.UtcNow + TimeSpan.FromSeconds(5);
}
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);
if (_transition is { Active: true } transition)
{
frame = transition.BlendFrame(frame);
transition.Tick(lastTick.Elapsed.TotalMilliseconds);
}
// Restarted EVERY frame (transition or not) — the old per-frame reset
// is what stops a transition that begins after idle from inheriting
// a giant ElapsedMs and completing instantly on its first tick.
lastTick.Restart();
renderSw.Stop();
renderTicks += renderSw.ElapsedTicks;
IFfmpegEncoder? encoder;
lock (_gate) encoder = _encoder;
if (encoder == null) break;
submitSw.Restart();
await encoder.SubmitFrameAsync(frame, ct);
submitSw.Stop();
submitTicks += submitSw.ElapsedTicks;
// Submit copied the bytes — everything from this tick is recyclable.
// Release AFTER submit, and the free-list Contains guard makes the
// Cut path (BlendFrame returns toFrame itself, aliasing scratch) safe.
ReleaseScratch(frame.BgraPixels);
ReleaseScratch(scratch);
statFrames++;
ReportStats();
// Advance the deadline; cost already spent is not slept again.
// Blew the frame budget: skip the wait AND the missed ticks —
// rebase the clock rather than bursting a catch-up pile
// (OBS rewinds its tick the same way; a burst would only
// queue stale frames into the encoder).
nextTick += intervalTicks;
var lag = nextTick - System.Diagnostics.Stopwatch.GetTimestamp();
if (lag <= 0)
{
nextTick = System.Diagnostics.Stopwatch.GetTimestamp() + intervalTicks;
lag = 0;
}
await _pacingDelay(
TimeSpan.FromSeconds(lag / (double)System.Diagnostics.Stopwatch.Frequency), ct);
}
}
}
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
{
lock (_gate) IsRunning = false;
}
}
private void OnHealthUpdated(object? sender, StreamHealth health) => HealthUpdated?.Invoke(this, health);
/// <summary>Render a scene, using the baked-crust optimization when a <see cref="SceneGraph"/>
/// 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). <paramref name="scratch"/> 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).</summary>
private VideoFrame RenderScene(
Scene scene,
CompositorOptions options,
VideoFrame? socialBarFrame,
int socialBarTop,
byte[]? scratch = null)
{
if (_sceneGraph == null)
return _compositor.Render(scene, _frameResolver, 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, _frameResolver, options, socialBarFrame, socialBarTop, scratch: scratch);
}
// No static base (first layer is dynamic or empty scene) — full render.
return _compositor.Render(scene, _frameResolver, 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();
}
}