Skip to main content

cairnq (Python)

SQLite-first, cross-language, storage-centered durable task runtime. The Python SDK. API and worker processes coordinate only through a shared SQLite file.

from cairnq import CairnQ, Worker

# Worker side — a handler always receives (ctx, payload).
worker = Worker.sqlite("tasks.db")

@worker.task                          # registered under the function name, "create_summary"
async def create_summary(ctx, payload):
    await ctx.progress(0.2, "reading")
    return {"summary": await llm.summarize(payload["text"])}

worker.serve()                        # blocking entry point; Ctrl-C closes cleanly

# API side (in your server) — submit returns immediately.
tasks = CairnQ.sqlite("tasks.db")
task = await tasks.submit("create_summary", {"text": text}, key=f"summary:{aid}")

@worker.task defaults the task name to the function's name. Pass a string for a dotted/namespaced name: @worker.task("summary.create").

Synchronous call (submit + wait):

from cairnq import TaskFailed, TaskTimeout

try:
    result = await tasks.call("create_summary", {"text": text}, wait_timeout_ms=10_000)
except TaskFailed as e:
    log(e.code, e.message, e.retryable)   # envelope fields, no e.error["code"] digging
except TaskTimeout as e:
    # The task keeps running — resume the wait instead of submitting again.
    result = await tasks.wait(e.task_id, timeout_ms=60_000)
    # …or tasks.wait_by_key(key), from a process that never held the id.

Inspect a task by id/key without memorizing status strings:

task = await tasks.get_by_key(key)
if task and task.succeeded:        # also .failed / .canceled / .running / .queued / .is_terminal
    use(task.result)

Optionally define a task once and share the symbol across both ends — the name lives in one place (no string drift), and call() is typed as the task's result:

from cairnq import TaskDef

summarize = TaskDef[dict, dict]("summarize")

@worker.task(summarize)            # registered under summarize.name
async def handle(ctx, payload): ...

result = await tasks.call(summarize, {"text": text})

Opt-in: every API still accepts a plain name string (cross-language callers use it).

Running it in production

worker = Worker.sqlite(
    "tasks.db",
    concurrency=4,            # handler calls at once; use max_in_flight_bytes to bound memory
    retry_backoff_ms=1_000,   # window doubles per attempt, capped at retry_backoff_max_ms (30s),
                              # jittered over its upper half; 0 disables
    on_error=lambda exc, info: log.warning("worker survived %s: %s", info, exc),
)

# Nothing else deletes rows, so give the client a retention policy — it sweeps
# terminal tasks in bounded batches for as long as the handle is open. A
# per-status mapping keeps each status on its own clock (statuses left out are
# never swept): spent results go in minutes, failures stay for diagnosis.
tasks = CairnQ.sqlite(
    "tasks.db",
    retention=Retention(older_than_ms={"succeeded": 300_000, "failed": 7 * 24 * 3600_000}),
)

A sync handler (def, not async def) is dispatched to a thread, so the usual shape around a blocking GPU or HTTP call keeps the worker's event loop — and with it every lease this worker holds — alive:

@worker.task("score")
def score(ctx, payload):
    return {"score": model.forward(payload["image"])}  # blocking, off the loop

A handler that does real side effects should bail out when it loses its lease — the task is already running on another worker and nothing it writes is recorded:

@worker.task("long.job")
async def long_job(ctx, payload):
    for chunk in chunks:
        if ctx.lost_lease or await ctx.canceled():
            return
        await process(chunk)

Multi-host

Same code, Postgres instead of the file — CairnQ.postgres(dsn) / Worker.postgres(dsn). Install with pip install cairnq[postgres].

schema puts cairnq's tables in a schema of their own: CairnQ.postgres(dsn, schema="cairnq") creates it if absent and sets search_path on every connection. Every process in a deployment must agree on it — a queue whose API and worker resolve to different schemas is two empty queues, and both sides come up healthy. cairnq refuses to connect where it can see that about to happen; the TypeScript SDK takes the same option and applies the same rule.

Sharing the application's connection

Given a PgExecutor instead of a DSN, cairnq runs inside a session the application already has — no second driver, no second pool:

from cairnq import CairnQ, PgExecutor

executor: PgExecutor = ...   # ~40 lines over your driver
tasks = CairnQ.postgres(executor)

An adapter passes rows through as its driver produced them; cairnq normalizes what the drivers disagree about. An executor cairnq was handed is never closed by cairnq.

That shared session is also what lets a task's settlement commit together with the rows the task produced:

@worker.task("render.document")
async def render_document(ctx, payload):
    rendered = await render(payload)

    async def write(session):
        await session.query("insert into pages (doc, n) values ($1, $2)", [...])
        return {"pages": len(rendered)}   # becomes the task's result

    return await ctx.succeed_in(write)

Without it the two are separate transactions, and a crash between them leaves work durable while the task still reads as running — on retry, recomputed. If the lease turns out to be gone, the settlement matches no row and the caller's writes roll back with it.

Watching

watch calls back when the tasks on a queue may have changed — for a dashboard that would otherwise poll:

stop = tasks.watch(on_signal, queues=["render"])

It is notify-accelerated polling, not an event log. On Postgres an idle watch costs nothing and a signal lands within milliseconds; where LISTEN is unavailable — a transaction-mode pooler, or SQLite, which has no channel — the timer alone still delivers poll signals. Treat a signal as "re-read now"; the truth is in stats() / list() / get().

The protocol (schema + canonical SQL) lives in ../cairnq-protocol and is shared verbatim with the TypeScript SDK. See ../cairnq-protocol/PROTOCOL.md.

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

cairnq-0.10.0.tar.gz (150.5 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

cairnq-0.10.0-py3-none-any.whl (125.6 kB view details)

Uploaded Python 3

File details

Details for the file cairnq-0.10.0.tar.gz.

File metadata

  • Download URL: cairnq-0.10.0.tar.gz
  • Upload date:
  • Size: 150.5 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for cairnq-0.10.0.tar.gz
Algorithm Hash digest
SHA256 520ce731683babc733419db24a445ebca07d9e8e7b9438a48ebed5fc4bf70b65
MD5 0ee1a7cf0771dbcdcbadbaeb8e2e3231
BLAKE2b-256 f83ad1af18a7f2590642e21dbe37c4fa0d41f74a151b388bbbc68e7389e37fe9

See more details on using hashes here.

Provenance

The following attestation bundles were made for cairnq-0.10.0.tar.gz:

Publisher: publish.yml on Jannchie/cairnq

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file cairnq-0.10.0-py3-none-any.whl.

File metadata

  • Download URL: cairnq-0.10.0-py3-none-any.whl
  • Upload date:
  • Size: 125.6 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for cairnq-0.10.0-py3-none-any.whl
Algorithm Hash digest
SHA256 d301d9d47fe34bac3b685cd4e465f56cf55e56191796f56c0f5016f6803b3f04
MD5 f2fd0802bb8f560e6ee8d087bdafb57e
BLAKE2b-256 c7f47748e2d8278b812032278e2ac252b2347319cc88e725c6ea83868a0fb39f

See more details on using hashes here.

Provenance

The following attestation bundles were made for cairnq-0.10.0-py3-none-any.whl:

Publisher: publish.yml on Jannchie/cairnq

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Release history Release notifications | RSS feed

0.14.0

2 files

0.13.0

2 files

0.12.0

2 files

0.11.0

2 files

This release

0.10.0 This release

2 files

0.9.0

2 files

0.8.0

2 files

0.7.0

2 files

0.6.0

2 files

0.5.0

2 files

0.4.0

2 files

0.3.0

2 files

0.2.0

2 files

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page