Skip to content

dura.engine

Defined in dura/engine.py. DurableEngine and the names below it are all re-exported from the top-level dura package.

DurableEngine

DurableEngine(db_path: str | Path, *, clock: Callable[[], datetime] | None = None, busy_timeout_ms: int = 5000)
Source code in dura/engine.py
def __init__(
    self,
    db_path: str | Path,
    *,
    clock: Callable[[], datetime] | None = None,
    busy_timeout_ms: int = 5000,
) -> None:
    self._db_path = str(db_path)
    self._clock = clock or (lambda: datetime.now(UTC))
    self._busy_timeout_ms = busy_timeout_ms
    self._local = threading.local()
    self._conn().executescript(_SCHEMA)

db_path property

db_path: str

Filesystem path of the engine database.

close

close() -> None

Close the connection for the calling thread, if any.

Source code in dura/engine.py
def close(self) -> None:
    """Close the connection for the calling thread, if any."""
    conn: sqlite3.Connection | None = getattr(self._local, "conn", None)
    if conn is not None:
        conn.close()
        self._local.conn = None

spawn_task

spawn_task(*, name: str, params: Any, idempotency_key: str | None = None, retry: RetryStrategy | None = None, max_attempts: int | None = None, available_after: timedelta | None = None, priority: int = 0) -> TaskRef

Create a task and its first run.

priority orders claiming: higher is claimed first (ties broken by availability, then insertion order). The default 0 is the lowest band; give maintenance and other latency-sensitive work a higher value so it is not starved behind a backlog of low-priority work. Retries inherit the run's priority.

If idempotency_key is supplied and a task already exists for it, no new task is created and the existing one is returned with created=False.

Source code in dura/engine.py
def spawn_task(
    self,
    *,
    name: str,
    params: Any,
    idempotency_key: str | None = None,
    retry: RetryStrategy | None = None,
    max_attempts: int | None = None,
    available_after: timedelta | None = None,
    priority: int = 0,
) -> TaskRef:
    """Create a task and its first run.

    ``priority`` orders claiming: higher is claimed first (ties broken by
    availability, then insertion order). The default 0 is the lowest band;
    give maintenance and other latency-sensitive work a higher value so it
    is not starved behind a backlog of low-priority work. Retries inherit
    the run's priority.

    If ``idempotency_key`` is supplied and a task already exists for it, no
    new task is created and the existing one is returned with
    ``created=False``.
    """
    now = self._now()
    available_at = now + (available_after or timedelta(0))
    delayed = available_at > now
    state = "sleeping" if delayed else "pending"
    task_id = _new_id()
    run_id = _new_id()

    with self._tx() as conn:
        if idempotency_key is not None:
            existing = conn.execute(
                "SELECT task_id, last_attempt_run, attempts "
                "FROM tasks WHERE idempotency_key = ?",
                (idempotency_key,),
            ).fetchone()
            if existing is not None:
                return TaskRef(
                    task_id=existing["task_id"],
                    run_id=existing["last_attempt_run"],
                    attempt=existing["attempts"],
                    created=False,
                )

        conn.execute(
            "INSERT INTO tasks (task_id, task_name, params, state, attempts, "
            "max_attempts, retry_strategy, enqueue_at, last_attempt_run, "
            "idempotency_key) VALUES (?, ?, ?, ?, 1, ?, ?, ?, ?, ?)",
            (
                task_id,
                name,
                _dumps(params),
                state,
                max_attempts,
                _dumps(retry.to_dict()) if retry else None,
                _fmt(now),
                run_id,
                idempotency_key,
            ),
        )
        conn.execute(
            "INSERT INTO runs (run_id, task_id, attempt, state, available_at, "
            "priority) VALUES (?, ?, 1, ?, ?, ?)",
            (run_id, task_id, state, _fmt(available_at), priority),
        )
    return TaskRef(task_id=task_id, run_id=run_id, attempt=1, created=True)

claim_task

claim_task(*, worker_id: str, timeout_secs: int = 120, max_priority: int | None = None) -> ClaimedTask | None

Atomically reserve the next available run, or return None.

Before claiming, runs whose lease has expired (a worker took them but never finished) are reset to pending so they can be picked up again. A reclaimed run keeps its attempt number: a crash is not a logical failure, so it does not consume a retry.

max_priority: when set, only claim runs whose priority is at or below this value. Used by lane workers dedicated to low-priority tasks.

