Repository navigation
fix(metrics): engine metrics from PROMETHEUS_MULTIPROC_DIR + graceful bind - #1075
ConnorLi96 wants to merge 9 commits into
Conversation
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughAdds Python ServeOrchestrator Prometheus multiprocess dir lifecycle for gRPC mode; introduces a Rust WorkerManager and LoadMonitor for worker telemetry, load polling, and engine-metrics aggregation (including a Python subprocess path); integrates engine metrics into /metrics; and adds timing/tracing and health endpoints. Changes
Sequence Diagram(s)sequenceDiagram
participant MetricsServer as SMG Metrics Server
participant WM as WorkerManager
participant HTTP as HTTP Worker
participant GRPC as gRPC Worker
participant Py as Python Subprocess
participant Dir as PROMETHEUS_MULTIPROC_DIR
MetricsServer->>WM: get_engine_metrics()
activate WM
WM->>HTTP: HTTP GET /metrics (fan-out for HTTP workers)
HTTP-->>WM: metrics text
Note over WM,GRPC: gRPC workers use multiprocess files
WM->>GRPC: identify gRPC workers
WM->>Py: spawn `python3` subprocess to read Dir and emit Prometheus text
activate Py
Py->>Dir: read multiprocess files
Py-->>WM: metrics text
deactivate Py
WM->>WM: aggregate_metrics(all texts)
WM-->>MetricsServer: aggregated metrics text
deactivate WM
MetricsServer-->>MetricsServer: render combined /metrics response
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
| Ok(text) if !text.trim().is_empty() => { | ||
| metric_packs.push(MetricPack { | ||
| labels: vec![], | ||
| metrics_text: text, | ||
| }); | ||
| } | ||
| Ok(_) => { | ||
| // No metrics available yet from gRPC workers — skip silently | ||
| } |
There was a problem hiding this comment.
🟡 Nit: The Ok(_) arm on line 345 is dead code. collect_prometheus_multiproc_metrics() already returns Err when the output is empty (lines 116-118), so the Ok(text) variant will always contain non-empty text, making the if !text.trim().is_empty() guard always true and the Ok(_) branch unreachable.
Consider simplifying:
| Ok(text) if !text.trim().is_empty() => { | |
| metric_packs.push(MetricPack { | |
| labels: vec![], | |
| metrics_text: text, | |
| }); | |
| } | |
| Ok(_) => { | |
| // No metrics available yet from gRPC workers — skip silently | |
| } | |
| Ok(text) => { | |
| metric_packs.push(MetricPack { | |
| labels: vec![], | |
| metrics_text: text, | |
| }); | |
| } |
| let output = tokio::process::Command::new("python3") | ||
| .args([ | ||
| "-c", | ||
| "import sys\n\ | ||
| from prometheus_client import CollectorRegistry, generate_latest\n\ | ||
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | ||
| registry = CollectorRegistry()\n\ | ||
| MultiProcessCollector(registry)\n\ | ||
| sys.stdout.buffer.write(generate_latest(registry))\n", | ||
| ]) | ||
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | ||
| .output() | ||
| .await | ||
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; |
There was a problem hiding this comment.
🟡 Nit: No timeout on the python3 subprocess. If the process hangs (e.g., corrupt .db file in the multiproc dir, or python3 not found and the OS stalls), the metrics endpoint request will hang indefinitely. Since this runs on every /metrics scrape, a stuck process could accumulate and exhaust resources.
Consider adding a timeout, e.g. via tokio::time::timeout:
let output = tokio::time::timeout(
std::time::Duration::from_secs(5),
tokio::process::Command::new("python3")
.args([...])
.env("PROMETHEUS_MULTIPROC_DIR", &dir)
.output(),
)
.await
.map_err(|_| "prometheus collector timed out after 5s".to_string())?
.map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?;| def _cleanup_prometheus_dir(self) -> None: | ||
| """Remove the temporary prometheus multiprocess directory and its .db files.""" | ||
| if self._prometheus_dir is None: | ||
| return | ||
| try: | ||
| shutil.rmtree(self._prometheus_dir) | ||
| logger.info("Cleaned up PROMETHEUS_MULTIPROC_DIR=%s", self._prometheus_dir) | ||
| except OSError as e: | ||
| logger.warning( | ||
| "Failed to clean up PROMETHEUS_MULTIPROC_DIR=%s: %s", | ||
| self._prometheus_dir, | ||
| e, | ||
| ) | ||
| self._prometheus_dir = None |
There was a problem hiding this comment.
🟡 Nit: After removing the directory, the PROMETHEUS_MULTIPROC_DIR environment variable still points to the now-deleted path. If any code (e.g., a straggling atexit handler or a library import) reads the env var after cleanup, it will reference a non-existent directory, which can cause confusing errors.
Consider also clearing the env var:
| def _cleanup_prometheus_dir(self) -> None: | |
| """Remove the temporary prometheus multiprocess directory and its .db files.""" | |
| if self._prometheus_dir is None: | |
| return | |
| try: | |
| shutil.rmtree(self._prometheus_dir) | |
| logger.info("Cleaned up PROMETHEUS_MULTIPROC_DIR=%s", self._prometheus_dir) | |
| except OSError as e: | |
| logger.warning( | |
| "Failed to clean up PROMETHEUS_MULTIPROC_DIR=%s: %s", | |
| self._prometheus_dir, | |
| e, | |
| ) | |
| self._prometheus_dir = None | |
| def _cleanup_prometheus_dir(self) -> None: | |
| """Remove the temporary prometheus multiprocess directory and its .db files.""" | |
| if self._prometheus_dir is None: | |
| return | |
| try: | |
| shutil.rmtree(self._prometheus_dir) | |
| logger.info("Cleaned up PROMETHEUS_MULTIPROC_DIR=%s", self._prometheus_dir) | |
| except OSError as e: | |
| logger.warning( | |
| "Failed to clean up PROMETHEUS_MULTIPROC_DIR=%s: %s", | |
| self._prometheus_dir, | |
| e, | |
| ) | |
| os.environ.pop("PROMETHEUS_MULTIPROC_DIR", None) | |
| self._prometheus_dir = None |
There was a problem hiding this comment.
Looks good! Clean implementation for gRPC metrics collection via prometheus multiproc and graceful metrics server bind failure.
3 minor nits flagged — all 🟡:
- Dead
Ok(_)arm in the metrics match (unreachable due to earlier empty check) - No timeout on the python3 subprocess for metrics collection
PROMETHEUS_MULTIPROC_DIRenv var not cleared on cleanup
None are blocking.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 7f618a4642
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| if text.trim().is_empty() { | ||
| return Err("no metrics available from gRPC workers yet".to_string()); |
There was a problem hiding this comment.
Treat empty gRPC multiprocess output as a non-error
When the multiprocess scrape has no samples yet (a normal startup state), this branch returns an Err, but get_engine_metrics logs all Err results as warnings. That makes the "skip silently" Ok(_) branch unreachable and produces warning spam on every scrape until metrics appear. Returning an empty Ok (or a distinct non-warning outcome) would preserve the intended quiet behavior.
Useful? React with 👍 / 👎.
| #[expect( | ||
| clippy::disallowed_methods, | ||
| reason = "no-op task for graceful degradation" | ||
| )] | ||
| return tokio::spawn(async {}); |
There was a problem hiding this comment.
Run Prometheus upkeep even on metrics bind failure
This early return bypasses the upkeep task started later in start_metrics_server, even though the recorder is already installed with an upkeep timeout. In the degraded path (port conflict), the router keeps serving traffic but never calls run_upkeep(), so histogram/metric maintenance is skipped for the lifetime of the process.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Code Review
This pull request enhances metrics collection by integrating Prometheus multi-process metrics for gRPC workers. The Python serve module now manages a temporary directory for Prometheus metrics, ensuring proper setup and cleanup. The Rust WorkerManager has been updated to collect these gRPC metrics by executing a Python subprocess. Additionally, the metrics HTTP server startup in Rust is made more resilient by gracefully handling port binding failures. A critical issue was identified in the inline Python script used for gRPC metrics collection, where an indentation error would prevent the script from executing correctly, thus hindering metrics collection.
| "import sys\n\ | ||
| from prometheus_client import CollectorRegistry, generate_latest\n\ | ||
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | ||
| registry = CollectorRegistry()\n\ | ||
| MultiProcessCollector(registry)\n\ | ||
| sys.stdout.buffer.write(generate_latest(registry))\n", |
There was a problem hiding this comment.
The inline Python script has an indentation error. Each line of the script after the first one starts with a space, which will cause a Python IndentationError: unexpected indent when the script is executed. This will prevent metrics from being collected from gRPC workers.
To fix this, please remove the leading spaces from each line of the Python script within the string literal.
"import sys\nfrom prometheus_client import CollectorRegistry, generate_latest\nfrom prometheus_client.multiprocess import MultiProcessCollector\nregistry = CollectorRegistry()\nMultiProcessCollector(registry)\nsys.stdout.buffer.write(generate_latest(registry))\n",|
Hi @ConnorLi96, this PR has merge conflicts that must be resolved before it can be merged. Please rebase your branch: git fetch origin main
git rebase origin/main
# resolve any conflicts, then:
git push --force-with-lease |
Signed-off-by: Scott Lee <scott@together.ai>
Signed-off-by: Scott Lee <scott@together.ai>
When the gRPC worker hasn't written any .db files yet (startup phase or no requests processed), generate_latest() returns empty bytes. Passing this empty string to parse_prometheus() triggers a parse error WARN log on every scrape interval. Add guards at both the collector function (return Err for empty output) and the call site (skip empty text silently) to prevent log spam. Signed-off-by: ConnorLi96 <ConnorLi96@users.noreply.github.com> Made-with: Cursor
When the metrics port is already in use, log an error and return a no-op handle instead of panicking, so the router can still operate. Signed-off-by: ConnorLi96 <ConnorLi96@users.noreply.github.com> Made-with: Cursor
The /metrics endpoint on the metrics port now returns both SMG's own metrics (smg_*) and engine metrics (vllm_*, sglang_*, nv_trt_*) in a single Prometheus text response. One scrape target per pod. Engine metrics deps (WorkerRegistry + reqwest::Client) are registered via OnceLock after AppContext init. Before registration, the handler returns SMG metrics only. No namespace collision — engine and SMG use distinct metric prefixes. The /engine_metrics endpoint on the main port stays as-is. Signed-off-by: ConnorLi96 <ConnorLi96@users.noreply.github.com> Made-with: Cursor
7f618a4 to
63eff15
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 63eff15f22
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| use crate::worker::{ | ||
| manager::{EngineMetricsResult, WorkerManager}, | ||
| registry::WorkerRegistry, |
There was a problem hiding this comment.
Wire /metrics to the new gRPC multiprocess collector
This import still routes /metrics through crate::worker::manager::WorkerManager, but the new gRPC multiprocess logic was added in model_gateway/src/core/worker_manager.rs and is never wired into the crate (there is no core module export). As a result, prometheus_handler continues to call the old get_engine_metrics implementation (HTTP fan-out to .../metrics), which fails for grpc:// workers and leaves engine metrics absent in gRPC mode.
Useful? React with 👍 / 👎.
| for resp in responses { | ||
| if let Ok(r) = resp.result { | ||
| if r.status().is_success() { | ||
| if let Ok(text) = r.text().await { | ||
| metric_packs.push(MetricPack { | ||
| labels: vec![("worker_addr".into(), resp.url)], | ||
| metrics_text: text, | ||
| }); | ||
| } | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
🟡 Nit: Three levels of silent error swallowing — network errors (resp.result), HTTP errors (r.status()), and body-read errors (r.text()) are all silently dropped. If a worker is consistently failing, the operator will see fewer metrics in the aggregated output but have no signal about which worker is down or why.
Consider adding at least a debug! log for failures so operators can diagnose missing metrics with RUST_LOG=debug:
| for resp in responses { | |
| if let Ok(r) = resp.result { | |
| if r.status().is_success() { | |
| if let Ok(text) = r.text().await { | |
| metric_packs.push(MetricPack { | |
| labels: vec![("worker_addr".into(), resp.url)], | |
| metrics_text: text, | |
| }); | |
| } | |
| } | |
| } | |
| } | |
| for resp in responses { | |
| match resp.result { | |
| Ok(r) if r.status().is_success() => { | |
| if let Ok(text) = r.text().await { | |
| metric_packs.push(MetricPack { | |
| labels: vec![("worker_addr".into(), resp.url)], | |
| metrics_text: text, | |
| }); | |
| } | |
| } | |
| Ok(r) => { | |
| debug!("Worker {} returned HTTP {} for /metrics", resp.url, r.status()); | |
| } | |
| Err(e) => { | |
| debug!("Failed to fetch /metrics from {}: {e}", resp.url); | |
| } | |
| } |
There was a problem hiding this comment.
Actionable comments posted: 3
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
bindings/python/src/smg/serve.py (1)
660-699:⚠️ Potential issue | 🟡 MinorAlways clean up
PROMETHEUS_MULTIPROC_DIR, even when no worker was recorded.Line 662 returns before
_cleanup_prometheus_dir()runs. If Line 607 created the temp directory and the firstPopenfails beforeself.workers.append(...), thefinally/atexitpath leaks the directory.Suggested fix
def _cleanup_workers(self) -> None: """SIGTERM all worker process groups, wait, then SIGKILL stragglers.""" if not self.workers: + self._cleanup_prometheus_dir() return🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@bindings/python/src/smg/serve.py` around lines 660 - 699, The early return in _cleanup_workers prevents _cleanup_prometheus_dir from running when self.workers is empty, leaking a previously created PROMETHEUS_MULTIPROC_DIR; update _cleanup_workers so it always calls self._cleanup_prometheus_dir() (e.g., remove the early return or call self._cleanup_prometheus_dir() before returning) so the temporary directory is removed even if no worker was recorded, keeping the rest of the SIGTERM/SIGKILL logic unchanged.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 114-119: The current code turns a cold-start empty scrape into
Err("no metrics available from gRPC workers yet"), causing repeated warning
logs; change the empty-output branch to return Ok(String::new()) instead of
Err(...) (i.e., when text.trim().is_empty() => Ok(String::new())), or
alternatively introduce a distinct sentinel error type and handle it silently at
the caller; update the callers that currently warn on any Err (the code that
iterates over the scraper results and logs warnings) to treat an empty String
(or the sentinel) as a non-error/no-op and skip logging so cold-start scrapes do
not spam warnings.
- Around line 556-595: The watch snapshot isn't pruned when all fetches fail
because the code returns early on group_loads.is_empty(); change
group_monitor_loop so that you always compute all_group_urls (from
workers.iter().map(|w| w.url().to_string())) and call tx.send_modify to remove
those URLs before the empty-check/continue, then only perform the
successful-path work (inserting group_loads into the map, calling
policy.update_loads and worker_load_manager.update_dp_loads) when group_loads is
non-empty; ensure worker_load_manager.update_dp_loads(&group_dp_loads) is not
called on failed ticks so the DP-rank cache remains unchanged.
- Around line 94-107: Wrap the tokio::process::Command invocation that currently
calls .output().await in a tokio::time::timeout using the existing
REQUEST_TIMEOUT constant: call tokio::time::timeout(REQUEST_TIMEOUT,
the_command.output()).await, then map the Timeout error to a descriptive Err
(e.g., "timed out running python3 prometheus collector") and unwrap the inner
Result from the command to preserve its original I/O error mapping (the existing
map_err("failed to run python3 prometheus collector: {e}") logic). Update the
code around the Command::new("python3") / .output().await block to perform this
two-level error handling so a hung subprocess or FS stall is bounded by
REQUEST_TIMEOUT.
---
Outside diff comments:
In `@bindings/python/src/smg/serve.py`:
- Around line 660-699: The early return in _cleanup_workers prevents
_cleanup_prometheus_dir from running when self.workers is empty, leaking a
previously created PROMETHEUS_MULTIPROC_DIR; update _cleanup_workers so it
always calls self._cleanup_prometheus_dir() (e.g., remove the early return or
call self._cleanup_prometheus_dir() before returning) so the temporary directory
is removed even if no worker was recorded, keeping the rest of the
SIGTERM/SIGKILL logic unchanged.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: da945273-c257-4551-b2f1-f6f482b42c26
📒 Files selected for processing (4)
bindings/python/src/smg/serve.pymodel_gateway/src/core/worker_manager.rsmodel_gateway/src/observability/metrics_server.rsmodel_gateway/src/server.rs
| let output = tokio::process::Command::new("python3") | ||
| .args([ | ||
| "-c", | ||
| "import sys\n\ | ||
| from prometheus_client import CollectorRegistry, generate_latest\n\ | ||
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | ||
| registry = CollectorRegistry()\n\ | ||
| MultiProcessCollector(registry)\n\ | ||
| sys.stdout.buffer.write(generate_latest(registry))\n", | ||
| ]) | ||
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | ||
| .output() | ||
| .await | ||
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
fd -type f -name "worker_manager.rs" | head -20Repository: lightseekorg/smg
Length of output: 230
🏁 Script executed:
if [ -f "model_gateway/src/core/worker_manager.rs" ]; then
wc -l "model_gateway/src/core/worker_manager.rs"
fiRepository: lightseekorg/smg
Length of output: 104
🏁 Script executed:
head -120 "model_gateway/src/core/worker_manager.rs" | tail -40Repository: lightseekorg/smg
Length of output: 1675
🏁 Script executed:
# Find the function containing the prometheus collector code
rg -A 5 "CollectorRegistry" model_gateway/src/core/worker_manager.rsRepository: lightseekorg/smg
Length of output: 460
🏁 Script executed:
# Search for REQUEST_TIMEOUT or similar timeout constants
rg "REQUEST_TIMEOUT|timeout.*=" model_gateway/src/ | head -20Repository: lightseekorg/smg
Length of output: 2022
🏁 Script executed:
# Check what function calls the Python subprocess and if it's on metrics path
rg -B 15 "python3.*prometheus_client" model_gateway/src/core/worker_manager.rs | head -50Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Find where collect_prometheus_multiproc_metrics is called
rg "collect_prometheus_multiproc_metrics" model_gateway/src/core/worker_manager.rsRepository: lightseekorg/smg
Length of output: 200
🏁 Script executed:
# Get the context around the function definition to see line numbers
cat -n model_gateway/src/core/worker_manager.rs | grep -A 35 "collect_prometheus_multiproc_metrics"Repository: lightseekorg/smg
Length of output: 3311
🏁 Script executed:
# Check if there are timeout patterns used elsewhere in this file
rg "tokio::time::timeout|Duration::from" model_gateway/src/core/worker_manager.rsRepository: lightseekorg/smg
Length of output: 285
🏁 Script executed:
# Look for where this function is used and what calls it
rg -B 5 -A 5 "collect_prometheus_multiproc_metrics" model_gateway/src/Repository: lightseekorg/smg
Length of output: 1739
🏁 Script executed:
# Get the function signature for the function containing line 338
sed -n '280,340p' model_gateway/src/core/worker_manager.rs | head -60Repository: lightseekorg/smg
Length of output: 2100
🏁 Script executed:
# Get more context to understand the function name
sed -n '250,340p' model_gateway/src/core/worker_manager.rsRepository: lightseekorg/smg
Length of output: 3148
🏁 Script executed:
# Look for "get_engine_metrics" or similar functions
rg "fn .*engine.*metric|pub.*fn.*metric" model_gateway/src/core/worker_manager.rs -A 3Repository: lightseekorg/smg
Length of output: 203
Add a timeout to the Python collector subprocess.
The subprocess call at lines 94–107 uses unbounded .output().await, which can block metrics scrapes indefinitely if python3 or the filesystem hangs. The codebase already defines REQUEST_TIMEOUT (5 seconds) and applies it to similar async operations elsewhere in this file (e.g., HTTP requests); use the same pattern here.
Suggested fix
+ let output = tokio::time::timeout(
+ REQUEST_TIMEOUT,
+ tokio::process::Command::new("python3")
+ .args([
+ "-c",
+ "import sys\n\
+ from prometheus_client import CollectorRegistry, generate_latest\n\
+ from prometheus_client.multiprocess import MultiProcessCollector\n\
+ registry = CollectorRegistry()\n\
+ MultiProcessCollector(registry)\n\
+ sys.stdout.buffer.write(generate_latest(registry))\n",
+ ])
+ .env("PROMETHEUS_MULTIPROC_DIR", &dir)
+ .output(),
+ )
+ .await
+ .map_err(|_| "python3 prometheus collector timed out".to_string())?
+ .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?;
- let output = tokio::process::Command::new("python3")
- .args([
- "-c",
- "import sys\n\
- from prometheus_client import CollectorRegistry, generate_latest\n\
- from prometheus_client.multiprocess import MultiProcessCollector\n\
- registry = CollectorRegistry()\n\
- MultiProcessCollector(registry)\n\
- sys.stdout.buffer.write(generate_latest(registry))\n",
- ])
- .env("PROMETHEUS_MULTIPROC_DIR", &dir)
- .output()
- .await
- .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?;📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| let output = tokio::process::Command::new("python3") | |
| .args([ | |
| "-c", | |
| "import sys\n\ | |
| from prometheus_client import CollectorRegistry, generate_latest\n\ | |
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | |
| registry = CollectorRegistry()\n\ | |
| MultiProcessCollector(registry)\n\ | |
| sys.stdout.buffer.write(generate_latest(registry))\n", | |
| ]) | |
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | |
| .output() | |
| .await | |
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; | |
| let output = tokio::time::timeout( | |
| REQUEST_TIMEOUT, | |
| tokio::process::Command::new("python3") | |
| .args([ | |
| "-c", | |
| "import sys\n\ | |
| from prometheus_client import CollectorRegistry, generate_latest\n\ | |
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | |
| registry = CollectorRegistry()\n\ | |
| MultiProcessCollector(registry)\n\ | |
| sys.stdout.buffer.write(generate_latest(registry))\n", | |
| ]) | |
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | |
| .output(), | |
| ) | |
| .await | |
| .map_err(|_| "python3 prometheus collector timed out".to_string())? | |
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/core/worker_manager.rs` around lines 94 - 107, Wrap the
tokio::process::Command invocation that currently calls .output().await in a
tokio::time::timeout using the existing REQUEST_TIMEOUT constant: call
tokio::time::timeout(REQUEST_TIMEOUT, the_command.output()).await, then map the
Timeout error to a descriptive Err (e.g., "timed out running python3 prometheus
collector") and unwrap the inner Result from the command to preserve its
original I/O error mapping (the existing map_err("failed to run python3
prometheus collector: {e}") logic). Update the code around the
Command::new("python3") / .output().await block to perform this two-level error
handling so a hung subprocess or FS stall is bounded by REQUEST_TIMEOUT.
| let text = String::from_utf8(output.stdout) | ||
| .map_err(|e| format!("prometheus collector output is not valid UTF-8: {e}"))?; | ||
| if text.trim().is_empty() { | ||
| return Err("no metrics available from gRPC workers yet".to_string()); | ||
| } | ||
| Ok(text) |
There was a problem hiding this comment.
The expected empty-output case still logs once per scrape.
Lines 116-118 turn the cold-start “no metrics yet” state into Err(...), and Lines 348-351 warn on every Err. That replaces the old parse spam with warning spam. Either return Ok(String::new()) for the empty case or distinguish that sentinel error and skip it silently at the call site.
Also applies to: 337-352
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/core/worker_manager.rs` around lines 114 - 119, The current
code turns a cold-start empty scrape into Err("no metrics available from gRPC
workers yet"), causing repeated warning logs; change the empty-output branch to
return Ok(String::new()) instead of Err(...) (i.e., when text.trim().is_empty()
=> Ok(String::new())), or alternatively introduce a distinct sentinel error type
and handle it silently at the caller; update the callers that currently warn on
any Err (the code that iterates over the scraper results and logs warnings) to
treat an empty String (or the sentinel) as a non-error/no-op and skip logging so
cold-start scrapes do not spam warnings.
| let results = future::join_all(futures).await; | ||
|
|
||
| // Collect successful loads | ||
| let mut group_loads: HashMap<String, WorkerLoadResponse> = HashMap::new(); | ||
| let mut group_dp_loads: HashMap<String, HashMap<isize, isize>> = HashMap::new(); | ||
| for (url, response) in results { | ||
| if let Some(load) = response { | ||
| group_loads.insert(url.clone(), load.clone()); | ||
| let dp_rank_loads = load.dp_rank_loads(); | ||
| group_dp_loads.insert(url, dp_rank_loads); | ||
| } | ||
| } | ||
|
|
||
| if group_loads.is_empty() { | ||
| debug!("No loads fetched for group {group_key}"); | ||
| continue; | ||
| } | ||
|
|
||
| debug!( | ||
| "Fetched loads from {}/{} workers in group {group_key}", | ||
| group_loads.len(), | ||
| workers.len() | ||
| ); | ||
|
|
||
| // Update policies with this group's loads | ||
| for policy in &power_of_two_policies { | ||
| policy.update_loads(&group_loads); | ||
| } | ||
| worker_load_manager.update_dp_loads(&group_dp_loads); | ||
|
|
||
| // Atomically merge into the shared watch channel. | ||
| // Remove all group URLs first to clear stale entries from workers | ||
| // that failed this tick, then insert successful loads. | ||
| let all_group_urls: Vec<String> = workers.iter().map(|w| w.url().to_string()).collect(); | ||
| tx.send_modify(|map| { | ||
| for url in &all_group_urls { | ||
| map.remove(url); | ||
| } | ||
| map.extend(group_loads); | ||
| }); |
There was a problem hiding this comment.
Prune the watch snapshot even when every fetch in a tick fails.
Line 569 exits before the group's URLs are removed from the shared watch map, so subscribers keep seeing stale per-worker loads after a full miss. The DP-rank cache can stay last-known-good, but the watch snapshot still needs to be cleared on failed ticks.
Suggested fix
- let results = future::join_all(futures).await;
+ let results = future::join_all(futures).await;
+ let all_group_urls: Vec<String> = workers.iter().map(|w| w.url().to_string()).collect();
// Collect successful loads
let mut group_loads: HashMap<String, WorkerLoadResponse> = HashMap::new();
let mut group_dp_loads: HashMap<String, HashMap<isize, isize>> = HashMap::new();
for (url, response) in results {
@@
}
if group_loads.is_empty() {
+ tx.send_modify(|map| {
+ for url in &all_group_urls {
+ map.remove(url);
+ }
+ });
debug!("No loads fetched for group {group_key}");
continue;
}
@@
- let all_group_urls: Vec<String> = workers.iter().map(|w| w.url().to_string()).collect();
tx.send_modify(|map| {
for url in &all_group_urls {
map.remove(url);
}
map.extend(group_loads);Based on learnings, group_monitor_loop should prune the watch snapshot for the group’s URLs on both successful and failed ticks, while leaving the DP-rank cache untouched on failed fetches.
📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| let results = future::join_all(futures).await; | |
| // Collect successful loads | |
| let mut group_loads: HashMap<String, WorkerLoadResponse> = HashMap::new(); | |
| let mut group_dp_loads: HashMap<String, HashMap<isize, isize>> = HashMap::new(); | |
| for (url, response) in results { | |
| if let Some(load) = response { | |
| group_loads.insert(url.clone(), load.clone()); | |
| let dp_rank_loads = load.dp_rank_loads(); | |
| group_dp_loads.insert(url, dp_rank_loads); | |
| } | |
| } | |
| if group_loads.is_empty() { | |
| debug!("No loads fetched for group {group_key}"); | |
| continue; | |
| } | |
| debug!( | |
| "Fetched loads from {}/{} workers in group {group_key}", | |
| group_loads.len(), | |
| workers.len() | |
| ); | |
| // Update policies with this group's loads | |
| for policy in &power_of_two_policies { | |
| policy.update_loads(&group_loads); | |
| } | |
| worker_load_manager.update_dp_loads(&group_dp_loads); | |
| // Atomically merge into the shared watch channel. | |
| // Remove all group URLs first to clear stale entries from workers | |
| // that failed this tick, then insert successful loads. | |
| let all_group_urls: Vec<String> = workers.iter().map(|w| w.url().to_string()).collect(); | |
| tx.send_modify(|map| { | |
| for url in &all_group_urls { | |
| map.remove(url); | |
| } | |
| map.extend(group_loads); | |
| }); | |
| let results = future::join_all(futures).await; | |
| let all_group_urls: Vec<String> = workers.iter().map(|w| w.url().to_string()).collect(); | |
| // Collect successful loads | |
| let mut group_loads: HashMap<String, WorkerLoadResponse> = HashMap::new(); | |
| let mut group_dp_loads: HashMap<String, HashMap<isize, isize>> = HashMap::new(); | |
| for (url, response) in results { | |
| if let Some(load) = response { | |
| group_loads.insert(url.clone(), load.clone()); | |
| let dp_rank_loads = load.dp_rank_loads(); | |
| group_dp_loads.insert(url, dp_rank_loads); | |
| } | |
| } | |
| if group_loads.is_empty() { | |
| tx.send_modify(|map| { | |
| for url in &all_group_urls { | |
| map.remove(url); | |
| } | |
| }); | |
| debug!("No loads fetched for group {group_key}"); | |
| continue; | |
| } | |
| debug!( | |
| "Fetched loads from {}/{} workers in group {group_key}", | |
| group_loads.len(), | |
| workers.len() | |
| ); | |
| // Update policies with this group's loads | |
| for policy in &power_of_two_policies { | |
| policy.update_loads(&group_loads); | |
| } | |
| worker_load_manager.update_dp_loads(&group_dp_loads); | |
| // Atomically merge into the shared watch channel. | |
| // Remove all group URLs first to clear stale entries from workers | |
| // that failed this tick, then insert successful loads. | |
| tx.send_modify(|map| { | |
| for url in &all_group_urls { | |
| map.remove(url); | |
| } | |
| map.extend(group_loads); | |
| }); |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/core/worker_manager.rs` around lines 556 - 595, The watch
snapshot isn't pruned when all fetches fail because the code returns early on
group_loads.is_empty(); change group_monitor_loop so that you always compute
all_group_urls (from workers.iter().map(|w| w.url().to_string())) and call
tx.send_modify to remove those URLs before the empty-check/continue, then only
perform the successful-path work (inserting group_loads into the map, calling
policy.update_loads and worker_load_manager.update_dp_loads) when group_loads is
non-empty; ensure worker_load_manager.update_dp_loads(&group_dp_loads) is not
called on failed ticks so the DP-rank cache remains unchanged.
| error = %e, | ||
| "PD backend request failed" | ||
| ); | ||
| error!("PD request transport error, both sides aborted: {e}"); |
There was a problem hiding this comment.
🟡 Nit: This push removed the comment that was here explaining why record_outcome is not called in this error path:
// Don't record_outcome here — the caller (execute_dual_dispatch)
// records outcomes from the response status after we return.
That comment documented a non-obvious design decision. Without it, a future reader might add a record_outcome call here (mirroring other error paths), which would cause double-counting with the caller. Consider restoring it above the return.
| pub(super) fn get_server_info(registry: &WorkerRegistry) -> Response { | ||
| let stats = registry.stats(); | ||
| let workers = external_workers(registry); | ||
| let workers = registry.get_all(); |
There was a problem hiding this comment.
🟡 Nit: This change removes the RuntimeType::External filter and drops the "external_workers" field from the /info JSON response. If any monitoring dashboards or scripts consume that field, they'll silently get null/missing. The behavioral change (health checks now cover all workers, not just external ones) seems intentional, but the field removal is a breaking API change worth noting in the PR description.
018f984 to
f4d1ecb
Compare
Add structured INFO logs for routing decisions and backend response times to both HTTP and gRPC router paths. Also fix health_generate to check all workers (not just External), so gRPC-connected backends (TRT-LLM, SGLang, vLLM) are visible. Signed-off-by: ConnorLi96 <ConnorLi96@users.noreply.github.com> Made-with: Cursor Signed-off-by: ConnorLi96 <ConnorLi96@users.noreply.github.com> Made-with: Cursor Signed-off-by: ConnorLi96 <ConnorLi96@users.noreply.github.com> Made-with: Cursor
| async fn health_generate(&self, _req: Request<Body>) -> Response { | ||
| let workers = self.worker_registry.get_all(); | ||
| if workers.is_empty() { | ||
| return (StatusCode::SERVICE_UNAVAILABLE, "No workers registered").into_response(); | ||
| } | ||
| let (healthy, unhealthy): (Vec<_>, Vec<_>) = workers.iter().partition(|w| w.is_healthy()); | ||
| if unhealthy.is_empty() { | ||
| ( | ||
| StatusCode::OK, | ||
| format!("OK - {} workers healthy", healthy.len()), | ||
| ) | ||
| .into_response() | ||
| } else { | ||
| let info: Vec<_> = unhealthy | ||
| .iter() | ||
| .map(|w| format!("{} ({})", w.model_id(), w.url())) | ||
| .collect(); | ||
| ( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| format!( | ||
| "{}/{} workers unhealthy: {}", | ||
| unhealthy.len(), | ||
| workers.len(), | ||
| info.join(", ") | ||
| ), | ||
| ) | ||
| .into_response() | ||
| } | ||
| } |
There was a problem hiding this comment.
🔴 Important: This PD router health check doesn't verify that both worker types (prefill and decode) are present and healthy. get_all() returns all workers regardless of type, so if all prefill workers are down (or none are registered) but decode workers are healthy, this will report "OK - N workers healthy" even though the router can't actually serve any requests (it needs at least one of each).
Compare with the HTTP PD router (http/pd_router.rs:1213) which selects a PD pair and tests both sides, and with the Debug impl just above this method (lines 364-377) which already filters by WorkerType::Prefill / WorkerType::Decode.
Consider checking both types:
| async fn health_generate(&self, _req: Request<Body>) -> Response { | |
| let workers = self.worker_registry.get_all(); | |
| if workers.is_empty() { | |
| return (StatusCode::SERVICE_UNAVAILABLE, "No workers registered").into_response(); | |
| } | |
| let (healthy, unhealthy): (Vec<_>, Vec<_>) = workers.iter().partition(|w| w.is_healthy()); | |
| if unhealthy.is_empty() { | |
| ( | |
| StatusCode::OK, | |
| format!("OK - {} workers healthy", healthy.len()), | |
| ) | |
| .into_response() | |
| } else { | |
| let info: Vec<_> = unhealthy | |
| .iter() | |
| .map(|w| format!("{} ({})", w.model_id(), w.url())) | |
| .collect(); | |
| ( | |
| StatusCode::SERVICE_UNAVAILABLE, | |
| format!( | |
| "{}/{} workers unhealthy: {}", | |
| unhealthy.len(), | |
| workers.len(), | |
| info.join(", ") | |
| ), | |
| ) | |
| .into_response() | |
| } | |
| } | |
| async fn health_generate(&self, _req: Request<Body>) -> Response { | |
| let prefill_workers = self.worker_registry.get_workers_filtered( | |
| None, | |
| Some(WorkerType::Prefill), | |
| Some(ConnectionMode::Grpc), | |
| None, | |
| false, | |
| ); | |
| let decode_workers = self.worker_registry.get_workers_filtered( | |
| None, | |
| Some(WorkerType::Decode), | |
| Some(ConnectionMode::Grpc), | |
| None, | |
| false, | |
| ); | |
| if prefill_workers.is_empty() || decode_workers.is_empty() { | |
| return ( | |
| StatusCode::SERVICE_UNAVAILABLE, | |
| format!( | |
| "Missing worker type: {} prefill, {} decode", | |
| prefill_workers.len(), | |
| decode_workers.len() | |
| ), | |
| ) | |
| .into_response(); | |
| } | |
| let workers = self.worker_registry.get_all(); | |
| let (healthy, unhealthy): (Vec<_>, Vec<_>) = workers.iter().partition(|w| w.is_healthy()); | |
| if unhealthy.is_empty() { | |
| ( | |
| StatusCode::OK, | |
| format!("OK - {} workers healthy", healthy.len()), | |
| ) | |
| .into_response() | |
| } else { | |
| let info: Vec<_> = unhealthy | |
| .iter() | |
| .map(|w| format!("{} ({})", w.model_id(), w.url())) | |
| .collect(); | |
| ( | |
| StatusCode::SERVICE_UNAVAILABLE, | |
| format!( | |
| "{}/{} workers unhealthy: {}", | |
| unhealthy.len(), | |
| workers.len(), | |
| info.join(", ") | |
| ), | |
| ) | |
| .into_response() | |
| } | |
| } |
| async fn health_generate(&self, _req: Request<Body>) -> Response { | ||
| let workers = self.worker_registry.get_all(); | ||
| if workers.is_empty() { | ||
| return (StatusCode::SERVICE_UNAVAILABLE, "No workers registered").into_response(); | ||
| } | ||
| let (healthy, unhealthy): (Vec<_>, Vec<_>) = workers.iter().partition(|w| w.is_healthy()); | ||
| if unhealthy.is_empty() { | ||
| ( | ||
| StatusCode::OK, | ||
| format!("OK - {} workers healthy", healthy.len()), | ||
| ) | ||
| .into_response() | ||
| } else { | ||
| let info: Vec<_> = unhealthy | ||
| .iter() | ||
| .map(|w| format!("{} ({})", w.model_id(), w.url())) | ||
| .collect(); | ||
| ( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| format!( | ||
| "{}/{} workers unhealthy: {}", | ||
| unhealthy.len(), | ||
| workers.len(), | ||
| info.join(", ") | ||
| ), | ||
| ) | ||
| .into_response() | ||
| } | ||
| } |
There was a problem hiding this comment.
🟡 Nit: This is byte-for-byte identical to GrpcPDRouter::health_generate. Consider extracting a shared helper (similar to openai/health.rs:health_generate) that both gRPC router impls can call, to avoid the two copies drifting apart over time.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f4d1ecb185
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| let smg_text = state.handle.render(); | ||
|
|
||
| let engine_text = if let Some(deps) = ENGINE_METRICS_DEPS.get() { | ||
| match WorkerManager::get_engine_metrics(&deps.worker_registry, &deps.client).await { |
There was a problem hiding this comment.
Restrict engine metrics scrape to local workers
This handler now calls WorkerManager::get_engine_metrics on every /metrics scrape, and that helper fans out GET /metrics to all registered workers (model_gateway/src/worker/manager.rs, fan_out(&workers, ...)). In OpenAI/IGW deployments, those workers can be external provider URLs, so each Prometheus scrape generates outbound requests (with bearer auth when configured) to upstream APIs that do not expose engine metrics, creating avoidable external traffic and possible rate-limit/credential-exposure risk. Filter to worker types/runtime modes that actually export local engine metrics before invoking the fan-out.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/routers/grpc/pd_router.rs`:
- Around line 392-417: The health check is using worker_registry.get_all() so
non-PD workers can affect PD health; modify the logic in this block to first
filter the workers list to only PD gRPC workers (e.g., those with the
PD/Prefill/Decode role or a protocol field indicating gRPC) before partitioning
by is_healthy(); keep the rest of the flow (partition, formatting using
is_healthy(), model_id(), url()) the same so only PD gRPC workers are counted
for OK vs SERVICE_UNAVAILABLE responses.
In `@model_gateway/src/routers/grpc/router.rs`:
- Around line 515-540: The health-check currently calls
self.worker_registry.get_all() and evaluates all workers; change it to only
consider gRPC-routable workers by filtering the collection before partitioning
so non-gRPC workers don't affect the gRPC router status. Update the code that
builds workers (from self.worker_registry.get_all()) to filter using the worker
predicate that indicates gRPC routability (e.g., a method like
is_grpc_routable() or supports_grpc() on the worker), then continue to use the
same partitioning (is_healthy()), model_id(), and url() logic on that filtered
list to produce the OK/503 response.
In `@model_gateway/src/routers/http/pd_router.rs`:
- Around line 647-669: The code currently logs "PD backend request completed" at
info! for any Ok(pd_result) even when prefill_resp or decode_resp is a 5xx;
change the Ok branch to check prefill_resp.status().is_server_error() ||
decode_resp.status().is_server_error() and, if true, emit warn! (include the
same fields: prefill_url = prefill.url(), decode_url = decode.url(),
prefill_status = prefill_resp.status().as_u16(), decode_status =
decode_resp.status().as_u16(), duration_ms = backend_duration_ms, streaming =
context.is_stream) with a message like "PD backend returned 5xx", otherwise keep
the info! "PD backend request completed"; locate this change in the match on
pd_result that assigns (prefill_resp, decode_resp).
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: e057761d-8eaa-4979-8526-2cde7de3fd15
📒 Files selected for processing (7)
model_gateway/src/routers/grpc/common/stages/request_execution.rsmodel_gateway/src/routers/grpc/common/stages/worker_selection.rsmodel_gateway/src/routers/grpc/pd_router.rsmodel_gateway/src/routers/grpc/router.rsmodel_gateway/src/routers/http/pd_router.rsmodel_gateway/src/routers/http/router.rsmodel_gateway/src/routers/openai/health.rs
| let workers = self.worker_registry.get_all(); | ||
| if workers.is_empty() { | ||
| return (StatusCode::SERVICE_UNAVAILABLE, "No workers registered").into_response(); | ||
| } | ||
| let (healthy, unhealthy): (Vec<_>, Vec<_>) = workers.iter().partition(|w| w.is_healthy()); | ||
| if unhealthy.is_empty() { | ||
| ( | ||
| StatusCode::OK, | ||
| format!("OK - {} workers healthy", healthy.len()), | ||
| ) | ||
| .into_response() | ||
| } else { | ||
| let info: Vec<_> = unhealthy | ||
| .iter() | ||
| .map(|w| format!("{} ({})", w.model_id(), w.url())) | ||
| .collect(); | ||
| ( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| format!( | ||
| "{}/{} workers unhealthy: {}", | ||
| unhealthy.len(), | ||
| workers.len(), | ||
| info.join(", ") | ||
| ), | ||
| ) | ||
| .into_response() |
There was a problem hiding this comment.
Scope PD health to PD gRPC workers only.
Line 392 uses get_all(), so unrelated workers can incorrectly flip PD health to 503. The PD health endpoint should evaluate only Prefill/Decode gRPC workers.
🔧 Proposed fix
- let workers = self.worker_registry.get_all();
+ let prefill_workers = self.worker_registry.get_workers_filtered(
+ None,
+ Some(WorkerType::Prefill),
+ Some(ConnectionMode::Grpc),
+ None,
+ false,
+ );
+ let decode_workers = self.worker_registry.get_workers_filtered(
+ None,
+ Some(WorkerType::Decode),
+ Some(ConnectionMode::Grpc),
+ None,
+ false,
+ );
+ let workers: Vec<_> = prefill_workers
+ .into_iter()
+ .chain(decode_workers.into_iter())
+ .collect();🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/grpc/pd_router.rs` around lines 392 - 417, The
health check is using worker_registry.get_all() so non-PD workers can affect PD
health; modify the logic in this block to first filter the workers list to only
PD gRPC workers (e.g., those with the PD/Prefill/Decode role or a protocol field
indicating gRPC) before partitioning by is_healthy(); keep the rest of the flow
(partition, formatting using is_healthy(), model_id(), url()) the same so only
PD gRPC workers are counted for OK vs SERVICE_UNAVAILABLE responses.
| let workers = self.worker_registry.get_all(); | ||
| if workers.is_empty() { | ||
| return (StatusCode::SERVICE_UNAVAILABLE, "No workers registered").into_response(); | ||
| } | ||
| let (healthy, unhealthy): (Vec<_>, Vec<_>) = workers.iter().partition(|w| w.is_healthy()); | ||
| if unhealthy.is_empty() { | ||
| ( | ||
| StatusCode::OK, | ||
| format!("OK - {} workers healthy", healthy.len()), | ||
| ) | ||
| .into_response() | ||
| } else { | ||
| let info: Vec<_> = unhealthy | ||
| .iter() | ||
| .map(|w| format!("{} ({})", w.model_id(), w.url())) | ||
| .collect(); | ||
| ( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| format!( | ||
| "{}/{} workers unhealthy: {}", | ||
| unhealthy.len(), | ||
| workers.len(), | ||
| info.join(", ") | ||
| ), | ||
| ) | ||
| .into_response() |
There was a problem hiding this comment.
Restrict gRPC health checks to gRPC-routable workers.
Line 515 currently inspects all registered workers. Unhealthy non-gRPC workers can incorrectly make this gRPC router report 503.
🔧 Proposed fix
- let workers = self.worker_registry.get_all();
+ let workers = self.worker_registry.get_workers_filtered(
+ None,
+ None,
+ Some(crate::worker::ConnectionMode::Grpc),
+ None,
+ false,
+ );🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/grpc/router.rs` around lines 515 - 540, The
health-check currently calls self.worker_registry.get_all() and evaluates all
workers; change it to only consider gRPC-routable workers by filtering the
collection before partitioning so non-gRPC workers don't affect the gRPC router
status. Update the code that builds workers (from
self.worker_registry.get_all()) to filter using the worker predicate that
indicates gRPC routability (e.g., a method like is_grpc_routable() or
supports_grpc() on the worker), then continue to use the same partitioning
(is_healthy()), model_id(), and url() logic on that filtered list to produce the
OK/503 response.
| let (prefill_response, decode_response) = match pd_result { | ||
| Ok((prefill_resp, decode_resp)) => (prefill_resp, decode_resp), | ||
| Ok((prefill_resp, decode_resp)) => { | ||
| info!( | ||
| target: "smg::upstream", | ||
| prefill_url = prefill.url(), | ||
| decode_url = decode.url(), | ||
| prefill_status = prefill_resp.status().as_u16(), | ||
| decode_status = decode_resp.status().as_u16(), | ||
| duration_ms = backend_duration_ms, | ||
| streaming = context.is_stream, | ||
| "PD backend request completed" | ||
| ); | ||
| (prefill_resp, decode_resp) | ||
| } | ||
| Err(e) => { | ||
| warn!( | ||
| target: "smg::upstream", | ||
| prefill_url = prefill.url(), | ||
| decode_url = decode.url(), | ||
| duration_ms = backend_duration_ms, | ||
| error = %e, | ||
| "PD backend request failed" | ||
| ); |
There was a problem hiding this comment.
Don't log PD 5xx responses as completed.
tokio::try_join! only tells us both HTTP sends returned a Response. If either prefill_resp.status() or decode_resp.status() is 5xx, this block still emits "PD backend request completed" at info!, even though the request is about to fail later in this method. That makes PD failure logs misleading and inconsistent with the regular HTTP router's 5xx warn! path in model_gateway/src/routers/http/router.rs Lines 350-373.
🔧 Suggested fix
- let (prefill_response, decode_response) = match pd_result {
- Ok((prefill_resp, decode_resp)) => {
- info!(
- target: "smg::upstream",
- prefill_url = prefill.url(),
- decode_url = decode.url(),
- prefill_status = prefill_resp.status().as_u16(),
- decode_status = decode_resp.status().as_u16(),
- duration_ms = backend_duration_ms,
- streaming = context.is_stream,
- "PD backend request completed"
- );
- (prefill_resp, decode_resp)
- }
+ let (prefill_response, decode_response) = match pd_result {
+ Ok((prefill_resp, decode_resp)) => {
+ let prefill_status = prefill_resp.status();
+ let decode_status = decode_resp.status();
+ let server_error =
+ prefill_status.is_server_error() || decode_status.is_server_error();
+
+ if server_error {
+ warn!(
+ target: "smg::upstream",
+ prefill_url = prefill.url(),
+ decode_url = decode.url(),
+ prefill_status = prefill_status.as_u16(),
+ decode_status = decode_status.as_u16(),
+ duration_ms = backend_duration_ms,
+ streaming = context.is_stream,
+ "PD backend request failed"
+ );
+ } else {
+ info!(
+ target: "smg::upstream",
+ prefill_url = prefill.url(),
+ decode_url = decode.url(),
+ prefill_status = prefill_status.as_u16(),
+ decode_status = decode_status.as_u16(),
+ duration_ms = backend_duration_ms,
+ streaming = context.is_stream,
+ "PD backend request completed"
+ );
+ }
+
+ (prefill_resp, decode_resp)
+ }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/http/pd_router.rs` around lines 647 - 669, The code
currently logs "PD backend request completed" at info! for any Ok(pd_result)
even when prefill_resp or decode_resp is a 5xx; change the Ok branch to check
prefill_resp.status().is_server_error() ||
decode_resp.status().is_server_error() and, if true, emit warn! (include the
same fields: prefill_url = prefill.url(), decode_url = decode.url(),
prefill_status = prefill_resp.status().as_u16(), decode_status =
decode_resp.status().as_u16(), duration_ms = backend_duration_ms, streaming =
context.is_stream) with a message like "PD backend returned 5xx", otherwise keep
the info! "PD backend request completed"; locate this change in the match on
pd_result that assigns (prefill_resp, decode_resp).
…kers
Two fixes to make /metrics on :9900 include sglang engine metrics when the
router is run with --connection-mode grpc:
1. model_gateway: in WorkerManager::get_engine_metrics, split workers by
connection mode. HTTP workers keep the existing fan_out path. gRPC
workers are scraped at http://{host}:{grpc_port+1}/metrics, which is the
sglang convention when --enable-metrics is set. Previously an HTTP GET
was fired at the gRPC port and silently failed, so the entire engine
metrics block ended up empty. Scrape errors are now logged instead of
swallowed. Adds grpc_worker_metrics_url helper + unit tests covering
grpc://, grpcs://, and @dp_rank suffix variants.
2. grpc_servicer.sglang.request_manager: gRPC mode launches only the
scheduler process (no TokenizerManager), so TokenizerMetricsCollector
was never initialized and request-level metrics (TTFT, e2e latency,
inter-token latency, prompt/generation token totals) never appeared.
GrpcRequestManager now creates its own TokenizerMetricsCollector when
server_args.enable_metrics is set and observes TTFT / inter-token
latency / e2e latency in _handle_batch_output, mirroring the logic in
TokenizerManager.
Verified end-to-end against a sglang worker in gRPC mode: sglang_* now
appears on :9900/metrics alongside smg_*, and counters increment on
inference traffic (5 requests: num_requests_total 5 -> 16; TTFT and
e2e_latency histograms populated).
Made-with: Cursor
|
Hi @ConnorLi96, the DCO sign-off check has failed. All commits must include a To fix existing commits: # Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-leaseTo sign off future commits automatically:
|
| let (host, port_and_suffix) = stripped.rsplit_once(':')?; | ||
| let port_str = port_and_suffix.split(['/', '@']).next()?; | ||
| let port: u16 = port_str.parse().ok()?; | ||
| Some(format!("http://{host}:{}/metrics", port + 1)) |
There was a problem hiding this comment.
🔴 Important: port + 1 will panic in debug builds (or silently wrap to 0 in release) when the gRPC worker is on port 65535. While 65535 is uncommon for gRPC, u16 arithmetic overflow is undefined behavior in safe Rust — it's a panic in debug and a silent wrap in release. Use checked_add so the function returns None (consistent with the other parse-failure paths) instead of panicking.
| Some(format!("http://{host}:{}/metrics", port + 1)) | |
| let port: u16 = port_str.parse().ok()?; | |
| let metrics_port = port.checked_add(1)?; | |
| Some(format!("http://{host}:{metrics_port}/metrics")) |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 0f93a4cfa6
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| for worker in &grpc_workers { | ||
| if let Some(metrics_url) = grpc_worker_metrics_url(worker.url()) { |
There was a problem hiding this comment.
Scrape gRPC metrics in parallel
This loop processes gRPC workers serially, and each iteration awaits an HTTP scrape with REQUEST_TIMEOUT (5s). If one or more worker metrics endpoints are slow/unreachable (a common startup or partial-failure case), total /metrics latency scales with worker count and can exceed Prometheus scrape timeouts, causing the whole scrape to fail. The gRPC branch should fan out requests concurrently like the HTTP branch.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
model_gateway/src/worker/manager.rs (1)
842-878:⚠️ Potential issue | 🟠 MajorScrape gRPC workers in parallel — sequential
awaitwill block/metrics.The HTTP branch uses
fan_outwithbuffer_unordered(MAX_CONCURRENT), but this loop awaits each gRPC worker serially with a 5sREQUEST_TIMEOUT. With N gRPC workers where some are slow/unreachable, total latency approachesN * 5s, which will trivially exceed a typical Prometheus scrape timeout (10s default) once you have more than ~2 unreachable workers, and will also stallprometheus_handlerinobservability/metrics_server.rs. Fan these out the same way as the HTTP path.♻️ Suggested fix: run gRPC scrapes concurrently
- // gRPC workers expose Prometheus metrics on a separate HTTP port - // (gRPC port + 1) when --enable-metrics is passed to sglang. - for worker in &grpc_workers { - if let Some(metrics_url) = grpc_worker_metrics_url(worker.url()) { - match client - .get(&metrics_url) - .timeout(REQUEST_TIMEOUT) - .send() - .await - { - Ok(r) if r.status().is_success() => { - if let Ok(text) = r.text().await { - metric_packs.push(MetricPack { - labels: vec![( - "worker_addr".into(), - worker.url().to_string(), - )], - metrics_text: text, - }); - } - } - Ok(r) => { - warn!( - "gRPC worker metrics endpoint {} returned {}", - metrics_url, - r.status() - ); - } - Err(e) => { - warn!( - "Failed to scrape gRPC worker metrics from {}: {e}", - metrics_url - ); - } - } - } - } + // gRPC workers expose Prometheus metrics on a separate HTTP port + // (gRPC port + 1) when --enable-metrics is passed to sglang. + let grpc_futures = grpc_workers.iter().filter_map(|worker| { + let metrics_url = grpc_worker_metrics_url(worker.url())?; + let client = client.clone(); + let worker_url = worker.url().to_string(); + Some(async move { + let result = client + .get(&metrics_url) + .timeout(REQUEST_TIMEOUT) + .send() + .await; + (worker_url, metrics_url, result) + }) + }); + let grpc_results: Vec<_> = stream::iter(grpc_futures) + .buffer_unordered(MAX_CONCURRENT) + .collect() + .await; + for (worker_url, metrics_url, result) in grpc_results { + match result { + Ok(r) if r.status().is_success() => { + if let Ok(text) = r.text().await { + metric_packs.push(MetricPack { + labels: vec![("worker_addr".into(), worker_url)], + metrics_text: text, + }); + } + } + Ok(r) => warn!( + "gRPC worker metrics endpoint {} returned {}", + metrics_url, + r.status() + ), + Err(e) => warn!( + "Failed to scrape gRPC worker metrics from {metrics_url}: {e}" + ), + } + }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/worker/manager.rs` around lines 842 - 878, The gRPC scraping loop currently awaits each request sequentially causing high total latency; replace the for-loop over grpc_workers with the same fan_out pattern used by the HTTP branch: build a stream over grpc_workers (using futures::stream::iter or similar), map each worker into an async closure that calls grpc_worker_metrics_url(worker.url()), does the client.get(...).timeout(REQUEST_TIMEOUT).send().await, processes the response into an Option<MetricPack> (using grpc_worker_metrics_url, client, REQUEST_TIMEOUT, MetricPack), then use .buffer_unordered(MAX_CONCURRENT) to run them concurrently and collect/for_each the resulting MetricPack options into metric_packs; ensure you preserve the existing log paths (warn on non-success and Err) and avoid shared-mutable push races by returning Option<MetricPack> from each future and collecting them into metric_packs after the concurrent execution.grpc_servicer/smg_grpc_servicer/sglang/request_manager.py (1)
804-863:⚠️ Potential issue | 🟡 MinorAborted/preempted requests are not reflected in finished-request metrics.
observe_one_finished_requestis only emitted in_handle_batch_outputon the success-finish branch (L718-736). Requests that terminate via_handle_abort_req(explicit client abort, scheduler-side priority preemption, queue-full, KV-cache pressure) setstate.finished = Trueand push an abort response to the queue, but never record an end-to-end finished observation. Depending on how downstream dashboards compare "started" vs. "finished" counts this will look like a leak, and abort-path latencies will be invisible.If the intent is "only count successful completions", consider at minimum emitting a distinct counter for aborted/preempted terminations (e.g., with a
finish_reasonlabel) so that totals reconcile. Otherwise, mirror the finished-request observation in_handle_abort_reqwith appropriate token counts (0/cumulative from state).🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py` around lines 804 - 863, The abort path in _handle_abort_req currently marks state.finished and pushes an abort_response but does not emit the finished-request metric (observe_one_finished_request) used in _handle_batch_output, so aborted/preempted requests are missing from finished-request metrics; update _handle_abort_req to call the same observation routine (or emit a distinct counter) immediately after marking the state finished and before/after putting to state.out_queue: include request_id/state.rid, finish_reason (use recv_obj.finished_reason or a synthesized {"type":"abort","message":"Abort before prefill"}), and token counts (use state.cumulative_prompt_tokens and state.cumulative_completion_tokens if available, else 0) so dashboards reconcile started vs finished counts; reference symbols: _handle_abort_req, _handle_batch_output, observe_one_finished_request, state.cumulative_prompt_tokens, state.cumulative_completion_tokens, rid_to_state.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py`:
- Around line 209-235: Initialize self._metrics_labels = {} before the try so
it's always defined, then in the try continue to assign labels and create
TokenizerMetricsCollector (using server_args, labels,
bucket_time_to_first_token, bucket_e2e_request_latency, and optional
bucket_inter_token_latency) and set self.metrics_collector and
self._metrics_labels on success; replace the broad except Exception: with
narrower handling—catch ImportError (log a warning with exc_info) and catch
TypeError (or other specific runtime/API-mismatch errors) separately so you log
the exception type and message via logger.warning, ensuring programming errors
aren't silently swallowed while preserving graceful degradation when
observability imports or APIs are missing.
In `@model_gateway/src/worker/manager.rs`:
- Around line 895-903: The function grpc_worker_metrics_url can overflow when
port == 65535; change its logic in grpc_worker_metrics_url to use a checked
addition (e.g., port.checked_add(1)) or explicitly return None if port == 65535
before formatting, so you never wrap to 0; update the code path that computes
port (port: u16 = port_str.parse().ok()?) to handle this check and return None
on overflow, and add a unit test asserting
grpc_worker_metrics_url("grpc://host:65535") -> None.
---
Outside diff comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py`:
- Around line 804-863: The abort path in _handle_abort_req currently marks
state.finished and pushes an abort_response but does not emit the
finished-request metric (observe_one_finished_request) used in
_handle_batch_output, so aborted/preempted requests are missing from
finished-request metrics; update _handle_abort_req to call the same observation
routine (or emit a distinct counter) immediately after marking the state
finished and before/after putting to state.out_queue: include
request_id/state.rid, finish_reason (use recv_obj.finished_reason or a
synthesized {"type":"abort","message":"Abort before prefill"}), and token counts
(use state.cumulative_prompt_tokens and state.cumulative_completion_tokens if
available, else 0) so dashboards reconcile started vs finished counts; reference
symbols: _handle_abort_req, _handle_batch_output, observe_one_finished_request,
state.cumulative_prompt_tokens, state.cumulative_completion_tokens,
rid_to_state.
In `@model_gateway/src/worker/manager.rs`:
- Around line 842-878: The gRPC scraping loop currently awaits each request
sequentially causing high total latency; replace the for-loop over grpc_workers
with the same fan_out pattern used by the HTTP branch: build a stream over
grpc_workers (using futures::stream::iter or similar), map each worker into an
async closure that calls grpc_worker_metrics_url(worker.url()), does the
client.get(...).timeout(REQUEST_TIMEOUT).send().await, processes the response
into an Option<MetricPack> (using grpc_worker_metrics_url, client,
REQUEST_TIMEOUT, MetricPack), then use .buffer_unordered(MAX_CONCURRENT) to run
them concurrently and collect/for_each the resulting MetricPack options into
metric_packs; ensure you preserve the existing log paths (warn on non-success
and Err) and avoid shared-mutable push races by returning Option<MetricPack>
from each future and collecting them into metric_packs after the concurrent
execution.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 51056492-142c-476f-b97c-0ab020b8708a
📒 Files selected for processing (2)
grpc_servicer/smg_grpc_servicer/sglang/request_manager.pymodel_gateway/src/worker/manager.rs
| self.metrics_collector = None | ||
| if server_args.enable_metrics: | ||
| try: | ||
| from sglang.srt.observability.metrics_collector import ( | ||
| TokenizerMetricsCollector, | ||
| ) | ||
|
|
||
| labels = { | ||
| "model_name": server_args.served_model_name, | ||
| } | ||
| self.metrics_collector = TokenizerMetricsCollector( | ||
| server_args=server_args, | ||
| labels=labels, | ||
| bucket_time_to_first_token=server_args.bucket_time_to_first_token, | ||
| bucket_e2e_request_latency=server_args.bucket_e2e_request_latency, | ||
| bucket_inter_token_latency=getattr( | ||
| server_args, "bucket_inter_token_latency", None | ||
| ), | ||
| ) | ||
| self._metrics_labels = labels | ||
| logger.info("TokenizerMetricsCollector initialized for gRPC request-level metrics") | ||
| except Exception: | ||
| logger.warning( | ||
| "Failed to initialize TokenizerMetricsCollector, " | ||
| "request-level metrics will be unavailable", | ||
| exc_info=True, | ||
| ) |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Metrics init: minor hardening — _metrics_labels default and narrower except.
Two small points:
self._metrics_labelsis only assigned inside the successfultrybranch. All current read sites are gated byself.metrics_collector is not None, so this is safe today, but any future reader that forgets the guard will hitAttributeError. Definingself._metrics_labels = {}before thetrymakes the invariant robust to refactors.except Exception:swallows everything includingKeyboardInterruptsubclasses (not in Py3, but) and, more importantly, masks programming errors likeTypeErrorfrom passing an unexpected kwarg (e.g.,bucket_inter_token_latencynot being accepted by older sglang versions). Consider logging the exception type at least, or catchingImportError+ a narrower runtime error separately so API-mismatch bugs surface during development.
Proposed diff
self.metrics_collector = None
+ self._metrics_labels: dict[str, str] = {}
if server_args.enable_metrics:
try:
from sglang.srt.observability.metrics_collector import (
TokenizerMetricsCollector,
)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py` around lines 209 -
235, Initialize self._metrics_labels = {} before the try so it's always defined,
then in the try continue to assign labels and create TokenizerMetricsCollector
(using server_args, labels, bucket_time_to_first_token,
bucket_e2e_request_latency, and optional bucket_inter_token_latency) and set
self.metrics_collector and self._metrics_labels on success; replace the broad
except Exception: with narrower handling—catch ImportError (log a warning with
exc_info) and catch TypeError (or other specific runtime/API-mismatch errors)
separately so you log the exception type and message via logger.warning,
ensuring programming errors aren't silently swallowed while preserving graceful
degradation when observability imports or APIs are missing.
| fn grpc_worker_metrics_url(worker_url: &str) -> Option<String> { | ||
| let stripped = worker_url | ||
| .strip_prefix("grpc://") | ||
| .or_else(|| worker_url.strip_prefix("grpcs://"))?; | ||
| let (host, port_and_suffix) = stripped.rsplit_once(':')?; | ||
| let port_str = port_and_suffix.split(['/', '@']).next()?; | ||
| let port: u16 = port_str.parse().ok()?; | ||
| Some(format!("http://{host}:{}/metrics", port + 1)) | ||
| } |
There was a problem hiding this comment.
Guard against port + 1 overflow on u16.
If port_str parses to 65535, port + 1 panics in debug and wraps to 0 in release, producing a bogus http://host:0/metrics. Unlikely in practice, but trivial to harden and avoids a hidden footgun if someone ever configures a high ephemeral port.
🛡️ Proposed fix
- let port: u16 = port_str.parse().ok()?;
- Some(format!("http://{host}:{}/metrics", port + 1))
+ let port: u16 = port_str.parse().ok()?;
+ let metrics_port = port.checked_add(1)?;
+ Some(format!("http://{host}:{metrics_port}/metrics"))Optionally add a test case asserting grpc_worker_metrics_url("grpc://host:65535") returns None.
📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| fn grpc_worker_metrics_url(worker_url: &str) -> Option<String> { | |
| let stripped = worker_url | |
| .strip_prefix("grpc://") | |
| .or_else(|| worker_url.strip_prefix("grpcs://"))?; | |
| let (host, port_and_suffix) = stripped.rsplit_once(':')?; | |
| let port_str = port_and_suffix.split(['/', '@']).next()?; | |
| let port: u16 = port_str.parse().ok()?; | |
| Some(format!("http://{host}:{}/metrics", port + 1)) | |
| } | |
| fn grpc_worker_metrics_url(worker_url: &str) -> Option<String> { | |
| let stripped = worker_url | |
| .strip_prefix("grpc://") | |
| .or_else(|| worker_url.strip_prefix("grpcs://"))?; | |
| let (host, port_and_suffix) = stripped.rsplit_once(':')?; | |
| let port_str = port_and_suffix.split(['/', '@']).next()?; | |
| let port: u16 = port_str.parse().ok()?; | |
| let metrics_port = port.checked_add(1)?; | |
| Some(format!("http://{host}:{metrics_port}/metrics")) | |
| } |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/worker/manager.rs` around lines 895 - 903, The function
grpc_worker_metrics_url can overflow when port == 65535; change its logic in
grpc_worker_metrics_url to use a checked addition (e.g., port.checked_add(1)) or
explicitly return None if port == 65535 before formatting, so you never wrap to
0; update the code path that computes port (port: u16 = port_str.parse().ok()?)
to handle this check and return None on overflow, and add a unit test asserting
grpc_worker_metrics_url("grpc://host:65535") -> None.
|
Hi @ConnorLi96, this PR has merge conflicts that must be resolved before it can be merged. Please rebase your branch: git fetch origin main
git rebase origin/main
# resolve any conflicts, then:
git push --force-with-lease |
…cs endpoint
Option A (sentinel round-trip): encode ':' to '__smgcolon48f__' before
passing text to openmetrics_parser (which rejects colons), then decode
the sentinel back to ':' after serialization. This preserves colons in
both metric names (sglang:num_running_reqs) and label values
(grpc://host:9001).
Option A chosen over B (passthrough-when-no-labels) because labels are
always injected by get_engine_metrics, making B a no-op. Option C
(raw text injection) was rejected as the largest diff with full
Prometheus syntax responsibility.
The sentinel uses only lowercase to avoid conflicts with
openmetrics_parser's case-sensitive keyword grammar (HELP/TYPE).
Before: sglang_num_running_reqs (colon clobbered to underscore at
metrics_aggregator.rs:19 via .replace(":", "_"))
After: sglang:num_running_reqs (original Prometheus namespace:name
format preserved)
Tracks: consolidation-doc §3.A-ish — colon preservation follow-up to #873
Signed-off-by: Connor Li <ConnorLi96@users.noreply.github.com>
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/worker/metrics_aggregator.rs`:
- Line 23: The current replace call (let encoded = metrics_text.replace(':',
COLON_SENTINEL);) is not input-safe because an input containing the literal
COLON_SENTINEL will be incorrectly decoded; fix by making the transformation
reversible: either (A) pre-escape any existing COLON_SENTINEL occurrences in
metrics_text (e.g., replace COLON_SENTINEL with an escaped variant) before
replacing ':' so round-trip decode distinguishes originals from encoded colons,
or (B) switch to a truly reversible encoding (e.g., percent-encoding or base64)
for metrics_text; alternatively, if you accept the negligible risk, add a clear
comment next to the COLON_SENTINEL constant and the use in encoded/replace
explaining the trust assumption and why collisions are acceptable. Ensure
references to metrics_text, COLON_SENTINEL, and the encoded variable are updated
accordingly.
- Around line 13-16: The comment above the COLON_SENTINEL constant uses the
`SAFETY:` marker but this is a safe-code invariant; change the comment marker to
`INVARIANT:` and keep the rest of the explanation verbatim so it documents the
round-trip assumption about encoding/decoding colons for `namespace:metric_name`
while reserving `SAFETY:` for true unsafe blocks; update the comment that
precedes `const COLON_SENTINEL: &str = "__smgcolon48f__";` accordingly.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: e6539697-1c95-47bc-875f-288083279259
📒 Files selected for processing (2)
model_gateway/src/worker/metrics_aggregator.rsmodel_gateway/tests/metrics_aggregator_test.rs
| // SAFETY: openmetrics_parser rejects colons in metric names. We encode colons to | ||
| // this placeholder before parsing and decode back after serialization, preserving | ||
| // the original `namespace:metric_name` colon format in the output. | ||
| const COLON_SENTINEL: &str = "__smgcolon48f__"; |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Use INVARIANT: instead of SAFETY: for this safe-code assumption.
Per repo convention, SAFETY: is reserved for unsafe blocks; safe-code assumptions should be documented with INVARIANT:. There is no unsafe here — this is a round-trip invariant over the input text.
✏️ Proposed change
-// SAFETY: openmetrics_parser rejects colons in metric names. We encode colons to
-// this placeholder before parsing and decode back after serialization, preserving
-// the original `namespace:metric_name` colon format in the output.
+// INVARIANT: openmetrics_parser rejects colons in metric names. We encode colons to
+// this placeholder before parsing and decode back after serialization, preserving
+// the original `namespace:metric_name` colon format in the output. The sentinel is
+// chosen to be extremely unlikely to collide with real metric text.
const COLON_SENTINEL: &str = "__smgcolon48f__";As per learnings, use the marker INVARIANT: to document assumptions in safe code; reserve SAFETY: for explaining why unsafe blocks are sound.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/worker/metrics_aggregator.rs` around lines 13 - 16, The
comment above the COLON_SENTINEL constant uses the `SAFETY:` marker but this is
a safe-code invariant; change the comment marker to `INVARIANT:` and keep the
rest of the explanation verbatim so it documents the round-trip assumption about
encoding/decoding colons for `namespace:metric_name` while reserving `SAFETY:`
for true unsafe blocks; update the comment that precedes `const COLON_SENTINEL:
&str = "__smgcolon48f__";` accordingly.
| let metrics_text = &metric_pack.metrics_text; | ||
| // openmetrics_parser doesn't handle colons in metric names; replace with underscores | ||
| let metrics_text = metrics_text.replace(":", "_"); | ||
| let encoded = metrics_text.replace(':', COLON_SENTINEL); |
There was a problem hiding this comment.
Sentinel round-trip is robust in practice but not input-safe.
metrics_text.replace(':', COLON_SENTINEL) followed by a reverse replace(COLON_SENTINEL, ":") will misbehave if the literal string __smgcolon48f__ ever appears in input metric text (e.g., inside a label value) — it would be decoded to : in the output. The probability is negligible given the sentinel shape, and upstream sources are trusted Prometheus exposition, so this is a minor caveat rather than a bug. Consider noting this assumption alongside the constant if you plan to take metric text from less-trusted sources later.
Also applies to: 40-40
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/worker/metrics_aggregator.rs` at line 23, The current
replace call (let encoded = metrics_text.replace(':', COLON_SENTINEL);) is not
input-safe because an input containing the literal COLON_SENTINEL will be
incorrectly decoded; fix by making the transformation reversible: either (A)
pre-escape any existing COLON_SENTINEL occurrences in metrics_text (e.g.,
replace COLON_SENTINEL with an escaped variant) before replacing ':' so
round-trip decode distinguishes originals from encoded colons, or (B) switch to
a truly reversible encoding (e.g., percent-encoding or base64) for metrics_text;
alternatively, if you accept the negligible risk, add a clear comment next to
the COLON_SENTINEL constant and the use in encoded/replace explaining the trust
assumption and why collisions are acceptable. Ensure references to metrics_text,
COLON_SENTINEL, and the encoded variable are updated accordingly.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 3f6eec0ebd
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| let (host, port_and_suffix) = stripped.rsplit_once(':')?; | ||
| let port_str = port_and_suffix.split(['/', '@']).next()?; | ||
| let port: u16 = port_str.parse().ok()?; | ||
| Some(format!("http://{host}:{}/metrics", port + 1)) |
There was a problem hiding this comment.
Guard gRPC metrics port increment against overflow
When a worker is configured on port 65535, port + 1 overflows in grpc_worker_metrics_url. In debug builds this can panic during /metrics scraping, and in release builds it wraps to 0, producing an invalid metrics URL and silently dropping that worker’s engine metrics. Using a checked increment (and returning None on overflow) avoids both crash and silent mis-scrape behavior.
Useful? React with 👍 / 👎.
Signed-off-by: Connor Li <ConnorLi96@users.noreply.github.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 9a9b6e0329
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| let _ = ENGINE_METRICS_DEPS.set(EngineMetricsDeps { | ||
| worker_registry, | ||
| client, | ||
| }); |
There was a problem hiding this comment.
Allow engine metrics deps to be refreshed per startup
register_engine_metrics_deps writes into a process-global OnceLock and ignores set failures, so only the first AppContext ever becomes visible to /metrics. In embedded/test scenarios that start the router more than once in the same process, subsequent servers will scrape engine metrics from stale registry/client state (and keep old Arcs alive) instead of their own workers. Use an updatable global (e.g., RwLock<Option<_>>) or explicit reset path rather than a one-shot lock here.
Useful? React with 👍 / 👎.
|
This pull request has been automatically marked as stale because it has not had any activity within 14 days. It will be automatically closed if no further activity occurs within 16 days. Leave a comment if you feel this pull request should remain open. Thank you! |
|
This pull request has been automatically closed due to inactivity. Please feel free to reopen if you intend to continue working on it. Thank you! |
Problem
PROMETHEUS_MULTIPROC_DIR.Solution
PROMETHEUS_MULTIPROC_DIR.dbfiles in gRPC modeChanges
bindings/python/src/smg/serve.py— setPROMETHEUS_MULTIPROC_DIRmodel_gateway/src/core/worker_manager.rs— metrics collection + empty output guardmodel_gateway/src/observability/metrics_server.rs— graceful bind failureCo-authored-by: Scott Lee scott@together.ai
Checklist:
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesMade with Cursor
Summary by CodeRabbit
New Features
Bug Fixes
Chores