From c04cac3da95ff77fe55d5bb0a1571e5a9f722831 Mon Sep 17 00:00:00 2001 From: Simo Lin Date: Tue, 14 Apr 2026 07:44:42 -0700 Subject: [PATCH] refactor(worker): canonical runtime API on WorkerRuntime MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Part B of the Worker trait trim outlined in .claude/plans/2026-04-14-worker-module-followup-cleanup.md (Finding 4, categories B + C). Follows the same Option 2 principle we landed on in Part A (PR #1127): canonical implementation on the underlying struct, trait surface preserved, callers unchanged, only delete truly dead methods. What changed - model_gateway/src/worker/worker.rs * Delete dead `Worker::reset_load` (grep confirmed 0 callers; trait had a default no-op body, BasicWorker had the only override). * Promote `WorkerRuntime::{status,set_status,revision,bump_revision}` from private to `pub fn` so external callers reach the runtime through a single accessor instead of being funneled through the trait each time. * Add canonical `pub fn` counter methods on `impl WorkerRuntime` as the single source of truth for runtime state: consecutive failures, consecutive successes, total pending probes, load counter, routing- key load, and processed-request counter. * Introduce `WorkerRuntime::try_decrement_load` which performs the saturating `fetch_update` and returns `bool`; the caller warns on underflow. Keeps the BasicWorker-specific side effects (tracing warn + metric update) at the forwarding layer. * Rewrite the 15 counter methods on `impl Worker for BasicWorker` so each one is a one-line forward to the corresponding `WorkerRuntime` method. `increment_load` / `decrement_load` still call `update_running_requests_metrics()` after the delegation so the metric export path is unchanged. `set_status` still updates `Metrics::set_worker_health` after the delegation. Why - Finding 4 Part A trimmed the metadata delegates. Part B does the equivalent consolidation for runtime state. Before this PR, every counter mutation went through a ~6-line block that reached into the WorkerRuntime's internal atomic fields; those fields had to be `pub` for the trait impl to touch them, and each piece of atomic bookkeeping (ordering + `+1` for post-increment) was duplicated at every call site in BasicWorker. - Collapsing the implementation onto WorkerRuntime lets us keep precisely one copy of each atomic ordering choice, and lets future `Worker` implementations (e.g. the GrpcWorker in bindings/golang that currently carries its own raw atomics) adopt the canonical runtime wholesale without duplicating the memory-ordering rules. - `reset_load` had no callers anywhere in the tree — keeping it in the trait just added noise. Deleting dead methods keeps the trait definition honest. How - Option 2 principle: trait methods stay (so callers still write `worker.load()`, `worker.increment_processed()`, etc. — no churn across routers/ and workflow/). The impl on BasicWorker becomes one-line forwards via `self.runtime.load().method()`. The canonical code lives in exactly one place: `impl WorkerRuntime`. - Category C (circuit breaker) is already Option 2-compliant — the six BasicWorker wrappers (`circuit_breaker_state`, `circuit_breaker_can_execute`, `record_circuit_breaker_outcome`, plus the `is_available` / `record_outcome` default impls and the `resilience` accessor) already forward to `CircuitBreaker::state()`, `.can_execute()`, `.record_outcome()` which own the canonical implementation in worker/circuit_breaker.rs. No changes needed. Verification - `cargo check -p smg` — clean. - `cargo clippy -p smg --all-targets -- -D warnings` — clean. - `cargo clippy -p smg-golang --all-targets -- -D warnings` — clean (GrpcWorker in bindings/golang still compiles against the trimmed trait surface). - `cargo test -p smg --lib` — 538 passed; 0 failed; 4 ignored. Refs: worker module deep refactor follow-up cleanup (Finding 4 Part B) Signed-off-by: Simo Lin --- model_gateway/src/worker/worker.rs | 168 ++++++++++++++++------------- 1 file changed, 95 insertions(+), 73 deletions(-) diff --git a/model_gateway/src/worker/worker.rs b/model_gateway/src/worker/worker.rs index fbbcdef623..4c88964daa 100644 --- a/model_gateway/src/worker/worker.rs +++ b/model_gateway/src/worker/worker.rs @@ -192,9 +192,6 @@ pub trait Worker: Send + Sync + fmt::Debug + 'static { /// Decrement the load counter fn decrement_load(&self); - /// Reset the load counter to 0 (for sync/recovery) - fn reset_load(&self) {} - /// Get the current routing-key load cardinality. fn routing_key_load(&self) -> usize; @@ -588,21 +585,97 @@ impl WorkerRuntime { } } - fn status(&self) -> WorkerStatus { + // ── Lifecycle status ──────────────────────────────────────────── + + pub fn status(&self) -> WorkerStatus { WorkerStatus::from_u8(self.status.load(Ordering::Acquire)) } - fn set_status(&self, status: WorkerStatus) { + pub fn set_status(&self, status: WorkerStatus) { self.status.store(status as u8, Ordering::Release); } - fn revision(&self) -> u64 { + pub fn revision(&self) -> u64 { self.revision.load(Ordering::Acquire) } - fn bump_revision(&self) -> u64 { + pub fn bump_revision(&self) -> u64 { self.revision.fetch_add(1, Ordering::AcqRel) + 1 } + + // ── Health-check counters ─────────────────────────────────────── + + pub fn consecutive_failures_increment(&self) -> usize { + self.consecutive_failures.fetch_add(1, Ordering::AcqRel) + 1 + } + + pub fn consecutive_failures_reset(&self) { + self.consecutive_failures.store(0, Ordering::Release); + } + + pub fn consecutive_successes_increment(&self) -> usize { + self.consecutive_successes.fetch_add(1, Ordering::AcqRel) + 1 + } + + pub fn consecutive_successes_reset(&self) { + self.consecutive_successes.store(0, Ordering::Release); + } + + pub fn total_pending_probes(&self) -> usize { + self.total_pending_probes.load(Ordering::Relaxed) + } + + pub fn total_pending_probes_increment(&self) -> usize { + self.total_pending_probes.fetch_add(1, Ordering::Relaxed) + 1 + } + + pub fn total_pending_probes_reset(&self) { + self.total_pending_probes.store(0, Ordering::Relaxed); + } + + // ── Load counter ──────────────────────────────────────────────── + + pub fn load(&self) -> usize { + self.load_counter.load(Ordering::Relaxed) + } + + pub fn increment_load(&self) { + self.load_counter.fetch_add(1, Ordering::Relaxed); + } + + /// Saturating decrement. Returns `true` if the counter was decremented, + /// `false` if it was already zero — callers can log when that happens. + pub fn try_decrement_load(&self) -> bool { + self.load_counter + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| { + current.checked_sub(1) + }) + .is_ok() + } + + // ── Routing-key load ──────────────────────────────────────────── + + pub fn routing_key_load(&self) -> usize { + self.worker_routing_key_load.value() + } + + pub fn increment_routing_key_load(&self, routing_key: &str) { + self.worker_routing_key_load.increment(routing_key); + } + + pub fn decrement_routing_key_load(&self, routing_key: &str) { + self.worker_routing_key_load.decrement(routing_key); + } + + // ── Processed-request counter ─────────────────────────────────── + + pub fn processed_requests(&self) -> usize { + self.processed_counter.load(Ordering::Relaxed) + } + + pub fn increment_processed(&self) { + self.processed_counter.fetch_add(1, Ordering::Relaxed); + } } /// Basic worker implementation @@ -739,78 +812,44 @@ impl Worker for BasicWorker { } fn consecutive_failures_increment(&self) -> usize { - self.runtime - .load() - .consecutive_failures - .fetch_add(1, Ordering::AcqRel) - + 1 + self.runtime.load().consecutive_failures_increment() } fn consecutive_failures_reset(&self) { - self.runtime - .load() - .consecutive_failures - .store(0, Ordering::Release); + self.runtime.load().consecutive_failures_reset(); } fn consecutive_successes_increment(&self) -> usize { - self.runtime - .load() - .consecutive_successes - .fetch_add(1, Ordering::AcqRel) - + 1 + self.runtime.load().consecutive_successes_increment() } fn consecutive_successes_reset(&self) { - self.runtime - .load() - .consecutive_successes - .store(0, Ordering::Release); + self.runtime.load().consecutive_successes_reset(); } fn total_pending_probes(&self) -> usize { - self.runtime - .load() - .total_pending_probes - .load(Ordering::Relaxed) + self.runtime.load().total_pending_probes() } fn total_pending_probes_increment(&self) -> usize { - self.runtime - .load() - .total_pending_probes - .fetch_add(1, Ordering::Relaxed) - + 1 + self.runtime.load().total_pending_probes_increment() } fn total_pending_probes_reset(&self) { - self.runtime - .load() - .total_pending_probes - .store(0, Ordering::Relaxed); + self.runtime.load().total_pending_probes_reset(); } fn load(&self) -> usize { - self.runtime.load().load_counter.load(Ordering::Relaxed) + self.runtime.load().load() } fn increment_load(&self) { - self.runtime - .load() - .load_counter - .fetch_add(1, Ordering::Relaxed); + self.runtime.load().increment_load(); self.update_running_requests_metrics(); } fn decrement_load(&self) { - let runtime = self.runtime.load(); - if runtime - .load_counter - .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| { - current.checked_sub(1) - }) - .is_err() - { + if !self.runtime.load().try_decrement_load() { tracing::warn!( worker_url = %self.metadata.spec.url, "Attempted to decrement load counter that is already at 0" @@ -819,41 +858,24 @@ impl Worker for BasicWorker { self.update_running_requests_metrics(); } - fn reset_load(&self) { - self.runtime.load().load_counter.store(0, Ordering::Relaxed); - self.update_running_requests_metrics(); - } - fn routing_key_load(&self) -> usize { - self.runtime.load().worker_routing_key_load.value() + self.runtime.load().routing_key_load() } fn increment_routing_key_load(&self, routing_key: &str) { - self.runtime - .load() - .worker_routing_key_load - .increment(routing_key); + self.runtime.load().increment_routing_key_load(routing_key); } fn decrement_routing_key_load(&self, routing_key: &str) { - self.runtime - .load() - .worker_routing_key_load - .decrement(routing_key); + self.runtime.load().decrement_routing_key_load(routing_key); } fn processed_requests(&self) -> usize { - self.runtime - .load() - .processed_counter - .load(Ordering::Relaxed) + self.runtime.load().processed_requests() } fn increment_processed(&self) { - self.runtime - .load() - .processed_counter - .fetch_add(1, Ordering::Relaxed); + self.runtime.load().increment_processed(); } fn metadata(&self) -> &WorkerMetadata {