diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 74b269d89..67eedf42c 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -2095,6 +2095,9 @@ impl VolumeServer for VolumeGrpcService { use tokio::io::AsyncWriteExt; let mut stream = request.into_inner(); + // Held while an EC .ecj is received, through the cleanup of a partial + // file below. Declared before the file so it is dropped after it. + let mut ecj_write: Option = None; // tokio::fs + BufWriter, as `drain_copy_stream_to_file` below already // does: the chunk writes and the final fsync are disk I/O and must not // run on the runtime worker that is also driving this stream. @@ -2244,6 +2247,20 @@ impl VolumeServer for VolumeGrpcService { } }; + // The mounted check above runs once; a volume can still + // mount on this journal while the stream writes it. + // Registered as a writer before the file is created, + // that mount cannot compact the journal and leave the + // rest of the stream in an unlinked inode. + if info.is_ec_volume && info.ext == ".ecj" && ecj_write.is_none() { + ecj_write = Some( + crate::storage::erasure_coding::ecj_registry::begin_ecj_write_async( + &path, + ) + .await, + ); + } + let f = tokio::fs::File::create(&path).await.map_err(|e| { Status::internal(format!("failed to create file: {}", e)) })?; @@ -8384,6 +8401,101 @@ mod tests { assert_eq!(written, payload, "file on disk does not match the payload"); } + /// A volume that mounts on a journal ReceiveFile is still writing must not + /// compact it: the rest of the stream would land in the replaced inode and + /// be gone at the next mount. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn receive_file_ecj_stream_blocks_mount_compaction() { + let (service, tmp) = make_local_service_with_volume("", None); + let dir = tmp.path().to_str().unwrap().to_string(); + let (port, _shutdown) = serve_source(service).await; + // An empty index is enough to mount; the store never sees this volume, + // so ReceiveFile's mounted check lets the .ecj through. + std::fs::write(format!("{}/4.ecx", dir), b"").unwrap(); + let ecj_path = format!("{}/4.ecj", dir); + + // A bloated prefix (100 ids, 4096 times over) the mount would compact, + // then an id only the second chunk carries. + let mut one = vec![0u8; 100 * NEEDLE_ID_SIZE]; + for (i, entry) in one.chunks_exact_mut(NEEDLE_ID_SIZE).enumerate() { + NeedleId(1000 + i as u64).to_bytes(entry); + } + let bloated = one.repeat(4096); + let mut tail = vec![0u8; NEEDLE_ID_SIZE]; + NeedleId(5000).to_bytes(&mut tail); + + let mut client = volume_server_pb::volume_server_client::VolumeServerClient::connect( + format!("http://127.0.0.1:{}", port), + ) + .await + .unwrap(); + let (tx, rx) = tokio::sync::mpsc::channel(4); + let call = tokio::spawn(async move { + client + .receive_file(tokio_stream::wrappers::ReceiverStream::new(rx)) + .await + }); + let send = |data| volume_server_pb::ReceiveFileRequest { data: Some(data) }; + tx.send(send(volume_server_pb::receive_file_request::Data::Info( + volume_server_pb::ReceiveFileInfo { + volume_id: 4, + ext: ".ecj".to_string(), + is_ec_volume: true, + file_size: (bloated.len() + tail.len()) as u64, + ..Default::default() + }, + ))) + .await + .unwrap(); + tx.send(send( + volume_server_pb::receive_file_request::Data::FileContent(bloated.clone()), + )) + .await + .unwrap(); + + // Mount once the first chunk is on disk and before the second is sent. + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + while std::fs::metadata(&ecj_path).map_or(0, |m| m.len()) < bloated.len() as u64 { + assert!( + std::time::Instant::now() < deadline, + "first chunk never reached disk" + ); + tokio::time::sleep(std::time::Duration::from_millis(5)).await; + } + let mount_dir = dir.clone(); + let mounted = tokio::task::spawn_blocking(move || { + crate::storage::erasure_coding::ec_volume::EcVolume::new( + &mount_dir, + &mount_dir, + "", + VolumeId(4), + ) + }) + .await + .unwrap() + .expect("mount during the stream"); + + tx.send(send( + volume_server_pb::receive_file_request::Data::FileContent(tail.clone()), + )) + .await + .unwrap(); + drop(tx); + let response = call.await.unwrap().unwrap().into_inner(); + assert_eq!(response.error, "", "ReceiveFile reported an error"); + drop(mounted); + + let mut want = bloated; + want.extend_from_slice(&tail); + let got = std::fs::read(&ecj_path).unwrap(); + assert_eq!( + got.len(), + want.len(), + "journal on disk is not what the stream sent" + ); + assert!(got == want, "journal on disk is not what the stream sent"); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[expect( clippy::await_holding_lock, diff --git a/seaweed-volume/src/storage/disk_location.rs b/seaweed-volume/src/storage/disk_location.rs index d588447e0..86925d2f8 100644 --- a/seaweed-volume/src/storage/disk_location.rs +++ b/seaweed-volume/src/storage/disk_location.rs @@ -18,7 +18,9 @@ use crate::storage::erasure_coding::ec_shard::{ DATA_SHARDS_COUNT, ERASURE_CODING_LARGE_BLOCK_SIZE, ERASURE_CODING_SMALL_BLOCK_SIZE, EcVolumeShard, ShardId, }; -use crate::storage::erasure_coding::ec_volume::{EcVolume, is_usable_ecx_file}; +use crate::storage::erasure_coding::ec_volume::{ + ECJ_COMPACT_TMP_EXT, EcVolume, is_usable_ecx_file, +}; use crate::storage::needle_map::NeedleMapKind; use crate::storage::super_block::SUPER_BLOCK_SIZE; use crate::storage::types::*; @@ -460,13 +462,16 @@ impl DiskLocation { let idx_base = volume_file_name(&self.idx_directory, collection, vid); const MAX_SHARD_COUNT: usize = 32; - // Remove index files from idx directory (.ecx, .ecj) + // Remove index files from idx directory (.ecx, .ecj, and a compaction + // tmp a crash may have left beside the .ecj) rm_if_present(format!("{}.ecx", idx_base))?; rm_if_present(format!("{}.ecj", idx_base))?; + rm_if_present(format!("{}{}", idx_base, ECJ_COMPACT_TMP_EXT))?; // Also try data directory in case .ecx/.ecj were created before -dir.idx was configured if self.idx_directory != self.directory { rm_if_present(format!("{}.ecx", base))?; rm_if_present(format!("{}.ecj", base))?; + rm_if_present(format!("{}{}", base, ECJ_COMPACT_TMP_EXT))?; } // Remove all EC shard files (.ec00 ~ .ec31) diff --git a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs index 2544887a7..7d8578de3 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs @@ -5,13 +5,14 @@ use std::collections::{HashMap, HashSet}; use std::fs::{self, File, OpenOptions}; -use std::io::{self, Write}; +use std::io::{self, BufWriter, Write}; use std::sync::RwLock; use std::time::{Instant, SystemTime, UNIX_EPOCH}; use crate::pb::master_pb; use crate::storage::erasure_coding::ec_locate; use crate::storage::erasure_coding::ec_shard::*; +use crate::storage::erasure_coding::ecj_registry::EcjHold; use crate::storage::io::read_exact_at; use crate::storage::io_error::IoErrorTracker; use crate::storage::needle::needle::{Needle, NeedleError, get_actual_size}; @@ -48,6 +49,91 @@ 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; +/// 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; + +/// Rewrite only when the journal is at least this many times larger than the +/// set it encodes. A journal holds one entry per delete, so a healthy one is +/// close to 1x; four times the set means most of the file is repeats. +const ECJ_COMPACT_RATIO: i64 = 4; + +/// Staging file for a compacted journal, next to the `.ecj` it replaces. +/// Listed with the other EC index files wherever those are removed, so a tmp +/// left by a crash between write and rename does not outlive its volume. +pub(crate) const ECJ_COMPACT_TMP_EXT: &str = ".ecj.compact.tmp"; + +/// Open `path` as a deletion journal append handle, creating it if absent. +/// Every handle `journal_delete` appends through comes from here, so a handle +/// reopened after compaction behaves exactly like the one opened at mount. +fn open_ecj_append(path: &str) -> io::Result { + open_volume_file( + OpenOptions::new() + .read(true) + .write(true) + .create(true) + .append(true), + path, + ) +} + +/// The filesystem steps of a mount's journal compaction: the reads that load +/// the set it compacts from, and the steps that publish the compacted file. +/// Production uses [`EcjFsOps::REAL`]; tests substitute their own steps to +/// cover the failure paths and races through the real mount. +#[derive(Clone, Copy)] +struct EcjFsOps { + read_at: fn(&File, &mut [u8], u64) -> io::Result<()>, + rename: fn(&str, &str) -> io::Result<()>, + fsync_dir: fn(&str) -> io::Result<()>, + reopen: fn(&str) -> io::Result, +} + +impl EcjFsOps { + const REAL: EcjFsOps = EcjFsOps { + read_at: read_exact_at, + rename: |from, to| fs::rename(from, to), + fsync_dir: fsync_ecj_dir, + reopen: open_ecj_append, + }; +} + +/// Fsync the directory holding `path`, failing when it cannot be opened. +/// The crate's `fsync_dir` treats an unopenable directory as success, which is +/// best-effort syncing; here the mount carries on after the rename and takes +/// deletes against the replacement, so an unsynced entry could lose them to a +/// power loss. A rename needs only write and search permission on the +/// directory, so it can succeed where opening the directory fails. Matches +/// Go's `util.FsyncDir`, which also skips the sync on Windows. +fn fsync_ecj_dir(path: &str) -> io::Result<()> { + #[cfg(windows)] + { + let _ = path; + Ok(()) + } + #[cfg(not(windows))] + { + let dir = match std::path::Path::new(path).parent() { + Some(p) if !p.as_os_str().is_empty() => p, + _ => std::path::Path::new("."), + }; + File::open(dir)?.sync_all() + } +} + +/// Why publishing a compacted journal failed, split by whether the volume +/// still has a working journal handle. Decided where the failure happens, not +/// inferred afterwards from `ecj_file`. +#[derive(Debug)] +enum EcjPublishError { + /// Nothing was published: the original journal is in place and the handle + /// to it is restored. The mount carries on with the uncompacted file. + Abandoned(io::Error), + /// The volume has no usable journal handle, so it must not mount: deletes + /// would fail, or land in an inode no longer at the journal's path. + HandleLost(io::Error), +} + /// Adds every whole needle id in the first `len` bytes of `ecj_file` to `ids`, /// reading `ECJ_LOAD_CHUNK_BYTES` at a time; a trailing partial record is /// ignored. @@ -55,6 +141,17 @@ pub(crate) fn read_ecj_ids( ecj_file: &File, len: u64, ids: &mut HashSet, +) -> io::Result<()> { + read_ecj_ids_with(read_exact_at, ecj_file, len, ids) +} + +/// [`read_ecj_ids`] with the positional read supplied by the caller, so the +/// mount can load through its injectable [`EcjFsOps`] steps. +fn read_ecj_ids_with( + read_at: fn(&File, &mut [u8], u64) -> io::Result<()>, + ecj_file: &File, + 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; @@ -62,7 +159,7 @@ pub(crate) fn read_ecj_ids( let mut want = std::cmp::min(ECJ_LOAD_CHUNK_BYTES as u64, len - off) as usize; want -= want % NEEDLE_ID_SIZE; // Positional read: the loader's handle is shared with journal appends. - read_exact_at(ecj_file, &mut buf[..want], off)?; + read_at(ecj_file, &mut buf[..want], off)?; for entry in buf[..want].as_chunks::().0 { ids.insert(NeedleId::from_bytes(entry)); } @@ -88,9 +185,13 @@ pub struct EcVolume { ecx_file: Option, ecx_file_size: i64, ecj_file: Option, - /// On-disk size of the .ecj deletion journal. Used only by IO helpers - /// (seek / set_len on partial writes) — the authoritative runtime - /// delete count comes from `deleted_needles.len()`. + /// This volume's registration as a holder of its `.ecj` path. Held for as + /// long as `ecj_file` may be open; see `ecj_registry`. + ecj_hold: Option, + /// On-disk size of the .ecj deletion journal: the rollback point for a + /// failed append, and the size mount-time compaction compares against the + /// id set and re-checks on disk before replacing the file. The runtime + /// delete count comes from `deleted_needles.len()`, not from this. ecj_file_size: i64, /// In-memory set of needle ids that have been deleted since the volume /// was encoded. .ecx is immutable at runtime — it only stores the @@ -424,6 +525,30 @@ pub(crate) fn is_usable_ecx_file(path: &str) -> bool { ecx_file_size(path).is_some_and(|size| size > 0) } +/// Stream sorted deleted ids to `tmp_path` and fsync it. Extracted so failure +/// paths (ENOSPC, read-only dir) are unit-testable without mounting a volume. +fn write_compacted_ecj_tmp(tmp_path: &str, ids: &[NeedleId]) -> io::Result<()> { + // Stream the ids out rather than materialising the whole encoded + // journal: the set and the sorted vector are already resident, and a + // third full-size buffer is a needless contiguous allocation at mount. + // Opened like every other volume file, so the journal this becomes has + // the same mode and open flags as one that was never compacted. + let file = open_volume_file( + OpenOptions::new().write(true).create(true).truncate(true), + tmp_path, + )?; + let mut w = BufWriter::new(file); + let mut buf = [0u8; NEEDLE_ID_SIZE]; + for id in ids { + id.to_bytes(&mut buf); + w.write_all(&buf)?; + } + // into_inner flushes; into_error keeps the flush's own ErrorKind. + w.into_inner() + .map_err(io::IntoInnerError::into_error)? + .sync_all() +} + impl EcVolume { /// Create a new EcVolume. Opens the .ecx index (required) and the .ecj journal. pub fn new( @@ -431,6 +556,18 @@ impl EcVolume { dir_idx: &str, collection: &str, volume_id: VolumeId, + ) -> io::Result { + Self::open_with(dir, dir_idx, collection, volume_id, EcjFsOps::REAL) + } + + /// [`Self::new`] with the journal-compaction filesystem steps supplied, so + /// tests can drive a failing publish through the real mount. + fn open_with( + dir: &str, + dir_idx: &str, + collection: &str, + volume_id: VolumeId, + ecj_ops: EcjFsOps, ) -> io::Result { // One load of the volume's `.vif`, used for both the shard config and // the version / dat-size fields below. @@ -493,6 +630,7 @@ impl EcVolume { ecx_file: None, ecx_file_size: 0, ecj_file: None, + ecj_hold: None, ecj_file_size: 0, deleted_needles: RwLock::new(HashSet::new()), disk_type: DiskType::default(), @@ -569,12 +707,20 @@ impl EcVolume { // Note: Go does NOT replay .ecj into .ecx at volume load (RebuildEcxFile // is only invoked from specific decode/rebuild gRPC handlers), so we // don't either. Tombstones from prior sessions were already written - // in-place in .ecx, and the journal grows monotonically until a - // decode/rebuild operation folds it in. + // in-place in .ecx. The journal only grows at runtime; the one rewrite + // at mount is `maybe_compact_ecj` below, which folds a bloated journal + // down to the set it encodes when no other holder or writer can reach + // the file. let ecj_base = crate::storage::volume::volume_file_name(&vol.ecx_actual_dir, collection, volume_id); let ecj_path = format!("{}.ecj", ecj_base); + // Register as a holder before touching the file: this waits out a + // compaction another disk's volume may be running on the same path, so + // the tail repair and the handle below see the final inode, and it + // stops any compaction from replacing the file under this handle. + vol.ecj_hold = Some(EcjHold::acquire(&ecj_path)); + // Repair a torn tail BEFORE the append handle exists. // // The file is a flat array of fixed-size records and the journal handle @@ -613,25 +759,30 @@ impl EcVolume { } } - let ecj_file = open_volume_file( - OpenOptions::new() - .read(true) - .write(true) - .create(true) - .append(true), - &ecj_path, - )?; + let ecj_file = open_ecj_append(&ecj_path)?; vol.ecj_file_size = ecj_file.metadata()?.len() as i64; vol.ecj_file = Some(ecj_file); // Seed the in-memory deleted set from the journal. - vol.load_deleted_needles_from_ecj()?; + vol.load_deleted_needles_from_ecj(ecj_ops.read_at)?; // Load the generation-0 EC bitrot checksum sidecar. Optional, except // when it contradicts the volume's own geometry — see // load_bitrot_for_generation. vol.load_active_bitrot_sidecar(&[])?; + // Last, once every check that can refuse the mount has passed: a + // volume the server declines to serve keeps its files as they were. + vol.maybe_compact_ecj(ecj_ops).map_err(|e| { + io::Error::new( + e.kind(), + format!( + "ec volume {}: .ecj compaction left no usable journal handle: {}", + volume_id.0, e + ), + ) + })?; + Ok(vol) } @@ -820,7 +971,10 @@ impl EcVolume { /// core with a 31 MB RSS — the set stays small because the ids repeat — /// and never opened its HTTP port, which made the master unregister every /// volume it held. Chunked reads cut that by ~`ECJ_LOAD_CHUNK_BYTES / 8`. - fn load_deleted_needles_from_ecj(&mut self) -> io::Result<()> { + fn load_deleted_needles_from_ecj( + &mut self, + read_at: fn(&File, &mut [u8], u64) -> io::Result<()>, + ) -> io::Result<()> { let ecj_file = match self.ecj_file.as_ref() { Some(f) => f, None => return Ok(()), @@ -833,7 +987,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(ecj_file, self.ecj_file_size as u64, &mut loaded)?; + read_ecj_ids_with(read_at, ecj_file, self.ecj_file_size as u64, &mut loaded)?; let mut set = self .deleted_needles @@ -842,6 +996,216 @@ impl EcVolume { set.extend(loaded); Ok(()) } + + /// Rewrite a bloated `.ecj` from the set just loaded out of it. + /// + /// The journal is semantically a SET of deleted needle ids, but it is + /// written as an append-only log that nothing dedupes or truncates. Several + /// paths append a peer's *entire* journal onto this one — + /// `VolumeEcShardsCopy` (`copy_ecj_file`), EC index recovery, and + /// `ec_decode`'s deliberate merge across holders — so a volume whose shards + /// are repeatedly balanced between two servers grows the file geometrically. + /// Observed in production: 1.51 TB and 1.30 TB for one volume holding ~100 + /// distinct ids, which wedged both owning volume servers at startup. + /// + /// Compaction is safe because the set IS the journal's meaning, provided + /// nothing else can write the file between the load and the rename. That is + /// what the `ecj_registry` reservation and the on-disk re-check establish: + /// this volume is the only holder of the path, no copy is writing to it, and + /// the file is still the inode and size that was loaded. Crash-safe via temp + /// file + fsync + rename + directory fsync. + /// + /// Returns `Err` only when the volume is left without a usable journal + /// handle, which must fail the mount. Every other failure leaves the + /// original journal in place and is logged here. + fn maybe_compact_ecj(&mut self, ops: EcjFsOps) -> io::Result<()> { + let ecj_path = self.ecj_file_name(); + let tmp_path = format!("{}{}", self.idx_base_name(), ECJ_COMPACT_TMP_EXT); + let wanted = self.ecj_needs_compaction()?; + let stale_tmp = std::path::Path::new(&tmp_path).exists(); + if !wanted && !stale_tmp { + return Ok(()); + } + let Some(_reservation) = self + .ecj_hold + .as_ref() + .and_then(EcjHold::try_begin_compaction) + else { + tracing::debug!( + volume_id = self.volume_id.0, + path = %ecj_path, + "skipping .ecj compaction: another holder or a copy can reach the journal", + ); + return Ok(()); + }; + // A tmp left by a crash between its write and the rename. Removed under + // the reservation, so it cannot be another holder's compaction in flight. + if stale_tmp { + let _ = fs::remove_file(&tmp_path); + } + if !wanted { + return Ok(()); + } + self.compact_ecj_reserved(&ecj_path, &tmp_path, ops) + } + + /// The part of [`Self::maybe_compact_ecj`] that runs under the path + /// reservation: re-check the file, write the compacted tmp, publish it. + /// Same error contract: `Err` only when no usable journal handle is left. + fn compact_ecj_reserved( + &mut self, + ecj_path: &str, + tmp_path: &str, + ops: EcjFsOps, + ) -> io::Result<()> { + match self.ecj_unchanged_since_load(ecj_path) { + Ok(true) => {} + Ok(false) => { + tracing::warn!( + volume_id = self.volume_id.0, + path = %ecj_path, + "skipping .ecj compaction: the journal changed on disk after it was loaded", + ); + return Ok(()); + } + Err(e) => { + tracing::warn!(volume_id = self.volume_id.0, error = %e, "compact .ecj: stat journal"); + return Ok(()); + } + } + let ids = self.sorted_deleted_ids()?; + tracing::warn!( + volume_id = self.volume_id.0, + collection = %self.collection, + on_disk_bytes = self.ecj_file_size, + unique_ids = ids.len(), + compacted_bytes = ids.len() * NEEDLE_ID_SIZE, + "compacting bloated .ecj deletion journal", + ); + if let Err(e) = write_compacted_ecj_tmp(tmp_path, &ids) { + // Leaving the partial temp file behind would pin its bytes for the + // life of the process — and the likeliest reason to land here is + // ENOSPC, where those bytes are exactly what is scarce. + let _ = fs::remove_file(tmp_path); + tracing::warn!(volume_id = self.volume_id.0, error = %e, "compact .ecj: write compacted journal"); + return Ok(()); + } + match self.publish_compacted_ecj(ecj_path, tmp_path, ops) { + Ok(()) => Ok(()), + Err(EcjPublishError::Abandoned(e)) => { + tracing::warn!(volume_id = self.volume_id.0, error = %e, "compact .ecj: journal left as it was"); + Ok(()) + } + Err(EcjPublishError::HandleLost(e)) => Err(e), + } + } + + /// Whether the loaded journal is bloated enough to rewrite: + /// `file_bytes >= max(ECJ_COMPACT_MIN_BYTES, ECJ_COMPACT_RATIO * set_bytes)`. + /// The 1 MiB floor keeps a small healthy journal from ever being rewritten. + /// Reads only the set's length, so the common no-op mount copies nothing. + fn ecj_needs_compaction(&self) -> io::Result { + let distinct = self + .deleted_needles + .read() + .map_err(|_| io::Error::other("deleted_needles lock poisoned"))? + .len(); + let compacted_len = (distinct as i64).saturating_mul(NEEDLE_ID_SIZE as i64); + Ok(self.ecj_file_size >= ECJ_COMPACT_MIN_BYTES + && self.ecj_file_size >= compacted_len.saturating_mul(ECJ_COMPACT_RATIO)) + } + + /// The deleted set, sorted so the rewritten file is deterministic: two + /// holders that compact the same set produce byte-identical journals, which + /// makes a difference between them meaningful rather than ordering noise. + fn sorted_deleted_ids(&self) -> io::Result> { + let mut ids: Vec = self + .deleted_needles + .read() + .map_err(|_| io::Error::other("deleted_needles lock poisoned"))? + .iter() + .copied() + .collect(); + ids.sort_unstable(); + Ok(ids) + } + + /// Whether the file at `ecj_path` is still the one this volume loaded: the + /// same inode as the open handle (Unix) and the size the set was read from. + /// A copy that appended, or replaced the file, after the load fails this, + /// and compacting then would drop what it wrote. + fn ecj_unchanged_since_load(&self, ecj_path: &str) -> io::Result { + let Some(held) = self.ecj_file.as_ref() else { + return Ok(false); + }; + let held = held.metadata()?; + let on_path = fs::metadata(ecj_path)?; + #[cfg(unix)] + { + use std::os::unix::fs::MetadataExt; + if held.dev() != on_path.dev() || held.ino() != on_path.ino() { + return Ok(false); + } + } + Ok(held.len() == on_path.len() && on_path.len() as i64 == self.ecj_file_size) + } + + /// Publish a compacted tmp file over the live journal: drop the handle + /// (Windows cannot rename over an open file), rename, fsync the directory, + /// reopen the append handle. + /// + /// A failed rename publishes nothing: the tmp is removed and the handle to + /// the original journal restored, giving [`EcjPublishError::Abandoned`]. If + /// that restore fails too, or anything fails after the rename, the volume + /// has no usable handle: [`EcjPublishError::HandleLost`], carrying every + /// error involved. + fn publish_compacted_ecj( + &mut self, + ecj_path: &str, + tmp_path: &str, + ops: EcjFsOps, + ) -> Result<(), EcjPublishError> { + // Drop our handle on the destination BEFORE replacing it. On Windows a + // rename over a file that is still open fails outright, which would + // make compaction a permanent no-op there; on Unix the handle would + // survive as an unlinked inode, which is worse than useless. + self.ecj_file = None; + if let Err(rename_err) = (ops.rename)(tmp_path, ecj_path) { + let _ = fs::remove_file(tmp_path); + return match (ops.reopen)(ecj_path) { + Ok(f) => { + self.ecj_file = Some(f); + Err(EcjPublishError::Abandoned(rename_err)) + } + Err(reopen_err) => Err(EcjPublishError::HandleLost(io::Error::new( + reopen_err.kind(), + format!( + "rename {} over {}: {}; reopening the original journal then failed: {}", + tmp_path, ecj_path, rename_err, reopen_err + ), + ))), + }; + } + // Syncing the temp file persisted its contents, not the directory entry + // that now points at it. Without this a power loss can restore the old + // bloated journal — and, worse, discard deletes that were acknowledged + // against the replacement in between. Matches the other replace-by-rename + // paths in this crate (ec_decoder, volume_idx_rebuild, volume_idx_repair). + let lost = |what: &str, e: io::Error| { + EcjPublishError::HandleLost(io::Error::new( + e.kind(), + format!("replaced {} but could not {}: {}", ecj_path, what, e), + )) + }; + (ops.fsync_dir)(ecj_path).map_err(|e| lost("fsync its directory", e))?; + // Reopen so later journal_delete appends land in the file just written. + let reopened = (ops.reopen)(ecj_path).map_err(|e| lost("reopen it", e))?; + let size = reopened.metadata().map_err(|e| lost("stat it", e))?.len(); + self.ecj_file_size = size as i64; + self.ecj_file = Some(reopened); + Ok(()) + } + /// Returns (file_count, delete_count) for this EC volume. Mirrors Go's /// `EcVolume.FileAndDeleteCount`: /// @@ -1589,15 +1953,7 @@ impl EcVolume { // in-memory deleted set (all of its contents are now materialized // in .ecx), and reset the cached size. fs::remove_file(&ecj_path)?; - let ecj_file = open_volume_file( - OpenOptions::new() - .read(true) - .write(true) - .create(true) - .append(true), - &ecj_path, - )?; - self.ecj_file = Some(ecj_file); + self.ecj_file = Some(open_ecj_append(&ecj_path)?); self.ecj_file_size = 0; if let Ok(mut set) = self.deleted_needles.write() { set.clear(); @@ -1871,6 +2227,7 @@ impl EcVolume { } self.ecx_file = None; self.ecj_file = None; + self.ecj_hold = None; } pub fn destroy(&mut self) { @@ -1889,6 +2246,7 @@ impl EcVolume { ); let _ = fs::remove_file(format!("{}.ecx", actual_base)); let _ = fs::remove_file(format!("{}.ecj", actual_base)); + let _ = fs::remove_file(format!("{}{}", actual_base, ECJ_COMPACT_TMP_EXT)); let _ = fs::remove_file(format!("{}.vif", actual_base)); // Also sweep the originally-configured idx dir in case stale files // exist there (ecx_file_name() / ecj_file_name() now resolve from @@ -1901,6 +2259,7 @@ impl EcVolume { ); let _ = fs::remove_file(format!("{}.ecx", idx_base)); let _ = fs::remove_file(format!("{}.ecj", idx_base)); + let _ = fs::remove_file(format!("{}{}", idx_base, ECJ_COMPACT_TMP_EXT)); let _ = fs::remove_file(format!("{}.vif", idx_base)); } if self.ecx_actual_dir != self.dir && self.dir_idx != self.dir { @@ -1911,6 +2270,7 @@ impl EcVolume { ); let _ = fs::remove_file(format!("{}.ecx", data_base)); let _ = fs::remove_file(format!("{}.ecj", data_base)); + let _ = fs::remove_file(format!("{}{}", data_base, ECJ_COMPACT_TMP_EXT)); let _ = fs::remove_file(format!("{}.vif", data_base)); } // Go's Destroy() also removes bitrot checksum sidecars so a later @@ -1929,6 +2289,7 @@ impl EcVolume { } self.ecx_file = None; self.ecj_file = None; + self.ecj_hold = None; } } @@ -2723,6 +3084,658 @@ mod tests { assert_eq!(deleted, ids); } + /// A journal of 1M records over 100 distinct ids must mount to those 100 + /// ids and be folded down to 100x8 bytes. The production failure: 1.51 TB + /// of ~100 distinct ids wedging startup. + #[test] + fn test_bloated_ecj_is_compacted_on_load() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + write_ecx_file(dir, "", VolumeId(37), &[]); + + // 100 unique ids x 10000 repeats x 8 B = 8 MiB on disk for 800 B of + // actual information — over both compaction thresholds. + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", VolumeId(37), &ids, 10_000); + + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(dir, "", VolumeId(37)) + ); + let before = std::fs::metadata(&ecj_path).unwrap().len(); + assert_eq!(before, 100 * 10_000 * NEEDLE_ID_SIZE as u64); + + let vol = EcVolume::new(dir, dir, "", VolumeId(37)).unwrap(); + + let deleted = vol.read_deleted_needles().unwrap(); + assert_eq!(deleted, ids, "compaction must preserve the deleted set"); + + let after = std::fs::metadata(&ecj_path).unwrap().len(); + assert_eq!( + after, + ids.len() as u64 * NEEDLE_ID_SIZE as u64, + "journal should have been folded down to one entry per id", + ); + assert!(after < before); + assert!(!std::path::Path::new(&format!("{}.compact.tmp", ecj_path)).exists()); + + drop(vol); + let vol2 = EcVolume::new(dir, dir, "", VolumeId(37)).unwrap(); + assert_eq!(vol2.read_deleted_needles().unwrap(), ids); + assert_eq!(std::fs::metadata(&ecj_path).unwrap().len(), after); + } + + /// Appends must still land in the file after a compaction reopened the + /// handle — otherwise deletes taken after mount would be lost on restart. + #[test] + fn test_journal_delete_after_compaction_persists() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + write_ecx_file( + dir, + "", + VolumeId(38), + &[(NeedleId(7), Offset::from_actual_offset(8), Size(10))], + ); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", VolumeId(38), &ids, 4096); + + let mut vol = EcVolume::new(dir, dir, "", VolumeId(38)).unwrap(); + vol.journal_delete(NeedleId(7)).unwrap(); + drop(vol); + + let vol2 = EcVolume::new(dir, dir, "", VolumeId(38)).unwrap(); + let deleted = vol2.read_deleted_needles().unwrap(); + assert!( + deleted.contains(&NeedleId(7)), + "post-compaction append was lost across remount: {:?}", + &deleted[..deleted.len().min(5)], + ); + assert_eq!(deleted.len(), ids.len() + 1); + } + + /// A healthy journal must be left byte-for-byte alone — compaction is for + /// pathology, not routine churn on every mount. + #[test] + fn test_healthy_ecj_is_not_rewritten() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + write_ecx_file(dir, "", VolumeId(39), &[]); + + // Duplicated 3x, but only 2.4 KB — under ECJ_COMPACT_MIN_BYTES. + let ids: Vec = (1..=100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", VolumeId(39), &ids, 3); + + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(dir, "", VolumeId(39)) + ); + let before_meta = std::fs::metadata(&ecj_path).unwrap(); + let before_bytes = std::fs::read(&ecj_path).unwrap(); + // Use mtime+inode proxy: metadata len + content equality. A rewrite + // would change mtime; content equality proves no-op. + let before_mtime = before_meta.modified().unwrap(); + + let vol = EcVolume::new(dir, dir, "", VolumeId(39)).unwrap(); + + let (_, delete_count) = vol.file_and_delete_count(); + assert_eq!(delete_count, ids.len() as u64); + assert_eq!(vol.read_deleted_needles().unwrap().len(), ids.len() * 3); + + let after_bytes = std::fs::read(&ecj_path).unwrap(); + assert_eq!( + before_bytes, after_bytes, + "small journal must not be rewritten" + ); + let after_mtime = std::fs::metadata(&ecj_path).unwrap().modified().unwrap(); + assert_eq!( + before_mtime, after_mtime, + "small journal mtime must be unchanged" + ); + } + + /// A torn tail plus an oversized journal must be repaired and compacted, + /// with no id lost or invented. + #[test] + fn test_torn_tail_plus_bloated_journal_is_repaired_and_compacted() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + write_ecx_file(dir, "", VolumeId(45), &[]); + + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", VolumeId(45), &ids, 4096); + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(dir, "", VolumeId(45)) + ); + // Append a torn tail. + { + use std::io::Write as _; + let mut f = std::fs::OpenOptions::new() + .append(true) + .open(&ecj_path) + .unwrap(); + f.write_all(&[0xAB, 0xCD, 0xEF]).unwrap(); + f.sync_all().unwrap(); + } + + let vol = EcVolume::new(dir, dir, "", VolumeId(45)).unwrap(); + assert_eq!(vol.read_deleted_needles().unwrap(), ids); + assert_eq!( + std::fs::metadata(&ecj_path).unwrap().len(), + ids.len() as u64 * NEEDLE_ID_SIZE as u64, + "torn tail should be truncated and journal compacted", + ); + } + + /// A journal reached through a shared index directory is compacted when + /// this volume is its only holder: the guard is who holds the path, not + /// which directory it lives in. + #[test] + fn test_shared_index_dir_journal_with_one_holder_is_compacted() { + let tmp = TempDir::new().unwrap(); + let data_dir = tmp.path().join("data"); + let idx_dir = tmp.path().join("idx"); + std::fs::create_dir_all(&data_dir).unwrap(); + std::fs::create_dir_all(&idx_dir).unwrap(); + let (data, idx) = (data_dir.to_str().unwrap(), idx_dir.to_str().unwrap()); + + write_ecx_file(idx, "", VolumeId(44), &[]); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(idx, "", VolumeId(44), &ids, 4096); + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(idx, "", VolumeId(44)) + ); + + let vol = EcVolume::new(data, idx, "", VolumeId(44)).unwrap(); + assert_eq!(vol.file_and_delete_count().1, ids.len() as u64); + assert_eq!( + std::fs::metadata(&ecj_path).unwrap().len(), + (ids.len() * NEEDLE_ID_SIZE) as u64, + ); + drop(vol); + + let vol2 = EcVolume::new(data, idx, "", VolumeId(44)).unwrap(); + assert_eq!(vol2.read_deleted_needles().unwrap(), ids); + } + + /// A shard mount can give disk B a volume whose `.ecj` is disk A's. When + /// A's own volume mounts later, it must not replace the file under B: + /// B's deletes would go to an unlinked inode and vanish at the next mount. + #[test] + fn test_sibling_holder_blocks_compaction() { + let tmp = TempDir::new().unwrap(); + let a_dir = tmp.path().join("a"); + let b_dir = tmp.path().join("b"); + std::fs::create_dir_all(&a_dir).unwrap(); + std::fs::create_dir_all(&b_dir).unwrap(); + let (a, b) = (a_dir.to_str().unwrap(), b_dir.to_str().unwrap()); + let vid = VolumeId(48); + + write_ecx_file( + a, + "", + vid, + &[(NeedleId(7), Offset::from_actual_offset(8), Size(10))], + ); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(a, "", vid, &ids, 1); + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(a, "", vid) + ); + + // B has no index of its own, so it resolves A's and holds A's journal. + let mut on_b = EcVolume::new(b, a, "", vid).unwrap(); + assert_eq!(on_b.ecj_file_name(), ecj_path); + + // The journal bloats (say, a copy appended to it) and A mounts. + write_bloated_ecj(a, "", vid, &ids, 4096); + let bloated = std::fs::metadata(&ecj_path).unwrap().len(); + let on_a = EcVolume::new(a, a, "", vid).unwrap(); + assert_eq!( + std::fs::metadata(&ecj_path).unwrap().len(), + bloated, + "compaction must not replace a journal another volume holds open", + ); + + // B's delete must land in the file at the journal's path. + on_b.journal_delete(NeedleId(7)).unwrap(); + assert_eq!(std::fs::metadata(&ecj_path).unwrap().len(), bloated + 8); + drop(on_b); + drop(on_a); + + // Sole holder now: the next mount compacts and keeps B's delete. + let vol = EcVolume::new(a, a, "", vid).unwrap(); + let deleted = vol.read_deleted_needles().unwrap(); + assert_eq!(deleted.len(), ids.len() + 1); + assert!(deleted.contains(&NeedleId(7))); + } + + /// A copy appending to the journal by path must not have its bytes + /// dropped by a compaction renaming over the file. + #[test] + fn test_active_copy_blocks_compaction() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let vid = VolumeId(49); + write_ecx_file(dir, "", vid, &[]); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", vid, &ids, 4096); + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(dir, "", vid) + ); + let bloated = std::fs::metadata(&ecj_path).unwrap().len(); + + let copy = crate::storage::erasure_coding::ecj_registry::begin_ecj_write(&ecj_path); + let vol = EcVolume::new(dir, dir, "", vid).unwrap(); + assert_eq!(std::fs::metadata(&ecj_path).unwrap().len(), bloated); + drop(vol); + drop(copy); + + let vol = EcVolume::new(dir, dir, "", vid).unwrap(); + assert_eq!( + std::fs::metadata(&ecj_path).unwrap().len(), + (ids.len() * NEEDLE_ID_SIZE) as u64, + ); + assert_eq!(vol.read_deleted_needles().unwrap(), ids); + } + + /// A journal that grew, or was replaced, after it was loaded must not be + /// compacted from the stale set. + #[test] + fn test_journal_changed_after_load_is_detected() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let vid = VolumeId(50); + write_ecx_file(dir, "", vid, &[]); + write_bloated_ecj(dir, "", vid, &[NeedleId(1)], 1); + + let vol = EcVolume::new(dir, dir, "", vid).unwrap(); + let ecj_path = vol.ecj_file_name(); + assert!(vol.ecj_unchanged_since_load(&ecj_path).unwrap()); + + { + let mut f = OpenOptions::new().append(true).open(&ecj_path).unwrap(); + f.write_all(&[0u8; NEEDLE_ID_SIZE]).unwrap(); + } + assert!(!vol.ecj_unchanged_since_load(&ecj_path).unwrap()); + + // Same size, different inode: a copy that replaced the file. + let replacement = format!("{}.replacement", ecj_path); + std::fs::write(&replacement, [0u8; NEEDLE_ID_SIZE]).unwrap(); + std::fs::rename(&replacement, &ecj_path).unwrap(); + #[cfg(unix)] + assert!(!vol.ecj_unchanged_since_load(&ecj_path).unwrap()); + } + + /// A compaction tmp left by a crash between its write and the rename is + /// removed at the next mount, even when the journal needs no compaction. + #[test] + fn test_stale_compaction_tmp_is_removed_at_mount() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let vid = VolumeId(51); + write_ecx_file(dir, "", vid, &[]); + write_bloated_ecj(dir, "", vid, &[NeedleId(1)], 1); + let tmp_path = format!( + "{}{}", + crate::storage::volume::volume_file_name(dir, "", vid), + ECJ_COMPACT_TMP_EXT + ); + std::fs::write(&tmp_path, [0u8; 64]).unwrap(); + + let mut vol = EcVolume::new(dir, dir, "", vid).unwrap(); + assert!(!std::path::Path::new(&tmp_path).exists()); + + // destroy() sweeps it too. + std::fs::write(&tmp_path, [0u8; 64]).unwrap(); + vol.destroy(); + assert!(!std::path::Path::new(&tmp_path).exists()); + } + + /// A mount the bitrot sidecar refuses must leave the journal as it was: + /// compaction runs only after every check that can fail the mount. + #[test] + fn test_refused_mount_does_not_compact() { + use crate::pb::volume_server_pb::{ChecksumAlgorithm, EcBitrotProtection}; + use crate::storage::erasure_coding::ec_bitrot; + use crate::storage::volume::{VifEcShardConfig, VifVolumeInfo}; + + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let vid = VolumeId(52); + let base = crate::storage::volume::volume_file_name(dir, "", vid); + write_ecx_file(dir, "", vid, &[]); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", vid, &ids, 4096); + std::fs::write( + format!("{}.vif", base), + serde_json::to_string(&VifVolumeInfo { + version: 3, + ec_shard_config: Some(VifEcShardConfig { + data_shards: 10, + parity_shards: 4, + ..Default::default() + }), + ..Default::default() + }) + .unwrap(), + ) + .unwrap(); + // A generation-0 sidecar recording a different layout: mount refused. + ec_bitrot::save_bitrot_sidecar( + &ec_bitrot::bitrot_sidecar_path(&base, 0), + &EcBitrotProtection { + algorithm: ChecksumAlgorithm::ChecksumCrc32c as i32, + block_size: ec_bitrot::DEFAULT_BITROT_BLOCK_SIZE as u32, + generation: 0, + ec_shard_config: Some(ec_bitrot::ec_shard_config(12, 4, 0)), + shards: Vec::new(), + encode_uuid: vec![0u8; 16], + }, + ) + .unwrap(); + let ecj_path = format!("{}.ecj", base); + let before = std::fs::read(&ecj_path).unwrap(); + + let err = EcVolume::new(dir, dir, "", vid) + .err() + .expect("a geometry-mismatched sidecar must refuse the mount"); + assert!(err.to_string().contains("refusing to serve"), "{}", err); + assert_eq!(std::fs::read(&ecj_path).unwrap(), before); + } + + /// Compaction must never leave a handle that does not refer to the + /// journal's pathname, and must not leave its temp file behind. + #[test] + fn test_compaction_leaves_a_usable_handle_and_no_temp_file() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + write_ecx_file( + dir, + "", + VolumeId(43), + &[(NeedleId(7), Offset::from_actual_offset(8), Size(10))], + ); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", VolumeId(43), &ids, 4096); + + let mut vol = EcVolume::new(dir, dir, "", VolumeId(43)).unwrap(); + assert!(vol.ecj_file.is_some(), "mount must leave a journal handle"); + + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(dir, "", VolumeId(43)) + ); + assert!(!std::path::Path::new(&format!("{}.compact.tmp", ecj_path)).exists()); + + vol.journal_delete(NeedleId(7)).unwrap(); + let on_disk = std::fs::metadata(&ecj_path).unwrap().len(); + assert_eq!( + on_disk, + ((ids.len() + 1) * NEEDLE_ID_SIZE) as u64, + "post-compaction append did not reach the journal's pathname", + ); + } + + /// Mount a volume `vid` over a bloated journal with the given compaction + /// filesystem steps. Returns the mount result, the journal path, and the + /// journal's bytes before the mount. + fn mount_bloated_with( + dir: &str, + vid: VolumeId, + ops: EcjFsOps, + ) -> (io::Result, String, Vec) { + write_ecx_file( + dir, + "", + vid, + &[(NeedleId(7), Offset::from_actual_offset(8), Size(10))], + ); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", vid, &ids, 4096); + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(dir, "", vid) + ); + let before = std::fs::read(&ecj_path).unwrap(); + ( + EcVolume::open_with(dir, dir, "", vid, ops), + ecj_path, + before, + ) + } + + fn failing_rename(_: &str, _: &str) -> io::Result<()> { + Err(io::Error::new( + io::ErrorKind::PermissionDenied, + "injected rename failure", + )) + } + + fn failing_reopen(_: &str) -> io::Result { + Err(io::Error::new( + io::ErrorKind::NotFound, + "injected reopen failure", + )) + } + + /// A failed rename publishes nothing: the mount succeeds on the original + /// journal, byte for byte, with no tmp left, and deletes still reach the + /// file at the journal's path. + #[test] + fn test_compaction_rename_failure_keeps_original_journal() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ops = EcjFsOps { + rename: failing_rename, + ..EcjFsOps::REAL + }; + let (vol, ecj_path, before) = mount_bloated_with(dir, VolumeId(46), ops); + let mut vol = vol.expect("a failed rename must not fail the mount"); + + assert_eq!(std::fs::read(&ecj_path).unwrap(), before); + assert!(!std::path::Path::new(&format!("{}.compact.tmp", ecj_path)).exists()); + assert_eq!(vol.file_and_delete_count().1, 100); + + vol.journal_delete(NeedleId(7)).unwrap(); + assert_eq!( + std::fs::metadata(&ecj_path).unwrap().len(), + before.len() as u64 + NEEDLE_ID_SIZE as u64, + "the restored handle must append to the original journal", + ); + } + + /// A failed rename whose handle restore also fails leaves no usable + /// handle: the mount fails, reports both errors, and says nothing was + /// compacted. The original journal is untouched. + #[test] + fn test_compaction_rename_and_restore_failure_is_mount_error() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ops = EcjFsOps { + rename: failing_rename, + reopen: failing_reopen, + ..EcjFsOps::REAL + }; + let (vol, ecj_path, before) = mount_bloated_with(dir, VolumeId(53), ops); + let msg = vol + .err() + .expect("mount must fail without a handle") + .to_string(); + + assert!(msg.contains("no usable journal handle"), "{}", msg); + assert!(msg.contains("injected rename failure"), "{}", msg); + assert!(msg.contains("injected reopen failure"), "{}", msg); + assert!(msg.contains("reopening the original journal"), "{}", msg); + assert_eq!(std::fs::read(&ecj_path).unwrap(), before); + assert!(!std::path::Path::new(&format!("{}.compact.tmp", ecj_path)).exists()); + } + + /// A reopen failure after the rename fails the mount: the journal was + /// replaced and the volume has no handle to it. + #[test] + fn test_compaction_post_rename_reopen_failure_is_mount_error() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ops = EcjFsOps { + reopen: failing_reopen, + ..EcjFsOps::REAL + }; + let (vol, ecj_path, _) = mount_bloated_with(dir, VolumeId(47), ops); + let msg = vol + .err() + .expect("mount must fail without a handle") + .to_string(); + + assert!(msg.contains("no usable journal handle"), "{}", msg); + assert!(msg.contains("could not reopen it"), "{}", msg); + assert!(msg.contains("injected reopen failure"), "{}", msg); + assert_eq!( + std::fs::metadata(&ecj_path).unwrap().len(), + 100 * NEEDLE_ID_SIZE as u64, + "the rename had published the compacted journal", + ); + } + + /// A directory fsync failure after the rename fails the mount too: the + /// replacement may not survive a power loss. + #[test] + fn test_compaction_post_rename_fsync_failure_is_mount_error() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ops = EcjFsOps { + fsync_dir: |_| Err(io::Error::other("injected fsync failure")), + ..EcjFsOps::REAL + }; + let (vol, _, _) = mount_bloated_with(dir, VolumeId(54), ops); + let msg = vol.err().expect("mount must fail").to_string(); + assert!(msg.contains("could not fsync its directory"), "{}", msg); + assert!(msg.contains("injected fsync failure"), "{}", msg); + } + + /// The directory sync after the rename must report a directory it cannot + /// open. The crate's best-effort `fsync_dir` returns Ok there, which would + /// let the mount take deletes against an entry that was never synced. + #[test] + fn test_fsync_ecj_dir_fails_when_directory_cannot_be_opened() { + let tmp = TempDir::new().unwrap(); + let missing = tmp.path().join("gone").join("1.ecj"); + let missing = missing.to_str().unwrap(); + #[cfg(not(windows))] + assert!(fsync_ecj_dir(missing).is_err()); + assert!(crate::storage::volume::fsync_dir(missing).is_ok()); + } + + /// The real mount over a directory it may write and search but not read: + /// the rename succeeds and the directory sync cannot open the directory, + /// so the mount fails rather than serving an unsynced replacement. + #[cfg(unix)] + #[test] + fn test_compaction_unreadable_directory_is_mount_error() { + use std::os::unix::fs::PermissionsExt; + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let set_mode = + |mode| std::fs::set_permissions(dir, std::fs::Permissions::from_mode(mode)).unwrap(); + write_ecx_file( + dir, + "", + VolumeId(55), + &[(NeedleId(7), Offset::from_actual_offset(8), Size(10))], + ); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", VolumeId(55), &ids, 4096); + + set_mode(0o300); + if File::open(dir).is_ok() { + // Root opens it regardless; nothing to test. + set_mode(0o755); + return; + } + let result = EcVolume::new(dir, dir, "", VolumeId(55)); + set_mode(0o755); + let msg = result + .err() + .expect("an unsynced replacement must fail the mount") + .to_string(); + assert!(msg.contains("could not fsync its directory"), "{}", msg); + } + + thread_local! { + /// The journal path and replacement bytes `rewrite_after_last_read` + /// writes; taken when it does. + static IN_PLACE_REWRITE: std::cell::RefCell)>> = + const { std::cell::RefCell::new(None) }; + } + + /// A load step that, once the last chunk is read, truncates and refills + /// the journal in place as a registered writer, start to finish, the way a + /// `ReceiveFile` would. + fn rewrite_after_last_read(f: &File, buf: &mut [u8], off: u64) -> io::Result<()> { + read_exact_at(f, buf, off)?; + let pending = IN_PLACE_REWRITE.with(|r| { + let mut r = r.borrow_mut(); + match r.as_ref() { + Some((_, bytes)) if off + buf.len() as u64 >= bytes.len() as u64 => r.take(), + _ => None, + } + }); + if let Some((path, bytes)) = pending { + let _write = crate::storage::erasure_coding::ecj_registry::begin_ecj_write(&path); + let mut out = OpenOptions::new().write(true).truncate(true).open(&path)?; + out.write_all(&bytes)?; + } + Ok(()) + } + + /// A `ReceiveFile` that runs during the mount's load and finishes before + /// compaction can leave the same inode at the same length with different + /// ids, which the inode-and-size re-check accepts. Compacting would then + /// overwrite the received ids with the stale set. + #[test] + fn test_in_place_rewrite_during_load_blocks_compaction() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let vid = VolumeId(56); + write_ecx_file( + dir, + "", + vid, + &[(NeedleId(7), Offset::from_actual_offset(8), Size(10))], + ); + let ids: Vec = (1000..1100).map(NeedleId).collect(); + write_bloated_ecj(dir, "", vid, &ids, 4096); + let ecj_path = format!( + "{}.ecj", + crate::storage::volume::volume_file_name(dir, "", vid) + ); + let mut one = vec![0u8; 100 * NEEDLE_ID_SIZE]; + for (i, entry) in one.chunks_exact_mut(NEEDLE_ID_SIZE).enumerate() { + NeedleId(2000 + i as u64).to_bytes(entry); + } + let rewrite = one.repeat(4096); + IN_PLACE_REWRITE.with(|r| *r.borrow_mut() = Some((ecj_path.clone(), rewrite.clone()))); + + let ops = EcjFsOps { + read_at: rewrite_after_last_read, + ..EcjFsOps::REAL + }; + let vol = EcVolume::open_with(dir, dir, "", vid, ops).unwrap(); + assert!( + IN_PLACE_REWRITE.with(|r| r.borrow().is_none()), + "the rewrite never ran" + ); + assert!( + std::fs::read(&ecj_path).unwrap() == rewrite, + "the received journal must not be compacted from the set loaded before it", + ); + drop(vol); + } + #[test] fn test_journal_delete_wrong_cookie() { let tmp = TempDir::new().unwrap(); diff --git a/seaweed-volume/src/storage/erasure_coding/ecj_registry.rs b/seaweed-volume/src/storage/erasure_coding/ecj_registry.rs new file mode 100644 index 000000000..426623875 --- /dev/null +++ b/seaweed-volume/src/storage/erasure_coding/ecj_registry.rs @@ -0,0 +1,326 @@ +//! Process-wide coordination of everything that touches one `.ecj` path. +//! +//! Mount-time compaction replaces a deletion journal with a new inode. That is +//! only safe while nothing else in this process can write the old one: +//! +//! - **Holders** are mounted `EcVolume`s with an append handle on the path. A +//! store holds one `EcVolume` per disk location, and a shard mount or a +//! cross-disk reconcile can point one disk's volume at another disk's +//! `.ecj`, so several holders of one path are normal. A holder that keeps +//! appending to a replaced inode acknowledges deletes that are gone at the +//! next mount. +//! - **Writers** append to or replace the path by name without holding it +//! open across calls: `ReceiveFile` of an EC `.ecj`, and the unmounted +//! append in `merge_ec_journal`, which `VolumeEcShardsCopy` and EC index +//! recovery funnel a peer's journal through. Bytes they write after the +//! compactor sized the journal would be dropped by the rename. +//! +//! Compaction therefore runs only while its caller is the sole holder and no +//! writer is active, and while it runs no holder may open the path and no +//! writer may start. Both wait instead; a compaction rewrites only the distinct +//! id set, so the wait is short. +//! +//! No writer active at the reservation is not enough: one that ran while the +//! holder loaded the journal, or after, and has finished may have rewritten it +//! in place to the same length (`ReceiveFile` truncates and refills), which +//! the inode-and-size re-check cannot see. So each writer bumps the path's +//! write generation as it starts, and a holder may compact only if no writer +//! was active when it registered and the generation has not moved since. +//! +//! Paths are keyed by their canonical parent directory, so two disk locations +//! that spell one directory differently still meet here. + +use std::collections::HashMap; +use std::path::{Path, PathBuf}; +use std::sync::{Condvar, LazyLock, Mutex, MutexGuard}; + +#[derive(Default)] +struct PathState { + holders: usize, + writers: usize, + compacting: bool, + /// Writers that have started on the path. Lives as long as the entry, + /// which a registered holder keeps. + write_gen: u64, +} + +impl PathState { + fn idle(&self) -> bool { + self.holders == 0 && self.writers == 0 && !self.compacting + } +} + +struct Registry { + paths: Mutex>, + changed: Condvar, +} + +static REGISTRY: LazyLock = LazyLock::new(|| Registry { + paths: Mutex::new(HashMap::new()), + changed: Condvar::new(), +}); + +fn lock() -> MutexGuard<'static, HashMap> { + // The critical sections only adjust counters and cannot panic midway, so + // a poisoned lock still guards consistent state. + REGISTRY.paths.lock().unwrap_or_else(|e| e.into_inner()) +} + +/// Canonical key for `path`: its resolved parent directory joined with the +/// file name. The file itself may not exist yet (a copy creates it), so only +/// the directory is resolved. +fn key_for(path: &str) -> PathBuf { + let p = Path::new(path); + let (Some(parent), Some(name)) = (p.parent(), p.file_name()) else { + return std::path::absolute(p).unwrap_or_else(|_| p.to_path_buf()); + }; + let parent = if parent.as_os_str().is_empty() { + Path::new(".") + } else { + parent + }; + let dir = std::fs::canonicalize(parent) + .or_else(|_| std::path::absolute(parent)) + .unwrap_or_else(|_| parent.to_path_buf()); + dir.join(name) +} + +/// Block until no compaction is running on `key`, then apply `f` to its state. +fn update_when_not_compacting(key: &Path, f: impl FnOnce(&mut PathState) -> R) -> R { + let mut paths = lock(); + while paths.get(key).is_some_and(|s| s.compacting) { + paths = REGISTRY + .changed + .wait(paths) + .unwrap_or_else(|e| e.into_inner()); + } + f(paths.entry(key.to_path_buf()).or_default()) +} + +fn release(key: &Path, f: impl FnOnce(&mut PathState)) { + let mut paths = lock(); + if let Some(state) = paths.get_mut(key) { + f(state); + if state.idle() { + paths.remove(key); + } + } + drop(paths); + REGISTRY.changed.notify_all(); +} + +/// A mounted `EcVolume`'s registration as a holder of its `.ecj`. Taken before +/// the journal is opened and released when dropped. +pub(crate) struct EcjHold { + key: PathBuf, + /// The path's write generation when the hold was taken, and whether a + /// writer was active then. Taken before the journal is opened and loaded, + /// so they cover every write the load might have missed. + write_gen: u64, + writer_at_start: bool, +} + +impl EcjHold { + /// Register as a holder of `ecj_path`, first waiting out any compaction in + /// progress so the handle opened afterwards is on the final inode. + pub(crate) fn acquire(ecj_path: &str) -> Self { + let key = key_for(ecj_path); + let (write_gen, writer_at_start) = update_when_not_compacting(&key, |s| { + s.holders += 1; + (s.write_gen, s.writers > 0) + }); + EcjHold { + key, + write_gen, + writer_at_start, + } + } + + /// Reserve the path for a compaction, or `None` when another holder or an + /// active writer could still reach the current inode, or when a writer + /// has run on the path since the hold was taken, so the journal may no + /// longer be what the holder loaded. + pub(crate) fn try_begin_compaction(&self) -> Option { + let mut paths = lock(); + let state = paths.get_mut(&self.key)?; + if state.holders != 1 || state.writers != 0 || state.compacting { + return None; + } + if self.writer_at_start || state.write_gen != self.write_gen { + return None; + } + state.compacting = true; + Some(EcjCompaction { + key: self.key.clone(), + }) + } +} + +impl Drop for EcjHold { + fn drop(&mut self) { + release(&self.key, |s| s.holders = s.holders.saturating_sub(1)); + } +} + +/// An exclusive reservation of a `.ecj` path for compaction. Holders and +/// writers wait until it is dropped. +pub(crate) struct EcjCompaction { + key: PathBuf, +} + +impl Drop for EcjCompaction { + fn drop(&mut self) { + release(&self.key, |s| s.compacting = false); + } +} + +/// An out-of-band writer (shard copy, index recovery, `ReceiveFile`) on a +/// `.ecj` path. +/// Compaction does not start while one is alive. +pub(crate) struct EcjWrite { + key: PathBuf, +} + +impl Drop for EcjWrite { + fn drop(&mut self) { + release(&self.key, |s| s.writers = s.writers.saturating_sub(1)); + } +} + +/// Register as a writer of `ecj_path`, waiting out any compaction in progress. +/// Blocks; async callers use [`begin_ecj_write_async`]. +pub(crate) fn begin_ecj_write(ecj_path: &str) -> EcjWrite { + let key = key_for(ecj_path); + update_when_not_compacting(&key, |s| { + s.writers += 1; + s.write_gen += 1; + }); + EcjWrite { key } +} + +/// [`begin_ecj_write`] for async handlers: the wait runs on the blocking pool +/// so a compaction in progress never stalls a runtime worker. +pub(crate) async fn begin_ecj_write_async(ecj_path: &str) -> EcjWrite { + let path = ecj_path.to_string(); + match tokio::task::spawn_blocking(move || begin_ecj_write(&path)).await { + Ok(write) => write, + // Only a panic inside the registry lands here, and it leaves no count + // behind; registering inline is still correct, merely blocking. + Err(_) => begin_ecj_write(ecj_path), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::mpsc; + use std::time::Duration; + use tempfile::TempDir; + + fn ecj(dir: &TempDir) -> String { + dir.path().join("1.ecj").to_str().unwrap().to_string() + } + + #[test] + fn sole_holder_may_compact() { + let dir = TempDir::new().unwrap(); + let hold = EcjHold::acquire(&ecj(&dir)); + assert!(hold.try_begin_compaction().is_some()); + } + + #[test] + fn second_holder_blocks_compaction() { + let dir = TempDir::new().unwrap(); + let a = EcjHold::acquire(&ecj(&dir)); + let b = EcjHold::acquire(&ecj(&dir)); + assert!(a.try_begin_compaction().is_none()); + drop(b); + assert!(a.try_begin_compaction().is_some()); + } + + #[test] + fn active_writer_blocks_compaction() { + let dir = TempDir::new().unwrap(); + let hold = EcjHold::acquire(&ecj(&dir)); + let w = begin_ecj_write(&ecj(&dir)); + assert!(hold.try_begin_compaction().is_none()); + drop(w); + // This hold loaded before the write; a later one may compact. + drop(hold); + let hold = EcjHold::acquire(&ecj(&dir)); + assert!(hold.try_begin_compaction().is_some()); + } + + /// A writer that ran after the hold was taken, or was already running + /// then, may have changed the journal the holder loaded, even though it + /// has finished by the time compaction asks. + #[test] + fn finished_writer_since_hold_blocks_compaction() { + let dir = TempDir::new().unwrap(); + let path = ecj(&dir); + + let hold = EcjHold::acquire(&path); + drop(begin_ecj_write(&path)); + assert!( + hold.try_begin_compaction().is_none(), + "a writer that started after the hold", + ); + drop(hold); + + let w = begin_ecj_write(&path); + let hold = EcjHold::acquire(&path); + drop(w); + assert!( + hold.try_begin_compaction().is_none(), + "a writer active when the hold was taken", + ); + drop(hold); + + let hold = EcjHold::acquire(&path); + assert!( + hold.try_begin_compaction().is_some(), + "no writer since the hold", + ); + } + + #[test] + fn differently_spelled_paths_share_one_key() { + let dir = TempDir::new().unwrap(); + std::fs::create_dir(dir.path().join("sub")).unwrap(); + let plain = ecj(&dir); + let dotted = dir + .path() + .join("sub") + .join("..") + .join("1.ecj") + .to_str() + .unwrap() + .to_string(); + let a = EcjHold::acquire(&plain); + let _b = EcjHold::acquire(&dotted); + assert!(a.try_begin_compaction().is_none()); + } + + #[test] + fn writer_waits_for_compaction_to_finish() { + let dir = TempDir::new().unwrap(); + let path = ecj(&dir); + let hold = EcjHold::acquire(&path); + let compaction = hold.try_begin_compaction().unwrap(); + + let (tx, rx) = mpsc::channel(); + let p = path.clone(); + let t = std::thread::spawn(move || { + let _w = begin_ecj_write(&p); + tx.send(()).unwrap(); + }); + assert!( + rx.recv_timeout(Duration::from_millis(100)).is_err(), + "a writer must not start while a compaction holds the path", + ); + drop(compaction); + rx.recv_timeout(Duration::from_secs(5)) + .expect("writer must proceed once the compaction ends"); + t.join().unwrap(); + } +} diff --git a/seaweed-volume/src/storage/erasure_coding/mod.rs b/seaweed-volume/src/storage/erasure_coding/mod.rs index 678b83961..4ed7ab004 100644 --- a/seaweed-volume/src/storage/erasure_coding/mod.rs +++ b/seaweed-volume/src/storage/erasure_coding/mod.rs @@ -10,6 +10,7 @@ pub mod ec_locate; pub mod ec_shard; pub mod ec_volume; pub mod ecj_merge; +pub(crate) mod ecj_registry; pub use ec_shard::{ DATA_SHARDS_COUNT, EcVolumeShard, MAX_SHARD_COUNT, MIN_TOTAL_DISKS, PARITY_SHARDS_COUNT, diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index b7d09c9d6..0d9332392 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -13,7 +13,9 @@ use crate::config::MinFreeSpace; use crate::pb::master_pb; use crate::storage::disk_location::DiskLocation; use crate::storage::erasure_coding::ec_shard::{EcVolumeShard, MAX_SHARD_COUNT, ShardId}; -use crate::storage::erasure_coding::ec_volume::{EcVolume, is_usable_ecx_file}; +use crate::storage::erasure_coding::ec_volume::{ + ECJ_COMPACT_TMP_EXT, EcVolume, is_usable_ecx_file, +}; use crate::storage::needle::needle::{Needle, get_actual_size}; use crate::storage::needle_map::NeedleMapKind; use crate::storage::super_block::{ReplicaPlacement, SUPER_BLOCK_SIZE}; @@ -1434,10 +1436,12 @@ impl Store { crate::storage::volume::volume_file_name(&loc.directory, collection, vid); let _ = std::fs::remove_file(format!("{}.ecx", idx_base)); let _ = std::fs::remove_file(format!("{}.ecj", idx_base)); + let _ = std::fs::remove_file(format!("{}{}", idx_base, ECJ_COMPACT_TMP_EXT)); // Also try data directory in case .ecx/.ecj were created before -dir.idx if loc.idx_directory != loc.directory { let _ = std::fs::remove_file(format!("{}.ecx", data_base)); let _ = std::fs::remove_file(format!("{}.ecj", data_base)); + let _ = std::fs::remove_file(format!("{}{}", data_base, ECJ_COMPACT_TMP_EXT)); } // A shard-only disk also drops its stale .vif (Go // removeEcSharedIndexFiles): a live .idx means this disk still diff --git a/seaweed-volume/src/storage/store_ec_journal.rs b/seaweed-volume/src/storage/store_ec_journal.rs index d31cd135f..6912a5dfa 100644 --- a/seaweed-volume/src/storage/store_ec_journal.rs +++ b/seaweed-volume/src/storage/store_ec_journal.rs @@ -83,6 +83,11 @@ fn merge_ec_journal_with( publish_to_journal_siblings(&store, primary, vid, &journal_path, ids); return Ok(added); } + // The path append registers as a writer so a mount compacting this + // journal cannot swap its inode underneath it (the write itself is + // already serialized with mounts by the store lock). + let _ecj_write = + crate::storage::erasure_coding::ecj_registry::begin_ecj_write(ecj_path); if let Some(added) = append_ecj_ids(ecj_path, &local, ids, size)? { return Ok(added); } diff --git a/weed/server/volume_grpc_copy.go b/weed/server/volume_grpc_copy.go index fdf24ab95..c48ab8abd 100644 --- a/weed/server/volume_grpc_copy.go +++ b/weed/server/volume_grpc_copy.go @@ -710,11 +710,17 @@ func (vs *VolumeServer) ReceiveFile(stream volume_server_pb.VolumeServer_Receive var targetFile *os.File var filePath string var bytesWritten uint64 + // Set while an EC .ecj is being received; released only after the file is + // closed and any partial copy removed. + var ecjWriteDone func() defer func() { if targetFile != nil { targetFile.Close() } + if ecjWriteDone != nil { + ecjWriteDone() + } }() for { @@ -803,6 +809,13 @@ func (vs *VolumeServer) ReceiveFile(stream volume_server_pb.VolumeServer_Receive // Create EC shard file path baseFileName := erasure_coding.EcShardBaseFileName(fileInfo.Collection, int(fileInfo.VolumeId)) filePath = util.Join(targetLocation.Directory, baseFileName+fileInfo.Ext) + // The mounted check above runs once; a volume can still mount on + // this journal while the stream writes it. Registered as a writer + // before the file is created, that mount cannot compact the + // journal and leave the rest of the stream in an unlinked inode. + if fileInfo.Ext == ".ecj" && ecjWriteDone == nil { + ecjWriteDone = erasure_coding.BeginEcjWrite(filePath) + } } else { // Regular volume file v := vs.store.GetVolume(needle.VolumeId(fileInfo.VolumeId)) diff --git a/weed/server/volume_grpc_copy_traversal_test.go b/weed/server/volume_grpc_copy_traversal_test.go index 89ada05e6..73fa4db58 100644 --- a/weed/server/volume_grpc_copy_traversal_test.go +++ b/weed/server/volume_grpc_copy_traversal_test.go @@ -30,15 +30,20 @@ func newTraversalTestStore(dir string) *storage.Store { } // fakeReceiveFileStream scripts a ReceiveFile request sequence and records the -// final response returned via SendAndClose. +// final response returned via SendAndClose. onRecv, when set, runs before +// request i is handed over, i.e. after the handler processed requests 0..i-1. type fakeReceiveFileStream struct { grpc.ServerStream - reqs []*volume_server_pb.ReceiveFileRequest - index int - resp *volume_server_pb.ReceiveFileResponse + reqs []*volume_server_pb.ReceiveFileRequest + index int + resp *volume_server_pb.ReceiveFileResponse + onRecv func(i int) } func (s *fakeReceiveFileStream) Recv() (*volume_server_pb.ReceiveFileRequest, error) { + if s.onRecv != nil { + s.onRecv(s.index) + } if s.index >= len(s.reqs) { return nil, io.EOF } diff --git a/weed/server/volume_grpc_erasure_coding.go b/weed/server/volume_grpc_erasure_coding.go index aef7162a9..aaee5c8af 100644 --- a/weed/server/volume_grpc_erasure_coding.go +++ b/weed/server/volume_grpc_erasure_coding.go @@ -735,13 +735,13 @@ func removeEcSharedIndexFiles(bName string, location *storage.DiskLocation, hasE dataBaseFilename := path.Join(location.Directory, bName) if hasEcxFile { // .ecx/.ecj may be in either dir depending on when -dir.idx was configured. - for _, p := range []string{indexBaseFilename + ".ecx", indexBaseFilename + ".ecj"} { + for _, p := range []string{indexBaseFilename + ".ecx", indexBaseFilename + ".ecj", indexBaseFilename + erasure_coding.EcjCompactTmpExt} { if err := removeFileIfExists(p); err != nil { return err } } if location.IdxDirectory != location.Directory { - for _, p := range []string{dataBaseFilename + ".ecx", dataBaseFilename + ".ecj"} { + for _, p := range []string{dataBaseFilename + ".ecx", dataBaseFilename + ".ecj", dataBaseFilename + erasure_coding.EcjCompactTmpExt} { if err := removeFileIfExists(p); err != nil { return err } @@ -805,10 +805,12 @@ func removeStaleEcArtifacts(dataBaseFileName, indexBaseFileName string, total in // .ecx/.ecj/.ecsum may sit in either dir depending on -dir.idx; clear both. record(removeFileIfExists(indexBaseFileName + ".ecx")) record(removeFileIfExists(indexBaseFileName + ".ecj")) + record(removeFileIfExists(indexBaseFileName + erasure_coding.EcjCompactTmpExt)) record(removeBitrotSidecars(indexBaseFileName)) if dataBaseFileName != indexBaseFileName { record(removeFileIfExists(dataBaseFileName + ".ecx")) record(removeFileIfExists(dataBaseFileName + ".ecj")) + record(removeFileIfExists(dataBaseFileName + erasure_coding.EcjCompactTmpExt)) record(removeBitrotSidecars(dataBaseFileName)) } diff --git a/weed/server/volume_grpc_receive_file_ecj_test.go b/weed/server/volume_grpc_receive_file_ecj_test.go new file mode 100644 index 000000000..bbcfba018 --- /dev/null +++ b/weed/server/volume_grpc_receive_file_ecj_test.go @@ -0,0 +1,88 @@ +package weed_server + +import ( + "bytes" + "os" + "path/filepath" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb" + "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" + "github.com/seaweedfs/seaweedfs/weed/storage/types" +) + +// A volume that mounts on a journal ReceiveFile is still writing must not +// compact it: the rest of the stream would land in the replaced inode and be +// gone at the next mount. +func TestReceiveFile_EcjStreamBlocksMountCompaction(t *testing.T) { + storeDir := t.TempDir() + // The store first, so its startup scan finds no volume to mount and + // ReceiveFile's mounted check lets the .ecj through. + vs := &VolumeServer{store: newTraversalTestStore(storeDir)} + + base := filepath.Join(storeDir, "4") + ecx := make([]byte, types.NeedleMapEntrySize) + types.NeedleIdToBytes(ecx[0:types.NeedleIdSize], 7) + types.OffsetToBytes(ecx[types.NeedleIdSize:types.NeedleIdSize+types.OffsetSize], types.ToOffset(8)) + types.SizeToBytes(ecx[types.NeedleIdSize+types.OffsetSize:], types.Size(10)) + if err := os.WriteFile(base+".ecx", ecx, 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(base+".vif", nil, 0o644); err != nil { + t.Fatal(err) + } + + // A bloated prefix (100 ids, 4096 times over) the mount would compact, + // then an id only the second chunk carries. + one := make([]byte, 100*types.NeedleIdSize) + for i := 0; i < 100; i++ { + types.NeedleIdToBytes(one[i*types.NeedleIdSize:], types.NeedleId(1000+i)) + } + bloated := bytes.Repeat(one, 4096) + tail := make([]byte, types.NeedleIdSize) + types.NeedleIdToBytes(tail, 5000) + + var mounted *erasure_coding.EcVolume + stream := &fakeReceiveFileStream{ + reqs: []*volume_server_pb.ReceiveFileRequest{ + infoReq(&volume_server_pb.ReceiveFileInfo{ + VolumeId: 4, + Ext: ".ecj", + IsEcVolume: true, + FileSize: uint64(len(bloated) + len(tail)), + }), + contentReq(bloated), + contentReq(tail), + }, + onRecv: func(i int) { + if i != 2 { + return + } + // The first chunk is on disk and the second is not yet sent. + ev, err := erasure_coding.NewEcVolume(types.HardDriveType, storeDir, storeDir, "", 4) + if err != nil { + t.Fatalf("mount during the stream: %v", err) + } + mounted = ev + }, + } + + if err := vs.ReceiveFile(stream); err != nil { + t.Fatalf("ReceiveFile: %v", err) + } + if stream.resp == nil || stream.resp.Error != "" { + t.Fatalf("ReceiveFile rejected the journal: %+v", stream.resp) + } + if mounted == nil { + t.Fatal("the mount hook never ran") + } + mounted.Close() + + got, err := os.ReadFile(base + ".ecj") + if err != nil { + t.Fatal(err) + } + if want := append(bloated, tail...); !bytes.Equal(got, want) { + t.Fatalf("journal on disk is %d bytes, want the %d the stream sent", len(got), len(want)) + } +} diff --git a/weed/storage/disk_location_ec.go b/weed/storage/disk_location_ec.go index dd3b0f4c7..be98e371d 100644 --- a/weed/storage/disk_location_ec.go +++ b/weed/storage/disk_location_ec.go @@ -639,10 +639,12 @@ func (l *DiskLocation) removeEcVolumeFiles(collection string, vid needle.VolumeI // EC loading for incomplete/missing shards on next startup removeFile(indexBaseFileName+".ecx", "EC index file") removeFile(indexBaseFileName+".ecj", "EC journal file") + removeFile(indexBaseFileName+erasure_coding.EcjCompactTmpExt, "EC journal compaction tmp") // Also try the data directory in case .ecx/.ecj were created before -dir.idx was configured if l.IdxDirectory != l.Directory { removeFile(baseFileName+".ecx", "EC index file (fallback)") removeFile(baseFileName+".ecj", "EC journal file (fallback)") + removeFile(baseFileName+erasure_coding.EcjCompactTmpExt, "EC journal compaction tmp (fallback)") } // Remove all EC shard files (.ec00 ~ .ec31) from data directory diff --git a/weed/storage/erasure_coding/ec_volume.go b/weed/storage/erasure_coding/ec_volume.go index 698b10431..7b6dd286f 100644 --- a/weed/storage/erasure_coding/ec_volume.go +++ b/weed/storage/erasure_coding/ec_volume.go @@ -1,9 +1,11 @@ package erasure_coding import ( + "bufio" "errors" "fmt" "os" + "path/filepath" "slices" "sync" "syscall" @@ -19,6 +21,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/storage/needle" "github.com/seaweedfs/seaweedfs/weed/storage/types" "github.com/seaweedfs/seaweedfs/weed/storage/volume_info" + "github.com/seaweedfs/seaweedfs/weed/util" ) var ( @@ -31,6 +34,21 @@ var ( // which keeps a bloated journal from spending mount in per-entry reads. const ecjLoadChunkBytes = 1 << 20 +// A .ecj smaller than this is never rewritten, however redundant. Below a +// megabyte the duplication costs nothing and the rewrite is pure churn. +const ecjCompactMinBytes = 1 << 20 + +// Rewrite only when the journal is at least this many times larger than the +// set it encodes. A journal holds one entry per delete, so a healthy one is +// close to 1x; four times the set means most of the file is repeats. +const ecjCompactRatio = 4 + +// EcjCompactTmpExt names the staging file for a compacted journal, next to the +// .ecj it replaces. Listed with the other EC index files wherever those are +// removed, so a tmp left by a crash between write and rename does not outlive +// its volume. +const EcjCompactTmpExt = ".ecj.compact.tmp" + type EcVolume struct { VolumeId needle.VolumeId Collection string @@ -51,19 +69,23 @@ type EcVolume struct { Version needle.Version ecjFile *os.File ecjFileAccessLock sync.Mutex - diskType types.DiskType - datFileSize int64 - ExpireAtSec uint64 //ec volume destroy time, calculated from the ec volume was created - ECContext *ECContext // EC encoding parameters + // ecjHold registers this volume as a holder of its .ecj path for as long + // as ecjFile may be open; see ecj_registry.go. + ecjHold *ecjHold + diskType types.DiskType + datFileSize int64 + ExpireAtSec uint64 //ec volume destroy time, calculated from the ec volume was created + ECContext *ECContext // EC encoding parameters // EncodeTsNs is the encode time (unix nanos) loaded from .vif; reads carry it // so a shard from a different encode run is rejected. 0 for pre-upgrade volumes. EncodeTsNs int64 // ecjFileSize mirrors the on-disk size of the .ecj deletion journal and - // is maintained under ecjFileAccessLock. It is only used by IO helpers - // (seek/truncate) — the authoritative runtime delete count comes from - // deletedNeedles. + // is maintained under ecjFileAccessLock: the write offset for appends, and + // the size mount-time compaction compares against the id set and re-checks + // on disk before replacing the file. The runtime delete count comes from + // deletedNeedles, not from this. ecjFileSize int64 // deletedNeedles is the in-memory set of needle ids that have been @@ -150,6 +172,13 @@ func statEcxSize(path string) (int64, error) { } func NewEcVolume(diskType types.DiskType, dir string, dirIdx string, collection string, vid needle.VolumeId) (ev *EcVolume, err error) { + return newEcVolumeWith(diskType, dir, dirIdx, collection, vid, realEcjFsOps) +} + +// newEcVolumeWith is NewEcVolume with the journal-compaction filesystem steps +// supplied, so tests can drive a failing load or publish through the real +// mount. +func newEcVolumeWith(diskType types.DiskType, dir string, dirIdx string, collection string, vid needle.VolumeId, ecjOps ecjFsOps) (ev *EcVolume, err error) { ev = &EcVolume{dir: dir, dirIdx: dirIdx, Collection: collection, VolumeId: vid, diskType: diskType} dataBaseFileName := EcShardFileName(collection, dir, int(vid)) @@ -204,7 +233,14 @@ func NewEcVolume(diskType types.DiskType, dir string, dirIdx string, collection ev.ecxCreatedAt = ecxFi.ModTime() // open ecj file and seed the in-memory deleted set from it. - if ev.ecjFile, err = backend.OpenVolumeFile(indexBaseFileName+".ecj", os.O_RDWR|os.O_CREATE); err != nil { + // + // Register as a holder first: this waits out a compaction another disk's + // volume may be running on the same path, so the handle below is on the + // final inode, and it stops any compaction from replacing the file under + // this handle. + ev.ecjHold = acquireEcjHold(indexBaseFileName + ".ecj") + if ev.ecjFile, err = openEcjFile(indexBaseFileName + ".ecj"); err != nil { + ev.Close() return nil, fmt.Errorf("cannot open ec volume journal %s.ecj: %v", indexBaseFileName, err) } if ecjFi, statErr := ev.ecjFile.Stat(); statErr == nil { @@ -219,15 +255,18 @@ func NewEcVolume(diskType types.DiskType, dir string, dirIdx string, collection whole := ev.ecjFileSize - ragged glog.Warningf("ec volume %d: truncating torn .ecj tail %d -> %d bytes", vid, ev.ecjFileSize, whole) if truncErr := ev.ecjFile.Truncate(whole); truncErr != nil { + ev.Close() return nil, fmt.Errorf("ec volume %d: repair torn .ecj tail: %w", vid, truncErr) } if syncErr := ev.ecjFile.Sync(); syncErr != nil { + ev.Close() return nil, fmt.Errorf("ec volume %d: sync .ecj after tail repair: %w", vid, syncErr) } ev.ecjFileSize = whole } ev.deletedNeedles = make(map[types.NeedleId]struct{}) - if loadErr := ev.loadDeletedNeedlesFromEcj(); loadErr != nil { + loadErr := ev.loadDeletedNeedlesFromEcj(ecjOps.readAt) + if loadErr != nil { glog.Warningf("ec volume %d: load deleted needles from .ecj: %v", vid, loadErr) } @@ -368,6 +407,14 @@ func NewEcVolume(diskType types.DiskType, dir string, dirIdx string, collection return nil, err } + // Fold a bloated journal back down to the set it encodes. Last, once every + // check that can refuse the mount has passed, so a volume the server + // declines to serve keeps its files as they were. + if compactErr := ev.compactEcjAfterLoad(loadErr, ecjOps); compactErr != nil { + ev.Close() + return nil, fmt.Errorf("ec volume %d: .ecj compaction left no usable journal handle: %w", vid, compactErr) + } + return } @@ -435,6 +482,10 @@ func (ev *EcVolume) Close() { _ = ev.ecjFile.Close() ev.ecjFile = nil } + if ev.ecjHold != nil { + ev.ecjHold.release() + ev.ecjHold = nil + } ev.ecjFileAccessLock.Unlock() if ev.ecxFile != nil { _ = ev.ecxFile.Sync() @@ -478,6 +529,7 @@ func (ev *EcVolume) Destroy() { for _, base := range ev.ecIndexBaseNames() { os.Remove(base + ".ecx") os.Remove(base + ".ecj") + os.Remove(base + EcjCompactTmpExt) } // The .vif is shared with a coexisting normal volume (e.g. mid-decode), so // only remove the active copy, not both. @@ -651,7 +703,7 @@ func (ev *EcVolume) markNeedleDeletedInMemory(needleId types.NeedleId) { // loadDeletedNeedlesFromEcj walks the .ecj journal and populates the // in-memory deleted set. Called once from NewEcVolume under the exclusive // ownership of the just-constructed (and not yet shared) EcVolume. -func (ev *EcVolume) loadDeletedNeedlesFromEcj() error { +func (ev *EcVolume) loadDeletedNeedlesFromEcj(readAt func(f *os.File, b []byte, off int64) (int, error)) error { if ev.ecjFile == nil || ev.ecjFileSize < int64(types.NeedleIdSize) { return nil } @@ -662,7 +714,7 @@ func (ev *EcVolume) loadDeletedNeedlesFromEcj() error { if want == 0 { break } - if _, err := ev.ecjFile.ReadAt(buf[:want], off); err != nil { + if _, err := readAt(ev.ecjFile, buf[:want], off); err != nil { return fmt.Errorf("read ecj at %d: %w", off, err) } for i := int64(0); i+int64(types.NeedleIdSize) <= want; i += int64(types.NeedleIdSize) { @@ -673,6 +725,246 @@ func (ev *EcVolume) loadDeletedNeedlesFromEcj() error { return nil } +// openEcjFile opens path as the deletion journal handle, creating it if +// absent. Both the mount and the reopen after compaction go through here, so +// the handle a compacted volume appends through behaves like the original. +func openEcjFile(path string) (*os.File, error) { + return backend.OpenVolumeFile(path, os.O_RDWR|os.O_CREATE) +} + +// ecjFsOps are the filesystem steps of a mount's journal compaction: the +// reads that load the set it compacts from, and the steps that publish the +// compacted file. Production uses realEcjFsOps; tests substitute failing steps +// to cover the failure paths through the real mount. +type ecjFsOps struct { + readAt func(f *os.File, b []byte, off int64) (int, error) + rename func(oldpath, newpath string) error + fsyncDir func(path string) error + reopen func(path string) (*os.File, error) +} + +var realEcjFsOps = ecjFsOps{ + readAt: (*os.File).ReadAt, + rename: os.Rename, + fsyncDir: func(p string) error { return util.FsyncDir(filepath.Dir(p)) }, + reopen: openEcjFile, +} + +// ecjHandleLostError marks a compaction failure that left the volume without a +// usable journal handle, so the mount must fail: deletes would error, or land +// in an inode no longer at the journal's path. Decided where the failure +// happens rather than inferred afterwards from ecjFile. +type ecjHandleLostError struct{ err error } + +func (e *ecjHandleLostError) Error() string { return e.err.Error() } +func (e *ecjHandleLostError) Unwrap() error { return e.err } + +// compactEcjAfterLoad compacts the journal unless its load failed. After a +// failed load the set holds only the part of the journal read before the +// error, and rewriting the file from it would delete the rest for good. +func (ev *EcVolume) compactEcjAfterLoad(loadErr error, ops ecjFsOps) error { + if loadErr != nil { + glog.Warningf("ec volume %d: not compacting .ecj: its load failed, so the in-memory set may be partial", ev.VolumeId) + return nil + } + return ev.maybeCompactEcj(ops) +} + +// maybeCompactEcj rewrites a bloated .ecj from the set just loaded out of it. +// +// The journal is semantically a SET of deleted needle ids, written as an +// append-only log that nothing dedupes. VolumeEcShardsCopy, EC index recovery +// and ec_decode's merge append a peer's whole journal onto this one, so a +// volume whose shards are balanced back and forth grows the file +// geometrically (1.51 TB for ~100 distinct ids in production). +// +// Compaction is safe because the set IS the journal's meaning, provided +// nothing else can write the file between the load and the rename. The +// ecj_registry reservation and the on-disk re-check establish that: this +// volume is the only holder of the path, no copy is writing to it, and the file +// is still the inode and size that was loaded. +// +// Returns an error only when the volume is left without a usable journal +// handle, which must fail the mount. Every other failure leaves the original +// journal in place and is logged here. +func (ev *EcVolume) maybeCompactEcj(ops ecjFsOps) error { + ecjPath := ev.FileName(".ecj") + tmpPath := EcShardFileName(ev.Collection, ev.ecxActualDir, int(ev.VolumeId)) + EcjCompactTmpExt + wanted := ev.ecjNeedsCompaction() + _, tmpStatErr := os.Stat(tmpPath) + staleTmp := tmpStatErr == nil + if (!wanted && !staleTmp) || ev.ecjHold == nil { + return nil + } + end, ok := ev.ecjHold.tryBeginCompaction() + if !ok { + glog.V(1).Infof("ec volume %d: skipping .ecj compaction: another holder or a copy can reach %s", ev.VolumeId, ecjPath) + return nil + } + defer end() + // A tmp left by a crash between its write and the rename. Removed under the + // reservation, so it cannot be another holder's compaction in flight. + if staleTmp { + _ = os.Remove(tmpPath) + } + if !wanted { + return nil + } + return ev.compactEcjReserved(ecjPath, tmpPath, ops) +} + +// compactEcjReserved is the part of maybeCompactEcj that runs under the path +// reservation: re-check the file, write the compacted tmp, publish it. Same +// error contract: an error only when no usable journal handle is left. +func (ev *EcVolume) compactEcjReserved(ecjPath, tmpPath string, ops ecjFsOps) error { + unchanged, err := ev.ecjUnchangedSinceLoad(ecjPath) + if err != nil { + glog.Warningf("ec volume %d: compact .ecj: stat journal: %v", ev.VolumeId, err) + return nil + } + if !unchanged { + glog.Warningf("ec volume %d: skipping .ecj compaction: %s changed on disk after it was loaded", ev.VolumeId, ecjPath) + return nil + } + ids := ev.sortedDeletedIds() + glog.Warningf("ec volume %d: compacting bloated .ecj deletion journal on-disk=%d unique=%d compacted=%d", + ev.VolumeId, ev.ecjFileSize, len(ids), len(ids)*types.NeedleIdSize) + if err := writeCompactedEcjTmp(tmpPath, ids); err != nil { + // A partial tmp would pin its bytes, and the likeliest cause here is + // ENOSPC, where those bytes are exactly what is scarce. + _ = os.Remove(tmpPath) + glog.Warningf("ec volume %d: compact .ecj: write compacted journal: %v", ev.VolumeId, err) + return nil + } + if err := ev.publishCompactedEcj(ecjPath, tmpPath, ops); err != nil { + var lost *ecjHandleLostError + if errors.As(err, &lost) { + return err + } + glog.Warningf("ec volume %d: compact .ecj: journal left as it was: %v", ev.VolumeId, err) + } + return nil +} + +// ecjNeedsCompaction reports whether the loaded journal is bloated enough to +// rewrite: file_bytes >= max(ecjCompactMinBytes, ecjCompactRatio * set_bytes). +// The 1 MiB floor keeps a small healthy journal from ever being rewritten. It +// reads only the set's length, so the common no-op mount copies nothing. +func (ev *EcVolume) ecjNeedsCompaction() bool { + ev.deletedNeedlesLock.RLock() + distinct := int64(len(ev.deletedNeedles)) + ev.deletedNeedlesLock.RUnlock() + compactedLen := distinct * int64(types.NeedleIdSize) + return ev.ecjFileSize >= ecjCompactMinBytes && ev.ecjFileSize >= compactedLen*ecjCompactRatio +} + +// sortedDeletedIds returns the deleted set sorted, so the rewritten file is +// deterministic and two holders compacting one set write identical bytes. +func (ev *EcVolume) sortedDeletedIds() []types.NeedleId { + ev.deletedNeedlesLock.RLock() + ids := make([]types.NeedleId, 0, len(ev.deletedNeedles)) + for id := range ev.deletedNeedles { + ids = append(ids, id) + } + ev.deletedNeedlesLock.RUnlock() + slices.Sort(ids) + return ids +} + +// ecjUnchangedSinceLoad reports whether the file at ecjPath is still the one +// this volume loaded: the same file as the open handle and the size the set +// was read from. A copy that appended, or replaced the file, after the load +// fails this, and compacting then would drop what it wrote. +func (ev *EcVolume) ecjUnchangedSinceLoad(ecjPath string) (bool, error) { + if ev.ecjFile == nil { + return false, nil + } + held, err := ev.ecjFile.Stat() + if err != nil { + return false, err + } + onPath, err := os.Stat(ecjPath) + if err != nil { + return false, err + } + return os.SameFile(held, onPath) && held.Size() == onPath.Size() && onPath.Size() == ev.ecjFileSize, nil +} + +// writeCompactedEcjTmp writes sorted ids to tmpPath through a buffer and +// fsyncs it. Opened like every other volume file, so the journal it becomes +// has the same mode and open flags as one that was never compacted. +func writeCompactedEcjTmp(tmpPath string, ids []types.NeedleId) error { + f, err := backend.OpenVolumeFile(tmpPath, os.O_WRONLY|os.O_CREATE|os.O_TRUNC) + if err != nil { + return err + } + w := bufio.NewWriterSize(f, ecjLoadChunkBytes) + var rec [types.NeedleIdSize]byte + for _, id := range ids { + types.NeedleIdToBytes(rec[:], id) + if _, err := w.Write(rec[:]); err != nil { + _ = f.Close() + return err + } + } + if err := w.Flush(); err != nil { + _ = f.Close() + return err + } + if err := f.Sync(); err != nil { + _ = f.Close() + return err + } + return f.Close() +} + +// publishCompactedEcj replaces the live journal with the compacted tmp file: +// drop the handle (Windows cannot rename over an open file), rename, fsync the +// directory, reopen the handle. +// +// A failed rename publishes nothing: the tmp is removed and the handle to the +// original journal restored, and the rename error is returned as is. If that +// restore fails too, or anything fails after the rename, the volume has no +// usable handle and the error is an *ecjHandleLostError carrying every error +// involved. +func (ev *EcVolume) publishCompactedEcj(ecjPath, tmpPath string, ops ecjFsOps) error { + if ev.ecjFile != nil { + _ = ev.ecjFile.Close() + ev.ecjFile = nil + } + if renameErr := ops.rename(tmpPath, ecjPath); renameErr != nil { + _ = os.Remove(tmpPath) + reopened, reopenErr := ops.reopen(ecjPath) + if reopenErr != nil { + return &ecjHandleLostError{fmt.Errorf("rename %s over %s: %w; reopening the original journal then failed: %w", + tmpPath, ecjPath, renameErr, reopenErr)} + } + ev.ecjFile = reopened + return renameErr + } + lost := func(what string, err error) error { + return &ecjHandleLostError{fmt.Errorf("replaced %s but could not %s: %w", ecjPath, what, err)} + } + // The tmp sync persisted contents, not the directory entry. Without this a + // power loss can restore the old journal and discard deletes acknowledged + // against the replacement. + if err := ops.fsyncDir(ecjPath); err != nil { + return lost("fsync its directory", err) + } + reopened, err := ops.reopen(ecjPath) + if err != nil { + return lost("reopen it", err) + } + fi, err := reopened.Stat() + if err != nil { + _ = reopened.Close() + return lost("stat it", err) + } + ev.ecjFile = reopened + ev.ecjFileSize = fi.Size() + return nil +} + func (ev *EcVolume) LocateEcShardNeedle(needleId types.NeedleId, version needle.Version) (offset types.Offset, size types.Size, intervals []Interval, err error) { // find the needle from ecx file diff --git a/weed/storage/erasure_coding/ec_volume_ecj_compact_test.go b/weed/storage/erasure_coding/ec_volume_ecj_compact_test.go new file mode 100644 index 000000000..5b233c0a6 --- /dev/null +++ b/weed/storage/erasure_coding/ec_volume_ecj_compact_test.go @@ -0,0 +1,415 @@ +package erasure_coding + +import ( + "errors" + "os" + "path/filepath" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/storage/types" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// compactTestIds are the 100 distinct ids the bloated journals below repeat. +func compactTestIds() []types.NeedleId { + ids := make([]types.NeedleId, 0, 100) + for i := 0; i < 100; i++ { + ids = append(ids, types.NeedleId(1000+i)) + } + return ids +} + +func compactEcjBytes(ids []types.NeedleId, repeats int) []byte { + one := make([]byte, len(ids)*types.NeedleIdSize) + for i, id := range ids { + types.NeedleIdToBytes(one[i*types.NeedleIdSize:], id) + } + b := make([]byte, 0, len(one)*repeats) + for i := 0; i < repeats; i++ { + b = append(b, one...) + } + return b +} + +// compactEcxEntry is one .ecx entry for a live needle, so it can be deleted. +func compactEcxEntry(id types.NeedleId) []byte { + b := make([]byte, types.NeedleMapEntrySize) + types.NeedleIdToBytes(b[0:types.NeedleIdSize], id) + types.OffsetToBytes(b[types.NeedleIdSize:types.NeedleIdSize+types.OffsetSize], types.ToOffset(8)) + types.SizeToBytes(b[types.NeedleIdSize+types.OffsetSize:], types.Size(10)) + return b +} + +// writeCompactTestVolume lays down .ecx (holding needle 7), .ecj and an empty +// .vif for volume vid in dir, returning the base file name. +func writeCompactTestVolume(t *testing.T, dir string, vid needle.VolumeId, ecj []byte) string { + t.Helper() + base := EcShardFileName("", dir, int(vid)) + require.NoError(t, os.WriteFile(base+".ecx", compactEcxEntry(7), 0644)) + require.NoError(t, os.WriteFile(base+".ecj", ecj, 0644)) + require.NoError(t, os.WriteFile(base+".vif", []byte{}, 0644)) + return base +} + +func fileSize(t *testing.T, path string) int64 { + t.Helper() + fi, err := os.Stat(path) + require.NoError(t, err) + return fi.Size() +} + +var errInjectedRename = errors.New("injected rename failure") +var errInjectedReopen = errors.New("injected reopen failure") + +func failingRename(string, string) error { return errInjectedRename } +func failingReopen(string) (*os.File, error) { return nil, errInjectedReopen } +func failingFsync(string) error { return errors.New("injected fsync failure") } +func withEcjOps(f func(*ecjFsOps)) ecjFsOps { ops := realEcjFsOps; f(&ops); return ops } +func mountCompactTest(dir string, vid needle.VolumeId, ops ecjFsOps) (*EcVolume, error) { + return newEcVolumeWith("hdd", dir, dir, "", vid, ops) +} + +// A load error leaves only part of the journal in the set. Compacting from it +// would delete the unread records from disk for good. The read fails inside the +// real mount, after the first chunk, so the mount goes on with a set missing +// the id that only the rest of the journal holds. +func TestEcjNotCompactedAfterLoadError(t *testing.T) { + dir := t.TempDir() + journal := append(compactEcjBytes(compactTestIds(), 4096), compactEcjBytes([]types.NeedleId{5000}, 1)...) + base := writeCompactTestVolume(t, dir, 60, journal) + + ops := withEcjOps(func(o *ecjFsOps) { + o.readAt = func(f *os.File, b []byte, off int64) (int, error) { + if off >= ecjLoadChunkBytes { + return 0, errors.New("injected read failure") + } + return f.ReadAt(b, off) + } + }) + ev, err := mountCompactTest(dir, 60, ops) + require.NoError(t, err) + require.False(t, ev.IsNeedleDeleted(5000), "the failed read must have left the set partial") + ev.Close() + + got, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, journal, got, "a partial set must never be written over the journal") + + // A full load compacts, and keeps the id only the tail held. + ev, err = mountCompactTest(dir, 60, realEcjFsOps) + require.NoError(t, err) + defer ev.Close() + assert.Equal(t, int64(101*types.NeedleIdSize), fileSize(t, base+".ecj")) + assert.True(t, ev.IsNeedleDeleted(5000)) +} + +// A shard mount can give disk B a volume whose .ecj is disk A's. When A's own +// volume mounts later it must not replace the file under B, whose deletes +// would go to an unlinked inode and vanish at the next mount. +func TestEcjSiblingHolderBlocksCompaction(t *testing.T) { + root := t.TempDir() + a, b := filepath.Join(root, "a"), filepath.Join(root, "b") + require.NoError(t, os.MkdirAll(a, 0755)) + require.NoError(t, os.MkdirAll(b, 0755)) + ids := compactTestIds() + base := writeCompactTestVolume(t, a, 61, compactEcjBytes(ids, 1)) + + // B has no index of its own, so it resolves A's and holds A's journal. + onB, err := NewEcVolume("hdd", b, a, "", 61) + require.NoError(t, err) + require.Equal(t, base+".ecj", onB.FileName(".ecj")) + + // The journal bloats (say, a copy appended to it) and A mounts. + require.NoError(t, os.WriteFile(base+".ecj", compactEcjBytes(ids, 4096), 0644)) + bloated := fileSize(t, base+".ecj") + onA, err := NewEcVolume("hdd", a, a, "", 61) + require.NoError(t, err) + assert.Equal(t, bloated, fileSize(t, base+".ecj"), "compaction must not replace a journal another volume holds open") + + // B's delete must land in the file at the journal's path. + require.NoError(t, onB.DeleteNeedleFromEcx(7)) + assert.Equal(t, bloated+types.NeedleIdSize, fileSize(t, base+".ecj")) + onB.Close() + onA.Close() + + // Sole holder now: the next mount compacts and keeps B's delete. + ev, err := NewEcVolume("hdd", a, a, "", 61) + require.NoError(t, err) + defer ev.Close() + assert.Equal(t, int64((len(ids)+1)*types.NeedleIdSize), fileSize(t, base+".ecj")) + assert.True(t, ev.IsNeedleDeleted(7)) +} + +// A copy appending to the journal by path must not have its bytes dropped by +// a compaction renaming over the file. +func TestEcjActiveCopyBlocksCompaction(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 62, compactEcjBytes(compactTestIds(), 4096)) + bloated := fileSize(t, base+".ecj") + + done := BeginEcjWrite(base + ".ecj") + ev, err := NewEcVolume("hdd", dir, dir, "", 62) + require.NoError(t, err) + assert.Equal(t, bloated, fileSize(t, base+".ecj")) + ev.Close() + done() + + ev, err = NewEcVolume("hdd", dir, dir, "", 62) + require.NoError(t, err) + defer ev.Close() + assert.Equal(t, int64(100*types.NeedleIdSize), fileSize(t, base+".ecj")) +} + +// A copy that starts while a compaction holds the path waits for it to end, +// then appends to the compacted file. +func TestEcjWriterWaitsForCompaction(t *testing.T) { + path := filepath.Join(t.TempDir(), "1.ecj") + hold := acquireEcjHold(path) + defer hold.release() + end, ok := hold.tryBeginCompaction() + require.True(t, ok) + + started := make(chan struct{}) + go func() { + done := BeginEcjWrite(path) + close(started) + done() + }() + select { + case <-started: + t.Fatal("a writer must not start while a compaction holds the path") + case <-time.After(100 * time.Millisecond): + } + end() + select { + case <-started: + case <-time.After(5 * time.Second): + t.Fatal("writer must proceed once the compaction ends") + } +} + +// A writer that ran after the hold was taken, or was already running then, +// may have changed the journal the holder loaded, even though it has finished +// by the time compaction asks. +func TestEcjFinishedWriterSinceHoldBlocksCompaction(t *testing.T) { + path := filepath.Join(t.TempDir(), "1.ecj") + // End a wrongly granted reservation, or the next writer waits on it forever. + refused := func(hold *ecjHold) bool { + end, ok := hold.tryBeginCompaction() + if ok { + end() + } + return !ok + } + + hold := acquireEcjHold(path) + BeginEcjWrite(path)() + assert.True(t, refused(hold), "a writer that started after the hold") + hold.release() + + done := BeginEcjWrite(path) + hold = acquireEcjHold(path) + done() + assert.True(t, refused(hold), "a writer active when the hold was taken") + hold.release() + + hold = acquireEcjHold(path) + defer hold.release() + end, ok := hold.tryBeginCompaction() + require.True(t, ok, "no writer since the hold") + end() +} + +// A ReceiveFile truncates the journal and refills it. One that runs during the +// mount's load and finishes before compaction can leave the same inode at the +// same length with different ids, which the inode-and-size re-check accepts; +// compacting would then overwrite the received ids with the stale set. +func TestEcjInPlaceRewriteDuringLoadBlocksCompaction(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 66, compactEcjBytes(compactTestIds(), 4096)) + received := make([]types.NeedleId, 0, 100) + for i := 0; i < 100; i++ { + received = append(received, types.NeedleId(2000+i)) + } + rewrite := compactEcjBytes(received, 4096) + + rewritten := false + ops := withEcjOps(func(o *ecjFsOps) { + o.readAt = func(f *os.File, b []byte, off int64) (int, error) { + n, err := f.ReadAt(b, off) + if !rewritten && off+int64(len(b)) >= int64(len(rewrite)) { + // The last chunk is read: a stream truncates and refills the + // journal in place, start to finish, before compaction. + rewritten = true + done := BeginEcjWrite(base + ".ecj") + w, openErr := os.OpenFile(base+".ecj", os.O_WRONLY|os.O_TRUNC, 0644) + require.NoError(t, openErr) + _, writeErr := w.Write(rewrite) + require.NoError(t, writeErr) + require.NoError(t, w.Close()) + done() + } + return n, err + } + }) + ev, err := mountCompactTest(dir, 66, ops) + require.NoError(t, err) + defer ev.Close() + require.True(t, rewritten) + + got, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, rewrite, got, "the received journal must not be compacted from the set loaded before it") +} + +// Two spellings of one directory must meet in the registry. +func TestEcjRegistryKeysResolveDirectories(t *testing.T) { + dir := t.TempDir() + require.NoError(t, os.Mkdir(filepath.Join(dir, "sub"), 0755)) + a := acquireEcjHold(filepath.Join(dir, "1.ecj")) + defer a.release() + b := acquireEcjHold(filepath.Join(dir, "sub", "..", "1.ecj")) + _, ok := a.tryBeginCompaction() + assert.False(t, ok) + b.release() + end, ok := a.tryBeginCompaction() + require.True(t, ok) + end() +} + +// A journal that grew, or was replaced, after it was loaded must not be +// compacted from the stale set. +func TestEcjChangedAfterLoadIsDetected(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 63, compactEcjBytes([]types.NeedleId{1}, 1)) + ev, err := NewEcVolume("hdd", dir, dir, "", 63) + require.NoError(t, err) + defer ev.Close() + + unchanged, err := ev.ecjUnchangedSinceLoad(base + ".ecj") + require.NoError(t, err) + assert.True(t, unchanged) + + f, err := os.OpenFile(base+".ecj", os.O_APPEND|os.O_WRONLY, 0644) + require.NoError(t, err) + _, err = f.Write(make([]byte, types.NeedleIdSize)) + require.NoError(t, err) + require.NoError(t, f.Close()) + unchanged, err = ev.ecjUnchangedSinceLoad(base + ".ecj") + require.NoError(t, err) + assert.False(t, unchanged) + + // Same size as loaded, different file: a copy that replaced the journal. + require.NoError(t, os.WriteFile(base+".replacement", make([]byte, types.NeedleIdSize), 0644)) + require.NoError(t, os.Rename(base+".replacement", base+".ecj")) + unchanged, err = ev.ecjUnchangedSinceLoad(base + ".ecj") + require.NoError(t, err) + assert.False(t, unchanged) +} + +// A failed rename publishes nothing: the mount succeeds on the original +// journal, byte for byte, with no tmp left, and deletes still reach the file +// at the journal's path. +func TestEcjRenameFailureKeepsOriginalJournal(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 64, compactEcjBytes(compactTestIds(), 4096)) + before, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + + ev, err := mountCompactTest(dir, 64, withEcjOps(func(o *ecjFsOps) { o.rename = failingRename })) + require.NoError(t, err, "a failed rename must not fail the mount") + defer ev.Close() + + after, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, before, after) + assert.NoFileExists(t, base+EcjCompactTmpExt) + + require.NoError(t, ev.DeleteNeedleFromEcx(7)) + assert.Equal(t, int64(len(before)+types.NeedleIdSize), fileSize(t, base+".ecj"), + "the restored handle must append to the original journal") +} + +// A failed rename whose handle restore also fails leaves no usable handle: +// the mount fails, reports both errors, and the journal is untouched. +func TestEcjRenameAndRestoreFailureIsMountError(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 65, compactEcjBytes(compactTestIds(), 4096)) + before, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + + _, err = mountCompactTest(dir, 65, withEcjOps(func(o *ecjFsOps) { + o.rename = failingRename + o.reopen = failingReopen + })) + require.Error(t, err) + assert.ErrorIs(t, err, errInjectedRename) + assert.ErrorIs(t, err, errInjectedReopen) + assert.Contains(t, err.Error(), "no usable journal handle") + assert.Contains(t, err.Error(), "reopening the original journal") + + after, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, before, after) + assert.NoFileExists(t, base+EcjCompactTmpExt) +} + +// A reopen failure after the rename fails the mount: the journal was replaced +// and the volume has no handle to it. +func TestEcjPostRenameReopenFailureIsMountError(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 66, compactEcjBytes(compactTestIds(), 4096)) + + _, err := mountCompactTest(dir, 66, withEcjOps(func(o *ecjFsOps) { o.reopen = failingReopen })) + require.Error(t, err) + assert.ErrorIs(t, err, errInjectedReopen) + assert.Contains(t, err.Error(), "could not reopen it") + assert.Equal(t, int64(100*types.NeedleIdSize), fileSize(t, base+".ecj"), "the rename had published the compacted journal") +} + +// A directory fsync failure after the rename fails the mount too. +func TestEcjPostRenameFsyncFailureIsMountError(t *testing.T) { + dir := t.TempDir() + writeCompactTestVolume(t, dir, 67, compactEcjBytes(compactTestIds(), 4096)) + + _, err := mountCompactTest(dir, 67, withEcjOps(func(o *ecjFsOps) { o.fsyncDir = failingFsync })) + require.Error(t, err) + assert.Contains(t, err.Error(), "could not fsync its directory") +} + +// A compaction tmp left by a crash between its write and the rename is removed +// at the next mount, even when the journal needs no compaction, and by +// Destroy. +func TestEcjStaleCompactionTmpRemoved(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 68, compactEcjBytes([]types.NeedleId{1}, 1)) + require.NoError(t, os.WriteFile(base+EcjCompactTmpExt, make([]byte, 64), 0644)) + + ev, err := NewEcVolume("hdd", dir, dir, "", 68) + require.NoError(t, err) + assert.NoFileExists(t, base+EcjCompactTmpExt) + + require.NoError(t, os.WriteFile(base+EcjCompactTmpExt, make([]byte, 64), 0644)) + ev.Destroy() + assert.NoFileExists(t, base+EcjCompactTmpExt) +} + +// A mount refused by a later check must leave the journal as it was: +// compaction runs only after every check that can fail the mount. +func TestEcjRefusedMountDoesNotCompact(t *testing.T) { + dir := t.TempDir() + base := writeCompactTestVolume(t, dir, 69, compactEcjBytes(compactTestIds(), 4096)) + require.NoError(t, os.WriteFile(base+".vif", []byte("{not a volume info"), 0644)) + before, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + + _, err = NewEcVolume("hdd", dir, dir, "", 69) + require.Error(t, err, "a malformed .vif must refuse the mount") + + after, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, before, after) +} diff --git a/weed/storage/erasure_coding/ec_volume_ecj_test.go b/weed/storage/erasure_coding/ec_volume_ecj_test.go index 7bcdb0757..d7643803a 100644 --- a/weed/storage/erasure_coding/ec_volume_ecj_test.go +++ b/weed/storage/erasure_coding/ec_volume_ecj_test.go @@ -3,6 +3,7 @@ package erasure_coding_test import ( "os" "testing" + "time" erasure_coding "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" "github.com/seaweedfs/seaweedfs/weed/storage/types" @@ -84,3 +85,83 @@ func TestEcjLoadsAcrossChunkBoundary(t *testing.T) { } assert.False(t, ev.IsNeedleDeleted(count+1)) } + +// A journal of 1M records over 100 distinct ids must mount to those 100 ids +// and be folded down to 100*8 bytes. Mirrors the Rust bloated-compaction test +// and the production 1.51 TB failure. +func TestEcjBloatedJournalCompactedOnMount(t *testing.T) { + dir := t.TempDir() + + const distinct = 100 + const repeats = 10_000 // 1M records = 8 MiB + ids := make([]types.NeedleId, 0, distinct) + for i := 0; i < distinct; i++ { + ids = append(ids, types.NeedleId(1000+i)) + } + ecj := make([]byte, 0, distinct*repeats*types.NeedleIdSize) + one := ecjBytes(ids...) + for i := 0; i < repeats; i++ { + ecj = append(ecj, one...) + } + + ev, base := mountEcVolume(t, dir, nil, ecj) + for _, id := range ids { + assert.True(t, ev.IsNeedleDeleted(id), "id %d", id) + } + ev.Close() + + fi, err := os.Stat(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, int64(distinct*types.NeedleIdSize), fi.Size(), "journal should have been folded down to one entry per id") + + // No temp file left behind, and remount is stable. + _, err = os.Stat(base + ".ecj.compact.tmp") + assert.True(t, os.IsNotExist(err)) + ev2, _ := mountEcVolume(t, dir, nil, nil) + defer ev2.Close() + for _, id := range ids { + assert.True(t, ev2.IsNeedleDeleted(id), "id %d", id) + } + fi2, err := os.Stat(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, int64(distinct*types.NeedleIdSize), fi2.Size()) +} + +// A healthy small journal must never be rewritten, however redundant. +func TestEcjHealthySmallJournalNotRewritten(t *testing.T) { + dir := t.TempDir() + + ids := make([]types.NeedleId, 0, 100) + for i := 1; i <= 100; i++ { + ids = append(ids, types.NeedleId(i)) + } + // Duplicated 3x, but only 2.4 KB — under ecjCompactMinBytes. + ecj := append(append(ecjBytes(ids...), ecjBytes(ids...)...), ecjBytes(ids...)...) + + // Lay the files down and take the baseline BEFORE mounting, so a rewrite + // during the mount shows up as a difference. The mtime is pushed into the + // past so a rewrite within the filesystem's timestamp granularity still + // changes it. + base := erasure_coding.EcShardFileName("", dir, 7) + require.NoError(t, os.WriteFile(base+".ecx", nil, 0644)) + require.NoError(t, os.WriteFile(base+".ecj", ecj, 0644)) + require.NoError(t, os.WriteFile(base+".vif", []byte{}, 0644)) + past := time.Now().Add(-time.Hour).Truncate(time.Second) + require.NoError(t, os.Chtimes(base+".ecj", past, past)) + before, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + + ev, err := erasure_coding.NewEcVolume("hdd", dir, dir, "", 7) + require.NoError(t, err) + for _, id := range ids { + assert.True(t, ev.IsNeedleDeleted(id), "id %d", id) + } + ev.Close() + + after, err := os.ReadFile(base + ".ecj") + require.NoError(t, err) + assert.Equal(t, before, after, "small journal must not be rewritten") + afterFi, err := os.Stat(base + ".ecj") + require.NoError(t, err) + assert.True(t, afterFi.ModTime().Equal(past), "small journal mtime must be unchanged, got %v", afterFi.ModTime()) +} diff --git a/weed/storage/erasure_coding/ecj_registry.go b/weed/storage/erasure_coding/ecj_registry.go new file mode 100644 index 000000000..f676aeccd --- /dev/null +++ b/weed/storage/erasure_coding/ecj_registry.go @@ -0,0 +1,183 @@ +package erasure_coding + +import ( + "path/filepath" + "sync" +) + +// Process-wide coordination of everything that touches one .ecj path. +// +// Mount-time compaction replaces a deletion journal with a new inode. That is +// only safe while nothing else in this process can write the old one: +// +// - Holders are mounted EcVolumes with a handle on the path. A store keeps one +// EcVolume per disk location, and a shard mount or a cross-disk reconcile +// can point one disk's volume at another disk's .ecj, so several holders of +// one path are normal. A holder that keeps appending to a replaced inode +// acknowledges deletes that are gone at the next mount. +// - Writers append to or replace the path by name without holding it open +// across calls: ReceiveFile of an EC .ecj, and the unmounted append in +// Store.MergeEcJournal, which VolumeEcShardsCopy and EC index recovery +// funnel a peer's journal through. +// Bytes they write after the compactor sized the journal would be dropped +// by the rename. +// +// Compaction therefore runs only while its caller is the sole holder and no +// writer is active, and while it runs no holder may open the path and no writer +// may start. Both wait instead; a compaction rewrites only the distinct id set, +// so the wait is short. +// +// No writer active at the reservation is not enough: one that ran while the +// holder loaded the journal, or after, and has finished may have rewritten it +// in place to the same length (ReceiveFile truncates and refills), which the +// inode-and-size re-check cannot see. So each writer bumps the path's write +// generation as it starts, and a holder may compact only if no writer was +// active when it registered and the generation has not moved since. +// +// Paths are keyed by their resolved parent directory, so two disk locations +// that spell one directory differently still meet here. + +type ecjPathState struct { + holders int + writers int + compacting bool + // writeGen counts writers that have started on the path. It survives as + // long as the entry does, and a registered holder keeps the entry. + writeGen uint64 +} + +var ecjPaths = struct { + sync.Mutex + changed *sync.Cond + state map[string]*ecjPathState +}{state: map[string]*ecjPathState{}} + +func init() { + ecjPaths.changed = sync.NewCond(&ecjPaths.Mutex) +} + +// ecjPathKey resolves the parent directory of path. The file itself may not +// exist yet (a copy creates it), so only the directory is resolved. +func ecjPathKey(path string) string { + dir, name := filepath.Split(path) + if dir == "" { + dir = "." + } + if abs, err := filepath.Abs(dir); err == nil { + dir = abs + } + if resolved, err := filepath.EvalSymlinks(dir); err == nil { + dir = resolved + } + return filepath.Join(dir, name) +} + +// ecjUpdateWhenNotCompacting blocks until no compaction is running on key, +// then applies f to its state. +func ecjUpdateWhenNotCompacting(key string, f func(*ecjPathState)) { + ecjPaths.Lock() + defer ecjPaths.Unlock() + for { + st := ecjPaths.state[key] + if st == nil { + st = &ecjPathState{} + ecjPaths.state[key] = st + } + if !st.compacting { + f(st) + return + } + ecjPaths.changed.Wait() + } +} + +func ecjRelease(key string, f func(*ecjPathState)) { + ecjPaths.Lock() + if st := ecjPaths.state[key]; st != nil { + f(st) + if st.holders == 0 && st.writers == 0 && !st.compacting { + delete(ecjPaths.state, key) + } + } + ecjPaths.Unlock() + ecjPaths.changed.Broadcast() +} + +// ecjHold is a mounted EcVolume's registration as a holder of its .ecj. +type ecjHold struct { + key string + once sync.Once + // The path's write generation when the hold was taken, and whether a + // writer was active then. Taken before the journal is opened and loaded, + // so they cover every write the load might have missed. + writeGen uint64 + writerAtStart bool +} + +// acquireEcjHold registers a holder of ecjPath, first waiting out any +// compaction in progress so the handle opened afterwards is on the final inode. +func acquireEcjHold(ecjPath string) *ecjHold { + h := &ecjHold{key: ecjPathKey(ecjPath)} + ecjUpdateWhenNotCompacting(h.key, func(st *ecjPathState) { + st.holders++ + h.writeGen = st.writeGen + h.writerAtStart = st.writers > 0 + }) + return h +} + +func (h *ecjHold) release() { + h.once.Do(func() { + ecjRelease(h.key, func(st *ecjPathState) { + if st.holders > 0 { + st.holders-- + } + }) + }) +} + +// tryBeginCompaction reserves the path for a compaction, or reports false when +// another holder or an active writer could still reach the current inode, or +// when a writer has run on the path since the hold was taken, so the journal +// may no longer be what the holder loaded. The returned func ends the +// reservation. +func (h *ecjHold) tryBeginCompaction() (end func(), ok bool) { + ecjPaths.Lock() + defer ecjPaths.Unlock() + st := ecjPaths.state[h.key] + if st == nil || st.holders != 1 || st.writers != 0 || st.compacting { + return nil, false + } + if h.writerAtStart || st.writeGen != h.writeGen { + return nil, false + } + st.compacting = true + var once sync.Once + return func() { + once.Do(func() { + ecjRelease(h.key, func(st *ecjPathState) { st.compacting = false }) + }) + }, true +} + +// BeginEcjWrite registers an out-of-band writer (shard copy, index recovery, +// ReceiveFile) of ecjPath, waiting out any compaction in progress. Compaction +// does not start until the returned func is called, so call it once the write, +// and any cleanup of a partial file, is done. +func BeginEcjWrite(ecjPath string) (done func()) { + key := ecjPathKey(ecjPath) + ecjUpdateWhenNotCompacting(key, func(st *ecjPathState) { + st.writers++ + st.writeGen++ + }) + var once sync.Once + return func() { + once.Do(func() { + ecjRelease(key, func(st *ecjPathState) { + if st.writers > 0 { + st.writers-- + } + }) + }) + } +} diff --git a/weed/storage/store_ec_journal.go b/weed/storage/store_ec_journal.go index c7e579252..a15a429ba 100644 --- a/weed/storage/store_ec_journal.go +++ b/weed/storage/store_ec_journal.go @@ -188,6 +188,12 @@ func (s *Store) mergeIntoMountedEcJournal(owner *DiskLocation, vid needle.Volume // the write reads the new records like any others. mounted reports that a // runtime appeared before the write; the caller then merges through it. func (s *Store) appendUnmountedEcJournal(owner *DiskLocation, vid needle.VolumeId, ecjPath string, local, ids map[types.NeedleId]struct{}, size int64, sync func(*erasure_coding.EcjAppend) error) (added int, mounted bool, err error) { + // The append and its sync, rollback or rewrite touch ecjPath through the + // open handle past the disk locks writeUnmountedEcJournal holds, so they + // run registered as a writer: a mount compacting this journal must not + // swap its inode underneath them. + done := erasure_coding.BeginEcjWrite(ecjPath) + defer done() pending, mounted, err := s.writeUnmountedEcJournal(owner, vid, ecjPath, local, ids, size) if pending == nil { return 0, mounted, err