432adfdaef
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).
429 lines
18 KiB
C#
429 lines
18 KiB
C#
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();
|
||
}
|
||
}
|