Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 39 additions & 7 deletions src/frontend/src/components/live-status.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/

Expand Down Expand Up @@ -37,13 +38,16 @@ 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;
let current: LiveSnapshot = EMPTY;
let source: EventSource | null = null;
let backoffIndex = 0;
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
let hiddenCloseTimer: ReturnType<typeof setTimeout> | null = null;
let pipOpen = false;
const listeners = new Set<(s: LiveSnapshot) => void>();

Expand Down Expand Up @@ -145,7 +149,8 @@ async function seed(): Promise<void> {

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;
Expand Down Expand Up @@ -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();
}
}
Expand Down Expand Up @@ -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);
}
Expand Down
42 changes: 42 additions & 0 deletions src/frontend/tests/unit/live-status.vitest.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,8 @@ describe('live-status module', () => {
const seed = Promise.withResolvers<Response>();
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);
Expand Down Expand Up @@ -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', () => {
Expand Down
49 changes: 49 additions & 0 deletions src/statichost/StaticHost/Extensions.cs
Original file line number Diff line number Diff line change
@@ -1,7 +1,48 @@
using Microsoft.AspNetCore.HttpOverrides;
using OpenTelemetry.Instrumentation.AspNetCore;

namespace Microsoft.Extensions.Hosting;

public static class Extensions
{
/// <summary>
/// Azure Front Door overwrites <c>X-Azure-ClientIP</c> with the client socket IP,
/// so it stays a single trustworthy entry even when App Service appends its own
/// <c>X-Forwarded-For</c> hop.
/// </summary>
public const string FrontDoorClientIpHeader = "X-Azure-ClientIP";

public static TBuilder AddForwardedClientIp<TBuilder>(this TBuilder builder) where TBuilder : IHostApplicationBuilder
{
builder.Services.Configure<ForwardedHeadersOptions>(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;
}

/// <summary>
/// Adds forwarded-header processing unless ASPNETCORE_FORWARDEDHEADERS_ENABLED
/// already added it, so a forwarded value is never consumed twice.
/// </summary>
public static IApplicationBuilder UseForwardedClientIp(this WebApplication app)
{
if (!app.Configuration.GetValue<bool>("FORWARDEDHEADERS_ENABLED"))
{
app.UseForwardedHeaders();
}

return app;
}

public static TBuilder AddServiceDefaults<TBuilder>(this TBuilder builder) where TBuilder : IHostApplicationBuilder
{
builder.ConfigureOpenTelemetry();
Expand All @@ -25,6 +66,14 @@ private static TBuilder ConfigureOpenTelemetry<TBuilder>(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<AspNetCoreTraceInstrumentationOptions>(static options =>
{
options.Filter = static context =>
!context.Request.Path.StartsWithSegments("/api/live/stream", StringComparison.OrdinalIgnoreCase);
});

builder.Services.AddOpenTelemetry()
.WithMetrics(metrics =>
{
Expand Down
20 changes: 14 additions & 6 deletions src/statichost/StaticHost/Live/LiveEndpoints.cs
Original file line number Diff line number Diff line change
Expand Up @@ -92,11 +92,17 @@ private static async Task<IResult> GetSnapshot(
internal static async Task StreamSse(
HttpContext context,
LiveStatusBroadcaster broadcaster,
IOptions<LiveStatusOptions> 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";
Expand All @@ -109,36 +115,38 @@ 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)
Comment thread
IEvangelist marked this conversation as resolved.
{
await Task.WhenAny(dataReady, heartbeatReady).ConfigureAwait(false);
if (dataReady.IsCompleted)
{
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);
Expand Down
9 changes: 9 additions & 0 deletions src/statichost/StaticHost/Live/LiveStatusOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,15 @@ public sealed class LiveStatusOptions
[Range(0, 5000)]
public int CoalesceWindowMs { get; set; } = 750;

/// <summary>
/// Upper bound for a single <c>/api/live/stream</c> 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.
/// </summary>
[Range(60, 24 * 60 * 60)]
public int StreamMaxLifetimeSeconds { get; set; } = 10 * 60;

/// <summary>
/// When <c>true</c>, expose the dev-only <c>POST /api/live/_dev/set</c>
/// endpoint that lets local devs (and Playwright) flip live state without
Expand Down
13 changes: 12 additions & 1 deletion src/statichost/StaticHost/Live/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down Expand Up @@ -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
Expand Down
4 changes: 4 additions & 0 deletions src/statichost/StaticHost/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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())
{
Expand Down
1 change: 1 addition & 0 deletions src/statichost/StaticHost/appsettings.json
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
"Live": {
"PublicBaseUrl": "https://aspire.dev",
"CoalesceWindowMs": 750,
"StreamMaxLifetimeSeconds": 600,
"EnableDevEndpoint": false,
"DevCommandSecret": "",
"Twitch": {
Expand Down
Loading
Loading