Files
LlamaCasty/Services/Encoder/FramePump.cs
T
gramps a11b15e444 feat(alerts): TASK 47 — alert box plays a video (built-in/custom clip) + read-time fade + message ticker
TASK 43's alert box grows a real video celebration. Per-alert IAlertClipDecoder
(ffmpeg bgra + f32le pipes, real-time paced, disposed at drain) plays the shipped
Assets/alert-default.mp4 (stamped into the Asset table at startup) unless the
creator picks their own file — path reference only, never stored in the DB; the
six AlertRenderer animations stay the fallback. ~0.3s fade rides the alpha
envelope on straight-source copies (EOF freeze-frames then fades out); audio
forwards to a new AudioMixer alert ring (8s, 48k stereo) drained at unity — no
duck, creator ruling — scaled by volume × fade. An auto-composed marquee ticker
('Funder — Super Chat · $10.00', 140px/s) scrolls top-of-frame via a
FramePump._alertTicker seam through Render/CompositeLayers, mixed into the cache
signature (dynamic overlay, never baked). New Stream Alerts section in LeftPanel.

Derivative-work references (how OBS/Streamlabs alert boxes do per-alert video):
- https://support.streamlabs.com/hc/en-us/articles/217741147-Setting-Up-Your-Streamlabs-Alerts (custom image/video per alert type + variations)
- https://obsproject.com/kb/stream-tutorial-2-alerts (alert overlay as an on-screen zone)
- https://streamlabs.com/content-hub/widgets/alert-box (per-event alert playback)

Good Dog: AlertLayerVideoTests drives a fake IAlertClipDecoder through the whole
lifecycle in one pass (custom path wins, decoder spawns/disposes, fade envelope
0→127→255, audio volume×fade, ticker scrolls, EOF fade-drain to idle). It caught
the clip branch of Advance not clearing _current before AdvanceToNext — the layer
stayed IsPlaying after drain (MyMistakes post-mortem).

Full vstest 319/319; clean build 0 warnings; scope check green.
2026-09-26 11:15:17 -07:00

