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, throwing X::LLM::Data::Pipeline::CheckpointDrift on 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; calling set dies. Reads are safe β€” but get returns 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 in finalize.

  • At-least-once. process-item may 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 pure process-item is 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 between finalize and the step-boundary checkpoint leaves the step not-yet-complete, so a resume re-runs finalize over the recovered results. It must derive provides purely 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-item receives :%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.

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.