Skip to main content

Pinion: a tiny pluggable job queue with retries.

Project description

Pinion

Pinion is a tiny, pluggable job queue and worker for Python. It provides a simple @task registry, an in-memory queue for quick starts, and a durable SQLite backend for cross-process work, plus a retry policy with exponential backoff.

Features

  • In-memory queue with thread-safe Condition coordination
  • Durable SQLite storage with atomic job claim (WAL) across processes
  • Pluggable Storage protocol (SPI) for custom backends
  • Task registry via @task decorator (case-insensitive names)
  • Worker loop with polling, retries, timeouts, and graceful stop/join
  • Exponential backoff retries with optional jitter and cap
  • Dead-letter queue (DLQ) after exhausted retries
  • Basic worker metrics (processed/succeeded/failed/retried/dead_lettered/reaped)
  • Job lifecycle tracking: PENDING, RUNNING, SUCCESS, FAILED

Requirements

  • Python 3.12+

Installation

  • From source (local dev): pip install -e .
  • CLI entry point installs as pinion

Quick Start

CLI

  • pinion prints a concise quickstart, admin command cheatsheet, and exits.
  • pinion --demo runs a tiny in-memory demo (adds 1+2).
  • SQLite admin commands: status, running, pending, dlq-list, dlq-replay, enqueue, worker.

Examples:

# Show queue summary for a DB
pinion status --db pinion.db

# Run a worker for 5s and import your task module for registration
pinion worker --db pinion.db --max-retries 2 --task-timeout 5 \
  --import your_project.tasks --run-seconds 5

# Enqueue a job by name with JSON args/kwargs
pinion enqueue add --db pinion.db --args '[1,2]'

Library usage (in-memory)

import threading, time
from pinion import task, Job, InMemoryStorage, Worker, RetryPolicy

@task()
def add(a: int, b: int) -> int:
    return a + b

storage = InMemoryStorage()
worker = Worker(storage, retry=RetryPolicy(jitter=False), task_timeout=2.0)
thread = threading.Thread(target=worker.run_forever, daemon=True)
thread.start()

storage.enqueue(Job("add", (1, 2)))   # args tuple
storage.enqueue(Job("BOOM"))           # case-insensitive lookup (if registered)

time.sleep(2.5)
worker.stop()
thread.join()
# Optional: access basic metrics
print(worker.metrics)

Library usage (SQLite)

import threading, time
from pinion import task, Job, Worker, RetryPolicy, SqliteStorage  # durable backend

@task("boom")
def fail() -> None:
    raise ValueError("kaboom")

storage = SqliteStorage("pinion.db")
worker = Worker(storage, retry=RetryPolicy(jitter=False), task_timeout=2.0)
t = threading.Thread(target=worker.run_forever, daemon=True)
t.start()

storage.enqueue(Job("fail"))

time.sleep(4.5)
worker.stop()
t.join()

# Inspect DLQ (SQLite backend) for permanently failed jobs
# rows: (id, func_name, args_json, kwargs_json, attempts, error, failed_at)
print(storage._conn.execute("SELECT * FROM dlq").fetchall())

Core Concepts

  • Job: encapsulates function name, args/kwargs, id, status, attempts, timestamps
  • Storage: SPI with enqueue, dequeue, mark_done, mark_failed, size, heartbeat, reap_stale, dead_letter
  • Task registry: mapping of case-insensitive names to callables via @task
  • Worker: pulls jobs, executes callables, applies retry policy and optional per-task timeouts
  • Retry policy: max_retries, base_delay, cap, optional jitter
  • DLQ: jobs moved to durable dead-letter storage after retries are exhausted
  • Metrics: basic counters available via worker.metrics

Extending Storage

Implement the Storage protocol to plug in your own backend (e.g., Redis, DB, file-based):

class MyStorage:
    def enqueue(self, job: Job) -> None: ...
    def dequeue(self, timeout: float | None = None) -> Job | None: ...
    def mark_done(self, job: Job) -> None: ...
    def mark_failed(self, job: Job, exc: Exception) -> None: ...
    def size(self) -> int: ...
    def heartbeat(self, job: Job) -> None: ...
    def reap_stale(self, visibility_timeout: float) -> int: ...
    def dead_letter(self, job: Job, exc: Exception) -> None: ...

