Where a Sourced.EventStore.Postgres.Subscriber has got to: one row per
subscriber name in sourced_checkpoints, holding the last sequence the
subscriber committed and a fingerprint of the query it committed it over.
There is nothing to hide behind a store API here. Sequences are global and
monotonic and reads are bounded by the watermark, so resuming is
subscribe(from: sequence + 1), and the row is an integer with a name. What
it does guard are the two things that would otherwise go wrong silently:
A position is only meaningful over the query it was computed for. A subscription carries an arbitrary query, so a named position could be resumed against a different slice than it was built over and skip every event the new query matches below it.
claim/4stores a fingerprint of the query on first use and refuses to resume under a different one.Two processes under one name would move one position.
advance/4moves the row only from the position the caller last saw, so the second of two subscribers to commit a batch finds its position gone and stops, rather than both projecting the same events.
Summary
Functions
Moves the position stored under name from from to to.
Loads the position stored under name, creating it at initial_position
on first use. Once the row exists the stored position wins and
initial_position is ignored.
A fingerprint of query that ignores the order its items, types and tags
are written in, so rewriting a query without changing what it matches keeps
its checkpoint.
Functions
@spec advance( repo :: module(), name :: String.t(), from :: non_neg_integer(), to :: pos_integer() ) :: :ok | {:error, :moved}
Moves the position stored under name from from to to.
Meant to run inside the transaction that projects the batch ending at to,
so the read model and its position commit together. Returns
{:error, :moved} when the stored position is no longer from, meaning
another process has advanced the same checkpoint.
@spec claim( repo :: module(), name :: String.t(), query :: Sourced.EventStore.Query.t() | nil, initial_position :: non_neg_integer() ) :: {:ok, non_neg_integer()} | {:error, :query_changed}
Loads the position stored under name, creating it at initial_position
on first use. Once the row exists the stored position wins and
initial_position is ignored.
Returns {:error, :query_changed} when the row was created over a query
with a different fingerprint than query; see fingerprint/1.
@spec fingerprint(Sourced.EventStore.Query.t() | nil) :: String.t()
A fingerprint of query that ignores the order its items, types and tags
are written in, so rewriting a query without changing what it matches keeps
its checkpoint.
Types are compared as given: an event module and the string it is stored as
are different queries here, even though Sourced.Middleware.Domain resolves
one to the other.