feat(smg): add WorkerSyncAdapter for v2 worker: CRDT namespace - #1313
Conversation
|
Warning You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again! |
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughAdds a new public Changes
Sequence Diagram(s)sequenceDiagram
participant WR as WorkerRegistry
participant WSA as WorkerSyncAdapter
participant CN as CrdtNamespace
participant Remote as RemoteMesh
rect rgba(100, 150, 200, 0.5)
Note over WR,Remote: Outbound: Local → Mesh
WR->>WSA: on_worker_changed(worker_id, state)
WSA->>WSA: bincode::serialize(state)
WSA->>CN: put("worker:{id}", bytes)
CN->>Remote: propagate update
end
rect rgba(150, 100, 200, 0.5)
Note over WR,Remote: Inbound: Mesh → Local
Remote->>CN: mesh update
CN->>WSA: subscription yields (key, Some(value))
WSA->>WSA: strip "worker:" → worker_id
WSA->>WSA: bincode::deserialize(value)
WSA->>WR: on_remote_worker_state(worker_id, state)
end
rect rgba(200, 100, 100, 0.5)
Note over WR,Remote: Removal / Tombstone
WR->>WSA: on_worker_removed(worker_id)
WSA->>CN: delete("worker:{id}")
CN->>Remote: propagate tombstone
CN->>WSA: subscription yields (key, None)
WSA->>WSA: treat as tombstone / ignore
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 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.
Clean, well-structured adapter. Error handling is appropriate (warn on ser/de failures, debug on tombstones), the echo-back idempotence from local put→subscribe is documented and safe, and tests cover the core paths including malformed payloads and wrong-prefix rejection. No issues found.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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/mesh/worker_sync.rs`:
- Around line 188-213: Extract a reusable async polling helper (e.g.,
poll_until) in the test module and replace the repeated for-loop in
start_routes_remote_state_into_registry with a call to that helper; implement
poll_until to accept a predicate (FnMut() -> bool) and a failure message,
perform the same retry/sleep loop (20 iterations, 10ms sleep) and panic with the
message if the predicate never becomes true, then call it as poll_until(||
registry.get_by_url("http://remote:8080").is_some(), "registry did not see the
remote worker").await; keep references to the same symbols
(start_routes_remote_state_into_registry, registry.get_by_url,
WorkerSyncAdapter::new) so the test logic is unchanged.
🪄 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: 41dc53e5-a929-4d78-a7d7-33dfc4aad57c
📒 Files selected for processing (3)
model_gateway/src/lib.rsmodel_gateway/src/mesh/mod.rsmodel_gateway/src/mesh/worker_sync.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 6fc23b10da
ℹ️ 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".
|
Hi @CatherineSue, this PR has merge conflicts that must be resolved before it can be merged. Please rebase your branch: git fetch origin main
git rebase origin/main
# resolve any conflicts, then:
git push --force-with-lease |
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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/mesh/adapters/worker_sync.rs`:
- Around line 73-76: The start() method currently only calls
self.workers.subscribe("") so it misses pre-existing CRDT entries; before
subscribing, read the current worker set from the CRDT and insert each into the
WorkerRegistry (or whatever registry field is used) to "backfill" existing
worker: entries, then proceed to call self.workers.subscribe("") to receive live
updates; locate start(), the workers subscribe call, and the registry insertion
logic (e.g., WorkerRegistry or self.registry, and functions that add workers)
and invoke the CRDT read API (e.g., get_all()/iter_entries()/snapshot) to
iterate existing entries and register them prior to subscribing.
🪄 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: ab2cde1b-5326-4dfb-bc98-6908e6366422
📒 Files selected for processing (3)
model_gateway/src/mesh/adapters/mod.rsmodel_gateway/src/mesh/adapters/worker_sync.rsmodel_gateway/src/mesh/mod.rs
Bridges the `worker:` CRDT namespace (MeshKV) and the in-process WorkerRegistry. Outbound APIs (`on_worker_changed`, `on_worker_removed`) bincode-serialise WorkerState and write to the namespace. An inbound task spawned by `start` routes remote updates through WorkerRegistry::on_remote_worker_state — the same sink the existing WorkerStateSubscriber implementation feeds — so URL deduplication, health promotion, and the Registered event fan-out stay unchanged. No server.rs wiring yet. Step 5 (dual-mode deployment) will wire the adapter alongside teardown of the v1 MeshSyncManager worker sync path. Third in the Step 4 adapter sequence; builds on the auto-registered config: prefix merged in the previous PR. Signed-off-by: Chang Su <chang.s.su@oracle.com>
Makes room for non-adapter mesh glue (bootstrap, shutdown, config helpers) to land alongside without mixing with the per-namespace adapters. Pure rename + module re-wire; no behaviour change. Signed-off-by: Chang Su <chang.s.su@oracle.com>
2f9d41e to
4e7e665
Compare
`start` previously only processed live subscription events, so a node that joined after workers had already gossiped would wait for the next unrelated write before its registry was populated — the v1 path explicitly replayed via `get_all_worker_states()` for exactly this reason. Subscribe first (so no replay/live race loses an event), then iterate the namespace and feed existing entries through the same `on_remote_worker_state` sink the live loop uses. Add a test that pre-seeds the namespace before `start` and asserts the registry sees the worker. Signed-off-by: Chang Su <chang.s.su@oracle.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 5d300b02f0
ℹ️ 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".
| this.backfill_existing(); | ||
| while let Some((key, value)) = sub.receiver.recv().await { |
There was a problem hiding this comment.
Drain subscription queue while backfilling
start subscribes and then runs backfill_existing() before the first recv, so the adapter does not consume subscription events during the whole replay window. In this codebase, subscription delivery is bounded and drop-on-full (CrdtNamespace::subscribe + SubscriberRegistry::notify in crates/mesh/src/kv.rs use an mpsc channel with try_send), so a busy cluster can overflow the queue while backfill is running; dropped worker updates are never replayed, leaving WorkerRegistry stale until a later write touches the same keys. This is most visible on startup when many workers are present and updates continue concurrently.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Fixed in cbce753. backfill_existing now runs on the caller's thread, outside the spawned task, so the recv loop is free to drain the channel concurrently on a multi-threaded runtime. The sink is idempotent on URL, so a key seen by both backfill and a live event just refreshes health. In practice the worker: channel capacity (1000) and the low update rate for worker state make the overflow very unlikely, but the reorder is free and tidier.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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/mesh/adapters/worker_sync.rs`:
- Around line 117-119: The backfill loop silently skips malformed keys when
doing let Some(worker_id) = key.strip_prefix(PREFIX).filter(|s| !s.is_empty())
else { continue; }; — change this to log a warning when a key does not start
with PREFIX or yields an empty suffix (i.e., when strip_prefix returns None or
the filtered empty string) so malformed keys are observable at startup; include
the offending key and mention PREFIX in the warning and still continue
processing subsequent keys.
🪄 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: 1321ad4a-e32e-4e6c-b5ff-ceb6707002c4
📒 Files selected for processing (1)
model_gateway/src/mesh/adapters/worker_sync.rs
Move `backfill_existing` out of the spawned recv task so the live loop starts draining immediately. `notify` uses `try_send` into a bounded mpsc (worker: capacity 1000), and a slow backfill inside the same task blocks the recv — on a busy startup that window could drop live updates the gossip path would not resend. Also mirror the live-path warning for malformed keys in `backfill_existing` so startup data-shape issues are observable. Signed-off-by: Chang Su <chang.s.su@oracle.com>
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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/mesh/adapters/worker_sync.rs`:
- Around line 130-143: Reject payloads whose embedded WorkerState.worker_id
disagrees with the authoritative key: in apply_incoming, after
bincode::deserialize::<WorkerState>(bytes) check that state.worker_id ==
worker_id and if not log a warning (include both ids) and return without calling
self.worker_registry.on_remote_worker_state(&state); likewise in
on_worker_changed verify state.worker_id == worker_id before bincode::serialize
and self.workers.put(&format!("{PREFIX}{worker_id}"), bytes), and if they differ
log a warning and skip publishing. Use the existing warn! call pattern and
reference apply_incoming, on_worker_changed, WorkerState, worker_registry,
workers.put and PREFIX to find the places to add these checks.
🪄 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: 149249af-f64c-4a0d-b829-d2e53de2f671
📒 Files selected for processing (1)
model_gateway/src/mesh/adapters/worker_sync.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: cbce753eb7
ℹ️ 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".
| } | ||
| debug!("WorkerSyncAdapter subscription closed"); | ||
| }); | ||
| self.backfill_existing(); |
There was a problem hiding this comment.
Yield before synchronous backfill to drain subscription queue
Fresh evidence for the earlier queue-drop concern: even though start() now calls tokio::spawn first, it immediately runs backfill_existing() synchronously, so on a single-thread Tokio runtime (including current_thread and 1-worker deployments) the spawned recv loop cannot poll until backfill returns. I checked crates/mesh/src/kv.rs, where subscriber delivery is bounded and drop-on-full via try_send, so a busy cluster can still lose worker updates during this window and leave WorkerRegistry stale until another write arrives.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Declining this one. Production runs on the multi-thread tokio runtime (default for #[tokio::main]), where the reorder in cbce753 already lets the recv loop drain on another worker thread while backfill runs, so the overflow window is closed.
On current_thread (1-worker test/dev runtimes), an extra tokio::spawn wrapper for backfill wouldn't actually help — backfill_existing is fully synchronous, so once the scheduler polls it the sync loop runs to completion without yielding regardless of whether it lives in its own task. A real fix for single-thread would require making backfill_existing async and sprinkling tokio::task::yield_now().await between iterations. Given worker state isn't a hot path (register / unregister / health toggle) and the channel capacity is 1000, the combination "single-thread runtime + sustained >200k worker-state writes/sec" needed to overflow isn't a realistic production scenario.
Happy to revisit if we ever ship on current_thread under real load.
Description
Problem
Mesh v2 routes worker state through the typed
worker:CRDT namespace onMeshKV, but the gateway has no adapter that writes to / reads from that namespace yet. Step 5's dual-mode deployment needs this bridge in place before the v1MeshSyncManagerworker-sync path can be torn down.Solution
Add
WorkerSyncAdapterundermodel_gateway/src/mesh/worker_sync.rs:on_worker_changed(worker_id, &state)/on_worker_removed(worker_id)— bincode-serialiseWorkerStateand callCrdtNamespace::put/deleteunderworker:{worker_id}.start(self: &Arc<Self>)— spawns a task that subscribes to theworker:prefix and routes non-tombstone events through the existingWorkerRegistry::on_remote_worker_statesink (the same entry point the currentWorkerStateSubscribertrait feeds), so URL-dedupe, health promotion, and theRegisteredevent fan-out are unchanged. Tombstones are logged at debug for now — the registry has no remote-remove hook yet, and wiring one belongs alongside the v1-mirror teardown in a later step.newasserts the namespace prefix isworker:so a mis-wired caller fails loudly at startup.No
server.rswiring in this PR. That comes with Step 5 dual-mode deployment together with shutting down the v1 path.Third in the Step 4 adapter sequence; builds on the auto-registered
config:prefix merged in the previous PR.Changes
model_gateway/src/mesh/module withWorkerSyncAdapter.pub mod meshfrommodel_gateway/src/lib.rs.WorkerRegistry, graceful handling of malformed payloads, and the prefix assertion.Test Plan
cargo test -p smg --lib mesh::worker_sync— 5 tests passing.cargo clippy -p smg --lib --tests -- -D warningscargo +nightly fmt --allChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Tests