mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
volume server: take the heartbeat's remaining store reads off the runtime (#11540)
* volume server: collect EC heartbeats and adjust volume max off the runtime The volume pass moved to the blocking pool, but the heartbeat task still called collect_ec_heartbeat and the following EC shard snapshot, and Store::maybe_adjust_volume_max, directly on its tokio worker. maybe_adjust_volume_max runs statvfs on every auto-sized disk and stats the .dat of every writable volume under the store read lock. All of them block the worker on the node-wide RwLock<Store> whenever a writer holds it or is queued, and every task sharing that worker stalls with it. Run the adjustment, on the pulse and after the master changes volume options, and the EC tick's heartbeat plus shard snapshot through off_runtime, like the volume pass. apply_master_volume_options now only reports whether the options changed; the loop adjusts off the runtime. What is collected and sent, and in what order, is unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume server: keep EC shard deltas out of the volume heartbeat's snapshot EC shard deltas are the diff between the store's shards and last_ec_shards, taken when volume_state_notify fires. But the volume tick and the options-changed heartbeat re-took last_ec_shards from the store too, and a volume heartbeat carries no shard list: a mount or unmount that landed while the notify was pending or the volume pass was collecting was absorbed into the baseline and never sent. The master only learned of it at the next EC tick, 17 pulses later. The EC tick likewise built its heartbeat and its baseline under two separate store guards, so a mount between them was lost the same way. A volume heartbeat now only takes out of the baseline the expired EC shards it reports deleted itself, so the next delta does not repeat them. The EC tick, and the initial EC heartbeat, build the full list and the baseline under one read guard, still on the blocking pool. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume server: take the heartbeat's remaining store reads off the runtime The heartbeat's volume and EC passes and the volume max adjustment already run on the blocking pool, but several reads of the node-wide RwLock<Store> were still taken directly on the heartbeat's tokio worker: the digest report reset before the first heartbeat, the duplicate-UUID directory lookup and the volume options a master response carries, the EC shard list a state notification is diffed against, and the deregistration heartbeat sent on stop and shutdown. The lock is writer-preferring, so with a writer holding or queued for it each of these parks the worker, and every task sharing that worker stalls with it. Run each through off_runtime, which now takes a closure so a pass can carry what it needs from the master's response. The notify branch's volume snapshot and EC read become one blocking pass, still under two guards in the same order. What is collected and sent, and in what order, is unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-authored-by: Chris Lu <chris.lu@gmail.com>
This commit is contained in:
1 file changed
+325
-63
@@ -373,6 +373,22 @@ fn diff_ec_shard_delta_messages(
|
||||
(new_ec_shards, deleted_ec_shards)
|
||||
}
|
||||
|
||||
/// A volume heartbeat carries no shard list, only the expired shards it
|
||||
/// deleted, so that is all it may take out of the delta baseline.
|
||||
fn forget_reported_ec_deletions(
|
||||
last_ec_shards: &mut HashMap<EcShardDeltaKey, master_pb::VolumeEcShardInformationMessage>,
|
||||
heartbeat: &master_pb::Heartbeat,
|
||||
) {
|
||||
last_ec_shards.retain(|(id, collection, disk_id, shard_id), _| {
|
||||
!heartbeat.deleted_ec_shards.iter().any(|deleted| {
|
||||
deleted.id == *id
|
||||
&& deleted.collection == *collection
|
||||
&& deleted.disk_id == *disk_id
|
||||
&& deleted.ec_index_bits & (1u32 << shard_id) != 0
|
||||
})
|
||||
});
|
||||
}
|
||||
|
||||
/// Perform one heartbeat session with a master server.
|
||||
async fn do_heartbeat(
|
||||
config: &HeartbeatConfig,
|
||||
@@ -393,24 +409,21 @@ async fn do_heartbeat(
|
||||
|
||||
let (tx, rx) = tokio::sync::mpsc::channel::<master_pb::Heartbeat>(32);
|
||||
|
||||
// This master may know nothing about this server, and has not yet said
|
||||
// whether it understands digests, so start from the whole list.
|
||||
state.store.read().unwrap().volume_report.reset();
|
||||
|
||||
// Keep track of what we sent, to generate delta updates
|
||||
let (initial_hb, initial_volumes) =
|
||||
off_runtime(config, state, collect_heartbeat_with_snapshot).await?;
|
||||
let (initial_hb, initial_volumes) = off_runtime(config, state, |config, state| {
|
||||
// This master may know nothing about this server, and has not yet said
|
||||
// whether it understands digests, so start from the whole list.
|
||||
state.store.read().unwrap().volume_report.reset();
|
||||
collect_heartbeat_with_snapshot(config, state)
|
||||
})
|
||||
.await?;
|
||||
let mut last_volumes: HashMap<u32, VolumeIdentity> = volume_identities(&initial_volumes);
|
||||
let mut last_ec_shards = {
|
||||
let store = state.store.read().unwrap();
|
||||
collect_ec_shard_delta_messages(&store)
|
||||
};
|
||||
|
||||
// Send initial heartbeats BEFORE calling send_heartbeat to avoid deadlock:
|
||||
// the server won't send response headers until it receives the first message,
|
||||
// but send_heartbeat().await waits for response headers.
|
||||
tx.send(initial_hb).await?;
|
||||
let initial_ec_hb = off_runtime(config, state, collect_ec_heartbeat).await?;
|
||||
let (initial_ec_hb, mut last_ec_shards) = ec_tick_pass(config, state).await?;
|
||||
tx.send(initial_ec_hb).await?;
|
||||
|
||||
let stream = tokio_stream::wrappers::ReceiverStream::new(rx);
|
||||
@@ -446,10 +459,11 @@ async fn do_heartbeat(
|
||||
// Match Go ordering: DuplicatedUuids first, then volume
|
||||
// options, then leader redirect.
|
||||
if !hb_resp.duplicated_uuids.is_empty() {
|
||||
let duplicate_dirs = {
|
||||
let store = state.store.read().unwrap();
|
||||
duplicate_directories(&store, &hb_resp.duplicated_uuids)
|
||||
};
|
||||
let uuids = hb_resp.duplicated_uuids.clone();
|
||||
let duplicate_dirs = off_runtime(config, state, move |_, state| {
|
||||
duplicate_directories(&state.store.read().unwrap(), &uuids)
|
||||
})
|
||||
.await?;
|
||||
error!(
|
||||
"Master reported duplicate volume directories: {:?}",
|
||||
duplicate_dirs
|
||||
@@ -460,19 +474,14 @@ async fn do_heartbeat(
|
||||
)
|
||||
.into());
|
||||
}
|
||||
let options_changed = {
|
||||
let s = state.store.read().unwrap();
|
||||
apply_master_volume_options(&s, &hb_resp)
|
||||
};
|
||||
let options_changed =
|
||||
volume_options_pass(config, state, &hb_resp).await?;
|
||||
if options_changed && maybe_adjust_volume_max(config, state).await? {
|
||||
let (adjusted_hb, adjusted_volumes) =
|
||||
off_runtime(config, state, collect_heartbeat_with_snapshot)
|
||||
.await?;
|
||||
last_volumes = volume_identities(&adjusted_volumes);
|
||||
last_ec_shards = {
|
||||
let store = state.store.read().unwrap();
|
||||
collect_ec_shard_delta_messages(&store)
|
||||
};
|
||||
forget_reported_ec_deletions(&mut last_ec_shards, &adjusted_hb);
|
||||
if tx.send(adjusted_hb).await.is_err() {
|
||||
return Ok(None);
|
||||
}
|
||||
@@ -501,10 +510,7 @@ async fn do_heartbeat(
|
||||
let (current_hb, current_volumes) =
|
||||
off_runtime(config, state, collect_heartbeat_with_snapshot).await?;
|
||||
last_volumes = volume_identities(¤t_volumes);
|
||||
last_ec_shards = {
|
||||
let store = state.store.read().unwrap();
|
||||
collect_ec_shard_delta_messages(&store)
|
||||
};
|
||||
forget_reported_ec_deletions(&mut last_ec_shards, ¤t_hb);
|
||||
if tx.send(current_hb).await.is_err() {
|
||||
return Ok(None);
|
||||
}
|
||||
@@ -525,12 +531,8 @@ async fn do_heartbeat(
|
||||
info!("Heartbeat stopping");
|
||||
return Ok(None);
|
||||
}
|
||||
let held_volumes = off_runtime(config, state, collect_volume_snapshot).await?;
|
||||
let (held_volumes, current_ec_shards) = notify_pass(config, state).await?;
|
||||
let current_volumes = volume_identities(&held_volumes);
|
||||
let current_ec_shards = {
|
||||
let store = state.store.read().unwrap();
|
||||
collect_ec_shard_delta_messages(&store)
|
||||
};
|
||||
|
||||
let mut new_vols = Vec::new();
|
||||
let mut del_vols = Vec::new();
|
||||
@@ -619,30 +621,41 @@ async fn send_deregister_heartbeat(
|
||||
state: &Arc<VolumeServerState>,
|
||||
tx: &tokio::sync::mpsc::Sender<master_pb::Heartbeat>,
|
||||
) {
|
||||
let empty = {
|
||||
let store = state.store.read().unwrap();
|
||||
// Deregister: no effective max computed, fall back to configured max.
|
||||
let (location_uuids, disk_tags) = collect_location_metadata(&store, &[]);
|
||||
master_pb::Heartbeat {
|
||||
id: store.id.clone(),
|
||||
ip: config.ip.clone(),
|
||||
port: config.port as u32,
|
||||
public_url: config.public_url.clone(),
|
||||
max_file_key: 0,
|
||||
data_center: config.data_center.clone(),
|
||||
rack: config.rack.clone(),
|
||||
has_no_volumes: true,
|
||||
has_no_ec_shards: true,
|
||||
grpc_port: config.grpc_port as u32,
|
||||
location_uuids,
|
||||
disk_tags,
|
||||
..Default::default()
|
||||
let empty = match off_runtime(config, state, deregister_heartbeat).await {
|
||||
Ok(empty) => empty,
|
||||
Err(e) => {
|
||||
warn!("Deregistration heartbeat not built: {}", e);
|
||||
return;
|
||||
}
|
||||
};
|
||||
let _ = tx.send(empty).await;
|
||||
tokio::time::sleep(Duration::from_millis(200)).await;
|
||||
}
|
||||
|
||||
fn deregister_heartbeat(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
) -> master_pb::Heartbeat {
|
||||
let store = state.store.read().unwrap();
|
||||
// Deregister: no effective max computed, fall back to configured max.
|
||||
let (location_uuids, disk_tags) = collect_location_metadata(&store, &[]);
|
||||
master_pb::Heartbeat {
|
||||
id: store.id.clone(),
|
||||
ip: config.ip.clone(),
|
||||
port: config.port as u32,
|
||||
public_url: config.public_url.clone(),
|
||||
max_file_key: 0,
|
||||
data_center: config.data_center.clone(),
|
||||
rack: config.rack.clone(),
|
||||
has_no_volumes: true,
|
||||
has_no_ec_shards: true,
|
||||
grpc_port: config.grpc_port as u32,
|
||||
location_uuids,
|
||||
disk_tags,
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
fn apply_metrics_push_settings(
|
||||
state: &VolumeServerState,
|
||||
address: &str,
|
||||
@@ -781,16 +794,54 @@ fn volume_identities(
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Runs a store pass on the blocking pool: it stats every volume's files.
|
||||
/// Runs a store pass on the blocking pool, so waiting for the store lock
|
||||
/// parks no runtime worker.
|
||||
async fn off_runtime<T: Send + 'static>(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
pass: fn(&HeartbeatConfig, &Arc<VolumeServerState>) -> T,
|
||||
pass: impl FnOnce(&HeartbeatConfig, &Arc<VolumeServerState>) -> T + Send + 'static,
|
||||
) -> Result<T, tokio::task::JoinError> {
|
||||
let (config, state) = (config.clone(), state.clone());
|
||||
tokio::task::spawn_blocking(move || pass(&config, &state)).await
|
||||
}
|
||||
|
||||
/// Applies the master's volume options; whether they changed.
|
||||
async fn volume_options_pass(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
hb_resp: &master_pb::HeartbeatResponse,
|
||||
) -> Result<bool, tokio::task::JoinError> {
|
||||
let hb_resp = hb_resp.clone();
|
||||
off_runtime(config, state, move |_, state| {
|
||||
apply_master_volume_options(&state.store.read().unwrap(), &hb_resp)
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// The held volumes and EC shards a state notification is diffed from.
|
||||
async fn notify_pass(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
) -> Result<
|
||||
(
|
||||
Vec<master_pb::VolumeInformationMessage>,
|
||||
HashMap<EcShardDeltaKey, master_pb::VolumeEcShardInformationMessage>,
|
||||
),
|
||||
tokio::task::JoinError,
|
||||
> {
|
||||
off_runtime(config, state, |config, state| {
|
||||
let volumes = collect_volume_snapshot(config, state);
|
||||
#[cfg(test)]
|
||||
{
|
||||
let store_id = state.store.read().unwrap().id.clone();
|
||||
read_phase_hook::park(read_phase_hook::Point::BeforeNotifyEcRead, &store_id);
|
||||
}
|
||||
let shards = collect_ec_shard_delta_messages(&state.store.read().unwrap());
|
||||
(volumes, shards)
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Go's MaybeAdjustVolumeMax: statvfs on every auto-sized disk, a stat per
|
||||
/// writable volume.
|
||||
async fn maybe_adjust_volume_max(
|
||||
@@ -803,7 +854,8 @@ async fn maybe_adjust_volume_max(
|
||||
.await
|
||||
}
|
||||
|
||||
/// The EC heartbeat and the shard list later deltas are diffed against.
|
||||
/// The EC heartbeat and the shard list later deltas are diffed against,
|
||||
/// under one guard so no mount lands between them unreported.
|
||||
async fn ec_tick_pass(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
@@ -815,9 +867,11 @@ async fn ec_tick_pass(
|
||||
tokio::task::JoinError,
|
||||
> {
|
||||
off_runtime(config, state, |config, state| {
|
||||
let heartbeat = collect_ec_heartbeat(config, state);
|
||||
let shards = collect_ec_shard_delta_messages(&state.store.read().unwrap());
|
||||
(heartbeat, shards)
|
||||
let store = state.store.read().unwrap();
|
||||
(
|
||||
collect_ec_heartbeat(config, &store),
|
||||
collect_ec_shard_delta_messages(&store),
|
||||
)
|
||||
})
|
||||
.await
|
||||
}
|
||||
@@ -932,6 +986,8 @@ mod read_phase_hook {
|
||||
BeforeVolumePass,
|
||||
/// Inside the volume pass, holding the store read lock.
|
||||
VolumePass,
|
||||
/// Between the notify branch's volume snapshot and its EC read.
|
||||
BeforeNotifyEcRead,
|
||||
}
|
||||
|
||||
type Park = (Point, String, Sender<()>, Receiver<()>);
|
||||
@@ -1356,12 +1412,8 @@ fn collect_live_ec_shards(
|
||||
}
|
||||
|
||||
/// Collect EC shard information into a Heartbeat message.
|
||||
fn collect_ec_heartbeat(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
) -> master_pb::Heartbeat {
|
||||
let store = state.store.read().unwrap();
|
||||
let ec_shards = collect_live_ec_shards(&store, true);
|
||||
fn collect_ec_heartbeat(config: &HeartbeatConfig, store: &Store) -> master_pb::Heartbeat {
|
||||
let ec_shards = collect_live_ec_shards(store, true);
|
||||
|
||||
let has_no = ec_shards.is_empty();
|
||||
master_pb::Heartbeat {
|
||||
@@ -2010,7 +2062,7 @@ mod tests {
|
||||
.unwrap();
|
||||
|
||||
let state = test_state_with_store(store);
|
||||
let heartbeat = collect_ec_heartbeat(&test_config(), &state);
|
||||
let heartbeat = collect_ec_heartbeat(&test_config(), &state.store.read().unwrap());
|
||||
|
||||
assert_eq!(heartbeat.ec_shards.len(), 1);
|
||||
assert!(!heartbeat.has_no_ec_shards);
|
||||
@@ -2302,6 +2354,126 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_volume_options_wait_for_the_store_off_the_runtime() {
|
||||
let state = test_state_with_store(Store::new(NeedleMapKind::InMemory));
|
||||
state
|
||||
.store
|
||||
.read()
|
||||
.unwrap()
|
||||
.volume_size_limit
|
||||
.store(1024, Ordering::Relaxed);
|
||||
let hb_resp = master_pb::HeartbeatResponse {
|
||||
volume_size_limit: 2048,
|
||||
preallocate: true,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let (release, released, writer) = hold_store_write_lock(&state);
|
||||
let apply = {
|
||||
let state = state.clone();
|
||||
tokio::spawn(async move { volume_options_pass(&test_config(), &state, &hb_resp).await })
|
||||
};
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let ran_while_held = !released.load(Ordering::SeqCst);
|
||||
drop(release);
|
||||
writer.join().unwrap();
|
||||
let changed = apply.await.unwrap().unwrap();
|
||||
|
||||
assert!(
|
||||
ran_while_held,
|
||||
"applying volume options parked the runtime on the store lock"
|
||||
);
|
||||
assert!(changed);
|
||||
let store = state.store.read().unwrap();
|
||||
assert!(store.get_preallocate());
|
||||
assert_eq!(store.volume_size_limit.load(Ordering::Relaxed), 2048);
|
||||
}
|
||||
|
||||
// The writer arrives between the volume snapshot and the EC read.
|
||||
#[tokio::test]
|
||||
async fn test_state_notification_waits_for_the_store_off_the_runtime() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let dir = temp_dir.path().to_str().unwrap();
|
||||
let mut store = reporting_store(dir, 1);
|
||||
store.id = "notify-off-runtime".to_string();
|
||||
let state = test_state_with_store(store);
|
||||
mount_test_ec_shard(&state, dir, "notify_off_runtime", 42);
|
||||
|
||||
let (entered_tx, entered_rx) = tokio::sync::oneshot::channel();
|
||||
let (resume_tx, resume_rx) = std::sync::mpsc::channel::<()>();
|
||||
read_phase_hook::arm(
|
||||
read_phase_hook::Point::BeforeNotifyEcRead,
|
||||
"notify-off-runtime",
|
||||
entered_tx,
|
||||
resume_rx,
|
||||
);
|
||||
let (held_tx, held_rx) = tokio::sync::oneshot::channel();
|
||||
let (release_tx, release_rx) = std::sync::mpsc::channel::<()>();
|
||||
let released = Arc::new(std::sync::atomic::AtomicBool::new(false));
|
||||
let writer = {
|
||||
let (state, released) = (state.clone(), released.clone());
|
||||
std::thread::spawn(move || {
|
||||
entered_rx.blocking_recv().unwrap();
|
||||
let guard = state.store.write().unwrap();
|
||||
drop(resume_tx);
|
||||
held_tx.send(()).unwrap();
|
||||
let _ = release_rx.recv_timeout(Duration::from_secs(3));
|
||||
released.store(true, Ordering::SeqCst);
|
||||
drop(guard);
|
||||
})
|
||||
};
|
||||
let pass = {
|
||||
let state = state.clone();
|
||||
tokio::spawn(async move { notify_pass(&test_config(), &state).await })
|
||||
};
|
||||
tokio::time::timeout(Duration::from_secs(10), held_rx)
|
||||
.await
|
||||
.expect("the notify pass never reached its EC read")
|
||||
.unwrap();
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let ran_while_held = !released.load(Ordering::SeqCst);
|
||||
drop(release_tx);
|
||||
writer.join().unwrap();
|
||||
let (volumes, shards) = pass.await.unwrap().unwrap();
|
||||
|
||||
assert!(
|
||||
ran_while_held,
|
||||
"the notify pass parked the runtime on the store lock"
|
||||
);
|
||||
assert_eq!(volumes.len(), 1);
|
||||
assert_eq!(shards.len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_deregister_heartbeat_waits_for_the_store_off_the_runtime() {
|
||||
let mut store = Store::new(NeedleMapKind::InMemory);
|
||||
store.id = "deregister-off-runtime".to_string();
|
||||
let state = test_state_with_store(store);
|
||||
let (tx, mut rx) = tokio::sync::mpsc::channel(1);
|
||||
|
||||
let (release, released, writer) = hold_store_write_lock(&state);
|
||||
let deregister = {
|
||||
let state = state.clone();
|
||||
tokio::spawn(
|
||||
async move { send_deregister_heartbeat(&test_config(), &state, &tx).await },
|
||||
)
|
||||
};
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let ran_while_held = !released.load(Ordering::SeqCst);
|
||||
drop(release);
|
||||
writer.join().unwrap();
|
||||
deregister.await.unwrap();
|
||||
let heartbeat = rx.recv().await.unwrap();
|
||||
|
||||
assert!(
|
||||
ran_while_held,
|
||||
"the deregistration heartbeat parked the runtime on the store lock"
|
||||
);
|
||||
assert_eq!(heartbeat.id, "deregister-off-runtime");
|
||||
assert!(heartbeat.has_no_volumes && heartbeat.has_no_ec_shards);
|
||||
}
|
||||
|
||||
// What the read pass decided on can go stale before the write lock is
|
||||
// taken: each action must re-check its target, not act on whatever now
|
||||
// holds the id.
|
||||
@@ -2446,6 +2618,96 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
fn mount_test_ec_shard(state: &VolumeServerState, dir: &str, collection: &str, id: u32) {
|
||||
std::fs::write(format!("{}/{}_{}.ec00", dir, collection, id), b"shard").unwrap();
|
||||
std::fs::write(format!("{}/{}_{}.ecx", dir, collection, id), [0u8; 16]).unwrap();
|
||||
state.store.write().unwrap().locations[0]
|
||||
.mount_ec_shards(VolumeId(id), collection, &[0], "")
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// One pass of the loop's volume tick, then of its notify branch: the
|
||||
/// EC delta the notify branch would send.
|
||||
fn ec_delta_after_volume_tick(
|
||||
state: &Arc<VolumeServerState>,
|
||||
last_ec_shards: &mut HashMap<EcShardDeltaKey, master_pb::VolumeEcShardInformationMessage>,
|
||||
) -> (
|
||||
master_pb::Heartbeat,
|
||||
Vec<master_pb::VolumeEcShardInformationMessage>,
|
||||
Vec<master_pb::VolumeEcShardInformationMessage>,
|
||||
) {
|
||||
let (heartbeat, _) = collect_heartbeat_with_snapshot(&test_config(), state);
|
||||
forget_reported_ec_deletions(last_ec_shards, &heartbeat);
|
||||
let current = collect_ec_shard_delta_messages(&state.store.read().unwrap());
|
||||
let (new_ec_shards, deleted_ec_shards) =
|
||||
diff_ec_shard_delta_messages(last_ec_shards, ¤t);
|
||||
(heartbeat, new_ec_shards, deleted_ec_shards)
|
||||
}
|
||||
|
||||
// A volume heartbeat carries no shard list, so a mount or unmount it
|
||||
// collects past must still go out as the notify branch's delta.
|
||||
#[test]
|
||||
fn test_ec_shard_change_before_a_volume_heartbeat_still_goes_out_as_a_delta() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let dir = temp_dir.path().to_str().unwrap();
|
||||
let state = test_state_with_store(reporting_store(dir, 1));
|
||||
mount_test_ec_shard(&state, dir, "ec_delta_kept", 91);
|
||||
mount_test_ec_shard(&state, dir, "ec_delta_unmounted", 92);
|
||||
let mut last_ec_shards = collect_ec_shard_delta_messages(&state.store.read().unwrap());
|
||||
|
||||
mount_test_ec_shard(&state, dir, "ec_delta_mounted", 93);
|
||||
state
|
||||
.store
|
||||
.write()
|
||||
.unwrap()
|
||||
.unmount_ec_shards(VolumeId(92), &[0]);
|
||||
let (heartbeat, new_ec_shards, deleted_ec_shards) =
|
||||
ec_delta_after_volume_tick(&state, &mut last_ec_shards);
|
||||
|
||||
assert!(heartbeat.ec_shards.is_empty() && !heartbeat.has_no_ec_shards);
|
||||
assert_eq!(
|
||||
new_ec_shards.iter().map(|s| s.id).collect::<Vec<_>>(),
|
||||
vec![93],
|
||||
"a mount collected past by the volume heartbeat was never sent"
|
||||
);
|
||||
assert_eq!(
|
||||
deleted_ec_shards.iter().map(|s| s.id).collect::<Vec<_>>(),
|
||||
vec![92],
|
||||
"an unmount collected past by the volume heartbeat was never sent"
|
||||
);
|
||||
}
|
||||
|
||||
// The volume heartbeat already told the master about the expired EC
|
||||
// volume it destroyed; the next delta must not repeat it.
|
||||
#[test]
|
||||
fn test_ec_volume_expired_by_a_volume_heartbeat_is_not_deleted_again() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let dir = temp_dir.path().to_str().unwrap();
|
||||
let state = test_state_with_store(reporting_store(dir, 0));
|
||||
mount_test_ec_shard(&state, dir, "ec_delta_kept", 94);
|
||||
mount_test_ec_shard(&state, dir, "ec_delta_expired", 95);
|
||||
let mut last_ec_shards = collect_ec_shard_delta_messages(&state.store.read().unwrap());
|
||||
state
|
||||
.store
|
||||
.write()
|
||||
.unwrap()
|
||||
.find_ec_volume_mut(VolumeId(95))
|
||||
.unwrap()
|
||||
.expire_at_sec = 1;
|
||||
|
||||
let (heartbeat, new_ec_shards, deleted_ec_shards) =
|
||||
ec_delta_after_volume_tick(&state, &mut last_ec_shards);
|
||||
|
||||
assert_eq!(heartbeat.deleted_ec_shards.len(), 1);
|
||||
assert_eq!(heartbeat.deleted_ec_shards[0].id, 95);
|
||||
assert!(new_ec_shards.is_empty());
|
||||
assert!(
|
||||
deleted_ec_shards.is_empty(),
|
||||
"an expired EC volume was reported deleted twice"
|
||||
);
|
||||
assert_eq!(last_ec_shards.len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_expired_ec_volume_gone_before_removal_is_not_reported_deleted() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
|
||||
Reference in new issue
Block a user