Public API surface for Scriba projections.
The five-line API
defmodule MyApp.Projections.Orders do
use Scriba.Projection,
name: "orders",
source: {Scriba.Source.Commanded, application: MyApp.CommandedApp},
target: {Scriba.Target.Ecto, repo: MyApp.Repo},
parallelism: 16
def handle(_event, _meta), do: :skip
end
{:ok, _pid} = Scriba.start_projection(MyApp.Projections.Orders)
Scriba.info(MyApp.Projections.Orders)
Scriba.pause(MyApp.Projections.Orders)
Scriba.resume(MyApp.Projections.Orders)
Scriba.stop(MyApp.Projections.Orders)
Scriba.list()Every public function accepts either:
- Module form —
Scriba.pause(MyApp.Projections.Orders). The module's__scriba_config__/0(generated byuse Scriba.Projection) provides name and version. Refactor- friendly: rename the module, the call site renames with it. - String form —
Scriba.pause("orders", 1). Operator-friendly when you only know the projection's name. Defaults to version 1 for the arity-1 calls.
Scriba.list/0 enumerates running projections via Scriba.Registry.
Summary
Functions
How many dead letters, of what kind, over what span.
Lists a projection's dead letters, newest first.
Returns a snapshot of the projection's runtime state. See Scriba.Info for
the shape.
See info/1. Takes explicit name (string) and version (integer).
Returns running projections as [%{name, version, state}].
Pauses a running projection — the source stops yielding new events. In-flight events already in Pipeline processors or batchers continue through their commit lifecycle.
See pause/1. Takes explicit name (string) and version (integer).
Clears a projection version's progress so it can start over.
Resumes a paused projection.
Starts a projection from its module's compile-time config.
Same as start_projection/1, but merges overrides into the module's
compile-time config before starting.
Stops a running, paused or halted projection. Waits for in-flight Broadway shutdown to complete before returning.
Functions
How many dead letters, of what kind, over what span.
%{total: 143,
by_error_kind: %{"commit:23505 (unique_violation)" => 140},
oldest: ~U[...], newest: ~U[...]}Accepts the same filters as dead_letters/2. The distribution is the
diagnosis: one kind on one stream is a poison event, one kind across every
stream is a schema problem that dead-lettering is papering over.
Lists a projection's dead letters, newest first.
Scriba.dead_letters(MyApp.Projections.Orders, limit: 10)
Scriba.dead_letters(MyApp.Projections.Orders, stream_id: "order-42", order: :asc)Takes a projection module, or a name and version with an explicit
repo:. The module form reads the repo from the projection's own target
config, so it works whether or not the projection is running — which is
the point: a halted or stopped projection is exactly when someone looks.
See Scriba.DeadLetter.list/3 for the options and the shape of a row.
@spec info(module() | String.t()) :: {:ok, Scriba.Info.t()} | {:error, :not_found}
Returns a snapshot of the projection's runtime state. See Scriba.Info for
the shape.
Accepts either a projection module or a (name, version) pair. The
arity-1 string form defaults to version 1.
Returns {:error, :not_found} if no Coordinator is registered for the
given identity.
When the projection has more than 1000
distinct streams, :stream_positions is :truncated rather than a
giant map; :safe_position is always populated regardless.
:safe_position is the minimum across cached streams — an introspection
figure, not a replay point. See Scriba.Info for why
it reads high and what it is not safe to build on.
The :status field is one of :initializing | :running | :paused | :draining | :halted | :stopped. :initializing is the window in which the
Coordinator is waiting for a Pipeline and its Broadway producer to register
— at startup, and again after a Pipeline restart. It then transitions
automatically to the state the projection was in when the Pipeline was
lost: :running, :paused or :halted.
:halted means a structural commit failure stopped the projection — it is
making no progress and needs a human. :halt_reason names the cause. This
is the field to poll if you want to detect a stopped projection without
subscribing to telemetry.
@spec info(String.t(), pos_integer()) :: {:ok, Scriba.Info.t()} | {:error, :not_found}
See info/1. Takes explicit name (string) and version (integer).
@spec list() :: [%{name: String.t(), version: pos_integer(), state: atom()}]
Returns running projections as [%{name, version, state}].
Enumerates Coordinators registered in Scriba.Registry; includes
projections in any state, including :stopped (terminal Coordinators
that haven't been terminated from the supervision tree). Empty list
when no projections are running.
Pauses a running projection — the source stops yielding new events. In-flight events already in Pipeline processors or batchers continue through their commit lifecycle.
Accepts either a projection module or a name string. Arity-1 string
form defaults to version 1; see pause/2 for explicit version.
Return values
:ok— pause signal sent to the source. The source'shandle_info/2will run on its own schedule; by the time this function returns the signal is in the source's mailbox but the source may not yet have flipped its internal flag. Operators should NOT assume "no commits possible" the instantpause/1returns.{:error, {:invalid_state, state}}— projection is not in a state where pause makes sense. The inner atom is the projection's current state::initializing— the Pipeline is being started or re-monitored; retry onceinfo/1reports a settled state, which after a Pipeline restart is whichever state the projection was in before.:paused— already paused. No idempotency — callers wanting "make sure this is paused" semantics checkinfo/1first or pattern-match this error case as success.:stopped— terminal state.:draining— stop in progress.:halted— terminal. A structural commit failure stopped the projection;stop/1is the way out once the cause is fixed.
@spec pause(String.t(), pos_integer()) :: :ok | {:error, {:invalid_state, atom()}}
See pause/1. Takes explicit name (string) and version (integer).
@spec reset( module() | String.t(), keyword() ) :: {:ok, map()} | {:error, {:running, atom()}} | {:error, :no_repo}
Clears a projection version's progress so it can start over.
Deletes its rows from scriba_positions and scriba_watermarks and drops
its entries from the position cache. After this the version has no memory
of what it applied, so starting it with start_from: :origin rebuilds from
the beginning.
It does not touch the read model. Scriba does not know which tables a handler writes — that is the handler's business — so truncating them is yours to do, in the same operation, or the rebuild will apply history on top of existing rows.
Dead letters are kept by default: they record what failed on the previous
attempt and are usually the reason for the rebuild. dead_letters: true
deletes them too.
Accepts a projection that is stopped or was never started, and refuses
anything else: clearing cursors under a live pipeline lets it commit
against a cache that no longer matches the table. stop/1 leaves the
Coordinator registered in :stopped, and that counts as stopped.
:ok = Scriba.stop(MyApp.Projections.OrdersV2)
{:ok, %{positions: 1_284, watermark: 1, dead_letters: 0}} =
Scriba.reset(MyApp.Projections.OrdersV2)Options
:dead_letters—truedeletes them too (defaultfalse).:repo— required for the name form, which has no config to read it from; optional for the module form, which takes it from the projection's target.:version— for the name form only (default1).
Returns
{:ok, %{positions: n, watermark: n, dead_letters: n}}— rows removed per table.dead_lettersis0unless you asked for them.{:error, {:running, state}}— the projection is live; stop it first.{:error, :no_repo}— no repo was given and none could be resolved.
Resumes a paused projection.
Accepts a module or name string; same dispatch semantics as pause/1.
Return values
:ok— resume signal sent. Same asynchrony caveat aspause/1: first commit follows whenever the source'shandle_info/2runs and Broadway's processor stage drains accumulated demand.{:error, {:invalid_state, state}}— projection is not paused.:runningis not idempotent here either.
@spec resume(String.t(), pos_integer()) :: :ok | {:error, {:invalid_state, atom()}}
See resume/1.
Starts a projection from its module's compile-time config.
Reads the projection's __scriba_config__/0 map (generated by
use Scriba.Projection) and starts a Scriba.Projection.Supervisor
under Scriba.Projections.Supervisor. Returns the per-projection
supervisor pid.
Returns immediately (does not wait for the Coordinator to transition
from :initializing to :running). Wait for the
[:scriba, :projection, :started] telemetry event or poll
Scriba.info/1 if you need a confirmation point.
Returns {:error, :already_started} if a projection with the same
(name, version) is already running.
Same as start_projection/1, but merges overrides into the module's
compile-time config before starting.
Useful for runtime tuning of :parallelism, :batch_size, :retry,
etc. per-environment without recompiling the module.
Rejects :name and :version overrides with ArgumentError.
Identity is compile-time: passing a different name/version at runtime
would silently create a different projection rather than override
the existing one's config, which is almost certainly a bug. If you
want a second projection with a different version, declare a second
module with version: 2 in its use.
Stops a running, paused or halted projection. Waits for in-flight Broadway shutdown to complete before returning.
Returns :ok on success, {:error, {:invalid_state, state}} if the
projection is :initializing, :stopped, or :draining. A projection
that is still starting cannot be stopped — wait for info/1 to report
:running.
The projection is not removed: its Coordinator stays registered in a
terminal :stopped state, so info/1 keeps answering and reset/2 can
still clear its cursors. Starting the same (name, version) again fails
with {:error, :already_started} until the supervision child is removed:
DynamicSupervisor.terminate_child(Scriba.Projections.Supervisor, pid)Keeping the row is deliberate — a stopped projection that vanished from
list/0 would be indistinguishable from one that never started.
@spec stop(String.t(), pos_integer()) :: :ok | {:error, {:invalid_state, atom()}}
See stop/1.