BlitzQ
API Reference

results

API reference for results.

Handles to enqueued tasks and result retrieval.

class Queue

A BlitzQ application.

Queue holds configuration and the task registry, and provides client operations (enqueue, results, inspection, cancellation). The same object is imported by web processes, which only publish, and by workers started with blitzq worker module:attribute.

Parameters

  • name str — Default queue for tasks without an explicit queue or routing rule.
  • redis_url str | None — Redis connection URL (rediss:// for TLS). Defaults to the BLITZQ_REDIS_URL environment variable, then redis://localhost:6379/0.
  • mode Literal['fast', 'reliable'] — "reliable" (default; Redis Streams, at-least-once, crash recovery) or "fast" (Redis lists, at-most-once, lowest overhead).
  • broker Broker | None — A custom :class:~blitzq.broker.base.Broker instead of Redis.
  • store_results bool — Store final state and return values (can be overridden per task).
  • result_ttl int — Seconds to keep task records.
  • track_state bool — Also record queued/scheduled/running/retrying transitions. Costs one extra write per transition.
  • routes Routes — Task-name glob patterns to queue names, or a callable.
  • visibility_timeout float — Reliable mode: seconds after which an unacknowledged message whose worker stopped renewing its lease is redelivered.
  • max_deliveries int — Reliable mode: dead-letter a message after this many deliveries (protects against messages that crash workers).

__init__(name: str = 'default', redis_url: str | None = None, mode: Literal['fast', 'reliable'] = 'reliable', broker: Broker | None = None, namespace: str = 'blitzq', serializer: Serializer | None = None, store_results: bool = True, result_ttl: int = 24 * 3600, track_state: bool = False, routes: Routes = None, default_retries: int = 0, default_retry_policy: RetryPolicy | None = None, default_timeout: float | None = None, visibility_timeout: float = 60.0, max_deliveries: int = 5, dead_letter_max: int = 100000, revoke_ttl: int = 3600, redis_options: dict[str, Any] | None = None) -> None

add_sync_wrapper(wrapper: SyncWrapper) -> None

Wrap every sync-task call inside its executor thread.

wrapper(call) must invoke call() and return its result. Used by integrations that need thread-local setup, such as Django connection management or a Flask application context.

cancel(task_id: str) -> bool

Cancel a task.

Returns True when the task was still scheduled (delayed or waiting for a retry) and has been removed - cancellation is then certain. Otherwise a revocation is recorded and False is returned: workers skip the task if they receive it after their next revocation sync (about one second) and cancel it if it is a running async task. Running thread/process tasks cannot be interrupted, and a task that already finished is unaffected.

cancel_sync(task_id: str) -> bool

close() -> None

Close async broker connections. Safe to call more than once.

close_sync() -> None

Close connections used by the synchronous API and stop its thread.

connect() -> None

Eagerly connect the async broker (optional; connections are lazy).

dead_letters(limit: int = 100, offset: int = 0) -> list[DeadLetter]

get_result(task_id: str, timeout: float | None = None) -> Any

Wait for a task's final state and return its result.

Waits on the broker's push notification when it has one (Redis: pub/sub on the task's record channel), falling back to a coarse poll otherwise. Requires result storage for the task. Raises TaskFailed for failed, dead-lettered or cancelled tasks and ResultTimeout after timeout seconds (None waits indefinitely).

get_result_sync(task_id: str, timeout: float | None = None) -> Any

inspect(task_id: str) -> TaskInfo | None

Everything known about a task, or None.

None means the task is queued or running without state tracking, its record expired, or the id is unknown.

inspect_sync(task_id: str) -> TaskInfo | None

new_task_id() -> str

on_shutdown(fn: Hook) -> Hook

Register a worker/scheduler shutdown hook (sync or async).

on_startup(fn: Hook) -> Hook

Register a worker/scheduler startup hook (sync or async).

periodic(schedule: Schedule | str | float | timedelta, tz: str | None = None, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, missed: MissedPolicy = 'run_once', name: str | None = None, queue: str | None = None, task_options: Any = {}) -> Callable[[Callable[..., Any]], Task[Any, Any]]

Register a periodic task (dispatched by blitzq scheduler).

schedule is a cron string ("*/5 * * * *", evaluated in tz), an interval in seconds or a timedelta, or a :class:Schedule. missed controls what happens to occurrences missed while no scheduler was running: "run_once" dispatches only the latest, "run_all" dispatches each (up to 100), "skip" dispatches none.

purge(queue: str | None = None) -> int

Delete waiting messages from queue (default queue if omitted).

queue_stats(queues: Sequence[str] | None = None) -> list[QueueStats]

retry(task_id: str) -> bool

Re-enqueue a dead-lettered task with a fresh attempt budget.

Returns False if no dead letter with this id exists.

retry_sync(task_id: str) -> bool

send(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]

Enqueue a task by name, without importing its code (producer-only apps).

send_sync(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]

status(task_id: str) -> TaskState | None

status_sync(task_id: str) -> TaskState | None

