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.1
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.1.tar.gz | 2.2 MB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| distributed_task_engine-0.1.1-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 2.3 MB
Release files / distributed_task_engine-0.1.1.tar.gz
| Download URL | distributed_task_engine-0.1.1.tar.gz |
|---|---|
| Size | 2.2 MB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
eb825ee7bfaadeb1c5f53e0ce0f09fd662dff809a7879f6230566e1283f08c88
|
|
BLAKE2b-256 checksum How to use checksums |
2b74a4b1cab4002ed3add05943f31cdf71d891d90b4c6265df2ae6cf4ed4152d
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
uv/0.12.7 {"installer":{"name":"uv","version":"0.12.7","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":null,"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}
|
Release files / distributed_task_engine-0.1.1-py3-none-any.whl
| Download URL | distributed_task_engine-0.1.1-py3-none-any.whl |
|---|---|
| Size | 141.9 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
12444bfca5a6541c6ee302355d91064c47baebc5e1ff7a7cb865b7903094cd03
|
|
BLAKE2b-256 checksum How to use checksums |
a9121a52f68adff037a19dd89b5c77753ea97ffc88585b2dd0af10955dab2219
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
uv/0.12.7 {"installer":{"name":"uv","version":"0.12.7","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":null,"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}
|