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