Spaces:
Paused
Paused
File size: 3,494 Bytes
28a08e7 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 | 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
|