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
351 changes: 208 additions & 143 deletions Services/Bridges/QueueStatusBridge.cs

Large diffs are not rendered by default.

40 changes: 37 additions & 3 deletions Services/Bridges/RealtimeBridge.cs
Original file line number Diff line number Diff line change
Expand Up @@ -11,16 +11,16 @@ namespace Craft.Services;
/// <summary>
/// Static publish surface for the realtime channel, so downstream PowerShell (and Craft's own C#) can
/// push job events without an HTTP round-trip — mirrors <see cref="AppLifecycleBridge"/> /
/// <see cref="SchedulerBridge"/>. Delivery, gating, GUID/size enforcement and storage live in
/// <see cref="SchedulerBridge"/>. Delivery, gating, job-id/size enforcement and storage live in
/// <see cref="RealtimeService"/>.
///
/// PS usage:
/// [Craft.Services.RealtimeBridge]::Publish($userId, $jobId, "start", @{ done = 0; total = 300 })
/// [Craft.Services.RealtimeBridge]::Publish($userId, $jobId, "update", @{ done = 142; total = 300 })
/// [Craft.Services.RealtimeBridge]::Publish($userId, $jobId, "end", @{ done = 300; total = 300 })
///
/// Only <c>userId</c> and <c>jobId</c> (a GUID) are required; everything else is optional. Every call is
/// best-effort and never throws back to the caller.
/// Only <c>userId</c> and <c>jobId</c> (a short token: 1-128 letters, digits, <c>-_.:</c>) are required; everything
/// else is optional. Every call is best-effort and never throws back to the caller.
/// </summary>
public static class RealtimeBridge
{
Expand All @@ -42,6 +42,40 @@ public static void Publish(string userId, string jobId, string mode, object? dat
public static void Publish(string userId, string jobId, string mode, object? data, string? urlHref, string? urlLabel) =>
Publish(userId, jobId, mode, data, urlHref, urlLabel, null, null);

/// <summary>
/// Grant <paramref name="userId"/> the events of <paramref name="jobId"/>, so later
/// <see cref="Notify(string)"/> calls reach them without knowing who they are. Call it from the request
/// that started the job, with that request's principal name.
/// PS usage: [Craft.Services.RealtimeBridge]::Watch($Request.Headers.'x-ms-client-principal-name', $JobId)
/// </summary>
public static void Watch(string userId, string jobId)
{
try { s_service?.Watch(userId, jobId, trackRun: false); }
catch { /* best-effort */ }
}

/// <summary>
/// <see cref="Watch"/>, and have Craft push the orchestrator run status for <paramref name="jobId"/>
/// (a queue id; matched like <see cref="QueueStatusBridge.GetRun"/>) until the run finishes.
/// </summary>
public static void WatchRun(string userId, string jobId)
{
try { s_service?.Watch(userId, jobId, trackRun: true); }
catch { /* best-effort */ }
}

/// <summary>Signal every user granted <paramref name="jobId"/> (mode update, no payload).</summary>
public static void Notify(string jobId) => Notify(jobId, null, null);

public static void Notify(string jobId, string? mode) => Notify(jobId, mode, null);

/// <summary>Publish to every user granted <paramref name="jobId"/> through <see cref="Watch"/> or <see cref="WatchRun"/>.</summary>
public static void Notify(string jobId, string? mode, object? data)
{
try { s_service?.Notify(jobId, mode, data); }
catch { /* best-effort */ }
}

/// <summary>Full form. <paramref name="mode"/> is start | update | end (default update).</summary>
public static void Publish(string userId, string jobId, string? mode, object? data,
string? urlHref, string? urlLabel, int? status, string? message)
Expand Down
5 changes: 3 additions & 2 deletions Services/Configuration/RealtimeSettings.cs
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,9 @@ public class RealtimeSettings
/// <summary>Heartbeat comment interval, seconds, to keep the stream alive through proxies. Default 20.</summary>
public int HeartbeatSeconds { get; set; } = 20;

/// <summary>TTL for a stored entry that never receives an <c>end</c> (crash backstop), minutes. Default 60.</summary>
public int EntryTtlMinutes { get; set; } = 60;
/// <summary>How long a job's last frame (an <c>end</c> included) is kept for reconnect replay, and a grant
/// with no activity is kept, minutes. Default 30.</summary>
public int EntryTtlMinutes { get; set; } = 30;

/// <summary>
/// Resolved enabled state. The <c>CRAFT_REALTIME_ENABLED</c> environment variable (true/1 or false/0)
Expand Down
2 changes: 1 addition & 1 deletion Services/Hosting/CraftHostBuilderExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -194,7 +194,7 @@ public static LogLevel AddCraftLogging(this WebApplicationBuilder builder)
return level;
}

private static readonly string[] second = new[] { "application/json", "text/json", "application/javascript", "text/javascript" };
private static readonly string[] second = new[] { "application/json", "text/json", "application/javascript", "text/javascript", "text/event-stream" };

/// <summary>
/// Resolves the on-the-fly compression level from the <c>CRAFT_API_COMPRESSION_LEVEL</c> env var
Expand Down
46 changes: 41 additions & 5 deletions Services/Hosting/Endpoints/RealtimeEndpoint.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using System.Text.Json;
using Craft.Auth;
using Craft.Configuration;
using Craft.Realtime;
using Microsoft.AspNetCore.Http.Features;

namespace Craft.Hosting.Endpoints;

Expand All @@ -11,6 +12,9 @@ namespace Craft.Hosting.Endpoints;
/// </summary>
public static class RealtimeEndpoint
{
/// <summary>The stream's path. Program.cs routes it through the /api compressor.</summary>
public const string Path = "/.craft/events";

/// <summary>
/// Maps the SSE endpoint on nodes that face a browser, when realtime is switched on.
/// </summary>
Expand Down Expand Up @@ -40,19 +44,20 @@ public static WebApplication MapCraftRealtimeEndpoint(
var realtime = app.Services.GetRequiredService<RealtimeService>();
var heartbeat = TimeSpan.FromSeconds(Math.Max(5, settings.Realtime.HeartbeatSeconds));

app.MapGet("/.craft/events", async (HttpContext ctx) =>
app.MapGet(Path, async (HttpContext ctx) =>
{
// Delivery is identity-gated: a stream is only ever fed this user's own job events.
var userId = ctx.Request.Headers["x-ms-client-principal-name"].ToString();
if (string.IsNullOrEmpty(userId)) { ctx.Response.StatusCode = 401; return; }
var userId = ResolveSignedInUser(ctx);
if (userId is null) { ctx.Response.StatusCode = 401; return; }

var (connId, conn) = realtime.Connect(userId);
if (conn is null) { ctx.Response.StatusCode = 503; return; } // over MaxConnections

ctx.Response.Headers["Content-Type"] = "text/event-stream";
ctx.Response.Headers["Cache-Control"] = "no-cache";
ctx.Response.Headers["X-Accel-Buffering"] = "no"; // stop nginx buffering the stream
ctx.Features.Get<IHttpResponseBodyFeature>()?.DisableBuffering();
// No DisableBuffering: every batch of frames is flushed below, which is all the stream needs, and
// the compressor then encodes a batch as one block instead of flushing after every write.

var ct = ctx.RequestAborted;
try
Expand Down Expand Up @@ -101,4 +106,35 @@ public static WebApplication MapCraftRealtimeEndpoint(
logger.LogInformation("[System] Realtime SSE endpoint: /.craft/events");
return app;
}

/// <summary>
/// The principal name of a signed-in user, read after <see cref="CraftAuthMiddleware"/> has normalised
/// the request. The principal must carry a real role, which rules out anonymous callers and app-only
/// API clients (normalised with no roles). Null when the caller is neither.
/// </summary>
internal static string? ResolveSignedInUser(HttpContext ctx)
{
var name = ctx.Request.Headers["x-ms-client-principal-name"].ToString();
var principal = ctx.Request.Headers["x-ms-client-principal"].ToString();
if (string.IsNullOrEmpty(name) || string.IsNullOrEmpty(principal)) return null;

try
{
using var doc = EasyAuthPrincipal.Decode(principal);
if (doc.RootElement.ValueKind != JsonValueKind.Object ||
!doc.RootElement.TryGetProperty("userRoles", out var roles) ||
roles.ValueKind != JsonValueKind.Array)
return null;

foreach (var role in roles.EnumerateArray())
if (role.ValueKind == JsonValueKind.String && role.GetString() is { Length: > 0 } r &&
!r.Equals("anonymous", StringComparison.OrdinalIgnoreCase))
return name;
}
catch
{
// Unreadable principal: not a signed-in user.
}
return null;
}
}
9 changes: 8 additions & 1 deletion Services/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -307,8 +307,12 @@ void RunInitialization()
});
}

// The realtime stream is API traffic too: same switch and compressor, each flushed frame batch sent at once.
if (apiCompressionEnabled)
app.UseWhen(IsEventsPath, ev => ev.UseResponseCompression());

if (compressionEnabled)
app.UseWhen(ctx => !IsApiPath(ctx), sf => sf.UseResponseCompression());
app.UseWhen(ctx => !IsApiPath(ctx) && !IsEventsPath(ctx), sf => sf.UseResponseCompression());

logger.LogInformation(
"[System] Compression — static: {Static} /api: {Api} level: {Level} (egress accounting: {Egress})",
Expand All @@ -322,6 +326,9 @@ void RunInitialization()
static bool IsApiPath(HttpContext c) =>
c.Request.Path.StartsWithSegments("/api", StringComparison.OrdinalIgnoreCase);

static bool IsEventsPath(HttpContext c) =>
c.Request.Path.Equals(RealtimeEndpoint.Path, StringComparison.OrdinalIgnoreCase);

// Nodes without the Http role do not short-circuit /api or auth paths: the HTTP endpoints simply aren't
// mapped (see the `if (capHttp)` blocks below), so those requests fall through to static file serving
// (a Frontend node can expose /api/me etc. from its own static dir) and finally to MapFallback (which
Expand Down
Loading
Loading