File size: 3,494 Bytes
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
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