Baida07 commited on
Commit
bd654eb
·
1 Parent(s): a0c8e17

sync: 175 file da Baida98/AI@d1881b9c (2026-08-25 07:39 UTC) [deploy-all] (#67)

Browse files

- sync: 175 file da Baida98/AI@d1881b9c (2026-08-25 07:39 UTC) [deploy-all] (b85ebc9eb1d4d6c1e63c7e9b2c341f62c891faf3)

agents/workflow_engine.py CHANGED
@@ -1,90 +1,112 @@
1
- import asyncio
2
  import logging
3
- import uuid
4
  import time
5
- from typing import List, Dict, Optional, Any
 
 
6
  from pydantic import BaseModel, Field
7
 
8
  _logger = logging.getLogger("agents.workflow_engine")
9
 
 
10
  class WorkflowStep(BaseModel):
11
  step_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
12
  tool_name: str
13
  args: Dict[str, Any]
14
- status: str = "pending" # pending, running, completed, failed
15
  result: Any = None
16
  error: Optional[str] = None
17
  started_at: Optional[float] = None
18
  finished_at: Optional[float] = None
19
 
 
20
  class Workflow(BaseModel):
21
  workflow_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
22
  name: str
23
  steps: List[WorkflowStep]
24
  status: str = "pending"
25
  created_at: float = Field(default_factory=time.time)
26
- metadata: Dict[str, Any] = {}
 
27
 
28
  class WorkflowExecutor:
29
  """
30
- ARCH-I4.3: Workflow Engine
31
- Coordina l'esecuzione di workflow persistenti orchestrati tramite l'Executor.
 
 
 
 
32
  """
33
- def __init__(self, kernel, executor):
 
34
  self.kernel = kernel
35
  self.executor = executor
36
  self.active_workflows: Dict[str, Workflow] = {}
37
 
 
 
 
 
38
  async def execute_workflow(self, workflow: Workflow) -> Workflow:
39
- """Esegue un workflow step-by-step."""
40
  self.active_workflows[workflow.workflow_id] = workflow
41
  workflow.status = "running"
42
- _logger.info(f"Avvio workflow: {workflow.name} ({workflow.workflow_id})")
43
 
44
  for step in workflow.steps:
45
  step.status = "running"
46
  step.started_at = time.time()
47
-
48
- _logger.info(f"Esecuzione step: {step.tool_name} in workflow {workflow.workflow_id}")
49
-
 
 
50
  try:
51
- # ARCH-I4.3 Integration: Usa il Kernel per risolvere e sottomettere il task
52
- # Risoluzione capability (ARCH-E3.2)
53
- res = await self.kernel.resolve_capability(step.tool_name)
54
-
55
- if res.get("status") == "resolved":
56
- worker = res["worker"]
57
  worker_id = worker.id if hasattr(worker, "id") else worker["id"]
58
- _logger.info(f"Step {step.tool_name} risolto su worker: {worker_id}")
59
-
60
- # Esecuzione via Executor (che ora usa il Kernel)
 
 
61
  result = await self.executor.run_tool(
62
  tool_name=step.tool_name,
63
- args=step.args,
64
- worker_hint=worker_id
65
  )
66
-
67
- step.result = result
68
- step.status = "completed"
69
  else:
70
- # Fallback all'esecuzione locale se nessun worker è trovato
71
- _logger.warning(f"Nessun worker per {step.tool_name}, provo esecuzione locale")
72
- result = await self.executor.run_tool(step.tool_name, step.args)
73
- step.result = result
74
- step.status = "completed"
75
-
76
- except Exception as e:
 
 
 
 
 
 
 
 
 
 
 
 
77
  step.status = "failed"
78
- step.error = str(e)
79
  workflow.status = "failed"
80
- _logger.error(f"Step {step.tool_name} fallito: {e}")
81
  break
82
-
83
- step.finished_at = time.time()
84
 
85
  if workflow.status == "running":
86
  workflow.status = "completed"
87
-
88
- _logger.info(f"Workflow {workflow.name} terminato con stato: {workflow.status}")
89
  return workflow
90
-
 
 
1
  import logging
 
2
  import time
3
+ import uuid
4
+ from typing import Any, Dict, List, Optional
5
+
6
  from pydantic import BaseModel, Field
7
 
8
  _logger = logging.getLogger("agents.workflow_engine")
9
 
10
+
11
  class WorkflowStep(BaseModel):
12
  step_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
13
  tool_name: str
14
  args: Dict[str, Any]
15
+ status: str = "pending" # pending, running, completed, failed
16
  result: Any = None
17
  error: Optional[str] = None
18
  started_at: Optional[float] = None
19
  finished_at: Optional[float] = None
20
 
21
+
22
  class Workflow(BaseModel):
23
  workflow_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
24
  name: str
25
  steps: List[WorkflowStep]
26
  status: str = "pending"
27
  created_at: float = Field(default_factory=time.time)
28
+ metadata: Dict[str, Any] = Field(default_factory=dict)
29
+
30
 
31
  class WorkflowExecutor:
32
  """
33
+ ARCH-I4.3: Workflow Engine.
34
+
35
+ Coordina workflow in-memory step-by-step tramite Kernel ed Executor. La
36
+ persistenza o il resume inter-processo non sono garantiti da questo motore;
37
+ i caller possono consultare lo stato del workflow corrente tramite
38
+ ``get_workflow``.
39
  """
40
+
41
+ def __init__(self, kernel: Any, executor: Any):
42
  self.kernel = kernel
43
  self.executor = executor
44
  self.active_workflows: Dict[str, Workflow] = {}
45
 
46
+ def get_workflow(self, workflow_id: str) -> Optional[Workflow]:
47
+ """Ritorna il workflow noto, inclusi gli stati terminali in memoria."""
48
+ return self.active_workflows.get(workflow_id)
49
+
50
  async def execute_workflow(self, workflow: Workflow) -> Workflow:
51
+ """Esegue un workflow step-by-step, mantenendo il fallback locale."""
52
  self.active_workflows[workflow.workflow_id] = workflow
53
  workflow.status = "running"
54
+ _logger.info("Avvio workflow: %s (%s)", workflow.name, workflow.workflow_id)
55
 
56
  for step in workflow.steps:
57
  step.status = "running"
58
  step.started_at = time.time()
59
+ _logger.info(
60
+ "Esecuzione step: %s in workflow %s",
61
+ step.tool_name,
62
+ workflow.workflow_id,
63
+ )
64
  try:
65
+ # ARCH-I4.3: il Kernel risolve la capability senza esporre
66
+ # l'infrastruttura al workflow.
67
+ resolution = await self.kernel.resolve_capability(step.tool_name)
68
+ if resolution.get("status") == "resolved":
69
+ worker = resolution["worker"]
 
70
  worker_id = worker.id if hasattr(worker, "id") else worker["id"]
71
+ _logger.info(
72
+ "Step %s risolto su worker: %s",
73
+ step.tool_name,
74
+ worker_id,
75
+ )
76
  result = await self.executor.run_tool(
77
  tool_name=step.tool_name,
78
+ inputs=step.args,
79
+ worker_hint=worker_id,
80
  )
 
 
 
81
  else:
82
+ # Nessun worker registrato: il comportamento storico resta
83
+ # l'esecuzione locale tramite lo stesso Executor.
84
+ _logger.warning(
85
+ "Nessun worker per %s, provo esecuzione locale",
86
+ step.tool_name,
87
+ )
88
+ result = await self.executor.run_tool(
89
+ tool_name=step.tool_name,
90
+ inputs=step.args,
91
+ )
92
+
93
+ step.result = result
94
+ if isinstance(result, dict) and result.get("success") is False:
95
+ step.status = "failed"
96
+ step.error = str(result.get("error", "Tool execution failed"))
97
+ workflow.status = "failed"
98
+ break
99
+ step.status = "completed"
100
+ except Exception as exc:
101
  step.status = "failed"
102
+ step.error = str(exc)
103
  workflow.status = "failed"
104
+ _logger.error("Step %s fallito: %s", step.tool_name, exc)
105
  break
106
+ finally:
107
+ step.finished_at = time.time()
108
 
109
  if workflow.status == "running":
110
  workflow.status = "completed"
111
+ _logger.info("Workflow %s terminato con stato: %s", workflow.name, workflow.status)
 
112
  return workflow
 
api/resolver.py CHANGED
@@ -1,11 +1,71 @@
1
- import time
2
  import logging
3
- from typing import List, Dict, Optional, Any
4
- from pydantic import BaseModel
 
 
 
 
5
  from .marketplace import WORKERS_REGISTRY, WorkerCapability
6
 
7
  _logger = logging.getLogger("api.resolver")
8
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
9
  class ResolverConstraints(BaseModel):
10
  min_version: Optional[str] = None
11
  max_cost: Optional[float] = None
@@ -14,62 +74,99 @@ class ResolverConstraints(BaseModel):
14
  require_gpu: bool = False
15
  min_priority: int = 100
16
 
 
 
 
 
 
 
 
 
17
  class CapabilityResolver:
18
  """
19
  ARCH-E3.2: Capability Resolver
20
  Mappa le capacità richieste dal Brain ai Worker disponibili tramite il Marketplace,
21
  scegliendo il migliore in base agli SLA.
22
  """
23
-
24
  @staticmethod
25
  async def resolve(
26
- capability: str,
27
- constraints: Optional[ResolverConstraints] = None
28
  ) -> Optional[WorkerCapability]:
29
  """
30
  Risolve una capability in un Worker specifico.
31
  Strategia:
32
  1. Filtra per capability supportata.
33
  2. Filtra per worker attivi (last_seen < 300s).
34
- 3. Applica constraints (versione, costo, latenza, GPU).
35
- 4. Ordina per (priority ASC, cost ASC, latency ASC).
 
36
  """
37
  now = int(time.time())
38
  candidates = []
39
-
40
  from .health_manager import health_manager
41
 
42
  for worker in WORKERS_REGISTRY.values():
43
  # 1. & 2. Filtro base + Health Check (ARCH-P5.1)
44
- is_alive = (now - worker.last_seen < 300)
45
  is_healthy = await health_manager.is_healthy(worker.id)
46
-
47
- if capability in worker.capabilities and is_alive and is_healthy:
48
- # 3. Applica constraints
49
- if constraints:
50
- if constraints.min_version and worker.version < constraints.min_version:
51
- continue
52
- if constraints.max_cost is not None and worker.cost > constraints.max_cost:
53
- continue
54
- if constraints.max_latency is not None and worker.latency > constraints.max_latency:
55
- continue
56
- if constraints.require_gpu and not worker.gpu:
 
 
 
 
57
  continue
58
-
59
- candidates.append(worker)
60
-
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
61
  if not candidates:
62
  _logger.warning(f"Nessun worker trovato per capability: {capability}")
63
  return None
64
-
65
- # 4. Ordinamento per SLA
66
  # Priorità: Priority (basso meglio), Cost (basso meglio), Latency (basso meglio)
67
- candidates.sort(key=lambda w: (w.priority, w.cost, w.latency))
68
-
69
  best_worker = candidates[0]
70
- _logger.info(f"Risolta capability '{capability}' su worker '{best_worker.id}' (score: p={best_worker.priority}, c={best_worker.cost}, l={best_worker.latency})")
71
-
 
 
 
 
 
 
72
  return best_worker
73
 
 
74
  # Singleton instance
75
  resolver = CapabilityResolver()
 
 
1
  import logging
2
+ import re
3
+ import time
4
+ from typing import Optional, Tuple
5
+
6
+ from pydantic import BaseModel, field_validator
7
+
8
  from .marketplace import WORKERS_REGISTRY, WorkerCapability
9
 
10
  _logger = logging.getLogger("api.resolver")
11
 
12
+ _SEMVER_IDENTIFIER = r"(?:0|[1-9]\d*|\d*[A-Za-z-][0-9A-Za-z-]*)"
13
+ _SEMVER_PATTERN = re.compile(
14
+ rf"^(?P<major>0|[1-9]\d*)\.(?P<minor>0|[1-9]\d*)\.(?P<patch>0|[1-9]\d*)"
15
+ rf"(?:-(?P<prerelease>{_SEMVER_IDENTIFIER}(?:\.{_SEMVER_IDENTIFIER})*))?"
16
+ rf"(?:\+(?P<build>[0-9A-Za-z-]+(?:\.[0-9A-Za-z-]+)*))?$"
17
+ )
18
+
19
+
20
+ def _parse_semver(version: str) -> Tuple[int, int, int, Optional[Tuple[str, ...]]]:
21
+ """Parsa una versione Semantic Versioning 2.0.0 senza dipendenze esterne."""
22
+ match = _SEMVER_PATTERN.fullmatch(version)
23
+ if not match:
24
+ raise ValueError(f"Invalid Semantic Version: {version!r}")
25
+
26
+ prerelease = match.group("prerelease")
27
+ return (
28
+ int(match.group("major")),
29
+ int(match.group("minor")),
30
+ int(match.group("patch")),
31
+ tuple(prerelease.split(".")) if prerelease else None,
32
+ )
33
+
34
+
35
+ def _compare_semver(left: str, right: str) -> int:
36
+ """Confronta due versioni SemVer, restituendo -1, 0 oppure 1."""
37
+ left_major, left_minor, left_patch, left_prerelease = _parse_semver(left)
38
+ right_major, right_minor, right_patch, right_prerelease = _parse_semver(right)
39
+
40
+ left_core = (left_major, left_minor, left_patch)
41
+ right_core = (right_major, right_minor, right_patch)
42
+ if left_core != right_core:
43
+ return -1 if left_core < right_core else 1
44
+
45
+ if left_prerelease is None and right_prerelease is None:
46
+ return 0
47
+ if left_prerelease is None:
48
+ return 1
49
+ if right_prerelease is None:
50
+ return -1
51
+
52
+ for left_identifier, right_identifier in zip(left_prerelease, right_prerelease):
53
+ if left_identifier == right_identifier:
54
+ continue
55
+
56
+ left_is_numeric = left_identifier.isdigit()
57
+ right_is_numeric = right_identifier.isdigit()
58
+ if left_is_numeric and right_is_numeric:
59
+ return -1 if int(left_identifier) < int(right_identifier) else 1
60
+ if left_is_numeric != right_is_numeric:
61
+ return -1 if left_is_numeric else 1
62
+ return -1 if left_identifier < right_identifier else 1
63
+
64
+ if len(left_prerelease) == len(right_prerelease):
65
+ return 0
66
+ return -1 if len(left_prerelease) < len(right_prerelease) else 1
67
+
68
+
69
  class ResolverConstraints(BaseModel):
70
  min_version: Optional[str] = None
71
  max_cost: Optional[float] = None
 
74
  require_gpu: bool = False
75
  min_priority: int = 100
76
 
77
+ @field_validator("min_version")
78
+ @classmethod
79
+ def validate_min_version(cls, value: Optional[str]) -> Optional[str]:
80
+ if value is not None:
81
+ _parse_semver(value)
82
+ return value
83
+
84
+
85
  class CapabilityResolver:
86
  """
87
  ARCH-E3.2: Capability Resolver
88
  Mappa le capacità richieste dal Brain ai Worker disponibili tramite il Marketplace,
89
  scegliendo il migliore in base agli SLA.
90
  """
91
+
92
  @staticmethod
93
  async def resolve(
94
+ capability: str,
95
+ constraints: Optional[ResolverConstraints] = None,
96
  ) -> Optional[WorkerCapability]:
97
  """
98
  Risolve una capability in un Worker specifico.
99
  Strategia:
100
  1. Filtra per capability supportata.
101
  2. Filtra per worker attivi (last_seen < 300s).
102
+ 3. Applica constraints (versione, costo, latenza, GPU, priorità).
103
+ 4. Se disponibile, preferisce la regione richiesta.
104
+ 5. Ordina per (priority ASC, cost ASC, latency ASC).
105
  """
106
  now = int(time.time())
107
  candidates = []
 
108
  from .health_manager import health_manager
109
 
110
  for worker in WORKERS_REGISTRY.values():
111
  # 1. & 2. Filtro base + Health Check (ARCH-P5.1)
112
+ is_alive = now - worker.last_seen < 300
113
  is_healthy = await health_manager.is_healthy(worker.id)
114
+ if capability not in worker.capabilities or not is_alive or not is_healthy:
115
+ continue
116
+
117
+ # 3. Applica constraints
118
+ if constraints:
119
+ if constraints.min_version:
120
+ try:
121
+ if _compare_semver(worker.version, constraints.min_version) < 0:
122
+ continue
123
+ except ValueError:
124
+ _logger.warning(
125
+ "Worker %s escluso: versione non valida per il vincolo SemVer (%r)",
126
+ worker.id,
127
+ worker.version,
128
+ )
129
  continue
130
+ if constraints.max_cost is not None and worker.cost > constraints.max_cost:
131
+ continue
132
+ if constraints.max_latency is not None and worker.latency > constraints.max_latency:
133
+ continue
134
+ if constraints.require_gpu and not worker.gpu:
135
+ continue
136
+ # Nel Marketplace una priorità più bassa è migliore; min_priority
137
+ # mantiene il nome del contratto esistente come soglia massima accettata.
138
+ if worker.priority > constraints.min_priority:
139
+ continue
140
+
141
+ candidates.append(worker)
142
+
143
+ if constraints and constraints.preferred_region:
144
+ regional_candidates = [
145
+ worker
146
+ for worker in candidates
147
+ if worker.region == constraints.preferred_region
148
+ ]
149
+ if regional_candidates:
150
+ candidates = regional_candidates
151
+
152
  if not candidates:
153
  _logger.warning(f"Nessun worker trovato per capability: {capability}")
154
  return None
155
+
156
+ # 5. Ordinamento per SLA
157
  # Priorità: Priority (basso meglio), Cost (basso meglio), Latency (basso meglio)
158
+ candidates.sort(key=lambda worker: (worker.priority, worker.cost, worker.latency))
 
159
  best_worker = candidates[0]
160
+ _logger.info(
161
+ "Risolta capability '%s' su worker '%s' (score: p=%s, c=%s, l=%s)",
162
+ capability,
163
+ best_worker.id,
164
+ best_worker.priority,
165
+ best_worker.cost,
166
+ best_worker.latency,
167
+ )
168
  return best_worker
169
 
170
+
171
  # Singleton instance
172
  resolver = CapabilityResolver()
api/workflows.py ADDED
@@ -0,0 +1,81 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """API in-memory per l'esecuzione esplicita di workflow tool-based."""
2
+
3
+ from __future__ import annotations
4
+
5
+ import asyncio
6
+ import logging
7
+ from typing import Any, Dict, List
8
+
9
+ from fastapi import APIRouter, Depends, HTTPException, status
10
+ from pydantic import BaseModel, Field
11
+
12
+ from .auth_guard import AuthRole, require_role
13
+ from agents.workflow_engine import Workflow, WorkflowExecutor, WorkflowStep
14
+
15
+ _logger = logging.getLogger("api.workflows")
16
+
17
+ router = APIRouter(
18
+ prefix="/api/workflows",
19
+ tags=["workflows"],
20
+ dependencies=[Depends(require_role(AuthRole.MACHINE))],
21
+ )
22
+
23
+
24
+ class WorkflowStartIn(BaseModel):
25
+ name: str = Field(min_length=1, max_length=200)
26
+ steps: List[WorkflowStep] = Field(min_length=1, max_length=100)
27
+ metadata: Dict[str, Any] = Field(default_factory=dict)
28
+
29
+
30
+ _workflow_executor: WorkflowExecutor | None = None
31
+ _workflow_tasks: Dict[str, asyncio.Task[None]] = {}
32
+
33
+
34
+ def _get_workflow_executor() -> WorkflowExecutor:
35
+ """Costruisce una sola istanza con i singleton Kernel/Executor esistenti."""
36
+ global _workflow_executor
37
+ if _workflow_executor is None:
38
+ from .kernel import kernel
39
+ from .state import _get_executor
40
+
41
+ executor = _get_executor()
42
+ if executor is None:
43
+ raise RuntimeError("Executor non disponibile")
44
+ _workflow_executor = WorkflowExecutor(kernel=kernel, executor=executor)
45
+ return _workflow_executor
46
+
47
+
48
+ async def _run_workflow(workflow: Workflow, executor: WorkflowExecutor) -> None:
49
+ """Esegue in background e conserva sempre uno stato terminale osservabile."""
50
+ try:
51
+ await executor.execute_workflow(workflow)
52
+ except Exception as exc: # defensive: il task non deve fallire silenziosamente
53
+ workflow.status = "failed"
54
+ workflow.metadata["runtime_error"] = str(exc)[:500]
55
+ _logger.exception("Workflow %s terminato con errore inatteso", workflow.workflow_id)
56
+ finally:
57
+ _workflow_tasks.pop(workflow.workflow_id, None)
58
+
59
+
60
+ @router.post("", response_model=Workflow, status_code=status.HTTP_202_ACCEPTED)
61
+ async def start_workflow(body: WorkflowStartIn) -> Workflow:
62
+ """Avvia un workflow esplicito senza bloccare la richiesta HTTP."""
63
+ try:
64
+ executor = _get_workflow_executor()
65
+ except RuntimeError as exc:
66
+ raise HTTPException(status_code=503, detail=str(exc)) from exc
67
+
68
+ workflow = Workflow(name=body.name, steps=body.steps, metadata=body.metadata)
69
+ executor.active_workflows[workflow.workflow_id] = workflow
70
+ task = asyncio.create_task(_run_workflow(workflow, executor))
71
+ _workflow_tasks[workflow.workflow_id] = task
72
+ return workflow
73
+
74
+
75
+ @router.get("/{workflow_id}", response_model=Workflow)
76
+ async def get_workflow(workflow_id: str) -> Workflow:
77
+ """Restituisce lo stato in-memory del workflow, inclusi gli esiti per step."""
78
+ workflow = _get_workflow_executor().get_workflow(workflow_id)
79
+ if workflow is None:
80
+ raise HTTPException(status_code=404, detail=f"Workflow {workflow_id} non trovato")
81
+ return workflow
main.py CHANGED
@@ -130,6 +130,7 @@ _ROUTER_MAP = {
130
  "research": "research",
131
  "agent_memory": "agent_memory",
132
  "agent": "agent",
 
133
  "exec": "exec",
134
  "vault": "vault",
135
  "browser": "browser",
 
130
  "research": "research",
131
  "agent_memory": "agent_memory",
132
  "agent": "agent",
133
+ "workflows": "workflows",
134
  "exec": "exec",
135
  "vault": "vault",
136
  "browser": "browser",
tests/test_capability_resolver.py ADDED
@@ -0,0 +1,173 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import asyncio
2
+ import sys
3
+ import time
4
+ import types
5
+ import unittest
6
+ from enum import IntEnum
7
+ from pathlib import Path
8
+
9
+ from pydantic import ValidationError
10
+
11
+ BACKEND_ROOT = Path(__file__).resolve().parents[1]
12
+ sys.path.insert(0, str(BACKEND_ROOT))
13
+
14
+
15
+ def _install_import_stubs() -> None:
16
+ auth_guard = types.ModuleType("api.auth_guard")
17
+
18
+ class AuthRole(IntEnum):
19
+ MACHINE = 1
20
+
21
+ auth_guard.AuthRole = AuthRole
22
+ auth_guard.require_role = lambda _role: (lambda: None)
23
+ sys.modules["api.auth_guard"] = auth_guard
24
+
25
+ health_module = types.ModuleType("api.health_manager")
26
+
27
+ class HealthyManager:
28
+ async def is_healthy(self, _worker_id: str) -> bool:
29
+ return True
30
+
31
+ health_module.health_manager = HealthyManager()
32
+ sys.modules["api.health_manager"] = health_module
33
+
34
+
35
+ _install_import_stubs()
36
+
37
+ from api.marketplace import WORKERS_REGISTRY, WorkerCapability
38
+ from api.resolver import CapabilityResolver, ResolverConstraints
39
+
40
+
41
+ class CapabilityResolverTests(unittest.IsolatedAsyncioTestCase):
42
+ def setUp(self):
43
+ self._original_registry = dict(WORKERS_REGISTRY)
44
+ WORKERS_REGISTRY.clear()
45
+
46
+ def tearDown(self):
47
+ WORKERS_REGISTRY.clear()
48
+ WORKERS_REGISTRY.update(self._original_registry)
49
+
50
+ def register_worker(self, **overrides) -> WorkerCapability:
51
+ values = {
52
+ "id": "worker",
53
+ "name": "Worker",
54
+ "version": "1.0.0",
55
+ "last_seen": int(time.time()),
56
+ "capabilities": ["vision"],
57
+ "cost": 1.0,
58
+ "latency": 100.0,
59
+ "region": "global",
60
+ "gpu": False,
61
+ "priority": 10,
62
+ }
63
+ values.update(overrides)
64
+ worker = WorkerCapability(**values)
65
+ WORKERS_REGISTRY[worker.id] = worker
66
+ return worker
67
+
68
+ async def test_uses_semver_not_lexicographic_ordering(self):
69
+ self.register_worker(id="v2", version="2.0.0", priority=20)
70
+ self.register_worker(id="v10", version="10.0.0", priority=10)
71
+
72
+ worker = await CapabilityResolver.resolve(
73
+ "vision",
74
+ ResolverConstraints(min_version="2.0.0"),
75
+ )
76
+
77
+ self.assertIsNotNone(worker)
78
+ self.assertEqual(worker.id, "v10")
79
+
80
+ async def test_excludes_prerelease_below_stable_minimum(self):
81
+ self.register_worker(id="candidate", version="2.0.0-rc.1", priority=1)
82
+ self.register_worker(id="stable", version="2.0.0", priority=10)
83
+
84
+ worker = await CapabilityResolver.resolve(
85
+ "vision",
86
+ ResolverConstraints(min_version="2.0.0"),
87
+ )
88
+
89
+ self.assertIsNotNone(worker)
90
+ self.assertEqual(worker.id, "stable")
91
+
92
+ def test_rejects_malformed_minimum_semver(self):
93
+ with self.assertRaises(ValidationError):
94
+ ResolverConstraints(min_version="2.0")
95
+
96
+ async def test_excludes_malformed_worker_version_only_when_constrained(self):
97
+ self.register_worker(id="malformed", version="not-a-version", priority=1)
98
+ self.register_worker(id="valid", version="2.0.0", priority=10)
99
+
100
+ worker = await CapabilityResolver.resolve(
101
+ "vision",
102
+ ResolverConstraints(min_version="2.0.0"),
103
+ )
104
+
105
+ self.assertIsNotNone(worker)
106
+ self.assertEqual(worker.id, "valid")
107
+
108
+ async def test_prefers_requested_region_when_available(self):
109
+ self.register_worker(id="global", region="global", priority=1)
110
+ self.register_worker(id="eu", region="eu-west", priority=20)
111
+
112
+ worker = await CapabilityResolver.resolve(
113
+ "vision",
114
+ ResolverConstraints(preferred_region="eu-west"),
115
+ )
116
+
117
+ self.assertIsNotNone(worker)
118
+ self.assertEqual(worker.id, "eu")
119
+
120
+ async def test_falls_back_when_requested_region_is_unavailable(self):
121
+ self.register_worker(id="global", region="global", priority=1)
122
+
123
+ worker = await CapabilityResolver.resolve(
124
+ "vision",
125
+ ResolverConstraints(preferred_region="eu-west"),
126
+ )
127
+
128
+ self.assertIsNotNone(worker)
129
+ self.assertEqual(worker.id, "global")
130
+
131
+ async def test_applies_priority_threshold_with_lower_values_preferred(self):
132
+ self.register_worker(id="allowed", priority=100, cost=10.0)
133
+ self.register_worker(id="excluded", priority=101, cost=0.0)
134
+
135
+ worker = await CapabilityResolver.resolve(
136
+ "vision",
137
+ ResolverConstraints(min_priority=100),
138
+ )
139
+
140
+ self.assertIsNotNone(worker)
141
+ self.assertEqual(worker.id, "allowed")
142
+
143
+ async def test_preserves_existing_cost_latency_and_gpu_constraints(self):
144
+ self.register_worker(
145
+ id="cpu",
146
+ cost=1.0,
147
+ latency=50.0,
148
+ gpu=False,
149
+ priority=1,
150
+ )
151
+ self.register_worker(
152
+ id="gpu",
153
+ cost=2.0,
154
+ latency=100.0,
155
+ gpu=True,
156
+ priority=10,
157
+ )
158
+
159
+ worker = await CapabilityResolver.resolve(
160
+ "vision",
161
+ ResolverConstraints(
162
+ max_cost=2.0,
163
+ max_latency=100.0,
164
+ require_gpu=True,
165
+ ),
166
+ )
167
+
168
+ self.assertIsNotNone(worker)
169
+ self.assertEqual(worker.id, "gpu")
170
+
171
+
172
+ if __name__ == "__main__":
173
+ unittest.main()
tests/test_workflow_integration.py ADDED
@@ -0,0 +1,175 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import ast
2
+ import asyncio
3
+ import sys
4
+ import types
5
+ import unittest
6
+ from enum import IntEnum
7
+ from pathlib import Path
8
+
9
+ from pydantic import ValidationError
10
+
11
+ BACKEND_ROOT = Path(__file__).resolve().parents[1]
12
+ sys.path.insert(0, str(BACKEND_ROOT))
13
+
14
+
15
+ def _install_auth_stub() -> None:
16
+ auth_guard = types.ModuleType("api.auth_guard")
17
+
18
+ class AuthRole(IntEnum):
19
+ MACHINE = 1
20
+
21
+ auth_guard.AuthRole = AuthRole
22
+ auth_guard.require_role = lambda _role: (lambda: None)
23
+ sys.modules["api.auth_guard"] = auth_guard
24
+
25
+
26
+ _install_auth_stub()
27
+
28
+ from agents.workflow_engine import Workflow, WorkflowExecutor, WorkflowStep
29
+ from api import workflows
30
+
31
+
32
+ class FakeKernel:
33
+ def __init__(self, resolutions):
34
+ self.resolutions = resolutions
35
+ self.calls = []
36
+
37
+ async def resolve_capability(self, tool_name):
38
+ self.calls.append(tool_name)
39
+ return self.resolutions.get(tool_name, {"status": "error"})
40
+
41
+
42
+ class FakeExecutor:
43
+ def __init__(self, results=None):
44
+ self.results = results or {}
45
+ self.calls = []
46
+
47
+ async def run_tool(self, tool_name, inputs, timeout=30.0, worker_hint=None):
48
+ self.calls.append({
49
+ "tool_name": tool_name,
50
+ "inputs": inputs,
51
+ "timeout": timeout,
52
+ "worker_hint": worker_hint,
53
+ })
54
+ result = self.results.get(tool_name, {"success": True, "output": tool_name})
55
+ if isinstance(result, Exception):
56
+ raise result
57
+ return result
58
+
59
+
60
+ class WorkflowExecutorTests(unittest.IsolatedAsyncioTestCase):
61
+ async def test_uses_resolved_worker_hint_and_records_completed_step(self):
62
+ kernel = FakeKernel({"search": {"status": "resolved", "worker": {"id": "worker-eu"}}})
63
+ executor = FakeExecutor()
64
+ engine = WorkflowExecutor(kernel=kernel, executor=executor)
65
+ workflow = Workflow(name="ricerca", steps=[WorkflowStep(tool_name="search", args={"query": "AI"})])
66
+
67
+ result = await engine.execute_workflow(workflow)
68
+
69
+ self.assertEqual(result.status, "completed")
70
+ self.assertEqual(result.steps[0].status, "completed")
71
+ self.assertIsNotNone(result.steps[0].started_at)
72
+ self.assertIsNotNone(result.steps[0].finished_at)
73
+ self.assertEqual(executor.calls[0]["worker_hint"], "worker-eu")
74
+
75
+ async def test_preserves_local_fallback_when_no_worker_is_resolved(self):
76
+ kernel = FakeKernel({})
77
+ executor = FakeExecutor()
78
+ engine = WorkflowExecutor(kernel=kernel, executor=executor)
79
+ workflow = Workflow(name="fallback", steps=[WorkflowStep(tool_name="local", args={"value": 1})])
80
+
81
+ result = await engine.execute_workflow(workflow)
82
+
83
+ self.assertEqual(result.status, "completed")
84
+ self.assertIsNone(executor.calls[0]["worker_hint"])
85
+
86
+ async def test_stops_after_unsuccessful_tool_result(self):
87
+ kernel = FakeKernel({})
88
+ executor = FakeExecutor({"first": {"success": False, "error": "denied"}})
89
+ engine = WorkflowExecutor(kernel=kernel, executor=executor)
90
+ workflow = Workflow(
91
+ name="errore",
92
+ steps=[
93
+ WorkflowStep(tool_name="first", args={}),
94
+ WorkflowStep(tool_name="second", args={}),
95
+ ],
96
+ )
97
+
98
+ result = await engine.execute_workflow(workflow)
99
+
100
+ self.assertEqual(result.status, "failed")
101
+ self.assertEqual(result.steps[0].status, "failed")
102
+ self.assertEqual(result.steps[0].error, "denied")
103
+ self.assertEqual(result.steps[1].status, "pending")
104
+ self.assertEqual([call["tool_name"] for call in executor.calls], ["first"])
105
+
106
+
107
+ class WorkflowApiTests(unittest.IsolatedAsyncioTestCase):
108
+ def setUp(self):
109
+ self.previous_executor = workflows._workflow_executor
110
+ self.previous_tasks = workflows._workflow_tasks.copy()
111
+ workflows._workflow_executor = None
112
+ workflows._workflow_tasks.clear()
113
+
114
+ async def asyncTearDown(self):
115
+ for task in list(workflows._workflow_tasks.values()):
116
+ task.cancel()
117
+ try:
118
+ await task
119
+ except asyncio.CancelledError:
120
+ pass
121
+ workflows._workflow_tasks.clear()
122
+ workflows._workflow_executor = self.previous_executor
123
+ workflows._workflow_tasks.update(self.previous_tasks)
124
+
125
+ def test_router_exposes_start_and_status_paths(self):
126
+ routes = {
127
+ (route.path, method)
128
+ for route in workflows.router.routes
129
+ for method in (getattr(route, "methods", set()) or set())
130
+ }
131
+
132
+ self.assertEqual(workflows.router.prefix, "/api/workflows")
133
+ self.assertIn(("/api/workflows", "POST"), routes)
134
+ self.assertIn(("/api/workflows/{workflow_id}", "GET"), routes)
135
+
136
+ def test_main_router_map_mounts_workflows(self):
137
+ tree = ast.parse((BACKEND_ROOT / "main.py").read_text(encoding="utf-8"))
138
+ router_map = next(
139
+ node.value
140
+ for node in tree.body
141
+ if isinstance(node, ast.Assign)
142
+ and any(isinstance(target, ast.Name) and target.id == "_ROUTER_MAP" for target in node.targets)
143
+ )
144
+ routes = {
145
+ key.value: value.value
146
+ for key, value in zip(router_map.keys, router_map.values)
147
+ if isinstance(key, ast.Constant) and isinstance(value, ast.Constant)
148
+ }
149
+
150
+ self.assertEqual(routes.get("workflows"), "workflows")
151
+
152
+ def test_start_contract_rejects_empty_steps(self):
153
+ with self.assertRaises(ValidationError):
154
+ workflows.WorkflowStartIn(name="vuoto", steps=[])
155
+
156
+ async def test_start_and_get_workflow_use_background_execution(self):
157
+ engine = WorkflowExecutor(kernel=FakeKernel({}), executor=FakeExecutor())
158
+ workflows._workflow_executor = engine
159
+ body = workflows.WorkflowStartIn(
160
+ name="API workflow",
161
+ steps=[WorkflowStep(tool_name="local", args={"n": 1})],
162
+ )
163
+
164
+ started = await workflows.start_workflow(body)
165
+ task = workflows._workflow_tasks[started.workflow_id]
166
+ await task
167
+ fetched = await workflows.get_workflow(started.workflow_id)
168
+
169
+ self.assertEqual(fetched.workflow_id, started.workflow_id)
170
+ self.assertEqual(fetched.status, "completed")
171
+ self.assertEqual(fetched.steps[0].status, "completed")
172
+
173
+
174
+ if __name__ == "__main__":
175
+ unittest.main()