ReactiveDag.Suspension (reactive_dag v0.17.0-rc.63)

Copy Markdown View Source

WHERE A CASCADE STOPPED — the library's only table, and the successor to the dirty queue.

A cascade runs to completion in one transaction, following the graph from whatever changed until it reaches something that cannot be done synchronously. Two things stop it: work too expensive to hold a transaction open, and work that needs a person. This table records those stopping points and nothing else.

The distinction from the queue it replaces is the whole design. A queue row said this cell needs recomputing — a conclusion, derived by walking the graph at mark time, which is stale by the time it is read. A suspension says this change reached here and could not continue — a fact, which stays true. Everything a resumption needs to know it derives when it runs.

The shape

suspension
  id          019a419f4          this row, never modified
  tenant      "red_hook_village"  whose graph
  waiting     "meeting_events"    the RESOURCE that stopped
  resource    "minutes_docs"      what changed
  row_uuid    019a3fc21            which row
  version_id  v-019a407b           and what moved
  reason      :expensive | :approval

waiting and resource are both resource names, and neither implies the other: one changed row may stop several resources, and they complete independently.

Immutable, and why that is the point

Rows are inserted and deleted, never updated. There is no unique constraint and no ON CONFLICT: a second change to the same stopping point writes a SECOND row.

What coalesces is the WORK, not the rows. A job reads every suspension at its point, does the expensive thing once, and discharges exactly the ids it read — so a suspension written during a nine-minute recompute is not in that list and survives for the next pass. An append-only table cannot lose a change that arrives mid-flight, and needs no revision counter that every future write path must remember to carry.

The cost is rows: a document changing fifty times before its resumption runs leaves fifty suspensions, all discharged together. That is a retention question, not a correctness one, and it is the trade taken deliberately — duplicate rows are cheap, a lost cascade is not.

A single mutable row per point, deleted only if unchanged since the job read it, would keep the count flat. It also makes correctness depend on every write path carrying a revision, and one that forgets loses a change with no error. Worth revisiting if volume becomes a real cost; not worth the invariant before then.

Configuration

config :reactive_dag, repo: MyApp.Repo, suspension_table: "my_suspensions"

Values are always parameterized. The table name — the one identifier SQL cannot parameterize — comes from config and is validated against an identifier grammar at read time, so a typo fails loudly rather than as a syntax error deep in a query.

Summary

Types

A stopping point: which tenant's graph, what stopped, and which row of what moved. This is what a resumption job carries — never a suspension id, never a version — so that a job queued at 12:00 and run at 12:05 acts on what is true at 12:05.

t()

Functions

Every suspension at a stopping point, oldest first.

Discharge suspensions BY ID, returning how many were removed.

Oban's job table, validated. Configurable only so the SQL above can be exercised against a real database without colliding with a host's own table.

Whether anything is suspended for this tenant.

The distinct stopping points with work outstanding, for one tenant.

Record that a cascade stopped, and return the new row's id.

Make stranded resumption jobs fetchable again, returning the job ids revived.

Run fun in a SAVEPOINT: if it returns {:error, reason}, everything it did is undone and {:error, reason} comes back — without disturbing the transaction around it.

Stopping points whose resumption job can never run.

The suspension table's name, validated.

The tenant a call is about: opts[:tenant], else "*" (untenanted).

Run fun in a transaction — the cascade's transaction.

Types

point()

@type point() :: %{
  tenant: String.t(),
  waiting: String.t(),
  resource: String.t(),
  row_uuid: String.t()
}

A stopping point: which tenant's graph, what stopped, and which row of what moved. This is what a resumption job carries — never a suspension id, never a version — so that a job queued at 12:00 and run at 12:05 acts on what is true at 12:05.

reason()

@type reason() :: :expensive | :approval

t()

@type t() :: %{
  id: String.t(),
  version_id: String.t(),
  reason: reason(),
  lap: non_neg_integer()
}

Functions

at(point)

@spec at(point()) :: [t()]

Every suspension at a stopping point, oldest first.

A resumption job's first act. An empty list is the ORDINARY outcome of a duplicate job — the work was already done and discharged — not an error.

Ordered by id, which is a UUIDv7 and therefore sorts by creation. That is not decorative: when several versions merge into one recompute, they must be applied in the order the changes happened.

discharge(ids)

@spec discharge([String.t()]) :: non_neg_integer()

Discharge suspensions BY ID, returning how many were removed.

By id, never by point — this is the line the append-only design rests on. A DELETE … WHERE tenant = … AND waiting = … would also remove suspensions written while the job was running, silently discarding changes nobody observed. Naming the ids the job actually read cannot do that.

Called in the same transaction as the resumption's writes, so a failed resumption leaves its suspensions in place and the work is retried.

oban_table()

@spec oban_table() :: String.t()

Oban's job table, validated. Configurable only so the SQL above can be exercised against a real database without colliding with a host's own table.

pending?(opts \\ [])

@spec pending?(keyword()) :: boolean()

Whether anything is suspended for this tenant.

points(opts \\ [])

