Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion README.rst
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,8 @@ Features
--------

- Define tasks using the ``@Task`` decorator
- Automatic execution order based on task dependencies
- Automatically resolves task execution order based on defined dependencies
- Supports both OR and AND dependency types
- Built-in logging and execution context for task outputs
- Task states tracking: ``PENDING``, ``SUCCESS``, ``FAILED``, ``SKIPPED``
- Fail-fast error handling for dependent tasks
Expand Down
124 changes: 100 additions & 24 deletions docs/documentation/task_dep.rst
Original file line number Diff line number Diff line change
@@ -1,50 +1,126 @@
Task Dependencies
===================
=================

In many workflows, tasks depend on the results of other tasks.
**Task Orchestrator** allows you to define dependencies easily
using the ``depends_on`` attribute.
**Pravaha** allows you to define task dependencies using the ``depends_on`` attribute.

This ensures that tasks are executed in the correct order and
prevents errors caused by missing prerequisites.
Defining dependencies ensures that tasks are executed in the correct order and
prevents failures caused by missing prerequisites.

---

Defining Dependencies
-----------------------
---------------------

You can specify dependencies as a list of task names when
defining a task:
Pravaha supports two types of dependencies:

1. **OR-based dependencies**
2. **AND-based dependencies**

---

OR-based Dependencies
---------------------

An **OR-based dependency** means that a task will execute if **at least one**
of its dependent tasks completes successfully.

Example
^^^^^^^

.. code-block:: python

from pravaha.core.task import Task
from pravaha.core.executor import TaskExecutor
from pravaha.core.task import Task
from pravaha.core.executor import TaskExecutor
from pravaha.dependency.dependency import Dependency

# Dependency(type="OR | AND", dependencies=list[str])
@Task("a", depends_on=Dependency("OR", dependencies=["b", "c"]))
def a():
print("a")

@Task("b")
def b():
pass

@Task(name="task1", depends_on=["task2"])
def task1():
print("Task1 executed")
@Task("c")
def c():
pass

@Task(name="task2")
def task2():
print("Task2 executed")
TaskExecutor.execute()

if __name__ == "__main__":
TaskExecutor.execute()
---

AND-based Dependencies
----------------------

An **AND-based dependency** means that a task will execute **only if all**
of its dependent tasks complete successfully.

AND-based dependencies can be defined in two ways.

### 1. Implicit AND dependency using a list

By default, passing a list to ``depends_on`` is treated as an **AND dependency**.

.. code-block:: python

from pravaha.core.task import Task
from pravaha.core.executor import TaskExecutor

@Task(name="task1", depends_on=["task2"])
def task1():
print("Task1 executed")

@Task(name="task2")
def task2():
print("Task2 executed")

if __name__ == "__main__":
TaskExecutor.execute()

Expected Output
-----------------
^^^^^^^^^^^^^^^

.. code-block:: text

Task2 executed
Task1 executed
Task2 executed
Task1 executed

---

### 2. Explicit AND dependency using ``Dependency``

You can also explicitly define an AND-based dependency using the ``Dependency`` object.

.. code-block:: python

from pravaha.core.task import Task
from pravaha.core.executor import TaskExecutor
from pravaha.dependency.dependency import Dependency

# Dependency(type="OR | AND", dependencies=list[str])
@Task("a", depends_on=Dependency("AND", dependencies=["b", "c"]))
def a():
print("a")

@Task("b")
def b():
pass

@Task("c")
def c():
pass

TaskExecutor.execute()

---

Notes
-------
-----

- **Task registration order does not matter** when dependencies are defined.
- The executor automatically determines the **correct execution order** based on the dependency graph.
- Circular dependencies will result in an error. Make sure your workflow graph is **acyclic**.
- The executor automatically determines the **correct execution order**
based on the dependency graph.
- Circular dependencies are not allowed. Ensure that your workflow graph
is **acyclic (DAG)**.
41 changes: 34 additions & 7 deletions pravaha/core/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,13 @@
from pravaha.enums.task_status import TaskStatus
from pravaha.validation.dag import DAGValidator
from pravaha.context.condition.context import ConditionContext
from datetime import datetime
from time import time, sleep
import os
from pravaha.core.task import ErrorInformation
from pravaha.utils.utilities import sort_task_on_the_basis_of_priority, filter_tasks_on_the_basis_of_tags
from pravaha.utils.dep_resolver import resolve_dependencies
from pravaha.exception.task import InvalidDependencyType
from datetime import datetime
from time import time, sleep
import os