Source code in dura/engine.py
def claim_task(
    self,
    *,
    worker_id: str,
    timeout_secs: int = 120,
    max_priority: int | None = None,
) -> ClaimedTask | None:
    """Atomically reserve the next available run, or return ``None``.

    Before claiming, runs whose lease has expired (a worker took them but
    never finished) are reset to ``pending`` so they can be picked up again.
    A reclaimed run keeps its attempt number: a crash is not a logical
    failure, so it does not consume a retry.

    ``max_priority``: when set, only claim runs whose priority is at or
    below this value. Used by lane workers dedicated to low-priority tasks.
    """
    now = self._now()
    now_str = _fmt(now)
    expires_str = _fmt(now + timedelta(seconds=timeout_secs))
    # A sentinel well above any real priority lets the general (unrestricted)
    # case share the same query as the lane case.
    priority_ceil = max_priority if max_priority is not None else 10_000

    with self._tx() as conn:
        conn.execute(
            "UPDATE runs SET state = 'pending', claimed_by = NULL, "
            "claim_expires_at = NULL "
            "WHERE state = 'running' AND claim_expires_at IS NOT NULL "
            "AND claim_expires_at <= ?",
            (now_str,),
        )

        # Pin runs_claimable: its (priority DESC, available_at) order matches
        # the ORDER BY, so the highest-priority due run is the first index
        # row. Without the hint the planner sorts the whole candidate set
        # (a 60k-row temp B-tree under a storm) unless ANALYZE has been run.
        row = conn.execute(
            "SELECT r.run_id, r.task_id, r.attempt, t.task_name, t.params "
            "FROM runs r INDEXED BY runs_claimable "
            "JOIN tasks t ON t.task_id = r.task_id "
            "WHERE r.state IN ('pending', 'sleeping') "
            "AND t.state IN ('pending', 'sleeping', 'running') "
            "AND r.available_at <= ? "
            "AND r.priority <= ? "
            "ORDER BY r.priority DESC, r.available_at, r.rowid LIMIT 1",
            (now_str, priority_ceil),
        ).fetchone()
        if row is None:
            return None

        conn.execute(
            "UPDATE runs SET state = 'running', claimed_by = ?, "
            "claim_expires_at = ?, started_at = COALESCE(started_at, ?) "
            "WHERE run_id = ?",
            (worker_id, expires_str, now_str, row["run_id"]),
        )
        conn.execute(
            "UPDATE tasks SET state = 'running', attempts = MAX(attempts, ?), "
            "first_started_at = COALESCE(first_started_at, ?), "
            "last_attempt_run = ? WHERE task_id = ?",
            (row["attempt"], now_str, row["run_id"], row["task_id"]),
        )

    return ClaimedTask(
        run_id=row["run_id"],
        task_id=row["task_id"],
        attempt=row["attempt"],
        name=row["task_name"],
        params=_loads(row["params"]),
    )

complete_run

complete_run(*, run_id: str, result: Any = None) -> None

Mark a run and its task as completed.

Source code in dura/engine.py
def complete_run(self, *, run_id: str, result: Any = None) -> None:
    """Mark a run and its task as completed."""
    now_str = _fmt(self._now())
    with self._tx() as conn:
        row = conn.execute(
            "SELECT task_id, state FROM runs WHERE run_id = ?", (run_id,)
        ).fetchone()
        self._guard_running(row, run_id)

        conn.execute(
            "UPDATE runs SET state = 'completed', completed_at = ?, result = ? "
            "WHERE run_id = ?",
            (now_str, _dumps(result), run_id),
        )
        conn.execute(
            "UPDATE tasks SET state = 'completed', completed_payload = ?, "
            "terminal_at = ?, last_attempt_run = ? WHERE task_id = ?",
            (_dumps(result), now_str, run_id, row["task_id"]),
        )
        conn.execute("DELETE FROM waits WHERE run_id = ?", (run_id,))

fail_run

fail_run(*, run_id: str, reason: dict[str, Any], retryable: bool = True) -> None

Mark a run as failed and schedule a retry if attempts remain.

retryable=False skips scheduling a retry regardless of attempts remaining, and fails the task outright. Use it for failures a retry can never fix, e.g. no handler is registered for the task's name (a deploy/config problem, not a transient one).

