Skip to content

Python API

Everything public is importable from the top-level package:

from ardiq import Ardiq, Task, Job, TaskResult, TaskInfo

The app: owns the Rust core, the task registry, and the wire codec.

Ardiq(
redis_url: str | None = None,
queue_name: str = "default",
priorities: list[str] | None = None,
*,
serializer: Callable[[Any], bytes] | None = None,
deserializer: Callable[[bytes], Any] | None = None,
cron_poll_s: float = 1.0,
**core_kwargs: Any,
)

Constructor arguments are documented in Configuration (core_kwargs covers concurrency, prefetch, idle_timeout_ms, result_ttl_ms, burst).

PropertyTypeDescription
worker_idstrThis worker’s id (set by the core).
default_prioritystrThe lane a task with no priority lands in.
burstboolRead/write; when True the loop exits once the queue drains.
taskslist[str]Names of the registered tasks.

Decorator that registers a task and returns a Task.

def task(
fn: Callable[..., Any] | None = None,
*,
name: str | None = None,
max_retries: int = 3,
backoff_ms: int = 0,
timeout: float | None = None,
priority: str | None = None,
unique: bool = False,
) -> Task

Usable bare (@app.task) or called (@app.task(max_retries=5)). See Defining tasks.

With unique=True, enqueuing a call identical to one already waiting or running returns that job instead of creating a second one. See Unique tasks.

Register a recurring task and return a Task; see Recurring tasks.

def cron(
spec: str | None = None,
*,
every: timedelta | float | None = None,
name: str | None = None,
max_retries: int = 3,
backoff_ms: int = 0,
timeout: float | None = None,
priority: str | None = None,
) -> Callable[..., Task]

Pass exactly one of spec (a 5-field cron expression, UTC) or every (seconds or a timedelta).

MethodReturnsDescription
await run()NoneStart the worker loop; runs until stop() or (in burst) the queue drains.
stop()NoneAsk the loop to shut down gracefully.
await send(name, *args, **kwargs)JobEnqueue by name, with no local registration; see Enqueuing by name.
ref(name, *, priority=None, unique=False)TaskA handle to a task registered elsewhere — enqueueable, not callable.
await queue_size()intNumber of jobs waiting across lanes.
await result(task_id, timeout=None)TaskResult | NoneFetch a result; with timeout (s) waits, else returns now-or-None.
await status(task_id)strqueued / scheduled / running / complete / not_found.
await info(task_id)TaskInfo | NoneSnapshot of an unfinished task, else None.
await abort(task_id)boolCancel a queued or running task; False if already finished.
await dead_letters(limit=100)list[DeadLetter]Tasks that failed for good, newest first; see Dead letter queue.
await dead_count()intNumber of tasks in the dead letter queue.
await replay(task_id)Job | NoneEnqueue a dead task again with the same id; None if it is not there.
await delete_dead(task_id)boolDrop a dead task without running it; False if it is not there.
lifespan(fn)decoratorRegister worker startup/shutdown; see Shared resources.
on_error(fn)decoratorRegister a failure hook; see Handling failures.
middleware(fn)decoratorWrap every attempt with an async (ctx, call_next); see Middleware.
on_enqueue(fn)decoratorRun a hook as each task is enqueued, to attach headers; see Middleware.
stateStateWorker-scoped resources set by the lifespan hook.

A registered task, returned by @app.task (or by app.ref, without a function). Call it to run inline; use its async methods to dispatch.

Generic as Task[**P, R] over the decorated function’s signature, so enqueue and inline calls are type-checked against it — see Type checking. A ref types as Task[..., Any].

MemberDescription
nameThe registered name.
fnThe underlying function, or None for a ref.
priorityThe task’s default priority lane (or None).
uniqueWhether an identical call already in flight is reused instead of enqueued again.
task(*args, **kwargs)Calling the Task runs fn inline, bypassing the queue. A ref raises TypeError.
await enqueue(*args, **kwargs)Dispatch to a worker; returns a Job.
options(...)Returns a bound task with per-call overrides; see below.
def options(
*,
task_id: str | None = None,
priority: str | None = None,
delay_ms: int = 0,
schedule_ms: int = 0,
expire_ms: int = 0,
unique: bool | None = None,
) -> _BoundTask

Returns an object with the same await enqueue(*args, **kwargs) method, carrying the overrides. See Enqueuing & scheduling.

