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
close ¶
Close the connection for the calling thread, if any.
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
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
complete_run ¶
Mark a run and its task as completed.
Source code in dura/engine.py
fail_run ¶
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
cancel_task ¶
Cancel a task and any of its non-terminal runs.
Source code in dura/engine.py
cancel_duplicate_tasks ¶
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
extend_claim ¶
Push a running lease forward (heartbeat for long steps).
Source code in dura/engine.py
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
get_checkpoint ¶
Return a stored checkpoint payload, or None if absent.
Source code in dura/engine.py
get_state ¶
Return the value stored at (namespace, key), or default.
Source code in dura/engine.py
has_state ¶
Whether any key exists under namespace (cheap emptiness check).
Source code in dura/engine.py
set_state ¶
Store value at (namespace, key) (upsert, last write wins).
Source code in dura/engine.py
set_state_many ¶
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
delete_state ¶
Delete (namespace, key). Returns True if it existed.
Source code in dura/engine.py
list_state ¶
Return all key -> value pairs in namespace (ordered by key).
Source code in dura/engine.py
update_state ¶
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
emit_event ¶
Emit a named signal (first write wins) and wake its waiters.
Source code in dura/engine.py
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
874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 | |
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
get_task ¶
get_task(task_id: str) -> TaskInfo
Return the current state of a task.
Source code in dura/engine.py
ready_run_count ¶
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
task_counts_by_state ¶
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
task_counts_by_name_and_state ¶
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
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
cleanup ¶
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
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 attemptsfixed- always waitbase_secondsexponential-base_seconds * factor ** (attempt - 1), capped atmax_seconds, or at 100xbase_secondswhenmax_secondsis unset. The delay is always capped: a task may retry forever (leavemax_attemptsunset), but the wait between attempts may not grow forever.
TaskRef
dataclass
¶
Returned by DurableEngine.spawn_task.
ClaimedTask
dataclass
¶
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).