Source code in dura/engine.py
def fail_run(
    self, *, run_id: str, reason: dict[str, Any], retryable: bool = True
) -> None:
    """Mark a run as failed and schedule a retry if attempts remain.

    ``retryable=False`` skips scheduling a retry regardless of attempts
    remaining, and fails the task outright. Use it for failures a retry
    can never fix, e.g. no handler is registered for the task's name (a
    deploy/config problem, not a transient one).
    """
    now = self._now()
    now_str = _fmt(now)
    with self._tx() as conn:
        row = conn.execute(
            "SELECT task_id, attempt, state, priority FROM runs WHERE run_id = ?",
            (run_id,),
        ).fetchone()
        self._guard_running(row, run_id)
        task_id = row["task_id"]
        attempt = row["attempt"]

        task = conn.execute(
            "SELECT retry_strategy, max_attempts FROM tasks WHERE task_id = ?",
            (task_id,),
        ).fetchone()

        conn.execute(
            "UPDATE runs SET state = 'failed', failed_at = ?, failure_reason = ? "
            "WHERE run_id = ?",
            (now_str, _dumps(reason), run_id),
        )

        next_attempt = attempt + 1
        max_attempts = task["max_attempts"]
        if retryable and (max_attempts is None or next_attempt <= max_attempts):
            strategy = (
                _loads(task["retry_strategy"]) if task["retry_strategy"] else None
            )
            delay = _retry_delay(strategy, attempt)
            available_at = now + timedelta(seconds=delay)
            run_state = "pending" if delay <= 0 else "sleeping"
            new_run_id = _new_id()
            conn.execute(
                "INSERT INTO runs (run_id, task_id, attempt, state, available_at, "
                "priority) VALUES (?, ?, ?, ?, ?, ?)",
                (
                    new_run_id,
                    task_id,
                    next_attempt,
                    run_state,
                    _fmt(available_at),
                    row["priority"],
                ),
            )
            conn.execute(
                "UPDATE tasks SET state = ?, attempts = MAX(attempts, ?), "
                "last_attempt_run = ? WHERE task_id = ?",
                (run_state, next_attempt, new_run_id, task_id),
            )
        else:
            conn.execute(
                "UPDATE tasks SET state = 'failed', attempts = MAX(attempts, ?), "
                "terminal_at = ?, last_attempt_run = ? WHERE task_id = ?",
                (attempt, now_str, run_id, task_id),
            )
        conn.execute("DELETE FROM waits WHERE run_id = ?", (run_id,))

cancel_task

cancel_task(task_id: str) -> None

Cancel a task and any of its non-terminal runs.

Source code in dura/engine.py
def cancel_task(self, task_id: str) -> None:
    """Cancel a task and any of its non-terminal runs."""
    now_str = _fmt(self._now())
    with self._tx() as conn:
        row = conn.execute(
            "SELECT state FROM tasks WHERE task_id = ?", (task_id,)
        ).fetchone()
        if row is None:
            raise TaskNotFound(f"Task {task_id} not found")
        if row["state"] in _TERMINAL_STATES:
            return
        conn.execute(
            "UPDATE tasks SET state = 'cancelled', terminal_at = ? "
            "WHERE task_id = ?",
            (now_str, task_id),
        )
        conn.execute(
            "UPDATE runs SET state = 'cancelled', claimed_by = NULL, "
            "claim_expires_at = NULL "
            "WHERE task_id = ? AND state NOT IN ('completed', 'failed', 'cancelled')",
            (task_id,),
        )
        conn.execute("DELETE FROM waits WHERE task_id = ?", (task_id,))

cancel_duplicate_tasks

cancel_duplicate_tasks(task_name: str) -> tuple[str | None, int]

Cancel all but the oldest active task with the given name.

Returns (surviving_task_id, count_cancelled). surviving_task_id is None when no non-terminal tasks with this name exist. Intended for use at startup to collapse duplicate self-perpetuating chains that were seeded across multiple restarts before the state-table guard was in place.

Source code in dura/engine.py
def cancel_duplicate_tasks(self, task_name: str) -> tuple[str | None, int]:
    """Cancel all but the oldest active task with the given name.

    Returns ``(surviving_task_id, count_cancelled)``.  ``surviving_task_id``
    is ``None`` when no non-terminal tasks with this name exist.  Intended
    for use at startup to collapse duplicate self-perpetuating chains that
    were seeded across multiple restarts before the state-table guard was in
    place.
    """
    rows = (
        self._conn()
        .execute(
            "SELECT task_id FROM tasks "
            "WHERE task_name = ? AND state NOT IN ('completed', 'failed', 'cancelled') "
            "ORDER BY enqueue_at ASC",
            (task_name,),
        )
        .fetchall()
    )
    if not rows:
        return None, 0
    ids = [r["task_id"] for r in rows]
    for dup_id in ids[1:]:
        self.cancel_task(task_id=dup_id)
    return ids[0], len(ids) - 1

extend_claim

extend_claim(*, run_id: str, by_secs: int) -> None

Push a running lease forward (heartbeat for long steps).

