Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -282,7 +282,8 @@ simulation-runs/sim-2026-04-02_14_30_45/
"scale_latency": 40,
"max_workers": 50,
"policy": "MaxPrefix", // RoundRobin, LeastLoaded, Random, MaxPrefix, Balanced
"periodic_infra_update_collection_time": 30,
"latency_percentile": "p95", // p90, p95, p99
"latency_window": 50, // last N completed requests for TTFT/ITL SLOs
"max_event_batch_size": 64
}
},
Expand Down
5 changes: 3 additions & 2 deletions configs/defaults.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,8 @@
"scale_latency": 40,
"max_workers": 50,
"policy": "MaxPrefix",
"periodic_infra_update_collection_time": 30,
"latency_percentile": "p95",
"latency_window": 50,
"max_event_batch_size": 64
}
},
Expand Down Expand Up @@ -56,7 +57,7 @@

"worker_params":{
"worker_local_queue_capacity": 1,
"periodic_infra_update_time": 30,
"periodic_infra_update_time": 5,
"kvcevent_coalesce_time": 30
},

Expand Down
5 changes: 3 additions & 2 deletions configs/defaults_otel.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,8 @@
"scale_latency": 40,
"max_workers": 50,
"policy": "MaxPrefix",
"periodic_infra_update_collection_time": 30,
"latency_percentile": "p95",
"latency_window": 50,
"max_event_batch_size": 64
}
},
Expand Down Expand Up @@ -59,7 +60,7 @@

"worker_params":{
"worker_local_queue_capacity": 1,
"periodic_infra_update_time": 30,
"periodic_infra_update_time": 5,
"kvcevent_coalesce_time": 30
},

Expand Down
5 changes: 3 additions & 2 deletions configs/hf.json
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,8 @@
"scale_latency": 40,
"max_workers": 50,
"policy": "MaxPrefix",
"periodic_infra_update_collection_time": 30,
"latency_percentile": "p95",
"latency_window": 50,
"max_event_batch_size": 64
}
},
Expand Down Expand Up @@ -55,7 +56,7 @@

"worker_params":{
"worker_local_queue_capacity": 1,
"periodic_infra_update_time": 30,
"periodic_infra_update_time": 5,
"kvcevent_coalesce_time": 30
},

Expand Down
5 changes: 4 additions & 1 deletion opal/core/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,5 +34,8 @@ class SystemEvent(OpalInfraEvent):
# 0 = min, 1 = max
load: float
ingress_queue_occupancy: float
mem_used: float
gpu_utilization: float
# kvc_utilization: fraction of GPU KV-cache blocks in use (1 - free/total), in [0, 1].
kvc_utilization: float = 0.0
# queue_depth: absolute in-flight request count on the worker (waiting + running).
queue_depth: int = 0
78 changes: 58 additions & 20 deletions opal/router/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
from opal.core.events import OpalInfraEvent, KVCEvent, SystemEvent
from opal.kvcache.kvbm import KVBM
from opal.core.request import LLMRequest
from opal.stats.metrics import MetricsSnapshot
from opal.worker.vllm_worker import LLMWorkerVLLMScheduler
from opal.utils.util import parse_bool, safe_process

Expand Down Expand Up @@ -39,6 +40,19 @@ def _get_policy_func(self):
f"{policy} no such routing policy. Supported policies are: RoundRobin, LeastLoaded, Random, MaxPrefix, Balanced"
)

def _get_latency_percentile(self) -> int:
name = str(self.opalConfig["router"]["router_params"]["latency_percentile"]).lower()
mapping = {"p90": 90, "p95": 95, "p99": 99}
if name not in mapping:
raise Exception(f"{name} is not a supported latency percentile. Supported: p90, p95, p99")
return mapping[name]

def _get_latency_window(self) -> int:
window = int(self.opalConfig["router"]["router_params"]["latency_window"])
if window <= 0:
raise Exception(f"latency_window must be a positive integer, got {window}")
return window

def __init__(self, opal_env, opal_config):
self.opal_env = opal_env
self.opalConfig = opal_config
Expand All @@ -51,12 +65,16 @@ def __init__(self, opal_env, opal_config):
self.results_queue = simpy.Store(self.sim_env)
# leave infinite capacity for this
self._event_queue = simpy.Store(self.sim_env)
self.periodic_infra_update_collection_time = self.opalConfig["router"]["router_params"][
"periodic_infra_update_collection_time"
]
self._kvbm = KVBM(self.opal_env)
self._worker_stats = None
self._policy_func = self._get_policy_func()
self._latency_percentile = self._get_latency_percentile()
self._latency_window = self._get_latency_window()

# Last SystemEvent reported by each worker and the latest
# MetricsSnapshot built from them.
self._latest_worker_metrics: dict[int, SystemEvent] = {}
self._latest_metrics: MetricsSnapshot | None = None

