StatifierPersistence.Retention (StatifierPersistence v0.21.0)

Copy Markdown View Source

Clearing what a finished execution leaves behind (ADR-0016).

An execution that has ended still stores its last position blob and, on an adapter that keeps one, its whole input log (ADR-0010) - the host's own event data, at rest, for the life of the store. prune/3 clears both for every execution that ended before a cutoff the host names, and keeps the execution row itself: its status, its failure, its metadata, its answer and its end stamp stay, so the drained query still counts it, a parent can still read a child's answer, and the execution id stays taken.

There is no clock here. The cutoff is a DateTime.t/0 the host computes from its own policy; nothing in this module takes a duration, defaults a window, or runs on a schedule (ADR-0012 decision 7). Which rows a host may delete on its own, and which it must not, is docs/retention.md.

Summary

Types

What prune/3 answers for single_batch: true: one batch's own counts plus whether another call may find more.

Options prune/3 accepts.

Functions

Clears the position blob and the input log of every execution that ended before cutoff, in batches, and answers what it cleared.

Types

batch_counts()

@type batch_counts() :: %{
  executions: non_neg_integer(),
  position_blobs: non_neg_integer(),
  inputs: non_neg_integer(),
  more?: boolean()
}

What prune/3 answers for single_batch: true: one batch's own counts plus whether another call may find more.

prune_opt()

@type prune_opt() ::
  {:batch_size, pos_integer()}
  | {:scope, StatifierPersistence.Storage.Adapter.prune_scope()}
  | {:single_batch, boolean()}

Options prune/3 accepts.

Functions

prune(store, cutoff, opts \\ [])

Clears the position blob and the input log of every execution that ended before cutoff, in batches, and answers what it cleared.

An execution is pruned when its status is :completed, :failed or :cancelled and its ended_at is set and strictly before cutoff. An execution with no end stamp is never touched, and neither is one whose row carries a stamp but was written back to a status that is not terminal. The execution row stays; see the moduledoc for what it keeps.

Answers {:ok, counts} summed over every batch: executions pruned, position_blobs nulled among them, and input log rows deleted. It is idempotent: a second call with the same cutoff answers zeros, because an execution with nothing left to clear is not selected again.

Each batch is one atomic unit in the adapter and commits on its own, unless a caller's transaction encloses the call (see below). So an {:error, reason} from a later batch leaves the earlier batches pruned, and calling again with the same cutoff carries on from where the failed batch stopped.

{:error, :execution_pruning_unsupported} for a store whose adapter does not declare the capability (StatifierPersistence.Storage.execution_pruning_supported?/1), before anything is read.

{:error, :unscoped_adapter} when :scope is given and the store's adapter cannot confine a batch to it - the in-memory adapter is one - before anything is cleared.

Raises ArgumentError when cutoff is not a DateTime.t/0 - a duration, a number of days or a Date included - when :batch_size is not a positive integer, and when :scope is not a non-empty keyword list of distinct columns with no nil value.

Options

  • :batch_size - at most how many executions one batch prunes. Defaults to 500. Smaller batches hold shorter transactions; the answer is the same.
  • :scope - a keyword list of column equalities, such as [tenant_id: "tenant-a"], that confines the prune to the rows holding every one of them. Every batch's selection, input log check, input log delete and position blob update carries the equalities, so a prune run inside one partition's transaction reads and writes no row of another. On StatifierPersistence.Storage.Ecto the columns are ones the host placed with :leading_columns. Left out, the prune covers the whole store, as it always has. nil is refused rather than read as IS NULL, because an equality with NULL matches no row.
  • :single_batch - false (the default) drains every batch as before. true runs exactly one batch and answers that batch's own counts plus more?: true when the batch took as many executions as batch_size: allows, so another call may find more; false when it took fewer, the same point where the default drain stops. A row another transaction holds locked (on Postgres) is skipped, not waited on, so it is left for a later call either way.

Inside a transaction of your own

Each batch is its own transaction only when nothing encloses it. On StatifierPersistence.Storage.Ecto a batch's transaction joins a caller's enclosing transaction, so prune/3 called inside one runs the whole drain - every batch - as one transaction that commits or rolls back with the caller's. :batch_size then bounds each batch's statements, not the transaction.

For one bounded transaction per batch - for example one that first sets a partition's context - call prune/3 with single_batch: true inside each of your transactions, and call again while more? is true:

def prune_tenant(store, cutoff, tenant_id) do
  batch =
    Repo.transaction(fn ->
      # set the tenant's context for this transaction here, then:
      opts = [scope: [tenant_id: tenant_id], single_batch: true]

      case Retention.prune(store, cutoff, opts) do
        {:ok, counts} -> counts
        {:error, reason} -> Repo.rollback(reason)
      end
    end)

  with {:ok, %{more?: true}} <- batch do
    prune_tenant(store, cutoff, tenant_id)
  end
end

It answers {:ok, counts} for the last batch, more?: false, or the first {:error, reason}, with every earlier batch committed.