Terminal / agents /workflow_engine.py
Baida07's picture
sync: 175 file da Baida98/AI@d1881b9c (2026-08-25 07:39 UTC) [deploy-all] (#67)
bd654eb
Raw
History Blame
4.34 kB
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