Repository navigation
feat(gateway): add --remove-unhealthy-workers - #714
Conversation
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
📝 WalkthroughWalkthroughA new boolean configuration, remove_unhealthy_workers, was added and propagated through CLI, config types, Python bindings, server startup, and the worker registry health-checker so unhealthy workers can be removed automatically when enabled. Documentation was updated to document the option. Changes
Sequence DiagramsequenceDiagram
participant Server
participant HealthChecker
participant WorkerRegistry
participant Storage
Server->>HealthChecker: start_health_checker(interval, remove_unhealthy=true)
HealthChecker->>WorkerRegistry: shallow_clone() (if removal enabled)
loop periodic
HealthChecker->>WorkerRegistry: perform health checks
WorkerRegistry-->>HealthChecker: report unhealthy workers
alt remove_unhealthy == true
HealthChecker->>WorkerRegistry: remove(unhealthy_worker)
WorkerRegistry->>Storage: persist removal / sync state
WorkerRegistry->>HealthChecker: confirm removal
else remove_unhealthy == false
HealthChecker->>WorkerRegistry: mark unhealthy (no removal)
end
end
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~25 minutes Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
📝 Coding Plan
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 |
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request enhances the gateway's worker health check mechanism by providing an option to automatically deregister workers that are deemed unhealthy. This feature improves the reliability and self-healing capabilities of the gateway, ensuring that only responsive and functional workers are part of the routing pool, which is especially beneficial for setups with frequently changing or transient worker instances. Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for GitHub and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
|
Hi @ekzhang, 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:
|
73140eb to
91aa03d
Compare
There was a problem hiding this comment.
Code Review
This pull request introduces a new --remove-unhealthy-workers option that automatically removes workers from the registry after they fail health checks. However, the current implementation poses a significant Denial of Service (DoS) risk as it can lead to the removal of all workers during transient failures, lacking built-in recovery or safeguards to maintain a minimum pool of workers. Additionally, while the implementation is well-structured and consistently applied, there is a suggestion to improve logging by handling the Result from the health check function.
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 (1)
model_gateway/src/core/worker_registry.rs (1)
746-766:⚠️ Potential issue | 🔴 CriticalRemove unhealthy workers by snapshot identity, not by URL.
This races with same-URL re-registration.
remove_by_url()clearsurl_to_idbefore it removes the worker, so a replacement that registers between the check and the delete can be removed or lose its URL mapping. That breaks the exact IGW re-registration flow this flag is meant for. Carry(WorkerId, Arc<dyn Worker>)through the health-check future, then remove only if the current registry entry is still the sameArc. Please also add a regression test for same-URL replacement during an in-flight health check.Possible fix sketch
- let workers: Vec<Arc<dyn Worker>> = workers_ref + let workers: Vec<(WorkerId, Arc<dyn Worker>)> = workers_ref .iter() - .map(|entry| entry.value().clone()) + .map(|entry| (entry.key().clone(), entry.value().clone())) .collect(); // Collect workers whose deadline has passed let due_workers: Vec<_> = workers .iter() - .filter(|w| !w.metadata().health_config.disable_health_check) + .filter(|(_, w)| !w.metadata().health_config.disable_health_check) .filter(|w| { + let (_, worker) = w; next_check - .get(w.url()) + .get(worker.url()) .is_some_and(|deadline| now >= *deadline) }) .cloned() .collect(); // Run due health checks in parallel and schedule the next deadline if !due_workers.is_empty() { for worker in &due_workers { + let (_, worker) = worker; let secs = worker.metadata().health_config.check_interval_secs; let secs = if secs > 0 { secs @@ let futs: Vec<_> = due_workers .into_iter() - .map(|w| async move { + .map(|(worker_id, w)| async move { let _ = w.check_health_async().await; - w + (worker_id, w) }) .collect(); let checked_workers = futures::future::join_all(futs).await; // Remove workers that transitioned to unhealthy if let Some(ref registry) = registry { - for worker in &checked_workers { - if !worker.is_healthy() { + for (worker_id, worker) in &checked_workers { + if !worker.is_healthy() + && registry + .get(worker_id) + .is_some_and(|current| Arc::ptr_eq(¤t, worker)) + { let url = worker.url().to_string(); tracing::warn!( worker_url = %url, "Removing unhealthy worker from registry" ); next_check.remove(&url); - registry.remove_by_url(&url); + registry.remove(worker_id); } } } }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/core/worker_registry.rs` around lines 746 - 766, The health-check currently maps over due_workers and only carries Arc<dyn Worker>, then calls registry.remove_by_url(&url) which races with a same-URL re-registration and can remove a new worker; change the future mapping to carry (WorkerId, Arc<dyn Worker>) through (e.g. map each w -> (w.id(), w) or similar) and after awaiting checked_workers, when a worker is unhealthy verify the registry still points to the exact same Arc (by fetching the current entry by URL or by ID and comparing pointer/ID equality) before removing; use registry.remove_by_id/remove_by_identity (or only call remove_by_url if the registry lookup confirms the stored Arc matches the captured Arc) and update next_check removal to operate on the WorkerId snapshot rather than unconditionally removing by URL; add a regression test that registers a replacement worker for the same URL mid-health-check to assert the new worker is not removed.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@bindings/python/src/lib.rs`:
- Line 782: The Python high-level API never exposes the Rust flag
remove_unhealthy_workers because RouterArgs (in the Python wrapper) lacks that
field and the router wrapper doesn't forward it; add a boolean field
remove_unhealthy_workers (default False) to RouterArgs and ensure the router
wrapper (the function that builds kwargs from RouterArgs) includes that key so
it gets passed through to the Rust binding (keep naming consistent with the Rust
binding and update any dict/serialization code that builds kwargs for the Rust
call).
In `@docs/concepts/reliability/health-checks.md`:
- Line 112: The docs entry for the flag `--remove-unhealthy-workers` is missing
that it disables automatic rejoin/self-healing; update the table row and the
duplicate reference entry to explicitly state that when
`--remove-unhealthy-workers` is true the worker is removed from the registry and
cannot be marked healthy again by later health checks (it will only return if it
re-registers), and ensure the wording references the service's "Self-Healing"
behavior to avoid contradiction.
---
Outside diff comments:
In `@model_gateway/src/core/worker_registry.rs`:
- Around line 746-766: The health-check currently maps over due_workers and only
carries Arc<dyn Worker>, then calls registry.remove_by_url(&url) which races
with a same-URL re-registration and can remove a new worker; change the future
mapping to carry (WorkerId, Arc<dyn Worker>) through (e.g. map each w ->
(w.id(), w) or similar) and after awaiting checked_workers, when a worker is
unhealthy verify the registry still points to the exact same Arc (by fetching
the current entry by URL or by ID and comparing pointer/ID equality) before
removing; use registry.remove_by_id/remove_by_identity (or only call
remove_by_url if the registry lookup confirms the stored Arc matches the
captured Arc) and update next_check removal to operate on the WorkerId snapshot
rather than unconditionally removing by URL; add a regression test that
registers a replacement worker for the same URL mid-health-check to assert the
new worker is not removed.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: d3f6d95c-ace0-4910-9c66-ca030005eb6a
📒 Files selected for processing (7)
bindings/python/src/lib.rsdocs/concepts/reliability/health-checks.mddocs/reference/configuration.mdmodel_gateway/src/config/types.rsmodel_gateway/src/core/worker_registry.rsmodel_gateway/src/main.rsmodel_gateway/src/server.rs
cb2d5a6 to
034d6a1
Compare
|
Hey @slin1237, let me know if you have any feedback on this PR |
|
@ekzhang |
|
Hi @slin1237, I double-checked that format passes |
This option automatically removes unhealthy workers from the gateway registry after failing the configured health check threshold. Especially useful for some setups in inference gateway mode, and you periodically re-register workers. cc @slin1237 Signed-off-by: Eric Zhang <ekzhang1@gmail.com>
Signed-off-by: Eric Zhang <ekzhang1@gmail.com>
034d6a1 to
4a8586e
Compare
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
model_gateway/src/core/worker_registry.rs (1)
703-730:⚠️ Potential issue | 🔴 CriticalMake unhealthy-worker removal identity-safe.
register()reuses the existingWorkerIdfor a URL and overwritesworkers[worker_id]. If a worker re-registers while its oldArcis still insidecheck_health_async(),remove_by_url()here will delete the fresh registration instead of the stale unhealthy instance. That race hits the exact “periodically re-registered workers” deployment this flag is meant for. Carry theWorkerIdthrough the snapshot and only remove when the registry still points to the sameArcthat was checked.🐛 Proposed fix
- let workers: Vec<Arc<dyn Worker>> = workers_ref - .iter() - .map(|entry| entry.value().clone()) - .collect(); + let workers: Vec<(WorkerId, Arc<dyn Worker>)> = workers_ref + .iter() + .map(|entry| (entry.key().clone(), entry.value().clone())) + .collect(); // Sync schedule with registry: add new workers, prune removed // and disabled ones so stale deadlines don't cause wakeups. let checkable_urls: std::collections::HashSet<String> = workers .iter() - .filter(|w| !w.metadata().health_config.disable_health_check) - .map(|w| w.url().to_string()) + .filter(|(_, w)| !w.metadata().health_config.disable_health_check) + .map(|(_, w)| w.url().to_string()) .collect(); next_check.retain(|url, _| checkable_urls.contains(url)); for url in &checkable_urls { next_check.entry(url.clone()).or_insert(now); } // Collect workers whose deadline has passed let due_workers: Vec<_> = workers .iter() - .filter(|w| !w.metadata().health_config.disable_health_check) - .filter(|w| { + .filter(|(_, w)| !w.metadata().health_config.disable_health_check) + .filter(|(_, w)| { next_check .get(w.url()) .is_some_and(|deadline| now >= *deadline) }) .cloned() .collect(); // Run due health checks in parallel and schedule the next deadline if !due_workers.is_empty() { - for worker in &due_workers { + for (_, worker) in &due_workers { let secs = worker.metadata().health_config.check_interval_secs; let secs = if secs > 0 { secs } else { default_interval_secs @@ } let futs: Vec<_> = due_workers .into_iter() - .map(|w| async move { + .map(|(worker_id, w)| async move { let _ = w.check_health_async().await; - w + (worker_id, w) }) .collect(); let checked_workers = futures::future::join_all(futs).await; // Remove workers that transitioned to unhealthy if let Some(ref registry) = registry { - for worker in &checked_workers { + for (worker_id, worker) in &checked_workers { if !worker.is_healthy() { let url = worker.url().to_string(); + let Some(current) = registry.get(worker_id) else { + continue; + }; + if !Arc::ptr_eq(¤t, worker) { + tracing::debug!( + worker_url = %url, + "Skipping removal because worker was replaced during health check" + ); + continue; + } tracing::warn!( worker_url = %url, "Removing unhealthy worker from registry" ); next_check.remove(&url); - registry.remove_by_url(&url); + registry.remove(worker_id); } } } }Also applies to: 746-766
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/core/worker_registry.rs` around lines 703 - 730, The snapshot currently clones only worker Arcs and uses url-based removal which races with re-registration; instead capture and carry the WorkerId alongside each Arc in the snapshot (e.g. produce Vec<(WorkerId, Arc<dyn Worker>)>) in the check_health_async() snapshot/collection logic used around the next_check/due_workers code and the similar block at the other site, and when a worker is found unhealthy verify the registry still points to the same Arc before removing: either use a remove_by_id or fetch registry entry by WorkerId and compare with Arc::ptr_eq to the snapshot Arc, only then call remove_by_url/remove_by_id. This ensures re-registered workers with the same WorkerId are not incorrectly removed.
♻️ Duplicate comments (1)
docs/concepts/reliability/health-checks.md (1)
112-112:⚠️ Potential issue | 🟡 MinorDocument that this opts out of self-healing.
With this flag enabled the worker is removed from the registry, so later health checks cannot mark it healthy again; it only returns if something re-registers it. Please spell that out here, and mirror the same wording in the duplicate reference / CLI descriptions, so this page does not contradict its own “Self-Healing” section.
📝 Suggested wording
-| `--remove-unhealthy-workers` | `false` | Remove workers after being marked unhealthy | +| `--remove-unhealthy-workers` | `false` | Remove workers from the registry after they are marked unhealthy. Use this only when workers are periodically re-registered, because removed workers do not automatically rejoin through later health checks |🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@docs/concepts/reliability/health-checks.md` at line 112, Clarify that the `--remove-unhealthy-workers` flag opts out of self-healing by explicitly stating that when enabled the worker is removed from the registry and cannot be marked healthy again by subsequent health checks (it will only return if re-registered); update the sentence for `--remove-unhealthy-workers` to include this wording and then mirror that exact phrasing in the duplicate reference and any CLI description strings to ensure consistency with the “Self-Healing” section and avoid contradiction.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Outside diff comments:
In `@model_gateway/src/core/worker_registry.rs`:
- Around line 703-730: The snapshot currently clones only worker Arcs and uses
url-based removal which races with re-registration; instead capture and carry
the WorkerId alongside each Arc in the snapshot (e.g. produce Vec<(WorkerId,
Arc<dyn Worker>)>) in the check_health_async() snapshot/collection logic used
around the next_check/due_workers code and the similar block at the other site,
and when a worker is found unhealthy verify the registry still points to the
same Arc before removing: either use a remove_by_id or fetch registry entry by
WorkerId and compare with Arc::ptr_eq to the snapshot Arc, only then call
remove_by_url/remove_by_id. This ensures re-registered workers with the same
WorkerId are not incorrectly removed.
---
Duplicate comments:
In `@docs/concepts/reliability/health-checks.md`:
- Line 112: Clarify that the `--remove-unhealthy-workers` flag opts out of
self-healing by explicitly stating that when enabled the worker is removed from
the registry and cannot be marked healthy again by subsequent health checks (it
will only return if re-registered); update the sentence for
`--remove-unhealthy-workers` to include this wording and then mirror that exact
phrasing in the duplicate reference and any CLI description strings to ensure
consistency with the “Self-Healing” section and avoid contradiction.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: d1fd0da5-a0d4-4972-9898-e55407911414
📒 Files selected for processing (8)
bindings/python/src/lib.rsbindings/python/src/smg/router_args.pydocs/concepts/reliability/health-checks.mddocs/reference/configuration.mdmodel_gateway/src/config/types.rsmodel_gateway/src/core/worker_registry.rsmodel_gateway/src/main.rsmodel_gateway/src/server.rs
Signed-off-by: Eric Zhang <ekzhang1@gmail.com>
Description
This option automatically removes unhealthy workers from the gateway registry after failing the configured health check threshold.
Especially useful for some setups in inference gateway mode, and you periodically re-register workers.
cc @slin1237
Problem
Solution
Changes
CLI option was added
Note: I updated the CLI options and documentation as well, it's documented in two places right now
Test Plan
Tests in
worker_registry.rsChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Documentation
Tests