perf(smg): mitigate the worker-sync lag-resync burst - #1667
Conversation
The 64-slot broadcast lagged under bulk startup registration and probe storms; every lag forces a full mesh resync. 1024 slots covers realistic worker counts at ~100 KB fixed cost. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
resync_local re-put every local worker on broadcast lag; each re-put mints a fresh Lamport op that invalidates every peer's send watermark for the key, bursting W ops to N peers even when nothing changed. The store is now consulted first: volatile load is ignored and specs compare as JSON values (WorkerSpec holds maps, so equal specs can differ byte-wise). Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Changed workers now skip the two JSON parses, and the WorkerState clone is gone — field-wise comparison replaces the zero-and-compare. Signed-off-by: Chang Su <8605658+CatherineSue@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 (1)
📝 WalkthroughWalkthroughWorkerSyncAdapter's local resync now avoids republishing worker state when the mesh store already holds an equivalent ChangesWorker State Management
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~20 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)
Comment |
There was a problem hiding this comment.
Code Review
This pull request optimizes the local worker resync process by skipping republishing when the stored state is already equivalent, and increases the worker registry's event broadcast channel capacity to 1024 to handle fleet-scale bursts. The review feedback suggests two key performance improvements: deferring the expensive worker_state_of serialization until after checking if the store matches, and refactoring store_matches to compare the worker's version/revision directly instead of performing expensive JSON deserialization and comparison of the specs.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| let state = worker_state_of(&id, &worker); | ||
| if !self.store_matches(&id, &state) { | ||
| self.on_worker_changed(id.as_str(), &state); | ||
| } |
There was a problem hiding this comment.
Defer the expensive worker_state_of call (which serializes the WorkerSpec to JSON) until after we have determined that the store does not match. This avoids unnecessary JSON serialization for all unchanged workers during a resync.
| let state = worker_state_of(&id, &worker); | |
| if !self.store_matches(&id, &state) { | |
| self.on_worker_changed(id.as_str(), &state); | |
| } | |
| if !self.store_matches(&id, &worker) { | |
| let state = worker_state_of(&id, &worker); | |
| self.on_worker_changed(id.as_str(), &state); | |
| } |
References
- Avoid heap allocations (like
Vec<u8>) within performance-critical comparison logic, especially during bulk processing like log compaction. Use allocation-free helpers or compare decoded values directly to reduce memory pressure.
There was a problem hiding this comment.
Deferring worker_state_of only pays if store_matches can decide "unchanged" without the spec bytes — i.e. under the version-only comparison in the sibling thread, which is unsound across restarts (revision() is an in-process counter that resets per boot). With semantic spec comparison kept, the common case (everything matches) is exactly the case that must compare specs, so building the state up front isn't wasted work in the steady state.
The cost worth cutting was the two JSON parses per matched worker: 857b747 short-circuits on spec byte equality, which holds in-process since re-serialising the same spec instance is byte-stable. The parse now only runs when bytes drift while scalars match (effectively once, after a restart).
| fn store_matches(&self, id: &WorkerId, state: &WorkerState) -> bool { | ||
| let Some(bytes) = self.workers.get(&format!("{PREFIX}{}", id.as_str())) else { | ||
| return false; | ||
| }; | ||
| let Ok(stored) = bincode::deserialize::<WorkerState>(&bytes) else { | ||
| return false; | ||
| }; | ||
| if stored.worker_id != state.worker_id | ||
| || stored.model_id != state.model_id | ||
| || stored.url != state.url | ||
| || stored.health != state.health | ||
| || stored.version != state.version | ||
| { | ||
| return false; | ||
| } | ||
| match ( | ||
| serde_json::from_slice::<serde_json::Value>(&stored.spec), | ||
| serde_json::from_slice::<serde_json::Value>(&state.spec), | ||
| ) { | ||
| (Ok(a), Ok(b)) => a == b, | ||
| _ => stored.spec == state.spec, | ||
| } | ||
| } |
There was a problem hiding this comment.
Since WorkerSpec is immutable for a given BasicWorker instance, and any update to a worker's specification (via replace or register_or_replace) always increments its monotonic version (revision) counter, comparing the version field is entirely sufficient to guarantee that the underlying spec has not changed.
By changing store_matches to take &Arc<dyn Worker> instead of &WorkerState, we can perform all necessary comparisons using cheap scalar fields and completely eliminate the expensive, CPU-bound JSON deserialization and semantic comparison of stored.spec and state.spec.
fn store_matches(&self, id: &WorkerId, worker: &Arc<dyn Worker>) -> bool {
let Some(bytes) = self.workers.get(&format!("{PREFIX}{}", id.as_str())) else {
return false;
};
let Ok(stored) = bincode::deserialize::<WorkerState>(&bytes) else {
return false;
};
stored.worker_id == id.as_str()
&& stored.model_id == worker.model_id()
&& stored.url == worker.url()
&& stored.health == worker.is_healthy()
&& stored.version == worker.revision()
}References
- Avoid heap allocations (like
Vec<u8>) within performance-critical comparison logic, especially during bulk processing like log compaction. Use allocation-free helpers or compare decoded values directly to reduce memory pressure.
There was a problem hiding this comment.
The version-only comparison is unsound across restarts: revision() is an in-process atomic that resets to 0 when the gateway restarts, and with deterministic ids (#1668) a restarted node re-registers the same worker_id. If the spec changed across the restart (config edit) while the republished revision number matches what the store holds from the prior incarnation, version-equality would false-match and suppress the re-assert — exactly the prior-incarnation shadow the reconcile pass exists to fix. So the semantic spec comparison stays as ground truth.
Adopted the spirit of this in 857b747: a byte-equality fast path now gates the JSON parses. Re-serialising the same in-process spec instance is byte-stable, so steady-state resync/reconcile short-circuits on memcmp and only falls through to the JSON comparison on key-order drift (i.e. across restarts, where the semantic check is the point).
There was a problem hiding this comment.
Clean, well-contained perf fix. store_matches correctly skips load (volatile/unread), uses JSON value comparison for spec to handle non-deterministic map ordering, and gates the expensive JSON parses behind cheap scalar checks. Deserialization failures conservatively fall through to re-put (the old behavior). Channel capacity bump (64→1024, ~100 KB fixed) is well-justified. Test covers both directions. No issues found.
0 🔴 Important · 0 🟡 Nit · 0 🟣 Pre-existing
…re_matches Re-serialising the same in-process spec instance is byte-stable, so the steady-state resync/reconcile comparison short-circuits on bytes and only falls through to the semantic JSON comparison on key-order drift (e.g. across restarts). Suggested by review on #1667. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Description
Problem
Two related burst amplifiers in the worker mesh sync, found by the #1661 perf review: (1) the registry's
WorkerEventbroadcast had 64 slots — bulk startup registration or a probe storm lags it, and every lag forces a full mesh resync; (2) that resync (resync_local) re-put every local worker unconditionally. Each re-put mints a fresh Lamport op, which invalidates every peer's per-key send watermark — so one lag event burst W ops to N peers (~16× amplification at 1,000 workers) even when nothing had changed, and a sustained status storm could re-trigger it repeatedly.Solution
store_matchesconsults the store before re-putting. Cheap scalar fields gate the comparison;loadis not compared (volatile, unread by importers); specs compare as JSON values rather than bytes, sinceWorkerSpecholds maps and two encodings of an identical spec can differ byte-wise. A skipped put is a watermark delta that never ships.Correctness containment: a lost publish or foreign tombstone leaves the store without (or with different) state, so
store_matchesreturns false and the resync still re-asserts — the lag-recovery role is preserved. False negatives merely re-put (the old behavior).Changes
model_gateway/src/worker/registry.rs— event channel capacity + sizing rationalemodel_gateway/src/mesh/adapters/worker_sync.rs—store_matches+ resync gating; test via the local-subscriber observable (put ⇒ notify, skip ⇒ empty channel)Test Plan
cargo test -p smg— suite green; newresync_skips_republish_when_state_unchangedasserts both directions (unchanged ⇒ skipped; health flip ⇒ republished).cargo clippy --all-targets --all-features -- -D warnings— clean.Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspasses🤖 Generated with Claude Code
Summary by CodeRabbit
Bug Fixes
Performance Improvements
Tests