task(fn: Callable[..., Any] | None = None, name: str | None = None, queue: str | None = None, retries: int | None = None, retry_policy: RetryPolicy | None = None, timeout: float | None = None, executor: Executor | None = None, store_result: bool | None = None, dead_letter: bool = True, priority: Priority = 'normal', rate_limit: RateLimit | str | None = None) -> Any

Register a task.

retries is the number of retries after the first attempt. executor defaults to "async" for coroutine functions and "thread" for regular functions; use "process" for CPU-bound functions (they must be importable module-level functions). dead_letter=False records terminal failures as failed instead of moving them to the dead-letter store.

priority ("high"/"normal"/"low") is this task's default priority within its queue; override per call with task.options(priority=...). Priority levels are physically separate broker queues that every worker checks in order (high, then normal, then low) while sharing the queue's overall concurrency

  • not a separate queue you need to remember to subscribe a worker to. See docs/architecture.md#task-priority.

rate_limit caps how often this task starts execution, shared across every worker (a Redis-backed token bucket, not a per-process counter): "10/s", "100/m" or "1000/hour", or a :class:~blitzq.RateLimit. A task over its limit is not executed and not counted as a retry; it is rescheduled for when a slot should be free. See docs/architecture.md#rate-limiting.

class TaskHandle

Reference to an enqueued task.

Async methods (result, info, status, cancel) must be awaited inside an event loop; *_sync variants block the calling thread and must not be used inside a running event loop.

__init__(task_id: str, app: Queue, queue: str) -> None

cancel() -> bool

cancel_sync() -> bool

info() -> TaskInfo | None

result(timeout: float | None = None) -> R

Wait for the task to finish and return its result.

Raises :class:~blitzq.exceptions.TaskFailed if it failed, was dead-lettered or cancelled, and :class:~blitzq.exceptions.ResultTimeout if no final state is visible within timeout seconds.

result_sync(timeout: float | None = None) -> R

status() -> TaskState | None

status_sync() -> TaskState | None

class TaskInfo

A snapshot of what BlitzQ knows about a task.

Records are written by producers (queued/scheduled when state tracking is enabled) and by workers (running when tracking is enabled, and the final state when results are stored). Readers may observe a slightly stale state during concurrent transitions; see docs/delivery_guarantees.md.

class TaskState

unwrap_result(info: TaskInfo) -> Any

Return info.result or raise TaskFailed for unsuccessful final states.

On this page

class Queue__init__(name: str = 'default', redis_url: str | None = None, mode: Literal['fast', 'reliable'] = 'reliable', broker: Broker | None = None, namespace: str = 'blitzq', serializer: Serializer | None = None, store_results: bool = True, result_ttl: int = 24 * 3600, track_state: bool = False, routes: Routes = None, default_retries: int = 0, default_retry_policy: RetryPolicy | None = None, default_timeout: float | None = None, visibility_timeout: float = 60.0, max_deliveries: int = 5, dead_letter_max: int = 100000, revoke_ttl: int = 3600, redis_options: dict[str, Any] | None = None) -> Noneadd_sync_wrapper(wrapper: SyncWrapper) -> Nonecancel(task_id: str) -> boolcancel_sync(task_id: str) -> boolclose() -> Noneclose_sync() -> Noneconnect() -> Nonedead_letters(limit: int = 100, offset: int = 0) -> list[DeadLetter]get_result(task_id: str, timeout: float | None = None) -> Anyget_result_sync(task_id: str, timeout: float | None = None) -> Anyinspect(task_id: str) -> TaskInfo | Noneinspect_sync(task_id: str) -> TaskInfo | Nonenew_task_id() -> stron_shutdown(fn: Hook) -> Hookon_startup(fn: Hook) -> Hookperiodic(schedule: Schedule | str | float | timedelta, tz: str | None = None, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, missed: MissedPolicy = 'run_once', name: str | None = None, queue: str | None = None, task_options: Any = {}) -> Callable[[Callable[..., Any]], Task[Any, Any]]purge(queue: str | None = None) -> intqueue_stats(queues: Sequence[str] | None = None) -> list[QueueStats]retry(task_id: str) -> boolretry_sync(task_id: str) -> boolsend(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]send_sync(task_name: str, args: Sequence[Any] = (), kwargs: dict[str, Any] | None = None, options: Any = {}) -> TaskHandle[Any]status(task_id: str) -> TaskState | Nonestatus_sync(task_id: str) -> TaskState | Nonetask(fn: Callable[..., Any] | None = None, name: str | None = None, queue: str | None = None, retries: int | None = None, retry_policy: RetryPolicy | None = None, timeout: float | None = None, executor: Executor | None = None, store_result: bool | None = None, dead_letter: bool = True, priority: Priority = 'normal', rate_limit: RateLimit | str | None = None) -> Anyclass TaskHandle__init__(task_id: str, app: Queue, queue: str) -> Nonecancel() -> boolcancel_sync() -> boolinfo() -> TaskInfo | Noneresult(timeout: float | None = None) -> Rresult_sync(timeout: float | None = None) -> Rstatus() -> TaskState | Nonestatus_sync() -> TaskState | Noneclass TaskInfoclass TaskStateunwrap_result(info: TaskInfo) -> Any