The address table's own two functions: reap/2, which removes address
rows whose execution finished longer ago than their horizon (ADR-0002,
sections 5 and 6), and by_execution/2, which reads the row naming one
execution.
The two handles it works on
reap/2 takes the router's configuration as the storage it works on,
and the configuration carries both handles it needs:
- the host's
:repo, with the:table_prefixand:prefixthat name the address table in it (StatifierRouter.Config.table/2): the rows it reads, stamps and deletes; - the
:store, the%StatifierPersistence.Storage{}the executions are kept in: the only thing it reads there is an execution's status, throughStatifierPersistence.Storage.fetch_execution/2.
A configuration without a :store is refused by the function head. The
bindings are the second argument, the host's current ones, and not the
configuration's: the horizon is computed from the bindings handed to each
reap.
What one reap does
For each address row it examines, in the order of the rows' ids:
- a row already stamped
terminal_seen_atis not read again: a terminal execution stays terminal; - for a row not yet stamped, it reads the execution's status. An
:activeexecution leaves the row as it is. A terminal one (:completed,:failedor:cancelled) is seen terminal now, and its row is stamped with the reap's time unless the next rule deletes it in the same reap; - a row whose execution is terminal is deleted once its horizon has
elapsed since
terminal_seen_at: whenterminal_seen_atplus the horizon is not later than the reap's time; - a row not yet stamped whose execution the store no longer holds -
fetch_execution/2answers{:error, :execution_not_found}- is deleted at this reap, whatever its horizon. The row's only use is to reach that execution, and a late event for its address can no longer be resolved to it and recorded as a drop, so nothing is kept by keeping the row. A row already stamped is not read again, so an execution removed after its row was stamped frees the row by its horizon rather than by this rule.
The horizon of a row is the longest dedupe horizon, horizon_ms, of any
enabled binding among bindings whose document is the row's document.
A document no enabled binding names has a horizon of zero. So disabling
or removing every binding for a document frees the rows of its finished
executions at the next reap, including a row that reap is the first to
see terminal; a later event for such an address opens a fresh execution
under :if_absent.
It deletes address rows only. It never deletes, alters or steps an execution, and it never writes the input log.
Its cost is bounded
Every row not yet stamped costs one status read, so a reap over every
row would grow with the number of live addresses. One call examines at
most :limit rows (1000 by default), those whose id is greater than
:after, and answers with next: the id of the last row it examined
when it examined a full :limit of them, and nil when it reached the
end of the table. A host sweeps the whole table by calling again with
after: next until next is nil. A host that only ever calls reap/2
with no options examines the first :limit rows each time, which frees
nothing behind them.
The rows are read in the order of their ids, and next is an id as the
table holds it: an integer under the default bigserial key, a string
under a text key a host built with the :primary_key option of
StatifierRouter.Migrations. "Greater" is the id column's own order,
so a text key sweeps in its collation's order, and a sortable id
sweeps roughly in insertion order. The sweep needs only that the order
is total and that a row keeps its id, which a primary key guarantees;
a row inserted during a sweep behind its cursor waits for the next one.
A cursor the id column cannot hold (a string under the default
bigserial key, an integer under a text key) is refused as
{:invalid_value, :after, value} only once the repo has prepared the
query and failed to bind it, so the repo's query telemetry event
([..., :query] under the repo's telemetry prefix) fires for that
call with an error result; releases before 0.8.0 refused a string
cursor before any query and emitted no event for it. A cursor that is
neither a positive integer nor a non-empty string is refused before
the repo is called, and no query event fires for it.
An {:error, reason} from fetch_execution/2 for any examined row ends
the call before it writes anything, and is returned. The one exception is
:execution_not_found, which the rules above make a deletion rather than
a refusal: without it one orphaned row would end every sweep that reaches
it, and the rows behind it would never be examined again.
This package runs no process to call it: the host schedules it, as it
schedules StatifierRouter.Dedupe.reap/2 (ADR-0002, section 6).
Summary
Types
What one reap did: how many rows it stamped terminal_seen_at on, how
many it deleted, and the cursor to continue from, or nil at the end of
the table.
Functions
The address row naming execution_id, or nil when no row names it.
Stamps and deletes the address rows under config as the module
documentation describes, with each row's horizon computed from
bindings. A host whose configuration gives a :bindings_resolver
hands it the bindings of every scope it routes, since a row's horizon
is read by document alone (ADR-0001, the Amendment of 2026-09-25).
Types
@type result() :: %{ stamped: non_neg_integer(), deleted: non_neg_integer(), next: StatifierRouter.Schema.Id.t() | nil }
What one reap did: how many rows it stamped terminal_seen_at on, how
many it deleted, and the cursor to continue from, or nil at the end of
the table.
Functions
@spec by_execution(StatifierRouter.Config.t(), String.t()) :: StatifierRouter.Schema.Address.t() | nil
The address row naming execution_id, or nil when no row names it.
ADR-0002, section 1's unique index is on (scope, document, key), and
the execution_id index StatifierRouter.Migrations.V01.up/1 adds is
not unique; what keeps the count at one row per execution is that each
create mints a fresh id and writes at most one row for it (ADR-0006,
section 1). This function reads the first row in id order rather than
asserting that invariant, so a second row cannot turn a read into a
raise.
An execution created under :always_new has no row at all (ADR-0002,
section 7), and so has neither a scope nor a key of its own: that is the
nil an execution-to-execution send is refused for as
unaddressed_sender (ADR-0006, section 6).
@spec reap(StatifierRouter.Config.t(), [StatifierRouter.Binding.t()], keyword()) :: {:ok, result()} | {:error, term()}
Stamps and deletes the address rows under config as the module
documentation describes, with each row's horizon computed from
bindings. A host whose configuration gives a :bindings_resolver
hands it the bindings of every scope it routes, since a row's horizon
is read by document alone (ADR-0001, the Amendment of 2026-09-25).
opts:
:now- aDateTimein UTC, the reap's time. Defaults toDateTime.utc_now/0.:limit- a positive integer, the most rows this call examines. Defaults to 1000.:after-nil, or an id of the table: examine only rows whose id is greater. Defaults tonil, the start of the table. The id is anextan earlier call answered: a positive integer under the default key, a string under a text key. One the table's id column cannot hold is refused as{:invalid_value, :after, value}.
Returns {:ok, result}, {:error, reason} from
StatifierPersistence.Storage.fetch_execution/2 other than
:execution_not_found, which deletes the row instead, or
{:error, reason} for a malformed option ({:invalid_opts, opts},
{:unknown_key, name}, {:invalid_value, name, value}).