Current section
Files
Jump to
Current section
Files
lib/continuum/runtime/journal/postgres.ex
defmodule Continuum.Runtime.Journal.Postgres do
@moduledoc """
Durable journal adapter backed by Postgres via Ecto.
Implements the `Continuum.Runtime.Journal` behaviour. Every append
operation is transactional and CAS-guarded by the lease state on the
run row. Appends lock the run row, validate the lease token, assign a
sequence number, and insert the event in one transaction.
The replay loop and engine code are identical whether this adapter or
`InMemory` is in use — the only difference is durability and the
fencing-token enforcement on writes.
"""
@behaviour Continuum.Runtime.Journal
import Ecto.Query
require Logger
alias Continuum.Runtime.{Instance, JournalError, Snapshotter}
alias Continuum.Schema.{ActivityResult, ActivityTask, Event, Run, Signal, Snapshot, Timer}
alias Continuum.Telemetry
@impl true
def start_run(%Instance{} = instance, run_id, workflow, input, opts \\ []) do
with_repo(instance, fn -> start_run_with_repo(run_id, workflow, input, opts) end)
end
defp start_run_with_repo(run_id, workflow, input, opts) do
metadata = workflow_metadata(workflow)
changeset =
%Run{}
|> Ecto.Changeset.change(%{
id: run_id,
workflow: metadata.workflow,
version_hash: metadata.version_hash,
namespace: normalize_namespace(Keyword.get(opts, :namespace, "default")),
state: "running",
input: encode_term(input),
attributes: normalize_attributes(Keyword.get(opts, :attributes, %{})),
correlation_id: run_id,
trace_context: Keyword.get(opts, :trace_context)
})
case Keyword.get(opts, :lease) do
nil ->
case repo().insert(changeset) do
{:ok, _} -> :ok
{:error, changeset} -> {:error, changeset}
end
lease_opts ->
insert_run_with_lease(changeset, run_id, lease_opts)
end
end
# Insert the run row already leased, in one transaction, so no concurrent
# dispatcher can claim the fresh row before the starting engine acquires its
# lease (the fresh-start steal race). The row only becomes visible with the
# lease fully set.
defp insert_run_with_lease(changeset, run_id, lease_opts) do
owner = Keyword.fetch!(lease_opts, :owner)
ttl_seconds = Keyword.get(lease_opts, :ttl_seconds, 30)
result =
repo().transaction(fn ->
with {:ok, _run} <-
repo().insert(Ecto.Changeset.change(changeset, %{lease_owner: owner})),
{:ok, %{rows: [[token]]}} <-
repo().query(
"""
UPDATE continuum_runs
SET lease_token = nextval('continuum_lease_token_seq'),
lease_expires_at = now() + make_interval(secs => $2)
WHERE id = $1::text::uuid
RETURNING lease_token
""",
[run_id, ttl_seconds]
) do
token
else
{:error, reason} -> repo().rollback(reason)
end
end)
case result do
{:ok, token} ->
Telemetry.execute([:continuum, :lease, :acquired], %{}, %{
run_id: run_id,
owner: owner,
lease_token: token
})
{:ok, %Continuum.Runtime.Lease{run_id: run_id, owner: owner, token: token}}
{:error, reason} ->
{:error, reason}
end
end
defp workflow_metadata(workflow) do
case Continuum.VersionRegistry.ensure_registered(workflow) do
{:ok, metadata} ->
%{workflow: metadata.workflow_string, version_hash: metadata.version_hash}
{:error, reason} ->
raise ArgumentError,
"expected #{inspect(workflow)} to use Continuum.Workflow before starting a durable run, got: #{inspect(reason)}"
end
end
defp normalize_attributes(nil), do: %{}
defp normalize_attributes(attributes) when is_map(attributes) do
case Jason.encode(attributes) do
{:ok, json} ->
Jason.decode!(json)
{:error, reason} ->
raise ArgumentError,
"expected :attributes to be JSON-encodable map data, got: #{inspect(reason)}"
end
end
defp normalize_attributes(other) do
raise ArgumentError, "expected :attributes to be a map, got: #{inspect(other)}"
end
defp normalize_namespace(nil), do: "default"
defp normalize_namespace(namespace) when is_binary(namespace) and byte_size(namespace) > 0,
do: namespace
defp normalize_namespace(other) do
raise ArgumentError, "expected :namespace to be a non-empty binary, got: #{inspect(other)}"
end
@impl true
def append!(%Instance{} = instance, run_id, event, lease_token) do
:ok = with_repo(instance, fn -> append_with_repo!(run_id, event, lease_token) end)
maybe_snapshot_after_event(instance, run_id, event, lease_token)
end
defp append_with_repo!(run_id, event, lease_token) do
{event_type, payload} = encode_event(event)
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
seq = event[:seq] || next_seq(run_id)
changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: seq,
event_type: event_type,
payload: payload,
inserted_at: DateTime.utc_now()
})
case repo().insert(changeset) do
{:ok, _} -> :ok
{:error, changeset} -> repo().rollback({:insert_failed, changeset})
end
end)
case result do
{:ok, :ok} ->
:ok
{:error, reason} ->
raise JournalError, op: :append!, reason: reason
end
end
@impl true
def load(%Instance{} = instance, run_id) do
with_repo(instance, fn -> load_with_repo(run_id) end)
end
defp load_with_repo(run_id) do
events =
repo().all(
from(e in Event,
where: e.run_id == ^run_id,
order_by: [asc: e.seq]
)
)
Enum.map(events, &decode_event/1)
end
@impl true
def load_with_snapshot(%Instance{} = instance, run_id, lease_token) do
with_repo(instance, fn -> load_with_snapshot_with_repo(run_id, lease_token) end)
end
defp load_with_snapshot_with_repo(run_id, lease_token) do
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
snapshot = latest_snapshot(run_id)
through_seq = if snapshot, do: snapshot.through_seq, else: -1
{snapshot, load_events_after(run_id, through_seq)}
end)
case result do
{:ok, value} ->
value
{:error, reason} ->
raise JournalError, op: :load_with_snapshot, reason: reason
end
end
@impl true
def take_snapshot!(%Instance{} = instance, %Continuum.Snapshot{} = snapshot) do
with_repo(instance, fn -> take_snapshot_with_repo!(snapshot) end)
end
defp take_snapshot_with_repo!(%Continuum.Snapshot{} = snapshot) do
repo().insert_all(
Snapshot,
[
%{
run_id: snapshot.run_id,
through_seq: snapshot.through_seq,
version_hash: snapshot.version_hash,
format_version: Continuum.Snapshot.format_version(),
payload: Continuum.Snapshot.encode(snapshot),
taken_at: snapshot.taken_at
}
],
on_conflict: :nothing,
conflict_target: [:run_id, :through_seq]
)
:ok
end
def schedule_activity!(%Instance{} = instance, run_id, event, task, lease_token) do
with_repo(instance, fn -> schedule_activity_with_repo!(run_id, event, task, lease_token) end)
end
defp schedule_activity_with_repo!(run_id, event, task, lease_token) do
{event_type, payload} = encode_event(event)
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
now = DateTime.utc_now()
event_changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: event.seq,
event_type: event_type,
payload: payload,
inserted_at: now
})
task_changeset =
%ActivityTask{}
|> Ecto.Changeset.change(%{
id: task.id,
run_id: run_id,
seq: task.seq,
mfa: encode_term(task),
attempt: 1,
state: "available"
})
with {:ok, _event} <- repo().insert(event_changeset),
{:ok, _task} <- repo().insert(task_changeset) do
:ok
else
{:error, changeset} -> repo().rollback({:activity_schedule_failed, changeset})
end
end)
case result do
{:ok, :ok} ->
:ok
{:error, reason} ->
raise JournalError, op: :schedule_activity!, reason: reason
end
end
@doc """
Start a child workflow run and journal `child_started` to the parent.
In one transaction (CAS-guarded by the parent's lease): insert the child run
row with `parent_run_id`/`parent_command_id`/`correlation_id` set and append
the `child_started` event to the parent's history. The child run is left
runnable for the dispatcher to claim.
"""
def start_child!(%Instance{} = instance, parent_run_id, child, lease_token) do
with_repo(instance, fn -> start_child_with_repo!(parent_run_id, child, lease_token) end)
end
defp start_child_with_repo!(parent_run_id, child, lease_token) do
metadata = workflow_metadata(child.workflow)
{event_type, payload} = encode_event(child.started_event)
now = DateTime.utc_now()
result =
repo().transaction(fn ->
parent = lock_and_load_run!(parent_run_id, lease_token)
enforce_child_depth!(parent_run_id)
child_changeset =
%Run{}
|> Ecto.Changeset.change(%{
id: child.child_run_id,
workflow: metadata.workflow,
version_hash: metadata.version_hash,
# Children inherit the parent's tenant scoping; without this a
# namespaced parent's children silently land in "default".
namespace: parent.namespace,
attributes: parent.attributes,
state: "running",
input: encode_term(child.input),
parent_run_id: parent_run_id,
parent_command_id: child.parent_command_id,
correlation_id: child.child_run_id,
trace_context: child.trace_context
})
event_changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: parent_run_id,
seq: child.started_event.seq,
event_type: event_type,
payload: payload,
inserted_at: now
})
with {:ok, _child} <- repo().insert(child_changeset),
{:ok, _event} <- repo().insert(event_changeset) do
:ok
else
{:error, changeset} -> repo().rollback({:start_child_failed, changeset})
end
end)
case result do
{:ok, :ok} ->
:ok
{:error, reason} ->
raise JournalError, op: :start_child!, reason: reason
end
end
# The cancel cascade walks descendants down to :max_child_depth. Creation
# must respect the same bound, or descendants deeper than the cascade can
# never be reached by Continuum.cancel — fail loudly at start_child time
# instead of silently outrunning the cascade.
defp enforce_child_depth!(parent_run_id) do
max_depth = max_child_depth()
sql = """
WITH RECURSIVE ancestors AS (
SELECT id, parent_run_id, 0 AS depth FROM continuum_runs WHERE id = $1
UNION ALL
SELECT p.id, p.parent_run_id, a.depth + 1
FROM continuum_runs p
JOIN ancestors a ON p.id = a.parent_run_id
WHERE a.depth < $2
)
SELECT max(depth) FROM ancestors
"""
case repo().query(sql, [Ecto.UUID.dump!(parent_run_id), max_depth]) do
{:ok, %{rows: [[parent_depth]]}} when is_integer(parent_depth) ->
if parent_depth + 1 > max_depth do
repo().rollback(
{:max_child_depth_exceeded, depth: parent_depth + 1, max_child_depth: max_depth}
)
else
:ok
end
_other ->
:ok
end
end
defp max_child_depth do
Application.get_env(:continuum, :max_child_depth, 10)
end
@doc """
Resolve a child's terminal state into the parent's history.
Locks the parent (CAS by lease). If the child run is terminal, appends the
matching `child_completed`/`child_failed`/`child_cancelled` event to the
parent and returns the decoded outcome; otherwise returns `:pending`.
"""
def await_child_terminal!(
%Instance{} = instance,
parent_run_id,
child_run_id,
command_id,
seq,
lease_token
) do
with_repo(instance, fn ->
await_child_terminal_with_repo!(parent_run_id, child_run_id, command_id, seq, lease_token)
end)
end
defp await_child_terminal_with_repo!(parent_run_id, child_run_id, command_id, seq, lease_token) do
result =
repo().transaction(fn ->
lock_and_validate_run!(parent_run_id, lease_token)
case child_terminal_state(child_run_id) do
{:completed, child_result} ->
event = %{
type: :child_completed,
child_run_id: child_run_id,
result: child_result,
command_id: command_id,
seq: seq
}
{:completed, child_result, insert_event!(parent_run_id, event)}
{:failed, error} ->
event = %{
type: :child_failed,
child_run_id: child_run_id,
error: error,
command_id: command_id,
seq: seq
}
{:failed, error, insert_event!(parent_run_id, event)}
{:cancelled} ->
event = %{
type: :child_cancelled,
child_run_id: child_run_id,
command_id: command_id,
seq: seq
}
{:cancelled, insert_event!(parent_run_id, event)}
:pending ->
:pending
end
end)
case result do
{:ok, value} ->
value
{:error, reason} ->
raise JournalError, op: :await_child_terminal!, reason: reason
end
end
defp child_terminal_state(child_run_id) do
# Follow a `continue_as_new` chain forward to its terminal run so a parent
# never sees an intermediate `{:continued, _}` marker as the child result.
terminal_id = follow_continued_chain(child_run_id)
case repo().one(
from(r in Run, where: r.id == ^terminal_id, select: {r.state, r.result, r.error})
) do
nil ->
:pending
{"completed", result, _error} ->
case decode_term(result) do
{:continued, _next_run_id} -> :pending
decoded -> {:completed, decoded}
end
{"cancelled", _result, _error} ->
{:cancelled}
# Classified by run *state*: a child that legitimately failed with the
# user error term :cancelled is a failure, not a cancellation.
{"failed", _result, error} ->
{:failed, decode_term(error)}
{_state, _result, _error} ->
:pending
end
end
@doc """
Resolve a run id to the live tip of its `continue_as_new` chain.
External callers hold the chain-root id; a run with no successor resolves
to itself. Used by signal delivery, cancel, and await so operations on a
continued run reach the current incarnation instead of the dead root.
"""
def resolve_chain_tip(%Instance{} = instance, run_id) do
with_repo(instance, fn -> follow_continued_chain(run_id) end)
end
defp follow_continued_chain(run_id) do
sql = """
WITH RECURSIVE chain AS (
SELECT id, 0 AS depth FROM continuum_runs WHERE id = $1
UNION ALL
SELECT c.id, ch.depth + 1
FROM continuum_runs c
JOIN chain ch ON c.continued_from_run_id = ch.id
)
SELECT id::text FROM chain ORDER BY depth DESC LIMIT 1
"""
case repo().query(sql, [Ecto.UUID.dump!(run_id)]) do
{:ok, %{rows: [[terminal_id]]}} -> terminal_id
_ -> run_id
end
end
@doc """
Complete the current run as `{:continued, next_run_id}` and insert the fresh
continuation run, in one lease-CAS-guarded transaction.
The new run carries `continued_from_run_id`, the chain's `correlation_id`
(the chain root's id), and any `parent_run_id`/`parent_command_id` so a
continued child stays a child.
"""
def continue_as_new!(
%Instance{} = instance,
run_id,
next_run_id,
next_input,
event,
lease_token
) do
with_repo(instance, fn ->
continue_as_new_with_repo!(run_id, next_run_id, next_input, event, lease_token)
end)
end
defp continue_as_new_with_repo!(run_id, next_run_id, next_input, event, lease_token) do
{event_type, payload} = encode_event(event)
now = DateTime.utc_now()
result =
repo().transaction(fn ->
# Lock our parent before our own row, the same parent-before-child
# order the cancel cascade takes (see lock_parent_first/1). A
# concurrent Continuum.cancel of an ancestor then serializes against
# us on the shared parent lock: it either clears our lease before we
# validate it (we roll back here) or observes the successor as a
# freshly-locked descendant of the parent. Without this the successor
# is parented to a row the cascade locked under only the root's lock
# and escapes cancellation (audit F3).
lock_parent_first(run_id)
run = repo().one(from(r in Run, where: r.id == ^run_id, lock: "FOR UPDATE"))
case run do
nil ->
repo().rollback({:run_not_found, run_id})
%Run{state: state} when state not in ["running", "suspended"] ->
repo().rollback({:run_not_active, state})
%Run{} = run ->
:ok = validate_lease!(run, lease_token)
correlation = run.correlation_id || run.id
{next_workflow, next_version_hash} = successor_workflow_version(run)
event_changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: event.seq,
event_type: event_type,
payload: payload,
inserted_at: now
})
next_changeset =
%Run{}
|> Ecto.Changeset.change(%{
id: next_run_id,
workflow: next_workflow,
version_hash: next_version_hash,
# The successor is the same logical process: keep its tenant
# scoping instead of resetting to the schema defaults.
namespace: run.namespace,
attributes: run.attributes,
state: "running",
input: encode_term(next_input),
correlation_id: correlation,
continued_from_run_id: run_id,
parent_run_id: run.parent_run_id,
parent_command_id: run.parent_command_id,
trace_context: run.trace_context,
# A durable cancel request addresses the logical chain, not
# one iteration: continuing must not let an acknowledged
# cancel evaporate before the owner's next heartbeat sees it.
cancel_requested_at: run.cancel_requested_at
})
with {:ok, _event} <- repo().insert(event_changeset),
{:ok, _next} <- repo().insert(next_changeset),
:ok <-
cas_update_run(run_id, lease_token, %{
state: "completed",
result: encode_term({:continued, next_run_id}),
correlation_id: correlation,
completed_at: now
}) do
# Re-parent unawaited live children to the successor: their
# child_started events live in the dead run's history, so the
# successor cannot await them, but cancelling the chain must
# still cascade into them.
repo().update_all(
from(r in Run,
where: r.parent_run_id == ^run_id and r.state in ["running", "suspended"]
),
set: [parent_run_id: next_run_id]
)
# Undelivered signals move to the successor's mailbox — external
# signalers address the logical chain, not this incarnation, and
# the predecessor can never consume them once terminal.
repo().update_all(
from(s in Signal, where: s.run_id == ^run_id and s.delivered == false),
set: [run_id: next_run_id]
)
correlation
else
{:error, changeset} -> repo().rollback({:continue_as_new_failed, changeset})
end
end
end)
case result do
{:ok, correlation} ->
correlation
{:error, reason} ->
raise JournalError, op: :continue_as_new!, reason: reason
end
end
# The successor starts with empty history, so any loaded version is
# replay-safe. Stamp it with the workflow's currently loaded version instead
# of pinning the predecessor's hash — a pinned chain never picks up deploys
# and goes unknown-version once the old module is gone.
defp successor_workflow_version(run) do
case Continuum.VersionRegistry.ensure_registered(Module.concat([run.workflow])) do
{:ok, metadata} -> {metadata.workflow_string, metadata.version_hash}
{:error, _reason} -> {run.workflow, run.version_hash}
end
end
@doc """
Schedule a compensation activity task.
Reuses the activity-task append path: the `compensation_scheduled` event and
the worker task are inserted under the run lease in one transaction. The task
carries `kind: :compensation` and `target_activity_id` so the worker journals
`compensation_completed`/`compensation_failed` on completion.
"""
def schedule_compensation!(%Instance{} = instance, run_id, event, task, lease_token) do
with_repo(instance, fn -> schedule_activity_with_repo!(run_id, event, task, lease_token) end)
end
def schedule_compensations!(%Instance{} = instance, run_id, scheduled, lease_token) do
with_repo(instance, fn ->
schedule_compensations_with_repo!(run_id, scheduled, lease_token)
end)
end
defp schedule_compensations_with_repo!(run_id, scheduled, lease_token) do
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
now = DateTime.utc_now()
Enum.each(scheduled, fn %{event: event, task: task} ->
{event_type, payload} = encode_event(event)
event_changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: event.seq,
event_type: event_type,
payload: payload,
inserted_at: now
})
task_changeset =
%ActivityTask{}
|> Ecto.Changeset.change(%{
id: task.id,
run_id: run_id,
seq: task.seq,
mfa: encode_term(task),
attempt: 1,
state: "available"
})
with {:ok, _event} <- repo().insert(event_changeset),
{:ok, _task} <- repo().insert(task_changeset) do
:ok
else
{:error, changeset} ->
repo().rollback({:compensation_batch_schedule_failed, changeset})
end
end)
:ok
end)
case result do
{:ok, :ok} ->
:ok
{:error, reason} ->
raise JournalError, op: :schedule_compensations!, reason: reason
end
end
def complete_compensation_task!(%Instance{} = instance, task, result, lease_token, opts \\ []) do
with_repo(instance, fn ->
complete_compensation_task_with_repo!(task, result, lease_token, opts)
end)
end
defp complete_compensation_task_with_repo!(task, result, lease_token, opts) do
idempotency = Keyword.get(opts, :idempotency)
tx_result =
repo().transaction(fn ->
lock_and_validate_active_run!(task.run_id, lease_token)
lock_and_validate_activity_task!(task)
committed_result = maybe_commit_idempotency_result(task, result, idempotency)
event = compensation_completed_event(task, committed_result)
with %{} <- insert_event!(task.run_id, event),
{1, _} <-
repo().update_all(
from(t in ActivityTask,
where:
t.id == ^task.id and t.run_id == ^task.run_id and t.state == "leased" and
t.lease_owner == ^task.lease_owner and t.attempt == ^task.attempt
),
set: [state: "completed", result: encode_term(committed_result)]
) do
:ok
else
{0, _} -> repo().rollback({:compensation_task_result_failed, :task_lease_mismatch})
end
end)
case tx_result do
{:ok, :ok} ->
Snapshotter.maybe_snapshot(task.instance, task.run_id, lease_token, __MODULE__)
:ok
{:error, reason} ->
raise JournalError, op: :complete_compensation_task!, reason: reason
end
end
defp compensation_completed_event(task, result) do
%{
type: :compensation_completed,
target_activity_id: task.target_activity_id,
result: result,
command_id: Map.get(task, :command_id),
seq: compensation_terminal_seq(task)
}
end
def fail_compensation_task!(%Instance{} = instance, task, error, lease_token) do
with_repo(instance, fn ->
event = %{
type: :compensation_failed,
target_activity_id: task.target_activity_id,
error: error,
attempt: task.attempt,
command_id: Map.get(task, :command_id),
seq: compensation_terminal_seq(task)
}
activity_task_result!(
task,
event,
[state: "discarded", error: encode_term(error)],
lease_token
)
end)
end
def get_activity_result(%Instance{} = instance, activity_module, idempotency_key) do
with_repo(instance, fn -> get_activity_result_with_repo(activity_module, idempotency_key) end)
end
defp get_activity_result_with_repo(activity_module, idempotency_key) do
activity_module = activity_module_key(activity_module)
case repo().one(
from(r in ActivityResult,
where:
r.activity_module == ^activity_module and r.idempotency_key == ^idempotency_key
)
) do
nil -> :miss
%ActivityResult{} = result -> {:ok, decode_term(result.result)}
end
end
def complete_activity_task!(%Instance{} = instance, task, result, lease_token, opts \\ []) do
with_repo(instance, fn ->
complete_activity_task_with_repo!(task, result, lease_token, opts)
end)
end
defp complete_activity_task_with_repo!(task, result, lease_token, opts) do
idempotency = Keyword.get(opts, :idempotency)
tx_result =
repo().transaction(fn ->
lock_and_validate_active_run!(task.run_id, lease_token)
lock_and_validate_activity_task!(task)
committed_result = maybe_commit_idempotency_result(task, result, idempotency)
event = activity_completed_event(task, committed_result)
with %{} <- insert_event!(task.run_id, event),
{1, _} <-
repo().update_all(
from(t in ActivityTask,
where:
t.id == ^task.id and t.run_id == ^task.run_id and t.state == "leased" and
t.lease_owner == ^task.lease_owner
),
set: [state: "completed", result: encode_term(committed_result)]
) do
:ok
else
{0, _} -> repo().rollback({:activity_task_result_failed, :task_lease_mismatch})
end
end)
case tx_result do
{:ok, :ok} ->
Snapshotter.maybe_snapshot(task.instance, task.run_id, lease_token, __MODULE__)
:ok
{:error, reason} ->
raise JournalError, op: :activity_task_result!, reason: reason
end
end
defp activity_completed_event(task, result) do
%{
type: :activity_completed,
mfa: task.mfa,
payload: result,
command_id: Map.get(task, :command_id),
seq: task.seq + 1
}
end
def fail_activity_task!(%Instance{} = instance, task, error, lease_token) do
with_repo(instance, fn -> fail_activity_task_with_repo!(task, error, lease_token) end)
end
defp fail_activity_task_with_repo!(task, error, lease_token) do
event = %{
type: :activity_failed,
mfa: task.mfa,
error: error,
attempt: task.attempt,
command_id: Map.get(task, :command_id),
seq: task.seq + 1
}
activity_task_result!(
task,
event,
[state: "discarded", error: encode_term(error)],
lease_token
)
end
def retry_activity_task!(%Instance{} = instance, task, error, backoff_ms, lease_token) do
with_repo(instance, fn ->
retry_activity_task_with_repo!(task, error, backoff_ms, lease_token)
end)
end
defp retry_activity_task_with_repo!(task, error, backoff_ms, lease_token) do
backoff_seconds = backoff_ms / 1_000
next_attempt = task.attempt + 1
encoded_error = encode_term(error)
result =
repo().transaction(fn ->
lock_and_validate_active_run!(task.run_id, lease_token)
lock_and_validate_activity_task!(task)
# available_at is computed on the database clock so the claim
# comparison (also DB time) measures the intended backoff regardless
# of app/DB clock skew.
query =
from(t in ActivityTask,
where:
t.id == ^task.id and t.run_id == ^task.run_id and t.state == "leased" and
t.lease_owner == ^task.lease_owner and t.attempt == ^task.attempt,
update: [
set: [
state: "available",
attempt: ^next_attempt,
available_at:
fragment("clock_timestamp() + make_interval(secs => ?)", ^backoff_seconds),
lease_owner: nil,
lease_expires_at: nil,
error: ^encoded_error
]
]
)
case repo().update_all(query, []) do
{1, _} -> :ok
{0, _} -> repo().rollback({:activity_task_retry_failed, :task_lease_mismatch})
end
end)
case result do
{:ok, :ok} ->
:ok
{:error, reason} ->
raise JournalError, op: :retry_activity_task!, reason: reason
end
end
def cancel_run!(%Instance{} = instance, run_id, lease_token) do
with_repo(instance, fn -> cancel_run_with_instance!(run_id, lease_token, instance) end)
end
defp cancel_run_with_instance!(run_id, lease_token, instance) do
result =
repo().transaction(fn ->
# Same lock order as the cascade (parent before child): take our own
# parent's row lock first so cancel and child-completion cannot
# deadlock AB-BA.
lock_parent_first(run_id)
lock_and_validate_active_run!(run_id, lease_token)
repo().update_all(
from(t in ActivityTask,
where: t.run_id == ^run_id and t.state in ["available", "leased"]
),
set: [
state: "discarded",
lease_owner: nil,
lease_expires_at: nil,
error: encode_term(:cancelled)
]
)
repo().update_all(
from(t in Timer, where: t.run_id == ^run_id and t.fired == false),
set: [fired: true]
)
# Cascade: cancel all in-flight descendant child runs, bounded by depth.
cancelled_descendants = cancel_descendants!(run_id)
case repo().update_all(
leased_run_query(run_id, lease_token),
set: [
state: "cancelled",
error: encode_term(:cancelled),
completed_at: DateTime.utc_now(),
next_wakeup_at: nil,
lease_owner: nil,
lease_token: nil,
lease_expires_at: nil
]
) do
{1, _} -> {maybe_wake_parent(run_id), cancelled_descendants}
{0, _} -> repo().rollback({:cancel_failed, :lease_mismatch})
end
end)
case result do
{:ok, {parent_run_id, cancelled_descendants}} ->
# Single broadcaster, single canonical state. Cascade-cancelled
# descendants broadcast too, so their awaiters don't block for the
# full timeout.
Continuum.Runtime.Engine.broadcast_run_finished(instance, run_id, :cancelled, :cancelled)
Enum.each(cancelled_descendants, fn descendant_id ->
Continuum.Runtime.Engine.broadcast_run_finished(
instance,
descendant_id,
:cancelled,
:cancelled
)
end)
wake_parent(instance, parent_run_id)
:ok
{:error, reason} ->
raise JournalError, op: :cancel_run!, reason: reason
end
end
# Returns the ids of the descendants that were actually flipped, so the
# caller can broadcast their termination after commit.
defp cancel_descendants!(run_id) do
case descendant_run_ids(run_id) do
[] ->
[]
descendant_ids ->
cancelled = encode_term(:parent_cancelled)
now = DateTime.utc_now()
repo().update_all(
from(t in ActivityTask,
where: t.run_id in ^descendant_ids and t.state in ["available", "leased"]
),
set: [state: "discarded", lease_owner: nil, lease_expires_at: nil, error: cancelled]
)
repo().update_all(
from(t in Timer, where: t.run_id in ^descendant_ids and t.fired == false),
set: [fired: true]
)
# Clear the lease so any live descendant engine fails its next write and
# stops cleanly — no post-cancel child events can be appended.
{_count, flipped} =
repo().update_all(
from(r in Run,
where: r.id in ^descendant_ids and r.state in ["running", "suspended"],
select: r.id
),
set: [
state: "cancelled",
error: cancelled,
completed_at: now,
next_wakeup_at: nil,
lease_owner: nil,
lease_token: nil,
lease_expires_at: nil
]
)
flipped || []
end
end
# Cancel's cascade locks parent rows before child rows; completion paths
# must take the same order (see complete!/fail!) or the two transactions
# deadlock AB-BA. Locking a missing/parentless run is a no-op.
defp lock_parent_first(run_id) do
case repo().one(from(r in Run, where: r.id == ^run_id, select: r.parent_run_id)) do
nil ->
:ok
parent_run_id ->
repo().one(
from(r in Run, where: r.id == ^parent_run_id, lock: "FOR UPDATE", select: r.id)
)
:ok
end
end
# Collect the run ids the cascade will cancel, locking each generation with
# FOR UPDATE before reading the next. The previous implementation snapshotted
# the whole subtree with an unlocked recursive CTE under only the root's
# lock; a concurrent continue_as_new!/start_child! could then commit a new
# descendant after the snapshot and before the final update_all, and that run
# escaped cancellation entirely (audit F3). Locking generation-by-generation
# forces any in-flight creator to serialize on its parent row: it either
# blocks until we cancel its parent (then fails its own lease/state
# validation) or has already committed and is seen by the next generation's
# locking read. Lock order is strictly parent-before-child, matching the rest
# of the cancel/completion paths, so no AB-BA deadlock is introduced.
defp descendant_run_ids(run_id) do
collect_locked_descendants(run_id, [run_id], 1, max_child_depth(), [])
end
defp collect_locked_descendants(root_run_id, parent_ids, depth, max_depth, acc) do
case lock_child_run_ids(parent_ids) do
[] ->
acc
# One generation past the bound (legacy data or a lowered
# :max_child_depth): these rows are NOT cancelled. Surface the
# truncation loudly instead of a silent partial cancel.
children when depth > max_depth ->
Logger.error(
"Continuum cancel cascade for run #{root_run_id} truncated at depth " <>
"#{max_depth}: #{length(children)}+ deeper descendant(s) keep running"
)
Telemetry.execute(
[:continuum, :run, :cancel_cascade_truncated],
%{count: length(children)},
%{run_id: root_run_id, max_child_depth: max_depth}
)
acc
children ->
collect_locked_descendants(
root_run_id,
children,
depth + 1,
max_depth,
acc ++ children
)
end
end
defp lock_child_run_ids([]), do: []
defp lock_child_run_ids(parent_ids) do
repo().all(
from(r in Run,
where: r.parent_run_id in ^parent_ids,
order_by: r.id,
lock: "FOR UPDATE",
select: r.id
)
)
end
defp maybe_wake_parent(run_id) do
case repo().one(from(r in Run, where: r.id == ^run_id, select: r.parent_run_id)) do
nil ->
nil
parent_run_id ->
repo().update_all(
from(r in Run, where: r.id == ^parent_run_id),
set: [next_wakeup_at: DateTime.utc_now() |> DateTime.truncate(:microsecond)]
)
repo().query("SELECT pg_notify('continuum_run_wake', $1)", [parent_run_id])
parent_run_id
end
end
defp wake_parent(_instance, nil), do: :ok
defp wake_parent(instance, parent_run_id),
do: Continuum.Runtime.Engine.wake(instance, parent_run_id)
defp run_in_transaction!(fun) do
case repo().transaction(fun) do
{:ok, value} ->
value
{:error, reason} ->
raise JournalError, op: :transaction, reason: reason
end
end
def schedule_timer!(%Instance{} = instance, run_id, event, timer, lease_token) do
with_repo(instance, fn -> schedule_timer_with_repo!(run_id, event, timer, lease_token) end)
end
defp schedule_timer_with_repo!(run_id, event, timer, lease_token) do
{event_type, payload} = encode_event(event)
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
event_changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: event.seq,
event_type: event_type,
payload: payload,
inserted_at: DateTime.utc_now()
})
timer_changeset =
%Timer{}
|> Ecto.Changeset.change(%{
id: timer.id,
run_id: run_id,
fires_at: timer.fires_at,
fired: false
})
with {:ok, _event} <- repo().insert(event_changeset),
{:ok, _timer} <- repo().insert(timer_changeset),
:ok <- notify_timer_armed_with_repo(run_id, timer.fires_at),
{1, _} <-
repo().update_all(
leased_run_query(run_id, lease_token),
set: [next_wakeup_at: timer.fires_at]
) do
:ok
else
{0, _} -> repo().rollback({:timer_schedule_failed, :lease_mismatch})
{:error, changeset} -> repo().rollback({:timer_schedule_failed, changeset})
end
end)
case result do
{:ok, :ok} ->
:ok
{:error, reason} ->
raise JournalError, op: :schedule_timer!, reason: reason
end
end
@doc false
def notify_timer_armed!(%Instance{} = instance, run_id, fires_at) do
with_repo(instance, fn -> notify_timer_armed_with_repo(run_id, fires_at) end)
end
def schedule_signal_await!(%Instance{} = instance, run_id, event, lease_token) do
with_repo(instance, fn -> schedule_signal_await_with_repo!(run_id, event, lease_token) end)
end
defp schedule_signal_await_with_repo!(run_id, event, lease_token) do
{event_type, payload} = encode_event(event)
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
now = DateTime.utc_now()
changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: event.seq,
event_type: event_type,
payload: payload,
inserted_at: now
})
with {:ok, _event} <- repo().insert(changeset),
:ok <- maybe_insert_signal_timeout_timer(run_id, event),
:ok <- maybe_set_signal_timeout_wakeup(run_id, event, lease_token) do
:ok
else
{:error, changeset} -> repo().rollback({:signal_await_failed, changeset})
{0, _} -> repo().rollback({:signal_await_failed, :lease_mismatch})
end
end)
case result do
{:ok, :ok} ->
:ok
{:error, reason} ->
raise JournalError, op: :schedule_signal_await!, reason: reason
end
end
def resolve_signal_await(%Instance{} = instance, run_id, await_event, lease_token) do
value =
with_repo(instance, fn ->
resolve_signal_await_with_repo(run_id, await_event, lease_token)
end)
maybe_snapshot_after_signal_resolution(instance, run_id, value, lease_token)
value
end
@doc false
def consume_pending_signal!(%Instance{} = instance, run_id, name, command_id, seq, lease_token) do
value =
with_repo(instance, fn ->
consume_pending_signal_with_repo!(run_id, name, command_id, seq, lease_token)
end)
maybe_snapshot_after_signal_resolution(instance, run_id, value, lease_token)
value
end
defp consume_pending_signal_with_repo!(run_id, name, command_id, seq, lease_token) do
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
case pending_signal(run_id, name) do
nil -> :none
%Signal{} = signal -> consume_signal_row!(run_id, name, signal, command_id, seq)
end
end)
case result do
{:ok, value} ->
value
{:error, reason} ->
raise JournalError, op: :consume_pending_signal!, reason: reason
end
end
defp resolve_signal_await_with_repo(run_id, await_event, lease_token) do
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
case signal_await_winner(run_id, await_event) do
:none -> consume_signal_or_timeout(run_id, await_event)
result -> result
end
end)
case result do
{:ok, value} ->
value
{:error, reason} ->
raise JournalError, op: :resolve_signal_await, reason: reason
end
end
def deliver_signal!(%Instance{} = instance, run_id, name, payload) do
with_repo(instance, fn -> deliver_signal_with_repo!(run_id, name, payload) end)
end
defp deliver_signal_with_repo!(run_id, name, payload) do
signal_name = Atom.to_string(name)
now = DateTime.utc_now()
result =
repo().transaction(fn ->
# Lock while walking the chain. If a concurrent continue_as_new wins,
# this waits, sees the completed predecessor, and re-locks the
# successor; if delivery wins, the later continue migrates the signal.
case lock_signal_delivery_tip(run_id) do
%Run{state: state} when state not in ["running", "suspended"] ->
repo().rollback(:run_terminal)
%Run{id: delivered_run_id} ->
changeset =
%Signal{}
|> Ecto.Changeset.change(%{
run_id: delivered_run_id,
name: signal_name,
payload: encode_term(payload),
delivered: false,
inserted_at: now
})
with {:ok, _signal} <- repo().insert(changeset),
{_count, _} <-
repo().update_all(
from(r in Run, where: r.id == ^delivered_run_id),
set: [next_wakeup_at: now]
),
{:ok, _} <-
repo().query("SELECT pg_notify('continuum_signal', $1)", [delivered_run_id]) do
delivered_run_id
else
{:error, reason} -> repo().rollback({:signal_delivery_failed, reason})
end
end
end)
case result do
{:ok, delivered_run_id} ->
Telemetry.execute([:continuum, :signal, :delivered], %{}, %{
run_id: delivered_run_id,
signal_name: name,
durable?: true
})
{:ok, delivered_run_id}
{:error, :not_found} ->
{:error, :not_found}
{:error, :run_terminal} ->
{:error, :run_terminal}
{:error, reason} ->
raise JournalError, op: :deliver_signal!, reason: reason
end
end
defp lock_signal_delivery_tip(run_id) do
case repo().one(from(r in Run, where: r.id == ^run_id, lock: "FOR UPDATE")) do
nil ->
repo().rollback(:not_found)
%Run{state: "completed", result: result} = run ->
case decode_term(result) do
{:continued, next_run_id} -> lock_signal_delivery_tip(next_run_id)
_other -> run
end
%Run{} = run ->
run
end
end
def consume_signal(%Instance{} = instance, run_id, name, lease_token) do
value = with_repo(instance, fn -> consume_signal_with_repo(run_id, name, lease_token) end)
maybe_snapshot_after_signal_resolution(instance, run_id, value, lease_token)
value
end
defp consume_signal_with_repo(run_id, name, lease_token) do
signal_name = Atom.to_string(name)
result =
repo().transaction(fn ->
lock_and_validate_run!(run_id, lease_token)
signal =
repo().one(
from(s in Signal,
where: s.run_id == ^run_id and s.name == ^signal_name and s.delivered == false,
order_by: [asc: s.inserted_at, asc: s.id],
limit: 1,
lock: "FOR UPDATE SKIP LOCKED"
)
)
case signal do
nil ->
:none
%Signal{} = signal ->
payload = decode_term(signal.payload)
event = %{
type: :signal_received,
name: name,
payload: payload,
seq: next_seq(run_id)
}
{event_type, event_payload} = encode_event(event)
with {:ok, _event} <-
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: event.seq,
event_type: event_type,
payload: event_payload,
inserted_at: DateTime.utc_now()
})
|> repo().insert(),
{1, _} <-
repo().update_all(
from(s in Signal, where: s.id == ^signal.id),
set: [delivered: true]
) do
{:ok, payload}
else
{0, _} -> repo().rollback({:signal_consume_failed, :already_delivered})
{:error, changeset} -> repo().rollback({:signal_consume_failed, changeset})
end
end
end)
case result do
{:ok, value} ->
value
{:error, reason} ->
raise JournalError, op: :consume_signal, reason: reason
end
end
def fire_timer!(%Instance{} = instance, run_id, timer_id, lease_token) do
:ok = with_repo(instance, fn -> fire_timer_with_repo!(run_id, timer_id, lease_token) end)
Snapshotter.maybe_snapshot(instance, run_id, lease_token, __MODULE__)
end
defp fire_timer_with_repo!(run_id, timer_id, lease_token) do
result =
repo().transaction(fn ->
# State and lease, not lease only: a cancel committing between the
# wheel's claim and this fire must roll back as {:run_not_active, _}
# instead of appending timer_fired to a terminal run's history.
lock_and_validate_active_run!(run_id, lease_token)
case timer_winner(run_id, timer_id) do
{:pending, timer_event, winner_seq} ->
event = %{
type: :timer_fired,
timer_id: timer_id,
command_id: Map.get(timer_event, :command_id),
seq: winner_seq
}
winner_event = insert_event!(run_id, event)
mark_timer_resolved(run_id, timer_id, lease_token)
{:ok, winner_event}
{:already_fired, winner_event} ->
mark_timer_resolved(run_id, timer_id, lease_token)
{:ok, winner_event}
{:already_resolved, _winner_event} ->
mark_timer_resolved(run_id, timer_id, lease_token)
:already_resolved
:not_found ->
repo().rollback({:timer_fire_failed, :not_found})
:mismatch ->
repo().rollback({:timer_fire_failed, :winner_mismatch})
end
end)
case result do
{:ok, _value} ->
:ok
{:error, reason} ->
raise JournalError, op: :fire_timer!, reason: reason
end
end
def clear_next_wakeup!(%Instance{} = instance, run_id, lease_token) do
with_repo(instance, fn -> clear_next_wakeup_with_repo!(run_id, lease_token) end)
end
defp clear_next_wakeup_with_repo!(run_id, lease_token) do
cas_update_run(run_id, lease_token, %{next_wakeup_at: nil})
end
defp activity_task_result!(task, event, task_updates, lease_token) do
result =
repo().transaction(fn ->
lock_and_validate_active_run!(task.run_id, lease_token)
lock_and_validate_activity_task!(task)
with %{} <- insert_event!(task.run_id, event),
{1, _} <-
repo().update_all(
from(t in ActivityTask,
where:
t.id == ^task.id and t.run_id == ^task.run_id and t.state == "leased" and
t.lease_owner == ^task.lease_owner and t.attempt == ^task.attempt
),
set: task_updates
) do
:ok
else
{0, _} -> repo().rollback({:activity_task_result_failed, :task_lease_mismatch})
end
end)
case result do
{:ok, :ok} ->
Snapshotter.maybe_snapshot(task.instance, task.run_id, lease_token, __MODULE__)
:ok
{:error, reason} ->
raise JournalError, op: :activity_task_result!, reason: reason
end
end
defp maybe_commit_idempotency_result(_task, result, nil), do: result
defp maybe_commit_idempotency_result(task, result, idempotency) do
activity_module = activity_module_key(Keyword.fetch!(idempotency, :module))
idempotency_key = Keyword.fetch!(idempotency, :key)
now = DateTime.utc_now()
{count, _} =
repo().insert_all(
ActivityResult,
[
%{
activity_module: activity_module,
idempotency_key: idempotency_key,
run_id: task.run_id,
seq: task.seq + 1,
result: encode_term(result),
completed_at: now
}
],
on_conflict: :nothing,
conflict_target: [:activity_module, :idempotency_key]
)
case count do
1 -> result
0 -> fetch_activity_result!(activity_module, idempotency_key)
end
end
defp fetch_activity_result!(activity_module, idempotency_key) do
case repo().one(
from(r in ActivityResult,
where:
r.activity_module == ^activity_module and r.idempotency_key == ^idempotency_key
)
) do
%ActivityResult{} = result ->
decode_term(result.result)
nil ->
repo().rollback({:activity_result_conflict_failed, {activity_module, idempotency_key}})
end
end
defp activity_module_key(module) when is_atom(module), do: Atom.to_string(module)
defp activity_module_key(module) when is_binary(module), do: module
defp lock_and_validate_active_run!(run_id, lease_token) do
run =
repo().one(
from(r in Run,
where: r.id == ^run_id,
lock: "FOR UPDATE"
)
)
case run do
nil ->
repo().rollback({:run_not_found, run_id})
%Run{state: state} = run when state in ["running", "suspended"] ->
validate_lease!(run, lease_token)
%Run{state: state} ->
repo().rollback({:run_not_active, state})
end
end
defp lock_and_validate_activity_task!(task) do
# Expiry is measured on the database clock — the same clock that wrote
# lease_expires_at — so app/DB skew cannot shrink or stretch the TTL.
result =
repo().one(
from(t in ActivityTask,
where: t.id == ^task.id,
lock: "FOR UPDATE",
select: {t, fragment("clock_timestamp()")}
)
)
case result do
nil ->
repo().rollback({:activity_task_not_found, task.id})
{db_task, db_now} ->
cond do
db_task.run_id != task.run_id ->
repo().rollback({:activity_task_run_mismatch, task.id})
db_task.state != "leased" ->
repo().rollback({:activity_task_not_leased, db_task.state})
db_task.lease_owner != task.lease_owner ->
repo().rollback(
{:activity_task_lease_mismatch,
expected: task.lease_owner, actual: db_task.lease_owner}
)
db_task.attempt != task.attempt ->
repo().rollback(
{:activity_task_attempt_mismatch, expected: task.attempt, actual: db_task.attempt}
)
is_nil(db_task.lease_expires_at) ->
repo().rollback({:activity_task_lease_missing_expiry, task.id})
DateTime.compare(db_task.lease_expires_at, db_now) == :lt ->
repo().rollback({:activity_task_lease_expired, task.id})
true ->
:ok
end
end
end
defp maybe_insert_signal_timeout_timer(run_id, %{
timeout_timer_id: timer_id,
timeout_at: fires_at
}) do
changeset =
%Timer{}
|> Ecto.Changeset.change(%{
id: timer_id,
run_id: run_id,
fires_at: fires_at,
fired: false
})
case repo().insert(changeset) do
{:ok, _timer} -> notify_timer_armed_with_repo(run_id, fires_at)
{:error, changeset} -> {:error, changeset}
end
end
defp maybe_insert_signal_timeout_timer(_run_id, _event), do: :ok
defp maybe_set_signal_timeout_wakeup(run_id, %{timeout_at: timeout_at}, lease_token) do
case repo().update_all(
leased_run_query(run_id, lease_token),
set: [next_wakeup_at: timeout_at]
) do
{1, _} -> :ok
other -> other
end
end
defp maybe_set_signal_timeout_wakeup(_run_id, _event, _lease_token), do: :ok
defp notify_timer_armed_with_repo(run_id, fires_at) do
payload = "#{run_id}|#{DateTime.to_iso8601(fires_at)}"
case repo().query("SELECT pg_notify('continuum_timer_armed', $1)", [payload]) do
{:ok, _} -> :ok
{:error, reason} -> {:error, reason}
end
end
defp signal_await_winner(run_id, await_event) do
winner_seq = await_event.seq + 1
await_name = await_event.name
timeout_timer_id = Map.get(await_event, :timeout_timer_id)
run_id
|> event_at(winner_seq)
|> case do
nil ->
:none
%{type: :signal_received, name: ^await_name, payload: payload} = winner_event ->
{:ok, payload, winner_event}
%{type: :timer_fired, timer_id: ^timeout_timer_id} = winner_event
when not is_nil(timeout_timer_id) ->
{:timeout, winner_event}
_other ->
repo().rollback({:signal_await_failed, :winner_mismatch})
end
end
defp consume_signal_or_timeout(run_id, await_event) do
case pending_signal(run_id, await_event.name) do
nil ->
maybe_timeout_signal_await(run_id, await_event)
%Signal{} = signal ->
{:ok, payload, winner_event} =
consume_signal_row!(
run_id,
await_event.name,
signal,
Map.get(await_event, :command_id),
await_event.seq + 1
)
mark_signal_timeout_resolved(run_id, await_event)
{:ok, payload, winner_event}
end
end
defp consume_signal_row!(run_id, name, %Signal{} = signal, command_id, seq) do
payload = decode_term(signal.payload)
winner_event =
insert_event!(run_id, %{
type: :signal_received,
name: name,
payload: payload,
command_id: command_id,
seq: seq
})
with {1, _} <-
repo().update_all(
from(s in Signal, where: s.id == ^signal.id and s.delivered == false),
set: [delivered: true]
) do
{:ok, payload, winner_event}
else
{0, _} -> repo().rollback({:signal_consume_failed, :already_delivered})
end
end
defp maybe_timeout_signal_await(
run_id,
%{timeout_timer_id: timer_id, timeout_at: timeout_at} = event
) do
if DateTime.compare(db_now(), timeout_at) in [:gt, :eq] do
winner_event =
insert_event!(run_id, %{
type: :timer_fired,
timer_id: timer_id,
command_id: Map.get(event, :command_id),
seq: event.seq + 1
})
mark_timer_resolved(run_id, timer_id, nil)
{:timeout, winner_event}
else
:none
end
end
defp maybe_timeout_signal_await(_run_id, _event), do: :none
defp pending_signal(run_id, name) do
signal_name = Atom.to_string(name)
repo().one(
from(s in Signal,
where: s.run_id == ^run_id and s.name == ^signal_name and s.delivered == false,
order_by: [asc: s.inserted_at, asc: s.id],
limit: 1,
lock: "FOR UPDATE SKIP LOCKED"
)
)
end
defp mark_signal_timeout_resolved(run_id, %{timeout_timer_id: timer_id}) do
mark_timer_resolved(run_id, timer_id, nil)
end
defp mark_signal_timeout_resolved(_run_id, _event), do: :ok
defp timer_winner(run_id, timer_id) do
events = load_events(run_id)
case Enum.find(events, &timer_owner?(&1, timer_id)) do
nil ->
:not_found
timer_event ->
winner_seq = timer_event.seq + 1
case Enum.find(events, &(&1.seq == winner_seq)) do
nil ->
{:pending, timer_event, winner_seq}
%{type: :timer_fired, timer_id: ^timer_id} = winner_event ->
{:already_fired, winner_event}
%{type: :signal_received} = winner_event when timer_event.type == :signal_awaited ->
{:already_resolved, winner_event}
_other ->
:mismatch
end
end
end
defp timer_owner?(%{type: :timer_started, timer_id: event_timer_id}, timer_id)
when event_timer_id == timer_id,
do: true
defp timer_owner?(%{type: :signal_awaited, timeout_timer_id: event_timer_id}, timer_id)
when event_timer_id == timer_id,
do: true
defp timer_owner?(_event, _timer_id), do: false
defp event_at(run_id, seq) do
repo().one(
from(e in Event,
where: e.run_id == ^run_id and e.seq == ^seq
)
)
|> case do
nil -> nil
event -> decode_event(event)
end
end
defp load_events(run_id) do
repo().all(
from(e in Event,
where: e.run_id == ^run_id,
order_by: [asc: e.seq]
)
)
|> Enum.map(&decode_event/1)
end
defp load_events_after(run_id, through_seq) do
repo().all(
from(e in Event,
where: e.run_id == ^run_id and e.seq > ^through_seq,
order_by: [asc: e.seq]
)
)
|> Enum.map(&decode_event/1)
end
defp latest_snapshot(run_id) do
# Newest *decodable* snapshot, not newest overall: in a mixed-version
# rolling deploy a newer node may have written a format this release
# cannot decode. Skipping it degrades to an older snapshot or full event
# replay (events are never pruned) instead of a crash/reclaim loop.
repo().one(
from(s in Snapshot,
where:
s.run_id == ^run_id and
s.format_version <= ^Continuum.Snapshot.format_version(),
order_by: [desc: s.through_seq],
limit: 1
)
)
|> case do
nil -> nil
snapshot -> decode_snapshot_or_skip(snapshot)
end
end
defp decode_snapshot_or_skip(%Snapshot{} = snapshot) do
decode_snapshot(snapshot)
rescue
error ->
Logger.warning(
"Continuum snapshot #{snapshot.id} for run #{snapshot.run_id} is undecodable, " <>
"falling back to event replay: #{Exception.message(error)}"
)
nil
end
defp insert_event!(run_id, event) do
{event_type, payload} = encode_event(event)
seq = event.seq || next_seq(run_id)
changeset =
%Event{}
|> Ecto.Changeset.change(%{
run_id: run_id,
seq: seq,
event_type: event_type,
payload: payload,
inserted_at: DateTime.utc_now()
})
case repo().insert(changeset) do
{:ok, event_record} -> decode_event(event_record)
{:error, changeset} -> repo().rollback({:event_insert_failed, changeset})
end
end
defp maybe_snapshot_after_event(instance, run_id, event, lease_token) do
if advancing_event?(Map.get(event, :type)) do
Snapshotter.maybe_snapshot(instance, run_id, lease_token, __MODULE__)
else
:ok
end
end
defp maybe_snapshot_after_signal_resolution(_instance, _run_id, :none, _lease_token), do: :ok
defp maybe_snapshot_after_signal_resolution(instance, run_id, _value, lease_token) do
Snapshotter.maybe_snapshot(instance, run_id, lease_token, __MODULE__)
end
defp advancing_event?(type) do
type in [
:side_effect,
:activity_completed,
:activity_failed,
:signal_received,
:timer_fired,
:patched,
:compensation_completed,
:compensation_failed
]
end
defp mark_timer_resolved(run_id, timer_id, lease_token) do
repo().update_all(
from(t in Timer, where: t.run_id == ^run_id and t.id == ^timer_id),
set: [fired: true]
)
run_query =
case lease_token do
nil -> from(r in Run, where: r.id == ^run_id)
token -> leased_run_query(run_id, token)
end
repo().update_all(run_query, set: [next_wakeup_at: nil])
:ok
end
@impl true
def suspend!(%Instance{} = instance, run_id, lease_token) do
with_repo(instance, fn -> suspend_with_repo!(run_id, lease_token) end)
end
defp suspend_with_repo!(run_id, lease_token) do
cas_update_active_run(run_id, lease_token, %{state: "suspended"})
end
@impl true
def complete!(%Instance{} = instance, run_id, result, lease_token) do
parent_run_id =
with_repo(instance, fn ->
parent =
run_in_transaction!(fn ->
lock_parent_first(run_id)
:ok =
cas_update_active_run(run_id, lease_token, %{
state: "completed",
result: encode_term(result),
completed_at: DateTime.utc_now()
})
maybe_wake_parent(run_id)
end)
Snapshotter.maybe_snapshot(instance, run_id, lease_token, __MODULE__)
Continuum.Runtime.Engine.broadcast_run_finished(instance, run_id, :completed, result)
parent
end)
wake_parent(instance, parent_run_id)
:ok
end
@impl true
def fail!(%Instance{} = instance, run_id, error, lease_token) do
parent_run_id =
with_repo(instance, fn ->
parent =
run_in_transaction!(fn ->
lock_parent_first(run_id)
:ok =
cas_update_active_run(run_id, lease_token, %{
state: "failed",
error: encode_term(error),
completed_at: DateTime.utc_now()
})
maybe_wake_parent(run_id)
end)
Snapshotter.maybe_snapshot(instance, run_id, lease_token, __MODULE__)
broadcast_failed(instance, run_id, error)
parent
end)
wake_parent(instance, parent_run_id)
:ok
end
@unknown_version_backoff_ms 5_000
@doc """
Release a run whose version is not loaded on this node.
Clears the lease and leaves the run `suspended` so another node that has
the version loaded can claim it — an unknown version is a per-node fact,
not a global one. Runs are no longer marked `stuck_unknown_version`;
`Continuum.VersionRegistry.upsert_instance/2` recovers legacy stuck rows.
The release also pushes `next_wakeup_at` a few seconds out: without the
backoff the same incapable node re-claims the run on its next poll (and
`NULLS FIRST` puts it at the head of every claim batch, starving runnable
work). Signals, timers, and parent wakes overwrite the backoff, and any
node — including a newly capable one — claims the run once it lapses.
"""
def release_unknown_version!(%Instance{} = instance, run_id, lease_token) do
backoff_until =
DateTime.add(DateTime.utc_now(), @unknown_version_backoff_ms, :millisecond)
with_repo(instance, fn ->
:ok =
cas_update_active_run(run_id, lease_token, %{
state: "suspended",
lease_owner: nil,
lease_token: nil,
lease_expires_at: nil,
next_wakeup_at: backoff_until
})
end)
end
@impl true
def get_run(%Instance{} = instance, run_id) do
with_repo(instance, fn -> get_run_with_repo(run_id) end)
end
defp get_run_with_repo(run_id) do
case repo().one(from(r in Run, where: r.id == ^run_id)) do
nil -> nil
run -> decode_run(run)
end
end
defp decode_snapshot(%Snapshot{payload: payload}) do
Continuum.Snapshot.decode(payload)
end
defp compensation_terminal_seq(%{parallel_batch?: true}), do: nil
defp compensation_terminal_seq(task), do: task.seq + 1
defp cas_update_run(run_id, lease_token, updates) do
query = leased_run_query(run_id, lease_token)
case repo().update_all(query, set: Map.to_list(updates)) do
{1, _} ->
:ok
{0, _} ->
raise JournalError, op: :cas_update_run, reason: {:cas_failed, run_id}
end
end
# Terminal and suspension transitions additionally require the run to still
# be active: a late fail!/suspend! must never flip a run that already
# reached completed/failed (e.g. a raise after complete! committed).
defp cas_update_active_run(run_id, lease_token, updates) do
query =
leased_run_query(run_id, lease_token)
|> where([r], r.state in ["running", "suspended"])
case repo().update_all(query, set: Map.to_list(updates)) do
{1, _} ->
:ok
{0, _} ->
raise JournalError, op: :cas_update_run, reason: {:cas_failed, run_id}
end
end
defp validate_lease!(%Run{lease_token: nil, lease_owner: nil}, nil), do: :ok
defp validate_lease!(%Run{lease_token: token}, token) when not is_nil(token), do: :ok
defp validate_lease!(%Run{} = run, lease_token) do
repo().rollback(
{:lease_mismatch,
expected: lease_token,
actual: %{lease_owner: run.lease_owner, lease_token: run.lease_token}}
)
end
defp lock_and_validate_run!(run_id, lease_token) do
run =
repo().one(
from(r in Run,
where: r.id == ^run_id,
lock: "FOR UPDATE"
)
)
case run do
nil -> repo().rollback({:run_not_found, run_id})
%Run{} = run -> validate_lease!(run, lease_token)
end
end
defp lock_and_load_run!(run_id, lease_token) do
run =
repo().one(
from(r in Run,
where: r.id == ^run_id,
lock: "FOR UPDATE"
)
)
case run do
nil ->
repo().rollback({:run_not_found, run_id})
%Run{} = run ->
:ok = validate_lease!(run, lease_token)
run
end
end
defp leased_run_query(run_id, nil) do
from(r in Run,
where: r.id == ^run_id and is_nil(r.lease_owner) and is_nil(r.lease_token)
)
end
defp leased_run_query(run_id, lease_token) do
from(r in Run,
where: r.id == ^run_id and r.lease_token == ^lease_token
)
end
defp next_seq(run_id) do
case repo().one(
from(e in Event,
where: e.run_id == ^run_id,
select: max(e.seq)
)
) do
nil -> 0
seq -> seq + 1
end
end
defp encode_event(%{type: type} = event) do
payload =
event
|> Map.delete(:type)
|> Map.delete(:seq)
|> encode_term()
{Atom.to_string(type), payload}
end
defp decode_event(%Event{event_type: event_type, payload: payload, seq: seq}) do
decoded = decode_term(payload)
type = String.to_atom(event_type)
decoded
|> Map.put(:type, type)
|> Map.put(:seq, seq)
|> atomize_keys_by_type(type)
end
defp atomize_keys_by_type(map, :side_effect) do
map
|> maybe_atomize(:kind)
end
defp atomize_keys_by_type(map, :activity_completed) do
map
|> maybe_decode_mfa()
end
defp atomize_keys_by_type(map, :signal_received) do
map
|> maybe_atomize(:name)
end
defp atomize_keys_by_type(map, _type), do: map
defp maybe_atomize(map, key) do
case Map.get(map, key) || Map.get(map, to_string(key)) do
nil -> map
val when is_binary(val) -> Map.put(map, key, String.to_atom(val))
_ -> map
end
end
defp maybe_decode_mfa(map) do
key = :mfa
str_key = "mfa"
case Map.get(map, key) || Map.get(map, str_key) do
[mod, fun, args] when is_binary(mod) and is_binary(fun) and is_list(args) ->
Map.put(map, key, {String.to_atom("Elixir." <> mod), String.to_atom(fun), args})
[mod, fun, args] when is_atom(mod) and is_atom(fun) and is_list(args) ->
Map.put(map, key, {mod, fun, args})
_ ->
map
end
end
defp decode_run(%Run{} = run) do
%{
run_id: run.id,
workflow: run.workflow,
state: String.to_atom(run.state),
result: decode_term(run.result),
error: decode_term(run.error),
input: decode_term(run.input),
attributes: run.attributes || %{},
namespace: run.namespace || "default",
version_hash: run.version_hash,
trace_context: run.trace_context
}
end
# Real database time (advances inside a transaction, unlike now()) — used
# where a comparison needs "current time" on the same clock that wrote the
# compared column.
defp db_now do
{:ok, %{rows: [[now]]}} = repo().query("SELECT clock_timestamp()")
now
end
defp encode_term(nil), do: nil
defp encode_term(term), do: :erlang.term_to_binary(term)
defp decode_term(nil), do: nil
defp decode_term(binary) when is_binary(binary), do: :erlang.binary_to_term(binary)
defp decode_term(other), do: other
defp broadcast_failed(_instance, _run_id, {_kind, _reason, stacktrace})
when is_list(stacktrace),
do: :ok
defp broadcast_failed(instance, run_id, error) do
Continuum.Runtime.Engine.broadcast_run_finished(instance, run_id, :failed, error)
end
defp with_repo(%Instance{} = instance, fun) do
previous = Process.get(:continuum_repo)
Process.put(:continuum_repo, instance.repo)
try do
fun.()
after
if is_nil(previous),
do: Process.delete(:continuum_repo),
else: Process.put(:continuum_repo, previous)
end
end
defp repo do
Process.get(:continuum_repo) || Application.fetch_env!(:continuum, :repo)
end
end