716a77f61a
Two defects made the producer 17x slow (37s record -> 2.1s/127-frame file, rawvideo stamps by arrival): FramePump slept the FULL interval after each render (period = render+submit+interval) and SceneCompositor did per-pixel float sampling + Math.Round blends over all 2.07M master pixels, scanning the whole destination per overlay (258ms avg render vs 1.5ms submit). Both solutions are established, not invented — researched before coding per the derivative-work rule: - deadline pacing: OBS libobs/media-io/video-io.c video_thread (nextTick += intervalTicks, sleep only the remainder, rebase on overrun, never burst) - row blits: libyuv pattern (BSD-3, chromium.googlesource.com/libyuv/libyuv) — 1:1 aligned identity fast path, per-pixel alpha branch, integer fixed-point blend, overlay clipped to the intersection rect, skip the dead black pre-fill when the backdrop covers ONE integration test: Pump_Paces_To_The_Deadline_Compensating_Render_Cost (lands after a fake-seam lesson: pacing fakes must await, not complete synchronously, or the pump loop runs inline on StartAsync and hangs vstest). Clean build 0 warnings; FramePumpTests 10/10, SceneCompositor/SceneGraph/ SocialBar/StretchMath 20/20. Docs same commit: ai.md pipeline section, TASKS.md TASK 18 (webcam-in-output + rename modal verified from take 3), MyMistakes recipe, HANDOFF rewritten. Take 4 pending on the user's machine.
391 lines
15 KiB
C#
391 lines
15 KiB
C#
using System.Net;
|
|
using System.Text;
|
|
using Xunit;
|
|
using ytLive.Models;
|
|
using ytLive.Services;
|
|
using ytLive.Services.Compositor;
|
|
using ytLive.Services.Encoder;
|
|
|
|
namespace ytLive.Tests;
|
|
|
|
/// <summary>
|
|
/// The live frame producer (TASK 4 ship step 5): composites the active scene at
|
|
/// the tier's FPS and paces frames into the encoder. The integration test drives
|
|
/// the full lifecycle against fakes — real SceneCompositor + real FramePump, fake
|
|
/// IFfmpegEncoder — proving the composite frame actually reaches the encoder and
|
|
/// that stop tears the pump down cleanly. The units pin the failure edges: the
|
|
/// no-URL skip (the TASK 5 seam), re-entrancy, and encoder death.
|
|
/// </summary>
|
|
public class FramePumpTests
|
|
{
|
|
private sealed class FakeEncoder : IFfmpegEncoder
|
|
{
|
|
public readonly List<VideoFrame> Frames = new();
|
|
public int StartCount;
|
|
public int StopCount;
|
|
public bool Disposed;
|
|
public EncoderOptions? LastOptions;
|
|
public Exception? StartError;
|
|
public TaskCompletionSource FrameArrived = new(TaskCreationOptions.RunContinuationsAsynchronously);
|
|
|
|
public event EventHandler<StreamHealth>? HealthUpdated;
|
|
public event EventHandler<string>? ProcessFailed;
|
|
|
|
public Task StartAsync(EncoderOptions options, CancellationToken cancellationToken = default)
|
|
{
|
|
StartCount++;
|
|
LastOptions = options;
|
|
if (StartError != null) throw StartError;
|
|
return Task.CompletedTask;
|
|
}
|
|
|
|
public Task SubmitFrameAsync(VideoFrame frame, CancellationToken cancellationToken = default)
|
|
{
|
|
Frames.Add(frame);
|
|
FrameArrived.TrySetResult();
|
|
return Task.CompletedTask;
|
|
}
|
|
|
|
public Task StopAsync(CancellationToken cancellationToken = default)
|
|
{
|
|
StopCount++;
|
|
return Task.CompletedTask;
|
|
}
|
|
|
|
public void Dispose() => Disposed = true;
|
|
|
|
public void RaiseProcessFailed(string message) => ProcessFailed?.Invoke(this, message);
|
|
|
|
public void RaiseHealth(StreamHealth health) => HealthUpdated?.Invoke(this, health);
|
|
}
|
|
|
|
private static Scene BackgroundScene()
|
|
{
|
|
var scene = new Scene { Name = "Live" };
|
|
scene.Elements.Add(new Source { Type = SourceType.DisplayCapture, IsBackground = true, CaptureKey = "monitor:0" });
|
|
return scene;
|
|
}
|
|
|
|
private static FramePump NewPump(FakeEncoder encoder, Func<EncoderOptions?>? options = null,
|
|
Func<Scene?>? scene = null, Func<SceneElement, VideoFrame?>? resolve = null,
|
|
List<string>? log = null,
|
|
Func<(VideoFrame? Frame, SocialBarPosition Position)>? socialBar = null)
|
|
{
|
|
return new FramePump(
|
|
sceneProvider: scene ?? (() => BackgroundScene()),
|
|
frameResolver: resolve ?? (_ => null),
|
|
compositorOptions: () => new CompositorOptions
|
|
{
|
|
SourceRectX = 0, SourceRectY = 0, SourceRectWidth = 64, SourceRectHeight = 48,
|
|
OutputWidth = 64, OutputHeight = 48,
|
|
},
|
|
encoderOptions: options ?? (() => new EncoderOptions
|
|
{
|
|
RtmpUrl = "rtmp://a.rtmp.youtube.com/live2/abc",
|
|
Width = 64, Height = 48, Fps = 60,
|
|
}),
|
|
encoderFactory: () => encoder,
|
|
log: log != null ? m => log.Add(m) : null,
|
|
pacingDelay: async (_, _) => await Task.Yield(), // deterministic: no real waits
|
|
socialBar: socialBar);
|
|
}
|
|
|
|
private static void AssertColor(VideoFrame frame, int x, int y, byte r, byte g, byte b)
|
|
{
|
|
var i = (y * frame.Width + x) * 4;
|
|
var br = frame.BgraPixels[i + 2];
|
|
var bg = frame.BgraPixels[i + 1];
|
|
var bb = frame.BgraPixels[i];
|
|
Assert.True(
|
|
Math.Abs(br - r) <= 2 && Math.Abs(bg - g) <= 2 && Math.Abs(bb - b) <= 2,
|
|
$"pixel ({x},{y}): expected rgb({r},{g},{b}), got rgb({br},{bg},{bb})");
|
|
}
|
|
|
|
private sealed class FakeHttpHandler : HttpMessageHandler
|
|
{
|
|
public FakeHttpHandler(string body) => Body = body;
|
|
|
|
public string Body;
|
|
public string? LastUri;
|
|
public int RequestCount;
|
|
|
|
protected override Task<HttpResponseMessage> SendAsync(
|
|
HttpRequestMessage request, CancellationToken cancellationToken)
|
|
{
|
|
RequestCount++;
|
|
LastUri = request.RequestUri?.ToString();
|
|
return Task.FromResult(new HttpResponseMessage(HttpStatusCode.OK)
|
|
{
|
|
Content = new StringContent(Body, Encoding.UTF8, "application/json"),
|
|
});
|
|
}
|
|
}
|
|
|
|
/// <summary>The ONE integration test for TASK 5: the real service (hermetic
|
|
/// HTTP) returns the reusable stream's RTMP URL, the provider seam hands it
|
|
/// to the pump, and the pump starts the encoder with that exact URL — the
|
|
/// chain that makes go-live actually push. The stream is reused (listed),
|
|
/// never recreated (no POST).</summary>
|
|
[Fact]
|
|
public async Task ReusableStream_Url_From_Service_Feeds_Encoder_Startup()
|
|
{
|
|
var handler = new FakeHttpHandler(
|
|
"""
|
|
{"items":[{"id":"S456","cdn":{"isReusable":true,"ingestionInfo":{"ingestionAddress":"rtmp://a.rtmp.youtube.com/live2","streamName":"KEY123"}}}]}
|
|
""");
|
|
var auth = new YouTubeAuthService("test-id", "test-secret");
|
|
auth.SetSession(new YouTubeChannel
|
|
{
|
|
AccessToken = "acc-123",
|
|
TokenExpiry = DateTime.UtcNow.AddHours(1),
|
|
});
|
|
var service = new YouTubeStreamService(auth, new HttpClient(handler));
|
|
|
|
var url = (await service.GetOrCreateReusableStreamAsync())?.RtmpUrl;
|
|
var encoder = new FakeEncoder();
|
|
using var pump = NewPump(encoder, options: () => url == null
|
|
? null
|
|
: new EncoderOptions { RtmpUrl = url, Width = 64, Height = 48, Fps = 60 });
|
|
|
|
await pump.StartAsync();
|
|
|
|
Assert.Equal("rtmp://a.rtmp.youtube.com/live2/KEY123", url);
|
|
Assert.Equal(1, encoder.StartCount);
|
|
Assert.Equal("rtmp://a.rtmp.youtube.com/live2/KEY123", encoder.LastOptions?.RtmpUrl);
|
|
Assert.Equal(1, handler.RequestCount); // reuse path: list only, no insert
|
|
Assert.NotNull(handler.LastUri);
|
|
Assert.Contains("liveStreams", handler.LastUri);
|
|
Assert.Contains("mine=true", handler.LastUri);
|
|
|
|
await pump.StopAsync();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Start_CompositesScene_FeedsEncoder_StopsCleanly()
|
|
{
|
|
var red = SceneCompositorTests.Solid(64, 48, 255, 0, 0);
|
|
var encoder = new FakeEncoder();
|
|
using var pump = NewPump(encoder,
|
|
resolve: e => e is Source { IsBackground: true } ? red : null);
|
|
|
|
await pump.StartAsync();
|
|
Assert.True(pump.IsRunning);
|
|
Assert.Equal(1, encoder.StartCount);
|
|
|
|
await encoder.FrameArrived.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
Assert.NotEmpty(encoder.Frames);
|
|
var frame = encoder.Frames[0];
|
|
Assert.Equal(64, frame.Width);
|
|
Assert.Equal(48, frame.Height);
|
|
AssertColor(frame, 0, 0, 255, 0, 0); // the background really was composited in
|
|
|
|
await pump.StopAsync();
|
|
Assert.False(pump.IsRunning);
|
|
Assert.Equal(1, encoder.StopCount);
|
|
Assert.True(encoder.Disposed);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Start_WithoutRtmpUrl_SkipsEncoder()
|
|
{
|
|
var encoder = new FakeEncoder();
|
|
var log = new List<string>();
|
|
using var pump = NewPump(encoder, options: () => null, log: log);
|
|
|
|
await pump.StartAsync();
|
|
|
|
Assert.False(pump.IsRunning);
|
|
Assert.Equal(0, encoder.StartCount);
|
|
Assert.Contains(log, m => m.Contains("no output configured"));
|
|
|
|
await pump.StopAsync(); // no-op after a skipped start
|
|
Assert.Equal(0, encoder.StopCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Start_WhileRunning_IsNoop()
|
|
{
|
|
var encoder = new FakeEncoder();
|
|
using var pump = NewPump(encoder);
|
|
|
|
await pump.StartAsync();
|
|
await pump.StartAsync();
|
|
|
|
Assert.Equal(1, encoder.StartCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Stop_WithoutStart_IsNoop()
|
|
{
|
|
var encoder = new FakeEncoder();
|
|
using var pump = NewPump(encoder);
|
|
|
|
await pump.StopAsync();
|
|
|
|
Assert.Equal(0, encoder.StopCount);
|
|
Assert.False(pump.IsRunning);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Start_WithSocialBarSeam_PlacesBarAtTopThenBottomEdge()
|
|
{
|
|
var red = SceneCompositorTests.Solid(64, 48, 255, 0, 0);
|
|
var bar = new VideoFrame(64, 8, new byte[64 * 8 * 4]);
|
|
Array.Fill(bar.BgraPixels, (byte)255); // opaque white strip
|
|
|
|
var encoder = new FakeEncoder();
|
|
var position = SocialBarPosition.Top;
|
|
using var pump = NewPump(encoder,
|
|
resolve: e => e is Source { IsBackground: true } ? red : null,
|
|
socialBar: () => (bar, position));
|
|
|
|
await pump.StartAsync();
|
|
await encoder.FrameArrived.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
Assert.NotEmpty(encoder.Frames);
|
|
AssertColor(encoder.Frames[0], 0, 0, 255, 255, 255); // bar at the top edge
|
|
|
|
// Flip to Bottom: the pump re-reads the seam each frame, so a later frame
|
|
// lands the bar at the bottom edge (SourceRectHeight - bar height) and the
|
|
// top corner clears back to background.
|
|
position = SocialBarPosition.Bottom;
|
|
VideoFrame? flipped = null;
|
|
var deadline = DateTime.UtcNow.AddSeconds(5);
|
|
while (DateTime.UtcNow < deadline && flipped == null)
|
|
{
|
|
for (var i = 1; i < encoder.Frames.Count; i++)
|
|
{
|
|
var f = encoder.Frames[i];
|
|
if (f.BgraPixels[2] == 255 && f.BgraPixels[1] == 0 && f.BgraPixels[0] == 0)
|
|
{
|
|
flipped = f;
|
|
break;
|
|
}
|
|
}
|
|
if (flipped == null) await Task.Delay(10);
|
|
}
|
|
Assert.NotNull(flipped);
|
|
AssertColor(flipped!, 0, 47, 255, 255, 255); // bar sits on the bottom edge
|
|
|
|
await pump.StopAsync();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Start_EncoderThrows_RaisesFailed_AndDisposes()
|
|
{
|
|
var encoder = new FakeEncoder { StartError = new InvalidOperationException("access denied") };
|
|
var failed = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
|
|
using var pump = NewPump(encoder);
|
|
pump.Failed += (_, m) => failed.TrySetResult(m);
|
|
|
|
await pump.StartAsync();
|
|
|
|
Assert.False(pump.IsRunning);
|
|
var message = await failed.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
Assert.Contains("access denied", message);
|
|
Assert.True(encoder.Disposed);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task ProcessDeath_StopsPump_AndRaisesFailed()
|
|
{
|
|
var encoder = new FakeEncoder();
|
|
var failed = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
|
|
using var pump = NewPump(encoder);
|
|
pump.Failed += (_, m) => failed.TrySetResult(m);
|
|
|
|
await pump.StartAsync();
|
|
encoder.RaiseProcessFailed("FFmpeg exited with code 1");
|
|
|
|
var message = await failed.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
Assert.Contains("1", message);
|
|
|
|
// The pump self-stops (fire-and-forget); poll the definitive teardown
|
|
// marker — the encoder's disposal — rather than the IsRunning flag, which
|
|
// StopAsync clears before the loop has fully drained.
|
|
var deadline = DateTime.UtcNow.AddSeconds(5);
|
|
while (!encoder.Disposed && DateTime.UtcNow < deadline)
|
|
await Task.Delay(10);
|
|
Assert.True(encoder.Disposed);
|
|
Assert.False(pump.IsRunning);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task HealthUpdated_ForwardsEncoderHealth()
|
|
{
|
|
var encoder = new FakeEncoder();
|
|
var health = new TaskCompletionSource<StreamHealth>(TaskCreationOptions.RunContinuationsAsynchronously);
|
|
using var pump = NewPump(encoder);
|
|
pump.HealthUpdated += (_, h) => health.TrySetResult(h);
|
|
|
|
await pump.StartAsync();
|
|
encoder.RaiseHealth(new StreamHealth { Status = StreamStatus.Streaming, FPS = 59.9 });
|
|
|
|
var h = await health.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
Assert.Equal(StreamStatus.Streaming, h.Status);
|
|
Assert.Equal(59.9, h.FPS, 1);
|
|
}
|
|
|
|
/// <summary>The ONE integration test for the take-3 starvation fix (2026-09-03):
|
|
/// the frame interval is a DEADLINE — render cost eats into it, never piles on
|
|
/// top (the OBS libobs video-io.c pacing pattern). The old loop slept the full
|
|
/// interval AFTER each render, so 30 wall-seconds of 258ms renders muxed into a
|
|
/// 2.1s 60fps time-lapse (rawvideo stamps by arrival). Here the fake render
|
|
/// costs ≥60ms of a 200ms interval and the pacing seam records every requested
|
|
/// wait: no request may reach the full interval (that was the bug), none may be
|
|
/// negative, and a machine slow enough to blow every deadline is still honest
|
|
/// (zero waits < interval passes too).</summary>
|
|
[Fact]
|
|
public async Task Pump_Paces_To_The_Deadline_Compensating_Render_Cost()
|
|
{
|
|
const int fps = 5; // 200ms interval — generous margin over the fake cost
|
|
const int renderCostMs = 60; // Thread.Sleep is a guaranteed lower bound
|
|
var interval = TimeSpan.FromSeconds(1d / fps);
|
|
|
|
var encoder = new FakeEncoder();
|
|
var scene = BackgroundScene();
|
|
var delays = new List<TimeSpan>();
|
|
using var pump = new FramePump(
|
|
sceneProvider: () => scene,
|
|
frameResolver: _ =>
|
|
{
|
|
System.Threading.Thread.Sleep(renderCostMs);
|
|
return null;
|
|
},
|
|
compositorOptions: () => new CompositorOptions
|
|
{
|
|
SourceRectX = 0, SourceRectY = 0, SourceRectWidth = 64, SourceRectHeight = 48,
|
|
OutputWidth = 64, OutputHeight = 48,
|
|
},
|
|
encoderOptions: () => new EncoderOptions
|
|
{
|
|
RtmpUrl = "rtmp://a.rtmp.youtube.com/live2/abc",
|
|
Width = 64, Height = 48, Fps = fps,
|
|
},
|
|
encoderFactory: () => encoder,
|
|
pacingDelay: async (d, ct) =>
|
|
{
|
|
lock (delays) delays.Add(d);
|
|
await Task.Delay(d, ct); // honors the request, like the real Task.Delay default
|
|
});
|
|
|
|
await pump.StartAsync();
|
|
var deadline = DateTime.UtcNow.AddSeconds(10);
|
|
while (DateTime.UtcNow < deadline)
|
|
{
|
|
lock (delays) if (delays.Count >= 3) break;
|
|
await Task.Delay(10);
|
|
}
|
|
await pump.StopAsync();
|
|
|
|
List<TimeSpan> got;
|
|
lock (delays) got = delays.ToList();
|
|
Assert.True(got.Count >= 3, $"pump paced only {got.Count} frames in 10s");
|
|
foreach (var d in got)
|
|
{
|
|
Assert.True(d >= TimeSpan.Zero, $"requested wait {d} is negative");
|
|
Assert.True(d < interval,
|
|
$"requested wait {d} reaches the full {interval} interval — render cost must subtract from the deadline (the take-3 bug)");
|
|
}
|
|
}
|
|
}
|