TEST-FRANKO / scheduler /task_executor.py
Franko Fišter
Added scheduling to data collection script
021b7e0
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