🗄️ 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)
| File | Size | Uploaded | |
|---|---|---|---|
| shqaff-0.0.2.tar.gz | 17.4 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|