feat(mesh): wire TreeSyncAdapter so cache-aware state syncs across nodes - #2119
Conversation
…te syncs across nodes Extend MeshAdapters::start to register the td:/tree:req:/tree:page: stream namespaces, construct a TreeSyncAdapter, and attach it to the running PolicyRegistry. The registry propagates the adapter (and the paired populate_hash_index flag) to every existing and future CacheAwarePolicy, matching the pattern already used for KvEventMonitor / LoadReceiver. CacheAwarePolicy::select_worker_with_tokens and select_worker_with_text publish a TreeDelta immediately after every hash_index write, keyed by the same hash_token_path(tokens) / hash_node_path(text) — so peers that repair against us land on the same tree node. CacheAwarePolicy::set_mesh_tree_sync is a single atomic setter that attaches the adapter AND flips populate_hash_index; the two fields cannot drift apart. PolicyRegistryTreeHandle dispatches per-model_id to whichever CacheAwarePolicy the registry has for that model, logging at debug when this node is not authoritative for a model so operators can distinguish legitimate empty results from repair storms. ClusterStatePeerList reads ClusterState and filters to alive non-self peers for TreeSyncAdapter's repair fan-out. Closes smg-project#1578. Signed-off-by: purp1e-ace <45795655+purp1e-ace@users.noreply.github.com>
ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (5)
📝 WalkthroughSummary by CodeRabbit
WalkthroughThe change wires ChangesMesh tree synchronization
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant Server
participant MeshAdapters
participant TreeSyncAdapter
participant PolicyRegistry
participant CacheAwarePolicy
Server->>MeshAdapters: start with ClusterState and PolicyRegistry
MeshAdapters->>TreeSyncAdapter: configure streams and start
MeshAdapters->>PolicyRegistry: attach TreeSyncAdapter
PolicyRegistry->>CacheAwarePolicy: set mesh tree synchronization
CacheAwarePolicy->>TreeSyncAdapter: publish TreeDelta
Possibly related PRs
Suggested labels: Suggested reviewers: 🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 2
🧹 Nitpick comments (4)
model_gateway/src/mesh/wiring.rs (3)
246-275: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win🟡 Nit: the
TreeHandlebridge has no direct test.
apply_known_remote_insert,open_repair_stream, andapply_repair_pageare new production code. The added tests cover registration, propagation, detach, peer filtering, and the outbound producer hook, but no test drives an inbound path throughPolicyRegistryTreeHandle.A test that registers a cache-aware policy, seeds the hash index, and calls
apply_known_remote_insertthrough the handle would pin the fallback contract (false/None/0) and would catch the resolution gap described in the comment on Lines 212-236.As per coding guidelines: "Run the pr-test-analyzer agent to verify that tests adequately cover new or changed functionality."
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@model_gateway/src/mesh/wiring.rs` around lines 246 - 275, Add a direct test for the PolicyRegistryTreeHandle bridge that registers a cache-aware policy, seeds the hash index, and invokes apply_known_remote_insert through the TreeHandle interface. Assert the delegated result and verify the fallback values false, None, and 0 for unresolved models across apply_known_remote_insert, open_repair_stream, and apply_repair_page, covering the resolution behavior in with_cache_aware.Source: Coding guidelines
365-374: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value🟡 Nit: this policy config duplicates the
cache_aware_policy_confighelper.The helper at Lines 455-466 defines the identical
PolicyConfig::CacheAwareliteral. Reuse it here so a threshold change updates one place.♻️ Proposed fix
- let policy_registry = Arc::new(PolicyRegistry::new(PolicyConfig::CacheAware { - cache_threshold: 0.5, - balance_abs_threshold: 32, - balance_rel_threshold: 1.5, - eviction_interval_secs: 60, - max_tree_size: 128, - block_size: 16, - balance_token_usage_threshold: 1.0, - overload_token_usage_threshold: 1.0, - })); + let policy_registry = Arc::new(PolicyRegistry::new(cache_aware_policy_config()));🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@model_gateway/src/mesh/wiring.rs` around lines 365 - 374, Replace the duplicated PolicyConfig::CacheAware literal in the policy_registry initialization with the existing cache_aware_policy_config helper, preserving the current PolicyRegistry construction while centralizing these threshold values.
29-44: 🚀 Performance & Scalability | 🔵 Trivial🟡 Nit: the three stream buffers are a fixed, unconfigurable memory commitment.
TD_BUFFER_BYTES(4 MiB),REPAIR_REQ_BUFFER_BYTES(256 KiB), andREPAIR_PAGE_BUFFER_BYTES(32 MiB) total roughly 36 MiB per node, reserved whenever mesh is enabled. The doc comments explain the sizing intent well, but an operator running many small gateway replicas cannot lower it, and an operator with large trees cannot raise it.Consider exposing these through the mesh server config, and emit a gauge for buffer occupancy so FIFO eviction on the
td:namespace (which silently degrades to the repair path) is observable.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@model_gateway/src/mesh/wiring.rs` around lines 29 - 44, Make TD_BUFFER_BYTES, REPAIR_REQ_BUFFER_BYTES, and REPAIR_PAGE_BUFFER_BYTES configurable through the mesh server configuration instead of fixed constants, preserving their current values as defaults. Use the configured values when initializing the corresponding td:, tree:req:, and tree:page: buffers, and add an occupancy gauge for each stream so FIFO eviction and buffer usage are observable.model_gateway/src/policies/cache_aware.rs (1)
279-286: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value🟡 Nit: the adapter store and the populate flag are not updated atomically.
The doc comment states both fields flip "in one atomic step". The write lock is released at Line 284 before the
storeat Line 285. Two concurrent callers (attach and detach) can interleave and leavemesh_tree_sync == Nonewithpopulate_hash_index == true. The hash index would then grow with no reader, which is the exact failure the flag prevents.Today only mesh startup calls this, so the window is not reachable. Either hold the write guard across the
store, or soften the doc comment to state the caller must serialize.♻️ Proposed fix
pub fn set_mesh_tree_sync(&self, adapter: Option<Arc<TreeSyncAdapter>>) { let populate = adapter.is_some(); - { - let mut guard = self.mesh_tree_sync.write(); - *guard = adapter; - } - self.populate_hash_index.store(populate, Ordering::Relaxed); + let mut guard = self.mesh_tree_sync.write(); + *guard = adapter; + // Store under the guard so no observer sees the pair split. + self.populate_hash_index.store(populate, Ordering::Relaxed); }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@model_gateway/src/policies/cache_aware.rs` around lines 279 - 286, Update set_mesh_tree_sync so the mesh_tree_sync write guard remains held while populate_hash_index is updated, making both fields change within the same synchronization window. Preserve the existing adapter presence logic and atomic flag ordering.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@model_gateway/src/mesh/wiring.rs`:
- Around line 212-236: Update with_cache_aware and inbound tree-delta handling
to resolve each delta through the policy owning its dispatch leg, including
PD/EPD prefill_policy, decode_policy, encode_policy, and the default policy when
no model entry exists, rather than consulting only model_policies. Normalize
model IDs before on_worker_added stores them so empty IDs become
UNKNOWN_MODEL_ID and match producer-published deltas.
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 1116-1124: Replace the hardcoded epoch: 0 values in both
sync_local_insert producer sites with a monotonic per-policy epoch counter,
ensuring each emitted TreeDelta receives the correct intra-batch ordering value
for both token and string paths. Update the surrounding policy state and call
sites as needed so the counter is shared and incremented consistently.
---
Nitpick comments:
In `@model_gateway/src/mesh/wiring.rs`:
- Around line 246-275: Add a direct test for the PolicyRegistryTreeHandle bridge
that registers a cache-aware policy, seeds the hash index, and invokes
apply_known_remote_insert through the TreeHandle interface. Assert the delegated
result and verify the fallback values false, None, and 0 for unresolved models
across apply_known_remote_insert, open_repair_stream, and apply_repair_page,
covering the resolution behavior in with_cache_aware.
- Around line 365-374: Replace the duplicated PolicyConfig::CacheAware literal
in the policy_registry initialization with the existing
cache_aware_policy_config helper, preserving the current PolicyRegistry
construction while centralizing these threshold values.
- Around line 29-44: Make TD_BUFFER_BYTES, REPAIR_REQ_BUFFER_BYTES, and
REPAIR_PAGE_BUFFER_BYTES configurable through the mesh server configuration
instead of fixed constants, preserving their current values as defaults. Use the
configured values when initializing the corresponding td:, tree:req:, and
tree:page: buffers, and add an occupancy gauge for each stream so FIFO eviction
and buffer usage are observable.
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 279-286: Update set_mesh_tree_sync so the mesh_tree_sync write
guard remains held while populate_hash_index is updated, making both fields
change within the same synchronization window. Preserve the existing adapter
presence logic and atomic flag ordering.
🪄 Autofix
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: CHILL
Plan: Pro Plus
Run ID: 5816f8d2-6c92-40ad-b401-d7e04876fe20
📒 Files selected for processing (5)
model_gateway/src/mesh/adapters/tree_sync.rsmodel_gateway/src/mesh/wiring.rsmodel_gateway/src/policies/cache_aware.rsmodel_gateway/src/policies/registry.rsmodel_gateway/src/server.rs
Broaden inbound resolution scope. `PolicyRegistryTreeHandle` now walks every distinct cache-aware policy the registry could dispatch requests through for a model (per-model → default → PD/EPD legs), deduplicated by Arc identity. Previously the bridge consulted only `model_policies`, so PD/EPD deployments and non-PD deployments that resolve `model_id` through the default policy never saw inbound deltas. Tighten the atomic pairing on `CacheAwarePolicy::set_mesh_tree_sync`. The `populate_hash_index` store now happens under the same `mesh_tree_sync` write guard as the adapter swap, so no observer can see the pair split. Clarify `TreeDelta.epoch` as a reserved slot. The receiver does not consult it — only `trace!`/`debug!` log lines read it — so hardcoding `0` at both producer sites is documented rather than a latent ordering bug. If a consumer ever wants intra-batch ordering, updating the doc + adding a counter is one change. Test additions in `mesh::wiring::tests`: - `bridge_returns_fallbacks_when_no_cache_aware_policy` — asserts `false` / `None` / `0` fallbacks when the chain has no cache-aware policy (non-cache-aware default, no per-model entries, no PD legs). - `bridge_dispatches_through_default_policy` — asserts the bridge reaches the default policy for models without a per-model entry, which is the single-model non-PD deployment shape. Also folded a duplicated `PolicyConfig::CacheAware` literal into the existing `cache_aware_policy_config` helper. Signed-off-by: purp1e-ace <45795655+purp1e-ace@users.noreply.github.com>
…x/1578-wire-tree-sync
Description
Problem
Cache-aware policy tree state does not synchronize across mesh nodes even
when
--enable-meshis on. Router B starts with an empty routing tree afterRouter A fails, losing all cache affinity.
path_hash_indexmetrics return-1; observed routing consistency after failover is-100%across all sevenwarm-up strategies tested in #1578.
Root cause: the pieces are all present —
TreeSyncAdapterimplementation,CacheAwarePolicy::apply_known_remote_insert/apply_repair_page/populate_hash_indexflag, hash-index generators — but the composition rootMeshAdapters::start(model_gateway/src/mesh/wiring.rs) only wires theworker:andrl:CRDT namespaces. Thetd:/tree:req:/tree:page:stream namespaces are never registered, no
TreeSyncAdapteris everconstructed, no producer-side hook fires, and
populate_hash_indexstaysfalseat its default.Closes #1578.
Solution
Extend
MeshAdapters::startto register the three tree-sync streamnamespaces, construct a
TreeSyncAdapter, and attach it to the runningPolicyRegistry. The registry propagates the adapter (and flipspopulate_hash_indexon) to every existing and futureCacheAwarePolicy,matching the pattern already used for
KvEventMonitor/LoadReceiver.Producer side:
CacheAwarePolicy::select_worker_with_tokensandselect_worker_with_textpublish aTreeDeltaright after everyhash_indexwrite, using the samehash_token_path(tokens)/hash_node_path(text)keys — so peers that request repair land on the sametree node.
Consumer side:
PolicyRegistryTreeHandleimplementsTreeHandleanddispatches per-
model_idto whicheverCacheAwarePolicythe registry hasfor that model.
ClusterStatePeerListimplementsPeerListby readingClusterStateand filtering to alive non-self peers.Changes
model_gateway/src/mesh/wiring.rs—MeshAdapters::startgrows twoparameters (
ClusterState,Arc<PolicyRegistry>); registerstd:(broadcast),
tree:req:(targeted),tree:page:(targeted) streamnamespaces; constructs
TreeSyncAdapterwithPolicyRegistryTreeHandle+ClusterStatePeerList; starts it before flippingset_mesh_tree_synconthe registry (order matters — see comment).
model_gateway/src/policies/cache_aware.rs— addsmesh_tree_syncfield and
set_mesh_tree_syncsetter that atomically attaches theadapter AND flips
populate_hash_index(single call, one write lock,one atomic store) so the paired invariant cannot drift; adds
sync_local_inserthelper (clones the adapter Arc out of the readguard before invoking
on_local_insert, so a future adapter path cannever deadlock against the policy read lock); wires the producer hook
into both
select_worker_*sites right after thehash_indexwrite.model_gateway/src/policies/registry.rs— addsmesh_tree_syncfield,set_mesh_tree_sync(propagates to default + PD + all model policies),maybe_inject_mesh_tree_synchelper, and injection increate_policy_from_typeso lazily-created per-model policies inheritthe adapter. The registry uses the single atomic setter throughout;
the standalone
set_populate_hash_indexonCacheAwarePolicyisgated
#[cfg(test)]because the pair now moves together inproduction.
model_gateway/src/mesh/adapters/tree_sync.rs— makes the three prefixconstants
pub(crate)so wiring can reference them; adds a#[cfg(test)] pub(crate) fn pending_delta_count_for_test(...)viewonto the outbound buffer so wiring integration tests can assert the
producer hook fired without waiting a gossip tick for drain.
model_gateway/src/server.rs— passeshandler.stateandapp_context.policy_registrytoMeshAdapters::start.The wiring flips
populate_hash_indexand the outbound adapter together —individually flipping only one would either OOM the gateway (index writes
with no readers, mesh off) or publish deltas the local node cannot resolve
on repair (adapter attached but index empty, mesh on).
Test Plan
Unit tests, added to
model_gateway/src/mesh/wiring.rs:start_wires_worker_inbound_end_to_end— existing test, extended withthe new
MeshAdapters::startsignature; confirmsworker:CRDT inboundloop still runs.
rl_namespace_uses_epoch_max_wins— existing, extended.tree_adapter_registers_and_flips_populate_flag— new. ConstructsMeshAdapterswith aPolicyRegistrydefaulted tocache_aware, thenfetches a per-model policy via
get_policy_or_default. Asserts thatpopulate_hash_indexreadstrueon the fetched policy — proves theadapter is attached and the flag propagates to lazily-created policies.
start_propagates_to_preexisting_model_policies— new. Creates aper-model cache-aware policy BEFORE wiring runs (so it lives in
model_policiesat attach time), then starts wiring and asserts thepopulate flag flipped on the same live policy Arc. Regression guard
for the propagation path that walks
model_policiesinset_mesh_tree_sync.detach_clears_populate_flag— new. After attach, callsset_mesh_tree_sync(None)and asserts the populate flag goes backoff. Guards the paired invariant that the atomic setter undoes both
flips together.
producer_hook_publishes_delta_on_select_worker— new. End-to-endregression guard for [Bug]: Mesh Cache-Aware Routing Tree Sync Mechanism Is Not Wired — Router B Never Receives Router A's Routing Tree #1578: builds a full
MeshAdapters, drives onestring request and one token request through
CacheAwarePolicy::select_worker, and asserts the outboundpending_deltasbuffer on theTreeSyncAdaptergrew for each. Deletingeither
sync_local_insertcall site fails this test.peer_list_reports_alive_only— new. Populates aClusterStatewithAlive + Alive + Down entries and asserts
alive_peers()excludes selfand the Down node.
start_panics_on_second_call— existing, extended.start_panics_on_colon_node_name— existing, extended.Manual reproduction of the #1578 scenario across two mesh nodes is NOT part
of this PR's test evidence — I don't have the multi-node cluster to run
that on. If the pre-merge reviewer wants that, I can coordinate a
follow-up.
Full run:
Note on
monitor.rs:744— a pre-existingclippy::unneeded_wildcard_patternlint appears on upstream
mainatc2cf59c7. I confirmed this by stashingmy changes and re-running clippy; the error persists. Out of scope for this
PR per the "one concern per PR" rule.
Checklist
cargo +nightly fmtpassescargo clippy -p smg --lib --tests -- -D warningspasses on thisdiff (a pre-existing
unneeded_wildcard_patternonmonitor.rs:744remains — see note in Test Plan; workspace
--all-featuresskippedbecause system OpenCV is not installed locally, per
CONTRIBUTING.mdfallback)