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 (
task-engine) and TypeScript / Node.js (@task-engine/sdk). - 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 (task-engine)
pip install 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 (@task-engine/sdk)
npm install @task-engine/sdk
# or
pnpm add @task-engine/sdk
Submit tasks and poll results with complete type safety:
import { TaskEngine } from "@task-engine/sdk";
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 "@task-engine/sdk";
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 web dashboard and API server
task-engine ui --port 8000
# 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.0
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.0.tar.gz | 2.2 MB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| distributed_task_engine-0.1.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 2.3 MB
Release files / distributed_task_engine-0.1.0.tar.gz
| Download URL | distributed_task_engine-0.1.0.tar.gz |
|---|---|
| Size | 2.2 MB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
29684cb1bac66dede9f8e6b66a10b609ba9254011d31a55302c7d671777472b8
|
|
BLAKE2b-256 checksum How to use checksums |
78333117c3fdcd82996967c25198910de956f94a59a4ce9a0f6c3fa83c9112d5
|
| 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.0-py3-none-any.whl
| Download URL | distributed_task_engine-0.1.0-py3-none-any.whl |
|---|---|
| Size | 141.6 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
d2e0c17148af5c33fcbe54f40ff0fbab5377d0855e8a33d9e64989052d584d23
|
|
BLAKE2b-256 checksum How to use checksums |
befc71d5f64ea8614708a57fe8e6c6c3de1ab59bc60ee43fcbb95ba41cf5be41
|
| 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}
|