File size: 3,478 Bytes
021b7e0
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
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