Task Engine
Event-Driven Distributed Task Processing Engine
Celery-like developer ergonomics with zero-extra-infrastructure on PostgreSQL, or enterprise-grade sub-millisecond throughput on Redis Streams.
Highlights
- Zero-Extra-Infrastructure Mode: Run background tasks, delayed jobs, and retries natively on PostgreSQL using ACID transactional guarantees and
FOR UPDATE SKIP LOCKED. No Redis or RabbitMQ required for small-to-medium deployments. - Enterprise High-Throughput Broker: Scale to millions of tasks with Redis Streams consumer groups (
XREADGROUP,XACK) with sub-millisecond dispatch. - Dual Language SDKs: Full-featured clients for Python (
distributed-task-engineon PyPI) and TypeScript / Node.js (@abdullah_mog/task-engineon npm). - Celery-like Ergonomics: Decorate functions with
@engine.taskor@task, dispatch with.delay()or.apply_async(), and track results viaAsyncResult. - Distributed Cron & Scheduling: Crontab and interval scheduling with automatic leader election via PostgreSQL advisory locks.
- Resilience & Fault Tolerance: Heartbeat leases (
lease_expires_at), automatic re-queueing of crashed worker tasks, exponential backoff retries with jitter, and dead-letter queue (DLQ) re-driving. - Multi-Tenancy & Security: First-class tenant isolation (
X-Tenant-Id), JWT & API Key authentication, and token bucket rate limiting. - Full Observability: Prometheus metrics (
/metrics), OpenTelemetry distributed tracing, structured JSON logs, and real-time WebSocket event feeds. - Real-time Web Console: Next.js 14 dashboard for queue monitoring, worker management, live metrics, and DLQ replay.
Quickstart
Python Quickstart (distributed-task-engine)
pip install distributed-task-engine
Define a background task in tasks.py:
from task_engine import TaskEngine, TaskContext
engine = TaskEngine.from_env()
@engine.task(queue="notifications", priority=3, max_retries=3)
async def send_welcome_email(payload: dict, context: TaskContext) -> dict:
print(f"Executing task {context.task_id} for tenant: {context.tenant_id}")
return {"status": "sent", "email": payload["email"]}
Start the worker:
task-engine worker -A tasks --queues notifications --concurrency 4
Submit and await results programmatically:
from tasks import send_welcome_email
result = send_welcome_email.delay({"email": "alice@example.com"})
print(f"Task ID: {result.id}")
output = result.get(timeout=30)
print(f"Task finished: {output}")
TypeScript / Node.js Quickstart (@abdullah_mog/task-engine)
npm install @abdullah_mog/task-engine
Alternative with pnpm:
pnpm add @abdullah_mog/task-engine
Submit tasks and poll results with complete type safety:
import { TaskEngine } from "@abdullah_mog/task-engine";
const engine = new TaskEngine({
apiUrl: process.env.TASK_ENGINE_API_URL || "http://localhost:8000",
apiKey: process.env.TASK_ENGINE_API_KEY,
tenantId: "default",
});
async function run() {
const result = await engine.submitTask({
taskType: "notifications.send_welcome_email",
payload: { email: "bob@example.com" },
queue: "notifications",
});
console.log(`Submitted: ${result.id}`);
const output = await result.get({ timeoutMs: 15000 });
console.log("Result:", output);
}
run();
Subscribe to real-time events over WebSockets:
import { TaskEngineWebSocket } from "@abdullah_mog/task-engine";
const ws = new TaskEngineWebSocket({
apiUrl: "http://localhost:8000",
channel: "tasks",
});
ws.on("task.completed", (event) => console.log("Completed:", event.data));
ws.connect();
CLI Command Reference
Start background worker process:
task-engine worker -A tasks --queues default,notifications --concurrency 10
Start distributed cron scheduler daemon:
task-engine scheduler -A tasks
Apply database migrations:
task-engine migrate
Start API backend server:
task-engine api start --port 8000
Start Web Operations Console (Flower equivalent):
pnpm --prefix frontend dev
Submit ad-hoc tasks from terminal:
task-engine submit notifications.send_welcome_email -p "{\"email\":\"user@test.com\"}"
Query task lifecycle and execution output:
task-engine status <task-id>
Documentation Hub
Comprehensive engineering guides, specifications, and runbooks:
- SDK Developer Guides:
- API Specifications:
- Infrastructure & Operations:
- Operational Runbooks:
- Architecture Decision Records:
- ADR Index
- ADR-001: Project Structure
- ADR-002: Configuration & Secrets
- ADR-003: Data Model & Persistence
- ADR-004: Core Services & Runtime
- ADR-005: API Surface & Schemas
- ADR-006: Redis Streams Broker Adapter
- ADR-007: Observability (Metrics, Tracing, Logging)
- ADR-008: Security, RBAC & Multi-Tenancy
- ADR-009: Test Strategy & Validation
Verification & Testing
Python Test Suite
uv run pytest tests/
uv run ruff check .
uv run mypy src/
TypeScript SDK Suite
cd packages/sdk-ts
pnpm test
pnpm run typecheck
pnpm run build
License
This project is licensed under the MIT License.
Metadata
Release files for distributed-task-engine 0.1.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 | |
|---|---|---|---|
| distributed_task_engine-0.1.2.tar.gz | 2.9 MB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| distributed_task_engine-0.1.2-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 3.8 MB
Release files / distributed_task_engine-0.1.2.tar.gz
| Download URL | distributed_task_engine-0.1.2.tar.gz |
|---|---|
| Size | 2.9 MB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
09c97f7ee688bc621e9d57b8ea85f795a4de6637e57264d6519ae8b5033b5f42
|
|
BLAKE2b-256 checksum How to use checksums |
2b680a98cbab497f77b0cde1de1a7e63db387631c30a02aacdf873b93d0e389b
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.14.3
|
Release files / distributed_task_engine-0.1.2-py3-none-any.whl
| Download URL | distributed_task_engine-0.1.2-py3-none-any.whl |
|---|---|
| Size | 881.2 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
dbbf4a3ee82f626e9994ff53d1421e3c4abd1811686ada16acf88aea8c85eab1
|
|
BLAKE2b-256 checksum How to use checksums |
a7bc9b3fb6c582f0339b5ed73bece8bcfcc9d8b32801d9d9825bf985c29b2866
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.14.3
|