Source code in dura/engine.py
def extend_claim(self, *, run_id: str, by_secs: int) -> None:
    """Push a running lease forward (heartbeat for long steps)."""
    if by_secs <= 0:
        raise ValueError("by_secs must be > 0")
    new_expiry = _fmt(self._now() + timedelta(seconds=by_secs))
    with self._tx() as conn:
        row = conn.execute(
            "SELECT state FROM runs WHERE run_id = ?", (run_id,)
        ).fetchone()
        if row is None:
            raise TaskNotFound(f"Run {run_id} not found")
        if row["state"] != "running":
            raise InvalidRunState(f"Run {run_id} is not running")
        conn.execute(
            "UPDATE runs SET claim_expires_at = ? WHERE run_id = ?",
            (new_expiry, run_id),
        )

checkpoint

checkpoint(*, task_id: str, step_name: str, fn: Callable[[], T], owner_run_id: str | None = None) -> T

Run fn once and persist its result, keyed to the task.

On the first call the result is computed and stored. On any later call for the same (task_id, step_name) - including a retry after a crash - the stored result is returned and fn is not invoked again.

fn runs outside the write transaction so a slow step (FTP, S3) does not hold the database lock.

Source code in dura/engine.py
def checkpoint(
    self,
    *,
    task_id: str,
    step_name: str,
    fn: Callable[[], T],
    owner_run_id: str | None = None,
) -> T:
    """Run ``fn`` once and persist its result, keyed to the task.

    On the first call the result is computed and stored. On any later call
    for the same ``(task_id, step_name)`` - including a retry after a crash -
    the stored result is returned and ``fn`` is not invoked again.

    ``fn`` runs outside the write transaction so a slow step (FTP, S3) does
    not hold the database lock.
    """
    conn = self._conn()
    existing = conn.execute(
        "SELECT state FROM checkpoints WHERE task_id = ? AND checkpoint_name = ?",
        (task_id, step_name),
    ).fetchone()
    if existing is not None:
        return _loads(existing["state"])

    result = fn()

    with self._tx() as conn:
        conn.execute(
            "INSERT INTO checkpoints (task_id, checkpoint_name, state, "
            "owner_run_id, updated_at) VALUES (?, ?, ?, ?, ?) "
            "ON CONFLICT (task_id, checkpoint_name) DO NOTHING",
            (task_id, step_name, _dumps(result), owner_run_id, _fmt(self._now())),
        )
        # Read back the authoritative value: if a concurrent run committed
        # first, everyone agrees on the same stored result.
        stored = conn.execute(
            "SELECT state FROM checkpoints "
            "WHERE task_id = ? AND checkpoint_name = ?",
            (task_id, step_name),
        ).fetchone()
    return _loads(stored["state"])

get_checkpoint

get_checkpoint(*, task_id: str, step_name: str) -> Any

Return a stored checkpoint payload, or None if absent.

Source code in dura/engine.py
def get_checkpoint(self, *, task_id: str, step_name: str) -> Any:
    """Return a stored checkpoint payload, or ``None`` if absent."""
    row = (
        self._conn()
        .execute(
            "SELECT state FROM checkpoints WHERE task_id = ? AND checkpoint_name = ?",
            (task_id, step_name),
        )
        .fetchone()
    )
    return _loads(row["state"]) if row is not None else None

get_state

get_state(*, namespace: str, key: str, default: Any = None) -> Any

Return the value stored at (namespace, key), or default.

Source code in dura/engine.py
def get_state(self, *, namespace: str, key: str, default: Any = None) -> Any:
    """Return the value stored at ``(namespace, key)``, or ``default``."""
    row = (
        self._conn()
        .execute(
            "SELECT value FROM state WHERE namespace = ? AND key = ?",
            (namespace, key),
        )
        .fetchone()
    )
    return _loads(row["value"]) if row is not None else default

has_state

has_state(namespace: str) -> bool

Whether any key exists under namespace (cheap emptiness check).

Source code in dura/engine.py
def has_state(self, namespace: str) -> bool:
    """Whether any key exists under ``namespace`` (cheap emptiness check)."""
    row = (
        self._conn()
        .execute("SELECT 1 FROM state WHERE namespace = ? LIMIT 1", (namespace,))
        .fetchone()
    )
    return row is not None

set_state

set_state(*, namespace: str, key: str, value: Any) -> None

Store value at (namespace, key) (upsert, last write wins).

Source code in dura/engine.py
def set_state(self, *, namespace: str, key: str, value: Any) -> None:
    """Store ``value`` at ``(namespace, key)`` (upsert, last write wins)."""
    with self._tx() as conn:
        conn.execute(
            "INSERT INTO state (namespace, key, value, updated_at) "
            "VALUES (?, ?, ?, ?) "
            "ON CONFLICT (namespace, key) DO UPDATE SET "
            "value = excluded.value, updated_at = excluded.updated_at",
            (namespace, key, _dumps(value), _fmt(self._now())),
        )

