diff --git a/seaweed-volume/src/main.rs b/seaweed-volume/src/main.rs index fb9e97106..96d77a0b0 100644 --- a/seaweed-volume/src/main.rs +++ b/seaweed-volume/src/main.rs @@ -371,6 +371,9 @@ async fn run( .to_string_lossy() .into_owned() }, + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }); // Load persisted state from disk if it exists (matches Go's State.Load on startup) diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 67eedf42c..f363648c3 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -869,6 +869,72 @@ impl VolumeGrpcService { } } +struct EcDecodeClaim<'a>(&'a VolumeServerState, VolumeId); + +impl Drop for EcDecodeClaim<'_> { + fn drop(&mut self) { + self.0 + .ec_decodes_in_flight + .lock() + .unwrap_or_else(|e| e.into_inner()) + .remove(&self.1); + } +} + +/// Keeps `vid` in `ec_decode_tail` while the decode publishes .idx and +/// compacts; local .ecj appenders wait that span out rather than commit a +/// delete the rebuilt index would miss. Dropping — including on panic — +/// lifts the marker and wakes the waiters. +struct EcDecodeTailGuard<'a>(&'a VolumeServerState, VolumeId); + +impl Drop for EcDecodeTailGuard<'_> { + fn drop(&mut self) { + self.0 + .ec_decode_tail + .lock() + .unwrap_or_else(|e| e.into_inner()) + .remove(&self.1); + self.0.ec_decode_tail_notify.notify_waiters(); + } +} + +/// Whether `vid`'s .ecj appends must wait for the decode publishing tail. +/// Only meaningful read under the store write lock: a decode can claim the +/// tail while a caller waits for the decoder's read lock, so membership +/// tested before acquiring the write lock is stale by commit time. +pub(crate) fn ec_decode_tail_contains(state: &VolumeServerState, vid: VolumeId) -> bool { + state + .ec_decode_tail + .lock() + .unwrap_or_else(|e| e.into_inner()) + .contains(&vid) +} + +/// Blocks a local .ecj append until `vid` leaves `ec_decode_tail`. Mirrors +/// Go's `EcVolume.ecjFileAccessLock`: decode holds it across journal +/// catch-up, .idx publication and compaction, so no committed delete falls +/// between the last catch_up and the .cpd/.cpx swap. Async wait — the +/// caller holds no lock while sleeping, so decode can never be deadlocked +/// by the append it is delaying. +pub(crate) async fn wait_ec_decode_tail(state: &Arc, vid: VolumeId) { + loop { + let notified = state.ec_decode_tail_notify.notified(); + tokio::pin!(notified); + // Register before testing the set so a tail that ends between the + // check and the await still wakes us. + notified.as_mut().enable(); + let in_tail = state + .ec_decode_tail + .lock() + .unwrap_or_else(|e| e.into_inner()) + .contains(&vid); + if !in_tail { + return; + } + notified.await; + } +} + #[tonic::async_trait] impl VolumeServer for VolumeGrpcService { // ---- Core volume operations ---- @@ -1530,10 +1596,31 @@ impl VolumeServer for VolumeGrpcService { let req = request.into_inner(); let vid = VolumeId(req.volume_id); + // A decode in flight may be mid-compaction on this volume's + // .dat/.idx; mounting across the .cpd/.cpx swap could load a mixed + // pair. Claim the vid for the mount rather than hold the set lock + // through it — the same exclusion for this vid, while unrelated + // mounts and decode claims stay unblocked. + // Retryable — the caller mounts after the decode RPC returns. + if !self + .state + .ec_decodes_in_flight + .lock() + .unwrap_or_else(|e| e.into_inner()) + .insert(vid) + { + return Err(Status::unavailable(format!( + "volume {} is being decoded", + req.volume_id + ))); + } + let _claim = EcDecodeClaim(&self.state, vid); + let mut store = self.state.store.write().unwrap(); store .mount_volume_by_id(vid, req.collection.as_deref()) .map_err(|e| Status::internal(e.to_string()))?; + drop(store); self.state.volume_state_notify.notify_one(); Ok(Response::new(volume_server_pb::VolumeMountResponse {})) @@ -3750,9 +3837,23 @@ impl VolumeServer for VolumeGrpcService { let vid = VolumeId(req.volume_id); let needle_id = NeedleId(req.file_key); + // A delete committed between the last journal catch_up and the + // compaction swap would be absent from the rebuilt .idx, so the + // tail membership must be checked under the store write lock: a + // decode can claim the tail while this delete waits for the + // decoder's read lock, and a check taken before acquiring it would + // be stale by the time the append commits. + let mut store = loop { + let store = self.state.store.write().unwrap(); + if !ec_decode_tail_contains(&self.state, vid) { + break store; + } + drop(store); + wait_ec_decode_tail(&self.state, vid).await; + }; + // Go's handler locates the needle first: absent fails the RPC so the // caller moves to the next holder; an existing tombstone is a no-op. - let mut store = self.state.store.write().unwrap(); if let Some(ec_vol) = store.find_ec_volume_mut(vid) { match ec_vol.find_needle_from_ecx(needle_id) { Ok(Some((_, size))) if size.is_deleted() => { @@ -3794,6 +3895,24 @@ impl VolumeServer for VolumeGrpcService { // marker, then mount. This server holds no EC shards for the vid, so there // is no in-place decode to run. if req.from_staged { + // An in-place decode in flight for this vid is rebuilding the + // same .dat/.idx the staged files would overwrite. Claim the + // vid for the adoption rather than hold the set lock through + // it — the same exclusion for this vid, while unrelated mounts + // and decode claims stay unblocked. + if !self + .state + .ec_decodes_in_flight + .lock() + .unwrap_or_else(|e| e.into_inner()) + .insert(vid) + { + return Err(Status::unavailable(format!( + "ec volume {} is already being decoded", + req.volume_id + ))); + } + let _claim = EcDecodeClaim(&self.state, vid); let want = DiskType::from_string(&req.disk_type); let base = { let store = self.state.store.read().unwrap(); @@ -3873,129 +3992,93 @@ impl VolumeServer for VolumeGrpcService { )); } - let store = self.state.store.read().unwrap(); - // Aggregate per-shard data dirs across all locations so the - // shard-presence check + decoder both see the union for - // cross-disk reconciled volumes (#9252). Mirrors Go's - // CollectEcShards. - let max_shard_count = crate::storage::erasure_coding::ec_shard::MAX_SHARD_COUNT; - let (ec_vol, shard_dirs) = store - .collect_ec_shard_dirs(vid, max_shard_count) - .ok_or_else(|| Status::not_found(format!("ec volume {} not found", req.volume_id)))?; + let job = { + let store = self.state.store.read().unwrap(); + // Aggregate per-shard data dirs across all locations so the + // shard-presence check + decoder both see the union for + // cross-disk reconciled volumes (#9252). Mirrors Go's + // CollectEcShards. + let max_shard_count = crate::storage::erasure_coding::ec_shard::MAX_SHARD_COUNT; + let (ec_vol, shard_dirs) = store + .collect_ec_shard_dirs(vid, max_shard_count) + .ok_or_else(|| { + Status::not_found(format!("ec volume {} not found", req.volume_id)) + })?; - if ec_vol.collection != req.collection { - return Err(Status::internal(format!( - "existing collection:{} unexpected input: {}", - ec_vol.collection, req.collection - ))); - } - - // Use EC context data shard count from the volume - let data_shards = ec_vol.data_shards as usize; - - // Validate data shard count range (matches Go's VolumeEcShardsToVolume) - if data_shards == 0 || data_shards > max_shard_count { - return Err(Status::invalid_argument(format!( - "invalid data shard count {} for volume {} (must be 1..{})", - data_shards, req.volume_id, max_shard_count - ))); - } - - // Check that all data shards are present somewhere on this server. - for (shard_id, dir) in shard_dirs[..data_shards].iter().enumerate() { - if dir.is_none() { + if ec_vol.collection != req.collection { return Err(Status::internal(format!( - "ec volume {} missing shard {}", - req.volume_id, shard_id + "existing collection:{} unexpected input: {}", + ec_vol.collection, req.collection ))); } - } - // Reconstruct the volume from EC shards. Use the EcVolume's - // own dir for the produced .dat (matches the volume's home - // disk) and its `ecx_actual_dir` for the .ecx lookup, while - // reading each shard from its real on-disk location. - let dat_dir = ec_vol.dir.clone(); - let ecx_dir = ec_vol.ecx_actual_dir().to_string(); - let idx_dir = ec_vol.dir_idx.clone(); - let collection = ec_vol.collection.clone(); - let vif_dat_file_size = ec_vol.dat_file_size; - let (large_block_size, small_block_size) = - (ec_vol.large_block_size(), ec_vol.small_block_size()); - // shard_dirs[i] is guaranteed Some for i in 0..data_shards by - // the check above; collect concrete dirs for the decoder. - let per_shard_dirs: Vec = shard_dirs[..data_shards] - .iter() - .map(|d| d.clone().unwrap()) - .collect(); - drop(store); + // Use EC context data shard count from the volume + let data_shards = ec_vol.data_shards as usize; - // Deletions journaled beside the .ecx, or collected by - // VolumeEcShardsCopy into the idx dir, count as deleted throughout. - let deleted = crate::storage::erasure_coding::ec_decoder::read_ecj_deletions( - &[&ecx_dir, &idx_dir], - &collection, - vid, - ) - .map_err(|e| Status::internal(format!("read ecj: {}", e)))?; - let has_live = crate::storage::erasure_coding::ec_decoder::has_live_needles( - &ecx_dir, - &collection, - vid, - &deleted, - ) - .map_err(|e| Status::internal(format!("HasLiveNeedles: {}", e)))?; - if !has_live { - return Err(Status::failed_precondition(format!( - "ec volume {} has no live entries", + // Validate data shard count range (matches Go's VolumeEcShardsToVolume) + if data_shards == 0 || data_shards > max_shard_count { + return Err(Status::invalid_argument(format!( + "invalid data shard count {} for volume {} (must be 1..{})", + data_shards, req.volume_id, max_shard_count + ))); + } + + // Check that all data shards are present somewhere on this server. + for (shard_id, dir) in shard_dirs[..data_shards].iter().enumerate() { + if dir.is_none() { + return Err(Status::internal(format!( + "ec volume {} missing shard {}", + req.volume_id, shard_id + ))); + } + } + + // Reconstruct the volume from EC shards. Use the EcVolume's + // own dir for the produced .dat (matches the volume's home + // disk) and its `ecx_actual_dir` for the .ecx lookup, while + // reading each shard from its real on-disk location. + // shard_dirs[i] is guaranteed Some for i in 0..data_shards by + // the check above; collect concrete dirs for the decoder. + EcDecodeJob { + vid, + dat_dir: ec_vol.dir.clone(), + ecx_dir: ec_vol.ecx_actual_dir().to_string(), + idx_dir: ec_vol.dir_idx.clone(), + collection: ec_vol.collection.clone(), + vif_dat_file_size: ec_vol.dat_file_size, + large_block_size: ec_vol.large_block_size() as usize, + small_block_size: ec_vol.small_block_size() as usize, + shard_dirs: shard_dirs[..data_shards] + .iter() + .map(|d| d.clone().unwrap()) + .collect(), + needle_map_kind: store.needle_map_kind, + } + }; + + // A dropped request leaves the blocking job running; keep the vid + // claimed until it finishes so a retry cannot race the in-flight + // decode on the same volume files. + if !self + .state + .ec_decodes_in_flight + .lock() + .unwrap() + .insert(vid) + { + return Err(Status::unavailable(format!( + "ec volume {} is already being decoded", req.volume_id ))); } - // Calculate .dat file size from .ecx entries (.ec00 lives on - // its own disk, .ecx on the index disk). - let dat_file_size = - crate::storage::erasure_coding::ec_decoder::find_dat_file_size_with_dirs( - &per_shard_dirs[0], - &ecx_dir, - &collection, - vid, - &deleted, - ) - .map_err(|e| Status::internal(format!("FindDatFileSize: {}", e)))?; - - // The shard block layout was fixed by the .dat size at encode time - // (recorded in .vif); deletions can shrink the live extent below a - // large-block row boundary, so the layout must not be derived from - // dat_file_size. The decoder infers the layout from the shard size - // when .vif does not record it. - // Write .dat file using block-interleaved reading from shards. - crate::storage::erasure_coding::ec_decoder::write_dat_file_from_shards( - &crate::storage::erasure_coding::ec_decoder::DatRebuild { - dat_dir: &dat_dir, - collection: &collection, - volume_id: vid, - dat_file_size, - encoded_dat_file_size: vif_dat_file_size, - data_shards, - shard_dirs: Some(&per_shard_dirs), - large_block_size: large_block_size as usize, - small_block_size: small_block_size as usize, - }, - ) - .map_err(|e| Status::internal(format!("WriteDatFile: {}", e)))?; - - // Write .idx from the .ecx wherever it lives, beside the .dat where - // the mount looks first (Go moves it there after the rebuild). - crate::storage::erasure_coding::ec_decoder::write_idx_file_from_ec_index_with_dirs( - &ecx_dir, - &dat_dir, - &collection, - vid, - &deleted, - dat_file_size, - ) - .map_err(|e| Status::internal(format!("WriteIdxFileFromEcIndex: {}", e)))?; + let state = self.state.clone(); + tokio::task::spawn_blocking(move || { + let _claim = EcDecodeClaim(&state, vid); + job.run(&state) + }) + .await + .map_err(|e| Status::internal(format!("decode ec volume {}: {}", vid, e)))??; // Go does NOT unmount EC shards or mount the volume here. // The caller (ec.balance / ec.decode) handles mount/unmount separately. @@ -6618,6 +6701,197 @@ fn get_disk_usage(path: &str) -> (u64, u64) { } } +/// What `VolumeEcShardsToVolume` needs to decode an EC volume, snapshotted +/// under the store lock so the decode itself runs without it. +struct EcDecodeJob { + vid: VolumeId, + dat_dir: String, + ecx_dir: String, + idx_dir: String, + collection: String, + vif_dat_file_size: i64, + large_block_size: usize, + small_block_size: usize, + /// Directory of each data shard. + shard_dirs: Vec, + needle_map_kind: crate::storage::needle_map::NeedleMapKind, +} + +impl EcDecodeJob { + fn run(self, state: &VolumeServerState) -> Result<(), Status> { + use crate::storage::erasure_coding::{ec_bitrot, ec_decoder}; + let EcDecodeJob { + vid, + dat_dir, + ecx_dir, + idx_dir, + collection, + vif_dat_file_size, + large_block_size, + small_block_size, + shard_dirs, + needle_map_kind, + } = self; + + // Deletions journaled beside the .ecx, or collected by + // VolumeEcShardsCopy into the idx dir, count as deleted throughout. + // The first pass runs unlocked — a slow journal must not stall every + // store writer — then a second pass runs with appends quiesced by the + // read lock. The rollback epoch proves whether it needs to be a full + // rescan: an unchanged epoch means no append was rolled back inside + // the window, so every id folded in was committed and catch_up only + // needs to pick up records appended meanwhile. + let epoch_before = crate::storage::erasure_coding::ec_volume::ECJ_ROLLBACK_EPOCH + .load(std::sync::atomic::Ordering::Acquire); + let mut deleted = ec_decoder::EcjDeletions::read(&[&ecx_dir, &idx_dir], &collection, vid) + .map_err(|e| Status::internal(format!("read ecj: {}", e)))?; + { + let _guard = state.store.read().unwrap(); + if epoch_before + == crate::storage::erasure_coding::ec_volume::ECJ_ROLLBACK_EPOCH + .load(std::sync::atomic::Ordering::Acquire) + { + deleted + .catch_up() + .map_err(|e| Status::internal(format!("read ecj: {}", e)))?; + } else { + // A rollback ran inside the window: the unlocked pass may + // have folded in bytes that were truncated away, or missed a + // record re-appended to the freed offset — re-read it all. + deleted + .rescan() + .map_err(|e| Status::internal(format!("read ecj: {}", e)))?; + } + } + let has_live = ec_decoder::has_live_needles(&ecx_dir, &collection, vid, &deleted.ids) + .map_err(|e| Status::internal(format!("HasLiveNeedles: {}", e)))?; + if !has_live { + return Err(Status::failed_precondition(format!( + "ec volume {} has no live entries", + vid + ))); + } + + // Calculate .dat file size from .ecx entries (.ec00 lives on + // its own disk, .ecx on the index disk). + let dat_file_size = ec_decoder::find_dat_file_size_with_dirs( + &shard_dirs[0], + &ecx_dir, + &collection, + vid, + &deleted.ids, + ) + .map_err(|e| Status::internal(format!("FindDatFileSize: {}", e)))?; + + // The shard block layout was fixed by the .dat size at encode time + // (recorded in .vif); deletions can shrink the live extent below a + // large-block row boundary, so the layout must not be derived from + // dat_file_size. The decoder infers the layout from the shard size + // when .vif does not record it. + // Write .dat file using block-interleaved reading from shards. + ec_decoder::write_dat_file_from_shards(&ec_decoder::DatRebuild { + dat_dir: &dat_dir, + collection: &collection, + volume_id: vid, + dat_file_size, + encoded_dat_file_size: vif_dat_file_size, + data_shards: shard_dirs.len(), + shard_dirs: Some(&shard_dirs), + large_block_size, + small_block_size, + }) + .map_err(|e| Status::internal(format!("WriteDatFile: {}", e)))?; + ec_decoder::verify_decoded_dat_file(&dat_dir, &collection, vid, dat_file_size) + .map_err(|e| Status::internal(format!("VerifyDecodedDatFile: {}", e)))?; + + // Publishing phase: local .ecj appends for this vid wait on + // ec_decode_tail_notify until the .idx is written and compaction + // done, so a delete committed mid-tail can never slip between + // catch_up and the .cpd/.cpx swap. Unlike a store lock held across + // all of it, this stalls writers for this volume only — an append + // that would commit now instead commits right after the tail, on a + // .ecj that survives like any post-decode journal record (Go holds + // EcVolume.ecjFileAccessLock over the same span). + state + .ec_decode_tail + .lock() + .unwrap_or_else(|e| e.into_inner()) + .insert(vid); + let _tail = EcDecodeTailGuard(state, vid); + + // Deletes journaled while the .dat was written. Journal appends hold + // the store write lock through their sync-or-truncate, so a read lock + // held from catch_up through the .idx publish guarantees every record + // read is committed — a rolled-back delete cannot leave a tombstone + // in the index — and no new journal record can be missed by the + // rebuilt .idx. + { + let _guard = state.store.read().unwrap(); + deleted + .catch_up() + .map_err(|e| Status::internal(format!("read ecj: {}", e)))?; + + // Write .idx from the .ecx wherever it lives, beside the .dat + // where the mount looks first (Go moves it there after the + // rebuild). + ec_decoder::write_idx_file_from_ec_index_with_dirs( + &ecx_dir, + &dat_dir, + &collection, + vid, + &deleted.ids, + dat_file_size, + ) + .map_err(|e| Status::internal(format!("WriteIdxFileFromEcIndex: {}", e)))?; + } + + // The EC generation is gone; a stale .ecsum must not pass for the + // protection of a later re-encode. + let mut sidecar_dirs = vec![&dat_dir]; + if ecx_dir != dat_dir { + sidecar_dirs.push(&ecx_dir); + } + for dir in sidecar_dirs { + let base = crate::storage::volume::volume_file_name(dir, &collection, vid); + if let Err(e) = ec_bitrot::remove_bitrot_sidecars(&base) { + tracing::warn!(volume_id = vid.0, error = %e, "remove bitrot sidecars of {base}"); + } + } + + // Drop the deleted needles — without the store lock, which must not + // be held through this rewrite: on a slow or large volume it would + // stall every store writer (journal appends, mounts) for unrelated + // volumes. This volume's own appends still wait on the tail set, and + // VolumeMount is held off by the decode claim, so the .cpd/.cpx swap + // cannot be raced. A failure is only logged, as in Go: the + // uncompacted .dat/.idx already make a complete volume. + if let Err(e) = crate::storage::store::Store::compact_volume_files( + &dat_dir, + &idx_dir, + &collection, + vid, + needle_map_kind, + ) { + // Benign only when the swap never started or was settled: a + // surviving .cpc marker means the commit was decided but the + // renames could not be reconciled, so .dat/.idx may be a + // mismatched pair — fail the decode and let the caller keep + // the shards rather than mount a corrupt volume. + let marker = format!( + "{}.cpc", + crate::storage::volume::volume_file_name(&dat_dir, &collection, vid) + ); + if std::path::Path::new(&marker).exists() { + return Err(Status::internal(format!( + "compact decoded volume {vid}: {e}" + ))); + } + tracing::error!(volume_id = vid.0, error = %e, "compact decoded volume"); + } + Ok(()) + } +} + /// Bytes compacted between two `VacuumVolumeCompact` progress reports. const COMPACT_REPORT_INTERVAL: i64 = 128 * 1024 * 1024; @@ -7093,6 +7367,9 @@ mod tests { security_file: String::new(), cli_white_list: vec![], state_file_path: String::new(), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }); ( @@ -7204,6 +7481,9 @@ mod tests { security_file: String::new(), cli_white_list: vec![], state_file_path: String::new(), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }); (VolumeGrpcService { state }, tmp) @@ -9368,6 +9648,9 @@ mod tests { security_file: String::new(), cli_white_list: vec![], state_file_path: String::new(), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }); (VolumeGrpcService { state }, tmp) @@ -11155,6 +11438,9 @@ mod tests { security_file: String::new(), cli_white_list: vec![], state_file_path: String::new(), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }); VolumeGrpcService { state } @@ -12080,9 +12366,22 @@ mod tests { assert_eq!(n.data, body); } + /// The decode's closing compaction fails when its .cpd cannot be created; + /// undo with [`unblock_decode_compaction`] before mounting. + fn block_decode_compaction(data: &str) { + std::fs::create_dir(format!("{data}/1.cpd")).unwrap(); + } + + fn unblock_decode_compaction(data: &str) { + std::fs::remove_dir(format!("{data}/1.cpd")).unwrap(); + assert!(!std::path::Path::new(&format!("{data}/1.cpx")).exists()); + } + /// ec.decode used to read the .ecx/.ecj from the data dir when writing the /// .idx, failing with NotFound after the .dat was already published, and /// ignored .ecj deletions when sizing the .dat (Go folds them in first). + /// Compaction is made to fail: the decode still succeeds, as in Go, so the + /// uncompacted .idx must stand on its own. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_volume_ec_shards_to_volume_reads_ecx_from_split_idx_dir() { // The lowest key (the first .idx row) and the .dat tail, each journaled @@ -12091,11 +12390,13 @@ mod tests { let (service, _tmp, data, idx, size_before_needle_3) = make_split_idx_ec_decode_service(&journal, false).await; let ecx = std::fs::read(format!("{idx}/1.ecx")).unwrap(); + block_decode_compaction(&data); service .volume_ec_shards_to_volume(ec_shards_to_volume_request()) .await .unwrap(); + unblock_decode_compaction(&data); // The journaled tail needle is not decoded. assert_eq!( @@ -12142,7 +12443,8 @@ mod tests { } /// A tail needle tombstoned in the .ecx itself (Go's RebuildEcxFile) is cut - /// from the .dat the same way, so its row must not survive either. + /// from the .dat the same way, so its row must not survive either, even + /// when compaction fails. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_volume_ec_shards_to_volume_drops_sealed_tail_tombstone() { let (service, _tmp, data, idx, size_before_needle_3) = @@ -12152,11 +12454,13 @@ mod tests { let size_at = 2 * NEEDLE_MAP_ENTRY_SIZE + NEEDLE_ID_SIZE + OFFSET_SIZE; TOMBSTONE_FILE_SIZE.to_bytes(&mut ecx[size_at..size_at + SIZE_SIZE]); std::fs::write(&ecx_path, &ecx).unwrap(); + block_decode_compaction(&data); service .volume_ec_shards_to_volume(ec_shards_to_volume_request()) .await .unwrap(); + unblock_decode_compaction(&data); assert_eq!( std::fs::metadata(format!("{data}/1.dat")).unwrap().len(), @@ -12169,6 +12473,47 @@ mod tests { assert_decoded_volume(&service, (2, 0), &[1, 2], &[3]); } + /// Go compacts the decoded volume (CompactVolumeFiles), so needles deleted + /// through the .ecj leave the .dat and .idx instead of waiting for a vacuum. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_volume_ec_shards_to_volume_compacts_the_decoded_volume() { + let (service, _tmp, data, idx, _) = make_split_idx_ec_decode_service(&[1], false).await; + let vifs = || [&data, &idx].map(|dir| std::fs::read(format!("{dir}/1.vif")).ok()); + let vifs_before = vifs(); + assert!(vifs_before.iter().any(Option::is_some)); + + service + .volume_ec_shards_to_volume(ec_shards_to_volume_request()) + .await + .unwrap(); + + let dat = std::fs::read(format!("{data}/1.dat")).unwrap(); + let deleted_body = b"split idx decode needle 1"; + assert!( + !dat.windows(deleted_body.len()).any(|w| w == deleted_body), + "the deleted needle must be compacted out of the .dat" + ); + assert_eq!(dat[4..6], 1u16.to_be_bytes(), "compaction revision"); + let idx = std::fs::read(format!("{data}/1.idx")).unwrap(); + let idx_rows: Vec<(NeedleId, Size)> = idx + .as_chunks::() + .0 + .iter() + .map(|row| { + let (key, _, size) = idx_entry_from_bytes(row); + (key, size) + }) + .collect(); + assert_eq!( + idx_rows.iter().map(|(key, _)| key.0).collect::>(), + [2, 3] + ); + assert!(idx_rows.iter().all(|(_, size)| !size.is_deleted())); + assert_eq!(vifs(), vifs_before, "the EC .vif is left alone"); + + assert_decoded_volume(&service, (2, 0), &[2, 3], &[1]); + } + /// Deletions only in the .ecj count toward "no live entries", as after Go's /// RebuildEcxFile, so the caller purges the shards instead of decoding. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] @@ -12183,6 +12528,303 @@ mod tests { assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}"); assert!(err.message().contains("has no live entries"), "{err:?}"); assert!(!std::path::Path::new(&format!("{data}/1.dat")).exists()); + assert!( + std::path::Path::new(&format!("{data}/1.ecsum")).exists(), + "a failed decode keeps the shards' bitrot sidecar" + ); + } + + /// A decode drops the bitrot sidecars beside the .dat and the .ecx, as Go + /// does, so a stale .ecsum cannot vouch for a later re-encode. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_volume_ec_shards_to_volume_removes_bitrot_sidecars() { + let (service, _tmp, data, idx, _) = make_split_idx_ec_decode_service(&[], false).await; + assert!( + std::path::Path::new(&format!("{data}/1.ecsum")).exists(), + "precondition: the encode wrote a generation-0 sidecar" + ); + let idx_sidecar = format!("{idx}/1.ecsum.v3"); + let other_volume = format!("{data}/2.ecsum"); + std::fs::write(&idx_sidecar, b"x").unwrap(); + std::fs::write(&other_volume, b"x").unwrap(); + + service + .volume_ec_shards_to_volume(ec_shards_to_volume_request()) + .await + .unwrap(); + + assert!(!std::path::Path::new(&format!("{data}/1.ecsum")).exists()); + assert!(!std::path::Path::new(&idx_sidecar).exists()); + assert!(std::path::Path::new(&other_volume).exists()); + } + + /// Only the .dat and .ecx dirs are swept: with the .ecx beside the shards, + /// the idx dir is not one of them. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_volume_ec_shards_to_volume_sweeps_sidecars_beside_dat_and_ecx_only() { + let (service, _tmp, data, idx, _) = make_split_idx_ec_decode_service(&[], true).await; + let idx_sidecar = format!("{idx}/1.ecsum"); + std::fs::write(&idx_sidecar, b"x").unwrap(); + + service + .volume_ec_shards_to_volume(ec_shards_to_volume_request()) + .await + .unwrap(); + + assert!(!std::path::Path::new(&format!("{data}/1.ecsum")).exists()); + assert!(std::path::Path::new(&idx_sidecar).exists()); + } + + /// Replaces the idx-dir .ecj with a FIFO: the decode's open of it blocks + /// until [`release_ecj_fifo`] opens the write end. + #[cfg(unix)] + fn make_ecj_fifo(idx: &str) -> String { + let fifo = format!("{idx}/1.ecj"); + std::fs::remove_file(&fifo).unwrap(); + let path = std::ffi::CString::new(fifo.clone()).unwrap(); + // SAFETY: `path` is a valid NUL-terminated string for the call. + assert_eq!(unsafe { libc::mkfifo(path.as_ptr(), 0o644) }, 0); + fifo + } + + /// Opens the write end once the decode has the read end open, releasing + /// it. Keep the returned handle so a later open does not block again. + #[cfg(unix)] + fn release_ecj_fifo(fifo: &str) -> std::fs::File { + use std::os::unix::fs::OpenOptionsExt; + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + loop { + match std::fs::OpenOptions::new() + .write(true) + .custom_flags(libc::O_NONBLOCK) + .open(fifo) + { + Ok(writer) => return writer, + // No reader yet. + Err(e) + if e.raw_os_error() == Some(libc::ENXIO) + && std::time::Instant::now() < deadline => + { + std::thread::sleep(std::time::Duration::from_millis(5)); + } + Err(e) => panic!("the decode never opened the .ecj: {e}"), + } + } + } + + /// The decode's file I/O used to run on the async runtime: parked on a + /// slow disk, it stalled every other task on that worker. + #[cfg(unix)] + #[tokio::test(flavor = "current_thread")] + async fn test_volume_ec_shards_to_volume_decodes_off_the_runtime() { + let (service, _tmp, data, idx, _) = make_split_idx_ec_decode_service(&[], true).await; + let fifo = make_ecj_fifo(&idx); + + let probed = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let releaser = { + let probed = probed.clone(); + let state = service.state.clone(); + std::thread::spawn(move || { + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while !probed.load(Ordering::SeqCst) && std::time::Instant::now() < deadline { + std::thread::sleep(std::time::Duration::from_millis(5)); + } + let probed_while_parked = probed.load(Ordering::SeqCst); + assert!( + state.store.try_write().is_ok(), + "the decode must not hold the store lock" + ); + (probed_while_parked, release_ecj_fifo(&fifo)) + }) + }; + + let (result, ()) = tokio::join!( + service.volume_ec_shards_to_volume(ec_shards_to_volume_request()), + async { probed.store(true, Ordering::SeqCst) } + ); + let (probed_while_parked, _writer) = releaser.join().unwrap(); + result.unwrap(); + assert!( + probed_while_parked, + "another task on the runtime must run while the decode is parked" + ); + assert!(std::path::Path::new(&format!("{data}/1.dat")).exists()); + } + + /// A delete journaled while the .dat is written still reaches the .idx: + /// the journals are read again just before it. + #[cfg(unix)] + #[tokio::test(flavor = "current_thread")] + async fn test_volume_ec_shards_to_volume_keeps_a_delete_journaled_mid_decode() { + let (service, _tmp, _data, idx, _) = make_split_idx_ec_decode_service(&[], true).await; + let fifo = make_ecj_fifo(&idx); + + let probed = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let deleter = { + let probed = probed.clone(); + let state = service.state.clone(); + std::thread::spawn(move || { + // The handler has taken its snapshot and handed off the decode. + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while !probed.load(Ordering::SeqCst) && std::time::Instant::now() < deadline { + std::thread::sleep(std::time::Duration::from_millis(5)); + } + // Held until the delete is journaled, so the decode cannot + // catch up before it; the release proves the decode is parked + // on the idx-dir .ecj, past its first read of the data-dir one. + let mut store = state.store.write().unwrap(); + let writer = release_ecj_fifo(&fifo); + store + .find_ec_volume_mut(VolumeId(1)) + .unwrap() + .journal_delete(NeedleId(2)) + .unwrap(); + writer + }) + }; + + let (result, ()) = tokio::join!( + service.volume_ec_shards_to_volume(ec_shards_to_volume_request()), + async { probed.store(true, Ordering::SeqCst) } + ); + let _writer = deleter.join().unwrap(); + result.unwrap(); + assert_decoded_volume(&service, (2, 0), &[1, 3], &[2]); + } + + /// A second decode while one is in flight is refused: the blocking job + /// outlives a dropped request and would race a retry on the same files. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_volume_ec_shards_to_volume_rejects_overlapping_decode() { + let (service, _tmp, _data, _idx, _) = make_split_idx_ec_decode_service(&[], false).await; + service + .state + .ec_decodes_in_flight + .lock() + .unwrap() + .insert(VolumeId(1)); + + let err = service + .volume_ec_shards_to_volume(ec_shards_to_volume_request()) + .await + .unwrap_err(); + assert_eq!(err.code(), tonic::Code::Unavailable); + } + + /// The blob-delete tail check must run under the store write lock: a + /// delete that passed an unlocked check could still be parked behind + /// the decoder's read lock when the decode claims the tail, then commit + /// its journal append after the rebuilt .idx — an acknowledged delete + /// the mount would never see. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_volume_ec_blob_delete_rechecks_tail_under_write_lock() { + let (service, _tmp, _data, idx, _) = make_split_idx_ec_decode_service(&[], false).await; + let state = service.state.clone(); + + // Park the delete on the store write lock, claimed in a thread so a + // std guard is never held across .await; the tail is claimed while + // it waits — the interleaving an unlocked check missed. + let claimed = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let held = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let releaser = { + let claimed = claimed.clone(); + let held = held.clone(); + let state = state.clone(); + std::thread::spawn(move || { + let guard = state.store.write().unwrap(); + held.store(true, Ordering::SeqCst); + while !claimed.load(Ordering::SeqCst) { + std::thread::sleep(std::time::Duration::from_millis(5)); + } + drop(guard); + }) + }; + // Spawn only once the write lock is held: on a slow runner the delete + // could otherwise acquire it first and commit before the tail exists. + while !held.load(Ordering::SeqCst) { + tokio::time::sleep(std::time::Duration::from_millis(5)).await; + } + let delete = tokio::spawn(async move { + service + .volume_ec_blob_delete(Request::new(volume_server_pb::VolumeEcBlobDeleteRequest { + volume_id: 1, + file_key: 1, + collection: String::new(), + version: 0, + })) + .await + }); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + state + .ec_decode_tail + .lock() + .unwrap() + .insert(VolumeId(1)); + claimed.store(true, Ordering::SeqCst); + releaser.join().unwrap(); + + // Under the lock the delete sees the tail and parks on the notify. + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + assert!( + !delete.is_finished(), + "a delete must not commit while its volume is in the decode tail" + ); + + state + .ec_decode_tail + .lock() + .unwrap() + .remove(&VolumeId(1)); + state.ec_decode_tail_notify.notify_waiters(); + delete.await.unwrap().unwrap(); + + let ecj = std::fs::read(format!("{idx}/1.ecj")).unwrap(); + assert_eq!( + ecj.len(), + NEEDLE_ID_SIZE, + "the delayed delete must land on the surviving journal" + ); + } + + /// A normal mount must exclude an in-flight decode for the same volume — + /// it could observe a half-swapped .dat/.idx — without holding the set + /// lock across the disk work, and the claim must not outlive the call. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_volume_mount_claims_the_vid_against_decodes() { + let (service, _tmp) = make_local_service_with_volume("mount_claim", None); + let mount_req = |volume_id| { + Request::new(volume_server_pb::VolumeMountRequest { + volume_id, + collection: Some(String::new()), + }) + }; + + service + .state + .ec_decodes_in_flight + .lock() + .unwrap() + .insert(VolumeId(1)); + let err = service.volume_mount(mount_req(1)).await.unwrap_err(); + assert_eq!(err.code(), tonic::Code::Unavailable); + service + .state + .ec_decodes_in_flight + .lock() + .unwrap() + .remove(&VolumeId(1)); + + // A mount that fails on disk still releases the claim it took. + let _ = service.volume_mount(mount_req(2)).await; + assert!( + service + .state + .ec_decodes_in_flight + .lock() + .unwrap() + .is_empty(), + "the mount claim must not outlive the call" + ); } // Among storage errors, the filer requeues a delete only on "is read only" diff --git a/seaweed-volume/src/server/handlers.rs b/seaweed-volume/src/server/handlers.rs index 38ef73332..4a8a223b9 100644 --- a/seaweed-volume/src/server/handlers.rs +++ b/seaweed-volume/src/server/handlers.rs @@ -4790,6 +4790,9 @@ mod tests { security_file: String::new(), cli_white_list: vec![], state_file_path: String::new(), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }) } diff --git a/seaweed-volume/src/server/heartbeat.rs b/seaweed-volume/src/server/heartbeat.rs index 9f7610c9f..cfa23c38f 100644 --- a/seaweed-volume/src/server/heartbeat.rs +++ b/seaweed-volume/src/server/heartbeat.rs @@ -1508,6 +1508,9 @@ mod tests { security_file: String::new(), cli_white_list: vec![], state_file_path: String::new(), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }) } diff --git a/seaweed-volume/src/server/store_ec.rs b/seaweed-volume/src/server/store_ec.rs index d54c8063f..35b76ca4d 100644 --- a/seaweed-volume/src/server/store_ec.rs +++ b/seaweed-volume/src/server/store_ec.rs @@ -419,13 +419,24 @@ async fn delete_on_ec_shard_holders( let mut last_err = None; if local_shards.contains(&shard_id) { - match journal_delete_local(state, target.vid, target.needle_id) { - Ok(()) => return Ok(true), - // Nothing was committed — the volume unmounted or remounted - // without the needle — so it is safe to fall back to other - // shard holders, unlike an RPC failure which may have landed. - Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(false), - Err(e) => last_err = Some(e), + // A decode in its publishing tail must not miss this delete. The + // tail membership is verified again under the store write lock + // inside journal_delete_local — WouldBlock means the decode claimed + // it in the gap after this wait — so wait and retry. + loop { + crate::server::grpc_server::wait_ec_decode_tail(state, target.vid).await; + match journal_delete_local(state, target.vid, target.needle_id) { + Ok(()) => return Ok(true), + Err(e) if e.kind() == io::ErrorKind::WouldBlock => continue, + // Nothing was committed — the volume unmounted or remounted + // without the needle — so it is safe to fall back to other + // shard holders, unlike an RPC failure which may have landed. + Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(false), + Err(e) => { + last_err = Some(e); + break; + } + } } } if let Some(addrs) = addrs { @@ -487,6 +498,15 @@ fn journal_delete_local( needle_id: NeedleId, ) -> io::Result<()> { let mut store = state.store.write().unwrap(); + // Membership is read under the write lock: a decode can claim the + // publishing tail while this call waited for the decoder's read lock, + // so a check taken earlier would be stale by commit time. + if crate::server::grpc_server::ec_decode_tail_contains(state, vid) { + return Err(io::Error::new( + io::ErrorKind::WouldBlock, + format!("ec volume {} is in decode publishing tail", vid.0), + )); + } let ecv = store.find_ec_volume_mut(vid).ok_or_else(|| { io::Error::new( io::ErrorKind::NotFound, diff --git a/seaweed-volume/src/server/volume_server.rs b/seaweed-volume/src/server/volume_server.rs index 725cc437e..cf848df16 100644 --- a/seaweed-volume/src/server/volume_server.rs +++ b/seaweed-volume/src/server/volume_server.rs @@ -114,6 +114,21 @@ pub struct VolumeServerState { pub cli_white_list: Vec, /// Path to state.pb file for persisting VolumeServerState across restarts. pub state_file_path: String, + /// Volumes with an EC decode in flight. A dropped request leaves the + /// blocking job running; this keeps a retry from racing it on the + /// same volume files. + pub ec_decodes_in_flight: + std::sync::Mutex>, + /// Volumes whose EC decode is in its publishing tail (journal catch-up, + /// .idx write, compaction). Local .ecj appenders wait on + /// `ec_decode_tail_notify` while their vid is listed, so no committed + /// delete falls between the last catch_up and the .cpd/.cpx swap — + /// the per-volume slice of Go's EcVolume.ecjFileAccessLock. + pub ec_decode_tail: + std::sync::Mutex>, + /// Wakes .ecj appenders waiting on `ec_decode_tail` when a decode's + /// publishing tail ends. + pub ec_decode_tail_notify: tokio::sync::Notify, } impl VolumeServerState { diff --git a/seaweed-volume/src/server/write_queue.rs b/seaweed-volume/src/server/write_queue.rs index e27ed486a..80389e802 100644 --- a/seaweed-volume/src/server/write_queue.rs +++ b/seaweed-volume/src/server/write_queue.rs @@ -229,6 +229,9 @@ mod tests { security_file: String::new(), cli_white_list: vec![], state_file_path: String::new(), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), }) } diff --git a/seaweed-volume/src/storage/erasure_coding/ec_decoder.rs b/seaweed-volume/src/storage/erasure_coding/ec_decoder.rs index 94126ebb6..474fcf5af 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_decoder.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_decoder.rs @@ -97,21 +97,97 @@ pub fn read_ecj_deletions( collection: &str, volume_id: VolumeId, ) -> io::Result> { - let mut ids = HashSet::new(); - for (i, dir) in dirs.iter().enumerate() { - if dirs[..i].contains(dir) { - continue; + Ok(EcjDeletions::read(dirs, collection, volume_id)?.ids) +} + +/// The ids [`read_ecj_deletions`] returns, plus how far each journal was read +/// so that ids journaled later can be added. +pub struct EcjDeletions { + pub ids: HashSet, + /// Each distinct journal path and the whole-record length read so far. + journals: Vec<(String, u64)>, +} + +impl EcjDeletions { + pub fn read(dirs: &[&str], collection: &str, volume_id: VolumeId) -> io::Result { + let mut journals: Vec<(String, u64)> = Vec::new(); + for dir in dirs { + let path = format!("{}.ecj", volume_file_name(dir, collection, volume_id)); + if !journals.iter().any(|(p, _)| *p == path) { + journals.push((path, 0)); + } } - let path = format!("{}.ecj", volume_file_name(dir, collection, volume_id)); - let file = match File::open(&path) { - Ok(file) => file, - Err(e) if e.kind() == io::ErrorKind::NotFound => continue, - Err(e) => return Err(e), + let mut deletions = EcjDeletions { + ids: HashSet::new(), + journals, }; - let len = file.metadata()?.len(); - read_ecj_ids(&file, len, &mut ids)?; + deletions.catch_up()?; + Ok(deletions) + } + + /// Adds the ids appended to each journal since the last read. A journal + /// that shrank is read again from the start — and the whole set rebuilt, + /// since ids already folded in from the truncated tail may have been a + /// rolled-back append. Non-regular journals (the FIFOs the tests stand + /// in for a blocking disk) stat empty and cannot be rolled back, so they + /// never count as shrunk. + pub fn catch_up(&mut self) -> io::Result<()> { + let mut shrank = false; + for (path, read_to) in &self.journals { + if *read_to == 0 { + continue; + } + match std::fs::metadata(path) { + Ok(m) => { + if m.is_file() && m.len() < *read_to { + shrank = true; + break; + } + } + // A journal that was read before and is now gone shrank to + // nothing (e.g. the volume was destroyed mid-scan) — the ids + // read from it no longer reflect committed content. + Err(e) if e.kind() == io::ErrorKind::NotFound => { + shrank = true; + break; + } + Err(e) => return Err(e), + } + } + if shrank { + self.ids.clear(); + for (_, read_to) in &mut self.journals { + *read_to = 0; + } + } + for (path, read_to) in &mut self.journals { + let file = match File::open(&*path) { + Ok(file) => file, + Err(e) if e.kind() == io::ErrorKind::NotFound => continue, + Err(e) => return Err(e), + }; + let len = file.metadata()?.len(); + if len < *read_to { + *read_to = 0; + } + read_ecj_ids(&file, *read_to, len, &mut self.ids)?; + *read_to = len - len % NEEDLE_ID_SIZE as u64; + } + Ok(()) + } + + /// Clears the set and re-reads every journal from the start. Run under + /// the caller's store read lock (no append in flight): the result is + /// then exactly the committed content — an earlier unlocked read may + /// have folded in bytes a rolled-back append later truncated, or missed + /// a record re-appended to the very offset a rollback freed. + pub fn rescan(&mut self) -> io::Result<()> { + self.ids.clear(); + for (_, read_to) in &mut self.journals { + *read_to = 0; + } + self.catch_up() } - Ok(ids) } /// What it takes to rebuild a volume's .dat from its EC data shards. @@ -312,6 +388,30 @@ pub fn write_dat_file_from_shards(spec: &DatRebuild<'_>) -> io::Result<()> { write_result } +/// Fails when the decoded `.dat` in `dat_dir` is shorter than the +/// `dat_file_size` bytes its EC index references: the caller deletes the +/// shards next, and they are the only other copy of the needles past the cut. +/// A longer file passes. +pub fn verify_decoded_dat_file( + dat_dir: &str, + collection: &str, + volume_id: VolumeId, + dat_file_size: i64, +) -> io::Result<()> { + let dat_path = format!("{}.dat", volume_file_name(dat_dir, collection, volume_id)); + let size = std::fs::metadata(&dat_path)?.len(); + if (size as i64) < dat_file_size { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + format!( + "decoded {} is {} bytes, short of the {} its ec index references", + dat_path, size, dat_file_size + ), + )); + } + Ok(()) +} + /// Write .idx file from .ecx index + .ecj deletion journal. /// /// See [`write_idx_file_from_ec_index_with_dirs`]; everything lives in `dir`. @@ -790,4 +890,77 @@ mod tests { let expected: HashSet = [1, 2, 3, 7, 9].into_iter().map(NeedleId).collect(); assert_eq!(ids, expected); } + + #[test] + fn test_verify_decoded_dat_file_rejects_a_short_dat() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let dat_path = format!("{dir}/1.dat"); + + let err = verify_decoded_dat_file(dir, "", VolumeId(1), 100).unwrap_err(); + assert_eq!(err.kind(), io::ErrorKind::NotFound); + + std::fs::write(&dat_path, vec![0u8; 99]).unwrap(); + let err = verify_decoded_dat_file(dir, "", VolumeId(1), 100).unwrap_err(); + assert!(err.to_string().contains("short of the 100"), "{err}"); + + std::fs::write(&dat_path, vec![0u8; 100]).unwrap(); + verify_decoded_dat_file(dir, "", VolumeId(1), 100).unwrap(); + std::fs::write(&dat_path, vec![0u8; 101]).unwrap(); + verify_decoded_dat_file(dir, "", VolumeId(1), 100).unwrap(); + } + + /// Ids appended after the first read, including the rest of a record torn + /// at that point, are picked up; a journal that shrank is read again. + #[test] + fn test_ecj_deletions_catch_up_reads_appended_ids() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ecj_path = format!("{dir}/1.ecj"); + let entry = |id: u64| { + let mut buf = [0u8; NEEDLE_ID_SIZE]; + NeedleId(id).to_bytes(&mut buf); + buf + }; + let append = |bytes: &[u8]| { + let mut f = std::fs::OpenOptions::new() + .create(true) + .append(true) + .open(&ecj_path) + .unwrap(); + f.write_all(bytes).unwrap(); + }; + let ids = |d: &EcjDeletions| { + let mut ids: Vec = d.ids.iter().map(|id| id.0).collect(); + ids.sort(); + ids + }; + + // No journal yet. + let mut deletions = EcjDeletions::read(&[dir, dir], "", VolumeId(1)).unwrap(); + assert!(deletions.ids.is_empty()); + + append(&entry(1)); + append(&entry(2)[..3]); + deletions.catch_up().unwrap(); + assert_eq!(ids(&deletions), [1]); + + append(&entry(2)[3..]); + append(&entry(3)); + deletions.catch_up().unwrap(); + assert_eq!(ids(&deletions), [1, 2, 3]); + + // A shrunk journal is rebuilt from its surviving content: ids folded + // in from the truncated tail may have been rolled back and must not + // linger as phantom tombstones. + std::fs::write(&ecj_path, entry(9)).unwrap(); + deletions.catch_up().unwrap(); + assert_eq!(ids(&deletions), [9]); + + // A journal removed since it was read is the extreme shrink: its + // earlier ids must go with it, not linger. + std::fs::remove_file(&ecj_path).unwrap(); + deletions.catch_up().unwrap(); + assert!(deletions.ids.is_empty()); + } } diff --git a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs index 7d8578de3..e47de96d5 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs @@ -49,6 +49,15 @@ pub(crate) struct ShardLocationCache { /// A multiple of `NEEDLE_ID_SIZE`; 1 MiB is 131072 entries per syscall. const ECJ_LOAD_CHUNK_BYTES: usize = 1 << 20; +/// Process-wide epoch bumped whenever a failed .ecj append truncates its +/// uncommitted tail. A decode that read the journal without the store lock +/// compares a sample taken before the read with one taken under the +/// quiescing lock: equal means no rollback ran in between, so every id it +/// folded in was committed and only incremental catch-up remains — the +/// full rescan is needed only when the epoch moved. +pub(crate) static ECJ_ROLLBACK_EPOCH: std::sync::atomic::AtomicU64 = + std::sync::atomic::AtomicU64::new(0); + /// A `.ecj` smaller than this is never rewritten, however redundant. Below a /// megabyte the duplication costs nothing and the rewrite is pure churn. const ECJ_COMPACT_MIN_BYTES: i64 = 1 << 20; @@ -134,15 +143,16 @@ enum EcjPublishError { HandleLost(io::Error), } -/// Adds every whole needle id in the first `len` bytes of `ecj_file` to `ids`, +/// Adds every whole needle id in bytes `from..len` of `ecj_file` to `ids`, /// reading `ECJ_LOAD_CHUNK_BYTES` at a time; a trailing partial record is -/// ignored. +/// ignored. `from` is a record boundary. pub(crate) fn read_ecj_ids( ecj_file: &File, + from: u64, len: u64, ids: &mut HashSet, ) -> io::Result<()> { - read_ecj_ids_with(read_exact_at, ecj_file, len, ids) + read_ecj_ids_with(read_exact_at, ecj_file, from, len, ids) } /// [`read_ecj_ids`] with the positional read supplied by the caller, so the @@ -150,11 +160,13 @@ pub(crate) fn read_ecj_ids( fn read_ecj_ids_with( read_at: fn(&File, &mut [u8], u64) -> io::Result<()>, ecj_file: &File, + from: u64, len: u64, ids: &mut HashSet, ) -> io::Result<()> { - let mut buf = vec![0u8; std::cmp::min(ECJ_LOAD_CHUNK_BYTES as u64, len) as usize]; - let mut off: u64 = 0; + let mut buf = + vec![0u8; std::cmp::min(ECJ_LOAD_CHUNK_BYTES as u64, len.saturating_sub(from)) as usize]; + let mut off: u64 = from; while off + NEEDLE_ID_SIZE as u64 <= len { let mut want = std::cmp::min(ECJ_LOAD_CHUNK_BYTES as u64, len - off) as usize; want -= want % NEEDLE_ID_SIZE; @@ -987,7 +999,7 @@ impl EcVolume { // held the `deleted_needles` write lock for the whole scan, which on a // bloated journal is the entire (unbounded) startup. let mut loaded: HashSet = HashSet::new(); - read_ecj_ids_with(read_at, ecj_file, self.ecj_file_size as u64, &mut loaded)?; + read_ecj_ids_with(read_at, ecj_file, 0, self.ecj_file_size as u64, &mut loaded)?; let mut set = self .deleted_needles @@ -2038,6 +2050,11 @@ impl EcVolume { let ecj_path = self.ecj_file_name(); let rollback = open_volume_file(OpenOptions::new().write(true), &ecj_path) .and_then(|f| f.set_len(prev_ecj_size as u64).and_then(|_| f.sync_all())); + // Bumped on attempt, not success: a failed rollback leaves the + // bytes in place and a decode re-reading for the epoch change + // simply sees them — harmless — while a successful one must + // never go unnoticed by an unlocked journal read. + ECJ_ROLLBACK_EPOCH.fetch_add(1, std::sync::atomic::Ordering::SeqCst); if let Err(trunc_err) = rollback { tracing::error!( volume_id = self.volume_id.0, diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index 0d9332392..bd0b1fb3f 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -22,6 +22,36 @@ use crate::storage::super_block::{ReplicaPlacement, SUPER_BLOCK_SIZE}; use crate::storage::types::*; use crate::storage::volume::{CompactionJob, VifVolumeInfo, Volume, VolumeError, VolumeSpec}; +/// Mirrors Go's ensureCompactVolumeSpace, per filesystem: the new .dat lands +/// next to the old one and the new .idx next to the old index, so when the +/// index directory is on another filesystem each disk answers for its own +/// share, while two directories on one filesystem must cover the sum. +fn ensure_compact_volume_space(v: &Volume, preallocate: u64) -> Result<(), VolumeError> { + let (data_bytes, index_bytes) = compaction_space_needed(v, preallocate); + let (dir, dir_idx) = (v.dir(), v.dir_idx()); + let check = |dir: &str, needed: u64| -> Result<(), VolumeError> { + let (_, free) = crate::storage::disk_location::get_disk_stats(dir); + if free < needed { + return Err(VolumeError::InsufficientSpace { + vid: v.id, + required: needed, + free, + }); + } + Ok(()) + }; + + if dir_idx.is_empty() || dir_idx == dir { + check(dir, data_bytes + index_bytes) + } else if !same_filesystem(dir, dir_idx) { + check(dir, data_bytes)?; + check(dir_idx, index_bytes) + } else { + check(dir, data_bytes + index_bytes)?; + check(dir_idx, index_bytes) + } +} + /// Top-level storage manager containing all disk locations and their volumes. pub struct Store { pub locations: Vec, @@ -1557,41 +1587,10 @@ impl Store { vid: VolumeId, preallocate: u64, ) -> Result, VolumeError> { - // Mirrors Go's ensureCompactVolumeSpace, per filesystem. - let (dir, dir_idx, data_bytes, index_bytes) = { - let (_, v) = self - .find_volume(vid) - .ok_or(VolumeError::VolumeNotFound(vid))?; - let (data_bytes, index_bytes) = compaction_space_needed(v, preallocate); - ( - v.dir().to_string(), - v.dir_idx().to_string(), - data_bytes, - index_bytes, - ) - }; - - let check = |dir: &str, needed: u64| -> Result<(), VolumeError> { - let (_, free) = crate::storage::disk_location::get_disk_stats(dir); - if free < needed { - return Err(VolumeError::InsufficientSpace { - vid, - required: needed, - free, - }); - } - Ok(()) - }; - - if dir_idx.is_empty() || dir_idx == dir { - check(&dir, data_bytes + index_bytes)?; - } else if !same_filesystem(&dir, &dir_idx) { - check(&dir, data_bytes)?; - check(&dir_idx, index_bytes)?; - } else { - check(&dir, data_bytes + index_bytes)?; - check(&dir_idx, index_bytes)?; - } + let (_, v) = self + .find_volume(vid) + .ok_or(VolumeError::VolumeNotFound(vid))?; + ensure_compact_volume_space(v, preallocate)?; let (_, v) = self .find_volume_mut(vid) @@ -1599,6 +1598,37 @@ impl Store { v.begin_compact_by_index() } + /// Rewrite the volume in `dir`/`dir_idx`, which is not mounted, with its + /// live needles only. Go's `Store.CompactVolumeFiles`. + pub fn compact_volume_files( + dir: &str, + dir_idx: &str, + collection: &str, + vid: VolumeId, + needle_map_kind: NeedleMapKind, + ) -> Result<(), VolumeError> { + let spec = VolumeSpec { + collection, + ..VolumeSpec::default() + }; + let mut v = Volume::new(dir, dir_idx, vid, needle_map_kind, &spec)?; + let mut compact = || -> Result<(), VolumeError> { + ensure_compact_volume_space(&v, 0)?; + v.compact_by_index(0, 0, |_| true)?; + v.commit_compact() + }; + let result = compact(); + if result.is_err() { + // A failed commit may have swapped only one of .dat/.idx; + // reconcile rolls a decided swap forward or removes orphan + // temp files before this volume can mount a mismatched pair. + let _ = v.reconcile_compact_state(); + let _ = v.cleanup_compact(); + } + v.close(); + result + } + /// Commit a completed compaction: swap files and reload. pub fn commit_compact_volume(&mut self, vid: VolumeId) -> Result<(bool, u64), VolumeError> { let (_, v) = self diff --git a/seaweed-volume/tests/http_integration.rs b/seaweed-volume/tests/http_integration.rs index 3ed07d003..723461d4e 100644 --- a/seaweed-volume/tests/http_integration.rs +++ b/seaweed-volume/tests/http_integration.rs @@ -108,6 +108,9 @@ fn build_test_state( file_size_limit_bytes: 0, maintenance_byte_per_second: 0, is_heartbeating: std::sync::atomic::AtomicBool::new(true), + ec_decodes_in_flight: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail: std::sync::Mutex::new(std::collections::HashSet::new()), + ec_decode_tail_notify: tokio::sync::Notify::new(), has_master: false, pre_stop_seconds: 0, volume_state_notify: tokio::sync::Notify::new(), diff --git a/weed/server/volume_grpc_erasure_coding.go b/weed/server/volume_grpc_erasure_coding.go index aaee5c8af..5405fe1e0 100644 --- a/weed/server/volume_grpc_erasure_coding.go +++ b/weed/server/volume_grpc_erasure_coding.go @@ -268,7 +268,23 @@ func (vs *VolumeServer) VolumeEcShardsRebuild(ctx context.Context, req *volume_s if !util.FileExists(indexBaseFileName+".ecx") && rebuildLocation.IdxDirectory != rebuildLocation.Directory { indexBaseFileName = path.Join(rebuildLocation.Directory, baseFileName) } - if err := erasure_coding.RebuildEcxFile(indexBaseFileName); err != nil { + if ev, found := rebuildLocation.FindEcVolume(needle.VolumeId(req.VolumeId)); found { + // A mounted volume appends runtime deletes through ev.ecjFile, so the + // fold runs under its journal lock on the journal's own base, and the + // handle is repointed: otherwise RebuildEcxFile's unlink strands later + // appends on the detached inode. + indexBaseFileName = ev.EcIndexBaseFileName() + unlockJournal := ev.LockDeletionJournal() + err := erasure_coding.RebuildEcxFile(indexBaseFileName) + if err == nil { + err = ev.ReopenDeletionJournal() + } + unlockJournal() + if err != nil { + recordEcRebuild("failure", time.Since(start)) + return nil, fmt.Errorf("RebuildEcxFile %s: %v", indexBaseFileName, err) + } + } else if err := erasure_coding.RebuildEcxFile(indexBaseFileName); err != nil { recordEcRebuild("failure", time.Since(start)) return nil, fmt.Errorf("RebuildEcxFile %s: %v", indexBaseFileName, err) } @@ -1138,6 +1154,11 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_ return vs.adoptStagedVolume(req) } + if _, loaded := vs.ecDecodesInFlight.LoadOrStore(req.VolumeId, struct{}{}); loaded { + return nil, status.Errorf(codes.Unavailable, "ec volume %d is already being decoded", req.VolumeId) + } + defer vs.ecDecodesInFlight.Delete(req.VolumeId) + // Collect all EC shards (NewEcVolume will load EC config from .vif into v.ECContext) // Use MaxShardCount (32) to support custom EC ratios up to 32 total shards tempShards := make([]string, erasure_coding.MaxShardCount) @@ -1168,18 +1189,50 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_ } } - dataBaseFileName, indexBaseFileName := v.DataBaseFileName(), v.IndexBaseFileName() - if !util.FileExists(indexBaseFileName + ".ecx") { - indexBaseFileName = dataBaseFileName + dataBaseFileName := v.DataBaseFileName() + // The fold and the index write must work on the same .ecj that runtime + // deletes append to through ecjFile — the volume's resolved index dir, + // which can differ from IndexBaseFileName when an .ecx copy exists in + // both the data and index directories. + indexBaseFileName := v.EcIndexBaseFileName() + + // Resolve the offline-compaction location before taking the journal + // lock: FindEcVolume takes the ec-volume map lock, and acquiring it inside + // the journal lock inverts DestroyEcVolume's map->journal order — a decode + // holding the journal lock while waiting on the map deadlocks a destroy + // holding the map lock while waiting on the journal. + var volumeLocation *storage.DiskLocation + for _, location := range vs.store.Locations { + if candidate, found := location.FindEcVolume(needle.VolumeId(req.VolumeId)); found && candidate == v { + volumeLocation = location + break + } + } + if volumeLocation == nil { + return nil, fmt.Errorf("ec volume %d location not found for offline compaction", req.VolumeId) } // Merge .ecj deletions into .ecx so that HasLiveNeedles and FindDatFileSize // see the full set of deleted needles. Without this, needles deleted after the // last ecx rebuild would still appear live, causing the decoded .dat to include // data that should be skipped and HasLiveNeedles to return a false positive. + // + // The journal lock is held across the fold: RebuildEcxFile unlinks .ecj, + // and a delete committed between its read and the unlink would land on the + // detached inode that ecjFile keeps open — synced, successful, and + // invisible to every path-based reader. ReopenDeletionJournal then points + // the handle back at a fresh journal, so deletes committed during the .dat + // rebuild stay durable and reach the index write below. + unlockJournal := v.LockDeletionJournal() if err := erasure_coding.RebuildEcxFile(indexBaseFileName); err != nil { + unlockJournal() return nil, fmt.Errorf("RebuildEcxFile %s: %v", indexBaseFileName, err) } + if err := v.ReopenDeletionJournal(); err != nil { + unlockJournal() + return nil, fmt.Errorf("reopen deletion journal %s: %v", indexBaseFileName, err) + } + unlockJournal() // If the EC index contains no live entries, decoding should be a no-op: // just allow the caller to purge EC shards and do not generate an empty normal volume. @@ -1215,6 +1268,11 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_ return nil, err } + // Runtime deletes serialize on the volume's journal lock; holding it from + // the journal-consuming index write through the offline compaction keeps a + // committed delete from slipping past the rebuilt .idx. + defer v.LockDeletionJournal()() + // write .idx file from .ecx and .ecj files if err := erasure_coding.WriteIdxFileFromEcIndex(indexBaseFileName); err != nil { return nil, fmt.Errorf("WriteIdxFileFromEcIndex %s: %v", v.IndexBaseFileName(), err) @@ -1227,17 +1285,6 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_ removeBitrotSidecars(indexBaseFileName) } - var volumeLocation *storage.DiskLocation - for _, location := range vs.store.Locations { - if candidate, found := location.FindEcVolume(needle.VolumeId(req.VolumeId)); found && candidate == v { - volumeLocation = location - break - } - } - if volumeLocation == nil { - return nil, fmt.Errorf("ec volume %d location not found for offline compaction", req.VolumeId) - } - if err := vs.store.CompactVolumeFiles( needle.VolumeId(req.VolumeId), v.Collection, @@ -1247,6 +1294,13 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_ 0, vs.compactionBytePerSecond, ); err != nil { + // Benign only when the swap never started or was settled: a surviving + // .cpc marker means the commit was decided but the renames could not + // be reconciled, so .dat/.idx may be a mismatched pair — fail the + // decode and let the caller keep the shards. + if util.FileExists(dataBaseFileName + ".cpc") { + return nil, fmt.Errorf("CompactVolumeFiles %s: %w", dataBaseFileName, err) + } glog.Errorf("CompactVolumeFiles %s: %v", dataBaseFileName, err) } diff --git a/weed/server/volume_server.go b/weed/server/volume_server.go index 1b6feeac4..737d48094 100644 --- a/weed/server/volume_server.go +++ b/weed/server/volume_server.go @@ -60,6 +60,11 @@ type VolumeServer struct { fileSizeLimitBytes int64 isHeartbeating bool stopChan chan bool + + // Volumes with an EC decode in flight. A cancelled request does not + // stop the handler; this keeps a retry from racing it on the same + // volume files. + ecDecodesInFlight sync.Map } func NewVolumeServer(adminMux, publicMux *http.ServeMux, ip string, diff --git a/weed/storage/disk_location.go b/weed/storage/disk_location.go index 51afa1f4c..5f468eb53 100644 --- a/weed/storage/disk_location.go +++ b/weed/storage/disk_location.go @@ -48,6 +48,11 @@ type DiskLocation struct { // erasure coding ecVolumes map[needle.VolumeId]*erasure_coding.EcVolume ecVolumesLock sync.RWMutex + // vids whose EcVolume was removed from ecVolumes but is still being + // destroyed off-lock; the channel closes when Destroy returns. A remount + // of the vid must wait for it, or the dying volume could unlink files + // the new one just opened. + ecVolumesDestroying map[needle.VolumeId]chan struct{} ecShardNotifyHandler func(collection string, vid needle.VolumeId, shardId erasure_coding.ShardId, ecVolume *erasure_coding.EcVolume) @@ -117,6 +122,7 @@ func NewDiskLocation(dir string, maxVolumeCount int32, minFreeSpace util.MinFree } location.volumes = make(map[needle.VolumeId]*Volume) location.ecVolumes = make(map[needle.VolumeId]*erasure_coding.EcVolume) + location.ecVolumesDestroying = make(map[needle.VolumeId]chan struct{}) location.closeCh = make(chan struct{}) go func() { location.CheckDiskSpace(config) @@ -474,8 +480,12 @@ func (l *DiskLocation) DeleteCollectionFromDiskLocation(collection string) (dele }() go func() { - for _, v := range delEcVolsMap { + for k, v := range delEcVolsMap { v.Destroy() + l.ecVolumesLock.Lock() + done := l.ecVolumesDestroying[k] + l.ecVolumesLock.Unlock() + close(done) } wg.Done() }() @@ -649,11 +659,20 @@ func (l *DiskLocation) Close() { l.volumesLock.Unlock() l.ecVolumesLock.Lock() - for _, ecVolume := range l.ecVolumes { - ecVolume.Close() + ecVolumes := make([]*erasure_coding.EcVolume, 0, len(l.ecVolumes)) + for vid, ecVolume := range l.ecVolumes { + ecVolumes = append(ecVolumes, ecVolume) + delete(l.ecVolumes, vid) } l.ecVolumesLock.Unlock() + // Close outside the write lock: EcVolume.Close takes the deletion-journal + // lock, which a running ec.decode can hold — closing under the map lock + // would stall every EC lookup and invert the map->journal lock order. + for _, ecVolume := range ecVolumes { + ecVolume.Close() + } + close(l.closeCh) return } diff --git a/weed/storage/disk_location_ec.go b/weed/storage/disk_location_ec.go index be98e371d..8205ede9c 100644 --- a/weed/storage/disk_location_ec.go +++ b/weed/storage/disk_location_ec.go @@ -24,6 +24,17 @@ var ( re = regexp.MustCompile(`\.ec\d{2,3}`) ) +// markEcDestroying records a pending destroy for vid and returns the channel +// closed when it finishes. Callers must hold ecVolumesLock for writing. +func (l *DiskLocation) markEcDestroying(vid needle.VolumeId) chan struct{} { + if l.ecVolumesDestroying == nil { + l.ecVolumesDestroying = make(map[needle.VolumeId]chan struct{}) + } + done := make(chan struct{}) + l.ecVolumesDestroying[vid] = done + return done +} + func (l *DiskLocation) FindEcVolume(vid needle.VolumeId) (*erasure_coding.EcVolume, bool) { l.ecVolumesLock.RLock() defer l.ecVolumesLock.RUnlock() @@ -37,13 +48,27 @@ func (l *DiskLocation) FindEcVolume(vid needle.VolumeId) (*erasure_coding.EcVolu func (l *DiskLocation) DestroyEcVolume(vid needle.VolumeId) { l.ecVolumesLock.Lock() - defer l.ecVolumesLock.Unlock() - ecVolume, found := l.ecVolumes[vid] + var done chan struct{} if found { - ecVolume.Destroy() + done = l.markEcDestroying(vid) delete(l.ecVolumes, vid) } + l.ecVolumesLock.Unlock() + + if !found { + return + } + // Destroy outside the write lock: EcVolume.Destroy's Close waits on the + // deletion-journal lock, which a running ec.decode can hold — under the + // map lock that wait would stall every EC lookup and invert the + // map->journal lock order. + ecVolume.Destroy() + + // The tombstone stays after close() as this vid's destroy generation: a + // remount compares it before and after opening files to spot a destroy + // that ran inside the window. The next successful mount clears it. + close(done) } // UnloadEcVolume drops the in-memory EcVolume for vid from this one disk without @@ -141,17 +166,57 @@ func (l *DiskLocation) LoadEcShard(collection string, vid needle.VolumeId, shard // (issue #9212). func (l *DiskLocation) loadEcShardWithIdxDir(collection string, vid needle.VolumeId, shardId erasure_coding.ShardId, idxDir string) (*erasure_coding.EcVolume, error) { - ecVolumeShard, err := erasure_coding.NewEcVolumeShard(l.DiskType, l.Directory, collection, vid, shardId) - if err != nil { - if err == os.ErrNotExist { - return nil, os.ErrNotExist + // A destroy-in-progress or already-finished one leaves a tombstone; it + // stays in the map until a mount clears it, so comparing the channel + // before and after opening the shard detects any destroy that could have + // unlinked the files in between — including one that ran to completion + // inside the window. + var ecVolumeShard *erasure_coding.EcVolumeShard + for { + l.ecVolumesLock.Lock() + gen := l.ecVolumesDestroying[vid] + l.ecVolumesLock.Unlock() + if gen != nil { + select { + case <-gen: + // Destroy already finished; the closed tombstone is the + // generation token the post-open recheck compares against. + default: + // Destroy in flight: its unlinks could detach anything this + // mount opens, so wait it out and take a fresh generation. + <-gen + continue + } } - return nil, fmt.Errorf("failed to create ec shard %d.%d: %w", vid, shardId, err) + + var err error + ecVolumeShard, err = erasure_coding.NewEcVolumeShard(l.DiskType, l.Directory, collection, vid, shardId) + if err != nil { + if err == os.ErrNotExist { + return nil, os.ErrNotExist + } + return nil, fmt.Errorf("failed to create ec shard %d.%d: %w", vid, shardId, err) + } + + l.ecVolumesLock.Lock() + if cur := l.ecVolumesDestroying[vid]; cur != gen { + // A destroy intervened while the shard was being opened; its + // unlink may have detached the file. Drop the handle, wait for + // the destroy, and retry on a clean slate. + l.ecVolumesLock.Unlock() + ecVolumeShard.Unmount() // release the gauge the constructor's Mount took + ecVolumeShard.Close() + if cur != nil { + <-cur + } + continue + } + break } - l.ecVolumesLock.Lock() defer l.ecVolumesLock.Unlock() ecVolume, found := l.ecVolumes[vid] if !found { + var err error ecVolume, err = erasure_coding.NewEcVolume(l.DiskType, l.Directory, idxDir, collection, vid) if err != nil { // Wrap with %w so MountEcShards / startup reconcile can use @@ -160,6 +225,9 @@ func (l *DiskLocation) loadEcShardWithIdxDir(collection string, vid needle.Volum return nil, fmt.Errorf("failed to create ec volume %d: %w", vid, err) } l.ecVolumes[vid] = ecVolume + // Fresh generation for this vid: the next destroy sets a new + // tombstone rather than matching the cleared one. + delete(l.ecVolumesDestroying, vid) } added, err := ecVolume.AddEcVolumeShard(ecVolumeShard) if err != nil { @@ -389,16 +457,21 @@ func (l *DiskLocation) loadEcShardsWithIdxDir(shards []string, collection string } func (l *DiskLocation) deleteEcVolumeById(vid needle.VolumeId) (e error) { - // Add write lock since we're modifying the ecVolumes map l.ecVolumesLock.Lock() - defer l.ecVolumesLock.Unlock() - ecVolume, ok := l.ecVolumes[vid] + var done chan struct{} + if ok { + done = l.markEcDestroying(vid) + delete(l.ecVolumes, vid) + } + l.ecVolumesLock.Unlock() + if !ok { return } + // Destroy outside the map lock — see DestroyEcVolume. ecVolume.Destroy() - delete(l.ecVolumes, vid) + close(done) return } @@ -410,8 +483,11 @@ func (l *DiskLocation) unmountEcVolumeByCollection(collectionName string) map[ne } } - for k, _ := range deltaVols { + for k := range deltaVols { delete(l.ecVolumes, k) + // Caller destroys these outside the lock; tombstone until then so a + // remount can't open files the destroy is about to unlink. + l.markEcDestroying(k) } return deltaVols } diff --git a/weed/storage/disk_location_ec_test.go b/weed/storage/disk_location_ec_test.go index 3ca3765f2..4c3e77a75 100644 --- a/weed/storage/disk_location_ec_test.go +++ b/weed/storage/disk_location_ec_test.go @@ -847,6 +847,113 @@ func TestLoadEcShardDuplicateReleasesTheNewShard(t *testing.T) { } } +// TestLoadEcShardWaitsForPendingDestroy: DestroyEcVolume deletes the map entry +// and destroys files outside ecVolumesLock, so a remount could otherwise open +// shard files mid-unlink. The ecVolumesDestroying tombstone must hold the mount +// until the destroy closes it, and the mount must then proceed and clear it. +func TestLoadEcShardWaitsForPendingDestroy(t *testing.T) { + dir := t.TempDir() + diskLocation := NewDiskLocation(dir, 10, util.MinFreeSpace{}, dir, types.HardDriveType, nil, stats.DefaultDiskIOProbeConfig()) + defer diskLocation.Close() + + if err := os.WriteFile(filepath.Join(dir, "125.ec00"), []byte("shard bytes"), 0o644); err != nil { + t.Fatalf("seed .ec00: %v", err) + } + if err := os.WriteFile(filepath.Join(dir, "125.ecx"), make([]byte, 16), 0o644); err != nil { + t.Fatalf("seed .ecx: %v", err) + } + + done := make(chan struct{}) + diskLocation.ecVolumesLock.Lock() + diskLocation.ecVolumesDestroying[needle.VolumeId(125)] = done + diskLocation.ecVolumesLock.Unlock() + + mounted := make(chan error, 1) + go func() { + _, err := diskLocation.LoadEcShard("", needle.VolumeId(125), erasure_coding.ShardId(0)) + mounted <- err + }() + + select { + case <-mounted: + t.Fatal("LoadEcShard returned while a destroy generation was still in flight") + case <-time.After(100 * time.Millisecond): + } + + close(done) + select { + case err := <-mounted: + if err != nil { + t.Fatalf("LoadEcShard after destroy completion: %v", err) + } + case <-time.After(5 * time.Second): + t.Fatal("LoadEcShard did not proceed after the destroy finished") + } + + if _, found := diskLocation.FindEcShard(needle.VolumeId(125), erasure_coding.ShardId(0)); !found { + t.Fatal("shard must be registered once the pending destroy is over") + } + diskLocation.ecVolumesLock.RLock() + _, tombstoneLeft := diskLocation.ecVolumesDestroying[needle.VolumeId(125)] + diskLocation.ecVolumesLock.RUnlock() + if tombstoneLeft { + t.Error("a successful remount must clear the vid's destroy tombstone") + } +} + +// TestDestroyEcVolumeClosesTombstone: once DestroyEcVolume returns, the +// tombstone is closed so waits unblock, and a fresh remount re-opens the +// regenerated files rather than anything the destroy unlinked. +func TestDestroyEcVolumeClosesTombstone(t *testing.T) { + dir := t.TempDir() + diskLocation := NewDiskLocation(dir, 10, util.MinFreeSpace{}, dir, types.HardDriveType, nil, stats.DefaultDiskIOProbeConfig()) + defer diskLocation.Close() + + seedShard := func() { + if err := os.WriteFile(filepath.Join(dir, "126.ec00"), []byte("shard bytes"), 0o644); err != nil { + t.Fatalf("seed .ec00: %v", err) + } + if err := os.WriteFile(filepath.Join(dir, "126.ecx"), make([]byte, 16), 0o644); err != nil { + t.Fatalf("seed .ecx: %v", err) + } + } + seedShard() + + if _, err := diskLocation.LoadEcShard("", needle.VolumeId(126), erasure_coding.ShardId(0)); err != nil { + t.Fatalf("initial LoadEcShard: %v", err) + } + diskLocation.DestroyEcVolume(needle.VolumeId(126)) + + if _, found := diskLocation.FindEcVolume(needle.VolumeId(126)); found { + t.Fatal("destroyed volume must not stay in ecVolumes") + } + diskLocation.ecVolumesLock.RLock() + done, ok := diskLocation.ecVolumesDestroying[needle.VolumeId(126)] + diskLocation.ecVolumesLock.RUnlock() + if !ok { + t.Fatal("destroy must leave a tombstone behind for remounts to compare against") + } + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("tombstone must be closed once DestroyEcVolume returns") + } + if util.FileExists(filepath.Join(dir, "126.ec00")) { + t.Fatal("destroy must have removed the shard file") + } + + seedShard() + if _, err := diskLocation.LoadEcShard("", needle.VolumeId(126), erasure_coding.ShardId(0)); err != nil { + t.Fatalf("remount after destroy: %v", err) + } + diskLocation.ecVolumesLock.RLock() + _, tombstoneLeft := diskLocation.ecVolumesDestroying[needle.VolumeId(126)] + diskLocation.ecVolumesLock.RUnlock() + if tombstoneLeft { + t.Error("remount must clear the stale destroy tombstone") + } +} + // TestLoadAllEcShardsSplitDirZeroSizedCleanup: the scan merges Directory and // IdxDirectory listings, so a stale zero-sized file in one directory and a // fresh same-named file in the other are different files behind one entry diff --git a/weed/storage/erasure_coding/ec_volume.go b/weed/storage/erasure_coding/ec_volume.go index 7b6dd286f..fb65f9732 100644 --- a/weed/storage/erasure_coding/ec_volume.go +++ b/weed/storage/erasure_coding/ec_volume.go @@ -587,6 +587,15 @@ func (ev *EcVolume) IndexBaseFileName() string { return EcShardFileName(ev.Collection, ev.dirIdx, int(ev.VolumeId)) } +// EcIndexBaseFileName returns the base path of the volume's live .ecx/.ecj +// pair — where NewEcVolume actually found them (ecxActualDir), which can +// differ from IndexBaseFileName when an .ecx copy exists in both the data +// and index directories. Anything that must agree with the live journal +// handle (ecjFile) uses this, not the configured index dir. +func (ev *EcVolume) EcIndexBaseFileName() string { + return EcShardFileName(ev.Collection, ev.ecxActualDir, int(ev.VolumeId)) +} + func (ev *EcVolume) ShardSize() uint64 { if len(ev.Shards) > 0 { return uint64(ev.Shards[0].Size()) diff --git a/weed/storage/erasure_coding/ec_volume_delete.go b/weed/storage/erasure_coding/ec_volume_delete.go index 44929935f..abb951d74 100644 --- a/weed/storage/erasure_coding/ec_volume_delete.go +++ b/weed/storage/erasure_coding/ec_volume_delete.go @@ -6,6 +6,7 @@ import ( "os" "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/storage/backend" "github.com/seaweedfs/seaweedfs/weed/storage/types" "github.com/seaweedfs/seaweedfs/weed/util" ) @@ -110,6 +111,41 @@ func (ev *EcVolume) appendJournalLocked(b []byte) error { return nil } +// LockDeletionJournal serializes the caller against runtime .ecj appends +// (DeleteNeedleFromEcx) until the returned func is called. Decode holds it +// while consuming the journal into the rebuilt index so a committed delete +// cannot slip past the publish. +func (ev *EcVolume) LockDeletionJournal() func() { + ev.ecjFileAccessLock.Lock() + return ev.ecjFileAccessLock.Unlock +} + +// ReopenDeletionJournal repoints ecjFile at the live .ecj path after +// RebuildEcxFile unlinked it. Without this, later DeleteNeedleFromEcx +// appends keep landing on the detached inode — synced, successful, and +// invisible to every path-based reader. Caller must hold the journal lock +// (see LockDeletionJournal). +func (ev *EcVolume) ReopenDeletionJournal() error { + if ev.ecjFile != nil { + _ = ev.ecjFile.Close() + ev.ecjFile = nil + } + ecjFile, err := backend.OpenVolumeFile(ev.FileName(".ecj"), os.O_RDWR|os.O_CREATE) + if err != nil { + return fmt.Errorf("reopen ec volume journal %s: %w", ev.FileName(".ecj"), err) + } + ev.ecjFile = ecjFile + // A successful RebuildEcxFile leaves the path unlinked, so this is a + // fresh empty file — but track whatever is actually there so the + // rollback truncate in DeleteNeedleFromEcx can never wipe real records. + if fi, statErr := ecjFile.Stat(); statErr == nil { + ev.ecjFileSize = fi.Size() + } else { + ev.ecjFileSize = 0 + } + return nil +} + func RebuildEcxFile(baseFileName string) error { if !util.FileExists(baseFileName + ".ecj") { diff --git a/weed/storage/store_vacuum.go b/weed/storage/store_vacuum.go index aee7843c1..f2e98d4ed 100644 --- a/weed/storage/store_vacuum.go +++ b/weed/storage/store_vacuum.go @@ -184,6 +184,12 @@ func (s *Store) CompactVolumeFiles(vid needle.VolumeId, collection string, locat } if err := tempVolume.CommitCompact(); err != nil { + // A failed commit may have swapped only one of .dat/.idx; reconcile + // rolls a decided swap forward or removes orphan temp files before + // this volume can mount with a mismatched pair. + if reconcileErr := tempVolume.reconcileCompactState(); reconcileErr != nil { + return fmt.Errorf("commit compact volume %d: %v (reconcile failed: %v)", vid, err, reconcileErr) + } if cleanupErr := tempVolume.cleanupCompact(); cleanupErr != nil { return fmt.Errorf("commit compact volume %d: %v (cleanup failed: %v)", vid, err, cleanupErr) }