Skip to content

fix: harden orchestration, limit thread usage, bound timers and thread usage, consolidate logging, new sequential orchestration mode - #41

Merged
Zacgoose merged 12 commits into
mainfrom
dev
Sep 8, 2026
Merged

Zacgoose merged 12 commits into
mainfrom
dev

Conversation

@Zacgoose

@Zacgoose Zacgoose commented Sep 8, 2026

Copy link
Copy Markdown
Contributor

This pull request introduces a new "sequential execution" mode for orchestrator runs, allowing tasks to be dispatched one at a time in order, rather than all at once in parallel. It also includes improvements to logging, memory usage, and background service robustness. The most important changes are grouped below.

Sequential Execution Mode for Orchestrator Runs:

  • Added a Sequential flag to orchestrator runs (OrchestratorRun, PendingOrchestration, and related APIs), enabling tasks to be dispatched one at a time in batch order instead of fanning out in parallel. This is controlled via a new Sequential property in the orchestrator input object and is now supported end-to-end in both PowerShell and C# layers. [1] [2] [3] [4] [5] [6] [7] [8] [9] [10] [11]

  • Updated the PowerShell runner service to allow pinning a worker for the duration of a sequential run, ensuring all tasks execute on the same worker without returning it to the pool between steps. [1] [2] [3]

Configuration and Performance Enhancements:

  • Added new orchestrator settings to control the status timer interval, enable/disable redrive backoff, and shed pending parameters to reduce memory usage when handling large backlogs.

Logging and Observability Improvements:

  • Changed per-task dispatch and completion logging from Write-Information to Write-Debug to reduce log noise, especially during crash recovery or high scale. [1] [2]

Background Service Robustness:

  • Improved the StatsHistoryService to gracefully handle cancellation during shutdown, preventing host faults due to background exceptions. [1] [2]

Internal API and Data Model Updates:

  • Extended task and status record models to track sequence order and additional metadata for sequential runs. [1] [2]

These changes collectively enable more flexible orchestrator execution patterns and improve the system's scalability, observability, and reliability.

Two ways a task could sit Pending forever with nothing to dispatch it.