839 lines
41 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 Func<VideoFrame?>? _alertTicker;
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);
}
// 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;
private long _outputIndex;
/// <summary>Captured by <see cref="ProbeRender"/> when a render exceeds ~20ms —
/// the 2026-09-12 half-speed-render probe. Appended (once) to the next stats or
/// stall line, then cleared.</summary>
private string _renderDetail = "";
/// <summary>C4 blit-on-change cache (2026-09-15): the FULL-render path (split=0
/// when the live-capture backdrop sits at element 0, or no SceneGraph) re-composites
/// the ENTIRE frame every tick even when no input changed — the ty-1841 take logged
/// a render stall on EVERY iteration (totalMs 21-44, one composite per ~30ms while
/// the content moved ~20 updates/s). The SceneGraph split can never help here: the
/// backdrop is element 0 and must stay dynamic (a baked capture goes stale), so the
/// cache lives at the pump: the last full-render output plus the input identity that
/// produced it. Unchanged identity → ONE BlockCopy (~3ms) reuses the composite;
/// changed identity → re-composite. The cache buffer is a SEPARATE long-lived array,
/// never the scratch pool — the caller burns the frame counter and recycles the
/// scratch AFTER RenderScene returns, so the cache is written from the rendered
/// scratch BEFORE returning (pre-burn, pre-recycle). The identity mirror is
/// <see cref="BuildFullRenderSignature"/>: it resolves the same frames the compositor
/// will (same resolver seam), so a changed frame (new capture, new webcam, new chat
/// raster) forces a re-render and an unchanged one never serves stale bytes. Internal
/// counters feed the Good Dog test + the 5s cache detail in the stats line.</summary>
private byte[]? _fullCachePixels;
private bool _fullCacheValid;
private ulong _lastFullSignature;
internal long CacheHits;
internal long CacheRenders;
/// <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,
Func<VideoFrame?>? alertTicker = 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;
_alertTicker = alertTicker;
_transition = transition;
_sceneGraph = sceneGraph;
}
public bool IsRunning { get; private set; }
/// <summary>The burned counter of the LAST frame submitted to the encoder. Stable
/// across a tick's multiple resolver passes (it advances only at burn/submit, which
/// happens AFTER render) — tests key per-tick content on it.</summary>
internal long OutputIndex => _outputIndex;
/// <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();
// 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, worstSubmit = 0;
int statFrames = 0;
int stalls = 0;
long lastCacheHits = 0, lastCacheRenders = 0;
var statsNext = DateTime.UtcNow + TimeSpan.FromSeconds(5);
void ReportStats()
{
if (DateTime.UtcNow < statsNext) return;
var target = 5d / interval.TotalSeconds; // frames expected per window
IFfmpegEncoder? encoder;
lock (_gate) encoder = _encoder;
var dropped = encoder == null ? 0 : encoder.DroppedFrames;
var detail = _renderDetail;
_renderDetail = "";
var cacheDetail = "";
if (statFrames > 0)
{
// C4 window detail: rendered composites (R) vs cache reuses (H) since the
// last report — hard numeric proof the blit-on-change cache engaged.
var hits = CacheHits - lastCacheHits;
var renders = CacheRenders - lastCacheRenders;
lastCacheHits = CacheHits;
lastCacheRenders = CacheRenders;
if (hits > 0 || renders > 0) cacheDetail = $", cache {renders}R/{hits}H";
}
_log?.Invoke(statFrames == 0
? "FramePump stats: NO frames produced in 5s (loop stalled?)" + (detail.Length > 0 ? " | " + detail : "")
: $"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, "
+ $"worst submit {worstSubmit / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F1}ms, "
+ $"dropped {dropped}, stalls {stalls}" + cacheDetail + (detail.Length > 0 ? " | " + detail : ""));
worstRender = 0;
worstSubmit = 0;
renderTicks = submitTicks = resolveTicks = waitTicks = 0;
statFrames = 0;
stalls = 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 iterStart = System.Diagnostics.Stopwatch.GetTimestamp();
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;
}
// The alert ticker is a per-frame overlay (it scrolls while an alert
// plays), NEVER baked into the static base. Lawn-mower seam: read it
// each tick and let the compositor pin it to the very top.
var tickerFrame = _alertTicker?.Invoke();
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, tickerFrame, 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;
// Deadline pacing, OBS duplicate-on-lag (slice 15): emit ONE frame per
// deadline slot — a fresh composite when the render kept up, a REPEAT
// of this iteration's latest composite for every slot the render
// overran. The OBS model never leaves a wall-time hole: libobs
// media-io/video-io.c runs on its own clock and duplicates the latest
// frame when the video thread lags, logging "lagged frames due to
// rendering lag/stalls" — never skipped time (docs.obsproject.com/
// backend-design: "If the video frame queue is full, it will duplicate
// the last frame"). Slice 10 chose the opposite for freshness: an
// overrun SKIPPED the missed slots, so a 60fps-authoring pump emitting
// one frame per 35ms overrun authored ~1.7x playback (ty-1726: 163
// video frames for 2.93s of audio) and ended the video 0.22s before the
// audio tail. Counting duplicates instead of skips keeps recording
// duration == wall duration (no acceleration, no audio tail cut) at the
// cost of a short judder during a stall — the accepted trade. The burst
// is microseconds (the enqueue never blocks, and the loop is clamped to
// deadlineNow) so it can't smear the way the take-9 blocking burst did.
var deadlineNow = System.Diagnostics.Stopwatch.GetTimestamp();
while (!ct.IsCancellationRequested && deadlineNow >= nextTick)
{
// Burned-in frame counter (slice 10 judge): a clock-independent
// pacing witness burned literally into the composite. Decode the
// recording and read the bottom-right strip: the number must
// advance +1 per authored frame and jump only by counted drops
// (queue overflow). Burned here — on the buffer finally submitted,
// which the enqueue's snapshot copies before the next iteration
// rewrites it — so every file frame carries its own index and a
// fast-render iteration that submitted nothing leaves no phantom
// gap (the slice-10 unconditional bump before the gate could).
_outputIndex++;
BurnFrameIndex(frame.BgraPixels, frame.Width, frame.Height, _outputIndex);
submitSw.Restart();
await encoder.SubmitFrameAsync(frame, ct);
submitSw.Stop();
statFrames++;
submitTicks += submitSw.ElapsedTicks;
if (submitSw.ElapsedTicks > worstSubmit) worstSubmit = submitSw.ElapsedTicks;
nextTick += intervalTicks;
}
// SubmitFrameAsync copied the bytes into the encoder's queue — the
// tick's buffers are recyclable once the enqueue snapshot 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;
// Stall logger (slice 10): an iteration spanning more than two full
// intervals is the old bug's fingerprint — name the stage instead of
// guessing. With the queue, submit should be ~1ms, so a stall here
// means RENDER or RESOLVE (the stage terms of the last take).
var iterWall = System.Diagnostics.Stopwatch.GetTimestamp() - iterStart;
if (iterWall > 2 * intervalTicks)
{
stalls++;
var detail = _renderDetail;
_renderDetail = "";
_log?.Invoke($"FramePump stall: iteration {iterWall / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F0}ms " +
$"(> 2× the {interval.TotalMilliseconds:F0}ms interval): worst render " +
$"{worstRender / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F0}ms, worst submit " +
$"{worstSubmit / (double)System.Diagnostics.Stopwatch.Frequency * 1000:F0}ms, dropped {encoder.DroppedFrames}" +
(detail.Length > 0 ? " | " + detail : ""));
}
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);
/// <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,
VideoFrame? tickerFrame,
byte[]? scratch = null,
Func<SceneElement, VideoFrame?>? resolver = null)
{
resolver ??= _frameResolver;
if (_sceneGraph == null)
{
var sw = System.Diagnostics.Stopwatch.StartNew();
var frame = RenderFull(scene, options, socialBarFrame, socialBarTop, tickerFrame, scratch, resolver);
ProbeRender("full-no-graph", scene, -1, sw.ElapsedMilliseconds, 0);
return frame;
}
var split = _sceneGraph.GetSplitPoint(scene);
if (split == scene.Elements.Count)
{
// Fully static scene: bake once, reuse.
var sw = System.Diagnostics.Stopwatch.StartNew();
var baked = _sceneGraph.GetBakedBase(scene, _frameResolver, _compositorOptions);
if (baked != null)
{
var stretched = StretchMath.BilinearScale(baked, options.OutputWidth, options.OutputHeight);
ProbeRender("fully-static", scene, split, sw.ElapsedMilliseconds, (int)sw.ElapsedMilliseconds);
return stretched;
}
}
var swBase = System.Diagnostics.Stopwatch.StartNew();
var baseFrame = _sceneGraph.GetBakedBase(scene, _frameResolver, _compositorOptions);
if (baseFrame != null)
{
var frame = SceneCompositor.CompositeLayers(
baseFrame, scene, split, resolver, options, socialBarFrame, socialBarTop, tickerFrame: tickerFrame, scratch: scratch);
ProbeRender("base+layers", scene, split, swBase.ElapsedMilliseconds, 0);
return frame;
}
// No static base (first layer is dynamic or empty scene) — full render.
var swFull = System.Diagnostics.Stopwatch.StartNew();
var frame2 = RenderFull(scene, options, socialBarFrame, socialBarTop, tickerFrame, scratch, resolver);
ProbeRender("full-render", scene, split, swFull.ElapsedMilliseconds, 0);
return frame2;
}
/// <summary>The full-render path (no SceneGraph, or no static base) wrapped in the
/// C4 blit-on-change cache: when every resolved input plus every element's
/// layout/visual bits match the last composite, reuse it with one BlockCopy instead
/// of re-compositing. Mirror the compositor's resolution <i>before</i> deciding so a
/// hit costs ~3ms of pure copy and a miss costs exactly the old full render.</summary>
private VideoFrame RenderFull(
Scene scene, CompositorOptions options, VideoFrame? socialBarFrame, int socialBarTop,
VideoFrame? tickerFrame,
byte[]? scratch, Func<SceneElement, VideoFrame?> resolver)
{
// 1:1 guard: the cache stores the OUTPUT bytes and scratch is source-sized —
// at 1:1 they are the same buffer size. Off-size tiers (a vertical 1080×1920
// tier, an H264-1080p tier) fall back to a per-tick full render; the deployed
// config is master-sized and caching a per-tick fresh scale buffer is the known
// vertical-tier follow-up (see ai.md), not silently baked here.
var useCache = scratch != null && options.SourceRectWidth == options.OutputWidth
&& options.SourceRectHeight == options.OutputHeight;
var signature = useCache
? BuildFullRenderSignature(scene, options, socialBarFrame, socialBarTop, tickerFrame, resolver)
: 0UL;
if (useCache && _fullCacheValid && signature == _lastFullSignature
&& _fullCachePixels!.Length == scratch!.Length)
{
Buffer.BlockCopy(_fullCachePixels, 0, scratch, 0, scratch.Length);
CacheHits++;
return new VideoFrame(options.SourceRectWidth, options.SourceRectHeight, scratch);
}
var rendered = _compositor.Render(scene, resolver, null, options, socialBarFrame, socialBarTop, scratch: scratch, tickerFrame: tickerFrame);
CacheRenders++;
if (useCache)
{
var pixels = rendered.BgraPixels;
if (_fullCachePixels == null || _fullCachePixels.Length != pixels.Length)
_fullCachePixels = new byte[pixels.Length];
// PRE-burn / PRE-recycle: the caller burns the counter and recycles this
// scratch after RenderScene returns — the cache keeps a clean copy, and the
// burn on the next hit writes to a FRESH scratch copy, never the cache.
Buffer.BlockCopy(pixels, 0, _fullCachePixels, 0, pixels.Length);
_lastFullSignature = signature;
_fullCacheValid = true;
}
else
{
_fullCacheValid = false; // size changed — never reuse a wrong-size buffer
}
return rendered;
}
/// <summary>Hash the complete set of inputs the compositor's FULL render consumes:
/// the crop/output dimensions, the social bar, and — in element order, mirroring
/// <c>SceneCompositor.Render</c>'s iteration — every element's layout/visual fields
/// plus the frame each one resolves (through the SAME resolver seam the render will
/// use). A changed frame (new capture, new webcam Epoch, new chat raster) changes the
/// hash → re-composite; an unchanged one reuses the cache. Buffer-array identity +
/// Epoch + CropBounds is the compositor's own paste-key shape — producers hand out
/// fresh immutable arrays, or recycled-ring arrays whose monotonic Epoch
/// distinguishes generations (the capture take-11 fix), so address-identity alone is
/// never trusted to key recycled content.</summary>
private static ulong BuildFullRenderSignature(
Scene scene, CompositorOptions options, VideoFrame? socialBarFrame, int socialBarTop,
VideoFrame? tickerFrame,
Func<SceneElement, VideoFrame?> resolver)
{
var h = 14695981039346656037UL; // FNV-1a offset basis
h = Mix(h, (ulong)options.SourceRectX);
h = Mix(h, (ulong)options.SourceRectY);
h = Mix(h, (ulong)options.SourceRectWidth);
h = Mix(h, (ulong)options.SourceRectHeight);
h = Mix(h, (ulong)options.OutputWidth);
h = Mix(h, (ulong)options.OutputHeight);
h = Mix(h, (ulong)socialBarTop);
if (socialBarFrame != null) h = MixFrame(h, socialBarFrame);
else h = Mix(h, 0xFFFFFFFFFFFFFFFFUL);
if (tickerFrame != null) h = MixFrame(h, tickerFrame);
else h = Mix(h, 0xFFFFFFFFFFFFFFFFUL);
var backdropResolved = false;
foreach (var element in scene.Elements)
{
h = Mix(h, (ulong)System.Runtime.CompilerServices.RuntimeHelpers.GetHashCode(element));
h = Mix(h, element.IsVisible ? 1UL : 0UL);
h = MixDouble(h, element.X);
h = MixDouble(h, element.Y);
h = MixDouble(h, element.Width);
h = MixDouble(h, element.Height);
h = MixDouble(h, element.Opacity);
h = Mix(h, (ulong)element.ClipShape);
h = Mix(h, element.IsMirrored ? 1UL : 0UL);
h = Mix(h, (ulong)element.BorderWidth);
h = MixDouble(h, element.BorderOpacity);
if (element.TryGetBorderColor(out var bcR, out var bcG, out var bcB))
h = Mix(h, (ulong)((bcR << 16) | (bcG << 8) | bcB));
else
h = Mix(h, 0UL);
// Mirror the compositor's per-element frame resolution EXACTLY, so the hash
// reacts to the same frames a render uses and never to a frame a render
// ignores (a false-positive change only wastes one render; a false NEGATIVE
// would serve stale bytes — that is what this mirror prevents).
VideoFrame? frame = null;
if (element is Source { IsBackground: true })
{
h = Mix(h, 1UL); // the live backdrop — first visible one is resolved, once
if (element.IsVisible && !backdropResolved)
{
frame = resolver(element);
backdropResolved = true;
}
}
else if (element is Source { Type: SourceType.TextOverlay })
{
h = Mix(h, 2UL); // skipped kind — never resolved by the compositor
}
else if (element is Source { Type: SourceType.Background })
{
h = Mix(h, 3UL); // background image — resolved even when invisible (Render does)
frame = resolver(element);
}
else
{
h = Mix(h, 4UL); // regular layer — resolved only when visible
if (element.IsVisible) frame = resolver(element);
}
if (frame != null) h = MixFrame(h, frame);
else h = Mix(h, 0UL);
}
return h;
}
private static ulong Mix(ulong h, ulong v) => (h ^ v) * 0x9E3779B97F4A7C15UL;
private static ulong MixDouble(ulong h, double v)
=> Mix(h, unchecked((ulong)BitConverter.DoubleToInt64Bits(v)));
private static ulong MixFrame(ulong h, VideoFrame frame)
{
h = Mix(h, (ulong)System.Runtime.CompilerServices.RuntimeHelpers.GetHashCode(frame.BgraPixels));
h = Mix(h, (ulong)frame.Epoch);
h = Mix(h, (ulong)frame.Width);
h = Mix(h, (ulong)frame.Height);
if (frame.CropBounds is { } cb)
{
h = Mix(h, (ulong)cb.X);
h = Mix(h, (ulong)cb.Y);
h = Mix(h, (ulong)cb.W);
h = Mix(h, (ulong)cb.H);
}
else
{
h = Mix(h, 0xFFFFFFFFFFFFFFFFUL);
}
return h;
}
/// <summary>2026-09-12 half-speed-render probe: the pump only produces ~27fps
/// (render ~35ms vs the 16.7ms deadline), and rawvideo muxes at the DECLARED fps,
/// so every take muxes at ~half its wall length (the truncation complaint). This
/// names WHICH compositor path ate the slow frame so the fix targets the real
/// stage. Records only when a frame takes ≥20ms (or carries the previous detail).</summary>
private void ProbeRender(string path, Scene scene, int split, long totalMs, int bakeMs)
{
if (totalMs < 20 && _renderDetail.Length == 0) return;
var dynamics = 0;
foreach (var e in scene.Elements)
if (e.Kind == ElementKind.Dynamic) dynamics++;
_renderDetail = $"render={path} split={split} elements={scene.Elements.Count} dynamic={dynamics}" +
(bakeMs > 0 ? $" bakeMs={bakeMs}" : "") + $" totalMs={totalMs}";
}
// Dot-matrix digits (5×7, one row per raster line, '1' = lit) burned into the
// bottom-right of every composite. Basic OCR-safe shapes, sized so the strip is
// a 36×8 white box in the corner of a 1920×1080 frame — readable with a zoomed
// player, invisible at normal size.
private static readonly string[][] DigitGlyphs =
{
new[] { "01110","10001","10001","10001","10001","10001","01110" }, // 0
new[] { "00100","01100","00100","00100","00100","00100","01110" }, // 1
new[] { "01110","10001","00001","00010","00100","01000","11111" }, // 2
new[] { "11111","00001","00010","00110","00001","10001","01110" }, // 3
new[] { "00010","00110","01010","10010","11111","00010","00010" }, // 4
new[] { "11111","10000","11110","00001","00001","10001","01110" }, // 5
new[] { "01110","10001","10000","11110","10001","10001","01110" }, // 6
new[] { "11111","00001","00010","00100","01000","01000","01000" }, // 7
new[] { "01110","10001","10001","01110","10001","10001","01110" }, // 8
new[] { "01110","10001","10001","01111","00001","10001","01110" }, // 9
};
private static void BurnFrameIndex(byte[] bgra, int width, int height, long index)
{
const int digitW = 5, digitH = 7, gap = 1, margin = 2, digitCount = 6;
var stripW = digitCount * (digitW + gap) - gap;
var left = width - margin - stripW;
var top = height - margin - digitH;
// Solid white box under the digits — the underlying scene can be anything.
for (var y = top; y < top + digitH; y++)
for (var x = left; x < left + stripW; x++)
{
var i = (y * width + x) * 4;
bgra[i] = 255;
bgra[i + 1] = 255;
bgra[i + 2] = 255;
}
var text = index.ToString("D6");
for (var d = 0; d < digitCount; d++)
{
var glyph = DigitGlyphs[text[d] - '0'];
for (var row = 0; row < digitH; row++)
for (var col = 0; col < digitW; col++)
{
if (glyph[row][col] != '1') continue;
var i = ((top + row) * width + (left + d * (digitW + gap) + col)) * 4;
bgra[i] = 0;
bgra[i + 1] = 0;
bgra[i + 2] = 0;
}
}
}
private void OnProcessFailed(object? sender, string message)
{
_log?.Invoke($"FramePump: encoder process failed: {message}");
Failed?.Invoke(this, message);
_ = StopAsync();
}
}