File size: 4,343 Bytes
b2b3de9
 
5f771b4
 
 
b2b3de9
 
 
 
5f771b4
b2b3de9
 
 
 
5f771b4
b2b3de9
 
 
 
 
5f771b4
b2b3de9
 
 
 
 
 
5f771b4
 
b2b3de9
 
 
5f771b4
 
 
 
 
 
b2b3de9
5f771b4
 
b2b3de9
 
 
 
5f771b4
 
 
 
b2b3de9
5f771b4
b2b3de9
 
5f771b4
b2b3de9
 
 
 
5f771b4
 
 
 
 
b2b3de9
5f771b4
 
 
 
 
b2b3de9
5f771b4
 
 
 
 
b2b3de9
 
5f771b4
 
b2b3de9
 
5f771b4
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
b2b3de9
5f771b4
b2b3de9
5f771b4
b2b3de9
5f771b4
 
b2b3de9
 
 
5f771b4
b2b3de9
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