BlitzQ
API Reference

client

API reference for client.

The application object: task registry, configuration and client operations.

class Broker

Abstract broker. See module docstring.

cancel(task_id: str, revoke_ttl: int) -> bool

Cancel a task.

Returns True if the task was scheduled and has been removed (cancellation is certain). Otherwise records a revocation that workers observe on their next revocation sync and returns False.

check_rate_limit(key: str, rate: float, capacity: float, now: float) -> float

Try to take one token from key's bucket (capacity tokens, refilling at rate tokens/second, shared across every caller).

Returns 0.0 if a token was taken (proceed), otherwise the number of seconds to wait before a token will be available (a token is not reserved for that wait; a task may need to check again after waiting if another caller took it first).

claim_periodic(name: str, occurrence: float, queue: str, data: bytes | None) -> bool

Atomically record occurrence as dispatched and publish data.

Succeeds only if no occurrence >= occurrence was recorded before, so any number of scheduler instances dispatch each occurrence at most once. data=None records the occurrence without publishing (used by the skip missed-run policy).

clone() -> Broker

Return a new, unconnected broker with identical configuration.

Used to give a separate event loop (such as the synchronous API's background thread) its own connections.

close() -> None

complete(completions: Sequence[Completion]) -> None

Apply a batch of completions (acks, records, retries, dead letters).

connect() -> None

Open connections eagerly. Brokers may also connect lazily.

dead_letter_count() -> int

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

Dead-letter records, newest first.

fetch(queue: str, count: int, timeout: float, consumer: str) -> list[Delivery]

Return up to count messages from queue.

timeout == 0 means do not block. Otherwise block up to timeout seconds waiting for at least one message.

get_dead_letter(task_id: str) -> bytes | None

get_record(task_id: str) -> bytes | None

heartbeat(queue: str, deliveries: Sequence[Delivery], consumer: str) -> list[Delivery]

Extend the lease on in-flight deliveries.

Returns the deliveries whose ownership was lost (another consumer recovered them). Brokers without leases return an empty list.

is_scheduled(task_id: str) -> float | None

Return the ETA if task_id is currently scheduled.

known_queues() -> set[str]

periodic_last() -> dict[str, float]

ping() -> bool

prepare_consumer(queues: Sequence[str], consumer: str) -> None

Prepare broker-side structures before fetching (e.g. consumer groups).

promote_due(now: float, limit: int) -> tuple[int, float | None]

Move due scheduled messages into their queues.

Returns (promoted_count, next_eta) where next_eta is the due time of the earliest remaining scheduled message, if any. Must be safe to call concurrently from many processes.

publish(requests: Sequence[PublishRequest]) -> None

Publish messages; those with eta in the future go to the schedule.

purge_dead_letters() -> int

purge_queue(queue: str) -> int

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

recover(queue: str, consumer: str, idle_timeout: float, count: int) -> list[Delivery]

Claim messages abandoned by other consumers for longer than idle_timeout.

register_worker(worker_id: str, info: bytes, ttl: int) -> None

replay_dead_letter(task_id: str, queue: str, data: bytes) -> bool

Atomically remove a dead letter, delete the task's stale record and publish data to queue.

Returns False if the dead letter no longer exists (for example it was replayed concurrently), in which case nothing is published.

revoked() -> set[str]

Currently active revocations.

scheduled(limit: int = 100, offset: int = 0) -> list[ScheduledEntry]

scheduled_count() -> int

set_record(task_id: str, data: bytes, ttl: int | None) -> None

unregister_worker(worker_id: str) -> None

wait_for_record(task_id: str, timeout: float) -> None

Block up to timeout seconds for a push notification that task_id's record or dead-letter entry may have changed.

Best effort, not exact: may return early with nothing to see (the caller must re-check the record), and may time out even though a change happened concurrently (for example a dropped pub/sub message during a reconnect). get_result() relies on this only to avoid polling tightly; it always re-reads the record itself afterwards, so a broker that never notifies (the default here) just falls back to coarser polling instead of being wrong.

workers() -> dict[str, bytes]

class ConfigurationError

Invalid configuration or API misuse.

class DeadLetter

A terminally failed message kept for inspection and manual replay.

class PeriodicTask

__init__(name: str, task: Task[Any, Any], schedule: Schedule, args: tuple[Any, ...] = (), kwargs: dict[str, Any] = dict(), missed: MissedPolicy = 'run_once', queue: str | None = None) -> None

class Portal

__init__(name: str = 'blitzq-sync') -> None

call(fn: Callable[..., Coroutine[Any, Any, T]], args: Any = (), what: str = 'this') -> T

stop(cleanup: Callable[[], Coroutine[Any, Any, Any]] | None = None) -> None

class PublishRequest

__init__(queue: str, task_id: str, data: bytes, eta: float | None = None, record: bytes | None = None, record_ttl: int | None = None) -> None

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 QueueStats

__init__(name: str, waiting: int = 0, in_progress: int | None = None, extra: dict[str, Any] = dict()) -> None

class RateLimit

A cap of count task starts per period_seconds, enforced globally.

The bucket holds up to count tokens (one burst's worth) and refills continuously at count / period_seconds tokens/second, so a limit of "100/m" allows a burst of 100 immediately after being idle, then settles to roughly one every 0.6s - it does not release exactly 100 once a minute.

__init__(count: float, period_seconds: float) -> None

parse(spec: str) -> RateLimit

Parse "100", "100/s", "100/m" or "100/hour".

class ResultTimeout

No final result became available within the requested wait time.

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

class RetryPolicy

Exponential backoff with optional jitter.

The delay before retry n (n = 1 for the first retry) is min(max_delay, initial_delay * backoff ** (n - 1)). With jitter enabled the delay is drawn uniformly from [delay / 2, delay] ("equal jitter"), which spreads retries of simultaneously failing tasks while keeping a guaranteed minimum wait.

retry_on lists exception types that trigger an automatic retry; dont_retry_on takes precedence and marks exceptions as permanent failures. An explicit :class:~blitzq.Retry raised by the task is always honoured while attempts remain.

__init__(initial_delay: float = 1.0, max_delay: float = 300.0, backoff: float = 2.0, jitter: bool = True, retry_on: tuple[type[BaseException], ...] = (Exception,), dont_retry_on: tuple[type[BaseException], ...] = ()) -> None

compute_delay(retry_number: int, rng: random.Random | None = None) -> float

Delay in seconds before retry number retry_number (1-based).

is_retryable(exc: BaseException) -> bool

class Router

Resolves the queue for a task name.

Precedence (highest first):

  1. queue= passed to task.options(...) at enqueue time
  2. queue= given to the @app.task decorator
  3. routing rules configured on the application (this class)
  4. the application's default queue

Rules are either a mapping of glob patterns to queue names, evaluated in insertion order (\{"app.images.*": "images"\}), or a callable returning a queue name or None.

__init__(routes: Routes, default: str) -> None

resolve(task_name: str, decorator_queue: str | None, call_queue: str | None) -> str

route(task_name: str) -> str | None

Queue selected by routing rules, or None if no rule matches.

class Schedule

next_after(ts: float) -> float

First occurrence strictly after epoch timestamp ts.

occurrences(after: float, until: float, limit: int = 10000) -> list[float]

Occurrences in (after, until], oldest first, at most limit (the latest kept).

class Serializer

Encodes envelopes, task records and results.

Parameters

  • format Literal['msgpack', 'json'] — "msgpack" (compact, default) or "json" (human-readable in Redis).
  • enc_hook Callable[[Any], Any] | None — Optional msgspec encode hook to convert otherwise unsupported argument types into supported ones (for example lambda o: o.to_dict()).
  • max_message_size int — Upper bound for an encoded message, enforced on both enqueue and receipt.

__init__(format: Literal['msgpack', 'json'] = 'msgpack', enc_hook: Callable[[Any], Any] | None = None, max_message_size: int = DEFAULT_MAX_MESSAGE_SIZE) -> None

check_value(value: Any) -> None

Raise SerializationError if value cannot be stored as a result.

decode_dead(data: bytes) -> DeadLetter

decode_envelope(data: bytes) -> Envelope

decode_info(data: bytes) -> TaskInfo

encode_dead(dead: DeadLetter) -> bytes

encode_envelope(env: Envelope) -> bytes

encode_info(info: TaskInfo) -> bytes

class Task

A function registered with a :class:~blitzq.Queue.

Calling the object runs the function directly in the current process. enqueue publishes it for a worker. For async def functions R is the awaited return type.

__init__(app: Queue, fn: Callable[P, Any], options: TaskOptions) -> None

enqueue(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]

Publish the task. Non-blocking for the event loop; returns once Redis accepted it.

enqueue_many(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]

Publish many calls in one pipelined round-trip.

Each item is a tuple of positional arguments.

enqueue_many_sync(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]

enqueue_sync(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]

Blocking variant of :meth:enqueue for synchronous code.

options(queue: str | None = None, delay: float | timedelta | None = None, eta: datetime | float | None = None, task_id: str | None = None, correlation_id: str | None = None, headers: Mapping[str, str] | None = None, timeout: float | None = None, priority: Priority | None = None) -> BoundTask[P, R]

Return a view of this task with per-call options applied.

delay (seconds or timedelta) and eta (aware datetime or epoch seconds) schedule the task for later. task_id sets an explicit id; enqueueing the same id twice creates two messages unless the first is still scheduled, in which case its schedule entry is replaced. priority ("high"/"normal"/"low") reorders this call within its queue's shared concurrency budget; see the priority parameter of @queue.task for how workers pick it up.

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 TaskOptions

__init__(name: str, queue: str | None, max_attempts: int, retry_policy: RetryPolicy, timeout: float | None, executor: Executor, store_result: bool, dead_letter: bool, priority: Priority, rate_limit: RateLimit | None) -> None

class TaskState

as_rate_limit(value: RateLimit | str | None) -> RateLimit | None

as_schedule(value: Schedule | str | float | timedelta, tz: str | None = None) -> Schedule

new_task_id() -> str

redis_broker(url: str, mode: Mode = 'reliable', options: Any = {}) -> Broker

Create a Redis broker for mode ("fast" or "reliable").

unwrap_result(info: TaskInfo) -> Any

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

On this page

class Brokercancel(task_id: str, revoke_ttl: int) -> boolcheck_rate_limit(key: str, rate: float, capacity: float, now: float) -> floatclaim_periodic(name: str, occurrence: float, queue: str, data: bytes | None) -> boolclone() -> Brokerclose() -> Nonecomplete(completions: Sequence[Completion]) -> Noneconnect() -> Nonedead_letter_count() -> intdead_letters(limit: int = 100, offset: int = 0) -> list[bytes]fetch(queue: str, count: int, timeout: float, consumer: str) -> list[Delivery]get_dead_letter(task_id: str) -> bytes | Noneget_record(task_id: str) -> bytes | Noneheartbeat(queue: str, deliveries: Sequence[Delivery], consumer: str) -> list[Delivery]is_scheduled(task_id: str) -> float | Noneknown_queues() -> set[str]periodic_last() -> dict[str, float]ping() -> boolprepare_consumer(queues: Sequence[str], consumer: str) -> Nonepromote_due(now: float, limit: int) -> tuple[int, float | None]publish(requests: Sequence[PublishRequest]) -> Nonepurge_dead_letters() -> intpurge_queue(queue: str) -> intqueue_stats(queues: Sequence[str]) -> list[QueueStats]recover(queue: str, consumer: str, idle_timeout: float, count: int) -> list[Delivery]register_worker(worker_id: str, info: bytes, ttl: int) -> Nonereplay_dead_letter(task_id: str, queue: str, data: bytes) -> boolrevoked() -> set[str]scheduled(limit: int = 100, offset: int = 0) -> list[ScheduledEntry]scheduled_count() -> intset_record(task_id: str, data: bytes, ttl: int | None) -> Noneunregister_worker(worker_id: str) -> Nonewait_for_record(task_id: str, timeout: float) -> Noneworkers() -> dict[str, bytes]class ConfigurationErrorclass DeadLetterclass PeriodicTask__init__(name: str, task: Task[Any, Any], schedule: Schedule, args: tuple[Any, ...] = (), kwargs: dict[str, Any] = dict(), missed: MissedPolicy = 'run_once', queue: str | None = None) -> Noneclass Portal__init__(name: str = 'blitzq-sync') -> Nonecall(fn: Callable[..., Coroutine[Any, Any, T]], args: Any = (), what: str = 'this') -> Tstop(cleanup: Callable[[], Coroutine[Any, Any, Any]] | None = None) -> Noneclass PublishRequest__init__(queue: str, task_id: str, data: bytes, eta: float | None = None, record: bytes | None = None, record_ttl: int | None = None) -> Noneclass 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 QueueStats__init__(name: str, waiting: int = 0, in_progress: int | None = None, extra: dict[str, Any] = dict()) -> Noneclass RateLimit__init__(count: float, period_seconds: float) -> Noneparse(spec: str) -> RateLimitclass ResultTimeout__init__(task_id: str, timeout: float | None, last_state: Any = None) -> Noneclass RetryPolicy__init__(initial_delay: float = 1.0, max_delay: float = 300.0, backoff: float = 2.0, jitter: bool = True, retry_on: tuple[type[BaseException], ...] = (Exception,), dont_retry_on: tuple[type[BaseException], ...] = ()) -> Nonecompute_delay(retry_number: int, rng: random.Random | None = None) -> floatis_retryable(exc: BaseException) -> boolclass Router__init__(routes: Routes, default: str) -> Noneresolve(task_name: str, decorator_queue: str | None, call_queue: str | None) -> strroute(task_name: str) -> str | Noneclass Schedulenext_after(ts: float) -> floatoccurrences(after: float, until: float, limit: int = 10000) -> list[float]class Serializer__init__(format: Literal['msgpack', 'json'] = 'msgpack', enc_hook: Callable[[Any], Any] | None = None, max_message_size: int = DEFAULT_MAX_MESSAGE_SIZE) -> Nonecheck_value(value: Any) -> Nonedecode_dead(data: bytes) -> DeadLetterdecode_envelope(data: bytes) -> Envelopedecode_info(data: bytes) -> TaskInfoencode_dead(dead: DeadLetter) -> bytesencode_envelope(env: Envelope) -> bytesencode_info(info: TaskInfo) -> bytesclass Task__init__(app: Queue, fn: Callable[P, Any], options: TaskOptions) -> Noneenqueue(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]enqueue_many(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]enqueue_many_sync(items: Iterable[tuple[Any, ...]]) -> list[TaskHandle[R]]enqueue_sync(args: P.args = (), kwargs: P.kwargs = {}) -> TaskHandle[R]options(queue: str | None = None, delay: float | timedelta | None = None, eta: datetime | float | None = None, task_id: str | None = None, correlation_id: str | None = None, headers: Mapping[str, str] | None = None, timeout: float | None = None, priority: Priority | None = None) -> BoundTask[P, R]class 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 TaskOptions__init__(name: str, queue: str | None, max_attempts: int, retry_policy: RetryPolicy, timeout: float | None, executor: Executor, store_result: bool, dead_letter: bool, priority: Priority, rate_limit: RateLimit | None) -> Noneclass TaskStateas_rate_limit(value: RateLimit | str | None) -> RateLimit | Noneas_schedule(value: Schedule | str | float | timedelta, tz: str | None = None) -> Schedulenew_task_id() -> strredis_broker(url: str, mode: Mode = 'reliable', options: Any = {}) -> Brokerunwrap_result(info: TaskInfo) -> Any