diff --git a/Services/Bridges/QueueStatusBridge.cs b/Services/Bridges/QueueStatusBridge.cs index 47db101..e7f6dc7 100644 --- a/Services/Bridges/QueueStatusBridge.cs +++ b/Services/Bridges/QueueStatusBridge.cs @@ -1,4 +1,5 @@ using System.Collections.Concurrent; +using System.Text; using System.Text.Json; using Craft.Orchestration; @@ -11,9 +12,10 @@ namespace Craft.Services; /// -/// Static bridge allowing PowerShell (Get-CIPPQueueData) to query orchestrator/job -/// progress without HTTP round-trips. Returns data in the shape the CIPP frontend expects. -/// PS usage: [Craft.Services.QueueStatusBridge]::GetRunStatus($Reference, $QueueId) +/// Run status for PowerShell and the realtime channel, without HTTP round-trips. App-neutral: a run's +/// counts, status and task list plus whatever label/link the app registered; the app shapes it for its own +/// UI. JSON is camelCase. +/// PS usage: [Craft.Services.QueueStatusBridge]::GetRun($QueueId) / ::GetRuns($Lookup) /// public static class QueueStatusBridge { @@ -21,11 +23,8 @@ public static class QueueStatusBridge private static OrchestratorService? s_orchestratorService; private static JobQueueStatusReader? s_queueReader; - /// - /// Maps QueueId (GUID) or Reference to friendly display metadata (Name, Link). - /// Populated by New-CippQueueEntry in CIPPNG mode. - /// - private static readonly ConcurrentDictionary s_queueMetadata = new(StringComparer.OrdinalIgnoreCase); + /// Display metadata the app registered for a queue id or reference. + private static readonly ConcurrentDictionary s_labels = new(StringComparer.OrdinalIgnoreCase); public static void Initialize(JobManager jobManager, OrchestratorService? orchestratorService = null, JobQueueStatusReader? queueReader = null) @@ -36,115 +35,207 @@ public static void Initialize(JobManager jobManager, OrchestratorService? orches } /// - /// Register friendly queue metadata from PowerShell (New-CippQueueEntry). - /// PS usage: [Craft.Services.QueueStatusBridge]::RegisterQueueMetadata($QueueId, $Name, $Link, $Reference) + /// Register a display label and link for a queue id (and its reference), returned with its runs. + /// PS usage: [Craft.Services.QueueStatusBridge]::RegisterQueueMetadata($QueueId, $Label, $Link, $Reference) /// - public static void RegisterQueueMetadata(string queueId, string name, string link, string reference) + public static void RegisterQueueMetadata(string queueId, string label, string link, string reference) { - var meta = new QueueMetadata { Name = name, Link = link ?? "", Reference = reference ?? "" }; + var meta = new RunLabel(label ?? "", link ?? ""); if (!string.IsNullOrEmpty(queueId)) - s_queueMetadata[queueId] = meta; + s_labels[queueId] = meta; if (!string.IsNullOrEmpty(reference)) - s_queueMetadata[reference] = meta; + s_labels[reference] = meta; } /// - /// Get queue/run status in the format expected by the CIPP frontend. - /// Looks up by run name (Reference) or returns all recent runs. - /// Returns a JSON string matching the Get-CIPPQueueData output shape. + /// Each run matching (a run name, a reference, or a queue id the run name ends + /// with), or every run of the last 3 hours when it is empty. JSON array of . /// - /// Optional run reference/name to filter by (maps to RunName in JobManager) - /// Optional queue ID (same as reference in Craft context) - /// JSON array of queue status objects - public static string GetRunStatus(string? reference = null, string? queueId = null) + public static string GetRuns(string? lookup = null) { if (s_jobManager == null) return "[]"; - - // PowerShell converts $null to "" when calling .NET string parameters, - // so treat empty strings the same as null. - var effectiveQueueId = string.IsNullOrEmpty(queueId) ? null : queueId; - var effectiveReference = string.IsNullOrEmpty(reference) ? null : reference; - var lookup = effectiveQueueId ?? effectiveReference; var summaries = GetMergedRunSummaries(); - if (!string.IsNullOrEmpty(lookup)) { - // Try exact match on run name first - var matched = summaries - .Where(s => s.Name.Equals(lookup, StringComparison.OrdinalIgnoreCase)) - .ToList(); + summaries = MatchRuns(summaries, lookup); + } + else + { + var cutoff = DateTime.UtcNow.AddHours(-3); + summaries = summaries.Where(s => s.StartedUtc == null || s.StartedUtc > cutoff).ToList(); + } + return JsonSerializer.Serialize(summaries.Select(s => ToStatus([Load(s)])).ToList(), s_json); + } - // If no exact match, try matching by Reference via orchestrator service - if (matched.Count == 0 && s_orchestratorService != null) - { - var runName = s_orchestratorService.FindRunByReference(lookup); - if (runName != null) - { - matched = summaries - .Where(s => s.Name.Equals(runName, StringComparison.OrdinalIgnoreCase)) - .ToList(); - } - } + /// + /// One for everything matches, chained and child runs + /// rolled up; JSON null when nothing matches yet. The same object the realtime channel pushes. + /// + public static string GetRun(string lookup) + { + if (s_jobManager == null || string.IsNullOrEmpty(lookup)) return "null"; + return GetRunRollups([lookup]).TryGetValue(lookup, out var run) ? run.Data.GetRawText() : "null"; + } - // Last resort: match run names ending with the lookup (QueueId is often the GUID suffix) - if (matched.Count == 0) + /// + /// A run as storage has it (which decides its status), its counts made to agree with its task list (see + /// ), and the task list. + /// + private static (JobRunSummary Storage, JobRunSummary Counts, List Tasks) Load(JobRunSummary s) + { + var tasks = GetTasks(s.Name); + return (s, Reconcile(s, tasks.Select(t => t.Status).ToList()), tasks); + } + + /// + /// The run's counts made to agree with its task list. Storage counts a task done only once its finish is + /// committed, a few seconds after the job manager marks it, so a tracker showing both would have its numbers + /// trail its per-task chips. A list holding every task gives the counts outright. A partial one (tasks not + /// yet claimed here, a fan-out bigger than the list, a parent whose total counts its child runs) can only + /// add: a task it shows finished is finished, so done counts take the higher of the two and the rest is + /// queued. Only the counts: the status stays storage's, which alone knows the run (and any child it is + /// waiting on) has finished, so a run never reports done early. + /// + internal static JobRunSummary Reconcile(JobRunSummary s, IReadOnlyCollection taskStatuses) + { + if (s.Total <= 0 || taskStatuses.Count == 0 || taskStatuses.Count > s.Total) return s; + int queued = 0, running = 0, completed = 0, failed = 0; + foreach (var status in taskStatuses) + { + switch (status) { - matched = summaries - .Where(s => s.Name.EndsWith(lookup, StringComparison.OrdinalIgnoreCase)) - .ToList(); + case "Queued": queued++; break; + case "Running": running++; break; + case "Failed": failed++; break; + default: completed++; break; } + } - summaries = matched; + var r = new JobRunSummary + { + Name = s.Name, + Reference = s.Reference, + Priority = s.Priority, + Total = s.Total, + StartedUtc = s.StartedUtc, + CompletedUtc = s.CompletedUtc + }; + if (taskStatuses.Count == s.Total) + { + (r.Queued, r.Running, r.Completed, r.Failed) = (queued, running, completed, failed); + return r; } - else + + r.Completed = Math.Max(s.Completed, completed); + r.Failed = Math.Max(s.Failed, failed); + r.Running = Math.Min(Math.Max(s.Running, running), s.Total - r.Completed - r.Failed); + r.Queued = s.Total - r.Completed - r.Failed - r.Running; + return r; + } + + /// + /// One status for one or more runs of the same queue: the first run's name, reference and label, the + /// summed counts, every run's tasks, and the status of the runs together as storage has them. + /// + private static RunStatus ToStatus(List<(JobRunSummary Storage, JobRunSummary Counts, List Tasks)> runs) + { + var first = runs[0].Counts; + var reference = s_orchestratorService?.GetRunReference(first.Name) ?? first.Name; + JobRunSummary counts = new(), storage = new(); + foreach (var (st, c, _) in runs) { - // Return only runs from last 3 hours (matches legacy behavior) - var cutoff = DateTime.UtcNow.AddHours(-3); - summaries = summaries - .Where(s => s.StartedUtc == null || s.StartedUtc > cutoff) + counts.Total += c.Total; + counts.Queued += c.Queued; + counts.Running += c.Running; + counts.Completed += c.Completed; + counts.Failed += c.Failed; + storage.Queued += st.Queued; + storage.Running += st.Running; + storage.Completed += st.Completed; + storage.Failed += st.Failed; + } + var label = FindLabel(reference, first.Name); + return new RunStatus + { + RunName = first.Name, + Reference = reference, + Label = label?.Label, + Link = label?.Link, + Status = DeriveStatus(storage), + Total = counts.Total, + Queued = counts.Queued, + Running = counts.Running, + Completed = counts.Completed, + Failed = counts.Failed, + StartedUtc = runs.Min(r => r.Counts.StartedUtc), + Tasks = runs.SelectMany(r => r.Tasks).ToList() + }; + } + + /// The label registered for a run: by its reference, its name, or the queue id its name ends with. + private static RunLabel? FindLabel(string reference, string runName) + { + if (s_labels.TryGetValue(reference, out var label) || s_labels.TryGetValue(runName, out label)) return label; + return runName.Length > 36 && Guid.TryParse(runName[^36..], out _) && s_labels.TryGetValue(runName[^36..], out label) + ? label + : null; + } + + /// + /// Runs for one queue id or reference: exact run name, then the orchestrator's reference, then the + /// run-name suffix (apps commonly name runs "Name-<QueueId>"). + /// + private static List MatchRuns(List summaries, string lookup) + { + var matched = summaries + .Where(s => s.Name.Equals(lookup, StringComparison.OrdinalIgnoreCase)) + .ToList(); + + if (matched.Count == 0 && s_orchestratorService?.FindRunByReference(lookup) is { } runName) + { + matched = summaries + .Where(s => s.Name.Equals(runName, StringComparison.OrdinalIgnoreCase)) .ToList(); } - var result = summaries.Select(s => + if (matched.Count == 0) { - var completedTasks = s.Completed + s.Failed; - var total = Math.Max(s.Total, 1); - var status = DeriveStatus(s); - var runReference = s_orchestratorService?.GetRunReference(s.Name) ?? s.Name; + matched = summaries + .Where(s => s.Name.EndsWith(lookup, StringComparison.OrdinalIgnoreCase)) + .ToList(); + } - // Look up friendly metadata by reference, then by run name, - // then by QueueId GUID suffix (run names follow "OrchestratorName-" pattern) - s_queueMetadata.TryGetValue(runReference, out var meta); - if (meta == null) - s_queueMetadata.TryGetValue(s.Name, out meta); - if (meta == null && s.Name.Length > 36) - { - var guidSuffix = s.Name[^36..]; - if (Guid.TryParse(guidSuffix, out _)) - s_queueMetadata.TryGetValue(guidSuffix, out meta); - } + return matched; + } - return new QueueStatusEntry - { - PartitionKey = "CippQueue", - RowKey = s.Name, - Name = meta?.Name ?? s.Name, - Link = meta?.Link ?? "", - Reference = runReference, - TotalTasks = s.Total, - CompletedTasks = completedTasks, - RunningTasks = s.Running, - FailedTasks = s.Failed, - PercentComplete = Math.Round(((double)completedTasks / total) * 100, 1), - PercentFailed = Math.Round(((double)s.Failed / total) * 100, 1), - PercentRunning = Math.Round(((double)s.Running / total) * 100, 1), - Tasks = GetTaskDetails(s.Name), - Status = status, - Timestamp = s.StartedUtc?.ToString("O") ?? DateTime.UtcNow.ToString("O") - }; - }).ToList(); + /// + /// What the realtime pump pushes for one queue id: its , and a signature of its + /// progress that changes only when the counts or a task's status do. + /// + internal sealed record RunRollup(string Status, JsonElement Data, string State); - return JsonSerializer.Serialize(result, s_jsonOptions); + /// The rolled-up status per queue id from one read of the run summaries. Ids with no run yet are left out. + internal static Dictionary GetRunRollups(IReadOnlyCollection ids) => + s_jobManager == null || ids.Count == 0 + ? new Dictionary(StringComparer.OrdinalIgnoreCase) + : RollUp(GetMergedRunSummaries(), ids); + + internal static Dictionary RollUp(List summaries, IReadOnlyCollection ids) + { + var result = new Dictionary(StringComparer.OrdinalIgnoreCase); + foreach (var id in ids) + { + var runs = MatchRuns(summaries, id); + if (runs.Count == 0) continue; + + var status = ToStatus(runs.Select(Load).ToList()); + var state = new StringBuilder() + .Append(status.Status).Append('|').Append(status.Total).Append('|').Append(status.Queued) + .Append('|').Append(status.Running).Append('|').Append(status.Completed).Append('|').Append(status.Failed); + foreach (var t in status.Tasks) state.Append('|').Append(t.Name).Append(':').Append(t.Status); + result[id] = new RunRollup(status.Status, JsonSerializer.SerializeToElement(status, s_json), state.ToString()); + } + return result; } /// @@ -173,11 +264,12 @@ private static List GetMergedRunSummaries() return s_jobManager?.GetRunSummaries() ?? []; } + /// Queued, Running, Completed or CompletedWithErrors. private static string DeriveStatus(JobRunSummary s) { if (s.Queued == 0 && s.Running == 0) { - return s.Failed > 0 ? "Completed (with errors)" : "Completed"; + return s.Failed > 0 ? "CompletedWithErrors" : "Completed"; } if (s.Running > 0 || s.Completed > 0 || s.Failed > 0) { @@ -186,77 +278,50 @@ private static string DeriveStatus(JobRunSummary s) return "Queued"; } - private static List GetTaskDetails(string runName) + /// A run's most recent tasks, each named by its job name without the run's own prefix. + private static List GetTasks(string runName) { if (s_jobManager == null) return []; - var jobs = s_jobManager.GetRunJobs(runName, limit: 100); - return jobs.Select(j => new TaskDetail + return s_jobManager.GetRunJobs(runName, limit: 100).Select(j => new RunTask { - Timestamp = (j.CompletedUtc ?? j.StartedUtc ?? j.QueuedUtc).ToString("O"), - Name = ExtractTaskDisplayName(j.Name, runName), - Status = j.Status + Name = j.Name.StartsWith(runName + "-", StringComparison.OrdinalIgnoreCase) ? j.Name[(runName.Length + 1)..] : j.Name, + Status = j.Status, + At = j.CompletedUtc ?? j.StartedUtc ?? j.QueuedUtc }).ToList(); } - /// - /// Extract a clean display name from the full job name. - /// Job names follow pattern: "{RunName}-{FunctionName}_{TenantFilter}" - /// We want just the tenant (or meaningful identifier) portion. - /// - private static string ExtractTaskDisplayName(string jobName, string runName) - { - // Strip the run name prefix (e.g. "GraphRequestOrchestrator-guid-") - var remaining = jobName; - if (jobName.StartsWith(runName + "-", StringComparison.OrdinalIgnoreCase)) - remaining = jobName[(runName.Length + 1)..]; + // ── Output models ── - // Strip function name prefix (e.g. "ListGraphRequestQueue_") — take everything after first underscore - var underscoreIdx = remaining.IndexOf('_'); - if (underscoreIdx >= 0 && underscoreIdx < remaining.Length - 1) - return remaining[(underscoreIdx + 1)..]; - - return remaining; - } - - // ── Output Models (match Get-CIPPQueueData shape) ── - - private sealed class QueueStatusEntry + /// A queue's runs as one status. Serialized camelCase. + internal sealed class RunStatus { - public string PartitionKey { get; set; } = ""; - public string RowKey { get; set; } = ""; - public string Name { get; set; } = ""; - public string Link { get; set; } = ""; + public string RunName { get; set; } = ""; public string Reference { get; set; } = ""; - public int TotalTasks { get; set; } - public int CompletedTasks { get; set; } - public int RunningTasks { get; set; } - public int FailedTasks { get; set; } - public double PercentComplete { get; set; } - public double PercentFailed { get; set; } - public double PercentRunning { get; set; } - public List Tasks { get; set; } = []; + public string? Label { get; set; } + public string? Link { get; set; } public string Status { get; set; } = ""; - public string Timestamp { get; set; } = ""; + public int Total { get; set; } + public int Queued { get; set; } + public int Running { get; set; } + public int Completed { get; set; } + public int Failed { get; set; } + public DateTime? StartedUtc { get; set; } + public List Tasks { get; set; } = []; } - private sealed class TaskDetail + internal sealed class RunTask { - public string Timestamp { get; set; } = ""; public string Name { get; set; } = ""; public string Status { get; set; } = ""; + public DateTime At { get; set; } } - private static readonly JsonSerializerOptions s_jsonOptions = new() + private sealed record RunLabel(string Label, string Link); + + private static readonly JsonSerializerOptions s_json = new() { - PropertyNamingPolicy = null, + PropertyNamingPolicy = JsonNamingPolicy.CamelCase, WriteIndented = false }; - - private sealed class QueueMetadata - { - public string Name { get; set; } = ""; - public string Link { get; set; } = ""; - public string Reference { get; set; } = ""; - } } diff --git a/Services/Bridges/RealtimeBridge.cs b/Services/Bridges/RealtimeBridge.cs index acbd176..5c87cce 100644 --- a/Services/Bridges/RealtimeBridge.cs +++ b/Services/Bridges/RealtimeBridge.cs @@ -11,7 +11,7 @@ namespace Craft.Services; /// /// 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 / -/// . Delivery, gating, GUID/size enforcement and storage live in +/// . Delivery, gating, job-id/size enforcement and storage live in /// . /// /// PS usage: @@ -19,8 +19,8 @@ namespace Craft.Services; /// [Craft.Services.RealtimeBridge]::Publish($userId, $jobId, "update", @{ done = 142; total = 300 }) /// [Craft.Services.RealtimeBridge]::Publish($userId, $jobId, "end", @{ done = 300; total = 300 }) /// -/// Only userId and jobId (a GUID) are required; everything else is optional. Every call is -/// best-effort and never throws back to the caller. +/// Only userId and jobId (a short token: 1-128 letters, digits, -_.:) are required; everything +/// else is optional. Every call is best-effort and never throws back to the caller. /// public static class RealtimeBridge { @@ -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); + /// + /// Grant the events of , so later + /// 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) + /// + public static void Watch(string userId, string jobId) + { + try { s_service?.Watch(userId, jobId, trackRun: false); } + catch { /* best-effort */ } + } + + /// + /// , and have Craft push the orchestrator run status for + /// (a queue id; matched like ) until the run finishes. + /// + public static void WatchRun(string userId, string jobId) + { + try { s_service?.Watch(userId, jobId, trackRun: true); } + catch { /* best-effort */ } + } + + /// Signal every user granted (mode update, no payload). + public static void Notify(string jobId) => Notify(jobId, null, null); + + public static void Notify(string jobId, string? mode) => Notify(jobId, mode, null); + + /// Publish to every user granted through or . + public static void Notify(string jobId, string? mode, object? data) + { + try { s_service?.Notify(jobId, mode, data); } + catch { /* best-effort */ } + } + /// Full form. is start | update | end (default update). public static void Publish(string userId, string jobId, string? mode, object? data, string? urlHref, string? urlLabel, int? status, string? message) diff --git a/Services/Configuration/RealtimeSettings.cs b/Services/Configuration/RealtimeSettings.cs index daa9c72..3074fc0 100644 --- a/Services/Configuration/RealtimeSettings.cs +++ b/Services/Configuration/RealtimeSettings.cs @@ -30,8 +30,9 @@ public class RealtimeSettings /// Heartbeat comment interval, seconds, to keep the stream alive through proxies. Default 20. public int HeartbeatSeconds { get; set; } = 20; - /// TTL for a stored entry that never receives an end (crash backstop), minutes. Default 60. - public int EntryTtlMinutes { get; set; } = 60; + /// How long a job's last frame (an end included) is kept for reconnect replay, and a grant + /// with no activity is kept, minutes. Default 30. + public int EntryTtlMinutes { get; set; } = 30; /// /// Resolved enabled state. The CRAFT_REALTIME_ENABLED environment variable (true/1 or false/0) diff --git a/Services/Hosting/CraftHostBuilderExtensions.cs b/Services/Hosting/CraftHostBuilderExtensions.cs index 810ae73..faa8e3b 100644 --- a/Services/Hosting/CraftHostBuilderExtensions.cs +++ b/Services/Hosting/CraftHostBuilderExtensions.cs @@ -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" }; /// /// Resolves the on-the-fly compression level from the CRAFT_API_COMPRESSION_LEVEL env var diff --git a/Services/Hosting/Endpoints/RealtimeEndpoint.cs b/Services/Hosting/Endpoints/RealtimeEndpoint.cs index 99e5a77..c61283a 100644 --- a/Services/Hosting/Endpoints/RealtimeEndpoint.cs +++ b/Services/Hosting/Endpoints/RealtimeEndpoint.cs @@ -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; @@ -11,6 +12,9 @@ namespace Craft.Hosting.Endpoints; /// public static class RealtimeEndpoint { + /// The stream's path. Program.cs routes it through the /api compressor. + public const string Path = "/.craft/events"; + /// /// Maps the SSE endpoint on nodes that face a browser, when realtime is switched on. /// @@ -40,11 +44,11 @@ public static WebApplication MapCraftRealtimeEndpoint( var realtime = app.Services.GetRequiredService(); 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 @@ -52,7 +56,8 @@ public static WebApplication MapCraftRealtimeEndpoint( 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()?.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 @@ -101,4 +106,35 @@ public static WebApplication MapCraftRealtimeEndpoint( logger.LogInformation("[System] Realtime SSE endpoint: /.craft/events"); return app; } + + /// + /// The principal name of a signed-in user, read after 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. + /// + 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; + } } diff --git a/Services/Program.cs b/Services/Program.cs index 414bd3d..d1f73a2 100644 --- a/Services/Program.cs +++ b/Services/Program.cs @@ -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})", @@ -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 diff --git a/Services/Realtime/RealtimeService.cs b/Services/Realtime/RealtimeService.cs index a4f9a7f..c538196 100644 --- a/Services/Realtime/RealtimeService.cs +++ b/Services/Realtime/RealtimeService.cs @@ -39,6 +39,14 @@ public sealed class RealtimeService : IDisposable private readonly ConcurrentDictionary> _conns = new(StringComparer.OrdinalIgnoreCase); + // Watch grants: jobId -> the users allowed to receive that job's events. Granted server-side by the + // app, from the authenticated request that started the job; the browser never names a job itself. + private readonly ConcurrentDictionary _watches = new(StringComparer.OrdinalIgnoreCase); + private readonly Timer? _runPump; + private int _pumping; + + private static readonly TimeSpan RunPumpInterval = TimeSpan.FromSeconds(3); + private long _seq; private int _connectionCount; @@ -56,6 +64,23 @@ public RealtimeService(CraftSettings settings, ILogger logger) var period = TimeSpan.FromMinutes(Math.Clamp(_cfg.EntryTtlMinutes, 1, 1440) / 2.0 + 0.5); _sweep = new Timer(_ => SweepExpired(), null, period, period); + _runPump = new Timer(_ => PumpRuns(), null, RunPumpInterval, RunPumpInterval); + } + + /// A watched run that has not appeared by then is never going to: stop reading run status for it. + internal TimeSpan RunStartGrace { get; set; } = TimeSpan.FromMinutes(10); + + /// Run-status source for the pump; swapped in tests. + internal Func, Dictionary> RunStatusSource { get; set; } = + QueueStatusBridge.GetRunRollups; + + private sealed class JobWatch + { + public readonly ConcurrentDictionary Users = new(StringComparer.OrdinalIgnoreCase); + public volatile bool TrackRun; + public string? LastRunState; + public readonly long CreatedTimestamp = Stopwatch.GetTimestamp(); + public long UpdatedTimestamp = Stopwatch.GetTimestamp(); } public bool Enabled => _enabled; @@ -90,19 +115,14 @@ public Connection(int capacity) // ── Publish ───────────────────────────────────────────────────────────────── /// - /// Publish a job event. and (a GUID) are - /// required; everything else is optional. Best-effort and non-throwing. + /// Publish a job event. and (see ) + /// are required; everything else is optional. Best-effort and non-throwing. /// public void Publish(string userId, string jobId, string? mode, object? data, string? urlHref, string? urlLabel, int? status, string? message) { if (!_enabled) return; - if (string.IsNullOrWhiteSpace(userId) || string.IsNullOrWhiteSpace(jobId)) return; - if (!Guid.TryParse(jobId, out _)) - { - _logger.LogWarning("[Realtime] Rejected non-GUID jobId '{JobId}'", jobId); - return; - } + if (string.IsNullOrWhiteSpace(userId) || !IsSafeJobId(jobId, _logger)) return; var m = NormalizeMode(mode); var outStatus = status; @@ -131,13 +151,10 @@ public void Publish(string userId, string jobId, string? mode, object? data, var frame = BuildFrame(jobId, m, seq, outStatus, outMessage, urlHref, urlLabel, dataJson); var key = Key(userId, jobId); - if (m == "end") - { - _matrix.TryRemove(key, out _); - } - else if (!oversized) + if (!oversized) { - // Store the delivered frame as the current message for reconnect replay. On oversize we + // Store the delivered frame as the current message for reconnect replay, "end" included so a tab + // that dropped before the end still gets the final state; the TTL sweep evicts it. On oversize we // keep the previous good frame (do not overwrite it with a truncated marker). if (_matrix.ContainsKey(key) || _matrix.Count < _cfg.MaxActiveJobs) { @@ -161,6 +178,108 @@ public void Publish(string userId, string jobId, string? mode, object? data, c.Enqueue(frame); } + // ── Watches ────────────────────────────────────────────────────────────────── + + /// + /// Grant the events of (see ). Call it from the + /// authenticated request that started the job, with that request's principal name. With + /// , Craft also pushes the job's orchestrator run status itself, matched the + /// same way as , until the run finishes. + /// + public void Watch(string userId, string jobId, bool trackRun) + { + if (!_enabled || string.IsNullOrWhiteSpace(userId) || !IsSafeJobId(jobId, _logger)) return; + + if (!_watches.TryGetValue(jobId, out var watch)) + { + if (_watches.Count >= _cfg.MaxActiveJobs) + { + _logger.LogWarning("[Realtime] MaxActiveJobs ({Max}) reached — not watching job {JobId}", + _cfg.MaxActiveJobs, jobId); + return; + } + watch = _watches.GetOrAdd(jobId, _ => new JobWatch()); + } + + watch.Users[userId] = 0; + watch.UpdatedTimestamp = Stopwatch.GetTimestamp(); + if (trackRun) watch.TrackRun = true; + } + + /// Publish to every user granted through . + public void Notify(string jobId, string? mode, object? data) + { + if (!_enabled || string.IsNullOrWhiteSpace(jobId) || !_watches.TryGetValue(jobId, out var watch)) return; + watch.UpdatedTimestamp = Stopwatch.GetTimestamp(); + foreach (var userId in watch.Users.Keys) + Publish(userId, jobId, mode, data, null, null, null, null); + } + + /// + /// Push the run status of every tracked watch whose counts changed since the last push. One read of + /// the run summaries serves all watches, so the cost does not grow with the number of open tabs. + /// + internal void PumpRuns() + { + if (Interlocked.Exchange(ref _pumping, 1) == 1) return; // a slow storage read is still running + try + { + var tracked = new List(); + foreach (var kv in _watches) + if (kv.Value.TrackRun) tracked.Add(kv.Key); + if (tracked.Count == 0) return; + + var rollups = RunStatusSource(tracked); + foreach (var jobId in tracked) + { + if (!_watches.TryGetValue(jobId, out var watch)) continue; + if (!rollups.TryGetValue(jobId, out var run)) + { + if (Stopwatch.GetElapsedTime(watch.CreatedTimestamp) <= RunStartGrace) continue; + // Tell the watcher to stop waiting rather than leave a tracker that no longer polls hanging. + watch.TrackRun = false; + Notify(jobId, "end", new Dictionary { ["status"] = "NotFound" }); + continue; + } + + if (run.State == watch.LastRunState) continue; + watch.LastRunState = run.State; + + var finished = run.Status is "Completed" or "CompletedWithErrors"; + if (finished) watch.TrackRun = false; + Notify(jobId, finished ? "end" : "update", run.Data); + } + } + catch (Exception ex) + { + _logger.LogWarning(ex, "[Realtime] Run status pump failed"); + } + finally + { + Volatile.Write(ref _pumping, 0); + } + } + + private const int MaxJobIdLength = 128; + + /// + /// A job id is any short token an app already uses (a GUID, BEC-20261008-ab12, a run name): 1-128 + /// letters, digits, -, _, . or :. Anything else is dropped, never thrown, and only + /// its length is logged, so a hostile id cannot reach a log line, a key separator or the frame. + /// + internal static bool IsSafeJobId(string? jobId, ILogger? logger = null) + { + if (jobId is { Length: > 0 and <= MaxJobIdLength }) + { + var safe = true; + foreach (var c in jobId) + if (!(char.IsAsciiLetterOrDigit(c) || c is '-' or '_' or '.' or ':')) { safe = false; break; } + if (safe) return true; + } + logger?.LogWarning("[Realtime] Rejected a job id that is not a short token ({Length} chars)", jobId?.Length ?? 0); + return false; + } + // ── SSE connection management ──────────────────────────────────────────────── /// Register a new SSE connection for a user. Returns null id/connection if over the cap. @@ -209,6 +328,9 @@ private void SweepExpired() foreach (var kv in _matrix) if (Stopwatch.GetElapsedTime(kv.Value.UpdatedTimestamp) > ttl) _matrix.TryRemove(kv.Key, out _); + foreach (var kv in _watches) + if (!kv.Value.TrackRun && Stopwatch.GetElapsedTime(kv.Value.UpdatedTimestamp) > ttl) + _watches.TryRemove(kv.Key, out _); } catch (Exception ex) { @@ -250,6 +372,8 @@ private static string BuildFrame(string jobId, string mode, long seq, int? statu { null => null, string or bool or int or long or double or float or decimal or DateTime or DateTimeOffset or Guid => v, + JsonElement => v, + PSObject { BaseObject: PSCustomObject } ps => NormalizeProperties(ps), PSObject ps => Normalize(ps.BaseObject), IDictionary d => NormalizeDict(d), IEnumerable e => NormalizeList(e), @@ -264,6 +388,15 @@ private static string BuildFrame(string jobId, string mode, long seq, int? statu return r; } + // A [pscustomobject] (and anything ConvertFrom-Json returns) has no CLR shape of its own: its note properties are the data. + private static Dictionary NormalizeProperties(PSObject ps) + { + var r = new Dictionary(StringComparer.Ordinal); + foreach (var p in ps.Properties) + r[p.Name] = Normalize(p.Value); + return r; + } + private static List NormalizeList(IEnumerable e) { var r = new List(); @@ -271,5 +404,9 @@ private static string BuildFrame(string jobId, string mode, long seq, int? statu return r; } - public void Dispose() => _sweep?.Dispose(); + public void Dispose() + { + _sweep?.Dispose(); + _runPump?.Dispose(); + } } diff --git a/docs/configuration.md b/docs/configuration.md index 853d93c..da7ec5b 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -672,10 +672,47 @@ endpoint is still only mapped by nodes carrying the **Http** or **Frontend** rol "MaxConnections": 1000, // max concurrent SSE streams "PerConnectionQueue": 256, // buffered frames per connection before the oldest is dropped "HeartbeatSeconds": 20, // keep-alive comment interval - "EntryTtlMinutes": 60 // backstop eviction for jobs that never send "end" + "EntryTtlMinutes": 30 // how long a job's last frame (an "end" included) stays for reconnect replay } ``` +A stream opens only for a signed-in user: the principal Craft's auth middleware resolved must carry a role +other than `anonymous`, so app-only API clients and anonymous callers get a 401. The browser never names a +job. The app grants a job to the user from the authenticated request that started it, and events for that +job then reach that user only: + +```powershell +$User = $Request.Headers.'x-ms-client-principal-name' +[Craft.Services.RealtimeBridge]::Watch($User, $JobId) # grant; the app publishes with Notify +[Craft.Services.RealtimeBridge]::WatchRun($User, $QueueId) # grant, and Craft pushes the run status itself +[Craft.Services.RealtimeBridge]::Notify($JobId) # from any worker: signal every granted user +``` + +`WatchRun` pushes what `QueueStatusBridge.GetRun($QueueId)` returns, every 3 seconds while the counts or a +task's status change, ending with an `end` frame when the run finishes. A run that has not appeared within 10 +minutes ends with `status: "NotFound"`. One read of the run summaries serves every watch. The status is +app-neutral (camelCase); shaping it for a UI is the app's job: + +```jsonc +{ + "runName": "SyncOrchestrator-6f1c2b8e-...", "reference": "...", "label": "Sync users", "link": "", + "status": "Running", // Queued | Running | Completed | CompletedWithErrors + "total": 16, "queued": 0, "running": 1, "completed": 14, "failed": 1, + "startedUtc": "2026-10-08T01:31:41Z", + "tasks": [ { "name": "Push-Sync_contoso.com", "status": "Completed", "at": "2026-10-08T01:31:45Z" } ] +} +``` + +A queue id that spans several runs (chained continuations, child runs named with the same id) is rolled up into +one status. `GetRuns($Lookup)` returns the same shape per run. Counts follow the task list when it holds every +task, so they never trail it; the status follows storage, so a run is never reported done before it and any +child it waits on have finished. + +The stream is compressed by the `/api` compressor under the same switch (`App:Api:Compression` / +`CRAFT_API_COMPRESSION`) and level: br or gzip as the browser accepts, one compression context per connection, so +the repeated keys and tenant names of successive frames compress against each other. Each batch of frames is +flushed as it is written, so compression never delays an event. + ### Frontend EasyAuth handles auth, redirects, and excluded paths at the App Service platform layer (see `Setup.UnauthenticatedClientAction` and `Setup.ExcludedPaths`). CRAFT only adds response headers EasyAuth doesn't touch — currently just CSP. diff --git a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 index 3d33224..ef1a48f 100644 --- a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 +++ b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 @@ -5,7 +5,7 @@ Author = 'CRAFT perf-harness' Description = 'Synthetic HTTP endpoints for load-testing CRAFT in http-only mode. Not for production.' PowerShellVersion = '7.2' - FunctionsToExport = @('Invoke-PerfPing', 'Invoke-PerfEcho', 'Invoke-PerfCpu', 'Invoke-PerfSleep', 'Invoke-PerfJson', 'Invoke-PerfFile', 'Invoke-PerfBgEnqueue', 'Push-PerfBg', 'Push-PerfBgLeaf', 'Invoke-PerfManyRuns', 'Push-PerfHold', 'Invoke-PerfThreads', 'Invoke-PerfTableOp', 'Push-PerfCheck', 'Invoke-PerfCheckCounts', 'Push-PerfSeq', 'Invoke-PerfSeqResult', 'Push-PerfSeqWorker', 'Invoke-PerfSeqWorkerEnqueue', 'Invoke-PerfSeqWorkerResult', 'Invoke-PerfThreadBreakdown', 'Invoke-PerfSeedRuns', 'Invoke-ListPerf', 'Invoke-PerfWhoami', 'Invoke-PerfTimerTick', 'Invoke-PerfTimerCount', 'Invoke-PerfPublish', 'Invoke-PerfAllocation', 'Invoke-PerfRuns', 'Push-PerfE2E', 'Push-PerfE2EPost', 'Invoke-PerfE2EStart', 'Invoke-PerfE2EState', 'Invoke-PerfE2ERuns', 'Invoke-PerfE2EBridge', 'Invoke-PerfE2ELegacy') + FunctionsToExport = @('Invoke-PerfPing', 'Invoke-PerfEcho', 'Invoke-PerfCpu', 'Invoke-PerfSleep', 'Invoke-PerfJson', 'Invoke-PerfFile', 'Invoke-PerfBgEnqueue', 'Push-PerfBg', 'Push-PerfBgLeaf', 'Invoke-PerfManyRuns', 'Push-PerfHold', 'Invoke-PerfThreads', 'Invoke-PerfTableOp', 'Push-PerfCheck', 'Invoke-PerfCheckCounts', 'Push-PerfSeq', 'Invoke-PerfSeqResult', 'Push-PerfSeqWorker', 'Invoke-PerfSeqWorkerEnqueue', 'Invoke-PerfSeqWorkerResult', 'Invoke-PerfThreadBreakdown', 'Invoke-PerfSeedRuns', 'Invoke-ListPerf', 'Invoke-PerfWhoami', 'Invoke-PerfTimerTick', 'Invoke-PerfTimerCount', 'Invoke-PerfPublish', 'Invoke-PerfWatchRun', 'Invoke-PerfAllocation', 'Invoke-PerfRuns', 'Push-PerfE2E', 'Push-PerfE2EPost', 'Invoke-PerfE2EStart', 'Invoke-PerfE2EState', 'Invoke-PerfE2ERuns', 'Invoke-PerfE2EBridge', 'Invoke-PerfE2ELegacy') CmdletsToExport = @() VariablesToExport = @() AliasesToExport = @() diff --git a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 index cb82228..109486b 100644 --- a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 +++ b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 @@ -563,6 +563,15 @@ function Invoke-PerfPublish { return @{ StatusCode = 200; Body = @{ ok = $true; endpoint = 'PerfPublish'; jobId = $jobId; mode = $mode; userId = $userId } } } +# Grant the caller a queue id's run status over the realtime channel (?jobId=). ok is false on an image +# without WatchRun, so the realtime orchestration check can skip instead of failing. +function Invoke-PerfWatchRun { + param($Request, $TriggerMetadata) + $Has = $null -ne [Craft.Services.RealtimeBridge].GetMethod('WatchRun') + if ($Has) { [Craft.Services.RealtimeBridge]::WatchRun([string]$Request.Headers.'x-ms-client-principal-name', [string]$Request.Query.jobId) } + return @{ StatusCode = 200; Body = @{ ok = $Has; jobId = [string]$Request.Query.jobId } } +} + # -- Orchestration e2e probes (scripts/run-e2e-orchestration.ps1) ----------------------------------------- # Each check gets its own namespace (ns). Tasks and PostExecutions record what they saw under that ns, in the # shared cache by default or in the E2EProbe table (sink=table) when the record has to survive a restart. @@ -800,7 +809,7 @@ function Invoke-PerfE2EBridge { startHasMaxConcurrency = ($Def -match 'MaxConcurrency'); startHasStopOnFailure = ($Def -match 'StopOnFailure') } } 'active' { $Body = @{ active = [Craft.Services.OrchestratorBridge]::IsRunActive($Name) } } - 'queue' { $Body = @{ entries = @([Craft.Services.QueueStatusBridge]::GetRunStatus($null, $Name) | ConvertFrom-Json) } } + 'queue' { $Body = @{ entries = @([Craft.Services.QueueStatusBridge]::GetRuns($Name) | ConvertFrom-Json) } } 'workers' { $Body = @{ busy = @([Craft.Services.WorkerMetricsBridge]::GetSnapshot().BgPool.Workers | Where-Object IsBusy | ForEach-Object { @{ id = $_.WorkerId; fn = $_.CurrentFunction } }) } diff --git a/perf-harness/scripts/run-e2e-orchestration.ps1 b/perf-harness/scripts/run-e2e-orchestration.ps1 index be0772e..8f5f68f 100644 --- a/perf-harness/scripts/run-e2e-orchestration.ps1 +++ b/perf-harness/scripts/run-e2e-orchestration.ps1 @@ -13,7 +13,7 @@ them. The restart check runs last: it restarts the SUT container. $OrchChecks (from run-e2e.ps1) limits the run to the named groups: - post fail child prio collide attrib cancel seq order status maxconc stopfail perf restart + post fail child prio collide attrib cancel seq order status maxconc stopfail perf realtime restart #> # Perf gates. Calibrated from 3 full runs on the reference dev box (combined role, BgPoolSize=4, cpus=2, @@ -397,9 +397,9 @@ Invoke-OrchCheck 'seq' { foreach ($B in @((Invoke-OrchBridge 'workers').busy)) { if ($B.fn -like "$Name-*") { [void]$Labels.Add($B.fn) } } $S = Get-OrchSummary $Ns; if ($S['sv'].ended -ge 4) { $true } } 60 250 - $Q = Wait-Orch { $E = @((Invoke-OrchBridge 'queue' $Name).entries)[0]; if ($E.Status -eq 'Completed') { $E } } 30 + $Q = Wait-Orch { $E = @((Invoke-OrchBridge 'queue' $Name).entries)[0]; if ($E.status -eq 'Completed') { $E } } 30 $E = $Q.value - Add-Result 'orch-seq' 'each-step-visible' ($Labels.Count -ge 3 -and $E.TotalTasks -eq 4 -and $E.CompletedTasks -eq 4 -and @($E.Tasks).Count -eq 4) '-' "worker labels seen=$($Labels.Count) (of 4 steps); queue total=$($E.TotalTasks) completed=$($E.CompletedTasks) tasks listed=$(@($E.Tasks).Count)" + Add-Result 'orch-seq' 'each-step-visible' ($Labels.Count -ge 3 -and $E.total -eq 4 -and $E.completed -eq 4 -and @($E.tasks).Count -eq 4) '-' "worker labels seen=$($Labels.Count) (of 4 steps); queue total=$($E.total) completed=$($E.completed) tasks listed=$(@($E.tasks).Count)" } # -- 9. Priority and start-order ------------------------------------------------------------------------------ @@ -589,6 +589,103 @@ Invoke-OrchCheck 'perf' { Add-Result 'orch-perf' 'memory-after-perf' ($M.rssMB -le $OrchGates.RssMB) "$($M.rssMB)MB" "rss=$($M.rssMB)MB workingSet=$($M.workingSetMB)MB heap=$($M.heapMB)MB committed=$($M.committedMB)MB; gate rss $($OrchGates.RssMB)MB" } +# -- 15. Realtime run status --------------------------------------------------------------------------------- +# A user granted a queue id (RealtimeBridge.WatchRun) gets Craft's run roll-up over /.craft/events: progress while +# the runs work, then one "end" once every task of every run carrying the id has ended, children and grandchildren +# included (a grandchild queued through the bridge, i.e. drained after the task that queued it returned). Nobody +# else sees those frames. +function Get-OrchPrincipal([string]$User) { + $Json = @{ identityProvider = 'aad'; userId = $User; userDetails = $User; userRoles = @('authenticated') } | ConvertTo-Json -Compress + @{ 'x-ms-client-principal-name' = $User; 'x-ms-client-principal' = [Convert]::ToBase64String([Text.Encoding]::UTF8.GetBytes($Json)) } +} + +function Start-OrchSse([string]$User, [string]$Path) { + $H = Get-OrchPrincipal $User + $ArgLine = "-N -s --compressed --max-time 300 -o `"$Path`" -H `"x-ms-client-principal-name: $User`" -H `"x-ms-client-principal: $($H['x-ms-client-principal'])`" $base/.craft/events" + Start-Process -FilePath $curl -ArgumentList $ArgLine -PassThru -NoNewWindow +} + +function Get-OrchSseFrames([string]$Path) { + $Txt = try { + $Sr = [IO.StreamReader]::new([IO.FileStream]::new($Path, 'Open', 'Read', 'ReadWrite')) + try { $Sr.ReadToEnd() } finally { $Sr.Dispose() } + } catch { '' } + return , @(foreach ($L in ($Txt -split "`n")) { if ($L.StartsWith('data: ')) { $L.Substring(6) | ConvertFrom-Json } }) +} + +function ConvertTo-OrchUnixMs([long]$Ticks) { [long](($Ticks - 621355968000000000) / 10000) } + +# The end frame alone must fill the queue tracker: every task listed, each in its final state. +function Test-OrchFrameTasks($Frame, [int]$Count) { + $T = @($Frame.data.tasks) + $T.Count -eq $Count -and @($T.Where({ $_.status -notin 'Completed', 'Failed' })).Count -eq 0 +} + +$RtNames = @('needs-signed-in-user', 'stream-compressed', 'fanout-progress', 'nested-ends-after-all', 'failure-ends-with-errors', 'other-user-sees-nothing') +Invoke-OrchCheck 'realtime' { + $Caps = try { Invoke-RestMethod "$base/API/PerfWatchRun" -Headers (Get-OrchPrincipal 'rt-probe') -TimeoutSec 20 } catch { $null } + if (-not $Caps.ok) { foreach ($N in $RtNames) { Add-Skip 'orch-realtime' $N 'image lacks RealtimeBridge.WatchRun' }; return } + + $Bare = Fetch "$base/.craft/events" @('--max-time', '3', '-H', 'x-ms-client-principal-name: rt-bare') + Add-Result 'orch-realtime' 'needs-signed-in-user' ($Bare.Code -eq 401) '-' "principal-name header alone -> HTTP $($Bare.Code)" + + $P = Get-OrchPrincipal 'rt-probe' + $Comp = Fetch "$base/.craft/events" @('--max-time', '2', '-H', 'Accept-Encoding: gzip', '-H', "x-ms-client-principal-name: rt-probe", '-H', "x-ms-client-principal: $($P['x-ms-client-principal'])") + Add-Result 'orch-realtime' 'stream-compressed' ($Comp.Headers -match '(?im)^content-encoding:\s*gzip') '-' "Accept-Encoding: gzip -> $((([regex]'(?im)^content-encoding:.*$').Match([string]$Comp.Headers).Value).Trim())" + + $Id = New-OrchId; $Ns = "rt-$Id" + $Fan = [guid]::NewGuid().ToString(); $Tree = [guid]::NewGuid().ToString(); $Bad = [guid]::NewGuid().ToString() + $Mine = New-TemporaryFile; $Theirs = New-TemporaryFile + $Streams = @((Start-OrchSse 'rt-alice' $Mine.FullName), (Start-OrchSse 'rt-mallory' $Theirs.FullName)) + try { + Start-Sleep -Milliseconds 1000 + foreach ($Q in $Fan, $Tree, $Bad) { $null = Invoke-RestMethod "$base/API/PerfWatchRun?jobId=$Q" -Headers (Get-OrchPrincipal 'rt-alice') -TimeoutSec 20 } + $null = Start-OrchRuns $Ns @( + @{ name = "E2ERtFan-$Fan"; label = 'fan'; tasks = 40; holdms = 1000 } + @{ name = "E2ERtP-$Tree"; label = 'rtP'; tasks = 3; holdms = 300 + overrides = @{ '0' = @{ child = @{ name = "E2ERtC-$Tree"; label = 'rtC'; tasks = 2; holdms = 1500 + overrides = @{ '0' = @{ child = @{ name = "E2ERtG-$Tree"; label = 'rtG'; tasks = 4; holdms = 2500; via = 'bridge' } } } } } } } + @{ name = "E2ERtBad-$Bad"; label = 'bad'; tasks = 4; holdms = 300; overrides = @{ '2' = @{ fail = $true } } } + ) + $W = Wait-Orch { + $F = Get-OrchSseFrames $Mine.FullName + if (@($F.Where({ $_.mode -eq 'end' -and $_.jobId -in @($Fan, $Tree, $Bad) }) | Select-Object -ExpandProperty jobId -Unique).Count -eq 3) { $true } + } 180 + Start-Sleep -Seconds 4 # window in which a second "end" or a late frame would show up + $Frames = Get-OrchSseFrames $Mine.FullName + $Rows = Get-OrchRecords $Ns + $Of = { param($Q) , @($Frames.Where({ $_.jobId -eq $Q })) } + $LastEnd = { param([string[]]$Labels) ConvertTo-OrchUnixMs (Get-OrchMax @(foreach ($L in $Labels) { (Get-OrchTasks $Rows $L).end })) } + + $F = & $Of $Fan; $Ends = @($F.Where({ $_.mode -eq 'end' })); $Ups = @($F.Where({ $_.mode -eq 'update' })) + $Done = @($Ups | ForEach-Object { [int]$_.data.completed + [int]$_.data.failed }) + $Mid = @($Ups.Where({ [int]$_.data.completed + [int]$_.data.failed -lt [int]$_.data.total })).Count + $E = $Ends | Select-Object -First 1 + $Ok = $W.ok -and $Ends.Count -eq 1 -and $Mid -ge 1 -and (($Done -join ',') -eq (@($Done | Sort-Object) -join ',')) -and + $E.data.status -eq 'Completed' -and [int]$E.data.completed -eq [int]$E.data.total -and [int]$E.data.total -eq 40 -and + $F[-1].mode -eq 'end' -and $E.ts -ge (& $LastEnd @('fan')) -and (Test-OrchFrameTasks $E 40) -and $E.data.runName + Add-Result 'orch-realtime' 'fanout-progress' $Ok "$($W.sec)s" "updates=$($Ups.Count) (mid-run $Mid) completed=[$($Done -join ',')] ends=$($Ends.Count) end=$($E.data.status) $($E.data.completed)/$($E.data.total) endTasks=$(@($E.data.tasks).Count) run=$($E.data.runName)" + + $F = & $Of $Tree; $Ends = @($F.Where({ $_.mode -eq 'end' })); $E = $Ends | Select-Object -First 1 + $Counts = @('rtP', 'rtC', 'rtG' | ForEach-Object { "$_=$((Get-OrchTasks $Rows $_).Count)" }) -join ' ' + $Lead = if ($E) { $E.ts - (& $LastEnd @('rtP', 'rtC', 'rtG')) } else { 'n/a' } + $Ok = $Ends.Count -eq 1 -and $Counts -eq 'rtP=3 rtC=2 rtG=4' -and $E.data.status -eq 'Completed' -and + [int]$E.data.completed -eq [int]$E.data.total -and [int]$E.data.total -ge 9 -and $Lead -ge 0 -and (Test-OrchFrameTasks $E 9) + Add-Result 'orch-realtime' 'nested-ends-after-all' $Ok "+${Lead}ms" "end minus last grandchild/child/parent task end=${Lead}ms ends=$($Ends.Count) tasks $Counts end=$($E.data.status) $($E.data.completed)/$($E.data.total) endTasks=$(@($E.data.tasks).Count) frames=$($F.Count)" + + $F = & $Of $Bad; $Ends = @($F.Where({ $_.mode -eq 'end' })); $E = $Ends | Select-Object -First 1 + $Ok = $Ends.Count -eq 1 -and $E.data.status -eq 'CompletedWithErrors' -and [int]$E.data.failed -eq 1 -and $E.ts -ge (& $LastEnd @('bad')) -and + (Test-OrchFrameTasks $E 4) -and @(@($E.data.tasks).Where({ $_.status -eq 'Failed' })).Count -eq 1 + Add-Result 'orch-realtime' 'failure-ends-with-errors' $Ok '-' "ends=$($Ends.Count) end=$($E.data.status) failed=$($E.data.failed) $($E.data.completed)/$($E.data.total)" + + $Leak = @((Get-OrchSseFrames $Theirs.FullName).Where({ $_.jobId -in @($Fan, $Tree, $Bad) })).Count + Add-Result 'orch-realtime' 'other-user-sees-nothing' ($Leak -eq 0) '-' "frames for alice's queues on mallory's stream=$Leak" + } finally { + $Streams | Stop-Process -Force -ErrorAction SilentlyContinue + Remove-Item $Mine, $Theirs -ErrorAction SilentlyContinue + } +} + # -- 10. Restart recovery + legacy table drop (restarts the SUT; keep last) ------------------------------------- Invoke-OrchCheck 'restart' { $Id = New-OrchId; $Ns = "rst-$Id"; $Name = "E2ERst-$Id" diff --git a/perf-harness/scripts/run-e2e.ps1 b/perf-harness/scripts/run-e2e.ps1 index b3c4420..def0022 100644 --- a/perf-harness/scripts/run-e2e.ps1 +++ b/perf-harness/scripts/run-e2e.ps1 @@ -132,7 +132,9 @@ try { foreach ($mode in 'start', 'update') { try { Invoke-RestMethod "$base/API/PerfPublish?jobId=$jobId&mode=$mode" -Headers @{ 'x-ms-client-principal-name' = 'ciuser' } -TimeoutSec 20 | Out-Null } catch {} } - $sseTxt = (& $curl '-N' '-s' '--max-time' '4' '-H' $hdr "$base/.craft/events" 2>$null | Out-String) + # The stream opens only for a signed-in user: a principal with a real role, as EasyAuth/Craft auth would pass on. + $ciPrincipal = 'x-ms-client-principal: ' + [Convert]::ToBase64String([Text.Encoding]::UTF8.GetBytes('{"identityProvider":"aad","userId":"ciuser","userDetails":"ciuser","userRoles":["authenticated"]}')) + $sseTxt = (& $curl '-N' '-s' '--max-time' '4' '-H' $hdr '-H' $ciPrincipal "$base/.craft/events" 2>$null | Out-String) try { Invoke-RestMethod "$base/API/PerfPublish?jobId=$jobId&mode=end" -Headers @{ 'x-ms-client-principal-name' = 'ciuser' } -TimeoutSec 20 | Out-Null } catch {} $sseOk = ($sseTxt -match [regex]::Escape($jobId)) -and ($sseTxt -match '"mode":"update"') Add-Result 'realtime' 'sse-deliver' $sseOk '-' "publish -> stored -> SSE replay delivered" diff --git a/tests/Craft.Tests/RealtimeCompressionTests.cs b/tests/Craft.Tests/RealtimeCompressionTests.cs new file mode 100644 index 0000000..cae6c22 --- /dev/null +++ b/tests/Craft.Tests/RealtimeCompressionTests.cs @@ -0,0 +1,106 @@ +using System.IO.Compression; +using System.Net; +using System.Text; +using Craft.Configuration; +using Craft.Hosting; +using Craft.Hosting.Endpoints; +using Craft.Realtime; +using Microsoft.AspNetCore.Builder; +using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.Hosting.Server; +using Microsoft.AspNetCore.Hosting.Server.Features; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +/// +/// The realtime stream goes through the /api compressor on real Kestrel: it negotiates br, and every frame +/// still reaches the browser the moment it is published, because the endpoint flushes and the compressor +/// passes flushes through. A buffered compressor would hold frames until the stream closes, which for SSE +/// is never. +/// +public class RealtimeCompressionTests +{ + private const string User = "alice@contoso.com"; + private const string Job = "6f1c2b8e-1d2a-4c1e-9f0a-3b2c1d4e5f60"; + + private sealed class CountingStream(Stream inner) : Stream + { + public long BytesRead; + public override int Read(byte[] buffer, int offset, int count) { var n = inner.Read(buffer, offset, count); BytesRead += n; return n; } + public override async ValueTask ReadAsync(Memory buffer, CancellationToken ct = default) + { + var n = await inner.ReadAsync(buffer, ct); + BytesRead += n; + return n; + } + public override bool CanRead => true; + public override bool CanSeek => false; + public override bool CanWrite => false; + public override long Length => throw new NotSupportedException(); + public override long Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); } + public override void Flush() { } + public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException(); + public override void SetLength(long value) => throw new NotSupportedException(); + public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException(); + } + + [Fact] + public async Task Stream_IsCompressed_AndEachFrameArrivesWhileOpen() + { + var settings = new CraftSettings { Realtime = new RealtimeSettings { Enabled = true, MaxMessageBytes = 4 * 1024 * 1024 } }; + var builder = WebApplication.CreateBuilder(); + builder.WebHost.ConfigureKestrel(o => o.Listen(IPAddress.Loopback, 0)); + builder.Logging.ClearProviders(); + builder.Services.AddCraftResponseCompression(); + builder.Services.AddSingleton(settings); + builder.Services.AddSingleton(); + var app = builder.Build(); + app.UseWhen(c => c.Request.Path.Equals(RealtimeEndpoint.Path, StringComparison.OrdinalIgnoreCase), + ev => ev.UseResponseCompression()); + app.MapCraftRealtimeEndpoint(CraftRoles.Resolve(settings, _ => null), settings, NullLogger.Instance); + var realtime = app.Services.GetRequiredService(); + + await app.StartAsync(); + try + { + var addr = app.Services.GetRequiredService().Features.Get()!.Addresses.First(); + using var handler = new HttpClientHandler { AutomaticDecompression = DecompressionMethods.None }; + using var client = new HttpClient(handler); + var req = new HttpRequestMessage(HttpMethod.Get, $"{addr}{RealtimeEndpoint.Path}"); + req.Headers.TryAddWithoutValidation("Accept-Encoding", "br"); + req.Headers.TryAddWithoutValidation("x-ms-client-principal-name", User); + req.Headers.TryAddWithoutValidation("x-ms-client-principal", + Convert.ToBase64String(Encoding.UTF8.GetBytes("{\"userDetails\":\"" + User + "\",\"userRoles\":[\"authenticated\"]}"))); + using var resp = await client.SendAsync(req, HttpCompletionOption.ResponseHeadersRead); + Assert.Equal("br", Assert.Single(resp.Content.Headers.ContentEncoding)); + + var wire = new CountingStream(await resp.Content.ReadAsStreamAsync()); + using var reader = new StreamReader(new BrotliStream(wire, CompressionMode.Decompress), Encoding.UTF8); + + // A queue entry the size a tenant fan-out produces. + var tasks = Enumerable.Range(0, 300) + .Select(i => (object?)new Dictionary { ["Name"] = $"tenant{i}.onmicrosoft.com", ["Status"] = "Completed" }) + .ToList(); + realtime.Watch(User, Job, trackRun: false); + realtime.Notify(Job, "update", new Dictionary { ["Status"] = "Running", ["Tasks"] = tasks }); + + using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + string? line; + do { line = await reader.ReadLineAsync(timeout.Token); } + while (line != null && !line.StartsWith("data: ", StringComparison.Ordinal)); + + Assert.NotNull(line); + Assert.Contains(Job, line); + Assert.Contains("tenant299.onmicrosoft.com", line); + Assert.True(wire.BytesRead < Encoding.UTF8.GetByteCount(line) / 4, + $"expected the frame compressed well below its {Encoding.UTF8.GetByteCount(line)} bytes, {wire.BytesRead} came over the wire"); + } + finally + { + await app.StopAsync(); + } + } +} diff --git a/tests/Craft.Tests/RealtimeRunStatusE2ETests.cs b/tests/Craft.Tests/RealtimeRunStatusE2ETests.cs new file mode 100644 index 0000000..4d3e4c8 --- /dev/null +++ b/tests/Craft.Tests/RealtimeRunStatusE2ETests.cs @@ -0,0 +1,188 @@ +using System.Text.Json; +using Craft.Configuration; +using Craft.Orchestration; +using Craft.Realtime; +using Craft.Services; +using Microsoft.Extensions.Logging.Abstractions; +using static Craft.Tests.OrchestrationHarness; + +namespace Craft.Tests; + +/// +/// End to end: real orchestrator, store and status reader, with the realtime run-status pump driven alongside +/// the work. A watched queue id must report progress while its runs work and send exactly one "end", only +/// once every run carrying the id has finished: a plain fan-out, a parent whose tasks queue child runs that +/// queue their own, and child runs that only reach storage after the task that queued them has returned +/// (PowerShell's QueueOrchestration is drained asynchronously). +/// +public class RealtimeRunStatusE2ETests +{ + private const string User = "alice@contoso.com"; + + private sealed class Watcher : IDisposable + { + public readonly RealtimeService Realtime = + new(new CraftSettings { Realtime = new RealtimeSettings { Enabled = true } }, NullLogger.Instance); + public readonly List Frames = []; + private readonly RealtimeService.Connection _conn; + + public Watcher(OrchestrationHarness h, string queueId) + { + var reader = new JobQueueStatusReader(NullLogger.Instance, h.Jobs, h.Store); + // The task list comes from the bridge's job manager. Tests in this class run one at a time. + QueueStatusBridge.Initialize(h.Jobs, h.Svc, reader); + Realtime.RunStatusSource = ids => + { + reader.GetAsync(TimeSpan.Zero).GetAwaiter().GetResult(); + return QueueStatusBridge.RollUp(reader.GetRunSummariesAsync().GetAwaiter().GetResult(), ids); + }; + _conn = Realtime.Connect(User).Conn!; + Realtime.Watch(User, queueId, trackRun: true); + } + + /// Pump once and keep what it sent. + public List Pump() + { + Realtime.PumpRuns(); + var fresh = new List(); + while (_conn.Reader.TryRead(out var frame)) + fresh.Add(JsonDocument.Parse(frame.Split("data: ")[1]).RootElement.Clone()); + Frames.AddRange(fresh); + return fresh; + } + + public IEnumerable Ends => Frames.Where(f => f.GetProperty("mode").GetString() == "end"); + + public void Dispose() => Realtime.Dispose(); + } + + private static int Completed(JsonElement f) => f.GetProperty("data").GetProperty("completed").GetInt32() + f.GetProperty("data").GetProperty("failed").GetInt32(); + + /// The end frame alone must be enough for the tracker: every task listed, all finished. + private static void AssertTaskList(JsonElement end, int tasks, int failed = 0) + { + var list = end.GetProperty("data").GetProperty("tasks").EnumerateArray().ToList(); + Assert.Equal(tasks, list.Count); + Assert.Equal(failed, list.Count(t => t.GetProperty("status").GetString() == "Failed")); + Assert.All(list, t => Assert.True(t.GetProperty("status").GetString() is "Completed" or "Failed")); + } + + /// Drive the work and the pump together; at every "end" assert that all are finished. + private static async Task DriveWatching(OrchestrationHarness h, Watcher w, string[] runs, int timeoutMs = 20_000) + { + var ok = await h.DriveUntil(async () => + { + foreach (var _ in w.Pump().Where(f => f.GetProperty("mode").GetString() == "end")) + foreach (var run in runs) + Assert.True(await h.Store.GetRunByNameAsync(run) is { IsFinished: true }, + $"'end' was sent while {run} was not finished"); + return w.Ends.Any(); + }, timeoutMs); + Assert.True(ok, "the watched queue never ended"); + w.Pump(); + } + + [Fact] + public async Task FanOut_ReportsProgressThenOneEnd() + { + await using var h = await CreateAsync(poolSize: 2); + var q = Guid.NewGuid().ToString(); + var gate = new SemaphoreSlim(0); + h.Svc.BeforeRun = _ => gate.WaitAsync(); + using var w = new Watcher(h, q); + + QueueStatusBridge.RegisterQueueMetadata(q, "Sync users", "/sync", ""); + Assert.True(await h.Start($"FanOut-{q}", Batch(12))); + // Let the tasks through a few at a time so the pump sees the run part-way. + for (var released = 0; released < 12; released += 3) + { + gate.Release(3); + await h.DriveUntil(() => Task.FromResult(h.Svc.Tasks.Count >= released + 3), 5_000); + await h.DriveUntil(() => Task.FromResult(false), 100); + w.Pump(); + } + await DriveWatching(h, w, [$"FanOut-{q}"]); + + var end = Assert.Single(w.Ends); + Assert.Equal("Completed", end.GetProperty("data").GetProperty("status").GetString()); + Assert.Equal(12, end.GetProperty("data").GetProperty("total").GetInt32()); + Assert.Equal(12, Completed(end)); + AssertTaskList(end, 12); + Assert.Equal("Sync users", end.GetProperty("data").GetProperty("label").GetString()); + Assert.Equal($"FanOut-{q}", end.GetProperty("data").GetProperty("runName").GetString()); + var progress = w.Frames.Where(f => f.GetProperty("mode").GetString() == "update").Select(Completed).ToList(); + Assert.True(progress.Count >= 2, $"expected progress updates, got [{string.Join(',', progress)}]"); + Assert.Equal(progress.OrderBy(c => c), progress); + Assert.Equal(w.Frames[^1], end); + } + + [Fact] + public async Task NestedChildRunsQueuedLate_HoldTheEndUntilTheLastOneFinishes() + { + await using var h = await CreateAsync(poolSize: 4); + var q = Guid.NewGuid().ToString(); + string parent = $"Parent-{q}", child = $"Child-{q}", grandchild = $"GrandChild-{q}"; + + // Queue a child the way a PowerShell task does: the parent link is registered at once, the run + // itself only lands once the bridge is drained, after the queuing task has returned. + void QueueLate(string from, string name, int tasks, string prefix) + { + var link = h.Svc.RegisterPendingChild(from, name); + Assert.NotNull(link); + _ = Task.Run(async () => + { + await Task.Delay(400); + await h.Svc.StartFromBatchAsync(name, Batch(tasks, prefix), 4, null, null, CancellationToken.None, + parentRunKey: link!.Value.ParentRunKey, childKey: link.Value.ChildKey); + }); + } + + h.Svc.Body = t => + { + var id = t["TenantFilter"].ToString(); + if (id == "p0") QueueLate(parent, child, 2, "c"); + if (id == "c0") QueueLate(child, grandchild, 5, "g"); + return "{}"; + }; + using var w = new Watcher(h, q); + + Assert.True(await h.Start(parent, Batch(3, "p"))); + await DriveWatching(h, w, [parent, child, grandchild]); + + var end = Assert.Single(w.Ends); + Assert.Equal("Completed", end.GetProperty("data").GetProperty("status").GetString()); + Assert.Equal(10, end.GetProperty("data").GetProperty("total").GetInt32()); + Assert.Equal(10, Completed(end)); + AssertTaskList(end, 10); + } + + [Fact] + public async Task FailedTasksInAChildRun_EndWithErrors() + { + await using var h = await CreateAsync(poolSize: 4); + var q = Guid.NewGuid().ToString(); + string parent = $"Parent-{q}", child = $"Child-{q}"; + h.Svc.Body = t => + { + var id = t["TenantFilter"].ToString()!; + if (id == "p0") + { + var link = h.Svc.RegisterPendingChild(parent, child); + h.Svc.StartFromBatchAsync(child, Batch(3, "c"), 4, null, null, CancellationToken.None, + parentRunKey: link!.Value.ParentRunKey, childKey: link.Value.ChildKey).GetAwaiter().GetResult(); + } + if (id == "c1") throw new InvalidOperationException("boom"); + return "{}"; + }; + using var w = new Watcher(h, q); + + Assert.True(await h.Start(parent, Batch(2, "p"))); + await DriveWatching(h, w, [parent, child]); + + var data = Assert.Single(w.Ends).GetProperty("data"); + Assert.Equal("CompletedWithErrors", data.GetProperty("status").GetString()); + Assert.Equal(5, data.GetProperty("total").GetInt32()); + Assert.Equal(1, data.GetProperty("failed").GetInt32()); + AssertTaskList(Assert.Single(w.Ends), 5, failed: 1); + } +} diff --git a/tests/Craft.Tests/RealtimeWatchTests.cs b/tests/Craft.Tests/RealtimeWatchTests.cs new file mode 100644 index 0000000..09c42cf --- /dev/null +++ b/tests/Craft.Tests/RealtimeWatchTests.cs @@ -0,0 +1,222 @@ +using System.Management.Automation; +using System.Text; +using System.Text.Json; +using Craft.Configuration; +using Craft.Hosting.Endpoints; +using Craft.Realtime; +using Craft.Services; +using Microsoft.AspNetCore.Http; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +public class RealtimeWatchTests +{ + private const string Job = "6f1c2b8e-1d2a-4c1e-9f0a-3b2c1d4e5f60"; + + private static RealtimeService NewService() => + new(new CraftSettings { Realtime = new RealtimeSettings { Enabled = true } }, NullLogger.Instance); + + private static List Drain(RealtimeService.Connection conn) + { + var frames = new List(); + while (conn.Reader.TryRead(out var f)) frames.Add(f); + return frames; + } + + [Fact] + public void Notify_ReachesOnlyGrantedUsers() + { + using var svc = NewService(); + var (_, alice) = svc.Connect("alice@contoso.com"); + var (_, mallory) = svc.Connect("mallory@contoso.com"); + + svc.Notify(Job, "update", null); // no grant yet + svc.Watch("alice@contoso.com", Job, trackRun: false); + svc.Notify(Job, "update", null); + + Assert.Single(Drain(alice!)); + Assert.Empty(Drain(mallory!)); + } + + [Theory] + [InlineData("6f1c2b8e-1d2a-4c1e-9f0a-3b2c1d4e5f60")] + [InlineData("BEC-20261008123000-ab12cd")] + [InlineData("Orchestrator_contoso.com:1")] + public void AnAppsOwnTokenIsAJobId(string jobId) + { + using var svc = NewService(); + var (_, alice) = svc.Connect("alice@contoso.com"); + svc.Watch("alice@contoso.com", jobId, trackRun: false); + + svc.Notify(jobId, "update", null); + + Assert.Contains($"\"jobId\":\"{jobId}\"", Assert.Single(Drain(alice!))); + } + + [Theory] + [InlineData("")] + [InlineData("a b")] + [InlineData("x\0y")] + [InlineData("line\r\ninjected")] + [InlineData("