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.4.tar.gz (11.7 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.4-py3-none-any.whl (15.8 kB view details)

Uploaded Python 3

File details

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

File metadata

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

File hashes

Hashes for pinion_queue-0.2.4.tar.gz
Algorithm Hash digest
SHA256 0cf918010d2650eeffa3e5f3cf7bda97e73576c88625507cbd0d95bcbc661ac0
MD5 48fca5c10a5353a6d862ae142d35ea1d
BLAKE2b-256 fa44908dd4e2f6182af40f69d7276edc12b68e02aba732698a97ae652169137b

See more details on using hashes here.

File details

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

File metadata

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

File hashes

Hashes for pinion_queue-0.2.4-py3-none-any.whl
Algorithm Hash digest
SHA256 153c99bc0b9d55dce4aed10d001b6e5fa0a2cb96f2cd213921284ba4eac9a58a
MD5 1582c3e9a76a50fa8149e88d91da9127
BLAKE2b-256 0062c21b4c0b3ab772f29de2c6eefae7a1efa080c965fa786f4dafc6ec11125d

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