From 0ff7794c54d20d8fa930822e62e60d7b3769b5cb Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Sat, 3 Oct 2026 04:08:19 +0300 Subject: [PATCH] volume: compact an oversized .ecj at mount, safely (Rust + Go) (#11555) * volume: compact an oversized .ecj at mount, safely (Rust + Go) Restore the mount-time compaction dropped from #11408, Rust + Go parity. A journal already bloated by repeated shard copies is folded down to the id set it encodes. - Trigger after load when file_records > max(threshold, 4x distinct), with a 1 MiB floor so small journals are never rewritten. The set is written to .ecj.compact.tmp + fsync, the handle dropped, renamed, the directory fsynced and the append handle reopened. A failure before the rename keeps the original journal and handle; a failure after it fails the mount. - Go never compacts after a failed journal load; the set would be partial and the rewrite would drop the unread records. - A per-path registry (ecj_registry.rs / ecj_registry.go) counts EcVolume holders and out-of-band writers of each .ecj. Compaction runs only when this volume is the sole holder and no copy is writing; holders and writers wait while one runs. This covers shared -dir.idx journals and cross-disk reconcile, where another EcVolume may hold the same journal. - VolumeEcShardsCopy and EC index recovery register as writers around their .ecj append and partial-file cleanup. - Under the reservation, re-check that the file on disk is still the inode and size that was loaded. - Publish errors are classified where they happen; a failed rename plus a failed restore reports both errors. - Compaction runs after the .vif / bitrot checks, so a refused mount leaves the journal untouched. - The tmp is opened like other volume files, removed at mount if a crash left it, and listed in every EC index cleanup path. Failure paths are tested through the real mount via injectable fs steps (open_with / newEcVolumeWith), plus sibling holders, active copies, changed-after-load, stale tmp cleanup, refused mounts and the Go load-error guard. Co-Authored-By: Claude Opus 5.5 (1M context) * volume: fail the mount when the compacted .ecj's directory cannot be synced The Rust mount synced the journal's directory after renaming the compacted file over it through the crate's best-effort fsync_dir, which returns Ok when the directory cannot be opened. A rename needs only write and search permission, so on a directory without read permission the replacement was published, never synced, and the mount went on taking deletes against it. Sync through a helper that propagates the open error, as Go's util.FsyncDir already does, so that case fails the mount like any other post-rename sync failure. Co-Authored-By: Claude Opus 5.5 (1M context) * volume: test the no-compaction-after-failed-load rule through the Go mount The test for it handed compactEcjAfterLoad an artificial error on a volume that had loaded cleanly, so it would not notice NewEcVolume dropping the real load error on the way to compaction. Make the journal read one of the injectable ecjFsOps steps and fail it inside the real mount, after the first chunk, on a journal whose last entry is an id the first chunk does not hold. The mount must leave the file byte for byte as it was; a clean remount then compacts and keeps that id. The Rust mount fails outright on a load error, so it has no equivalent path. Co-Authored-By: Claude Opus 5.5 (1M context) * volume: register ReceiveFile's .ecj writes with the journal registry ReceiveFile refuses a mounted EC volume only once, when the info message arrives, then creates the .ecj and streams chunks into it. A volume that mounted on that journal mid-stream could find a bloated prefix, pass the inode-and-size re-check and rename a compacted file over it; the rest of the stream then went to the unlinked inode and was lost. Register the path as a writer before the file is created, in both the Go and Rust handlers, and hold it until the file is closed and any partial copy removed, as the shard-copy and index-recovery appends already do. Co-Authored-By: Claude Opus 5.5 (1M context) * volume: skip .ecj compaction when a writer ran since the journal was loaded Compaction checked only that no writer was active at the reservation, and that the file was still the loaded inode at the loaded size. A ReceiveFile truncates and refills the journal in place, so one that ran during the mount's load, or after it, and finished before the reservation could leave different ids at the same length; compaction then wrote the stale set over them. Give each path a write generation that every writer bumps as it starts. A holder records it, and whether a writer was active, when it registers, which is before it opens and loads the journal. It may compact only if no writer was active then and the generation has not moved. Same rule in Go and Rust; the journal read becomes an injectable step in Rust as it is in Go, so both test the in-place rewrite through the real mount. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: match the ReadOnly(VolumeId) variant in write_volume_needles #11543 matched VolumeError::ReadOnly as a unit variant in Store::write_volume_needles, and #11544 changed it to ReadOnly(VolumeId) in the same merge window. Each passed CI on its own, but master no longer compiles the Rust volume server. Carry the volume id through. Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) Co-authored-by: Chris Lu Co-authored-by: Chris Lu --- seaweed-volume/src/server/grpc_server.rs | 112 ++ seaweed-volume/src/storage/disk_location.rs | 9 +- .../src/storage/erasure_coding/ec_volume.rs | 1067 ++++++++++++++++- .../storage/erasure_coding/ecj_registry.rs | 326 +++++ .../src/storage/erasure_coding/mod.rs | 1 + seaweed-volume/src/storage/store.rs | 6 +- .../src/storage/store_ec_journal.rs | 5 + weed/server/volume_grpc_copy.go | 13 + .../server/volume_grpc_copy_traversal_test.go | 13 +- weed/server/volume_grpc_erasure_coding.go | 6 +- .../volume_grpc_receive_file_ecj_test.go | 88 ++ weed/storage/disk_location_ec.go | 2 + weed/storage/erasure_coding/ec_volume.go | 314 ++++- .../ec_volume_ecj_compact_test.go | 415 +++++++ .../erasure_coding/ec_volume_ecj_test.go | 81 ++ weed/storage/erasure_coding/ecj_registry.go | 183 +++ weed/storage/store_ec_journal.go | 6 + 17 files changed, 2600 insertions(+), 47 deletions(-) create mode 100644 seaweed-volume/src/storage/erasure_coding/ecj_registry.rs create mode 100644 weed/server/volume_grpc_receive_file_ecj_test.go create mode 100644 weed/storage/erasure_coding/ec_volume_ecj_compact_test.go create mode 100644 weed/storage/erasure_coding/ecj_registry.go 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