Scriba (Scriba v0.2.0)

Copy Markdown View Source

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 by use 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 or paused projection. Waits for in-flight Broadway shutdown to complete before returning.

Functions

dead_letter_stats(module_or_name, opts \\ [])

@spec dead_letter_stats(
  module() | String.t(),
  keyword()
) :: map()

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.

dead_letters(module_or_name, opts \\ [])

@spec dead_letters(
  module() | String.t(),
  keyword()
) :: [map()]

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.

info(module_or_name)

@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 brief window between Coordinator start and Broadway producer registration; transitions to :running automatically.

: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.

info(name, version)

@spec info(String.t(), pos_integer()) :: {:ok, Scriba.Info.t()} | {:error, :not_found}

See info/1. Takes explicit name (string) and version (integer).

list()

@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.

pause(module_or_name)

@spec pause(module() | String.t()) :: :ok | {:error, {:invalid_state, atom()}}

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's handle_info/2 will 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 instant pause/1 returns.

  • {:error, {:invalid_state, state}} — projection is not in a state where pause makes sense. The inner atom is the projection's current state:

    • :initializing — engine starting up; retry once info/1 reports :running.
    • :paused — already paused. No idempotency — callers wanting "make sure this is paused" semantics check info/1 first 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/1 is the way out once the cause is fixed.

pause(name, version)

@spec pause(String.t(), pos_integer()) :: :ok | {:error, {:invalid_state, atom()}}

See pause/1. Takes explicit name (string) and version (integer).

reset(module_or_name, opts \\ [])

@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 — true deletes them too (default false).
  • :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 (default 1).

Returns

  • {:ok, %{positions: n, watermark: n, dead_letters: n}} — rows removed per table. dead_letters is 0 unless 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.

resume(module_or_name)

@spec resume(module() | String.t()) :: :ok | {:error, {:invalid_state, atom()}}

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 as pause/1: first commit follows whenever the source's handle_info/2 runs and Broadway's processor stage drains accumulated demand.

  • {:error, {:invalid_state, state}} — projection is not paused. :running is not idempotent here either.

resume(name, version)

@spec resume(String.t(), pos_integer()) :: :ok | {:error, {:invalid_state, atom()}}

See resume/1.

start_projection(module)

@spec start_projection(module()) :: {:ok, pid()} | {:error, :already_started | term()}

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.

start_projection(module, overrides)

@spec start_projection(
  module(),
  keyword()
) :: {:ok, pid()} | {:error, :already_started | term()}

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.

stop(module_or_name)

@spec stop(module() | String.t()) :: :ok | {:error, {:invalid_state, atom()}}

Stops a running or paused 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.

stop(name, version)

@spec stop(String.t(), pos_integer()) :: :ok | {:error, {:invalid_state, atom()}}

See stop/1.