Scriba.Source.Commanded (Scriba v0.2.0)

Copy Markdown View Source

Broadway producer that subscribes to a Commanded event store and yields events as %Scriba.Event{}-wrapped Broadway.Messages.

Usage

source: {Scriba.Source.Commanded, application: MyApp.CommandedApp}

Options

  • :application (required) — the user's Commanded.Application module.
  • :subscription_name (default "scriba") — name passed to Commanded for durable subscription tracking.
  • :start_from (default :origin) — where to start reading: :origin, :current, or a specific event number.
  • :buffer_size, :concurrency_limit, :partition_by — forwarded verbatim to the event store adapter's subscription. Unset means the adapter's own default applies. See "Subscription buffer and throughput".

Subscription buffer and throughput

:buffer_size is how many events the store will send before it requires an acknowledgement. EventStore's default is 1, and that default, not :parallelism, is what bounds a projection's catch-up rate: Scriba acknowledges after the batch commits, so a batcher waiting on a single in-flight event waits out its full :batch_timeout before acking and releasing the next one.

Measured against a real EventStore, 5,000 events over 100 streams:

:buffer_sizeThroughput10M events
unset (adapter default, 1)9.1 events/sec12.7 days
5006,002 events/sec28 minutes

Short runs read lower — 500 events reach roughly 2,200/sec, because startup is a larger share of the measurement.

Scriba sets no default of its own — the adapter's applies unless configured. Raising it trades memory and redelivered-work-after-a-crash for throughput: up to :buffer_size events are held in flight, and an unclean restart replays whatever had not been acknowledged.

source: {Scriba.Source.Commanded,
         application: MyApp.CommandedApp,
         buffer_size: 500}

Acknowledgement watermark

Two properties of this producer make a larger buffer safe.

Acknowledgement is issued by this process, not by the Broadway batch processor that committed the batch. An event store may resolve the acking subscriber from self() and silently ignore an ack from anywhere else, which stalls the subscription for good once its buffer fills.

And only the longest gapless run of committed events is acknowledged. Acks are prefix acks — acking event 7 acks everything up to 7 — while batches commit out of source order whenever :parallelism exceeds 1. A handler still working on event 5 must therefore hold the watermark at 4, however many later events have committed; otherwise a crash in that window loses event 5 with no dead letter, no cursor anomaly and no log line, because the store believes it was delivered and nothing redelivers it.

Optional dependency (§11)

This module is the only place in Scriba that depends on :commanded. To honor optional: true, every Commanded reference is dynamic:

  • Module atoms are built from string literals (:"Elixir.Commanded.X"), so the compiler does not record them as static module dependencies.
  • Calls go through apply/3.

Result: Scriba (and this module) compile cleanly even when :commanded is absent. start_link/1 raises a clear error in that case; child_spec/1 is always safe.

Watermark persistence

The producer also records how far the projection has got. ack_contiguous/1 already computes the highest gapless committed position in order to acknowledge safely, so persisting that number costs a write rather than a second calculation: it goes to scriba_watermarks (see Scriba.Watermark), throttled to at most one write a second while events are in flight and flushed immediately once the in-flight queue drains, so an idle projection does not sit on a stale number.

The Pipeline supplies the repo and projection identity through a :scriba_watermark option it injects into every source's opts. A source that does not compute a contiguous position ignores it, and a failed write is logged rather than raised — the watermark is observability, and losing one should not cost the projection.

Standby: what happens when another subscriber holds the name

A persistent subscription admits one subscriber. Rather than failing, a producer that is refused stands by: it starts, retries in the background, and acquires the subscription when the holder releases it. On a multi-node deployment that is a warm standby — one node projects, the others wait, and a failover needs nobody's intervention.

The curve has three phases, because two different failures share it. The first five attempts are milliseconds apart (50ms to 800ms), for the case where a producer died deliberately to force a replay and the store has not yet processed the DOWN. The next thirty are a second apart: a killed supervision tree holds its registered names until its slowest in-flight handler returns, so recovery can take longer than the fast attempts cover. Only after that does the cadence settle to about a minute with jitter, so standbys that started together do not retry in lockstep.

[:scriba, :source, :standby] fires on every failed attempt and [:scriba, :source, :subscribed] when the subscription is acquired, which is how a takeover is observable. A standby's projection reports :running — its pipeline is up and healthy — so the telemetry, not the status, is what distinguishes the node doing the work from the ones waiting.

Configuration errors are not retried. An application that is not running or a store that cannot be reached raises, because retrying forever would hide it.

Pause/resume memory caveat

pause/1 sets a paused: true flag — handle_demand/2 returns no messages while paused, accumulating it in the state's :demand. The Commanded subscription, however, keeps pushing events into the source's pending :queue regardless of pause state (we don't unsubscribe).

Memory grows during pause, bounded by how many events the upstream EventStore delivers in the pause window. For operator-driven pauses (seconds to minutes on low-volume projections), this is fine. For long pauses on high-throughput projections, this is a documented sharp edge.

The cleaner alternative — unsubscribe on pause, re-subscribe on resume from the current cursor — is v0.5 hardening territory and changes EventStore subscription state in non-trivial ways. Out of scope for now.