Repository navigation
feat(worker): add Draining state with workflow-driven drain - #1491
Conversation
Selecting a worker that has just received a `RemoveWorker` job opens a narrow window where in-flight requests can be sent to a worker that is about to be torn down. Add a transitional `WorkerStatus::Draining` so policies stop selecting the worker the moment removal starts, and a new `DrainWorkersStep` in the worker_removal workflow that holds the worker in `Draining` for a configurable settle window before the existing remove-from-* steps run. Drain semantics live in the workflow, not service discovery: any `Job::RemoveWorker` (K8s pod deletion, --remove-unhealthy-workers, manual API) gets uniform draining. The step skips the sleep entirely when no worker in the batch is Ready, so cleanup of broken workers stays fast. Per-worker overrides follow the existing `HealthCheckUpdate` pattern: default lives on `HealthCheckConfig::drain_settle_secs` (5s, exposed as --drain-settle-secs), and any worker may override via its `WorkerSpec::health.drain_settle_secs`. The step takes the max across the batch for a single sleep. Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (2)
📝 WalkthroughWalkthroughThis PR implements a graceful worker draining mechanism by introducing a ChangesWorker Draining with Configurable Settle Period
Sequence Diagram(s)sequenceDiagram
participant Client
participant Workflow
participant DrainStep
participant Registry
participant Sleep
Client->>Workflow: request RemoveWorker(workers)
Workflow->>DrainStep: execute(context)
DrainStep->>Registry: transition_ready_to_draining(id, rev)
Registry-->>DrainStep: transition result
DrainStep->>Sleep: sleep(max_drain_settle_secs)
Sleep-->>DrainStep: wake
DrainStep-->>Workflow: return Success
Workflow->>Client: completion / result
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~25 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Warning Review ran into problems🔥 ProblemsGit: Failed to clone repository. Please run the Comment |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f89b15d64d
ℹ️ 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".
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@bindings/python/src/lib.rs`:
- Around line 789-790: The new parameter drain_settle_secs was inserted
mid-signature of _Router which breaks positional-arg callers; update the API by
either moving drain_settle_secs to the end of the _Router/__init__ parameter
list or make new parameters keyword-only (add a * before it in the signature) so
existing positional usage (e.g., Router/Robot wrapping) remains compatible;
adjust the _Router constructor signature and any internal callers (e.g.,
Robot(_Router(...))) accordingly and add a brief note/test ensuring external
positional calls continue to work.
In `@model_gateway/src/workflow/job_queue.rs`:
- Around line 397-402: The timeout for the RemoveWorker path is currently built
using only the global default (router_config.health_check.drain_settle_secs)
which can undercut longer per-worker overrides; update the timeout calculation
where timeout_duration is set in job_queue.rs (the RemoveWorker handling) to use
the worker's effective drain_settle_secs when present (e.g., take max(global
router_config.health_check.drain_settle_secs, worker.override.drain_settle_secs)
or otherwise use the worker-config value) so the Duration::from_secs uses the
per-worker override if larger; ensure you reference the worker config / job
payload that contains the override when computing the final timeout.
🪄 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: 3e73c7b6-5833-4754-a1f6-5ba4b7e094a9
📒 Files selected for processing (14)
bindings/python/src/lib.rscrates/protocols/src/worker.rsmodel_gateway/Cargo.tomlmodel_gateway/src/config/types.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/worker/builder.rsmodel_gateway/src/worker/manager.rsmodel_gateway/src/worker/registry.rsmodel_gateway/src/worker/worker.rsmodel_gateway/src/workflow/job_queue.rsmodel_gateway/src/workflow/steps/local/drain_workers.rsmodel_gateway/src/workflow/steps/local/mod.rstui/src/app.rs
| /// callers that need to invoke `transition_status_if_revision` with | ||
| /// the current worker revision. | ||
| pub fn get_id_by_url(&self, url: &str) -> Option<WorkerId> { | ||
| self.url_to_id.get(url).map(|id| id.clone()) |
There was a problem hiding this comment.
🟡 Nit: Clippy preference — .map(|id| id.clone()) can be .cloned().
| self.url_to_id.get(url).map(|id| id.clone()) | |
| self.url_to_id.get(url).map(|id| id.clone()) |
(Actually, DashMap's Ref doesn't impl the right trait for .cloned() on the Option<Ref> — disregard if that's the case here.)
| WorkerStatus::Failed | WorkerStatus::Draining => { | ||
| // Terminal for the health-state machine. Failed is removed | ||
| // by `--remove-unhealthy-workers`; Draining is removed by | ||
| // the discovery drain timer once in-flight requests settle. | ||
| } |
There was a problem hiding this comment.
🟡 Nit: Draining is correctly terminal in the state machine, but launch_due_probes (line ~399) only short-circuits for Failed — Draining workers still get health-probed every interval even though compute_next_status always returns None for them.
Consider adding Draining alongside Failed in the early-exit check in launch_due_probes to skip pointless probes during the drain window:
if launched_status == WorkerStatus::Failed || launched_status == WorkerStatus::Draining {
next_check.remove(&worker_id);
...(For Draining you'd skip the removal push since drain-initiated removal handles that.)
There was a problem hiding this comment.
Clean, well-structured addition of the Draining worker state. The workflow-driven drain step, config plumbing across CLI/Python/protocol layers, and comprehensive tests all look solid.
Summary: 0 🔴 Important · 3 🟡 Nit · 0 🟣 Pre-existing
Nits:
job_queue.rs: Timeout forwait_for_completionuses only the globaldrain_settle_secs— per-worker overrides could cause spurious timeout errors.registry.rs: Minor style —.map(|id| id.clone())→.cloned().manager.rs:Drainingworkers still get health-probed even thoughcompute_next_statustreats them as terminal. Consider skipping probes likeFaileddoes.
There was a problem hiding this comment.
Code Review
This pull request introduces a Draining status for workers and a configurable drain_settle_secs period to allow in-flight requests to finish before a worker is removed from the registry. Key changes include the addition of the DrainWorkersStep to the worker removal workflow, updates to the health check configuration, and new CLI parameters. Review feedback highlights a potential timeout issue in the job queue when per-worker overrides exceed the global default and suggests refactoring the use of '0' as a special value for disabling the settle period to align with project conventions.
- bindings/python: move `drain_settle_secs` to the end of the `_Router` parameter list so external positional callers aren't broken by the mid-signature insertion (CodeRabbit, Codex). - workflow/job_queue: cap the `RemoveWorker` caller wait at `30 + max(global_drain_settle_secs, 600)` so a per-worker override larger than the global default no longer produces a spurious timeout error from the queue while the workflow is still draining (Codex, CodeRabbit, Claude, Gemini). - worker/manager: short-circuit `launch_due_probes` for `Draining` workers the same way `Failed` is handled — `compute_next_status` is already a no-op for them, and the drain workflow owns removal so the probe pass should not push another `RemoveWorker` candidate (Claude). - workflow/drain_workers: drop the `max_drain_secs == 0` short-circuit; `sleep(Duration::ZERO)` is a no-op, and removing the special-case keeps the step from inventing a new "0 disables" convention (Gemini). Signed-off-by: Simo Lin <25425177+slin1237@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: c5f0987402
ℹ️ 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".
| // Timeout is generous so the per-worker drain_settle_secs sleep | ||
| // can complete even when several workers stack their settle | ||
| // windows. The step itself caps by taking max() across them. | ||
| .with_timeout(Duration::from_secs(600)) |
There was a problem hiding this comment.
Avoid timing out valid drain windows
When any worker (or the global CLI/config) sets drain_settle_secs above 600, DrainWorkersStep sleeps that full value but the workflow engine wraps this step in this fixed 600s timeout (tokio::time::timeout in the workflow engine marks it failed). In that configuration the worker has already been moved to Draining, the downstream removal steps never run, and the worker can remain stuck out of service instead of being removed. Even with the new caller wait timeout, the step timeout itself still needs to be derived from or capped against the configured drain window.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@model_gateway/src/workflow/steps/local/drain_workers.rs`:
- Around line 46-64: The loop currently uses the snapshot `worker` (from
`workers_to_remove`) to decide drain eligibility and revision; instead, fetch
the live registry entry before deciding and transitioning: use
`app_context.worker_registry.get_id_by_url(url)` (as already present) then
retrieve the current registry entry (e.g. via the registry's
get/get_by_id/get_entry method) and use that entry's status and revision for the
`WorkerStatus::Ready` check and for the call to
`transition_status_if_revision(&worker_id, revision, WorkerStatus::Draining)`;
only proceed if the registry's status is `Ready`, and read
`health_config.drain_settle_secs` from the registry entry's metadata for
updating `max_drain_secs`, so decisions are based on current registry state
rather than the snapshot `worker`.
🪄 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: 22a44e95-5caf-4636-9087-c86909a656db
📒 Files selected for processing (4)
bindings/python/src/lib.rsmodel_gateway/src/worker/manager.rsmodel_gateway/src/workflow/job_queue.rsmodel_gateway/src/workflow/steps/local/drain_workers.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 5d15749ad9
ℹ️ 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".
| // No `with_timeout`: the step's purpose is to sleep for the | ||
| // resolved `drain_settle_secs` (which can be set per-worker | ||
| // via `WorkerSpec::health.drain_settle_secs`). A static | ||
| // workflow-level timeout would be set at definition time | ||
| // without visibility into runtime config and would | ||
| // pre-emptively fail the workflow for legitimately long | ||
| // drain windows, leaving workers stuck in `Draining`. |
There was a problem hiding this comment.
🟡 Nit: This is now the only step in the entire codebase without a with_timeout. Every other step — including variable-duration ones — sets an upper bound (e.g., mcp_registration uses 7200s, tokenizer_registration uses 300s). Removing the timeout entirely means a misconfigured drain_settle_secs (or a future bug in the sleep path) would hang the workflow indefinitely, leaving workers stuck in Draining.
The previous 600s cap was too rigid, but the fix could be a generous safety-net timeout rather than none at all — e.g., with_timeout(Duration::from_secs(3600)) gives a 1-hour ceiling that no realistic drain window should hit, while still preventing unbounded hangs.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@model_gateway/src/workflow/steps/local/drain_workers.rs`:
- Around line 297-305: The test incorrectly shares the same Arc-backed worker
between snapshot and live registry so calling
worker.set_status(WorkerStatus::Ready) mutates both; update the test to make the
snapshot contain a distinct worker instance (e.g. create a second Arc via
build_worker or clone a new worker with the same id/status) instead of reusing
the live Arc so the snapshot stays Pending while the live worker is set to
Ready; adjust the snapshot creation (the snapshot variable) while leaving
make_app_context and the live worker unchanged, and ensure any other occurrences
(around lines with make_context) use the independent snapshot instance.
🪄 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: 69cdb44d-8ebb-4ff8-9293-03d386a6e987
📒 Files selected for processing (2)
model_gateway/src/workflow/steps/local/drain_workers.rsmodel_gateway/src/workflow/steps/local/mod.rs
- workflow/drain_workers: read status / revision / drain_settle_secs from the live registry entry, not the snapshot captured by find_workers_to_remove. A worker that became `Ready` between the snapshot and this step (eg. a probe just succeeded) would otherwise bypass the transition and start serving traffic up until the registry-removal step ran. Adds a regression test that flips a Pending → Ready transition between snapshot and execute and confirms the step still drains it (CodeRabbit). - workflow/local/mod: drop the static `with_timeout(600s)` on the drain step. Workflow-level timeouts are set at definition time so they cannot track the runtime `drain_settle_secs`, and a per-worker override above 600s would otherwise pre-emptively fail the workflow with the worker stuck in `Draining` (Codex). Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
5d15749 to
250f74b
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 250f74bfc2
ℹ️ 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".
| // `MAX_DRAIN_WAIT_SECS` (large enough for realistic | ||
| // overrides), and at least the global default if it's | ||
| // higher. 30s on top covers the other removal steps. | ||
| const MAX_DRAIN_WAIT_SECS: u64 = 600; |
There was a problem hiding this comment.
Derive removal wait from the actual drain window
Fresh evidence after the earlier timeout review is that this version still uses a fixed 600s floor here instead of the per-worker value. When a worker sets health.drain_settle_secs above 600 while the global setting remains lower, DrainWorkersStep sleeps the full per-worker value, but wait_for_completion returns Workflow timeout after 630s and the RemoveWorker job is reported failed even though the drain is valid and still in progress. The wait budget needs to be computed from the workers being removed, or the configured per-worker drain window must be capped to match this timeout.
Useful? React with 👍 / 👎.
Summary
WorkerStatus::Draining = 4so policies stop selecting a worker the momentRemoveWorkerstarts, closing the window where a request could be routed to a worker that's about to be torn down.DrainWorkersStepworkflow step (inserted betweenfind_workers_to_removeandremove_from_policy_registry) transitions each Ready worker in the removal batch toDraining(revision-checked, so a same-URL replace isn't clobbered) and sleeps the maxdrain_settle_secsacross the batch. Non-Ready workers (Pending/NotReady/Failed) skip both transition and sleep — keeps--remove-unhealthy-workersfast.Job::RemoveWorker(K8s pod deletion,--remove-unhealthy-workers, manual API) gets uniform drain behaviour.--drain-settle-secs(default 5s) onHealthCheckConfig::drain_settle_secs, with per-worker overrides throughWorkerSpec::health.drain_settle_secs(mirrors the existingHealthCheckUpdatepattern). Set to0to skip draining entirely.WorkerRegistry::get_id_by_urlhelper added so the step can calltransition_status_if_revisionwithout scanning all workers.30 + drain_settle_secsto accommodate the drain window.Test plan
cargo +nightly fmt --all(silent)cargo clippy --all-targets --all-features -- -D warnings(clean)cargo test— 3518 passed / 0 failed across the workspace.crates/protocols):try_from_u8(4),from_u8(4),Display→"draining", serde round-trip Draining ↔"draining",is_routable()false for Draining, parametrized routability check.crates/protocols): defaultdrain_settle_secs == 5;HealthCheckUpdateoverride;is_emptyincludes the new field; existing JSON without the field deserializes (serde default).get_healthy_worker_indicesparametrized — Ready included, Pending/NotReady/Failed/Draining excluded.get_id_by_url(registered + missing);transition_statusto Draining emitsWorkerEvent::StatusChanged(proves mesh sync forwards Draining the same way as Ready/NotReady).tokio::time::advance): empty list → no-op; only-non-Ready → no transition + no sleep; single Ready → immediate Draining + sleep until window elapses; multiple workers → uses maxdrain_settle_secs;drain_settle_secs == 0→ still transitions to Draining (so policies stop) but skips the sleep.Summary by CodeRabbit
New Features
Behavior Changes