import asyncio import logging import uuid import time from typing import List, Dict, Optional, Any from pydantic import BaseModel, Field _logger = logging.getLogger("agents.workflow_engine") class WorkflowStep(BaseModel): step_id: str = Field(default_factory=lambda: str(uuid.uuid4())) tool_name: str args: Dict[str, Any] status: str = "pending" # pending, running, completed, failed result: Any = None error: Optional[str] = None started_at: Optional[float] = None finished_at: Optional[float] = None class Workflow(BaseModel): workflow_id: str = Field(default_factory=lambda: str(uuid.uuid4())) name: str steps: List[WorkflowStep] status: str = "pending" created_at: float = Field(default_factory=time.time) metadata: Dict[str, Any] = {} class WorkflowExecutor: """ ARCH-I4.3: Workflow Engine Coordina l'esecuzione di workflow persistenti orchestrati tramite l'Executor. """ def __init__(self, kernel, executor): self.kernel = kernel self.executor = executor self.active_workflows: Dict[str, Workflow] = {} async def execute_workflow(self, workflow: Workflow) -> Workflow: """Esegue un workflow step-by-step.""" self.active_workflows[workflow.workflow_id] = workflow workflow.status = "running" _logger.info(f"Avvio workflow: {workflow.name} ({workflow.workflow_id})") for step in workflow.steps: step.status = "running" step.started_at = time.time() _logger.info(f"Esecuzione step: {step.tool_name} in workflow {workflow.workflow_id}") try: # ARCH-I4.3 Integration: Usa il Kernel per risolvere e sottomettere il task # Risoluzione capability (ARCH-E3.2) res = await self.kernel.resolve_capability(step.tool_name) if res.get("status") == "resolved": worker = res["worker"] worker_id = worker.id if hasattr(worker, "id") else worker["id"] _logger.info(f"Step {step.tool_name} risolto su worker: {worker_id}") # Esecuzione via Executor (che ora usa il Kernel) result = await self.executor.run_tool( tool_name=step.tool_name, args=step.args, worker_hint=worker_id ) step.result = result step.status = "completed" else: # Fallback all'esecuzione locale se nessun worker รจ trovato _logger.warning(f"Nessun worker per {step.tool_name}, provo esecuzione locale") result = await self.executor.run_tool(step.tool_name, step.args) step.result = result step.status = "completed" except Exception as e: step.status = "failed" step.error = str(e) workflow.status = "failed" _logger.error(f"Step {step.tool_name} fallito: {e}") break step.finished_at = time.time() if workflow.status == "running": workflow.status = "completed" _logger.info(f"Workflow {workflow.name} terminato con stato: {workflow.status}") return workflow