1. Resolver dropped tasks on a task-script-path cache miss. ResolveTaskWorkAsync
   read _taskScriptPaths first and returned null on a miss, so the JobManager
   marked the job Skipped and the pump deleted its queue row — leaving the task
   Pending in the run with no queue row and nothing to re-dispatch it. But
   _taskScriptPaths is written only by DispatchPendingTasksAsync, and the pump
   (a BackgroundService) begins claiming persisted queue rows at host start,
   before ResumeInterruptedRunsAsync has re-dispatched a resumed run; the
   rehydration branch also re-adds a run to _activeRuns without a cached path.
   Resolve the run first, then rebuild the path from the persisted
   run.TaskScriptName exactly as the resume path does (the ScriptRepository is
   loaded before the pump's first claim) and cache it. Only an empty
   TaskScriptName with no naming-convention match is genuinely unrunnable.

2. Orphan re-drive trusted the index, which can outlive its queue row.
   RedrivePendingTasksAsync decided "orphaned" from GetQueuedTaskIdsAsync, which
   reads the index table only, while the pump claims from the queue table only.
   An index row can survive with no matching claimable queue row (a removal
   deletes the index side first; a run dispatched under an in-memory build
   re-enters with index rows and no queue rows; an owned row with no LeaseUntil
   is excluded by the claim filter). The task is then invisible to the pump but
   reported queued by the index, so neither resume nor the watchdog re-enqueues
   it. Add GetDispatchableTaskIdsAsync, which verifies each candidate against the
   queue table (a point read apiece, so the caller passes the small aged-Pending
   candidate set) and returns only tasks the pump can still claim; re-drive
   anything else. RequeueToTable's upsert re-writes both rows.

Tests: OrchestratorTaskPathRehydrationTests, JobQueueDispatchableTests.
…unts

Thousands of concurrently-live runs walked the managed heap into the GC hard
limit and OOM-crashed the host (exit 139), while the per-run status machinery
flooded the log and re-read storage on every tick. Four changes, each measured
on a new run-count-axis harness, all behind Orchestrator config flags that
default to the new behaviour (ShedPendingParameters, RedriveBackoff,
StatusTimerIntervalSeconds):

- Shed pending-task Parameters. A live run pinned every task's Parameters
  payload in _activeRuns for the run's whole lifetime, so a large pending
  backlog was the dominant retained heap. The Tasks table is already
  authoritative, so drop the in-memory payload once a task is persisted and
  enqueued, and rehydrate it from storage at dispatch. Measured: 3000 runs of
  4 tasks with a 10KB payload under a 256MB heap cap went from heap climbing to
  244MB then exit 139, to a flat ~33MB with no OOM. Shedding happens before the
  rows become claimable so it cannot race a concurrent dispatch. Fail-closed: a
  missing row at dispatch (distinct from a present-but-empty payload) fails the
  task explicitly rather than running it with no parameters — restoring, for the
  one case shedding introduced, the resilience a resident payload used to give.

- Back off the per-run re-drive. It verified against the queue table every tick
  even when nothing was orphaned, the dominant per-tick storage + deserialization
  cost. Once a run verifies clean the interval doubles to a 15-minute cap and
  snaps back to the base interval the moment an orphan appears. Measured: 3.5x
  fewer re-drive storage reads over a fixed window.

- Collapse the status-log flood. LogRunStatus walked the task list four times
  and emitted a line every tick unconditionally; now it makes one pass and skips
  the line when the counts are unchanged, keeping a slow heartbeat so a
  long-lived run still shows it is alive. Measured: 7.4x fewer status log lines.

- Replace the per-run System.Threading.Timer with a single status/re-drive sweep
  over _activeRuns. At high live-run counts that was one Timer object per run and
  a continuous drizzle of fire-and-forget callbacks onto the thread pool; one
  sweep is a single scheduling source that allocates nothing per run.

The status/re-drive tick cadence is now configurable (StatusTimerIntervalSeconds,
default 60), also used to slow the loop on a constrained deployment. 682/682
tests pass; a 2000-task fan-out still completes with every task terminal.
run-oom.ps1 tests fan-out WIDTH (one run of N tasks); run-manyruns.ps1 exercises
the axis it never did — the number of concurrently-live RUNS — under a GC heap
cap and a constrained CPU/thread budget, reporting heap, Gen2 GCs, thread-pool
depth, HTTP-SLOW, status-log volume and re-drive storage reads as functions of
that count. It also A/Bs the new orchestrator flags in a single image.

New PerfApi endpoints back it and the failure-mode exploration:
- PerfManyRuns: create M live runs of hold-open (or payload-carrying) tasks.
- PerfThreads: thread-pool counters + the cumulative re-drive storage-read count.
- PerfCheck / PerfCheckCounts: a per-task marker that confirms a payload survived
  the shed -> rehydrate round trip (detects any parameter loss).
- PerfTableOp: count/list/delete orchestrator table rows (a row, a partition, or
  a whole table) while runs are live, to probe deletion resilience.

docker-compose.bg.yml gains matching env knobs (LOG_LEVEL, STATUS_INTERVAL,
REDRIVE_BACKOFF, SHED_PARAMS).
Add an opt-in Sequential run mode: a run's tasks execute ONE AT A TIME in
payload order, the way Azure Durable Functions sequences an orchestration,
rather than the default fan-out that enqueues every task up front and drains
them in parallel. This is the mechanism an ordered workflow needs so a step
that must follow another (e.g. convert-to-shared strictly after the grants
before it) stops racing it.

Mechanism (chaining, no worker pinning): only the current task — the lowest
OrchestratorTaskItem.Sequence still Pending — is ever enqueued; the next is
enqueued when the current reaches a terminal state, so the durable queue never
holds more than one of the run's tasks. It still runs on whatever worker is
free.

- OrchestratorRun.Sequential and OrchestratorTaskItem.Sequence (payload index,
  stamped at batch parse), both persisted — including through the status
  writer's Replace-mode task writes, so a resumed run keeps its order.
- DispatchPendingTasksAsync enqueues only the current task (covering first
  dispatch and resume), and only when it lacks a queue row — never the next
  one, which would put two of the run's tasks in flight.
- On terminal, AdvanceSequentialAsync enqueues the next task; it advances past
  a FAILED step too, so one failure cannot strand the rest Pending forever.
- RedrivePendingTasksAsync is restricted for sequential runs: it never re-drives
  while a task is Running, and otherwise only re-enqueues the current task if
  its own queue row is genuinely gone — the not-yet-reached tasks deliberately
  have no row and must not be treated as orphaned.
- Parameter shedding is skipped for sequential runs (small, and the advance
  path does not rehydrate).

Threaded through OrchestratorBridge and Start-CraftOrchestrator
(InputObject.Sequential); off by default, so fan-out behaviour is unchanged.
Pinned by OrchestratorSequentialTests across persistence, the dispatch gate,
the advance step, and the re-drive restriction.
PerfManyRuns gains a seq flag (sets InputObject.Sequential) and a PerfSeq task
that records the order tasks start in and the maximum concurrency observed into
a shared tally, read via PerfSeqResult. This lets the harness show a sequential
run executing in payload order at concurrency 1 even with idle workers, against
a fan-out control that interleaves — the end-to-end counterpart to the unit
tests that pin the mechanism.
PerfThreadBreakdown exposes WorkerMetricsBridge.GetMemoryBreakdown()'s OS-thread
tally — total thread count and a breakdown by ThreadState, plus processor count
and the PS worker-pool sizes. Lets the harness answer "where are Craft's threads"
at a point in time and confirm the count stays bounded as the concurrently-live
run count climbs (the axis the per-run status timers used to grow it on).
PerfSeedRuns writes run + task + counter rows straight into the orchestrator
tables (Azure Table batch transactions), so a high live-run-count test is not
gated by Start-CraftOrchestrator's per-run enqueue path. Restart the container
after seeding and ResumeInterruptedRunsAsync materialises them into live runs:
5000 runs seed in ~20s and resume in seconds, versus minutes through the batch
path. Used to compare thread usage of the per-run-timer scheduler against the
single sweep at run counts the enqueue path could not reach.
…o Debug

Crash recovery emitted several Info lines PER RUN — "Found interrupted run",
"Resuming interrupted run: N pending", "Released N stale claims", "Dispatched N
tasks" — and Invoke-CraftTask wrote two Info lines PER TASK ("Dispatching task
to X" / "Completed task X"). On a crash-looping instance that replayed thousands
of lines on every restart, since a restart resumes every run and re-dispatches
all its pending tasks.

- ResumeInterruptedRunsAsync now tallies run outcomes and emits ONE aggregate
  line ("Crash recovery: resumed N run(s) (P pending re-dispatched), ...")
  instead of the per-run lines, which drop to Debug. Genuine problems (a run
  that cannot resume, a post-execution abandoned) still log at Warning/Error as
  they happen and are counted in the summary.
- DispatchPendingTasksAsync takes a quiet flag: its per-run "Dispatched N tasks"
  line stays Info for a normal orchestration start (one line) but drops to Debug
  when called from recovery, where the aggregate covers it.
- Invoke-CraftTask's per-task dispatch/complete lines move to Debug — two Info
  lines per task was the bulk of the volume, worst during recovery. Push-*
  functions still log their own output at Info.

Measured: resuming 2000 runs went from ~6000 Info lines to 1. Separate from the
per-tick status flood already collapsed by LogRunStatus's skip-when-unchanged.
…or sample

The Sequential option was documented in the .PARAMETER block but the only usage
example showed a fan-out. Add a second .EXAMPLE demonstrating Sequential = $true
with an ordered, Durable-style batch (an offboarding sequence where a later step
must not race the ones before it), and label the existing example as the
default fan-out.
A sequential run now runs as ONE driver on ONE pinned worker instead of chaining
task-by-task through the queue. The single entry row that dispatch enqueues drives the
whole run inline: check a background worker out once, run every step on it in payload
(Sequence) order, reclaim it once at the end. The run keeps the same worker from start to
finish and never returns to the pool between steps to be re-scheduled elsewhere.

- PowerShellRunnerService.ExecuteScript / ExecuteScriptWithOutput gain an optional
  pinnedWorker: when supplied the script runs on that worker and it is NOT reclaimed there,
  so the driver owns the worker's lifecycle across all steps. CheckoutBackgroundWorker /
  ReclaimBackgroundWorker expose the checkout/reclaim pair. InvokeAsync still resets the
  runspace after every step, so steps stay isolated on the shared worker.
- OrchestratorService.BuildSequentialRunWork replaces the AdvanceSequentialAsync chaining.
  ResolveTaskWorkAsync routes a sequential run's entry claim to it. Failure policy is
  best-effort: a failing step is recorded and the driver moves on, so one bad step cannot
  strand the rest; a cancelled run marks the remaining steps Cancelled; a duplicate entry
  row is a no-op while a driver is active (TryAdd guard). _activeSequentialDrivers registers
  the driver for its whole life so the re-drive leaves the run's deliberately row-less steps
  alone, yet still restarts the driver if the entry row is lost.
- The three PowerShell interactions are behind internal virtual seams so the driver loop is
  unit-tested without a worker pool. OrchestratorSequentialTests is reworked to the pinned
  model: ordering, best-effort, cancellation, single checkout/reclaim, duplicate no-op,
  routing, and the re-drive restriction. Full suite green (701).
Push-PerfSeqWorker records (run, idx, worker) per step under a unique key, and
Invoke-PerfSeqWorkerEnqueue starts N sequential (or fan-out) runs holding per step so
several overlap. Proves the pinned driver live: every step of a sequential run lands on the
same worker in payload order, and concurrently-running sequential runs each pin their own
worker — while a fan-out (seq=false) run spreads its steps across workers out of order.
Wrap the startup delay and timer loop in a cancellation-aware try/catch so expected shutdowns do not fault the host. This preserves per-sample warning logging while avoiding BackgroundServiceExceptionBehavior.StopHost during normal service stop.
@Zacgoose Zacgoose changed the title fix: harden orchestration, limit thread usage, bound timers and thread usage, consolidate logging fix: harden orchestration, limit thread usage, bound timers and thread usage, consolidate logging, new sequential orchestration mode Sep 8, 2026
@Zacgoose
Zacgoose merged commit f65a978 into main Sep 8, 2026
7 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants