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; /// /// 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. /// public class FramePumpTests { private sealed class FakeEncoder : IFfmpegEncoder { public readonly List Frames = new(); public readonly List Backings = 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? HealthUpdated; public event EventHandler? 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) { // Snapshot the bytes like the production FfmpegEncoder does (WriteAsync to // stdin copies before returning). Holding the reference would race the // pump's scratch pool legitimately reusing the buffer one tick later. Frames.Add(new VideoFrame(frame.Width, frame.Height, (byte[])frame.BgraPixels.Clone())); Backings.Add(frame.BgraPixels); // the pre-copy buffer — pool-reuse witness 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); public int DroppedFrames { get; set; } } 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? options = null, Func? scene = null, Func? resolve = null, List? 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 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"), }); } } /// 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). [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(); 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(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(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(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); } /// 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). [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(); 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 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)"); } } /// The ONE integration test for the slice-15 pacing fix (2026-09-14): /// a render that overruns the deadline must NOT skip the missed slots — a skip /// authors fast playback: ty-1726 rendered one frame per 35ms overrun into a /// 60fps container (163 frames for 2.93s of audio) and cut the audio tail off. /// OBS fills an overrun's slots by DUPLICATING the newest frame (libobs /// media-io/video-io.c; docs.obsproject.com/backend-design: "If the video frame /// queue is full, it will duplicate the last frame") so recording duration stays /// equal to wall duration. Here a 60fps pump renders every frame at ~35ms cost /// (a genuine overrun): the slice-10 keep-fresh loop authored ~1 per 35ms; the /// slot-filling loop must emit a frame for ~every 16.6ms of wall time. [Fact] public async Task Pump_Overrun_Renders_EmitsEverySlot_NotSkipped() { const int fps = 60; var encoder = new FakeEncoder(); using var pump = NewPump(encoder, resolve: _ => { System.Threading.Thread.Sleep(35); return null; }); await pump.StartAsync(); await Task.Delay(7000); // ~420 slots of wall time await pump.StopAsync(); // Slot-faithful: ~1 frame per 16.6ms -> ~415 in the window. Skipping (the // slice-10 behavior): ~1 per 35ms -> ~200. 0.65 * fps * (elapsed seconds) // is a middle line a skipped pump can never cross, a filling one only fails // on a pathological run. Assert.True(encoder.Frames.Count >= (int)(7.0 * fps * 0.65), $"pump emitted {encoder.Frames.Count} frames in ~7s wall at {fps}fps — " + $"{7.0 / encoder.Frames.Count * fps:F1}x playback — an overrun must " + "DUPLICATE the newest frame per missed slot (OBS), not skip it"); } /// The ONE integration test for the take-4 scratch pool: the pump recycles /// its master buffers across frames (the 8.3MB-per-tick LOH churn that cost GC /// stalls inside "render") while EVERY frame's content stays correct — stale bytes /// from an earlier tick would show as a wrong color. The backdrop alternates /// red/blue per tick (submit snapshots pin the color), and a repeated backing-array /// identity proves the pool actually recycles (an inert pool keeps the pixel /// assertions green while leaving the churn in place). [Fact] public async Task Pump_Pools_ScratchBuffers_Across_Frames_Without_Stale_Pixels() { var red = SceneCompositorTests.Solid(64, 48, 255, 0, 0); var blue = SceneCompositorTests.Solid(64, 48, 0, 0, 255); var flip = 0; var encoder = new FakeEncoder(); using var pump = NewPump(encoder, resolve: _ => flip++ % 2 == 0 ? red : blue); await pump.StartAsync(); var deadline = DateTime.UtcNow.AddSeconds(10); while (DateTime.UtcNow < deadline && encoder.Frames.Count < 6) await Task.Delay(10); await pump.StopAsync(); Assert.True(encoder.Frames.Count >= 6, $"only {encoder.Frames.Count} frames produced"); for (var i = 0; i < encoder.Frames.Count; i++) { if (i % 2 == 0) AssertColor(encoder.Frames[i], 0, 0, 255, 0, 0); else AssertColor(encoder.Frames[i], 0, 0, 0, 0, 255); } Assert.True(encoder.Backings.Distinct().Count() < encoder.Backings.Count, "every frame rode a fresh buffer — the scratch pool is inert"); } /// One integration test for the take-8 off-UI fix: the pump loop must NOT /// run on the starting thread's SynchronizationContext (StartAsync fires from a UI /// command handler; the old loop inherited it, so every frame competed with the /// live preview for the dispatcher — the "wait 10ms after a 22ms render" that the /// quantum stats exposed). An inline-pumping sync context makes the OLD code run /// its resolver on the starting thread; the Task.Run'd loop must never do so. [Fact] public async Task Pump_Produces_OffTheStartingContext() { var startThread = Environment.CurrentManagedThreadId; var resolverThreads = new List(); var firstFrame = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var encoder = new FakeEncoder(); using var pump = NewPump(encoder, resolve: _ => { lock (resolverThreads) resolverThreads.Add(Environment.CurrentManagedThreadId); firstFrame.TrySetResult(); return null; }); var prev = SynchronizationContext.Current; SynchronizationContext.SetSynchronizationContext(new InlineSyncContext()); try { await pump.StartAsync(); await firstFrame.Task.WaitAsync(TimeSpan.FromSeconds(5)); } finally { SynchronizationContext.SetSynchronizationContext(prev); await pump.StopAsync(); } Assert.NotEmpty(resolverThreads); Assert.All(resolverThreads, id => Assert.True(id != startThread, $"resolver ran on the starting thread ({startThread}) — the loop is back on the caller context")); } private sealed class InlineSyncContext : SynchronizationContext { public override void Post(SendOrPostCallback d, object? state) => d(state); public override void Send(SendOrPostCallback d, object? state) => d(state); public override SynchronizationContext CreateCopy() => this; } }