The contiguous global position a projection has reached.
Per-stream cursors (scriba_positions) answer "where is this stream?".
They cannot answer "where is this projection?", because a minimum across
them counts only streams the projection has written to, and a maximum
counts work that may sit above an event still in flight. The watermark is
the highest position P such that every event up to P has been
accounted for — committed, skipped, or dead-lettered — with no gap below it.
That single number is what makes three things possible: how far behind a projection is, how far a rebuild has got, and where a replica could resume without missing anything.
It lags reality, and only in the safe direction
The source computes the watermark when it acknowledges (see
Scriba.Source.Commanded), and persists it outside the commit
transaction. A crash between a commit and the write leaves the stored
watermark behind what was actually applied, never ahead. Resuming from a
stale watermark redelivers events that already committed, which dedup
absorbs; resuming from one that ran ahead would skip events that never did.
Only one of those is recoverable, so the write is deliberately not atomic
with the commit.
Writes are also throttled — a projection committing thousands of events a
second does not need thousands of watermark rows a second — so the stored
value trails the in-memory one by up to Scriba.Source.Commanded's write
interval even while everything is healthy.
occurred_at
Alongside the position, the row carries the occurred_at of the event at
that position, when the source knows it. now() - occurred_at is the
projection's lag in time, which is what operators alert on, and it needs no
knowledge of where the event store's head is — something Commanded's
adapter behaviour does not expose.
Summary
Functions
Reads the stored watermark, or nil when the projection has never written
one — a projection that has not yet committed anything, or one whose source
does not report a watermark.
Lag in milliseconds: how long ago the event at the watermark happened.
Records position as the projection's watermark, if it is ahead of what is
already stored.
Types
@type projection() :: %{name: String.t(), version: pos_integer()}
Functions
@spec get(module(), projection()) :: %{position: non_neg_integer(), occurred_at: DateTime.t() | nil} | nil
Reads the stored watermark, or nil when the projection has never written
one — a projection that has not yet committed anything, or one whose source
does not report a watermark.
@spec lag_ms(module(), projection()) :: non_neg_integer() | nil
Lag in milliseconds: how long ago the event at the watermark happened.
nil when there is no watermark yet, or when the source did not report an
occurred_at. This measures time, not events — an idle projection that is
fully caught up reports the age of the last event it saw, which is the
number an operator wants when asking "is anything still flowing?".
@spec put(module(), projection(), non_neg_integer(), DateTime.t() | nil) :: :ok
Records position as the projection's watermark, if it is ahead of what is
already stored.
GREATEST in SQL rather than a read-then-write: two producers for the same
projection should never both be running, but if one lingers through a
handover its late write must not move the watermark backwards.