Spaces:
Running
Running
| 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 | |