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 at 1, no gaps).

  • step-started precedes its item events, which precede that step's step-completed/step-failed (item events arrive from 0.4.0).

  • An item step's first item event is a progress activation snapshot (from 0.5.1): it is emitted before the first item-started, so a consumer can render 0/N (fresh) or done/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

Item-engine
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-item receives 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.0 behaviour.

  • 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-retryable that 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 attempts 1 and one attempt-history entry β€” 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-done are 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-id is a stable lineage identifier: run mints it once and every resume adopts the checkpoint's run-id, so a dead record and a later requeued record 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-item is at-least-once; recording is exactly-once per checkpoint lineage. A pure process-item is 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

Item-step
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 resume if the checkpoint already exists (so it composes with an external re-invoker that restarts the process after a crash), otherwise run with the seeded Context. Every later attempt is resume.

  • 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), and X::LLM::Data::Pipeline::CheckpointDrift (the same checkpoint will drift again). Also X::LLM::Data::Pipeline::ItemsDead unless :retry-dead β€” without requeuing, a re-run reproduces the identical dead set, so retrying is futile.

  • Retried: everything else. Each retry emits a run-retry event (attempt, max-attempts, delay, error, exception), backs off via &.schedule-after, then resumes. When the run-retry budget is spent the last error is rethrown unchanged.

  • :retry-dead composes: when True, every attempt requeues dead items first, and ItemsDead becomes 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); only fail-on-dead steps 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 + "" + .key)
       | map(last) | map(select(.kind == "dead"))
       | .[] | {step, key, attempts, error: .error.message}' run.dlq.jsonl

or in Raku with JSONL::Reader (dedupe by last record per (step, key)):


use JSONL::Reader;
my %last;
%last{"{.<step>}\0{.<key>}"} = $_ for JSONL::Reader.new(:path($dlq)).list.map(*.value);
for %last.values.grep(*.<kind> eq 'dead') -> %d {
    say "%d<step>/%d<key>: %d<error><message> (%d<attempts> attempts)";
    with %d<inference> { say "  model summary: {.<summary>}" }   # the raw inference cause
}

Should I use :retry-dead?

  • Transient item failures (backend hiccup, timeout) β†’ :retry-dead so a later attempt reprocesses them once the upstream recovers.

  • Poison items (malformed input, a bug the model can't get past) β†’ leave :retry-dead off; they stay journaled for offline triage and never burn the run-retry budget. Fix the cause, then a one-off resume(:retry-dead) drains them.

SEE ALSO

LLM::Data::Pipeline::Event for the full event taxonomy and JSON shapes; LLM::Data::Pipeline::Step::Items for the item-step contract; LLM::Data::Pipeline::RetryPolicy and LLM::Data::Pipeline::Exceptions.

AUTHOR

Matt Doughty <[email protected]>

COPYRIGHT AND LICENSE

Copyright 2026 Matt Doughty

This library is free software; you can redistribute it and/or modify it under the Artistic License 2.0.

LLM::Data::Pipeline v0.7.1

Generic pipeline framework for composing and running sequential data processing steps with checkpointing

Authors

  • Matt Doughty

License

Artistic-2.0

Dependencies

JSON::Fast:ver<0.19>:auth<cpan:TIMOTIMO>JSONL:ver<0.1.3+>:auth<zef:apogee>Digest::SHA256::Native:ver<1.0.0+>:auth<zef:bduggan>

Test Dependencies

Provides

  • LLM::Data::Pipeline
  • LLM::Data::Pipeline::Context
  • LLM::Data::Pipeline::Event
  • LLM::Data::Pipeline::Exceptions
  • LLM::Data::Pipeline::Plan
  • LLM::Data::Pipeline::RetryPolicy
  • LLM::Data::Pipeline::Runner
  • LLM::Data::Pipeline::Step
  • LLM::Data::Pipeline::Step::Items

The Camelia image is copyright 2009 by Larry Wall. "Raku" is a trademark of the Yet Another Society. All rights reserved.

Built with Podlite β€” the markup and publishing tools behind this site.