Files
LlamaCasty/Services/Encoder/FramePump.cs
T
gramps 27bf74389d diag(build): per-build GUID stamp + resolve/blit timing split (take-6 attribution failure)
Take 6 measured render 35-41ms — WORSE than take 5's 25.5 — and the run could
not be attributed to a binary: exe mtime != build contents (incremental builds
serve stale exes; a source edit without rebuild is a silent old binary). Three
takes of a perf saga had been judged against builds nobody could prove.

- ytLive.csproj GenerateBuildStamp target: fresh GUID per compile (writes
  obj/BuildStamp.g.cs -> Helpers/BuildStamp.Id/BuiltLocal). Deliberately defeats
  incremental lies: every 'dotnet build' recompiles the app project.
- Wordmark shows the id as a superscript (TopBar.xaml, x:Static, 9px grey
  BaselineAlignment=Superscript); startup.log records 'Build <id> (compiled
  <time>)' so every take is cross-readable with the visible UI.
- FramePump stats split the tick: 'avg render Xms (resolve Y), avg submit Z' —
  the resolver is timed separately (wrapper resolver on per-tick paths; bake
  keeps the raw one) so take 7 names the hot half of 'render' with data.
- Fixed a latent transition-clock bug found on the way: lastTick now restarts
  every frame (the branch rework had restarted it only during transitions,
  letting a transition begun after idle complete instantly on its first Tick).

ONE integration test family: BuildStampTests (unit: shape) +
BuildStampDisplayTests (RealApp, namescoped FindName on TopBar proves the
wordmark SHOWS the id). 47/47 per-class green, clean build 0 warnings. Docs
same commit. User's top-bar/session spec (re-sent twice) + settled Q&A
decisions folded into HANDOFF Unit B — next work unit after take 7 verdict.
2026-09-04 11:26:46 -07:00

452 lines
19 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();
// 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;
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");
renderTicks = submitTicks = resolveTicks = 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.
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) so a transition's first
// Tick sees per-frame time, not the pump's whole uptime.
lastTick.Restart();
// 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,
Func<SceneElement, VideoFrame?>? 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();
}
}