mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-01 12:16:07 +00:00
volume server: collect EC heartbeats and adjust volume max off the runtime (#11532)
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>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
5ece8dd63c
commit
cc1ec48151
@@ -297,7 +297,7 @@ fn apply_master_volume_options(store: &Store, hb_resp: &master_pb::HeartbeatResp
|
||||
volume_opts_changed = true;
|
||||
}
|
||||
|
||||
volume_opts_changed && store.maybe_adjust_volume_max()
|
||||
volume_opts_changed
|
||||
}
|
||||
|
||||
type EcShardDeltaKey = (u32, String, u32, u32);
|
||||
@@ -410,7 +410,8 @@ async fn do_heartbeat(
|
||||
// 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?;
|
||||
tx.send(collect_ec_heartbeat(config, state)).await?;
|
||||
let initial_ec_hb = off_runtime(config, state, collect_ec_heartbeat).await?;
|
||||
tx.send(initial_ec_hb).await?;
|
||||
|
||||
let stream = tokio_stream::wrappers::ReceiverStream::new(rx);
|
||||
let mut response_stream = client.send_heartbeat(stream).await?.into_inner();
|
||||
@@ -459,11 +460,11 @@ async fn do_heartbeat(
|
||||
)
|
||||
.into());
|
||||
}
|
||||
let changed = {
|
||||
let options_changed = {
|
||||
let s = state.store.read().unwrap();
|
||||
apply_master_volume_options(&s, &hb_resp)
|
||||
};
|
||||
if changed {
|
||||
if options_changed && maybe_adjust_volume_max(config, state).await? {
|
||||
let (adjusted_hb, adjusted_volumes) =
|
||||
off_runtime(config, state, collect_heartbeat_with_snapshot)
|
||||
.await?;
|
||||
@@ -496,10 +497,7 @@ async fn do_heartbeat(
|
||||
}
|
||||
|
||||
_ = volume_tick.tick() => {
|
||||
{
|
||||
let s = state.store.read().unwrap();
|
||||
s.maybe_adjust_volume_max();
|
||||
}
|
||||
maybe_adjust_volume_max(config, state).await?;
|
||||
let (current_hb, current_volumes) =
|
||||
off_runtime(config, state, collect_heartbeat_with_snapshot).await?;
|
||||
last_volumes = volume_identities(¤t_volumes);
|
||||
@@ -513,11 +511,8 @@ async fn do_heartbeat(
|
||||
}
|
||||
|
||||
_ = ec_tick.tick() => {
|
||||
let current_ec_hb = collect_ec_heartbeat(config, state);
|
||||
last_ec_shards = {
|
||||
let store = state.store.read().unwrap();
|
||||
collect_ec_shard_delta_messages(&store)
|
||||
};
|
||||
let (current_ec_hb, current_ec_shards) = ec_tick_pass(config, state).await?;
|
||||
last_ec_shards = current_ec_shards;
|
||||
if tx.send(current_ec_hb).await.is_err() {
|
||||
return Ok(None);
|
||||
}
|
||||
@@ -796,6 +791,37 @@ async fn off_runtime<T: Send + 'static>(
|
||||
tokio::task::spawn_blocking(move || pass(&config, &state)).await
|
||||
}
|
||||
|
||||
/// Go's MaybeAdjustVolumeMax: statvfs on every auto-sized disk, a stat per
|
||||
/// writable volume.
|
||||
async fn maybe_adjust_volume_max(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
) -> Result<bool, tokio::task::JoinError> {
|
||||
off_runtime(config, state, |_, state| {
|
||||
state.store.read().unwrap().maybe_adjust_volume_max()
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// The EC heartbeat and the shard list later deltas are diffed against.
|
||||
async fn ec_tick_pass(
|
||||
config: &HeartbeatConfig,
|
||||
state: &Arc<VolumeServerState>,
|
||||
) -> Result<
|
||||
(
|
||||
master_pb::Heartbeat,
|
||||
HashMap<EcShardDeltaKey, master_pb::VolumeEcShardInformationMessage>,
|
||||
),
|
||||
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)
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Collect volume information into a Heartbeat message.
|
||||
fn collect_heartbeat_with_snapshot(
|
||||
config: &HeartbeatConfig,
|
||||
@@ -2177,6 +2203,105 @@ mod tests {
|
||||
assert_eq!(volumes.len(), 2);
|
||||
}
|
||||
|
||||
/// Holds the store write lock on another thread until `release` is
|
||||
/// dropped, or for 3s so a parked runtime fails the test instead of
|
||||
/// hanging it. The flag turns true just before the lock is let go.
|
||||
fn hold_store_write_lock(
|
||||
state: &Arc<VolumeServerState>,
|
||||
) -> (
|
||||
std::sync::mpsc::Sender<()>,
|
||||
Arc<std::sync::atomic::AtomicBool>,
|
||||
std::thread::JoinHandle<()>,
|
||||
) {
|
||||
let (held_tx, held_rx) = std::sync::mpsc::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 || {
|
||||
let guard = state.store.write().unwrap();
|
||||
held_tx.send(()).unwrap();
|
||||
let _ = release_rx.recv_timeout(Duration::from_secs(3));
|
||||
released.store(true, Ordering::SeqCst);
|
||||
drop(guard);
|
||||
})
|
||||
};
|
||||
held_rx.recv().unwrap();
|
||||
(release_tx, released, writer)
|
||||
}
|
||||
|
||||
// The heartbeat task shares a worker with other tasks; on this
|
||||
// single-threaded runtime, a pass that waits for the store on the worker
|
||||
// stops everything else until the writer lets go.
|
||||
#[tokio::test]
|
||||
async fn test_ec_tick_pass_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, 0);
|
||||
std::fs::write(format!("{}/ec_tick_off_runtime_41.ec00", dir), b"shard").unwrap();
|
||||
std::fs::write(format!("{}/ec_tick_off_runtime_41.ecx", dir), [0u8; 16]).unwrap();
|
||||
store.locations[0]
|
||||
.mount_ec_shards(VolumeId(41), "ec_tick_off_runtime", &[0], "")
|
||||
.unwrap();
|
||||
let state = test_state_with_store(store);
|
||||
|
||||
let (release, released, writer) = hold_store_write_lock(&state);
|
||||
let pass = {
|
||||
let state = state.clone();
|
||||
tokio::spawn(async move { ec_tick_pass(&test_config(), &state).await })
|
||||
};
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let ran_while_held = !released.load(Ordering::SeqCst);
|
||||
drop(release);
|
||||
writer.join().unwrap();
|
||||
let (heartbeat, shards) = pass.await.unwrap().unwrap();
|
||||
|
||||
assert!(
|
||||
ran_while_held,
|
||||
"the EC pass parked the runtime on the store lock"
|
||||
);
|
||||
assert_eq!(heartbeat.ec_shards.len(), 1);
|
||||
assert_eq!(shards.len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_volume_max_adjustment_waits_for_the_store_off_the_runtime() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let mut store = reporting_store(temp_dir.path().to_str().unwrap(), 1);
|
||||
// Auto-sized, so the adjustment stats the disk and the volume.
|
||||
store.locations[0].original_max_volume_count = 0;
|
||||
store.locations[0]
|
||||
.max_volume_count
|
||||
.store(0, Ordering::Relaxed);
|
||||
store
|
||||
.volume_size_limit
|
||||
.store(1024 * 1024, Ordering::Relaxed);
|
||||
let state = test_state_with_store(store);
|
||||
|
||||
let (release, released, writer) = hold_store_write_lock(&state);
|
||||
let adjust = {
|
||||
let state = state.clone();
|
||||
tokio::spawn(async move { maybe_adjust_volume_max(&test_config(), &state).await })
|
||||
};
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let ran_while_held = !released.load(Ordering::SeqCst);
|
||||
drop(release);
|
||||
writer.join().unwrap();
|
||||
let changed = adjust.await.unwrap().unwrap();
|
||||
|
||||
assert!(
|
||||
ran_while_held,
|
||||
"the adjustment parked the runtime on the store lock"
|
||||
);
|
||||
assert!(changed);
|
||||
assert!(
|
||||
state.store.read().unwrap().locations[0]
|
||||
.max_volume_count
|
||||
.load(Ordering::Relaxed)
|
||||
>= 1
|
||||
);
|
||||
}
|
||||
|
||||
// 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.
|
||||
@@ -2532,7 +2657,7 @@ mod tests {
|
||||
|
||||
assert!(store.get_preallocate());
|
||||
assert_eq!(store.volume_size_limit.load(Ordering::Relaxed), 2048);
|
||||
assert!(!changed);
|
||||
assert!(changed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
Reference in New Issue
Block a user