Sourced. EventStore. Postgres. Subscriber behaviour
(sourced_postgres v0.3.0)
Copy Markdown
View Source
A subscriber that persists the last sequence that it processed.
Sourced.EventStore.subscribe/2 delivers at-most-once and leaves the
subscriber's position in the subscriber's memory, so a subscriber that
restarts has nothing to resume from. This module keeps that position in a
Sourced.EventStore.Postgres.Checkpoint row instead, and commits it in the
same transaction as whatever handle_events/2 writes.
Usage
defmodule MyApp.OrderProjector do
use Sourced.EventStore.Postgres.Subscriber
@impl true
def handle_events(events, state) do
Enum.each(events, &MyApp.Orders.project/1)
{:ok, state}
end
endStart it under a supervisor, after the store it reads from:
children = [
MyApp.Repo,
{Sourced.EventStore, store},
{MyApp.OrderProjector, store: store, name: "orders", query: [%{tags: ["order"]}]}
]handle_events/2 runs inside a transaction on the store's repo, so writes
it makes through that repo commit with the checkpoint. Returning
{:error, reason} or raising rolls the transaction back and stops the
subscriber, and its supervisor starts it again from the last committed
position. Anything with an effect outside that transaction — sending an
email, calling another service — happens once per delivery of the batch
rather than once per event, and needs its own idempotency.
Handlers that do not write to the repo
For a handler whose work is not in the database at all the transaction buys
nothing — an HTTP call cannot be rolled back — and costs a pooled connection
held open for however long that work takes. Override transaction?/2 to
return false and the handler runs outside any transaction, with the
checkpoint advanced in one short statement of its own once the handler has
returned {:ok, _}:
@impl true
def transaction?(_events, _state), do: falseIt is asked per batch, with the events about to be handled, so a handler can also decide by what is in them.
That is at-least-once delivery: a crash between the handler finishing and the checkpoint committing redelivers the batch on the next start, so the handler has to be idempotent. Only use it for handlers that do not write through the store's repo, since one that does would silently lose the atomicity above.
Options
:store— required, theSourced.EventStoreto subscribe to.:name— required, the name the position is stored under, unique per read model.:query— theSourced.EventStore.Queryto subscribe with. Defaults to every event.:start_from— where a new name starts::origin(the default) for every event ever stored,:latestfor only what is appended from now on, or a sequence for the first event to deliver. Applies only the first time the name is seen; once a position exists it is resumed from regardless, or:latestwould skip events on every restart.
The whole list is passed on to init/1.
What the name guarantees
A checkpoint is a position over one query, advanced by one process:
The query's fingerprint is stored with the position on first use, and starting under the same name with a different query raises rather than resuming, since the new query could match events below the stored position that were never delivered. To change what a read model subscribes to, rebuild it under a new name — or delete the row, if the change is known not to matter.
Two subscribers under one name is a configuration error, and is detected rather than tolerated: each moves the checkpoint only from the position it last committed, so the second of the two to commit a batch finds its position gone and stops with
{:checkpoint_moved, name}, its own writes for that batch rolling back with it.
The subscription underneath
The subscriber monitors the process delivering its subscription (see
Sourced.EventStore.Subscription) and, when that process goes down, stops
with {:subscription_down, reason} rather than resubscribing in place. Its
supervisor starting it again is the same recovery every other failure gets —
claim the checkpoint, subscribe from it, catch up — and one path for all of
them is worth more than sparing a restart on an event that is rare to begin
with: the store's listening connection reconnects on its own without its
process dying, so this is a crash path, not a network one.
Summary
Callbacks
Handles a batch of events, in ascending sequence order, inside the transaction that commits their position.
Called once when the subscriber starts, with the options it was started
with. Returns the state handle_events/2 is first called with.
Whether handle_events/2 runs inside the transaction that commits the
batch's position. Asked once per batch, with the events about to be handled.
Functions
Returns a specification to start this module under a supervisor.
Starts a subscriber running the callbacks of module. use-ing this module
defines a start_link/1 that calls this.
Types
@type option() :: {:store, Sourced.EventStore.t()} | {:name, String.t()} | {:query, Sourced.EventStore.Query.t()} | {:start_from, start_from()} | {atom(), term()}
@type start_from() :: :origin | :latest | non_neg_integer()
Callbacks
@callback handle_events(events :: [Sourced.StoredEvent.t()], state :: term()) :: {:ok, state :: term()} | {:error, reason :: term()}
Handles a batch of events, in ascending sequence order, inside the transaction that commits their position.
Called once when the subscriber starts, with the options it was started
with. Returns the state handle_events/2 is first called with.
Defaults to {:ok, opts}.
@callback transaction?(events :: [Sourced.StoredEvent.t()], state :: term()) :: boolean()
Whether handle_events/2 runs inside the transaction that commits the
batch's position. Asked once per batch, with the events about to be handled.
Defaults to true. See "Handlers that do not write to the repo" above for
when to return false.
Functions
Returns a specification to start this module under a supervisor.
See Supervisor.
@spec start_link(module(), [option()]) :: GenServer.on_start()
Starts a subscriber running the callbacks of module. use-ing this module
defines a start_link/1 that calls this.