| """
|
| 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
|
| """
|
|
|
| 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
|
|
|
|
|
| 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
|
| """
|
|
|
| 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
|
| )
|
|
|
|
|
| gpu_util = self._normalize_gpu_utilization(gpu_utilization)
|
|
|
|
|
| avg_sla = 0.0
|
| if sla_metrics:
|
| avg_sla = sla_metrics.get("avg_urgency", 0.0)
|
|
|
|
|
| normalized_latency = self._normalize(
|
| avg_job_latency_ms,
|
| self.max_latency_ms
|
| )
|
|
|
|
|
| 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
|
| """
|
|
|
| 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.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"],
|
| )
|
|
|
|
|
|
|
| _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",
|
| ]
|
|
|