refactor(mesh): extract NamespaceCrdtEngine + LwwEngine + transitional EpochMaxWins - #1539
Conversation
📝 WalkthroughWalkthroughCrdtOrMap is refactored into a prefix-key router that delegates per-key operations to pluggable NamespaceCrdtEngine implementations. Two engines are added: LwwEngine (default) and EpochMaxWinsLegacyEngine. Aggregate queries, GC, export, and merge operate by composing per-engine results and bucketizing incoming ops by longest-prefix-match. ChangesEngine Router Architecture
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related issues
Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 77327d171c
ℹ️ 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".
There was a problem hiding this comment.
Code Review
This pull request refactors CrdtOrMap into a router that delegates operations to prefix-specific engines via a new NamespaceCrdtEngine trait, isolating strategy-specific logic to prevent regression bugs. Review feedback identifies an issue with engine prioritization in all_engines that could lead to data duplication in all() and keys(), and suggests using HashSet for unique key collection. Reviewers also recommended performance optimizations, such as zero-copy operation handling and more efficient state retrieval methods. Additionally, a discrepancy between the put_local documentation and its implementation was noted, along with a potential violation of the monotonic property in the generation counter.
| pub fn all(&self) -> BTreeMap<String, Vec<u8>> { | ||
| let mut all = BTreeMap::new(); | ||
| for engine in self.all_engines() { | ||
| for key in engine.keys() { | ||
| if let Some(value) = engine.get(&key) { | ||
| all.insert(key, value); | ||
| } | ||
| } | ||
| } | ||
| all | ||
| } |
There was a problem hiding this comment.
The implementation of all() is inefficient as it performs a full scan of keys via engine.keys() followed by individual engine.get() calls for each key. This results in all() method to the NamespaceCrdtEngine trait to allow engines to return their full state more efficiently (e.g., by iterating over their internal KvStore directly).
There was a problem hiding this comment.
Declining. all() is called from tests and an HTTP admin endpoint that hand-formats responses; N is small in practice (~tens of keys per engine, low call frequency). Adding all() to the engine trait makes the trait surface bigger now for a perf win that doesn't matter at SMG scale. Worth revisiting if a future caller iterates the full map on the hot path; deferring for now.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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 `@crates/mesh/src/crdt_kv/engine/mod.rs`:
- Around line 35-37: The put_local trait docs currently state rejected writes
return None but actual engines return the current live value on rejection;
update the doc for the put_local method in the engine trait to reflect the real
contract: specify that put_local returns Some(previous_live_bytes) when the
write is accepted, returns Some(current_live_bytes) when the write is rejected
(i.e., the value that prevented the write), and only returns None for cases with
no well-defined previous value (e.g., per-point shard updates); reference the
put_local method in the engine trait and ensure wording matches existing engine
implementations or else adjust implementations to match the updated contract.
🪄 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: c6e8f953-4b70-44cd-a112-e33ab3d025a3
📒 Files selected for processing (7)
crates/mesh/src/crdt_kv/crdt.rscrates/mesh/src/crdt_kv/engine/epoch_max_wins.rscrates/mesh/src/crdt_kv/engine/lww.rscrates/mesh/src/crdt_kv/engine/mod.rscrates/mesh/src/crdt_kv/kv_store.rscrates/mesh/src/crdt_kv/mod.rscrates/mesh/src/crdt_kv/operation.rs
💤 Files with no reviewable changes (1)
- crates/mesh/src/crdt_kv/kv_store.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: ebd9875fe1
ℹ️ 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".
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/crdt_kv/engine/mod.rs (1)
68-76:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winRelax the
apply_remote_opscontract to match current engines.The trait says engines “apply only the post-canonicalisation result to live state”, but
crates/mesh/src/crdt_kv/engine/epoch_max_wins.rsLines 373-404 intentionally replays the raw incoming batch after compaction. Please document the required end state instead of prescribing a replay algorithm the current impl does not follow.🤖 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 `@crates/mesh/src/crdt_kv/engine/mod.rs` around lines 68 - 76, The trait doc for apply_remote_ops is too prescriptive about replaying only the post-canonicalisation result; update the documentation on fn apply_remote_ops(Operation) in engine/mod.rs to instead state the required end-state contract (the engine must merge the incoming ops into its operation log, perform canonicalisation/compaction/tombstone collapse, and ensure live state is consistent with the canonicalised log), allow implementations (e.g. the epoch_max_wins implementation) to replay raw incoming batches if they choose, and retain the ownership requirement (takes Vec<Operation> so the engine can move the batch without cloning); do not change the function signature, just relax and rephrase the comment to describe the required end state and ownership semantics rather than prescribing a specific replay algorithm.
🤖 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 `@crates/mesh/src/crdt_kv/crdt.rs`:
- Around line 62-84: register_merge_strategy currently installs a new engine for
a prefix but leaves any preexisting keys in the default engine stranded; update
register_merge_strategy (and helpers around EngineHandle/default_engine) to
detect and handle existing data for the new prefix: either (A) reject late
registration with an assertion/error if any keys or log entries matching the
prefix exist in default_engine (scan its keys/entries for prefix matches) or (B)
migrate matching state by iterating default_engine's entries/log for that prefix
and applying them into the newly created EngineHandle (e.g., via
engine.merge/ingest APIs) before replacing the engines list; pick one behavior
and implement the check/migration atomically while holding the engines write
lock so lookups are not raced.
---
Outside diff comments:
In `@crates/mesh/src/crdt_kv/engine/mod.rs`:
- Around line 68-76: The trait doc for apply_remote_ops is too prescriptive
about replaying only the post-canonicalisation result; update the documentation
on fn apply_remote_ops(Operation) in engine/mod.rs to instead state the required
end-state contract (the engine must merge the incoming ops into its operation
log, perform canonicalisation/compaction/tombstone collapse, and ensure live
state is consistent with the canonicalised log), allow implementations (e.g. the
epoch_max_wins implementation) to replay raw incoming batches if they choose,
and retain the ownership requirement (takes Vec<Operation> so the engine can
move the batch without cloning); do not change the function signature, just
relax and rephrase the comment to describe the required end state and ownership
semantics rather than prescribing a specific replay algorithm.
🪄 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: 210cfd01-b4b7-44a5-a7bc-2ee32e551033
📒 Files selected for processing (4)
crates/mesh/src/crdt_kv/crdt.rscrates/mesh/src/crdt_kv/engine/epoch_max_wins.rscrates/mesh/src/crdt_kv/engine/lww.rscrates/mesh/src/crdt_kv/engine/mod.rs
…l EpochMaxWins Step 2 of the per-strategy engine refactor (.claude/docs/mesh/crdt-kv-namespace-engine-design.md). No behavior change; preserves the full public surface of `CrdtOrMap` and all 149 mesh lib tests. Why this shape The old `CrdtOrMap` was one shared store - one `KvStore`, one `ValueMetadata` map, one `key_locks`, one `LamportClock`, one `OperationLog` - plus a per-prefix `MergeStrategy` table that branched at every entry point. Every strategy-specific invariant had to be threaded through every shared call site, and every bug in PR #1469 traced back to that seam. What changed - `crdt_kv::engine::NamespaceCrdtEngine` is the new boundary: a single trait that owns local writes, reads, replication, GC, and the in-engine state machine for one namespace's CRDT strategy. State (store, metadata, locks, clock, log) is engine-internal; the trait surface stays byte-oriented at the boundary. - `LwwEngine` implements LWW with its own self-contained state. The old LWW code (`record_insert_metadata`, `record_remove_metadata`, `compact_key_metadata`, `gc_tombstones`, LWW branches of `apply_insert_locked` / `apply_remove`) moves verbatim into this engine. - `EpochMaxWinsLegacyEngine` is a transitional wrapper that lifts the current inline EpochMaxWins code into the trait shape with its own state. PR #3 will replace it with a real `RateLimitEngine` that holds typed `RateLimitShard` state directly without the LWW-shaped `ValueMetadata` layer. - `CrdtOrMap` becomes a thin router: a sorted prefix table maps each key to its engine via longest-prefix-match. Unregistered keys fall through to a built-in default LWW engine so callers that never call `register_merge_strategy` (notably the in-crate tests) keep working. `crdt.rs` shrinks from 808 lines to ~240. What remains for follow-ups - PR #3: replace `EpochMaxWinsLegacyEngine` with `RateLimitEngine` (typed shard state, no LWW-shaped metadata, drops the `_with_strategy` siblings on `OperationLog`). - PR #4 (cosmetic): split `replica.rs` into `replica.rs` + `clock.rs`; introduce associated types on the trait if needed; promote `OperationLog::{merge,compact}_with_strategy` to per-engine ops. - Lifecycle hooks (drain/shutdown) intentionally deferred - CRDT state is purely in-memory today, no need for them yet. Verification - `cargo test -p smg-mesh --lib` - 149 tests pass. - `cargo clippy --workspace --all-targets --all-features -- -D warnings` - clean. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
…y_remote_ops Three follow-ups to the engine refactor based on the bot review batch. register_merge_strategy is one-shot. The old "upsert" branch was dead code copied from the pre-refactor merge_strategies upsert; nobody ever hit it because MeshKV::configure_crdt_prefix already asserts "each prefix configured exactly once" at the public boundary. Drop the replacement branch and assert in CrdtOrMap directly. Closes the codex P1 (register-after-write orphaning the data already routed under that prefix) and the gemini findings that flowed from replacement (generation could decrease, keys could duplicate, all_engines order ambiguity) - all reduced to "what if replacement happens?" and the answer is now "it panics." NamespaceCrdtEngine::put_local doc said "None on rejection" but both engines return Some(current_live_bytes) when rejecting an older (timestamp, replica_id). Four reviewers flagged the mismatch. The return-the-current-value behaviour is more useful than None (caller sees what is actually live without a follow-up get), so align the doc to the impl rather than the other way around. NamespaceCrdtEngine::apply_remote_ops now takes Vec<Operation> by value instead of &[Operation]. The router already builds owned per-engine buckets, so passing them by value lets the engine move the batch into its operation log via OperationLog::from_operations without an extra Vec allocation. Per-op clones inside merge_with_strategy remain (unavoidable - merge_with_strategy stores its own owned copies); the savings is the outer Vec allocation only. 149 mesh lib tests still pass; cargo clippy --workspace clean. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Two follow-ups to PR #1539 reviews on commit ebd9875. Share a single LamportClock across all engines (codex P1). Each engine used to construct its own LamportClock::new(), so one replica could emit ops with the same (replica_id, timestamp) across different engines - mixed-config peers (haven't registered the second prefix, route both keys into one LwwEngine) then deduplicate by op-id in OperationLog::merge_with_strategy and silently drop one of them, diverging from sender. CrdtOrMap now owns an Arc<LamportClock> and passes it to every engine on construction. This is what we had pre-refactor, what spec mesh-v2 §2.2 describes, and what design issue #1540 already recommends - just not what the implementation had. Drop stale doc comment on register_merge_strategy (claude nit). Earlier edit appended new wording above the old block instead of replacing it; rustdoc was rendering two contradictory paragraphs as one. Coderabbit's "late prefix registration orphans data" finding is real in principle (lower-level CrdtOrMap::insert allows it) but not reachable through the public API: CrdtNamespace::put asserts the key matches its prefix, and getting a CrdtNamespace for "worker:" requires calling configure_crdt_prefix("worker:") first. Defensive assert was declined on the same "don't design for hypothetical" basis as the earlier coderabbit and gemini findings. 149 mesh lib tests still pass; cargo clippy clean. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
External callers route through `CrdtNamespace::put` / `delete`, which assert the key starts with the namespace's registered prefix. The namespace handle itself only exists after `MeshKV::configure_crdt_prefix` (which calls `register_merge_strategy`) has registered the prefix's engine. So forcing all external writes through `CrdtNamespace` makes "write before register" structurally unreachable from outside the crate - the late-registration data-orphaning scenario coderabbit and codex flagged can't happen via the public API. Drops just the two mutating entry points. `CrdtOrMap` itself stays publicly re-exported (it appears in some return types) and its reads (`get`, `keys`, etc.) remain public. The crate-internal router methods (`merge`, `get_operation_log`, `gc_tombstones`) await downstream callers as the gossip wiring lands; they're tested today but not production-wired yet. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
ebd9875 to
f443e2d
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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 `@crates/mesh/src/crdt_kv/engine/epoch_max_wins.rs`:
- Around line 369-408: In apply_remote_ops, avoid allocating an extra vector by
removing let mut to_apply = ops.clone(); instead sort ops in place (make ops
mutable), then build the OperationLog from ops.clone() and reuse the sorted ops
for applying; i.e., call ops.sort_by_key(...) on the incoming Vec, pass a single
clone into OperationLog::from_operations, and iterate over the (now-sorted) ops
for clock updates and apply_remote_insert/apply_remote_remove so only one clone
remains.
🪄 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: 85dd28c1-50e3-47ee-b1f0-d48c825a3756
📒 Files selected for processing (7)
crates/mesh/src/crdt_kv/crdt.rscrates/mesh/src/crdt_kv/engine/epoch_max_wins.rscrates/mesh/src/crdt_kv/engine/lww.rscrates/mesh/src/crdt_kv/engine/mod.rscrates/mesh/src/crdt_kv/kv_store.rscrates/mesh/src/crdt_kv/mod.rscrates/mesh/src/crdt_kv/operation.rs
💤 Files with no reviewable changes (1)
- crates/mesh/src/crdt_kv/kv_store.rs
Description
Problem
The old
CrdtOrMapwas one shared store — oneKvStore, oneValueMetadatamap, onekey_locks, oneLamportClock, oneOperationLog— plus a per-prefixMergeStrategytable that branched at every entry point (apply_insert_locked,apply_remove,merge,compact_with_strategy,merge_with_strategy,snapshot_and_truncate). Every strategy-specific invariant had to be threaded through every shared call site, and every bug fixed in #1469 traced back to that seam.Solution
Extract
NamespaceCrdtEngineas the architectural boundary. Each registered prefix gets its own engine that owns its full CRDT state machine (live store, metadata, key locks, clock, operation log, conflict-resolution rules).CrdtOrMapbecomes a thin router that matches each key to the right engine by longest-prefix-match and delegates.This PR implements step 2 of the migration plan: extract the trait, implement
LwwEngine, and lift the existing EpochMaxWins logic into a transitionalEpochMaxWinsLegacyEngineso the trait surface is exercised end-to-end. A follow-up PR will replace the legacy wrapper with a realRateLimitEnginethat holds typedRateLimitShardstate without LWW-shaped metadata.Design tracking issue: #1540 has the full design — problem, alternatives considered, what each layer owns, wire-format compatibility, the remote-apply footgun the engine API seals, and the full migration plan.
No behavior change in this PR. All 149 mesh lib tests pass without modification. Public surface of
CrdtOrMapis preserved.Changes
New module
crates/mesh/src/crdt_kv/engine/mod.rs—NamespaceCrdtEnginetrait. Byte-oriented at the boundary (per design): the engine can use typed state internally (e.g.RateLimitShard) butput_local/get/export_opsdeal inVec<u8>andOperation. Methods cover local writes, reads, replication (export_ops/apply_remote_ops), and GC. State (store, metadata, locks, clock, log) is engine-internal — the trait deliberately doesn't exposeValueMetadata.lww.rs—LwwEngine. Self-contained state: ownKvStore, ownValueMetadatamap, ownkey_locks, ownLamportClock, ownOperationLog. The old LWW code (record_insert_metadata,record_remove_metadata,compact_key_metadata,gc_tombstones, LWW branches ofapply_insert_locked/apply_remove) moves verbatim into this engine.epoch_max_wins.rs—EpochMaxWinsLegacyEngine. Transitional wrapper that lifts the current inline EpochMaxWins code into the trait shape with its own state. A follow-up will replace it with a realRateLimitEnginethat holds typedRateLimitShardstate directly without the LWW-shapedValueMetadatalayer.Router rewrite
crates/mesh/src/crdt_kv/crdt.rs—CrdtOrMaprewritten as a router. Sorted prefix table (Arc<RwLock<Arc<[(String, EngineHandle)]>>>) maps each key to its engine via longest-prefix-match. Unregistered keys fall through to a built-in defaultLwwEngineso callers that never callregister_merge_strategy(notably the in-crate tests) keep working unchanged. File shrinks from 808 lines to ~240.Small supporting changes
OperationLog::from_operations(Vec<Operation>) -> Self— crate-private constructor used by the router to concatenate per-engine logs back into a singleOperationLogfor gossip export.KvStore::allremoved — dead code after the router stopped using it; engines aggregate viakeys()+get().Out of scope (deliberate follow-ups)
EpochMaxWinsLegacyEnginewith a realRateLimitEngine(typed shard state, no LWW-shaped metadata, drop the_with_strategysiblings onOperationLog).replica.rsintoreplica.rs+clock.rs; introduce associated types on the trait if needed.Test Plan
cargo test -p smg-mesh --lib— 149 tests pass (same as pre-refactor).cargo clippy --workspace --all-targets --all-features -- -D warnings— clean.cargo build -p smg— clean.Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit