Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
cf24560
fix(orchestrator): bound re-drive verifications to one per run and ei…
Zacgoose Oct 5, 2026
16f687d
perf(storage): find split-entity part rows by RowKey range and only t…
Zacgoose Oct 5, 2026
1464c73
feat(queue): claim oldest run first within a priority band
Zacgoose Oct 5, 2026
c98a091
fix(queue): key runs by start time and keep queue reads flat at scale
Zacgoose Oct 5, 2026
5055813
feat(orchestration): keep all run and task state in storage
Zacgoose Oct 5, 2026
e636885
feat(orchestration): let runs of one name overlap unless the caller o…
Zacgoose Oct 5, 2026
d00b4b7
feat(runtime): pass AllowCollision and the parent run key from Start-…
Zacgoose Oct 5, 2026
a4d72cf
feat(orchestration): report runs skipped because a run of that name i…
Zacgoose Oct 5, 2026
4f902e1
feat(orchestration): add MaxConcurrency and sequential StopOnFailure,…
Zacgoose Oct 5, 2026
ceda47b
fix(jobs): finish the claim of a job cancelled just after the dispatc…
Zacgoose Oct 5, 2026
8667f17
test(orchestration): pin end-to-end table calls per task and per smal…
Zacgoose Oct 5, 2026
7cc9a3b
fix(powershell): keep a call's output when the previous call on the w…
Zacgoose Oct 5, 2026
d0e22d5
feat(orchestration): one instance works the queue, with repairable in…
Zacgoose Oct 5, 2026
edc5a87
test(e2e): add orchestration engine regression section
Zacgoose Oct 5, 2026
97e87ec
fix(status): derive every waiting-in-storage count from live holdings…
Zacgoose Oct 5, 2026
1854b76
fix(orchestration): make bulk cancel fast and exact, and tighten the …
Zacgoose Oct 5, 2026
bbdf0bb
fix(orchestration): never run a claim handed back at shutdown, and lo…
Zacgoose Oct 6, 2026
f584d4b
test(perf): add orchestration benchmark for comparing engines on Azur…
Zacgoose Oct 6, 2026
f400c87
perf(orchestration): finish a sequential step and claim the next in o…
Zacgoose Oct 6, 2026
28e742b
fix(orchestration): treat queue-id suffixed run names as one name for…
Zacgoose Oct 6, 2026
e04011c
fix(orchestration): show every step of a sequential run in worker sta…
Zacgoose Oct 6, 2026
c296552
feat(auth): keep the token's claims on the normalised principal
Zacgoose Oct 6, 2026
2a30801
feat(memory): trim memory on a timer and read memory detail from the …
Zacgoose Oct 6, 2026
0b1db91
style(tests): sort usings in OrchestratorBridgeLineageTests
Zacgoose Oct 6, 2026
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
30 changes: 28 additions & 2 deletions Runtime/CraftRuntime/Start-CraftOrchestrator.ps1
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,13 @@ function Start-CraftOrchestrator {
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).
- AllowCollision (bool) — optional, default $true: runs of one name stack up side by side.
$false skips this run while another run of the same name is still
going (recurring work that must not pile up).
- MaxConcurrency (int) — optional, default 0 (no limit): at most this many of the run's tasks
run at once. Not used with Sequential.
- StopOnFailure (bool) — optional, Sequential only: the first failed step cancels the steps
after it. By default a sequential run carries on past a failure.

