Sourced.EventStore.Postgres.Checkpoint (sourced_postgres v0.3.0)

Copy Markdown View Source

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/4 stores 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/4 moves 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

advance(repo, name, from, to)

@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.

claim(repo, name, query, initial_position \\ 0)

@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.

fingerprint(query)

@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.