feat(watch) #300: prefix watch, watcher limits, and revision field - #380
Conversation
- WatchRegistry: add prefix DashMap with prefix_segments() O(depth) dispatch - register_prefix(): slash-boundary validation, LimitExceeded on cap reached - WatchConfig: add max_watcher_count (default unlimited); watcher_buffer_size 10→256 - Proto: WatchRequest.prefix: bool, WatchResponse.revision: uint64 - GrpcClient + EmbeddedClient: add watch_prefix() - 5 new integration tests in watch_events_embedded.rs - service-discovery examples: add --prefix mode, live HashMap registry, verified CLI output - watch-feature.md: Prefix Watch section, updated config + architecture diagram
handle_put() silently returned 500 for NotLeader. Add is_not_leader()
to catch both error variants; return 503 with {"error":"not_leader"}.
README: show find-leader-first pattern via /primary health check.
Also raise zombie test deadline 20s→35s (split-vote flakiness on slow CI).
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: defaults Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (15)
💤 Files with no reviewable changes (2)
✅ Files skipped from review due to trivial changes (4)
🚧 Files skipped from review as they are similar to previous changes (4)
📝 WalkthroughWalkthroughThis PR extends the watch subsystem with slash-delimited prefix watches, propagates a Raft-applied-index ChangesWatch Prefix + Revision + Limits Feature
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes 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: 10
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
d-engine-core/src/watch/mod.rs (1)
139-147:⚠️ Potential issue | 🟠 Major | ⚡ Quick winPreserve revision on CANCELED events instead of forcing
0.Line 146 hardcodes
revision: 0, which loses the event anchor clients need afterCANCELED. Use the triggering event’s revision when constructing the sentinel.Proposed fix
-pub(crate) fn make_cancel_event(key: bytes::Bytes) -> WatchEvent { +pub(crate) fn make_cancel_event( + key: bytes::Bytes, + revision: u64, +) -> WatchEvent { @@ - revision: 0, + revision, } }// d-engine-core/src/watch/manager.rs let _ = watcher .sender .try_send(crate::watch::make_cancel_event(event.key.clone(), event.revision));🤖 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/mod.rs` around lines 139 - 147, The CANCELED sentinel currently hardcodes revision: 0 in make_cancel_event, which loses the triggering event anchor; update make_cancel_event to accept a revision parameter (e.g., fn make_cancel_event(key: bytes::Bytes, revision: i64) -> WatchEvent) and set WatchEvent.revision to that value while keeping value empty, event_type = WatchEventType::Canceled and error = ErrorCode::WatchBufferOverflow; then update callers (e.g., the watcher.sender.try_send invocation in manager code) to pass the original event.revision into make_cancel_event so the CANCELED event preserves the triggering revision.examples/service-discovery-standalone/admin.rs (1)
97-105:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winUpdate the printed index-key hint to match the new path format.
The message still suggests
services/{name}_index, but the code now reads/services/{name}_index.Suggested patch
- println!("Workaround: Maintain an index key like 'services/{name}_index'"); + println!("Workaround: Maintain an index key like '/services/{name}_index'");🤖 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/admin.rs` around lines 97 - 105, Update the user-facing hint in the Commands::List match arm so the printed index key matches the actual path used by the code: change the suggested index key string from "services/{name}_index" to "/services/{name}_index" (the code constructs index_key via let index_key = format!("/services/{name}_index");).
🧹 Nitpick comments (1)
d-engine-client/src/grpc_client_test.rs (1)
1592-1605: ⚡ Quick winAssert
revisionin the watch success test as part of the new response contract.Now that
revisionis explicit in fixtures, add assertions for both events so this test also guards client-side field propagation.🤖 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 - 1605, Update the "watch success" test to assert the WatchResponse.revision field for both events so the test verifies client-side propagation; after locating the generated responses (the two WatchResponse instances with WatchEventType::Put and WatchEventType::Delete), add assertions that their .revision equals the expected fixture value (0 in the shown fixtures) for both the put and delete responses, using the same response variables or indexing used in the test.
🤖 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/watch/manager_test.rs`:
- Around line 811-812: Replace brittle
tokio::time::sleep(Duration::from_millis(50)).await-based waits with a
deterministic readiness check: instead of sleeping, poll or await a specific
condition (e.g., that the watcher is registered in the registry or that a
recv/try_recv returns a value) with a bounded timeout. Concretely, replace
occurrences of tokio::time::sleep(...) in manager_test.rs with a small loop that
repeatedly checks the watcher registration/registry visibility or attempts to
receive (using tokio::time::timeout + short interval retries) until success or
deadline, and factor this into a helper like wait_until_registered or
recv_with_deadline and use it at all listed sites (811, 838-839, 863-864,
890-891, 917-918, 947-948, 1010-1011, 1074-1075) to eliminate flakiness.
In `@d-engine-core/src/watch/manager.rs`:
- Around line 238-254: Atomically reserve a watcher slot using
self.total_count.fetch_update (or compare_exchange loop) instead of the current
separate check and later fetch_add: attempt to increment total_count only if
current < self.max_watcher_count and return
Err(WatchError::LimitExceeded(self.max_watcher_count)) if it fails; after a
successful atomic reserve proceed to create the mpsc::channel, Watcher { id,
sender } and insert into self.prefix or self.exact; remove the later
self.total_count.fetch_add(1, ...) so you don't double-count, and if insertion
can fail (unlikely) ensure you undo the reserved increment by decrementing
total_count.
In `@d-engine-server/src/network/grpc/grpc_raft_service.rs`:
- Around line 498-504: The current branch that handles prefix registrations
always maps registry.register_prefix errors to Status::invalid_argument; change
the error mapping for registry.register_prefix so that when the registry returns
a LimitExceeded (or equivalent watcher-limit) error you map it to
Status::resource_exhausted(e.to_string())—mirroring the exact-watch path used
for registry.register—while preserving invalid_argument for genuine argument
errors; update the is_prefix branch around register_prefix to inspect the error
kind (e.g., LimitExceeded) and return resource_exhausted for that case and
invalid_argument otherwise.
In `@d-engine-server/tests/watch_and_subscriptions/watch_events_embedded.rs`:
- Around line 445-446: Replace the brittle fixed sleep calls used in the prefix
integration tests with a deterministic readiness/await helper (e.g., create and
use a helper like await_registration_ready_or_event that polls/awaits the
specific condition) instead of sleep(Duration::from_millis(50)). Update each
occurrence referenced (the sleep at around 50ms and the 200ms negative-assertion
windows) to use that helper with a bounded timeout (e.g., 500–1000ms) and an
explicit failure if the condition isn't met; adjust calls in the tests that rely
on startup/registration jitter (the sleeps around lines noted) to call the
helper and lengthen the negative assertion window to be more robust (e.g., from
200ms to ~500ms).
In `@examples/service-discovery-embedded/README.md`:
- Line 15: Three fenced code blocks in the README are missing language
specifiers; update each opening triple backtick to include a language (use
"text") for the ASCII-art and log examples so markdownlint MD040 is satisfied.
Specifically, change the block that begins with the ASCII box
(┌────────────────...), the block containing the "[exact-watch] config changed:
/config/payment-service/timeout..." log, and the block containing the
"[prefix-watch] node up: /services/payment/node1 → 10.0.0.1:8080..." log so
their opening fences are ```text instead of just ```.
In `@examples/service-discovery-standalone/README.md`:
- Line 15: Add explicit fenced-code languages (e.g., text) to the three code
blocks in README.md that are currently untyped: the ASCII art block starting
with "┌────────────────────────────────────────────────────────┐", the block
beginning with "=== Current State ===", and the block showing "[PUT ]
/services/payment/node1 = 10.0.0.1:8080 (revision=3)"; update each opening
triple backtick to include the language (```text) so markdownlint MD040 warnings
are resolved.
- Around line 26-30: The architecture diagram uses inconsistent service key
paths (/svc/pay/node1 and /services/pay/) that conflict with the rest of the
docs/scripts which use /services/payment/..., so update the diagram entries that
show "--key /svc/pay/node1" and "--key /services/pay/" to the correct paths
(e.g., "--key /services/payment/node1" and "--key /services/payment/" or the
appropriate /services/payment/... prefix) so the diagram matches the
documented/scripted /services/payment/... hierarchy.
In `@examples/service-discovery-standalone/watcher.rs`:
- Around line 164-188: The current loop uses ? on
client.watch_prefix(prefix).await which makes setup failures propagate instead
of retrying, and it doesn't treat tonic::Status with Code::Cancelled as a
transient reconnect; change the setup to explicitly match the result of
client.watch_prefix(prefix).await (instead of map_err + ?): if Ok(stream)
proceed, if Err(e) inspect e (e.g., if it is tonic::Status with Code::Cancelled
or other transient errors) then log with context (include e), sleep
(tokio::time::sleep) and continue the loop to reconnect; only propagate/return
on fatal/non-retryable errors. Keep the existing stream.next() error handling
but ensure Err(e) from the stream that is a Cancelled also breaks to the
reconnect path (reuse the same log/sleep pattern). Refer to client.watch_prefix,
prefix, stream, stream.next(), registry, and tokio::time::sleep when applying
changes.
In `@examples/three-nodes-embedded/README.md`:
- Around line 114-146: The README currently documents leader-only writes (health
check /primary and followers returning 503 {"error":"not_leader"}) but later
text still claims writes can be sent to any node and will be forwarded; update
that section to match the actual behavior: change the statement that writes can
be sent to any node and forwarded to instead state that writes must be sent to
the leader and followers will return 503 with {"error":"not_leader"}, and adjust
the example or port references (the /primary health check and /kv POST examples)
so both places consistently show finding the leader via /primary and sending
POSTs only to the leader (no forwarding). Ensure the README references the
/primary health endpoint and the /kv write/read examples (/kv POST and /kv/:key
GET) consistently.
In `@examples/three-nodes-embedded/src/main.rs`:
- Around line 222-225: The response currently returns the debug-formatted
internal error using format!("{e:?}") in the handler (producing
StatusCode::INTERNAL_SERVER_ERROR and Json(...)), which can leak internals;
change the handler to return a stable, generic error message in the Json body
(e.g., {"error":"internal server error"}) and move the detailed logging of e to
server-side logs using your logger/tracing before returning the response; ensure
you update the code paths that build the Json response (the block creating
StatusCode::INTERNAL_SERVER_ERROR and Json(...)) to stop including
format!("{e:?}") and instead return the generic string while logging e
separately.
---
Outside diff comments:
In `@d-engine-core/src/watch/mod.rs`:
- Around line 139-147: The CANCELED sentinel currently hardcodes revision: 0 in
make_cancel_event, which loses the triggering event anchor; update
make_cancel_event to accept a revision parameter (e.g., fn
make_cancel_event(key: bytes::Bytes, revision: i64) -> WatchEvent) and set
WatchEvent.revision to that value while keeping value empty, event_type =
WatchEventType::Canceled and error = ErrorCode::WatchBufferOverflow; then update
callers (e.g., the watcher.sender.try_send invocation in manager code) to pass
the original event.revision into make_cancel_event so the CANCELED event
preserves the triggering revision.
In `@examples/service-discovery-standalone/admin.rs`:
- Around line 97-105: Update the user-facing hint in the Commands::List match
arm so the printed index key matches the actual path used by the code: change
the suggested index key string from "services/{name}_index" to
"/services/{name}_index" (the code constructs index_key via let index_key =
format!("/services/{name}_index");).
---
Nitpick comments:
In `@d-engine-client/src/grpc_client_test.rs`:
- Around line 1592-1605: Update the "watch success" test to assert the
WatchResponse.revision field for both events so the test verifies client-side
propagation; after locating the generated responses (the two WatchResponse
instances with WatchEventType::Put and WatchEventType::Delete), add assertions
that their .revision equals the expected fixture value (0 in the shown fixtures)
for both the put and delete responses, using the same response variables or
indexing used in the test.
🪄 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: 0cffa624-3ef8-45a6-980b-7f25224d1573
⛔ Files ignored due to path filters (2)
d-engine-proto/src/generated/d_engine.client.rsis excluded by!**/generated/**examples/quick-start-standalone/go.sumis excluded by!**/*.sum
📒 Files selected for processing (25)
CHANGELOG.mdd-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/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/network/grpc/grpc_raft_service.rsd-engine-server/src/node/builder.rsd-engine-server/tests/watch_and_subscriptions/watch_events_embedded.rsd-engine-server/tests/watch_and_subscriptions/watch_membership_embedded.rsd-engine/src/docs/server_guide/watch-feature.mdexamples/quick-start-standalone/go.modexamples/service-discovery-embedded/README.mdexamples/service-discovery-embedded/server.rsexamples/service-discovery-standalone/Makefileexamples/service-discovery-standalone/README.mdexamples/service-discovery-standalone/admin.rsexamples/service-discovery-standalone/watcher.rsexamples/three-nodes-embedded/README.mdexamples/three-nodes-embedded/src/main.rs
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
Codecov Report❌ Patch coverage is 📢 Thoughts on this report? Let us know! |
…s, CANCELED break, docs
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (1)
examples/three-nodes-embedded/src/main.rs (1)
222-225:⚠️ Potential issue | 🟠 Major | ⚡ Quick winDo not return raw internal error messages in 500 responses.
Line 224 currently reflects internal error text back to clients. Return a stable generic payload and log details server-side instead.
Proposed fix
- Err(e) => ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(serde_json::json!({ "error": format!("{e}") })), - ), + Err(e) => { + tracing::error!("PUT /kv failed: {e}"); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({ "error": "internal_server_error" })), + ) + },🤖 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/three-nodes-embedded/src/main.rs` around lines 222 - 225, The handler's Err(e) arm currently echoes internal error text to clients; change the response payload to a stable generic message (e.g., Json(serde_json::json!({ "error": "internal server error" }))) while logging the actual error server-side using your logger (e.g., tracing::error! or log::error!) before returning (keep the StatusCode::INTERNAL_SERVER_ERROR and the Json wrapper). Locate the Err(e) match arm that returns (StatusCode::INTERNAL_SERVER_ERROR, Json(...)) and replace format!("{e}") with the generic string and add a logging call like tracing::error!("handler_name failed: {:?}", e) (or similar) just before returning.
🤖 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-server/tests/watch_and_subscriptions/watch_events_grpc_standalone.rs`:
- Around line 361-373: The test currently only checks that client.watch_prefix
calls (result and result2) return Err; update both assertions to verify the gRPC
status is INVALID_ARGUMENT by extracting the error as a tonic::Status (or
equivalent Status type used in the test harness) and asserting status.code() ==
tonic::Code::InvalidArgument (or the equivalent enum/value). Specifically, after
calling client.watch_prefix(b"/config") and client.watch_prefix(b"config/"),
convert the returned Err into a Status and assert the status code equals
InvalidArgument so the test enforces the exact gRPC contract.
In `@examples/three-nodes-embedded/src/main.rs`:
- Around line 201-207: The current is_not_leader function brittlely parses
free-form error message text from ClientApiError; change it to use the typed
discriminator by returning e.code() == ErrorCode::NotLeader instead of matching
on message contents; update the is_not_leader function to call
ClientApiError::code() and compare to ErrorCode::NotLeader (and add an import
for ErrorCode if missing) so leader detection uses the structured error code.
---
Duplicate comments:
In `@examples/three-nodes-embedded/src/main.rs`:
- Around line 222-225: The handler's Err(e) arm currently echoes internal error
text to clients; change the response payload to a stable generic message (e.g.,
Json(serde_json::json!({ "error": "internal server error" }))) while logging the
actual error server-side using your logger (e.g., tracing::error! or
log::error!) before returning (keep the StatusCode::INTERNAL_SERVER_ERROR and
the Json wrapper). Locate the Err(e) match arm that returns
(StatusCode::INTERNAL_SERVER_ERROR, Json(...)) and replace format!("{e}") with
the generic string and add a logging call like tracing::error!("handler_name
failed: {:?}", e) (or similar) just before returning.
🪄 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: 9c574619-9bcb-458a-95ff-03897fe7e378
📒 Files selected for processing (11)
d-engine-core/src/watch/manager.rsd-engine-core/src/watch/manager_test.rsd-engine-core/src/watch/mod.rsd-engine-server/src/network/grpc/grpc_raft_service.rsd-engine-server/tests/watch_and_subscriptions/watch_events_embedded.rsd-engine-server/tests/watch_and_subscriptions/watch_events_grpc_standalone.rsexamples/service-discovery-embedded/README.mdexamples/service-discovery-standalone/README.mdexamples/service-discovery-standalone/watcher.rsexamples/three-nodes-embedded/README.mdexamples/three-nodes-embedded/src/main.rs
✅ Files skipped from review due to trivial changes (2)
- examples/three-nodes-embedded/README.md
- examples/service-discovery-embedded/README.md
🚧 Files skipped from review as they are similar to previous changes (6)
- d-engine-core/src/watch/mod.rs
- examples/service-discovery-standalone/README.md
- d-engine-server/tests/watch_and_subscriptions/watch_events_embedded.rs
- examples/service-discovery-standalone/watcher.rs
- d-engine-core/src/watch/manager.rs
- d-engine-core/src/watch/manager_test.rs
- Export WatchEventType/WatchEvent/WatchError/WatcherHandle via d_engine:: (watch feature) - Add __test_support feature exposing state_machine_test / storage_engine_test - service-discovery-standalone: replace d-engine-client dep with d-engine facade - service-discovery-embedded: remove direct d-engine-core dep - sled-cluster: remove direct d-engine-proto dep - Docs: fix all internal-crate imports to use d_engine::
What Does This PR Do?
Extends the Watch API with prefix watch (
watch_prefix()), a hard cap on watcher count (max_watcher_count), and arevisionfield on every event so callers can detect missed events after reconnection.Type:
Why Is This Needed?
For features: Link to approved issue #300
Without prefix watch, an API gateway monitoring N service nodes must register N separate exact-key watchers. One
watch_prefix("/services/payment/")replaces all of them — any node joining, updating, or leaving fires a single stream. The revision field enables gap detection on reconnect.Checklist
Required:
make testpassesIf changing APIs:
Testing
How tested:
manager_test.rs— prefix registration, dispatch, LimitExceeded, InvalidPrefix, prefix_segments decompositionwatch_events_embedded.rs— prefix fires for child keys, does not fire for parent key, root prefix, revision increases monotonically, handle drop cleans up watchermake run-watcher-prefix+make register-node1/2/3)For bug fixes:
test_embedded_prefix_watch_does_not_fire_for_parent_key)Does This Follow d-engine's Principles?
Reviewer Notes
Architecture: Two separate
DashMaps (exact + prefix).prefix_segments(key)decomposes an event key into its slash-terminated prefix candidates — O(depth) lookups, no linear scan over all registered prefixes.Proto backward compat:
WatchRequest.prefixdefaults tofalse(existing exact-key behavior unchanged).WatchResponse.revisiondefaults to0(old clients ignore it).Also includes:
fix(example) #300—three-nodes-embeddednow returns 503 with{"error":"not_leader"}instead of 500 for follower writes; README updated to show find-leader-first pattern.Estimated review complexity:
Summary by CodeRabbit
New Features
Configuration
Improvements
Tests