.EXAMPLE
# Fan-out (default): every task is queued up front and drained in parallel by the worker pool.
Expand Down Expand Up @@ -68,6 +75,13 @@ function Start-CraftOrchestrator {

$OrchestratorName = $InputObject.OrchestratorName ?? 'UnnamedOrchestrator'

# Collisions off: a run of this name that is still going wins, and this one is skipped up front.
$AllowCollision = $InputObject.AllowCollision -ne $false
if (-not $AllowCollision -and [Craft.Services.OrchestratorBridge]::IsRunActive($OrchestratorName)) {
Write-Warning "Craft: Skipped orchestrator '$OrchestratorName' - a run with this name is still active"
return "Craft-$OrchestratorName-Skipped"
}

# QueueFunction pattern: call the function first to generate batch items
if (-not $InputObject.Batch -and $InputObject.QueueFunction) {
$QueueFuncName = "Push-$($InputObject.QueueFunction.FunctionName)"
Expand Down Expand Up @@ -140,11 +154,20 @@ function Start-CraftOrchestrator {
# Lineage: pass the enclosing run explicitly. The bridge's own ambient read is null for calls
# made from the pipeline thread — which is exactly where this function runs — so without this
# a parent run would finalize (and dispatch its PostExecution) before its child runs complete.
$ParentRunName = $OpContext.RunName
# RunKey names the exact run when several runs share a name.
$ParentRunName = $OpContext.RunKey ?? $OpContext.RunName

# 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)
$MaxConcurrency = [int]($InputObject.MaxConcurrency ?? 0)
$StopOnFailure = [bool]($InputObject.StopOnFailure)
if ($Sequential -and $MaxConcurrency -gt 0) {
Write-Warning "Craft: MaxConcurrency is ignored for '$OrchestratorName': a sequential run already runs one step at a time"
}
if ($StopOnFailure -and -not $Sequential) {
Write-Warning "Craft: StopOnFailure is ignored for '$OrchestratorName': it applies to sequential runs only"
}

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(
Expand All @@ -155,7 +178,10 @@ function Start-CraftOrchestrator {
$PostExecParametersJson,
$InputObject.Reference,
$ParentRunName,
$Sequential
$Sequential,
$AllowCollision,
$MaxConcurrency,
$StopOnFailure
)
return "Craft-$OrchestratorName"
}
20 changes: 20 additions & 0 deletions Services/Auth/EasyAuthPrincipal.cs
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,26 @@ public static JsonDocument Decode(string headerValue) =>
public static string Encode<T>(T principal) =>
Convert.ToBase64String(Encoding.UTF8.GetBytes(JsonSerializer.Serialize(principal)));

/// <summary>
/// Encodes the normalised SWA-format principal for an EasyAuth <paramref name="source"/>, keeping the
/// token's original <c>claims</c> so the hosted app can read any of them (e.g. <c>azp</c>, <c>scp</c>).
/// </summary>
/// <remarks>
/// The claims are those of the token EasyAuth validated, and <c>userRoles</c> marks the result as
/// already normalised (<see cref="NeedsTransform"/>), so it is never transformed twice.
/// </remarks>
public static string EncodeNormalised(
JsonElement source, string identityProvider, string userId, string userDetails, IReadOnlyList<string> userRoles)
{
var claims = source.ValueKind == JsonValueKind.Object &&
source.TryGetProperty("claims", out var sourceClaims) &&
sourceClaims.ValueKind == JsonValueKind.Array
? sourceClaims.Clone()
: JsonDocument.Parse("[]").RootElement.Clone();

return Encode(new { identityProvider, userId, userDetails, userRoles, claims });
}

/// <summary>
/// Pulls the identity claims out of an EasyAuth principal.
/// </summary>
Expand Down
3 changes: 3 additions & 0 deletions Services/Bridges/JobRecord.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,9 @@ public class JobRecord
public string Id { get; set; } = string.Empty;
public string Name { get; set; } = string.Empty;
public string? RunName { get; set; }

/// <summary>The outing of the run this job belongs to, when it came from a stored run.</summary>
public string? RunKey { get; set; }
public int Priority { get; set; }
public string Status { get; set; } = "Queued";
public DateTime QueuedUtc { get; set; }
Expand Down
164 changes: 97 additions & 67 deletions Services/Bridges/OrchestratorBridge.cs
Original file line number Diff line number Diff line change
Expand Up @@ -23,19 +23,27 @@ public static class OrchestratorBridge

public static void Initialize(OrchestratorService service) => s_service = service;

/// <param name="allowCollision">True (the default) lets runs of one name stack up; false skips this run
/// while another run of the same name is unfinished.</param>
/// <param name="maxConcurrency">At most this many of the run's tasks run at once; 0 (the default) is no
/// limit. Ignored for a sequential run.</param>
/// <param name="stopOnFailure">Sequential runs only: the first failed step cancels the rest instead of the
/// run carrying on (the default).</param>
public static void QueueOrchestration(string name, string batchJson, int priority,
string? postExecFunctionName = null, string? postExecParametersJson = null,
string? reference = null, string? parentRunName = null, bool sequential = false)
string? reference = null, string? parentRunName = null, bool sequential = false, bool allowCollision = true,
int maxConcurrency = 0, bool stopOnFailure = 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
// character would register a child link no live run ever matches.
name = TableKeys.Sanitize(name);
parentRunName = ResolveParentRunName(name, parentRunName);
var gated = RegisterPendingChild(parentRunName, name);
var child = RegisterPendingChild(parentRunName, name);
s_pending.Enqueue(new PendingOrchestration(name, batchJson, priority,
postExecFunctionName, postExecParametersJson, parentRunName, reference,
PendingChildRegistered: gated, Sequential: sequential));
Sequential: sequential, ParentRunKey: child?.ParentRunKey, ChildKey: child?.ChildKey,
AllowCollision: allowCollision, MaxConcurrency: maxConcurrency, StopOnFailure: stopOnFailure));
}

