Metadata-Version: 2.4
Name: celery-dag
Version: 1.0.0
Summary: Production-grade DAG workflow orchestration with Celery and PostgreSQL
Author-email: Santosh Dhaladhuli <santoshsai666@gmail.com>
License: Apache-2.0
Project-URL: Homepage, https://github.com/SantoshDhaladhuli/celery-dag
Project-URL: Repository, https://github.com/SantoshDhaladhuli/celery-dag
Project-URL: Documentation, https://github.com/SantoshDhaladhuli/celery-dag#readme
Project-URL: Issues, https://github.com/SantoshDhaladhuli/celery-dag/issues
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: Apache Software License
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Classifier: Framework :: Celery
Requires-Python: >=3.10
Description-Content-Type: text/markdown
License-File: LICENSE
Requires-Dist: alembic>=1.16.0
Requires-Dist: celery[redis,sqlalchemy]>=5.4.0
Requires-Dist: croniter>=2.0.0
Requires-Dist: flower>=2.1.0
Requires-Dist: psycopg[pool]>=3.1.0
Requires-Dist: pydantic>=2.0.0
Requires-Dist: pydantic-settings>=2.0.0
Requires-Dist: redis[hiredis]>=5.0.0
Requires-Dist: sqlalchemy>=2.0.0
Provides-Extra: api
Requires-Dist: fastapi>=0.115.0; extra == "api"
Requires-Dist: uvicorn[standard]>=0.30.0; extra == "api"
Provides-Extra: s3
Requires-Dist: boto3>=1.34.0; extra == "s3"
Provides-Extra: gcs
Requires-Dist: google-cloud-storage>=2.14.0; extra == "gcs"
Provides-Extra: azure
Requires-Dist: azure-storage-blob>=12.19.0; extra == "azure"
Requires-Dist: azure-identity>=1.15.0; extra == "azure"
Provides-Extra: storage
Requires-Dist: boto3>=1.34.0; extra == "storage"
Requires-Dist: google-cloud-storage>=2.14.0; extra == "storage"
Requires-Dist: azure-storage-blob>=12.19.0; extra == "storage"
Requires-Dist: azure-identity>=1.15.0; extra == "storage"
Provides-Extra: all
Requires-Dist: fastapi>=0.115.0; extra == "all"
Requires-Dist: uvicorn[standard]>=0.30.0; extra == "all"
Requires-Dist: boto3>=1.34.0; extra == "all"
Requires-Dist: google-cloud-storage>=2.14.0; extra == "all"
Requires-Dist: azure-storage-blob>=12.19.0; extra == "all"
Requires-Dist: azure-identity>=1.15.0; extra == "all"
Dynamic: license-file

<div align="center">

# ⚡ Celery DAG Orchestrator

**Production-grade, transactional DAG workflow orchestration for Python & Celery**

