Event
NAME
LLM::Data::Pipeline::Event - Typed, immutable events emitted by the pipeline Runner
SYNOPSIS
use LLM::Data::Pipeline::Event;
use LLM::Data::Pipeline::Runner;
my LLM::Data::Pipeline::Runner $runner .= new(
on-event => -> LLM::Data::Pipeline::Event $e {
given $e {
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";
}
}
}
);
$runner.run($plan, $ctx);
# Every event is also a JSON-safe Hash, ready for a log line:
$runner.events.tap(-> $e { say to-json($e.to-hash) });
DESCRIPTION
Every observable moment in a run is an immutable object composing the
LLM::Data::Pipeline::Event role. The role supplies three envelope fields ā
UInt $.seqā strictly increasing per-run sequence, starting at1.Instant $.atā wall-clock time the event was emitted.Str $.run-idā the run that produced it (stable within a run).
ā plus two methods: kind (a stable kebab-case string discriminator) and
to-hash (a JSON-safe Hash of every field). In to-hash, at is rendered
as a POSIX epoch-seconds number (fractional), so the whole hash round-trips
cleanly through JSON::Fast:
my %h = $event.to-hash;
my $json = to-json(%h); # never throws on a Stage-2 event
my %back = from-json($json); # %back<kind>, %back<seq>, %back<run-id>, %back<at>, ā¦
EVENT TAXONOMY
Every kind in the taxonomy is emitted as of this release (0.5.1). Earlier
releases emitted a growing subset; the Status column records when each kind
first shipped.
Nearly all events carry a per-run seq starting at 1. The sole exception is
run-retry, a supervision event emitted by run-until-done between
inner run/resume attempts: it carries seq 0 to mark that it sits outside
any single run's totally-ordered sequence, and is delivered via &.on-event
(each inner run's events() Supply has already ended in done).
| Kind | Class | Payload fields | Status |
|---|---|---|
| run-started | RunStarted | plan-size, resumed, checkpoint-path | emitted (0.2.0) |
| step-skipped | StepSkipped | step, ordinal, total, description | emitted (0.2.0) |
| step-started | StepStarted | step, ordinal, total, description | emitted (0.2.0) |
| step-completed | StepCompleted | step, ordinal, total, description, duration, items-done?, items-dead? | emitted (0.2.0) |
| checkpoint-written | CheckpointWritten | path, trigger | emitted (0.2.0) |
| run-completed | RunCompleted | duration, steps-run, steps-skipped, items-done, items-dead | emitted (0.2.0) |
| run-failed | RunFailed | step, error, exception | emitted (0.2.0) |
| step-failed | StepFailed | step, attempt, error, exception, will-retry, retry-delay | emitted (0.3.0) |
| step-retry | StepRetry | step, attempt, delay, error | emitted (0.3.0) |
| item-started | ItemStarted | step, key, attempt | emitted (0.4.0) |
| item-completed | ItemCompleted | step, key, attempt, duration | emitted (0.4.0) |
| item-failed | ItemFailed | step, key, attempt, error, exception | emitted (0.4.0) |
| item-dead-lettered | ItemDeadLettered | step, key, attempts, error | emitted (0.4.0) |
| item-requeued | ItemRequeued | step, key | emitted (0.4.0) |
| progress | Progress | step, done, dead, total, in-flight, pending | emitted (0.4.0; also at activation from 0.5.1) |
| run-cancelled | RunCancelled | step, items-done, items-dead | emitted (0.4.0) |
| telemetry | Telemetry | step, key, stage, data | emitted (0.4.0) |
| run-retry | RunRetry | attempt, max-attempts, delay, error, exception | emitted (0.5.0) |
The step-completed items-done/items-dead fields are present only for
Step::Items steps (omitted for plain steps). checkpoint-written trigger
is one of 'step', 'item', 'dead-letter', 'cancel', or ā since
0.6.0 ā 'partial' (an item saved sub-unit progress from inside its own
work; see LLM::Data::Pipeline::Runner's "Resumable items"). 'partial' and
'item' writes are coalesced by checkpoint-every / checkpoint-interval;
the other three never are.
Sub-unit progress itself is deliberately not an event: a partial can be the
whole of an item's generated output, which has no business on the event stream
or in the dead-letter journal. What a consumer wants for "this item is still
moving" is the step's own telemetry stream.
progress is emitted once at item-step activation as well (since 0.5.1):
before the step's first item-started, carrying in-flight 0, the done
and dead counts rehydrated from the checkpoint, and pending for everything
not yet recorded. So a fresh stage reports 0/N and a resumed one done/N the
moment it activates, instead of nothing until its first item finishes:
when LLM::Data::Pipeline::Event::Progress {
# Fires immediately on activation, then after every item terminal.
say sprintf('%s %d/%d (%d dead, %d in flight)',
$e.step, $e.done, $e.total, $e.dead, $e.in-flight);
}
A step skipped as already-complete on resume never activates, so it emits no
progress at all ā step-skipped is the only thing a consumer sees for it.
All duration and delay fields are seconds as a number; checkpoint-path
is Str and may be the (undefined) type object when no checkpoint path was
given, which serializes to JSON null.
SEE ALSO
LLM::Data::Pipeline::Runner for how and when these events are emitted and the per-run ordering guarantees.
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.