ALLM.Pipeline.Dsl.Stage (allm_pipeline v0.1.0)

Copy Markdown View Source

One compiled stage of an ALLM.Pipeline declaration.

__pipeline__(:stages) returns a list of these — structs, not bare atoms, so a caller reads the field it means (& &1.name, & &1.concurrency) rather than re-deriving it.

The struct is built at RUNTIME by the function @before_compile generates, not stored in a module attribute: every hook field holds a real function, and Module.put_attribute/3 rejects those. The DSL accumulates quoted AST during the module body and splices it into that generated function, where a :hook_name atom becomes a local capture (&hook_name/2) — which is why a hook may be defp (Phase 4 D6).

Fields

FieldMeaning
namethe stage's atom name, unique within the pipeline
kind:stage (runs once) or :fan_out (runs once per item)
stepthe ALLM.Pipeline.Step module, or nil for the escape-hatch form
bodythe escape-hatch / per-item body, arity 2 (ctx, subject), or nil
inputarity-2 (ctx, subject) hook building the Step's input struct
overfan_out only: arity-1 (prev_output) hook returning the item list
skip_when{:opt, key, default} / arity-1 (ctx) hook, or nil
carryfield names captured into the context's carry map — stage only, see below
parent:source_stage (default) or :per_item — see the moduledoc of ALLM.Pipeline
concurrencynil (inherit the pipeline's), a pos_integer(), or {:opt, key, default}
on_errorkind: :stage: :fail_run (default) or :continue. It governs a whole-stage failure; a fan_out has none (its items fail individually into its output), so declaring it there is a compile error

carry:stage only

The keys are captured from the stage's own output — not from its subject, which is the previous stage's output — and merged into the carry map for every later stage. That is what makes a carried value survive any number of skipped stages between producer and consumer (D4). Runtime.apply_result/3 performs the capture; runtime_test.exs's "carry" describe pins it.

carry: is not available on a fan_out (removed in Phase 4.5.3): a fan-out has N items and one successor, so there was no non-arbitrary value to propagate, and the option only ever captured into each item's own context and nowhere else — a silent-bug shape. To read a per-item value downstream, filter the fan-out's [ALLM.Pipeline.Dsl.Item.t()] output with ALLM.Pipeline.Dsl.Item.ok_items/1 in the next stage and read the field off each item. Dsl.assert_carry_placement!/4 rejects carry: on a fan_out at compile time.

A key the subject does not have

carry: is validated at compile time only as a literal list of atoms (Dsl.validate_carry!/3) — whether a key is a field of whatever the stage produces is knowable only at run time. A key that is not is dropped: the carry map is unchanged and every later Context.carried(ctx, key) returns its default. Runtime.capture/3 logs a Logger.warning naming the key, the stage and the subject's struct when that happens (user decision, 2026-08-21 — warn rather than raise, since raising would change framework behaviour nothing gates). It is a detector, not a gate: the run still succeeds, so a new carry: is verified by asserting the value ARRIVES in the consuming stage, never by the declaration compiling. runtime_test.exs's "a carry: key the subject does not have is dropped LOUDLY" pins the message.

The commonest way to hit it is a typo in a field name, or a carry: key that belongs to a different stage's output than the one declaring it.

Summary

Types

concurrency_spec()

@type concurrency_spec() :: pos_integer() | {:opt, atom(), pos_integer()}

skip_spec()

@type skip_spec() ::
  {:opt, atom(), term()} | (ALLM.Pipeline.Context.t() -> as_boolean(term()))

t()

@type t() :: %ALLM.Pipeline.Dsl.Stage{
  body: (ALLM.Pipeline.Context.t(), term() -> term()) | nil,
  carry: [atom()],
  concurrency: concurrency_spec() | nil,
  input: (ALLM.Pipeline.Context.t(), term() -> struct()) | nil,
  kind: :stage | :fan_out,
  name: atom(),
  on_error: :fail_run | :continue,
  over: (term() -> [term()]) | nil,
  parent: :source_stage | :per_item,
  skip_when: skip_spec() | nil,
  step: module() | nil
}