refactor(mesh): typed RateLimitEngine, drop EpochMaxWinsLegacyEngine - #1549
Conversation
…Engine The transitional engine introduced in #1539 mirrored the LWW layout (`Vec<ValueMetadata>` per key + separate `KvStore` byte map + per-key `Mutex` table) even though EpochMaxWins state is already richer: the `RateLimitShard` carries its own tombstone_version and live-points frontier inline. `RateLimitEngine` holds `DashMap<String, ShardEntry>` where each entry is a typed `RateLimitState` (`Live(shard)` or `Tombstone(version)`) plus a local `tombstoned_at` instant for GC. DashMap's entry API supplies the per-key atomicity the per-key `Mutex` table used to provide, so the mutex map and its cleanup plumbing are gone. The tombstoned_at clock is never refreshed on a tombstone-to-tombstone merge, preserving the PR #1469 codex P2 fix. Typed primitives in `crdt_kv/epoch_max_wins.rs` (`RateLimitState`, `state_from_insert_value`, `RateLimitState::merge`, new `encode_live`) are promoted to `pub(super)` so the engine operates on shards directly without byte round-trips on the put/apply paths. The byte-bounded `merge_live_value` / `apply_tombstone` wrappers (and the now-unused `LiveMerge` / `TombstoneApply` shapes) are deleted - their sole callers were the legacy engine, with two test helpers updated to the typed API. Net: -240 lines. All 145 mesh tests pass unchanged. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
|
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:
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (2)
📝 WalkthroughWalkthroughRemoves the legacy EpochMaxWinsLegacyEngine, adds a new RateLimitEngine and re-exports it, routes MergeStrategy::EpochMaxWins to RateLimitEngine, promotes typed RateLimitState/RateLimitShard helpers, refactors merge plumbing, and adds a tombstone GC test. ChangesEngine Migration and Integration
Sequence DiagramsequenceDiagram
participant Client
participant RateLimitEngine
participant DashMap
participant OperationLog
participant LamportClock
Client->>RateLimitEngine: put_local(key, value)
RateLimitEngine->>RateLimitEngine: merge into RateLimitState
RateLimitEngine->>OperationLog: append Operation
RateLimitEngine->>Client: return previous value
Client->>RateLimitEngine: apply_remote_ops(ops)
RateLimitEngine->>RateLimitEngine: sort ops by (timestamp, replica_id)
RateLimitEngine->>RateLimitEngine: compact with EpochMaxWins
RateLimitEngine->>LamportClock: update per op
RateLimitEngine->>DashMap: merge each op into state
RateLimitEngine->>Client: return (generation only incremented on change)
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Code Review
This pull request replaces the transitional EpochMaxWinsLegacyEngine with a direct, typed RateLimitEngine for the EpochMaxWins CRDT. This removes the legacy LWW-shaped metadata layer, separate key-value store, and per-key mutex table, consolidating them into a single DashMap holding typed RateLimitState per key. The review feedback suggests two performance optimizations: first, consolidating double lookups in put_local into a single lookup using the DashMap entry API; second, avoiding an expensive clone of the operations vector in apply_remote_ops by reordering the state updates and log merging.
There was a problem hiding this comment.
Clean, well-executed refactoring. The typed RateLimitEngine correctly preserves all CRDT invariants (tombstone GC grace clock, per-point merge semantics, compacted-snapshot replay) while removing ~224 lines of LWW-shaped indirection.
One nit flagged inline: get() now re-serializes on every call instead of returning cached bytes. Worth watching if rl: reads turn out to be hot, but not blocking.
Summary: 0 🔴 Important · 1 🟡 Nit · 0 🟣 Pre-existing
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 61c98237f7
ℹ️ 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: 2
🤖 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/rate_limit.rs`:
- Around line 162-181: put_local and delete_local return incorrect/unsafe
"previous" bytes because current_encoded is called outside the per-key lock and
merge_insert/merge_remove only return a boolean; change merge_insert and
merge_remove to return a structured enum/result that (while holding the DashMap
entry lock) includes: whether the previous value is well-defined, the previous
live bytes if any, and whether the merge changed state or produced a malformed
payload; then update put_local to call merge_insert and derive its returned
Option<Vec<u8>> from that structured outcome (avoiding calling current_encoded
before the lock, and returning None when EpochMaxWins semantics make the
previous value undefined), and update delete_local similarly to return the
previous live bytes when merge_remove transitioned Live->Tombstone; keep
existing side effects (creating Operation via Operation::insert/remove,
append_op, generation.fetch_add, debug logging, and RateLimitVersion usage)
triggered only when the structured result indicates a change.
In `@crates/mesh/src/crdt_kv/epoch_max_wins.rs`:
- Around line 326-330: The current match branch that handles
(decode_shard(local), decode_shard(remote)) and then calls
RateLimitState::Live(local_shard).merge(RateLimitState::Live(remote_shard)) must
not silently fall through when the typed merge yields a non-live result; change
the else arm to fail fast (panic! or .expect with a clear message) so callers
always receive a Live shard. Locate the match arm that binds local_shard and
remote_shard and replace the current else handling with an explicit panic/expect
stating that the merge produced a non-live/tombstone result for the given
shards.
🪄 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: 166f4459-d136-43a3-99ef-eae19d587f0c
📒 Files selected for processing (5)
crates/mesh/src/crdt_kv/crdt.rscrates/mesh/src/crdt_kv/engine/epoch_max_wins.rscrates/mesh/src/crdt_kv/engine/mod.rscrates/mesh/src/crdt_kv/engine/rate_limit.rscrates/mesh/src/crdt_kv/epoch_max_wins.rs
💤 Files with no reviewable changes (1)
- crates/mesh/src/crdt_kv/engine/epoch_max_wins.rs
`update_entry` previously kept `tombstoned_at` on every tombstone -> tombstone transition, conflating two distinct cases: - An older dominated Remove arriving late at an already-tombstoned key. The GC clock must NOT reset, otherwise stale gossip from a lagging peer pins the clock indefinitely. - A newer winning Remove that supersedes the existing tombstone version. The GC clock MUST reset so the new tombstone gets its full grace period. The transitional engine refreshed `created_at` on the version-advance case by clearing and re-pushing metadata; this regression was introduced when the typed engine collapsed both branches into "keep tombstoned_at as-is." Without the refresh, a newer winning tombstone inherits the prior tombstone's already-aged clock, GC drops the entry immediately, and a delayed insert older than the winning tombstone but newer than any remaining local frontier can resurrect the key. `update_entry` now compares the winning tombstone version before vs after the merge and refreshes only when it advances. The older-delayed-remove preservation path stays intact (same-version merge leaves the clock alone). Add a symmetric regression test next to the existing `test_epoch_max_wins_older_delayed_remove_preserves_tombstone_age` that exercises the advance path and asserts GC does not collect immediately. Also strip stray planning metadata (PR-number refs) from comments on the touched lines. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 056cae8444
ℹ️ 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
🤖 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/tests.rs`:
- Around line 675-688: The test's timing window is too tight around
replica.merge(&newer) and the subsequent call to
replica.gc_tombstones_with_grace(Duration::from_millis(50)), causing flaky
failures; change the test to avoid relying on precise wall-clock ordering by
either increasing the grace gap (e.g., use much larger sleep durations around
thread::sleep(Duration::from_millis(60)) and the grace argument) or replace the
single sleep + immediate gc_tombstones_with_grace call with a deadline/polling
loop that repeatedly calls replica.gc_tombstones_with_grace until the expected
condition is stable (or a timeout is reached), ensuring the merge of newer
OperationLog and the refreshed tombstone state are observed reliably.
🪄 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: c1b898ac-7479-4685-9da5-77d8719fb59e
📒 Files selected for processing (2)
crates/mesh/src/crdt_kv/engine/rate_limit.rscrates/mesh/src/crdt_kv/tests.rs
`merge_insert` / `merge_remove` now return a `MergeOutcome` captured under the per-key entry lock, recording whether state changed, the encoded prior live bytes (if the prior state was `Live`), and whether the new state is live. `put_local` and `delete_local` derive their return values from that outcome, closing three contract gaps: - `put_local` previously sampled `previous` via `current_encoded(key)` before taking the entry lock, so concurrent same-key writers could observe a "previous" that no longer matched what the merge displaced. The sample now happens inside the entry guard. - `put_local` returned `Some(prior_bytes)` for every accepted insert, including per-point frontier updates where the prior shard is not a well-defined "previous." The trait says return `None` for that case (key remains alive with a different shard, no clean displaced value). The function now returns `None` for the per-point and vacant-insert paths, `Some(prior_live)` only when the merge transitions `Live -> Tombstone` (rare, via compacted incoming with a tombstone bound that kills the existing frontier), and `Some(current_live)` on rejected or malformed input. - `delete_local` returned `None` unconditionally, even when the delete killed the last live points. The trait says return prior live bytes in that case. The function now returns `Some(prior_live)` exactly when the entry transitions `Live -> Tombstone`, and `None` for deletes that arrive at vacant or already-tombstoned keys or leave live points behind (lower-version tombstone). Four regression tests in `crdt_kv/tests.rs` pin the new behavior: per-point update returns `None`, malformed put returns current live bytes, delete that kills the last live points returns prior live bytes, delete on a never-seen key returns `None`. The restructure also collapses the prior two-lookup pattern in `put_local` (`current_encoded` + entry lookup) into a single entry lookup, addressing the related lookup-discipline nit on the same code path. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
`test_epoch_max_wins_newer_tombstone_refreshes_tombstone_age` asserted `removed == 0` with only a 50ms grace, which is sensitive to scheduler stalls between the `replica.merge(&newer)` call and the `gc_tombstones_with_grace` check - if more than 50ms elapsed between them, the new tombstone would age past grace and the assertion would flip to 1. The companion older-delayed-remove test (which asserts `removed == 1`) is symmetric in form but not in flakiness, since a longer scheduler stall only ages the tombstone further past its grace window without flipping the assertion. Widen the first-tombstone sleep to 250ms and the GC grace to 150ms so the second tombstone's GC age has ~150ms of headroom before the assertion becomes wrong. The first tombstone is still well past grace (250 > 150), so the only-the-newer-tombstone behavior remains testable. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: aee9a4aa21
ℹ️ 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".
`merge_insert` previously decoded the incoming payload before taking the DashMap entry guard, so a malformed put bailed out before any per-key serialisation. `put_local` then read current live bytes via a separate `current_encoded` lookup, leaving a window where a concurrent valid same-key write could commit between the failed decode and the fallback read. The returned "current" bytes were correct in absolute terms but did not reflect the value at this call's serialisation point, contradicting the per-key atomic read/modify/return shape the previous engine had with its dedicated key-lock map. Move `state_from_insert_value` inside the entry guard. A malformed payload now produces a no-change `MergeOutcome` carrying the prior live bytes (sampled under the same lock), indistinguishable to `put_local` from a dominated/idempotent rejection. `merge_insert` no longer returns `Option<MergeOutcome>` - malformed is just another rejection. Callers simplify accordingly. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
The `#[cfg(test)] fn merge` helper used by the unit tests below it silently fell back to `local.to_vec()` when the typed merge yielded a non-`Live` result (the `let-else` guard returned the local bytes). Every caller below assumes a live merge result, so a future regression that makes the merge produce a tombstone-only state would have looked like a test pass against the unchanged local bytes. Panic with a descriptive message instead. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Description
Problem
The EpochMaxWins engine introduced in #1539 was deliberately transitional. It mirrored the LWW layout (
Vec<ValueMetadata>per key + separateKvStorebyte map + per-keyMutextable) even though EpochMaxWins state is already richer: eachRateLimitShardcarries its owntombstone_versionandlive_pointsfrontier inline.That duplication forced every put/apply to encode the typed shard → bytes, hand bytes to the byte-bounded merge primitive, decode bytes back to typed state internally, and re-encode. It also meant carrying an LWW-shaped metadata vec whose only unique information was the
created_attimestamp for tombstone GC.Solution
Replace the engine with one that holds typed
RateLimitStateper key directly.DashMap<String, ShardEntry>, whereShardEntryis the typed state plus a singletombstoned_at: Option<Instant>for GC grace. DashMap's entry API gives equivalent per-key atomicity, so the per-keyMutextable and its cleanup plumbing are gone.The byte boundary stays at the
NamespaceCrdtEnginetrait — engines are still byte-oriented externally so the router (CrdtOrMap) stays strategy-agnostic. Internally the engine operates on typed shards without round-trips, matching what the comment inengine/mod.rs:30-31already promised the boundary was for.Changes
crates/mesh/src/crdt_kv/engine/rate_limit.rs—RateLimitEngineholdingDashMap<String, ShardEntry>.crates/mesh/src/crdt_kv/engine/epoch_max_wins.rs— the transitional legacy engine.crates/mesh/src/crdt_kv/epoch_max_wins.rs: typed primitives promoted topub(super)(RateLimitState,RateLimitShard,state_from_insert_value,RateLimitState::merge, newRateLimitState::encode_live). Dead byte-bounded wrappers (merge_live_value,apply_tombstone,LiveMerge,TombstoneApply,state_from_stored_value,decode_stored_live) and their two test helpers updated to use the typed API.crates/mesh/src/crdt_kv/crdt.rs: registration site now instantiatesRateLimitEngineforMergeStrategy::EpochMaxWins.crates/mesh/src/crdt_kv/engine/mod.rs: re-export updated; trait doc comment refers toRateLimitState/RateLimitEngine.Net: +345 / −569 ≈ -224 lines.
Tombstone GC invariant preserved
The
tombstoned_atclock is set only on the live → tombstone transition. Tombstone → tombstone merges (a late delayed Remove arriving at a key that's already tombstoned) keep the originaltombstoned_at, so the grace clock cannot be restarted by a stale Remove. This is the PR #1469 codex P2 fix, captured directly inupdate_entry.Trait surface unchanged
NamespaceCrdtEngineis identical.put_local/delete_local/apply_remote_opssemantics are preserved end-to-end, verified by the existing 16+ EpochMaxWins tests incrdt_kv/tests.rspassing unchanged.Test Plan
cargo test -p smg-mesh --lib— 145 passed, 0 failed (including the EpochMaxWins suite that exercises the engine end-to-end throughCrdtOrMap)cargo clippy -p smg-mesh --all-targets— cleancargo check --workspace— cleanChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Refactor
Tests