Skip to content

feat #327, #328: cluster membership streaming (EmbeddedEngine + GrpcClient watch_membership) - #366

Merged
JoshuaChi merged 6 commits into
mainfrom
feature/327-328-watch-membership
May 5, 2026
Merged

JoshuaChi merged 6 commits into
mainfrom
feature/327-328-watch-membership

Conversation

@JoshuaChi

@JoshuaChi JoshuaChi commented May 5, 2026 •

Copy link
Copy Markdown
Contributor

What Does This PR Do?

Adds real-time cluster membership change notifications for both embedded and standalone (gRPC)
modes. Callers subscribe once and receive a snapshot on every committed ConfChange — no polling.

Type:


Why Is This Needed?

For features: Issues #327 and #328

Operators and upper-layer tooling (service discovery, load balancers, health dashboards) need
to react to cluster topology changes (node join, promotion, removal) without polling
get_cluster_metadata. A push-based membership stream enables zero-latency reactions and
eliminates unnecessary RPC traffic.


Checklist

Required:

  • make test passes
  • Added tests for new code
  • Commits squashed to 1-2 logical units

If changing APIs:

  • Updated relevant docs (CHANGELOG.md)
  • Explained why complexity is justified

Testing

How tested:

  • Unit tests (server): test_watch_membership_returns_unavailable_when_node_not_ready,
    test_watch_membership_yields_current_snapshot_then_sentinel_on_sender_drop
    — verify the gRPC handler's not-ready guard and the mark_changed() + UNAVAILABLE sentinel behavior.

  • Unit tests (client): test_watch_membership_returns_err_when_server_rejects,
    test_watch_membership_receives_snapshots_in_order,
    test_watch_membership_empty_stream_closes_cleanly
    — verify GrpcClient::watch_membership() against a mock gRPC server.

  • Integration tests (embedded): 7 tests in watch_membership_embedded.rs covering
    initial snapshot delivery, node join/promotion, zombie no-op, multi-subscriber fanout,
    and committed_index monotonicity.

  • Integration test (standalone/gRPC): watch_membership_standalone.rs — boots a real
    2-node cluster, opens a gRPC stream, joins a 3rd learner node, asserts the stream yields
    a snapshot with learners=[3] and committed_index > 0.

All 13 watch_membership tests pass: cargo nextest run --all-features -E 'test(watch_membership)'


Does This Follow d-engine's Principles?

  • Solves a real problem for most users (not just my edge case)
  • Keeps implementation simple
  • Doesn't bloat the API surface

Reviewer Notes

4 commits, 3 concerns to focus on:

  1. mark_changed() before WatchStream::new() (grpc_raft_service.rs): This forces the
    first stream item to be the current snapshot, matching watch::Receiver::borrow() semantics.
    Requires tokio ≥ 1.37 (project uses 1.51.1). Alternative would be an extra initial send()
    which adds complexity; mark_changed() is idiomatic.

  2. UNAVAILABLE sentinel (grpc_raft_service.rs): The stream is chained with a
    futures::stream::once(Err(Status::unavailable(...))) so clients get a clean signal instead
    of a silent stream close. Clients should reconnect and resubscribe on receiving this error.

  3. GrpcClient::watch_membership() is NOT in ClientApi trait: Membership watch is
    standalone-only (the embedded equivalent uses EmbeddedEngine::watch_membership() which
    returns a watch::Receiver). Putting both under one trait would force type-system gymnastics;
    keeping them separate preserves simplicity.

Estimated review complexity:

  • Deep (> 300 lines)

Summary by CodeRabbit

  • New Features

    • Added cluster membership streaming: clients can now subscribe to real-time membership snapshots through both embedded and gRPC APIs, with immediate initial snapshot delivery and updates on each committed membership change.
  • Important Fixes

    • Zombie detection now emits warnings only; unreachable nodes are no longer automatically removed and require deliberate operator action for removal.

…t API

- New `MembershipSnapshot` struct (members, learners, committed_index)
- `EmbeddedEngine::watch_membership()` returns a `watch::Receiver<MembershipSnapshot>`
  that fires on every Raft-consensus-backed conf change (AddNode, Promote, BatchRemove)
- Notification reaches all nodes (leader, follower, learner) via the existing
  CommitHandler::apply_config_change path — no extra machinery
- 7 integration tests covering: initial snapshot, join, zombie warn-only,
  learner promotion, all-nodes notification, multiple subscribers, monotonic index
…or peer removed

The inner reconnect loop in process_batch called open_replication_stream in a tight
loop without checking whether the owning LeaderState had been dropped. When a peer
was removed or the leader stepped down, the worker kept retrying indefinitely.

Fix: check task_rx.is_closed() at the top of the reconnect loop. The task channel
sender lives inside LeaderState; once it drops, is_closed() returns true and the
worker returns immediately instead of spinning forever.