set_state_many

set_state_many(*, namespace: str, items: Mapping[str, Any]) -> int

Upsert many key -> value entries under namespace in one transaction. Much faster than a loop of set_state for bulk imports (one commit, not one per key). Returns the number written.

Source code in dura/engine.py
def set_state_many(self, *, namespace: str, items: Mapping[str, Any]) -> int:
    """Upsert many ``key -> value`` entries under ``namespace`` in one
    transaction. Much faster than a loop of ``set_state`` for bulk
    imports (one commit, not one per key). Returns the number written.
    """
    now = _fmt(self._now())
    rows = [(namespace, key, _dumps(value), now) for key, value in items.items()]
    with self._tx() as conn:
        conn.executemany(
            "INSERT INTO state (namespace, key, value, updated_at) "
            "VALUES (?, ?, ?, ?) "
            "ON CONFLICT (namespace, key) DO UPDATE SET "
            "value = excluded.value, updated_at = excluded.updated_at",
            rows,
        )
    return len(rows)

delete_state

delete_state(*, namespace: str, key: str) -> bool

Delete (namespace, key). Returns True if it existed.

Source code in dura/engine.py
def delete_state(self, *, namespace: str, key: str) -> bool:
    """Delete ``(namespace, key)``. Returns ``True`` if it existed."""
    with self._tx() as conn:
        cur = conn.execute(
            "DELETE FROM state WHERE namespace = ? AND key = ?",
            (namespace, key),
        )
        return cur.rowcount > 0

list_state

list_state(namespace: str) -> dict[str, Any]

Return all key -> value pairs in namespace (ordered by key).

Source code in dura/engine.py
def list_state(self, namespace: str) -> dict[str, Any]:
    """Return all ``key -> value`` pairs in ``namespace`` (ordered by key)."""
    rows = (
        self._conn()
        .execute(
            "SELECT key, value FROM state WHERE namespace = ? ORDER BY key",
            (namespace,),
        )
        .fetchall()
    )
    return {row["key"]: _loads(row["value"]) for row in rows}

update_state

update_state(*, namespace: str, key: str, fn: Callable[[Any], Any], default: Any = None) -> Any

Atomic read-modify-write of (namespace, key).

fn receives the current value (or default if absent) and returns the new value, all inside the write transaction so concurrent workers cannot interleave. Keep fn pure and fast - no I/O - since it runs while the database write lock is held. Returns the new value.

Source code in dura/engine.py
def update_state(
    self,
    *,
    namespace: str,
    key: str,
    fn: Callable[[Any], Any],
    default: Any = None,
) -> Any:
    """Atomic read-modify-write of ``(namespace, key)``.

    ``fn`` receives the current value (or ``default`` if absent) and returns
    the new value, all inside the write transaction so concurrent workers
    cannot interleave. Keep ``fn`` pure and fast - no I/O - since it runs
    while the database write lock is held. Returns the new value.
    """
    with self._tx() as conn:
        row = conn.execute(
            "SELECT value FROM state WHERE namespace = ? AND key = ?",
            (namespace, key),
        ).fetchone()
        current = _loads(row["value"]) if row is not None else default
        new_value = fn(current)
        conn.execute(
            "INSERT INTO state (namespace, key, value, updated_at) "
            "VALUES (?, ?, ?, ?) "
            "ON CONFLICT (namespace, key) DO UPDATE SET "
            "value = excluded.value, updated_at = excluded.updated_at",
            (namespace, key, _dumps(new_value), _fmt(self._now())),
        )
    return new_value

emit_event

emit_event(*, event_name: str, payload: Any = None) -> None

Emit a named signal (first write wins) and wake its waiters.

Source code in dura/engine.py
def emit_event(self, *, event_name: str, payload: Any = None) -> None:
    """Emit a named signal (first write wins) and wake its waiters."""
    now = self._now()
    now_str = _fmt(now)
    with self._tx() as conn:
        already = conn.execute(
            "SELECT 1 FROM events WHERE event_name = ?", (event_name,)
        ).fetchone()
        if already is not None:
            return  # first emit wins; ignore later ones

        conn.execute(
            "INSERT INTO events (event_name, payload, emitted_at) VALUES (?, ?, ?)",
            (event_name, _dumps(payload), now_str),
        )

        waiters = conn.execute(
            "SELECT run_id, task_id FROM waits "
            "WHERE event_name = ? AND (timeout_at IS NULL OR timeout_at > ?)",
            (event_name, now_str),
        ).fetchall()
        for waiter in waiters:
            conn.execute(
                "UPDATE runs SET state = 'pending', available_at = ?, "
                "event_payload = ?, wake_event = NULL, claimed_by = NULL, "
                "claim_expires_at = NULL "
                "WHERE run_id = ? AND state = 'sleeping'",
                (now_str, _dumps(payload), waiter["run_id"]),
            )
            conn.execute(
                "UPDATE tasks SET state = 'pending' WHERE task_id = ?",
                (waiter["task_id"],),
            )
        conn.execute("DELETE FROM waits WHERE event_name = ?", (event_name,))

