Repository navigation
feat(scan) #378: add scan_prefix API for zero-race-window watch reconnection - #382
Conversation
…nection
Add linearizable prefix scan across all storage backends and client interfaces,
enabling clients to reconnect watches without missing events.
Protocol:
- Add KvEntry, ScanRequest, ScanResponse messages to client_api.proto
- Add HandleClientScan RPC to RaftClientService
Core:
- Add ScanResult{entries, revision} to StateMachine trait (default impl returns error)
- Add ClientCmd::Scan variant; leader serves immediately from state machine,
followers reject with failed_precondition("Not leader")
- Add scan_prefix() to ClientApi trait
Storage backends:
- RocksDB: set_iterate_upper_bound(prefix_successor) — O(results), no full scan
- File: HashMap::iter().filter(starts_with) — O(total keys), unavoidable for HashMap
- Sled: sled::Tree::scan_prefix() — O(results), native B-tree support
Server/Client:
- Add handle_client_scan gRPC handler in grpc_raft_service
- Add scan_prefix() to EmbeddedClient and GrpcClient
- Re-export ScanResult from d-engine-server and d-engine crates
Tests:
- 4 RocksDB unit tests covering prefix isolation, upper-bound exclusion,
empty prefix, and revision correctness
Docs/Examples:
- Update watch-feature.md with real Watch->Scan->filter reconnect pattern
- Update service-discovery-standalone watcher.rs to use scan_prefix()
…ne trait impl impl StateMachine for RocksDBStateMachine was missing scan_prefix delegation. The trait default impl always returns an error; any call through an SM: StateMachine bound (all production paths via SMOF<T>) silently failed. The pub fn scan_prefix at impl RocksDBStateMachine was unreachable through the trait — only direct struct dispatch (e.g. unit tests) could hit it.
|
Warning Rate limit exceeded
You’ve run out of usage credits. Purchase more in the billing tab. ⌛ How to resolve this issue?After the wait time has elapsed, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans have higher rate limits than the trial, open-source and free plans. In all cases, we re-allow further reviews after a brief timeout. Please see our FAQ for further information. ℹ️ Review info⚙️ Run configurationConfiguration used: defaults Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (3)
📝 WalkthroughWalkthroughAdds a linearizable prefix-scan API end-to-end: proto messages and RPC, ChangesPrefix Scan Implementation
Sequence Diagram(s)sequenceDiagram
participant Client
participant GrpcClient
participant RaftLeader
participant StateMachine
participant WatchStream
Note over Client,WatchStream: Zero-race-window reconnection (watch-first, scan-second)
Client->>WatchStream: register prefix watch
Client->>GrpcClient: scan_prefix(prefix)
GrpcClient->>RaftLeader: HandleClientScan(ScanRequest)
RaftLeader->>StateMachine: scan_prefix(prefix)
StateMachine-->>RaftLeader: ScanResult{entries, scan_revision}
RaftLeader-->>GrpcClient: ScanResponse{entries, revision}
GrpcClient-->>Client: ScanResult
Note over Client: Initialize local registry from entries
loop Drain watch events
WatchStream-->>Client: WatchEvent{key, value, revision}
alt revision <= scan_revision
Note over Client: Discard (already included)
else
Note over Client: Apply to registry
end
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related issues
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 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: 5
🤖 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-client/src/grpc_client.rs`:
- Around line 535-542: The current map_err on the handle_client_scan RPC
collapses all tonic gRPC Statuses into ErrorCode::Uncategorized; change the
map_err closure for client.handle_client_scan(request) to inspect the
tonic::Status (status.code() and any metadata/trailers) and map to the
appropriate d_engine_proto::error::ErrorCode variants (e.g., NotLeader,
NetworkError, Timeout, etc.), preserving status.message() and any
leader/redirect info in the Business error payload; ensure the mapping logic
mirrors other RPCs’ error mapping helpers (use the same helper or inline logic)
so callers can detect NotLeader/network semantics for retries and leader
routing.
In `@d-engine-proto/proto/client/client_api.proto`:
- Around line 109-119: The ScanRequest's consistency_policy is ignored: update
handle_client_scan to copy consistency_policy into ClientCmd::Scan so the
command carries the requested policy, and in leader_state's push_client_cmd (the
branch handling ClientCmd::Scan and the call to
ctx.state_machine().scan_prefix(&prefix)) enforce the same
noop-commit/read-index logic used by LinearizableRead—i.e., if the effective
policy is linearizable (or omitted/default), reject with the same
"LeaderNotReady: noop not committed" until noop has committed or perform the
equivalent read-index check; if the policy is eventual, proceed immediately.
Ensure error messages and command flow mirror LinearizableRead for consistency
and that ctx.state_machine().scan_prefix is only called after the required
verification.
In `@d-engine-server/src/api/embedded_client.rs`:
- Around line 528-529: The current map_err call flattens gRPC status into a
generic server_error string (result.map_err(|status| server_error(format!("RPC
error: {}", status.message())))), losing status.code() classification; change
this to preserve and propagate the original tonic::Status (or create/return a
ServerError variant that wraps the tonic::Status) instead of formatting only the
message—e.g., replace the closure with one that returns Err(status) or
ServerError::Rpc(status) (and adjust the function return type/signature and
callers accordingly) so callers can inspect status.code() for
retries/leader-handling.
In `@d-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine.rs`:
- Around line 598-612: The prefix successor logic is wrong for empty prefixes
and when the last byte(s) are 0xFF; instead of returning an empty ScanResult for
empty prefix or doing a single byte wrapping add, compute the correct exclusive
upper bound by walking upper from the end, popping trailing 0xFF bytes until you
find a byte < 0xFF, increment that byte and truncate upper to that position+1,
and only call opts.set_iterate_upper_bound(upper) when such an upper exists; if
all bytes were 0xFF (or the original prefix is empty) leave no upper bound so
the iterator scans to the end; update the code around variables
prefix/upper/last_mut and the ReadOptions set_iterate_upper_bound usage and
remove the early return that used last_applied_index/ScanResult for empty
prefix.
In `@examples/service-discovery-standalone/watcher.rs`:
- Around line 164-167: run_prefix_watch currently propagates errors from
client.watch_prefix(...) causing the watcher to exit on transient failures;
instead, inside the reconnect loop surrounding the watch, catch/map the error
from client.watch_prefix(prefix).await, log or record the failure, wait/backoff
briefly, and then continue the loop to retry establishing the watch; only exit
the loop on explicit shutdown, not on a watch registration error. Locate the
call to client.watch_prefix(...) and move its error handling into the reconnect
loop (retrying) rather than returning the anyhow error from run_prefix_watch.
🪄 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: a56df6da-16fd-4503-9ee3-97f7cadbfedb
⛔ Files ignored due to path filters (1)
d-engine-proto/src/generated/d_engine.client.rsis excluded by!**/generated/**
📒 Files selected for processing (22)
d-engine-client/src/grpc_client.rsd-engine-client/src/mock_rpc.rsd-engine-core/src/client/client_api.rsd-engine-core/src/event.rsd-engine-core/src/raft_role/leader_state.rsd-engine-core/src/raft_role/role_state.rsd-engine-core/src/storage/state_machine.rsd-engine-core/src/test_utils/mock/mock_rpc.rsd-engine-proto/proto/client/client_api.protod-engine-server/src/api/embedded_client.rsd-engine-server/src/lib.rsd-engine-server/src/network/grpc/grpc_raft_service.rsd-engine-server/src/storage/adaptors/file/file_state_machine.rsd-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine.rsd-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine_test.rsd-engine-server/src/test_utils/mock/mock_rpc.rsd-engine/src/docs/server_guide/watch-feature.mdexamples/service-discovery-embedded/README.mdexamples/service-discovery-embedded/server.rsexamples/service-discovery-standalone/README.mdexamples/service-discovery-standalone/watcher.rsexamples/sled-cluster/src/sled_state_machine.rs
Codecov Report❌ Patch coverage is 📢 Thoughts on this report? Let us know! |
…, 0xFF prefix boundary grpc_client.rs: use ClientApiError::from(status) for scan RPC failures instead of hardcoded Uncategorized, consistent with all other RPCs and preserving NotLeader/ network semantics for caller retry logic. proto + generated: remove consistency_policy from ScanRequest — field was dead code (handler extracted only prefix, policy was silently dropped). Documented in #381 where scan consistency guarantees will be addressed together with LinearizableRead. rocksdb_state_machine.rs: fix prefix_successor carry propagation for 0xFF byte suffix. Previous wrapping_add(1) turned 0xFF → 0x00, producing an upper bound lexicographically below the prefix — RocksDB iterator returned 0 rows silently. Now trims trailing 0xFF bytes before incrementing; returns None (no upper bound, starts_with guard used instead) when all bytes are 0xFF. watcher.rs: watch_prefix setup failure in reconnect loop now continues to next iteration instead of propagating error out of the function — consistent with how scan errors are handled and correctly demonstrates the resilient reconnect pattern. New test: test_scan_prefix_0xff_suffix_carry_propagation verifies the fix. client_read_test.rs: add RED regression test for #381 (LinearizableRead served without quorum in minority partition) — documents the bug, expected to fail until #381 is fixed.
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 `@d-engine-core/src/raft_role/leader_state_test/client_read_test.rs`:
- Line 2174: Replace the single-line assert with a rustfmt-friendly multi-line
form: update the assert! that checks !state.cluster_metadata.single_voter (the
precondition) to use the block style with each argument on its own line, e.g.
call assert!(\n !state.cluster_metadata.single_voter,\n "precondition:
multi-voter cluster",\n); so rustfmt accepts it.
In `@d-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine_test.rs`:
- Around line 296-297: The second assert_eq! is a single long line that fails
rustfmt; update the assertion comparing result.entries[0].0 to
Bytes::copy_from_slice(key) to a multi-line assert_eq! (put actual and expected
on separate lines, each on its own line, with proper trailing comma and optional
message) so rustfmt can format it; locate the assert using symbols assert_eq!,
result.entries, Bytes::copy_from_slice and key and rewrite that assertion into
the multi-line style.
🪄 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: f7ef5a3a-63fb-4198-b319-434f7542bf98
⛔ Files ignored due to path filters (1)
d-engine-proto/src/generated/d_engine.client.rsis excluded by!**/generated/**
📒 Files selected for processing (6)
d-engine-client/src/grpc_client.rsd-engine-core/src/raft_role/leader_state_test/client_read_test.rsd-engine-proto/proto/client/client_api.protod-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine.rsd-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine_test.rsexamples/service-discovery-standalone/watcher.rs
🚧 Files skipped from review as they are similar to previous changes (2)
- d-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine.rs
- examples/service-discovery-standalone/watcher.rs
There was a problem hiding this comment.
🧹 Nitpick comments (3)
d-engine-core/src/raft_role/leader_state_test/client_read_test.rs (1)
2207-2213: ⚡ Quick winDifferentiate
EmptyvsClosedin the no-response assertion.
resp_rx.try_recv().is_err()also passes when the sender is dropped/closed. For this regression, assert specifically that the channel is empty.Proposed tightening
+ use crate::maybe_clone_oneshot::TryRecvError; + assert!( - resp_rx.try_recv().is_err(), + matches!(resp_rx.try_recv(), Err(TryRecvError::Empty)), "BUG `#381`: LinearizableRead was served immediately in Phase 3 without any \ quorum confirmation. A partitioned leader can return stale data. \ The read must be deferred until handle_append_result confirms quorum \ (same as the LeaseRead expired path in process_lease_read)." );🤖 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/raft_role/leader_state_test/client_read_test.rs` around lines 2207 - 2213, The current assertion uses resp_rx.try_recv().is_err(), which conflates the channel being empty with it being closed; instead assert specifically that try_recv returned an Err of the "Empty" variant. Replace the generic is_err check with a check that matches Err(TryRecvError::Empty) (e.g., using matches!(resp_rx.try_recv(), Err(TryRecvError::Empty))) and add the appropriate import for TryRecvError if needed so the test ensures the receiver is empty rather than closed when verifying no response.d-engine-client/src/grpc_client_test.rs (1)
1706-1708: ⚡ Quick winAssert returned scan payload contents, not just count/revision.
test_scan_prefix_successshould also verify the returned keys/values to catch response-mapping regressions that still preservelen()andrevision.Proposed assertion expansion
assert_eq!(result.revision, 7); assert_eq!(result.entries.len(), 2); + assert_eq!(result.entries[0].0, Bytes::from("/services/node1")); + assert_eq!(result.entries[0].1, Bytes::from("10.0.0.1")); + assert_eq!(result.entries[1].0, Bytes::from("/services/node2")); + assert_eq!(result.entries[1].1, Bytes::from("10.0.0.2"));🤖 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 1706 - 1708, Update the test_scan_prefix_success assertions to validate the actual returned scan payload contents, not just count and revision: after asserting result.revision and result.entries.len(), assert the expected keys and values in result.entries (e.g., compare entries[i].key and entries[i].value or construct an expected map and assert equality) to catch response-mapping regressions; use the test name test_scan_prefix_success and the result variable (and its entries field) to locate where to add these assertions.d-engine-server/src/api/embedded_client_test.rs (1)
347-350: ⚡ Quick winAssert key→value pairs, not just key presence.
This test can pass even if values are mismatched. Please also validate returned values per key.
Proposed test hardening
let keys: Vec<_> = result.entries.iter().map(|(k, _)| k.clone()).collect(); assert!(keys.contains(&bytes::Bytes::from("/services/node1"))); assert!(keys.contains(&bytes::Bytes::from("/services/node2"))); assert!(keys.contains(&bytes::Bytes::from("/services/node3"))); + let by_key: std::collections::HashMap<_, _> = + result.entries.iter().cloned().collect(); + assert_eq!( + by_key.get(&bytes::Bytes::from("/services/node1")), + Some(&bytes::Bytes::from("10.0.0.1")) + ); + assert_eq!( + by_key.get(&bytes::Bytes::from("/services/node2")), + Some(&bytes::Bytes::from("10.0.0.2")) + ); + assert_eq!( + by_key.get(&bytes::Bytes::from("/services/node3")), + Some(&bytes::Bytes::from("10.0.0.3")) + );🤖 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-server/src/api/embedded_client_test.rs` around lines 347 - 350, The test currently only checks for presence of keys by collecting keys from result.entries into the keys variable; instead build a mapping from key→value (e.g., iterate result.entries into a HashMap or look up values per key) and assert each key maps to the expected value (use the same bytes::Bytes::from("/services/nodeX") keys and compare the corresponding value bytes to the expected value for node1/node2/node3). Update the assertions to check value equality (not just key containment) using result.entries (or the constructed map) and the expected bytes for each service.
🤖 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.
Nitpick comments:
In `@d-engine-client/src/grpc_client_test.rs`:
- Around line 1706-1708: Update the test_scan_prefix_success assertions to
validate the actual returned scan payload contents, not just count and revision:
after asserting result.revision and result.entries.len(), assert the expected
keys and values in result.entries (e.g., compare entries[i].key and
entries[i].value or construct an expected map and assert equality) to catch
response-mapping regressions; use the test name test_scan_prefix_success and the
result variable (and its entries field) to locate where to add these assertions.
In `@d-engine-core/src/raft_role/leader_state_test/client_read_test.rs`:
- Around line 2207-2213: The current assertion uses resp_rx.try_recv().is_err(),
which conflates the channel being empty with it being closed; instead assert
specifically that try_recv returned an Err of the "Empty" variant. Replace the
generic is_err check with a check that matches Err(TryRecvError::Empty) (e.g.,
using matches!(resp_rx.try_recv(), Err(TryRecvError::Empty))) and add the
appropriate import for TryRecvError if needed so the test ensures the receiver
is empty rather than closed when verifying no response.
In `@d-engine-server/src/api/embedded_client_test.rs`:
- Around line 347-350: The test currently only checks for presence of keys by
collecting keys from result.entries into the keys variable; instead build a
mapping from key→value (e.g., iterate result.entries into a HashMap or look up
values per key) and assert each key maps to the expected value (use the same
bytes::Bytes::from("/services/nodeX") keys and compare the corresponding value
bytes to the expected value for node1/node2/node3). Update the assertions to
check value equality (not just key containment) using result.entries (or the
constructed map) and the expected bytes for each service.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: ead48c0b-4dc1-4a69-ab06-c140101f73b2
📒 Files selected for processing (11)
d-engine-client/src/grpc_client_test.rsd-engine-client/src/mock_rpc.rsd-engine-client/src/mock_rpc_service.rsd-engine-core/src/raft_role/follower_state_test.rsd-engine-core/src/raft_role/leader_state_test/client_read_test.rsd-engine-core/src/storage/state_machine_test.rsd-engine-server/src/api/embedded_client_test.rsd-engine-server/src/network/grpc/grpc_raft_service_test.rsd-engine-server/src/storage/adaptors/file/file_state_machine_test.rsd-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine.rsd-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine_test.rs
🚧 Files skipped from review as they are similar to previous changes (1)
- d-engine-server/src/storage/adaptors/rocksdb/rocksdb_state_machine.rs
What Does This PR Do?
Adds a
scan_prefix(prefix)RPC that returns all key-value pairs under a namespace prefix as a linearizable snapshot with a Raft revision. Enables the zero-race-window Watch reconnection pattern: Watch first → scan snapshot → drain events withrevision > scan_revision.Type:
scan_prefix) #378 approved)Why Is This Needed?
For features: Link to approved issue #378
Without
scan_prefix, clients cannot recover missed state after a Watch stream drops. Re-registering a Watch from "now" leaves a gap for every write that occurred during the outage — silently stale routing tables, missed config updates, phantom services.scan_prefix+ the zero-race-window pattern eliminates this gap at the client level with no server-side changes to the Watch protocol.Checklist
Required:
make testpassesIf changing APIs:
watch-feature.md, example READMEs,watcher.rs)scan_prefix) #378)Testing
How tested:
rocksdb_state_machine_test.rs— 128 lines covering prefix scan correctness, empty prefix, key boundary (prefix successor), and revision consistencyscan-watchworkload — verifies no-gap / no-phantom / revision-monotonicity underFAULTS=partition,FAULTS=kill,partition,FAULTS=none, and buffer-overflow (RATE=200) scenarios; all four passmake testfull 9-workload Jepsen suite passesFor bug fixes: N/A
Does This Follow d-engine's Principles?
Reviewer Notes
Two commits in this PR:
feat(scan) #378— the API itself: proto,StateMachinetrait, RocksDB + file + sled-cluster impls,EmbeddedClient,GrpcClient, gRPC server handler,watch-feature.md+ example updatesfix(scan) #378— bug found during Jepsen testing:impl StateMachine for RocksDBStateMachinewas missing thescan_prefixdelegation; the trait default always returns error, so all production call paths throughSM: StateMachine(i.e.SMOF<T>) silently failed while direct struct dispatch in unit tests passedRocksDB implementation note: Uses
set_iterate_upper_bound(prefix_successor(prefix))+ forward iterator — noprefix_extractorneeded, block-level truncation gives O(results) not O(total keys).Estimated review complexity:
Summary by CodeRabbit
New Features
Documentation
Examples
Tests