Also hide leader_id in ClusterConf responses until noop commits: before the noop
entry commits, the leader has not confirmed quorum readiness, so exposing its id
too early causes clients to route to it before it is ready to serve.
…via BatchRemove

Previously handle_zombie_node called execute_request_immediately(BatchRemove) when a
node exceeded the connection failure threshold. This was too aggressive: a node that
is temporarily restarting would be permanently ejected from the cluster.

Membership changes are high-risk Raft consensus operations. The framework should
detect and report, not decide. Now handle_zombie_node only emits warn!(Zombie detected)
with the node_id and status. Removal remains a deliberate operator or upper-layer
decision.

Peer failure/success telemetry channels (grpc_transport → health_monitor → zombie bridge)
are preserved so the warning still fires reliably when the threshold is crossed.
@coderabbitai

coderabbitai Bot commented May 5, 2026 •

Copy link
Copy Markdown

Warning

Rate limit exceeded

@JoshuaChi has exceeded the limit for the number of commits that can be reviewed per hour. Please wait 6 minutes and 35 seconds before requesting another review.

To keep reviews running without waiting, you can enable usage-based add-on for your organization. This allows additional reviews beyond the hourly cap. Account admins can enable it under billing.

⌛ How to resolve this issue?

After the wait time has elapsed, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

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 configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro

Run ID: 6047158e-874f-4b22-99ce-855567465294

📥 Commits

Reviewing files that changed from the base of the PR and between afd2dda and db18994.

