diff --git a/src/frontend/src/components/live-status.ts b/src/frontend/src/components/live-status.ts index cf0ca1f91..699224716 100644 --- a/src/frontend/src/components/live-status.ts +++ b/src/frontend/src/components/live-status.ts @@ -3,8 +3,9 @@ * * Wires the header `.live-btn` icon and dispatches a typed * `aspire:live-change` CustomEvent on `document` whenever the snapshot - * actually changes. Reconnects on errors with exponential backoff and - * forces a reconnect when the tab becomes visible again. A missing snapshot + * actually changes. Reconnects on errors with jittered exponential backoff, + * releases the stream after a tab stays hidden, and reconnects when the tab + * becomes visible again. A missing snapshot * endpoint disables live updates until a full page reload. */ @@ -37,6 +38,8 @@ const EMPTY: LiveSnapshot = { }; const BACKOFF_MS = [1_000, 2_000, 5_000, 15_000, 30_000]; +// Hidden tabs release their stream after this grace period; visibility reconnects. +const HIDDEN_CLOSE_DELAY_MS = 60_000; let started = false; let unavailable = false; @@ -44,6 +47,7 @@ let current: LiveSnapshot = EMPTY; let source: EventSource | null = null; let backoffIndex = 0; let reconnectTimer: ReturnType | null = null; +let hiddenCloseTimer: ReturnType | null = null; let pipOpen = false; const listeners = new Set<(s: LiveSnapshot) => void>(); @@ -145,7 +149,8 @@ async function seed(): Promise { function scheduleReconnect(): void { if (unavailable || reconnectTimer) return; - const delay = BACKOFF_MS[Math.min(backoffIndex, BACKOFF_MS.length - 1)]; + // Randomize between 50% and 150% so tabs dropped together don't reconnect together. + const delay = BACKOFF_MS[Math.min(backoffIndex, BACKOFF_MS.length - 1)] * (0.5 + Math.random()); backoffIndex = Math.min(backoffIndex + 1, BACKOFF_MS.length - 1); reconnectTimer = setTimeout(() => { reconnectTimer = null; @@ -204,13 +209,39 @@ function closeSource(): void { } } +function clearReconnectTimer(): void { + if (reconnectTimer) { + clearTimeout(reconnectTimer); + reconnectTimer = null; + } +} + +function clearHiddenCloseTimer(): void { + if (hiddenCloseTimer) { + clearTimeout(hiddenCloseTimer); + hiddenCloseTimer = null; + } +} + +function scheduleHiddenClose(): void { + if (hiddenCloseTimer) return; + hiddenCloseTimer = setTimeout(() => { + hiddenCloseTimer = null; + if (document.visibilityState !== 'hidden') return; + clearReconnectTimer(); + closeSource(); + }, HIDDEN_CLOSE_DELAY_MS); +} + function onVisibilityChange(): void { + if (document.visibilityState === 'hidden') { + scheduleHiddenClose(); + return; + } + clearHiddenCloseTimer(); if (document.visibilityState === 'visible' && !source) { backoffIndex = 0; - if (reconnectTimer) { - clearTimeout(reconnectTimer); - reconnectTimer = null; - } + clearReconnectTimer(); connect(); } } @@ -267,6 +298,7 @@ export function init(): void { started = true; void seed(); connect(); + if (document.visibilityState === 'hidden') scheduleHiddenClose(); document.addEventListener('visibilitychange', onVisibilityChange); document.addEventListener('aspire:live-pip-change', onLivePipChange); } diff --git a/src/frontend/tests/unit/live-status.vitest.test.ts b/src/frontend/tests/unit/live-status.vitest.test.ts index fa1ac8236..cf0da49f4 100644 --- a/src/frontend/tests/unit/live-status.vitest.test.ts +++ b/src/frontend/tests/unit/live-status.vitest.test.ts @@ -111,6 +111,8 @@ describe('live-status module', () => { const seed = Promise.withResolvers(); const fetch = vi.fn(() => seed.promise); const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + // Center the reconnect jitter so delays equal the base backoff. + vi.spyOn(Math, 'random').mockReturnValue(0.5); vi.stubGlobal('document', document); vi.stubGlobal('window', { location: { pathname: '/docs/' } }); vi.stubGlobal('EventSource', MockEventSource); @@ -168,6 +170,46 @@ describe('live-status module', () => { expect(sources).toHaveLength(4); expect(warn).not.toHaveBeenCalled(); }); + + it('jitters reconnect delays around the base backoff', async () => { + const { sources } = await setup(); + vi.mocked(Math.random).mockReturnValue(0); + sources[0].onerror?.(); + await vi.advanceTimersByTimeAsync(499); + expect(sources).toHaveLength(1); + await vi.advanceTimersByTimeAsync(1); + expect(sources).toHaveLength(2); + }); + + it('closes the stream in hidden tabs and reconnects when visible', async () => { + const { document, sources } = await setup(); + document.visibilityState = 'hidden'; + document.dispatchEvent(new Event('visibilitychange')); + sources[0].onerror?.(); + await vi.advanceTimersByTimeAsync(59_999); + expect(sources).toHaveLength(2); + expect(sources[1].close).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(sources[1].close).toHaveBeenCalledOnce(); + await vi.advanceTimersByTimeAsync(300_000); + expect(sources).toHaveLength(2); + + document.visibilityState = 'visible'; + document.dispatchEvent(new Event('visibilitychange')); + expect(sources).toHaveLength(3); + }); + + it('keeps the stream when a tab becomes visible within the grace period', async () => { + const { document, sources } = await setup(); + document.visibilityState = 'hidden'; + document.dispatchEvent(new Event('visibilitychange')); + await vi.advanceTimersByTimeAsync(30_000); + document.visibilityState = 'visible'; + document.dispatchEvent(new Event('visibilitychange')); + await vi.advanceTimersByTimeAsync(60_000); + expect(sources).toHaveLength(1); + expect(sources[0].close).not.toHaveBeenCalled(); + }); }); it('delivers the current snapshot to a new subscriber synchronously', () => { diff --git a/src/statichost/StaticHost/Extensions.cs b/src/statichost/StaticHost/Extensions.cs index 2c02c4a8b..c541dd330 100644 --- a/src/statichost/StaticHost/Extensions.cs +++ b/src/statichost/StaticHost/Extensions.cs @@ -1,7 +1,48 @@ +using Microsoft.AspNetCore.HttpOverrides; +using OpenTelemetry.Instrumentation.AspNetCore; + namespace Microsoft.Extensions.Hosting; public static class Extensions { + /// + /// Azure Front Door overwrites X-Azure-ClientIP with the client socket IP, + /// so it stays a single trustworthy entry even when App Service appends its own + /// X-Forwarded-For hop. + /// + public const string FrontDoorClientIpHeader = "X-Azure-ClientIP"; + + public static TBuilder AddForwardedClientIp(this TBuilder builder) where TBuilder : IHostApplicationBuilder + { + builder.Services.Configure(static options => + { + options.ForwardedHeaders |= ForwardedHeaders.XForwardedFor; + options.ForwardedForHeaderName = FrontDoorClientIpHeader; + options.ForwardLimit = 1; + // Front Door and App Service front-end addresses are not fixed. A request that + // bypasses Front Door can spoof this header, so use the value for telemetry and + // diagnostics, not for authorization. + options.KnownIPNetworks.Clear(); + options.KnownProxies.Clear(); + }); + + return builder; + } + + /// + /// Adds forwarded-header processing unless ASPNETCORE_FORWARDEDHEADERS_ENABLED + /// already added it, so a forwarded value is never consumed twice. + /// + public static IApplicationBuilder UseForwardedClientIp(this WebApplication app) + { + if (!app.Configuration.GetValue("FORWARDEDHEADERS_ENABLED")) + { + app.UseForwardedHeaders(); + } + + return app; + } + public static TBuilder AddServiceDefaults(this TBuilder builder) where TBuilder : IHostApplicationBuilder { builder.ConfigureOpenTelemetry(); @@ -25,6 +66,14 @@ private static TBuilder ConfigureOpenTelemetry(this TBuilder builder) options.ConnectionString = TelemetryConstants.AzureMonitorConnectionString; }); + // Long-lived SSE requests would otherwise dominate request duration telemetry + // (App Insights "Server response time") with connection lifetimes, not latency. + builder.Services.Configure(static options => + { + options.Filter = static context => + !context.Request.Path.StartsWithSegments("/api/live/stream", StringComparison.OrdinalIgnoreCase); + }); + builder.Services.AddOpenTelemetry() .WithMetrics(metrics => { diff --git a/src/statichost/StaticHost/Live/LiveEndpoints.cs b/src/statichost/StaticHost/Live/LiveEndpoints.cs index 98a0ac86c..84940362d 100644 --- a/src/statichost/StaticHost/Live/LiveEndpoints.cs +++ b/src/statichost/StaticHost/Live/LiveEndpoints.cs @@ -92,11 +92,17 @@ private static async Task GetSnapshot( internal static async Task StreamSse( HttpContext context, LiveStatusBroadcaster broadcaster, + IOptions options, TimeProvider time, CancellationToken cancellationToken) { await broadcaster.RefreshAsync(cancellationToken).ConfigureAwait(false); + // End every stream within a jittered bound so request durations stay finite + // and clients that reconnect after a deployment don't stay synchronized. + var maxLifetime = TimeSpan.FromSeconds(options.Value.StreamMaxLifetimeSeconds); + var lifetimeDuration = maxLifetime * (1 - (Random.Shared.NextDouble() * 0.2)); + context.Response.StatusCode = StatusCodes.Status200OK; context.Response.Headers.ContentType = "text/event-stream"; context.Response.Headers["X-Accel-Buffering"] = "no"; @@ -109,14 +115,15 @@ internal static async Task StreamSse( var (reader, unsubscribe) = broadcaster.Subscribe(); using var _ = unsubscribe; - using var pending = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + using var lifetime = new CancellationTokenSource(lifetimeDuration, time); + using var pending = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, lifetime.Token); using var heartbeat = new PeriodicTimer(TimeSpan.FromSeconds(15), time); var dataReady = reader.WaitToReadAsync(pending.Token).AsTask(); var heartbeatReady = heartbeat.WaitForNextTickAsync(pending.Token).AsTask(); try { - while (!cancellationToken.IsCancellationRequested) + while (!pending.IsCancellationRequested) { await Task.WhenAny(dataReady, heartbeatReady).ConfigureAwait(false); if (dataReady.IsCompleted) @@ -124,21 +131,22 @@ internal static async Task StreamSse( if (!await dataReady.ConfigureAwait(false)) break; while (reader.TryRead(out var next)) { - await context.Response.Body.WriteAsync(next.Frame, cancellationToken).ConfigureAwait(false); + await context.Response.Body.WriteAsync(next.Frame, pending.Token).ConfigureAwait(false); } - await context.Response.Body.FlushAsync(cancellationToken).ConfigureAwait(false); + await context.Response.Body.FlushAsync(pending.Token).ConfigureAwait(false); dataReady = reader.WaitToReadAsync(pending.Token).AsTask(); } else { if (!await heartbeatReady.ConfigureAwait(false)) break; - await context.Response.WriteAsync(":hb\n\n", cancellationToken).ConfigureAwait(false); - await context.Response.Body.FlushAsync(cancellationToken).ConfigureAwait(false); + await context.Response.WriteAsync(":hb\n\n", pending.Token).ConfigureAwait(false); + await context.Response.Body.FlushAsync(pending.Token).ConfigureAwait(false); heartbeatReady = heartbeat.WaitForNextTickAsync(pending.Token).AsTask(); } } } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { /* client disconnected */ } + catch (OperationCanceledException) when (lifetime.IsCancellationRequested) { /* max lifetime reached; client reconnects */ } finally { await pending.CancelAsync().ConfigureAwait(false); diff --git a/src/statichost/StaticHost/Live/LiveStatusOptions.cs b/src/statichost/StaticHost/Live/LiveStatusOptions.cs index adb9868f1..91e0e57b7 100644 --- a/src/statichost/StaticHost/Live/LiveStatusOptions.cs +++ b/src/statichost/StaticHost/Live/LiveStatusOptions.cs @@ -38,6 +38,15 @@ public sealed class LiveStatusOptions [Range(0, 5000)] public int CoalesceWindowMs { get; set; } = 750; + /// + /// Upper bound for a single /api/live/stream connection. Each stream ends + /// after a random lifetime between 80% and 100% of this value, so long-lived + /// tabs reconnect periodically instead of holding one request open for days + /// and reconnects are spread out. Defaults to 10 minutes. + /// + [Range(60, 24 * 60 * 60)] + public int StreamMaxLifetimeSeconds { get; set; } = 10 * 60; + /// /// When true, expose the dev-only POST /api/live/_dev/set /// endpoint that lets local devs (and Playwright) flip live state without diff --git a/src/statichost/StaticHost/Live/README.md b/src/statichost/StaticHost/Live/README.md index c736b3854..85925609b 100644 --- a/src/statichost/StaticHost/Live/README.md +++ b/src/statichost/StaticHost/Live/README.md @@ -215,6 +215,7 @@ The SSE endpoint reports offline unless a local simulation sets live state. "Live": { "PublicBaseUrl": "https://aspire.dev", "CoalesceWindowMs": 750, + "StreamMaxLifetimeSeconds": 600, "EnableDevEndpoint": false, "DevCommandSecret": "", "Twitch": { @@ -527,7 +528,17 @@ project or `Aspire.Dev.slnx` trigger that job. described above. - Reconciliation timers are the safety net for missed individual webhooks. - SSE heartbeats every 15 s defeat proxy idle-timeouts; the client uses - exponential backoff with a `visibilitychange`-aware reconnect. + jittered exponential backoff (50–150% of each step) with a + `visibilitychange`-aware reconnect. +- Each SSE stream ends after a random 80–100% of `StreamMaxLifetimeSeconds` + (default 10 minutes), and the client reconnects. Without this cap, one + request could stay open for days and inflate request-duration telemetry + until the next restart. +- A tab that stays hidden for 60 s closes its stream; the stream reconnects + when the tab becomes visible and receives the current snapshot. +- `/api/live/stream` is excluded from ASP.NET Core request traces, so App + Insights server response time reflects page and API latency, not stream + lifetimes. - A `404` from the initial `/api/live/` snapshot request closes the event stream and cancels reconnects until a full page reload. Tab visibility changes and Astro navigation do not restart it. This avoids repeated missing-endpoint diff --git a/src/statichost/StaticHost/Program.cs b/src/statichost/StaticHost/Program.cs index 06f736b60..9463687b5 100644 --- a/src/statichost/StaticHost/Program.cs +++ b/src/statichost/StaticHost/Program.cs @@ -3,6 +3,7 @@ var builder = WebApplication.CreateBuilder(args); builder.AddServiceDefaults(); +builder.AddForwardedClientIp(); // Production supplies the vault reference; local Aspire runs intentionally do not. if (builder.Configuration.GetConnectionString("secrets") is not null) @@ -26,6 +27,9 @@ await using var app = builder.Build(); +// Restore the Front Door client IP before telemetry or any middleware reads it. +app.UseForwardedClientIp(); + // Only enable HSTS in production if (!app.Environment.IsDevelopment()) { diff --git a/src/statichost/StaticHost/appsettings.json b/src/statichost/StaticHost/appsettings.json index bd5964e7e..99c6f3634 100644 --- a/src/statichost/StaticHost/appsettings.json +++ b/src/statichost/StaticHost/appsettings.json @@ -9,6 +9,7 @@ "Live": { "PublicBaseUrl": "https://aspire.dev", "CoalesceWindowMs": 750, + "StreamMaxLifetimeSeconds": 600, "EnableDevEndpoint": false, "DevCommandSecret": "", "Twitch": { diff --git a/tests/StaticHost.Tests/Live/LiveStreamTests.cs b/tests/StaticHost.Tests/Live/LiveStreamTests.cs index 3f89150ec..7a2a25b1b 100644 --- a/tests/StaticHost.Tests/Live/LiveStreamTests.cs +++ b/tests/StaticHost.Tests/Live/LiveStreamTests.cs @@ -13,7 +13,8 @@ public async Task StreamSse_SendsHeartbeatsThenUpdatesAndStopsOnDisconnect() context.Response.Body = body; var streaming = LiveStatusEndpointRouteBuilderExtensions.StreamSse( - context, broadcaster, time, cancellation.Token); + context, broadcaster, Options.Create(new LiveStatusOptions { StreamMaxLifetimeSeconds = 24 * 60 * 60 }), + time, cancellation.Token); Assert.StartsWith("event: state\n", await body.ReadWriteAsync()); for (var i = 0; i < 100; i++) @@ -33,6 +34,63 @@ public async Task StreamSse_SendsHeartbeatsThenUpdatesAndStopsOnDisconnect() Assert.False(body.Writes.Reader.TryRead(out _)); } + [Fact] + public async Task StreamSse_EndsWithinJitteredMaxLifetime() + { + var time = new FakeTimeProvider(DateTimeOffset.UnixEpoch); + using var broadcaster = LiveTestHelpers.CreateBroadcaster(timeProvider: time); + using var body = new RecordingStream(); + var context = new DefaultHttpContext(); + context.Response.Body = body; + + var streaming = LiveStatusEndpointRouteBuilderExtensions.StreamSse( + context, broadcaster, Options.Create(new LiveStatusOptions { StreamMaxLifetimeSeconds = 600 }), + time, CancellationToken.None); + Assert.StartsWith("event: state\n", await body.ReadWriteAsync()); + + // The jittered lifetime is never shorter than 80% of the maximum. + for (var elapsed = 15; elapsed < 480; elapsed += 15) + { + time.Advance(TimeSpan.FromSeconds(15)); + Assert.Equal(":hb\n\n", await body.ReadWriteAsync()); + } + Assert.False(streaming.IsCompleted); + + // Advance to the full 600-second maximum. + time.Advance(TimeSpan.FromSeconds(600 - 465)); + await streaming.WaitAsync(TimeSpan.FromSeconds(5)); + } + + [Fact] + public async Task StreamSse_MaxLifetimeInterruptsBlockedWrites() + { + var time = new FakeTimeProvider(DateTimeOffset.UnixEpoch); + using var broadcaster = LiveTestHelpers.CreateBroadcaster(timeProvider: time); + using var body = new StalledStream(); + var context = new DefaultHttpContext(); + context.Response.Body = body; + + var streaming = LiveStatusEndpointRouteBuilderExtensions.StreamSse( + context, broadcaster, Options.Create(new LiveStatusOptions { StreamMaxLifetimeSeconds = 60 }), + time, CancellationToken.None); + await body.Stalled.Task.WaitAsync(TimeSpan.FromSeconds(5)); + Assert.False(streaming.IsCompleted); + + time.Advance(TimeSpan.FromSeconds(60)); + await streaming.WaitAsync(TimeSpan.FromSeconds(5)); + } + + private sealed class StalledStream : MemoryStream + { + public TaskCompletionSource Stalled { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + + public override async ValueTask WriteAsync(ReadOnlyMemory buffer, CancellationToken cancellationToken = default) + { + Stalled.TrySetResult(); + await Task.Delay(Timeout.Infinite, cancellationToken); + } + } + private sealed class RecordingStream : MemoryStream { public Channel Writes { get; } = Channel.CreateUnbounded();