Scriba.Watermark (Scriba v0.2.2)

Copy Markdown View Source

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

projection()

@type projection() :: %{name: String.t(), version: pos_integer()}

Functions

get(repo, map)

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

lag_ms(repo, projection)

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

put(repo, map, position, occurred_at \\ nil)

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