Repository navigation
feat(watch) #379: prev_kv + progress heartbeat - #385
Conversation
## Production code - WatchEvent: add `prev_value: Bytes` field; add `WatchEventType::Progress` variant - WatchRegistry: `register()` / `register_prefix()` take `prev_kv: bool`; expose `prev_kv_watcher_count_arc()` for shared atomic counter - DefaultStateMachineHandler: when `prev_kv_watcher_count > 0`, read old values from state machine before `apply_chunk` and populate `prev_value` in broadcast events - WatchDispatcher: `heartbeat_interval_ms` param; tokio interval with ±10% jitter broadcasts `Progress` events to all watchers; `broadcast_progress()` fan-out iterates registry maps directly (key-routing would miss watchers on empty key) - EmbeddedClient: add `watch_with_options(key, prev_kv)` and `watch_prefix_with_options(prefix, prev_kv)` as explicit prev_kv entry points - Proto: add `prev_value` field to `WatchResponse`; add `prev_kv` field to `WatchRequest` - Config: add `heartbeat_interval_ms: u64` to `WatchConfig` - NodeBuilder: wire `prev_kv_watcher_count` Arc into `DefaultStateMachineHandler::new()` ## Tests (RED → GREEN) - manager_test: 6 new tests — prev_kv isolation per watcher, watcher count tracking, progress heartbeat delivery, heartbeat disabled produces no Progress events - default_state_machine_handler_test: 3 new tests — skip get when count=0, read and broadcast prev_value when count=1, stop reading after count drops to zero - grpc_client_test: add `prev_value` field to WatchResponse fixture - embedded_client_test: guard both test mods with `#[cfg(all(test, feature="rocksdb"))]` ## API compatibility fixes (callers broken by new fields / signatures) - benches/state_machine.rs: `prev_value`, `register(_, false)`, WatchDispatcher 5-arg - tests/embedded_client: `WatchEventType::X as i32` → `WatchEventType::X` - tests/watch_events_embedded: same - examples/service-discovery-embedded: same + explicit Canceled arm - examples/service-discovery-standalone: same + Progress arm for exhaustive match - integration/mod.rs: DefaultStateMachineHandler::new() 7th arg ## Docs - watch-feature.md: replace "No prev_kv / No heartbeat (planned)" with feature docs; add `heartbeat_interval_ms` config entry; fix `WatchEventType::CANCELED` casing (×3); fix `Canceled as i32` in reconnection pattern; add Progress to Best Practices - embedded_client.rs: full rustdoc for `watch_with_options` and `watch_prefix_with_options` — prev_kv semantics, performance note, empty-prev_value cases, code examples - service-discovery-embedded: Demo 3 — prev_kv audit log pattern (put/delete/canceled/progress)
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: defaults Review profile: CHILL Plan: Pro Run ID: ⛔ Files ignored due to path filters (1)
📒 Files selected for processing (8)
✅ Files skipped from review due to trivial changes (1)
🚧 Files skipped from review as they are similar to previous changes (5)
📝 WalkthroughWalkthroughThis PR adds optional previous-value reporting (WatchRequest.prev_kv → WatchResponse.prev_value) and periodic Progress heartbeats. Proto, watch manager dispatcher, state-machine handler, client/server wiring, tests, docs, examples, benchmarks, and config (heartbeat_interval_ms) are updated to support opt-in prev_kv and heartbeat delivery. ChangesWatch prev_kv and heartbeat progress
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related issues
Possibly related PRs
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)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 8
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (3)
d-engine-client/src/grpc_client.rs (1)
217-228:⚠️ Potential issue | 🟠 Major | ⚡ Quick winExpose
prev_kvin the public gRPC watch API.
watch/watch_prefixhardcodeprev_kv: false, and there is no*_with_optionsvariant here. That makesprev_kveffectively inaccessible throughGrpcClient.Suggested shape
pub async fn watch( &self, key: impl AsRef<[u8]>, ) -> ClientApiResult<tonic::Streaming<WatchResponse>> { + self.watch_with_options(key, false).await +} + +pub async fn watch_with_options( + &self, + key: impl AsRef<[u8]>, + prev_kv: bool, +) -> ClientApiResult<tonic::Streaming<WatchResponse>> { let client_inner = self.client_inner.load(); let request = WatchRequest { client_id: client_inner.client_id, key: Bytes::copy_from_slice(key.as_ref()), prefix: false, - prev_kv: false, + prev_kv, }; ... }pub async fn watch_prefix( &self, prefix: impl AsRef<[u8]>, ) -> ClientApiResult<tonic::Streaming<WatchResponse>> { + self.watch_prefix_with_options(prefix, false).await +} + +pub async fn watch_prefix_with_options( + &self, + prefix: impl AsRef<[u8]>, + prev_kv: bool, +) -> ClientApiResult<tonic::Streaming<WatchResponse>> { let client_inner = self.client_inner.load(); let request = WatchRequest { client_id: client_inner.client_id, key: Bytes::copy_from_slice(prefix.as_ref()), prefix: true, - prev_kv: false, + prev_kv, }; ... }Also applies to: 249-260
🤖 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 `@d-engine-client/src/grpc_client.rs` around lines 217 - 228, The public gRPC watch API currently hardcodes WatchRequest.prev_kv = false in the watch and watch_prefix implementations (see methods watch and watch_prefix in grpc_client.rs), making prev_kv inaccessible; add an option to set prev_kv by either adding new methods watch_with_options and watch_prefix_with_options that accept a boolean prev_kv (or a small options struct) and use that value when constructing WatchRequest.prev_kv, or by adding an optional prev_kv parameter to the existing watch/watch_prefix signatures; update the request construction in both locations to use the provided prev_kv and keep existing convenience methods delegating to the new variants with prev_kv=false for backwards compatibility.d-engine-core/src/watch/manager_test.rs (1)
651-658:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winAssert registration succeeds before dropping the handle.
reg.register(k, false)returns aResult, but this test drops it without checking. If registration regresses and starts returningErr, the final watcher count still stays at 0 and the cleanup path never gets exercised.Suggested fix
let handle = tokio::spawn(async move { - let watcher = reg.register(k, false); + let watcher = reg + .register(k, false) + .expect("concurrent registration should succeed"); tokio::time::sleep(Duration::from_millis(10)).await; drop(watcher); // Yield to ensure drop completes tokio::task::yield_now().await; });🤖 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 `@d-engine-core/src/watch/manager_test.rs` around lines 651 - 658, The test currently calls reg.register(k, false) and immediately drops the returned value without checking the Result; change the spawned task to assert the registration succeeded by handling the Result from reg.register(k, false) (e.g., unwrap or assert_ok) into a watcher variable before the sleep and drop so any Err will fail the test and exercise the cleanup path (reference: reg.register and the watcher variable created in the async move block).d-engine-core/src/state_machine_handler/default_state_machine_handler.rs (1)
239-245:⚠️ Potential issue | 🟠 Major | ⚡ Quick winAdvance
last_appliedbefore emitting watch events.Progress heartbeats read
last_applied, but this code broadcasts data events first and only updateslast_appliedlater. That allows a client to observePUT/DELETEat revisionNand then aProgressevent at revision< N, which violates the monotonic revision contract. Move thelast_appliedupdate/notify ahead of the watch broadcast, or clamp heartbeat revisions to the highest revision already emitted.🤖 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 `@d-engine-core/src/state_machine_handler/default_state_machine_handler.rs` around lines 239 - 245, The watch broadcast currently emits events (via broadcast_watch_events called where apply_result and watch_event_tx are checked) before advancing/notifying last_applied, which can let Progress heartbeats report lower revisions; move the logic that updates/sets last_applied and sends its notification (the code that advances last_applied/notify in DefaultStateMachineHandler) to execute before the block that calls broadcast_watch_events (or, alternatively, clamp heartbeat/Progress revisions to the highest revision emitted), so that last_applied is monotonically advanced prior to calling broadcast_watch_events(&chunk, results, tx, prev_values.as_deref()) when apply_result is Ok.
🧹 Nitpick comments (2)
examples/service-discovery-standalone/watcher.rs (1)
254-254: ⚡ Quick winAvoid registry re-print on heartbeat
Progressevents.Returning
falseon Line 254 still flows into unconditionalprint_registry(...)in the caller, so heartbeat mode can spam logs. Consider gating printouts toPut/Deletemutations only.🤖 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 `@examples/service-discovery-standalone/watcher.rs` at line 254, The handler currently returns false for WatchEventType::Progress but the caller still calls print_registry(...) unconditionally, causing heartbeat spamming; update the caller to only call print_registry when the event represents a mutation (e.g., WatchEventType::Put or WatchEventType::Delete) or change the handler to return an explicit flag/enum indicating "mutated" vs "heartbeat" and have the caller check that before invoking print_registry; reference WatchEventType::Progress, WatchEventType::Put/WatchEventType::Delete, and print_registry to locate the logic to gate the print.d-engine-client/src/grpc_client_test.rs (1)
1592-1607: ⚡ Quick winAssert
prev_valuein the watch stream test.The fixture now includes
prev_value, but the test never validates it. Adding explicit assertions will lock in the new response contract.Suggested test assertions
let ev1 = stream.next().await.expect("expected first event").expect("no error"); assert_eq!(ev1.key, Bytes::from("lock")); assert_eq!(ev1.value, Bytes::from("owner-a")); + assert_eq!(ev1.prev_value, Bytes::new()); assert_eq!(ev1.event_type, WatchEventType::Put as i32); // Receive second event: Delete let ev2 = stream.next().await.expect("expected second event").expect("no error"); assert_eq!(ev2.key, Bytes::from("lock")); + assert_eq!(ev2.prev_value, Bytes::new()); assert_eq!(ev2.event_type, WatchEventType::Delete as i32);🤖 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 `@d-engine-client/src/grpc_client_test.rs` around lines 1592 - 1607, The watch stream test creates WatchResponse fixtures including a prev_value but never asserts it; update the test that consumes these WatchResponse entries (in grpc_client_test.rs) to assert the prev_value field for both events: for the Put event assert prev_value == Bytes::new() or the expected previous contents, and for the Delete event assert prev_value == Bytes::new() (or the appropriate previous value used in the fixture). Locate the WatchResponse structures and add assertions against response.prev_value (and use the same Bytes constructors as in the fixture: Bytes::from(...) or Bytes::new()) alongside the existing assertions for key/value/event_type to lock the contract.
🤖 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
`@d-engine-core/src/state_machine_handler/default_state_machine_handler_test.rs`:
- Around line 2111-2148: The test currently never exercises the 1→0 transition
because prev_kv_count is initialized to 0; change it to drive two applies:
initialize prev_kv_count = Arc::new(AtomicUsize::new(1)), configure the
MockStateMachine to expect a get() (or otherwise allow one read) on the first
apply and expect_apply_chunk for both applies, then after the first
handler.apply_chunk(...) set prev_kv_count.store(0, Ordering::SeqCst) and call
handler.apply_chunk(...) again and assert no panic / no get() on the second
apply. Update references:
test_apply_chunk_stops_reading_prev_value_after_count_drops_to_zero,
prev_kv_count, DefaultStateMachineHandler::new, handler.apply_chunk,
MockStateMachine.expect_get / expect_apply_chunk, and insert_entry.
In `@d-engine-core/src/state_machine_handler/default_state_machine_handler.rs`:
- Around line 643-679: read_prev_values currently reads every key from sm alone
so later entries in the same chunk see the original state rather than any
earlier in-chunk mutations; fix it by simulating an in-chunk overlay: iterate
chunk entries in order in read_prev_values, for each entry decode the
WriteCommand (as you already do), determine the key (from Operation::Insert /
Delete / CompareAndSwap), compute prev_value by first checking an in-memory
HashMap overlay (key -> bytes::Bytes) and if missing falling back to
sm.get(key).ok().flatten().unwrap_or_default(), then push that prev_value
(Some(bytes) for Insert/Delete/CAS, None for Noop/Config) and update the overlay
to the post-apply value for that operation (Insert -> ins.value, Delete -> empty
Bytes, CompareAndSwap -> the operation’s new value) so subsequent entries see
earlier changes; use the same types referenced in this function (Entry,
WriteCommand, Operation, Payload, sm.get) when locating code to change.
In `@d-engine-core/src/watch/manager_test.rs`:
- Around line 1335-1349: The test drops the broadcast sender returned by
setup_watch_system_with_heartbeat immediately (via the anonymous first tuple
element), which can shut down the dispatcher and make the test vacuously pass;
preserve the sender so the dispatcher remains alive by capturing the first tuple
element (e.g. bind broadcast_tx or _broadcast_tx) instead of ignoring it, keep
that variable in scope until after the watcher is registered and the sleep
completes, and only drop it at the end of the test (or let it naturally go out
of scope) so the assertion truly verifies heartbeat_interval_ms = 0 behavior for
register/receiver interactions.
In `@d-engine-core/src/watch/manager.rs`:
- Around line 493-513: The code currently always constructs a
tokio::time::interval_at and does Instant::now() +
Duration::from_millis(u64::MAX) when heartbeat_enabled is false, which can
overflow; instead, change the logic so you do not build the interval at all when
heartbeats are disabled: make the heartbeat variable optional (e.g.,
Option<tokio::time::Interval>) and only create the interval_at using
first_tick_ms and base_ms inside the heartbeat_enabled branch; when disabled set
heartbeat = None (or another safe sentinel) and update any downstream uses of
heartbeat to handle the None case. Ensure you reference the existing symbols
heartbeat_enabled, first_tick_ms, base_ms, and tokio::time::interval_at when
making the change.
In `@d-engine-proto/proto/client/client_api.proto`:
- Around line 206-210: Change the bytes prev_value field to a presence-bearing
type and propagate that through the WatchEvent API: replace the scalar "bytes
prev_value = 6;" with a wrapped type (e.g. google.protobuf.BytesValue) or use
proto3 "optional bytes prev_value" so callers can distinguish "absent" vs
"present-empty". Update the proto import to include the wrapper
(google/protobuf/wrappers.proto) if using BytesValue, update the WatchEvent
message and any other messages or RPCs that expose prev_value to use the new
type, and update all producer/consumer code that constructs or reads WatchEvent
to check the wrapper’s presence (.value vs null/has_value or optional presence)
rather than relying on empty bytes. Ensure generated code and tests are updated
accordingly.
In `@d-engine-server/src/api/embedded_client.rs`:
- Around line 462-463: Update the misleading rustdoc that mentions a `prev_kv`
parameter for the `watch()` function: remove or reword the sentence describing
`prev_kv` and instead point users to `watch_with_options` (which accepts that
option) or document how to obtain previous values via `watch_with_options`.
Modify the doc comment adjacent to the `watch()` declaration (the doc block
guarded by #[cfg(feature = "watch")]) to reference `watch_with_options` by name
and remove the incorrect parameter description so the docs accurately reflect
the `watch()` signature.
In `@d-engine/src/docs/server_guide/watch-feature.md`:
- Around line 208-209: Clarify that the `error == WATCH_BUFFER_OVERFLOW`
detection applies only to gRPC/proto transport events: update the lines
referencing `WatchEventType::Canceled` and `ErrorCode::WATCH_BUFFER_OVERFLOW`
(e.g., the text mentioning WatchEventType::Canceled and
ErrorCode::WATCH_BUFFER_OVERFLOW) to state that embedded watchers may not expose
an `error` field and instead match cancellation only via
`WatchEventType::Canceled`, and instruct clients using gRPC/proto to check
`error == WATCH_BUFFER_OVERFLOW` while embedded/in-process watchers should rely
on the canceled event and re-sync via the Read API and re-register accordingly.
- Around line 220-223: Update the example to use the new API: replace the
deprecated watch_with_prev_kv call with the new watch_with_options invocation
and pass prev_kv as an option (e.g., call watch_with_options(..., prev_kv)) so
the snippet compiles when copy-pasted; keep the rest of the example (the watcher
variable and the event loop using watcher.receiver_mut().recv().await) unchanged
and ensure the example references the watch_with_options symbol instead of
watch_with_prev_kv.
---
Outside diff comments:
In `@d-engine-client/src/grpc_client.rs`:
- Around line 217-228: The public gRPC watch API currently hardcodes
WatchRequest.prev_kv = false in the watch and watch_prefix implementations (see
methods watch and watch_prefix in grpc_client.rs), making prev_kv inaccessible;
add an option to set prev_kv by either adding new methods watch_with_options and
watch_prefix_with_options that accept a boolean prev_kv (or a small options
struct) and use that value when constructing WatchRequest.prev_kv, or by adding
an optional prev_kv parameter to the existing watch/watch_prefix signatures;
update the request construction in both locations to use the provided prev_kv
and keep existing convenience methods delegating to the new variants with
prev_kv=false for backwards compatibility.
In `@d-engine-core/src/state_machine_handler/default_state_machine_handler.rs`:
- Around line 239-245: The watch broadcast currently emits events (via
broadcast_watch_events called where apply_result and watch_event_tx are checked)
before advancing/notifying last_applied, which can let Progress heartbeats
report lower revisions; move the logic that updates/sets last_applied and sends
its notification (the code that advances last_applied/notify in
DefaultStateMachineHandler) to execute before the block that calls
broadcast_watch_events (or, alternatively, clamp heartbeat/Progress revisions to
the highest revision emitted), so that last_applied is monotonically advanced
prior to calling broadcast_watch_events(&chunk, results, tx,
prev_values.as_deref()) when apply_result is Ok.
In `@d-engine-core/src/watch/manager_test.rs`:
- Around line 651-658: The test currently calls reg.register(k, false) and
immediately drops the returned value without checking the Result; change the
spawned task to assert the registration succeeded by handling the Result from
reg.register(k, false) (e.g., unwrap or assert_ok) into a watcher variable
before the sleep and drop so any Err will fail the test and exercise the cleanup
path (reference: reg.register and the watcher variable created in the async move
block).
---
Nitpick comments:
In `@d-engine-client/src/grpc_client_test.rs`:
- Around line 1592-1607: The watch stream test creates WatchResponse fixtures
including a prev_value but never asserts it; update the test that consumes these
WatchResponse entries (in grpc_client_test.rs) to assert the prev_value field
for both events: for the Put event assert prev_value == Bytes::new() or the
expected previous contents, and for the Delete event assert prev_value ==
Bytes::new() (or the appropriate previous value used in the fixture). Locate the
WatchResponse structures and add assertions against response.prev_value (and use
the same Bytes constructors as in the fixture: Bytes::from(...) or Bytes::new())
alongside the existing assertions for key/value/event_type to lock the contract.
In `@examples/service-discovery-standalone/watcher.rs`:
- Line 254: The handler currently returns false for WatchEventType::Progress but
the caller still calls print_registry(...) unconditionally, causing heartbeat
spamming; update the caller to only call print_registry when the event
represents a mutation (e.g., WatchEventType::Put or WatchEventType::Delete) or
change the handler to return an explicit flag/enum indicating "mutated" vs
"heartbeat" and have the caller check that before invoking print_registry;
reference WatchEventType::Progress, WatchEventType::Put/WatchEventType::Delete,
and print_registry to locate the logic to gate the print.
🪄 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: defaults
Review profile: CHILL
Plan: Pro
Run ID: ddcf0594-21a1-433c-9d21-59f37668d976
⛔ Files ignored due to path filters (1)
d-engine-proto/src/generated/d_engine.client.rsis excluded by!**/generated/**
📒 Files selected for processing (20)
d-engine-client/src/grpc_client.rsd-engine-client/src/grpc_client_test.rsd-engine-core/src/config/raft.rsd-engine-core/src/state_machine_handler/default_state_machine_handler.rsd-engine-core/src/state_machine_handler/default_state_machine_handler_test.rsd-engine-core/src/watch/manager.rsd-engine-core/src/watch/manager_test.rsd-engine-core/src/watch/mod.rsd-engine-proto/proto/client/client_api.protod-engine-server/benches/state_machine.rsd-engine-server/src/api/embedded_client.rsd-engine-server/src/api/embedded_client_test.rsd-engine-server/src/network/grpc/grpc_raft_service.rsd-engine-server/src/node/builder.rsd-engine-server/src/test_utils/integration/mod.rsd-engine-server/tests/embedded_client/embedded_client_operations.rsd-engine-server/tests/watch_and_subscriptions/watch_events_embedded.rsd-engine/src/docs/server_guide/watch-feature.mdexamples/service-discovery-embedded/server.rsexamples/service-discovery-standalone/watcher.rs
Codecov Report❌ Patch coverage is 📢 Thoughts on this report? Let us know! |
- MissedTickBehavior::Skip on heartbeat interval: Burst mode would re-deliver N identical-revision Progress events after a stall, misleading clients into thinking the stream was alive during the gap - WatchEvent.prev_value: Option<Bytes> — None means "not requested", Some(empty) means "key didn't exist"; eliminates the ambiguity of a bare Bytes::new() standing for two distinct states - read_prev_values overlay map: two writes to the same key in one apply_chunk batch now produce correct sequential prev_values; previously both would read the pre-batch SM state - Jitter seed: replace ThreadId byte-sum with DefaultHasher mixing ThreadId + SystemTime; covers k8s rolling-restart where subsec_nanos alone is near-identical across pods started in the same millisecond - Upgrade astral-tokio-tar 0.6.1 -> 0.6.2 (RUSTSEC-2026-0145: PAX header desynchronization allows file smuggling during extraction) - Test: test_apply_chunk_stops_reading_prev_value_after_count_drops_to_zero now exercises the actual 1->0 transition instead of starting at 0 - Docs/examples: watch-feature.md and service-discovery example updated for Option<Bytes> semantics, correct API names, and gRPC-scoped WATCH_BUFFER_OVERFLOW detection note
What Does This PR Do?
Adds
prev_kvsupport to watch registrations (each mutation event carries the value the key held before the write) and a configurable progress heartbeat that periodically signals all active watchers even when no writes occur.Type:
Why Is This Needed?
For features: Issue #379
prev_kveliminates the extra Read round-trip developers currently need when they want to know what a key held before it changed (audit logs, distributed lock handoff, CAS-based state machines). Without it, the pattern iswatch → event arrives → get(key)with an unavoidable race window.Progress heartbeat lets clients distinguish a healthy-but-idle watch stream from a silently broken connection, removing the need for application-level keepalives.
Checklist
Required:
make testpassesIf changing APIs:
watch-feature.md, rustdoc onwatch_with_options/watch_prefix_with_options)Testing
How tested:
prev_kv_watcher_countatomic tracks registration/deregistration correctly; progress heartbeat delivered at configured interval; no Progress events when heartbeat disabledget()when count=0 (mockall panics if called unexpectedly); reads and broadcastsprev_valuewhen count=1; stops reading after count drops back to zeroembedded_client_operationsandwatch_events_embeddedsuites pass with updated APIFor bug fixes: n/a
For performance improvements: n/a — prev_kv read is off by default; only active when at least one watcher opts in
Does This Follow d-engine's Principles?
Arc<AtomicUsize>counter; no per-event branching when no prev_kv watchers existwatch()/watch_prefix()default toprev_kv=false; new surface iswatch_with_options/watch_prefix_with_optionsonlyReviewer Notes
Key design decision to review:
prev_kv_watcher_countis a sharedArc<AtomicUsize>betweenWatchRegistryandDefaultStateMachineHandler. The handler reads it on everyapply_chunk— zero cost when count=0, oneget()per write batch when count>0. Alternative (per-event flag in broadcast message) was rejected because it would require reading old value even for non-watched keys.Progress heartbeat routing:
broadcast_progress()cannot use the normal key-routing path (Progress has an empty key, which routes to nobody). It snapshots all registered keys first, then dispatches to each. DashMap shard locks are released between snapshot and dispatch to avoid potential deadlock.Estimated review complexity:
Summary by CodeRabbit
New Features
Documentation
Examples & Tests