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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 20 additions & 11 deletions Services/Orchestration/BackgroundTaskLimiter.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -97,12 +97,18 @@ public class BackgroundTaskLimiter : IDisposable
/// <summary>
/// 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 <see cref="HttpPressureConcurrency"/>.
/// 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.
/// </summary>
public int HttpPressureThreshold { get; }

/// <summary>
/// BG concurrency while HTTP-throttled. Default: half the ceiling, minimum 2.
/// </summary>
public int HttpPressureConcurrency { get; }

/// <summary>
/// How long HTTP pressure must be sustained before throttling BG tasks.
/// Default: 10 seconds.
Expand Down Expand Up @@ -132,10 +138,13 @@ public BackgroundTaskLimiter(ILogger<BackgroundTaskLimiter> 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(
Expand Down Expand Up @@ -164,9 +173,9 @@ public BackgroundTaskLimiter(ILogger<BackgroundTaskLimiter> 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);
}

/// <summary>
Expand Down Expand Up @@ -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)
Expand Down
19 changes: 18 additions & 1 deletion Services/Orchestration/JobQueuePump.cs
Original file line number Diff line number Diff line change
Expand Up @@ -43,12 +43,17 @@ public class JobQueuePump : BackgroundService
/// left instead of re-reading every row every tick.</summary>
private readonly Dictionary<string, DateTime> _leaseExpiry = new(StringComparer.Ordinal);

/// <summary>Completes when startup recovery is done; nothing is claimed before it. Null (tests, a host
/// without an orchestrator) claims immediately.</summary>
private readonly Task? _claimGate;

public JobQueuePump(ILogger<JobQueuePump> 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.
Expand Down Expand Up @@ -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)
Expand Down
41 changes: 38 additions & 3 deletions Services/Orchestration/OrchestratorService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,19 @@ public class OrchestratorService : IJobDescriptorStateWriter
?.Name;
}

/// <summary>
/// Completes once startup recovery has finished — or been abandoned, see <see cref="MarkRecoveryDone"/>.
/// <see cref="JobQueuePump"/> claims nothing before it. A claim taken earlier rehydrates the run into
/// <c>_activeRuns</c> 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.
/// </summary>
public Task RecoveryDone => _recoveryDone.Task;
private readonly TaskCompletionSource _recoveryDone = new(TaskCreationOptions.RunContinuationsAsynchronously);

/// <summary>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.</summary>
public void MarkRecoveryDone() => _recoveryDone.TrySetResult();

private static readonly JsonSerializerOptions s_jsonOptions = new()
{
WriteIndented = true,
Expand Down Expand Up @@ -1094,7 +1107,7 @@ private Func<CancellationToken, Task> 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
{
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -1471,6 +1505,7 @@ private Func<CancellationToken, Task> 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.
Expand Down
7 changes: 7 additions & 0 deletions Services/Orchestration/OrchestratorTaskItem.cs
Original file line number Diff line number Diff line change
Expand Up @@ -26,4 +26,11 @@ public class OrchestratorTaskItem
/// Sequence order. Non-sequential runs leave it 0 and ignore it.
/// </summary>
public int Sequence { get; set; }

/// <summary>
/// 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 <c>OrchestratorService.ResolveTaskWorkAsync</c>.
/// </summary>
internal bool OwnedHere;
}
61 changes: 35 additions & 26 deletions Services/Orchestration/SchedulerService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
55 changes: 55 additions & 0 deletions tests/Craft.Tests/LimiterHttpPressureTests.cs
Original file line number Diff line number Diff line change
@@ -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;

/// <summary>
/// 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.
/// </summary>
public class LimiterHttpPressureTests
{
private static BackgroundTaskLimiter NewLimiter(int httpPoolSize, int bgPoolSize, Dictionary<string, string?>? extra = null)
{
var settings = new CraftSettings();
settings.Worker.HttpPoolSize = httpPoolSize;
settings.Worker.BgPoolSize = bgPoolSize;
var values = new Dictionary<string, string?>
{
["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<ScriptRepository>.Instance, settings);
var pool = new PowerShellWorkerPool(repo, NullLogger<PowerShellWorkerPool>.Instance, config, settings);
return new BackgroundTaskLimiter(NullLogger<BackgroundTaskLimiter>.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);
}
}
3 changes: 2 additions & 1 deletion tests/Craft.Tests/OrchestratorFinalizedRunTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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]);
Expand Down
Loading
Loading