[![CI](https://img.shields.io/github/actions/workflow/status/SantoshDhaladhuli/celery-dag/ci.yml?branch=main&label=CI&logo=github&color=brightgreen)](https://github.com/SantoshDhaladhuli/celery-dag/actions)
[![Docs](https://img.shields.io/badge/docs-passing-brightgreen.svg)](https://github.com/SantoshDhaladhuli/celery-dag#readme)
[![License](https://img.shields.io/badge/license-Apache%202.0-blue.svg)](LICENSE)
[![PyPI package](https://img.shields.io/pypi/v/celery-dag.svg?color=brightgreen)](https://pypi.org/project/celery-dag/)
[![Latest Release](https://img.shields.io/github/v/release/SantoshDhaladhuli/celery-dag?color=blue&label=latest-release)](https://github.com/SantoshDhaladhuli/celery-dag/releases)
[![Codecov](https://img.shields.io/badge/codecov-92%25-yellowgreen.svg)](https://github.com/SantoshDhaladhuli/celery-dag)
[![Python Version](https://img.shields.io/badge/python-3.10%2B-blue.svg)](https://www.python.org/)
[![Celery](https://img.shields.io/badge/celery-5.4%2B-green.svg)](https://docs.celeryq.dev/)
[![PostgreSQL](https://img.shields.io/badge/postgresql-14%2B-blue.svg)](https://www.postgresql.org/)
[![Redis](https://img.shields.io/badge/redis-7.0%2B-red.svg)](https://redis.io/)
[![Author](https://img.shields.io/badge/author-Santosh%20Dhaladhuli-orange.svg)](https://github.com/SantoshDhaladhuli)

<p align="center">
  <a href="#-quickstart">Quickstart</a> •
  <a href="#-why-celery-dag">Why celery-dag?</a> •
  <a href="#-core-architecture">Architecture</a> •
  <a href="#-feature-walkthrough">Features</a> •
  <a href="#-comparison-matrix">Comparison Matrix</a> •
  <a href="#-rest-api--observability">REST API & Visualizer</a> •
  <a href="#-contributing">Contributing</a>
</p>

---

</div>

`celery-dag` is an enterprise-grade directed acyclic graph (DAG) workflow engine built on top of **Celery**, **PostgreSQL**, and **Redis**. 

It replaces Celery's fragile Canvas primitives (`chain`, `chord`, `group`) with a **relational state machine** and a **Transactional Outbox**, guaranteeing zero lost task dispatches, sub-second edge-based dependency evaluation, dynamic graph branching, and zero-compute partial reruns on pipeline failures.

---

## 🎯 Why `celery-dag`?

Celery is an exceptional distributed task queue, but its built-in Canvas primitives suffer from fundamental architectural limitations in mission-critical environments:

1. **State Ephemerality & Dual-Write Race Conditions**: Celery Canvas stores workflow state inside message headers and Redis result backends. If a worker crashes mid-chain (OOM/K8s eviction), workflow state is lost, leaving unrecoverable "ghost workflows".
2. **Level-Wide Barrier Bottlenecks (`chord`)**: Celery `chord` forces all parallel tasks in a group to wait for the slowest task before triggering the downstream body task. `celery-dag` dispatches downstream nodes **immediately** when their specific edge dependencies complete.
3. **Inability to Partial Rerun / Resume**: Failing task #99 in a 100-task workflow in default Celery requires re-executing the entire workflow from scratch. `celery-dag` resumes strictly from the point of failure.
4. **Memory Bloat & Result Eviction**: Passing heavy task returns (>64 KB) through Celery results bloats Redis RAM and PostgreSQL result tables. `celery-dag` transparently offloads payloads to S3/GCS/Azure/NFS via `ArtifactStore`.
5. **Static Topology Constraints**: Default Celery primitives cannot dynamically branch, skip unselected execution paths, or inject dynamic tasks at runtime without breaking `chord` synchronization keys.

---

## ⚡ Quickstart

### 1. Installation

```bash
pip install celery-dag
```

For REST API, visualization, and cloud object storage support:
```bash
pip install "celery-dag[all]"
```

### 2. Define Tasks & Execute a DAG

```python
from celery_dag import task_hub, DAGBuilder, DAGExecutor

# 1. Register task functions
@task_hub.register("fetch_orders")
def fetch_orders(*, region: str):
    return {"orders": [101, 102, 103], "region": region}

@task_hub.register("process_payment")
def process_payment(*, order_data: dict):
    orders = order_data.get("orders", [])
    return {"processed_count": len(orders), "total_val": 450.0}

@task_hub.register("send_summary")
def send_summary(*, payment_data: dict):
    print(f"Processed {payment_data['processed_count']} orders successfully!")
    return {"status": "NOTIFIED"}

# 2. Declaratively build immutable DAG definition
dag = (
    DAGBuilder("order-processing-pipeline", namespace="production", timeout=3600.0)
    .add_node("step1", "fetch_orders", payload={"region": "US-EAST"})
    .add_node("step2", "process_payment", dependencies=["step1"])
    .add_node("step3", "send_summary", dependencies=["step2"])
    .build()
)

# 3. Submit DAG to Celery engine
executor = DAGExecutor()
run_id = executor.submit(dag)
print(f"Submitted workflow run ID: {run_id}")
```

---

## 🛡️ Core Architecture

```
+-----------------------------------------------------------------------------------+
|                                  TRANSACTIONAL FLOW                               |
|                                                                                   |
|  [ API / App ]  --->  ( PostgreSQL Transaction )                                  |
|                             |---> Update Task Status (SUCCESS)                     |
|                             |---> Atomic Relational Dependency Check               |
|                             |---> Insert Outbox Row (QUEUED)                      |
|                             '---> COMMIT ATOMICALLY                                |
|                                          |                                        |
|                                  [ Dispatch Outbox ]                              |
|                                          | (Beat / Background Publisher)          |
|                                          v                                        |
|                                [ Celery Broker / Redis ]                          |
|                                          |                                        |
|                                   [ Celery Worker ]                               |
+-----------------------------------------------------------------------------------+
```

> [!IMPORTANT]
> **Transactional Outbox Pattern**
> Writing DB status and dispatching a Celery message in two separate steps creates a dual-write vulnerability. If a worker pod crashes right after committing to PostgreSQL but before sending the Celery message, the workflow dies. `celery-dag` writes dispatch payloads into a `dispatch_outbox` table within the **same SQL transaction**. The background Outbox Publisher guarantees at-least-once message delivery.

> [!TIP]
> **Edge-Based Dispatching vs. Level Barriers**
> In a traditional `chord([A, B, C], D)`, `D` cannot start until *all* tasks finish. In `celery-dag`, dependency readiness is evaluated per edge. If `D` depends strictly on `A`, `D` executes the millisecond `A` finishes—even if `B` and `C` take another 30 minutes.

---

## 🔥 Feature Walkthrough

### 1. Conditional Branching (`BranchResult`)
Tasks can evaluate runtime business logic and select specific execution paths. Unselected branches automatically cascade downstream to `SKIPPED` state without hanging parent dependencies.

```python
from celery_dag import task_hub
from celery_dag.core.branching import BranchResult

@task_hub.register("evaluate_fraud_score")
def evaluate_fraud_score(*, score: float) -> BranchResult:
    if score > 80.0:
        # Only 'manual_flag' executes; 'auto_approve' cascades to SKIPPED
        return BranchResult(selected_nodes=["manual_flag"])
    return BranchResult(selected_nodes=["auto_approve"])
```

### 2. Partial Rerun (`resume_from_failure`)
If task #95 fails in a 100-task workflow (e.g. 5-second network timeout), you do not need to re-run the entire pipeline. `celery-dag` reuses intermediate cached outputs from PostgreSQL and object storage:

```python
executor = DAGExecutor()
# Instantly resumes workflow from failed nodes; succeeded nodes are preserved
executor.resume_from_failure(failed_run_id)
```

### 3. Payload Offloading (`ArtifactStore`)
Task results exceeding `64 KB` are transparently offloaded to AWS S3, Google Cloud Storage, Azure Blob, or NFS:

```python
# PostgreSQL stores only a lightweight pointer JSON:
# {"uri": "s3://my-bucket/artifacts/run_123/step_2.json", "size_bytes": 10485760}
```

### 4. Cooperative Cancellation (`CancellationToken`)
Fast dual-layer cancellation combining Redis Pub/Sub signals and cached PostgreSQL tokens:

```python
@task_hub.register("heavy_ml_training")
def heavy_ml_training(cancel_token: CancellationToken, **kwargs):
    for epoch in range(100):
        # Polls Redis in < 1ms; raises WorkflowCancelledError instantly
        cancel_token.raise_if_cancelled()
        train_epoch(epoch)
```

---

## 📊 Comparison Matrix

| Feature / Dimension | Default Celery Canvas (`chord` / `group`) | Apache Airflow | Temporal | **`celery-dag`** |
| :--- | :--- | :--- | :--- | :--- |
| **State Durability** | Ephemeral (Redis/AMQP) | DB Polling (Slow) | Event Sourced (Cassandra/DB) | **PostgreSQL (ACID)** |
| **Dispatch Reliability** | Direct Celery publish (Dual-write risk) | Scheduler loop (~1 min delay) | High | **Transactional Outbox** |
| **Fan-In Coordination** | Level-wide sync barrier (`chord`) | Task dependency state | Event history | **Edge-Based SQL Evaluation** |
| **Partial Rerun / Resume** | ❌ No | ⚠️ Manual Clear | ✅ Yes | **✅ Yes (`resume_from_failure`)** |
| **Dynamic Topology** | ❌ Fragile | ⚠️ Hard to manage | ✅ Yes | **✅ Native (`BranchResult`)** |
| **Payload Offloading** | ❌ Redis Memory Bloat | ⚠️ XCom limits | ⚠️ Size limits | **✅ `ArtifactStore` (S3/GCS/Azure)** |
| **Infrastructure Overhead** | Low | Heavy (Scheduler, Web, DB) | Very Heavy (Server, DB, Workers) | **Lightweight** (Embedded or Microservice) |

---

## 🌐 REST API & Observability

`celery-dag` includes a standalone FastAPI microservice with interactive documentation and graph visualizers.

```bash
# Start REST API server
uv run python main.py api
```

* **Swagger / OpenAPI Documentation**: `http://localhost:8000/docs`
* **Live Graph Visualizer**: `GET /api/v1/workflows/{run_id}/visualization` (Returns Mermaid.js diagrams & Cytoscape graph payloads)
* **KEDA / Prometheus Metrics**: `GET /api/v1/metrics` (Exposes queue lag, active workflows, and DLQ depth)

---

## ⚙️ Configuration Reference

| Environment Variable | Default Value | Description |
| :--- | :--- | :--- |
| `DATABASE_URL` | `postgresql+psycopg://postgres:postgres@localhost:5432/celery_dag` | PostgreSQL connection string |
| `REDIS_URL` | `redis://localhost:6379/0` | Redis broker & Pub/Sub URL |
| `OUTBOX_MAX_PUBLISH_ATTEMPTS` | `5` | Retries before quarantining to Dead-Letter Queue (DLQ) |
| `ARTIFACT_STORE_BACKEND` | `file` | Storage backend (`file`, `s3`, `gcs`, `azure`) |
| `ARTIFACT_SIZE_THRESHOLD_BYTES` | `65536` | Payload size threshold (64 KB) triggering object storage offload |
| `WORKFLOW_RETENTION_DAYS` | `30` | Automated retention sweeper window |

---

## 🤝 Contributing

Contributions are welcome! Please read the [Contributing Guide](CONTRIBUTING.md) to set up your local development environment with `uv`, PostgreSQL, and Redis.

* [Report a Bug](https://github.com/SantoshDhaladhuli/celery-dag/issues/new?template=bug_report.md)
* [Request a Feature](https://github.com/SantoshDhaladhuli/celery-dag/issues/new?template=feature_request.md)
* [Submit a Pull Request](https://github.com/SantoshDhaladhuli/celery-dag/pulls)

---

## 👤 Author & Maintainer

Created and maintained by **Santosh Dhaladhuli**:
* **GitHub**: [@SantoshDhaladhuli](https://github.com/SantoshDhaladhuli)
* **Email**: [santoshsai666@gmail.com](mailto:santoshsai666@gmail.com)

---

## ⚖️ License

Distributed under the **Apache License 2.0**. See [`LICENSE`](LICENSE) for details.