await_event

await_event(*, run_id: str, task_id: str, step_name: str, event_name: str, timeout_secs: int | None = None) -> tuple[bool, Any]

Suspend the current run until event_name fires or the timeout.

Returns (should_suspend, payload):

  • (False, payload) - the event already fired (or fired while we slept); the worker continues.
  • (False, None) - the wait timed out; the worker continues.
  • (True, None) - the run has been parked; the worker must return without completing or failing the run. It will be re-claimed when the event fires or the timeout elapses.

The outcome is checkpointed, so re-execution after a resume resolves the same step without parking again.

Source code in dura/engine.py
def await_event(
    self,
    *,
    run_id: str,
    task_id: str,
    step_name: str,
    event_name: str,
    timeout_secs: int | None = None,
) -> tuple[bool, Any]:
    """Suspend the current run until ``event_name`` fires or the timeout.

    Returns ``(should_suspend, payload)``:

    * ``(False, payload)`` - the event already fired (or fired while we
      slept); the worker continues.
    * ``(False, None)``    - the wait timed out; the worker continues.
    * ``(True, None)``     - the run has been parked; the worker must return
      without completing or failing the run. It will be re-claimed when the
      event fires or the timeout elapses.

    The outcome is checkpointed, so re-execution after a resume resolves the
    same step without parking again.
    """
    now = self._now()
    with self._tx() as conn:
        # 1. Resolved on a previous pass?
        checkpoint = conn.execute(
            "SELECT state FROM checkpoints "
            "WHERE task_id = ? AND checkpoint_name = ?",
            (task_id, step_name),
        ).fetchone()
        if checkpoint is not None:
            return False, _loads(checkpoint["state"])

        run = conn.execute(
            "SELECT wake_event, event_payload FROM runs WHERE run_id = ?",
            (run_id,),
        ).fetchone()
        if run is None:
            raise TaskNotFound(f"Run {run_id} not found")

        # 2. Woken by emit_event: payload was delivered onto our run.
        if run["event_payload"] is not None:
            payload = _loads(run["event_payload"])
            self._write_checkpoint(conn, task_id, step_name, payload, run_id, now)
            conn.execute(
                "UPDATE runs SET event_payload = NULL, wake_event = NULL "
                "WHERE run_id = ?",
                (run_id,),
            )
            return False, payload

        # 3. Event already emitted before we asked.
        event = conn.execute(
            "SELECT payload FROM events WHERE event_name = ?", (event_name,)
        ).fetchone()
        if event is not None:
            payload = _loads(event["payload"])
            self._write_checkpoint(conn, task_id, step_name, payload, run_id, now)
            return False, payload

        # 4. Woken by timeout (we were already waiting on this event).
        if run["wake_event"] == event_name:
            conn.execute(
                "UPDATE runs SET wake_event = NULL WHERE run_id = ?", (run_id,)
            )
            conn.execute(
                "DELETE FROM waits WHERE run_id = ? AND step_name = ?",
                (run_id, step_name),
            )
            self._write_checkpoint(conn, task_id, step_name, None, run_id, now)
            return False, None

        # 5. First encounter: register the wait and park the run.
        if timeout_secs is None:
            timeout_at = None
            available_at = _INFINITY
        else:
            deadline = now + timedelta(seconds=timeout_secs)
            timeout_at = _fmt(deadline)
            available_at = timeout_at
        conn.execute(
            "INSERT INTO waits (run_id, step_name, task_id, event_name, timeout_at) "
            "VALUES (?, ?, ?, ?, ?) "
            "ON CONFLICT (run_id, step_name) DO UPDATE SET "
            "event_name = excluded.event_name, timeout_at = excluded.timeout_at",
            (run_id, step_name, task_id, event_name, timeout_at),
        )
        conn.execute(
            "UPDATE runs SET state = 'sleeping', wake_event = ?, "
            "available_at = ?, claimed_by = NULL, claim_expires_at = NULL "
            "WHERE run_id = ?",
            (event_name, available_at, run_id),
        )
        conn.execute(
            "UPDATE tasks SET state = 'sleeping' WHERE task_id = ?", (task_id,)
        )
        return True, None