/// <summary>
Expand All @@ -49,20 +57,29 @@ public static void QueueOrchestration(string name, string batchJson, int priorit
///
/// The file is owned by the orchestrator from this point: it is deleted once parsed.
/// </summary>
/// <param name="allowCollision">True (the default) lets runs of one name stack up; false skips this run
/// while another run of the same name is unfinished.</param>
/// <param name="maxConcurrency">At most this many of the run's tasks run at once; 0 (the default) is no
/// limit. Ignored for a sequential run.</param>
/// <param name="stopOnFailure">Sequential runs only: the first failed step cancels the rest instead of the
/// run carrying on (the default).</param>
public static void QueueOrchestrationFromFile(string name, string batchFilePath, int priority,
string? postExecFunctionName = null, string? postExecParametersJson = null,
string? reference = null, string? parentRunName = null, bool sequential = false)
string? reference = null, string? parentRunName = null, bool sequential = false, bool allowCollision = true,
int maxConcurrency = 0, bool stopOnFailure = false)
{
name = TableKeys.Sanitize(name);
parentRunName = ResolveParentRunName(name, parentRunName);
var gated = RegisterPendingChild(parentRunName, name);
var child = RegisterPendingChild(parentRunName, name);
s_pending.Enqueue(new PendingOrchestration(name, string.Empty, priority,
postExecFunctionName, postExecParametersJson, parentRunName, reference, batchFilePath,
PendingChildRegistered: gated, Sequential: sequential));
Sequential: sequential, ParentRunKey: child?.ParentRunKey, ChildKey: child?.ChildKey,
AllowCollision: allowCollision, MaxConcurrency: maxConcurrency, StopOnFailure: stopOnFailure));
}

