Repository navigation
Conversation
|
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:
📝 WalkthroughWalkthroughSwitches mesh tree sync from legacy TreeState to kv_index::snapshot::TreeSnapshot across mesh and model-gateway: adds "tree_snapshot" payload/type, snapshot-based materialization, delta apply, checkpointing, compatibility shims for legacy "tree_state", and API re-exports. Changes
Sequence Diagram(s)sequenceDiagram
participant Client
participant Controller
participant SyncManager
participant Stores
participant PolicyRegistry
Client->>Controller: send PolicyState (policy_type="tree_snapshot", config bytes)
Controller->>SyncManager: deserialize to TreeSnapshot
Controller->>SyncManager: apply_remote_tree_snapshot(model_id, snapshot, version, actor)
SyncManager->>Stores: store snapshot bytes in `tree_configs` and advance `tree_version`
SyncManager->>PolicyRegistry: notify subscribers with &TreeSnapshot
PolicyRegistry->>PolicyRegistry: merge snapshot into local trees (cache-aware)
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
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 unit tests (beta)
Comment |
There was a problem hiding this comment.
Code Review
This pull request transitions the mesh synchronization logic from a flat list of operations to a compact radix tree snapshot format. Key changes include updating the synchronization manager, incremental update collector, and gossip service to support the new TreeSnapshot type while providing backward compatibility for the legacy TreeState format. The review feedback highlights several areas for improvement: addressing potential race conditions by capturing versions during state materialization, reducing logic duplication between the incremental collector and sync manager, and optimizing performance by avoiding redundant serialization of tree snapshots. Additionally, it is recommended to log deserialization failures as warnings to improve system observability.
| let serialized = match snapshot.to_bytes() { | ||
| Ok(bytes) => bytes, | ||
| Err(err) => { | ||
| debug!(error = %err, model_id = %model_id, "Failed to serialize remote tree state"); | ||
| debug!(error = %err, model_id = %model_id, "Failed to serialize remote tree snapshot"); | ||
| return; | ||
| } | ||
| }; |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 00624ec329
ℹ️ 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 tree = self | ||
| .string_trees | ||
| .entry(model_id.to_string()) | ||
| .or_insert_with(|| Arc::new(Tree::new())); | ||
| tree.merge_snapshot(snapshot); |
There was a problem hiding this comment.
Preserve token ops in remote snapshot apply path
apply_remote_tree_snapshot now merges only into string_trees, so remote mesh updates never populate token_trees. In this commit, gossip routes tree updates through the snapshot subscriber path, which means TreeKey::Tokens inserts from other nodes are effectively dropped; gRPC/token-based cache-aware routing will miss cross-node cache affinity and behave as cold on every node.
Useful? React with 👍 / 👎.
| let current_version = self.stores.tree_version(&key); | ||
| if version > current_version { | ||
| entry.insert(serialized); |
There was a problem hiding this comment.
Make snapshot version check atomic with config writes
This compares version against tree_versions and writes tree_configs based on that check, but tree_versions is only advanced later (outside the entry lock). With concurrent remote full-state updates, a lower-version snapshot can overwrite a higher-version one if it runs between the higher write and the later advance_tree_version, leaving stale snapshot bytes paired with a high version counter and causing state regression until another sync fixes it.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 7
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@crates/mesh/src/incremental.rs`:
- Around line 291-325: Add a short clarifying comment inside the full-state
fallback block where pending ops are replayed (the loop matching
super::tree_ops::TreeOperation::Insert and matching
super::tree_ops::TreeKey::Text) stating that TreeKey::Tokens inserts are
intentionally skipped here because token-tree snapshot format is not yet
supported and token synchronization is deferred to the later phase described in
sync.rs; mention that the version is still advanced and this omission is
intentional to avoid being mistaken for a bug or data loss.
In `@crates/mesh/src/ping_server.rs`:
- Around line 145-150: The no-sync_manager fallback currently writes raw bytes
into tree_configs without updating tree_versions/tree_generation or performing
version de-dup, causing stale-version re-gossip; modify the fallback at the tree
snapshot handling paths (where stores.tree_version(&key) is read and where raw
writes occur into tree_configs) to call the same versioned write/bookkeeping
routine used by MeshSyncManager::apply_remote_tree_snapshot() (or extract that
logic into a shared helper) so that tree_versions, tree_generation, and dedup
checks are updated atomically when storing the snapshot.
In `@crates/mesh/src/sync.rs`:
- Around line 486-501: The current materialize_tree_snapshot logic treats a
missing tree_configs entry as an empty tree which is incorrect when there are
pending ops; update the branch handling None so it first checks
tree_version()==0 and whether self.stores.tree_ops_pending contains_key(key): if
version==0 and no pending ops, return an empty kv_index::Tree as before,
otherwise materialize the pending-only tree (apply operations from
self.stores.tree_ops_pending for key to produce a kv_index::Tree) and return
that snapshot/size; apply the same guard/fix to the other vacant-case usages
referenced (the code paths around the logic at the other locations involving
tree_configs absence) so remote snapshots/deltas are compared/merged against the
materialized pending-only local tree instead of being treated as empty.
- Around line 718-735: In apply_remote_tree_state_compat, do not silently ignore
unsupported TreeOperation or TreeKey variants (e.g., TreeOperation::Remove or
TreeKey::Tokens); validate the entire tree_state.operations first and if any op
is unsupported, return an error or early-fail (or at minimum log and skip
advancing the version) instead of building a partial tree and calling
apply_remote_tree_snapshot; specifically update the loop in
apply_remote_tree_state_compat (and the analogous loop around delta handling) to
detect non-Insert(Text) or non-Text keys, propagate a failure to the caller (or
avoid calling apply_remote_tree_snapshot/applying delta.new_version) so the
remote peer can retry or the conflict can be handled safely.
- Around line 477-525: The checkpoint currently computes pending_count =
pending.len() then only materializes Insert(Text) ops, causing unhandled ops
(Remove and TreeKey::Tokens) to be drained and lost; fix by making pending_count
reflect only the ops actually serialized into the snapshot (increment
pending_count inside the loop when handling TreeOperation::Insert with
TreeKey::Text) instead of using pending.len(), and update the other identical
block (the similar materialization logic around the 877-913 region) to use the
same approach so only applied/serialized ops are counted and subsequently
drained.
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 982-985: The test currently asserts mesh_sync snapshot node_count
but not that the policy was actually restored; change the assertion to validate
the restored policy behavior by calling the policy's
restore_tree_state_from_mesh() (or invoking set_mesh_sync() then
restore_tree_state_from_mesh()) and then asserting against the policy tree or a
selection outcome (e.g., call select_worker(...) or inspect PolicyTree methods)
to ensure the restored policy contains the expected entries (shared prefix
"test_text_" and two leaf suffixes) rather than just checking
mesh_sync.get_tree_snapshot(); locate functions restore_tree_state_from_mesh,
set_mesh_sync, select_worker, and the PolicyTree/Policy struct in this file to
update the test accordingly.
- Around line 337-341: The DashMap entry guard for string_trees is held while
calling Tree::merge_snapshot, blocking other concurrent operations; fix by using
string_trees.entry(model_id.to_string()).or_insert_with(||
Arc::new(Tree::new())) to obtain the Arc<Tree>, clone that Arc (e.g., let
tree_arc = entry_ref.clone()), drop the map guard immediately, and then call
tree_arc.merge_snapshot(snapshot) so merge_snapshot runs without holding the
DashMap guard; reference symbols: string_trees, model_id, Arc<Tree>,
merge_snapshot, entry, or_insert_with.
🪄 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: a0f39856-778a-412c-a90f-544f8395c6fe
📒 Files selected for processing (8)
crates/mesh/src/controller.rscrates/mesh/src/incremental.rscrates/mesh/src/lib.rscrates/mesh/src/ping_server.rscrates/mesh/src/sync.rsmodel_gateway/src/policies/cache_aware.rsmodel_gateway/src/policies/registry.rsmodel_gateway/tests/mesh_integration_test.rs
| pub fn apply_remote_tree_state_compat( | ||
| &self, | ||
| model_id: String, | ||
| tree_state: TreeState, | ||
| actor: Option<String>, | ||
| ) { | ||
| let version = tree_state.version; | ||
| // Build a local Tree from the operations, then snapshot it | ||
| let tree = kv_index::Tree::new(); | ||
| for op in &tree_state.operations { | ||
| if let TreeOperation::Insert(insert_op) = op { | ||
| if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key { | ||
| tree.insert(text, &insert_op.tenant); | ||
| } | ||
| } | ||
| } | ||
| let snapshot = tree.snapshot(); | ||
| self.apply_remote_tree_snapshot(model_id, snapshot, version, actor); |
There was a problem hiding this comment.
Compat and delta paths can't silently acknowledge unsupported tree ops.
Lines 727-732 and 757-763 only replay Insert(Text), yet the caller still advances to tree_state.version/delta.new_version. A peer can therefore send TreeOperation::Remove or TreeKey::Tokens, get a successful apply, and still diverge permanently from this node because the skipped ops will never be retried once the version moves forward.
Also applies to: 755-765
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@crates/mesh/src/sync.rs` around lines 718 - 735, In
apply_remote_tree_state_compat, do not silently ignore unsupported TreeOperation
or TreeKey variants (e.g., TreeOperation::Remove or TreeKey::Tokens); validate
the entire tree_state.operations first and if any op is unsupported, return an
error or early-fail (or at minimum log and skip advancing the version) instead
of building a partial tree and calling apply_remote_tree_snapshot; specifically
update the loop in apply_remote_tree_state_compat (and the analogous loop around
delta handling) to detect non-Insert(Text) or non-Text keys, propagate a failure
to the caller (or avoid calling apply_remote_tree_snapshot/applying
delta.new_version) so the remote peer can retry or the conflict can be handled
safely.
| let snapshot = mesh_sync.get_tree_snapshot("model1").unwrap(); | ||
| // 2 inserts with shared prefix "test_text_" → at least 3 snapshot nodes | ||
| // (root, shared prefix, and 2 leaf suffixes) | ||
| assert!(snapshot.node_count() >= 3); |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Assert the restored policy behavior, not the mesh snapshot again.
These lines only re-read state from mesh_sync, so a no-op restore_tree_state_from_mesh() would still pass. Please assert against the restored policy tree or a select_worker() outcome so this test actually covers the set_mesh_sync() hydration path.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/policies/cache_aware.rs` around lines 982 - 985, The test
currently asserts mesh_sync snapshot node_count but not that the policy was
actually restored; change the assertion to validate the restored policy behavior
by calling the policy's restore_tree_state_from_mesh() (or invoking
set_mesh_sync() then restore_tree_state_from_mesh()) and then asserting against
the policy tree or a selection outcome (e.g., call select_worker(...) or inspect
PolicyTree methods) to ensure the restored policy contains the expected entries
(shared prefix "test_text_" and two leaf suffixes) rather than just checking
mesh_sync.get_tree_snapshot(); locate functions restore_tree_state_from_mesh,
set_mesh_sync, select_worker, and the PolicyTree/Policy struct in this file to
update the test accordingly.
- Fix pending_count drain bug: materialize_tree_snapshot now tracks _applied_count separately from total_pending. Remove/Token ops in the pending buffer are logged when skipped but still drained on checkpoint (they cannot be retried and would block the buffer forever). - Fix DashMap guard held during merge_snapshot: clone Arc<Tree> before merging to release the shard guard, avoiding blocking concurrent routing. - Fix ping_server no-sync_manager fallback: advance tree_versions when storing snapshot/tree_state bytes directly to tree_configs. - Fix materialize_tree_snapshot return type: now returns (TreeSnapshot, u64, usize) with version captured under same access window, eliminating TOCTOU between materialization and checkpoint version comparison. - Add warn! log for snapshot deserialization failures (was silent return). - Add debug! log for Remove ops skipped in apply_remote_tree_state_compat. - Add clarifying comments about token key deferral in incremental.rs. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 312523bc32
ℹ️ 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 current_version = self.stores.tree_version(&key); | ||
|
|
||
| // Version checks | ||
| if delta.base_version > current_version || current_version >= delta.new_version { |
There was a problem hiding this comment.
Validate delta versions against tree_configs atomically
In apply_remote_tree_delta, the acceptance check reads current_version from tree_versions while the actual payload write happens to tree_configs, and tree_versions is only advanced after releasing the tree_configs entry lock. With concurrent remote updates, a lower-version delta can pass this check during that window and overwrite a newer snapshot/delta that was already written to tree_configs, leaving stale bytes paired with a higher version counter and causing later sync decisions to run on regressed state.
Useful? React with 👍 / 👎.
| let current = self.stores.tree_version(&key); | ||
| if our_version >= current { | ||
| entry.insert(serialized); |
There was a problem hiding this comment.
Gate checkpoint writes on persisted config version
checkpoint_tree_states compares our_version against tree_versions instead of the version represented by the tree_configs entry it is about to overwrite. Because remote apply paths write tree_configs first and only then call advance_tree_version, checkpoint can race in that gap, see an old counter value, and overwrite a newly written higher-version snapshot with stale checkpoint data; once the counter is advanced, the stale bytes remain hidden behind a high version number.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 4
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
crates/mesh/src/ping_server.rs (1)
137-158:⚠️ Potential issue | 🟠 MajorSnapshot requests still skip pending tree ops.
tree_configsonly holds the checkpointed snapshot; newer inserts can still be sitting intree_ops_pending. Wrapping those raw bytes withstores.tree_version(&key)can therefore advertise versionN+kwhile shipping versionNdata, so a joining peer can miss those inserts permanently. Please build these entries from the materialized current snapshot instead of rawtree_configsbytes.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/ping_server.rs` around lines 137 - 158, The code currently uses raw checkpoint bytes from stores.tree_configs (entry.value()) which can be stale because newer inserts live in stores.tree_ops_pending; replace that by materializing the up-to-date snapshot for each key: decode the checkpoint bytes, fetch and apply pending ops from stores.tree_ops_pending for the same key to produce the current tree snapshot bytes, then wrap those bytes into PolicyState (still using stores.tree_version(&key) for the authoritative version) before serializing and pushing into entries; use existing symbols stores.tree_configs, stores.tree_ops_pending, stores.tree_version, and PolicyState to locate and update the logic where entries.push((key, serialized)) is produced.crates/mesh/src/sync.rs (1)
850-858: 🧹 Nitpick | 🔵 TrivialDelta rejected when tree exists only in pending ops buffer.
When
tree_configsis vacant buttree_ops_pendingcontains operations for this model, a delta withbase_version > 0is rejected (lines 852-858). This is correct since deltas require a stored snapshot as base.However, this creates a timing dependency: if
checkpoint_tree_stateshasn't run yet, valid deltas matching the pending state's version will be rejected. Consider documenting this behavior or logging when this edge case occurs.📝 Proposed documentation/logging
Entry::Vacant(entry) => { // No existing config — new tree from delta. if delta.base_version > 0 { debug!( - "Skipped remote tree delta: model={} (base_version={}, new_version={}, no local state)", + "Skipped remote tree delta: model={} (base_version={}, new_version={}, no local snapshot — checkpoint may be pending)", model_id, delta.base_version, delta.new_version ); return; }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/sync.rs` around lines 850 - 858, When handling Entry::Vacant in the tree_configs path, detect the edge case where tree_ops_pending contains operations for the same model_id and delta.base_version > 0 and emit an informative log (or document) instead of silently rejecting: update the vacant-branch logic around the `if delta.base_version > 0` check to query `tree_ops_pending` for the model and log that a pending state exists and that `checkpoint_tree_states` must run before applying this delta (include model_id, delta.base_version, delta.new_version and a short hint). This preserves the current rejection behavior but records the timing dependency for easier debugging.
♻️ Duplicate comments (2)
crates/mesh/src/ping_server.rs (1)
1071-1114:⚠️ Potential issue | 🟠 MajorVersion-check the no-
sync_managertree writes before replacing bytes.Both fallback branches overwrite
stores.tree_configsfirst and only then calladvance_tree_version(). Sinceadvance_tree_version()only keeps the max counter, a stale payload can roll back the stored snapshot bytes while the local version stays high.tree_generationis also untouched, soIncrementalUpdateCollectorwill not re-gossip the replacement. Funnel these writes through the same versioned tree-store path as the sync-manager flow, or at least guard on a winningpolicy_state.versionand bump tree generation on success.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/ping_server.rs` around lines 1071 - 1114, The fallback branches that write directly to stores.tree_configs then call stores.advance_tree_version() can overwrite a newer payload; change them to first check/compare policy_state.version against the current stored version (the same versioning path used by sync_manager.apply_remote_tree_snapshot / apply_remote_tree_state_compat), only replace the bytes when policy_state.version wins, and then call stores.advance_tree_version() and increment the tree_generation (so IncrementalUpdateCollector will re-gossip) as a single atomic logical update; alternatively route these no-sync_manager writes through the same versioned update helper used by the sync_manager path to ensure consistent version checks and generation bumps instead of unconditionally inserting bytes then advancing the version.model_gateway/src/policies/cache_aware.rs (1)
983-989:⚠️ Potential issue | 🟡 MinorAssert the restored policy state, not the mesh snapshot again.
These assertions still mostly re-check the source
mesh_syncdata. A brokenapply_remote_tree_snapshot()that only createsstring_trees["model1"]would still pass. Please assert on the hydrated policy tree contents or a routing outcome afterset_mesh_sync().🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/policies/cache_aware.rs` around lines 983 - 989, The test currently re-checks mesh_sync (mesh_sync.get_tree_snapshot) instead of verifying the policy was actually hydrated; change the assertions to validate the restored policy state after set_mesh_sync() / apply_remote_tree_snapshot(): inspect policy.string_trees.get("model1")'s tree contents (e.g., node_count, expected keys/prefixes, or exact entries) or perform a routing lookup using the policy to verify queries for "test_text_*" route to the expected leaves; replace the mesh_snapshot assertions with assertions against the policy's hydrated tree structure or a routing outcome to ensure apply_remote_tree_snapshot() populated the policy, not just created the key.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@crates/mesh/src/incremental.rs`:
- Around line 324-329: The code currently calls
snapshot.to_bytes().unwrap_or_default() when building PolicyState (variables:
snapshot, PolicyState, config, current_version), which silently converts
serialization failures into an empty config and advances the version; instead,
handle the Result from snapshot.to_bytes() explicitly: if to_bytes() returns
Err, do not construct/send the full_state or mark the version as sent — return
or propagate the error (or log and retry) so the failure is surfaced and
retried; replace unwrap_or_default() with proper error handling that aborts this
update path on serialization failure.
In `@crates/mesh/src/sync.rs`:
- Around line 517-547: Remove the dead _applied_count bookkeeping: delete the
declaration let mut _applied_count = 0; and the increment _applied_count += 1
inside the TreeOperation::Insert branch (where tree.insert(text,
&insert_op.tenant) is called), since only total_pending is returned and
_applied_count is never used; if you want to document intent, replace the
removed increments with a short comment near the loop noting we only track
total_pending (not applied insertions) for snapshot materialization.
- Around line 713-720: The Vacant branch for Entry::Vacant currently inserts the
remote serialized snapshot into tree_configs without verifying local progress;
update this branch to fetch the local tree_version (and/or inspect
tree_ops_pending for this model_id) and only insert the remote snapshot if
remote version > local tree_version (match the logic used in the Occupied
branch), otherwise skip inserting the snapshot (or drop it) to avoid storing a
stale snapshot inconsistent with advance_tree_version/sync_tree_operation state.
- Around line 782-791: The closure apply_ops_to_tree only handles
TreeOperation::Insert(Text) and silently ignores Remove and TreeKey::Tokens;
update apply_ops_to_tree to explicitly match and log skipped operations (e.g.,
when op is TreeOperation::Remove or when Insert has
super::tree_ops::TreeKey::Tokens) using the module's debug/trace logger before
continuing so dropped ops are visible; keep the existing tree.insert(text,
&insert_op.tenant) for Text inserts and still return tree.snapshot() at the end.
---
Outside diff comments:
In `@crates/mesh/src/ping_server.rs`:
- Around line 137-158: The code currently uses raw checkpoint bytes from
stores.tree_configs (entry.value()) which can be stale because newer inserts
live in stores.tree_ops_pending; replace that by materializing the up-to-date
snapshot for each key: decode the checkpoint bytes, fetch and apply pending ops
from stores.tree_ops_pending for the same key to produce the current tree
snapshot bytes, then wrap those bytes into PolicyState (still using
stores.tree_version(&key) for the authoritative version) before serializing and
pushing into entries; use existing symbols stores.tree_configs,
stores.tree_ops_pending, stores.tree_version, and PolicyState to locate and
update the logic where entries.push((key, serialized)) is produced.
In `@crates/mesh/src/sync.rs`:
- Around line 850-858: When handling Entry::Vacant in the tree_configs path,
detect the edge case where tree_ops_pending contains operations for the same
model_id and delta.base_version > 0 and emit an informative log (or document)
instead of silently rejecting: update the vacant-branch logic around the `if
delta.base_version > 0` check to query `tree_ops_pending` for the model and log
that a pending state exists and that `checkpoint_tree_states` must run before
applying this delta (include model_id, delta.base_version, delta.new_version and
a short hint). This preserves the current rejection behavior but records the
timing dependency for easier debugging.
---
Duplicate comments:
In `@crates/mesh/src/ping_server.rs`:
- Around line 1071-1114: The fallback branches that write directly to
stores.tree_configs then call stores.advance_tree_version() can overwrite a
newer payload; change them to first check/compare policy_state.version against
the current stored version (the same versioning path used by
sync_manager.apply_remote_tree_snapshot / apply_remote_tree_state_compat), only
replace the bytes when policy_state.version wins, and then call
stores.advance_tree_version() and increment the tree_generation (so
IncrementalUpdateCollector will re-gossip) as a single atomic logical update;
alternatively route these no-sync_manager writes through the same versioned
update helper used by the sync_manager path to ensure consistent version checks
and generation bumps instead of unconditionally inserting bytes then advancing
the version.
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 983-989: The test currently re-checks mesh_sync
(mesh_sync.get_tree_snapshot) instead of verifying the policy was actually
hydrated; change the assertions to validate the restored policy state after
set_mesh_sync() / apply_remote_tree_snapshot(): inspect
policy.string_trees.get("model1")'s tree contents (e.g., node_count, expected
keys/prefixes, or exact entries) or perform a routing lookup using the policy to
verify queries for "test_text_*" route to the expected leaves; replace the
mesh_snapshot assertions with assertions against the policy's hydrated tree
structure or a routing outcome to ensure apply_remote_tree_snapshot() populated
the policy, not just created the key.
🪄 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: d1f4f023-ee27-4cd9-b0f0-7311d24e07d4
📒 Files selected for processing (4)
crates/mesh/src/incremental.rscrates/mesh/src/ping_server.rscrates/mesh/src/sync.rsmodel_gateway/src/policies/cache_aware.rs
| let mut _applied_count = 0; | ||
| let total_pending; | ||
| if let Some(pending) = self.stores.tree_ops_pending.get(key) { | ||
| pending_count = pending.len(); | ||
| total_pending = pending.len(); | ||
| for op in pending.iter() { | ||
| tree_state.add_operation(op.clone()); | ||
| match op { | ||
| TreeOperation::Insert(insert_op) => { | ||
| if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key { | ||
| tree.insert(text, &insert_op.tenant); | ||
| _applied_count += 1; | ||
| } | ||
| // Token keys: deferred to future phase (not counted) | ||
| } | ||
| TreeOperation::Remove(_) => { | ||
| debug!("Skipping Remove op during snapshot materialization for {key}"); | ||
| // Not counted — must not be drained | ||
| } | ||
| } | ||
| } | ||
| } else { | ||
| total_pending = 0; | ||
| } | ||
|
|
||
| // Return total_pending (not applied_count) because checkpoint | ||
| // drains from the front of the Vec. If all ops up to total_pending | ||
| // are Insert(Text), this is correct. If some are Remove/Tokens, | ||
| // we must still drain total_pending to advance past them — they | ||
| // cannot be retried and would block the buffer forever. | ||
| // The trade-off: Remove/Token ops are lost on checkpoint. This is | ||
| // acceptable because Remove is unused and Tokens are deferred. | ||
| Some((tree.snapshot(), version, total_pending)) |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
_applied_count is computed but never used — consider removing dead code.
The _applied_count variable (lines 517, 526) tracks how many Insert(Text) ops were actually applied, but the function returns total_pending instead. The underscore prefix suppresses the unused warning, but the computation serves no purpose.
If the intent is to document the difference between applied and total ops, a comment would suffice. Otherwise, remove the dead code.
♻️ Proposed cleanup
- // Only count ops actually serialized into the snapshot.
- // Remove and Token ops are NOT representable in snapshot format —
- // they must NOT be drained from pending or they will be lost.
- let mut _applied_count = 0;
let total_pending;
if let Some(pending) = self.stores.tree_ops_pending.get(key) {
total_pending = pending.len();
for op in pending.iter() {
match op {
TreeOperation::Insert(insert_op) => {
if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key {
tree.insert(text, &insert_op.tenant);
- _applied_count += 1;
}
// Token keys: deferred to future phase (not counted)
}
TreeOperation::Remove(_) => {
debug!("Skipping Remove op during snapshot materialization for {key}");
- // Not counted — must not be drained
}
}
}🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@crates/mesh/src/sync.rs` around lines 517 - 547, Remove the dead
_applied_count bookkeeping: delete the declaration let mut _applied_count = 0;
and the increment _applied_count += 1 inside the TreeOperation::Insert branch
(where tree.insert(text, &insert_op.tenant) is called), since only total_pending
is returned and _applied_count is never used; if you want to document intent,
replace the removed increments with a short comment near the loop noting we only
track total_pending (not applied insertions) for snapshot materialization.
| let apply_ops_to_tree = |tree: &kv_index::Tree, ops: &[TreeOperation]| { | ||
| for op in ops { | ||
| if let TreeOperation::Insert(insert_op) = op { | ||
| if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key { | ||
| tree.insert(text, &insert_op.tenant); | ||
| } | ||
| } | ||
| } | ||
| tree.snapshot() | ||
| }; |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Delta apply silently drops Remove/Token ops without logging.
The apply_ops_to_tree helper (lines 782-791) only processes Insert(Text) operations, silently skipping Remove and TreeKey::Tokens. Unlike apply_remote_tree_state_compat which logs skipped Remove ops (line 752-756), this path provides no visibility into dropped operations.
For consistency and operational observability, consider adding debug logging for skipped operations:
🔍 Proposed logging for skipped ops
let apply_ops_to_tree = |tree: &kv_index::Tree, ops: &[TreeOperation]| {
for op in ops {
- if let TreeOperation::Insert(insert_op) = op {
- if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key {
- tree.insert(text, &insert_op.tenant);
+ match op {
+ TreeOperation::Insert(insert_op) => {
+ if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key {
+ tree.insert(text, &insert_op.tenant);
+ }
+ // Token keys silently skipped — deferred to future phase
+ }
+ TreeOperation::Remove(remove_op) => {
+ debug!(
+ model_id = %model_id,
+ tenant = %remove_op.tenant,
+ "Skipping Remove op in delta apply"
+ );
}
}
}
tree.snapshot()
};📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| let apply_ops_to_tree = |tree: &kv_index::Tree, ops: &[TreeOperation]| { | |
| for op in ops { | |
| if let TreeOperation::Insert(insert_op) = op { | |
| if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key { | |
| tree.insert(text, &insert_op.tenant); | |
| } | |
| } | |
| } | |
| tree.snapshot() | |
| }; | |
| let apply_ops_to_tree = |tree: &kv_index::Tree, ops: &[TreeOperation]| { | |
| for op in ops { | |
| match op { | |
| TreeOperation::Insert(insert_op) => { | |
| if let super::tree_ops::TreeKey::Text(ref text) = insert_op.key { | |
| tree.insert(text, &insert_op.tenant); | |
| } | |
| // Token keys silently skipped — deferred to future phase | |
| } | |
| TreeOperation::Remove(remove_op) => { | |
| debug!( | |
| model_id = %model_id, | |
| tenant = %remove_op.tenant, | |
| "Skipping Remove op in delta apply" | |
| ); | |
| } | |
| } | |
| } | |
| tree.snapshot() | |
| }; |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@crates/mesh/src/sync.rs` around lines 782 - 791, The closure
apply_ops_to_tree only handles TreeOperation::Insert(Text) and silently ignores
Remove and TreeKey::Tokens; update apply_ops_to_tree to explicitly match and log
skipped operations (e.g., when op is TreeOperation::Remove or when Insert has
super::tree_ops::TreeKey::Tokens) using the module's debug/trace logger before
continuing so dropped ops are visible; keep the existing tree.insert(text,
&insert_op.tenant) for Text inserts and still return tree.snapshot() at the end.
The collector was combining tree deltas, tree full-state snapshots, and policy entries into a single batch under StoreType::Policy. With diverse 20k-char prompts, the full-state snapshot alone is 28 MB, exceeding the 10 MB gRPC limit. This caused the ENTIRE batch to be rejected every cycle — including the small (~1 MB) deltas that would have fit. Since mark_sent was never called, versions never advanced, creating an infinite retry loop where no tree data was ever synchronized. Fix: - Split tree entries into per-key batches in collect_all_updates() so each tree key gets independent size checking. Small deltas pass through even when a full-state snapshot is oversized. - Phase 2 now skips keys that have pending ops (already covered by Phase 1 deltas), preventing redundant full-state snapshots from being added to every cycle. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (1)
crates/mesh/src/incremental.rs (1)
324-329:⚠️ Potential issue | 🟠 MajorHandle
snapshot.to_bytes()failures explicitly.
unwrap_or_default()can sendpolicy_type = "tree_snapshot"withconfig = []andversion = current_version. Once that batch is marked sent, this collector stops retrying the real snapshot for that key.♻️ Suggested fix
let snapshot = tree.snapshot(); + let snapshot_bytes = match snapshot.to_bytes() { + Ok(bytes) => bytes, + Err(error) => { + debug!( + "Skipping full tree snapshot fallback for {} — failed to serialize snapshot: {}", + key, error + ); + continue; + } + }; let full_state = PolicyState { model_id, policy_type: "tree_snapshot".to_string(), - config: snapshot.to_bytes().unwrap_or_default(), + config: snapshot_bytes, version: current_version, };🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/incremental.rs` around lines 324 - 329, The code currently uses snapshot.to_bytes().unwrap_or_default() when constructing PolicyState (in the block creating full_state after tree.snapshot()), which can silently send an empty config and stop retries; replace unwrap_or_default with explicit error handling: call snapshot.to_bytes() and if it Err, log/return that error (or propagate it up) so the batch is not marked sent, and only construct PolicyState (with policy_type "tree_snapshot" and config set to the bytes) when to_bytes() succeeds; ensure this handling touches the snapshot.to_bytes() call site and the surrounding logic that marks batches as sent so failed serializations can be retried.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@crates/mesh/src/incremental.rs`:
- Around line 477-480: Splitting StoreType::Policy into multiple batches causes
premature advancement of generation because collector.mark_sent(*store_type,
updates) advances last_scanned.policy and last_scanned.tree on the first
successful batch; modify the logic in the batch collection/return path (the
functions that yield the (store_type, updates) tuples consumed by controller.rs)
so that generation/last_scanned is only advanced after all policy batches for a
single collection succeed, or alternatively include explicit metadata with each
returned tuple (e.g., is_final_policy_batch or batch_index/total_batches) so
collector.mark_sent can detect the final policy batch before advancing
generations; update collector.mark_sent (and any callers) to accept and act on
that metadata and ensure the same fix is applied to the other affected location
(lines referenced 492-517) to preserve retry semantics.
- Around line 312-329: The code builds a tree_snapshot PolicyState using
current_version even when pending contains operations that the snapshot format
won't represent; scan the pending buffer before snapshotting and if any op is
not TreeOperation::Insert with TreeKey::Text (i.e., matches
TreeOperation::Remove or TreeKey::Tokens), do not emit a tree_snapshot at
current_version — instead abort/refuse to create the PolicyState (or return an
error) so we don't advance the advertised version without delivering those ops;
implement this check around the loop that iterates pending (referencing pending,
super::tree_ops::TreeOperation::Insert, super::tree_ops::TreeOperation::Remove,
super::tree_ops::TreeKey::Text, TreeKey::Tokens, tree.snapshot(), PolicyState,
current_version, full_state) and only construct full_state when the buffer
contains exclusively representable Insert(Text) entries.
---
Duplicate comments:
In `@crates/mesh/src/incremental.rs`:
- Around line 324-329: The code currently uses
snapshot.to_bytes().unwrap_or_default() when constructing PolicyState (in the
block creating full_state after tree.snapshot()), which can silently send an
empty config and stop retries; replace unwrap_or_default with explicit error
handling: call snapshot.to_bytes() and if it Err, log/return that error (or
propagate it up) so the batch is not marked sent, and only construct PolicyState
(with policy_type "tree_snapshot" and config set to the bytes) when to_bytes()
succeeds; ensure this handling touches the snapshot.to_bytes() call site and the
surrounding logic that marks batches as sent so failed serializations can be
retried.
🪄 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: 7b05ed77-72d6-4284-86bf-fca91af51902
📒 Files selected for processing (1)
crates/mesh/src/incremental.rs
| /// Collect all incremental updates across all stores. | ||
| /// | ||
| /// Tree updates are split into per-key batches so a single oversized | ||
| /// snapshot doesn't block delivery of small deltas or policy updates. |
There was a problem hiding this comment.
Splitting StoreType::Policy into multiple batches breaks retry semantics.
crates/mesh/src/controller.rs:424-520 calls collector.mark_sent(*store_type, updates) for each returned tuple. After the first successful StoreType::Policy batch, mark_sent() advances both last_scanned.policy and last_scanned.tree; any later policy/tree batch from the same collection that hits size limits or backpressure will not be re-collected until a fresh generation bump. Please defer the generation advance until all policy batches from one collection succeed, or return enough metadata for mark_sent() to know when it is handling the final policy batch.
Also applies to: 492-517
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@crates/mesh/src/incremental.rs` around lines 477 - 480, Splitting
StoreType::Policy into multiple batches causes premature advancement of
generation because collector.mark_sent(*store_type, updates) advances
last_scanned.policy and last_scanned.tree on the first successful batch; modify
the logic in the batch collection/return path (the functions that yield the
(store_type, updates) tuples consumed by controller.rs) so that
generation/last_scanned is only advanced after all policy batches for a single
collection succeed, or alternatively include explicit metadata with each
returned tuple (e.g., is_final_policy_batch or batch_index/total_batches) so
collector.mark_sent can detect the final policy batch before advancing
generations; update collector.mark_sent (and any callers) to accept and act on
that metadata and ensure the same fix is applied to the other affected location
(lines referenced 492-517) to preserve retry semantics.
The full-state fallback builds the entire tree from snapshot + all pending ops and sends it as a single entry. With diverse 20k-char prompts this produces 35 MB entries that exceed the 10 MB gRPC limit every cycle, blocking all tree sync. Remove the full-state fallback entirely. Trees sync via: - Deltas (small, ~1 MB per cycle) for incremental updates - Initial snapshot exchange (ping_server) for new peer join - Phase 2 (tree_configs-only entries) for remote-only trees When a version gap prevents delta construction, skip the key. The next checkpoint folds pending ops into tree_configs, resetting the delta window for the next cycle. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
1 similar comment
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
f363a2d to
60aa15d
Compare
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
60aa15d to
0a05345
Compare
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
- Fix pending_count drain bug: materialize_tree_snapshot now tracks _applied_count separately from total_pending. Remove/Token ops in the pending buffer are logged when skipped but still drained on checkpoint (they cannot be retried and would block the buffer forever). - Fix DashMap guard held during merge_snapshot: clone Arc<Tree> before merging to release the shard guard, avoiding blocking concurrent routing. - Fix ping_server no-sync_manager fallback: advance tree_versions when storing snapshot/tree_state bytes directly to tree_configs. - Fix materialize_tree_snapshot return type: now returns (TreeSnapshot, u64, usize) with version captured under same access window, eliminating TOCTOU between materialization and checkpoint version comparison. - Add warn! log for snapshot deserialization failures (was silent return). - Add debug! log for Remove ops skipped in apply_remote_tree_state_compat. - Add clarifying comments about token key deferral in incremental.rs. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
The collector was combining tree deltas, tree full-state snapshots, and policy entries into a single batch under StoreType::Policy. With diverse 20k-char prompts, the full-state snapshot alone is 28 MB, exceeding the 10 MB gRPC limit. This caused the ENTIRE batch to be rejected every cycle — including the small (~1 MB) deltas that would have fit. Since mark_sent was never called, versions never advanced, creating an infinite retry loop where no tree data was ever synchronized. Fix: - Split tree entries into per-key batches in collect_all_updates() so each tree key gets independent size checking. Small deltas pass through even when a full-state snapshot is oversized. - Phase 2 now skips keys that have pending ops (already covered by Phase 1 deltas), preventing redundant full-state snapshots from being added to every cycle. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
The full-state fallback builds the entire tree from snapshot + all pending ops and sends it as a single entry. With diverse 20k-char prompts this produces 35 MB entries that exceed the 10 MB gRPC limit every cycle, blocking all tree sync. Remove the full-state fallback entirely. Trees sync via: - Deltas (small, ~1 MB per cycle) for incremental updates - Initial snapshot exchange (ping_server) for new peer join - Phase 2 (tree_configs-only entries) for remote-only trees When a version gap prevents delta construction, skip the key. The next checkpoint folds pending ops into tree_configs, resetting the delta window for the next cycle. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
0a05345 to
d9fb223
Compare
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
When mark_sent fails (batch rejected for size), last_sent_version never advances. Next cycle collects ALL unsent ops into one delta. With 50 rps × 60s = 3000 ops × 20k chars = 60 MB delta — far exceeding the 10 MB gRPC limit. This creates an infinite growth loop. Fix: cap delta to MAX_OPS_PER_DELTA (200) ops per cycle. With 20k char ops, 200 × 20k = 4 MB which fits in 10 MB. Remaining ops are sent in subsequent cycles as mark_sent advances the version window. Also updates two tests that relied on full-state fallback for new peers after buffer trim. Without full-state fallback, new peers with version gaps receive the tree via the initial snapshot exchange (ping_server). Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
crates/mesh/src/ping_server.rs (1)
137-158:⚠️ Potential issue | 🟠 MajorSnapshot responses are still built from stale tree blobs, and they can still overflow the transport limit.
Line 142 copies the last checkpointed
tree_configsbytes, while Line 147 advertisesstores.tree_version(&key), which can already include newer ops still sitting intree_ops_pendingor even keys that have notree_configsentry yet. That lets a joiner store versionNwith bytes fromN-k, then drop the follow-up delta as stale. Because these entries are now multi-megabyte"tree_snapshot"payloads, feeding them into the later count-based chunker can also recreate >10 MBSnapshotChunks. Please build snapshot replies from a materialized(snapshot, version)view overtree_configs + tree_ops_pending, and split tree entries by serialized byte size before chunking.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/ping_server.rs` around lines 137 - 158, The current code reads from stores.tree_configs and uses stores.tree_version(&key) which can be newer than the serialized checkpoint bytes, causing joiners to receive stale tree_snapshot payloads and oversized chunks; change the logic that builds the "tree_snapshot" PolicyState so it materializes a coherent (snapshot, version) pair by merging stores.tree_configs with in-memory tree_ops_pending (apply pending ops on top of the checkpoint to produce the up-to-date snapshot and its corresponding version) instead of using the raw tree_configs bytes and stores.tree_version separately, and after creating the PolicyState and serializing it, split large tree entries by serialized byte size (not by count) before passing them to the later chunker to ensure no SnapshotChunk exceeds the transport limit; refer to stores.tree_configs, tree_ops_pending, stores.tree_version, PolicyState and the "tree_snapshot" policy_type to locate and update this logic.
♻️ Duplicate comments (4)
crates/mesh/src/incremental.rs (1)
454-475:⚠️ Potential issue | 🔴 CriticalSplit policy batches can still be lost after the first successful send.
The sender still calls
mark_sent(StoreType::Policy, updates)per batch. If the first policy batch succeeds and a later tree batch from the same collection hits size limits or channel backpressure,last_scanned.policy/treeis already advanced, so that unsent key is skipped on the next collection until another generation bump happens. Please defer generation advancement until the final policy batch succeeds, or attach batch metadata somark_sent()can tell when it is handling the last batch.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/incremental.rs` around lines 454 - 475, The current split of policy and tree updates (in the loop that builds policy_updates and tree_batches) advances progress via mark_sent(StoreType::Policy, updates) per batch which can move last_scanned.policy/tree prematurely and cause later unsent tree batches to be skipped; change the send/ack logic so generation advancement is deferred until the final policy-related batch for that collection succeeds—either by bundling metadata with each batch indicating "is_last_for_generation" or by aggregating all policy/tree batches for a given generation and only calling mark_sent once after all batches are confirmed; update the code paths that call mark_sent and any logic that updates last_scanned.policy/last_scanned.tree to check this metadata or only advance after the final batch succeeds (refer to the tree_batches/policy_updates construction and the mark_sent(StoreType::Policy, ...) calls).model_gateway/src/policies/cache_aware.rs (1)
983-989:⚠️ Potential issue | 🟡 MinorThis test still doesn't prove the restore happened.
tree.is_some()only shows that a local entry exists, andmesh_sync.get_tree_snapshot("model1")re-reads the source mesh state. An implementation that creates an empty local tree would still pass here. Please assert restored routing behavior or inspect the restored local tree contents directly.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/policies/cache_aware.rs` around lines 983 - 989, The test currently only checks mesh_sync.get_tree_snapshot("model1") and tree.is_some(), which can pass without actual local restore; change the assertions to inspect the restored local tree contents or routing behavior: unwrap the local tree returned by policy.string_trees.get("model1") (e.g., let tree = policy.string_trees.get("model1").unwrap()) and assert its snapshot/node/leaf contents include the expected shared prefix and two leaf suffixes (or check node_count on the unwrapped tree), or alternatively assert routing by invoking the policy's routing method for "model1" with known keys and verifying it resolves to the restored leaves; replace the weak assert with these stronger checks referencing policy.string_trees.get and mesh_sync.get_tree_snapshot where appropriate.crates/mesh/src/sync.rs (1)
713-719:⚠️ Potential issue | 🟠 MajorA vacant
tree_configsslot no longer means there is no local tree.After this refactor, a model can live only in
tree_ops_pendingwhiletree_configsis still vacant. These branches still treat that case as empty local state, so a stale remote snapshot/delta can be inserted over a newer pending-only tree;advance_tree_version()then leaves the version counter ahead of the stored bytes. Please gate the vacant case ontree_version == 0and no pending ops, or materialize the pending-only local tree before deciding whether to accept the remote payload.Also applies to: 850-869
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/sync.rs` around lines 713 - 719, The vacant-branch logic for inserting a remote snapshot into tree_configs must not assume absence of a local tree: before inserting serialized into the Entry::Vacant branch, check that tree_version == 0 and that tree_ops_pending for model_id is empty; if either condition fails, materialize the pending-only local tree (or reject/merge the remote payload) so you don't overwrite a newer pending-only state and desync advance_tree_version(); apply the same guard/behavior to the corresponding vacant-case at the other site handling model_id (the block analogous to lines ~850-869) and reference tree_configs, tree_ops_pending, tree_version, advance_tree_version(), serialized, model_id, and version when making the change.crates/mesh/src/ping_server.rs (1)
1071-1088:⚠️ Potential issue | 🟠 MajorThe no-
sync_managerfallback still overwrites newer tree bytes without a version guard.Both fallback branches call
stores.tree_configs.insert(...)unconditionally and only thenadvance_tree_version(...). If an older peer sends version 8 after this node already has version 10 locally, the version counter stays at 10 but the stored snapshot bytes roll back to 8, which leaves later materialization/re-gossip inconsistent. Reuse the same guarded write path asMeshSyncManager::apply_remote_tree_snapshot()here.Also applies to: 1089-1114
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@crates/mesh/src/ping_server.rs` around lines 1071 - 1088, The fallback path that writes policy_state.config into stores.tree_configs unconditionally can regress the stored snapshot when an older version arrives; change the no-sync_manager branch to perform the same guarded update used by MeshSyncManager::apply_remote_tree_snapshot(): read the current stored version for the given key, compare it to policy_state.version, and only replace tree_configs and call stores.advance_tree_version(&key, policy_state.version) when the incoming version is newer (or otherwise permitted by the same tie-break logic used in apply_remote_tree_snapshot), ensuring identical version-check semantics; apply the same guarded-write change to the other occurrence referenced (lines 1089–1114).
🤖 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/policies/cache_aware.rs`:
- Around line 251-263: The mesh restore currently only replays
get_all_tree_snapshots() into string_trees via apply_remote_tree_snapshot(),
which leaves token_trees uninitialized and breaks select_worker_with_tokens();
update restore_tree_state_from_mesh (and the similar block at the other
occurrence) to also restore token state: either extend
apply_remote_tree_snapshot to return or emit token-related data and merge that
into token_trees, or add a new
apply_remote_token_snapshot/get_all_token_snapshots path that reconstructs
token_trees from the mesh snapshots; ensure the code updates the same shared
structures used by select_worker_with_tokens() so cache-affinity state is
restored after mesh recovery.
---
Outside diff comments:
In `@crates/mesh/src/ping_server.rs`:
- Around line 137-158: The current code reads from stores.tree_configs and uses
stores.tree_version(&key) which can be newer than the serialized checkpoint
bytes, causing joiners to receive stale tree_snapshot payloads and oversized
chunks; change the logic that builds the "tree_snapshot" PolicyState so it
materializes a coherent (snapshot, version) pair by merging stores.tree_configs
with in-memory tree_ops_pending (apply pending ops on top of the checkpoint to
produce the up-to-date snapshot and its corresponding version) instead of using
the raw tree_configs bytes and stores.tree_version separately, and after
creating the PolicyState and serializing it, split large tree entries by
serialized byte size (not by count) before passing them to the later chunker to
ensure no SnapshotChunk exceeds the transport limit; refer to
stores.tree_configs, tree_ops_pending, stores.tree_version, PolicyState and the
"tree_snapshot" policy_type to locate and update this logic.
---
Duplicate comments:
In `@crates/mesh/src/incremental.rs`:
- Around line 454-475: The current split of policy and tree updates (in the loop
that builds policy_updates and tree_batches) advances progress via
mark_sent(StoreType::Policy, updates) per batch which can move
last_scanned.policy/tree prematurely and cause later unsent tree batches to be
skipped; change the send/ack logic so generation advancement is deferred until
the final policy-related batch for that collection succeeds—either by bundling
metadata with each batch indicating "is_last_for_generation" or by aggregating
all policy/tree batches for a given generation and only calling mark_sent once
after all batches are confirmed; update the code paths that call mark_sent and
any logic that updates last_scanned.policy/last_scanned.tree to check this
metadata or only advance after the final batch succeeds (refer to the
tree_batches/policy_updates construction and the mark_sent(StoreType::Policy,
...) calls).
In `@crates/mesh/src/ping_server.rs`:
- Around line 1071-1088: The fallback path that writes policy_state.config into
stores.tree_configs unconditionally can regress the stored snapshot when an
older version arrives; change the no-sync_manager branch to perform the same
guarded update used by MeshSyncManager::apply_remote_tree_snapshot(): read the
current stored version for the given key, compare it to policy_state.version,
and only replace tree_configs and call stores.advance_tree_version(&key,
policy_state.version) when the incoming version is newer (or otherwise permitted
by the same tie-break logic used in apply_remote_tree_snapshot), ensuring
identical version-check semantics; apply the same guarded-write change to the
other occurrence referenced (lines 1089–1114).
In `@crates/mesh/src/sync.rs`:
- Around line 713-719: The vacant-branch logic for inserting a remote snapshot
into tree_configs must not assume absence of a local tree: before inserting
serialized into the Entry::Vacant branch, check that tree_version == 0 and that
tree_ops_pending for model_id is empty; if either condition fails, materialize
the pending-only local tree (or reject/merge the remote payload) so you don't
overwrite a newer pending-only state and desync advance_tree_version(); apply
the same guard/behavior to the corresponding vacant-case at the other site
handling model_id (the block analogous to lines ~850-869) and reference
tree_configs, tree_ops_pending, tree_version, advance_tree_version(),
serialized, model_id, and version when making the change.
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 983-989: The test currently only checks
mesh_sync.get_tree_snapshot("model1") and tree.is_some(), which can pass without
actual local restore; change the assertions to inspect the restored local tree
contents or routing behavior: unwrap the local tree returned by
policy.string_trees.get("model1") (e.g., let tree =
policy.string_trees.get("model1").unwrap()) and assert its snapshot/node/leaf
contents include the expected shared prefix and two leaf suffixes (or check
node_count on the unwrapped tree), or alternatively assert routing by invoking
the policy's routing method for "model1" with known keys and verifying it
resolves to the restored leaves; replace the weak assert with these stronger
checks referencing policy.string_trees.get and mesh_sync.get_tree_snapshot where
appropriate.
🪄 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: 2d9686d1-cc81-4601-aa26-4f9f809beec1
📒 Files selected for processing (9)
crates/mesh/src/controller.rscrates/mesh/src/incremental.rscrates/mesh/src/lib.rscrates/mesh/src/ping_server.rscrates/mesh/src/sync.rsdeploy/helm/smg/templates/_helpers.tplmodel_gateway/src/policies/cache_aware.rsmodel_gateway/src/policies/registry.rsmodel_gateway/tests/mesh_integration_test.rs
| fn restore_tree_state_from_mesh(&self) { | ||
| let tree_states = { | ||
| let snapshots = { | ||
| let guard = self.mesh_sync.read(); | ||
| guard | ||
| .as_ref() | ||
| .map(|mesh_sync| mesh_sync.get_all_tree_states()) | ||
| .map(|mesh_sync| mesh_sync.get_all_tree_snapshots()) | ||
| }; | ||
|
|
||
| if let Some(tree_states) = tree_states { | ||
| for tree_state in &tree_states { | ||
| // Use the merge path (not replace) so concurrent live updates | ||
| // arriving via subscriber callbacks are not overwritten. | ||
| self.apply_remote_tree_state(&tree_state.model_id, tree_state); | ||
| if let Some(snapshots) = snapshots { | ||
| for (model_id, snapshot) in &snapshots { | ||
| self.apply_remote_tree_snapshot(model_id, snapshot); | ||
| } | ||
| } |
There was a problem hiding this comment.
Snapshot recovery is string-tree only.
After this change, mesh restore replays only get_all_tree_snapshots(), and apply_remote_tree_snapshot() only merges into string_trees. select_worker_with_tokens() still depends on token_trees, so a node that restarts or falls back to snapshot recovery loses gRPC cache-affinity state until fresh traffic repopulates it. Please keep a token-tree recovery path instead of treating the snapshot as a full tree restore.
Also applies to: 332-345
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/policies/cache_aware.rs` around lines 251 - 263, The mesh
restore currently only replays get_all_tree_snapshots() into string_trees via
apply_remote_tree_snapshot(), which leaves token_trees uninitialized and breaks
select_worker_with_tokens(); update restore_tree_state_from_mesh (and the
similar block at the other occurrence) to also restore token state: either
extend apply_remote_tree_snapshot to return or emit token-related data and merge
that into token_trees, or add a new
apply_remote_token_snapshot/get_all_token_snapshots path that reconstructs
token_trees from the mesh snapshots; ensure the code updates the same shared
structures used by select_worker_with_tokens() so cache-affinity state is
restored after mesh recovery.
| if let Ok(snapshot) = | ||
| kv_index::snapshot::TreeSnapshot::from_bytes( | ||
| &policy_state.config, | ||
| ) | ||
| { | ||
| sync_manager | ||
| .apply_remote_tree_snapshot( | ||
| policy_state | ||
| .model_id | ||
| .clone(), | ||
| snapshot, | ||
| policy_state.version, | ||
| actor, | ||
| ); | ||
| } |
There was a problem hiding this comment.
🟡 Nit: Silent failure — when TreeSnapshot::from_bytes fails here, the error is silently dropped (if let Ok(...)). In controller.rs (line 674), the same case logs a warn!. Consider logging for consistency so operators can diagnose corrupted snapshots in this code path too.
Remove the entire raw-operation sync path (tree_ops_pending buffer + TreeStateDelta encoding). This was the root cause of the memory leak and 28+ MB message sizes: every request buffered the full prompt text (up to 80k chars) as a TreeOperation, and the collector sent these raw ops as deltas over the wire. New approach: - sync_tree_operation() just bumps the version counter (O(1), zero alloc) - checkpoint_tree_states() gets snapshots directly from the radix trees via TreeSnapshotProvider (implemented by CacheAwarePolicy/PolicyRegistry) - Collector sends compact TreeSnapshot bytes from tree_configs - No raw operations buffered, no raw ops on the wire The TreeSnapshotProvider trait bridges the crate boundary: smg-mesh defines it, model_gateway implements it on CacheAwarePolicy (which owns the radix trees). The provider is registered in server.rs where mesh is wired up. Removed: - tree_ops_pending buffer writes from hot path - Phase 1 delta collection (TreeStateDelta from pending ops) - 9 delta-specific tests (tested dead code path) - Full-state fallback from collector Changed: - sync.rs: sync_tree_operation ignores the operation, just bumps version - sync.rs: checkpoint_tree_states uses TreeSnapshotProvider - sync.rs: added TreeSnapshotProvider trait + registration - incremental.rs: simplified to only send tree_configs snapshots - lib.rs: re-export TreeSnapshotProvider - server.rs: register PolicyRegistry as TreeSnapshotProvider - cache_aware.rs: implement TreeSnapshotProvider (snapshot string_trees) - registry.rs: implement TreeSnapshotProvider (delegate to cache_aware) Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
|
|
||
| use anyhow::Result; | ||
| use futures::Stream; | ||
| use kv_index::RadixTree; |
There was a problem hiding this comment.
🟡 Nit: RadixTree is imported but never used in this file. This will likely trigger an unused_imports warning.
| use kv_index::RadixTree; |
|
|
||
| use std::sync::Arc; | ||
|
|
||
| use kv_index::RadixTree; |
There was a problem hiding this comment.
🟡 Nit: RadixTree is imported but never used in this test file. This will trigger an unused_imports warning.
| use kv_index::RadixTree; |
| if let Ok(snapshot) = | ||
| kv_index::snapshot::TreeSnapshot::from_bytes( | ||
| &policy_state.config, | ||
| ) | ||
| { | ||
| sync_manager | ||
| .apply_remote_tree_snapshot( | ||
| policy_state | ||
| .model_id | ||
| .clone(), | ||
| snapshot, | ||
| policy_state.version, | ||
| actor, | ||
| ); | ||
| } |
There was a problem hiding this comment.
🟡 Nit: Deserialization errors for tree_snapshot are silently swallowed here — the if let Ok(snapshot) discards the Err case with no log. Compare with controller.rs (line 673–678) which correctly logs a warn! on failure. This inconsistency makes debugging gossip deserialization issues harder in production.
Consider adding a log on the Err case:
match kv_index::snapshot::TreeSnapshot::from_bytes(&policy_state.config) {
Ok(snapshot) => { sync_manager.apply_remote_tree_snapshot(...); }
Err(e) => {
log::warn!("Failed to deserialize tree snapshot for model {}: {e}", policy_state.model_id);
}
}| if policy_state.policy_type == "tree_snapshot" { | ||
| // New compact snapshot format | ||
| if let Some(ref sync_manager) = sync_manager { | ||
| if let Ok(snapshot) = kv_index::snapshot::TreeSnapshot::from_bytes( | ||
| &policy_state.config | ||
| ) { | ||
| sync_manager.apply_remote_tree_snapshot( | ||
| policy_state.model_id.clone(), | ||
| snapshot, | ||
| policy_state.version, | ||
| Some(entry.actor.clone()), | ||
| ); | ||
| } | ||
| } else { | ||
| // No sync manager — store bytes and advance version. | ||
| stores.tree_configs.insert(key.clone(), policy_state.config.clone()); | ||
| stores.advance_tree_version(&key, policy_state.version); | ||
| } |
There was a problem hiding this comment.
🟡 Nit: Same silent swallow pattern — if let Ok(snapshot) on line 1074 discards the Err without logging. The anti-entropy path should log deserialization failures the same way controller.rs does. Corrupted snapshots arriving via anti-entropy would be invisible in production logs.
| impl smg_mesh::TreeSnapshotProvider for PolicyRegistry { | ||
| fn snapshot_all_trees(&self) -> Vec<(String, smg_mesh::TreeSnapshot)> { | ||
| // Collect snapshots from the default policy's cache_aware trees. | ||
| if let Some(cache_aware) = self | ||
| .default_policy | ||
| .as_any() | ||
| .downcast_ref::<CacheAwarePolicy>() | ||
| { | ||
| return cache_aware.snapshot_all_trees(); | ||
| } | ||
| Vec::new() | ||
| } |
There was a problem hiding this comment.
🔴 Important: snapshot_all_trees() only collects snapshots from default_policy. In PD (prefill/decode disaggregation) mode, the prefill and decode policies have their own CacheAwarePolicy instances with independent string_trees. Trees inserted into those policies will never be checkpointed to mesh.
Compare with apply_remote_tree_snapshot (lines 527–571) which correctly fans out to all four policy locations (model-specific, default, prefill, decode). The checkpoint path should do the inverse — collect from all policies, not just default.
For example:
fn snapshot_all_trees(&self) -> Vec<(String, smg_mesh::TreeSnapshot)> {
let mut snapshots = Vec::new();
let mut seen = std::collections::HashSet::new();
// Helper to collect from a cache_aware policy
let mut collect = |policy: &dyn LoadBalancingPolicy| {
if let Some(ca) = policy.as_any().downcast_ref::<CacheAwarePolicy>() {
for (model_id, snap) in ca.snapshot_all_trees() {
if seen.insert(model_id.clone()) {
snapshots.push((model_id, snap));
}
}
}
};
collect(&*self.default_policy);
if let Some(p) = self.prefill_policy.get() { collect(&**p); }
if let Some(p) = self.decode_policy.get() { collect(&**p); }
snapshots
}| // Version check and insert happen under the same shard lock. | ||
| // Read tree_version INSIDE the entry match arm so no concurrent | ||
| // apply_remote_tree_snapshot can advance the version between | ||
| // our read and our write. | ||
| let applied = match self.stores.tree_configs.entry(key.clone()) { | ||
| Entry::Occupied(mut entry) => { | ||
| let current_version = TreeState::from_bytes(entry.get()) | ||
| .ok() | ||
| .map(|ts| ts.version) | ||
| .unwrap_or(0); | ||
| if tree_state.version > current_version { | ||
| let current_version = self.stores.tree_version(&key); |
There was a problem hiding this comment.
🟡 Nit: The comment says "Read tree_version INSIDE the entry match arm so no concurrent apply_remote_tree_snapshot can advance the version between our read and our write" — but tree_version() reads from a separate AtomicU64 in tree_versions, which is NOT guarded by the DashMap shard lock on tree_configs. A concurrent advance_tree_version can still race between the tree_version read (line 688) and the entry.insert (line 690).
The old code had the same gap (it read version from the blob bytes under the shard lock, which was truly atomic). In practice this is benign — worst case, a slightly-stale snapshot wins and the next checkpoint corrects it — but the comment is misleading. Consider updating the comment to reflect the actual guarantee.
| // Only count ops actually serialized into the snapshot. | ||
| // Remove and Token ops are NOT representable in snapshot format — | ||
| // they must NOT be drained from pending or they will be lost. | ||
| let mut _applied_count = 0; |
There was a problem hiding this comment.
🟡 Nit: _applied_count is computed but never read (the leading underscore suppresses the warning but the variable itself serves no purpose). Consider removing it entirely — the comment on line 505–507 documents the intent, and the actual return uses total_pending.
| pub fn checkpoint_tree_states(&self) { | ||
| use dashmap::mapref::entry::Entry; | ||
|
|
||
| let keys: Vec<String> = self | ||
| .stores | ||
| .tree_ops_pending | ||
| .iter() | ||
| .map(|e| e.key().clone()) | ||
| .collect(); | ||
| let provider = self.tree_snapshot_provider.read().clone(); | ||
| let Some(provider) = provider else { | ||
| return; | ||
| }; | ||
|
|
||
| for key in keys { | ||
| let model_id = key.strip_prefix("tree:").unwrap_or(&key); | ||
| if let Some((tree_state, pending_count)) = self.materialize_tree_state(&key, model_id) { | ||
| if let Ok(serialized) = tree_state.to_bytes() { | ||
| // Write to tree_configs only if our materialized version | ||
| // >= the current config version. A concurrent remote | ||
| // update may have written a newer entry between | ||
| // materialize_tree_state and now. | ||
| let inserted = match self.stores.tree_configs.entry(key.clone()) { | ||
| Entry::Occupied(mut entry) => { | ||
| let current = TreeState::from_bytes(entry.get()) | ||
| .ok() | ||
| .map(|ts| ts.version) | ||
| .unwrap_or(0); | ||
| if tree_state.version >= current { | ||
| entry.insert(serialized); | ||
| true | ||
| } else { | ||
| // A concurrent remote update wrote a newer | ||
| // version — skip our stale checkpoint. | ||
| false | ||
| } | ||
| } | ||
| Entry::Vacant(entry) => { | ||
| entry.insert(serialized); | ||
| true | ||
| } | ||
| }; | ||
|
|
||
| // Only drain pending ops if the checkpoint was actually | ||
| // written. If a concurrent remote update won the version | ||
| // race, our materialized data is stale — draining would | ||
| // lose ops that haven't been persisted anywhere. | ||
| if inserted { | ||
| self.stores.tree_ops_pending.alter(&key, |_, mut ops| { | ||
| if ops.len() <= pending_count { | ||
| ops = Vec::new(); | ||
| } else { | ||
| ops.drain(..pending_count); | ||
| } | ||
| ops | ||
| }); | ||
| } | ||
| } | ||
| for (model_id, snapshot) in provider.snapshot_all_trees() { | ||
| let key = tree_state_key(&model_id); | ||
| if let Ok(serialized) = snapshot.to_bytes() { | ||
| self.stores.tree_configs.insert(key.clone(), serialized); | ||
| // Ensure version is current so the collector picks it up. | ||
| self.stores.bump_tree_version(&key); | ||
| } | ||
| } |
There was a problem hiding this comment.
🟡 Nit: checkpoint_tree_states unconditionally overwrites tree_configs without checking whether the provider's snapshot is newer than what's already stored. A concurrent apply_remote_tree_snapshot from a peer (which does version-check) could store a newer version, and then this checkpoint overwrites it with a potentially older local snapshot. The bump_tree_version call also increments unconditionally, so the stale data gets a newer version number — making it impossible for the peer's version to win on subsequent syncs.
Consider adding a version comparison (or at least a generation check) before overwriting, similar to the old checkpoint_tree_states which compared versions before inserting.
Two critical fixes for the two-layer sync protocol: 1. TenantInsert.node_path was storing the full prompt text (80k chars) making it no smaller than raw TreeOperation. Changed to node_path_hash: u64 (8 bytes). The receiver looks up nodes by hash; unknown hashes are buffered until the next structure snapshot. Similarly, TenantEvict now uses node_path_hash (0 = global evict). 2. sync_tree_operation was STILL pushing raw operations to tree_ops_pending (the memory leak source). Removed entirely. The tenant delta buffer is now the only sync data on the hot path. Also: - hash_node_path() utility for consistent path hashing - TreeStateSubscriber::apply_tenant_delta is now a required method (no default impl since hash can't reconstruct text) - Removed 8 dead delta-encoding tests - Updated 6 tests to match new behavior Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Layer 2: Periodic structure snapshots for convergence - TreeSnapshotProvider trait on CacheAwarePolicy provides snapshots directly from radix trees — no raw ops, no pending buffer - Collector sends LZ4-compressed TreeSnapshot every 30 rounds (~30s) via "tree_structure_snapshot" policy_type - Receiver decompresses and calls merge_snapshot() for three-case merge - ~2.2 MB raw → ~300 KB compressed per snapshot - Replaces old Phase 1/2 (tree_ops_pending scan + tree_configs fallback) Layer 3: Hash-based idle skip - Already implemented: tree_generation atomic is only bumped on changes - When idle (no requests), tree_changed = false → zero sync traffic Wire changes: - lz4_flex dependency added to smg-mesh - TreeSnapshotProvider trait + registration on MeshSyncManager - Controller wires provider to collector on stream handler spawn - ping_server handles "tree_structure_snapshot" with LZ4 decompress - TreeStateSubscriber gains apply_structure_snapshot method - CacheAwarePolicy + PolicyRegistry implement snapshot provider + subscriber Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
- Fix Issue 8 (CRITICAL): Snapshot counter was reset to 0 in Phase 0 before Layer 2 checked it, so snapshots never fired. Now uses a separate models_needing_snapshot HashSet passed from Phase 0 to Layer 2. - Fix Issue 1 (CRITICAL): Tenant delta inserts were all no-ops because hash→node lookup was unimplemented. Added truncated node_path (first 2000 chars) to TenantInsert alongside the hash. Receiver uses node_path for insert_text. 2000 chars × 200 rps = 400 KB/round. - Fix Issue 4 (CRITICAL): TreeSnapshot bytes stored in tree_configs caused TreeState::from_bytes() failures. Removed tree_configs.insert from apply_remote_structure_snapshot — snapshot is applied directly to subscribers via merge_snapshot. - Fix Issue 3 (HIGH): Replaced std DefaultHasher (unstable across Rust versions) with xxh3 from xxhash-rust for deterministic cross-version node path hashing. - Fix Issue 6: Layer 2 now filters snapshots to only models in models_needing_snapshot instead of snapshotting all models. - Fix Issue 7: apply_structure_snapshot now delegates to all policies (model, default, prefill, decode) matching apply_tenant_delta pattern. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…leak The ping_server snapshot exchange cloned tree_configs bytes (TreeState format, tens of MB with 80k-char prompts) on every ping round (~1/s). This caused continuous memory growth to 9+ GB even after traffic stopped. Tree data is now synced exclusively via: - Layer 1: tenant deltas every gossip round (~400 KB) - Layer 2: LZ4-compressed TreeSnapshot every 30 rounds (~300 KB) The legacy tree_configs snapshot exchange path is removed. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…ree-sync Applies all changes from commits 878f0c1..0733dc1 (the other engineer's implementation of the two-layer sync design): Layer 1 (every 1s): Tenant deltas with node_path_hash (8 bytes) + path_hash_index for receiver-side resolution. No prompt text on wire. Layer 2 (every 30s): LZ4-compressed TreeSnapshot via TreeSnapshotProvider. TreeSnapshot/from_snapshot/merge_snapshot restored from PR #974. Key differences from our previous approach: - TenantInsert uses ONLY node_path_hash (no truncated node_path text) - Receiver resolves hash via path_hash_index (populated on local inserts) - export_tree_state on TreeStateSubscriber for checkpoint (no tree_ops_pending) - path_hash_index stores matched prefix (~200 chars), not full prompt - blake3 hash instead of xxh3 for cross-platform determinism 792 tests passing (168 kv-index + 174 mesh + 450 gateway). Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
The ping_server snapshot exchange cloned tree_configs bytes on every ping round (~1/s per peer reconnect). With large trees from 200 rps 80k-char traffic, this caused continuous memory growth to 6+ GB on the receiver side even after traffic stopped. Tree data is synced via Layer 1 (tenant deltas) and Layer 2 (periodic compressed snapshots). The legacy tree_configs path in the snapshot exchange is removed. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
checkpoint_tree_states called export_tree_state() every 10 gossip rounds, which walked the entire radix tree and reconstructed TreeState with full prompt text as TreeKey::Text for every leaf. With 2000 entries of 80k chars, this allocated ~160 MB every 10 seconds, causing 2+ GB memory usage that kept growing. Tree data is synced via Layer 1 (tenant deltas, 8-byte hashes) and Layer 2 (LZ4-compressed TreeSnapshot every 30 rounds). checkpoint is no longer needed — tree_configs is not used for tree data transmission. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
C1: checkpoint_tree_states now exports compact TreeSnapshot (shared prefix edges) from subscribers via new export_tree_snapshot trait method, stores in tree_configs. Phase 2 sends as tree_state_lz4 with LZ4 compression. No full prompt text reconstruction. C2: controller.rs incremental stream handler now handles tenant_delta and tree_state_lz4 message types (was falling through to CRDT policy store). Mirrors ping_server.rs logic. C3: select_worker_min_load text path now populates path_hash_index so tenant deltas from the imbalanced-load path are resolvable on receiver side. H3: hash_node_path and hash_token_path return 1 instead of 0 when hash collides with GLOBAL_EVICTION_HASH (0 = global eviction sentinel). Also: TreeState::from_snapshot() conversion for receivers that get TreeSnapshot bytes in tree_state_lz4 payloads. materialize_tree_state handles both TreeState and TreeSnapshot bytes in tree_configs. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…y leak) select_worker_min_load was inserting text.to_string() (the full 80k+ char prompt) into path_hash_index. At 200 rps under load imbalance, this accumulated 16 MB/s until the next eviction clear(). This is the same class of memory leak we fixed everywhere else. Fix: skip path_hash_index population in the min-load path (no match result available to extract the matched prefix). Layer 2 snapshots handle convergence for these entries. Also: version comparison in apply_remote_tree_operation now falls back to the authoritative atomic tree_version counter when tree_configs holds TreeSnapshot bytes (which can't be deserialized as TreeState). Prevents stale remote state from overwriting newer local state. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
checkpoint_tree_states was calling export_tree_snapshot + to_bytes + tree_configs.insert every 10 gossip rounds unconditionally. With a large tree (~50 MB snapshot), this caused continuous memory growth from allocation churn even when no requests were arriving. Add last_checkpoint_gen counter to StateStores. Checkpoint now compares tree_generation against last checkpoint — if unchanged, returns immediately. This eliminates all per-round allocations when idle. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Phase 2 in the collector was reading tree_configs (31 MB TreeSnapshot), LZ4-compressing, serializing, pushing to updates — only to have the controller reject it for exceeding 10 MB. Since mark_sent was never called, last_sent_version stayed at 0 and the cycle repeated every round (~1/s), allocating 31 MB of garbage per second. Fix: check compressed size before adding to updates. If > 8 MB, skip and advance last_sent_version (via write lock upgrade) to stop the retry loop. The oversized snapshot is effectively dropped — Layer 1 tenant deltas handle steady-state sync, and the tree will eventually shrink after eviction, at which point the snapshot will fit. Also: checkpoint_tree_states skips when tree_generation unchanged. GC clears tree_configs to force re-export after eviction. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
… logs
Two changes:
1. sync_insert_hash: new lightweight sync path that accepts a
pre-computed hash + tenant instead of TreeKey::Text(text.to_string()).
The main text routing path (select_worker_with_text) now uses this,
avoiding 80k+ String allocation on every request. At 200 rps that was
16 MB/s of allocator churn causing glibc malloc fragmentation
(RSS grows and never returns to OS).
The token path and min-load path still use sync_insert_operation with
TreeKey since they need different handling.
2. Memory debug instrumentation:
- Controller GC cycle logs: tree_configs entries/bytes, tree_versions,
tenant_delta_inserts/evictions counts, policy/worker CRDT sizes
- Eviction task logs: string_tree/token_tree model counts,
path_hash_index entry count
These fire every 60s (GC) and every eviction_interval_secs (300s
default), providing visibility into memory usage without code changes.
Refs: #983
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Add INFO-level logs at every allocation/lifecycle point: Controller: - spawn_sync_stream_handler: peer name + active connection count - Handler exit: peer name - IncrementalUpdateCollector creation: per-peer - GC cycle: CRDT policy oplog length, sync_connections size, retry count Collector: - Phase 0 summary: tenant delta update count + total bytes - Phase 2 summary: tree_configs update count + total bytes - Oversized snapshot skip: promoted to INFO with structured fields Sync: - sync_tree_insert_hash: every 1000 calls, log total + buffer sizes - sync_tree_operation: every 1000 calls, same - apply_remote_tenant_delta: every 100 calls, log applied/skipped/resolved/unresolved Eviction task: - Tree model counts + path_hash_index size These logs fire periodically (not every request) to minimize overhead while providing full visibility into memory usage patterns. Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
… per cycle The debug log I added called tree.snapshot() every eviction cycle to log node_count and edge_bytes. With 3500+ unique 80k-char prompt leaves, snapshot() cloned 170 MB of edge text into a temporary Vec. Over 10 minutes this caused ~850 MB of allocator fragmentation that glibc never returned to the OS. This was the 1.7 GB vs 270 Mi difference between mesh ON and mesh OFF. The snapshot call only existed in the mesh ON test image (added in snapshot-21), not in the mesh OFF baseline image (snapshot-19). Replace with lightweight model count log (no tree walk, no allocation). Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
checkpoint_tree_states called export_tree_snapshot → tree.snapshot() every 10 gossip rounds. With 3500+ unique 80k-char prompt leaves, this allocated 170 MB of edge text into a Vec<SnapshotNode>, serialized it, and stored in tree_configs — every 10 seconds. The serialized bytes were then read by the collector, LZ4-compressed to 34 MB, found to exceed the 8 MB limit, and skipped. The entire 170 MB allocation was wasted work causing allocator fragmentation. Tree data syncs via Layer 1 (tenant deltas, ~50 bytes per insert). Layer 2 (full tree snapshots) is deferred to future work — options include chunked snapshots or incremental tree diffs. Also removed the tree.snapshot() call from the eviction debug log (same 170 MB allocation issue). Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…nitoring Remove: - Static atomic counters in sync.rs (TREE_INSERT_CALL_COUNT, etc) - Per-1000-call periodic logs in sync_tree_insert_hash/sync_tree_operation - Per-100-call periodic logs in apply_remote_tenant_delta - sync_connections/retry count logs in controller GC cycle Downgrade to DEBUG: - spawn_sync_stream_handler/collector created/handler exited logs - Phase 0/Phase 2 summary logs in incremental.rs (fire every round) Keep at INFO (fires every 60s, useful for ops): - Mesh memory: tree_configs, tree_versions, tenant_inserts, policy/worker CRDT - CRDT policy oplog length - Tree memory: model counts, path_hash_index size Refs: #983 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Summary
"tree_state"are converted to snapshot on receiveWhat changed
crates/mesh/src/sync.rsTreeStateSubscribertrait:apply_remote_tree_state(&TreeState)→apply_remote_tree_snapshot(&TreeSnapshot)apply_remote_tree_snapshot(): stores snapshot bytes intree_configswith atomic version dedupapply_remote_tree_state_compat(): backward compat — converts old TreeState → Tree → TreeSnapshotmaterialize_tree_snapshot(): reconstructs Tree from stored snapshot + pending opscheckpoint_tree_states(): serializes as TreeSnapshot instead of TreeStatesnapshot_to_tree_state(): backward compat conversion forget_tree_state()callerscrates/mesh/src/incremental.rs"tree_snapshot"policy_typecrates/mesh/src/ping_server.rs+controller.rs"tree_snapshot"(new) +"tree_state"(backward compat)tree_configs"tree_snapshot"policy_type with version from atomic mapmodel_gateway/src/policies/cache_aware.rsapply_remote_tree_snapshot()callstree.merge_snapshot()(3-case merge algorithm)restore_tree_state_from_mesh()usesget_all_tree_snapshots()with model_idmodel_gateway/src/policies/registry.rsPolicyRegistryimplements newapply_remote_tree_snapshottrait methodWhy
With 20k-char prompts and 2048 cached entries, TreeState serialized to ~40 MB (each insert stored the full prompt text). This exceeded the 10 MB gRPC message limit, causing full-state sync to fail every cycle. TreeSnapshot preserves the radix tree structure where shared prefixes are stored once, reducing wire size by 10-20x.
Design decisions
"tree_state"handled by converting to snapshot on receive. New peers send"tree_snapshot".TreeStateDeltastill usesVec<TreeOperation>(deltas are small).TreeSnapshothas no embedded version. Versions live in atomictree_versionsDashMap, read inside entry shard lock to avoid TOCTOU.Test Plan
cargo test -p kv-index --lib— 168 passedcargo test -p smg-mesh --lib— 177 passedcargo test -p smg --lib— 450 passedcargo test -p smg --test mesh_integration_test— 8 passedcargo clippy --all-targets --all-features -- -D warnings— cleancargo +nightly fmt --all— cleanradix_tree_benchmark) will run on this PRChecklist
Summary by CodeRabbit
New Features
Performance
Compatibility
Chores