wait_for_event

wait_for_event(*, run_id: str, task_id: str, step_name: str, event_name: str, timeout_secs: int | None = None) -> Any

Ergonomic wrapper over await_event for workflow handlers.

Returns the event payload (or None on timeout) when the run may proceed. When the run is parked, it raises WorkflowSuspended, which unwinds the handler so the worker loop skips completion. Use this inside handlers; use await_event when you need the raw tuple.

Source code in dura/engine.py
def wait_for_event(
    self,
    *,
    run_id: str,
    task_id: str,
    step_name: str,
    event_name: str,
    timeout_secs: int | None = None,
) -> Any:
    """Ergonomic wrapper over ``await_event`` for workflow handlers.

    Returns the event payload (or ``None`` on timeout) when the run may
    proceed. When the run is parked, it raises ``WorkflowSuspended``,
    which unwinds the handler so the worker loop skips completion. Use this
    inside handlers; use ``await_event`` when you need the raw tuple.
    """
    should_suspend, payload = self.await_event(
        run_id=run_id,
        task_id=task_id,
        step_name=step_name,
        event_name=event_name,
        timeout_secs=timeout_secs,
    )
    if should_suspend:
        raise WorkflowSuspended(task_id)
    return payload

get_task

get_task(task_id: str) -> TaskInfo

Return the current state of a task.

Source code in dura/engine.py
def get_task(self, task_id: str) -> TaskInfo:
    """Return the current state of a task."""
    conn = self._conn()
    row = conn.execute(
        "SELECT state, attempts, completed_payload, last_attempt_run "
        "FROM tasks WHERE task_id = ?",
        (task_id,),
    ).fetchone()
    if row is None:
        raise TaskNotFound(f"Task {task_id} not found")

    result = None
    failure_reason = None
    if row["state"] == "completed":
        result = _loads(row["completed_payload"])
    elif row["state"] == "failed" and row["last_attempt_run"]:
        run = conn.execute(
            "SELECT failure_reason FROM runs WHERE run_id = ?",
            (row["last_attempt_run"],),
        ).fetchone()
        if run is not None:
            failure_reason = _loads(run["failure_reason"])

    return TaskInfo(
        task_id=task_id,
        state=row["state"],
        attempts=row["attempts"],
        result=result,
        failure_reason=failure_reason,
    )

ready_run_count

ready_run_count() -> int

Number of runs claimable right now (the work backlog).

Mirrors claim_task's candidacy: runs in pending/sleeping whose task is not terminal and whose available_at has arrived. Useful as a queue-depth gauge.

Source code in dura/engine.py
def ready_run_count(self) -> int:
    """Number of runs claimable right now (the work backlog).

    Mirrors claim_task's candidacy: runs in ``pending``/``sleeping`` whose
    task is not terminal and whose ``available_at`` has arrived. Useful as a
    queue-depth gauge.
    """
    now = _fmt(self._now())
    row = (
        self._conn()
        .execute(
            "SELECT COUNT(*) AS n FROM runs r JOIN tasks t ON t.task_id = r.task_id "
            "WHERE r.state IN ('pending', 'sleeping') "
            "AND t.state IN ('pending', 'sleeping', 'running') "
            "AND r.available_at <= ?",
            (now,),
        )
        .fetchone()
    )
    return row["n"]

task_counts_by_state

task_counts_by_state() -> dict[str, int]

Number of tasks in each state, e.g. {"pending": 3, "failed": 1}.

States with no tasks are omitted rather than reported as zero.

Source code in dura/engine.py
def task_counts_by_state(self) -> dict[str, int]:
    """Number of tasks in each state, e.g. ``{"pending": 3, "failed": 1}``.

    States with no tasks are omitted rather than reported as zero.
    """
    rows = self._conn().execute(
        "SELECT state, COUNT(*) AS n FROM tasks GROUP BY state"
    )
    return {row["state"]: row["n"] for row in rows}

task_counts_by_name_and_state

task_counts_by_name_and_state() -> dict[str, dict[str, int]]

Task counts grouped by task_name, then by state.

e.g. {"send_file": {"pending": 2, "completed": 5}}. Names and states with no matching tasks are omitted.

Source code in dura/engine.py
def task_counts_by_name_and_state(self) -> dict[str, dict[str, int]]:
    """Task counts grouped by ``task_name``, then by ``state``.

    e.g. ``{"send_file": {"pending": 2, "completed": 5}}``. Names and
    states with no matching tasks are omitted.
    """
    rows = self._conn().execute(
        "SELECT task_name, state, COUNT(*) AS n FROM tasks "
        "GROUP BY task_name, state"
    )
    counts: dict[str, dict[str, int]] = {}
    for row in rows:
        counts.setdefault(row["task_name"], {})[row["state"]] = row["n"]
    return counts

