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
7 changes: 5 additions & 2 deletions Runtime/CraftRuntime/Invoke-CraftTask.ps1
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand All @@ -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"
}
31 changes: 29 additions & 2 deletions Runtime/CraftRuntime/Start-CraftOrchestrator.ps1
Original file line number Diff line number Diff line change
Expand Up @@ -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 = @(
Expand All @@ -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 = '[email protected]'
Sequential = $true
Batch = @(
@{ FunctionName = 'RevokeSessions'; User = '[email protected]' }
@{ FunctionName = 'GrantMailboxAccess'; User = '[email protected]'; Delegate = '[email protected]' }
@{ FunctionName = 'ConvertToShared'; User = '[email protected]' }
)
}

.FUNCTIONALITY
Internal
#>
Expand Down Expand Up @@ -120,15 +142,20 @@ 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,
$Priority,
$PostExecFunctionName,
$PostExecParametersJson,
$InputObject.Reference,
$ParentRunName
$ParentRunName,
$Sequential
)
return "Craft-$OrchestratorName"
}
14 changes: 7 additions & 7 deletions Services/Bridges/OrchestratorBridge.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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));
}

/// <summary>
Expand All @@ -51,14 +51,14 @@ public static void QueueOrchestration(string name, string batchJson, int priorit
/// </summary>
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));
}

/// <summary>
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
{
Expand Down Expand Up @@ -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<PendingPlannerRun> s_pendingPlanners = new();

Expand Down
24 changes: 24 additions & 0 deletions Services/Configuration/OrchestratorSettings.cs
Original file line number Diff line number Diff line change
Expand Up @@ -103,4 +103,28 @@ public class OrchestratorSettings
/// periodic sweep; the startup pass, which follows crash recovery, still runs.
/// </summary>
public int CleanupIntervalHours { get; set; } = 4;

/// <summary>
/// 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.
/// </summary>
public int StatusTimerIntervalSeconds { get; set; } = 60;

/// <summary>
/// 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.
/// </summary>
public bool RedriveBackoff { get; set; } = true;

/// <summary>
/// Whether a task sheds its <c>Parameters</c> 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.
/// </summary>
public bool ShedPendingParameters { get; set; } = true;
}
44 changes: 26 additions & 18 deletions Services/Hosting/StatsHistoryService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 9 additions & 0 deletions Services/Orchestration/OrchestratorRun.cs
Original file line number Diff line number Diff line change
Expand Up @@ -28,4 +28,13 @@ public class OrchestratorRun
public int PostExecAttemptCount { get; set; }

public string? ParentRunName { get; set; }

/// <summary>
/// Sequential execution mode. When true the run's tasks are dispatched ONE AT A TIME, in ascending
/// <see cref="OrchestratorTaskItem.Sequence"/> (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.
/// </summary>
public bool Sequential { get; set; }
}
Loading
Loading