FlowExtra.Pipeline (flowextra v0.6.0)

Copy Markdown View Source

Defines macros for pipeline creating.

Initialization and options

Two init/1 callbacks run when a pipeline starts, both in the starting caller, before any worker process exists — never per request:

  • the pipeline module's own init/1 receives the options passed to start/1/supervised_start/2 and returns the pipeline options;
  • each module pipe's init/1 receives its options (pipeline options merged with the pipe's declared opts:) and returns the options that pipe's call/2,3 sees. On the asynchronous engine it runs once per declared replica (count: 3 prepares three times, one per stage process); the synchronous engine executes one representative per stage and prepares once per declared occurrence.

init/1 is configuration preparation: it must return a map (module pipes) or a map/keyword list (pipeline), validated before any topology starts, and it is not a worker-resource lifecycle callback — use a process-owned mechanism for resources that must be recreated with each worker. Supervisor-driven restarts reuse the prepared options; only a new explicit start runs init/1 again.

Summary

Functions

The absolute local monotonic deadline (milliseconds) for a call timeout; :infinity maps to nil — no deadline. Local to one BEAM node.

Whether a packet's deadline has passed (nil never expires).

The call prelude, shared by every pipeline module: arms the revocable reply alias (FX-005), then hands the WHOLE packet to the admission owner (FX-001) — reservation and forwarding are one transaction inside the owner, so no caller death can strand a reserved permit without its packet, and the submission's settling wait is bounded by this call's own deadline (one budget, never reset, checked again at dequeue). The consumer incarnation the owner forwarded to is returned for the caller's monitor: one incarnation for monitor and submission. Every refusal path revokes the alias before it raises — cleanup is guaranteed on exceptional exits.

Remaining budget for an absolute deadline — :infinity when there is none.

Delivers a finished packet to the caller: a packet the pipeline itself expired raises the caller's own timeout — the caller cannot tell (and need not) whether its wait or the pipeline noticed the deadline first.

Validates module-pipe init/1 output: a module pipe's options must be a map by the time they reach a stage or the sync walker.

Validates pipeline-level init/1 output at admission: whatever init returns must be a map or keyword list — the forms the pipe walker can turn into pipe options. Shared by both engines so the refusal is loud, early, and names the culprit.

Types

t()

@type t() :: %FlowExtra.Pipeline{
  in_name: term(),
  module: module(),
  out_name: term(),
  owner_name: term() | nil,
  parent: pid() | nil,
  sup_name: term()
}

Functions

deadline(timeout)

@spec deadline(timeout()) :: integer() | nil

The absolute local monotonic deadline (milliseconds) for a call timeout; :infinity maps to nil — no deadline. Local to one BEAM node.

The deadline bounds the call's admission wait too: past its budget the caller still allows 25 ms for a late admission acknowledgment before reporting the outcome (FlowExtra.Admission slack), so a reply already on its way is heard rather than turned into uncertainty.

error_pipe(atom, options \\ [opts: [], count: 1])

(macro)

expired?(deadline)

@spec expired?(integer() | nil) :: boolean()

Whether a packet's deadline has passed (nil never expires).

pipe(atom, options \\ [opts: [], count: 1])

(macro)

prepare_call(pipeline, owner_name, ip, deadline, reply_alias)

@spec prepare_call(t(), term(), FlowExtra.IP.t(), integer() | nil, reference()) ::
  {pid(), reference()}

The call prelude, shared by every pipeline module: arms the revocable reply alias (FX-005), then hands the WHOLE packet to the admission owner (FX-001) — reservation and forwarding are one transaction inside the owner, so no caller death can strand a reserved permit without its packet, and the submission's settling wait is bounded by this call's own deadline (one budget, never reset, checked again at dequeue). The consumer incarnation the owner forwarded to is returned for the caller's monitor: one incarnation for monitor and submission. Every refusal path revokes the alias before it raises — cleanup is guaranteed on exceptional exits.

remaining(deadline)

@spec remaining(integer() | nil) :: timeout()

Remaining budget for an absolute deadline — :infinity when there is none.

unwrap!(ip, pipeline)

@spec unwrap!(FlowExtra.IP.t(), t()) :: :ok

Delivers a finished packet to the caller: a packet the pipeline itself expired raises the caller's own timeout — the caller cannot tell (and need not) whether its wait or the pipeline noticed the deadline first.

validate_module_init!(module, opts)

@spec validate_module_init!(module(), term()) :: :ok

Validates module-pipe init/1 output: a module pipe's options must be a map by the time they reach a stage or the sync walker.

validate_opts!(pipeline_module, opts)

@spec validate_opts!(module(), term()) :: map() | keyword()

Validates pipeline-level init/1 output at admission: whatever init returns must be a map or keyword list — the forms the pipe walker can turn into pipe options. Shared by both engines so the refusal is loud, early, and names the culprit.