Runner
NAME
LLM::Data::Pipeline::Runner - Executes a Plan against a Context, emitting a typed event stream
SYNOPSIS
use LLM::Data::Pipeline::Runner;
use LLM::Data::Pipeline::Event;
my LLM::Data::Pipeline::Runner $runner .= new(
on-event => -> LLM::Data::Pipeline::Event $e {
given $e {
when LLM::Data::Pipeline::Event::RunStarted {
say "run {$e.run-id} started ({$e.plan-size} steps, resumed={$e.resumed})";
}
when LLM::Data::Pipeline::Event::StepStarted {
say "β {$e.step} ({$e.ordinal}/{$e.total})";
}
when LLM::Data::Pipeline::Event::StepCompleted {
say "β {$e.step} in {$e.duration.round(0.01)}s";
}
when LLM::Data::Pipeline::Event::RunCompleted {
say "done: {$e.steps-run} run, {$e.steps-skipped} skipped";
}
}
}
);
my $ctx = $runner.run($plan, $ctx, :checkpoint-path('checkpoint.json'.IO));
# Resume: completed steps arrive as step-skipped, the rest execute.
my $ctx2 = $runner.resume($plan, 'checkpoint.json'.IO);
DESCRIPTION
The Runner drives a Plan over a Context, validating dependencies first, then executing each step in order with a JSON checkpoint written after every successful step (temp-file + rename for atomicity). Everything observable is surfaced as a typed LLM::Data::Pipeline::Event.
Observing a run
&.on-event (primary)
A synchronous callback, invoked once per event on the run thread in strict
seq order:
my @log;
my $runner = LLM::Data::Pipeline::Runner.new(
on-event => -> $e { @log.push($e.to-hash) },
);
$runner.run($plan, $ctx);
Because delivery is synchronous, ordering and backpressure are free: the next
step never starts until your handler returns. Handler exceptions are shielded
(reported via note) so a broken observer can neither abort the run nor cause
later events to be dropped.
method events(-- Supply)> (secondary)
A Supply mirror for reactive consumers:
my $runner = LLM::Data::Pipeline::Runner.new;
$runner.events.tap(
-> $e { say $e.kind },
done => { say 'stream closed' },
);
$runner.run($plan, $ctx); # tap BEFORE running
The Supply is backed by a Supplier; taps run synchronously on the run
thread (the same caveat as &.on-event). Tap before calling
run/resume, or you will miss events already emitted. The Supply receives
done when the run finishes β on normal completion and on failure (the
run-failed event is emitted first, then the exception is rethrown, then
done fires). A consumer wanting isolation from the run thread should bridge
with .Channel. Each run closes its Supplier; call events again for a fresh
Supply before the next run.
&.on-step (deprecated)
The original string-callback API is retained as a thin shim over the event
stream and fires identically to previous releases β 'skip', 'start',
'complete' β driven off step-skipped/step-started/step-completed:
my $runner = LLM::Data::Pipeline::Runner.new(
on-step => -> Str:D $name, Str:D $event { say "$name: $event" },
);
Deprecated β it exposes only step transitions and no run/checkpoint/item
detail. It will be removed at 1.0; migrate to &.on-event.
Ordering guarantees
Events form a per-run total order by
seq(contiguous, starting at1, no gaps).step-startedprecedes its item events, which precede that step'sstep-completed/step-failed(item events arrive from0.4.0).An item step's first item event is a
progressactivation snapshot (from0.5.1): it is emitted before the firstitem-started, so a consumer can render0/N(fresh) ordone/N(resumed) the moment the stage activates. A step skipped as already-complete on resume never activates and emits none.Cross-item ordering is defined only by
seqβ do not infer any other relationship between events of different items.
run-id is stable for the lifetime of a single run/resume call and
differs between calls (minted per invocation). A validation failure throws
before run-started is emitted, so a stream that opens with run-started
means validation has already passed.
Step retries
Each plain step is executed under the &.step-retry
RetryPolicy. The default is a single
attempt:
# Default: try each step once; on failure the run aborts and the step's
# ORIGINAL exception propagates unchanged (pre-0.3.0 contract).
my $runner = LLM::Data::Pipeline::Runner.new;
# Opt into retries for every plain step:
my $resilient = LLM::Data::Pipeline::Runner.new(
step-retry => LLM::Data::Pipeline::RetryPolicy.new(
:max-attempts(3), :base-delay(5), :max-delay(120),
),
);
With max-attempts at 1, a failing step behaves exactly as before: no
step-failed event is emitted, run-failed fires, and the original
exception is rethrown (so domain types such as
X::LLM::Data::Inference::Exhausted still flow out, and existing CATCH
handlers keep matching).
With max-attempts above 1, each failed attempt emits step-failed
(will-retry True, retry-delay = the backoff); the Runner then waits that
long (via &.schedule-after, blocking the run thread) and emits step-retry
before the next attempt. When the budget is spent it emits a final
step-failed (will-retry False), then run-failed, then throws
X::LLM::Data::Pipeline::StepExhausted β a Pipeline-typed wrapper carrying the
step name, attempt count, and the last-error (the final attempt's real
exception). Exhausted plain steps always abort the run; skipping them would be
unsound against provides/requires.
Retry implies idempotency. A failed attempt's Context mutations are not
rolled back β the same Context, with whatever partial writes the failed attempt
made, is handed to the next attempt. A step that opts into retries must be
idempotent or tolerant of its own partial writes; a step that cannot promise
this must keep max-attempts at 1.
Injectable time
&.now (default { now }) supplies the clock used for run and step
durations, and &.schedule-after (default
< -> Real $s, &cb { Promise.in($s).then(&cb) } >) supplies the retry timer.
Both are injectable so tests can run a virtual clock β capturing the requested
backoff delays and returning an already-kept Promise β with zero real sleeps.
Event at timestamps always use the real wall clock.
PARALLEL ITEM STEPS (Step::Items)
A Step::Items is processed by the Runner's
item engine rather than a single execute.
The single-coordinator model
The run thread is the coordinator and owns ALL mutable state (results,
attempts, dead set, in-flight set, checkpoint, DLQ). It spins up degree worker
start blocks that pull jobs from a work Channel, run process-item under
a try, and message results back through a single inbox Channel. Workers
never touch state, files, or events. Retries and telemetry re-enter through
the inbox, so every state transition happens in one place β the inbox loop.
There are no locks: event seq ordering and single-writer checkpoint/DLQ fall
out for free. :degree(1) is the sequential mode on the same code path.
Dispatch is throttled to at most degree items in flight; the rest wait in a
ready queue. This bounds concurrency and lets cancellation actually stop pending
work (only the β€ degree in-flight items must drain).
The Context is frozen for the duration of item processing: process-item
runs on a worker thread and any set throws. Reads are safe (nothing mutates
the data). The engine thaws before finalize, then writes the reserved keys
"{name}/items" (keyβresult Hash) and "{name}/dead" (sorted List).
Knobs
| Knob | Default | Meaning |
|---|---|---|
| degree | 4 | Worker parallelism (per-step override via degree) |
| item-retry | 3 tries | Per-item RetryPolicy (per-step override via item-retry) |
| is-cancelled | (none) | Cooperative cancel hook, polled initially + after each message |
| checkpoint-every | 1 | Coalesce item checkpoints to 1 per N item terminals |
| checkpoint-interval | 0.0 | Minimum seconds between coalesced item checkpoints |
Dead-letter, step-boundary, and cancel checkpoints are never coalesced away.
Item-terminal and 'partial' (sub-unit progress) checkpoints are: both go
through the same gate, so the two knobs bound the total checkpoint write rate
of a step however chatty its items are.
Checkpoint v2
Written atomically (temp file + rename) as:
{ "version": 2, "run-id": "β¦", "completed-steps": [...], "context": {...},
"step-state": { "<step>": {
"fingerprint": {"count": N, "keys-sha256": "β¦", "items-sha256": "β¦"},
"done": {"<key>": true}, "results": {"<key>": ...},
"attempts": {"<key>": 2}, "dead": {"<key>": true},
"partials": {"<key>": {...}}, "started-at": "β¦" } },
"updated-at": "β¦" }
In-progress item results live in step-state.results while the Context is
frozen (this is what makes mid-step checkpoints useful for resume); they are
promoted to the reserved Context key at finalize and pruned. done /
attempts / dead are kept after completion (:retry-dead needs dead).
partials holds the sub-unit progress of items that are still in progress
(see #Resumable items) and is empty for a completed step. Legacy v1
checkpoints (no version key) are still accepted on resume, as are v2
checkpoints written before 0.6.0 (an absent partials reads as "no item
has saved any" β every item simply starts from the beginning). The format is
still v2: partials are purely additive, so a 0.6.0 checkpoint also loads
into an older Runner, which ignores the key.
At activation the engine materializes items() once and records a
fingerprint: count, keys-sha256 = SHA-256 of to-json of the key list,
and items-sha256 = SHA-256 of to-json of the materialized items. Resuming
into an in-progress step recomputes all three and throws
X::LLM::Data::Pipeline::CheckpointDrift on any mismatch β so items-sha256
catches value-level drift that stable keys (e.g. index keys) cannot. Duplicate
keys from item-key abort at activation.
Resumable items (sub-unit progress)
Some items are not one unit of work. A step that rewrites an 89-turn chapter
makes 89 sequential model calls for one item, and before 0.6.0 a failure
on turn 60 threw all 60 finished turns away: the item retried, and
process-item started again at turn 1. The same was true of a resume β a
mid-item kill lost the whole item.
An item step can now save its sub-unit progress as it goes and be handed it back on the next attempt. Two pieces, both opt-in:
method partial-sink(:$step!, :$key!)mints a thread-safe closure for one item. Call it with a JSON-safe Hash after each completed sub-unit; it is a Channel send, so it is safe from the worker thread.process-itemreceives the item's last saved Hash as the named argument:%partial.%()means "nothing saved β start fresh".
method process-item($ctx, $item, Str:D $key, :%partial --> Any) {
my &save = $runner.partial-sink(:step(self.name), :$key);
# Resume where the last attempt stopped. An empty partial (or one whose
# shape you don't recognise) means start from the beginning.
my @done = (%partial<turns> // []).list.Array;
my Int $i = (%partial<next> // 0).Int;
for @($item<turns>)[$i ..^ *] -> $turn {
@done.push(write-turn($turn)); # the expensive model call
$i++;
save({ turns => @done, next => $i }); # after each completed sub-unit
}
@done;
}
Cadence. A partial rides the ordinary checkpoint coalescing gate with
< trigger => 'partial' >, unforced β so checkpoint-every /
checkpoint-interval throttle it exactly like item terminals (the default of
1 persists every save). Coalescing is safe here in a way that is worth stating:
an in-process retry is handed the in-memory partial and never reads the
checkpoint, so throttling can only widen the redo window of a process kill,
never that of a retry.
Clearing. The coordinator drops an item's partial the moment the item
becomes terminal β on success, on dead-lettering, on a :retry-dead requeue,
and when :retry-dead reopens a completed step. A completed step's
step-state therefore carries < "partials": {} >, and a requeued item
always re-runs from the start (its saved sub-units are the ones that walked
into the failure). Saves for a key that is already terminal, or minted by a
different step, are dropped with a note.
Cancellation keeps partials. A cancel drains the in-flight items and their saves are still accepted β that is the whole payoff: a run cancelled 60 turns into a chunk resumes at turn 61.
Normalization (and what a partial must be). Every partial is normalized
through JSON at receipt (from-json(to-json(β¦))), so what a retry gets from
memory is identical to what a resume gets off disk. A partial that cannot
round-trip through JSON is a bug in the step and fails loudly at its first
save rather than at some later checkpoint write. There is deliberately no
per-partial input digest: the step's item-list fingerprint and the upstream
steps' completion already pin the inputs, which makes a partial exactly as
trustworthy as a rehydrated item result β and results have never carried one
either.
Checkpoint size. While a step is mid-flight its partials can roughly double
the checkpoint (the in-progress prose is stored alongside the finished items'
results). Both are pruned at finalize, so the completed-run checkpoint is
unchanged in size. A step whose partials are genuinely large can trade
durability for size with checkpoint-every.
At-least-once narrows, it does not go away. process-item is still
at-least-once, but the unit of re-execution becomes the sub-unit boundary:
sub-units whose partial was persisted are not re-run, and sub-units after the
last save are. Any external side effect inside process-item must still be
idempotent.
Item-retry advice (item-retryable)
The item RetryPolicy answers "how many attempts may this item have?".
It cannot answer "would another attempt do anything?", because only the
failure knows that. Some deaths are deterministic: an LLM completion
cut off by max_tokens, a schema violation in a fixed input, a
malformed record. Re-running them re-issues the identical work and buys
the identical death β three times the latency and spend on the way to
the same dead-letter record.
A failure can therefore advise the engine, by exposing a method on the thrown exception:
class X::MyStep::Unfixable is Exception {
method message(--> Str) { 'the input cannot produce a result' }
#| Advice to the item engine: do not spend further attempts.
method item-retryable(--> Bool:D) { False }
}
The engine probes it once per failed attempt, on the worker thread that caught the exception, and ships the resulting Bool to the coordinator alongside the error message. An advised-False failure dead-letters immediately, on that attempt, whatever the policy's remaining budget:
Absent means retryable. The probe is
.?-based and defaults to True, so every exception that has never heard of this β which is almost all of them β keeps the exact pre-0.7.0behaviour.Duck-typed, both ways. Pipeline calls a method by name and imports nothing; the exception's library imports no Pipeline type. Anything can opt in, including
X::LLM::Data::Inference::Truncated, which declares it False for exactly the reason above.Advice cannot malfunction into a dead letter. A hostile or buggy
item-retryablethat throws is read as retryable, and a non-Bool return is coerced (.so) on the worker rather than reaching the coordinator.The attempt count stays true. A first-attempt non-retryable death records
attempts1 and oneattempt-historyentry β the budget it declined to spend is not counted as spent.Only the item engine reads it. Whole-step retries (#Step retries) and
run-until-doneare unaffected; a step that wants to stop a run rethrows, as it always has.
The exception object itself deliberately never crosses into the coordinator: the inbox carries the stringified type name, the message, and this Bool, so nothing thread-hostile is shared and the DLQ record shape is unchanged.
Dead-letter queue (DLQ)
foo.checkpoint.json β foo.dlq.jsonl, appended one record at a time
(open/append/close β crash-safe). No checkpoint path β no DLQ file (dead items
are still reported in-memory). Record shape:
{ "schema": 1, "kind": "dead"|"requeued", "run-id": "β¦", "step": "β¦", "key": "β¦",
"item-digest": "sha256β¦", "attempts": N,
"attempt-history": [{"attempt":1,"at":β¦,"error":"β¦","duration":β¦}],
"error": {"exception":"β¦","message":"β¦"},
"inference": {"attempts":β¦,"summary":"β¦"}, // optional, from telemetry-sink
"first-attempt-at": β¦, "dead-at": β¦ }
The inference block is the last telemetry Hash with < stage => 'exhausted' >
for that (step, key); wire it via method telemetry-sink(:$step!, :$key), whose
returned closure only does a Channel send (thread-safe). Pipeline never imports
any inference types β the coupling is this one documented Hash shape.
Note the two attempt fields differ in scope: attempts is the cumulative
attempt count carried in the checkpoint across resumes, whereas
attempt-history lists only the attempts made by the current process lineage
since the last resume. Per-attempt history is deliberately not persisted in
step-state (it would bloat the checkpoint); after a resume the history restarts
while attempts keeps counting.
Consistency rules
The checkpoint is authoritative for control flow; the DLQ is a forensic, append-only journal.
Write order on exhaustion is DLQ first, then checkpoint. A crash between the two yields a duplicate dead record on the next run; consumers dedupe by the last record per (run-id, step, key). This dedupe is sound only because
run-idis a stable lineage identifier:runmints it once and everyresumeadopts the checkpoint'srun-id, so adeadrecord and a laterrequeuedrecord for the same item share a key and collapse correctly. (A per-invocation run-id would split them and break dedupe.)A DLQ append failure is fatal to the run β the journal is the one artifact this feature exists to produce and is never silently dropped.
process-itemis at-least-once; recording is exactly-once per checkpoint lineage. A pureprocess-itemis effectively exactly-once; any external side effect must be idempotent. A step that saves partials narrows the re-execution unit to the sub-unit boundary (see #Resumable items) but does not change this rule.
resume(:retry-dead)
Requeues every dead item (clearing it from dead/attempts, appending a
requeued journal record, emitting item-requeued) after verifying its digest
against the journal (CheckpointDrift on mismatch). Reopening a completed
step removes it AND every later step from completed-steps and clears their
item state β downstream consumed its provides, so soundness wins over the
tempting single-step re-run β while rehydrating the reopened step's successes
from its reserved Context key so only the formerly-dead items re-run.
Failure modes
| Situation | | Outcome |
|---|---|
| Item attempt throws, budget left | | item-failed, backoff, item-retry, re-run |
| Throw advises item-retryable False | | item-failed, then DLQ 'dead' at once β remaining budget skipped |
| Item exhausts its retry budget | | DLQ 'dead' record, item-dead-lettered, checkpoint; run continues |
| Step finishes with dead items | | run-completed with items-dead > 0 (unless fail-on-dead) |
| fail-on-dead step has dead items | | step checkpointed complete, then X::β¦::ItemsDead throws |
| Context mutated in process-item | | the set throws β the item fails (then retries/dead-letters) |
| Duplicate item-key at activation | | run aborts (die) |
| Item set changed shape on resume | | X::β¦::CheckpointDrift |
| is-cancelled turns true | | drain in-flight, checkpoint, run-cancelled, X::β¦::Cancelled |
| DLQ append fails | | fatal β the exception propagates out of the run |
Cancellation latency
is-cancelled is observed at the top of the inbox loop, i.e. when the next
inbox message arrives. Once observed, no further items are dispatched (including
items whose backoff timer fires β their wake-retry is ignored and they stay
pending for resume), and the coordinator drains only the β€ degree already
in-flight items. If the only outstanding work is a retry timer (no items in
flight), cancellation is not acted on until that timer fires and delivers its
message β so the worst-case cancel latency is roughly the pending backoff delay,
which is bounded by the item RetryPolicy's max-delay.
RUN UNTIL DONE
method run-until-done(Plan:D, Context:D, IO::Path:D :$checkpoint-path!,
RetryPolicy:D :$run-retry = β¦, Bool :$retry-dead = False -- Context:D)> supervises
a pipeline to completion across whole-run retries β the outer ring above per-item
and per-step retries.
my $ctx = $runner.run-until-done(
$plan, $seed-ctx,
checkpoint-path => 'run.checkpoint.json'.IO,
run-retry => LLM::Data::Pipeline::RetryPolicy.new(
:max-attempts(5), :base-delay(30), :max-delay(600)), # patient, minutes-scale
);
Attempt 1 is
resumeif the checkpoint already exists (so it composes with an external re-invoker that restarts the process after a crash), otherwiserunwith the seeded Context. Every later attempt isresume.Not retried (rethrown immediately): plan-validation failures (checked up front β a malformed plan cannot be fixed by retrying),
X::LLM::Data::Pipeline::Cancelled(the caller's own intent), andX::LLM::Data::Pipeline::CheckpointDrift(the same checkpoint will drift again). AlsoX::LLM::Data::Pipeline::ItemsDeadunless:retry-deadβ without requeuing, a re-run reproduces the identical dead set, so retrying is futile.Retried: everything else. Each retry emits a
run-retryevent (attempt,max-attempts,delay,error,exception), backs off via&.schedule-after, thenresumes. When therun-retrybudget is spent the last error is rethrown unchanged.:retry-dead composes: when True, every attempt requeues dead items first, and
ItemsDeadbecomes retryable. Default False keeps poison items dead across retries so they do not burn the run-retry budget β that is what the DLQ is for.Completing with dead items is a success (
run-completed,items-dead> 0); onlyfail-on-deadsteps escalate that to a failure.The attempt count is process-local β a fresh invocation starts with a fresh budget; nothing about it is persisted.
OPERATIONAL RECIPES
run-until-done vs an external supervisor
run-until-done retries in-process: ideal for transient upstream trouble (a
flaky backend, a rate limit) where the process itself is healthy. It does not
survive the process dying. For that, pair it with an external supervisor
(systemd Restart=on-failure, a Kubernetes Job with backoffLimit, a cron
re-invoker): because attempt 1 resumes an existing checkpoint, simply
re-invoking the same command after a crash picks up exactly where it left off.
The two compose β in-process retries for transient errors, the external
supervisor for process death.
DLQ triage
Dead items are journaled next to the checkpoint as < <name>.dlq.jsonl >. To
see what is currently dead (last record per item wins), with jq:
# Latest record per (step,key); show only those still 'dead'.
jq -s 'group_by(.step + "