toro 🐂
An async-first, Redis-backed job queue for Python. Every state transition is
an atomic Lua script; producing and processing are asyncio end to end.
pip install toro-queue # the import name is `toro`
Installed as
toro-queueon PyPI (the nametorowas taken), but youimport toro. See the docs for the architecture, the reliability model, and the detailed guides.
Pairs with matador, a live web dashboard for your queues.
Why toro
- Async-native. Enqueue and process with
async/await- no thread pools, no sync bridge. A natural fit for FastAPI, aiohttp, or any asyncio app. - Atomic by construction. Claims, retries, promotions and finishes are Lua scripts, so a job can't be lost or double-committed between two round trips.
- At-least-once delivery. Per-job locks + a background mark-and-sweep recover jobs from workers that crashed - without the visibility-timeout double-delivery trap of some other queues.
- Typed. Ships
py.typed; the public API is fully annotated.
Features
| Enqueue | delayed jobs, global priorities (FIFO within a band) |
| Retries | fixed or exponential backoff, capped attempts |
| Schedules | repeatable cron and fixed-interval (every) jobs |
| Flows | parent/child job trees: fan-out/fan-in, failure policies, flow-aware retry |
| Rate limiting | queue-wide token bucket shared across all workers |
| Global concurrency | one cap on jobs active at once, across every worker process |
| Dedup | custom (idempotent) job ids + a throttle window ({id, ttl}) |
| Auto-removal | keep the last N and/or finished-within-age completed/failed |
| Reliability | per-job locks, lock renewal, stalled-job recovery |
| Observability | progress, per-job logs, lifecycle events, await result() |
| Lifecycle | pause / resume, graceful shutdown that drains in-flight jobs |
| Dashboard | matador - a live web UI |
Quick start
import asyncio
from toro import Queue, Worker
async def main():
queue = Queue("emails")
await queue.add("welcome", {"to": "ada@example.com"})
async def process(job):
print("sending", job.data)
return {"ok": True}
worker = Worker("emails", process, concurrency=8)
worker.on("completed", lambda job, result: print("done", job.id))
await worker.run()
asyncio.run(main())
A taste of the options
# Priorities, delay, and retry-with-backoff
await queue.add("report", data, priority=10, delay=5000,
attempts=5, backoff={"type": "exponential", "delay": 1000})
# Idempotent custom id (a second add with the same id is ignored)
await queue.add("charge", data, job_id="order-1234")
# A repeatable schedule (cron or every-N-ms); "run now" with trigger_scheduler
await queue.add_scheduler("nightly-rollup", cron="0 0 * * *")
# A flow: children run first (fan-out), the parent runs on their results (fan-in)
from toro import FlowChild as c
report = await queue.add_flow("report", {"q": 3},
children=[c("fetch", {"shard": i}) for i in range(3)])
# Queue-wide rate limit: at most 100 jobs / second across every worker
worker = Worker("emails", process, rate_limit={"max": 100, "duration": 1000})
# Wait for a result from the producer side
job = await queue.add("resize", {"src": "a.png"})
print(await job.result(timeout=30))
Flows
A flow enqueues a parent and its children as one atomic tree. The children run first (fan-out, nested arbitrarily); the parent parks until every child has settled, then runs and reads their results (fan-in). One primitive covers fan-out/fan-in and chained steps, with per-child failure policies and flow-aware retry that recovers a whole failed flow in one shot. Full guide: docs/flows.md.
Develop
Managed with uv; the Astral toolchain throughout.
uv sync # venv + deps + dev group
uv run ruff check . # lint (strict: select = ALL)
uv run ruff format . # format
uv run ty check # type check
uv run pytest -m "unit or integration" # tests (integration needs Redis on :6379)
uv run python examples/basic.py
The suite is a pyramid - -m unit (fast, no Redis), -m integration (Redis),
and -m load (the open-loop benchmark harness in tests/load/).
License
Release files for toro-queue 0.6.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| toro_queue-0.6.0.tar.gz | 159.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| toro_queue-0.6.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 209.7 kB
Release files / toro_queue-0.6.0.tar.gz
| Download URL | toro_queue-0.6.0.tar.gz |
|---|---|
| Size | 159.7 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
66a7cedae0a499cadd2c8ee52718875deadc489bd112271152e214c1047e0924
|
|
BLAKE2b-256 checksum How to use checksums |
c586480f1d46eec3f0165aa45a5e597c6799a1a1cf1de59302a16d4b9359d79e
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 20, 2026.
Transparency logRelease files / toro_queue-0.6.0-py3-none-any.whl
| Download URL | toro_queue-0.6.0-py3-none-any.whl |
|---|---|
| Size | 50.0 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
9808b387a71fc902512149be4c3cbbefd4696cb7234156f2ac6352a1851efaa2
|
|
BLAKE2b-256 checksum How to use checksums |
c4a1a10edda50748f20e54f7be8ebdc79cb39b8cfe9a48c2a3660e08c4933c35
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 20, 2026.
Transparency log