Repository navigation
perf #407: coalesce AppendEntries and decouple Raft event processing layers - #411
Conversation
|
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 (2)
✅ Files skipped from review due to trivial changes (1)
🚧 Files skipped from review as they are similar to previous changes (1)
📝 WalkthroughWalkthroughThe PR introduces a drain-based buffering architecture in the Raft main loop, adding ChangesAppendEntries Batching and Merge Pipeline
Sequence Diagram(s)sequenceDiagram
participant gRPC as gRPC AppendEntries Handler
participant Channel as raft_event_tx
participant Raft as Raft Main Loop
participant Drain as Drain Helpers
participant Merge as merge_append_entries
participant Role as RaftRoleState
participant gRPCResp as Response Handlers
gRPC->>Channel: RaftEvent::AppendEntries(req, vec![resp_tx])
Channel->>Raft: tokio::select! fires raft arm
Raft->>Raft: push_back to buffered_raft_event
Raft->>Drain: drain_raft_events (up to max_batch_size)
Drain->>Raft: try_recv additional events into buffer
Raft->>Merge: process_raft_events: merge_append_entries(...)
Merge->>Merge: coalesce contiguous same-term AppendEntries
Merge->>Merge: accumulate entries and senders Vec
Merge->>Raft: rewrite buffered_raft_event with merged events
Raft->>Role: handle_raft_event(merged AppendEntries, Vec<senders>)
Role->>Role: handle_append_entries_request_workflow(..., senders)
Role->>gRPCResp: iterate senders: send response
gRPCResp->>gRPC: oneshot::Sender::send result
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related issues
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: 4
🤖 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_test/merge_append_entries_tests.rs`:
- Around line 855-856: The test at line 855-856 with the loop containing the
tuple [(9u64, 5, 6), (14, 5, 5)] is validating term-boundary behavior rather
than max_merge_entries cap enforcement because the two test cases use different
terms (6 and 5), causing merge_append_entries() to stop at the term mismatch
before evaluating the max_merge_entries guard. Fix this by modifying both test
cases to use the same term value so that the max_merge_entries cap enforcement
is actually tested rather than the term boundary logic. Ensure the prev and
count values remain as parameters to properly exercise the cap behavior while
keeping the term consistent across both iterations.
In `@d-engine-core/src/raft.rs`:
- Around line 388-393: The error handling in the `process_role_events` method
currently swallows all errors returned by `handle_role_event`, including fatal
errors like `RoleEvent::FatalError`, which prevents the node from properly
exiting when a fatal condition occurs. Refactor the error handling in
`process_role_events` to distinguish between fatal and non-fatal errors,
matching the fatal/non-fatal error handling pattern already implemented in the
`process_raft_events` method. Fatal errors should be propagated by returning
them from `process_role_events`, while non-fatal errors should continue to be
logged and ignored.
- Around line 397-416: The merge_append_entries() function is only called once
before the while loop begins in the process_raft_events method, which means if
the buffered queue contains non-AppendEntries events (like VoteRequest)
interspersed with AppendEntries events, the later consecutive AppendEntries runs
will not be merged. Move the merge_append_entries() call inside the while loop
so it executes before each event is processed, ensuring that all consecutive
runs of AppendEntries entries are merged regardless of their position in the
queue, not just the prefix entries.
- Around line 446-448: The merge guard condition for RaftEvent::AppendEntries is
incomplete and does not validate the log-matching property of Raft. The current
condition checks next_prev == req.prev_log_index and term == req.term, but it is
missing a check for prev_log_term consistency. Add an additional check to the
merge guard condition to verify that the log term at prev_log_index matches
req.prev_log_term, ensuring that the log-matching checks are preserved across
merged boundaries and preventing inconsistent append entries requests from being
accepted as part of a merged operation.
🪄 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: 32f85a0b-75ec-4ec0-8c65-59bf2919f168
📒 Files selected for processing (22)
benches/standalone-bench/src/main.rsd-engine-core/src/config/raft.rsd-engine-core/src/event.rsd-engine-core/src/raft.rsd-engine-core/src/raft_role/candidate_state.rsd-engine-core/src/raft_role/candidate_state_test.rsd-engine-core/src/raft_role/follower_state_test.rsd-engine-core/src/raft_role/leader_state.rsd-engine-core/src/raft_role/leader_state_test/become_follower_test.rsd-engine-core/src/raft_role/leader_state_test/client_write_test.rsd-engine-core/src/raft_role/leader_state_test/commit_index_test.rsd-engine-core/src/raft_role/leader_state_test/event_handling_test.rsd-engine-core/src/raft_role/leader_state_test/replication_test.rsd-engine-core/src/raft_role/learner_state_test.rsd-engine-core/src/raft_role/role_state.rsd-engine-core/src/raft_test/drain_based_batch_architecture_tests.rsd-engine-core/src/raft_test/merge_append_entries_tests.rsd-engine-core/src/raft_test/mod.rsd-engine-core/src/storage/buffered_raft_log.rsd-engine-server/src/network/grpc/grpc_raft_service.rsd-engine-server/src/storage/adaptors/rocksdb/rocksdb_storage_engine.rsd-engine/src/docs/performance/throughput-optimization-guide.md
💤 Files with no reviewable changes (1)
- d-engine-core/src/raft_test/mod.rs
Codecov Report❌ Patch coverage is 📢 Thoughts on this report? Let us know! |
Leader now advances next_index immediately when firing AppendEntries (prev_log_index + entries.len() + 1) instead of waiting for ACK, eliminating the stop-and-wait LEGACY resend on every round. update_peer_index replaces update_peer_indexes with two distinct paths: - Success: max(current, peer_match+1) — never regress speculative value - Conflict: retreat to follower hint, floor-guarded by match_index+1 Add 7 unit tests covering speculative advance, success preservation, conflict rollback, and floor guard correctness.
…layers - Buffer raft/role events in VecDeque; merge contiguous same-term AEs, absorb heartbeats via max(leader_commit_index), cap via max_merge_entries - AppendEntries sender → Vec to broadcast across merged events - Add merge_append_entries_tests.rs - Fix is_write_durable() comments: flush_wal(true) = batched Level 3, not unrecoverable - Update throughput-optimization-guide.md: correct MemFirst durability semantics
…n loop merge_append_entries was called once before the loop; AE runs following a non-AE event were never coalesced. matches!() guard fires only when front is AppendEntries, skipping push/pop cost for non-AE events. Adds process_raft_events_tests.rs: drain-loop tests asserting via mock times() — times(1) vs times(2) on handle_append_entries is the invariant.
…rge_entries test
e5ec77a to
2742334
Compare
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_test/process_raft_events_tests.rs`:
- Around line 237-239: The comment on the expect_handle_append_entries() call
incorrectly references "merged dispatch" which doesn't occur in this test—this
appears to be copy-pasted from another test. Replace the misleading comment with
an accurate description explaining that the 2 calls correspond to 2 append
entries that are each dispatched independently without any merging, reflecting
the actual behavior of VR isolation in this specific test.
- Around line 247-252: The inline comments describing entry ranges in the test
do not match the actual AppendEntries requests being created. Update the comment
above the first queue.push_back call with make_ae_event to correctly reflect
that it creates entries 10-14 (not 10-19) based on the make_request call with
prev=9 and 5 entries. Similarly, update the comment above the second
queue.push_back call with make_ae_event to correctly reflect that it creates
entries 15-19 (not 20-29) based on the make_request call with prev=14 and 5
entries. These comments were copy-pasted from an earlier test and need
correction to accurately describe the actual entry ranges in this 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: 1223c5ac-c04c-47e0-a121-121e0b9ab14e
📒 Files selected for processing (23)
benches/standalone-bench/src/main.rsd-engine-core/src/config/raft.rsd-engine-core/src/event.rsd-engine-core/src/raft.rsd-engine-core/src/raft_role/candidate_state.rsd-engine-core/src/raft_role/candidate_state_test.rsd-engine-core/src/raft_role/follower_state_test.rsd-engine-core/src/raft_role/leader_state.rsd-engine-core/src/raft_role/leader_state_test/become_follower_test.rsd-engine-core/src/raft_role/leader_state_test/client_write_test.rsd-engine-core/src/raft_role/leader_state_test/commit_index_test.rsd-engine-core/src/raft_role/leader_state_test/event_handling_test.rsd-engine-core/src/raft_role/leader_state_test/replication_test.rsd-engine-core/src/raft_role/learner_state_test.rsd-engine-core/src/raft_role/role_state.rsd-engine-core/src/raft_test/drain_based_batch_architecture_tests.rsd-engine-core/src/raft_test/merge_append_entries_tests.rsd-engine-core/src/raft_test/mod.rsd-engine-core/src/raft_test/process_raft_events_tests.rsd-engine-core/src/storage/buffered_raft_log.rsd-engine-server/src/network/grpc/grpc_raft_service.rsd-engine-server/src/storage/adaptors/rocksdb/rocksdb_storage_engine.rsd-engine/src/docs/performance/throughput-optimization-guide.md
💤 Files with no reviewable changes (1)
- d-engine-core/src/raft_test/mod.rs
✅ Files skipped from review due to trivial changes (4)
- benches/standalone-bench/src/main.rs
- d-engine-core/src/storage/buffered_raft_log.rs
- d-engine-core/src/raft_test/drain_based_batch_architecture_tests.rs
- d-engine-core/src/raft_role/leader_state_test/commit_index_test.rs
🚧 Files skipped from review as they are similar to previous changes (16)
- d-engine-core/src/raft_role/leader_state_test/event_handling_test.rs
- d-engine-server/src/network/grpc/grpc_raft_service.rs
- d-engine-core/src/config/raft.rs
- d-engine-server/src/storage/adaptors/rocksdb/rocksdb_storage_engine.rs
- d-engine-core/src/raft_role/follower_state_test.rs
- d-engine-core/src/raft_role/leader_state_test/client_write_test.rs
- d-engine-core/src/raft_role/role_state.rs
- d-engine-core/src/raft_role/candidate_state.rs
- d-engine-core/src/raft_role/candidate_state_test.rs
- d-engine-core/src/raft_role/learner_state_test.rs
- d-engine-core/src/raft_role/leader_state_test/replication_test.rs
- d-engine-core/src/raft_role/leader_state.rs
- d-engine-core/src/event.rs
- d-engine/src/docs/performance/throughput-optimization-guide.md
- d-engine-core/src/raft.rs
- d-engine-core/src/raft_test/merge_append_entries_tests.rs
What Does This PR Do?
Introduces a
VecDequebuffer between theselect!loop and Raft event processing, then addsmerge_append_entries()to coalesce consecutive same-termAppendEntriesevents — including heartbeat absorption viamax(leader_commit_index)— before dispatch. Also corrects durability comments and docs that misrepresentedMemFirstas "not power-loss safe."Type:
Why Is This Needed?
For features: Issue #407
Under high replication load, the
select!loop received oneAppendEntriesevent per iteration. Each event triggered a full follower append cycle (log write +durable_indexadvance + response), serialising what could be a single batched write. Role events (heartbeats, votes) arriving between AEs also prevented natural batching. Decouplingselect!from processing via buffered VecDeques allowsmerge_append_entries()to coalesce the leading run of contiguous AEs — including heartbeat absorption — into one dispatch, reducing fsync round-trips and follower response round-trips per unit of replicated data.Checklist
Required:
make testpassesIf changing APIs:
Testing
How tested:
merge_append_entries_tests.rs— 14 tests covering: empty buffer, single event, basic 3-way merge, metadata preservation, non-contiguous stop, non-AE FIFO stop, heartbeat absorption (mid/end/consecutive/all-HB), sender collection order, dropped receiver, term ordering (newer-first stop),max_merge_entriescapmake testsuiteFor performance improvements:
Does This Follow d-engine's Principles?
Reviewer Notes
AppendEntriessecond field changes fromMaybeCloneOneshotSender→Vec<MaybeCloneOneshotSender>— all call sites updated. Themerge_append_entries()method ispub(crate)only, no public API change.Durability comment fix:
is_write_durable() = false+flush()=flush_wal(true)/sync_all()→ committed data is disk-durable (batched Level 3).throughput-optimization-guide.mdand inline comments corrected accordingly.drain_based_batch_architecture_tests.rshas pre-existing compile errors (orphaned before this branch, now surfaced via#[path]) — tracked separately, not part of this PR.Estimated review complexity:
Summary by CodeRabbit
New Features
max_merge_entriesto cap consecutiveAppendEntriescoalescing.AppendEntries(including heartbeats).next_indexadvancement before append-result feedback.Bug Fixes
next_indexregressions on stale conflict acknowledgements with refined clamp/update logic.AppendEntriesresponses to multiple recipients by logging per-recipient send failures without aborting processing.Documentation