Spaces:
Running
Running
File size: 2,884 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 | import time
import logging
from typing import List, Dict, Optional, Any
from pydantic import BaseModel
from .marketplace import WORKERS_REGISTRY, WorkerCapability
_logger = logging.getLogger("api.resolver")
class ResolverConstraints(BaseModel):
min_version: Optional[str] = None
max_cost: Optional[float] = None
max_latency: Optional[float] = None
preferred_region: Optional[str] = None
require_gpu: bool = False
min_priority: int = 100
class CapabilityResolver:
"""
ARCH-E3.2: Capability Resolver
Mappa le capacità richieste dal Brain ai Worker disponibili tramite il Marketplace,
scegliendo il migliore in base agli SLA.
"""
@staticmethod
async def resolve(
capability: str,
constraints: Optional[ResolverConstraints] = None
) -> Optional[WorkerCapability]:
"""
Risolve una capability in un Worker specifico.
Strategia:
1. Filtra per capability supportata.
2. Filtra per worker attivi (last_seen < 300s).
3. Applica constraints (versione, costo, latenza, GPU).
4. Ordina per (priority ASC, cost ASC, latency ASC).
"""
now = int(time.time())
candidates = []
from .health_manager import health_manager
for worker in WORKERS_REGISTRY.values():
# 1. & 2. Filtro base + Health Check (ARCH-P5.1)
is_alive = (now - worker.last_seen < 300)
is_healthy = await health_manager.is_healthy(worker.id)
if capability in worker.capabilities and is_alive and is_healthy:
# 3. Applica constraints
if constraints:
if constraints.min_version and worker.version < constraints.min_version:
continue
if constraints.max_cost is not None and worker.cost > constraints.max_cost:
continue
if constraints.max_latency is not None and worker.latency > constraints.max_latency:
continue
if constraints.require_gpu and not worker.gpu:
continue
candidates.append(worker)
if not candidates:
_logger.warning(f"Nessun worker trovato per capability: {capability}")
return None
# 4. Ordinamento per SLA
# Priorità: Priority (basso meglio), Cost (basso meglio), Latency (basso meglio)
candidates.sort(key=lambda w: (w.priority, w.cost, w.latency))
best_worker = candidates[0]
_logger.info(f"Risolta capability '{capability}' su worker '{best_worker.id}' (score: p={best_worker.priority}, c={best_worker.cost}, l={best_worker.latency})")
return best_worker
# Singleton instance
resolver = CapabilityResolver()
|