StatifierRouter.Addresses (StatifierRouter v0.9.2)

Copy Markdown View Source

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:

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_at is not read again: a terminal execution stays terminal;
  • for a row not yet stamped, it reads the execution's status. An :active execution leaves the row as it is. A terminal one (:completed, :failed or :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: when terminal_seen_at plus the horizon is not later than the reap's time;
  • a row not yet stamped whose execution the store no longer holds - fetch_execution/2 answers {: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

result()

@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

by_execution(config, execution_id)

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

reap(config, bindings, opts \\ [])

@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 - a DateTime in UTC, the reap's time. Defaults to DateTime.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 to nil, the start of the table. The id is a next an 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}).