Structured concurrency and parallelism for asyncio, threading, and multiprocessing using queues.
Project description
osiiso
Structured task queues for Python across asyncio, threads, and processes.
osiiso gives you one compact queue API for three execution backends:
AsyncQueuefor coroutine-heavy I/O and async integrations (SharedQueueadds thread-safe submission onto one shared event loop).ThreadQueuefor blocking I/O, synchronous SDKs, filesystem work, and SQLite writes.ProcessQueuefor CPU-heavy work that benefits from separate subprocesses.
It is dependency-free at runtime, typed with py.typed, and built around a predictable workflow: submit tasks, apply options, run the queue, then inspect handles and a RunSummary.
Contents
- Why osiiso
- Installation
- Choose a queue
- Quick start
- Core concepts
- Task options
- Results and errors
- Examples
- Documentation
- Development
- Community
Why osiiso
- Shared API across async, thread, and process execution.
- Priority scheduling where lower priority numbers run first.
- Retries with optional delay and exponential backoff.
- Per-task timeouts and queue-level run timeouts.
- Built-in rate limiting (
rate=/burst=) across all backends. - Delayed and absolute-time scheduling that never blocks a worker.
- A persistent process pool — subprocesses are reused across tasks, with crash detection and automatic respawn.
- Graceful shutdown that drains outstanding work, with
must_completeprotection. - Batch workflows with
submit(),map(), andgroup(). - Resumable fan-out via
Checkpoint— a killed run picks up where it left off. - Awaitable async handles and blocking sync handles, with
add_done_callback(). - Completion-order iteration via
as_completed()(async) anditer_completed()(sync). - One-shot helpers:
amap(),tmap(), andpmap()return ordered values. - Structured
RunSummaryand immutableTaskResultrecords withmetadatapassthrough. - Lifecycle hooks for
on_start,on_complete, andon_retry. - Worker
initializersupport for threads and subprocesses. - Optional
uvloopintegration throughosiiso.run().
Installation
pip install osiiso
With optional uvloop support:
pip install "osiiso[uvloop]"
The project targets Python 3.13 and newer.
Choose a queue
| Workload | Queue | Good for |
|---|---|---|
| Coroutine-based I/O | AsyncQueue |
HTTP clients, async databases, websockets, API fan-out |
| Coroutine I/O fed from threads | SharedQueue |
One event loop serving producers on other threads |
| Blocking synchronous work | ThreadQueue |
File operations, blocking SDKs, SQLite writes, sync integrations |
| CPU-heavy functions | ProcessQueue |
Ranking, parsing, scoring, transformations, analytics |
The queues intentionally look similar, so work can move between execution models with minimal changes.
Quick start
import asyncio
import osiiso
async def fetch(name: str) -> str:
await asyncio.sleep(0.1)
return f"fetched {name}"
async def main():
async with osiiso.AsyncQueue(workers=4) as q:
q.submit(fetch, "users", priority=0)
q.submit(fetch, "posts", retries=2, retry_delay=0.25, timeout=5)
summary = await q.run(strict=True)
return summary.values
print(osiiso.run(main()))
Core concepts
submit()
Use submit() for one task. It returns a handle immediately.
handle = q.submit(fetch_user, "ada", retries=3, timeout=10, name="fetch-user")
Async handles are awaitable:
result = await handle
value = handle.value()
Thread and process handles are blocking:
result = handle.wait(timeout=5)
value = handle.value()
map()
Use map() for one callable over many inputs.
q.map(download, urls, retries=2, group_id="downloads")
q.map(add, [(1, 2), (3, 4), (5, 6)], name="add")
q.map(request, [{"method": "GET", "url": "https://example.com"}])
Tuple entries are unpacked as positional arguments. Mapping entries are passed as keyword arguments.
group()
Use group() for a named batch, especially when tasks have different callables.
group = q.group(
[
(extract, "db"),
(transform, raw_records),
(load, destination),
],
group_id="etl-batch-1",
)
summary = q.run()
values = group.values()
For AsyncQueue, use await group.wait() and await group.values().
values() raises ExecutionError if any task failed or was cancelled, so a
returned tuple always lines up 1:1 with the inputs. Iterate a group in
completion order with group.as_completed().
One-shot helpers
When all you need is "run this function over these inputs", skip the queue
ceremony. Each helper builds a queue, runs it, and returns values in input
order, raising ExecutionError on any failure:
pages = await osiiso.amap(fetch, urls, workers=8, retries=2) # AsyncQueue
sizes = osiiso.tmap(stat_file, paths, workers=8) # ThreadQueue
scores = osiiso.pmap(rank, datasets, workers=4) # ProcessQueue
Resumable runs
Long fan-outs fail partway. Pass a Checkpoint and the inputs that already
succeeded are not submitted again — their stored values are handed straight
back, so results still line up 1:1 with the input:
from osiiso import Checkpoint, ThreadQueue
with Checkpoint("scrape.sqlite") as cp, ThreadQueue(workers=8, rate=5) as q:
grp = q.group(fetch, urls, checkpoint=cp, retries=3)
q.run()
pages = grp.values()
Kill it at 30k of 50k URLs and run it again: only the missing 20k are fetched. Only successes are recorded, so failed and cancelled tasks retry next run.
This is completion tracking keyed by input, not a durable task queue — the callable is never persisted, and nothing recovers work that was queued but never started. See Resumable Runs.
Rate limiting
Every queue accepts rate (maximum task attempts per second) and burst
(how many attempts may start back-to-back after an idle period):
async with osiiso.AsyncQueue(workers=8, rate=10, burst=3) as q:
q.map(call_api, payloads, retries=2)
await q.run()
Bound tasks
Bind a callable to a queue with @q.task().
async with osiiso.AsyncQueue(workers=4) as q:
@q.task(retries=2, retry_delay=0.25, name="fetch")
async def fetch(url: str) -> str:
return await client.get(url)
fetch("https://example.com")
fetch.map(["https://example.org", "https://example.net"])
summary = await q.run(strict=True)
Queue examples
AsyncQueue
import asyncio
import osiiso
async def fetch(name: str) -> str:
await asyncio.sleep(0.1)
return f"fetched {name}"
async def main():
async with osiiso.AsyncQueue(workers=4) as q:
q.map(fetch, ["users", "posts", "comments"], retries=2, timeout=5)
summary = await q.run(strict=True)
print(summary.values)
osiiso.run(main())
ThreadQueue
import time
import osiiso
def resize(path: str) -> str:
time.sleep(0.1)
return f"resized {path}"
with osiiso.ThreadQueue(workers=4) as q:
q.map(resize, ["a.png", "b.png", "c.png"], name="resize")
summary = q.run(strict=True)
print(summary.values)
ProcessQueue
Each worker owns a persistent subprocess that is reused across tasks, so spawn cost is paid once per worker instead of once per task. Timeouts and cancellation terminate the subprocess (it is respawned for the next task), and a crashed worker is reported as that task's failure while the pool recovers automatically.
Tasks and arguments must be pickleable, and standard spawn rules apply:
guard your script's entry point with if __name__ == "__main__":.
import osiiso
def score(n: int) -> int:
return sum(i * i for i in range(n))
if __name__ == "__main__":
with osiiso.ProcessQueue(workers=4) as q:
q.map(score, [10_000, 20_000, 30_000], name="score")
summary = q.run(strict=True)
print(summary.values)
Task options
Task behavior can be configured inline or through an immutable TaskOptions object.
from osiiso import TaskOptions
retrying = TaskOptions(retries=3, retry_delay=0.5, backoff=2, timeout=10)
urgent = retrying.replace(priority=0, name="urgent-api-call")
q.submit(fetch, url, opts=urgent)
q.submit(fetch, other_url, retries=3, retry_delay=0.5, backoff=2)
| Option | Default | Meaning |
|---|---|---|
priority |
3 |
Lower numbers run first. |
must_complete |
False |
Protects a task during graceful shutdown. |
timeout |
None |
Per-task timeout in seconds. |
retries |
0 |
Retry attempts after the first failure. |
retry_delay |
0.0 |
Delay before the first retry. |
backoff |
1.0 |
Multiplier applied after each retry. |
delay |
None |
Run after this many seconds. |
run_at |
None |
Run at an absolute epoch timestamp. |
name |
None |
Custom result and hook name. |
group_id |
None |
Group label for summaries. |
detached |
False |
Task still runs, but its result is excluded from the RunSummary; observe it via its handle. |
metadata |
None |
Arbitrary user data carried onto the handle and the TaskResult. |
TaskOptions validates invalid combinations immediately. For example, delay and run_at are mutually exclusive, negative retries are rejected, and unknown submit options raise TypeError.
Results and errors
Every run() returns a RunSummary.
summary.ok
summary.succeeded
summary.failed
summary.cancelled
summary.timed_out
summary.values
summary.errors
summary.by_task_id()
summary.by_name()
summary.by_group()
summary.raise_for_errors()
summary.display()
Use strict=True when failures should raise ExecutionError after the run finishes:
summary = await q.run(strict=True)
Each task result is stored as a TaskResult with task id, name, status, value, exception, attempts, priority, timing, group id, and cancellation metadata.
Lifecycle and policies
Queues support finite and long-running modes:
q = osiiso.AsyncQueue(mode="finite", fail_policy="continue", on_timeout="complete")
mode="finite"runs pending work and exits.mode="infinite"keeps workers alive until shutdown or timeout.fail_policy="continue"records failures and keeps processing.fail_policy="fail_first"cancels remaining work after the first failure (must_completetasks are spared).on_timeout="complete"letsmust_completetasks finish when a run times out.on_timeout="cancel"cancels everything when a run times out.
Leaving the context manager (or calling shutdown()) drains all outstanding
work — including scheduled tasks — before stopping workers. Use
shutdown(force=True) (or raise inside the with block) to cancel
everything immediately instead.
Bounded queues (size=N) cap outstanding tasks: sync queues block submit()
until space frees (natural backpressure), while AsyncQueue.submit() raises
QueueFullError.
Hooks give you a simple integration point for logging, metrics, and tracing:
def completed(result: osiiso.TaskResult) -> None:
print(result.name, result.status, result.duration)
q = osiiso.ThreadQueue(on_complete=completed)
Examples
Run the compact feature gallery:
uv run python examples/feature_gallery.py
Run the complete Hacker News style showcase:
uv run python -m examples.hackernews_showcase --limit 6
Use the live Hacker News API:
uv run python -m examples.hackernews_showcase --limit 20 --online
The showcase uses all three backends:
AsyncQueuefetches feeds, items, and users.ThreadQueuepersists records into SQLite.ProcessQueueranks stories and computes keywords.
Documentation
Full documentation is available at ichinga-samuel.github.io/osiiso.
Run the docs locally:
python -m pip install -e ".[docs]"
mkdocs serve
Build the docs strictly:
mkdocs build --strict
Development
Install the development dependencies:
python -m pip install -e ".[dev]"
Run tests:
uv run pytest
Run Ruff:
uv run ruff check .
Build the package:
python -m build
Build the docs with the docs extra:
uv run --extra docs mkdocs build --strict
Community
- Read the contribution guide before opening larger pull requests.
- Check the changelog for release history and upcoming changes.
- Use support guidance for questions, bug reports, and feature requests.
- Report vulnerabilities through the security policy, not public issues.
- Follow the code of conduct when participating in project spaces.
Project layout
.
|-- src/osiiso/ # Library source
|-- tests/ # Unit tests for async, thread, process, and options behavior
|-- docs/ # MkDocs documentation
|-- examples/feature_gallery.py # Compact API showcase
|-- examples/hackernews_showcase # Complete multi-backend example project
|-- pyproject.toml # Package metadata and tool configuration
`-- mkdocs.yml # Documentation site configuration
Public API
from osiiso import (
AsyncQueue,
ThreadQueue,
ProcessQueue,
TaskOptions,
TaskHandle,
SyncTaskHandle,
TaskGroup,
SyncTaskGroup,
TaskResult,
RunSummary,
ExecutionError,
ClosedError,
QueueFullError,
OsiisoError,
run,
amap,
tmap,
pmap,
as_completed,
iter_completed,
)
License
osiiso is released under the MIT License.
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file osiiso-1.1.0.tar.gz.
File metadata
- Download URL: osiiso-1.1.0.tar.gz
- Upload date:
- Size: 78.3 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: uv/0.9.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
cdaef99fb369307e747ac8e056bc98ae4eae7e4f7293e934bef06c8a42dd7d99
|
|
| MD5 |
6552924a5617c46de80f1332d21c5b44
|
|
| BLAKE2b-256 |
8c9526ecaabe7e2e54d0b2f312dc02f9a93bec34eefa9055e603f692a0d00488
|
File details
Details for the file osiiso-1.1.0-py3-none-any.whl.
File metadata
- Download URL: osiiso-1.1.0-py3-none-any.whl
- Upload date:
- Size: 53.8 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: uv/0.9.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
496ba9629b549c0875d2b5b15ef559f9aed495a1606ef44b66f96a5294df9eda
|
|
| MD5 |
f223c94686762c3638d9ca83d658a4e6
|
|
| BLAKE2b-256 |
d9e269223e5a9309589da7d0b217d1bb1bd7a1770deeaf5e66c8b0eb3dc9b9a7
|