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