Repository navigation
refactor #339: event-driven zombie detection and O(1) stale learner deadline - #364
Conversation
…eadline Zombie detection: replace get_zombie_candidates() polling with Sender<u32> signal. RaftMembership::new() returns (membership, Receiver<u32>); builder.rs spawns minimal bridge task → RoleEvent::ZombieDetected → handle_zombie_detected. HealthMonitor has zero dependency on core types. Stale learner: replace membership_maintenance_interval polling with stale_check_deadline: Option<Instant>. tick() does O(1) VecDeque-front comparison; refresh_stale_deadline/check_and_purge_stale_front replace run_periodic_maintenance. 7 TDD tests added. Remove membership_maintenance_interval and stale_check_interval config fields. Remove test_mem_first_loses_unflushed_data_on_crash: assumed Level 1 semantics, invalid for Level 2 (WAL replay guarantees process-crash recovery).
📝 WalkthroughWalkthroughRemoved periodic membership polling and config knobs; replaced zombie-candidate polling with channel-based zombie signals and RoleEvent::ZombieDetected; leader stale-learner handling changed from periodic scans to a front-deadline mechanism; updated traits, tests, builders, and health-monitor/membership constructors accordingly. Changes
Sequence Diagram(s)sequenceDiagram
autonumber
actor Net as Network
participant HM as HealthMonitor
participant RB as RaftMembership (owner)
participant NB as NodeBuilder Task
participant RL as Raft Role Loop
participant RS as Raft (handle_role_event)
participant LD as LeaderState
participant CN as Consensus
Net->>HM: record_failure(node_id)
HM->>NB: zombie_tx.send(node_id)
NB->>NB: if membership.is_zombie_valid(node_id)
NB->>RL: send RoleEvent::ZombieDetected(node_id)
RL->>RS: forward RoleEvent::ZombieDetected(node_id)
RS->>LD: handle_zombie_detected(node_id)
alt node promotable/non-readonly
LD->>CN: propose BatchRemove(node_id)
CN-->>LD: result
else readonly
LD-->>RS: skip
end
sequenceDiagram
autonumber
participant APP as LeaderState
participant Q as PendingPromotions(FIFO)
participant TM as Time
participant MB as Membership
participant CN as Consensus
Note over APP,Q: stale_check_deadline = None
APP->>Q: enqueue(learner ready_since)
APP->>APP: refresh_stale_deadline(front.ready_since + threshold)
TM-->>APP: tick(now < deadline)
APP-->>MB: (no calls)
TM-->>APP: tick(now >= deadline)
APP->>APP: check_and_purge_stale_front()
loop while front.expired
APP->>Q: pop_front()
APP->>MB: get_node_status(node)
alt promotable/non-readonly
APP->>CN: propose BatchRemove(node)
else readonly
APP-->>APP: skip
end
end
APP->>APP: refresh_stale_deadline()
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 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: 3
🧹 Nitpick comments (2)
d-engine-server/src/node/builder.rs (1)
404-412: Consider adding tracing for zombie detection events.The bridging task silently forwards zombie signals. Adding a trace log would improve observability for debugging cluster membership issues in production.
📡 Suggested tracing instrumentation
let role_tx_for_zombie = role_tx.clone(); tokio::spawn(async move { let mut zombie_rx = zombie_rx; while let Some(node_id) = zombie_rx.recv().await { + tracing::debug!(node_id, "Forwarding ZombieDetected signal to role event loop"); let _ = role_tx_for_zombie.send(RoleEvent::ZombieDetected(node_id)); } + tracing::debug!("Zombie signal bridge task exiting (membership dropped)"); });🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@d-engine-server/src/node/builder.rs` around lines 404 - 412, Add tracing inside the spawned task that bridges zombie signals so each forwarded event is logged for observability: inside the async move closure that uses role_tx_for_zombie and zombie_rx, emit a trace-level log (e.g. with tracing::trace!) when a node_id is received before sending RoleEvent::ZombieDetected(node_id), including the node_id and any contextual info; ensure the tracing crate is available in the module and keep the send behavior unchanged.d-engine-core/src/raft_role/leader_state_test/stale_learner_deadline_test.rs (1)
94-115: Test 3 uses real time sleep withoutstart_paused = true.While the test logic is sound (verifying FIFO ordering where front stays oldest), using
tokio::time::sleepwithout paused time could introduce minor flakiness under extreme system load.Consider using paused time for determinism:
⏱️ Optional: Use paused time for determinism
-#[tokio::test] +#[tokio::test(start_paused = true)] async fn test_stale_check_deadline_unchanged_on_second_enqueue() { let threshold = Duration::from_secs(30); let (mut leader, _ctx, _role_tx, _raft_tx) = setup_with_threshold(threshold).await; // Enqueue first learner let t0 = Instant::now(); leader.pending_promotions.push_back(PendingPromotion::new(10, t0)); leader.refresh_stale_deadline(threshold); let first_deadline = leader.stale_check_deadline.unwrap(); // Enqueue second learner slightly later - tokio::time::sleep(Duration::from_millis(10)).await; + tokio::time::advance(Duration::from_millis(10)).await; leader.pending_promotions.push_back(PendingPromotion::new(11, Instant::now()));🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@d-engine-core/src/raft_role/leader_state_test/stale_learner_deadline_test.rs` around lines 94 - 115, The test test_stale_check_deadline_unchanged_on_second_enqueue currently uses real time sleep which can be flaky; change the test attribute to #[tokio::test(start_paused = true)] and replace tokio::time::sleep(Duration::from_millis(10)).await with tokio::time::advance(Duration::from_millis(10)).await so the simulated time moves deterministically; keep the rest of the logic (pushing PendingPromotion, calling leader.refresh_stale_deadline, and asserting leader.stale_check_deadline) unchanged so the front-of-queue oldest behavior is still validated.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@d-engine-core/src/raft_role/mod.rs`:
- Around line 370-384: handle_zombie_detected currently drops ZombieDetected
events unless self is RaftRole::Leader, which loses one-shot detections; instead
persist or replay them across role transitions: when handle_zombie_detected
receives a zombie id and self is not RaftRole::Leader, append the id into a
persistent/struct field (e.g., pending_zombie_ids on the role state or parent)
rather than no-op, and when transitioning to leader (BecomeLeader) drain that
buffer and call LeaderState::handle_zombie_node (or revalidate with
RaftHealthMonitor::record_failure) for each id; also ensure duplicates are
deduplicated or revalidated to avoid repeated processing and clear the buffer
after replay.
In `@d-engine-server/src/network/health_monitor.rs`:
- Around line 55-65: The queued zombie notification can be stale; change the
signaling to include a generation/token so removal is revalidated: when you call
zombie_tx.try_send(node_id) from the failure-detection code, send a small struct
(node_id, gen) instead of just node_id, maintain a per-node generation counter
in the same map used by failure_counts (or an accompanying map) and increment or
bump that generation in record_success (and when clearing failure_counts) so any
outstanding ZombieDetected can be ignored unless the received gen matches the
current generation for that node; update the consumer that handles
ZombieDetected/BatchRemove to check the gen against the current map before
proposing removal.
- Around line 55-58: The current use of zombie_tx.try_send(node_id) can drop the
message on Full or Closed and thus lose the one-time zombie signal when
new_count == self.zombie_threshold; replace this lossy send with a reliable send
that retries or awaits until it succeeds and handle the Closed case explicitly:
on Full, await/send again (e.g., use async Sender::send(node_id).await or loop
retry on try_send until Ok) and on Closed log/report the error so the zombie
notification for the node_id from the new_count == self.zombie_threshold branch
is not silently dropped.
---
Nitpick comments:
In
`@d-engine-core/src/raft_role/leader_state_test/stale_learner_deadline_test.rs`:
- Around line 94-115: The test
test_stale_check_deadline_unchanged_on_second_enqueue currently uses real time
sleep which can be flaky; change the test attribute to
#[tokio::test(start_paused = true)] and replace
tokio::time::sleep(Duration::from_millis(10)).await with
tokio::time::advance(Duration::from_millis(10)).await so the simulated time
moves deterministically; keep the rest of the logic (pushing PendingPromotion,
calling leader.refresh_stale_deadline, and asserting
leader.stale_check_deadline) unchanged so the front-of-queue oldest behavior is
still validated.
In `@d-engine-server/src/node/builder.rs`:
- Around line 404-412: Add tracing inside the spawned task that bridges zombie
signals so each forwarded event is logged for observability: inside the async
move closure that uses role_tx_for_zombie and zombie_rx, emit a trace-level log
(e.g. with tracing::trace!) when a node_id is received before sending
RoleEvent::ZombieDetected(node_id), including the node_id and any contextual
info; ensure the tracing crate is available in the module and keep the send
behavior unchanged.
🪄 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: 39256cc7-2925-4a9a-9b43-20e7aa71583c
📒 Files selected for processing (21)
config/base/raft.tomld-engine-core/benches/leader_state_bench.rsd-engine-core/src/config/raft.rsd-engine-core/src/event.rsd-engine-core/src/membership.rsd-engine-core/src/raft.rsd-engine-core/src/raft_role/leader_state.rsd-engine-core/src/raft_role/leader_state_test/membership_change_test.rsd-engine-core/src/raft_role/leader_state_test/mod.rsd-engine-core/src/raft_role/leader_state_test/replication_test.rsd-engine-core/src/raft_role/leader_state_test/stale_learner_deadline_test.rsd-engine-core/src/raft_role/mod.rsd-engine-core/src/storage/buffered_raft_log_test/flush_strategy_test.rsd-engine-core/src/test_utils/mock/mock_raft_builder.rsd-engine-server/src/membership/raft_membership.rsd-engine-server/src/membership/raft_membership_test.rsd-engine-server/src/network/health_monitor.rsd-engine-server/src/network/health_monitor_test.rsd-engine-server/src/node/builder.rsd-engine-server/src/test_utils/integration/mod.rsd-engine-server/src/test_utils/mock/mock_node_builder.rs
💤 Files with no reviewable changes (8)
- d-engine-core/src/test_utils/mock/mock_raft_builder.rs
- d-engine-core/benches/leader_state_bench.rs
- d-engine-server/src/test_utils/mock/mock_node_builder.rs
- d-engine-core/src/storage/buffered_raft_log_test/flush_strategy_test.rs
- config/base/raft.toml
- d-engine-core/src/raft_role/leader_state_test/replication_test.rs
- d-engine-core/src/config/raft.rs
- d-engine-core/src/membership.rs
Codecov Report❌ Patch coverage is 📢 Thoughts on this report? Let us know! |
9e519a9 to
beb931a
Compare
…isory, document unavoidable duplicates
There was a problem hiding this comment.
🧹 Nitpick comments (1)
d-engine-server/src/node/builder.rs (1)
380-393: Consider documenting that pre-built memberships disable zombie detection.When a membership is injected via
self.membership, the dummy channel(_tx, rx)is created and_txis immediately dropped. This causes the bridge task to exit immediately (recv returnsNone), silently disabling zombie detection.If this is intentional (e.g., for testing or custom health monitoring), consider adding a doc comment or log message to make this behavior explicit. Otherwise, callers who inject memberships expecting full functionality may be surprised.
📝 Suggested documentation improvement
- // Pre-built membership has no zombie_rx; create a dummy closed channel. + // Pre-built membership has no zombie_rx; create a dummy closed channel. + // NOTE: This disables zombie detection for externally-provided memberships. + // The bridge task will exit immediately when _tx drops. |m| { let (_tx, rx) = tokio::sync::mpsc::channel(1); (m, rx) },🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@d-engine-server/src/node/builder.rs` around lines 380 - 393, When an injected membership is used (the self.membership.take() branch producing (membership_inner, zombie_rx)), we currently create a one-shot dummy channel (_tx, rx) and drop _tx which causes zombie detection to be disabled; make this explicit by either adding a doc comment to the constructor/method that accepts an injected membership (documenting that pre-built RaftMembership disables zombie detection) and/or emitting a debug/info log when the dummy channel is created so callers know zombie detection is intentionally disabled; reference the injection site (self.membership.take()), the resulting variables membership_inner and zombie_rx, and the RaftMembership::new() path when adding the docs/logs.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Nitpick comments:
In `@d-engine-server/src/node/builder.rs`:
- Around line 380-393: When an injected membership is used (the
self.membership.take() branch producing (membership_inner, zombie_rx)), we
currently create a one-shot dummy channel (_tx, rx) and drop _tx which causes
zombie detection to be disabled; make this explicit by either adding a doc
comment to the constructor/method that accepts an injected membership
(documenting that pre-built RaftMembership disables zombie detection) and/or
emitting a debug/info log when the dummy channel is created so callers know
zombie detection is intentionally disabled; reference the injection site
(self.membership.take()), the resulting variables membership_inner and
zombie_rx, and the RaftMembership::new() path when adding the docs/logs.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: 68bb6896-ae04-439d-b657-eeaa9399b0b7
⛔ Files ignored due to path filters (2)
Cargo.lockis excluded by!**/*.lockexamples/client-usage-standalone/Cargo.lockis excluded by!**/*.lock
📒 Files selected for processing (10)
d-engine-core/Cargo.tomld-engine-core/src/raft_role/leader_state.rsd-engine-core/src/raft_role/mod.rsd-engine-core/src/raft_role/raft_role_test.rsd-engine-core/src/raft_role/role_state.rsd-engine-server/src/membership/raft_membership.rsd-engine-server/src/network/health_monitor.rsd-engine-server/src/network/health_monitor_test.rsd-engine-server/src/node/builder.rsdeny.toml
✅ Files skipped from review due to trivial changes (1)
- d-engine-core/Cargo.toml
🚧 Files skipped from review as they are similar to previous changes (2)
- d-engine-core/src/raft_role/mod.rs
- d-engine-core/src/raft_role/leader_state.rs
What Does This PR Do?
Replaces two polling loops in membership maintenance with event-driven / deadline-based mechanisms: zombie node detection now fires on failure threshold via channel signal; stale learner expiry uses a single
Option<Instant>deadline against the VecDeque front instead of scanning on a fixed interval.Type:
Why Is This Needed?
Issue #339. The previous implementation polled
get_zombie_candidates()and scannedpending_promotionson everymembership_maintenance_intervaltick regardless of activity. This PR eliminates both hot-path polls:RaftHealthMonitoremitsSender<u32>when failure count hits threshold; a bridge task forwards it asRoleEvent::ZombieDetected.HealthMonitorgains zero dependency on core types.stale_check_deadline: Option<Instant>is set tofront().ready_since + thresholdon enqueue.tick()does oneInstantcomparison; no queue scan, no periodic wakeup when idle.Checklist
make testpasses (1492 tests)stale_learner_deadline_test.rs)If changing APIs:
RaftMembership::new()returns(membership, Receiver<u32>)— all call sites updatedmembership_maintenance_intervalandstale_check_intervalconfig fields —raft.tomlupdatedTesting
make test, 1492/1492)Does This Follow d-engine's Principles?
Reviewer Notes
RaftMembership::new()return type change is the widest blast radius — 25+ call sites updated mechanically. Bridge task inbuilder.rsis intentionally minimal (no shutdown signal; lifecycle follows sender drop naturally).test_mem_first_loses_unflushed_data_on_crashremoved: was testing Level 1 semantics (unflushed = lost on crash); d-engine uses Level 2 (WAL replay guarantees process-crash recovery), making the test both semantically wrong and racy.Estimated review complexity:
Summary by CodeRabbit
New Features
Refactor
Tests
Chores