refactor(worker): extract health loop into WorkerManager - #1105
Conversation
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
📝 WalkthroughWalkthroughConsolidates per-worker runtime state into Changes
Sequence Diagram(s)sequenceDiagram
participant Server as Server
participant Manager as WorkerManager
participant Registry as WorkerRegistry
participant Worker as BasicWorker
participant Queue as JobQueue
Server->>Manager: WorkerManager::start(registry, config, job_queue)
activate Manager
Manager->>Registry: reconcile_snapshot()
Registry-->>Manager: Vec<WorkerDescriptor>
Manager->>Manager: schedule probe deadlines
loop Concurrent probes
Manager->>Worker: check_health_async()
Worker-->>Manager: Ok() / Err(HealthCheckFailed)
alt success
Manager->>Registry: transition_status_if_revision(worker_id, expected_rev, Ready)
Registry-->>Manager: Some((old,new)) / None
else failure
Manager->>Manager: update counters, compute next status
alt becomes Failed and remove_unhealthy
Manager->>Queue: Job::RemoveWorker{url, expected_revision: Some(rev)}
else
Manager->>Registry: transition_status_if_revision(..., NotReady/Pending)
end
end
end
deactivate Manager
Estimated code review effort🎯 4 (Complex) | ⏱️ ~75 minutes Possibly related PRs
Suggested labels
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 docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Code Review
This pull request refactors the worker health check and lifecycle management system by introducing a WorkerManager to centralize the health check loop and state machine logic. It also implements a WorkerRuntime struct to manage shared mutable state, allowing for state preservation and revision tracking during worker replacements. Review feedback suggests pinning the shutdown notification future in the WorkerManager loop to prevent missed signals and notes that skipping Failed workers during reconciliation could lead to stale registry entries.
| if descriptor.status == WorkerStatus::Failed { | ||
| // Startup reconcile and lagged rebuild must be side-effect-free: | ||
| // do not reschedule already-failed workers for probing or removal. | ||
| next_check.remove(&descriptor.worker_id); | ||
| return; | ||
| } |
There was a problem hiding this comment.
Skipping Failed workers during reconcile (startup or lag recovery) prevents them from being removed by the WorkerManager if remove_unhealthy is enabled. While the intent is to keep reconcile side-effect-free, this may leave stale Failed workers in the registry indefinitely if they were added via mesh sync or were already present during a lag event.
| Ok(WorkerEvent::Replaced { worker_id, new, .. }) => { | ||
| schedule_worker_at( | ||
| &mut next_check, | ||
| worker_id, | ||
| new.status(), | ||
| &new.metadata().health_config, | ||
| &config, | ||
| tokio::time::Instant::now(), | ||
| true, | ||
| ); | ||
| } |
There was a problem hiding this comment.
🟡 Nit: When a Replaced event arrives while an old probe is still in-flight for the same worker, schedule_worker_at(..., immediate: true) sets the next_check deadline to now. However, queue_due_probes skips workers that are in in_flight, so the immediate re-probe can't be launched. The stale past-deadline entry then causes sleep_until to resolve instantly on every loop iteration until the old probe completes or times out — a brief busy-loop (bounded by the health check timeout).
Consider also clearing the worker from in_flight here — the old probe will be correctly discarded by revision fencing anyway:
| Ok(WorkerEvent::Replaced { worker_id, new, .. }) => { | |
| schedule_worker_at( | |
| &mut next_check, | |
| worker_id, | |
| new.status(), | |
| &new.metadata().health_config, | |
| &config, | |
| tokio::time::Instant::now(), | |
| true, | |
| ); | |
| } | |
| Ok(WorkerEvent::Replaced { worker_id, new, .. }) => { | |
| in_flight.remove(&worker_id); | |
| schedule_worker_at( | |
| &mut next_check, | |
| worker_id, | |
| new.status(), | |
| &new.metadata().health_config, | |
| &config, | |
| tokio::time::Instant::now(), | |
| true, | |
| ); | |
| } |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 2c48a7ec87
ℹ️ 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 sleep_until = next_check | ||
| .values() | ||
| .min() | ||
| .copied() |
There was a problem hiding this comment.
Exclude in-flight workers from sleep deadline selection
The loop picks sleep_until from all next_check entries, including workers already in in_flight. If a probe takes longer than its interval (common when timeout > interval or endpoints are slow), that worker’s deadline stays in the past, queue_due_probes skips it because it is in-flight, and tokio::time::sleep_until(sleep_until) fires immediately on every iteration, creating a tight CPU spin until the probe finishes. This can peg a core under degraded workers and starve other tasks; compute the next wakeup from non-in-flight workers (or clamp past deadlines when all due workers are in-flight).
Useful? React with 👍 / 👎.
| ) { | ||
| let now = tokio::time::Instant::now(); | ||
| for descriptor in registry.reconcile_snapshot() { | ||
| schedule_descriptor_at(next_check, descriptor, config, now, false); |
There was a problem hiding this comment.
Schedule reconciled workers for immediate probing
During startup and lag recovery, reconcile_from_registry calls schedule_descriptor_at(..., false), which defers every existing worker’s first probe by a full interval. Because health-checked workers start as Pending, this delays them becoming routable after startup (and similarly delays recovery after Lagged rebuilds) even when they are healthy. Reconciled workers should be scheduled immediately (while still skipping side effects for already-Failed workers) to avoid avoidable availability gaps.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 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/service_discovery.rs`:
- Around line 641-644: The RemoveWorker jobs are being enqueued with
expected_revision: None which allows a stale remove to match a newly replaced
worker; before creating Job::RemoveWorker (the instances at the locations that
set expected_revision: None), read the current worker's revision (snapshot the
worker.revision or equivalent field for old_url) and set expected_revision:
Some(current_revision) when constructing Job::RemoveWorker so the existing
find_workers_to_remove logic can ignore stale queued removals; apply the same
change at all three enqueue sites that currently use expected_revision: None and
ensure AddWorker continues to use the replacement revision semantics.
In `@model_gateway/src/worker/manager.rs`:
- Around line 252-269: The loop can hot-spin because sleep_until is set to an
already-expired deadline when probes are saturated; update the logic that
computes sleep_until (which currently uses next_check.min()) to detect when the
earliest deadline <= now and in_flight.len() >= MAX_CONCURRENT_HEALTH_PROBES (or
when queue_due_probes() could not schedule more), and in that case choose a
future backoff (e.g., now +
Duration::from_secs(config.default_check_interval_secs) or a short jittered
backoff) instead of the expired time so the tokio::time::sleep_until branch will
actually sleep; make this change around the code that computes sleep_until and
the async select using probes, next_check, in_flight, and
MAX_CONCURRENT_HEALTH_PROBES to prevent busy-waiting when probe concurrency is
saturated.
In `@model_gateway/src/worker/registry.rs`:
- Around line 991-1000: The current logic allocates an entry in
worker_mutation_locks for worker_id before checking whether the worker actually
exists, causing orphaned Mutex entries when workers are missing; change the flow
in the helper(s) that use worker_mutation_locks (the block creating lock via
worker_mutation_locks.entry(...).or_insert_with(...).clone()) so you first probe
self.workers.get(worker_id) and only then create/clone the lock for extant
workers, or if you must create the lock early then remove the just-created entry
when self.workers.get(worker_id) returns None (and before returning None); apply
the same fix to the other occurrence around the revision check (the second
helper at the 1028–1038 region).
In `@model_gateway/src/worker/worker.rs`:
- Around line 151-167: The Worker trait currently provides default no-op
implementations for revision() and inherit_shared_state_from(), allowing new
impl Worker to opt out of replacement semantics; remove the default bodies to
make revision(&self) -> u64 and inherit_shared_state_from(&self, _other: &dyn
Worker) -> bool required methods on the Worker trait so every implementor must
explicitly handle revision fencing and shared-state adoption; update any
existing impls (e.g., BasicWorker) to implement these methods accordingly to
preserve current behavior.
🪄 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: 6ee84468-852d-4913-a768-58ac831271e2
📒 Files selected for processing (16)
bindings/golang/src/policy.rsmodel_gateway/src/observability/metrics_ws/collectors.rsmodel_gateway/src/policies/manual.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/server.rsmodel_gateway/src/service_discovery.rsmodel_gateway/src/worker/builder.rsmodel_gateway/src/worker/circuit_breaker.rsmodel_gateway/src/worker/manager.rsmodel_gateway/src/worker/registry.rsmodel_gateway/src/worker/service.rsmodel_gateway/src/worker/worker.rsmodel_gateway/src/workflow/job_queue.rsmodel_gateway/src/workflow/steps/local/find_workers_to_remove.rsmodel_gateway/src/workflow/steps/local/mod.rsmodel_gateway/src/workflow/steps/local/update_worker_properties.rs
| /// Get the current monotonic worker revision. | ||
| /// | ||
| /// Same-URL `replace()` increments the revision so stale probe outcomes | ||
| /// can be discarded without mutating the newly installed worker object. | ||
| fn revision(&self) -> u64 { | ||
| 0 | ||
| } | ||
|
|
||
| /// Set the worker's lifecycle status. | ||
| fn set_status(&self, status: WorkerStatus); | ||
|
|
||
| /// Adopt shared mutable runtime state from a previous worker object. | ||
| /// | ||
| /// Used by same-URL `replace()` so in-flight traffic and counters remain | ||
| /// attached to a single shared runtime across the old and new objects. | ||
| fn inherit_shared_state_from(&self, _other: &dyn Worker) -> bool { | ||
| false |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
set -euo pipefail
python - <<'PY'
import pathlib
import re
def iter_impls(text: str):
for match in re.finditer(r'impl[^{]*\bWorker\b[^{]*\bfor\b[^{]*\{', text):
header = match.group(0)[:-1].strip()
start = match.end() - 1 # points at the opening brace
depth = 0
i = start
while i < len(text):
ch = text[i]
if ch == '{':
depth += 1
elif ch == '}':
depth -= 1
if depth == 0:
yield header, text[start + 1:i]
break
i += 1
for path in sorted(pathlib.Path('.').rglob('*.rs')):
text = path.read_text()
for header, body in iter_impls(text):
has_revision = re.search(r'\bfn\s+revision\s*\(', body) is not None
has_inherit = re.search(r'\bfn\s+inherit_shared_state_from\s*\(', body) is not None
print(f"{path}: {header}")
print(f" revision override: {has_revision}")
print(f" inherit_shared_state_from override: {has_inherit}")
PYRepository: lightseekorg/smg
Length of output: 320
🏁 Script executed:
# Check where replace() is called and whether GrpcWorker instances participate
rg -n "\.replace\(" model_gateway/src/core/worker_registry.rs | head -20Repository: lightseekorg/smg
Length of output: 141
🏁 Script executed:
# Check the WorkerRegistry context to see what worker types flow through replace()
rg -B5 -A5 "replace\(" model_gateway/src/core/worker_registry.rs | head -40Repository: lightseekorg/smg
Length of output: 141
🏁 Script executed:
# Check if GrpcWorker is only used in a specific context that excludes it from replace()
rg -n "GrpcWorker" --type rustRepository: lightseekorg/smg
Length of output: 831
🏁 Script executed:
# Find the actual structure and locate worker_registry
find . -name "worker_registry*" -o -name "*registry*" | grep -E "\.(rs|toml)$" | head -20Repository: lightseekorg/smg
Length of output: 541
🏁 Script executed:
# Check the directory structure
ls -la model_gateway/src/core/ 2>/dev/null || ls -la model_gateway/src/ 2>/dev/null || find . -type d -name "model_gateway" | head -5Repository: lightseekorg/smg
Length of output: 1035
🏁 Script executed:
# Search for where replace() is actually called on workers
rg -n "\.replace\(" --type rust | grep -i workerRepository: lightseekorg/smg
Length of output: 758
🏁 Script executed:
# Check the actual worker registry implementation
head -100 model_gateway/src/worker/registry.rsRepository: lightseekorg/smg
Length of output: 3867
🏁 Script executed:
# Check what Worker types are registered/replaced in the registry
rg "impl Worker" model_gateway/src/worker/ --type rustRepository: lightseekorg/smg
Length of output: 832
🏁 Script executed:
# Look at the replace() method signature and see which types it operates on
rg -B3 -A10 "fn replace\(" model_gateway/src/worker/registry.rsRepository: lightseekorg/smg
Length of output: 701
🏁 Script executed:
# Check if GrpcWorker instances are registered in the WorkerRegistry
rg "GrpcWorker" model_gateway/ --type rustRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Look at the full replace() method to understand what happens
rg -A30 "pub fn replace\(" model_gateway/src/worker/registry.rsRepository: lightseekorg/smg
Length of output: 1373
🏁 Script executed:
# Check what happens if inherit_shared_state_from returns false in replace()
rg -B5 -A15 "inherit_shared_state_from" model_gateway/src/worker/registry.rsRepository: lightseekorg/smg
Length of output: 846
🏁 Script executed:
# Check if GrpcWorker is in bindings and whether it's integrated into the main registry
cat bindings/golang/src/policy.rs | head -120Repository: lightseekorg/smg
Length of output: 3786
🏁 Script executed:
# Look at the actual inherit_shared_state_from implementation in BasicWorker
rg -B2 -A10 "fn inherit_shared_state_from" model_gateway/src/worker/worker.rsRepository: lightseekorg/smg
Length of output: 1090
🏁 Script executed:
# Check the Worker trait definition to see current signatures
rg -B5 -A20 "trait Worker" model_gateway/src/worker/worker.rs | head -80Repository: lightseekorg/smg
Length of output: 1716
🏁 Script executed:
# Check if GrpcWorker is used with WorkerRegistry or any replace() calls
rg "GrpcWorker" . --type rust -A3 -B3 | grep -E "(registry|replace|Registry|Replace)"Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Look at where GrpcWorker is instantiated and used
rg "GrpcWorker::new" . --type rust -B5 -A10Repository: lightseekorg/smg
Length of output: 1133
🏁 Script executed:
# Verify the trait definition - look at the actual default implementations
sed -n '151,250p' model_gateway/src/worker/worker.rsRepository: lightseekorg/smg
Length of output: 3774
🏁 Script executed:
# Check if GrpcWorker could ever be passed through WorkerRegistry.replace()
# by looking at whether bindings/golang is integrated with the main registry
rg "use.*Worker" bindings/golang/src/policy.rs | head -20Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check the actual location and scope of GrpcWorker usage - is it only in Go FFI?
rg "GrpcWorker" bindings/golang/src/ --type rust | wc -lRepository: lightseekorg/smg
Length of output: 62
🏁 Script executed:
# Verify the trait signature - confirm revision() and inherit_shared_state_from() are truly default implementations
rg -B2 -A5 "fn revision\(\|fn inherit_shared_state_from\(" model_gateway/src/worker/worker.rs | grep -A5 "fn revision\|fn inherit_shared_state_from"Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Verify GrpcWorker is isolated in bindings - check if it's ever used in model_gateway
rg "GrpcWorker" model_gateway/ --type rustRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Confirm bindings/golang is separate by checking if there's any integration between the Go bindings and WorkerRegistry
rg "bindings" model_gateway/src/worker/registry.rs model_gateway/src/worker/manager.rs --type rustRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check if GrpcWorker documents or explains why it doesn't need inherit_shared_state_from()
rg -B10 -A10 "impl Worker for GrpcWorker" bindings/golang/src/policy.rsRepository: lightseekorg/smg
Length of output: 504
🏁 Script executed:
# Verify the exact implementations in GrpcWorker
sed -n '109,250p' bindings/golang/src/policy.rs | head -80Repository: lightseekorg/smg
Length of output: 2041
🏁 Script executed:
# Understand the use case: is GrpcWorker ever created through a discovery/registration flow that could lead to replace()?
rg -B5 -A15 "MultiWorkerClientHandle" bindings/golang/src/policy.rs | head -40Repository: lightseekorg/smg
Length of output: 1750
🏁 Script executed:
# Final verification: confirm GrpcWorker is NEVER registered through any discovery or registry path
# by checking if there's any code path from service discovery to GrpcWorker
rg "service_discovery\|discovery" bindings/golang/ --type rustRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check where replace() is actually called and what types it operates on
rg -B5 "registry.replace" model_gateway/src/ --type rustRepository: lightseekorg/smg
Length of output: 2240
🏁 Script executed:
# Confirm BasicWorker is the only type that flows through WorkerRegistry
rg "BasicWorker" model_gateway/src/worker/registry.rs model_gateway/src/worker/manager.rs model_gateway/src/service_discovery.rs --type rust | head -20Repository: lightseekorg/smg
Length of output: 2000
🏁 Script executed:
# Check if there's any documentation or reasoning for why revision() and inherit_shared_state_from() are defaults
rg -B10 "fn revision\(\)" model_gateway/src/worker/worker.rs | head -20Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Verify all impl Worker blocks in the entire codebase to make absolutely certain
rg "impl.*Worker\s" --type rust | grep -v "impl WorkerType\|impl Worker for\|impl WorkerService\|impl WorkerMetadata\|impl WorkerRuntime\|impl WorkerLoad\|impl WorkerRouting"Repository: lightseekorg/smg
Length of output: 747
🏁 Script executed:
# Check if there are comments about GrpcWorker's purpose or constraints
rg -B5 "FFI workers\|Go SDK\|bindings" bindings/golang/src/policy.rs | head -30Repository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Final check: confirm GrpcWorker is never passed to any function expecting WorkerRegistry behavior
rg "register\|manager\|discovery" bindings/golang/src/ --type rustRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Verify the comment in GrpcWorker confirms it's isolated from the state machine
sed -n '180,250p' bindings/golang/src/policy.rsRepository: lightseekorg/smg
Length of output: 2086
🏁 Script executed:
# Final confirmation: check if there are any other Worker implementations
rg "impl Worker for" --type rustRepository: lightseekorg/smg
Length of output: 183
Make the replacement hooks mandatory for Worker implementations.
With revision() defaulting to 0 and inherit_shared_state_from() defaulting to false, a new impl Worker can still compile without participating in revision fencing or shared-runtime adoption. While the current codebase is safe (only BasicWorker flows through WorkerRegistry::replace()), this design allows future implementations to silently opt out of the replacement protocol. Any new impl Worker that omits these overrides would reset status/load/circuit-breaker state on same-URL replacement, reintroducing the race this PR eliminates.
Suggested direction: Make both methods required trait methods (remove default implementations).
- fn revision(&self) -> u64 {
- 0
- }
+ fn revision(&self) -> u64;
-
- fn inherit_shared_state_from(&self, _other: &dyn Worker) -> bool {
- false
- }
+ fn inherit_shared_state_from(&self, other: &dyn Worker) -> bool;🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/worker/worker.rs` around lines 151 - 167, The Worker trait
currently provides default no-op implementations for revision() and
inherit_shared_state_from(), allowing new impl Worker to opt out of replacement
semantics; remove the default bodies to make revision(&self) -> u64 and
inherit_shared_state_from(&self, _other: &dyn Worker) -> bool required methods
on the Worker trait so every implementor must explicitly handle revision fencing
and shared-state adoption; update any existing impls (e.g., BasicWorker) to
implement these methods accordingly to preserve current behavior.
Addresses two gaps in PR 7 (3c3b52a) that landed the WorkerManager extraction but left the event-driven schedule vulnerable to missing workers in two scenarios. ## Changes model_gateway/src/worker/registry.rs: - `on_remote_worker_state()`: after `register_inner()` inserts a mesh-imported worker, broadcast `WorkerEvent::Registered` on the local event bus so in-process subscribers pick it up. The outgoing mesh sync is still skipped (preserves the existing version-bump loop protection); the event is a purely local broadcast and does not re-enter the mesh. The existing-worker update path remains silent — it goes through `set_healthy()`, and local probes reconcile the status on the next tick. - New test `test_mesh_imported_worker_emits_registered_event`: verifies the new-worker mesh import path emits `Registered`. - New test `test_mesh_worker_state_update_is_silent_on_existing_worker`: regression guard so future changes that accidentally emit events on the health-update path are caught. model_gateway/src/worker/manager.rs: - `WorkerManager::start()`: run `reconcile_from_registry()` synchronously on the caller's thread before spawning the health loop task. Subscribe first, then snapshot, so any registration that lands after the subscribe call is either captured in the synchronous snapshot or delivered as a `Registered` event through the broadcast buffer — no third possibility regardless of async scheduling. - `run_health_loop()`: now takes `next_check` as a parameter instead of building it inline. The internal `reconcile_from_registry()` call is removed; lag-recovery rebuild is the only in-loop reconcile path. - Added explanatory doc comments on both the subscribe-then-snapshot ordering and the run_health_loop contract so future readers don't re-introduce the race. - New test `test_reconcile_from_registry_captures_pending_and_ready_workers`: positive complement to the existing skips-failed test. - New test `test_worker_manager_start_is_deterministic_with_preexisting_workers`: end-to-end start/shutdown with a pre-registered worker. ## Why Before these fixes: 1. Mesh-imported workers registered via `on_remote_worker_state` went through `register_inner()`, which does not broadcast events. Pre-PR 7 this was fine because the health checker loop iterated `self.workers` directly. Post-PR 7 the WorkerManager relies on the event stream, so mesh-imported workers added dynamically never entered the health schedule. With `--remove-unhealthy-workers` enabled, a broken mesh worker could never reach `Failed` locally (set_healthy(false) only transitions Ready → NotReady) and therefore could not be auto-removed. 2. `WorkerManager::start()` previously deferred `reconcile_from_registry` into the spawned task. Whether it ran before or after the mesh replay loop in `server.rs` depended on async task scheduling. Under certain scheduling orders or in a single-threaded runtime the reconcile could run before mesh replay and miss the replayed workers entirely — the inverse symptom of the same underlying event gap. ## How Both fixes preserve the PR 7 invariant that the registry remains a pure collection and the state machine lives in WorkerManager. No business logic moves back into the registry. The event emission in `on_remote_worker_state` uses the existing `event_tx` already owned by the registry, and the synchronous reconcile reuses the existing `reconcile_from_registry` helper. ## Test plan - cargo check -p smg (clean) - cargo clippy -p smg --all-targets -- -D warnings (clean) - cargo test -p smg --lib (531 passed, 4 ignored, 0 failed) - 4 new tests covering both fixes pass Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
♻️ Duplicate comments (2)
model_gateway/src/worker/manager.rs (1)
283-300:⚠️ Potential issue | 🟠 MajorAvoid hot-spinning when probe concurrency is saturated.
If
in_flightis already atMAX_CONCURRENT_HEALTH_PROBESand the earliest deadline innext_checkis overdue,sleep_untilis set to an already-expired instant. Thesleep_untilbranch then wakes immediately and this loop can burn CPU until some probe completes.Proposed fix
- let sleep_until = next_check - .values() - .min() - .copied() - .unwrap_or_else(|| now + Duration::from_secs(config.default_check_interval_secs)); + let has_probe_capacity = in_flight.len() < MAX_CONCURRENT_HEALTH_PROBES; + let sleep_until = next_check + .values() + .filter(|deadline| has_probe_capacity || **deadline > now) + .min() + .copied() + .unwrap_or_else(|| now + Duration::from_secs(config.default_check_interval_secs));🤖 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 283 - 300, The loop can hot-spin when in_flight is at MAX_CONCURRENT_HEALTH_PROBES and sleep_until is already past; fix by clamping sleep_until to a short-future instant when concurrency is saturated. In the select around probes.next() and tokio::time::sleep_until(sleep_until), check if in_flight.len() >= MAX_CONCURRENT_HEALTH_PROBES (or use the existing MAX_CONCURRENT_HEALTH_PROBES constant) and if sleep_until <= now then set sleep_until = now + Duration::from_millis(...) (a small backoff) before entering tokio::select!, so the sleep branch yields a real await instead of immediate wakeups.model_gateway/src/worker/registry.rs (1)
991-1000:⚠️ Potential issue | 🟠 MajorDon't allocate per-worker mutation locks for missing workers.
Both helpers create a
worker_mutation_locksentry before provingself.workersstill containsworker_id. Every stale probe completion or delayed removal job that arrives after removal can therefore leave behind an orphan lock entry, and this refactor explicitly makes those stale async paths normal.Proposed fix
- let lock = self - .worker_mutation_locks - .entry(worker_id.clone()) - .or_insert_with(|| Arc::new(parking_lot::Mutex::new(()))) - .clone(); + let (lock, created_lock) = match self.worker_mutation_locks.entry(worker_id.clone()) { + Entry::Occupied(entry) => (entry.get().clone(), false), + Entry::Vacant(entry) => { + let lock = Arc::new(parking_lot::Mutex::new(())); + entry.insert(lock.clone()); + (lock, true) + } + }; let _guard = lock.lock(); - let worker = self.workers.get(worker_id)?.clone(); + let Some(worker) = self.workers.get(worker_id).map(|entry| entry.clone()) else { + if created_lock { + self.worker_mutation_locks.remove(worker_id); + } + return None; + };Apply the same pattern in both
apply_if_revision()andtransition_status_inner().Also applies to: 1028-1037
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/worker/registry.rs` around lines 991 - 1000, The code currently allocates a per-worker lock via worker_mutation_locks.entry(...) before verifying the worker actually exists, causing orphan locks for missing workers; update apply_if_revision() and transition_status_inner() to first check self.workers.get(worker_id) (clone the Worker or its revision) and validate expected_revision, and only then perform worker_mutation_locks.entry(...). In short: move the self.workers lookup and revision check before creating/locking the Arc<Mutex> so the entry is only created for existing workers (use the existing worker variable from the map lookup to guide the lock acquisition).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Duplicate comments:
In `@model_gateway/src/worker/manager.rs`:
- Around line 283-300: The loop can hot-spin when in_flight is at
MAX_CONCURRENT_HEALTH_PROBES and sleep_until is already past; fix by clamping
sleep_until to a short-future instant when concurrency is saturated. In the
select around probes.next() and tokio::time::sleep_until(sleep_until), check if
in_flight.len() >= MAX_CONCURRENT_HEALTH_PROBES (or use the existing
MAX_CONCURRENT_HEALTH_PROBES constant) and if sleep_until <= now then set
sleep_until = now + Duration::from_millis(...) (a small backoff) before entering
tokio::select!, so the sleep branch yields a real await instead of immediate
wakeups.
In `@model_gateway/src/worker/registry.rs`:
- Around line 991-1000: The code currently allocates a per-worker lock via
worker_mutation_locks.entry(...) before verifying the worker actually exists,
causing orphan locks for missing workers; update apply_if_revision() and
transition_status_inner() to first check self.workers.get(worker_id) (clone the
Worker or its revision) and validate expected_revision, and only then perform
worker_mutation_locks.entry(...). In short: move the self.workers lookup and
revision check before creating/locking the Arc<Mutex> so the entry is only
created for existing workers (use the existing worker variable from the map
lookup to guide the lock acquisition).
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: bae45188-5d58-437f-b3b0-8548438db461
📒 Files selected for processing (2)
model_gateway/src/worker/manager.rsmodel_gateway/src/worker/registry.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 29893d99f2
ℹ️ 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 descriptor.status == WorkerStatus::Failed { | ||
| // Startup reconcile and lagged rebuild must be side-effect-free: | ||
| // do not reschedule already-failed workers for probing or removal. | ||
| next_check.remove(&descriptor.worker_id); | ||
| return; |
There was a problem hiding this comment.
Queue removal for failed workers seen on replace events
schedule_descriptor_at drops every Failed descriptor from next_check, and this helper is also used for WorkerEvent::Replaced. When a same-URL replace preserves runtime status (including Failed), the replacement event removes that worker from scheduling instead of triggering a new removal path; if the prior removal job is discarded as stale by revision fencing, no further probe/removal is ever queued and the failed worker can remain indefinitely in the registry. This breaks remove_unhealthy cleanup for the replace-while-failed case.
Useful? React with 👍 / 👎.
…#1105) Signed-off-by: key4ng <rukeyang@gmail.com>
Description
Problem
WorkerRegistrystill owned the background health-check loop and unhealthy-worker removal path, which mixed collection concerns with lifecycle control. That left same-URL replacement vulnerable to stale probe outcomes and stale removal jobs, and it also reset live runtime state during replace.Solution
Move health-check lifecycle ownership into
WorkerManager, makecheck_health_async()a pure probe, and preserve mutable runtime state across same-URL replacement with an internal sharedWorkerRuntimeplus revision fencing.Changes
WorkerRegistryintoWorkerManagerand wire it throughserver.rsWorkerRegistrya pure collection/mutation layer that emits events and provides revision-gated transition helperscargo clippy --all-targets --all-features -- -D warningsabsolute-path lint in the Go bindingTest Plan
cargo check -p smg --libcargo test -p smg worker::registry::tests -- --nocapturecargo test -p smg worker::manager::tests -- --nocapturecargo test -p smg worker::worker::tests -- --nocapturecargo test -p smg policies::manual::tests -- --nocapturecargo clippy -p smg --lib -- -D warningscargo clippy --all-targets --all-features -- -D warningscargo +nightly fmt --all --checkChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
Refactor
Bug Fixes
Observability