diff --git a/seaweed-volume/src/server/heartbeat.rs b/seaweed-volume/src/server/heartbeat.rs index dbfc81bd3..1460ab991 100644 --- a/seaweed-volume/src/server/heartbeat.rs +++ b/seaweed-volume/src/server/heartbeat.rs @@ -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( 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, +) -> Result { + 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, +) -> Result< + ( + master_pb::Heartbeat, + HashMap, + ), + 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, + ) -> ( + std::sync::mpsc::Sender<()>, + Arc, + 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]