Spaces:
Runtime error
Runtime error
| 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 |