Skip to content

feat(smg): MeshAdapters composition root — register CRDT engines before gossip, start inbound sync - #1626

Merged
slin1237 merged 6 commits into
mainfrom
chang/mesh-adapters-wiring
Jun 10, 2026
Merged

slin1237 merged 6 commits into
mainfrom
chang/mesh-adapters-wiring

Conversation

@CatherineSue

@CatherineSue CatherineSue commented Jun 10, 2026 •

Copy link
Copy Markdown
Member

Description

Problem

No mesh sync adapter runs in production: every *Adapter::new call is #[cfg(test)]-only, so the CRDT gossip machinery (#1570/#1586/#1592) moves data between mesh nodes but nothing in the gateway produces or consumes it. The worker:/rl: CRDT namespaces are never registered either — and server.rs spawns mesh_server.start() at build time, so even registering them later would race gossip: a remote op arriving for an unregistered prefix merges through the default last-writer-wins engine with the wrong semantics (rl: needs epoch-max-wins).

Solution

Add MeshAdapters — a thin composition root (mesh/wiring.rs) that registers the worker: (LWW) and rl: (EpochMaxWins) namespaces, constructs the existing adapters against the gateway's registries, and starts their inbound sync loops in one call. server.rs builds the mesh server without starting it, calls MeshAdapters::start once the worker registry exists, and only then spawns gossip — engine registration always precedes remote traffic.

The detail adapters (WorkerSyncAdapter, RateLimitSyncAdapter, TreeSyncAdapter) are unchanged — this PR only calls their existing public APIs. Mesh-off is represented by the absence of the whole struct (one Option at the AppState level, no per-adapter fields).

Changes

  • model_gateway/src/mesh/wiring.rs (new): MeshAdapters::start(mesh_kv, node_name, worker_registry) + worker() / rate_limit() accessors (the seams for the follow-up outbound-publish and rate-limit-middleware PRs). Single start (no separate new) because the adapters' start() methods are not idempotent and duplicate prefix registration panics — double-construction fails fast at startup.
  • model_gateway/src/server.rs: mesh build block no longer spawns; new block after app_context creation constructs MeshAdapters then spawns gossip; AppState.mesh_adapters: Option<Arc<MeshAdapters>>; stale "not yet started from server.rs" comment updated.
  • Test/bench AppState literals gain mesh_adapters: None.

Behavior changes

  • Inbound worker sync goes live: remote worker: states now register workers on this node (inert in practice until the outbound PR makes nodes publish).
  • A mesh self_name containing : now panics at startup (the rate-limit shard-key invariant surfacing at wiring time) — fail-fast by design.

Test Plan

  • cargo test -p smg — full suite green (mesh-off path: mesh_server_config = None ⇒ identical behavior to today).
  • cargo clippy -p smg --all-targets --all-features -- -D warnings — clean.

New tests in wiring.rs:

  • start_wires_worker_inbound_end_to_end — put through the adapter echoes back through the registered namespace into the WorkerRegistry (proves registration + live inbound loop).
  • rl_namespace_uses_epoch_max_wins — sync_counter(epoch=2, 5) then (epoch=1, 100) aggregates to 5; a mis-registered LWW engine would give 100 (observably pins the merge strategy).
  • start_panics_on_second_call / start_panics_on_colon_node_name — fail-fast invariants.

Follow-up PR (outbound): registry WorkerOrigin tracking, mesh-id adoption, echo suppression, WorkerEvent → mesh publish loop, inbound tombstone → remove.

Checklist
  • cargo +nightly fmt passes
  • cargo clippy --all-targets --all-features -- -D warnings passes
  • (Optional) Documentation updated
  • (Optional) Please join us on Slack #sig-smg to discuss, review, and merge PRs

🤖 Generated with Claude Code

Summary by CodeRabbit

  • Chores

    • New mesh adapters initialization integrated into gateway startup; mesh server spawn is deferred until adapters are started. App state now holds an optional mesh adapters handle.
  • Bug Fixes

    • Mesh server name validation now rejects empty or invalid names (e.g., containing ':') during config build.
  • Tests

    • Updated tests and benchmarks to populate the new state field and validate adapter and config behavior.

Single owner that registers the worker: (last-writer-wins) and rl:
(epoch-max-wins) CRDT namespaces, constructs the existing sync adapters,
and starts their inbound loops in one call. Mesh-off is the absence of
the whole struct, so per-adapter Option fields never leak into server
state. The adapters themselves are unchanged.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
…Adapters

Previously server.rs spawned `mesh_server.start()` at build time and never
registered the `worker:`/`rl:` CRDT namespaces at all, so no adapter ran in
production. Registering engines after gossip is live is also a race: a
remote op arriving for an unregistered prefix merges through the default
last-writer-wins engine with the wrong semantics (`rl:` needs
epoch-max-wins).