⛔ Files ignored due to path filters (1)
  • Cargo.lock is excluded by !**/*.lock
📒 Files selected for processing (5)
  • Makefile
  • d-engine-core/src/raft_role/leader_state_test/worker_lifecycle_test.rs
  • d-engine-proto/proto/client/client_api.proto
  • d-engine-server/src/membership/raft_membership.rs
  • d-engine-server/tests/watch_and_subscriptions/watch_membership_embedded.rs
📝 Walkthrough

Walkthrough

This PR implements cluster membership streaming capabilities through a new WatchMembership RPC and embedded API, enabling clients to receive real-time snapshots of committed cluster membership. It also changes zombie node detection from automatic removal to warning-only behavior and refines leader visibility in cluster membership queries.

Changes

Cluster Membership Streaming and Zombie Detection Updates

Layer / File(s) Summary
Proto Definitions
d-engine-proto/proto/client/client_api.proto
Introduces WatchMembershipRequest, MembershipSnapshot messages, and WatchMembership streaming RPC to RaftClientService.
Membership Data Structure
d-engine-server/src/membership/membership_snapshot.rs, d-engine-server/src/membership/mod.rs
Adds MembershipSnapshot struct with members, learners, and committed_index fields for point-in-time cluster state snapshots.
Core Raft Logic Changes
d-engine-core/src/raft_role/leader_state.rs
(1) Derives current leader for cluster conf only when noop is committed via noop_log_id.and(...). (2) Exits replication worker reconnect loop when task channel closes. (3) Changes handle_zombie_node from proposing BatchRemove to emitting warning logs only.
Membership Snapshot Management
d-engine-server/src/membership/raft_membership.rs
Adds tokio::watch channel for membership snapshots; notify_config_applied() publishes snapshots on commit; new APIs subscribe_membership(), on_peer_stream_success(), on_peer_stream_failed() for subscription and health reporting.
Node and Embedded Engine Wiring
d-engine-server/src/node/mod.rs, d-engine-server/src/api/embedded.rs, d-engine-server/src/lib.rs
Node<T> gains membership_rx field and membership_change_notifier() method; EmbeddedEngine stores receiver and exposes watch_membership() API; MembershipSnapshot re-exported at crate root.
gRPC Transport and Service Integration
d-engine-server/src/network/grpc/grpc_transport.rs, d-engine-server/src/network/grpc/grpc_raft_service.rs
GrpcTransport adds peer_failure_tx/peer_success_tx channels to report stream health; grpc_raft_service implements watch_membership RPC that streams snapshots with immediate initial delivery and terminal UNAVAILABLE on shutdown.
Node Builder Bridges
d-engine-server/src/node/builder.rs
Wires peer health channels from transport through membership health monitor via spawned bridge tasks.
Client APIs
d-engine-client/src/grpc_client.rs
Adds GrpcClient::watch_membership() method that sends WatchMembershipRequest and returns tonic::Streaming<MembershipSnapshot>.
Mock RPC Infrastructure
d-engine-client/src/mock_rpc.rs, d-engine-client/src/mock_rpc_service.rs, d-engine-core/src/test_utils/mock/mock_rpc.rs, d-engine-server/src/test_utils/mock/mock_rpc.rs, d-engine-server/src/test_utils/mock/mock_node_builder.rs
Extends mock RPC services with WatchMembershipStream type and watch_membership() implementations; adds mock server helpers and updates node builders to wire membership receivers.
Core Logic Tests
d-engine-core/src/raft_role/follower_state_test.rs, d-engine-core/src/raft_role/leader_state_test/event_handling_test.rs, d-engine-core/src/raft_role/leader_state_test/membership_change_test.rs, d-engine-core/src/raft_role/leader_state_test/worker_lifecycle_test.rs
Tests validate noop-based leader visibility in cluster conf, zombie-detection warn-only behavior, and replication worker early exit on channel close.
Client/RPC Tests
d-engine-client/src/grpc_client_test.rs, d-engine-server/src/network/grpc/grpc_raft_service_test.rs
Tests verify GrpcClient::watch_membership() error handling, stream ordering, and server-side watch_membership RPC readiness checks and stream termination.
Integration Tests
d-engine-server/tests/watch_and_subscriptions/mod.rs, d-engine-server/tests/watch_and_subscriptions/watch_membership_embedded.rs, d-engine-server/tests/watch_and_subscriptions/watch_membership_standalone.rs
End-to-end tests covering embedded and standalone modes; verify initial snapshots, firing on membership changes, concurrent subscribers, monotonic committed_index, and zombie-detection warn behavior.
Test Helpers and Stability Improvements
d-engine-server/tests/common/mod.rs, d-engine-server/tests/cas_operations/leader_failover_cas_standalone.rs, d-engine-server/tests/cas_operations/snapshot_recovery_standalone.rs, d-engine-server/tests/cluster_lifecycle/scale_single_to_three_node_embedded.rs, d-engine-server/tests/failover_and_recovery/leader_failover_standalone.rs
Adds wait_for_stable_leader() and create_rejoin_node_config() helpers to improve test stability during leader transitions; updates existing tests to use new helpers instead of fixed sleeps.
Documentation
CHANGELOG.md
Documents "Cluster membership streaming" and "Zombie detection no longer auto-removes unreachable nodes" features in v0.2.4 release notes.

Sequence Diagram(s)

sequenceDiagram
    participant Client as gRPC/Embedded Client
    participant EmbeddedEngine as EmbeddedEngine /<br>GrpcService
    participant Node as Node<br>(Raft)
    participant RaftMembership as RaftMembership
    participant MembershipNotifier as Membership<br>Watch Channel

    Client->>EmbeddedEngine: watch_membership() /<br>GrpcClient::watch_membership()
    activate EmbeddedEngine
    EmbeddedEngine->>Node: membership_change_notifier()
    activate Node
    Node-->>EmbeddedEngine: watch::Receiver<MembershipSnapshot>
    deactivate Node
    EmbeddedEngine-->>Client: Receiver / Streaming<MembershipSnapshot>
    deactivate EmbeddedEngine

    activate Client
    Note over Client: Await first snapshot
    Client->>MembershipNotifier: borrow() / next()
    activate MembershipNotifier
    MembershipNotifier-->>Client: Initial MembershipSnapshot
    deactivate MembershipNotifier
    deactivate Client

    Node->>RaftMembership: ConfChange committed
    activate RaftMembership
    RaftMembership->>MembershipNotifier: send_replace(snapshot)
    activate MembershipNotifier
    MembershipNotifier->>Client: Emit snapshot update
    deactivate MembershipNotifier
    deactivate RaftMembership
    Client->>Client: changed() fires, read updated snapshot
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~45 minutes

Possibly related issues

Possibly related PRs

  • deventlab/d-engine#206: Extends v0.2.0 work on membership/watch APIs with MembershipSnapshot, WatchMembership RPCs, and client/server/watch wiring.
  • deventlab/d-engine#175: Related through proto API surface modifications (proto/client/client_api.proto) for membership watch support.
  • deventlab/d-engine#144: Related through gRPC transport layer modifications and replication/peer-stream machinery updates.

Poem

🐰 Snapshots flutter through the mesh,
Members dance, committed fresh,
Zombies warn instead of die,
Watch them stream and watch them fly!
— A Rabbit's Ode to Membership Streams 🐇✨

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly and specifically describes the main change: adding cluster membership streaming functionality to both EmbeddedEngine and GrpcClient, directly matching the file changes and PR objectives.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feature/327-328-watch-membership

Tip

💬 Introducing Slack Agent: The best way for teams to turn conversations into code.

Slack Agent is built on CodeRabbit's deep understanding of your code, so your team can collaborate across the entire SDLC without losing context.

  • Generate code and open pull requests
  • Plan features and break down work
  • Investigate incidents and troubleshoot customer tickets together
  • Automate recurring tasks and respond to alerts with triggers
  • Summarize progress and report instantly

Built for teams:

  • Shared memory across your entire org—no repeating context
  • Per-thread sandboxes to safely plan and execute work
  • Governance built-in—scoped access, auditability, and budget controls

One agent for your entire SDLC. Right inside Slack.

👉 Get started


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.

❤️ Share

Comment @coderabbitai help to get the list of available commands and usage tips.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 5

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
d-engine-server/tests/watch_and_subscriptions/mod.rs (1)

6-9: ⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Module doc comment is stale — update to list the two new submodules.

watch_membership_embedded and watch_membership_standalone are declared but omitted from the ## Test Coverage section.

📝 Proposed fix
 //! ## Test Coverage
 //!
 //! - `watch_events_embedded.rs` - In-process key change subscriptions (embedded mode)
 //! - `watch_events_grpc_standalone.rs` - gRPC streaming watch events (standalone mode)
 //! - `watch_performance_gate_embedded.rs` - Watch latency and throughput benchmarks (embedded mode)
+//! - `watch_membership_embedded.rs` - Committed membership change notifications (embedded mode)
+//! - `watch_membership_standalone.rs` - gRPC membership snapshot streaming (standalone mode)
🤖 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/tests/watch_and_subscriptions/mod.rs` around lines 6 - 9,
Update the module doc comment at the top of tests/watch_and_subscriptions/mod.rs
to reflect the current submodules: replace the stale entries (e.g.,
watch_events_embedded.rs, watch_events_grpc_standalone.rs,
watch_performance_gate_embedded.rs) with the actual submodules now declared
(include watch_membership_embedded and watch_membership_standalone) and add both
names to the "## Test Coverage" section so those two new submodules are listed;
ensure the submodule list and coverage section match the declared mod statements
(watch_membership_embedded, watch_membership_standalone).
🧹 Nitpick comments (7)
d-engine-server/tests/common/mod.rs (2)

595-614: ⚡ Quick win

wait_for_stable_leader can loop indefinitely — add an explicit deadline.

client.refresh(None).await.ok() silently discards all errors. If the cluster is permanently unreachable (or a test bug prevents recovery), refresh keeps failing silently, get() keeps returning ConnectionTimeout, and the function sleeps 100 ms then loops — forever. There is no escape hatch.

Contrast this with the Phase 5 loop in leader_failover_cas_standalone.rs which uses refresh(None).await? (propagates errors) and is therefore bounded by the 30 s cluster_ready_timeout. wait_for_stable_leader has no equivalent bound.

A CI job will eventually time it out, but without a clear test-failure message, making regressions hard to diagnose.

⏱️ Proposed fix: wrap in an explicit timeout
-pub async fn wait_for_stable_leader(client: &Client) -> Result<(), ClientApiError> {
-    loop {
-        client.refresh(None).await.ok();
-        match client.get(b"__stability_probe__").await {
-            Ok(_) => return Ok(()),
-            Err(ClientApiError::Business {
-                code: ErrorCode::StaleOperation,
-                ..
-            }) => continue,
-            Err(ClientApiError::Network {
-                code: ErrorCode::ConnectionTimeout,
-                ..
-            }) => {
-                tokio::time::sleep(Duration::from_millis(100)).await;
-                continue;
-            }
-            Err(e) => return Err(e),
-        }
-    }
-}
+pub async fn wait_for_stable_leader(client: &Client) -> Result<(), ClientApiError> {
+    const STABILITY_TIMEOUT: Duration = Duration::from_secs(30);
+    tokio::time::timeout(STABILITY_TIMEOUT, async {
+        loop {
+            client.refresh(None).await.ok();
+            match client.get(b"__stability_probe__").await {
+                Ok(_) => return Ok(()),
+                Err(ClientApiError::Business {
+                    code: ErrorCode::StaleOperation,
+                    ..
+                }) => continue,
+                Err(ClientApiError::Network {
+                    code: ErrorCode::ConnectionTimeout,
+                    ..
+                }) => {
+                    tokio::time::sleep(Duration::from_millis(100)).await;
+                    continue;
+                }
+                Err(e) => return Err(e),
+            }
+        }
+    })
+    .await
+    .map_err(|_elapsed| {
+        std::io::Error::new(
+            std::io::ErrorKind::TimedOut,
+            "wait_for_stable_leader: no stable leader within 30s",
+        )
+        .into()
+    })?
+}
🤖 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/tests/common/mod.rs` around lines 595 - 614, The
wait_for_stable_leader function currently swallows refresh errors via
client.refresh(None).await.ok() and can loop forever on repeated
ConnectionTimeouts; change it to enforce an explicit deadline (e.g., use
tokio::time::timeout or track Instant::now() + duration) and propagate refresh
errors instead of ignoring them (replace the .ok() with ?.await or check the
Result), returning a clear timeout/err result if the deadline is exceeded;
update references in wait_for_stable_leader and the loop that matches
client.get(...) so the function returns an Err with a descriptive error when the
overall timeout is reached.

600-611: ⚡ Quick win

Add handling for NotLeader and LeaderChanged errors to the retry loop.

The retry set currently covers only StaleOperation (4002) and ConnectionTimeout (1001), but linearizable reads on non-leader nodes return NotLeader (4001, Business layer) and leadership changes return LeaderChanged (1003, Network layer). Both errors will escape via Err(e) => return Err(e) and fail callers during failover settling, which contradicts the function's purpose.

Add match arms to retry on:

  • ClientApiError::Business { code: ErrorCode::NotLeader, .. }
  • ClientApiError::Network { code: ErrorCode::LeaderChanged, .. }

with similar backoff (e.g., 100ms sleep) as the ConnectionTimeout case.

🤖 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/tests/common/mod.rs` around lines 600 - 611, The retry loop
currently matches only StaleOperation and ConnectionTimeout; add two new match
arms to also handle ClientApiError::Business { code: ErrorCode::NotLeader, .. }
and ClientApiError::Network { code: ErrorCode::LeaderChanged, .. } so they
behave like the retry cases: await
tokio::time::sleep(Duration::from_millis(100)).await and then continue the loop
rather than returning Err(e). Update the match in the function containing the
shown match on Err(ClientApiError::...) to include these two arms so
linearizable reads and leader changes are retried with the same backoff.
d-engine-server/src/membership/membership_snapshot.rs (2)

3-7: 💤 Low value

Optional: rustdoc intra-doc link for EmbeddedEngine::watch_membership.

The bracketed reference [EmbeddedEngine::watch_membership] is unlikely to resolve from within the membership module since EmbeddedEngine lives in crate::api. With #![warn(missing_docs)] enabled at the crate level you may want to also enable broken_intra_doc_links (or use an explicit path) so this link doesn't silently render as plain text.

Suggested fix
-/// Delivered via [`EmbeddedEngine::watch_membership`] whenever a `ConfChange`
+/// Delivered via [`crate::api::EmbeddedEngine::watch_membership`] whenever a `ConfChange`
🤖 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/membership/membership_snapshot.rs` around lines 3 - 7,
Doc link to EmbeddedEngine::watch_membership will not resolve from the
membership module; update the intra-doc link in membership_snapshot.rs to use
the full path (e.g. `crate::api::EmbeddedEngine::watch_membership`) or enable
broken intra-doc links at the crate root (add `#![warn(broken_intra_doc_links)]`
or `#![allow(broken_intra_doc_links)]` as appropriate) so the bracketed
reference resolves correctly; locate the bracketed link in the membership
snapshot doc comment and replace it with the explicit path or add the
crate-level attribute to fix the link.

35-46: 💤 Low value

Consider documenting the Default semantics for committed_index.

MembershipSnapshot::default() yields committed_index = 0, which is also the value used to seed the watch channel in mock_node_builder.rs (Lines 321, 370). The doc comment states the index is "strictly monotonically increasing", but a real ConfChange could in principle commit at index 0 (or, more realistically, the first published snapshot may share 0 with the seed). Subscribers using committed_index <= last_applied as an idempotency guard need to know whether 0 is a valid sentinel "no snapshot yet" value or a legitimate snapshot index. A short note in the rustdoc would prevent confusion.

🤖 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/membership/membership_snapshot.rs` around lines 35 - 46,
Update the rustdoc for MembershipSnapshot::committed_index to state the Default
semantics clearly: indicate that MembershipSnapshot::default() sets
committed_index = 0 and that 0 is used as the sentinel "no snapshot yet" value
(and matches the seed used by the watch channel), or alternatively document that
0 can be a valid index and recommend a specific idempotency check (e.g., use
strict < or > comparisons rather than <=). Modify the comment on the
committed_index field in the MembershipSnapshot struct to explicitly describe
which convention is followed so subscribers using committed_index <=
last_applied know how to interpret 0.
d-engine-core/src/raft_role/leader_state.rs (2)

1966-1990: 💤 Low value

Reconnect loop guard looks correct; consider also bailing on shutdown.

The new task_rx.is_closed() check at the top of the reconnect loop correctly prevents the worker from spinning on open_replication_stream after the leader has stepped down (verified by test_replication_worker_exits_when_handle_dropped). One small consideration: if the underlying transport keeps failing fast (Err immediately) with a long max_delay_ms, the worker now sleeps via tokio::time::sleep without a select on a shutdown signal — is_closed() is only re-checked once per backoff cycle. That's acceptable since dropping the handle still ensures eventual exit, but a tokio::select! over sleep and a closed-channel future would shave off up to max_delay_ms of idle time on step-down. Fine to leave as-is given the existing test coverage.

🤖 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.rs` around lines 1966 - 1990, The
reconnect loop currently sleeps with tokio::time::sleep between retries, which
only re-checks task_rx.is_closed() once per backoff cycle; update the retry
sleep to use tokio::select! so the worker can exit immediately on shutdown:
replace the tokio::time::sleep(Duration::from_millis(backoff_ms)).await call
with a tokio::select! that awaits either the sleep future or task_rx.closed()
(or task_rx.is_closed()’s async closed notification) and break/return early if
the channel is closed; keep the existing backoff_ms doubling logic and error
logging around transport.open_replication_stream(peer_id, membership.clone(),
response_compress_enabled).

3620-3646: ⚡ Quick win

Warn-only zombie handling is intentional; consider emitting a metric for observability.

The shift from auto-BatchRemove to warn-only is well-justified by the comment (a node failing N attempts may simply be restarting). Since operators are now expected to act on these warnings, it would help to also emit a counter (similar to membership.stale_learner_removed at Line 3607) so dashboards/alerts can surface persistent zombies without scraping logs:

Suggested addition
         warn!(
             node_id,
             ?status,
             "Zombie detected: node is persistently unreachable — manual intervention may be required"
         );
+        metrics::counter!(
+            "membership.zombie_detected",
+            &[("node_id", node_id.to_string())]
+        )
+        .increment(1);

         Ok(())

Also note that _role_tx is now an unused parameter on this pub async method; since it's still required by the RaftRoleState::handle_zombie_detected trait method signature, keeping it (with the leading underscore) is the right call.

🤖 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.rs` around lines 3620 - 3646, Add an
observability counter increment when a zombie is detected: inside pub async fn
handle_zombie_node (the function shown) after confirming the node exists (i.e.,
where you currently call warn!), call the same membership metrics increment used
for membership.stale_learner_removed to increment a new counter (e.g.
"membership.zombie_detected" or similar) via the membership/metrics API (use
ctx.membership().<metrics increment helper> or the same helper used by
stale_learner_removed) and then emit the warn log; keep the unused _role_tx
parameter as-is.
d-engine-server/src/network/grpc/grpc_raft_service.rs (1)

528-558: 💤 Low value

Remove redundant mark_changed() call and use the public accessor for consistency.

The tokio_stream::wrappers::WatchStream::new (in 0.1.16) yields the current value on its first poll automatically via an internal borrow_and_update() call. The rx.mark_changed() on line 539 is redundant; WatchStream handles the initial snapshot delivery without it. The comment incorrectly attributes this behavior to mark_changed().

Additionally, per the pattern already used in d-engine-server/src/api/embedded.rs, use the public accessor self.membership_change_notifier() instead of directly accessing self.membership_rx. This maintains consistency with the embedded API and provides a single chokepoint if the storage type ever changes.

♻️ Suggested cleanup
-        let mut rx = self.membership_rx.clone();
-        // Deliver the current snapshot immediately; subsequent items arrive on each ConfChange.
-        rx.mark_changed();
+        // WatchStream::new yields the current value on first poll, so subscribers always
+        // receive the snapshot at subscribe time, then one item per committed ConfChange.
+        let rx = self.membership_change_notifier();
🤖 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/network/grpc/grpc_raft_service.rs` around lines 528 -
558, Remove the redundant rx.mark_changed() call in watch_membership and switch
to using the public accessor membership_change_notifier() instead of directly
cloning self.membership_rx: locate the watch_membership method, remove the
mark_changed() invocation (it’s unnecessary because
tokio_stream::wrappers::WatchStream::new yields the current value on first
poll), replace let mut rx = self.membership_rx.clone() with let mut rx =
self.membership_change_notifier().clone() (or the exact accessor call used
elsewhere), and keep the rest of the stream mapping and chain logic unchanged.
🤖 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/worker_lifecycle_test.rs`:
- Around line 425-427: The timeout only unwraps the Result from
tokio::time::timeout but not the inner Option from first_attempt_rx.recv(), so
add a second expect (or assert_some) after the timeout to fail if recv()
returned None; specifically, change the await chain around
tokio::time::timeout(std::time::Duration::from_secs(5),
first_attempt_rx.recv()).await.expect("timed out waiting for first reconnect
attempt") to also unwrap the Option (e.g., .expect("first reconnect signal not
sent") on the value returned) so the test fails if the channel was closed
without sending.

In `@d-engine-proto/proto/client/client_api.proto`:
- Around line 179-181: The doc comment for the MembershipSnapshot/ConfChange
event is inconsistent: it mentions "Remove" but the code/tests use
"BatchRemove"; update the comment near MembershipSnapshot and ConfChange to use
the exact operation name "BatchRemove" (or list both as "Remove (BatchRemove)"
if you want backward clarity) so client-facing docs match the actual API
operation names such as AddNode, Promote, and BatchRemove.

In `@d-engine-server/src/membership/raft_membership.rs`:
- Around line 444-460: When AddLearner re-applies and finds an existing learner
(the if let Some(existing) guard.nodes branch), avoid silently returning Ok(())
if the incoming address or status differs; instead compare incoming
address/status with existing.address and existing.status and emit a warn!
(including self.node_id, node_id, existing.address, existing.status and the
incoming values) before the early return so operators see stale-address or
status mismatches; keep the early return to preserve idempotence but add that
warning in the AddLearner handling code path.

In `@d-engine-server/tests/watch_and_subscriptions/watch_membership_embedded.rs`:
- Line 471: The test doc comments incorrectly state Promotable is value 2 while
the code constant STATUS_PROMOTABLE: i32 = 1 (matching
d_engine_proto::common::NodeStatus::Promotable); update the doc comment in
watch_membership_embedded.rs (the block around the test that currently says
"Promotable (2)" and "status=2") to "Promotable (1)" and "status=1" (also apply
the same correction to the other comment block at the referenced second
occurrence around lines 535-537) so the comments match the STATUS_PROMOTABLE
constant and the proto enum.
- Around line 141-154: Update the comment above with_fast_zombie to remove the
stale reference to "tests 3 and 7" and the incorrect claim about node removal:
state that it is used only by test_watch_membership_zombie_warns_without_removal
(test 3), that test_watch_membership_committed_index_monotonically_increasing
(test 7) does not use this override and instead builds TOML via
make_toml/node_toml, and clarify that with_fast_zombie triggers immediate zombie
warnings only (no automatic BatchRemove or node removal) per current design.

---

Outside diff comments:
In `@d-engine-server/tests/watch_and_subscriptions/mod.rs`:
- Around line 6-9: Update the module doc comment at the top of
tests/watch_and_subscriptions/mod.rs to reflect the current submodules: replace
the stale entries (e.g., watch_events_embedded.rs,
watch_events_grpc_standalone.rs, watch_performance_gate_embedded.rs) with the
actual submodules now declared (include watch_membership_embedded and
watch_membership_standalone) and add both names to the "## Test Coverage"
section so those two new submodules are listed; ensure the submodule list and
coverage section match the declared mod statements (watch_membership_embedded,
watch_membership_standalone).

---

Nitpick comments:
In `@d-engine-core/src/raft_role/leader_state.rs`:
- Around line 1966-1990: The reconnect loop currently sleeps with
tokio::time::sleep between retries, which only re-checks task_rx.is_closed()
once per backoff cycle; update the retry sleep to use tokio::select! so the
worker can exit immediately on shutdown: replace the
tokio::time::sleep(Duration::from_millis(backoff_ms)).await call with a
tokio::select! that awaits either the sleep future or task_rx.closed() (or
task_rx.is_closed()’s async closed notification) and break/return early if the
channel is closed; keep the existing backoff_ms doubling logic and error logging
around transport.open_replication_stream(peer_id, membership.clone(),
response_compress_enabled).
- Around line 3620-3646: Add an observability counter increment when a zombie is
detected: inside pub async fn handle_zombie_node (the function shown) after
confirming the node exists (i.e., where you currently call warn!), call the same
membership metrics increment used for membership.stale_learner_removed to
increment a new counter (e.g. "membership.zombie_detected" or similar) via the
membership/metrics API (use ctx.membership().<metrics increment helper> or the
same helper used by stale_learner_removed) and then emit the warn log; keep the
unused _role_tx parameter as-is.

In `@d-engine-server/src/membership/membership_snapshot.rs`:
- Around line 3-7: Doc link to EmbeddedEngine::watch_membership will not resolve
from the membership module; update the intra-doc link in membership_snapshot.rs
to use the full path (e.g. `crate::api::EmbeddedEngine::watch_membership`) or
enable broken intra-doc links at the crate root (add
`#![warn(broken_intra_doc_links)]` or `#![allow(broken_intra_doc_links)]` as
appropriate) so the bracketed reference resolves correctly; locate the bracketed
link in the membership snapshot doc comment and replace it with the explicit
path or add the crate-level attribute to fix the link.
- Around line 35-46: Update the rustdoc for MembershipSnapshot::committed_index
to state the Default semantics clearly: indicate that
MembershipSnapshot::default() sets committed_index = 0 and that 0 is used as the
sentinel "no snapshot yet" value (and matches the seed used by the watch
channel), or alternatively document that 0 can be a valid index and recommend a
specific idempotency check (e.g., use strict < or > comparisons rather than <=).
Modify the comment on the committed_index field in the MembershipSnapshot struct
to explicitly describe which convention is followed so subscribers using
committed_index <= last_applied know how to interpret 0.

In `@d-engine-server/src/network/grpc/grpc_raft_service.rs`:
- Around line 528-558: Remove the redundant rx.mark_changed() call in
watch_membership and switch to using the public accessor
membership_change_notifier() instead of directly cloning self.membership_rx:
locate the watch_membership method, remove the mark_changed() invocation (it’s
unnecessary because tokio_stream::wrappers::WatchStream::new yields the current
value on first poll), replace let mut rx = self.membership_rx.clone() with let
mut rx = self.membership_change_notifier().clone() (or the exact accessor call
used elsewhere), and keep the rest of the stream mapping and chain logic
unchanged.

In `@d-engine-server/tests/common/mod.rs`:
- Around line 595-614: The wait_for_stable_leader function currently swallows
refresh errors via client.refresh(None).await.ok() and can loop forever on
repeated ConnectionTimeouts; change it to enforce an explicit deadline (e.g.,
use tokio::time::timeout or track Instant::now() + duration) and propagate
refresh errors instead of ignoring them (replace the .ok() with ?.await or check
the Result), returning a clear timeout/err result if the deadline is exceeded;
update references in wait_for_stable_leader and the loop that matches
client.get(...) so the function returns an Err with a descriptive error when the
overall timeout is reached.
- Around line 600-611: The retry loop currently matches only StaleOperation and
ConnectionTimeout; add two new match arms to also handle
ClientApiError::Business { code: ErrorCode::NotLeader, .. } and
ClientApiError::Network { code: ErrorCode::LeaderChanged, .. } so they behave
like the retry cases: await tokio::time::sleep(Duration::from_millis(100)).await
and then continue the loop rather than returning Err(e). Update the match in the
function containing the shown match on Err(ClientApiError::...) to include these
two arms so linearizable reads and leader changes are retried with the same
backoff.
🪄 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: 08a4b1c7-8c9b-4550-bc20-30d991d1afe9

📥 Commits

Reviewing files that changed from the base of the PR and between 8c2e998 and afd2dda.

⛔ Files ignored due to path filters (3)
  • d-engine-proto/src/generated/d_engine.client.rs is excluded by !**/generated/**
  • examples/single-node-expansion/Cargo.lock is excluded by !**/*.lock
  • examples/three-nodes-embedded/Cargo.lock is excluded by !**/*.lock
📒 Files selected for processing (33)
  • CHANGELOG.md
  • d-engine-client/src/grpc_client.rs
  • d-engine-client/src/grpc_client_test.rs
  • d-engine-client/src/mock_rpc.rs
  • d-engine-client/src/mock_rpc_service.rs
  • d-engine-core/src/raft_role/follower_state_test.rs
  • d-engine-core/src/raft_role/leader_state.rs
  • d-engine-core/src/raft_role/leader_state_test/event_handling_test.rs
  • d-engine-core/src/raft_role/leader_state_test/membership_change_test.rs
  • d-engine-core/src/raft_role/leader_state_test/worker_lifecycle_test.rs
  • d-engine-core/src/test_utils/mock/mock_rpc.rs
  • d-engine-proto/proto/client/client_api.proto
  • d-engine-server/src/api/embedded.rs
  • d-engine-server/src/lib.rs
  • d-engine-server/src/membership/membership_snapshot.rs
  • d-engine-server/src/membership/mod.rs
  • d-engine-server/src/membership/raft_membership.rs
  • d-engine-server/src/membership/raft_membership_test.rs
  • d-engine-server/src/network/grpc/grpc_raft_service.rs
  • d-engine-server/src/network/grpc/grpc_raft_service_test.rs
  • d-engine-server/src/network/grpc/grpc_transport.rs
  • d-engine-server/src/node/builder.rs
  • d-engine-server/src/node/mod.rs
  • d-engine-server/src/test_utils/mock/mock_node_builder.rs
  • d-engine-server/src/test_utils/mock/mock_rpc.rs
  • d-engine-server/tests/cas_operations/leader_failover_cas_standalone.rs
  • d-engine-server/tests/cas_operations/snapshot_recovery_standalone.rs
  • d-engine-server/tests/cluster_lifecycle/scale_single_to_three_node_embedded.rs
  • d-engine-server/tests/common/mod.rs
  • d-engine-server/tests/failover_and_recovery/leader_failover_standalone.rs
  • d-engine-server/tests/watch_and_subscriptions/mod.rs
  • d-engine-server/tests/watch_and_subscriptions/watch_membership_embedded.rs
  • d-engine-server/tests/watch_and_subscriptions/watch_membership_standalone.rs

Comment thread d-engine-core/src/raft_role/leader_state_test/worker_lifecycle_test.rs Outdated
Comment thread d-engine-proto/proto/client/client_api.proto
Comment thread d-engine-server/src/membership/raft_membership.rs
Comment thread d-engine-server/tests/watch_and_subscriptions/watch_membership_embedded.rs Outdated
@JoshuaChi

Copy link
Copy Markdown
Contributor Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented May 5, 2026

Copy link
Copy Markdown
✅ Actions performed

Review triggered.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@JoshuaChi
JoshuaChi merged commit 741dd4f into main May 5, 2026
9 checks passed
@JoshuaChi
JoshuaChi deleted the feature/327-328-watch-membership branch May 5, 2026 04:35
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant