Command Bus abstraction over PostgreSQL + PGMQ for reliable async command processing
Project description
rcmd - Reliable Commands
A Python library providing Command Bus abstraction over PostgreSQL + PGMQ.
Installation
pip install reliable-cmd
Overview
Command Bus enables reliable command processing with:
- At-least-once delivery via PGMQ visibility timeout
- Transactional guarantees - commands sent atomically with business data
- Retry policies with configurable backoff
- Troubleshooting queue for failed commands with operator actions
- Audit trail for all state transitions
Requirements
- Python 3.11+
- PostgreSQL 15+ with PGMQ extension
- uv package manager (recommended) or pip
Quick Start
# Clone the repository
git clone https://github.com/your-org/commandbus.git
cd commandbus
# Install dependencies (uses uv)
make install-dev
# Start PostgreSQL with PGMQ
make docker-up
# Run tests
make test
Alternative: Using pip with venv
# Create and activate virtual environment
python -m venv .venv
source .venv/bin/activate # On Windows: .venv\Scripts\activate
# Install dependencies
pip install -e ".[dev,e2e]"
# Install pre-commit hooks
pre-commit install
Developer Guide
This section covers how to set up command handlers and configure workers for your domain.
1. Define Command Handlers
Use the @handler decorator to mark methods as command handlers. Handlers are organized in classes with constructor-injected dependencies:
from psycopg_pool import AsyncConnectionPool
from commandbus import Command, HandlerContext, handler
class OrderHandlers:
"""Handlers for order domain commands."""
def __init__(self, pool: AsyncConnectionPool) -> None:
"""Inject dependencies via constructor."""
self._pool = pool
@handler(domain="orders", command_type="CreateOrder")
async def handle_create_order(
self, cmd: Command, ctx: HandlerContext
) -> dict[str, Any]:
"""Handle CreateOrder command.
Args:
cmd: The command with command_id and data
ctx: Handler context (currently provides metadata)
Returns:
Result dict stored in command record
"""
order_data = cmd.data
# Process the order...
return {"status": "created", "order_id": str(cmd.command_id)}
@handler(domain="orders", command_type="CancelOrder")
async def handle_cancel_order(
self, cmd: Command, ctx: HandlerContext
) -> dict[str, Any]:
"""Handle CancelOrder command."""
# Cancel logic...
return {"status": "cancelled"}
2. Handle Errors
Use built-in exception types to control retry behavior:
from commandbus.exceptions import PermanentCommandError, TransientCommandError
@handler(domain="orders", command_type="ProcessPayment")
async def handle_payment(self, cmd: Command, ctx: HandlerContext) -> dict[str, Any]:
try:
result = await payment_gateway.process(cmd.data)
return {"status": "paid", "transaction_id": result.id}
except PaymentDeclined as e:
# Permanent failure - no retry, moves to troubleshooting queue
raise PermanentCommandError(
code="PAYMENT_DECLINED",
message=str(e)
)
except GatewayTimeout as e:
# Transient failure - will be retried according to policy
raise TransientCommandError(
code="GATEWAY_TIMEOUT",
message=str(e)
)
3. Register Handlers and Create Worker
Create a composition root that wires up dependencies and registers handlers:
from psycopg_pool import AsyncConnectionPool
from commandbus import HandlerRegistry, RetryPolicy, Worker
async def create_pool() -> AsyncConnectionPool:
pool = AsyncConnectionPool(
conninfo="postgresql://localhost:5432/mydb", # configure auth as needed
min_size=2,
max_size=10,
)
await pool.open()
return pool
def create_registry(pool: AsyncConnectionPool) -> HandlerRegistry:
"""Create registry and register all handlers."""
# Create handler instances with dependencies
order_handlers = OrderHandlers(pool)
inventory_handlers = InventoryHandlers(pool)
# Register handlers - decorator metadata is used for routing
registry = HandlerRegistry()
registry.register_instance(order_handlers)
registry.register_instance(inventory_handlers)
return registry
def create_worker(pool: AsyncConnectionPool) -> Worker:
"""Create worker with retry policy."""
registry = create_registry(pool)
retry_policy = RetryPolicy(
max_attempts=3,
backoff_schedule=[10, 60, 300], # seconds between retries
)
return Worker(
pool=pool,
domain="orders",
registry=registry,
retry_policy=retry_policy,
visibility_timeout=30, # seconds before message redelivery
)
async def run_worker() -> None:
"""Main entry point."""
pool = await create_pool()
try:
worker = create_worker(pool)
await worker.run(
concurrency=4, # concurrent command handlers
poll_interval=1.0, # seconds between queue polls
)
finally:
await pool.close()
if __name__ == "__main__":
import asyncio
asyncio.run(run_worker())
4. Send Commands
Use the CommandBus to send commands:
from uuid import uuid4
from commandbus import CommandBus
async def create_order(bus: CommandBus, order_data: dict) -> UUID:
command_id = uuid4()
await bus.send(
domain="orders",
command_type="CreateOrder",
command_id=command_id,
data=order_data,
max_attempts=3, # optional, overrides retry policy
)
return command_id
For high-throughput scenarios, use batch sending:
from commandbus.models import SendRequest
requests = [
SendRequest(
domain="orders",
command_type="CreateOrder",
command_id=uuid4(),
data={"product_id": "123", "quantity": 1},
)
for _ in range(1000)
]
result = await bus.send_batch(requests)
print(f"Sent {result.total_commands} commands in {result.chunks_processed} chunks")
E2E Test Application
The repository includes an end-to-end test application with a web UI for testing command processing behaviors.
Prerequisites
- Docker and Docker Compose
- Python 3.11+ with dependencies installed (see Quick Start)
Running the E2E Application
1. Start the database:
make docker-up
2. Start the web UI:
make e2e-app
The web UI is available at http://localhost:5001
3. Start workers (in a separate terminal):
cd tests/e2e
python -m app.worker
To run multiple workers for load testing:
cd tests/e2e
for i in {1..4}; do
python -m app.worker &
done
Test Behaviors
The E2E UI supports various command behaviors for testing:
| Behavior | Description |
|---|---|
| No-Op | Returns immediately, for throughput benchmarking |
| Success | Completes successfully after optional delay |
| Fail Permanent | Fails with PermanentCommandError, moves to TSQ |
| Fail Transient | Fails with TransientCommandError, retries |
| Fail Transient Then Succeed | Fails N times, then succeeds |
| Timeout | Simulates slow execution |
Bulk Generation
For load testing, use the bulk generation form:
- Select behavior type (No-Op recommended for pure throughput tests)
- Set count (up to 1,000,000)
- Set execution time (0ms for maximum throughput)
- Click "Generate Bulk Commands"
Monitoring
The E2E UI provides:
- Dashboard: Real-time status counts and throughput metrics
- Commands: List and filter commands by status
- Troubleshooting Queue: View and action failed commands
- Audit Trail: Full event history per command
Documentation
- Implementation Spec - Detailed design and API
- Architecture Decisions - ADRs explaining key choices
- Contributing - How to contribute
License
MIT
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file reliable_cmd-0.1.0.tar.gz.
File metadata
- Download URL: reliable_cmd-0.1.0.tar.gz
- Upload date:
- Size: 33.6 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
63f2fa62dfd47c909847c1c7bd68efcc3627403b375ecba6576b3addebcd6207
|
|
| MD5 |
02c2db40b3d2214919e0aef0e8d1f6d9
|
|
| BLAKE2b-256 |
3fa56637dc39f8c12c160f7c933f75b51c2c022cd9791a79305a77726e90f3e7
|
Provenance
The following attestation bundles were made for reliable_cmd-0.1.0.tar.gz:
Publisher:
publish.yml on FreeSideNomad/rcmd
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
reliable_cmd-0.1.0.tar.gz -
Subject digest:
63f2fa62dfd47c909847c1c7bd68efcc3627403b375ecba6576b3addebcd6207 - Sigstore transparency entry: 791251149
- Sigstore integration time:
-
Permalink:
FreeSideNomad/rcmd@a90c8caf546a8057a86d6fdd3a03873fb85869dc -
Branch / Tag:
refs/tags/v0.1.0 - Owner: https://github.com/FreeSideNomad
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish.yml@a90c8caf546a8057a86d6fdd3a03873fb85869dc -
Trigger Event:
push
-
Statement type:
File details
Details for the file reliable_cmd-0.1.0-py3-none-any.whl.
File metadata
- Download URL: reliable_cmd-0.1.0-py3-none-any.whl
- Upload date:
- Size: 41.9 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
7287ce943b62538927f2169bf969b8ec3973a5d517acc796b436bd758beea988
|
|
| MD5 |
894057d97dc3d8423e7f718e95700fb8
|
|
| BLAKE2b-256 |
c4df9f3962b5835b3823923072f2ccbcd4d47dc91e7b0bb36f03f91a99c9080e
|
Provenance
The following attestation bundles were made for reliable_cmd-0.1.0-py3-none-any.whl:
Publisher:
publish.yml on FreeSideNomad/rcmd
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
reliable_cmd-0.1.0-py3-none-any.whl -
Subject digest:
7287ce943b62538927f2169bf969b8ec3973a5d517acc796b436bd758beea988 - Sigstore transparency entry: 791251155
- Sigstore integration time:
-
Permalink:
FreeSideNomad/rcmd@a90c8caf546a8057a86d6fdd3a03873fb85869dc -
Branch / Tag:
refs/tags/v0.1.0 - Owner: https://github.com/FreeSideNomad
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish.yml@a90c8caf546a8057a86d6fdd3a03873fb85869dc -
Trigger Event:
push
-
Statement type: