mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
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) <noreply@anthropic.com> * 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) <noreply@anthropic.com> * 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) <noreply@anthropic.com> * 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) <noreply@anthropic.com> * 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) <noreply@anthropic.com> * 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) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com> Co-authored-by: Chris Lu <chris.lu@gmail.com>
This commit is contained in:
17 files changed
+2600
-47
No files matched your search
@@ -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<crate::storage::erasure_coding::ecj_registry::EcjWrite> = 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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
File diff suppressed because it is too large.
Load diff
@@ -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<HashMap<PathBuf, PathState>>,
|
||||
changed: Condvar,
|
||||
}
|
||||
|
||||
static REGISTRY: LazyLock<Registry> = LazyLock::new(|| Registry {
|
||||
paths: Mutex::new(HashMap::new()),
|
||||
changed: Condvar::new(),
|
||||
});
|
||||
|
||||
fn lock() -> MutexGuard<'static, HashMap<PathBuf, PathState>> {
|
||||
// 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<R>(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<EcjCompaction> {
|
||||
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();
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
@@ -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--
|
||||
}
|
||||
})
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
Reference in new issue
Block a user