A projection engine for Elixir event-sourced systems. Pairs with
Commanded for the event-sourcing side; replaces the role
commanded_ecto_projections
used to play before it stopped being actively maintained.
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(%OrderPlaced{} = event, _meta) do
{:insert, %OrderReadModel{id: event.order_id, status: "pending"}}
end
def handle(%OrderShipped{order_id: id}, _meta) do
{:update, OrderReadModel, [id: id], set: [status: "shipped"]}
end
def handle(_, _), do: :skip
end
{:ok, _} = Scriba.start_projection(MyApp.Projections.Orders)That's the API. The rest is operational scaffolding you get for free.
Already running commanded_ecto_projections? Read
MIGRATION.md. Short version: your existing cursor
carries over — both libraries track Commanded's global event_number,
so you hand your last_seen_event_number to Scriba as :start_from and
cut over in place. No read-model rebuild, no maintenance window.
Status
v0.2. The architectural contract is frozen
(SCRIBA_ARCHITECTURE.md) and the operational
primitives — telemetry, dead-letter routing and inspection, retry policy,
real pause/resume, lag reporting, multi-node standby — are in place.
Property tests cover per-stream ordering, cursor/read-model consistency,
and resume-from-cursor after a restart; fault injection against real
Postgres covers the failure taxonomy; a separate harness (bench/)
exercises acknowledgement, standby takeover and the watermark against a
real event store. Effectively-once under injected crash schedules is
enforced structurally rather than by a property test — see architecture
§10.
What's deliberately out of scope:
- A LiveView dashboard.
broadway_dashboardalready renders Scriba projections — see "Seeing a projection in LiveDashboard" below. - Throughput metrics of Scriba's own — Broadway's batch telemetry and the
per-event
:stopevents already carry the rate. - Online rebuild / shadow targets / atomic swap (v0.3).
- Sources other than Commanded; targets other than Ecto/Postgres (v0.4).
- Multi-target fan-out (v0.4).
See SCRIBA_ARCHITECTURE.md §2 for the full
in-scope / out-of-scope split.
Why this exists
You're running Commanded. You need read models. The library you'd
have reached for —
commanded_ecto_projections
— has been unmaintained for years and has known sharp edges around
crash recovery, per-stream ordering across parallelism, and lag
visibility.
Scriba is an opinionated rewrite of that role with three principles:
- Correctness over throughput. Position update and read-model
write commit atomically inside a single
Ecto.Multi. There is no configuration that lets you turn this off, because you should not want to. - Per-stream ordering is preserved. All events for a given
stream_idroute to the same processor via modular hashing of the stream id. Within a processor, events are serial. Across streams, events are parallel. - Sharp edges are documented, not hidden. Dead-lettered events advance the cursor (skip-and-continue, not block-the-projection). The handler return contract is six return shapes, not a DSL. Lag is a telemetry-consumer concern, not an engine feature.
Migrating an existing projector is a mechanical rewrite —
project %Event{}, fn multi -> ... end becomes def handle(%Event{}, meta)
returning a tagged tuple. MIGRATION.md covers the
rewrite, the Ecto.Multi differences, the callbacks with no equivalent
(after_update/3, schema_prefix/1, consistency: :strong), and cursor
carry-over.
Installation
defp deps do
[
{:scriba, "~> 0.1"},
# Optional — needed only if you use Scriba.Source.Commanded,
# which is the only source shipped in v0.1.
{:commanded, "~> 1.4"}
]
end:commanded is Scriba's one optional dependency: nothing in the engine
references it statically, so Scriba compiles without it, and
Scriba.Source.Commanded.start_link/1 raises with instructions if you
configure the Commanded source without adding it.
Everything else arrives transitively and is not optional —
:ecto_sql and :postgrex back the position cursor and dead-letter
tables (Scriba.Position, Scriba.DeadLetter, Scriba.Migrations), not
merely Scriba.Target.Ecto; :broadway is the pipeline runtime;
:telemetry and :jason are used throughout. You do not list them
yourself, but they will be in your dependency tree.
Scriba is a Broadway topology
Worth knowing before you read the failure-modes section below, which is
written in Broadway's vocabulary: each projection is a
Broadway pipeline. The source is a Broadway
producer, :parallelism sets the processor concurrency, :batch_size /
:batch_timeout configure the batcher, and event acknowledgement runs
through a Broadway acknowledger that Scriba implements against your event
store. You never write Broadway code — but when this README says "the
batch is marked failed", that is Broadway.Message.failed/2, and Broadway
is where the retry-on-redelivery behaviour comes from.
Catch-up throughput depends on :buffer_size
The event store decides how many events it will send before it requires an
acknowledgement, and its default is one. Scriba acknowledges after the batch
commits, so with a single event in flight the batcher waits out its whole
:batch_timeout before releasing the next one — which bounds catch-up far
below anything :parallelism can affect. Measured against a real EventStore,
5,000 events over 100 streams: 9.1 events/sec at the adapter default,
6,002 events/sec with buffer_size: 500.
source: {Scriba.Source.Commanded,
application: MyApp.CommandedApp,
buffer_size: 500}Scriba sets no default of its own. Raising the buffer trades memory and
redelivered-work-after-a-crash for throughput — see
Scriba.Source.Commanded for the full table and the tradeoff.
Add a migration to your repo for Scriba's tables:
defmodule MyApp.Repo.Migrations.AddScribaTables do
use Ecto.Migration
def up, do: Scriba.Migrations.up()
def down, do: Scriba.Migrations.down()
endThis creates three tables in your read-model database:
scriba_positions— per-stream cursor (one row per{projection, version, stream_id}).scriba_dead_letters— failed events for inspection / manual replay.scriba_watermarks— the contiguous global position per projection.
Upgrading from Scriba 0.1.x, whose schema had the first two, means a second migration that names the version you already have:
def up, do: Scriba.Migrations.up(from: 1)
def down, do: Scriba.Migrations.down(to: 1)Start projections during your application boot. The idiomatic pattern
is a small Task in your supervision tree that calls
Scriba.start_projection/1 once Ecto and Commanded are up:
# Bank.Application or equivalent
children = [
MyApp.Repo,
MyApp.CommandedApp,
MyApp.Projections.Starter # Task that calls Scriba.start_projection
]See examples/bank/lib/bank/projections/starter.ex
for the working pattern.
Core concepts
name and version
A projection's identity is (name, version). name is the stable
logical label ("orders"). version is an integer (default 1) you
bump when you want to run a new projection side-by-side with the old
one during a cutover.
# orders v1 — the current production projection
defmodule MyApp.Projections.OrdersV1 do
use Scriba.Projection, name: "orders", version: 1, ...
end
# orders v2 — running alongside, populating a new read model
defmodule MyApp.Projections.OrdersV2 do
use Scriba.Projection, name: "orders", version: 2, ...
endDo not encode the version in the name (name: "orders_v2"). The
macro warns at compile time when it sees that pattern, because
commanded_ecto_projections historically used it and it makes
side-by-side versioning awkward. See architecture §5.
Handler return shapes
Your handle(event, meta) clauses must return one of these six values:
| Return | Effect |
|---|---|
:skip | Event acknowledged, no read-model write, cursor does not advance for that stream. |
{:insert, schema_struct} | Ecto.Multi.insert/3 |
{:update, schema_module, filter_keyword, [set: keyword]} | Ecto.Multi.update_all/4 filtered by the keyword |
{:delete, schema_module, filter_keyword} | Ecto.Multi.delete_all/3 |
{:multi, %Ecto.Multi{}} | Merged into the batch's Multi — escape hatch for :inc, complex queries, etc. |
{:error, reason} | Routes to dead-letter (after retries, if enabled). |
Raising an exception is also valid; it's converted to a dead-letter
entry whose error_kind is the exception's module name. Per
architecture §4.2.
The meta map
The second argument to handle/2:
%{
id: "event-uuid",
stream_id: "aggregate-uuid",
type: "OrderPlaced",
position: 42,
metadata: %{correlation_id: "..."},
occurred_at: ~U[2026-05-14 12:00:00Z]
}Use :id for idempotency keys (globally unique). Use :position for
ordering within Scriba's internal accounting; don't use it as a
stable identifier across event-store rebuilds.
Per-stream ordering
Events with the same stream_id route to the same processor —
:erlang.phash2(stream_id, parallelism), plain modular hashing rather
than a consistent-hash ring — so different events on the same stream are
never processed concurrently. Events across streams are parallel,
bounded by :parallelism.
This is the invariant that makes "balance += amount" projections correct without explicit locking.
Atomic position commit
Every successful batch commits the user's read-model writes AND the
per-stream cursor advances in one Ecto.Multi transaction. There
is no observable state where the read model advanced but the cursor
didn't — or vice versa.
This is the property that makes crash recovery work: on Coordinator restart, the source is told to resume from the durably committed cursor, and Scriba's source-side dedup skips any events the upstream re-delivers below that cursor.
Operational features
Telemetry
Fifteen events fire — from the Pipeline, the Coordinator, position-cache
init, and the source. The full surface table is in Scriba.Telemetry's moduledoc
and in architecture §6.3. Highlights:
[:scriba, :projection, :event, :start | :stop | :exception]
[:scriba, :projection, :event, :skipped]
[:scriba, :projection, :batch, :stop]
[:scriba, :projection, :dead_letter]
[:scriba, :projection, :started | :paused | :resumed]
[:scriba, :projection, :cache_initialized]
[:scriba, :projection, :lag]
[:scriba, :source, :standby | :subscribed]
[:scriba, :projection, :halted]
[:scriba, :source, :batch, :failed]:event :stop fires once per successful handler invocation —
counting these gives you exact "events processed" without doing
position-arithmetic across streams. :dead_letter fires once per
dead-lettered event, with projection: %{name, version} plus position,
stream_id, event_type and error_kind metadata for alerting.
[:scriba, :source, :batch, :failed] means a batch did not commit and
nothing in it was acknowledged. For a transient failure the source then
restarts to replay from its last durable checkpoint; for a structural one it
deliberately stays stopped, because no replay can fix a missing column.
Isolated occurrences are normal under transient database trouble; a sustained
stream of them means no progress.
[:scriba, :projection, :halted] is the page. The projection hit a
structural failure — a column that doesn't exist, a missing privilege — and
stopped on purpose, because replaying it would loop forever and
dead-lettering it would destroy a batch over a fixable deploy-ordering
mistake. Nothing is lost; nothing proceeds either. The metadata carries the
SQLSTATE.
Scriba attaches no handlers of its own — it emits and gets out of the
way, so it never competes with your observability stack. You write your
own with :telemetry.attach_many/4 in your application's start/2.
Dead-letter routing
A failing event — handler returned {:error, _} or raised, AFTER
retry exhaustion — gets a row inserted into scriba_dead_letters
atomically with the cursor advance. The projection does not
block on bad events. This is a deliberate sharp edge, documented at
architecture §9.2:
"When an event is dead-lettered, the position advances past it. The alternative — blocking the projection until the bad event is resolved — stops every subsequent event for one bad one, and does it silently."
If you want block-on-failure semantics for a specific projection,
you build that on top: subscribe to [:scriba, :projection, :dead_letter] telemetry, page someone, manually replay from the
dead-letter table once they've fixed the underlying issue.
Retry policy
Default: 3 attempts with exponential backoff (100ms, 1s, 10s) before dead-letter. Configurable per projection:
use Scriba.Projection,
...
retry: [max_attempts: 5, backoff: [100, 500, 2000, 10_000]]Or retry: false to opt out (one attempt, immediate dead-letter on
failure).
The retry loop is in-handler Process.sleep — see architecture §9.1
and §9 for why this is the right primitive rather
than Process.send_after (Broadway's processor model). Each retry
attempt re-invokes :telemetry.span/3, so per-attempt
:event :start / :event :stop / :event :exception events fire.
Operators counting :event :start per event_id can see retry
activity without a dedicated :retry event.
Pause / resume
Scriba.pause(MyApp.Projections.Orders) signals the source to stop
yielding new events. The Pipeline tree stays alive; in-flight events
finish their commit lifecycle. Scriba.resume(...) reverses the
signal. Both fire telemetry; both are honest about asynchrony (the
source-pause signal is send/2, so by the time pause/1 returns the
signal is in the source's mailbox but the source's handle_info
may not yet have run).
pause on :paused and resume on :running return {:error, {:invalid_state, _}} — they are deliberately not idempotent.
Silent idempotency hides bugs. If you want "make sure this is
paused" semantics, check Scriba.info/1 first or pattern-match the
matching-state error case as success.
Failure modes worth knowing
Handler raises
Re-raise is caught by Scriba, the exception is logged as a
:event :exception telemetry event, and the event enters the retry
loop. After retry exhaustion, the original exception's struct and
stacktrace are recorded in scriba_dead_letters with error_kind
equal to the exception module name (e.g. "Elixir.ArgumentError").
Multi transaction fails
What happens depends on why, and Scriba reads that from the SQLSTATE
rather than guessing (Scriba.Failure). Guessing from "did some events
succeed?" is wrong in both directions: resource pressure fails
non-uniformly, so a partial success looks deterministic when it isn't;
and handler code deployed ahead of its migration fails uniformly, so it
looks transient when it very much isn't.
Transient (connection loss, deadlock, serialization failure, resource exhaustion, cancelled query). Nothing in the batch is acknowledged — including the events that succeeded, because event-store acks are prefix-acks and cannot express a gap. The source stops its producer, the subscription rewinds to its last durable checkpoint, and the batch is redelivered. Source-side dedup filters whatever did commit.
Integrity (unique violation, NOT NULL, foreign key, check constraint,
numeric overflow). Deterministic and specific to one event, so replaying
it forever is pointless. The batch is re-applied one transaction per
event: the offending event lands in scriba_dead_letters with
error_kind "commit:23505 (unique_violation)" or similar, the rest
commit, and the projection keeps moving.
Structural (undefined column or table, insufficient privilege, and
anything Scriba cannot classify). Neither response is safe —
dead-lettering would destroy a batch over a fixable deploy-ordering
mistake, replaying would loop forever — so the projection halts and
says so via [:scriba, :projection, :halted] and a log line naming the
SQLSTATE. Nothing is acknowledged and no cursor moves, so nothing is
lost. It resumes when you fix the cause and restart.
Multi failures do not retry through the per-event retry policy. The retry layer wraps the handler call, not the Multi commit. If your DB is intermittently failing, you want it to recover at the DB level — not for individual events to retry-then-dead-letter against a sick database.
Per-stream ordering survives all three: in the per-event pass a stream stops at its first unresolved event rather than skipping past it.
Verification status. The replay path — the source refusing to acknowledge, killing its producer, rewinding the subscription, backing off and redelivering — is exercised by the suite through the
Scriba.Test.Sourcedouble, which requeues failed messages in position order, but not against a real event store: conservation across a database outage (events delivered == read-model rows + dead letters + skipped) is unverified there. Of the failure classification,:integrityand:structuralare verified against real Postgres (test/property_db/e3_fault_injection_test.exs);:transientis covered by unit tests against constructedPostgrex.Errorstructs, including the57P01adocker stopemits.
Source redelivers events the projection has already committed
This happens after Pipeline restart — Commanded's subscription
resumes from its acked position, which may be behind Scriba's
durably committed position. Source-side dedup at the Pipeline's
handle_message/3 (architecture §3.4) catches these: events whose
position is at or below the committed cursor for their stream are
returned as :skip without invoking the handler.
test/property_db/pd3_cursor_resume_test.exs verifies this against real
Postgres over 100 iterations: after a clean stop, a projection restarted
with :start_from set to the committed cursor processes nothing it has
already applied. The integration-test side, where the source replays from
zero and dedup absorbs it, is in
test/scriba/projection/pipeline_test.exs.
Example app
examples/bank/
is a self-contained Mix project
that demonstrates the full path: real Commanded
(Commanded.EventStore.Adapters.InMemory for fast iteration), real
Ecto, real read model. The projection module itself is under 70 lines
including comments.
cd examples/bank
mix deps.get
mix bank.setup # creates the database, runs migrations
mix bank.demo # opens 3 accounts, dispatches 50 random ops, prints balances
The demo's wait-for-completion uses a telemetry counter on
:event :stop — exactly the pattern you'd reach for in your own
test code. See
examples/bank/lib/mix/tasks/bank.demo.ex.
The example app is not included in the Hex package tarball — these links go to GitHub. Clone the repo to run it.
Seeing a projection in LiveDashboard
A projection is a Broadway topology, so
broadway_dashboard renders one
with no work from Scriba — it discovers pipelines through
Broadway.all_running/0 and handles Scriba's {:via, Registry, ...} names:
# deps
{:broadway_dashboard, "~> 0.4"}
# router
live_dashboard "/dashboard",
additional_pages: [broadway: BroadwayDashboard]It shows the producer, processors and batchers, their concurrency, and
successful/failed counts per stage. Scriba ships no dashboard of its own:
the ecosystem distributes operational UIs as companion packages, and this
one already works. bench/test/broadway_dashboard_spike_test.exs keeps that
claim honest.
How far along, and how far behind
Scriba.info/1 reports two numbers an operator can alert on:
{:ok, info} = Scriba.info(MyApp.Projections.Orders)
info.watermark # 48_213 — every event up to here is accounted for
info.lag_ms # 1_240 — the event at that position happened 1.2s agoThe watermark is the contiguous global position: every event at or below it has been committed, skipped or dead-lettered, with no gap underneath. That is the number a replica could resume from, and the one that says how far a rebuild has got. Per-stream cursors cannot answer either question — a minimum across them ignores streams the projection never wrote to, and a maximum counts work sitting above an event still in flight.
Lag is measured from the event's own timestamp rather than from the event store's head, which Commanded's adapter behaviour does not expose. An idle, fully caught-up projection therefore reports the age of the last event it saw, which is what you want when asking whether anything is still flowing.
Both are written outside the commit transaction and throttled to roughly one
write a second while events are in flight, flushing immediately once the
projection catches up. They can lag what was applied; they cannot run ahead
of it. See Scriba.Watermark for why that direction is the safe one.
Running on more than one node
Every node runs the same supervision tree, so every node tries to start the projection — and a persistent subscription admits one subscriber. Scriba treats that as normal: one node acquires the subscription and projects, the others stand by, retrying about once a minute, and take over when the holder goes away. No leader election, no extra dependency, nothing to configure.
[:scriba, :source, :standby] # waiting; another subscriber holds the name
[:scriba, :source, :subscribed] # acquired it, including after a takeoverA standby's projection reports :running — its pipeline is up and healthy,
it simply has no subscription yet — so those two events, not info/1, are
what tell you which node is doing the work.
The same behaviour covers a cutover from commanded_ecto_projections: start
Scriba while the old projector still holds the name, and it picks up the
moment you stop it.
Reading the dead-letter table
Dead letters outlive the projection that produced them, so these read from the table rather than from a running process — which is the point, since a halted or stopped projection is when you go looking.
Scriba.dead_letter_stats(MyApp.Projections.Orders)
#=> %{total: 143,
# by_error_kind: %{"commit:23505 (unique_violation)" => 140,
# "Elixir.ArgumentError" => 3},
# oldest: ~U[2026-09-16 09:12:03Z], newest: ~U[2026-09-16 11:40:55Z]}
Scriba.dead_letters(MyApp.Projections.Orders, limit: 10)
Scriba.dead_letters(MyApp.Projections.Orders, stream_id: "order-42", order: :asc)The distribution is the diagnosis. One kind on one stream is a poison event.
One kind spread across every stream is a schema or handler problem that
dead-lettering is papering over — and if a whole batch fails that way at
once, Scriba.Circuit halts the projection rather than draining the stream
into the table.
A row carries :id, :position, :stream_id, :event_type,
:error_kind, :error_message, :occurred_at and the serialized
:event_data. :id orders rows that share a timestamp and is the natural
paging key alongside :limit and :offset. There is
no replay function: :event_data records what failed rather than a value
that can be re-dispatched, so replaying means reading the event from the
source by :position. Filters, paging and ordering are in
Scriba.DeadLetter.list/3.
Testing your projections
Scriba.Testing runs a projection's handle/2 clauses and commits the
results through its configured target, without starting a pipeline — so a
test asserts on read-model rows rather than on handler return values.
test "a deposit increases the balance" do
Scriba.Testing.project(MyApp.Projections.Balances, [
%AccountOpened{account_id: "acc-1"},
%Deposited{account_id: "acc-1", amount_cents: 500}
])
assert Repo.get(Balance, "acc-1").balance_cents == 500
endproject/3 returns a Scriba.Testing.Result with what committed, what the
handler skipped, what it failed on, and what it returned that the target
cannot apply. Events run in one transaction, in order, and per-stream cursors
advance exactly as they would in production:
Scriba.Testing.project(MyProjection, [
{%Deposited{}, stream_id: "acc-1"},
{%Deposited{}, stream_id: "acc-2"}
])Scriba.Testing.handle/3 calls a single clause with no database at all, for
asserting the shape a handler returns — including :skip for event types the
projection ignores.
It exercises the handler and the commit, not the pipeline around them: no retries, no dead-letter routing, no dedup, no telemetry. Handler failures come back to the caller instead of being routed, so a test can assert on them directly.
Running Scriba's own tests
| Command | What runs | Requires |
|---|---|---|
mix test.fast | Non-property tests | — |
mix test.property_db | Real-Postgres property tests in test/property_db/ | SCRIBA_TEST_DB_* env vars |
mix test.all | Everything ExUnit will run with the current environment | — for always-runnable parts; env vars for property_db |
mix test.all is the canonical full-suite command. Which command
produced a result matters: "fast suite: 99/0" and "full suite: 99/0"
describe different coverage, and a bare "99 tests, 0 failures" does not
say whether the real-Postgres tests ran at all.
Real-Postgres property tests
Property tests in test/property_db/ need a live Postgres. Set all five
of:
| Variable | Example |
|---|---|
SCRIBA_TEST_DB_HOST | localhost |
SCRIBA_TEST_DB_PORT | 5433 |
SCRIBA_TEST_DB_NAME | scriba_test |
SCRIBA_TEST_DB_USER | postgres |
SCRIBA_TEST_DB_PASS | postgres |
Local convenience: copy .env.local.example to .env.local and fill
in real values. config/test.exs auto-loads it; .env.local is
gitignored. Shell environment variables override .env.local.
Policy is all-or-nothing-or-error: all five set → tests run; all
five unset → tests excluded with a startup message; any subset
partially set → config/test.exs raises at config load. Partial
configuration is treated as a misconfiguration, not a graceful
degrade.
The database must already exist; mix test does not create it.
Migrations run in test_helper.exs against an existing connection.
docker compose up -d starts a Postgres 16 on port 5433 with
scriba_test already created, matching the example values above.
Documentation
MIGRATION.md— migrating fromcommanded_ecto_projections:project/2→handle/2,Ecto.Multidifferences, and how to carry your existing cursor across so you cut over in place instead of rebuilding read models.REBUILDING.md— rebuilding a read model from history: the(name, version)side-by-side procedure, watching progress, cutting over, and the two things that bite (side effects replay; dead letters are not replayed).SCRIBA_ARCHITECTURE.md— the architectural contract. Read this before opening a PR that changes engine behavior.examples/bank/README.md— example app walkthrough.Scriba.Telemetrymoduledoc — the full telemetry event surface.
Versioning
Scriba follows semver. The public API is start_projection, pause,
resume, stop, info, list, dead_letters, dead_letter_stats and
reset in Scriba, plus the macro at Scriba.Projection and the
test-time helpers in Scriba.Testing. The six lifecycle functions frozen
at v0.1.0 have not changed; 0.2.0 added the last three, which are
additive.
0.2.0 adds a table (scriba_watermarks), so upgrading from 0.1.x means one
migration — Scriba.Migrations.up(from: 1). Nothing else breaks.
Scriba.Target and Scriba.Source behaviours are not frozen, and will
widen if other sources or targets are built. Custom adapter authors should
pin against a specific minor version.
Only Commanded and Ecto/Postgres ship today, and nothing else is being prepared for speculatively — an interface with one implementation behind it encodes that implementation's assumptions. If you need another source or target, open an issue: it gets built with you, and the second implementation is what reveals the right shape. One constraint is worth knowing up front, because it is the guarantee rather than a detail: a target commits read-model rows, cursor advances and dead letters in a single transaction, so a store that cannot do that cannot provide effectively-once delivery.
License
Apache-2.0. See LICENSE.
Contributing
Issues and PRs welcome at https://github.com/thatsme/scriba.
Architecture-affecting changes should reference the relevant section
of SCRIBA_ARCHITECTURE.md; deliberate departures from the contract
require documentation in the PR explaining why.