From 818650484064fe30058f36efb4bd8ecb712f5433 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:00:30 +0000 Subject: [PATCH 1/8] fix: lower the source generator's Roslyn pin to the net10.0 GA compiler Microsoft.CodeAnalysis.CSharp 5.6.0 made the generator unloadable (CS9057) on any SDK shipping an older compiler, including current distro packages of the .NET 10 SDK. Pin to 5.0.0 - the GA compiler - so standard installs build; the generator only uses ForAttributeWithMetadataName (Roslyn 4.3+). Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- src/Spangle.SourceGenerator/Spangle.SourceGenerator.csproj | 6 ++++-- .../Spangle.SourceGenerator.Tests.csproj | 4 ++-- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/src/Spangle.SourceGenerator/Spangle.SourceGenerator.csproj b/src/Spangle.SourceGenerator/Spangle.SourceGenerator.csproj index 2264334..4c946ce 100644 --- a/src/Spangle.SourceGenerator/Spangle.SourceGenerator.csproj +++ b/src/Spangle.SourceGenerator/Spangle.SourceGenerator.csproj @@ -13,8 +13,10 @@ - - + + + diff --git a/tests/Spangle.SourceGenerator.Tests/Spangle.SourceGenerator.Tests.csproj b/tests/Spangle.SourceGenerator.Tests/Spangle.SourceGenerator.Tests.csproj index 4a4adc1..492be68 100644 --- a/tests/Spangle.SourceGenerator.Tests/Spangle.SourceGenerator.Tests.csproj +++ b/tests/Spangle.SourceGenerator.Tests/Spangle.SourceGenerator.Tests.csproj @@ -14,8 +14,8 @@ - - + + From 337c34c0bfc27a8416e76ccbbcb30302b2cc2a61 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:06:49 +0000 Subject: [PATCH 2/8] feat: TLS everywhere it was missing - HTTPS delivery/management, RTMPS, metadata token Every listener gains an optional Tls block (PKCS#12 or PEM cert/key): - Http.Tls: HLS/DASH delivery over HTTPS - Management.Tls: keeps the Bearer token off the wire in cleartext - Rtmp.Tls: RTMPS via the TLS connection middleware ahead of the handler Options validation fails fast on an enabled Tls block with no certificate. The metadata injection endpoint on the public delivery port gets its own optional Bearer token (Http.MetadataInjectionToken); the constant-time comparison moves to a shared TokenGate used by the management gate too. Verified live: HTTPS 200 on both ports, plaintext refused, TLSv1.3 handshake on 1935, 401 without the metadata token / through with it. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- .../SpangleServiceCollectionExtensions.cs | 7 +++ .../Management/SpangleManagementEndpoints.cs | 15 +----- .../SpangleMediaServerOptions.cs | 52 ++++++++++++++++++- src/Spangle.Extensions.Kestrel/TokenGate.cs | 32 ++++++++++++ .../WebHostBuilderSpangleExtensions.cs | 23 ++++++-- src/Spangle.MediaServer/Program.cs | 7 +++ src/Spangle.MediaServer/spanglesettings.yaml | 16 +++++- 7 files changed, 131 insertions(+), 21 deletions(-) create mode 100644 src/Spangle.Extensions.Kestrel/TokenGate.cs diff --git a/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs b/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs index ebb2365..22da364 100644 --- a/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs +++ b/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs @@ -27,6 +27,13 @@ public static void AddSpangle(this IServiceCollection services) && IPAddress.IsLoopback(address); }, "Management.Token is required when Management.BindAddress is not a loopback address") + .Validate(static options => + { + // a TLS block without a certificate is a misconfiguration, not "plaintext please" + static bool Ok(TlsOptions tls) => !tls.Enabled || !string.IsNullOrEmpty(tls.CertificatePath); + return Ok(options.Rtmp.Tls) && Ok(options.Http.Tls) && Ok(options.Management.Tls); + }, + "Tls.CertificatePath is required wherever Tls.Enabled is set") .ValidateOnStart(); services.AddSingleton(); services.AddSingleton(); diff --git a/src/Spangle.Extensions.Kestrel/Management/SpangleManagementEndpoints.cs b/src/Spangle.Extensions.Kestrel/Management/SpangleManagementEndpoints.cs index c945020..02b54f8 100644 --- a/src/Spangle.Extensions.Kestrel/Management/SpangleManagementEndpoints.cs +++ b/src/Spangle.Extensions.Kestrel/Management/SpangleManagementEndpoints.cs @@ -2,8 +2,6 @@ using System.Reflection; using System.Runtime.CompilerServices; using System.Runtime.InteropServices; -using System.Security.Cryptography; -using System.Text; using System.Text.Json; using Microsoft.Extensions.Options; using Spangle.Spinner; @@ -107,10 +105,7 @@ static async (string streamKey, HttpRequest request, TimedMetadataHub hub) => return null; // loopback-only bind is enforced by options validation } - const string prefix = "Bearer "; - string auth = ctx.Request.Headers.Authorization.ToString(); - if (auth.StartsWith(prefix, StringComparison.Ordinal) - && FixedTimeEquals(auth.AsSpan(prefix.Length).Trim(), opt.Token)) + if (TokenGate.Matches(ctx.Request.Headers.Authorization.ToString(), opt.Token)) { return null; } @@ -118,14 +113,6 @@ static async (string streamKey, HttpRequest request, TimedMetadataHub hub) => return Results.Unauthorized(); } - private static bool FixedTimeEquals(ReadOnlySpan provided, string expected) - { - Span providedUtf8 = provided.Length <= 64 ? stackalloc byte[256] : new byte[provided.Length * 4]; - int written = Encoding.UTF8.GetBytes(provided, providedUtf8); - byte[] expectedUtf8 = Encoding.UTF8.GetBytes(expected); - return CryptographicOperations.FixedTimeEquals(providedUtf8[..written], expectedUtf8); - } - private static async IAsyncEnumerable> StatsFeedAsync( ManagementStatsService stats, [EnumeratorCancellation] CancellationToken ct) { diff --git a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs index 5d4c5fa..b8a3937 100644 --- a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs +++ b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs @@ -1,4 +1,5 @@ using System.ComponentModel.DataAnnotations; +using System.Security.Cryptography.X509Certificates; namespace Spangle.Extensions.Kestrel; @@ -37,6 +38,40 @@ public class ManagementOptions /// (Authorization: Bearer ...). Mandatory for non-loopback binds. /// public string? Token { get; set; } + + /// + /// TLS for the management port. Without it the Bearer token travels in + /// cleartext, so anything beyond loopback should turn this on. + /// + public TlsOptions Tls { get; set; } = new(); +} + +/// +/// TLS for one listener. Plaintext until is set, then +/// is a PKCS#12/PFX file (with +/// if the file has one) — or a PEM +/// certificate when points at the PEM private key. +/// +public class TlsOptions +{ + public bool Enabled { get; set; } + + public string? CertificatePath { get; set; } + + public string? CertificatePassword { get; set; } + + public string? KeyPath { get; set; } + + internal X509Certificate2 LoadCertificate() + { + if (string.IsNullOrEmpty(CertificatePath)) + { + throw new InvalidOperationException("Tls.CertificatePath is required when Tls.Enabled"); + } + return string.IsNullOrEmpty(KeyPath) + ? X509CertificateLoader.LoadPkcs12FromFile(CertificatePath, CertificatePassword) + : X509Certificate2.CreateFromPemFile(CertificatePath, KeyPath); + } } public class RtmpOptions : MediaProtocolOptions @@ -55,6 +90,9 @@ public class RtmpOptions : MediaProtocolOptions /// metadata carried in the HLS output. Adds one spinner hop to the pipeline. /// public bool TimedMetadata { get; set; } = true; + + /// RTMPS: TLS on the RTMP listener (publishers connect with rtmps://) + public TlsOptions Tls { get; set; } = new(); } public class SrtOptions : MediaProtocolOptions @@ -75,10 +113,20 @@ public class HttpOptions /// /// Enables POST /api/streams/{key}/metadata: injects timed ID3 metadata into a - /// live session. The endpoint has no authentication of its own — protect it at - /// the network level or front it with your own middleware. + /// live session. Set to require a Bearer + /// token; without one, protect the endpoint at the network level instead. /// public bool MetadataInjection { get; set; } = true; + + /// + /// Bearer token required on metadata injection requests when set + /// (Authorization: Bearer ...). The endpoint lives on the public + /// delivery port, so set this anywhere the port is reachable by viewers. + /// + public string? MetadataInjectionToken { get; set; } + + /// HTTPS for the delivery port (HLS/DASH and the test player) + public TlsOptions Tls { get; set; } = new(); } public class HlsOptions : MediaProtocolOptions diff --git a/src/Spangle.Extensions.Kestrel/TokenGate.cs b/src/Spangle.Extensions.Kestrel/TokenGate.cs new file mode 100644 index 0000000..5037bdf --- /dev/null +++ b/src/Spangle.Extensions.Kestrel/TokenGate.cs @@ -0,0 +1,32 @@ +using System.Security.Cryptography; +using System.Text; + +namespace Spangle.Extensions.Kestrel; + +/// +/// Constant-time Bearer token comparison, shared by every token-gated endpoint +/// (the management surface and metadata injection). +/// +public static class TokenGate +{ + private const string Prefix = "Bearer "; + + /// Checks an Authorization header value against the expected token. + public static bool Matches(string authorizationHeader, string expected) + { + ArgumentNullException.ThrowIfNull(authorizationHeader); + ArgumentNullException.ThrowIfNull(expected); + return authorizationHeader.StartsWith(Prefix, StringComparison.Ordinal) + && FixedTimeEquals(authorizationHeader.AsSpan(Prefix.Length).Trim(), expected); + } + + private static bool FixedTimeEquals(ReadOnlySpan provided, string expected) + { + Span providedUtf8 = provided.Length <= 64 + ? stackalloc byte[256] + : new byte[Encoding.UTF8.GetMaxByteCount(provided.Length)]; + int written = Encoding.UTF8.GetBytes(provided, providedUtf8); + byte[] expectedUtf8 = Encoding.UTF8.GetBytes(expected); + return CryptographicOperations.FixedTimeEquals(providedUtf8[..written], expectedUtf8); + } +} diff --git a/src/Spangle.Extensions.Kestrel/WebHostBuilderSpangleExtensions.cs b/src/Spangle.Extensions.Kestrel/WebHostBuilderSpangleExtensions.cs index 5171ad8..e057cd1 100644 --- a/src/Spangle.Extensions.Kestrel/WebHostBuilderSpangleExtensions.cs +++ b/src/Spangle.Extensions.Kestrel/WebHostBuilderSpangleExtensions.cs @@ -27,14 +27,20 @@ public static void ConfigureSpangle(this KestrelServerOptions options) { var opt = options.ApplicationServices.GetRequiredService>(); options.ListenAnyIP(opt.Value.Rtmp.Port, - listenOptions => { listenOptions.UseConnectionHandler(); }); + listenOptions => + { + // RTMPS: the TLS middleware wraps the raw connection before the handler sees it + ApplyTls(listenOptions, opt.Value.Rtmp.Tls); + listenOptions.UseConnectionHandler(); + }); // Explicit Listen* calls override URL-based configuration, so HTTP must be explicit too - options.ListenAnyIP(opt.Value.Http.Port); + options.ListenAnyIP(opt.Value.Http.Port, listenOptions => ApplyTls(listenOptions, opt.Value.Http.Tls)); // The management surface (console + control API) never shares the delivery port ManagementOptions management = opt.Value.Management; if (management.Enabled) { - options.Listen(System.Net.IPAddress.Parse(management.BindAddress), management.Port); + options.Listen(System.Net.IPAddress.Parse(management.BindAddress), management.Port, + listenOptions => ApplyTls(listenOptions, management.Tls)); } var loggerFactory = options.ApplicationServices.GetService(); if (loggerFactory != null) @@ -42,4 +48,15 @@ public static void ConfigureSpangle(this KestrelServerOptions options) SpangleLogManager.SetLoggerFactory(loggerFactory); } } + + [System.Diagnostics.CodeAnalysis.SuppressMessage("Reliability", "CA2000:Dispose objects before losing scope", + Justification = "The certificate must stay alive as long as the listener - the process lifetime")] + private static void ApplyTls(ListenOptions listenOptions, TlsOptions tls) + { + if (!tls.Enabled) + { + return; + } + listenOptions.UseHttps(tls.LoadCertificate()); + } } diff --git a/src/Spangle.MediaServer/Program.cs b/src/Spangle.MediaServer/Program.cs index a8c1642..c8ee622 100644 --- a/src/Spangle.MediaServer/Program.cs +++ b/src/Spangle.MediaServer/Program.cs @@ -92,9 +92,16 @@ if (options.Http.MetadataInjection) { var metadataHub = app.Services.GetRequiredService(); + string? metadataToken = options.Http.MetadataInjectionToken; app.MapPost("/api/streams/{streamKey}/metadata", async (string streamKey, HttpRequest request) => { + if (!string.IsNullOrEmpty(metadataToken) + && !TokenGate.Matches(request.Headers.Authorization.ToString(), metadataToken)) + { + request.HttpContext.Response.Headers.WWWAuthenticate = "Bearer"; + return Results.Unauthorized(); + } using var body = await JsonDocument.ParseAsync(request.Body, cancellationToken: request.HttpContext.RequestAborted).ConfigureAwait(false); if (!body.RootElement.TryGetProperty("name", out JsonElement nameElement) || nameElement.GetString() is not { Length: > 0 } name) diff --git a/src/Spangle.MediaServer/spanglesettings.yaml b/src/Spangle.MediaServer/spanglesettings.yaml index 125c97c..949f7ca 100644 --- a/src/Spangle.MediaServer/spanglesettings.yaml +++ b/src/Spangle.MediaServer/spanglesettings.yaml @@ -4,14 +4,23 @@ Spangle: Port: 1935 # AudioOnlyFallbackMs: 3000 # audio with no video in sight becomes audio-only (0 = off) # TimedMetadata: true # AMF data events -> timed ID3 in the HLS output + # Tls: # RTMPS (publishers connect with rtmps://) + # Enabled: true + # CertificatePath: cert.pfx # PKCS#12; or a PEM cert with KeyPath + # CertificatePassword: "" + # KeyPath: cert.key # set for a PEM cert/key pair Srt: Enabled: true Port: 9998 # Passphrase: "correct horse battery staple" # optional wire encryption (10-79 bytes) Http: Port: 8080 - # MetadataInjection: true # POST /api/streams/{key}/metadata -> timed ID3 in the - # # output. No built-in auth: protect it at the network level. + # MetadataInjection: true # POST /api/streams/{key}/metadata -> timed ID3 in the + # # output. Token below, or protect it at the network level. + # MetadataInjectionToken: "" # Bearer token for metadata injection on this public port + # Tls: # HTTPS for HLS/DASH delivery (same shape as Rtmp.Tls) + # Enabled: true + # CertificatePath: cert.pfx Management: # The control API (and web console) listens on its own port, never on the # delivery port. Binding beyond loopback requires Token. @@ -19,6 +28,9 @@ Spangle: Port: 8081 # BindAddress: 127.0.0.1 # e.g. 0.0.0.0 to manage remotely (Token required) # Token: "" # Bearer token for every /api/manage request + # Tls: # HTTPS here keeps the token off the wire in cleartext + # Enabled: true + # CertificatePath: cert.pfx Hls: Enabled: true # Storage: Memory # "Memory" (default): the live window is served from process From 6bf5a46422ffe978af12717a41c773c59afc23f4 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:09:08 +0000 Subject: [PATCH 3/8] feat: config-driven publish allowlist - Publish.AllowedStreamNames Out of the box the only publish policy was allow-all. Listing names under Publish.AllowedStreamNames now switches the built-in authorizer to an exact-match (case-sensitive) allowlist over the raw published name (RTMP publish name / SRT streamid), keeping last-wins takeover for valid names so reconnects still work. A custom IPublishAuthorizer in DI still overrides the built-in policy entirely. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- .../AllowListPublishAuthorizer.cs | 25 ++++++++++ .../SpangleServiceCollectionExtensions.cs | 12 ++++- .../SpangleMediaServerOptions.cs | 15 ++++++ src/Spangle.MediaServer/spanglesettings.yaml | 4 ++ .../AllowListPublishAuthorizerTests.cs | 48 +++++++++++++++++++ 5 files changed, 102 insertions(+), 2 deletions(-) create mode 100644 src/Spangle.Core/AllowListPublishAuthorizer.cs create mode 100644 tests/Spangle.Core.Tests/AllowListPublishAuthorizerTests.cs diff --git a/src/Spangle.Core/AllowListPublishAuthorizer.cs b/src/Spangle.Core/AllowListPublishAuthorizer.cs new file mode 100644 index 0000000..9c8da5c --- /dev/null +++ b/src/Spangle.Core/AllowListPublishAuthorizer.cs @@ -0,0 +1,25 @@ +using System.Collections.Frozen; + +namespace Spangle; + +/// +/// Allows only the configured stream names, matched exactly against the raw name the +/// publisher presented (RTMP publish name / SRT streamid). Contested names keep the +/// last-wins takeover policy of , so a reconnect +/// by a holder of a valid name still just works. +/// +public sealed class AllowListPublishAuthorizer(IEnumerable allowedStreamNames) : IPublishAuthorizer +{ + private readonly FrozenSet _allowed = allowedStreamNames.ToFrozenSet(StringComparer.Ordinal); + + public ValueTask AuthorizeAsync(PublishRequest request, CancellationToken ct) + { + ArgumentNullException.ThrowIfNull(request); + return new ValueTask( + !_allowed.Contains(request.StreamName) + ? PublishDecision.Deny + : request.ExistingSession is null + ? PublishDecision.Allow + : PublishDecision.Takeover); + } +} diff --git a/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs b/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs index 22da364..7fcd299 100644 --- a/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs +++ b/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs @@ -40,8 +40,16 @@ public static void AddSpangle(this IServiceCollection services) services.AddSingleton(); services.AddSingleton(); // apps replace this to implement their own publish policy (deny lists, - // key validation, first-wins, ...); the default is allow-all + last-wins - services.TryAddSingleton(); + // key validation, first-wins, ...); the built-in policy is allow-all + + // last-wins, or an exact-match allowlist when Publish.AllowedStreamNames is set + services.TryAddSingleton(static provider => + { + PublishOptions publish = + provider.GetRequiredService>().Value.Publish; + return publish.AllowedStreamNames.Count > 0 + ? new AllowListPublishAuthorizer(publish.AllowedStreamNames) + : new DefaultPublishAuthorizer(); + }); // Memory (default) serves the live window without touching disk; // File makes the output an on-disk archive as well services.TryAddSingleton(static provider => diff --git a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs index b8a3937..0c98d85 100644 --- a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs +++ b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs @@ -15,6 +15,21 @@ public class SpangleMediaServerOptions public HlsOptions Hls { get; set; } = new(); public HttpOptions Http { get; set; } = new(); public ManagementOptions Management { get; set; } = new(); + public PublishOptions Publish { get; set; } = new(); +} + +/// +/// Publish authorization policy. The default stays allow-all + last-wins; listing +/// stream names here switches the built-in authorizer to an exact-match allowlist. +/// A custom registered in DI always wins over both. +/// +public class PublishOptions +{ + /// + /// Stream names (RTMP publish names / SRT streamids) allowed to publish. + /// Empty keeps the allow-all policy. + /// + public IList AllowedStreamNames { get; } = []; } /// diff --git a/src/Spangle.MediaServer/spanglesettings.yaml b/src/Spangle.MediaServer/spanglesettings.yaml index 949f7ca..cafcb68 100644 --- a/src/Spangle.MediaServer/spanglesettings.yaml +++ b/src/Spangle.MediaServer/spanglesettings.yaml @@ -31,6 +31,10 @@ Spangle: # Tls: # HTTPS here keeps the token off the wire in cleartext # Enabled: true # CertificatePath: cert.pfx + # Publish: + # AllowedStreamNames: # exact-match allowlist for publishers (RTMP name / + # - my-stream-key # SRT streamid); empty = allow all. A custom + # - another-key # IPublishAuthorizer in DI overrides this entirely. Hls: Enabled: true # Storage: Memory # "Memory" (default): the live window is served from process diff --git a/tests/Spangle.Core.Tests/AllowListPublishAuthorizerTests.cs b/tests/Spangle.Core.Tests/AllowListPublishAuthorizerTests.cs new file mode 100644 index 0000000..87533a0 --- /dev/null +++ b/tests/Spangle.Core.Tests/AllowListPublishAuthorizerTests.cs @@ -0,0 +1,48 @@ +using System.Net; + +namespace Spangle.Tests; + +public class AllowListPublishAuthorizerTests +{ + private static PublishRequest Request(string name, ExistingSessionInfo? existing = null) => new() + { + Protocol = "TEST", + StreamName = name, + StreamKey = StreamKeys.Sanitize(name), + RemoteEndPoint = new IPEndPoint(IPAddress.Loopback, 1), + ExistingSession = existing, + }; + + [Fact] + public async Task ListedNameIsAllowed() + { + var authorizer = new AllowListPublishAuthorizer(["live/x", "live/y"]); + (await authorizer.AuthorizeAsync(Request("live/x"), CancellationToken.None)) + .Should().Be(PublishDecision.Allow); + } + + [Fact] + public async Task UnlistedNameIsDenied() + { + var authorizer = new AllowListPublishAuthorizer(["live/x"]); + (await authorizer.AuthorizeAsync(Request("live/z"), CancellationToken.None)) + .Should().Be(PublishDecision.Deny); + } + + [Fact] + public async Task MatchIsExactAndCaseSensitive() + { + var authorizer = new AllowListPublishAuthorizer(["live/x"]); + (await authorizer.AuthorizeAsync(Request("LIVE/X"), CancellationToken.None)) + .Should().Be(PublishDecision.Deny, "a stream name is a credential; matching must be exact"); + } + + [Fact] + public async Task ContestedListedNameKeepsLastWinsTakeover() + { + var authorizer = new AllowListPublishAuthorizer(["live/x"]); + var existing = new ExistingSessionInfo { Id = "A", StartedAt = DateTimeOffset.UtcNow }; + (await authorizer.AuthorizeAsync(Request("live/x", existing), CancellationToken.None)) + .Should().Be(PublishDecision.Takeover, "a reconnect by a valid key holder must not be blocked"); + } +} From b3ce5f01a746207cd88053a4a6f3a93cdc95ef1a Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:12:57 +0000 Subject: [PATCH 4/8] feat: ended streams no longer live forever - Hls.EndedStreamTtlSeconds MemoryHLSStorage kept every ended stream's final window until the same key published again, so distinct short-lived stream names accumulated without bound. Streams now carry a last-write stamp; EvictIdleStreams (new IEvictingHLSStorage capability) frees any stream idle past the TTL unless the publish registry still owns its key, releasing blocked LL-DASH readers on the way out. A hosted sweeper runs it a few times per TTL. Default TTL 300s; 0 restores keep-until-republish. File storage is an archive by design and is never cleaned. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- src/Spangle.Core/PublishSessionRegistry.cs | 3 + src/Spangle.Core/Transport/HLS/HLSStorage.cs | 77 +++++++++++++++++-- .../SpangleServiceCollectionExtensions.cs | 1 + .../HlsEvictionService.cs | 39 ++++++++++ .../SpangleMediaServerOptions.cs | 8 ++ src/Spangle.MediaServer/spanglesettings.yaml | 2 + .../HLS/MemoryStorageEvictionTests.cs | 73 ++++++++++++++++++ 7 files changed, 198 insertions(+), 5 deletions(-) create mode 100644 src/Spangle.Extensions.Kestrel/HlsEvictionService.cs create mode 100644 tests/Spangle.Core.Tests/Transport/HLS/MemoryStorageEvictionTests.cs diff --git a/src/Spangle.Core/PublishSessionRegistry.cs b/src/Spangle.Core/PublishSessionRegistry.cs index 2ec35fa..7d13edb 100644 --- a/src/Spangle.Core/PublishSessionRegistry.cs +++ b/src/Spangle.Core/PublishSessionRegistry.cs @@ -168,6 +168,9 @@ internal void Release(string streamName, string sessionId) } } + /// True while a publish session owns the given (sanitized) stream key. + public bool IsLive(string streamKey) => _sessions.ContainsKey(streamKey); + /// /// Snapshots every live publish session for monitoring. Codec and byte counters /// come from the session's receiver context when it is still alive. diff --git a/src/Spangle.Core/Transport/HLS/HLSStorage.cs b/src/Spangle.Core/Transport/HLS/HLSStorage.cs index ea2abb8..08e3b86 100644 --- a/src/Spangle.Core/Transport/HLS/HLSStorage.cs +++ b/src/Spangle.Core/Transport/HLS/HLSStorage.cs @@ -73,6 +73,21 @@ internal LiveBlobReader(Func read) } } +/// +/// Optional capability of a storage backend: dropping streams nobody writes anymore. +/// Only the memory backend implements this — file storage is an archive by design +/// and is never cleaned. +/// +public interface IEvictingHLSStorage +{ + /// + /// Frees every stream without a write for at least , + /// except those claims a publisher still owns. + /// Returns the number of streams evicted. + /// + int EvictIdleStreams(TimeSpan idleFor, Func isLive); +} + /// The storage area of a single stream (one playlist and its media blobs). public interface IHLSStreamStorage { @@ -94,9 +109,10 @@ public interface IHLSStreamStorage /// /// Keeps the live window in process memory. Blobs the playlist trims are freed, so /// a stream holds roughly its sliding window (a few MB); after the stream ends its -/// final window stays servable until the same stream key publishes again. +/// final window stays servable until it is evicted (see ) +/// or the same stream key publishes again. /// -public sealed class MemoryHLSStorage : IHLSStorage +public sealed class MemoryHLSStorage : IHLSStorage, IEvictingHLSStorage { private readonly ConcurrentDictionary _streams = new(StringComparer.Ordinal); @@ -106,6 +122,29 @@ public IHLSStreamStorage GetStream(string streamKey) => public bool TryGetStream(string streamKey, out IHLSStreamStorage stream) => _streams.TryGetValue(streamKey, out stream!); + public int EvictIdleStreams(TimeSpan idleFor, Func isLive) + { + ArgumentNullException.ThrowIfNull(isLive); + var evicted = 0; + long now = Environment.TickCount64; + foreach ((string key, IHLSStreamStorage stream) in _streams) + { + if (stream is not StreamStorage storage + || now - storage.LastWriteTicks < idleFor.TotalMilliseconds + || isLive(key)) + { + continue; + } + // value-checked removal: a concurrent re-publish replaced the entry, leave it be + if (_streams.TryRemove(new KeyValuePair(key, stream))) + { + storage.Drop(); + evicted++; + } + } + return evicted; + } + public override string ToString() => "memory"; private sealed class StreamStorage : IHLSStreamStorage, ILiveBlobStreamStorage @@ -113,8 +152,29 @@ private sealed class StreamStorage : IHLSStreamStorage, ILiveBlobStreamStorage private readonly ConcurrentDictionary _blobs = new(StringComparer.Ordinal); private readonly ConcurrentDictionary _growing = new(StringComparer.Ordinal); private volatile string? _playlist; + private long _lastWriteTicks = Environment.TickCount64; + + internal long LastWriteTicks => Volatile.Read(ref _lastWriteTicks); - public void WriteBlob(string name, ReadOnlySpan content) => _blobs[name] = content.ToArray(); + private void Touch() => Volatile.Write(ref _lastWriteTicks, Environment.TickCount64); + + /// Frees everything after eviction; blocked LL-DASH readers are released. + internal void Drop() + { + foreach (GrowingBlob blob in _growing.Values) + { + blob.Complete(); + } + _growing.Clear(); + _blobs.Clear(); + _playlist = null; + } + + public void WriteBlob(string name, ReadOnlySpan content) + { + Touch(); + _blobs[name] = content.ToArray(); + } public void DeleteBlob(string name) { @@ -125,7 +185,11 @@ public void DeleteBlob(string name) } } - public void PublishPlaylist(string text) => _playlist = text; + public void PublishPlaylist(string text) + { + Touch(); + _playlist = text; + } public bool TryReadBlob(string name, out ReadOnlyMemory content) { @@ -142,8 +206,11 @@ public bool TryReadBlob(string name, out ReadOnlyMemory content) // ---- growing blobs (LL-DASH) ---- - public void AppendBlob(string name, ReadOnlySpan content) => + public void AppendBlob(string name, ReadOnlySpan content) + { + Touch(); _growing.GetOrAdd(name, static _ => new GrowingBlob()).Append(content); + } public void CompleteBlob(string name) { diff --git a/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs b/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs index 7fcd299..ba8029f 100644 --- a/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs +++ b/src/Spangle.Extensions.Kestrel/DependencyInjection/SpangleServiceCollectionExtensions.cs @@ -60,6 +60,7 @@ public static void AddSpangle(this IServiceCollection services) : new MemoryHLSStorage(); }); services.AddHostedService(); + services.AddHostedService(); // Management surface: log capture, delivery counters, and the stats join services.AddSingleton(); diff --git a/src/Spangle.Extensions.Kestrel/HlsEvictionService.cs b/src/Spangle.Extensions.Kestrel/HlsEvictionService.cs new file mode 100644 index 0000000..bf01376 --- /dev/null +++ b/src/Spangle.Extensions.Kestrel/HlsEvictionService.cs @@ -0,0 +1,39 @@ +using Microsoft.Extensions.Options; +using Spangle.Transport.HLS; +using ZLogger; + +namespace Spangle.Extensions.Kestrel; + +/// +/// Frees ended streams from evicting storage backends (memory) once they have been +/// idle for Hls.EndedStreamTtlSeconds. Without this, every distinct stream key +/// ever published would hold its final window in memory for the process lifetime. +/// +public sealed class HlsEvictionService( + IHLSStorage storage, + PublishSessionRegistry sessions, + IOptions options, + ILogger logger) : BackgroundService +{ + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + int ttlSeconds = options.Value.Hls.EndedStreamTtlSeconds; + if (ttlSeconds <= 0 || storage is not IEvictingHLSStorage evicting) + { + return; // eviction disabled, or the backend (file archive) never evicts + } + + var ttl = TimeSpan.FromSeconds(ttlSeconds); + // sweeping a few times per TTL keeps the overshoot small without busy-looping + var period = TimeSpan.FromSeconds(Math.Clamp(ttlSeconds / 4.0, 5.0, 60.0)); + using var timer = new PeriodicTimer(period); + while (await timer.WaitForNextTickAsync(stoppingToken).ConfigureAwait(false)) + { + int evicted = evicting.EvictIdleStreams(ttl, sessions.IsLive); + if (evicted > 0) + { + logger.ZLogInformation($"Freed {evicted} ended stream(s) idle for {ttl}"); + } + } + } +} diff --git a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs index 0c98d85..8bb171a 100644 --- a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs +++ b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs @@ -183,6 +183,14 @@ public class HlsOptions : MediaProtocolOptions /// Target duration of LL-HLS partial segments in seconds [Range(0.1, 5.0)] public double PartTargetDuration { get; set; } = 0.5; + + /// + /// How long an ended stream's final window stays servable from memory storage, + /// in seconds, before it is freed. 0 keeps every ended stream until the same + /// key publishes again — memory then grows with the number of distinct stream + /// keys ever published. File storage is an archive and is never cleaned. + /// + [Range(0, 604_800)] public int EndedStreamTtlSeconds { get; set; } = 300; } public abstract class MediaProtocolOptions diff --git a/src/Spangle.MediaServer/spanglesettings.yaml b/src/Spangle.MediaServer/spanglesettings.yaml index cafcb68..e21479e 100644 --- a/src/Spangle.MediaServer/spanglesettings.yaml +++ b/src/Spangle.MediaServer/spanglesettings.yaml @@ -48,3 +48,5 @@ Spangle: # TsPassthrough: true # SRT→TS-HLS re-segments the source TS as-is (half the # # container work); set false to force demux+remux, e.g. # # when MediaFrame spinner plugins must run on SRT sessions + # EndedStreamTtlSeconds: 300 # memory storage frees an ended stream's final window + # # after this long (0 = keep until the key publishes again) diff --git a/tests/Spangle.Core.Tests/Transport/HLS/MemoryStorageEvictionTests.cs b/tests/Spangle.Core.Tests/Transport/HLS/MemoryStorageEvictionTests.cs new file mode 100644 index 0000000..44ff9d4 --- /dev/null +++ b/tests/Spangle.Core.Tests/Transport/HLS/MemoryStorageEvictionTests.cs @@ -0,0 +1,73 @@ +using Spangle.Transport.HLS; + +namespace Spangle.Tests.Transport.HLS; + +public class MemoryStorageEvictionTests +{ + [Fact] + public void IdleEndedStreamIsEvicted() + { + var storage = new MemoryHLSStorage(); + IHLSStreamStorage stream = storage.GetStream("ended"); + stream.WriteBlob("seg0.ts", [1, 2, 3]); + stream.PublishPlaylist("#EXTM3U"); + + int evicted = storage.EvictIdleStreams(TimeSpan.Zero, static _ => false); + + evicted.Should().Be(1); + storage.TryGetStream("ended", out _).Should().BeFalse("the final window must be freed"); + } + + [Fact] + public void LiveStreamSurvivesEvenWhenIdle() + { + var storage = new MemoryHLSStorage(); + storage.GetStream("live").WriteBlob("seg0.ts", [1]); + + int evicted = storage.EvictIdleStreams(TimeSpan.Zero, static key => key == "live"); + + evicted.Should().Be(0); + storage.TryGetStream("live", out _).Should().BeTrue("a publisher still owns the key"); + } + + [Fact] + public void RecentWriteDefersEviction() + { + var storage = new MemoryHLSStorage(); + storage.GetStream("fresh").WriteBlob("seg0.ts", [1]); + + int evicted = storage.EvictIdleStreams(TimeSpan.FromHours(1), static _ => false); + + evicted.Should().Be(0); + storage.TryGetStream("fresh", out _).Should().BeTrue("the stream wrote within the TTL"); + } + + [Fact] + public async Task EvictionReleasesBlockedLiveBlobReaders() + { + var storage = new MemoryHLSStorage(); + IHLSStreamStorage stream = storage.GetStream("lldash"); + var live = (ILiveBlobStreamStorage)stream; + live.AppendBlob("seg0.m4s", [1, 2]); + live.TryOpenLiveBlob("seg0.m4s", out LiveBlobReader reader).Should().BeTrue(); + (await reader.ReadNextAsync(CancellationToken.None)).Should().NotBeNull(); + + storage.EvictIdleStreams(TimeSpan.Zero, static _ => false).Should().Be(1); + + // the reader must terminate rather than hang on a freed blob + using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); + (await reader.ReadNextAsync(cts.Token)).Should().BeNull(); + } + + [Fact] + public void RepublishAfterEvictionStartsClean() + { + var storage = new MemoryHLSStorage(); + storage.GetStream("re").WriteBlob("seg0.ts", [1]); + storage.EvictIdleStreams(TimeSpan.Zero, static _ => false); + + IHLSStreamStorage fresh = storage.GetStream("re"); + fresh.TryReadBlob("seg0.ts", out _).Should().BeFalse("eviction freed the old window"); + fresh.Playlist.Should().BeNull(); + } +} From f941b2829ac6b3ab62ad6fdeed1b778d0b859e2f Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:14:58 +0000 Subject: [PATCH 5/8] chore: retire the last TODO - Rtmp.Bandwidth comes from configuration The window-acknowledgement/peer-bandwidth value announced during the RTMP connect sequence was hardcoded at 1.5 Mbps with a TODO to make it configurable; it now binds from Rtmp.Bandwidth (same default). Also gives the NALU-length-size guard in NALFileFormat a message naming the actual offending value instead of a bare InvalidDataException. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- src/Spangle.Core/Codecs/NALFileFormat.cs | 3 ++- src/Spangle.Core/Transport/Rtmp/RtmpReceiverContext.cs | 5 ++++- src/Spangle.Extensions.Kestrel/RtmpConnectionHandler.cs | 1 + src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs | 6 ++++++ 4 files changed, 13 insertions(+), 2 deletions(-) diff --git a/src/Spangle.Core/Codecs/NALFileFormat.cs b/src/Spangle.Core/Codecs/NALFileFormat.cs index c8acf4c..3a19c69 100644 --- a/src/Spangle.Core/Codecs/NALFileFormat.cs +++ b/src/Spangle.Core/Codecs/NALFileFormat.cs @@ -33,7 +33,8 @@ public static uint ReadNALULength(ref ReadOnlySequence buff, int lengthSiz 1 => BufferMarshal.As(lenBuff), 2 => BufferMarshal.As(lenBuff).HostValue, 4 => BufferMarshal.As(lenBuff).HostValue, - _ => throw new InvalidDataException(), + _ => throw new InvalidDataException( + $"NALU length size must be 1, 2 or 4 bytes; the decoder configuration declared {lengthSize}"), }; buff = buff.Slice(lenBuff.End); diff --git a/src/Spangle.Core/Transport/Rtmp/RtmpReceiverContext.cs b/src/Spangle.Core/Transport/Rtmp/RtmpReceiverContext.cs index ece3c4a..2ffe149 100644 --- a/src/Spangle.Core/Transport/Rtmp/RtmpReceiverContext.cs +++ b/src/Spangle.Core/Transport/Rtmp/RtmpReceiverContext.cs @@ -32,7 +32,10 @@ public sealed class RtmpReceiverContext : ReceiverContextBase + /// Announced to the peer as WindowAcknowledgementSize / SetPeerBandwidth during + /// the connect sequence. Hosts override this from configuration (Rtmp.Bandwidth). + /// public uint Bandwidth = 1500000; /// diff --git a/src/Spangle.Extensions.Kestrel/RtmpConnectionHandler.cs b/src/Spangle.Extensions.Kestrel/RtmpConnectionHandler.cs index b3b10c7..b1c3390 100644 --- a/src/Spangle.Extensions.Kestrel/RtmpConnectionHandler.cs +++ b/src/Spangle.Extensions.Kestrel/RtmpConnectionHandler.cs @@ -26,6 +26,7 @@ public override async Task OnConnectedAsync(ConnectionContext connection) var ct = cts.Token; var receiver = connection.CreateRtmpReceiverContext(ct); + receiver.Bandwidth = options.Value.Rtmp.Bandwidth; var hlsOptions = options.Value.Hls; var segmentFormat = hlsOptions.SegmentFormat.ToLowerInvariant() switch { diff --git a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs index 8bb171a..e01ac71 100644 --- a/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs +++ b/src/Spangle.Extensions.Kestrel/SpangleMediaServerOptions.cs @@ -106,6 +106,12 @@ public class RtmpOptions : MediaProtocolOptions /// public bool TimedMetadata { get; set; } = true; + /// + /// Announced to publishers as WindowAcknowledgementSize / SetPeerBandwidth (bytes) + /// during the connect sequence. + /// + [Range(1, uint.MaxValue)] public uint Bandwidth { get; set; } = 1_500_000; + /// RTMPS: TLS on the RTMP listener (publishers connect with rtmps://) public TlsOptions Tls { get; set; } = new(); } From 5fd077238cd49e9f37f37cc96f407aa942c974bf Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:20:48 +0000 Subject: [PATCH 6/8] test: direct coverage for the SRT 0x47 resync and the Bearer token gate TryResync widens to internal (InternalsVisibleTo already covers the test assembly) so the alignment-loss paths get exercised without a live SRT socket: garbage skip to a verified boundary, payload 0x47 rejection, keep-tail-for-more-data, and full-garbage drop. TokenGate gets the match/mismatch/scheme/whitespace/long-token matrix. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- .../Transport/SRT/SRTReceiverContext.cs | 2 +- .../Management/TokenGateTests.cs | 57 +++++++++++++ .../Transport/SRT/SrtResyncTests.cs | 83 +++++++++++++++++++ 3 files changed, 141 insertions(+), 1 deletion(-) create mode 100644 tests/Spangle.Core.Tests/Management/TokenGateTests.cs create mode 100644 tests/Spangle.Core.Tests/Transport/SRT/SrtResyncTests.cs diff --git a/src/Spangle.Core/Transport/SRT/SRTReceiverContext.cs b/src/Spangle.Core/Transport/SRT/SRTReceiverContext.cs index 42d6b48..3cea251 100644 --- a/src/Spangle.Core/Transport/SRT/SRTReceiverContext.cs +++ b/src/Spangle.Core/Transport/SRT/SRTReceiverContext.cs @@ -163,7 +163,7 @@ private static byte PeekByte(in ReadOnlySequence buff) /// found; distinguishes "a candidate exists but /// there aren't enough bytes to verify it yet" (keep the tail for the next read). /// - private static bool TryResync(ref ReadOnlySequence buff, out bool needMoreData) + internal static bool TryResync(ref ReadOnlySequence buff, out bool needMoreData) { ReadOnlySequence search = buff.Slice(1); Span next = stackalloc byte[1]; // hoisted: stackalloc in a loop grows the stack per iteration diff --git a/tests/Spangle.Core.Tests/Management/TokenGateTests.cs b/tests/Spangle.Core.Tests/Management/TokenGateTests.cs new file mode 100644 index 0000000..f54d662 --- /dev/null +++ b/tests/Spangle.Core.Tests/Management/TokenGateTests.cs @@ -0,0 +1,57 @@ +using Spangle.Extensions.Kestrel; + +namespace Spangle.Tests.Management; + +public class TokenGateTests +{ + [Fact] + public void CorrectTokenMatches() + { + TokenGate.Matches("Bearer sekrit", "sekrit").Should().BeTrue(); + } + + [Fact] + public void WrongTokenIsRejected() + { + TokenGate.Matches("Bearer wrong", "sekrit").Should().BeFalse(); + } + + [Fact] + public void MissingSchemeIsRejected() + { + TokenGate.Matches("sekrit", "sekrit").Should().BeFalse(); + } + + [Fact] + public void SchemeIsCaseSensitive() + { + TokenGate.Matches("bearer sekrit", "sekrit").Should().BeFalse(); + } + + [Fact] + public void EmptyHeaderIsRejected() + { + TokenGate.Matches("", "sekrit").Should().BeFalse(); + } + + [Fact] + public void SurroundingWhitespaceAroundTheTokenIsTolerated() + { + TokenGate.Matches("Bearer sekrit ", "sekrit").Should().BeTrue(); + } + + [Fact] + public void PrefixOfTheTokenIsRejected() + { + TokenGate.Matches("Bearer sekri", "sekrit").Should().BeFalse(); + TokenGate.Matches("Bearer sekrit2", "sekrit").Should().BeFalse(); + } + + [Fact] + public void LongTokensTakeTheHeapPathAndStillMatch() + { + string token = new('x', 300); // beyond the stackalloc fast path + TokenGate.Matches($"Bearer {token}", token).Should().BeTrue(); + TokenGate.Matches($"Bearer {token}y", token).Should().BeFalse(); + } +} diff --git a/tests/Spangle.Core.Tests/Transport/SRT/SrtResyncTests.cs b/tests/Spangle.Core.Tests/Transport/SRT/SrtResyncTests.cs new file mode 100644 index 0000000..1bdaf9f --- /dev/null +++ b/tests/Spangle.Core.Tests/Transport/SRT/SrtResyncTests.cs @@ -0,0 +1,83 @@ +using System.Buffers; +using Spangle.Containers.M2TS; +using Spangle.Transport.SRT; + +namespace Spangle.Tests.Transport.SRT; + +/// +/// The 0x47-resync loop: after alignment loss, the receiver must latch onto a sync +/// byte only when another sync byte sits exactly one packet later, so a 0x47 inside +/// a payload cannot fool it. +/// +public class SrtResyncTests +{ + private const int PacketSize = 188; // M2TSWriter.PacketSize + + private static byte[] Packet(byte fill) + { + var packet = new byte[PacketSize]; + Array.Fill(packet, fill); + packet[0] = 0x47; + return packet; + } + + [Fact] + public void ResyncSkipsGarbageToTheVerifiedBoundary() + { + byte[] garbage = [0xFF, 0x00, 0x12, 0x34, 0xFF]; + var buff = new ReadOnlySequence([.. garbage, .. Packet(0x00), .. Packet(0x01)]); + + SRTReceiverContext.TryResync(ref buff, out bool needMoreData).Should().BeTrue(); + + needMoreData.Should().BeFalse(); + buff.Length.Should().Be(2 * PacketSize); + buff.FirstSpan[0].Should().Be(0x47); + } + + [Fact] + public void PayloadSyncByteDoesNotFoolTheResync() + { + // a fake 0x47 at index 1: one packet after it lands inside packet bytes + // that are all zero, so the candidate fails verification + byte[] prefix = [0xFF, 0x47, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF]; + var buff = new ReadOnlySequence([.. prefix, .. Packet(0x00), .. Packet(0x00)]); + + SRTReceiverContext.TryResync(ref buff, out bool needMoreData).Should().BeTrue(); + + needMoreData.Should().BeFalse(); + buff.Length.Should().Be(2 * PacketSize, "the fake sync in the garbage must be skipped"); + buff.FirstSpan[0].Should().Be(0x47); + } + + [Fact] + public void UnverifiableCandidateWaitsForMoreData() + { + // a sync byte exists, but fewer than PacketSize+1 bytes follow it + byte[] tail = new byte[20]; + tail[0] = 0xFF; + tail[5] = 0x47; + var buff = new ReadOnlySequence(tail); + + SRTReceiverContext.TryResync(ref buff, out bool needMoreData).Should().BeFalse(); + + needMoreData.Should().BeTrue("the candidate needs the next read to verify"); + buff.FirstSpan[0].Should().Be(0x47, "the tail from the candidate on must be kept"); + } + + [Fact] + public void PureGarbageIsDroppedEntirely() + { + byte[] garbage = new byte[300]; // no 0x47 anywhere + Array.Fill(garbage, (byte)0xAA); + var buff = new ReadOnlySequence(garbage); + + SRTReceiverContext.TryResync(ref buff, out bool needMoreData).Should().BeFalse(); + needMoreData.Should().BeFalse("nothing in the buffer is worth keeping"); + } + + [Fact] + public void PacketSizeConstantMatchesTheMuxer() + { + M2TSWriter.PacketSize.Should().Be(PacketSize); + } +} From 8fedaf520ff19991b8a794a2748b5df338c51f7d Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:20:48 +0000 Subject: [PATCH 7/8] docs: README roadmap catches up with TLS, publish allowlist, ended-stream TTL Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- README.md | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 1ac5293..0a9e838 100644 --- a/README.md +++ b/README.md @@ -45,8 +45,8 @@ $ ffmpeg -re -f lavfi -i testsrc2=size=640x360:rate=30 -f lavfi -i sine=frequenc ``` Configuration lives in `src/Spangle.MediaServer/spanglesettings.yaml` -(ports, segment format TS/fMP4, LL-HLS, SRT passphrase), overridable with -`SMS_`-prefixed environment variables. +(ports, segment format TS/fMP4, LL-HLS, SRT passphrase, TLS per listener, +publish allowlist), overridable with `SMS_`-prefixed environment variables. Family ------ @@ -83,6 +83,18 @@ Roadmap errors (CI builds Release, so the gate holds there too). The deliberate exceptions live in `.editorconfig`, each with a written reason — spec-mirroring enums, numbered protocol-flow files, wire-struct fields +- [x] TLS on every listener that lacked it: HTTPS for HLS/DASH delivery + (`Http.Tls`), HTTPS for the management port so the Bearer token never + travels in cleartext (`Management.Tls`), and RTMPS (`Rtmp.Tls`) — + PKCS#12 or PEM cert/key, plaintext stays the default +- [x] Publish allowlist out of the box: `Publish.AllowedStreamNames` switches + the built-in authorizer from allow-all to an exact-match allowlist + (a custom `IPublishAuthorizer` in DI still overrides everything); the + metadata-injection endpoint on the public port takes an optional Bearer + token (`Http.MetadataInjectionToken`) +- [x] Ended streams no longer hold memory forever: memory storage frees a + stream's final window after `Hls.EndedStreamTtlSeconds` (default 300; + 0 restores keep-until-republish), skipping keys a publisher still owns ### Mid term — features From a472f89f0ad7452203c32a3cbac7dd5d9abace49 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 12 Jul 2026 02:24:03 +0000 Subject: [PATCH 8/8] test: the CMAF packager and the TS-mode segmenter get dedicated suites CmafPackager was only exercised incidentally; it now has a structural ISO-BMFF box parser and 9 tests: init-segment box trees for H.264+AAC / video-only (avc1 with verbatim avcC, mp4a with the ASC inside esds, mdhd timescales), the unmapped-codec rejection, and fragment structure - styp/ moof/mdat order, trun data offsets into the shared mdat, signed composition offsets, sequence numbers and tfdt base times advancing across fragments. The TS-mode HLSSegmenter (the passthrough variant was tested, this one was not) gets 5: keyframe-after-target cutting with held PSI tables, window trimming with physical blob deletion, handover export without ENDLIST, resumed continuation with EXT-X-DISCONTINUITY, and broken-input rejection. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01QwUovURiQeAM8ctCu2gjAa --- .../Containers/ISOBMFF/CmafPackagerTests.cs | 319 ++++++++++++++++++ .../Transport/HLS/HLSSegmenterTests.cs | 202 +++++++++++ 2 files changed, 521 insertions(+) create mode 100644 tests/Spangle.Core.Tests/Containers/ISOBMFF/CmafPackagerTests.cs create mode 100644 tests/Spangle.Core.Tests/Transport/HLS/HLSSegmenterTests.cs diff --git a/tests/Spangle.Core.Tests/Containers/ISOBMFF/CmafPackagerTests.cs b/tests/Spangle.Core.Tests/Containers/ISOBMFF/CmafPackagerTests.cs new file mode 100644 index 0000000..4c0aa83 --- /dev/null +++ b/tests/Spangle.Core.Tests/Containers/ISOBMFF/CmafPackagerTests.cs @@ -0,0 +1,319 @@ +using System.Buffers.Binary; +using System.Text; +using Spangle.Containers.ISOBMFF; + +namespace Spangle.Tests.Containers.ISOBMFF; + +/// +/// CmafPackager builds fragmented MP4: one init segment (ftyp+moov) and media +/// segments (styp+moof+mdat). These tests parse the emitted box structure back +/// out of the raw bytes and verify layout, sample entries, trun bookkeeping and +/// the moof-relative data offsets. +/// +public class CmafPackagerTests +{ + private static readonly byte[] s_avcConfig = [0x01, 0x64, 0x00, 0x1F, 0xFF, 0xE1, 0x00, 0x04]; + private static readonly byte[] s_aacConfig = [0x12, 0x10]; // AAC-LC, 44.1 kHz, stereo + + private static CmafPackager CreateH264AacPackager() => new( + new CmafVideoTrack { Codec = VideoCodec.H264, ConfigRecord = s_avcConfig, Width = 640, Height = 360 }, + new CmafAudioTrack { Codec = AudioCodec.AAC, Config = s_aacConfig, SampleRate = 44100, ChannelCount = 2 }); + + // ---- init segment ---- + + [Fact] + public void InitSegmentForH264AndAacHasExpectedBoxTree() + { + byte[] init = CreateH264AacPackager().BuildInitSegment(); + + List top = ParseBoxes(init, 0, init.Length); + top.Select(b => b.Type).Should().Equal("ftyp", "moov"); + + // ftyp: major brand iso6 with the CMAF structural brand + Box ftyp = top[0]; + FourCcAt(init, ftyp.PayloadStart).Should().Be("iso6"); + Encoding.ASCII.GetString(init, ftyp.PayloadStart, ftyp.PayloadEnd - ftyp.PayloadStart) + .Should().Contain("cmfc"); + + List moov = ParseBoxes(init, top[1].PayloadStart, top[1].PayloadEnd); + moov.Select(b => b.Type).Should().Equal("mvhd", "trak", "trak", "mvex"); + + // mvex declares both tracks, in track-id order + List mvex = ParseBoxes(init, moov[3].PayloadStart, moov[3].PayloadEnd); + mvex.Select(b => b.Type).Should().Equal("trex", "trex"); + ReadUInt32(init, mvex[0].PayloadStart + 4).Should().Be(CmafPackager.VideoTrackId); + ReadUInt32(init, mvex[1].PayloadStart + 4).Should().Be(CmafPackager.AudioTrackId); + } + + [Fact] + public void InitSegmentCarriesAvc1SampleEntryWithVerbatimAvcC() + { + byte[] init = CreateH264AacPackager().BuildInitSegment(); + Box videoTrak = TrakBoxes(init)[0]; + + Box mdhd = Descend(init, videoTrak, "mdia")[0]; + mdhd.Type.Should().Be("mdhd"); + ReadUInt32(init, mdhd.PayloadStart + 12).Should().Be(90000u, "video uses the 90 kHz PTS timescale"); + + Box stsd = StsdOf(init, videoTrak); + ReadUInt32(init, stsd.PayloadStart + 4).Should().Be(1u, "one sample entry"); + + List entries = ParseBoxes(init, stsd.PayloadStart + 8, stsd.PayloadEnd); + entries.Should().ContainSingle().Which.Type.Should().Be("avc1"); + Box avc1 = entries[0]; + + // VisualSampleEntry: width/height after 6 reserved + 2 dri + 16 pre_defined/reserved + ReadUInt16(init, avc1.PayloadStart + 24).Should().Be(640); + ReadUInt16(init, avc1.PayloadStart + 26).Should().Be(360); + + // the codec config box follows the 78-byte fixed VisualSampleEntry fields + List configBoxes = ParseBoxes(init, avc1.PayloadStart + 78, avc1.PayloadEnd); + configBoxes.Should().ContainSingle().Which.Type.Should().Be("avcC"); + init.AsSpan(configBoxes[0].PayloadStart, configBoxes[0].PayloadEnd - configBoxes[0].PayloadStart) + .ToArray().Should().Equal(s_avcConfig, "the avcC record is carried verbatim"); + } + + [Fact] + public void InitSegmentCarriesMp4aSampleEntryWithEsdsHoldingTheAsc() + { + byte[] init = CreateH264AacPackager().BuildInitSegment(); + Box audioTrak = TrakBoxes(init)[1]; + + Box mdhd = Descend(init, audioTrak, "mdia")[0]; + ReadUInt32(init, mdhd.PayloadStart + 12).Should().Be(44100u, "the audio timescale is the sample rate"); + + Box stsd = StsdOf(init, audioTrak); + List entries = ParseBoxes(init, stsd.PayloadStart + 8, stsd.PayloadEnd); + entries.Should().ContainSingle().Which.Type.Should().Be("mp4a"); + Box mp4a = entries[0]; + + // AudioSampleEntry: channelcount after 6 reserved + 2 dri + 8 reserved + ReadUInt16(init, mp4a.PayloadStart + 16).Should().Be(2); + ReadUInt32(init, mp4a.PayloadStart + 24).Should().Be(44100u << 16, "samplerate is 16.16 fixed"); + + // esds follows the 28-byte fixed AudioSampleEntry fields + List esds = ParseBoxes(init, mp4a.PayloadStart + 28, mp4a.PayloadEnd); + esds.Should().ContainSingle().Which.Type.Should().Be("esds"); + + // DecSpecificInfo (tag 0x05) carries the AudioSpecificConfig verbatim + ReadOnlySpan dsi = [0x05, (byte)s_aacConfig.Length, .. s_aacConfig]; + init.AsSpan(esds[0].PayloadStart, esds[0].PayloadEnd - esds[0].PayloadStart) + .IndexOf(dsi).Should().BeGreaterThan(0); + } + + [Fact] + public void VideoOnlyInitSegmentHasOneTrackAndOneTrex() + { + var packager = new CmafPackager( + new CmafVideoTrack { Codec = VideoCodec.H264, ConfigRecord = s_avcConfig, Width = 640, Height = 360 }, + audio: null); + byte[] init = packager.BuildInitSegment(); + + List top = ParseBoxes(init, 0, init.Length); + List moov = ParseBoxes(init, top[1].PayloadStart, top[1].PayloadEnd); + moov.Select(b => b.Type).Should().Equal("mvhd", "trak", "mvex"); + List mvex = ParseBoxes(init, moov[2].PayloadStart, moov[2].PayloadEnd); + mvex.Should().ContainSingle(); + ReadUInt32(init, mvex[0].PayloadStart + 4).Should().Be(CmafPackager.VideoTrackId); + } + + [Fact] + public void UnmappedVideoCodecThrowsNotSupported() + { + var packager = new CmafPackager( + new CmafVideoTrack { Codec = VideoCodec.VP9, ConfigRecord = [], Width = 640, Height = 360 }, + audio: null); + + Action act = () => packager.BuildInitSegment(); + act.Should().Throw().WithMessage("*VP9*"); + } + + // ---- media segments ---- + + [Fact] + public void FragmentEmitsStypMoofMdatInOrderWithCorrectTruns() + { + CmafPackager packager = CreateH264AacPackager(); + + CmafSample[] videoSamples = + [ + MakeSample([0xA0, 0xA1, 0xA2, 0xA3], duration: 3000, compositionOffset: 3000, isSync: true), + MakeSample([0xB0, 0xB1], duration: 3000, compositionOffset: -1500, isSync: false), + ]; + CmafSample[] audioSamples = [MakeSample([0xC0, 0xC1, 0xC2], duration: 1024, compositionOffset: 0, isSync: true)]; + + using var stream = new MemoryStream(); + packager.BuildFragment(90_000, videoSamples, 44_100, audioSamples, stream); + byte[] fragment = stream.ToArray(); + + List top = ParseBoxes(fragment, 0, fragment.Length); + top.Select(b => b.Type).Should().Equal("styp", "moof", "mdat"); + Box moof = top[1]; + Box mdat = top[2]; + + List moofChildren = ParseBoxes(fragment, moof.PayloadStart, moof.PayloadEnd); + moofChildren.Select(b => b.Type).Should().Equal("mfhd", "traf", "traf"); + ReadUInt32(fragment, moofChildren[0].PayloadStart + 4).Should().Be(1u, "the first fragment is sequence 1"); + + // video traf: tfhd/tfdt/trun with per-sample fields and moof-relative data offset + List videoTraf = ParseBoxes(fragment, moofChildren[1].PayloadStart, moofChildren[1].PayloadEnd); + videoTraf.Select(b => b.Type).Should().Equal("tfhd", "tfdt", "trun"); + ReadUInt32(fragment, videoTraf[0].PayloadStart + 4).Should().Be(CmafPackager.VideoTrackId); + ReadUInt64(fragment, videoTraf[1].PayloadStart + 4).Should().Be(90_000u, "tfdt carries the video base time"); + + int trun = videoTraf[2].PayloadStart; + ReadUInt32(fragment, trun).Should().Be(0x01_000F01u, "trun v1 with data-offset|duration|size|flags|cts"); + ReadUInt32(fragment, trun + 4).Should().Be(2u, "sample count"); + int moofStart = moof.PayloadStart - 8; + ReadUInt32(fragment, trun + 8).Should().Be((uint)(mdat.PayloadStart - moofStart), + "data_offset points at the mdat payload, relative to the moof start"); + ReadUInt32(fragment, trun + 12).Should().Be(3000u); + ReadUInt32(fragment, trun + 16).Should().Be(4u, "size of the first sample"); + ReadUInt32(fragment, trun + 20).Should().Be(0x02000000u, "keyframe sample flags"); + ReadUInt32(fragment, trun + 24).Should().Be(3000u, "composition offset of the first sample"); + ReadUInt32(fragment, trun + 32).Should().Be(2u, "size of the second sample"); + ReadUInt32(fragment, trun + 36).Should().Be(0x01010000u, "non-sync sample flags"); + BinaryPrimitives.ReadInt32BigEndian(fragment.AsSpan(trun + 40)).Should().Be(-1500, + "trun v1 composition offsets are signed"); + + // audio traf: its data offset skips the video bytes + List audioTraf = ParseBoxes(fragment, moofChildren[2].PayloadStart, moofChildren[2].PayloadEnd); + audioTraf.Select(b => b.Type).Should().Equal("tfhd", "tfdt", "trun"); + ReadUInt32(fragment, audioTraf[0].PayloadStart + 4).Should().Be(CmafPackager.AudioTrackId); + ReadUInt64(fragment, audioTraf[1].PayloadStart + 4).Should().Be(44_100u); + int audioTrun = audioTraf[2].PayloadStart; + ReadUInt32(fragment, audioTrun).Should().Be(0x00_000701u, "trun v0 with data-offset|duration|size|flags"); + ReadUInt32(fragment, audioTrun + 8).Should().Be((uint)(mdat.PayloadStart - moofStart + 6), + "the audio payload starts after the 4+2 video bytes"); + + // mdat holds video samples then audio samples, verbatim + fragment.AsSpan(mdat.PayloadStart, mdat.PayloadEnd - mdat.PayloadStart).ToArray() + .Should().Equal([0xA0, 0xA1, 0xA2, 0xA3, 0xB0, 0xB1, 0xC0, 0xC1, 0xC2]); + } + + [Fact] + public void SequenceNumbersIncrementAndBaseTimesAdvanceAcrossFragments() + { + CmafPackager packager = CreateH264AacPackager(); + CmafSample[] video = [MakeSample([0x01], duration: 3000, compositionOffset: 0, isSync: true)]; + CmafSample[] audio = [MakeSample([0x02], duration: 1024, compositionOffset: 0, isSync: true)]; + + (uint Sequence, ulong VideoTfdt, ulong AudioTfdt) BuildAndInspect(ulong videoBase, ulong audioBase) + { + using var stream = new MemoryStream(); + packager.BuildFragment(videoBase, video, audioBase, audio, stream); + byte[] fragment = stream.ToArray(); + + List top = ParseBoxes(fragment, 0, fragment.Length); + List moof = ParseBoxes(fragment, top[1].PayloadStart, top[1].PayloadEnd); + List videoTraf = ParseBoxes(fragment, moof[1].PayloadStart, moof[1].PayloadEnd); + List audioTraf = ParseBoxes(fragment, moof[2].PayloadStart, moof[2].PayloadEnd); + return (ReadUInt32(fragment, moof[0].PayloadStart + 4), + ReadUInt64(fragment, videoTraf[1].PayloadStart + 4), + ReadUInt64(fragment, audioTraf[1].PayloadStart + 4)); + } + + BuildAndInspect(0, 0).Should().Be((1u, 0ul, 0ul)); + BuildAndInspect(180_000, 88_200).Should().Be((2u, 180_000ul, 88_200ul)); + BuildAndInspect(360_000, 176_400).Should().Be((3u, 360_000ul, 176_400ul)); + } + + [Fact] + public void AudioTrafIsOmittedWhenNoAudioSamplesArrive() + { + CmafPackager packager = CreateH264AacPackager(); + + using var stream = new MemoryStream(); + packager.BuildFragment(0, [MakeSample([0x01], 3000, 0, isSync: true)], 0, [], stream); + byte[] fragment = stream.ToArray(); + + List top = ParseBoxes(fragment, 0, fragment.Length); + List moof = ParseBoxes(fragment, top[1].PayloadStart, top[1].PayloadEnd); + moof.Select(b => b.Type).Should().Equal("mfhd", "traf"); + List traf = ParseBoxes(fragment, moof[1].PayloadStart, moof[1].PayloadEnd); + ReadUInt32(fragment, traf[0].PayloadStart + 4).Should().Be(CmafPackager.VideoTrackId); + } + + [Fact] + public void AudioOnlyFragmentHasASingleAudioTraf() + { + var packager = new CmafPackager(video: null, + new CmafAudioTrack { Codec = AudioCodec.AAC, Config = s_aacConfig, SampleRate = 44100, ChannelCount = 2 }); + + using var stream = new MemoryStream(); + packager.BuildFragment(0, [], 22_050, [MakeSample([0x0A, 0x0B], 1024, 0, isSync: true)], stream); + byte[] fragment = stream.ToArray(); + + List top = ParseBoxes(fragment, 0, fragment.Length); + top.Select(b => b.Type).Should().Equal("styp", "moof", "mdat"); + List moof = ParseBoxes(fragment, top[1].PayloadStart, top[1].PayloadEnd); + moof.Select(b => b.Type).Should().Equal("mfhd", "traf"); + List traf = ParseBoxes(fragment, moof[1].PayloadStart, moof[1].PayloadEnd); + ReadUInt32(fragment, traf[0].PayloadStart + 4).Should().Be(CmafPackager.AudioTrackId); + ReadUInt64(fragment, traf[1].PayloadStart + 4).Should().Be(22_050u); + fragment.AsSpan(top[2].PayloadStart, top[2].PayloadEnd - top[2].PayloadStart).ToArray() + .Should().Equal([0x0A, 0x0B]); + } + + // ======================================================================= + + private readonly record struct Box(string Type, int PayloadStart, int PayloadEnd); + + /// Parses a run of ISO-BMFF boxes; they must tile [start, end) exactly. + private static List ParseBoxes(ReadOnlySpan data, int start, int end) + { + var boxes = new List(); + int pos = start; + while (pos < end) + { + (pos + 8).Should().BeLessThanOrEqualTo(end, "a box header must fit in the remaining range"); + var size = (int)BinaryPrimitives.ReadUInt32BigEndian(data[pos..]); + size.Should().BeGreaterThanOrEqualTo(8, "compact (32-bit, typed) box headers are expected"); + (pos + size).Should().BeLessThanOrEqualTo(end, "a box must not overrun its container"); + boxes.Add(new Box(FourCcAt(data, pos + 4), pos + 8, pos + size)); + pos += size; + } + return boxes; + } + + private static string FourCcAt(ReadOnlySpan data, int offset) => + Encoding.ASCII.GetString(data.Slice(offset, 4)); + + private static uint ReadUInt32(ReadOnlySpan data, int offset) => + BinaryPrimitives.ReadUInt32BigEndian(data[offset..]); + + private static ushort ReadUInt16(ReadOnlySpan data, int offset) => + BinaryPrimitives.ReadUInt16BigEndian(data[offset..]); + + private static ulong ReadUInt64(ReadOnlySpan data, int offset) => + BinaryPrimitives.ReadUInt64BigEndian(data[offset..]); + + private static List TrakBoxes(byte[] init) + { + List top = ParseBoxes(init, 0, init.Length); + return ParseBoxes(init, top[1].PayloadStart, top[1].PayloadEnd) + .Where(b => b.Type == "trak").ToList(); + } + + /// The children of the first child of . + private static List Descend(byte[] data, Box parent, string childType) + { + Box child = ParseBoxes(data, parent.PayloadStart, parent.PayloadEnd).First(b => b.Type == childType); + return ParseBoxes(data, child.PayloadStart, child.PayloadEnd); + } + + private static Box StsdOf(byte[] init, Box trak) + { + List mdia = Descend(init, trak, "mdia"); + List stbl = Descend(init, mdia.First(b => b.Type == "minf"), "stbl"); + return stbl.First(b => b.Type == "stsd"); + } + + private static CmafSample MakeSample(byte[] data, uint duration, int compositionOffset, bool isSync) => new() + { + Data = data, + Duration = duration, + CompositionOffset = compositionOffset, + IsSync = isSync, + }; +} diff --git a/tests/Spangle.Core.Tests/Transport/HLS/HLSSegmenterTests.cs b/tests/Spangle.Core.Tests/Transport/HLS/HLSSegmenterTests.cs new file mode 100644 index 0000000..035deb4 --- /dev/null +++ b/tests/Spangle.Core.Tests/Transport/HLS/HLSSegmenterTests.cs @@ -0,0 +1,202 @@ +using System.Buffers; +using Spangle.Containers.M2TS; +using Spangle.Transport.HLS; + +namespace Spangle.Tests.Transport.HLS; + +/// +/// The TS segmenter consumes our own muxer's output (PAT/PMT written right before +/// every keyframe), cuts at the first keyframe whose PCR is at least the target +/// duration past the segment start, and maintains a sliding-window playlist. +/// +public class HLSSegmenterTests +{ + private static readonly byte[] s_keyAu = [0x00, 0x00, 0x00, 0x01, 0x65, 0x11, 0x22, 0x33]; + private static readonly byte[] s_pAu = [0x00, 0x00, 0x00, 0x01, 0x41, 0x44, 0x55]; + + /// Writes PAT+PMT followed by a keyframe PES, the way the live muxer does. + private static void WriteKeyframeGroup(M2TSWriter muxer, ArrayBufferWriter ts, ulong pts) + { + muxer.WriteProgramTables(ts); + muxer.WritePes(ts, M2TSWriter.PidVideo, M2TSWriter.StreamIdVideo, s_keyAu, + pts, dts: null, randomAccess: true, withPcr: true); + } + + private static void WritePFrame(M2TSWriter muxer, ArrayBufferWriter ts, ulong pts) + { + muxer.WritePes(ts, M2TSWriter.PidVideo, M2TSWriter.StreamIdVideo, s_pAu, + pts, dts: null, randomAccess: false, withPcr: true); + } + + private static void Feed(HLSSegmenter segmenter, ReadOnlySpan ts) + { + for (var i = 0; i < ts.Length; i += M2TSWriter.PacketSize) + { + segmenter.ProcessPacket(ts.Slice(i, M2TSWriter.PacketSize)); + } + } + + [Fact] + public void CutsAtTheFirstKeyframePastTheTargetDuration() + { + IHLSStreamStorage storage = new MemoryHLSStorage().GetStream("test"); + var segmenter = new HLSSegmenter(storage, targetDuration: 2.0); + + // Keyframes every second: the 1s boundary must NOT cut, the 2s one must + var muxer = new M2TSWriter { VideoCodec = VideoCodec.H264 }; + var ts = new ArrayBufferWriter(); + WriteKeyframeGroup(muxer, ts, pts: 0); + WritePFrame(muxer, ts, pts: 45_000); + WriteKeyframeGroup(muxer, ts, pts: 90_000); // 1.0 s: below target, no cut + WriteKeyframeGroup(muxer, ts, pts: 180_000); // 2.0 s: cut here + WritePFrame(muxer, ts, pts: 225_000); + + Feed(segmenter, ts.WrittenSpan); + segmenter.Complete(); + + string playlist = storage.Playlist!; + playlist.Should().Contain("#EXTINF:2.000,\nseg00000.ts", + "the segment closes at the first keyframe at/after the 2s target, not at the 1s keyframe"); + playlist.Should().Contain("#EXTINF:0.500,\nseg00001.ts", "Complete flushes the remainder"); + playlist.Should().Contain("#EXT-X-MEDIA-SEQUENCE:0"); + playlist.Should().EndWith("#EXT-X-ENDLIST\n"); + + // Both segments exist, are whole packets, and start with PAT+PMT then a video packet + foreach (string name in (string[])["seg00000.ts", "seg00001.ts"]) + { + storage.TryReadBlob(name, out ReadOnlyMemory segMemory).Should().BeTrue(); + byte[] seg = segMemory.ToArray(); + (seg.Length % M2TSWriter.PacketSize).Should().Be(0); + PidOf(seg, 0).Should().Be(M2TSWriter.PidPat, $"{name} must start with a PAT"); + PidOf(seg, 1).Should().Be(M2TSWriter.PidPmt); + PidOf(seg, 2).Should().Be(M2TSWriter.PidVideo); + } + + // The tables held while deciding the no-cut at 1s must stay in segment 0: + // it contains all three PATs written before the cut boundary's tables + storage.TryReadBlob("seg00000.ts", out ReadOnlyMemory seg0).Should().BeTrue(); + CountPid(seg0.ToArray(), M2TSWriter.PidPat).Should().Be(2, "the 0s and 1s tables belong to segment 0"); + storage.TryReadBlob("seg00001.ts", out ReadOnlyMemory seg1).Should().BeTrue(); + CountPid(seg1.ToArray(), M2TSWriter.PidPat).Should().Be(1, "the 2s tables lead segment 1"); + } + + [Fact] + public void TrimsTheWindowAndDeletesOldSegmentsFromStorage() + { + IHLSStreamStorage storage = new MemoryHLSStorage().GetStream("test"); + var segmenter = new HLSSegmenter(storage, targetDuration: 2.0, windowSize: 2); + + // Keyframes 2.5s apart: every boundary cuts; 4 groups + a tail -> 4 segments + var muxer = new M2TSWriter { VideoCodec = VideoCodec.H264 }; + var ts = new ArrayBufferWriter(); + for (var i = 0; i < 4; i++) + { + WriteKeyframeGroup(muxer, ts, pts: (ulong)i * 225_000); + } + WritePFrame(muxer, ts, pts: 720_000); // flushed by Complete as the 4th segment + + Feed(segmenter, ts.WrittenSpan); + segmenter.Complete(); + + string playlist = storage.Playlist!; + playlist.Should().Contain("#EXT-X-MEDIA-SEQUENCE:2", "two segments fell out of the window"); + playlist.Should().NotContain("seg00000.ts").And.NotContain("seg00001.ts"); + playlist.Should().Contain("seg00002.ts").And.Contain("seg00003.ts"); + + storage.TryReadBlob("seg00000.ts", out _).Should().BeFalse("trimmed segments are deleted from storage"); + storage.TryReadBlob("seg00001.ts", out _).Should().BeFalse(); + storage.TryReadBlob("seg00002.ts", out _).Should().BeTrue(); + storage.TryReadBlob("seg00003.ts", out _).Should().BeTrue(); + } + + [Fact] + public void ExportHandoverFlushesWithoutEndListAndCarriesTheWindow() + { + IHLSStreamStorage storage = new MemoryHLSStorage().GetStream("test"); + var segmenter = new HLSSegmenter(storage, targetDuration: 2.0); + + var muxer = new M2TSWriter { VideoCodec = VideoCodec.H264 }; + var ts = new ArrayBufferWriter(); + WriteKeyframeGroup(muxer, ts, pts: 0); + WriteKeyframeGroup(muxer, ts, pts: 225_000); // cuts segment 0 at 2.5s + WritePFrame(muxer, ts, pts: 270_000); + + Feed(segmenter, ts.WrittenSpan); + HLSPlaylistHandover handover = segmenter.ExportHandover(); + + handover.Sequence.Should().Be(2, "the remainder was flushed as the second segment"); + handover.Window.Should().HaveCount(2); + handover.Window[0].Name.Should().Be("seg00000.ts"); + handover.Window[0].Duration.Should().BeApproximately(2.5, 0.001); + handover.Window[1].Name.Should().Be("seg00001.ts"); + + storage.Playlist.Should().NotContain("#EXT-X-ENDLIST", "a takeover keeps the playlist live"); + storage.TryReadBlob("seg00001.ts", out _).Should().BeTrue(); + } + + [Fact] + public void ResumedSegmenterContinuesTheSequenceWithADiscontinuity() + { + IHLSStreamStorage storage = new MemoryHLSStorage().GetStream("test"); + var first = new HLSSegmenter(storage, targetDuration: 2.0); + + var muxer = new M2TSWriter { VideoCodec = VideoCodec.H264 }; + var ts = new ArrayBufferWriter(); + WriteKeyframeGroup(muxer, ts, pts: 0); + WriteKeyframeGroup(muxer, ts, pts: 225_000); + Feed(first, ts.WrittenSpan); + HLSPlaylistHandover handover = first.ExportHandover(); + + // the successor session restarts its own timeline at zero + var second = new HLSSegmenter(storage, targetDuration: 2.0, resume: handover); + var muxer2 = new M2TSWriter { VideoCodec = VideoCodec.H264 }; + var ts2 = new ArrayBufferWriter(); + WriteKeyframeGroup(muxer2, ts2, pts: 0); + WriteKeyframeGroup(muxer2, ts2, pts: 225_000); + Feed(second, ts2.WrittenSpan); + second.Complete(); + + string playlist = storage.Playlist!; + playlist.Should().Contain("seg00002.ts", "the media sequence continues after the takeover"); + playlist.Should().Contain("#EXT-X-DISCONTINUITY", "players must expect a timestamp jump"); + playlist.Should().Contain("#EXT-X-ENDLIST"); + } + + [Fact] + public void BrokenPacketStreamIsRejected() + { + IHLSStreamStorage storage = new MemoryHLSStorage().GetStream("test"); + var segmenter = new HLSSegmenter(storage, targetDuration: 2.0); + + var badSync = new byte[M2TSWriter.PacketSize]; + badSync[0] = 0x48; + Action wrongSync = () => segmenter.ProcessPacket(badSync); + wrongSync.Should().Throw(); + + var truncated = new byte[M2TSWriter.PacketSize - 1]; + truncated[0] = 0x47; + Action shortPacket = () => segmenter.ProcessPacket(truncated); + shortPacket.Should().Throw(); + } + + // ======================================================================= + + private static ushort PidOf(ReadOnlySpan ts, int packetIndex) + { + int o = packetIndex * M2TSWriter.PacketSize; + return (ushort)(((ts[o + 1] & 0x1F) << 8) | ts[o + 2]); + } + + private static int CountPid(ReadOnlySpan ts, ushort pid) + { + var count = 0; + for (var i = 0; i < ts.Length / M2TSWriter.PacketSize; i++) + { + if (PidOf(ts, i) == pid) + { + count++; + } + } + return count; + } +}