self.num_workers = self.opalConfig["simulation"]["num_workers"]
self._worker_cls = LLMWorkerVLLMScheduler
Expand Down Expand Up @@ -149,6 +167,25 @@ def _per_second_stats(self):
stats.add_per_unit_gpu_utilization(utilization)
self.log.debug(f"Breaking the per second stats loop at {self.sim_env.now}")

def _pool_metrics_snapshot(self, stats):
"""Create a MetricsSnapshot from the last SystemEvent each worker pushed,
plus SLO percentiles over the configured sliding window of completions.

Called when process_events() drains a batch that contains SystemEvents.
"""
ttft, itl = stats.recent_ttft_itl(self._latency_percentile, self._latency_window)
snapshot = MetricsSnapshot(
timestamp=self.sim_env.now,
queue_depth_per_worker={wid: ev.queue_depth for wid, ev in self._latest_worker_metrics.items()},
kvc_util_per_worker={wid: ev.kvc_utilization for wid, ev in self._latest_worker_metrics.items()},
percentile=self._latency_percentile,
window=self._latency_window,
ttft_secs=ttft,
itl_secs=itl,
)
self._latest_metrics = snapshot
stats.add_metrics_snapshot(snapshot.to_dict())

def _policy_leastloaded(self, req: LLMRequest):
queue_size = min(self._outstanding_requests_per_worker.values())
worker = next((k for k, v in self._outstanding_requests_per_worker.items() if v == queue_size), None)
Expand Down Expand Up @@ -245,38 +282,39 @@ def queue_events(self, elist: list[OpalInfraEvent], delay: float = 0):
for elem in elist:
yield self._event_queue.put(elem)

def _ingest_event(self, e: OpalInfraEvent, kvbm_events: list[KVCEvent], systems_events: list[SystemEvent]):
if isinstance(e, KVCEvent):
kvbm_events.append(e)
elif isinstance(e, SystemEvent):
systems_events.append(e)
# Store last reported SystemEvent from this worker. Used to create telemetry snapshot
self._latest_worker_metrics[e.worker_id] = e
else:
raise Exception(f"Unknown event type {type(e)}")

def process_events(self):
max_event_batch_size = self.opalConfig["router"]["router_params"]["max_event_batch_size"]

while not self.opal_env.are_we_done():
kvbm_events = []
systems_events = []

# Collect events up to batch size
e = yield self._event_queue.get()
self._ingest_event(e, kvbm_events, systems_events)

while len(self._event_queue.items) > 0 and not self.opal_env.are_we_done():
# Stop if either batch is full
if len(kvbm_events) >= max_event_batch_size or len(systems_events) >= max_event_batch_size:
break

e = yield self._event_queue.get()
if isinstance(e, KVCEvent):
kvbm_events.append(e)
elif isinstance(e, SystemEvent):
systems_events.append(e)
else:
raise Exception(f"Unknown event type {type(e)}")

# Process batches if we have any events
self._ingest_event(e, kvbm_events, systems_events)

if kvbm_events or systems_events:
self._kvbm.process_kvc_events(kvbm_events)
self._kvbm.process_system_events(systems_events)

# If queue still has items, continue immediately
if len(self._event_queue.items) > 0:
continue

# Otherwise sleep until next periodic check
yield self.sim_env.timeout(self.periodic_infra_update_collection_time)
if systems_events:
stats = self.opal_env.workload_orchestrator.get_active_stage_stats()
self._pool_metrics_snapshot(stats)

def shutdown(self):
self._stats_request_allocated_per_worker = dict(sorted(self._stats_request_allocated_per_worker.items()))
Expand Down
48 changes: 48 additions & 0 deletions opal/stats/metrics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
# SPDX-License-Identifier: Apache-2.0
from __future__ import annotations
from dataclasses import dataclass, field


@dataclass
class MetricsSnapshot:
"""A point-in-time view of cluster telemetry as seen by the router.

Built when workers push SystemEvent telemetry.
Worker fields are the last report from each worker (delayed by
periodic_infra_update_time). TTFT/ITL are the configured percentile over a
sliding window of recent completed requests.

Units:
- queue_depth_per_worker: absolute in-flight request count per worker.
- kvc_util_per_worker: fraction in [0, 1].
- ttft_secs / itl_secs: seconds; -1 until `window` requests have completed.
"""

timestamp: float = 0.0
queue_depth_per_worker: dict[int, int] = field(default_factory=dict)
kvc_util_per_worker: dict[int, float] = field(default_factory=dict)
percentile: int = 95 # 90, 95, or 99
window: int = 50 # last N completed requests used for ttft/itl
ttft_secs: float = -1.0
itl_secs: float = -1.0

@property
def max_queue_depth(self) -> int:
return max(self.queue_depth_per_worker.values(), default=0)

