Items
NAME
LLM::Data::Pipeline::Step::Items - A pipeline step whose work is a set of independently-processed items
SYNOPSIS
use LLM::Data::Pipeline::Step::Items;
class TagChunks does LLM::Data::Pipeline::Step::Items {
method name(--> Str:D) { 'tag-chunks' }
method description(--> Str:D) { 'Tag every story chunk' }
method requires(--> List) { ('chunks',) }
method provides(--> List) { ('tag-chunks/items', 'tags') }
# Deterministic given the Context; JSON-safe items.
method items($ctx --> Iterable) { $ctx.get('chunks').list }
# Content-derived keys survive reordering of the chunk list.
method item-key(Int:D $i, $chunk --> Str:D) { $chunk<id> }
# Worker thread: no Context mutation, at-least-once.
method process-item($ctx, $chunk, Str:D $key --> Any) {
my $task = build-task($chunk);
$task.on-call-complete($runner.telemetry-sink(:step<tag-chunks>, :$key));
$task.execute; # returns a JSON-safe result
}
# Coordinator thread: assemble provides from the collected results.
method finalize($ctx --> Nil) {
my %items = $ctx.get('tag-chunks/items');
$ctx.set('tags', %items.values.flat.unique.sort.list);
}
}
DESCRIPTION
A Step::Items is a Step whose body is a set of
items processed independently and in parallel by the Runner's item engine.
Instead of execute, you implement items, process-item, and (usually)
finalize. The Runner branches on does Step::Items and never calls
execute (it dies if you do).
The contract
Determinism.
items($ctx)must return the same items in the same order for the same Context. It is called once per activation and materialized; a resumed in-progress step recomputes it and compares a fingerprint (count + SHA-256 of the keys) against the checkpoint, throwingX::LLM::Data::Pipeline::CheckpointDrifton a mismatch.JSON-safe items. Each item is digested (SHA-256 of its JSON) for the dead-letter journal and must round-trip through JSON.
Unique keys.
item-key($i, $item)must be unique across the set. Duplicates are detected at activation and abort the run. The default key is the 0-based index (stable only if the list order is stable); override with a content-derived key for reorder-stable resume.No Context mutation in
process-item. It runs on a worker thread while the Context is frozen; callingsetdies. Reads are safe β butgetreturns the live reference, so mutating a fetched value in place (pushing to a fetched Array, assigning into a fetched Hash) is NOT caught by the frozen flag and is a data race under parallelism β undefined behavior. Treat everything you read as read-only; assemble all outputs infinalize.At-least-once.
process-itemmay run more than once for a key (retry, or resume after a crash between processing and checkpoint). Recording is exactly-once per checkpoint lineage, so a pureprocess-itemis effectively exactly-once; anything with external side effects must be idempotent. Saving partials (below) narrows the unit that gets re-run to the sub-unit boundary; it does not remove the requirement.Idempotent
finalize. A crash betweenfinalizeand the step-boundary checkpoint leaves the step not-yet-complete, so a resume re-runsfinalizeover the recovered results. It must deriveprovidespurely from the reserved keys / results, never accumulate onto prior state.
Resumable items: saving sub-unit progress (0.6.0)
An item that is one unit of work needs nothing here. An item that is a run of sequential sub-units β twenty model calls to rewrite twenty turns of one chunk β should save its progress as it goes, or a failure on sub-unit 19 throws eighteen finished ones away and a resume re-runs the lot.
The contract is two halves and it is entirely opt-in:
process-itemreceives:%partialβ whatever this item last saved, or%()for "nothing saved, start fresh". Treat a partial whose shape you do not recognise (an older release's, a hand-edited checkpoint) exactly like%(): fall back to a fresh start rather than dying, since a malformed partial must never be able to make an item unprocessable.Runner.partial-sink(:$step!, :$key!)mints the save closure for that item. Call it with a JSON-safe Hash after each completed sub-unit β not before one, not batched β so what is saved is always a prefix of the work that is really done.
method process-item($ctx, $chunk, Str:D $key, :%partial --> Any) {
my &save = $runner.partial-sink(:step(self.name), :$key);
my @turns = (%partial<turns> // []).list.Array; # already-written turns
my Int $at = (%partial<next> // 0).Int; # where to resume
for @($chunk<beats>)[$at ..^ *] -> $beat {
@turns.push(write-turn($beat)); # the expensive part
$at++;
save({ turns => @turns, next => $at });
}
@turns;
}
Two properties make this sound, and both are the step's responsibility:
The resumed work must be deterministic given the partial. Anything the loop derives from earlier sub-units (a transcript, a handover line, a running index) must be re-derivable from the partial plus the item β not carried only in a local variable of the attempt that died. If it cannot be re-derived, save it.
A save is a promise the sub-unit is finished. The engine hands the partial back verbatim on the next attempt and never re-runs what it covers.
The Runner drops an item's partial the moment the item is recorded done or
dead-lettered, and :retry-dead requeues an item without one. Cancellation,
by contrast, keeps partials β that is the point of them. See
LLM::Data::Pipeline::Runner's "Resumable items" for the persistence cadence,
the JSON normalization rule, and the checkpoint-size trade-off.
Reserved Context keys
Before finalize the engine writes two reserved keys into the (thawed)
Context, so finalize and downstream steps can read them β declare whichever
you consume in provides:
"{name}/items"β Hash of item key β result (successful items only)."{name}/dead"β sorted List of dead-lettered item keys.
SEE ALSO
LLM::Data::Pipeline::Runner for the parallel engine, checkpoint v2, and DLQ; LLM::Data::Pipeline::RetryPolicy for per-item retries.
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.