Build the mesh server without starting it, construct MeshAdapters (which
registers both namespaces and starts the inbound sync loops) once the
worker registry exists, and only then spawn gossip. AppState carries the
adapters as a single optional handle for the follow-up outbound and
rate-limit middleware wiring; mesh-off behavior is unchanged.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
@CatherineSue
CatherineSue requested a review from slin1237 as a code owner June 10, 2026 04:17
@coderabbitai

coderabbitai Bot commented Jun 10, 2026 •

Copy link
Copy Markdown

Review Change Stack

Note

Reviews paused

It 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 reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: 92432f71-561e-48ee-8d64-0c654ba43c85

📥 Commits

Reviewing files that changed from the base of the PR and between 2e0997a and 344e0e8.

📒 Files selected for processing (3)
  • model_gateway/src/config/mod.rs
  • model_gateway/src/config/validation.rs
  • model_gateway/src/main.rs

📝 Walkthrough

Walkthrough

Adds MeshAdapters composition root and re-exports it from mesh; updates server startup to start adapters before spawning the mesh server and stores the adapters handle in AppState; updates tests and a benchmark to initialize the new AppState field; adds mesh-server-name config validation and tests.

Changes

Mesh Adapters Wiring and Server Integration

Layer / File(s) Summary
MeshAdapters composition root and merge-strategy tests
model_gateway/src/mesh/wiring.rs
Introduces MeshAdapters that builds and starts WorkerSyncAdapter (namespace worker: with LastWriterWins) and RateLimitSyncAdapter (namespace rl: with EpochMaxWins), exposes accessors, and adds Tokio tests for inbound propagation, epoch-max semantics, and panic cases.
Module exports for MeshAdapters
model_gateway/src/mesh/mod.rs
Adds public wiring module and re-exports MeshAdapters from it.
Server startup integration with explicit initialization order
model_gateway/src/server.rs
Imports MeshAdapters, adds AppState.mesh_adapters: Option<Arc<MeshAdapters>>, refactors mesh setup to return (mesh_server, mesh_handler), starts MeshAdapters::start(...) before spawning mesh_server.start(), and stores the started adapters in AppState.
Test fixture updates for AppState mesh_adapters field
model_gateway/tests/common/test_app.rs, model_gateway/tests/wasm_test.rs, model_gateway/benches/wasm_middleware_latency.rs
Initializes mesh_adapters: None in AppState struct literals used by test helpers, a wasm test, and a benchmark to reflect the updated AppState shape.
Mesh server-name config validation and tests
model_gateway/src/config/validation.rs, model_gateway/src/config/mod.rs, model_gateway/src/main.rs
Adds validate_mesh_server_name, re-exports it, calls it during mesh-server-name config construction, and adds unit tests verifying invalid/valid names.

Sequence Diagram

sequenceDiagram
    participant StartupFn as startup()
    participant MeshHandler as mesh_handler
    participant MeshAdapters
    participant WorkerRegistry
    participant MeshServer as mesh_server

    StartupFn->>MeshHandler: build mesh_server, mesh_handler
    MeshHandler-->>StartupFn: (mesh_server, mesh_handler)
    StartupFn->>StartupFn: build AppContext
    alt mesh_handler exists
        StartupFn->>MeshAdapters: start(mesh_kv, node_name, worker_registry)
        MeshAdapters->>MeshAdapters: configure worker: LastWriterWins
        MeshAdapters->>MeshAdapters: configure rl: EpochMaxWins
        MeshAdapters->>WorkerRegistry: register inbound sync loop
        MeshAdapters-->>StartupFn: Arc~MeshAdapters~
    end
    alt mesh_server exists
        StartupFn->>MeshServer: spawn mesh_server.start()
    end
Loading

Estimated Code Review Effort

🎯 4 (Complex) | ⏱️ ~45 minutes

Possibly Related PRs

  • lightseekorg/smg#1327: Both PRs introduce and expose the RateLimitSyncAdapter for the v2 rl: CRDT namespace.
  • lightseekorg/smg#1313: Related WorkerSyncAdapter work that MeshAdapters now constructs and starts.
  • lightseekorg/smg#1164: Changes to MeshKV and MergeStrategy APIs that MeshAdapters uses for namespace configuration.

Suggested Labels

mesh

Suggested Reviewers

  • slin1237
  • claude
  • tonyluj
  • llfl

Poem

🐇 I stitched the mesh with nimble paws,
Workers whisper without a pause,
Epochs climb and writers win,
Adapters hum; the syncs begin.
Hoppity hops—tests pass applause!

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title accurately describes the main change: introducing MeshAdapters as a composition root that registers CRDT engines and starts inbound sync loops before gossip.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch chang/mesh-adapters-wiring

Comment @coderabbitai help to get the list of available commands and usage tips.

