diff --git a/Runtime/CraftRuntime/Invoke-CraftTask.ps1 b/Runtime/CraftRuntime/Invoke-CraftTask.ps1 index 1881264..275d3b7 100644 --- a/Runtime/CraftRuntime/Invoke-CraftTask.ps1 +++ b/Runtime/CraftRuntime/Invoke-CraftTask.ps1 @@ -23,7 +23,10 @@ function Invoke-CraftTask { $PushFunction = "Push-$FunctionName" $TenantLabel = $Item.TenantFilter ?? $Item.QueueName ?? $Item.defaultDomainName ?? $(if ($Item.Tenant -is [string]) { $Item.Tenant } elseif ($Item.Tenant.defaultDomainName) { $Item.Tenant.defaultDomainName } else { 'unknown' }) - Write-Information "Dispatching task to $PushFunction for $TenantLabel" + # Per-task dispatch/complete are Debug: at scale (and especially during crash recovery, when every + # pending task re-dispatches) two Info lines per task flooded the log. Push-* functions still log their + # own meaningful output at Info; enable Debug to correlate an individual task to its tenant. + Write-Debug "Dispatching task to $PushFunction for $TenantLabel" $Result = & $PushFunction -Item ([PSCustomObject]$Item) @@ -32,5 +35,5 @@ function Invoke-CraftTask { ConvertTo-Json -InputObject @($Result) -Depth 20 -Compress } - Write-Information "Completed task $PushFunction for $TenantLabel" + Write-Debug "Completed task $PushFunction for $TenantLabel" } diff --git a/Runtime/CraftRuntime/Start-CraftOrchestrator.ps1 b/Runtime/CraftRuntime/Start-CraftOrchestrator.ps1 index 0efd7f7..7e7daac 100644 --- a/Runtime/CraftRuntime/Start-CraftOrchestrator.ps1 +++ b/Runtime/CraftRuntime/Start-CraftOrchestrator.ps1 @@ -24,8 +24,15 @@ function Start-CraftOrchestrator { - FunctionName (string) — Push-{FunctionName} is called with aggregated results - Parameters (object) — extra parameters forwarded to the post-exec function - SkipLog (bool) — optional; suppress run logging + - Sequential (bool) — optional; run the batch ONE TASK AT A TIME in payload order + (Durable-style sequencing) instead of fanning out in parallel. + The whole run is PINNED to a single worker: it starts on one + worker and runs every step on it to completion, without going + back to the pool between steps. A step that fails is recorded and + the run carries on with the next (best-effort). .EXAMPLE + # Fan-out (default): every task is queued up front and drained in parallel by the worker pool. Start-CraftOrchestrator -InputObject @{ OrchestratorName = 'MyDataCollection' Batch = @( @@ -35,6 +42,21 @@ function Start-CraftOrchestrator { PostExecution = @{ FunctionName = 'AggregateResults' } } + .EXAMPLE + # Sequential: run ordered steps ONE AT A TIME, in payload order (Durable-style), all on ONE pinned + # worker. Each step starts only after the previous one finishes, so a step that must follow another + # (here ConvertToShared after the mailbox grants that precede it) cannot race it. Once the run starts + # it keeps the same worker until every step is done. + Start-CraftOrchestrator -InputObject @{ + OrchestratorName = 'Offboard-jdoe@contoso.com' + Sequential = $true + Batch = @( + @{ FunctionName = 'RevokeSessions'; User = 'jdoe@contoso.com' } + @{ FunctionName = 'GrantMailboxAccess'; User = 'jdoe@contoso.com'; Delegate = 'manager@contoso.com' } + @{ FunctionName = 'ConvertToShared'; User = 'jdoe@contoso.com' } + ) + } + .FUNCTIONALITY Internal #> @@ -120,7 +142,11 @@ function Start-CraftOrchestrator { # a parent run would finalize (and dispatch its PostExecution) before its child runs complete. $ParentRunName = $OpContext.RunName - Write-Information "Craft: Queuing orchestrator '$OrchestratorName' ($TaskCount tasks, P$Priority$(if ($PostExecFunctionName) { ", PostExec: $PostExecFunctionName" })$(if ($ParentRunName) { ", Parent: $ParentRunName" }))" + # Sequential mode: PowerShell marshals absent/$false to $false. When set, the orchestrator queues the + # batch one task at a time in payload order rather than fanning out. + $Sequential = [bool]($InputObject.Sequential) + + Write-Information "Craft: Queuing orchestrator '$OrchestratorName' ($TaskCount tasks, P$Priority$(if ($Sequential) { ', Sequential' })$(if ($PostExecFunctionName) { ", PostExec: $PostExecFunctionName" })$(if ($ParentRunName) { ", Parent: $ParentRunName" }))" [Craft.Services.OrchestratorBridge]::QueueOrchestrationFromFile( $OrchestratorName, $BatchPath, @@ -128,7 +154,8 @@ function Start-CraftOrchestrator { $PostExecFunctionName, $PostExecParametersJson, $InputObject.Reference, - $ParentRunName + $ParentRunName, + $Sequential ) return "Craft-$OrchestratorName" } diff --git a/Services/Bridges/OrchestratorBridge.cs b/Services/Bridges/OrchestratorBridge.cs index b75978b..17ccc6d 100644 --- a/Services/Bridges/OrchestratorBridge.cs +++ b/Services/Bridges/OrchestratorBridge.cs @@ -25,7 +25,7 @@ public static class OrchestratorBridge public static void QueueOrchestration(string name, string batchJson, int priority, string? postExecFunctionName = null, string? postExecParametersJson = null, - string? reference = null, string? parentRunName = null) + string? reference = null, string? parentRunName = null, bool sequential = false) { // Sanitized here as well as at run creation so the child-run registration below // records the SAME name the service ends up creating — a raw name with a table-illegal @@ -35,7 +35,7 @@ public static void QueueOrchestration(string name, string batchJson, int priorit var gated = RegisterPendingChild(parentRunName, name); s_pending.Enqueue(new PendingOrchestration(name, batchJson, priority, postExecFunctionName, postExecParametersJson, parentRunName, reference, - PendingChildRegistered: gated)); + PendingChildRegistered: gated, Sequential: sequential)); } /// @@ -51,14 +51,14 @@ public static void QueueOrchestration(string name, string batchJson, int priorit /// public static void QueueOrchestrationFromFile(string name, string batchFilePath, int priority, string? postExecFunctionName = null, string? postExecParametersJson = null, - string? reference = null, string? parentRunName = null) + string? reference = null, string? parentRunName = null, bool sequential = false) { name = TableKeys.Sanitize(name); parentRunName = ResolveParentRunName(name, parentRunName); var gated = RegisterPendingChild(parentRunName, name); s_pending.Enqueue(new PendingOrchestration(name, string.Empty, priority, postExecFunctionName, postExecParametersJson, parentRunName, reference, batchFilePath, - PendingChildRegistered: gated)); + PendingChildRegistered: gated, Sequential: sequential)); } /// @@ -112,7 +112,7 @@ public static void DrainPending() if (s_service == null) { DiscardUndispatchable(p); continue; } s_service.StartFromBatchAsync(p.Name, p.BatchJson, p.Priority, p.PostExecFunctionName, p.PostExecParametersJson, CancellationToken.None, - p.ParentRunName, p.Reference, p.BatchFilePath) + p.ParentRunName, p.Reference, p.BatchFilePath, p.Sequential) .GetAwaiter().GetResult(); } catch (Exception ex) @@ -144,7 +144,7 @@ public static async Task DrainPendingAsync() if (s_service == null) { DiscardUndispatchable(p); continue; } await s_service.StartFromBatchAsync(p.Name, p.BatchJson, p.Priority, p.PostExecFunctionName, p.PostExecParametersJson, CancellationToken.None, - p.ParentRunName, p.Reference, p.BatchFilePath); + p.ParentRunName, p.Reference, p.BatchFilePath, p.Sequential); } catch (Exception ex) { @@ -187,7 +187,7 @@ private static void DiscardUndispatchable(PendingOrchestration p) public record PendingOrchestration(string Name, string BatchJson, int Priority, string? PostExecFunctionName, string? PostExecParametersJson, string? ParentRunName, string? Reference = null, string? BatchFilePath = null, - bool PendingChildRegistered = false); + bool PendingChildRegistered = false, bool Sequential = false); private static readonly ConcurrentQueue s_pendingPlanners = new(); diff --git a/Services/Configuration/OrchestratorSettings.cs b/Services/Configuration/OrchestratorSettings.cs index 437c14b..8c97e36 100644 --- a/Services/Configuration/OrchestratorSettings.cs +++ b/Services/Configuration/OrchestratorSettings.cs @@ -103,4 +103,28 @@ public class OrchestratorSettings /// periodic sweep; the startup pass, which follows crash recovery, still runs. /// public int CleanupIntervalHours { get; set; } = 4; + + /// + /// The per-run status/re-drive tick cadence (seconds, default 60). Each live run has a timer firing at + /// this interval to log status, re-drive orphaned tasks, and re-check completion. It is also the floor + /// of the re-drive backoff. Raising it cuts per-run overhead at high live-run counts; the perf harness + /// lowers it to exercise the backoff in compressed time. Minimum 1s. + /// + public int StatusTimerIntervalSeconds { get; set; } = 60; + + /// + /// Whether the per-run re-drive backs off geometrically once it has verified a run has no orphaned + /// tasks (default true). When false the re-drive verifies against storage on every tick — the old + /// behaviour, retained as a safety switch and for A/B measurement of the backoff's effect. + /// + public bool RedriveBackoff { get; set; } = true; + + /// + /// Whether a task sheds its Parameters payload from the in-memory run graph once it is durably + /// persisted and enqueued, rehydrating it from the Tasks table at dispatch (default true). This bounds + /// the retained memory of a large pending backlog — thousands of runs each holding every task's payload + /// is what drives the live-set toward the GC heap ceiling — at the cost of one point read per task at + /// dispatch. False keeps the payload resident the whole time (the old behaviour), for A/B or safety. + /// + public bool ShedPendingParameters { get; set; } = true; } diff --git a/Services/Hosting/StatsHistoryService.cs b/Services/Hosting/StatsHistoryService.cs index 15b6f9c..738c511 100644 --- a/Services/Hosting/StatsHistoryService.cs +++ b/Services/Hosting/StatsHistoryService.cs @@ -61,32 +61,40 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken) // Load persisted history from disk LoadFromDisk(); - // Wait for worker pool to be ready before starting collection - await Task.Delay(TimeSpan.FromSeconds(10), stoppingToken); + try + { + // Wait for worker pool to be ready before starting collection + await Task.Delay(TimeSpan.FromSeconds(10), stoppingToken); - _logger.LogInformation("[StatsHistory] Started — sampling every {Interval}s, retaining {Days} days", - SampleIntervalSeconds, RetentionDays); + _logger.LogInformation("[StatsHistory] Started — sampling every {Interval}s, retaining {Days} days", + SampleIntervalSeconds, RetentionDays); - using var timer = new PeriodicTimer(TimeSpan.FromSeconds(SampleIntervalSeconds)); + using var timer = new PeriodicTimer(TimeSpan.FromSeconds(SampleIntervalSeconds)); - while (await timer.WaitForNextTickAsync(stoppingToken)) - { - try + while (await timer.WaitForNextTickAsync(stoppingToken)) { - var point = CollectSample(); - AppendPointToDisk(point); + try + { + var point = CollectSample(); + AppendPointToDisk(point); - _ticksSinceCompact++; - if (_ticksSinceCompact >= CompactEveryNTicks) + _ticksSinceCompact++; + if (_ticksSinceCompact >= CompactEveryNTicks) + { + CompactFile(); + _ticksSinceCompact = 0; + } + } + catch (Exception ex) { - CompactFile(); - _ticksSinceCompact = 0; + _logger.LogWarning(ex, "[StatsHistory] Sample collection failed"); } } - catch (Exception ex) - { - _logger.LogWarning(ex, "[StatsHistory] Sample collection failed"); - } + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + // Host shutting down (the Task.Delay or the timer await was cancelled). Expected — + // swallow it so we don't fault the host under BackgroundServiceExceptionBehavior.StopHost. } // Final compaction on shutdown so retention is applied before we exit diff --git a/Services/Orchestration/OrchestratorRun.cs b/Services/Orchestration/OrchestratorRun.cs index 851ab92..9808b28 100644 --- a/Services/Orchestration/OrchestratorRun.cs +++ b/Services/Orchestration/OrchestratorRun.cs @@ -28,4 +28,13 @@ public class OrchestratorRun public int PostExecAttemptCount { get; set; } public string? ParentRunName { get; set; } + + /// + /// Sequential execution mode. When true the run's tasks are dispatched ONE AT A TIME, in ascending + /// (batch) order: only the current task is ever enqueued, + /// and the next is enqueued when it reaches a terminal state. Runs on any free worker (no pinning) — + /// the durable queue simply never holds more than one of this run's tasks at once. The default (false) + /// is the fan-out behaviour: every task is enqueued up front and drained in parallel by the pool. + /// + public bool Sequential { get; set; } } diff --git a/Services/Orchestration/OrchestratorService.cs b/Services/Orchestration/OrchestratorService.cs index 28ce4d9..3d714b2 100644 --- a/Services/Orchestration/OrchestratorService.cs +++ b/Services/Orchestration/OrchestratorService.cs @@ -3,6 +3,7 @@ using System.Text.Json; using System.Text.Json.Serialization; using Craft.Configuration; +using Craft.PowerShellHost; using Craft.Services; using Craft.Storage; @@ -59,9 +60,47 @@ public class OrchestratorService : IJobDescriptorStateWriter /// drain runs in the background while the enqueuing task is marked terminal immediately. /// private readonly ConcurrentDictionary _pendingChildRuns = new(); - private readonly ConcurrentDictionary _runStatusTimers = new(); private readonly ConcurrentDictionary _cancelledRuns = new(); + /// + /// Last status line emitted per run — the (completed, failed, running, pending) tuple and when. Lets + /// skip re-emitting an identical line every 60s for a run that has not + /// changed (the dominant log volume at scale — thousands of runs parked at "0 running / N pending"), + /// while a slow heartbeat still proves a long-lived run is alive. Dropped at finalize. + /// + private readonly ConcurrentDictionary _lastStatusLog = new(); + + /// + /// Per-run re-drive backoff: when the storage verification in may + /// next run, and the interval it grew to. The re-drive is a watchdog for the rare orphaned-Pending task; + /// in steady state it reads storage and finds nothing, so once it does it backs off geometrically instead + /// of paying a full index read (+ a point read per candidate) on every 60s tick for every live run — the + /// dominant per-tick storage cost at scale. Snaps back to the base interval the moment it finds an orphan. + /// Dropped at finalize. + /// + private readonly ConcurrentDictionary _redriveBackoff = new(); + + /// Per-run status/re-drive tick cadence, from Orchestrator:StatusTimerIntervalSeconds. + private readonly TimeSpan _statusInterval; + + /// Re-drive backoff floor — one tick. The interval grows from here to . + private readonly TimeSpan _redriveBase; + + /// Whether the re-drive backoff is active, from Orchestrator:RedriveBackoff. + private readonly bool _redriveBackoffEnabled; + + /// Whether pending tasks shed their Parameters payload, from Orchestrator:ShedPendingParameters. + private readonly bool _shedParameters; + + private static long _redriveStorageReads; + + /// + /// Count of storage verifications the re-drive has performed (a + /// call: one index-partition read + a point read per candidate). Instrumentation for the perf harness — + /// the backoff's whole purpose is to hold this down at high live-run counts. + /// + public static long RedriveStorageReads => Interlocked.Read(ref _redriveStorageReads); + /// /// Resolved task-script path per run. One entry per RUN (not per task), so a 738-task fan-out costs /// one string reference instead of 738 captured ones. Populated at dispatch, dropped at finalize. @@ -75,6 +114,12 @@ public class OrchestratorService : IJobDescriptorStateWriter /// private const int MaxPostExecAttempts = 3; + /// + /// How often an UNCHANGED run still emits a status line, so a long-lived run proves it is alive without + /// logging the identical line on every 60s tick. A real status change always logs immediately. + /// + private static readonly TimeSpan StatusHeartbeat = TimeSpan.FromMinutes(10); + /// /// Runs whose finalize has been claimed, so it happens once. Claimed in CheckRunCompletion, /// released on the deferral/failure paths there, and cleared in DispatchPendingTasksAsync when a @@ -83,6 +128,15 @@ public class OrchestratorService : IJobDescriptorStateWriter /// private readonly ConcurrentDictionary _finalizingRuns = new(); + /// + /// Sequential runs whose single driver job is currently executing. A sequential run's steps all run on + /// one pinned worker inside one driver, and the not-yet-run steps deliberately have no queue row — so + /// while a driver is active the re-drive must not treat those steps as orphaned and enqueue them (which + /// would spawn a second driver). Set when the driver starts, cleared when it finishes; in-memory only, + /// so after a crash it is empty and the re-drive/resume correctly re-triggers the driver. + /// + private readonly ConcurrentDictionary _activeSequentialDrivers = new(); + /// /// Get the Reference for a given run name, or null if not found/no reference set. /// @@ -127,6 +181,14 @@ public OrchestratorService( _writer = writer; _settings = settings; + // Status/re-drive tick cadence. Configurable so a constrained deployment can slow it (fewer ticks = + // less per-run overhead at high live-run counts) and so the perf harness can compress it to exercise + // the re-drive backoff quickly. The backoff floor is one tick. + _statusInterval = TimeSpan.FromSeconds(Math.Max(1, _settings.Orchestrator.StatusTimerIntervalSeconds)); + _redriveBase = _statusInterval; + _redriveBackoffEnabled = _settings.Orchestrator.RedriveBackoff; + _shedParameters = _settings.Orchestrator.ShedPendingParameters; + // The queue holds descriptors; this is how they become work again at dispatch time, and how // operator changes to a queued task are made durable. _jobManager.SetWorkResolver(ResolveTaskWorkAsync); @@ -144,6 +206,14 @@ public void PriorityChanged(JobDescriptor descriptor, int newPriority) { if (!TryFindLive(descriptor, out _, out var task)) return; + // Rehydrate a shed payload before the durable write below (Replace mode) overwrites the stored + // ParametersJson with null. This is a rare, operator-initiated path, so the blocking read is fine. + if (_shedParameters && task.Parameters == null) + { + var p = _store.GetTaskParametersAsync(descriptor.RunName, task.Id).GetAwaiter().GetResult() ?? []; + lock (_lock) { task.Parameters ??= p; } + } + lock (_lock) task.Priority = newPriority; _writer.QueueTask(descriptor.RunName, task); } @@ -416,6 +486,12 @@ public async Task ResumeInterruptedRunsAsync(CancellationToken ct) if (reattached > 0) _logger.LogInformation("[Scheduler] Reattached {Count} in-flight child runs to their parents", reattached); + // Recovery emits ONE aggregate line, not a handful per run: at scale a crash-loop replayed thousands + // of per-run "Found/Resuming/Released/Dispatched" lines on every restart. Per-run detail is kept at + // Debug; the counts below carry the summary. Genuine problems (unresumable, post-exec abandoned) still + // log at their own level as they happen. + var resumed = 0; var pendingTotal = 0; var postExecResumed = 0; + var staleReleased = 0; var finalizedNow = 0; var unresumable = 0; var postExecGaveUp = 0; foreach (var runName in summaries.Select(s => s.Name)) { try @@ -444,19 +520,21 @@ public async Task ResumeInterruptedRunsAsync(CancellationToken ct) await _store.UpsertRunAsync(run); await _store.CleanupRunAsync(run.Name); await _queue.RemoveRunAsync(run.Name, ct); + postExecGaveUp++; continue; } - _logger.LogInformation( + _logger.LogDebug( "[Scheduler] Resuming interrupted PostExecution for run: {Name} (PostExecStatus={Status}, attempt {Attempt}/{Max})", run.Name, run.PostExecStatus, run.PostExecAttemptCount + 1, MaxPostExecAttempts); DispatchPostExecution(run); + postExecResumed++; continue; } if (run.Status != "Running") continue; - _logger.LogInformation("[Scheduler] Found interrupted run: {Name}", run.Name); + _logger.LogDebug("[Scheduler] Found interrupted run: {Name}", run.Name); // Use the stored task script name, fall back to naming convention var taskPath = !string.IsNullOrEmpty(run.TaskScriptName) @@ -466,6 +544,7 @@ public async Task ResumeInterruptedRunsAsync(CancellationToken ct) { _logger.LogWarning("[Scheduler] Cannot resume {Name}: task script not found (tried {Script})", run.Name, run.TaskScriptName ?? $"Invoke-{run.Name}Task"); + unresumable++; continue; } @@ -507,9 +586,12 @@ public async Task ResumeInterruptedRunsAsync(CancellationToken ct) { var released = await _queue.ReleaseRunClaimsAsync(run.Name, ct); if (released > 0) - _logger.LogInformation( + { + staleReleased += released; + _logger.LogDebug( "[Scheduler] Released {Count} stale claim(s) held by the previous process for {Name}", released, run.Name); + } } catch (Exception ex) { @@ -517,12 +599,14 @@ public async Task ResumeInterruptedRunsAsync(CancellationToken ct) _logger.LogWarning(ex, "[Scheduler] Could not release stale claims for {Name}", run.Name); } - _logger.LogInformation("[Scheduler] Resuming interrupted run {Name}: {Pending} pending", run.Name, pending); - await DispatchPendingTasksAsync(run, taskPath, run.Priority, ct); + _logger.LogDebug("[Scheduler] Resuming interrupted run {Name}: {Pending} pending", run.Name, pending); + resumed++; pendingTotal += pending; + await DispatchPendingTasksAsync(run, taskPath, run.Priority, ct, quiet: true); } else { await FinalizeRunAsync(run); + finalizedNow++; } } catch (Exception ex) @@ -538,6 +622,15 @@ public async Task ResumeInterruptedRunsAsync(CancellationToken ct) } } + if (resumed + finalizedNow + postExecResumed + unresumable + postExecGaveUp > 0) + _logger.LogInformation( + "[Scheduler] Crash recovery: resumed {Resumed} run(s) ({Pending} pending tasks re-dispatched), " + + "{PostExec} post-execution(s), {Finalized} finalized, {Stale} stale claim(s) released" + + "{Unresumable}{GaveUp}", + resumed, pendingTotal, postExecResumed, finalizedNow, staleReleased, + unresumable > 0 ? $", {unresumable} unresumable" : "", + postExecGaveUp > 0 ? $", {postExecGaveUp} post-exec abandoned" : ""); + // First retention pass, now that every run that could be resumed is back in _activeRuns and so // exempt from the abandoned-run rule. The scheduler keeps it going on an interval from here. try @@ -620,6 +713,46 @@ public async Task RunRetentionLoopAsync(CancellationToken ct) } } + /// + /// One loop that ticks every live run's status / re-drive / completion-recheck, at + /// . It replaces the per-run that used to do this: at + /// high live-run counts that meant one Timer object per run and a continuous stream of fire-and-forget + /// callbacks onto the thread pool (~M/interval per second), whereas one sweep over + /// is a single scheduling source that allocates nothing per run. The per-run work stays cheap — + /// skips an unchanged line, is gated by + /// its backoff and returns synchronously when backed off, and + /// short-circuits — so ticking thousands of runs in one pass is fast. A throw for one run is logged and + /// neither stops the sweep nor takes the host down (a Timer callback that threw would have crashed it). + /// Started once by , alongside the retention loop. + /// + public async Task RunStatusSweepLoopAsync(CancellationToken ct) + { + using var timer = new PeriodicTimer(_statusInterval); + try + { + while (await timer.WaitForNextTickAsync(ct)) + { + foreach (var run in _activeRuns.Values) + { + try + { + LogRunStatus(run); + RedrivePendingTasks(run); + lock (_lock) { CheckRunCompletion(run); } + } + catch (Exception ex) + { + _logger?.LogWarning(ex, "[Scheduler] Run status tick failed for {Name}", run?.Name); + } + } + } + } + catch (OperationCanceledException) + { + // Host shutdown. + } + } + /// /// Start an orchestrator run from a pre-built batch. /// Called by OrchestratorBridge.DrainPending() when PowerShell's Start-CIPPOrchestrator @@ -634,7 +767,8 @@ public async Task RunRetentionLoopAsync(CancellationToken ct) /// public async Task StartFromBatchAsync(string name, string batchJson, int priority, string? postExecFunctionName, string? postExecParametersJson, CancellationToken ct, - string? parentRunName = null, string? reference = null, string? batchFilePath = null) + string? parentRunName = null, string? reference = null, string? batchFilePath = null, + bool sequential = false) { // The batch file is this method's to dispose of, on EVERY path — including the two // "already running, skipping" returns below, which never look at it. Those are the common @@ -643,7 +777,7 @@ public async Task StartFromBatchAsync(string name, string batchJson, int priorit try { await StartFromBatchCoreAsync(name, batchJson, priority, postExecFunctionName, - postExecParametersJson, parentRunName, reference, batchFilePath, ct); + postExecParametersJson, parentRunName, reference, batchFilePath, sequential, ct); } finally { @@ -660,7 +794,7 @@ await StartFromBatchCoreAsync(name, batchJson, priority, postExecFunctionName, private async Task StartFromBatchCoreAsync(string name, string batchJson, int priority, string? postExecFunctionName, string? postExecParametersJson, - string? parentRunName, string? reference, string? batchFilePath, CancellationToken ct) + string? parentRunName, string? reference, string? batchFilePath, bool sequential, CancellationToken ct) { // Run names become PartitionKeys verbatim, and batch names carry user-typed task names // ("Alert on Entra ID P1/P2 …"). An illegal key character 400s every write for the run — @@ -722,6 +856,10 @@ private async Task StartFromBatchCoreAsync(string name, string batchJson, int pr return; } + // Stamp payload order. Sequential dispatch reads this to run tasks one at a time in the order + // they were submitted; fan-out ignores it. The parser preserves batch order, so the index is it. + for (var i = 0; i < tasks.Count; i++) tasks[i].Sequence = i; + var genericTaskFunc = _settings.Orchestrator.GenericTaskFunction; if (string.IsNullOrEmpty(genericTaskFunc)) { @@ -746,7 +884,8 @@ private async Task StartFromBatchCoreAsync(string name, string batchJson, int pr TaskScriptName = genericTaskFunc, PostExecFunctionName = postExecFunctionName, PostExecParametersJson = postExecParametersJson, - ParentRunName = parentRunName + ParentRunName = parentRunName, + Sequential = sequential }; await _store.UpsertRunAsync(run); @@ -772,7 +911,8 @@ private async Task StartFromBatchCoreAsync(string name, string batchJson, int pr } } - private async Task DispatchPendingTasksAsync(OrchestratorRun run, string taskPath, int priority, CancellationToken ct) + private async Task DispatchPendingTasksAsync(OrchestratorRun run, string taskPath, int priority, + CancellationToken ct, bool quiet = false) { _activeRuns.TryAdd(run.Name, run); // Registered before anything is enqueued — the resolver reads it on the dispatch side. @@ -783,39 +923,12 @@ private async Task DispatchPendingTasksAsync(OrchestratorRun run, string taskPat // this, the finalize claim from their previous outing would strand the next one forever. _finalizingRuns.TryRemove(run.Name, out _); - // Start periodic status timer (every 60s) for this run - if (!_runStatusTimers.ContainsKey(run.Name)) - { - // CheckRunCompletion is re-run here on purpose. It is normally driven by task transitions, - // but a run whose finalize was deferred because storage still showed work outstanding has no - // transitions left to retrigger it - without this periodic re-check that deferral would be - // permanent, which is a worse failure than the premature finalize it exists to prevent. - var timer = new Timer(_ => - { - // A System.Threading.Timer callback that throws crashes the process. This periodic - // maintenance tick must never take the host down on a transient error — a dependency - // disposed during shutdown, a race on run state — so it logs and waits for the next tick. - try - { - LogRunStatus(run); - RedrivePendingTasks(run); - lock (_lock) { CheckRunCompletion(run); } - } - catch (Exception ex) - { - _logger?.LogWarning(ex, "[Scheduler] Run status tick failed for {Name}", run?.Name); - } - }, - null, TimeSpan.FromSeconds(60), TimeSpan.FromSeconds(60)); - if (!_runStatusTimers.TryAdd(run.Name, timer)) - { - // Lost the ContainsKey→TryAdd race (concurrent dispatch of the same run — startup - // resume vs a scheduler tick). An active periodic Timer is rooted by the runtime's - // timer queue, so an undisposed loser would fire — and pin this run graph through its - // closure — for the process lifetime. - timer.Dispose(); - } - } + // No per-run timer: the single RunStatusSweepLoopAsync ticks every run in _activeRuns (which this + // run was just added to). One sweep replaces what used to be a System.Threading.Timer per run — at + // high live-run counts that was thousands of Timer objects and a continuous drizzle of fire-and-forget + // callbacks onto the thread pool. The sweep still runs LogRunStatus + RedrivePendingTasks + + // CheckRunCompletion for each run at the same cadence (the last re-run on purpose: a finalize deferred + // while storage showed work outstanding has no task transition left to retrigger it). var pending = run.Tasks.Where(t => t.Status == "Pending").ToList(); @@ -848,7 +961,41 @@ private async Task DispatchPendingTasksAsync(OrchestratorRun run, string taskPat alreadyQueued = []; } - var toQueue = pending.Where(t => !alreadyQueued.Contains(t.Id)).ToList(); + List toQueue; + if (run.Sequential) + { + // Sequential mode: a single entry row starts the ONE pinned driver that then runs every step of + // the run inline (see BuildSequentialRunWork), so exactly ONE queue row ever exists for the run — + // the lowest-Sequence task still Pending — and the later steps get no rows at all. Enqueue that + // one task, and only if it does not already have a row. A second row for the same run would let a + // second driver start and break the pinning/ordering. This covers both first dispatch (enqueues + // Sequence 0 to start the driver) and resume (re-enqueues the current step only if its row is + // gone, so a fresh driver picks the run back up where it left off). + var current = pending.OrderBy(t => t.Sequence).FirstOrDefault(); + toQueue = current != null && !alreadyQueued.Contains(current.Id) ? [current] : []; + } + else + { + toQueue = pending.Where(t => !alreadyQueued.Contains(t.Id)).ToList(); + } + + // Shed the Parameters payload BEFORE the tasks become claimable. The caller has already persisted + // them (UpsertTaskBatchAsync), so the Tasks table is authoritative; the live graph keeps each task + // object (its identity + Status drive completion tracking) but drops the payload it does not need + // while it waits, and BuildTaskWork rehydrates it from storage at dispatch. This is what bounds the + // retained memory of a large pending backlog — thousands of runs each holding every task's payload is + // what walks the live-set into the GC heap ceiling. Ordering matters: shedding AFTER the enqueue + // could null a task the pump had already claimed and whose BuildTaskWork had just rehydrated it, so + // shed here, before EnqueueBatchAsync makes the rows claimable. Only toQueue (the rows enqueued in + // THIS call) is touched — a task already queued from a prior call may be mid-dispatch. Sequential + // runs are exempt: they are small (a handful of ordered steps) and their pinned driver holds every + // step's payload in the live graph while it runs, so shedding would buy no memory and only add reads. + if (_shedParameters && !run.Sequential) + { + lock (_lock) + foreach (var t in toQueue) + if (t.Status == "Pending") t.Parameters = null!; + } // One batched write per priority bucket rather than one per task. The queue is the backlog now; // the JobManager only ever sees the batch JobQueuePump claims from it. @@ -856,24 +1003,244 @@ await _queue.EnqueueBatchAsync(run.Name, toQueue.Select(t => (t.Id, t.Priority ?? priority)).ToList(), DateTime.UtcNow, ct); + // quiet = called from crash recovery, where a per-run line per resumed run is the flood the + // aggregate summary replaces — drop to Debug. A normal orchestration start logs it at Info (one line). + var level = quiet ? LogLevel.Debug : LogLevel.Information; if (toQueue.Count == pending.Count) { - _logger.LogInformation("[Scheduler] Dispatched {Count} tasks for {Name} at P{Priority}", + _logger.Log(level, "[Scheduler] Dispatched {Count} tasks for {Name} at P{Priority}", toQueue.Count, run.Name, priority); } else { - _logger.LogInformation( + _logger.Log(level, "[Scheduler] Dispatched {Count} tasks for {Name} at P{Priority} ({Existing} already queued)", toQueue.Count, run.Name, priority, pending.Count - toQueue.Count); } } /// - /// Enqueue one task BY IDENTITY. The JobManager holds only (runName, taskId, priority) — no run - /// graph, no task, no script path, no service reference — and calls back into - /// at dispatch time to rebuild the work. + /// Build the work for a SEQUENTIAL run: a single delegate that checks out ONE background worker and + /// runs every step of the run on it, in payload (Sequence) order, one at a time, then reclaims the + /// worker once at the end. This is what "pin one worker for the whole run" means — the run starts on a + /// worker and stays there until it finishes, never going back to the pool between steps to be + /// re-scheduled onto a different one. The run is dispatched as a SINGLE queue row (the entry task; see + /// DispatchPendingTasksAsync), so exactly one JobManager slot and one worker are held for the run's + /// whole duration and the mid-run steps never get their own rows. + /// + /// Failure policy is best-effort: a step that throws is recorded Failed and the driver moves on to the + /// next, so one bad step cannot strand the rest and the run ends CompletedWithErrors. Each step runs + /// through InvokeAsync, whose finally resets the runspace after success, failure OR cancellation, so + /// steps stay isolated on the shared worker. Registered in for + /// its whole life so the re-drive leaves the not-yet-reached (deliberately row-less) steps alone while + /// the driver is progressing. + /// + private Func BuildSequentialRunWork(OrchestratorRun run, string taskPath) + { + return async (jobCt) => + { + // One driver per run. Registered for the driver's whole life so the re-drive does not mistake + // the run's row-less pending steps for orphans and enqueue them. TryAdd (not an unconditional + // set) also closes the duplicate-entry-row race: if two rows for the same run are claimed before + // either driver marks the entry step Running, the loser here does nothing and the winner runs + // every step. Because this guard is OUTSIDE the try/finally, the loser never runs the finally + // that would otherwise remove the winner's registration. + if (!_activeSequentialDrivers.TryAdd(run.Name, true)) + { + _logger.LogDebug( + "[Scheduler] Sequential run {Run} already has an active driver — skipping duplicate entry", run.Name); + return; + } + PowerShellWorker? worker = null; + var workerFaulted = false; + try + { + // One worker for the whole run: checked out once here, reclaimed once in the finally. + worker = CheckoutSequentialWorker(jobCt); + + while (true) + { + OrchestratorTaskItem? task; + lock (_lock) + { + task = run.Tasks.Where(t => t.Status == "Pending") + .OrderBy(t => t.Sequence) + .FirstOrDefault(); + } + if (task == null) break; // every step is terminal — the run is done + + // Run cancelled while we were working through it: mark this and every remaining step + // Cancelled and stop. + if (_cancelledRuns.ContainsKey(run.Name)) + { + CancelRemainingSequentialTasks(run); + break; + } + + // Rehydrate a shed payload. Sequential runs are exempt from shedding (see + // DispatchPendingTasksAsync), so this is normally a no-op — kept for a resumed run whose + // in-memory payload was dropped. A missing row means the durable state was lost: fail + // that step closed rather than run it blank, then carry on with the next. + if (_shedParameters && task.Parameters == null) + { + var rehydrated = await _store.GetTaskParametersAsync(run.Name, task.Id, jobCt); + if (rehydrated == null) + { + FailTaskTerminally(run, task, + "Parameters could not be rehydrated at dispatch — the Tasks-table row is missing. " + + "The task's payload was shed from memory and storage no longer has it."); + continue; + } + lock (_lock) { task.Parameters ??= rehydrated; } + } + + lock (_lock) { task.Status = "Running"; } + // Durable "Running" marker — awaited before the invoke, same as the parallel path. + try + { + await _writer.MarkRunningAsync(run.Name, task, jobCt); + } + catch (MarkerNotPersistedException ex) + { + // The marker never landed, so storage still has this step Pending. Unlike the + // parallel path we cannot just give a slot back and let a re-drive retry the one + // task — we own the worker and the whole run — and running later steps out of order + // is not allowed. Put the step back to Pending and STOP the driver. The entry job + // then completes with pending work left and no driver active, so the re-drive + // re-enqueues the current step and a fresh driver resumes the run here. + lock (_lock) { task.Status = "Pending"; } + _logger.LogWarning(ex, + "[Scheduler] Sequential run {Run}: could not persist the Running marker for {Task} — " + + "leaving the run for the re-drive to resume", run.Name, task.Id); + break; + } + + try + { + // Run the step on the pinned worker. Only a PostExecution run needs the output + // captured and stored; otherwise the seam returns empty and nothing is stored. + var output = await RunSequentialStepAsync(run, task, taskPath, worker); + if (!string.IsNullOrEmpty(run.PostExecFunctionName) + && !string.IsNullOrEmpty(output) + && !_writer.TryQueueResult(run.Name, task.Id, output)) + { + await _store.StoreResultAsync(run.Name, task.Id, output); + } + + lock (_lock) + { + task.Status = "Completed"; + task.CompletedUtc = DateTime.UtcNow; + task.Parameters = null!; + CheckRunCompletion(run); + } + PersistTaskAndRunAsync(run, task); + _logger.LogDebug("[Scheduler] Sequential task completed: {TaskId}", task.Id); + } + catch (OperationCanceledException) when (jobCt.IsCancellationRequested) + { + // App shutting down mid-step. Leave this step Running (its durable marker is written) + // for resume on next startup, and let the exception abort the driver. + _logger.LogInformation( + "[Scheduler] Sequential run {Run} interrupted by shutdown at {TaskId}", run.Name, task.Id); + throw; + } + catch (Exception ex) + { + lock (_lock) + { + task.Status = "Failed"; + task.LastError = ex.Message; + task.CompletedUtc = DateTime.UtcNow; + task.Parameters = null!; + CheckRunCompletion(run); + } + PersistTaskAndRunAsync(run, task); + _logger.LogError(ex, + "[Scheduler] Sequential task failed: {TaskId} — continuing with the next step", task.Id); + // best-effort: fall through to the next Pending step + } + } + } + catch (OperationCanceledException) when (jobCt.IsCancellationRequested) + { + // Shutdown (checkout or a step was cancelled). The pipeline may have been Stop()'d, so treat + // the worker as faulted on reclaim; rethrow so the JobManager marks the entry job Cancelled. + workerFaulted = true; + throw; + } + catch (Exception ex) + { + // An unexpected driver-level failure — not a per-step task error, which is handled inline. + workerFaulted = true; + _logger.LogError(ex, "[Scheduler] Sequential driver for {Run} failed", run.Name); + throw; + } + finally + { + ReclaimSequentialWorker(worker, workerFaulted); + _activeSequentialDrivers.TryRemove(run.Name, out _); + } + }; + } + + // ── sequential-driver seams ─────────────────────────────────────────────────────────────────────── + // The driver's loop logic (ordering, best-effort, cancellation, marker recovery) is exercised by unit + // tests through a subclass that overrides these three methods, so the tests need no PowerShell worker + // pool. Production runs the real pool: one worker checked out for the whole run, each step invoked on + // it, reclaimed once at the end. Keep the checkout/run/reclaim split — the whole point is one checkout + // and one reclaim around many step invocations. + + /// Check out the single worker a sequential run is pinned to. Virtual for tests. + internal virtual PowerShellWorker? CheckoutSequentialWorker(CancellationToken ct) + => _psRunner.CheckoutBackgroundWorker(ct); + + /// Return the pinned worker once the run is done. No-op for a null worker (checkout failed). + /// Virtual for tests. + internal virtual void ReclaimSequentialWorker(PowerShellWorker? worker, bool faulted) + { + if (worker != null) _psRunner.ReclaimBackgroundWorker(worker, faulted: faulted); + } + + /// Run one sequential step on the pinned worker and return its captured output (empty when the + /// run has no PostExecution and so needs no result). Virtual for tests. InvokeAsync resets the runspace + /// afterwards, so successive steps stay isolated on the shared worker. + internal virtual async Task RunSequentialStepAsync( + OrchestratorRun run, OrchestratorTaskItem task, string taskPath, PowerShellWorker? worker) + { + var parameters = new Dictionary + { + { "TaskJson", JsonSerializer.Serialize(task.Parameters, s_jsonOptions) } + }; + if (!string.IsNullOrEmpty(run.PostExecFunctionName)) + return await _psRunner.ExecuteScriptWithOutput(taskPath, parameters, pinnedWorker: worker); + await _psRunner.ExecuteScript(taskPath, parameters, pinnedWorker: worker); + return string.Empty; + } + + /// + /// Mark every still-Pending or Running step of a cancelled sequential run Cancelled, in one lock, and + /// persist them. Called by the driver when it notices the run was cancelled between steps. /// + private void CancelRemainingSequentialTasks(OrchestratorRun run) + { + List remaining; + lock (_lock) + { + remaining = run.Tasks.Where(t => t.Status is "Pending" or "Running").ToList(); + foreach (var t in remaining) + { + t.Status = "Cancelled"; + t.LastError = "Cancelled by user"; + t.CompletedUtc = DateTime.UtcNow; + t.Parameters = null!; + } + CheckRunCompletion(run); + } + foreach (var t in remaining) _writer.QueueTask(run.Name, t); + _writer.QueueRun(run); + } + /// /// Put one task back on the durable queue. Fire-and-forget because every caller is on a lock or a /// timer callback, and a failure is recoverable: the task is still Pending in storage, so the next @@ -950,13 +1317,6 @@ private void FailTaskTerminally(OrchestratorRun run, OrchestratorTaskItem task, /// private async Task?> ResolveTaskWorkAsync(JobDescriptor descriptor, CancellationToken ct) { - if (!_taskScriptPaths.TryGetValue(descriptor.RunName, out var taskPath)) - { - _logger.LogWarning("[Orchestrator] No task script path known for run {Run} — dropping {Task}", - descriptor.RunName, descriptor.TaskId); - return null; - } - OrchestratorRun? run; OrchestratorTaskItem? task; @@ -1005,6 +1365,35 @@ private void FailTaskTerminally(OrchestratorRun run, OrchestratorTaskItem task, return null; } + // Resolve the task script (a PowerShell function name) for this run. The steady-state path finds + // it in _taskScriptPaths, cached by DispatchPendingTasksAsync when the run was dispatched. A MISS + // is NOT a reason to drop the task. The pump is a BackgroundService that begins claiming persisted + // queue rows at host start, whereas the only writer of _taskScriptPaths does not run for a resumed + // run until ResumeInterruptedRunsAsync reaches it — and that waits on the worker pool first — and + // the rehydration branch above re-adds a run to _activeRuns without a cached path either. In both + // windows the run record still carries TaskScriptName, so rebuild the path from it exactly as the + // resume path does (ScriptRepository is fully loaded before the pump's first claim) and cache it, + // so sibling tasks of the same run cost nothing. Only an empty TaskScriptName with no + // naming-convention match is genuinely unrunnable. Returning null on the miss instead lets the + // JobManager mark the job Skipped and the pump delete its queue row, permanently dropping a task + // that is still Pending in the run — with no queue row left, nothing re-dispatches it. + if (!_taskScriptPaths.TryGetValue(run.Name, out var taskPath)) + { + taskPath = !string.IsNullOrEmpty(run.TaskScriptName) + ? _psRunner.FindScript(run.TaskScriptName) + : FindTaskScript(run.Name); + + if (string.IsNullOrEmpty(taskPath)) + { + _logger.LogWarning( + "[Orchestrator] No task script for run {Run} (TaskScriptName={Script}) — dropping {Task}", + run.Name, run.TaskScriptName, descriptor.TaskId); + return null; + } + + _taskScriptPaths[run.Name] = taskPath; + } + // Already terminal (e.g. cancelled, or completed by a previous attempt while queued), or already // executing. // @@ -1020,6 +1409,14 @@ private void FailTaskTerminally(OrchestratorRun run, OrchestratorTaskItem task, return null; } + // Sequential runs are dispatched as a single entry row; that one claim drives the WHOLE run on one + // pinned worker (BuildSequentialRunWork ignores which step this descriptor named and works through + // every Pending step in order). The "Running" guard above already stops a duplicate row from + // starting a second driver once the entry step is marked Running, and the driver's own TryAdd closes + // the remaining pre-mark race. + if (run.Sequential) + return BuildSequentialRunWork(run, taskPath); + return BuildTaskWork(run, task, taskPath); } @@ -1028,6 +1425,31 @@ private Func BuildTaskWork(OrchestratorRun run, Orchest return async (jobCt) => { + // Rehydrate the Parameters payload shed while this task waited in the backlog. Done BEFORE + // any status write — MarkRunningAsync snapshots Parameters and every task-row write is + // Replace, so a null payload here would overwrite the stored one. One point read, only for a + // task actually being dispatched. The ??= keeps a value another dispatch attempt already set. + if (_shedParameters && task.Parameters == null) + { + // GetTaskParametersAsync returns null only when the Tasks row itself is GONE (a + // present-but-empty payload comes back as an empty dictionary). A missing row means this + // task's durable state was lost out from under a live run — the shed dropped the in-memory + // copy on the promise that storage still had it. Running now would invoke the task with NO + // parameters (its FunctionName and inputs both live in the payload), which for a real task + // is worse than not running it. Fail closed instead of executing blank. Before shedding + // the payload was resident, so a deleted row could not affect an in-flight dispatch; this + // guard restores that safety for the one case shedding introduced. + var rehydrated = await _store.GetTaskParametersAsync(run.Name, task.Id, jobCt); + if (rehydrated == null) + { + FailTaskTerminally(run, task, + "Parameters could not be rehydrated at dispatch — the Tasks-table row is missing. " + + "The task's payload was shed from memory and storage no longer has it."); + return; + } + lock (_lock) { task.Parameters ??= rehydrated; } + } + // Check if run was cancelled while this job was queued if (_cancelledRuns.ContainsKey(run.Name)) { @@ -1146,6 +1568,7 @@ private Func BuildTaskWork(OrchestratorRun run, Orchest CheckRunCompletion(run); } PersistTaskAndRunAsync(run, task); + _logger.LogError(ex, "[Scheduler] Task failed: {TaskId}", task.Id); throw; // Let JobManager also track the failure } @@ -1244,6 +1667,14 @@ private void DeferTask(OrchestratorRun run, OrchestratorTaskItem task, Exception /// private static readonly TimeSpan RedriveAge = TimeSpan.FromMinutes(5); + /// + /// Re-drive backoff bounds. The first verification for a run runs at the status-timer cadence; each time + /// it confirms nothing orphaned the interval doubles up to , so a run stuck for + /// hours costs a handful of index reads rather than one per minute. The cap bounds how long a genuinely + /// orphaned task can wait to be caught (worst case ~RedriveMax), which the watchdog trades for the cost. + /// + private static readonly TimeSpan RedriveMax = TimeSpan.FromMinutes(15); + /// /// Re-queue tasks that are Pending in memory but that nothing owns — no queued job, no running job. /// @@ -1280,6 +1711,14 @@ private void DeferTask(OrchestratorRun run, OrchestratorTaskItem task, Exception private async Task RedrivePendingTasksAsync(OrchestratorRun run) { var now = DateTime.UtcNow; + + // Backoff gate. Once the verification below has confirmed a run has nothing orphaned, it need not run + // again for a while: orphaning is caused by specific rare events (a removed/expired queue row, a + // crash/migration), not something that spontaneously arises every 60s. Skipping here avoids the whole + // tick body — the candidates scan AND the storage read — for a run that verified clean, which for a + // large stuck backlog is nearly every run on nearly every tick. + if (_redriveBackoffEnabled && _redriveBackoff.TryGetValue(run.Name, out var st) && now < st.NextUtc) return; + List candidates; lock (_lock) @@ -1290,15 +1729,55 @@ private async Task RedrivePendingTasksAsync(OrchestratorRun run) .Where(t => !_deferrals.TryGetValue(DeferralKey(run.Name, t.Id), out var s) || now - s.LastUtc >= RedriveAge) .ToList(); + + // Sequential run: one pinned driver runs every step, so the ONLY queue row that ever exists is + // the entry row that started the driver — the not-yet-reached steps deliberately have none. + // Applying the generic "Pending with no queue row = orphaned" rule to them would re-drive them + // all and spawn a second driver. The run is progressing whenever a driver is registered for it + // OR its entry job is still queued/running (that job's identity is one of this run's tasks) — + // re-drive nothing in either case. The driver registration closes the gap the entry-job check + // alone leaves open between steps (no task queued, none marked Running for an instant). Only when + // neither holds is the driver truly gone (never started, or died with the process): re-enqueue + // ONLY the current step (lowest Sequence still Pending) so a fresh driver resumes the run. + if (run.Sequential) + { + var driverActive = _activeSequentialDrivers.ContainsKey(run.Name) + || run.Tasks.Any(t => _jobManager.IsQueuedOrRunning($"{run.Name}-{t.Id}")); + if (driverActive) + { + candidates.Clear(); + } + else + { + var next = candidates.OrderBy(t => t.Sequence).FirstOrDefault(); + candidates = next != null ? [next] : []; + } + } } - if (candidates.Count == 0) return; + if (candidates.Count == 0) + { + // No candidates to verify (all Pending tasks are queued/running, or none are Pending). Don't grow + // the backoff — a run mid-drain legitimately produces no candidates and should stay responsive — + // just clear any prior backoff so the next real candidate is checked promptly. + _redriveBackoff.TryRemove(run.Name, out _); + return; + } - // Storage decides. A Pending task that still has a row is waiting its turn, not orphaned. - HashSet stillQueued; + // Storage decides — but the queue TABLE decides, not the index. Asking the index (the old + // GetQueuedTaskIdsAsync here) reports a task queued whenever its index row exists, and an index + // row can outlive the queue row it points at. Such a task is invisible to the pump yet looks + // "queued" to this check, so it is never re-driven and its run stalls indefinitely with the task + // Pending — this watchdog keeps ticking and finds nothing orphaned. GetDispatchableTaskIdsAsync + // verifies each candidate against the queue table (one point read apiece; the candidate set is + // small), returning only tasks the pump can actually still claim. Anything else is a ghost to + // re-enqueue. + HashSet dispatchable; try { - stillQueued = await _queue.GetQueuedTaskIdsAsync(run.Name); + Interlocked.Increment(ref _redriveStorageReads); + dispatchable = await _queue.GetDispatchableTaskIdsAsync( + run.Name, candidates.Select(t => t.Id).ToList()); } catch (Exception ex) { @@ -1308,9 +1787,25 @@ private async Task RedrivePendingTasksAsync(OrchestratorRun run) return; } - var orphaned = candidates.Where(t => !stillQueued.Contains(t.Id)).ToList(); - if (orphaned.Count == 0) return; + var orphaned = candidates.Where(t => !dispatchable.Contains(t.Id)).ToList(); + if (orphaned.Count == 0) + { + // Verified clean: grow the interval (double, capped) so this run's next storage read is further + // out. A run stuck for hours thus costs O(log) reads, not one per minute. + if (_redriveBackoffEnabled) + { + var next = _redriveBackoff.TryGetValue(run.Name, out var cur) + ? TimeSpan.FromTicks(Math.Min(cur.Interval.Ticks * 2, RedriveMax.Ticks)) + : _redriveBase; + _redriveBackoff[run.Name] = (now + next, next); + } + return; + } + // Found orphans — something is wrong with this run's queue rows, so snap back to close watch and + // re-drive them. + if (_redriveBackoffEnabled) + _redriveBackoff[run.Name] = (now + _redriveBase, _redriveBase); foreach (var task in orphaned) { // Clear the exhausted counter, or DeferTask would abandon it again on its first attempt. @@ -1319,27 +1814,50 @@ private async Task RedrivePendingTasksAsync(OrchestratorRun run) } _logger.LogWarning( - "[Scheduler] Re-drove {Count} orphaned Pending task(s) in {Run} — no queue row and not queued or running", + "[Scheduler] Re-drove {Count} orphaned Pending task(s) in {Run} — no runnable queue row and not queued or running", orphaned.Count, run.Name); } private void LogRunStatus(OrchestratorRun run) { - var elapsed = DateTime.UtcNow - run.StartedUtc; - int completed, failed, running, pending; + // Nothing consumes this Info line at a higher level, and the flood of them is itself a measured + // cost, so do no work at all when Info is disabled. + if (!_logger.IsEnabled(LogLevel.Information)) return; + + // One pass, not four Count(predicate) calls. Enumerable.Count over the List boxes an enumerator per + // call, and this runs on every run's 60s timer — four boxed enumerators × M runs per minute. + int completed = 0, failed = 0, running = 0, pending = 0; lock (_lock) { - completed = run.Tasks.Count(t => t.Status == "Completed"); - failed = run.Tasks.Count(t => t.Status == "Failed"); - running = run.Tasks.Count(t => t.Status == "Running"); - pending = run.Tasks.Count(t => t.Status == "Pending"); + foreach (var t in run.Tasks) + { + switch (t.Status) + { + case "Completed": completed++; break; + case "Failed": failed++; break; + case "Running": running++; break; + case "Pending": pending++; break; + } + } } - var memSnapshot = BackgroundTaskLimiter.GetMemorySnapshot(); + + // Skip the line — and the string format, the nine boxed args, and the memory snapshot it needs — + // when nothing has changed since the last tick. A run parked at "0 running / N pending" for hours + // re-emitted the identical line every 60s (M of them per minute at scale). Log on a real transition, + // plus a slow heartbeat so a long-lived run still shows it is alive. + var now = DateTime.UtcNow; + if (_lastStatusLog.TryGetValue(run.Name, out var prev) + && prev.C == completed && prev.F == failed && prev.R == running && prev.P == pending + && now - prev.LoggedUtc < StatusHeartbeat) + return; + _lastStatusLog[run.Name] = (completed, failed, running, pending, now); + + var elapsed = now - run.StartedUtc; _logger.LogInformation( "[Scheduler] Run {Name} T+{Elapsed:F1}min: {Completed}/{Total} done {Running} running {Pending} pending {Failed} failed jobs={Active}a/{Queued}q {Memory}", run.Name, elapsed.TotalMinutes, completed, run.Tasks.Count, running, pending, failed, _jobManager.ActiveCount, _jobManager.QueuedCount, - memSnapshot); + BackgroundTaskLimiter.GetMemorySnapshot()); } private void CheckRunCompletion(OrchestratorRun run) @@ -1524,11 +2042,12 @@ private async Task FinalizeRunCoreAsync(OrchestratorRun run) _writer.QueueRun(run); await _writer.FlushAsync(); + // Removed from _activeRuns first, so the status sweep stops ticking it before its per-run maps go. _activeRuns.TryRemove(run.Name, out _); _cancelledRuns.TryRemove(run.Name, out _); _taskScriptPaths.TryRemove(run.Name, out _); - _runStatusTimers.TryRemove(run.Name, out var timer); - timer?.Dispose(); + _lastStatusLog.TryRemove(run.Name, out _); + _redriveBackoff.TryRemove(run.Name, out _); _finalizeDeferrals.TryRemove(run.Name, out _); // Deferral and re-queue tracking is keyed per task and nothing else removes entries for tasks // that ended without passing through their happy-path cleanup — without this sweep the residue diff --git a/Services/Orchestration/OrchestratorStatusWriter.cs b/Services/Orchestration/OrchestratorStatusWriter.cs index 99008be..1208e8b 100644 --- a/Services/Orchestration/OrchestratorStatusWriter.cs +++ b/Services/Orchestration/OrchestratorStatusWriter.cs @@ -75,7 +75,7 @@ public OrchestratorStatusWriter(OrchestratorTableStore store, ILogger run + "" + task; private static TaskStatusWrite Snap(string run, OrchestratorTaskItem t) => new( run, t.Id, t.Status, JsonSerializer.Serialize(t.Parameters, s_json), t.AttemptCount, t.LastError, - t.CompletedUtc, t.Priority); + t.CompletedUtc, t.Priority, t.Sequence); /// Persist the pre-invoke "Running" marker durably before the task runs. Under the barrier it is /// batched with other concurrently-starting tasks (N tasks → ~1 transaction) yet still lands before the diff --git a/Services/Orchestration/OrchestratorTaskItem.cs b/Services/Orchestration/OrchestratorTaskItem.cs index 6eae8fc..504595c 100644 --- a/Services/Orchestration/OrchestratorTaskItem.cs +++ b/Services/Orchestration/OrchestratorTaskItem.cs @@ -19,4 +19,11 @@ public class OrchestratorTaskItem /// priority when ResumeInterruptedRunsAsync re-queues the task. /// public int? Priority { get; set; } + + /// + /// Position of this task in the batch as submitted (0-based). Only meaningful for a run marked + /// , where tasks are dispatched one at a time in ascending + /// Sequence order. Non-sequential runs leave it 0 and ignore it. + /// + public int Sequence { get; set; } } diff --git a/Services/Orchestration/SchedulerService.cs b/Services/Orchestration/SchedulerService.cs index d85d952..a8e84d6 100644 --- a/Services/Orchestration/SchedulerService.cs +++ b/Services/Orchestration/SchedulerService.cs @@ -98,6 +98,10 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken) // forget is deliberate: the loop handles its own failures and ends with the stopping token. _ = _orchestrator.RunRetentionLoopAsync(stoppingToken); + // One status/re-drive sweep over all live runs, replacing the per-run timers. Same fire-and-forget + // contract: it handles its own per-run errors and ends with the stopping token. + _ = _orchestrator.RunStatusSweepLoopAsync(stoppingToken); + while (!stoppingToken.IsCancellationRequested) { var now = DateTimeOffset.UtcNow; diff --git a/Services/PowerShellHost/PowerShellRunnerService.cs b/Services/PowerShellHost/PowerShellRunnerService.cs index 6863f25..f162254 100644 --- a/Services/PowerShellHost/PowerShellRunnerService.cs +++ b/Services/PowerShellHost/PowerShellRunnerService.cs @@ -452,16 +452,19 @@ private void DrainBridgesInBackground() /// /// Execute a script by function name (for scheduler / orchestrator). No HTTP context needed. - /// Runs on the background pool. + /// Runs on the background pool. When is supplied the script runs on + /// THAT worker and it is NOT returned to the pool here — the caller owns its lifecycle (see + /// for the same contract and why a sequential run uses it). /// - public async Task ExecuteScript(string functionName, Dictionary? parameters = null) + public async Task ExecuteScript(string functionName, Dictionary? parameters = null, + PowerShellWorker? pinnedWorker = null) { var prof = DispatchProfiler.Enabled; var totalStart = prof ? Stopwatch.GetTimestamp() : 0; long checkoutTicks = 0, invokeTicks = 0; var checkoutStart = prof ? Stopwatch.GetTimestamp() : 0; var sw = Stopwatch.StartNew(); - var worker = _pool.CheckoutBackground(CancellationToken.None); + var worker = pinnedWorker ?? _pool.CheckoutBackground(CancellationToken.None); if (prof) checkoutTicks = Stopwatch.GetTimestamp() - checkoutStart; // Set invocation context — inherits RunName and Priority from parent OperationContext if set @@ -570,7 +573,9 @@ public async Task ExecuteScript(string functionName, Dictionary? if (onInfo != null) worker.Streams.Information.DataAdded -= onInfo; if (onDebug != null) worker.Streams.Debug.DataAdded -= onDebug; if (onVerbose != null) worker.Streams.Verbose.DataAdded -= onVerbose; - _pool.Reclaim(worker, isHttp: false, faulted: exceptionOccurred); + // A pinned worker is owned by the caller (sequential driver) — it reclaims once, after the + // whole run. Only reclaim here when we checked the worker out ourselves. + if (pinnedWorker == null) _pool.Reclaim(worker, isHttp: false, faulted: exceptionOccurred); } // BG dispatch profiling (checkout + invoke + total; marshal/extract N/A for a script call). @@ -586,14 +591,22 @@ public async Task ExecuteScript(string functionName, Dictionary? /// Execute a script on the background pool and capture its output stream as a string. /// Used by OrchestratorService for planner scripts that return JSON task lists. /// - public async Task ExecuteScriptWithOutput(string functionName, Dictionary? parameters = null) + /// + /// Run one script and return its output. When is supplied the script + /// runs on THAT worker and it is NOT returned to the pool here — the caller owns its lifecycle. This is + /// how a sequential orchestration keeps one worker for its whole run: check a worker out once, invoke + /// each step on it (InvokeAsync still resets the runspace per step, so steps stay isolated), reclaim + /// once at the end. Passing null preserves the original checkout-per-call, reclaim-in-finally behaviour. + /// + public async Task ExecuteScriptWithOutput(string functionName, Dictionary? parameters = null, + PowerShellWorker? pinnedWorker = null) { var prof = DispatchProfiler.Enabled; var totalStart = prof ? Stopwatch.GetTimestamp() : 0; long checkoutTicks = 0, invokeTicks = 0; var checkoutStart = prof ? Stopwatch.GetTimestamp() : 0; var sw = Stopwatch.StartNew(); - var worker = _pool.CheckoutBackground(CancellationToken.None); + var worker = pinnedWorker ?? _pool.CheckoutBackground(CancellationToken.None); if (prof) checkoutTicks = Stopwatch.GetTimestamp() - checkoutStart; // Set invocation context — inherits RunName and Priority from parent OperationContext if set @@ -692,10 +705,19 @@ public async Task ExecuteScriptWithOutput(string functionName, Dictionar if (onInfo != null) worker.Streams.Information.DataAdded -= onInfo; if (onDebug != null) worker.Streams.Debug.DataAdded -= onDebug; if (onVerbose != null) worker.Streams.Verbose.DataAdded -= onVerbose; - _pool.Reclaim(worker, isHttp: false); + // A pinned worker is owned by the caller (sequential driver) — it reclaims once, after the + // whole run. Only reclaim here when we checked the worker out ourselves. + if (pinnedWorker == null) _pool.Reclaim(worker, isHttp: false); } } + /// Check out a background worker to pin across a run's steps. Pair with . + public PowerShellWorker CheckoutBackgroundWorker(CancellationToken ct = default) => _pool.CheckoutBackground(ct); + + /// Return a worker taken with . + public void ReclaimBackgroundWorker(PowerShellWorker worker, bool faulted = false) + => _pool.Reclaim(worker, isHttp: false, faulted: faulted); + /// /// Find a script by command name. Checks ScriptRepository (standalone files) /// first, then falls back to checking if it exists as a module function. diff --git a/Services/Storage/JobQueueStore.cs b/Services/Storage/JobQueueStore.cs index d895f17..0dd5f93 100644 --- a/Services/Storage/JobQueueStore.cs +++ b/Services/Storage/JobQueueStore.cs @@ -728,6 +728,60 @@ public async Task> GetQueuedTaskIdsAsync(string runName, Cancell return ids; } + /// + /// Of (all belonging to ), the ones the pump + /// can still dispatch: they have a queue row that EXISTS and that the claim filter will match — now, + /// because it is free, or later, because it holds a lease that will lapse. + /// + /// This is the queue-table counterpart to , which answers purely + /// from the index. The index is what makes "does this run still have queued work" a single-partition + /// read, but it can OUTLIVE the queue rows it points at, and then it lies: it reports a task queued + /// that no pump will ever run. That divergence is not hypothetical — + /// + /// a removal deletes the index row first, so a crash in between (or a + /// that deleted queue rows before its index partition) can leave the + /// opposite; + /// a run left Pending under a build that dispatched into memory rather than this + /// queue re-enters here with index rows and no queue rows; + /// a row owned with NO LeaseUntil is excluded by + /// forever, so it sits with an index entry advertising it. + /// + /// A task in any of those states is invisible to the pump AND reported "queued" by the index, so the + /// re-drive that trusts the index never re-enqueues it and the run stalls indefinitely with it, its + /// watchdog never firing. Verifying against the queue table costs one point read per id, so the caller + /// passes a SMALL candidate set (the re-drive's aged-Pending tasks), never the whole run. + /// + public async Task> GetDispatchableTaskIdsAsync( + string runName, IReadOnlyCollection taskIds, CancellationToken ct = default) + { + var result = new HashSet(StringComparer.Ordinal); + if (taskIds.Count == 0) return result; + + var wanted = taskIds as HashSet ?? new HashSet(taskIds, StringComparer.Ordinal); + var now = DateTimeOffset.UtcNow; + + // The index carries the bucket + queue row key to address each row directly — one partition read, + // then a point read only for the ids the caller asked about. + foreach (var e in await ReadIndexAsync(runName, ct)) + { + if (!wanted.Contains(e.TaskId)) continue; + + var row = await _store.GetAsync(_queueTable, e.Bucket, e.QueueRowKey, ct); + if (row == null) continue; // index points at a queue row that is gone — a ghost + + // Owned with no lease is what the server-side claim filter cannot match (it is neither + // Owner eq '' nor LeaseUntil lt now), so the pump would never dispatch it however long it + // waits — a ghost as surely as a missing row. A free row, or one under a lease live or + // lapsed, the pump will get. + if (!string.IsNullOrEmpty(row.GetString("Owner")) && row.GetDateTimeOffset("LeaseUntil") == null) + continue; + + result.Add(e.TaskId); + } + + return result; + } + /// /// Bring the queue tables up to , once per storage account. The marker row /// written at the end is checked first, so every later start is a single point read. diff --git a/Services/Storage/OrchestratorTableStore.cs b/Services/Storage/OrchestratorTableStore.cs index b096f41..1e207fc 100644 --- a/Services/Storage/OrchestratorTableStore.cs +++ b/Services/Storage/OrchestratorTableStore.cs @@ -118,7 +118,10 @@ public async Task> WriteRunStatusBatchAsync(IReadOnlyList< ["PostExecAttemptCount"] = run.PostExecAttemptCount, ["Reference"] = run.Reference, ["ParentRunName"] = run.ParentRunName, - ["TaskCount"] = run.Tasks.Count + ["TaskCount"] = run.Tasks.Count, + // Persisted so a resumed sequential run keeps advancing one task at a time (0/1 — StoreRow has + // no bool reader). Absent on older rows reads as 0 = the fan-out default. + ["Sequential"] = run.Sequential ? 1 : 0 } }; @@ -145,7 +148,8 @@ public async Task> WriteRunStatusBatchAsync(IReadOnlyList< // (so FindRunByReference could not see it) and a null ParentRunName (so its finalize // never re-checked the parent). Absent on rows written before this existed. Reference = runRow.GetString("Reference"), - ParentRunName = runRow.GetString("ParentRunName") + ParentRunName = runRow.GetString("ParentRunName"), + Sequential = (runRow.GetInt32("Sequential") ?? 0) == 1 }; var tasks = new List(); @@ -178,6 +182,7 @@ public async Task> WriteRunStatusBatchAsync(IReadOnlyList< LastError = taskRow.GetString("LastError"), // Absent on rows written before per-task priority existed — null means "inherit the run's". Priority = taskRow.GetInt32("Priority"), + Sequence = taskRow.GetInt32("Sequence") ?? 0, CompletedUtc = taskRow.GetDateTimeOffset("CompletedUtc")?.UtcDateTime }); } @@ -186,6 +191,28 @@ public async Task> WriteRunStatusBatchAsync(IReadOnlyList< return run; } + /// + /// Read one task's Parameters from the Tasks table (a single point read + deserialize), matching the + /// deserialization uses. For the pending-Parameters shedding path: the live + /// graph keeps the task object but drops its Parameters payload while it waits, and this rehydrates them + /// at dispatch. Null if the row or its ParametersJson is missing. + /// + public async Task?> GetTaskParametersAsync( + string runName, string taskId, CancellationToken ct = default) + { + var row = await _store.GetAsync(_tasksTable, runName, taskId, ct); + var parametersJson = row?.GetString("ParametersJson"); + if (string.IsNullOrEmpty(parametersJson)) return null; + try + { + return JsonSerializer.Deserialize>(parametersJson, s_jsonOptions) ?? []; + } + catch + { + return []; + } + } + /// List all known run names. public async Task> ListRunsAsync() { @@ -461,6 +488,7 @@ public async Task> WriteTaskStatusBatchAsync(IReadOnlyList ["AttemptCount"] = task.AttemptCount, ["LastError"] = task.LastError, ["Priority"] = task.Priority, + ["Sequence"] = task.Sequence, ["CompletedUtc"] = task.CompletedUtc.HasValue ? new DateTimeOffset(task.CompletedUtc.Value, TimeSpan.Zero) : (DateTimeOffset?)null @@ -476,6 +504,7 @@ public async Task> WriteTaskStatusBatchAsync(IReadOnlyList ["AttemptCount"] = w.AttemptCount, ["LastError"] = w.LastError, ["Priority"] = w.Priority, + ["Sequence"] = w.Sequence, ["CompletedUtc"] = w.CompletedUtc.HasValue ? new DateTimeOffset(w.CompletedUtc.Value, TimeSpan.Zero) : (DateTimeOffset?)null diff --git a/Services/Storage/TaskStatusWrite.cs b/Services/Storage/TaskStatusWrite.cs index de0cf21..6c22012 100644 --- a/Services/Storage/TaskStatusWrite.cs +++ b/Services/Storage/TaskStatusWrite.cs @@ -9,4 +9,4 @@ namespace Craft.Storage; /// time the task moved to Running. /// public record TaskStatusWrite(string RunName, string TaskId, string Status, string? ParametersJson, - int AttemptCount, string? LastError, DateTime? CompletedUtc, int? Priority); + int AttemptCount, string? LastError, DateTime? CompletedUtc, int? Priority, int Sequence = 0); diff --git a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psd1 index 3e0b5f9..95a728f 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-PerfBgEnqueue', 'Push-PerfBg', 'Push-PerfBgLeaf', 'Invoke-ListPerf', 'Invoke-PerfWhoami', 'Invoke-PerfTimerTick', 'Invoke-PerfTimerCount', 'Invoke-PerfPublish', 'Invoke-PerfAllocation', 'Invoke-PerfRuns') + FunctionsToExport = @('Invoke-PerfPing', 'Invoke-PerfEcho', 'Invoke-PerfCpu', 'Invoke-PerfSleep', 'Invoke-PerfJson', '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') 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 c896a50..cdf4f8f 100644 --- a/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 +++ b/perf-harness/api-harness/API/Modules/PerfApi/PerfApi.psm1 @@ -105,6 +105,316 @@ function Push-PerfBgLeaf { return @{ ok = $true; idx = $Item.idx } } +# Run-COUNT axis driver (run-manyruns.ps1). The OOM harness above tests one run of N tasks (fan-out +# WIDTH); this creates MANY separate runs that stay live (fan-out COUNT), reproducing the production shape +# where thousands of scheduled runs sit in _activeRuns — each pinning a per-run 60s status Timer + its +# task graph, and each re-reading the queue index on every tick. Memory is bounded in run-width but scales +# with run-COUNT, which this exercises. +# +# Each run is a small batch of PerfHold tasks that sleep for the whole observation window, so with a tiny +# BG pool only a few are ever claimed and the rest sit Pending — the run never finalizes, so it stays in +# _activeRuns with its timer firing every 60s. Called in chunks by the harness to avoid a long single HTTP +# request. Query: +# runs=M how many runs to create THIS call (default 200) +# tasks=K tasks per run (default 4) — K > (claimable) keeps each run with Pending work forever +# holdms=H per-task sleep so a claimed task never completes during the window (default 3600000 = 1h) +# paramkb=P per-task payload string in KB (default 0) — inflates each OrchestratorTaskItem.Parameters +# so the retained run graph is production-weight (real runs carry tenant/audit data, not no-ops). +# Retained memory then scales as M x K x paramkb, which is the axis that reaches the heap ceiling. +# prefix=P run-name prefix (default PerfHold) so the harness can group a wave +function Invoke-PerfManyRuns { + param($Request, $TriggerMetadata) + $runs = 200; if ($Request.Query.runs) { $runs = [int]$Request.Query.runs } + $tasks = 4; if ($Request.Query.tasks) { $tasks = [int]$Request.Query.tasks } + $holdms = 3600000; if ($Request.Query.holdms) { $holdms = [int]$Request.Query.holdms } + $paramkb = 0; if ($Request.Query.paramkb) { $paramkb = [int]$Request.Query.paramkb } + $func = 'PerfHold'; if ($Request.Query.func) { $func = [string]$Request.Query.func } # PerfHold | PerfCheck + $prefix = 'PerfHold'; if ($Request.Query.prefix) { $prefix = [string]$Request.Query.prefix } + $seq = ([string]$Request.Query.seq -eq 'true') # sequential: one task at a time, payload order + $payload = if ($paramkb -gt 0) { 'x' * ($paramkb * 1024) } else { $null } + + $created = 0 + for ($r = 0; $r -lt $runs; $r++) { + $batch = @(for ($i = 0; $i -lt $tasks; $i++) { + # marker = 'm'+idx lets Push-PerfCheck confirm the payload survived shed→rehydrate at dispatch. + $item = @{ FunctionName = $func; idx = $i; holdms = $holdms; marker = ('m' + $i) } + if ($payload) { $item['payload'] = $payload } + $item + }) + Start-CraftOrchestrator -InputObject @{ + OrchestratorName = "$prefix-$([guid]::NewGuid().ToString('N').Substring(0, 10))" + Batch = $batch + Sequential = $seq + } | Out-Null + $created++ + } + return @{ StatusCode = 200; Body = @{ ok = $true; endpoint = 'PerfManyRuns'; created = $created; tasksPerRun = $tasks; holdms = $holdms; paramkb = $paramkb; sequential = $seq } } +} + +# Table manipulation for failure-mode exploration: delete/inspect orchestrator table rows WHILE runs are +# live, to see whether Craft survives losing state under it. Uses the same storage connection the app uses. +# op=count rows in a partition (needs table[,pk]) +# op=list first rows of a partition (table[,pk]) — RowKeys only +# op=deleteRow delete one entity (table,pk,rk) +# op=deletePart delete every row in a partition (table,pk) — e.g. a run's Tasks or Queue-index partition +# op=tables list table names +# Tables (prefix PerfBgOrch): PerfBgOrchRuns (pk 'Run'), PerfBgOrchTasks (pk=runName, rk=taskId, + Counter), +# PerfBgOrchResults, PerfBgOrchQueue, PerfBgOrchQueueIndex. +function Invoke-PerfTableOp { + param($Request, $TriggerMetadata) + $op = [string]$Request.Query.op; if (-not $op) { $op = 'count' } + $table = [string]$Request.Query.table + $pk = [string]$Request.Query.pk + $rk = [string]$Request.Query.rk + try { + $svc = [Azure.Data.Tables.TableServiceClient]::new($env:AzureWebJobsStorage) + if ($op -eq 'tables') { + $names = @($svc.Query() | ForEach-Object { $_.Name }) + return @{ StatusCode = 200; Body = @{ ok = $true; op = $op; tables = $names } } + } + $tc = $svc.GetTableClient($table) + switch ($op) { + 'deleteRow' { + $tc.DeleteEntity($pk, $rk) | Out-Null + return @{ StatusCode = 200; Body = @{ ok = $true; op = $op; table = $table; deleted = "$pk/$rk" } } + } + 'deletePart' { + $filter = "PartitionKey eq '$pk'" + $n = 0 + foreach ($e in $tc.Query[Azure.Data.Tables.TableEntity]($filter)) { + $tc.DeleteEntity($e.PartitionKey, $e.RowKey) | Out-Null; $n++ + } + return @{ StatusCode = 200; Body = @{ ok = $true; op = $op; table = $table; pk = $pk; deleted = $n } } + } + 'deleteAll' { + $n = 0 + foreach ($e in $tc.Query[Azure.Data.Tables.TableEntity]("PartitionKey gt ''")) { + $tc.DeleteEntity($e.PartitionKey, $e.RowKey) | Out-Null; $n++ + } + return @{ StatusCode = 200; Body = @{ ok = $true; op = $op; table = $table; deleted = $n } } + } + 'list' { + $filter = if ($pk) { "PartitionKey eq '$pk'" } else { "PartitionKey gt ''" } + $rows = @() + foreach ($e in $tc.Query[Azure.Data.Tables.TableEntity]($filter)) { + $rows += @{ pk = $e.PartitionKey; rk = $e.RowKey } + if ($rows.Count -ge 25) { break } + } + return @{ StatusCode = 200; Body = @{ ok = $true; op = $op; table = $table; rows = $rows } } + } + default { + $filter = if ($pk) { "PartitionKey eq '$pk'" } else { "PartitionKey gt ''" } + $n = 0 + foreach ($e in $tc.Query[Azure.Data.Tables.TableEntity]($filter)) { $n++ } + return @{ StatusCode = 200; Body = @{ ok = $true; op = 'count'; table = $table; pk = $pk; count = $n } } + } + } + } catch { + return @{ StatusCode = 500; Body = @{ ok = $false; op = $op; error = "$_" } } + } +} + +# Pre-seed live runs directly into the orchestrator tables, bypassing Start-CraftOrchestrator's per-run +# enqueue cost — so a HIGH-scale thread comparison (per-run timers vs the single sweep) is not gated by how +# fast the batch/planner path can create runs. Writes the Runs rows (Status=Running), each run's Tasks rows +# (Status=Pending) + the "!!run-counter" row, in Azure Table batch transactions. RESTART the container after +# seeding: ResumeInterruptedRunsAsync reads the Running run rows and resumes them into live _activeRuns +# entries (each with a per-run status timer on the old build, or joined to the sweep on the new one). +# Query: runs=M (default 2000), tasks=K (default 1), holdms=H, prefix=P, tableprefix=T (default PerfBgOrch). +function Invoke-PerfSeedRuns { + param($Request, $TriggerMetadata) + $runs = 2000; if ($Request.Query.runs) { $runs = [int]$Request.Query.runs } + $tasks = 1; if ($Request.Query.tasks) { $tasks = [int]$Request.Query.tasks } + $holdms = 3600000; if ($Request.Query.holdms) { $holdms = [int]$Request.Query.holdms } + $prefix = 'seed'; if ($Request.Query.prefix) { $prefix = [string]$Request.Query.prefix } + $tp = 'PerfBgOrch'; if ($Request.Query.tableprefix) { $tp = [string]$Request.Query.tableprefix } + try { + $svc = [Azure.Data.Tables.TableServiceClient]::new($env:AzureWebJobsStorage) + $rc = $svc.GetTableClient("${tp}Runs"); $rc.CreateIfNotExists() | Out-Null + $tc = $svc.GetTableClient("${tp}Tasks"); $tc.CreateIfNotExists() | Out-Null + $now = [DateTimeOffset]::UtcNow + $upsert = [Azure.Data.Tables.TableTransactionActionType]::UpsertReplace + $runBatch = [System.Collections.Generic.List[Azure.Data.Tables.TableTransactionAction]]::new() + $created = 0 + for ($i = 0; $i -lt $runs; $i++) { + $name = "$prefix-$i" + $r = [Azure.Data.Tables.TableEntity]::new('Run', $name) + $r['Status'] = 'Running'; $r['Priority'] = [int]4; $r['StartedUtc'] = $now + $r['TaskScriptName'] = 'Invoke-CraftTask'; $r['TaskCount'] = [int]$tasks; $r['Sequential'] = [int]0 + $runBatch.Add([Azure.Data.Tables.TableTransactionAction]::new($upsert, $r)) + if ($runBatch.Count -eq 100) { $rc.SubmitTransaction($runBatch) | Out-Null; $runBatch.Clear() } + + $taskBatch = [System.Collections.Generic.List[Azure.Data.Tables.TableTransactionAction]]::new() + for ($j = 0; $j -lt $tasks; $j++) { + $t = [Azure.Data.Tables.TableEntity]::new($name, "${name}_t$j") + $t['Status'] = 'Pending' + $t['ParametersJson'] = "{`"FunctionName`":`"PerfHold`",`"idx`":$j,`"holdms`":$holdms}" + $t['AttemptCount'] = [int]0; $t['Sequence'] = [int]$j + $taskBatch.Add([Azure.Data.Tables.TableTransactionAction]::new($upsert, $t)) + } + $cnt = [Azure.Data.Tables.TableEntity]::new($name, '!!run-counter') + $cnt['Remaining'] = [int]$tasks; $cnt['Total'] = [int]$tasks + $taskBatch.Add([Azure.Data.Tables.TableTransactionAction]::new($upsert, $cnt)) + $tc.SubmitTransaction($taskBatch) | Out-Null + $created++ + } + if ($runBatch.Count -gt 0) { $rc.SubmitTransaction($runBatch) | Out-Null } + return @{ StatusCode = 200; Body = @{ ok = $true; seeded = $created; tasksPerRun = $tasks + note = 'restart the container to resume these into live runs' } } + } catch { + return @{ StatusCode = 500; Body = @{ ok = $false; error = "$_" } } + } +} + +# Parameters-integrity task: verifies the payload the run was created with survived the shed→rehydrate round +# trip. Increments a shared 'ok' counter when its marker parameter is present and correct, 'lost' when it is +# missing/empty (payload lost — a shedding race, or the Tasks row was deleted before dispatch so rehydration +# read nothing). Read the tallies via /API/PerfCheckCounts. holdms lets it sit Pending like PerfHold. +function Push-PerfCheck { + param($Item) + $cache = [Craft.Services.PowerShellRunnerService]::GetSharedCache('PerfCheck') + $marker = [string]$Item.marker + $key = if ($marker -and $marker -eq ('m' + $Item.idx)) { 'ok' } else { 'lost' } + # Interlocked-ish: the shared cache is concurrent; a coarse increment is fine for a tally. + $n = 0; if ($cache[$key]) { $n = [int]$cache[$key] } + $cache[$key] = $n + 1 + if ($key -eq 'lost') { + $ln = 0; if ($cache['lostSample']) { $ln = [int]$cache['lostSample'] } + $cache['lastLost'] = "idx=$($Item.idx) marker='$marker'" + } + if ($Item.holdms -and [int]$Item.holdms -gt 0) { Start-Sleep -Milliseconds ([int]$Item.holdms) } + return @{ ok = $true; idx = $Item.idx; check = $key } +} + +function Invoke-PerfCheckCounts { + param($Request, $TriggerMetadata) + $cache = [Craft.Services.PowerShellRunnerService]::GetSharedCache('PerfCheck') + return @{ StatusCode = 200; Body = @{ ok = $true + okCount = [int]$cache['ok']; lostCount = [int]$cache['lost']; lastLost = [string]$cache['lastLost'] } } +} + +# The hold-open task: sleeps holdms so a claimed task never reaches a terminal state during the test, which +# is what keeps its run live in _activeRuns (and its 60s status timer firing). No allocation, no fan-out — +# this axis is about run COUNT, not per-task work. +function Push-PerfHold { + param($Item) + $ms = 3600000; if ($Item.holdms) { $ms = [int]$Item.holdms } + Start-Sleep -Milliseconds $ms + return @{ ok = $true; idx = $Item.idx } +} + +# Sequential-mode probe: records the ORDER tasks start in and the MAX concurrency observed, into a shared +# cache. A sequential run should show order = payload order (0,1,2,...) and maxActive = 1 (one at a time); +# a fan-out run shows interleaved order and maxActive > 1. Read via /API/PerfSeqResult. +function Push-PerfSeq { + param($Item) + $c = [Craft.Services.PowerShellRunnerService]::GetSharedCache('PerfSeq') + $c['order'] = "$($c['order'])$($Item.idx)," + $a = [int]$c['active'] + 1; $c['active'] = $a + if ($a -gt [int]$c['maxActive']) { $c['maxActive'] = $a } + if ($Item.holdms -and [int]$Item.holdms -gt 0) { Start-Sleep -Milliseconds ([int]$Item.holdms) } + $c['active'] = [int]$c['active'] - 1 + return @{ ok = $true; idx = $Item.idx } +} + +# Live OS-thread breakdown via the C# bridge: total thread count and a tally by ThreadState, plus the +# processor count and PS worker-pool size. Answers "where are Craft's threads" at a point in time — +# bounded PS pool + thread-pool workers, and whether anything is growing under load. +function Invoke-PerfThreadBreakdown { + param($Request, $TriggerMetadata) + $b = [Craft.Services.WorkerMetricsBridge]::GetMemoryBreakdown() + $states = @{} + foreach ($k in $b.ThreadStates.Keys) { $states[$k] = $b.ThreadStates[$k] } + return @{ StatusCode = 200; Body = @{ ok = $true + threadCount = $b.ThreadCount; processorCount = $b.ProcessorCount + httpWorkers = $b.HttpWorkers; bgWorkers = $b.BgWorkers; threadStates = $states } } +} + +function Invoke-PerfSeqResult { + param($Request, $TriggerMetadata) + $c = [Craft.Services.PowerShellRunnerService]::GetSharedCache('PerfSeq') + return @{ StatusCode = 200; Body = @{ ok = $true + order = [string]$c['order']; maxActive = [int]$c['maxActive'] } } +} + +# ── Worker-pinning probe ──────────────────────────────────────────────────────────────────────────── +# Proves the pinned sequential driver: every step of a sequential run runs on the SAME worker, and +# concurrently-running sequential runs each pin their OWN worker. Each step records (run, idx, worker) under +# a unique key so distinct-key writes stay concurrency-safe on the synchronized Hashtable. The worker id is +# the per-invoke stamped $global:CraftOperationContext.WorkerId ("W"); the run name is carried on the item +# so grouping never depends on RunName propagation into the task context. Read via /API/PerfSeqWorkerResult. +function Push-PerfSeqWorker { + param($Item) + $ctx = Get-Variable -Name 'CraftOperationContext' -Scope Global -ValueOnly -ErrorAction SilentlyContinue + $worker = if ($ctx -and $ctx.WorkerId) { [string]$ctx.WorkerId } else { 'W?' } + $run = [string]$Item.run + $c = [Craft.Services.PowerShellRunnerService]::GetSharedCache('PerfSeqWorker') + $c["$run|$([int]$Item.idx)"] = "$worker@$([DateTime]::UtcNow.Ticks)" + if ($Item.holdms -and [int]$Item.holdms -gt 0) { Start-Sleep -Milliseconds ([int]$Item.holdms) } + return @{ ok = $true; run = $run; idx = $Item.idx; worker = $worker } +} + +# Start N sequential (or fan-out, with seq=false) runs of K steps each, holding holdms per step so several +# runs are in flight at once — the case that shows each sequential run keeps its own single worker. Each step +# carries its run name. Clears the shared cache first so a run of the harness starts clean. +function Invoke-PerfSeqWorkerEnqueue { + param($Request, $TriggerMetadata) + $runs = 4; if ($Request.Query.runs) { $runs = [int]$Request.Query.runs } + $steps = 5; if ($Request.Query.steps) { $steps = [int]$Request.Query.steps } + $holdms = 500; if ($Request.Query.holdms) { $holdms = [int]$Request.Query.holdms } + $seq = -not ([string]$Request.Query.seq -eq 'false') # default sequential; seq=false → fan-out contrast + $c = [Craft.Services.PowerShellRunnerService]::GetSharedCache('PerfSeqWorker'); $c.Clear() + $names = @() + for ($r = 0; $r -lt $runs; $r++) { + $name = "SeqW$r-$([guid]::NewGuid().ToString('N').Substring(0, 6))" + $names += $name + $batch = @(for ($i = 0; $i -lt $steps; $i++) { @{ FunctionName = 'PerfSeqWorker'; run = $name; idx = $i; holdms = $holdms } }) + Start-CraftOrchestrator -InputObject @{ OrchestratorName = $name; Batch = $batch; Sequential = $seq } | Out-Null + } + return @{ StatusCode = 200; Body = @{ ok = $true; endpoint = 'PerfSeqWorkerEnqueue' + runs = $runs; steps = $steps; holdms = $holdms; sequential = $seq; names = $names } } +} + +function Invoke-PerfSeqWorkerResult { + param($Request, $TriggerMetadata) + $c = [Craft.Services.PowerShellRunnerService]::GetSharedCache('PerfSeqWorker') + $rows = @() + foreach ($k in @($c.Keys)) { + $parts = ([string]$k) -split '\|', 2 + $vp = ([string]$c[$k]) -split '@', 2 + $rows += @{ run = $parts[0]; idx = [int]$parts[1]; worker = $vp[0]; ticks = [long]$vp[1] } + } + return @{ StatusCode = 200; Body = @{ ok = $true; endpoint = 'PerfSeqWorkerResult'; count = $rows.Count; rows = @($rows) } } +} + +# Thread-pool + process-thread telemetry, for the "thread constrained" half of the many-runs harness. The +# per-run timers fire their re-drive as fire-and-forget work onto the .NET thread pool, so PendingWorkItemCount +# climbing (work queued faster than threads drain it) is the thread-starvation signal that inflates the +# client-side wall-time of otherwise-fast table reads. +function Invoke-PerfThreads { + param($Request, $TriggerMetadata) + $maxW = 0; $maxIo = 0; $minW = 0; $minIo = 0; $availW = 0; $availIo = 0 + [System.Threading.ThreadPool]::GetMaxThreads([ref]$maxW, [ref]$maxIo) | Out-Null + [System.Threading.ThreadPool]::GetMinThreads([ref]$minW, [ref]$minIo) | Out-Null + [System.Threading.ThreadPool]::GetAvailableThreads([ref]$availW, [ref]$availIo) | Out-Null + $proc = [System.Diagnostics.Process]::GetCurrentProcess() + return @{ StatusCode = 200; Body = @{ ok = $true; endpoint = 'PerfThreads' + threadPool = @{ + threadCount = [System.Threading.ThreadPool]::ThreadCount + pendingWorkItems = [long][System.Threading.ThreadPool]::PendingWorkItemCount + completedWorkItems = [long][System.Threading.ThreadPool]::CompletedWorkItemCount + maxWorker = $maxW; maxIo = $maxIo + minWorker = $minW; minIo = $minIo + busyWorker = ($maxW - $availW); busyIo = ($maxIo - $availIo) + } + process = @{ + osThreadCount = $proc.Threads.Count + # Cumulative re-drive storage verifications (index read + point reads). ②'s backoff holds this down. + redriveReads = [long][Craft.Orchestration.OrchestratorService]::RedriveStorageReads + } + } } +} + # Worker/queue allocation snapshot — the harness's downstream wrapper around the CRAFT bridge, standing # in for what a real app (e.g. CIPP) does: CRAFT exposes the data as [Craft.Services.WorkerMetricsBridge], # the app wraps whichever fields it wants into its own endpoint. Returns the shape run-orch.ps1 and the diff --git a/perf-harness/docker-compose.bg.yml b/perf-harness/docker-compose.bg.yml index 0b92b7a..4e429f6 100644 --- a/perf-harness/docker-compose.bg.yml +++ b/perf-harness/docker-compose.bg.yml @@ -58,8 +58,16 @@ services: - App__Orchestrator__BatchStatusWrites=${BATCH_WRITES:-true} - App__Orchestrator__DurableRunningBarrier=${DURABLE_BARRIER:-true} - App__Orchestrator__StatusFlushIntervalMs=${FLUSH_MS:-25} - - CRAFT_LOG_LEVEL=Warning - - Logging__LogLevel__Default=Warning + # Per-run status/re-drive tick cadence. run-manyruns.ps1 lowers it to compress the re-drive backoff. + - App__Orchestrator__StatusTimerIntervalSeconds=${STATUS_INTERVAL:-60} + # Re-drive backoff (② ). run-manyruns.ps1 flips it to A/B the backoff's effect at a fixed interval. + - App__Orchestrator__RedriveBackoff=${REDRIVE_BACKOFF:-true} + # Pending-Parameters shedding (retained-memory fix). run-manyruns.ps1 flips it to A/B the memory effect. + - App__Orchestrator__ShedPendingParameters=${SHED_PARAMS:-true} + # Default Warning keeps the other bg-harness runs quiet; run-manyruns.ps1 sets LOG_LEVEL=Information + # to reproduce production's Info-level per-run status logging (the log flood is part of the cost under test). + - CRAFT_LOG_LEVEL=${LOG_LEVEL:-Warning} + - Logging__LogLevel__Default=${LOG_LEVEL:-Warning} - CRAFT_DISPATCH_TIMING=${DISPATCH_TIMING:-} # Reused pipeline thread default on; run-bg.ps1 -NoReuseThread sets false for A/B (must be a valid bool). - App__Worker__ReuseRunspaceThread=${REUSE_THREAD:-true} diff --git a/perf-harness/scripts/run-manyruns.ps1 b/perf-harness/scripts/run-manyruns.ps1 new file mode 100644 index 0000000..9c65cb2 --- /dev/null +++ b/perf-harness/scripts/run-manyruns.ps1 @@ -0,0 +1,221 @@ +<# +.SYNOPSIS + Run-COUNT axis harness: prove (or, after the fix, disprove) that the per-run status/re-drive tick + storm makes memory + GC + thread-pool pressure scale with the number of concurrently-live runs. + +.DESCRIPTION + oom-analysis.md proved memory is bounded on the fan-out WIDTH axis: one run of 20,000 tasks peaks at + ~64 MB because the backlog lives in the Queue table and the JobManager holds an O(batch) buffer. This + harness exercises the axis that was never tested — the number of concurrently-live RUNS. + + Each live run pins a per-run System.Threading.Timer that fires every 60s (OrchestratorService), and on + every tick it: (a) LogRunStatus — walks the run's task list 4x and logs a status line, unconditionally; + (b) RedrivePendingTasksAsync — reads the queue index + a point-read per pending candidate, every tick, + even when nothing is orphaned. With M live runs that is ~M/60 callbacks/second of fire-and-forget work + on the thread pool plus M pinned run graphs. This reproduces the production shape (thousands of runs + parked at "0 running / N pending" for hours) that climbs to the GC heap ceiling and OOM-crashes (exit + 139), where the fast (~8 ms) index reads show up as multi-second [HTTP-SLOW] purely because the thread + is frozen by continuous full GCs / starved of the thread pool. + + Flow: bring CRAFT up (Http+Background + Azurite) at Information log level (so the per-run status flood is + part of the measured cost) under an optional GC heap hard limit and a constrained CPU/thread budget; + create M runs of K hold-open tasks (they sleep the whole window so no run finalizes → M stays live); + then OBSERVE for -WatchSec across several 60s tick waves, sampling heap / gc2 / thread-pool depth, and + finally scrape the container log for [HTTP-SLOW], tick failures, and (if it died) the exit code. + +.EXAMPLE + # Baseline repro against the current image — climb to the heap ceiling and crash. + pwsh scripts\run-manyruns.ps1 -Runs 2500 -TasksPerRun 4 -HeapLimitMB 400 -Cpus 1 -Label baseline + + # After the fix (rebuild craft:local first) — same M, heap should stay flat, no HTTP-SLOW, no crash. + pwsh scripts\run-manyruns.ps1 -Runs 2500 -TasksPerRun 4 -HeapLimitMB 400 -Cpus 1 -Label fixed +#> +[CmdletBinding()] +param( + [string]$SutImage = 'craft:local', + [string]$Label = 'manyruns', + [int]$Runs = 2500, # M — concurrently-live runs (the independent variable) + [int]$TasksPerRun = 4, # K — tasks per run; K > claimable keeps each run permanently Pending + [int]$ParamKB = 0, # per-task payload KB — makes the retained run graph production-weight + [int]$HoldMs = 3600000, # per-task sleep (1h) — longer than the test, so no run finalizes + [int]$ChunkRuns = 200, # runs created per PerfManyRuns HTTP call (avoids a long single request) + [int]$BgPool = 4, + [double]$Cpus = 1, # constrain CPU → small thread pool (the "thread constrained" half) + [int]$HeapLimitMB = 400, # 0 = unconstrained; set just above the live-graph baseline to force the ceiling + [string]$LogLevel = 'Information', + [int]$StatusIntervalSec = 60, # per-run status/re-drive tick cadence; lower to compress the backoff + [ValidateSet('true','false')][string]$RedriveBackoff = 'true', # ②: 'false' A/Bs the pre-backoff behaviour + [ValidateSet('true','false')][string]$ShedParameters = 'true', # retained-memory fix: 'false' A/Bs the old resident-payload behaviour + [int]$WatchSec = 300, # observe several tick waves after enqueue + [int]$PollMs = 1500, + [int]$Port = 5298, + [int]$ReadyTimeoutSec = 240, + [int]$EnqueueTimeoutSec = 120, + [switch]$KeepUp +) +$ErrorActionPreference = 'Stop' +$here = Split-Path -Parent $MyInvocation.MyCommand.Path +$root = Split-Path -Parent $here +$compose = Join-Path $root 'docker-compose.bg.yml' +$resultsDir = Join-Path $root 'results' +New-Item -ItemType Directory -Force $resultsDir | Out-Null +$stamp = Get-Date -Format 'yyyyMMdd-HHmmss' +$base = "http://127.0.0.1:$Port" +$sutContainer = 'craft-perf-bg-sut' + +function Info($m){ Write-Host "[manyruns] $m" -ForegroundColor Cyan } +function Warn($m){ Write-Host "[manyruns] $m" -ForegroundColor Yellow } +function Get-Alloc { try { Invoke-RestMethod "$base/API/PerfAllocation" -TimeoutSec 5 } catch { $null } } +function Get-Threads { try { Invoke-RestMethod "$base/API/PerfThreads" -TimeoutSec 5 } catch { $null } } + +$env:SUT_IMAGE=$SutImage; $env:SUT_PORT="$Port"; $env:SUT_CPUS="$Cpus"; $env:BG_POOL="$BgPool" +$env:GC_HEAP_LIMIT_MB = if($HeapLimitMB -gt 0){ "$HeapLimitMB" } else { '' } +$env:LOG_LEVEL = $LogLevel +# JobQueueBatchSize/PollIntervalMs bind as int — the compose's ${JOB_BATCH:-} yields an EMPTY string when +# unset, which the config binder rejects at startup (JobQueuePump ctor). Give them real values like run-oom. +$env:JOB_BATCH = "$BgPool"; $env:JOB_POLL_MS = "1000" +$env:STATUS_INTERVAL = "$StatusIntervalSec" +$env:REDRIVE_BACKOFF = $RedriveBackoff +$env:SHED_PARAMS = $ShedParameters +# Nothing about this test depends on the limiter ramp — the tick storm is C# thread-pool work, not BG-pool +# work. Leave the BG defaults; the hold-open tasks only need a couple of slots. + +Info "image=$SutImage runs=$Runs tasksPerRun=$TasksPerRun bgPool=$BgPool cpus=$Cpus heapLimitMB=$(if($HeapLimitMB -gt 0){$HeapLimitMB}else{'none'}) logLevel=$LogLevel watchSec=$WatchSec" +Info "compose up ..." +docker compose -f $compose up -d 2>&1 | Out-Host +if ($LASTEXITCODE -ne 0){ throw "compose up failed" } + +$result = [ordered]@{ label=$Label; timestamp=$stamp; runs=$Runs; tasksPerRun=$TasksPerRun; paramKB=$ParamKB; holdMs=$HoldMs + bgPool=$BgPool; cpus=$Cpus; heapLimitMB=$HeapLimitMB; logLevel=$LogLevel; statusIntervalSec=$StatusIntervalSec; redriveBackoff=$RedriveBackoff; shedParameters=$ShedParameters; watchSec=$WatchSec } +try { + Info "waiting for ready (timeout ${ReadyTimeoutSec}s) ..." + $ready=$false; $dl=(Get-Date).AddSeconds($ReadyTimeoutSec) + while((Get-Date) -lt $dl){ + try { if((Invoke-RestMethod "$base/healthz" -TimeoutSec 5).status -eq 'ready'){ $ready=$true; break } } catch {} + Start-Sleep -Seconds 2 + } + if(-not $ready){ throw 'SUT never became ready' } + + # Baseline heap (retry until the PS pool answers with a real reading). + Start-Sleep -Seconds 2 + $b=$null + for($i=0; $i -lt 30; $i++){ $a=Get-Alloc; if($a -and ([double]$a.memory.heapMB) -gt 0){ $b=$a.memory; break }; Start-Sleep -Milliseconds 500 } + if(-not $b){ throw 'could not read a baseline memory sample from /API/PerfAllocation' } + $baseHeap=[double]$b.heapMB; $baseUsed=[double]$b.containerUsedMB; $gcLimit=[double]$b.gcHeapLimitMB + $th0=Get-Threads + Info ("baseline: heap={0}MB containerUsed={1}MB gcHeapLimit={2}MB threadPool={3} osThreads={4}" -f ` + $baseHeap,$baseUsed,$gcLimit,$(if($th0){$th0.threadPool.threadCount}else{'?'}),$(if($th0){$th0.process.osThreadCount}else{'?'})) + $result.baselineHeapMB=$baseHeap; $result.baselineContainerUsedMB=$baseUsed; $result.gcHeapLimitMB=$gcLimit + + # ── Create M live runs (chunked) ──────────────────────────────────────────── + Info "creating $Runs runs x $TasksPerRun tasks (chunks of $ChunkRuns) ..." + $t0=Get-Date; $created=0; $enqCrashed=$false + while($created -lt $Runs){ + $chunk=[math]::Min($ChunkRuns, $Runs-$created) + $u="$base/API/PerfManyRuns?runs=$chunk&tasks=$TasksPerRun&holdms=$HoldMs¶mkb=$ParamKB&prefix=$Label" + $resp = & curl.exe -s --max-time $EnqueueTimeoutSec $u 2>$null + $ok=$null; try { $ok=($resp | ConvertFrom-Json).created } catch {} + if($null -eq $ok){ + # A crash mid-enqueue is itself a result (the graph didn't even fit to build) — record and stop. + if(-not (Get-Alloc)){ $enqCrashed=$true; Warn "SUT unreachable during enqueue at created=$created — crashed while building the run set"; break } + Warn "enqueue chunk returned no count (resp: $resp) — retrying"; Start-Sleep -Seconds 2; continue + } + $created+=[int]$ok + $a=Get-Alloc + Info (" created={0,6}/{1} qtotal={2,7} heap={3}MB gc2={4}" -f ` + $created,$Runs,$(if($a){[int]$a.queue.total}else{'?'}),$(if($a){[double]$a.memory.heapMB}else{'?'}),$(if($a){[int]$a.memory.gc2}else{'?'})) + } + $enqSec=[math]::Round(((Get-Date)-$t0).TotalSeconds,1) + $result.runsCreated=$created; $result.enqueueSec=$enqSec + + # ── Observe: watch the 60s tick waves drive heap / gc2 / thread-pool ───────── + Info "observing tick waves for ${WatchSec}s ..." + $samples=New-Object System.Collections.ArrayList + $peakHeap=$baseHeap; $peakUsed=$baseUsed; $peakGc2=[int]$b.GC2 + $baseGc2=[int]$b.GC2; $peakPending=0; $peakOsThreads=0; $peakTpThreads=0 + $unreachable=0; $maxUnreachable=0; $crashed=$false + $redriveReadsFirst=$null; $redriveReadsLast=$null + $lastLog=Get-Date; $wStart=Get-Date; $wdl=(Get-Date).AddSeconds($WatchSec) + $firstGc2=$null; $firstGc2At=$null; $lastGc2=$null; $lastGc2At=$null + while((Get-Date) -lt $wdl){ + $a=Get-Alloc + if(-not $a){ + $unreachable++; if($unreachable -gt $maxUnreachable){$maxUnreachable=$unreachable} + if($unreachable -ge 40){ $crashed=$true; Warn "SUT unreachable ~20s — process crashed (OOM repro)"; break } + Start-Sleep -Milliseconds 500; continue + } + $unreachable=0 + $th=Get-Threads + $t=[math]::Round(((Get-Date)-$wStart).TotalSeconds,1) + $heap=[double]$a.memory.heapMB; $used=[double]$a.memory.containerUsedMB; $gc2=[int]$a.memory.gc2 + if($heap -gt $peakHeap){$peakHeap=$heap}; if($used -gt $peakUsed){$peakUsed=$used}; if($gc2 -gt $peakGc2){$peakGc2=$gc2} + if($null -eq $firstGc2){ $firstGc2=$gc2; $firstGc2At=Get-Date } + $lastGc2=$gc2; $lastGc2At=Get-Date + $pending=if($th){[long]$th.threadPool.pendingWorkItems}else{0} + $osThreads=if($th){[int]$th.process.osThreadCount}else{0} + $tpThreads=if($th){[int]$th.threadPool.threadCount}else{0} + if($th -and $null -ne $th.process.redriveReads){ if($null -eq $redriveReadsFirst){$redriveReadsFirst=[long]$th.process.redriveReads}; $redriveReadsLast=[long]$th.process.redriveReads } + if($pending -gt $peakPending){$peakPending=$pending} + if($osThreads -gt $peakOsThreads){$peakOsThreads=$osThreads} + if($tpThreads -gt $peakTpThreads){$peakTpThreads=$tpThreads} + [void]$samples.Add([pscustomobject]@{ t=$t; heapMB=$heap; usedMB=$used; gc2=$gc2 + qtotal=[int]$a.queue.total; bgBusy=[int]$a.pool.bgBusy + tpPending=$pending; tpThreads=$tpThreads; osThreads=$osThreads }) + if(((Get-Date)-$lastLog).TotalSeconds -ge 10){ + $lastLog=Get-Date + Info ("t={0,6}s heap={1}MB used={2}MB gc2={3} tpPending={4} tpThreads={5} osThreads={6} qtotal={7}" -f ` + $t,$heap,$used,$gc2,$pending,$tpThreads,$osThreads,[int]$a.queue.total) + } + Start-Sleep -Milliseconds $PollMs + } + + # gc2 rate over the observation window = the GC-thrash signal (full compacting collections / minute). + $gc2Rate=$null + if($firstGc2 -ne $null -and $lastGc2At -gt $firstGc2At){ + $mins=((($lastGc2At)-($firstGc2At)).TotalMinutes) + if($mins -gt 0){ $gc2Rate=[math]::Round((($lastGc2-$firstGc2)/$mins),1) } + } + + # ── Scrape the container log for the symptoms + exit code ───────────────────── + $log = (& docker logs $sutContainer 2>&1) -join "`n" + $httpSlow = ([regex]::Matches($log,'HTTP-SLOW')).Count + $tickFail = ([regex]::Matches($log,'status tick failed')).Count + $redriveSkip = ([regex]::Matches($log,'skipping re-drive')).Count + $oomLines = ([regex]::Matches($log,'OutOfMemoryException')).Count + # Per-run status lines emitted ("[Scheduler] Run T+...") — the flood ① is meant to collapse. + $statusLines = ([regex]::Matches($log,'\] Run .*T\+')).Count + $exitCode=$null + try { $exitCode=[int](& docker inspect $sutContainer --format '{{.State.ExitCode}}' 2>$null) } catch {} + $isRunning=$null + try { $isRunning=(& docker inspect $sutContainer --format '{{.State.Running}}' 2>$null) } catch {} + + $result.enqueueCrashed=$enqCrashed; $result.crashed=($crashed -or $enqCrashed) + $result.peakHeapMB=$peakHeap; $result.peakContainerUsedMB=$peakUsed + $result.heapGrowthMB=[math]::Round($peakHeap-$baseHeap,1) + $result.baseGc2=$baseGc2; $result.peakGc2=$peakGc2; $result.gc2PerMin=$gc2Rate + $result.peakThreadPoolPending=$peakPending; $result.peakThreadPoolThreads=$peakTpThreads; $result.peakOsThreads=$peakOsThreads + $redriveReadsWindow = if($null -ne $redriveReadsFirst -and $null -ne $redriveReadsLast){ $redriveReadsLast-$redriveReadsFirst } else { $null } + $result.httpSlowCount=$httpSlow; $result.tickFailCount=$tickFail; $result.redriveSkipCount=$redriveSkip; $result.oomExceptionCount=$oomLines; $result.statusLogLines=$statusLines + $result.redriveStorageReadsInWindow=$redriveReadsWindow + $result.containerExitCode=$exitCode; $result.containerRunning=$isRunning; $result.maxUnreachableStreak=$maxUnreachable + $result.samples=$samples + + Write-Host "" + Write-Host "===== run-count axis: $Label ($created runs x $TasksPerRun, heapLimit=$(if($HeapLimitMB -gt 0){"${HeapLimitMB}MB"}else{'none'}), cpus=$Cpus) =====" -ForegroundColor Yellow + Write-Host (" baseline heap : {0} MB gc heap hard limit: {1}" -f $baseHeap,$(if($gcLimit -gt 0){"$gcLimit MB"}else{'(none)'})) -ForegroundColor Gray + Write-Host (" PEAK heap : {0} MB (growth {1} MB) peak container used: {2} MB" -f $peakHeap,$result.heapGrowthMB,$peakUsed) -ForegroundColor Gray + Write-Host (" Gen2 full GCs : {0} -> {1} rate: {2}/min <-- GC-thrash signal" -f $baseGc2,$peakGc2,$(if($null -ne $gc2Rate){$gc2Rate}else{'?'})) -ForegroundColor Gray + Write-Host (" thread-pool pending : peak {0} tpThreads peak {1} osThreads peak {2} <-- thread starvation" -f $peakPending,$peakTpThreads,$peakOsThreads) -ForegroundColor Gray + Write-Host (" [HTTP-SLOW] log lines : {0} tick-failures: {1} OOMException: {2}" -f $httpSlow,$tickFail,$oomLines) -ForegroundColor $(if($httpSlow -gt 0 -or $oomLines -gt 0){'Yellow'}else{'Gray'}) + Write-Host (" status log lines : {0} <-- per-run 'T+' flood (① collapses this)" -f $statusLines) -ForegroundColor Gray + Write-Host (" re-drive storage reads: {0} in window (backoff={1}) <-- ② holds this down" -f $(if($null -ne $redriveReadsWindow){$redriveReadsWindow}else{'?'}),$RedriveBackoff) -ForegroundColor Gray + Write-Host (" process crashed : {0} exitCode={1} running={2}" -f $result.crashed,$exitCode,$isRunning) -ForegroundColor $(if($result.crashed){'Red'}else{'Green'}) + + $jsonOut = Join-Path $resultsDir "$Label-m$Runs-h$HeapLimitMB-$stamp.json" + ($result | ConvertTo-Json -Depth 6) | Set-Content $jsonOut -Encoding utf8 + Info "wrote $jsonOut" +} +finally { + if($KeepUp){ Warn "leaving up (-KeepUp). down: docker compose -f `"$compose`" down -v" } + else { Info "tearing down ..."; docker compose -f $compose down -v 2>&1 | Out-Null } +} diff --git a/tests/Craft.Tests/JobQueueDispatchableTests.cs b/tests/Craft.Tests/JobQueueDispatchableTests.cs new file mode 100644 index 0000000..3821a71 --- /dev/null +++ b/tests/Craft.Tests/JobQueueDispatchableTests.cs @@ -0,0 +1,106 @@ +using Craft.Configuration; +using Craft.Storage; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +/// +/// The index table answers "does this run still have queued work" from a single partition, which is +/// what keeps the re-drive and finalize paths off a full-table scan. But the index can OUTLIVE the +/// queue rows it points at, and when it does it lies: it reports a task queued that no pump will ever +/// claim, so an orphan re-drive that trusts the index skips it and the run stalls indefinitely with +/// the task Pending — resume logs "Dispatched 0 tasks (N already queued)" and the orphan watchdog +/// never fires, because GetQueuedTaskIdsAsync (index-only) keeps reporting the task queued while the +/// queue table has no runnable row for it. Clearing the stale index row is the only thing that lets +/// the watchdog see it as orphaned again. +/// +/// is the fix: it verifies each candidate +/// against the queue TABLE, returning only tasks the pump can actually still claim. These tests pin +/// the three divergence shapes it has to catch. +/// +public class JobQueueDispatchableTests +{ + private static readonly TimeSpan Lease = TimeSpan.FromMinutes(20); + private const string QueueTable = "OrchestratorQueue"; + + private static (JobQueueStore Queue, RunRemainingCounterTests.ConditionalStore Backing) NewQueue() + { + var backing = new RunRemainingCounterTests.ConditionalStore(); + var queue = new JobQueueStore(NullLogger.Instance, new CraftSettings(), backing); + return (queue, backing); + } + + [Fact] + public async Task NormallyQueuedTasksAreAllDispatchable() + { + var (queue, _) = NewQueue(); + await queue.InitializeAsync(); + await queue.EnqueueBatchAsync("R", [("a", 4), ("b", 4), ("c", 4)], DateTime.UtcNow); + + var dispatchable = await queue.GetDispatchableTaskIdsAsync("R", ["a", "b", "c"]); + + Assert.Equal(3, dispatchable.Count); + Assert.Contains("a", dispatchable); + Assert.Contains("b", dispatchable); + Assert.Contains("c", dispatchable); + } + + [Fact] + public async Task AnIndexRowWithNoQueueRowIsNotDispatchable_ButTheIndexStillListsIt() + { + var (queue, backing) = NewQueue(); + await queue.InitializeAsync(); + await queue.EnqueueBatchAsync("R", [("a", 4), ("b", 4)], DateTime.UtcNow); + + // Delete ONLY b's queue row, leaving its index row — the exact divergence a crash between the + // two deletes, or a run carried in from a pre-pump build, leaves behind. + await backing.DeleteAsync(QueueTable, JobQueueStore.Bucket(4), JobQueueStore.BuildRowKey("R", "b")); + + // The index — what the old re-drive trusted — still reports both as queued. + var indexView = await queue.GetQueuedTaskIdsAsync("R"); + Assert.Contains("a", indexView); + Assert.Contains("b", indexView); + + // The queue-verified view sees b for the ghost it is. + var dispatchable = await queue.GetDispatchableTaskIdsAsync("R", ["a", "b"]); + Assert.Contains("a", dispatchable); + Assert.DoesNotContain("b", dispatchable); + } + + [Fact] + public async Task AQueueRowOwnedWithNoLeaseIsNotDispatchable() + { + var (queue, backing) = NewQueue(); + await queue.InitializeAsync(); + await queue.EnqueueBatchAsync("R", [("a", 4)], DateTime.UtcNow); + + // Owned with no LeaseUntil: neither "Owner eq ''" nor "LeaseUntil lt now", so the claim filter + // can never match it and the pump will never dispatch it — a ghost as surely as a missing row. + var key = JobQueueStore.BuildRowKey("R", "a"); + var row = await backing.GetAsync(QueueTable, JobQueueStore.Bucket(4), key); + Assert.NotNull(row); + row!["Owner"] = "dead-instance"; + row["LeaseUntil"] = (DateTimeOffset?)null; + await backing.UpsertAsync(QueueTable, row); + + var dispatchable = await queue.GetDispatchableTaskIdsAsync("R", ["a"]); + Assert.Empty(dispatchable); + } + + [Fact] + public async Task AQueueRowUnderALeaseStaysDispatchable() + { + var (queue, _) = NewQueue(); + await queue.InitializeAsync(); + await queue.EnqueueBatchAsync("R", [("a", 4)], DateTime.UtcNow); + + // A live claim: owned under a lease. The pump (or its lapse) will run it, so the re-drive must + // NOT treat it as orphaned — doing so would reset a row a worker is actively holding and could + // run the task a second time. + var claimed = await queue.ClaimBatchAsync("worker-a", 8, Lease); + Assert.Single(claimed); + + var dispatchable = await queue.GetDispatchableTaskIdsAsync("R", ["a"]); + Assert.Contains("a", dispatchable); + } +} diff --git a/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs b/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs index 83beb9e..6be6fc6 100644 --- a/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs +++ b/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs @@ -234,7 +234,8 @@ public async Task DispatchingARunAgain_ClearsAPreviousFinalizeClaim() { var (svc, _) = NewService(); Set(svc, "_finalizingRuns", new ConcurrentDictionary()); - Set(svc, "_runStatusTimers", new ConcurrentDictionary()); + // No _runStatusTimers any more: the per-run status timer was replaced by a single sweep loop + // (RunStatusSweepLoopAsync), so DispatchPendingTasksAsync no longer creates or tracks a timer. Claims(svc)["recurring-run"] = true; var run = new OrchestratorRun { Name = "recurring-run", Status = "Running", StartedUtc = DateTime.UtcNow }; @@ -244,7 +245,7 @@ public async Task DispatchingARunAgain_ClearsAPreviousFinalizeClaim() { var mi = typeof(OrchestratorService).GetMethod("DispatchPendingTasksAsync", BindingFlags.NonPublic | BindingFlags.Instance)!; - await (Task)mi.Invoke(svc, [run, "Invoke-CraftTask", 4, CancellationToken.None])!; + await (Task)mi.Invoke(svc, [run, "Invoke-CraftTask", 4, CancellationToken.None, false])!; } catch { /* expected: no queue on this instance */ } diff --git a/tests/Craft.Tests/OrchestratorSequentialTests.cs b/tests/Craft.Tests/OrchestratorSequentialTests.cs new file mode 100644 index 0000000..f5824e2 --- /dev/null +++ b/tests/Craft.Tests/OrchestratorSequentialTests.cs @@ -0,0 +1,527 @@ +using System.Collections.Concurrent; +using System.Reflection; +using Craft.Configuration; +using Craft.Orchestration; +using Craft.PowerShellHost; +using Craft.Services; +using Craft.Storage; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +/// +/// Pins SEQUENTIAL orchestrator mode. +/// +/// A run marked runs its tasks ONE AT A TIME, in ascending +/// (payload) order, ON A SINGLE PINNED WORKER: one entry row +/// is dispatched, that one claim drives the WHOLE run inline (checkout one worker → run every step on it → +/// reclaim once), and the not-yet-reached steps never get their own queue row. The default (false) is the +/// existing fan-out: every task enqueued up front and drained in parallel by the pool. +/// +/// These fix the contract at every seam it touches: +/// - persistence — the flag and the per-task order survive a round trip AND a status-write rewrite +/// (Replace mode erases any column the write omits), so a resumed run keeps its order; +/// - dispatch — a sequential run enqueues only its single entry row; fan-out enqueues all; +/// - driver — every step runs, in Sequence order, on ONE worker; a failing step does not strand the +/// rest (best-effort); a cancelled run marks the remaining steps Cancelled; a duplicate +/// entry row is a no-op while a driver is already active; +/// - re-drive — the watchdog leaves a run alone while its driver is active (or its entry job is still +/// queued/running), yet still restarts the driver if the entry row is lost. +/// +/// The driver's loop logic is exercised through , a subclass that overrides the +/// three PowerShell seams (checkout / run-step / reclaim), so these tests need no worker pool. The actual +/// pinning to a live worker is proven separately by live validation against the dev backend. +/// +public class OrchestratorSequentialTests +{ + private const string TaskFunc = "Invoke-CraftTask"; + + // ─── seam subclass ────────────────────────────────────────────────────────────────────────────── + // Overrides only the three PowerShell interactions. Everything else — ordering, best-effort, marker + // writes, completion, cancellation — runs the real OrchestratorService code. + + private sealed class SeqDriver : OrchestratorService + { + // OrchestratorService has only a parameterized constructor; a subclass must chain to it to compile. + // Never actually invoked — the harness builds this via GetUninitializedObject — so the nulls are safe. + private SeqDriver() : base(null!, null!, null!, null!, null!, null!, null!, null!) { } + + public List Order = new(); // step ids in the order the driver ran them + public HashSet FailIds = new(); // ids whose step throws (best-effort test) + public int Checkouts; // must be exactly 1 for a run that starts + public int Reclaims; // must pair 1:1 with Checkouts + public string StepOutput = string.Empty; // what a step "returns" (PostExecution capture test) + public Action? OnStep; // test hook, runs synchronously inside a step + + internal override PowerShellWorker? CheckoutSequentialWorker(CancellationToken ct) + { + Checkouts++; + return null; // step execution is faked, so the worker handle is never dereferenced + } + + internal override void ReclaimSequentialWorker(PowerShellWorker? worker, bool faulted) => Reclaims++; + + internal override Task RunSequentialStepAsync( + OrchestratorRun run, OrchestratorTaskItem task, string taskPath, PowerShellWorker? worker) + { + Order.Add(task.Id); + OnStep?.Invoke(task); + if (FailIds.Contains(task.Id)) + throw new InvalidOperationException("boom " + task.Id); + return Task.FromResult(StepOutput); + } + } + + // ─── harness ──────────────────────────────────────────────────────────────────────────────────── + + private sealed record Harness(SeqDriver Svc, OrchestratorTableStore Store, JobQueueStore Queue, JobManager Jobs); + + private static async Task NewHarnessAsync() + { + var settings = new CraftSettings + { + Orchestrator = { TablePrefix = "seq" + Guid.NewGuid().ToString("N")[..8] } + }; + // Direct, synchronous status writes — no batching barrier or drain loop to stand up in a unit test. + settings.Orchestrator.BatchStatusWrites = false; + + var backing = new RunRemainingCounterTests.ConditionalStore(); + var store = new OrchestratorTableStore(NullLogger.Instance, settings, backing); + var queue = new JobQueueStore(NullLogger.Instance, settings, backing); + await store.InitializeAsync(); + await queue.InitializeAsync(); + var writer = new OrchestratorStatusWriter(store, NullLogger.Instance, settings); + + // Build the service without its constructor (which drags in the PowerShell runner and the worker + // pool) and set only the fields the dispatch / driver / re-drive paths read — the same shape the + // other orchestrator unit tests use. + var svc = (SeqDriver)System.Runtime.CompilerServices.RuntimeHelpers.GetUninitializedObject(typeof(SeqDriver)); + svc.Order = new(); + svc.FailIds = new(); + Set(svc, "_logger", NullLogger.Instance); + Set(svc, "_store", store); + Set(svc, "_queue", queue); + Set(svc, "_writer", writer); + Set(svc, "_lock", new object()); + Set(svc, "_activeRuns", new ConcurrentDictionary()); + Set(svc, "_taskScriptPaths", new ConcurrentDictionary()); + Set(svc, "_finalizingRuns", new ConcurrentDictionary()); + Set(svc, "_finalizeDeferrals", new ConcurrentDictionary()); + Set(svc, "_childRuns", new ConcurrentDictionary>()); + Set(svc, "_cancelledRuns", new ConcurrentDictionary()); + Set(svc, "_activeSequentialDrivers", new ConcurrentDictionary()); + Set(svc, "_requeueFailures", new ConcurrentDictionary()); + Set(svc, "_deferrals", NewFieldDict(svc, "_deferrals")); + Set(svc, "_redriveBackoff", NewFieldDict(svc, "_redriveBackoff")); + Set(svc, "_shedParameters", false); + Set(svc, "_redriveBackoffEnabled", false); // pin the sequential logic, not the backoff timing + Set(svc, "_redriveBase", TimeSpan.FromSeconds(60)); + Set(svc, "_settings", settings); + + // A JobManager with an empty job map, so IsQueuedOrRunning answers false (nothing dispatched here). + var jm = (JobManager)System.Runtime.CompilerServices.RuntimeHelpers.GetUninitializedObject(typeof(JobManager)); + var jobsField = typeof(JobManager).GetField("_jobs", BindingFlags.NonPublic | BindingFlags.Instance)!; + jobsField.SetValue(jm, Activator.CreateInstance(jobsField.FieldType)); + Set(svc, "_jobManager", jm); + + return new Harness(svc, store, queue, jm); + } + + private static object NewFieldDict(object svc, string field) + { + var t = typeof(OrchestratorService).GetField(field, BindingFlags.NonPublic | BindingFlags.Instance)!.FieldType; + return Activator.CreateInstance(t)!; + } + + private static void Set(object target, string field, object? value) => + typeof(OrchestratorService).GetField(field, BindingFlags.NonPublic | BindingFlags.Instance)! + .SetValue(target, value); + + private static T Get(object target, string field) => + (T)typeof(OrchestratorService).GetField(field, BindingFlags.NonPublic | BindingFlags.Instance)! + .GetValue(target)!; + + private static Task Invoke(OrchestratorService svc, string method, params object[] args) => + (Task)typeof(OrchestratorService).GetMethod(method, BindingFlags.NonPublic | BindingFlags.Instance)! + .Invoke(svc, args)!; + + private static Task Dispatch(Harness h, OrchestratorRun run) => + Invoke(h.Svc, "DispatchPendingTasksAsync", run, TaskFunc, run.Priority, CancellationToken.None, false); + private static Task Redrive(Harness h, OrchestratorRun run) => Invoke(h.Svc, "RedrivePendingTasksAsync", run); + + /// Run the sequential driver to completion. Pre-seeds _finalizingRuns so the background + /// finalize the last step schedules cannot mutate the run under the test's assertions. + private static async Task DriveAsync(Harness h, OrchestratorRun run, string taskPath = TaskFunc, + CancellationToken ct = default) + { + Get>(h.Svc, "_finalizingRuns")[run.Name] = true; + var mi = typeof(OrchestratorService).GetMethod("BuildSequentialRunWork", + BindingFlags.NonPublic | BindingFlags.Instance)!; + var work = (Func)mi.Invoke(h.Svc, [run, taskPath])!; + try { await work(ct); } + catch (TargetInvocationException ex) { throw ex.InnerException ?? ex; } + } + + private static OrchestratorRun MakeRun(string name, int count, bool sequential) + { + var run = new OrchestratorRun + { + Name = name, + Status = "Running", + Priority = 4, + Sequential = sequential, + StartedUtc = DateTime.UtcNow, + TaskScriptName = TaskFunc + }; + for (var i = 0; i < count; i++) + run.Tasks.Add(new OrchestratorTaskItem + { + Id = $"{name}_t{i}", + Status = "Pending", + Sequence = i, + Parameters = new Dictionary { ["FunctionName"] = "Push-Noop", ["idx"] = i } + }); + return run; + } + + private static async Task> QueuedAsync(Harness h, string run) => + (await h.Queue.GetQueuedTaskIdsAsync(run)).OrderBy(x => x, StringComparer.Ordinal).ToList(); + + /// Re-drive re-enqueues via a fire-and-forget Task.Run, so poll for the row to land. + private static async Task> WaitQueuedAsync(Harness h, string run, int expected) + { + for (var i = 0; i < 100; i++) + { + var ids = await h.Queue.GetQueuedTaskIdsAsync(run); + if (ids.Count >= expected) return ids.OrderBy(x => x, StringComparer.Ordinal).ToList(); + await Task.Delay(20); + } + return await QueuedAsync(h, run); + } + + private static string St(OrchestratorRun run, int i) => run.Tasks.Single(t => t.Sequence == i).Status; + + // ─── persistence ──────────────────────────────────────────────────────────────────────────────── + + [Fact] + public async Task Sequential_And_Sequence_SurviveCreationRoundTrip() + { + var h = await NewHarnessAsync(); + var run = MakeRun("persist-seq", 3, sequential: true); + await h.Store.UpsertRunAsync(run); + await h.Store.UpsertTaskBatchAsync(run.Name, run.Tasks); + + var loaded = await h.Store.GetRunAsync(run.Name); + + Assert.NotNull(loaded); + Assert.True(loaded!.Sequential); + Assert.Equal([0, 1, 2], + loaded.Tasks.OrderBy(t => t.Sequence).Select(t => t.Sequence).ToArray()); + } + + [Fact] + public async Task FanOutRun_RoundTripsSequentialFalse() + { + var h = await NewHarnessAsync(); + var run = MakeRun("persist-fanout", 2, sequential: false); + await h.Store.UpsertRunAsync(run); + await h.Store.UpsertTaskBatchAsync(run.Name, run.Tasks); + + var loaded = await h.Store.GetRunAsync(run.Name); + + Assert.False(loaded!.Sequential); + } + + [Fact] + public async Task Sequence_SurvivesStatusWriteRewrite() + { + // The status writer rewrites a task row on every transition with Replace semantics, so Sequence + // must be part of that write or a resumed run would read every task as Sequence 0 and lose order. + var h = await NewHarnessAsync(); + var run = MakeRun("persist-statuswrite", 2, sequential: true); + await h.Store.UpsertRunAsync(run); + await h.Store.UpsertTaskBatchAsync(run.Name, run.Tasks); + + // Simulate a terminal status write (as the coalescing writer does) for the second task. + var t1 = run.Tasks[1]; + await h.Store.WriteTaskStatusBatchAsync( + [ + new TaskStatusWrite(run.Name, t1.Id, "Completed", "{}", 0, null, DateTime.UtcNow, null, t1.Sequence) + ]); + + var loaded = await h.Store.GetRunAsync(run.Name); + var reloaded = loaded!.Tasks.Single(t => t.Id == t1.Id); + Assert.Equal(1, reloaded.Sequence); // not erased to 0 by the Replace write + Assert.Equal("Completed", reloaded.Status); + } + + // ─── dispatch gate ────────────────────────────────────────────────────────────────────────────── + + [Fact] + public async Task Sequential_Dispatch_EnqueuesOnlyTheEntryRow() + { + var h = await NewHarnessAsync(); + var run = MakeRun("disp-seq", 5, sequential: true); + + await Dispatch(h, run); + + var queued = await QueuedAsync(h, run.Name); + Assert.Equal(["disp-seq_t0"], queued); // only Sequence 0, despite five pending tasks + } + + [Fact] + public async Task FanOut_Dispatch_EnqueuesEveryTask() + { + var h = await NewHarnessAsync(); + var run = MakeRun("disp-fanout", 5, sequential: false); + + await Dispatch(h, run); + + var queued = await QueuedAsync(h, run.Name); + Assert.Equal(5, queued.Count); + } + + [Fact] + public async Task Sequential_Dispatch_OnResume_EnqueuesEntryOnly_NeverASecondRow() + { + // Resume shape: tasks 0,1 done; task 2 is the reached task and its queue row SURVIVED the restart; + // task 3 has not been reached. Dispatch must not enqueue task 3 (a second row would let a second + // driver start) — it must leave the already-queued entry alone. + var h = await NewHarnessAsync(); + var run = MakeRun("disp-resume", 4, sequential: true); + run.Tasks[0].Status = "Completed"; + run.Tasks[1].Status = "Completed"; + await h.Queue.EnqueueBatchAsync(run.Name, [(run.Tasks[2].Id, 4)], DateTime.UtcNow); + + await Dispatch(h, run); + + var queued = await QueuedAsync(h, run.Name); + Assert.Equal(["disp-resume_t2"], queued); // task 3 NOT enqueued + } + + [Fact] + public async Task Sequential_Dispatch_OnResume_ReEnqueuesCurrentStep_WhenItsRowIsGone() + { + // The reached step's queue row was lost. Dispatch re-enqueues exactly it, to restart the driver. + var h = await NewHarnessAsync(); + var run = MakeRun("disp-resume-gone", 3, sequential: true); + run.Tasks[0].Status = "Completed"; + + await Dispatch(h, run); + + var queued = await QueuedAsync(h, run.Name); + Assert.Equal(["disp-resume-gone_t1"], queued); + } + + // ─── driver ───────────────────────────────────────────────────────────────────────────────────── + + [Fact] + public async Task Driver_RunsEveryStepInSequenceOrder_OnExactlyOneWorker() + { + var h = await NewHarnessAsync(); + var run = MakeRun("drv-order", 4, sequential: true); + run.Tasks.Reverse(); // insertion order 3,2,1,0 — only the Sequence sort can recover 0,1,2,3 + + await DriveAsync(h, run); + + Assert.Equal(["drv-order_t0", "drv-order_t1", "drv-order_t2", "drv-order_t3"], h.Svc.Order); + Assert.All(run.Tasks, t => Assert.Equal("Completed", t.Status)); + Assert.Equal(1, h.Svc.Checkouts); // one worker for the whole run + Assert.Equal(1, h.Svc.Reclaims); // reclaimed once, at the end + } + + [Fact] + public async Task Driver_BestEffort_ContinuesPastAFailingStep() + { + var h = await NewHarnessAsync(); + var run = MakeRun("drv-besteffort", 4, sequential: true); + h.Svc.FailIds.Add("drv-besteffort_t1"); // the second step throws + + await DriveAsync(h, run); + + // Every step is still attempted, in order, and the failure does not strand the rest. + Assert.Equal(["drv-besteffort_t0", "drv-besteffort_t1", "drv-besteffort_t2", "drv-besteffort_t3"], + h.Svc.Order); + Assert.Equal("Completed", St(run, 0)); + Assert.Equal("Failed", St(run, 1)); + Assert.Equal("Completed", St(run, 2)); + Assert.Equal("Completed", St(run, 3)); + Assert.Equal(1, h.Svc.Checkouts); // still one worker — a step failure does not re-grab a worker + Assert.Equal(1, h.Svc.Reclaims); + } + + [Fact] + public async Task Driver_Cancellation_MarksRemainingStepsCancelled_AndStops() + { + var h = await NewHarnessAsync(); + var run = MakeRun("drv-cancel", 4, sequential: true); + var cancelled = Get>(h.Svc, "_cancelledRuns"); + h.Svc.OnStep = t => { if (t.Sequence == 1) cancelled[run.Name] = true; }; // cancel while step 1 runs + + await DriveAsync(h, run); + + // Steps 0 and 1 completed; the driver noticed the cancellation before step 2 and stopped there. + Assert.Equal(["drv-cancel_t0", "drv-cancel_t1"], h.Svc.Order); + Assert.Equal("Completed", St(run, 0)); + Assert.Equal("Completed", St(run, 1)); + Assert.Equal("Cancelled", St(run, 2)); + Assert.Equal("Cancelled", St(run, 3)); + Assert.Equal(1, h.Svc.Reclaims); // the pinned worker is still reclaimed on the way out + } + + [Fact] + public async Task Driver_DuplicateEntry_IsANoOp_WhileAnotherDriverIsActive() + { + // A duplicate entry row for a run that already has an active driver must not start a second one. + var h = await NewHarnessAsync(); + var run = MakeRun("drv-dup", 3, sequential: true); + Get>(h.Svc, "_activeSequentialDrivers")[run.Name] = true; + + await DriveAsync(h, run); + + Assert.Empty(h.Svc.Order); // ran nothing + Assert.Equal(0, h.Svc.Checkouts); // never grabbed a worker + Assert.All(run.Tasks, t => Assert.Equal("Pending", t.Status)); // left for the real driver + } + + [Fact] + public async Task Driver_ClearsItsRegistration_WhenDone() + { + var h = await NewHarnessAsync(); + var run = MakeRun("drv-cleanup", 2, sequential: true); + + await DriveAsync(h, run); + + // The run must not stay registered, or its re-drive would be suppressed forever. + Assert.False(Get>(h.Svc, "_activeSequentialDrivers") + .ContainsKey(run.Name)); + } + + [Fact] + public async Task Driver_PostExecutionRun_CapturesAndStoresStepOutput() + { + var h = await NewHarnessAsync(); + var run = MakeRun("drv-postexec", 2, sequential: true); + run.PostExecFunctionName = "Push-Aggregate"; // capture path + h.Svc.StepOutput = "{\"ok\":true}"; + + await DriveAsync(h, run); + + var results = await h.Store.GetResultsAsync(run.Name); + Assert.Equal(2, results.Length); + Assert.All(results, r => Assert.Contains("ok", r)); + } + + [Fact] + public async Task ResolveTaskWork_ForSequentialRun_ReturnsAndRunsTheDriver() + { + // Routing: a claimed entry row for a sequential run resolves to the pinned driver (not the parallel + // per-task work), and invoking it drives the whole run on one worker. + var h = await NewHarnessAsync(); + var run = MakeRun("resolve-seq", 3, sequential: true); + Get>(h.Svc, "_activeRuns")[run.Name] = run; + Get>(h.Svc, "_taskScriptPaths")[run.Name] = TaskFunc; + Get>(h.Svc, "_finalizingRuns")[run.Name] = true; + + var mi = typeof(OrchestratorService).GetMethod("ResolveTaskWorkAsync", + BindingFlags.NonPublic | BindingFlags.Instance)!; + var task = (Task)mi.Invoke(h.Svc, [new JobDescriptor(run.Name, run.Tasks[0].Id, 4), CancellationToken.None])!; + await task; + var work = (Func?)task.GetType().GetProperty("Result")!.GetValue(task); + + Assert.NotNull(work); + await work!(CancellationToken.None); + + Assert.Equal(["resolve-seq_t0", "resolve-seq_t1", "resolve-seq_t2"], h.Svc.Order); + Assert.Equal(1, h.Svc.Checkouts); + } + + // ─── re-drive restriction ─────────────────────────────────────────────────────────────────────── + + [Fact] + public async Task Redrive_Sequential_DoesNothing_WhileADriverIsActive() + { + // The driver runs every step inline; the not-yet-reached steps deliberately have no queue row. The + // watchdog must not treat them as orphaned and enqueue them — that would spawn a second driver. + var h = await NewHarnessAsync(); + var run = MakeRun("rd-driver", 3, sequential: true); + Get>(h.Svc, "_activeSequentialDrivers")[run.Name] = true; + + await Redrive(h, run); + await Task.Delay(100); // give any (erroneous) fire-and-forget requeue time to land + + Assert.Empty(await QueuedAsync(h, run.Name)); + } + + [Fact] + public async Task Redrive_Sequential_DoesNothing_WhileTheEntryJobIsStillQueued() + { + // No driver registered yet, but the entry job is still Queued/Running in the JobManager (the window + // between the pump claiming the entry row and the driver registering). The watchdog must still leave + // the run alone. + var h = await NewHarnessAsync(); + var run = MakeRun("rd-entryjob", 3, sequential: true); + var jobs = Get>(h.Jobs, "_jobs"); + jobs[$"{run.Name}-{run.Tasks[0].Id}"] = new JobRecord + { + Id = $"{run.Name}-{run.Tasks[0].Id}", + Name = TaskFunc, + RunName = run.Name, + Priority = 4, + Status = "Running", + QueuedUtc = DateTime.UtcNow + }; + + await Redrive(h, run); + await Task.Delay(100); + + Assert.Empty(await QueuedAsync(h, run.Name)); + } + + private static T Get(JobManager jm, string field) => + (T)typeof(JobManager).GetField(field, BindingFlags.NonPublic | BindingFlags.Instance)!.GetValue(jm)!; + + [Fact] + public async Task Redrive_Sequential_NoOp_WhenTheEntryStepStillHasItsRow() + { + // Nothing running, entry step (Sequence 0) is queued and waiting; later steps have no row. The entry + // is dispatchable (row present) so not orphaned, and the later steps must not be touched. + var h = await NewHarnessAsync(); + var run = MakeRun("rd-waiting", 3, sequential: true); + await h.Queue.EnqueueBatchAsync(run.Name, [(run.Tasks[0].Id, 4)], DateTime.UtcNow); + + await Redrive(h, run); + await Task.Delay(100); + + Assert.Equal(["rd-waiting_t0"], await QueuedAsync(h, run.Name)); // unchanged; t1/t2 not enqueued + } + + [Fact] + public async Task Redrive_Sequential_ReDrivesOnlyTheCurrentStep_WhenStalled() + { + // Driver gone and the entry row lost: current step (Sequence 0) has no queue row and no driver is + // active. The watchdog re-enqueues exactly the current step — and none of the not-yet-reached ones — + // so a fresh driver resumes the run. + var h = await NewHarnessAsync(); + var run = MakeRun("rd-stalled", 3, sequential: true); + + await Redrive(h, run); + + var queued = await WaitQueuedAsync(h, run.Name, 1); + Assert.Equal(["rd-stalled_t0"], queued); // only the current; t1/t2 stay unqueued + } + + [Fact] + public async Task Redrive_FanOut_ReEnqueuesAllOrphanedTasks() + { + // The default fan-out behaviour is unchanged: every orphaned (Pending, no queue row) task is + // re-driven, not just the first. + var h = await NewHarnessAsync(); + var run = MakeRun("rd-fanout", 3, sequential: false); + + await Redrive(h, run); + + var queued = await WaitQueuedAsync(h, run.Name, 3); + Assert.Equal(3, queued.Count); + } +} diff --git a/tests/Craft.Tests/OrchestratorTaskPathRehydrationTests.cs b/tests/Craft.Tests/OrchestratorTaskPathRehydrationTests.cs new file mode 100644 index 0000000..0684b09 --- /dev/null +++ b/tests/Craft.Tests/OrchestratorTaskPathRehydrationTests.cs @@ -0,0 +1,123 @@ +using System.Collections.Concurrent; +using System.Reflection; +using Craft.Configuration; +using Craft.Orchestration; +using Craft.PowerShellHost; +using Craft.Services; +using Craft.Storage; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +/// +/// The resolver must rebuild a task's script from the run record when the in-memory path cache misses, +/// instead of dropping the task. +/// +/// The cache (_taskScriptPaths) is written only by DispatchPendingTasksAsync. The JobQueuePump, a +/// BackgroundService, starts claiming persisted queue rows at host start — before the scheduler has +/// even waited for the worker pool, let alone run ResumeInterruptedRunsAsync, which is what dispatches +/// (and so caches the path for) a resumed run. In that window every claimed row resolved to a cache +/// miss and the task was dropped: the resolver returned null, the JobManager marked the job Skipped, +/// and the pump deleted the queue row — permanently, for a task still Pending in a run the pump would +/// never see re-queued. +/// +/// The run record persists TaskScriptName, and the ScriptRepository is fully loaded before the pump's +/// first claim, so the path is rebuildable from storage exactly as the resume path rebuilds it. These +/// tests pin that: a cache miss on a live run rehydrates and caches the path; only a run with no +/// resolvable script at all is still dropped. +/// +public class OrchestratorTaskPathRehydrationTests +{ + private static (OrchestratorService Svc, OrchestratorTableStore Store, ConcurrentDictionary Paths) + NewService(params string[] knownModuleFunctions) + { + var settings = new CraftSettings + { + Orchestrator = { TablePrefix = "rehyd" + Guid.NewGuid().ToString("N")[..6] } + }; + var config = new ConfigurationBuilder().AddInMemoryCollection([]).Build(); + + var repo = new ScriptRepository(NullLogger.Instance, settings); + // IsModuleFunction short-circuits on a non-null _moduleFunctionNames, so FindScript resolves + // these names without a ScriptBlock or a module on disk. + typeof(ScriptRepository).GetField("_moduleFunctionNames", BindingFlags.NonPublic | BindingFlags.Instance)! + .SetValue(repo, new HashSet(knownModuleFunctions, StringComparer.OrdinalIgnoreCase)); + + var pool = new PowerShellWorkerPool(repo, NullLogger.Instance, config, settings); + var runner = new PowerShellRunnerService(NullLogger.Instance, pool, repo, settings); + var store = new OrchestratorTableStore(NullLogger.Instance, settings, + new RunRemainingCounterTests.ConditionalStore()); + var paths = new ConcurrentDictionary(); + + var svc = (OrchestratorService)System.Runtime.CompilerServices.RuntimeHelpers + .GetUninitializedObject(typeof(OrchestratorService)); + Set(svc, "_logger", NullLogger.Instance); + Set(svc, "_store", store); + Set(svc, "_psRunner", runner); + Set(svc, "_activeRuns", new ConcurrentDictionary()); + Set(svc, "_taskScriptPaths", paths); // deliberately EMPTY — the pump-before-resume window + Set(svc, "_lock", new object()); + return (svc, store, paths); + } + + private static void Set(object target, string field, object value) => + typeof(OrchestratorService).GetField(field, BindingFlags.NonPublic | BindingFlags.Instance)! + .SetValue(target, value); + + private static async Task ResolveAsync(OrchestratorService svc, string runName, string taskId) + { + var mi = typeof(OrchestratorService).GetMethod("ResolveTaskWorkAsync", + BindingFlags.NonPublic | BindingFlags.Instance)!; + try + { + var task = (Task)mi.Invoke(svc, [new JobDescriptor(runName, taskId, 4), CancellationToken.None])!; + await task; + return task.GetType().GetProperty("Result")!.GetValue(task); + } + catch (TargetInvocationException ex) + { + throw ex.InnerException ?? ex; + } + } + + private static async Task SeedRunAsync(OrchestratorTableStore store, string name, string taskScriptName) + { + await store.InitializeAsync(); + await store.UpsertRunAsync(new OrchestratorRun + { + Name = name, + Status = "Running", + Priority = 4, + StartedUtc = DateTime.UtcNow, + TaskScriptName = taskScriptName, + Tasks = [new OrchestratorTaskItem { Id = "task-0", Status = "Pending" }] + }); + await store.UpsertTaskAsync(name, new OrchestratorTaskItem { Id = "task-0", Status = "Pending" }); + } + + [Fact] + public async Task CacheMissOnALiveRun_RehydratesThePathFromTaskScriptName_AndCachesIt() + { + var (svc, store, paths) = NewService("Invoke-CraftTask"); + await SeedRunAsync(store, "AuditLogProcessV2-contoso.com", "Invoke-CraftTask"); + + var work = await ResolveAsync(svc, "AuditLogProcessV2-contoso.com", "task-0"); + + Assert.NotNull(work); // previously: dropped, because _taskScriptPaths had no entry yet + Assert.Equal("Invoke-CraftTask", paths["AuditLogProcessV2-contoso.com"]); + } + + [Fact] + public async Task ARunWithNoResolvableTaskScript_IsStillDropped() + { + // Empty TaskScriptName, and the naming-convention fallback (Invoke-Task) resolves to nothing + // because the repo knows no such function. This is the only case that should still drop. + var (svc, store, _) = NewService(/* no known functions */); + await SeedRunAsync(store, "mysteryrun", taskScriptName: ""); + + var work = await ResolveAsync(svc, "mysteryrun", "task-0"); + + Assert.Null(work); + } +}