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