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 at 1.

  • 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).

Pipeline
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.

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.