feat(mesh): gossip CRDT operation log over the wire (d-3a) - #1570
Conversation
Wires the CRDT data path onto the existing sync_stream transport so the
`worker:`, `rl:`, and `config:` namespaces converge across nodes, and
delivers remote changes to local subscribers. Until now CRDT writes only
mutated the local op-log; `CrdtOrMap::merge` was reachable only from
tests and the gossip loop carried stream traffic only.
Wire format:
- New `CrdtOp { key, value, tombstone, timestamp, replica_id }` and
`CrdtBatch { ops }` proto messages + `CRDT_BATCH` StreamMessageType.
`CrdtOp` mirrors the in-crate `Operation`; `replica_id` is the
ReplicaId UUID in text form. Field-by-field (not an opaque bincode
blob) for a language-agnostic schema.
- `transport/crdt_batch.rs`: `build_crdt_batch` (op-log -> CrdtBatch),
`wrap_crdt_batch` (StreamMessage envelope), `dispatch_crdt_batch`
(decode -> merge). Ops with an unparsable replica_id are dropped.
Producer/consumer:
- `RoundBatch` gains a `crdt_ops` snapshot, filled in
`collect_round_batch` from the store's op-log. Both per-peer senders
(gossip_controller dialed streams, gossip_service accepted streams)
broadcast the snapshot each round alongside the stream batch. Merge is
idempotent by op-id, so re-sending seen ops is a no-op; per-peer
watermark filtering is a follow-up.
- Both inbound dispatch blocks gain a `CrdtBatch` arm routing to
`dispatch_crdt_batch`.
Subscriber fan-out + value-shape alignment (migration step 7):
- `NamespaceCrdtEngine::apply_remote_ops` now returns `Vec<CrdtChange>`;
`CrdtOrMap::merge` concatenates per-engine changes; `MeshKV::merge_crdt_ops`
fires `SubscriberRegistry::notify` per change with the canonical
post-merge value (matching `get`). `CrdtNamespace::put` likewise now
notifies the canonical post-insert value instead of the raw caller
payload, so local-write and remote-merge subscribers see one shape
(the `rl:` adapter no longer needs to special-case the raw payload).
- Change detection emits a `CrdtChange` iff the observable `get(key)`
value actually changes (snapshot before/after the apply loop), not on
mere op acceptance. This suppresses spurious events: a newer-version
insert rewriting byte-identical bytes, and a tombstone for an
already-absent key (Tombstone->Tombstone or vacant->Tombstone, both
encoding to None). Generation/log semantics are unchanged.
Tests: codec round-trip + envelope shape; end-to-end (no gRPC)
convergence, remote-merge subscriber fan-out, idempotent re-delivery
fires nothing, `rl:` canonical-shard shape, and tombstone->None;
`CrdtOrMap::merge` change-detection contract (new value, byte-identical
no-op, never-seen tombstone no-op, kill-live emits None). 165 mesh tests
pass; gateway builds and its mesh-adapter tests pass; clippy + fmt clean.
Not in this PR (follow-ups): per-peer CRDT watermark (bandwidth/at-least-once
retry), initial snapshot-on-join, and wiring the adapters into server.rs.
Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
|
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 (6)
📝 WalkthroughWalkthroughAdds per-key change reporting to CRDT remote merges, introduces wire/batching for CRDT ops over gossip, and integrates merge delivery into MeshKV to notify subscribers with canonical post-merge values. ChangesCRDT Operation Broadcasting and Merge Change Reporting
Sequence Diagram(s)sequenceDiagram
participant SenderMeshKV as Sender MeshKV
participant RoundBatch as RoundBatch
participant Gossip as Gossip Sender
participant Transport as crdt_batch (wrap/build)
participant RecvGossip as Receiver Gossip
participant ReceiverMeshKV as Receiver MeshKV
participant Subscriber as Subscriber
SenderMeshKV->>RoundBatch: collect_round_batch() (includes crdt_ops)
RoundBatch-->>Gossip: RoundBatch
Gossip->>Transport: build_crdt_batches(ops)
Transport-->>Gossip: CrdtBatch frames
Gossip->>RecvGossip: send StreamMessage(CrdtBatch)
RecvGossip->>Transport: dispatch_crdt_batch()
Transport->>ReceiverMeshKV: merge_crdt_ops(ops)
ReceiverMeshKV->>ReceiverMeshKV: store.merge + snapshot per-key before/after
ReceiverMeshKV->>Subscriber: notify Vec<CrdtChange>
Estimated code review effort🎯 4 (Complex) | ⏱️ ~50 minutes Possibly related issues
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)
Warning Review ran into problems🔥 ProblemsStopped waiting for pipeline failures after 30000ms. One of your pipelines takes longer than our 30000ms fetch window to run, so review may not consider pipeline-failure results for inline comments if any failures occurred after the fetch window. Increase the timeout if you want to wait longer or run a Comment |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 4f77e884d5
ℹ️ 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".
| return None; | ||
| } | ||
| Some(CrdtBatch { | ||
| ops: ops.iter().map(op_to_proto).collect(), |
There was a problem hiding this comment.
Split CRDT batches before sending them over gRPC
When the CRDT op-log snapshot contains a large value or enough keys to exceed the 10 MiB MAX_MESSAGE_SIZE used by the tonic clients/servers, this constructs a single unbounded CrdtBatch frame. Unlike the stream path, which chunks values below the gRPC cap, this frame will be rejected during encode/decode and can close/fail sync_stream, leaving that peer unable to receive CRDT updates until the log shrinks. Please split/chunk CRDT batches or enforce the transport limit before wrapping them in a StreamMessage.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Fixed in 2d7ad42. build_crdt_batches now splits the op-log snapshot into frames each estimated to stay under MAX_STREAM_CHUNK_BYTES (= MAX_MESSAGE_SIZE − 64 KiB envelope margin, the same budget the stream path uses), and both senders loop over the returned batches. A single op larger than the budget is emitted alone (best-effort); worker:/rl:/config: values are far below the cap, so in practice this only bounds the op count per frame.
There was a problem hiding this comment.
Code Review
This pull request implements the CRDT-over-gossip data path, enabling the broadcasting and merging of CRDT operation logs (CrdtBatch) between peers, and updating subscribers only when observable values change. The review feedback highlights several improvement opportunities: avoiding unnecessary heap allocations in the LWW and Rate Limit engines by checking contains_key with a borrowed string slice instead of using entry with an owned string; offloading the CPU-bound dispatch_crdt_batch operation to tokio::task::spawn_blocking to prevent blocking the async executor; and adding defensive checks to ensure large serialized batches do not exceed the gRPC MAX_MESSAGE_SIZE limit.
| let mut before: std::collections::HashMap<String, Option<Vec<u8>>> = | ||
| std::collections::HashMap::new(); | ||
| for op in &unseen { | ||
| before | ||
| .entry(op.key().to_string()) | ||
| .or_insert_with(|| self.store.get(op.key())); | ||
| } |
There was a problem hiding this comment.
In Rust, HashMap::entry takes an owned key, which means op.key().to_string() will allocate a new String on every single iteration of the loop, even if the key is already present in the before map. Since unseen can contain multiple operations for the same key (especially during bulk merges or heavy gossip rounds), this results in unnecessary heap allocations.
To avoid these allocations, check if the key is already present using contains_key with a borrowed &str before inserting.
| let mut before: std::collections::HashMap<String, Option<Vec<u8>>> = | |
| std::collections::HashMap::new(); | |
| for op in &unseen { | |
| before | |
| .entry(op.key().to_string()) | |
| .or_insert_with(|| self.store.get(op.key())); | |
| } | |
| let mut before: std::collections::HashMap<String, Option<Vec<u8>>> = | |
| std::collections::HashMap::new(); | |
| for op in &unseen { | |
| if !before.contains_key(op.key()) { | |
| before.insert(op.key().to_string(), self.store.get(op.key())); | |
| } | |
| } |
There was a problem hiding this comment.
Fixed in 2d7ad42 — switched to if !before.contains_key(op.key()) { before.insert(op.key().to_string(), self.store.get(op.key())) } so the String is only allocated for keys not yet snapshotted.
| let mut before: std::collections::HashMap<String, Option<Vec<u8>>> = | ||
| std::collections::HashMap::new(); | ||
| for op in &ops { | ||
| before | ||
| .entry(op.key().to_string()) | ||
| .or_insert_with(|| self.current_encoded(op.key())); | ||
| } |
There was a problem hiding this comment.
Similar to the LWW engine, calling before.entry(op.key().to_string()) allocates a new String on every iteration of the loop even if the key already exists in the map. Since ops can contain duplicate keys, we can avoid these unnecessary heap allocations by checking contains_key first.
| let mut before: std::collections::HashMap<String, Option<Vec<u8>>> = | |
| std::collections::HashMap::new(); | |
| for op in &ops { | |
| before | |
| .entry(op.key().to_string()) | |
| .or_insert_with(|| self.current_encoded(op.key())); | |
| } | |
| let mut before: std::collections::HashMap<String, Option<Vec<u8>>> = | |
| std::collections::HashMap::new(); | |
| for op in &ops { | |
| if !before.contains_key(op.key()) { | |
| before.insert(op.key().to_string(), self.current_encoded(op.key())); | |
| } | |
| } |
There was a problem hiding this comment.
Fixed in 2d7ad42, same contains_key guard as the LWW engine.
| pub fn dispatch_crdt_batch(mesh_kv: &MeshKV, batch: CrdtBatch) { | ||
| let ops: Vec<Operation> = batch.ops.into_iter().filter_map(proto_to_op).collect(); | ||
| if ops.is_empty() { | ||
| return; | ||
| } | ||
| mesh_kv.merge_crdt_ops(ops); | ||
| } |
There was a problem hiding this comment.
The dispatch_crdt_batch function is called synchronously on the async executor thread. However, merging a CrdtBatch (which can contain up to AUTO_COMPACT_THRESHOLD of 10,000 operations) is a highly CPU-bound operation involving sorting, compaction, multiple map lookups, and lock acquisitions. Running this synchronously can block the async executor thread for tens of milliseconds, leading to latency spikes or heartbeat timeouts on other active connections.
Consider offloading this CPU-bound merge operation to tokio::task::spawn_blocking to keep the async executor responsive.
References
- Offload CPU-bound operations, such as large-scale bincode serialization or complex tree walks, to
tokio::task::spawn_blockingto prevent blocking the async executor threads.
There was a problem hiding this comment.
Deferring this one, with reasoning:
-
Not CRDT-specific.
dispatch_stream_batch(the sibling, already in production) also runs synchronously on the executor and handles larger traffic (tree pages up to 64 MB, chunked). If sync dispatch were a problem it would already be one there. The right fix is a single cross-cutting "non-blocking receive path" change covering both stream and CRDT dispatch, not a CRDT-only bolt-on. -
Cost only bites out-of-spec. The merge is O(batch + log). In the design's supported regime (~30 keys) the op-log is tiny and the merge is microseconds. The "tens of ms" only materializes at the ~10K-op regime, which is already out of the v2 design envelope.
-
spawn_blockingisn't free. It requires the work to beSend + 'static(cloning anArcof state into the closure — a real restructure), pulls a runtime dependency into a currently runtime-agnostic codec layer, and converts a serialized executor path into a concurrent one against the sameCrdtOrMap(safe given the engines' internal locking, but a semantics change deserving its own review). For the common small-batch case it also adds latency. -
No measurement justifies it yet — it's a hypothesis for an out-of-spec regime, and the broadcast-full-op-log behavior is interim (the watermark follow-up shrinks batches to deltas anyway).
Noting it as a candidate for a future cross-cutting non-blocking-dispatch change, gated on an actual profile during the d-3 rollout.
| pub fn build_crdt_batch(ops: &[Operation]) -> Option<CrdtBatch> { | ||
| if ops.is_empty() { | ||
| return None; | ||
| } | ||
| Some(CrdtBatch { | ||
| ops: ops.iter().map(op_to_proto).collect(), | ||
| }) | ||
| } |
There was a problem hiding this comment.
Since CrdtBatch currently broadcasts the full op-log snapshot (which can grow up to AUTO_COMPACT_THRESHOLD of 10,000 operations) without any chunking or size validation, a large op-log could easily produce a serialized message that exceeds the gRPC MAX_MESSAGE_SIZE limit. If this happens, the gRPC stream will be abruptly terminated with a "message too large" error, disrupting gossip convergence.
While delta-based watermark filtering is planned as a follow-up, consider adding a defensive check or logging warning if the serialized batch size approaches MAX_MESSAGE_SIZE to prevent silent stream failures.
References
- When chunking data for a transport protocol, reserve a margin from the maximum message size for protocol envelope overhead (headers, metadata) to prevent the total serialized message from being rejected.
There was a problem hiding this comment.
Addressed in 2d7ad42 — went with splitting rather than just a warning: build_crdt_batches bounds each frame below MAX_MESSAGE_SIZE (reusing MAX_STREAM_CHUNK_BYTES with its 64 KiB envelope margin), so a large full-op-log broadcast can never produce an oversized frame that tears down the stream. New test build_crdt_batches_splits_over_budget covers it.
| /// Receiver-side dispatch for a `CrdtBatch`: decode each op and merge the batch | ||
| /// into the local CRDT store. Ops with an unparsable `replica_id` are skipped. | ||
| /// Merge is idempotent by op-id, so a batch the node has already absorbed is a | ||
| /// no-op. (Subscriber notification on remote merge is a follow-up — d-3a-2.) |
There was a problem hiding this comment.
🟡 Nit: This parenthetical says subscriber notification is a follow-up (d-3a-2), but merge_crdt_ops (called via mesh_kv.merge_crdt_ops(ops) on line 98) already fires subscriber_registry.notify for every changed key. The comment is stale — this PR implements the fan-out.
| /// no-op. (Subscriber notification on remote merge is a follow-up — d-3a-2.) | |
| /// no-op. |
There was a problem hiding this comment.
Good catch — fixed in 2d7ad42. The doc now states subscriber fan-out on remote merge is implemented here (via merge_crdt_ops), and that an already-absorbed batch is a no-op that fires no event. The d-3a-2 reference was stale leftover from before the fan-out landed in this PR.
- Split the per-round CRDT op-log snapshot into size-bounded frames (`build_crdt_batches`, capped at MAX_STREAM_CHUNK_BYTES) instead of one unbounded `CrdtBatch`. A large op-log (broadcast in full each round until per-peer watermark filtering lands) could otherwise serialize to a frame above MAX_MESSAGE_SIZE, which tonic rejects on encode/decode and which would tear down the sync_stream. Both senders now loop over the returned batches. (codex P2 + gemini) - Avoid a per-op `String` allocation in both engines' change-detection snapshot: check `contains_key(&str)` before inserting, since a batch may repeat a key. (gemini) - Fix the stale `dispatch_crdt_batch` doc comment: subscriber fan-out on remote merge is implemented in this PR (via `merge_crdt_ops`), not a follow-up. (claude) Tests: `build_crdt_batches` empty/single + over-budget split coverage. 166 mesh tests pass; clippy + fmt clean. Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Description
Problem
The mesh-v2 gap analysis (
.claude/docs/mesh/mesh-v2-gap-analysis.md) flagged the d-3 cutover as the single largest gap: CRDT operations never reached the wire.CrdtNamespace::put/deleteonly mutated the local op-log,CrdtOrMap::mergewas reachable only from tests, and the gossip loop carried stream traffic (td:/tree:*) only. So theworker:,rl:, andconfig:CRDT namespaces were silent across nodes, and remote CRDT changes never reached local subscribers.Solution
Wire the CRDT data path onto the existing
sync_streamtransport (this is d-3a; the adapter wiring inserver.rsand rate-limit enforcement are follow-ups). Two halves, in one PR:getreturns.Changes
Wire format (
crates/mesh/src/proto/gossip.proto)CrdtOp { key, value, tombstone, timestamp, replica_id }+CrdtBatch { ops }+CRDT_BATCHStreamMessageType+ acrdt_batcharm in theStreamMessageoneof.CrdtOpmirrors the in-crateOperation;replica_idis theReplicaIdUUID in text form. Field-by-field (not an opaque bincode blob) for a language-agnostic schema.Codec (
crates/mesh/src/transport/crdt_batch.rs, new)build_crdt_batch(op-log →CrdtBatch),wrap_crdt_batch(envelope),dispatch_crdt_batch(decode → merge). Ops with an unparsablereplica_idare dropped rather than poisoning the merge.Producer / consumer
RoundBatchgains acrdt_opssnapshot, filled incollect_round_batchfrom the store's op-log. Both per-peer senders (gossip_controllerdialed streams,gossip_serviceaccepted streams) broadcast the snapshot each round alongside the stream batch. Merge is idempotent by op-id, so re-sending seen ops is a no-op — per-peer watermark filtering is a follow-up.CrdtBatcharm →dispatch_crdt_batch.Fan-out + step 7 (
engine/mod.rs,engine/lww.rs,engine/rate_limit.rs,crdt.rs,kv.rs)NamespaceCrdtEngine::apply_remote_opsnow returnsVec<CrdtChange>;CrdtOrMap::mergeconcatenates per-engine changes;MeshKV::merge_crdt_opsfiresSubscriberRegistry::notifyper change with the canonical post-merge value (matchingget).CrdtNamespace::putnow notifies the canonical post-insert value instead of the raw caller payload, so local-write and remote-merge subscribers see one shape (therl:adapter no longer needs to special-case the raw payload; its module doc is updated).CrdtChangeiff the observableget(key)value actually changes (snapshot before/after the apply loop), not on mere op acceptance. This suppresses spurious events: a newer-version insert rewriting byte-identical bytes, and a tombstone for an already-absent key (Tombstone→Tombstoneorvacant→Tombstone, both encoding toNone). Generation/log semantics are unchanged.Verification
An adversarial multi-agent pass reviewed the change-detection for missed/spurious/wrong-value events. It found no missed events and no convergence bugs, and two P2 spurious-event classes — both fixed by the observable-value diff above. (Two findings were correct-as-is: a frontier reshuffle that changes the encoded shard bytes correctly fires, and a benign concurrency TOCTOU where the emitted value is always a valid
getsnapshot and converges.)Test Plan
cargo test -p smg-mesh --lib— 165 passed, 0 failed (+15 over baseline)tests/crdt_integration.rs): convergence, remote-merge subscriber fan-out, idempotent re-delivery fires nothing,rl:canonical-shard shape, tombstone→NoneCrdtOrMap::mergechange-detection contract: new value, byte-identical no-op, never-seen tombstone no-op, kill-live emitsNonecargo build -p smg+cargo test -p smg --lib mesh::adapters— pass (downstream unaffected)cargo clippy -p smg-mesh --all-targets+cargo fmt— cleanNet: +635 / −21 across 18 files.
Follow-ups (not in this PR)
worker:/rl:/tree) intoserver.rsand connect outbound change paths (d-3 PR 2/3); rate-limit middleware enforcement (d-3b).Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Tests