/// <summary>
/// Resolve the parent run of a queued orchestration. The explicit argument wins — PowerShell
/// Resolve the parent run of a queued orchestration: a run key (exact — runs of one name can overlap) or a
/// run name (the newest outing). The explicit argument wins — PowerShell
/// callers MUST pass it (read from the stamped $global:CraftOperationContext), because the
/// ambient fallback cannot work for them: the pipeline runs on the runspace's reused thread,
/// whose frozen ExecutionContext never sees the per-invocation AsyncLocal (see
Expand All @@ -75,7 +92,7 @@ public static void QueueOrchestrationFromFile(string name, string batchFilePath,
private static string? ResolveParentRunName(string name, string? parentRunName)
{
if (string.IsNullOrEmpty(parentRunName))
parentRunName = OperationContext.Current?.RunName;
parentRunName = OperationContext.Current?.RunKey ?? OperationContext.Current?.RunName;
if (string.IsNullOrEmpty(parentRunName))
return null;
// Sanitized like the child name: the parent was created under its sanitized name, and the
Expand All @@ -88,76 +105,87 @@ public static void QueueOrchestrationFromFile(string name, string batchFilePath,
}

/// <summary>
/// Register the child link at ENQUEUE time — while the parent task's script is still executing,
/// so the parent cannot pass its completion check before the gate exists. Registering after
/// StartFromBatchAsync (the old shape) loses that race for the parent's LAST task: the drain
/// runs in a background Task.Run while the enqueuing task is marked terminal immediately, so
/// the parent would finalize — and dispatch PostExecution — before its child was visible.
/// Returns whether a gate was taken, so the drain releases exactly what was registered.
/// Make the parent wait for this child, at ENQUEUE time — while the parent's task is still executing, so
/// the parent cannot reach its barrier first. The returned keys travel with the queued run: the child
/// fills the placeholder when it finishes, and a child that is never created releases it on the drain.
/// </summary>
private static bool RegisterPendingChild(string? parentRunName, string childName) =>
!string.IsNullOrEmpty(parentRunName) &&
s_service?.TryRegisterPendingChildRun(parentRunName, childName) == true;
private static (string ParentRunKey, string ChildKey)? RegisterPendingChild(string? parentRunName, string childName) =>
string.IsNullOrEmpty(parentRunName) ? null : s_service?.RegisterPendingChild(parentRunName, childName);

/// <summary>
/// Synchronous drain — blocks until all pending orchestrations are started.
/// Safe to call from any context (no SynchronizationContext on background workers).
/// Whether a run of this name is unfinished or already queued here to start, counting outings that carry a
/// queue id suffix (<c>Name-{guid}</c>) as the same name. Lets a caller that does not want overlapping runs
/// (<c>allowCollision: false</c>) skip, and say so, before building the batch. The start itself checks
/// again, so a run that appears in between is still skipped.
/// PS usage: <c>[Craft.Services.OrchestratorBridge]::IsRunActive($name)</c>.
/// </summary>
public static bool IsRunActive(string name)
{
var family = WorkStore.CollisionFamily(TableKeys.Sanitize(name));
if (s_pending.Any(p => WorkStore.CollisionFamily(TableKeys.Sanitize(p.Name)) == family)) return true;
return s_service != null && Task.Run(() => s_service.IsRunActiveAsync(name)).GetAwaiter().GetResult();
}

private static readonly System.Text.Json.JsonSerializerOptions s_inspectJson = new() { WriteIndented = true };

/// <summary>
/// Why a run is or is not moving, as JSON: counts and mode, whether the scheduler can see it, its claims
/// and who holds them, child runs it waits for, its aggregation, the instance lock, and a diagnosis.
/// Takes a run key, or a run name (every unfinished run of it, else the latest).
/// PS usage: <c>[Craft.Services.OrchestratorBridge]::InspectRun('MailboxRules_contoso.com')</c>.
/// </summary>
public static string InspectRun(string nameOrKey) => s_service == null
? "{\"error\":\"orchestrator not initialised\"}"
: System.Text.Json.JsonSerializer.Serialize(Task.Run(() => s_service.InspectRunAsync(nameOrKey)).GetAwaiter().GetResult(), s_inspectJson);

/// <summary>
/// Rebuild the Ready and Finished indexes from the active-run list now, as the pump does at startup: relists
/// unfinished runs, retires finished ones, removes runs whose creation never finished. Returns a JSON summary.
/// PS usage: <c>[Craft.Services.OrchestratorBridge]::RepairIndexes()</c>.
/// </summary>
public static string RepairIndexes() => s_service == null
? "{\"error\":\"orchestrator not initialised\"}"
: System.Text.Json.JsonSerializer.Serialize(Task.Run(() => s_service.RepairIndexesAsync()).GetAwaiter().GetResult(), s_inspectJson);

/// <summary>Synchronous drain — blocks until all pending orchestrations are started.</summary>
public static void DrainPending()
{
while (s_pending.TryDequeue(out var p))
{
try
{
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.Sequential)
.GetAwaiter().GetResult();
}
catch (Exception ex)
{
s_service?._logger.LogError(ex, "[Orchestrator] DrainPending failed for {Name}", p.Name);
}
finally
{
// The enqueue-time gate lifts on EVERY path once the start attempt is over: a
// started child is in _activeRuns by now (which takes over blocking the parent),
// and one that failed to start must stop blocking — a leaked gate would defer the
// parent's finalize forever, re-checked every 60s for the process lifetime.
if (p.PendingChildRegistered)
s_service?.ReleasePendingChildRun(p.Name);
}
}
Task.Run(() => StartAsync(p)).GetAwaiter().GetResult();
DrainPendingPlanners();
}

/// <summary>
/// Async drain — preferred from async call sites (PostExec lambdas, ExecuteScript).
/// </summary>
/// <summary>Async drain — preferred from async call sites (PostExec, ExecuteScript).</summary>
public static async Task DrainPendingAsync()
{
while (s_pending.TryDequeue(out var p))
await StartAsync(p);
await DrainPendingPlannersAsync();
}

private static async Task StartAsync(PendingOrchestration p)
{
var created = false;
try
{
try
{
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.Sequential);
}
catch (Exception ex)
{
s_service?._logger.LogError(ex, "[Orchestrator] DrainPending failed for {Name}", p.Name);
}
finally
if (s_service == null) { DiscardUndispatchable(p); return; }
created = await s_service.StartFromBatchAsync(p.Name, p.BatchJson, p.Priority,
p.PostExecFunctionName, p.PostExecParametersJson, CancellationToken.None,
p.ParentRunName, p.Reference, p.BatchFilePath, p.Sequential, p.ParentRunKey, p.ChildKey, p.AllowCollision,
p.MaxConcurrency, p.StopOnFailure);
}
catch (Exception ex)
{
s_service?._logger.LogError(ex, "[Orchestrator] DrainPending failed for {Name}", p.Name);
}
finally
{
if (!created && p.ParentRunKey != null && p.ChildKey != null && s_service != null)
{
// See DrainPending: the gate lifts whatever the outcome of the start attempt.
if (p.PendingChildRegistered)
s_service?.ReleasePendingChildRun(p.Name);
try { await s_service.AbandonPendingChildAsync(p.ParentRunKey, p.ChildKey); }
catch (Exception ex) { s_service._logger.LogWarning(ex, "[Orchestrator] Could not release {Name} from its parent", p.Name); }
}
}
await DrainPendingPlannersAsync();
}

/// <summary>
Expand All @@ -179,15 +207,17 @@ private static void DiscardUndispatchable(PendingOrchestration p)

/// <summary>
/// A queued run. Exactly one of <paramref name="BatchJson"/> and <paramref name="BatchFilePath"/>
/// carries the batch; the file path wins when both are set.
/// <paramref name="PendingChildRegistered"/> records whether enqueue took a pending-child gate on
/// the orchestrator, so the drain releases exactly the gates that were taken — releasing on a
/// refused registration could lift a gate held by ANOTHER queued entry of the same child name.
/// carries the batch; the file path wins when both are set. <paramref name="ParentRunKey"/> and
/// <paramref name="ChildKey"/> are set when the parent was made to wait for this run.
/// </summary>
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 Sequential = false);
string? Reference = null, string? BatchFilePath = null, bool Sequential = false,
string? ParentRunKey = null, string? ChildKey = null, bool AllowCollision = true, int MaxConcurrency = 0,
bool StopOnFailure = false)
{
public bool PendingChildRegistered => ChildKey != null;
}

private static readonly ConcurrentQueue<PendingPlannerRun> s_pendingPlanners = new();

Expand Down
Loading
Loading