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.

Python Version TypeScript License: MIT Architecture


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.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 (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:


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)

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

Built distribution (wheel)

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

Release history Release notifications | RSS feed

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

This release

0.1.0 This release

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