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: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,15 @@ All notable changes to this project will be documented in this file.
- Stream terminates with `UNAVAILABLE` on server shutdown — callers reconnect and resubscribe
- `MembershipSnapshot` carries `members`, `learners`, and `committed_index` (idempotency key for deduplication after reconnect)

- **Watch: prefix watch and watcher limits** (#300):
- `watch_prefix(prefix)` API for both embedded and gRPC clients — one watcher covers an entire key namespace (e.g. `/services/payment/`) without registering a per-key watcher
- Prefix semantics: prefix must start and end with `/`; matching uses slash-boundary decomposition so `/services/payment/node1` matches `/services/payment/` but not `/services/payment`
- `WatchConfig.max_watcher_count` — hard cap on total active watchers (exact + prefix); `register()` / `watch_prefix()` return `LimitExceeded` error once reached (default: effectively unlimited)
- `WatchConfig.watcher_buffer_size` default raised from 10 → 256, reducing spurious `CANCELED` events under burst load
- Every `WatchEvent` carries `revision: u64` (Raft applied index) — use it to detect missed events after reconnection
- Architecture: two separate `DashMap`s (exact and prefix) — O(1) exact lookup, O(depth) prefix dispatch; no linear scan
- Proto: added `prefix: bool` to `WatchRequest`; added `revision: uint64` to `WatchResponse`

- **Unified RocksDB option** (#295): opt-in via `storage.unified_db = true`.
Uses a single RocksDB instance with 4 column families instead of two separate instances.
- Reduces memory RSS and file descriptor usage
Expand Down
31 changes: 31 additions & 0 deletions d-engine-client/src/grpc_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -221,6 +221,7 @@ impl GrpcClient {
let request = WatchRequest {
client_id: client_inner.client_id,
key: Bytes::copy_from_slice(key.as_ref()),
prefix: false,
};

// Watch can connect to any node (leader or follower)
Expand All @@ -237,6 +238,36 @@ impl GrpcClient {
}
}
}

/// Watch all keys under a path prefix.
///
/// `prefix` must start with '/' and end with '/', e.g. `b"/services/"`.
/// Returns a stream of events for any key whose path begins with the prefix.
pub async fn watch_prefix(
&self,
prefix: impl AsRef<[u8]>,
) -> ClientApiResult<tonic::Streaming<WatchResponse>> {
let client_inner = self.client_inner.load();

let request = WatchRequest {
client_id: client_inner.client_id,
key: Bytes::copy_from_slice(prefix.as_ref()),
prefix: true,
};

let mut client = self.make_client().await?;

match client.watch(request).await {
Ok(response) => {
debug!("Prefix watch stream established");
Ok(response.into_inner())
}
Err(status) => {
error!("Prefix watch request failed: {:?}", status);
Err(status.into())
}
}
}
}

// ==================== Core ClientApi Trait Implementation ====================
Expand Down
2 changes: 2 additions & 0 deletions d-engine-client/src/grpc_client_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1594,12 +1594,14 @@ mod watch_tests {
value: Bytes::from("owner-a"),
event_type: WatchEventType::Put as i32,
error: 0,
revision: 0,
},
WatchResponse {
key: Bytes::from("lock"),
value: Bytes::new(),
event_type: WatchEventType::Delete as i32,
error: 0,
revision: 0,
},
];

Expand Down
1 change: 1 addition & 0 deletions d-engine-client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ mod utils;
pub use builder::*;
pub use config::*;
pub use d_engine_core::client::{ClientApi, ClientApiError, ClientApiResult};
pub use d_engine_proto::error::ErrorCode;
pub use grpc_client::*;
pub(crate) use pool::*;

Expand Down
22 changes: 20 additions & 2 deletions d-engine-core/src/config/raft.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1190,8 +1190,9 @@ fn default_client_compression() -> bool {
/// ```toml
/// [raft.watch]
/// event_queue_size = 10240
/// watcher_buffer_size = 10
/// watcher_buffer_size = 256
/// enable_metrics = false
/// max_watcher_count = 10000
/// ```
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct WatchConfig {
Expand Down Expand Up @@ -1242,6 +1243,16 @@ pub struct WatchConfig {
/// **Default**: false (minimal overhead in production)
#[serde(default = "default_enable_watch_metrics")]
pub enable_metrics: bool,

/// Hard cap on total active watchers (exact + prefix combined).
///
/// `register()` and `register_prefix()` return `WatchError::LimitExceeded`
/// once this count is reached. Leave at the default (effectively unlimited)
/// or set to a finite value to prevent runaway watcher growth.
///
/// **Default**: i64::MAX (effectively unlimited)
#[serde(default = "default_max_watcher_count")]
pub max_watcher_count: usize,
}

impl Default for WatchConfig {
Expand All @@ -1250,6 +1261,7 @@ impl Default for WatchConfig {
event_queue_size: default_event_queue_size(),
watcher_buffer_size: default_watcher_buffer_size(),
enable_metrics: default_enable_watch_metrics(),
max_watcher_count: default_max_watcher_count(),
}
}
}
Expand Down Expand Up @@ -1294,13 +1306,19 @@ const fn default_event_queue_size() -> usize {
}

const fn default_watcher_buffer_size() -> usize {
10
256
}

const fn default_enable_watch_metrics() -> bool {
false
}

const fn default_max_watcher_count() -> usize {
// i64::MAX — the config crate uses signed 64-bit internally, so usize::MAX overflows it.
// This value is effectively unlimited for any realistic deployment.
9_223_372_036_854_775_807
}

/// Performance metrics configuration
///
/// Controls emission of observability metrics. Disabling metrics reduces overhead
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -658,12 +658,14 @@ where
value: insert.value,
event_type: WatchEventType::Put as i32,
error: 0,
revision: entry.index,
}),
Some(Operation::Delete(delete)) => Some(WatchResponse {
key: delete.key,
value: bytes::Bytes::new(),
event_type: WatchEventType::Delete as i32,
error: 0,
revision: entry.index,
}),
Some(Operation::CompareAndSwap(cas)) => {
// Only broadcast if CAS actually mutated the value.
Expand All @@ -674,6 +676,7 @@ where
value: cas.new_value,
event_type: WatchEventType::Put as i32,
error: 0,
revision: entry.index,
})
} else {
None
Expand Down
Loading
Loading