StatifierPersistence.Storage.Ecto (StatifierPersistence v0.21.0)

Copy Markdown View Source

The Ecto StatifierPersistence.Storage.Adapter: the storage contract over the schemas a host generates with use StatifierPersistence.Ecto (ADR-0002), against the tables the versioned migrations helper creates. Requires the optional ecto_sql dependency (ADR-0005).

defmodule MyApp.Persistence do
  use StatifierPersistence.Ecto, repo: MyApp.Repo
end

{:ok, store} =
  StatifierPersistence.Storage.new(
    StatifierPersistence.Storage.Ecto,
    persistence: MyApp.Persistence
  )

Options init/1 accepts:

  • :persistence - required, a module that called use StatifierPersistence.Ecto. The repo, the schema modules, and the table names all come from its resolved configuration, so this adapter adds no knobs of its own (ADR-0002 decision 3).
  • :input_log_cap - the per-execution cap on ADR-0010's input log: :infinity (the default, no cap) or a positive integer. It counts entries, and past it append_input/3 refuses and closes the execution's log with a marker (decision 6). It has no bounded default, because one would be this package silently truncating a host's log.
  • :sandbox - when true, isolate/1 checks out an Ecto.Adapters.SQL.Sandbox connection: the hook a test suite (this package's conformance suite included) uses to wrap each test in its own transaction. Default false, and isolate/1 is then a no-op.

Engine identities (content_hash, session_id, execution_id) are stored verbatim in text columns and blobs in bytea columns, so both round-trip byte-identically (ADR-0002 decision 1, ADR-0003 decision 1). The identity guard lives in StatifierPersistence.Storage, above this adapter like above every other one (ADR-0003 decision 2); nothing here decodes a blob.

insert_execution/2's :execution_exists refusal rides the V01 unique index on execution_id - one atomic insert, never a check-then-insert. A backend failure a callback cannot observe as a value (the database down, a timeout) raises the driver's own exception rather than being flattened into a default (this package's errors-are-events rule).

Summary

Functions

Appends one input at the execution's next ordinal (the optional StatifierPersistence.Storage.Adapter.append_input/3).

Counts the executions on content_hash per stored arm (the optional StatifierPersistence.Storage.Adapter.count_executions_by_content_hash/2, ADR-0012 decision 3).

Fetches the chart stored under content_hash, or :chart_not_found.

Fetches the execution stored under execution_id, or :execution_not_found.

Fetches the position stored for session_id, or :position_not_found.

Reads the tombstone on content_hash, or nil for a hash that has none (the optional StatifierPersistence.Storage.Adapter.fetch_retired_info/2).

Resolves the :persistence host module into the handle every other callback takes: the host's repo, its three generated schema modules, and its executions table name (for the unique-constraint mapping).

Inserts execution_record, refusing a duplicate execution_id with {:error, :execution_exists}.

Per-test isolation (the optional StatifierPersistence.Storage.Adapter.isolate/1): checks out an Ecto.Adapters.SQL.Sandbox connection when this handle was built with sandbox: true, and is a no-op otherwise.

Lists the ids of the :active executions on content_hash (the optional StatifierPersistence.Storage.Adapter.list_active_execution_ids_by_content_hash/2, ADR-0012 decision 4).

Lists the ids of the executions on content_hash whose status is one of statuses (the optional StatifierPersistence.Storage.Adapter.list_execution_ids_by_content_hash/3, ADR-0017 decision 8), in ascending execution id.

Lists the executions whose stored metadata contains every key/value pair in metadata (ADR-0006 decision 3's equality-match list helper).

Lists an execution's whole log in ascending seq (the optional StatifierPersistence.Storage.Adapter.list_inputs/2), or :execution_not_found for an execution this adapter does not hold.

Runs fun under per-execution mutual exclusion for execution_id (the optional StatifierPersistence.Storage.Adapter.lock_execution/3, ADR-0004 decision 5 as amended 2026-08-22).

Prunes one batch of finished executions (the optional StatifierPersistence.Storage.Adapter.prune_executions/4, ADR-0016) in one transaction.

Retires the chart on content_hash: the counts and the tombstone in one transaction (the optional StatifierPersistence.Storage.Adapter.retire_chart/3, ADR-0012 decisions 5 and 6).

Stores chart_record, idempotent on its content_hash: an insert with on_conflict: :nothing against the unique index, so a repeated save of the same hash neither duplicates the row nor rewrites it.

Stores position_record under its session_id, overwriting any position already stored for that session: an upsert replacing the record columns (and updated_at) on the unique index.

Declares whether the store this adapter is pointed at can be tombstoned (the optional StatifierPersistence.Storage.Adapter.supports_chart_retirement?/1, ADR-0012 decision 6).

Declares the drained query (the optional StatifierPersistence.Storage.Adapter.supports_content_hash_query?/1): this adapter counts the executions table's rows on one content_hash, on every Ecto backend.

Declares outcome support (the optional StatifierPersistence.Storage.Adapter.supports_execution_outcome?/1): this adapter stores an execution's answer in the V03 outcome_blob column, under the configured :blob_type like every other blob column.

Declares execution pruning (the optional StatifierPersistence.Storage.Adapter.supports_execution_pruning?/1, ADR-0016) on every Ecto backend: the batch is a select, a delete and an update over columns V05 and V08 give every backend.

Declares input log support (the optional StatifierPersistence.Storage.Adapter.supports_input_log?/1): this adapter keeps ADR-0010's log in the V05 inputs table, on every Ecto backend.

Declares metadata support (the optional StatifierPersistence.Storage.Adapter.supports_metadata?/1): this adapter stores an execution's metadata in the V02 jsonb column (ADR-0006 decision 3), on a Postgres repo.

Declares the tree migration unit (the optional StatifierPersistence.Storage.Adapter.supports_tree_migration?/1, ADR-0015 decision 3): its writes are UPDATEs of existing columns, and the linkage pin lives in the existing metadata column, so it needs no schema version.

Overwrites the execution stored under execution_record's execution_id with the full record, or refuses with :execution_not_found.

Writes a tree migration's re-pins and parks in one transaction (the optional StatifierPersistence.Storage.Adapter.write_tree_migration/2, ADR-0015 decision 3).

Functions

append_input(opts, execution_id, map)

Appends one input at the execution's next ordinal (the optional StatifierPersistence.Storage.Adapter.append_input/3).

The ordinal is the execution's current maximum plus one, read and written inside the exclusion the caller already holds (StatifierPersistence.Executions' serialized unit). The V05 unique index on (execution_id, seq), not the read, is what makes denseness true: a lost race fails the write with {:adapter, :seq_conflict} rather than duplicating an ordinal.

Past the configured input_log_cap: the log closes itself. The cap's last slot is written as a row with a nil input_blob - decision 6's marker - and this call and every later one for that execution return {:error, :input_log_full}. The step that produced the input is not failed by it: the log records its own truncation instead (ADR-0010 decision 5).

count_executions_by_content_hash(opts, content_hash)

Counts the executions on content_hash per stored arm (the optional StatifierPersistence.Storage.Adapter.count_executions_by_content_hash/2, ADR-0012 decision 3).

One grouped count against the V07 index on executions(content_hash): no row is loaded and no blob is read, which is the point of a callback that answers "what is running on this chart" rather than "which executions are".

The database answers only the arms it holds rows in, so the grouped result is folded onto a map of zeros: every key is present for every hash, and a hash this store has never seen answers zeros.

children is a second query, because it counts something the executions table's own content_hash column does not hold: the linkage pins naming this hash whose parent execution is :active or :needs_migration (ADR-0012 decision 1, as ADR-0014 decision 4 reads it). It is a containment probe on the reserved metadata key joined to the parent row by execution_id, and containment is the operator V03's jsonb_path_ops GIN index on metadata serves.

Off Postgres it is 0, and that is the count rather than a gap: supports_metadata?/1 is false there, so StatifierPersistence.Storage.insert_execution/5 refuses a metadata: option at open and no linkage reaches the table through a supported door.

fetch_chart(opts, content_hash)

Fetches the chart stored under content_hash, or :chart_not_found.

A tombstoned hash answers {:error, {:chart_retired, info}} instead, carrying who retired it and when (ADR-0012 decision 6): the row is still there and its blobs are nil, and a record with nil blobs is not a chart this adapter holds.

fetch_execution(opts, execution_id)

Fetches the execution stored under execution_id, or :execution_not_found.

fetch_position(opts, session_id)

Fetches the position stored for session_id, or :position_not_found.

fetch_retired_info(opts, content_hash)

Reads the tombstone on content_hash, or nil for a hash that has none (the optional StatifierPersistence.Storage.Adapter.fetch_retired_info/2).

One SELECT of retired_at and retired_by on the unique index, restricted to a tombstoned row, so neither chart blob crosses the wire and the cost does not grow with the chart.

init(opts)

Resolves the :persistence host module into the handle every other callback takes: the host's repo, its three generated schema modules, and its executions table name (for the unique-constraint mapping).

Refuses a module that never called use StatifierPersistence.Ecto with {:error, {:adapter, {:not_a_persistence_host, module}}}. Makes no database call: reachability surfaces on first use, per call site.

insert_execution(opts, execution_record)

Inserts execution_record, refusing a duplicate execution_id with {:error, :execution_exists}.

The refusal is the V01 unique index on execution_id speaking: the insert carries a unique_constraint/3 on that index's name, so two concurrent inserts of one execution_id cannot both return :ok and no separate existence check ever runs.

metadata is stored in the V02 jsonb column, NULL for the empty map. jsonb holds only JSON-representable values, which makes term narrower here than in Elixir (ADR-0006 decision 3): a tuple, a pid, a reference, an atom, a struct, or a binary that is not valid UTF-8 has no jsonb form. This adapter refuses such a map at open with {:error, :metadata_unsupported} - the failure shape ADR-0006 decision 3 leaves to the implementation - rather than letting the encoder raise from inside a transaction or, worse, storing something that is not what the caller handed over. Refusing at open is the principle the decision already sets for an adapter that cannot store a map; a value it cannot store is the same answer at a finer grain.

isolate(opts)

Per-test isolation (the optional StatifierPersistence.Storage.Adapter.isolate/1): checks out an Ecto.Adapters.SQL.Sandbox connection when this handle was built with sandbox: true, and is a no-op otherwise.

list_active_execution_ids_by_content_hash(opts, content_hash)

Lists the ids of the :active executions on content_hash (the optional StatifierPersistence.Storage.Adapter.list_active_execution_ids_by_content_hash/2, ADR-0012 decision 4).

One column, one arm, under the same V07 index on executions(content_hash) the grouped count uses: the ids are what a pin source is handed as its context, and nothing else about those executions is read.

list_execution_ids_by_content_hash(opts, content_hash, statuses)

Lists the ids of the executions on content_hash whose status is one of statuses (the optional StatifierPersistence.Storage.Adapter.list_execution_ids_by_content_hash/3, ADR-0017 decision 8), in ascending execution id.

One column under the same V07 index on executions(content_hash) the :active listing uses, with the status set as an IN list. The ids are sorted in Elixir rather than by an ORDER BY, so the order is the binary order of the ids whatever collation the backend's text column carries.

list_execution_states_by_metadata(opts, metadata)

The indexed status projection over a metadata match (the optional StatifierPersistence.Storage.Adapter.list_execution_states_by_metadata/2).

The same jsonb containment predicate list_executions_by_metadata/2 issues - which V03's GIN jsonb_path_ops index serves - with a three-column select: in place of the whole row. No blob column is read, which is the point: a fan-out of N children asks this question N times, and the listing would move N identity and position blobs each time.

child_index is extracted from this package's own reserved linkage namespace inside metadata, and is nil for a matched execution carrying no linkage.

Takes the same non-empty string-keyed map, with the same ArgumentError for anything else.

Off Postgres this refuses with {:error, :metadata_unsupported} rather than issuing SQL the backend cannot parse - the same answer supports_metadata?/1 already gives the facade, given directly to a caller who reached the callback itself.

list_executions_by_metadata(opts, metadata)

Lists the executions whose stored metadata contains every key/value pair in metadata (ADR-0006 decision 3's equality-match list helper).

Equality match on all pairs is the whole query surface: no ranges, no partial matches, no containment operators exposed to the caller, and no ordering guarantee. A host needing more than that queries its own column directly - ADR-0002's configurable table names already make that a supported thing to do.

The query is one jsonb containment predicate, which a GIN index on the column serves directly; V02 ships no index, because which pairs a host queries by is the host's call (ADR-0006 decision 4).

metadata must be a non-empty map of string keys: a zero-pair "contains every given pair" matches every execution with any metadata at all, which is a caller bug far more often than a request, so it raises ArgumentError rather than answering it.

StatifierPersistence.Storage.Ecto.list_executions_by_metadata(
  store.opts,
  %{"tenant_id" => "acct_01H8X"}
)

Returns records in fetch_execution/2's shape.

Off Postgres this refuses with {:error, :metadata_unsupported} rather than issuing SQL the backend cannot parse - the same answer supports_metadata?/1 already gives the facade, given directly to a caller who reached the callback itself.

list_inputs(opts, execution_id)

Lists an execution's whole log in ascending seq (the optional StatifierPersistence.Storage.Adapter.list_inputs/2), or :execution_not_found for an execution this adapter does not hold.

One index-ordered read of the V05 unique index. No filter, no range, no limit: the whole log is what a replay consumes and the cap is what bounds it (ADR-0010 decision 2).

lock_execution(opts, execution_id, fun)

@spec lock_execution(
  StatifierPersistence.Storage.Adapter.opts(),
  StatifierPersistence.Storage.Adapter.execution_id(),
  (-> result)
) :: {:ok, result} | {:error, StatifierPersistence.Storage.Adapter.error()}
when result: term()

Runs fun under per-execution mutual exclusion for execution_id (the optional StatifierPersistence.Storage.Adapter.lock_execution/3, ADR-0004 decision 5 as amended 2026-08-22).

Everything happens inside one transaction that spans fun. It takes pg_advisory_xact_lock(hashtextextended(execution_id, 0)) first - unconditional per-execution exclusion whether or not the execution row exists yet - and then SELECT ... FOR UPDATE on the execution row when it does, keeping the row itself locked against every other writer for the rest of the transaction. Both locks are transaction-scoped, so any exit from fun releases them: a normal return commits, and a raise rolls back and propagates to the caller with nothing leaked.

prune_executions(opts, cutoff, limit, scope)

Prunes one batch of finished executions (the optional StatifierPersistence.Storage.Adapter.prune_executions/4, ADR-0016) in one transaction.

The batch is selected on the V08 index on executions(ended_at), oldest first, and its input log rows are deleted and its position blobs nulled by id. On Postgres the selection locks its rows with FOR UPDATE SKIP LOCKED: a row another transaction is writing is left for the next call rather than waited on, and two prunes running at once take disjoint batches. Off Postgres the backend's own write lock serialises the transaction.

A scope confines every statement - the selection, the input log check inside it, the input log delete and the position blob update - to the rows holding each of its column equalities. Its columns are the host's :leading_columns, which the generated schemas do not declare, so a scoped batch queries the two tables by name, under the schemas' own prefix, and reads the package columns it needs by the same names. A column that is not one of the host's :leading_columns raises ArgumentError before any statement runs.

This transaction joins a caller's own when there is one. Nothing here rolls back: every statement either runs or raises the driver's own exception.

retire_chart(opts, content_hash, retirement)

Retires the chart on content_hash: the counts and the tombstone in one transaction (the optional StatifierPersistence.Storage.Adapter.retire_chart/3, ADR-0012 decisions 5 and 6).

The transaction is what the callback's contract asks for, and one thing inside it is worth naming, because it is what closes the race the record cares about. The tombstone is not an update the counts authorise; it is a single conditional UPDATE that re-asserts every one of them in its own WHERE - no :active or :needs_migration execution row on the hash, no position row on it, no durable-child pin naming it, and the row not already retired. The counts taken above it are what a refusal reports; the UPDATE is what decides. So there is no interval between the count and the write for an execution to be created in: an execution visible when the statement runs is in its NOT EXISTS, and one committed after it is after the tombstone. A statement that matches no row is read back as the refusal it is - the counts are taken again and reported, or the retired arm is answered if a concurrent retirement won.

This transaction joins a caller's own when there is one, and it takes no per-execution lock: the two contracts the README's "Writing inside a caller's transaction" and "Delivering while a step is in flight" sections state are untouched by it.

No row: :chart_not_found. Already retired: the retired arm, never a second tombstone.

save_chart(opts, chart_record)

Stores chart_record, idempotent on its content_hash: an insert with on_conflict: :nothing against the unique index, so a repeated save of the same hash neither duplicates the row nor rewrites it.

A tombstoned hash is refused with {:error, {:chart_retired, info}} and not revived (ADR-0012 decision 6). What keeps a tombstoned row from being revived is the unique index, not the read: the insert never rewrites an existing row, so no interleaving puts the bytes back. The read is what turns a save of a retired hash into the retired arm rather than a silent :ok, and it answers for the row as it stood when it ran. The read and the insert are one transaction, but under Postgres's default READ COMMITTED each statement reads its own snapshot, so a retirement committing between the two leaves this save answering :ok - which is the answer the save would have had in the order it was read in, a save followed by a retirement, and the row stays tombstoned. The transaction joins a caller's own when there is one, which is what keeps the README's "Writing inside a caller's transaction" contract intact.

save_position(opts, position_record)

Stores position_record under its session_id, overwriting any position already stored for that session: an upsert replacing the record columns (and updated_at) on the unique index.

supports_chart_retirement?(opts)

@spec supports_chart_retirement?(StatifierPersistence.Storage.Adapter.opts()) ::
  boolean()

Declares whether the store this adapter is pointed at can be tombstoned (the optional StatifierPersistence.Storage.Adapter.supports_chart_retirement?/1, ADR-0012 decision 6).

It asks the store, not the backend. The retirement writes NULL into identity_blob and chart_blob and writes retired_at and retired_by, so what has to be true is that those four columns are there and the two blob columns are nullable - which is exactly what migration V07 arranges, and which V07 can only arrange on Postgres, because ecto_sqlite3 raises from the modify/3 that drops a NOT NULL.

Asking the catalog rather than the adapter module keeps the answer true for the host V07's moduledoc sends elsewhere: one that altered the two columns itself, in a migration of its own on a backend that is not Postgres, has a store that can be retired against and this predicate says so. One query, on a call a host makes once per retirement.

The probe reads the columns' presence and the two blobs' nullability, not their types. A host that adds the tombstone columns itself with types other than V07's - retired_at as text, say - is answered true here, and the retirement's UPDATE then raises the database's own error rather than refusing at open. V07 writes the types this adapter writes; a hand-built store has to match them.

supports_content_hash_query?(opts)

@spec supports_content_hash_query?(StatifierPersistence.Storage.Adapter.opts()) ::
  boolean()

Declares the drained query (the optional StatifierPersistence.Storage.Adapter.supports_content_hash_query?/1): this adapter counts the executions table's rows on one content_hash, on every Ecto backend.

Unconditional, unlike supports_metadata?/1. The query is an equality predicate on a text column with a GROUP BY over another, which every backend this package tests parses, and V07 indexes the column it filters on for every backend too.

supports_execution_outcome?(opts)

@spec supports_execution_outcome?(StatifierPersistence.Storage.Adapter.opts()) ::
  boolean()

Declares outcome support (the optional StatifierPersistence.Storage.Adapter.supports_execution_outcome?/1): this adapter stores an execution's answer in the V03 outcome_blob column, under the configured :blob_type like every other blob column.

supports_execution_pruning?(opts)

@spec supports_execution_pruning?(StatifierPersistence.Storage.Adapter.opts()) ::
  boolean()

Declares execution pruning (the optional StatifierPersistence.Storage.Adapter.supports_execution_pruning?/1, ADR-0016) on every Ecto backend: the batch is a select, a delete and an update over columns V05 and V08 give every backend.

supports_input_log?(opts)

@spec supports_input_log?(StatifierPersistence.Storage.Adapter.opts()) :: boolean()

Declares input log support (the optional StatifierPersistence.Storage.Adapter.supports_input_log?/1): this adapter keeps ADR-0010's log in the V05 inputs table, on every Ecto backend.

Unconditional, unlike supports_metadata?/1, which answers false off Postgres because the containment SQL its listings issue does not parse there. Nothing in the input log needs a Postgres-only feature - a table, four columns and a unique index (ADR-0010 decision 9) - so there is nothing here for a backend to decline.

supports_metadata?(opts)

@spec supports_metadata?(StatifierPersistence.Storage.Adapter.opts()) :: boolean()

Declares metadata support (the optional StatifierPersistence.Storage.Adapter.supports_metadata?/1): this adapter stores an execution's metadata in the V02 jsonb column (ADR-0006 decision 3), on a Postgres repo.

On any other Ecto adapter this answers false, which is ADR-0006 decision 3's refusal-at-open arm rather than a new one. The capability that record defines is the column and the equality-match list helper, and the helper is Postgres-only SQL: both list_executions_by_metadata/2 and list_execution_states_by_metadata/2 are jsonb containment (@>) with a -> ... ->> extraction, which a non-Postgres backend does not parse. Declaring the capability true there would strand a durable subchart or a fan-out at the far end of a listing that cannot answer, with children already started that nothing could then settle (sp-11w).

The two listings consult this answer themselves, so a caller holding the raw callback gets the same {:error, :metadata_unsupported} the facade gives rather than a raise from the driver (sp-4eo). V03's metadata index is skipped on the same adapters, for the same reason; sp-5lm tracks this surface.

supports_retired_info?(opts)

@spec supports_retired_info?(StatifierPersistence.Storage.Adapter.opts()) :: boolean()

Declares the narrow tombstone read (the optional StatifierPersistence.Storage.Adapter.supports_retired_info?/1).

Always true, and unlike supports_chart_retirement?/1 it does not ask the store: the read selects retired_at and retired_by, which fetch_chart/2 and save_chart/2 already read on every store this adapter serves. A store that can never carry a tombstone answers nil for every hash, cheaply, which is the answer a create needs.

supports_tree_migration?(opts)

@spec supports_tree_migration?(StatifierPersistence.Storage.Adapter.opts()) ::
  boolean()

Declares the tree migration unit (the optional StatifierPersistence.Storage.Adapter.supports_tree_migration?/1, ADR-0015 decision 3): its writes are UPDATEs of existing columns, and the linkage pin lives in the existing metadata column, so it needs no schema version.

update_execution(opts, execution_record)

Overwrites the execution stored under execution_record's execution_id with the full record, or refuses with :execution_not_found.

One update_all/3 keyed on execution_id: the match count is the existence check, so refusal and overwrite are a single statement.

metadata is not in the set: list, and that is the documented exception to the full overwrite: the map is write-once (ADR-0006 decision 1 grants it at create and grants no way to change it), so the stored column is left exactly as insert_execution/2 wrote it and the given record's metadata is ignored.

outcome_blob is the second exception, and it joins the set: list only when the given record carries one: a nil leaves the stored column alone, so an ordinary step of an execution that has already answered does not erase its answer.

ended_at is the third, and it is decided in the same statement: a record carrying a stamp writes COALESCE(ended_at, stamp), so a row that already holds one keeps it, and a record carrying nil leaves the column alone. No read comes first, so two writers racing to end one execution cannot both stamp it.

write_tree_migration(opts, writes)

Writes a tree migration's re-pins and parks in one transaction (the optional StatifierPersistence.Storage.Adapter.write_tree_migration/2, ADR-0015 decision 3).

Every execution the writes name is read first, and one that is not stored is {:error, :execution_not_found} before any write is made: that refusal writes nothing, so it is returned as it is and nothing is rolled back, as retire_chart/3 returns its refusals. Reached inside a caller's own transaction - the per-execution lock's is one - the transaction joins it, so the refusal reaches the caller as the adapter's own reason and the enclosing transaction is left open.

Each write is then one update_all/3 keyed on execution_id. A write that still matches no row rolls the transaction back with rollback/1, so the writes before it are undone; inside a caller's own transaction that aborts the enclosing one too, which then answers its own rollback: a failure after a write must not return an error the enclosing transaction would then commit (the callback's contract).