diff --git a/seaweed-volume/src/server/heartbeat.rs b/seaweed-volume/src/server/heartbeat.rs index 1460ab991..9f7610c9f 100644 --- a/seaweed-volume/src/server/heartbeat.rs +++ b/seaweed-volume/src/server/heartbeat.rs @@ -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, + 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::(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 = 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, tx: &tokio::sync::mpsc::Sender, ) { - 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, +) -> 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( config: &HeartbeatConfig, state: &Arc, - pass: fn(&HeartbeatConfig, &Arc) -> T, + pass: impl FnOnce(&HeartbeatConfig, &Arc) -> T + Send + 'static, ) -> Result { 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, + hb_resp: &master_pb::HeartbeatResponse, +) -> Result { + 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, +) -> Result< + ( + Vec, + HashMap, + ), + 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, @@ -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, -) -> 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, + last_ec_shards: &mut HashMap, + ) -> ( + master_pb::Heartbeat, + Vec, + Vec, + ) { + 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![93], + "a mount collected past by the volume heartbeat was never sent" + ); + assert_eq!( + deleted_ec_shards.iter().map(|s| s.id).collect::>(), + 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();