diff --git a/Services/Orchestration/BackgroundTaskLimiter.cs b/Services/Orchestration/BackgroundTaskLimiter.cs index 5b196d0..0e39401 100644 --- a/Services/Orchestration/BackgroundTaskLimiter.cs +++ b/Services/Orchestration/BackgroundTaskLimiter.cs @@ -16,8 +16,8 @@ namespace Craft.Orchestration; /// BG scale-down: when queue drains to 0 active + 0 waiting, scales back to baseline. /// /// HTTP pressure throttle: if the number of busy HTTP workers meets or exceeds -/// HttpPressureThreshold for a sustained period, BG concurrency is reduced to 2 -/// to give HTTP maximum CPU headroom. When HTTP pressure drops, BG concurrency +/// HttpPressureThreshold for a sustained period, BG concurrency is reduced to +/// HttpPressureConcurrency (default half the ceiling) to give HTTP CPU headroom. When HTTP pressure drops, BG concurrency /// restores to baseline. /// /// - Ceiling is capped to BgPoolSize (the real bottleneck). @@ -97,12 +97,18 @@ public class BackgroundTaskLimiter : IDisposable /// /// Number of busy HTTP workers that triggers BG throttling. /// When HttpPoolSize - HttpAvailable >= this value for HttpPressureSeconds, - /// BG concurrency drops to 2. - /// Default: half of HttpPoolSize (e.g. 2 on a 4-worker pool). + /// BG concurrency drops to . + /// Default: two thirds of HttpPoolSize (e.g. 4 on a 6-worker pool). Half the pool tripped on routine + /// API-client traffic and held BG at 2 for hours, starving the scheduled cache runs. /// Set to 0 to disable HTTP pressure throttling. /// public int HttpPressureThreshold { get; } + /// + /// BG concurrency while HTTP-throttled. Default: half the ceiling, minimum 2. + /// + public int HttpPressureConcurrency { get; } + /// /// How long HTTP pressure must be sustained before throttling BG tasks. /// Default: 10 seconds. @@ -132,10 +138,13 @@ public BackgroundTaskLimiter(ILogger logger, IConfigurati ScaleUpAfter = TimeSpan.FromSeconds( configuration.GetValue("BackgroundScaleUpAfterSeconds", 15)); - // HTTP pressure: when this many HTTP workers are busy, throttle BG to 1 + // HTTP pressure: when this many HTTP workers are busy, throttle BG to HttpPressureConcurrency var httpPoolSize = Math.Max(1, settings.Worker.HttpPoolSize); HttpPressureThreshold = configuration.GetValue("BackgroundHttpPressureThreshold", - Math.Max(1, httpPoolSize / 2)); + Math.Max(1, httpPoolSize * 2 / 3)); + HttpPressureConcurrency = Math.Clamp( + configuration.GetValue("BackgroundHttpPressureConcurrency", CeilingConcurrency / 2), + Math.Min(2, CeilingConcurrency), CeilingConcurrency); // How long HTTP pressure must persist before throttling HttpPressureAfter = TimeSpan.FromSeconds( @@ -164,9 +173,9 @@ public BackgroundTaskLimiter(ILogger logger, IConfigurati _monitorTimer = new Timer(MonitorCallback, null, TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10)); _logger.LogInformation("[System] Limiter init: baseline={Base} ceiling={Ceiling} scaleAfter={ScaleAfter}s " + - "burstToCeiling={Burst} overSubscribe={Over} httpPressureThreshold={HttpThreshold} httpPressureAfter={HttpAfter}s cpus={Cpus}", + "burstToCeiling={Burst} overSubscribe={Over} httpPressureThreshold={HttpThreshold} httpPressureConcurrency={HttpConc} httpPressureAfter={HttpAfter}s cpus={Cpus}", BaseConcurrency, CeilingConcurrency, ScaleUpAfter.TotalSeconds, _burstToCeiling, _overSubscribe, - HttpPressureThreshold, HttpPressureAfter.TotalSeconds, Environment.ProcessorCount); + HttpPressureThreshold, HttpPressureConcurrency, HttpPressureAfter.TotalSeconds, Environment.ProcessorCount); } /// @@ -459,12 +468,12 @@ private void CheckHttpPressure() var duration = DateTime.UtcNow - _httpPressureSince.Value; if (duration >= HttpPressureAfter) { - // Sustained HTTP pressure — throttle BG to minimum (2) + // Sustained HTTP pressure — throttle BG to HttpPressureConcurrency _logger.LogInformation("[System] Limiter: HTTP pressure detected ({Busy}/{Total} workers busy for {Sec}s), " + - "throttling BG to 2", httpBusy, _pool.HttpPoolSize, duration.TotalSeconds); + "throttling BG to {Target}", httpBusy, _pool.HttpPoolSize, duration.TotalSeconds, HttpPressureConcurrency); _httpThrottled = true; _queuePressureSince = null; // reset BG scale-up tracking - ScaleDown(2, "HTTP pressure"); + ScaleDown(HttpPressureConcurrency, "HTTP pressure"); } } else if (!underPressure && _httpThrottled) diff --git a/Services/Orchestration/JobQueuePump.cs b/Services/Orchestration/JobQueuePump.cs index 76b3ca5..a5e030c 100644 --- a/Services/Orchestration/JobQueuePump.cs +++ b/Services/Orchestration/JobQueuePump.cs @@ -43,12 +43,17 @@ public class JobQueuePump : BackgroundService /// left instead of re-reading every row every tick. private readonly Dictionary _leaseExpiry = new(StringComparer.Ordinal); + /// Completes when startup recovery is done; nothing is claimed before it. Null (tests, a host + /// without an orchestrator) claims immediately. + private readonly Task? _claimGate; + public JobQueuePump(ILogger logger, JobQueueStore queue, JobManager jobs, - IConfiguration configuration, CraftSettings settings) + IConfiguration configuration, CraftSettings settings, OrchestratorService? orchestrator = null) { _logger = logger; _queue = queue; _jobs = jobs; + _claimGate = orchestrator?.RecoveryDone; // Identifies this instance's claims. The container id is stable for the life of the process and // distinct per instance, which is exactly the scope a lease needs. @@ -86,6 +91,18 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken) _owner, _batchSize, _lowWater, _lease.TotalSeconds, _pollInterval.TotalMilliseconds, _idlePollInterval.TotalMilliseconds); + // No claim before startup recovery has run. A row claimed earlier rehydrates its run from storage — + // stale Running markers from the previous process included — into the live graph ahead of recovery, + // whose reset then lands on a copy that loses the _activeRuns race. Rows enqueued meanwhile (timer or + // HTTP-started runs) just wait in the table and are claimed on the first cycle after. + if (_claimGate is { IsCompleted: false }) + { + _logger.LogInformation("[JobQueuePump] Waiting for startup recovery before claiming"); + try { await _claimGate.WaitAsync(stoppingToken); } + catch (OperationCanceledException) { return; } + _logger.LogInformation("[JobQueuePump] Startup recovery done — claiming"); + } + var idleTicks = 0; while (!stoppingToken.IsCancellationRequested) diff --git a/Services/Orchestration/OrchestratorService.cs b/Services/Orchestration/OrchestratorService.cs index 3d714b2..bc7843f 100644 --- a/Services/Orchestration/OrchestratorService.cs +++ b/Services/Orchestration/OrchestratorService.cs @@ -155,6 +155,19 @@ public class OrchestratorService : IJobDescriptorStateWriter ?.Name; } + /// + /// Completes once startup recovery has finished — or been abandoned, see . + /// claims nothing before it. A claim taken earlier rehydrates the run into + /// _activeRuns ahead of recovery, recovery's own copy then loses the TryAdd, and the live graph + /// keeps the dead process's stale Running markers instead of recovery's reset. + /// + public Task RecoveryDone => _recoveryDone.Task; + private readonly TaskCompletionSource _recoveryDone = new(TaskCreationOptions.RunContinuationsAsynchronously); + + /// Open the claim gate. Called from a finally, so a recovery that throws or never runs + /// (shutdown mid-startup, storage down) still releases the pump rather than wedging it. + public void MarkRecoveryDone() => _recoveryDone.TrySetResult(); + private static readonly JsonSerializerOptions s_jsonOptions = new() { WriteIndented = true, @@ -1094,7 +1107,7 @@ private Func BuildSequentialRunWork(OrchestratorRun run lock (_lock) { task.Parameters ??= rehydrated; } } - lock (_lock) { task.Status = "Running"; } + lock (_lock) { task.Status = "Running"; task.OwnedHere = true; } // Durable "Running" marker — awaited before the invoke, same as the parallel path. try { @@ -1402,11 +1415,32 @@ private void FailTaskTerminally(OrchestratorRun run, OrchestratorTaskItem task, // Intune collection that was re-claimed four minutes in and ran twice. This does not block crash // recovery: ResumeInterruptedRunsAsync flips interrupted tasks from Running back to Pending // before re-dispatching them, so a task that genuinely needs re-running never reaches here as - // Running. Within a live process, Running means a worker has it. + // Running. + // + // Running only means "a worker HERE has it" when this process wrote it (OwnedHere). A Running read + // from storage is another process's pre-invoke marker, and it reaches the live graph whenever a run + // is rehydrated — by this resolver, or by the pump winning the _activeRuns race with recovery at + // startup, or from a container that outlived this one's recovery. Dropping those left the task + // Running forever: nothing re-drives Running, so the run never finalized. Reaching here means this + // process holds the task's queue claim (only the pump enqueues descriptors, only for rows it + // claimed), so the previous owner is gone: count the interrupted attempt exactly as recovery does, + // and run it. + var poisoned = false; lock (_lock) { - if (task.Status is "Completed" or "Failed" or "Cancelled" or "Running") + if (task.Status is "Completed" or "Failed" or "Cancelled") return null; + if (task.Status == "Running") + { + if (task.OwnedHere) return null; + poisoned = ++task.AttemptCount >= 3; + if (!poisoned) task.Status = "Pending"; + } + } + if (poisoned) + { + FailTaskTerminally(run, task, $"Cancelled {task.AttemptCount} times by host interruption"); + return null; } // Sequential runs are dispatched as a single entry row; that one claim drives the WHOLE run on one @@ -1471,6 +1505,7 @@ private Func BuildTaskWork(OrchestratorRun run, Orchest lock (_lock) { task.Status = "Running"; + task.OwnedHere = true; } // Pre-script "Running" write is awaited — the durability marker for crash recovery. Batched // across concurrently-starting tasks by the status writer, but still durable before the invoke. diff --git a/Services/Orchestration/OrchestratorTaskItem.cs b/Services/Orchestration/OrchestratorTaskItem.cs index 504595c..3df1994 100644 --- a/Services/Orchestration/OrchestratorTaskItem.cs +++ b/Services/Orchestration/OrchestratorTaskItem.cs @@ -26,4 +26,11 @@ public class OrchestratorTaskItem /// Sequence order. Non-sequential runs leave it 0 and ignore it. /// public int Sequence { get; set; } + + /// + /// Set when a worker in THIS process marks the task Running. Never persisted, so a Running status + /// rehydrated from storage (another process's pre-invoke marker) reads false: it is not proof that + /// anything here is executing the task. See OrchestratorService.ResolveTaskWorkAsync. + /// + internal bool OwnedHere; } diff --git a/Services/Orchestration/SchedulerService.cs b/Services/Orchestration/SchedulerService.cs index a8e84d6..2a10fee 100644 --- a/Services/Orchestration/SchedulerService.cs +++ b/Services/Orchestration/SchedulerService.cs @@ -58,40 +58,49 @@ public SchedulerService( protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - _logger.LogInformation("[Scheduler] Waiting for worker pool to be ready"); + // Every exit from startup — recovery done, recovery threw, storage never came up, shutdown mid-way — + // opens the claim gate, so the pump can never be left waiting on a recovery that is not coming. try { - await Task.Run(() => _pool.WaitForBgReady(Timeout.InfiniteTimeSpan), stoppingToken); - } - catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) - { - return; - } - _logger.LogInformation("[Scheduler] Worker pool ready — service starting"); + _logger.LogInformation("[Scheduler] Waiting for worker pool to be ready"); + try + { + await Task.Run(() => _pool.WaitForBgReady(Timeout.InfiniteTimeSpan), stoppingToken); + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + return; + } + _logger.LogInformation("[Scheduler] Worker pool ready — service starting"); - LoadConfig(); - ResolveTimezone(); + LoadConfig(); + ResolveTimezone(); - // Seed all tasks with "now" so we never catch up on missed runs from before startup. - // Only cron ticks that occur AFTER this moment will fire. - var startup = DateTimeOffset.UtcNow; - foreach (var task in _tasks) - { - _lastRun[task.Id] = startup; - } + // Seed all tasks with "now" so we never catch up on missed runs from before startup. + // Only cron ticks that occur AFTER this moment will fire. + var startup = DateTimeOffset.UtcNow; + foreach (var task in _tasks) + { + _lastRun[task.Id] = startup; + } - // Wait for the storage backend to accept connections before the first orchestrator store access. - // Avoids a startup error when storage becomes reachable a moment after the app. - await _storageHealth.WaitUntilReadyAsync(TimeSpan.FromSeconds(60), stoppingToken); + // Wait for the storage backend to accept connections before the first orchestrator store access. + // Avoids a startup error when storage becomes reachable a moment after the app. + await _storageHealth.WaitUntilReadyAsync(TimeSpan.FromSeconds(60), stoppingToken); - // Resume any orchestrator runs that were interrupted by a previous crash - try - { - await _orchestrator.ResumeInterruptedRunsAsync(stoppingToken); + // Resume any orchestrator runs that were interrupted by a previous crash + try + { + await _orchestrator.ResumeInterruptedRunsAsync(stoppingToken); + } + catch (Exception ex) + { + _logger.LogError(ex, "[Scheduler] Failed to resume interrupted orchestrator runs"); + } } - catch (Exception ex) + finally { - _logger.LogError(ex, "[Scheduler] Failed to resume interrupted orchestrator runs"); + _orchestrator.MarkRecoveryDone(); } // Retention sweeps for the rest of the process lifetime; recovery ran the first one. Fire and diff --git a/tests/Craft.Tests/LimiterHttpPressureTests.cs b/tests/Craft.Tests/LimiterHttpPressureTests.cs new file mode 100644 index 0000000..345d182 --- /dev/null +++ b/tests/Craft.Tests/LimiterHttpPressureTests.cs @@ -0,0 +1,55 @@ +using System.Globalization; +using Craft.Configuration; +using Craft.Orchestration; +using Craft.PowerShellHost; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Craft.Tests; + +/// +/// HTTP-pressure defaults. Half the HTTP pool (3 of 6) tripped on routine API-client traffic and the +/// fixed throttle target of 2 then held an 8-slot BG pool at 2 for hours, so a scheduled cache run of +/// ~5.5k tasks could not finish inside a day. +/// +public class LimiterHttpPressureTests +{ + private static BackgroundTaskLimiter NewLimiter(int httpPoolSize, int bgPoolSize, Dictionary? extra = null) + { + var settings = new CraftSettings(); + settings.Worker.HttpPoolSize = httpPoolSize; + settings.Worker.BgPoolSize = bgPoolSize; + var values = new Dictionary + { + ["BackgroundMaxConcurrency"] = bgPoolSize.ToString(CultureInfo.InvariantCulture), + }; + foreach (var kv in extra ?? []) values[kv.Key] = kv.Value; + var config = new ConfigurationBuilder().AddInMemoryCollection(values).Build(); + var repo = new ScriptRepository(NullLogger.Instance, settings); + var pool = new PowerShellWorkerPool(repo, NullLogger.Instance, config, settings); + return new BackgroundTaskLimiter(NullLogger.Instance, config, settings, pool); + } + + [Theory] + [InlineData(6, 8, 4, 4)] // hosted default profile + [InlineData(2, 2, 1, 2)] // Basic 1-cpu profile: target never exceeds the ceiling + [InlineData(16, 24, 10, 12)] + public void DefaultsScaleWithPools(int http, int bg, int threshold, int target) + { + using var limiter = NewLimiter(http, bg); + Assert.Equal(threshold, limiter.HttpPressureThreshold); + Assert.Equal(target, limiter.HttpPressureConcurrency); + } + + [Fact] + public void ConfigOverridesAreHonouredAndClampedToTheCeiling() + { + using var limiter = NewLimiter(6, 8, new() + { + ["BackgroundHttpPressureThreshold"] = "5", + ["BackgroundHttpPressureConcurrency"] = "50", + }); + Assert.Equal(5, limiter.HttpPressureThreshold); + Assert.Equal(8, limiter.HttpPressureConcurrency); + } +} diff --git a/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs b/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs index 6be6fc6..f519476 100644 --- a/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs +++ b/tests/Craft.Tests/OrchestratorFinalizedRunTests.cs @@ -290,7 +290,8 @@ public async Task DescriptorForATaskThatIsAlreadyRunning_IsDropped() Status = "Running", Priority = 4, StartedUtc = DateTime.UtcNow, - Tasks = [new OrchestratorTaskItem { Id = "task-0", Status = "Running" }] + // OwnedHere: a worker in THIS process is executing it — the state dispatch leaves behind. + Tasks = [new OrchestratorTaskItem { Id = "task-0", Status = "Running", OwnedHere = true }] }; await store.UpsertRunAsync(run); await store.UpsertTaskAsync("live-run", run.Tasks[0]); diff --git a/tests/Craft.Tests/OrchestratorStaleRunningTests.cs b/tests/Craft.Tests/OrchestratorStaleRunningTests.cs new file mode 100644 index 0000000..2d6c5fb --- /dev/null +++ b/tests/Craft.Tests/OrchestratorStaleRunningTests.cs @@ -0,0 +1,207 @@ +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; + +/// +/// A "Running" status this process did not write must not be trusted as "a worker here has it". +/// +/// The resolver drops a descriptor whose task reads Running, which is right for a duplicate queue row +/// claimed while THIS process is executing the task. But Running also reaches the live graph from +/// storage: the durable pre-invoke marker another process wrote before it died. Two production paths +/// put it there, both seen on a hosted instance across an App Service container swap (old and new +/// containers overlap for tens of seconds to minutes, sharing one storage account): +/// +/// A. The old container creates or advances a run after the new container's startup recovery has +/// already passed. It dies holding claims whose tasks are marked Running. The new container's pump +/// claims one of the run's rows, ResolveTaskWorkAsync rehydrates the run from storage with +/// those markers, and once the dead claims' leases lapse and those rows are claimed, the resolver +/// drops each one as "already running" — Skipped, row deleted. Nothing ever re-drives a Running +/// task, and the run can never finalize (observed: 505/508 for 31 hours). +/// +/// B. The pump starts claiming at host start, before ResumeInterruptedRunsAsync (which waits +/// for the worker pool). Its rehydrated copy goes into _activeRuns first; recovery then +/// loads its OWN copy, flips Running→Pending on that copy and in storage, releases the dead claims +/// and calls DispatchPendingTasksAsync, whose _activeRuns.TryAdd loses to the pump's copy. +/// The live graph keeps the stale Running; the released row is claimed and dropped exactly as in A. +/// +/// Holding the queue claim is the ownership proof: only this process's pump enqueues descriptor jobs, +/// and it only does so for a row it has claimed. So a Running task that no worker in this process +/// started is an interrupted task, and is handled the way recovery handles one — attempt counted +/// (poison bound kept) and run. +/// +public class OrchestratorStaleRunningTests +{ + private const string TaskFunc = "Invoke-CraftTask"; + + private sealed record Harness(OrchestratorService Svc, OrchestratorTableStore Store, JobQueueStore Queue, + ConcurrentDictionary Active); + + private static async Task NewHarnessAsync() + { + var settings = new CraftSettings { Orchestrator = { TablePrefix = "stale" + Guid.NewGuid().ToString("N")[..8] } }; + settings.Orchestrator.BatchStatusWrites = false; + var config = new ConfigurationBuilder().AddInMemoryCollection([]).Build(); + + 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); + + var repo = new ScriptRepository(NullLogger.Instance, settings); + typeof(ScriptRepository).GetField("_moduleFunctionNames", BindingFlags.NonPublic | BindingFlags.Instance)! + .SetValue(repo, new HashSet([TaskFunc], StringComparer.OrdinalIgnoreCase)); + var pool = new PowerShellWorkerPool(repo, NullLogger.Instance, config, settings); + var runner = new PowerShellRunnerService(NullLogger.Instance, pool, repo, settings); + + var svc = (OrchestratorService)System.Runtime.CompilerServices.RuntimeHelpers + .GetUninitializedObject(typeof(OrchestratorService)); + var active = new ConcurrentDictionary(); + Set(svc, "_logger", NullLogger.Instance); + Set(svc, "_store", store); + Set(svc, "_queue", queue); + Set(svc, "_writer", writer); + Set(svc, "_psRunner", runner); + Set(svc, "_settings", settings); + Set(svc, "_lock", new object()); + Set(svc, "_activeRuns", active); + Set(svc, "_taskScriptPaths", new ConcurrentDictionary()); + Set(svc, "_finalizingRuns", new ConcurrentDictionary()); + Set(svc, "_finalizeDeferrals", new ConcurrentDictionary()); + Set(svc, "_childRuns", new ConcurrentDictionary>()); + Set(svc, "_recoveringChildren", new ConcurrentDictionary()); + Set(svc, "_pendingChildRuns", new ConcurrentDictionary()); + Set(svc, "_cancelledRuns", new ConcurrentDictionary()); + Set(svc, "_activeSequentialDrivers", new ConcurrentDictionary()); + Set(svc, "_requeueFailures", new ConcurrentDictionary()); + Set(svc, "_deferrals", NewFieldValue(svc, "_deferrals")); + Set(svc, "_redriveBackoff", NewFieldValue(svc, "_redriveBackoff")); + Set(svc, "_shedParameters", false); + 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, active); + } + + private static object NewFieldValue(object svc, string field) => + Activator.CreateInstance(typeof(OrchestratorService) + .GetField(field, BindingFlags.NonPublic | BindingFlags.Instance)!.FieldType)!; + + 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 run, string taskId) + { + var mi = typeof(OrchestratorService).GetMethod("ResolveTaskWorkAsync", BindingFlags.NonPublic | BindingFlags.Instance)!; + try + { + var task = (Task)mi.Invoke(svc, [new JobDescriptor(run, taskId, 4), CancellationToken.None])!; + await task; + return task.GetType().GetProperty("Result")!.GetValue(task); + } + catch (TargetInvocationException ex) { throw ex.InnerException ?? ex; } + } + + /// Storage as another process left it: its in-flight task carries the durable Running marker. + private static async Task SeedAsync(Harness h, string name, int staleAttempts = 0) + { + var tasks = new List + { + new() { Id = "t-stale", Status = "Running", AttemptCount = staleAttempts, Parameters = new() { ["FunctionName"] = "Push-Noop" } }, + new() { Id = "t-next", Status = "Pending", Parameters = new() { ["FunctionName"] = "Push-Noop" } }, + new() { Id = "t-other", Status = "Pending", Parameters = new() { ["FunctionName"] = "Push-Noop" } }, + }; + await h.Store.UpsertRunAsync(new OrchestratorRun + { + Name = name, + Status = "Running", + Priority = 4, + StartedUtc = DateTime.UtcNow, + TaskScriptName = TaskFunc, + Tasks = tasks + }); + foreach (var t in tasks) await h.Store.UpsertTaskAsync(name, t); + } + + // ── Mode A ───────────────────────────────────────────────────────────────────────────────────── + + [Fact] + public async Task ModeA_RunningMarkerFromADeadProcess_IsRunWhenThisProcessClaimsItsRow() + { + var h = await NewHarnessAsync(); + await SeedAsync(h, "run-a"); + + // The dead process's lease lapsed; this process claimed the row. Previously: null (dropped forever). + var work = await ResolveAsync(h.Svc, "run-a", "t-stale"); + + Assert.NotNull(work); + var task = h.Active["run-a"].Tasks.Single(t => t.Id == "t-stale"); + Assert.Equal(1, task.AttemptCount); // the interrupted attempt is counted, as recovery counts it + } + + [Fact] + public async Task ModeA_PoisonBoundHolds_AThirdInterruptedAttemptFailsTheTask() + { + var h = await NewHarnessAsync(); + await SeedAsync(h, "run-poison", staleAttempts: 2); + + Assert.Null(await ResolveAsync(h.Svc, "run-poison", "t-stale")); + + var task = h.Active["run-poison"].Tasks.Single(t => t.Id == "t-stale"); + Assert.Equal("Failed", task.Status); + } + + // ── Mode B ───────────────────────────────────────────────────────────────────────────────────── + + [Fact] + public async Task ModeB_PumpClaimBeforeRecovery_LeavesAStaleRunningInTheLiveGraph_WhichMustStillRun() + { + var h = await NewHarnessAsync(); + await SeedAsync(h, "run-b"); + + // 1. Host start: the pump claims a sibling row before recovery has run. The rehydrated copy — + // carrying the dead process's Running marker — becomes the live graph. + Assert.NotNull(await ResolveAsync(h.Svc, "run-b", "t-next")); + var live = h.Active["run-b"]; + + // 2. Recovery runs on its own copy: flips t-stale to Pending in storage, releases the claims and + // re-dispatches — but its TryAdd loses to the pump's copy. + await h.Svc.ResumeInterruptedRunsAsync(CancellationToken.None); + Assert.Same(live, h.Active["run-b"]); + + // 3. The released row is claimed. Previously: dropped as "already running", never run again. + Assert.NotNull(await ResolveAsync(h.Svc, "run-b", "t-stale")); + } + + // ── The guard's real job is kept ─────────────────────────────────────────────────────────────── + + [Fact] + public async Task ADuplicateRowForATaskThisProcessIsExecuting_IsStillDropped() + { + var h = await NewHarnessAsync(); + await SeedAsync(h, "run-dup"); + await h.Store.UpsertTaskAsync("run-dup", + new OrchestratorTaskItem { Id = "t-stale", Status = "Pending", Parameters = new() { ["FunctionName"] = "Push-Noop" } }); + + // Start the task the way dispatch does, up to the point it is running on a worker here. + var work = (Func)(await ResolveAsync(h.Svc, "run-dup", "t-stale"))!; + var task = h.Active["run-dup"].Tasks.Single(t => t.Id == "t-stale"); + _ = Task.Run(() => work(CancellationToken.None)); // marks Running, then blocks on worker checkout + for (var i = 0; i < 200 && task.Status != "Running"; i++) await Task.Delay(10); + Assert.Equal("Running", task.Status); + + // A second row for the same task, claimed mid-flight, must not start another copy. + Assert.Null(await ResolveAsync(h.Svc, "run-dup", "t-stale")); + } +} diff --git a/tests/Craft.Tests/StartupClaimGateTests.cs b/tests/Craft.Tests/StartupClaimGateTests.cs new file mode 100644 index 0000000..7f9a9ec --- /dev/null +++ b/tests/Craft.Tests/StartupClaimGateTests.cs @@ -0,0 +1,140 @@ +using System.Reflection; +using System.Runtime.CompilerServices; +using Craft.Configuration; +using Craft.Endpoints; +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 pump claims nothing until startup recovery is done, and every way startup can end opens the gate. +/// +/// A claim taken before recovery rehydrates its run from storage — the previous process's Running markers +/// included — into the live graph. Recovery then resets those markers on its own copy, which loses the +/// _activeRuns race, so the live graph keeps the stale Running (see OrchestratorStaleRunningTests, +/// mode B). Gating the first claim on recovery removes the race; the resolver's ownership guard stays as +/// defense in depth. +/// +public class StartupClaimGateTests +{ + /// An orchestrator carrying only the gate. Field initializers do not run on an uninitialized + /// object, so the gate is installed by hand; every other field stays null. + private static OrchestratorService GatedOrchestrator() + { + var svc = (OrchestratorService)RuntimeHelpers.GetUninitializedObject(typeof(OrchestratorService)); + typeof(OrchestratorService).GetField("_recoveryDone", BindingFlags.NonPublic | BindingFlags.Instance)! + .SetValue(svc, new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously)); + typeof(OrchestratorService).GetField("_logger", BindingFlags.NonPublic | BindingFlags.Instance)! + .SetValue(svc, NullLogger.Instance); + return svc; + } + + private static (JobQueuePump Pump, JobManager Jobs) NewPump(OrchestratorService orchestrator) + { + var settings = new CraftSettings(); + settings.Worker.BgPoolSize = 4; + var config = new ConfigurationBuilder().AddInMemoryCollection(new Dictionary + { + ["JobQueuePollIntervalMs"] = "100", + }).Build(); + + var queue = new JobQueueStore(NullLogger.Instance, settings, new RunRemainingCounterTests.ConditionalStore()); + queue.InitializeAsync().GetAwaiter().GetResult(); + queue.EnqueueBatchAsync("StandardsApply", + Enumerable.Range(0, 20).Select(i => ($"task-{i:D3}", 4)).ToList(), DateTime.UtcNow).GetAwaiter().GetResult(); + + var repo = new ScriptRepository(NullLogger.Instance, settings); + var pool = new PowerShellWorkerPool(repo, NullLogger.Instance, config, settings); + var limiter = new BackgroundTaskLimiter(NullLogger.Instance, config, settings, pool); + var jobs = new JobManager(NullLogger.Instance, settings, limiter); + return (new JobQueuePump(NullLogger.Instance, queue, jobs, config, settings, orchestrator), jobs); + } + + private static async Task WaitUntil(Func condition, int timeoutMs = 5000) + { + var deadline = Environment.TickCount64 + timeoutMs; + while (Environment.TickCount64 < deadline) + { + if (condition()) return true; + await Task.Delay(20); + } + return condition(); + } + + [Fact] + public async Task PumpClaimsNothingBeforeRecovery_AndClaimsNormallyAfter() + { + var orchestrator = GatedOrchestrator(); + var (pump, jobs) = NewPump(orchestrator); + + await pump.StartAsync(CancellationToken.None); + try + { + await Task.Delay(500); // five poll intervals with rows waiting + Assert.Equal(0, jobs.QueuedCount); + + orchestrator.MarkRecoveryDone(); + + Assert.True(await WaitUntil(() => jobs.QueuedCount > 0), "the pump did not claim once recovery was done"); + } + finally + { + await Task.WhenAny(pump.StopAsync(CancellationToken.None), Task.Delay(3000)); + } + } + + private static SchedulerService NewScheduler(OrchestratorService orchestrator, bool poolReady) + { + var settings = new CraftSettings(); + var config = new ConfigurationBuilder().AddInMemoryCollection([]).Build(); + var repo = new ScriptRepository(NullLogger.Instance, settings); + var pool = new PowerShellWorkerPool(repo, NullLogger.Instance, config, settings); + if (poolReady) pool.Initialize(enableHttp: false, enableBg: false); // signals ready, builds nothing + var limiter = new BackgroundTaskLimiter(NullLogger.Instance, config, settings, pool); + var runner = new PowerShellRunnerService(NullLogger.Instance, pool, repo, settings); + var health = new StorageHealthMonitor(new RunRemainingCounterTests.ConditionalStore(), NullLogger.Instance); + return new SchedulerService(NullLogger.Instance, runner, limiter, orchestrator, + new JobManager(NullLogger.Instance, settings, limiter), settings, pool, health, + NativeScheduledTasks.Empty, null!); + } + + [Fact] + public async Task ARecoveryThatThrows_StillOpensTheGate() + { + // The gated orchestrator has no store, so ResumeInterruptedRunsAsync throws on its first line. + var orchestrator = GatedOrchestrator(); + var scheduler = NewScheduler(orchestrator, poolReady: true); + using var cts = new CancellationTokenSource(); + + await scheduler.StartAsync(cts.Token); + try + { + Assert.True(await Task.WhenAny(orchestrator.RecoveryDone, Task.Delay(10_000)) == orchestrator.RecoveryDone, + "a failed recovery left the claim gate shut — the pump would never claim"); + } + finally + { + cts.Cancel(); + try { await scheduler.StopAsync(CancellationToken.None); } catch { /* the gutted orchestrator may fault the loop */ } + } + } + + [Fact] + public async Task ShutdownBeforeTheWorkerPoolIsReady_StillOpensTheGate() + { + var orchestrator = GatedOrchestrator(); + var scheduler = NewScheduler(orchestrator, poolReady: false); + using var cts = new CancellationTokenSource(); + + await scheduler.StartAsync(cts.Token); + cts.Cancel(); + try { await scheduler.StopAsync(CancellationToken.None); } catch { /* cancellation */ } + + Assert.True(await Task.WhenAny(orchestrator.RecoveryDone, Task.Delay(5_000)) == orchestrator.RecoveryDone); + } +}