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
8 changes: 8 additions & 0 deletions app/crowdb-chunk-kv-server/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -408,6 +408,12 @@ impl ChunkKvService {
let Ok(durable_bytes) = partition.estimated_bytes() else {
continue;
};
let summary = partition
.observe_tree()
.ok()
.and_then(|observation| observation.summary);
let (logical_bytes, logical_metrics_exact) =
summary.map_or((0, false), |summary| (summary.live_logical_bytes, summary.exact));
let live_byte_samples =
if max_samples >= 2 && snapshot.lifecycle == crowdb_chunk_kv::PartitionLifecycle::Serving {
self.load_sampling.samples(*id, snapshot.ownership_epoch)
Expand All @@ -420,6 +426,8 @@ impl ChunkKvService {
low: snapshot.partition_id.low,
},
durable_bytes,
logical_bytes,
logical_metrics_exact,
live_byte_samples,
independently_recoverable: self.independently_recoverable.load().contains(&(
Id128 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ struct PlanningState {
healthy: HashMap<u64, (InstanceValue, ChunkKvExtra)>,
partition_loads: HashMap<(u64, Id128), ChunkKvPartitionLoad>,
partition_bytes: HashMap<(u64, Id128), u64>,
partition_metrics_exact: HashMap<(u64, Id128), bool>,
hosted_epochs: HashMap<(u64, Id128), u64>,
active_partitions: HashSet<Id128>,
busy_owners: HashSet<u64>,
Expand Down Expand Up @@ -144,6 +145,15 @@ async fn planning_state(
.map(move |load| ((*id, load.partition_id), effective_bytes(load)))
})
.collect();
let partition_metrics_exact = healthy
.iter()
.flat_map(|(id, (_, extra))| {
extra
.partition_loads
.iter()
.map(move |load| ((*id, load.partition_id), load.logical_metrics_exact))
})
.collect();
let hosted_epochs = healthy
.iter()
.flat_map(|(id, (_, extra))| {
Expand All @@ -158,6 +168,7 @@ async fn planning_state(
healthy,
partition_loads,
partition_bytes,
partition_metrics_exact,
hosted_epochs,
active_partitions: HashSet::new(),
busy_owners: HashSet::new(),
Expand Down Expand Up @@ -230,7 +241,12 @@ async fn plan_split(
continue;
};
let effective_bytes = state.partition_bytes[&(entry.owner.instance_id, entry.partition_id)];
if (count_shortfall || effective_bytes > policy.target_partition_bytes)
let exact_size = state
.partition_metrics_exact
.get(&(entry.owner.instance_id, entry.partition_id))
.copied()
.unwrap_or(false);
if (count_shortfall || exact_size && effective_bytes > policy.target_partition_bytes)
&& eligible(entry, state)
&& cooled_down(entry, &state.last_changed_ms, policy, now_ms)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,9 @@ pub(super) fn partition_load<'a>(
}

pub(super) fn effective_bytes(load: &ChunkKvPartitionLoad) -> u64 {
if load.logical_metrics_exact {
return load.logical_bytes;
}
load.durable_bytes.max(
load.live_byte_samples
.iter()
Expand Down
21 changes: 20 additions & 1 deletion app/crowdb-kv-server/src/mgmt/group_ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
//! readiness, and async operation polling.

use axum::extract::{Path, Query, State};
use axum::http::StatusCode;
use axum::http::{HeaderMap, StatusCode};
use axum::response::IntoResponse;
use axum::Json;
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -447,11 +447,30 @@ pub(super) async fn join_group_via_snapshot(
pub(super) async fn remove_group(
State(state): State<RegistryArc>,
Path((sid, gid)): Path<(u64, u64)>,
headers: HeaderMap,
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
let store = state
.get_store(sid)
.ok_or_else(|| err_json(StatusCode::NOT_FOUND, format!("store {sid} not found")))?;

let group = store.get_group(gid).ok_or_else(|| {
err_json(
StatusCode::NOT_FOUND,
format!("group {gid} not found in store {sid}"),
)
})?;
let membership_guard = store.membership_guard(gid);
let _membership_guard = membership_guard.lock().await;
let expected_epoch = super::replica_ops::check_expected_epoch(&headers, group.membership_epoch())?;
if expected_epoch == Some(group.membership_epoch()) {
return Err(err_json(
StatusCode::CONFLICT,
format!(
"membership epoch {} already has a different configuration",
group.membership_epoch()
),
));
}
info!(s = sid, g = gid, "removing PxGroup via management API");
if !store.remove_group(gid) {
return Err(err_json(
Expand Down
105 changes: 103 additions & 2 deletions app/crowdb-kv-server/src/mgmt/replica_ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
//! Replica management endpoints: list, add, remove, batch-add remote replicas.

use axum::extract::{Path, State};
use axum::http::StatusCode;
use axum::http::{HeaderMap, StatusCode};
use axum::Json;
use tracing::{debug, info};

Expand Down Expand Up @@ -42,7 +42,6 @@ pub(super) async fn list_remote_replicas(
format!("group {gid} not found in store {sid}"),
)
})?;

let remotes: Vec<RemoteReplicaInfo> = group
.remote_replica_info()
.into_iter()
Expand Down Expand Up @@ -74,6 +73,7 @@ pub(super) async fn list_remote_replicas(
pub(super) async fn add_remote_replicas(
State(state): State<RegistryArc>,
Path((sid, gid)): Path<(u64, u64)>,
headers: HeaderMap,
Json(remotes): Json<Vec<RemoteReplicaInfo>>,
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
let store = state
Expand All @@ -85,6 +85,9 @@ pub(super) async fn add_remote_replicas(
format!("group {gid} not found in store {sid}"),
)
})?;
let membership_guard = store.membership_guard(gid);
let _membership_guard = membership_guard.lock().await;
let expected_epoch = check_expected_epoch(&headers, group.membership_epoch())?;

let local_id = group.local_replica().id;
for r in &remotes {
Expand All @@ -98,6 +101,18 @@ pub(super) async fn add_remote_replicas(
));
}
}
if expected_epoch == Some(group.membership_epoch()) {
if remotes_match_existing(&group, &remotes) {
return Ok(StatusCode::OK);
}
return Err(err_json(
StatusCode::CONFLICT,
format!(
"membership epoch {} already has a different configuration",
group.membership_epoch()
),
));
}

debug!(
s = sid,
Expand All @@ -120,6 +135,10 @@ pub(super) async fn add_remote_replicas(
.map(|r| (r.replica_id, r.endpoint.clone(), r.voting))
.collect();
let new_group = rebuild_group_with_new_remotes(&group, &new_remotes);
let new_group = new_group;
if let Some(epoch) = expected_epoch {
new_group.set_membership_epoch(epoch);
}
new_group.local_replica().set_endpoint(
store
.listen_addr()
Expand Down Expand Up @@ -161,6 +180,7 @@ pub(super) async fn add_remote_replicas(
pub(super) async fn remove_remote_replica(
State(state): State<RegistryArc>,
Path((sid, gid, rid)): Path<(u64, u64, u64)>,
headers: HeaderMap,
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
let store = state
.get_store(sid)
Expand All @@ -171,6 +191,9 @@ pub(super) async fn remove_remote_replica(
format!("group {gid} not found in store {sid}"),
)
})?;
let membership_guard = store.membership_guard(gid);
let _membership_guard = membership_guard.lock().await;
let expected_epoch = check_expected_epoch(&headers, group.membership_epoch())?;

let local_id = group.local_replica().id;
if rid == local_id {
Expand All @@ -183,11 +206,23 @@ pub(super) async fn remove_remote_replica(
// Check if remote exists
let exists = group.remote_replica_info().iter().any(|(id, _, _)| *id == rid);
if !exists {
if expected_epoch == Some(group.membership_epoch()) {
return Ok(StatusCode::OK);
}
return Err(err_json(
StatusCode::NOT_FOUND,
format!("remote replica {rid} not found in group {gid}"),
));
}
if expected_epoch == Some(group.membership_epoch()) {
return Err(err_json(
StatusCode::CONFLICT,
format!(
"membership epoch {} already has a different configuration",
group.membership_epoch()
),
));
}

info!(
s = sid,
Expand All @@ -211,6 +246,9 @@ pub(super) async fn remove_remote_replica(
.collect(),
);
new_group.remove_remote_replica(rid);
if let Some(epoch) = expected_epoch {
new_group.set_membership_epoch(epoch);
}
let current_term = group.local_replica().current_term_snapshot();
if new_group.quorum() == 1 {
new_group.local_replica().become_leader();
Expand Down Expand Up @@ -260,6 +298,7 @@ pub(super) async fn remove_remote_replica(
pub(super) async fn batch_add_remote_replicas(
State(state): State<RegistryArc>,
Path((sid, gid)): Path<(u64, u64)>,
headers: HeaderMap,
Json(topology): Json<TopologyResponse>,
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
let store = state
Expand All @@ -271,6 +310,9 @@ pub(super) async fn batch_add_remote_replicas(
format!("group {gid} not found in store {sid}"),
)
})?;
let membership_guard = store.membership_guard(gid);
let _membership_guard = membership_guard.lock().await;
let expected_epoch = check_expected_epoch(&headers, group.membership_epoch())?;

let local_id = group.local_replica().id;
let mut new_remotes = Vec::new();
Expand All @@ -295,6 +337,18 @@ pub(super) async fn batch_add_remote_replicas(
info!(s = sid, g = gid, "batch add remotes: no new remotes to add");
return Ok(StatusCode::OK);
}
if expected_epoch == Some(group.membership_epoch()) {
if remotes_match_existing(&group, &new_remotes) {
return Ok(StatusCode::OK);
}
return Err(err_json(
StatusCode::CONFLICT,
format!(
"membership epoch {} already has a different configuration",
group.membership_epoch()
),
));
}

debug!(
s = sid,
Expand All @@ -317,6 +371,10 @@ pub(super) async fn batch_add_remote_replicas(
.map(|r| (r.replica_id, r.endpoint.clone(), r.voting))
.collect();
let new_group = rebuild_group_with_new_remotes(&group, &remotes_tuple);
let new_group = new_group;
if let Some(epoch) = expected_epoch {
new_group.set_membership_epoch(epoch);
}
new_group.local_replica().set_endpoint(
store
.listen_addr()
Expand All @@ -339,3 +397,46 @@ pub(super) async fn batch_add_remote_replicas(
);
Ok(StatusCode::OK)
}

const MEMBERSHIP_EPOCH_HEADER: &str = "x-crowdb-membership-epoch";

pub(super) fn check_expected_epoch(
headers: &HeaderMap,
actual: u64,
) -> Result<Option<u64>, (StatusCode, Json<ErrorResponse>)> {
let Some(value) = headers.get(MEMBERSHIP_EPOCH_HEADER) else {
return Ok(None);
};
let expected = value
.to_str()
.ok()
.and_then(|value| value.parse::<u64>().ok())
.ok_or_else(|| err_json(StatusCode::BAD_REQUEST, "invalid membership epoch header"))?;
// A freshly-created replacement replica has no persisted peers yet, so
// it can legitimately join an already-advanced authority epoch. The
// endpoint callers pass the complete successor from Group 0; retain a
// narrow bootstrap allowance for a newly-created replica (the supported
// membership size is at most seven voters) while still rejecting
// arbitrary jumps such as stale requests with epoch 99.
if expected != actual
&& expected != actual.saturating_add(1)
&& !(actual == 0 && (2..=8).contains(&expected))
{
return Err(err_json(
StatusCode::CONFLICT,
format!("membership epoch conflict: expected {expected}, actual {actual}"),
));
}
Ok(Some(expected))
}

fn remotes_match_existing(group: &crowdb_kv::cluster::group::PxGroup, remotes: &[RemoteReplicaInfo]) -> bool {
let existing = group.remote_replica_info();
remotes.iter().all(|requested| {
existing.iter().any(|(id, endpoint, voting)| {
*id == requested.replica_id
&& *endpoint == requested.endpoint.as_str()
&& *voting == requested.voting
})
})
}
Loading
Loading