@spec points(keyword()) :: [
  %{
    point: point(),
    reason: reason(),
    count: pos_integer(),
    oldest: DateTime.t()
  }
]

The distinct stopping points with work outstanding, for one tenant.

Non-consuming, and aggregated: one entry per point per reason, with how many suspensions have accumulated there and how long the oldest has waited. This is what an operator's view reads — "twelve things waiting on a person" — and what surfaces a point whose resumption keeps failing, since its count climbs while its oldest recedes.

record(point, version_id, reason, lap \\ 0)

@spec record(point(), String.t(), reason(), non_neg_integer()) :: String.t()

Record that a cascade stopped, and return the new row's id.

Called from INSIDE the cascade's transaction, so a rolled-back cascade leaves no suspension: the change it describes never happened.

version_id may be "*" when the change could not be attributed — a source that cannot say which of its items moved, or a write path that records no diff. Resumption then recomputes the whole cell: expensive, correct, and honest. This gives the whole-cell marker a principled meaning it has never had — not "something happened somewhere", but "this suspension could not be narrowed".

lap is how many consecutive resumption-driven trips around a declared feedback loop led to this suspension — 0 (the default, and the value for every suspension outside a loop) means the work arrived here from an external change. The column exists because a loop through a suspending cell ends every CASCADE cleanly: no in-memory counter survives the suspend → commit → resume chain, so the count must ride on the one thing that does — this row. ReactiveDag.Cascade refuses to record a lap past its budget, which is what bounds the chain.

revive(opts \\ [])

@spec revive(keyword()) :: [integer()]

Make stranded resumption jobs fetchable again, returning the job ids revived.

Raises max_attempts above attempt — the same expression Oban's own retry_all_jobs/2 uses — and clears the backoff so the queue picks the job up on its next poll. Idempotent: a job already fetchable does not match.

Safe to run blind. A resumption whose suspensions were discharged in the meantime finds nothing at step 1 and exits, which is its ordinary duplicate path.

savepoint(fun)

@spec savepoint((-> result)) :: result | {:error, term()} when result: term()

Run fun in a SAVEPOINT: if it returns {:error, reason}, everything it did is undone and {:error, reason} comes back — without disturbing the transaction around it.

This is what lets one fallible cell fail inside a cascade that is otherwise committing. The branch that failed stops; every other branch carries on and commits.

It must RETURN its failure, not raise

An exception inside a nested transaction aborts the OUTER one: Postgres marks the connection failed and every later statement errors until the whole thing rolls back. A savepoint only isolates a failure that arrives as a value.

That is why ReactiveDag.Source's poll/1 returns {:error, reason} "contained, not raised" — the contract was already the right shape for this.

A repo without transaction/2 runs fun directly: no savepoint, so a failure is not isolated.

stranded(opts \\ [])

@spec stranded(keyword()) :: [map()]

Stopping points whose resumption job can never run.

A job that exhausts its attempts is normally discarded, which is visible. But Oban returns a failed job to available with a backoff BEFORE checking attempts, so a final attempt that fails leaves attempt == max_attempts in state available — and the fetch query requires attempt < max_attempts (Oban.Engines.Basic). The job is then unfetchable and undiscarded: it will never run and nothing reports it as failed.

That would merely be a stalled point, except the resumption worker's uniqueness is states: :incomplete, which INCLUDES available. So the stranded job also dedups every future enqueue for its point. The suspension stays outstanding, no job will ever discharge it, and the queue looks healthy.

Observed in production: a resumption whose third attempt hit DBConnection.ConnectionError while the machine was being replaced.

Returns a point plus its job_id for each. Repair with revive/1. Note that Oban.retry_job/1 does NOT fix these — it skips jobs already in available (retry_all_jobs/2), so it reports success and changes nothing.

table()

@spec table() :: String.t()

The suspension table's name, validated.

tenant(opts)

@spec tenant(keyword()) :: String.t()

The tenant a call is about: opts[:tenant], else "*" (untenanted).

"*" rather than nil because a null never equals itself. Every read here matches on the tenant, and tenant = NULL matches nothing — an untenanted suspension would be written and then never found, which is the quiet failure of tenancy in its worst form: a resumption that finds no work reports SUCCESS. "*" is also already this library's spelling for "the whole thing", so the vocabulary is not new.

Normalised in ONE place so a caller passing nil, omitting the option, or passing an atom all mean the same row.

transaction(fun)

@spec transaction((-> result)) :: result when result: term()

Run fun in a transaction — the cascade's transaction.

Bounded, unlike the drain's it replaces. The drain wrapped its recompute in timeout: :infinity and justified it with "a recompute legitimately runs for minutes" — which is exactly the condition that made a nine-minute extraction hold a connection until the database killed it.

A cascade transaction contains only fast work by construction: anything slow declared itself so and became a suspension. An infinite timeout here would hide the one bug this design exists to remove — a cell declared cheap that is not. Configure with:

config :reactive_dag, cascade_timeout: 30_000

A repo without transaction/2 runs fun directly. The library only ever needed query!/2, and requiring a new capability would break such a host on upgrade for a guarantee it may not want — at the cost that its cascades are not atomic.