Scriba.Source.Commanded (Scriba v0.1.3)

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, 500 events over 50 streams:

:buffer_sizeThroughput10M events
unset (adapter default, 1)9.0 events/sec12.9 days
5002,183 events/sec76 minutes

Longer runs amortise startup and go faster still — 5,000 events at buffer_size: 500 reaches roughly 4,700/sec on the same machine.

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.

Pause/resume memory caveat (v0.1)

pause/1 sets a paused: true flag — handle_demand/2 returns no messages while paused, accumulating demand into pending_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 v0.1.