from typing import List, Dict, Any, Optional from scheduler.tasks.base_task import BaseTask, TaskStatus, TaskResult from scheduler.task_registry import TaskRegistry import asyncio import logging logger = logging.getLogger(__name__) class TaskExecutor: def __init__(self, registry: TaskRegistry): self.registry = registry self.execution_context: Dict[str, Any] = {} async def execute_task_chain(self, task_ids: List[str], context: Dict[str, Any] = None) -> Dict[str, TaskResult]: """Execute tasks in dependency order""" if context: self.execution_context.update(context) tasks = [self.registry.get_task(task_id) for task_id in task_ids] tasks = [task for task in tasks if task is not None] if not tasks: logger.warning("No valid tasks found to execute") return {} # Sort tasks by dependencies (topological sort) sorted_tasks = self._topological_sort(tasks) completed_task_ids = [] results = {} logger.info(f"Starting execution chain with {len(sorted_tasks)} tasks") for task in sorted_tasks: if task.can_execute(completed_task_ids): # Add results from previous tasks to context task_context = self.execution_context.copy() task_context['previous_results'] = results result = await task.run(task_context) results[task.task_id] = result if result.status == TaskStatus.COMPLETED: completed_task_ids.append(task.task_id) # Add task result data to global context for next tasks if result.data: self.execution_context.update(result.data) elif result.status == TaskStatus.FAILED: logger.error(f"Task chain stopped due to failure in task: {task.task_id}") break else: logger.warning(f"Skipping task {task.task_id} - dependencies not met") results[task.task_id] = TaskResult(TaskStatus.SKIPPED, error="Dependencies not satisfied") logger.info(f"Task chain execution completed. {len(completed_task_ids)} tasks succeeded") return results def _topological_sort(self, tasks: List[BaseTask]) -> List[BaseTask]: """Sort tasks based on dependencies using topological sort""" task_dict = {task.task_id: task for task in tasks} visited = set() temp_visited = set() result = [] def visit(task_id: str): if task_id in temp_visited: raise ValueError(f"Circular dependency detected involving task: {task_id}") if task_id in visited: return temp_visited.add(task_id) task = task_dict.get(task_id) if task: for dep_id in task.dependencies: if dep_id in task_dict: visit(dep_id) visited.add(task_id) result.append(task) temp_visited.remove(task_id) for task in tasks: if task.task_id not in visited: visit(task.task_id) return result