dequeue should block until timeout or a job is available, mark the job RUNNING, and increment attempts before returning the job.

Design Notes

  • InMemoryStorage uses a Condition for coordinating producers/consumers.
  • SqliteStorage uses WAL mode and an atomic claim (BEGIN IMMEDIATE + UPDATE) to safely select a PENDING job across processes. Access is serialized with a lock.
  • Worker uses a JobExecution context manager to mark success/failure and provides a thread-based per-task timeout option.
  • Retries are scheduled by re-enqueuing the same job after a computed delay.
  • Registry keys are normalized to lowercase for case-insensitive task names.
  • DLQ persists final failures (SQLite backend has a dlq table).

Limitations

  • In-memory storage is ephemeral; jobs are not persisted across restarts.
  • SQLite backend is local to a machine; horizontal scaling requires a different backend.
  • Thread-based timeouts cannot kill Python threads; long-running tasks should be cooperative or run in separate processes.
  • No result storage/return channel; tasks handle their own outputs.
  • Minimal inspection/metrics API.

Breaking Changes (since 0.1.x)

  • Storage SPI: added dead_letter(job, exc); custom backends must implement it.

Release and Publishing

Pinion targets Python 3.12+ and is published as pinion-queue.

Suggested versioning for this release: bump to 0.2.0.

  1. Update version
  • Edit pyproject.toml: set version = "0.2.0".
  • Edit pinion/__init__.py: set __version__ = "0.2.0".
  1. Build distributions
python -m pip install --upgrade pip build twine
python -m build
twine check dist/*
  1. Publish to PyPI
# Set PYPI token in environment (from your PyPI account)
export TWINE_USERNAME="__token__"
export TWINE_PASSWORD="pypi-***your-token***"

twine upload dist/*
  1. Tag the release
git tag v0.2.0
git push --tags

If pinion-queue is not available on PyPI, choose an alternative name or organization namespace.

Project Layout

  • pinion/queue.py — core queue, worker, storages, demo tasks
  • pinion/cli.py — simple CLI demo (pinion)

Upgrading

  • pip: pip install -U pinion-queue
  • pipx: pipx upgrade pinion-queue
  • Poetry: poetry update pinion-queue

The CLI performs a lightweight update check (with a short timeout) and prints a hint if a newer version is available. Disable via PINION_NO_UPDATE_CHECK=1 or pinion --no-update-check.

See the Changelog for release notes: CHANGELOG.md.


Pinion aims to be a tiny, understandable foundation you can extend with a real storage backend and operational features as needed.

Project details


Download files

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

Source Distribution

pinion_queue-0.2.6.tar.gz (11.9 kB view details)

Uploaded Source

Built Distribution

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

pinion_queue-0.2.6-py3-none-any.whl (16.0 kB view details)

Uploaded Python 3

File details

Details for the file pinion_queue-0.2.6.tar.gz.

File metadata

  • Download URL: pinion_queue-0.2.6.tar.gz
  • Upload date:
  • Size: 11.9 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.13.3

File hashes

Hashes for pinion_queue-0.2.6.tar.gz
Algorithm Hash digest
SHA256 9b108b1ccf812b347303b92afa5c5c95519f6d9ab82441926d6f5374ece67b5b
MD5 66e80911d522d38b92db252fb80cbf78
BLAKE2b-256 191f87ea0eb025e77cd9677fde51702f25e4f5ede8f9141740ed81120872cfac

See more details on using hashes here.

File details

Details for the file pinion_queue-0.2.6-py3-none-any.whl.

File metadata

  • Download URL: pinion_queue-0.2.6-py3-none-any.whl
  • Upload date:
  • Size: 16.0 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.13.3

File hashes

Hashes for pinion_queue-0.2.6-py3-none-any.whl
Algorithm Hash digest
SHA256 6641169ce1df56efd06ad9ea20283de02d9e37dbfa9086a6a68021ebd03df514
MD5 af47dcdf7e2054e11db4a34f610bb5ad
BLAKE2b-256 84a8c3912f6ae153cbb06f7e5a4f92e8f872b6c50326d7eac5ee1208f7a0f75d

See more details on using hashes here.

Supported by

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