await add.options(priority="high", delay_ms=5000).enqueue(2, 3)
def prepare(*args, **kwargs) -> PreparedTask

The call .enqueue would make, held back so enqueue_many can send a batch in one round trip. Available on .options(...) too. Arguments are checked against the task’s signature exactly as .enqueue checks them.

async def enqueue_many(tasks: Iterable[PreparedTask]) -> list[Job]

Sends prepared tasks in one round trip and returns their Jobs in order. Any iterable works, including a generator. Priorities are validated across the whole batch before anything is sent. See Enqueuing in bulk.

jobs = await queue.enqueue_many(charge.prepare(oid) for oid in order_ids)

An immutable handle to an enqueued task — just the app plus an id.

MemberReturnsDescription
appArdiqThe owning app.
idstrThe job id.
await result(timeout=None)TaskResult | NoneFetch the result; with timeout (s) waits, raising TimeoutError.
await status()strCurrent status.
await info()TaskInfo | NoneSnapshot if unfinished, else None.
await abort()boolCancel the task; False if already finished. See Aborting tasks.

A NamedTuple describing a finished task.

FieldTypeDescription
successboolWhether the task returned (vs failed after retries).
valueAnyReturn value on success; error repr on failure.
triesintNumber of attempts.
enqueue_timeintEpoch ms when enqueued.
startintEpoch ms when execution started.
finishintEpoch ms when execution finished.
duration_msintProperty: finish - start.
abortedboolProperty: the task was cancelled rather than failing on its own.

A namespace holding worker-scoped resources, reachable as app.state. Populated by the @app.lifespan hook — either from the mapping it yields or by direct assignment. Reading an attribute that was never set raises AttributeError naming it.

app.state.db # set by the lifespan hook

A NamedTuple snapshot of an unfinished task (queued, scheduled, or running).

FieldTypeDescription
task_idstrThe job id.
fn_namestrRegistered task name.
argstuplePositional arguments.
kwargsdictKeyword arguments.
enqueue_timeintEpoch ms when enqueued.
triesintAttempts so far.
statusstrCurrent status.
scheduled_atint | NoneEpoch ms if waiting in the delayed set, else None.

A NamedTuple for a task in the dead letter queue.

FieldTypeDescription
task_idstrThe job id, kept by a replay.
fn_namestrRegistered task name.
argstuplePositional arguments.
kwargsdictKeyword arguments.
prioritystrThe lane it ran in, and the one a replay uses.
errorstrThe failure, as its TaskResult.value reported it.
triesintAttempts made before it gave up.
enqueue_timeintEpoch ms when enqueued.
failed_atintEpoch ms when it failed for good.

A NamedTuple handed to an @app.middleware for each attempt. See Middleware.

FieldTypeDescription
task_idstrThe job id.
namestrThe task’s registered name.
triesintThe attempt being made, counting from 1.
argstuplePositional arguments.
kwargsdictKeyword arguments.
headersdictHeaders attached at enqueue; empty if none.

A NamedTuple handed to an @app.on_enqueue hook for each task being enqueued.

FieldTypeDescription
task_idstrThe id the task will have.
namestrThe task’s name.
argstuplePositional arguments.
kwargsdictKeyword arguments.
headersdictFill it to send entries with the task.

A NamedTuple handed to every @app.on_error hook; see Handling failures.

FieldTypeDescription
namestrThe task’s registered name.
task_idstrThe job id.
excBaseExceptionThe exception the attempt raised.
triesintThe attempt that just failed, counting from 1.
will_retryboolWhether another attempt is coming.

A NamedTuple describing the task running right now, returned by current_task() — None outside a worker. See Knowing which task you are.

FieldTypeDescription
task_idstrThe job id.
namestrThe task’s registered name.
triesintThe attempt in progress, counting from 1.

BrokerError → ArdiqError → RuntimeError. BrokerError means Redis was unreachable; ArdiqError is anything else the core raises. See When the broker itself fails.

An exception a task raises to run again, optionally after a delay of its choosing. It respects max_retries; see Retrying on demand.

raise Retry("rate limited", delay_ms=30_000)
ArgumentTypeDefaultDescription
messagestr"retry requested"Why, recorded as the error if the retries run out.
delay_msint | NoneNoneWait this long before the next attempt; None uses the task’s backoff.

inline(app) is an async context manager. While it is open, app runs each task as soon as it is enqueued, in memory, with no Redis. See Testing.

async with inline(app):
job = await add.enqueue(2, 3)

Made bytay.dev