recent_failures

recent_failures(*, limit: int = 20) -> list[FailureInfo]

The most recent failed runs, newest first, with their reason.

Includes every failed attempt, not just a task's latest one: a task retried three times and failed each time contributes three entries.

Source code in dura/engine.py
def recent_failures(self, *, limit: int = 20) -> list[FailureInfo]:
    """The most recent failed runs, newest first, with their reason.

    Includes every failed attempt, not just a task's latest one: a task
    retried three times and failed each time contributes three entries.
    """
    rows = self._conn().execute(
        "SELECT r.run_id, r.task_id, t.task_name, r.attempt, r.failed_at, "
        "r.failure_reason FROM runs r JOIN tasks t ON t.task_id = r.task_id "
        "WHERE r.state = 'failed' ORDER BY r.failed_at DESC LIMIT ?",
        (limit,),
    )
    return [
        FailureInfo(
            run_id=row["run_id"],
            task_id=row["task_id"],
            task_name=row["task_name"],
            attempt=row["attempt"],
            failed_at=row["failed_at"],
            failure_reason=_loads(row["failure_reason"]),
        )
        for row in rows
    ]

cleanup

cleanup(ttl: timedelta = timedelta(days=30)) -> int

Delete terminal tasks (and their rows) older than ttl.

Durable state is deliberately left untouched: it is meant to outlive the tasks that wrote it. Returns the number of tasks removed.

Source code in dura/engine.py
def cleanup(self, ttl: timedelta = timedelta(days=30)) -> int:
    """Delete terminal tasks (and their rows) older than ``ttl``.

    Durable ``state`` is deliberately left untouched: it is meant to outlive
    the tasks that wrote it. Returns the number of tasks removed.
    """
    cutoff = _fmt(self._now() - ttl)
    with self._tx() as conn:
        ids = [
            r["task_id"]
            for r in conn.execute(
                "SELECT task_id FROM tasks "
                "WHERE state IN ('completed', 'failed', 'cancelled') "
                "AND terminal_at IS NOT NULL AND terminal_at < ?",
                (cutoff,),
            ).fetchall()
        ]
        for task_id in ids:
            conn.execute("DELETE FROM waits WHERE task_id = ?", (task_id,))
            conn.execute("DELETE FROM checkpoints WHERE task_id = ?", (task_id,))
            conn.execute("DELETE FROM runs WHERE task_id = ?", (task_id,))
            conn.execute("DELETE FROM tasks WHERE task_id = ?", (task_id,))
    # VACUUM cannot run inside a transaction; run it once afterwards.
    if ids:
        self._conn().execute("VACUUM")
    return len(ids)

Data classes

RetryStrategy dataclass

RetryStrategy(*, kind: str = 'none', base_seconds: float = 30.0, factor: float = 2.0, max_seconds: float | None = None, jitter_factor: float = 0.2)

How a failed task should be retried.

  • none - no delay between attempts
  • fixed - always wait base_seconds
  • exponential - base_seconds * factor ** (attempt - 1), capped at max_seconds, or at 100x base_seconds when max_seconds is unset. The delay is always capped: a task may retry forever (leave max_attempts unset), but the wait between attempts may not grow forever.

TaskRef dataclass

TaskRef(*, task_id: str, run_id: str, attempt: int, created: bool)

Returned by DurableEngine.spawn_task.

ClaimedTask dataclass

ClaimedTask(*, run_id: str, task_id: str, attempt: int, name: str, params: Any)

A task handed to a worker by DurableEngine.claim_task.

TaskInfo dataclass

TaskInfo(*, task_id: str, state: str, attempts: int, result: Any = None, failure_reason: Any = None)

Read model returned by DurableEngine.get_task.

FailureInfo dataclass

FailureInfo(*, run_id: str, task_id: str, task_name: str, attempt: int, failed_at: str, failure_reason: Any)

One failed run, as returned by DurableEngine.recent_failures.

Exceptions

EngineError

Bases: Exception

Base class for engine errors.

TaskNotFound

Bases: EngineError

Raised when a referenced task or run does not exist.

TaskCancelledError

Bases: EngineError

Raised when an operation is attempted on a cancelled task.

InvalidRunState

Bases: EngineError

Raised when a run is not in a state that permits the operation.

WorkflowSuspended

Bases: EngineError

Raised by DurableEngine.wait_for_event when a run parks itself.

The worker loop catches this and leaves the run alone: it is now sleeping and will be re-claimed when the awaited event fires or the timeout elapses. A handler must let this propagate (do not catch it).