Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 3 additions & 6 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,9 @@ All notable changes to this project will be documented in this file.

### Added

- **ReadActor fast path for Eventual/LeaseRead** (#392): Dedicated `ReadActor` task serves `EventualConsistency` and `LeaseRead` without entering the Raft loop, eliminating channel contention under high read concurrency.
- `Arc<SM>` is owned exclusively by ReadActor — guarantees RocksDB LOCK release before Raft shutdown on `stop()`, fixing `test_snapshot_recovery_embedded`
- `ReadLease.revoke()`: atomic lease invalidation on leader demotion (replaces `invalidate()`)
- `EmbeddedClient` routes Eventual/Lease → ReadActor; falls back to `cmd_tx` on `LeaseInvalid`/`SmStopped`
- New `[raft.read_actor]` config section: `channel_capacity` (default 512), `max_drain` (default 100)
- **Benchmark (local embedded, 100 concurrent clients)**: Lease Read +10.9%, Eventual Read +13.1% vs v0.2.4; write throughput unchanged (see `benches/reports/v0.2.5/`)
- **AppendEntries coalescing + speculative replication** (#407): Consecutive same-term `AppendEntries` in the receive buffer are merged before dispatch (heartbeats absorbed via `max(leader_commit_index)`); leader advances `next_index` speculatively on send, eliminating stop-and-wait per replication round.

- **ReadActor fast path for Eventual/LeaseRead** (#392): Dedicated `ReadActor` serves `EventualConsistency` and `LeaseRead` off the Raft loop, eliminating channel contention under high read concurrency. Lease Read +10.9%, Eventual Read +13.1% vs v0.2.4 (100 concurrent clients, local embedded).

### Fixed

Expand Down
15 changes: 8 additions & 7 deletions benches/standalone-bench/src/main.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,10 @@
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
use std::time::Duration;
use std::time::Instant;

use clap::Parser;
use clap::Subcommand;
use d_engine::Client;
Expand All @@ -9,12 +16,6 @@ use hdrhistogram::Histogram;
use rand::Rng;
use rand::SeedableRng;
use rand::distributions::Alphanumeric;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
use std::time::Duration;
use std::time::Instant;
use tokio::sync::Semaphore;

#[derive(Parser, Debug, Clone)]
Expand Down Expand Up @@ -209,7 +210,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {

// Initialize the connection pool
let endpoints = cli.endpoints.clone();
println!("Initializing client connection with: {:?}", &endpoints);
println!("Initializing client connection with: {:?}", endpoints);
let client_pool = ClientPool::new(endpoints, cli.conns)
.await
.expect("Failed to create client pool");
Expand Down
17 changes: 17 additions & 0 deletions d-engine-core/src/config/raft.rs
Original file line number Diff line number Diff line change
Expand Up @@ -333,12 +333,19 @@ pub struct BatchingConfig {
/// **Default**: 100
#[serde(default = "default_max_batch_size")]
pub max_batch_size: usize,

/// Maximum total entries to accumulate in a single merge_append_entries pass.
/// Prevents unbounded memory growth when the buffer has many large batches.
/// Default: 1000
#[serde(default = "default_max_merge_entries")]
pub max_merge_entries: usize,
}

impl Default for BatchingConfig {
fn default() -> Self {
Self {
max_batch_size: default_max_batch_size(),
max_merge_entries: default_max_merge_entries(),
}
}
}
Expand All @@ -350,6 +357,11 @@ impl BatchingConfig {
"batching.max_batch_size must be > 0".into(),
)));
}
if self.max_merge_entries == 0 {
return Err(Error::Config(ConfigError::Message(
"batching.max_merge_entries must be > 0".into(),
)));
}
Ok(())
}
}
Expand All @@ -360,6 +372,11 @@ fn default_append_interval() -> u64 {
fn default_max_batch_size() -> usize {
100
}

fn default_max_merge_entries() -> usize {
1000
}

fn default_entries_per_replication() -> u64 {
100
}
Expand Down
7 changes: 6 additions & 1 deletion d-engine-core/src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,7 @@ pub enum RaftEvent {

AppendEntries(
AppendEntriesRequest,
MaybeCloneOneshotSender<std::result::Result<AppendEntriesResponse, Status>>,
Vec<MaybeCloneOneshotSender<std::result::Result<AppendEntriesResponse, Status>>>,
),

// Response snapshot stream from Leader
Expand Down Expand Up @@ -193,6 +193,11 @@ pub enum RaftEvent {
},
}

// SAFETY: tonic::Streaming variants are only accessed on the thread that created them;
// no shared references to snapshot variants cross thread boundaries.
// Full decoupling tracked in #409.
unsafe impl Sync for RaftEvent {}

#[cfg(test)]
#[cfg_attr(test, derive(Debug, Clone))]
#[allow(unused)]
Expand Down
Loading
Loading