Spaces:
Running
Running
File size: 4,343 Bytes
28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 28a08e7 bd654eb 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 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 | 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
|