class TaskExecutor:
"""
Expand Down Expand Up @@ -53,11 +54,27 @@ def _execute_helper(cls, task: Task):
if task.state in [TaskStatus.SUCCESS, TaskStatus.FAILED, TaskStatus.SKIPPED]:
return

for dep_name in task.depends_on:
dep_task = cls.tasks[dep_name]
dep_type = task.depends_on.type

if dep_type not in ['AND', 'OR', ""]:
raise InvalidDependencyType(f"Invalid dependency type: {dep_type} with task: {task.name}")

for dep_name in task.depends_on.dependencies:
dep_task = cls.tasks[dep_name]
cls._execute_helper(dep_task)
if dep_type == "OR" and dep_task.state == TaskStatus.SUCCESS:
break

can_we_proceed = True

if dep_type == "OR":
if not any(cls.tasks[dep].state == TaskStatus.SUCCESS for dep in task.depends_on.dependencies):
can_we_proceed = False
elif dep_type == "AND":
if any(cls.tasks[dep].state in [TaskStatus.FAILED, TaskStatus.SKIPPED] for dep in task.depends_on.dependencies):
can_we_proceed = False

if any(cls.tasks[dep].state in [TaskStatus.FAILED, TaskStatus.SKIPPED] for dep in task.depends_on):
if not can_we_proceed:
task.state = TaskStatus.SKIPPED
return

Expand All @@ -68,7 +85,17 @@ def _execute_helper(cls, task: Task):
return

# Preparing inputs for dependency outputs.
inputs = [cls.ExecutionContext.get(dep_name) for dep_name in task.depends_on]

inputs = []

if dep_type == "OR":
for dep_name in task.depends_on.dependencies:
dep_task = cls.tasks[dep_name]
if dep_task.state == TaskStatus.SUCCESS:
inputs.append(cls.ExecutionContext.get(dep_name))

else:
inputs = [cls.ExecutionContext.get(dep_name) for dep_name in task.depends_on.dependencies]

# Executing task along with checking for retry.
# Implementing retry logic.
Expand Down
18 changes: 16 additions & 2 deletions pravaha/core/task.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
from pravaha.enums.task_status import TaskStatus
from pravaha.retry.policy import RetryPolicy
from pravaha.enums.task_priority import TaskPriority
from pravaha.dependency.dependency import Dependency
from typing import Union, List, Optional

class Task:
"""
Expand All @@ -22,9 +24,8 @@ class Task:
tag (str): Tag to the task.
"""

def __init__(self, name, depends_on=None, retries: RetryPolicy=None, condition=None, priority=TaskPriority.NORMAL, tag=None):
def __init__(self, name, depends_on: Optional[Union[Dependency, List[str]]] = None, retries: Optional[RetryPolicy]=None, condition=None, priority=TaskPriority.NORMAL, tag=None):
self.name = name
self.depends_on = depends_on or []
self.retries = retries
self.function_ref = None
self.state = TaskStatus.PENDING # Default state
Expand All @@ -38,6 +39,19 @@ def __init__(self, name, depends_on=None, retries: RetryPolicy=None, condition=N
self.condition = condition
self.priority = priority
self.tag = tag
self.depends_on = self._normalize_dependency(depends_on)

@staticmethod
def _normalize_dependency(depends_on: Optional[Union[Dependency, List[str]]]) -> Dependency:

if depends_on is None:
return Dependency()
elif isinstance(depends_on, Dependency):
return depends_on
elif isinstance(depends_on, list):
return Dependency("AND", depends_on)
else:
raise TypeError(f"depends_on must be dependency or List[str]")

def __call__(self, original_function):

Expand Down
Empty file added pravaha/dependency/__init__.py
Empty file.
10 changes: 10 additions & 0 deletions pravaha/dependency/dependency.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
class Dependency:
def __init__(self, type: str = "", dependencies: list[str] = []):
self.type = type
self.dependencies = dependencies

def get_type(self):
return self.name

def get_dependencies(self):
return self.dependencies
3 changes: 3 additions & 0 deletions pravaha/exception/task.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,4 +2,7 @@ class TaskFailedError(Exception):
pass

class TaskNotFoundError(Exception):
pass

class InvalidDependencyType(Exception):
pass
2 changes: 1 addition & 1 deletion pravaha/utils/dep_resolver.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ def dfs(task_name: str):
visited.add(task_name)
task: Task = registry[task_name]

for dep_name in task.depends_on:
for dep_name in task.depends_on.dependencies:
dfs(dep_name)

resolved[task_name] = task # Adding resolved task to dictionary.
Expand Down
4 changes: 2 additions & 2 deletions pravaha/utils/execution_plan.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ def dry_run():

tasks = Registry.get_task()
for task in tasks.values():
if task.depends_on is not None and task.depends_on != []:
final_output.append(f"{task.name} - depends on ({','.join(task.depends_on)})")
if task.depends_on is not None and task.depends_on.dependencies != []:
final_output.append(f"{task.name} - depends on ({','.join(task.depends_on.dependencies)}) with dependency-type: {task.depends_on.type}")
else:
final_output.append(f"{task.name} - no dependencies ")

Expand Down
4 changes: 2 additions & 2 deletions pravaha/validation/dag.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ def validate(tasks: dict):

# Missing dependency check.
for task in tasks.values():
for dep in task.depends_on:
for dep in task.depends_on.dependencies:
if dep not in tasks:
raise MissingDependencyError(f"Task '{task.name}' depends on missing task '{dep}'")

Expand All @@ -19,7 +19,7 @@ def dfs(task_name, path):
path_visited.add(task_name)
path.append(task_name)

for dep in tasks[task_name].depends_on:
for dep in tasks[task_name].depends_on.dependencies:
if dep not in visited:
dfs(dep, path)
elif dep in path_visited:
Expand Down
Loading