Skip to main content

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.

PyPI version npm version Python Version TypeScript License: MIT GitHub release


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-engine on PyPI) and TypeScript / Node.js (@abdullah_mog/task-engine on npm).
  • Celery-like Ergonomics: Decorate functions with @engine.task or @task, dispatch with .delay() or .apply_async(), and track results via AsyncResult.
  • 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:


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)

Source distribution for distributed-task-engine 0.1.2
File Size Uploaded
distributed_task_engine-0.1.2.tar.gz 2.9 MB Details

Built distribution (wheel)

Table of built distributions (wheels) for distributed-task-engine 0.1.2
File Interpreter ABI Platform
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

Release history Release notifications | RSS feed

0.1.3

2 release files

This release

0.1.2 This release

2 release files

0.1.1

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page