aegislm / backend /autotune /state_encoder.py
ACA050's picture
Upload 28 files
3e4aa3e verified
Raw
History Blame Contribute Delete
11.1 kB
"""
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",
]