using System.Net; using ErsatzTV.Infrastructure.Streaming; using Microsoft.Extensions.Logging; using NSubstitute; using NUnit.Framework; using Shouldly; namespace ErsatzTV.Infrastructure.Tests.Streaming; [TestFixture] public class HttpRemoteStreamProberTests { private const string Url = "http://localhost:8409/media/jellyfin/abc123"; [Test] public async Task Should_Report_Unavailable_On_404_From_The_Media_Server() { // a media-server 404 arrives after our /media/... endpoint redirected, so the response's // final request uri is the media server's, not the probe url HttpRemoteStreamProber prober = ProberReturning( HttpStatusCode.NotFound, finalUri: "http://jellyfin:8096/Videos/abc123/stream?static=true"); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeFalse(); } // ersatztv#473 review finding: our OWN /media/{provider}/... endpoint 404s when the media source // is unconfigured or momentarily missing. Failing closed there would blank every item on that // source, which is exactly what the fail-open contract exists to prevent. [Test] public async Task Should_Fail_Open_On_404_That_Was_Not_Redirected() { HttpRemoteStreamProber prober = ProberReturning(HttpStatusCode.NotFound, finalUri: Url); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeTrue(); } // a plex key can contain spaces/unicode; pin that an un-redirected 404 on such a url still fails // OPEN. (This passes against a naive string comparison too - Uri.ToString() unescapes - so it // guards the behaviour, not the implementation choice.) [Test] public async Task Should_Fail_Open_On_404_For_An_Unredirected_Url_Needing_Escaping() { const string plexUrl = "http://localhost:8409/media/plex/1/library/parts/1/a file.mkv"; HttpRemoteStreamProber prober = ProberReturning(HttpStatusCode.NotFound, finalUri: plexUrl); bool result = await prober.IsAvailable(plexUrl, CancellationToken.None); result.ShouldBeTrue(); } [TestCase(HttpStatusCode.OK)] [TestCase(HttpStatusCode.PartialContent)] [TestCase(HttpStatusCode.NoContent)] public async Task Should_Report_Available_On_Success(HttpStatusCode statusCode) { HttpRemoteStreamProber prober = ProberReturning(statusCode); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeTrue(); } // the fail-open contract: a probe that cannot answer must never block a tune that would // otherwise have worked. these cases exist so a future refactor can't silently invert it. [TestCase(HttpStatusCode.InternalServerError)] [TestCase(HttpStatusCode.BadGateway)] [TestCase(HttpStatusCode.Unauthorized)] [TestCase(HttpStatusCode.Forbidden)] public async Task Should_Fail_Open_On_Other_Status_Codes(HttpStatusCode statusCode) { HttpRemoteStreamProber prober = ProberReturning(statusCode); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeTrue(); } // a server that ignores `Range: bytes=0-0` answers 200 with the WHOLE FILE. The probe must not // read it -- buffering a video on the streaming hot path would be far worse than the aborted // socket the drain was added to avoid. (Review finding against the first fix commit.) [Test] public async Task Should_Not_Read_The_Body_When_The_Server_Ignores_The_Range_Request() { var body = new TrackingStream(64 * 1024 * 1024); var response = new HttpResponseMessage(HttpStatusCode.OK) { Content = new StreamContent(body) }; var prober = new HttpRemoteStreamProber( new StubHttpClientFactory(new FixedResponseHttpMessageHandler(response)), Substitute.For>()); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeTrue(); body.BytesRead.ShouldBe(0); } // the counterpart: when the server DID honour the range, the one byte is read so the connection // goes back to the pool rather than being aborted [Test] public async Task Should_Drain_The_Single_Byte_When_The_Server_Honours_The_Range_Request() { var body = new TrackingStream(1); var response = new HttpResponseMessage(HttpStatusCode.PartialContent) { Content = new StreamContent(body) }; var prober = new HttpRemoteStreamProber( new StubHttpClientFactory(new FixedResponseHttpMessageHandler(response)), Substitute.For>()); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeTrue(); body.BytesRead.ShouldBe(1); } [Test] public async Task Should_Fail_Open_On_Transport_Failure() { var prober = new HttpRemoteStreamProber( new StubHttpClientFactory(new ThrowingHttpMessageHandler(new HttpRequestException("no route to host"))), Substitute.For>()); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeTrue(); } [Test] public async Task Should_Fail_Open_On_Timeout() { var prober = new HttpRemoteStreamProber( new StubHttpClientFactory(new ThrowingHttpMessageHandler(new TaskCanceledException("timed out"))), Substitute.For>()); bool result = await prober.IsAvailable(Url, CancellationToken.None); result.ShouldBeTrue(); } // caller cancellation (shutdown / client disconnect) is a genuine signal, NOT a probe failure -- // swallowing it would let the handler go on building an ffmpeg command on a dead token. [Test] public async Task Should_Propagate_Caller_Cancellation() { HttpRemoteStreamProber prober = ProberReturning(HttpStatusCode.OK); using var cts = new CancellationTokenSource(); await cts.CancelAsync(); await Should.ThrowAsync(() => prober.IsAvailable(Url, cts.Token)); } private static HttpRemoteStreamProber ProberReturning(HttpStatusCode statusCode, string finalUri = null) => new( new StubHttpClientFactory(new StatusCodeHttpMessageHandler(statusCode, finalUri)), Substitute.For>()); private sealed class StubHttpClientFactory(HttpMessageHandler handler) : IHttpClientFactory { public HttpClient CreateClient(string name) => new(handler, disposeHandler: false); } private sealed class StatusCodeHttpMessageHandler(HttpStatusCode statusCode, string finalUri = null) : HttpMessageHandler { protected override Task SendAsync( HttpRequestMessage request, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); // HttpClient rewrites RequestMessage.RequestUri to the final hop when it follows a // redirect; finalUri lets a test stand in for "the media server answered this". if (finalUri is not null) { request.RequestUri = new Uri(finalUri); } return Task.FromResult(new HttpResponseMessage(statusCode) { RequestMessage = request }); } } private sealed class FixedResponseHttpMessageHandler(HttpResponseMessage response) : HttpMessageHandler { protected override Task SendAsync( HttpRequestMessage request, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); response.RequestMessage = request; return Task.FromResult(response); } } /// A readable stream that records how many bytes were actually pulled from it. private sealed class TrackingStream(long length) : Stream { public int BytesRead { get; private set; } public override bool CanRead => true; public override bool CanSeek => false; public override bool CanWrite => false; public override long Length => length; public override long Position { get => BytesRead; set => throw new NotSupportedException(); } public override int Read(byte[] buffer, int offset, int count) { if (BytesRead >= length) { return 0; } int toRead = (int)Math.Min(count, length - BytesRead); Array.Clear(buffer, offset, toRead); BytesRead += toRead; return toRead; } public override void Flush() { } public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException(); public override void SetLength(long value) => throw new NotSupportedException(); public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException(); } private sealed class ThrowingHttpMessageHandler(Exception exception) : HttpMessageHandler { protected override Task SendAsync( HttpRequestMessage request, CancellationToken cancellationToken) => Task.FromException(exception); } }