@github-actions github-actions Bot added benchmarks Benchmark changes tests Test changes model-gateway Model gateway crate changes labels Jun 10, 2026

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces the composition root for the gateway's mesh sync adapters via the new MeshAdapters struct, wiring up the inbound sync loops for worker and rate-limit adapters during server startup. It also ensures that CRDT namespaces are registered before starting the mesh server gossip to avoid incorrect merge semantics. The review feedback suggests simplifying the MeshAdapters struct by deriving Debug directly, and increasing the retry duration in the end-to-end integration test to prevent potential flakiness in CI environments.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread model_gateway/src/mesh/wiring.rs Outdated
Comment thread model_gateway/src/mesh/wiring.rs Outdated

@claude claude Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Clean, focused composition root. The build→register→start ordering in server.rs is correct — CRDT namespaces are registered before gossip starts, preventing wrong-merge-strategy issues. Tests cover the key invariants (e2e inbound sync, epoch-max-wins strategy, fail-fast on double-call and invalid node names). No issues found.

Both fields implement Debug so the manual impl was equivalent to the
derive. The inbound e2e test polls up to 1s instead of 200ms to absorb
slow CI schedulers; success still returns on the first hit.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 5323d3ee02

ℹ️ 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 rl_ns = mesh_kv.configure_crdt_prefix("rl:", MergeStrategy::EpochMaxWins);
let worker = WorkerSyncAdapter::new(worker_ns, worker_registry);
let rate_limit = RateLimitSyncAdapter::new(rl_ns, node_name);
worker.start();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Keep unhealthy mesh-imported workers unroutable

When a peer publishes a worker: state with health = false and its serialized WorkerSpec has health.disable_health_check = true, starting this adapter imports that state through WorkerRegistry::on_remote_worker_state. That path builds from the spec, BasicWorkerBuilder defaults disabled-health-check workers to Ready, and it only overrides the status when state.health is true, so an explicitly unhealthy remote worker becomes routable. Please force a non-ready status for false-health imports before registering, or skip those imports.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Valid — confirmed: BasicWorkerBuilder defaults disable_health_check workers to Ready and on_remote_worker_state never forces un-ready on health=false. Deferring the fix to the outbound-sync follow-up PR, which rewrites on_remote_worker_state anyway (origin tracking + id adoption) and will carry the else-branch fix plus a regression test. The path is unreachable until that PR lands: nothing publishes worker: states in production yet, so no unhealthy import can occur on this PR alone.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
model_gateway/src/mesh/wiring.rs (1)

100-105: 🧹 Nitpick | 🔵 Trivial | 💤 Low value

Consider extracting the retry timeout as a constant.

The timeout increase (20 → 100 iterations, 1s total) is reasonable for async processing in test environments. For maintainability, consider extracting const TEST_POLL_TIMEOUT_MS: u64 = 1000; and const TEST_POLL_INTERVAL_MS: u64 = 10; so the retry count is computed as TIMEOUT / INTERVAL.

🤖 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 100 - 105, Extract the
hardcoded retry loop into named constants and compute the iteration count from
them: define const TEST_POLL_TIMEOUT_MS: u64 = 1000 and const
TEST_POLL_INTERVAL_MS: u64 = 10 (or similar names), replace the literal 100 with
(TEST_POLL_TIMEOUT_MS / TEST_POLL_INTERVAL_MS) and replace
Duration::from_millis(10) with Duration::from_millis(TEST_POLL_INTERVAL_MS) in
the async retry around registry.get_by_url("http://remote:8080") so the timeout
and interval are configurable and the intent is clear.
🤖 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.

Outside diff comments:
In `@model_gateway/src/mesh/wiring.rs`:
- Around line 100-105: Extract the hardcoded retry loop into named constants and
compute the iteration count from them: define const TEST_POLL_TIMEOUT_MS: u64 =
1000 and const TEST_POLL_INTERVAL_MS: u64 = 10 (or similar names), replace the
literal 100 with (TEST_POLL_TIMEOUT_MS / TEST_POLL_INTERVAL_MS) and replace
Duration::from_millis(10) with Duration::from_millis(TEST_POLL_INTERVAL_MS) in
the async retry around registry.get_by_url("http://remote:8080") so the timeout
and interval are configurable and the intent is clear.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: ac608993-8daf-447a-935c-40a29d20f735

📥 Commits

Reviewing files that changed from the base of the PR and between d62e3cf and 5323d3e.

📒 Files selected for processing (1)
  • model_gateway/src/mesh/wiring.rs

Comments referencing follow-up work go stale the moment that work lands
and are easy to forget to update; describe only what exists.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: cea5619dc0

ℹ️ 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".

Comment thread model_gateway/src/server.rs
A user-supplied --mesh-server-name that is empty or contains ':' (the
rate-limit shard-key separator) previously parsed fine and then panicked
at adapter construction during startup. Reject it in
build_mesh_server_config with a ConfigError instead; the adapter assert
stays as a backstop.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Bin entrypoints stay test-free: the check moves to a pub
validate_mesh_server_name helper tested beside the other config
validators; main.rs keeps a one-line call.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

benchmarks Benchmark changes model-gateway Model gateway crate changes tests Test changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants