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()