Skip to main content

🗄️ shqaff

A lightweight, PostgreSQL-backed task queue for Python. Uses SQLAlchemy for storage and a finite state machine for reliable task lifecycle management.

Features

  • PostgreSQL-backed persistence with row-level locking (SELECT ... FOR UPDATE SKIP LOCKED)
  • Finite state machine for task transitions: pending -> in_progress -> done | failed
  • Configurable retry logic with max retry limits
  • Simple consumer registration pattern
  • Minimal dependencies: SQLAlchemy, psycopg2, transitions

Installation

pip install shqaff

For development:

pip install shqaff[dev]

Quick Start

1. Start PostgreSQL

Using Docker:

docker run -d --name shqaff-db \
  -e POSTGRES_DB=shqaff \
  -e POSTGRES_USER=postgres \
  -e POSTGRES_PASSWORD=postgres \
  -p 5432:5432 \
  postgres:15

2. Define a Consumer

A consumer is a class that processes tasks. Subclass Consumer and implement the name property and run() method:

from shqaff.consumer import Consumer
from shqaff.registry import register_consumer

class EmailConsumer(Consumer):

    @property
    def name(self) -> str:
        return "send_email"

    def run(self, payload: dict) -> None:
        to = payload["to"]
        subject = payload["subject"]
        body = payload["body"]
        print(f"Sending email to {to}: {subject}")
        # ... actual email sending logic here

register_consumer(EmailConsumer)

3. Create Tasks

from shqaff.db import init_db, SessionLocal
from shqaff.producer import create_task

# Initialize the database tables
init_db()

db = SessionLocal()

# Enqueue a task
create_task(
    db=db,
    task_name="welcome_email",
    consumer="send_email",
    payload={
        "to": "user@example.com",
        "subject": "Welcome!",
        "body": "Thanks for signing up.",
    },
    max_retries=3,
)

4. Process Tasks

from shqaff.event_loop import process_tasks

# Poll the database and process pending tasks
# This runs an infinite loop with a configurable interval
process_tasks(db=db, poll_interval=2, batch_size=10)

Single-batch Processing

For cron-style execution, process one batch and exit:

from shqaff.event_loop import process_once

process_once(db=db, batch_size=10)

Full Working Example

from dataclasses import asdict, dataclass

from shqaff.consumer import Consumer
from shqaff.db import init_db, SessionLocal
from shqaff.event_loop import process_tasks
from shqaff.producer import create_task
from shqaff.registry import register_consumer


@dataclass
class OrderPayload:
    order_id: int
    items: list[str]


class OrderConsumer(Consumer):

    @property
    def name(self) -> str:
        return "process_order"

    def run(self, payload: dict) -> None:
        order = OrderPayload(**payload)
        print(f"Processing order #{order.order_id}: {order.items}")


if __name__ == "__main__":
    init_db()
    register_consumer(OrderConsumer)

    db = SessionLocal()
    create_task(
        db=db,
        task_name="new_order",
        consumer="process_order",
        payload=asdict(OrderPayload(order_id=42, items=["widget", "gadget"])),
    )

    # Process tasks every 2 seconds
    process_tasks(db=db, poll_interval=2)

Repository Layer

TaskRepository is a thin, engine-agnostic CRUD layer over the task model. It holds no connection of its own — it operates on the Session you pass in, so it works identically against PostgreSQL in production and SQLite in tests:

from shqaff import SessionLocal, TaskRepository, TaskStatus

repo = TaskRepository(SessionLocal())

task = repo.create(task_name="welcome_email", consumer="send_email", payload={"to": "a@b.c"})
repo.get(task.id)
repo.list(status=TaskStatus.PENDING, limit=20)
repo.count(status=TaskStatus.PENDING)
repo.set_status(task.id, TaskStatus.DONE)
repo.delete(task.id)

Command-Line Interface

Installing the package exposes a shqaff command with full CRUD over the queue:

shqaff init-db
shqaff create --task-name welcome_email --consumer send_email --payload '{"to": "a@b.c"}'
shqaff list --status pending
shqaff get 1
shqaff update 1 --status done
shqaff delete 1

The database URL is resolved from the --database-url option, the SHQAFF_DATABASE_URL environment variable, or the SHAQAFF_DB_* configuration variables, in that order.

Task Lifecycle

Tasks follow a strict state machine:

pending ──start()──> in_progress ──succeed()──> done
                          │
                          └──fail()──> failed
  • pending: Task is queued, waiting to be picked up
  • in_progress: A consumer is currently executing the task
  • done: Task completed successfully
  • failed: Task exceeded max retries and is permanently failed

When a task raises an exception during processing, it is retried (reset to pending) until max_retries is reached, at which point it transitions to failed.

Configuration

Database connection is configured via environment variables:

Variable Default Description
SHAQAFF_DB_HOST localhost PostgreSQL host
SHAQAFF_DB_PORT 5432 PostgreSQL port
SHAQAFF_DB_NAME shqaff Database name
SHAQAFF_DB_USER postgres Database user
SHAQAFF_DB_PASS (empty) Database password

Docker

Run the full stack with Docker Compose:

docker compose up -d db       # Start PostgreSQL
docker compose up shqaff      # Run the demo app
docker compose run test       # Run the test suite

Testing

The suite spans unit, integration, fuzz (Hypothesis) and end-to-end layers. The repository, CLI and lifecycle tests run against in-memory SQLite and need no PostgreSQL; the legacy DB-backed tests (test_smoke, test_event_loop, test_retry_logic, test_fail_after_max_retries) require a running instance.

pip install -e ".[dev]"

# Runs everywhere, no database required:
pytest tests/test_repository.py tests/test_repository_fuzz.py \
       tests/test_cli.py tests/test_e2e.py tests/test_package.py \
       tests/test_task_fsm.py tests/test_invalid_fsm_transition.py

# Full suite (needs a running PostgreSQL instance):
pytest tests/

License

MIT

Metadata

Release files for shqaff 0.0.2

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for shqaff 0.0.2
File Size Uploaded
shqaff-0.0.2.tar.gz 17.4 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for shqaff 0.0.2
File Interpreter ABI Platform
shqaff-0.0.2-py3-none-any.whl Python 3 none any Details

Total release size: 30.3 kB

Release files / shqaff-0.0.2.tar.gz

Download URL shqaff-0.0.2.tar.gz
Size 17.4 kB
Tags Source
SHA-256 checksum
How to use checksums
0420b3a32cf3ed610acfbf97ed42526d3d4c78a230d3e4981cf6b988c24ac999
BLAKE2b-256 checksum
How to use checksums
09b41eac80fe8b33236ab79c5f8304dd8a10436514c2d887ee3c81e30c18234f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.12.5

Release files / shqaff-0.0.2-py3-none-any.whl

Download URL shqaff-0.0.2-py3-none-any.whl
Size 13.0 kB
Tags Python 3
SHA-256 checksum
How to use checksums
ca49d58e4279f4a6d71b581f26e3613e1ded877058c3f4e5eeab44031480ec33
BLAKE2b-256 checksum
How to use checksums
06c3cf3cd9a17d0b067221fab4363a9f6d7f99aa35d34694162e9fc0faf9eb81
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.12.5

Release history Release notifications | RSS feed

This release

0.0.2 This release

2 release files

0.0.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page