""" State Encoder for Autotune Module Encodes system metrics into state vectors for the RL agent. Provides normalization and feature extraction from scheduler metrics. """ from typing import Dict, Any, Optional import logging from .schemas import SystemState, QueueLevel logger = logging.getLogger(__name__) class StateEncoder: """ Encodes system metrics into state vectors for RL agent. State vector: [Q_h, Q_m, Q_l, U_gpu, SLA_avg, CostPressure, DriftScore, FailureRate] Provides: - Metric normalization (0-1 range) - Queue length encoding - Resource utilization encoding - SLA and cost pressure encoding - Drift and failure rate encoding """ def __init__( self, max_queue_high: int = 1000, max_queue_medium: int = 5000, max_queue_low: int = 10000, max_gpu_count: int = 100, max_tenant_count: int = 1000, max_latency_ms: float = 60000.0, ): """ Initialize state encoder. Args: max_queue_high: Maximum high priority queue size max_queue_medium: Maximum medium priority queue size max_queue_low: Maximum low priority queue size max_gpu_count: Maximum GPU count in cluster max_tenant_count: Maximum tenant count max_latency_ms: Maximum expected latency in ms """ # Normalization bounds self.max_queue_high = max_queue_high self.max_queue_medium = max_queue_medium self.max_queue_low = max_queue_low self.max_gpu_count = max_gpu_count self.max_tenant_count = max_tenant_count self.max_latency_ms = max_latency_ms # Default values for missing metrics self.defaults = { "queue_high": 0.0, "queue_medium": 0.0, "queue_low": 0.0, "gpu_utilization": 0.5, "avg_sla_urgency": 0.0, "cost_pressure": 0.0, "drift_score": 0.0, "failure_rate": 0.0, "tenant_count": 1, "avg_job_latency_ms": 1000.0, } def encode( self, queue_lengths: Dict[str, int], gpu_utilization: float, sla_metrics: Optional[Dict[str, float]] = None, cost_pressure: float = 0.0, drift_score: float = 0.0, failure_rate: float = 0.0, tenant_count: int = 1, avg_job_latency_ms: float = 1000.0, ) -> SystemState: """ Encode raw metrics into SystemState. Args: queue_lengths: Dictionary with 'high', 'medium', 'low' queue sizes gpu_utilization: GPU utilization (0-1 or percentage 0-100) sla_metrics: Optional SLA metrics dictionary cost_pressure: Cost pressure metric (0-1) drift_score: Model drift severity (0-1) failure_rate: Job failure rate (0-1) tenant_count: Number of active tenants avg_job_latency_ms: Average job latency in ms Returns: SystemState with normalized values """ # Normalize queue lengths queue_high = self._normalize( queue_lengths.get("high", 0), self.max_queue_high ) queue_medium = self._normalize( queue_lengths.get("medium", 0), self.max_queue_medium ) queue_low = self._normalize( queue_lengths.get("low", 0), self.max_queue_low ) # Normalize GPU utilization gpu_util = self._normalize_gpu_utilization(gpu_utilization) # Extract SLA metrics avg_sla = 0.0 if sla_metrics: avg_sla = sla_metrics.get("avg_urgency", 0.0) # Normalize latency normalized_latency = self._normalize( avg_job_latency_ms, self.max_latency_ms ) # Create system state state = SystemState( queue_high=queue_high, queue_medium=queue_medium, queue_low=queue_low, gpu_utilization=gpu_util, avg_sla_urgency=avg_sla, cost_pressure=cost_pressure, drift_score=drift_score, failure_rate=failure_rate, tenant_count=tenant_count, avg_job_latency_ms=avg_job_latency_ms, ) return state def encode_from_scheduler( self, scheduler_policy: Any, gpu_utilization: float = 0.5, sla_metrics: Optional[Dict[str, float]] = None, cost_pressure: float = 0.0, drift_score: float = 0.0, failure_rate: float = 0.0, tenant_count: int = 1, avg_job_latency_ms: float = 1000.0, ) -> SystemState: """ Encode system state from scheduler policy engine. Args: scheduler_policy: SchedulingPolicyEngine instance gpu_utilization: Current GPU utilization sla_metrics: Optional SLA metrics cost_pressure: Cost pressure drift_score: Drift score failure_rate: Failure rate tenant_count: Tenant count avg_job_latency_ms: Average job latency Returns: SystemState with normalized values """ # Get queue sizes from scheduler queue_sizes = scheduler_policy.get_all_queue_sizes() queue_lengths = { "high": queue_sizes.get("high", 0), "medium": queue_sizes.get("medium", 0), "low": queue_sizes.get("low", 0), } return self.encode( queue_lengths=queue_lengths, gpu_utilization=gpu_utilization, sla_metrics=sla_metrics, cost_pressure=cost_pressure, drift_score=drift_score, failure_rate=failure_rate, tenant_count=tenant_count, avg_job_latency_ms=avg_job_latency_ms, ) def encode_from_metrics( self, metrics: Dict[str, Any], ) -> SystemState: """ Encode system state from metrics dictionary. Expects metrics to have: - queue_lengths: Dict[str, int] - gpu_utilization: float - avg_sla_urgency: float (optional) - cost_pressure: float - drift_score: float - failure_rate: float - tenant_count: int (optional) - avg_job_latency_ms: float (optional) Args: metrics: Dictionary of system metrics Returns: SystemState with normalized values """ queue_lengths = metrics.get("queue_lengths", {}) gpu_utilization = metrics.get("gpu_utilization", 0.5) cost_pressure = metrics.get("cost_pressure", 0.0) drift_score = metrics.get("drift_score", 0.0) failure_rate = metrics.get("failure_rate", 0.0) tenant_count = metrics.get("tenant_count", 1) avg_job_latency_ms = metrics.get("avg_job_latency_ms", 1000.0) sla_metrics = None if "avg_sla_urgency" in metrics: sla_metrics = {"avg_urgency": metrics["avg_sla_urgency"]} return self.encode( queue_lengths=queue_lengths, gpu_utilization=gpu_utilization, sla_metrics=sla_metrics, cost_pressure=cost_pressure, drift_score=drift_score, failure_rate=failure_rate, tenant_count=tenant_count, avg_job_latency_ms=avg_job_latency_ms, ) def _normalize(self, value: float, max_value: float) -> float: """ Normalize value to [0, 1] range. Args: value: Raw value max_value: Maximum expected value Returns: Normalized value in [0, 1] """ if max_value <= 0: return 0.0 normalized = value / max_value return min(1.0, max(0.0, normalized)) def _normalize_gpu_utilization(self, value: float) -> float: """ Normalize GPU utilization to [0, 1] range. Handles both decimal (0-1) and percentage (0-100) formats. Args: value: GPU utilization value Returns: Normalized value in [0, 1] """ # If value > 1, assume it's in percentage form if value > 1.0: value = value / 100.0 return min(1.0, max(0.0, value)) def get_state_bounds(self) -> Dict[str, tuple]: """ Get the bounds for each state dimension. Returns: Dictionary mapping state field to (min, max) tuple """ return { "queue_high": (0.0, 1.0), "queue_medium": (0.0, 1.0), "queue_low": (0.0, 1.0), "gpu_utilization": (0.0, 1.0), "avg_sla_urgency": (0.0, 1.0), "cost_pressure": (0.0, 1.0), "drift_score": (0.0, 1.0), "failure_rate": (0.0, 1.0), } def get_default_state(self) -> SystemState: """ Get default system state. Returns: SystemState with default values """ return SystemState( queue_high=self.defaults["queue_high"], queue_medium=self.defaults["queue_medium"], queue_low=self.defaults["queue_low"], gpu_utilization=self.defaults["gpu_utilization"], avg_sla_urgency=self.defaults["avg_sla_urgency"], cost_pressure=self.defaults["cost_pressure"], drift_score=self.defaults["drift_score"], failure_rate=self.defaults["failure_rate"], tenant_count=self.defaults["tenant_count"], avg_job_latency_ms=self.defaults["avg_job_latency_ms"], ) # Global encoder instance _state_encoder: Optional[StateEncoder] = None def get_state_encoder( max_queue_high: int = 1000, max_queue_medium: int = 5000, max_queue_low: int = 10000, ) -> StateEncoder: """ Get or create global StateEncoder instance. Args: max_queue_high: Max high priority queue size max_queue_medium: Max medium priority queue size max_queue_low: Max low priority queue size Returns: StateEncoder instance """ global _state_encoder if _state_encoder is None: _state_encoder = StateEncoder( max_queue_high=max_queue_high, max_queue_medium=max_queue_medium, max_queue_low=max_queue_low, ) return _state_encoder __all__ = [ "StateEncoder", "get_state_encoder", ]