Terminal / agents /workflow_engine.py
Baida-A
Initial clean deploy (Reverse Proxy removed)
28a08e7
Raw
History Blame
3.49 kB
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