import logging import time import uuid from typing import Any, Dict, List, Optional 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] = Field(default_factory=dict) class WorkflowExecutor: """ ARCH-I4.3: Workflow Engine. Coordina workflow in-memory step-by-step tramite Kernel ed Executor. La persistenza o il resume inter-processo non sono garantiti da questo motore; i caller possono consultare lo stato del workflow corrente tramite ``get_workflow``. """ def __init__(self, kernel: Any, executor: Any): self.kernel = kernel self.executor = executor self.active_workflows: Dict[str, Workflow] = {} def get_workflow(self, workflow_id: str) -> Optional[Workflow]: """Ritorna il workflow noto, inclusi gli stati terminali in memoria.""" return self.active_workflows.get(workflow_id) async def execute_workflow(self, workflow: Workflow) -> Workflow: """Esegue un workflow step-by-step, mantenendo il fallback locale.""" self.active_workflows[workflow.workflow_id] = workflow workflow.status = "running" _logger.info("Avvio workflow: %s (%s)", workflow.name, workflow.workflow_id) for step in workflow.steps: step.status = "running" step.started_at = time.time() _logger.info( "Esecuzione step: %s in workflow %s", step.tool_name, workflow.workflow_id, ) try: # ARCH-I4.3: il Kernel risolve la capability senza esporre # l'infrastruttura al workflow. resolution = await self.kernel.resolve_capability(step.tool_name) if resolution.get("status") == "resolved": worker = resolution["worker"] worker_id = worker.id if hasattr(worker, "id") else worker["id"] _logger.info( "Step %s risolto su worker: %s", step.tool_name, worker_id, ) result = await self.executor.run_tool( tool_name=step.tool_name, inputs=step.args, worker_hint=worker_id, ) else: # Nessun worker registrato: il comportamento storico resta # l'esecuzione locale tramite lo stesso Executor. _logger.warning( "Nessun worker per %s, provo esecuzione locale", step.tool_name, ) result = await self.executor.run_tool( tool_name=step.tool_name, inputs=step.args, ) step.result = result if isinstance(result, dict) and result.get("success") is False: step.status = "failed" step.error = str(result.get("error", "Tool execution failed")) workflow.status = "failed" break step.status = "completed" except Exception as exc: step.status = "failed" step.error = str(exc) workflow.status = "failed" _logger.error("Step %s fallito: %s", step.tool_name, exc) break finally: step.finished_at = time.time() if workflow.status == "running": workflow.status = "completed" _logger.info("Workflow %s terminato con stato: %s", workflow.name, workflow.status) return workflow