@property
def max_kvc_util(self) -> float:
return max(self.kvc_util_per_worker.values(), default=0.0)

def to_dict(self) -> dict:
return {
"timestamp": self.timestamp,
"queue_depth_per_worker": dict(self.queue_depth_per_worker),
"kvc_util_per_worker": dict(self.kvc_util_per_worker),
"percentile": self.percentile,
"window": self.window,
"ttft_secs": self.ttft_secs,
"itl_secs": self.itl_secs,
"max_queue_depth": self.max_queue_depth,
"max_kvc_util": self.max_kvc_util,
}
23 changes: 23 additions & 0 deletions opal/stats/stage_statistics.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ def __init__(self):
self.per_unit_req_done = []
self.per_unit_workers = []
self.per_unit_gpu_utilization = []
# List of all MetricsSnapshot objects recorded during the simulation.
self.metrics_snapshots = []
while val < max_bin:
self.bins.append(val)
val *= 2
Expand Down Expand Up @@ -59,6 +61,7 @@ def to_dict(self):
"per_unit_throughput": self.per_unit_req_done,
"per_unit_workers": self.per_unit_workers,
"per_unit_gpu_utilization": self.per_unit_gpu_utilization,
"metrics_snapshots": self.metrics_snapshots,
# numpy array -> list
"latencies": self.latencies.tolist(),
"stage_time_start": self.stage_time_start,
Expand All @@ -85,6 +88,7 @@ def from_dict(cls, data):
obj.per_unit_req_done = data["per_unit_throughput"]
obj.per_unit_workers = data["per_unit_workers"]
obj.per_unit_gpu_utilization = data["per_unit_gpu_utilization"]
obj.metrics_snapshots = data.get("metrics_snapshots", [])
obj.stage_time_end = data["stage_time_end"]
obj.stage_time_start = data["stage_time_start"]
obj.kvc_tier_tokens = defaultdict(int, data.get("kvc_tier_tokens", {}))
Expand Down Expand Up @@ -226,6 +230,25 @@ def add_per_unit_workers(self, worker_count: int):
def add_per_unit_workdone(self, done: int):
self.per_unit_req_done.append(done)

def add_metrics_snapshot(self, snapshot: dict):
self.metrics_snapshots.append(snapshot)

def recent_ttft_itl(self, percentile: int, window: int) -> Tuple[float, float]:
"""Return (ttft_secs, itl_secs) at `percentile` over the last `window` completed requests.

Returns (-1.0, -1.0) until at least `window` requests have finished.
"""
if window <= 0:
raise ValueError(f"window must be a positive integer, got {window}")
if len(self.raw_ttft_values) < window:
return -1.0, -1.0
recent_ttft = self.raw_ttft_values[-window:]
ttft = float(np.percentile(recent_ttft, percentile))
recent_decode = self.raw_decode_values[-window:]
_, all_itls = self._calculate_itl_tpot(recent_decode)
itl = float(np.percentile(all_itls, percentile)) if len(all_itls) > 0 else -1.0
return ttft, itl

def sample_workdone_per_K(self, k: int = 1):
return sample_series_K(self.per_unit_req_done, k)

Expand Down
20 changes: 15 additions & 5 deletions opal/webserver/index.html
Original file line number Diff line number Diff line change
Expand Up @@ -325,9 +325,19 @@ <h1>OPAL config builder</h1>
</div>

<div class="field">
<label>periodic_infra_update_collection_time</label>
<input type="number" name="router.router_params.periodic_infra_update_collection_time" value="30" />
<div class="hint">Virtual seconds between router-side infra-event drains (KV events from workers).</div>
<label>latency_percentile</label>
<select name="router.router_params.latency_percentile">
<option>p90</option>
<option selected>p95</option>
<option>p99</option>
</select>
<div class="hint">Percentile used for snapshot TTFT/ITL SLOs. Computed over <code>latency_window</code> recent completions.</div>
</div>

<div class="field">
<label>latency_window</label>
<input type="number" name="router.router_params.latency_window" value="50" />
<div class="hint">Number of most recent completed requests used for TTFT/ITL percentiles. Snapshots store -1 until this many requests have finished.</div>
</div>

<div class="field">
Expand Down Expand Up @@ -366,8 +376,8 @@ <h3 class="subhead">worker_params</h3>

<div class="field">
<label>periodic_infra_update_time</label>
<input type="number" name="worker.worker_params.periodic_infra_update_time" value="30" />
<div class="hint">Virtual seconds between worker→router status pushes (gpu_util, queue depth, mem_used).</div>
<input type="number" name="worker.worker_params.periodic_infra_update_time" value="5" />
<div class="hint">Virtual seconds between worker→router status pushes (queue depth, kvc_util). Snapshots are taken when these arrive.</div>
</div>

<div class="field">
Expand Down
Loading