Spaces:
Running
Running
| 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. | |
| """ | |
| 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() | |