mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 14:31:57 +02:00
volume: merge .ecj as a set union on EC shard copy + index recovery (Rust+Go) (#11554)
* volume: merge .ecj as a set union on EC shard copy + index recovery (Rust+Go) An EC volume's deletion journal is a set of needle ids, but shard copy and index recovery appended the peer's whole journal, doubling the file on every ec_balance round trip. Fold the peer's ids in as a union instead: only ids the local journal lacks are appended. - The journal is never replaced. A mounted EcVolume merges a peer's ids through its live handle under the lock deletes take (Go MergeJournal / Rust merge_journal), wherever its journal lives. - An unmounted journal gets only the missing ids appended while mounts are excluded; the delta is read outside the lock and re-read if the journal changed. - The source .ecj streams into memory as an id set: no staging files, chunked reads, memory proportional to distinct ids. - Go and Rust agree that a source journal exists when it sends a modified time or any bytes. A missing source stays a no-op. - Rust runs every merge in spawn_blocking and shares one receive/merge path between shard copy and index recovery. The decode path and the journal format are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume: route .ecj merges to the runtime that holds the journal open Disks sharing one index directory all resolved as the journal's owner, so the last one won and a sibling's mounted runtime was skipped: the merge appended behind its open handle and the sibling kept serving the peer's deleted needles until remount. Callers now name the receiving disk by its data directory; the merge goes through that disk's runtime, else a sibling runtime whose journal is the target file. In Go the unmounted append now holds every disk's EC lock (in location order) while it rechecks for a mount, so a sibling mounting from this disk's index during the unlocked read is merged through instead. In Rust a mount that lands during the read is merged through directly and its added count returned, rather than discarded and reported as zero. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume: sync merged .ecj records outside the disks' EC locks The unmounted merge held every disk's EC read lock across its fsync, so a slow sync on one disk held off mounts on all of them, along with the EC reads queued behind those mounts. Mounts only need to be excluded while the records are written: the write now happens under the locks and the fsync after they are released, since a later mount reads the written records from the page cache. A failed fsync rolls back only if nothing has mounted the journal or appended to it since the write. A merge through a mounted volume now keeps only that volume's disk locked across its fsync. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume: roll back an unsynced .ecj merge through a volume mounted mid-sync If a volume mounted after the unmounted merge wrote its records but before the fsync failed, the rollback kept the records because the journal was now open, leaving ids in the volume's deleted set that may never reach disk; a retried merge then saw them and synced nothing. The rollback now goes through that volume the way its own failed journal fsync does: truncate back and drop the ids from the in-memory set, so a retry appends and syncs them again. It still keeps the records if the volume journaled since, as truncating would lose that delete. No fsync runs under the disk locks. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume: decide .ecj merge rollback from the journal's actual length Two runtimes can hold one journal (cross-disk mounts). The rollback of an unsynced merge checked one runtime's cached ecjFileSize, which another runtime's appends leave stale, so it could truncate a delete that runtime had already synced. The rollback now holds every holder's journal lock and truncates only if the file's actual length is still the append's end, then updates each holder's size and deleted set. Otherwise later records follow the merged ones, so they stay and are rewritten in place and synced outside the locks, rather than left possibly not durable. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume: keep unsynced .ecj merge ids out of mounted deleted sets When a merge's fsync failed, later records blocked the rollback, and the rewrite-and-sync failed as well, the merged ids stayed in every mounted volume's deleted set without being shown durable, so a retried merge saw them as present and synced nothing. They now leave those sets while the records stay in the file, matching DeleteNeedleFromEcx, which publishes an id only after its record syncs. The merge returns the error and a retry appends and syncs them again. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume: publish merged .ecj ids to every holder of the journal Two runtimes can journal into the same file when disks share an index directory. The merge went through only the first holder, leaving a sibling's in-memory deleted set without the ids, so it could keep serving a needle the peer deleted until it remounted. Every holder of the journal now gets the merged ids, in Go and in the volume server. Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * Publish merged .ecj ids to the journal actually written mountedEcJournal prefers the receiving disk's own runtime for the vid, whose journal may live in its data directory while the copied records name a sibling's journal in the index directory. Publishing by the requested ecjPath then marked a holder of a different file deleted on records that file never persisted, resurrecting the needles on remount. Publish by the picked runtime's journal path instead. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com> Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
16 files changed
+2292
-64
No files matched your search
@@ -17,6 +17,7 @@ use crate::pb::master_pb;
|
||||
use crate::pb::volume_server_pb;
|
||||
use crate::pb::volume_server_pb::volume_server_server::VolumeServer;
|
||||
use crate::storage::erasure_coding::ec_shard::{DATA_SHARDS_COUNT, ShardId, shard_id_try_from};
|
||||
use crate::storage::erasure_coding::ecj_merge::EcjIdDecoder;
|
||||
use crate::storage::needle::needle::{self, Needle};
|
||||
use crate::storage::types::*;
|
||||
use crate::storage::volume::VolumeSpec;
|
||||
@@ -3239,7 +3240,10 @@ impl VolumeServer for VolumeGrpcService {
|
||||
}
|
||||
}
|
||||
|
||||
// Copy .ecj file if requested
|
||||
// Copy .ecj file if requested. The journal is a *set* of ids: merge
|
||||
// the source's into the local one as a union, never append it whole,
|
||||
// or every balance round trip doubles it. A source without one
|
||||
// is not an error.
|
||||
if req.copy_ecj_file {
|
||||
let copy_req = volume_server_pb::CopyFileRequest {
|
||||
volume_id: req.volume_id,
|
||||
@@ -3261,19 +3265,26 @@ impl VolumeServer for VolumeGrpcService {
|
||||
))
|
||||
})?
|
||||
.into_inner();
|
||||
|
||||
let file_path = {
|
||||
let base =
|
||||
crate::storage::volume::volume_file_name(&dest_idx_dir, &req.collection, vid);
|
||||
format!("{}.ecj", base)
|
||||
};
|
||||
let file = tokio::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.append(true)
|
||||
.open(&file_path)
|
||||
.await
|
||||
.map_err(|e| Status::internal(format!("create {}: {}", file_path, e)))?;
|
||||
drain_copy_stream_to_file(&mut stream, file, &file_path, ".ecj").await?;
|
||||
let (ids, found) = receive_ecj_ids(&mut stream).await.map_err(|e| {
|
||||
Status::internal(format!(
|
||||
"VolumeEcShardsCopy volume {} copy .ecj: {}",
|
||||
vid, e
|
||||
))
|
||||
})?;
|
||||
if found {
|
||||
let ecj_path = format!(
|
||||
"{}.ecj",
|
||||
crate::storage::volume::volume_file_name(&dest_idx_dir, &req.collection, vid)
|
||||
);
|
||||
merge_ecj_ids(&self.state, vid, dest_dir.clone(), ecj_path, ids)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Status::internal(format!(
|
||||
"VolumeEcShardsCopy volume {} merge .ecj: {}",
|
||||
vid, e
|
||||
))
|
||||
})?;
|
||||
}
|
||||
}
|
||||
|
||||
// Copy .vif file if requested
|
||||
@@ -5960,6 +5971,50 @@ async fn drain_copy_stream_to_file(
|
||||
}
|
||||
}
|
||||
|
||||
/// Decode a CopyFile stream of an `.ecj` into its distinct ids without staging
|
||||
/// it on disk. `found` is false only when the source has no journal, which it
|
||||
/// signals with neither a modified time nor any bytes; an empty journal still
|
||||
/// carries its modified time.
|
||||
pub(crate) async fn receive_ecj_ids(
|
||||
stream: &mut tonic::Streaming<volume_server_pb::CopyFileResponse>,
|
||||
) -> std::io::Result<(std::collections::HashSet<NeedleId>, bool)> {
|
||||
let mut decoder = EcjIdDecoder::default();
|
||||
let mut found = false;
|
||||
while let Some(chunk) = stream
|
||||
.message()
|
||||
.await
|
||||
.map_err(|e| std::io::Error::other(format!("recv .ecj: {}", e)))?
|
||||
{
|
||||
found |= chunk.modified_ts_ns != 0 || !chunk.file_content.is_empty();
|
||||
decoder.push(&chunk.file_content);
|
||||
}
|
||||
Ok((decoder.into_ids(), found))
|
||||
}
|
||||
|
||||
/// Merge received `.ecj` ids into vid's local journal at `ecj_path` on the
|
||||
/// disk whose data directory is `data_dir`, off the async runtime (the merge
|
||||
/// reads, appends and fsyncs). Shared by shard copy and index recovery.
|
||||
pub(crate) async fn merge_ecj_ids(
|
||||
state: &std::sync::Arc<super::volume_server::VolumeServerState>,
|
||||
vid: VolumeId,
|
||||
data_dir: String,
|
||||
ecj_path: String,
|
||||
ids: std::collections::HashSet<NeedleId>,
|
||||
) -> std::io::Result<usize> {
|
||||
let state = std::sync::Arc::clone(state);
|
||||
tokio::task::spawn_blocking(move || {
|
||||
crate::storage::store_ec_journal::merge_ec_journal(
|
||||
&state.store,
|
||||
vid,
|
||||
&data_dir,
|
||||
&ecj_path,
|
||||
&ids,
|
||||
)
|
||||
})
|
||||
.await
|
||||
.map_err(|e| std::io::Error::other(format!("join .ecj merge: {}", e)))?
|
||||
}
|
||||
|
||||
/// One file of a volume copy: what to ask the source for and where it lands.
|
||||
#[derive(Clone, Copy)]
|
||||
struct CopyFileSpec<'a> {
|
||||
|
||||
@@ -1753,7 +1753,7 @@ async fn fetch_ec_index_from_one_peer(
|
||||
.await
|
||||
.map_err(|e| io::Error::other(format!("copy .ecx: {}", e)))?
|
||||
.into_inner();
|
||||
drain_copy_stream(stream, ecx_path, false).await?;
|
||||
drain_copy_stream(stream, ecx_path).await?;
|
||||
|
||||
let meta =
|
||||
fs::metadata(ecx_path).map_err(|e| io::Error::other(format!("stat copied .ecx: {}", e)))?;
|
||||
@@ -1766,15 +1766,30 @@ async fn fetch_ec_index_from_one_peer(
|
||||
)));
|
||||
}
|
||||
|
||||
// .ecj is the source peer's deletion journal (appended); .vif carries EC
|
||||
// params. Both are best-effort: a missing .ecj is recreated at mount and a
|
||||
// missing .vif falls back to default EC parameters. A failed .ecj append
|
||||
// leaves a partial file, so drop it.
|
||||
// .ecj is the source peer's deletion journal; .vif carries EC params. Both
|
||||
// are best-effort: a missing .ecj is recreated at mount and a missing .vif
|
||||
// falls back to default EC parameters. The journal is a *set*: merge the
|
||||
// peer's ids into any local ones as a union instead of appending, so
|
||||
// a volume bounced between servers cannot double its journal. The merge
|
||||
// only appends whole records, so a failure leaves nothing to clean up.
|
||||
match client.copy_file(copy_req(".ecj", true)).await {
|
||||
Ok(resp) => {
|
||||
if let Err(e) = drain_copy_stream(resp.into_inner(), ecj_path, true).await {
|
||||
let mut stream = resp.into_inner();
|
||||
let merged = match crate::server::grpc_server::receive_ecj_ids(&mut stream).await {
|
||||
Ok((ids, true)) => crate::server::grpc_server::merge_ecj_ids(
|
||||
state,
|
||||
m.vid,
|
||||
m.data_dir.clone(),
|
||||
ecj_path.to_string(),
|
||||
ids,
|
||||
)
|
||||
.await
|
||||
.map(|_| ()),
|
||||
Ok((_, false)) => Ok(()),
|
||||
Err(e) => Err(e),
|
||||
};
|
||||
if let Err(e) = merged {
|
||||
tracing::warn!(volume_id = m.vid.0, peer = %peer, "copy .ecj: {}", e);
|
||||
let _ = fs::remove_file(ecj_path);
|
||||
}
|
||||
}
|
||||
Err(e) => tracing::warn!(volume_id = m.vid.0, peer = %peer, "copy .ecj: {}", e),
|
||||
@@ -1782,7 +1797,7 @@ async fn fetch_ec_index_from_one_peer(
|
||||
|
||||
match client.copy_file(copy_req(".vif", true)).await {
|
||||
Ok(resp) => {
|
||||
if let Err(e) = drain_copy_stream(resp.into_inner(), vif_path, false).await {
|
||||
if let Err(e) = drain_copy_stream(resp.into_inner(), vif_path).await {
|
||||
tracing::warn!(volume_id = m.vid.0, peer = %peer, "copy .vif: {}", e);
|
||||
}
|
||||
}
|
||||
@@ -1792,22 +1807,14 @@ async fn fetch_ec_index_from_one_peer(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Drain a CopyFile stream into a local file, appending or truncating.
|
||||
/// Drain a CopyFile stream into a local file, truncating it first.
|
||||
async fn drain_copy_stream(
|
||||
mut stream: tonic::Streaming<crate::pb::volume_server_pb::CopyFileResponse>,
|
||||
dest_path: &str,
|
||||
append: bool,
|
||||
) -> io::Result<()> {
|
||||
use std::io::Write;
|
||||
let mut file = if append {
|
||||
fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.append(true)
|
||||
.open(dest_path)
|
||||
} else {
|
||||
fs::File::create(dest_path)
|
||||
}
|
||||
.map_err(|e| io::Error::other(format!("create {}: {}", dest_path, e)))?;
|
||||
let mut file = fs::File::create(dest_path)
|
||||
.map_err(|e| io::Error::other(format!("create {}: {}", dest_path, e)))?;
|
||||
while let Some(chunk) = stream
|
||||
.message()
|
||||
.await
|
||||
|
||||
@@ -1642,23 +1642,33 @@ impl EcVolume {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let mut buf = [0u8; NEEDLE_ID_SIZE];
|
||||
needle_id.to_bytes(&mut buf);
|
||||
self.append_journal(&buf)?;
|
||||
if let Ok(mut set) = self.deleted_needles.write() {
|
||||
set.insert(needle_id);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Append whole records to `.ecj` and sync them. On any failure the file is
|
||||
/// truncated back to the pre-append length so the on-disk journal and
|
||||
/// `deleted_needles` cannot drift.
|
||||
fn append_journal(&mut self, records: &[u8]) -> io::Result<()> {
|
||||
let prev_ecj_size = self.ecj_file_size;
|
||||
let append_result: io::Result<()> = {
|
||||
let ecj_file = self
|
||||
.ecj_file
|
||||
.as_mut()
|
||||
.ok_or_else(|| io::Error::other("ecj file not open"))?;
|
||||
let mut buf = [0u8; NEEDLE_ID_SIZE];
|
||||
needle_id.to_bytes(&mut buf);
|
||||
ecj_file.write_all(&buf).and_then(|_| ecj_file.sync_all())
|
||||
ecj_file
|
||||
.write_all(records)
|
||||
.and_then(|_| ecj_file.sync_all())
|
||||
};
|
||||
|
||||
match append_result {
|
||||
Ok(()) => {
|
||||
self.ecj_file_size += NEEDLE_ID_SIZE as i64;
|
||||
if let Ok(mut set) = self.deleted_needles.write() {
|
||||
set.insert(needle_id);
|
||||
}
|
||||
self.ecj_file_size += records.len() as i64;
|
||||
Ok(())
|
||||
}
|
||||
Err(e) => {
|
||||
@@ -1669,20 +1679,12 @@ impl EcVolume {
|
||||
// lacks FILE_WRITE_DATA, so set_len through it fails with
|
||||
// ERROR_ACCESS_DENIED and the rollback would silently not
|
||||
// happen.
|
||||
let ecj_path = format!(
|
||||
"{}.ecj",
|
||||
crate::storage::volume::volume_file_name(
|
||||
&self.ecx_actual_dir,
|
||||
&self.collection,
|
||||
self.volume_id,
|
||||
)
|
||||
);
|
||||
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()));
|
||||
if let Err(trunc_err) = rollback {
|
||||
tracing::error!(
|
||||
volume_id = self.volume_id.0,
|
||||
needle_id = needle_id.0,
|
||||
truncate_error = %trunc_err,
|
||||
"failed to truncate ecj after append failure"
|
||||
);
|
||||
@@ -1828,6 +1830,32 @@ impl EcVolume {
|
||||
Ok(needles)
|
||||
}
|
||||
|
||||
/// Fold a peer's deletion journal into this mounted volume: append only
|
||||
/// the ids not already deleted here, in one write and one fsync, then
|
||||
/// publish them into the in-memory set — the same commit order as
|
||||
/// `journal_delete`. `&mut self` serializes it with deletes under the store
|
||||
/// lock, and the live handle is appended to in place, so no delete can land
|
||||
/// in an orphaned file. Returns how many ids were added.
|
||||
pub fn merge_journal(&mut self, ids: &HashSet<NeedleId>) -> io::Result<usize> {
|
||||
let delta = super::ecj_merge::ecj_delta(ids, |id| self.is_needle_deleted(*id));
|
||||
if delta.is_empty() {
|
||||
return Ok(0);
|
||||
}
|
||||
self.append_journal(&super::ecj_merge::encode_ecj_ids(&delta))?;
|
||||
if let Ok(mut set) = self.deleted_needles.write() {
|
||||
set.extend(delta.iter().copied());
|
||||
}
|
||||
Ok(delta.len())
|
||||
}
|
||||
|
||||
/// Marks ids deleted in memory only, for runtimes sharing a journal
|
||||
/// another runtime already appended them to.
|
||||
pub fn publish_merged_ids(&self, ids: &HashSet<NeedleId>) {
|
||||
if let Ok(mut set) = self.deleted_needles.write() {
|
||||
set.extend(ids.iter().copied());
|
||||
}
|
||||
}
|
||||
|
||||
// ---- Lifecycle ----
|
||||
|
||||
pub fn close(&mut self) {
|
||||
@@ -2514,6 +2542,71 @@ mod tests {
|
||||
assert_eq!((fc, dc), (2, 2));
|
||||
}
|
||||
|
||||
/// A mounted volume merges a peer's journal through its own handle: only
|
||||
/// the missing ids are appended, the file keeps its inode, the in-memory
|
||||
/// set follows, and a later delete lands in the same live file.
|
||||
#[test]
|
||||
fn test_merge_journal_appends_missing_ids_in_place() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let entries: Vec<_> = (1..=5)
|
||||
.map(|id| {
|
||||
(
|
||||
NeedleId(id),
|
||||
Offset::from_actual_offset(8 * id as i64),
|
||||
Size(100),
|
||||
)
|
||||
})
|
||||
.collect();
|
||||
write_ecx_file(dir, "", VolumeId(1), &entries);
|
||||
|
||||
let mut vol = EcVolume::new(dir, dir, "", VolumeId(1)).unwrap();
|
||||
vol.journal_delete(NeedleId(1)).unwrap();
|
||||
let ecj_path = vol.ecj_file_name();
|
||||
#[cfg(unix)]
|
||||
let inode = {
|
||||
use std::os::unix::fs::MetadataExt;
|
||||
std::fs::metadata(&ecj_path).unwrap().ino()
|
||||
};
|
||||
|
||||
let peer: HashSet<NeedleId> = [1, 2, 3].into_iter().map(NeedleId).collect();
|
||||
assert_eq!(vol.merge_journal(&peer).unwrap(), 2);
|
||||
assert_eq!(
|
||||
vol.merge_journal(&peer).unwrap(),
|
||||
0,
|
||||
"a repeated merge appends nothing"
|
||||
);
|
||||
for id in 1..=3 {
|
||||
assert!(vol.is_needle_deleted(NeedleId(id)), "id {}", id);
|
||||
}
|
||||
assert_eq!(vol.file_and_delete_count().1, 3);
|
||||
|
||||
vol.journal_delete(NeedleId(4)).unwrap();
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::MetadataExt;
|
||||
assert_eq!(
|
||||
std::fs::metadata(&ecj_path).unwrap().ino(),
|
||||
inode,
|
||||
"the journal must never be replaced under the open handle"
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
vol.read_deleted_needles().unwrap(),
|
||||
(1..=4).map(NeedleId).collect::<Vec<_>>()
|
||||
);
|
||||
vol.close();
|
||||
|
||||
let vol = EcVolume::new(dir, dir, "", VolumeId(1)).unwrap();
|
||||
for id in 1..=4 {
|
||||
assert!(
|
||||
vol.is_needle_deleted(NeedleId(id)),
|
||||
"id {} survives remount",
|
||||
id
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Write a raw `.ecj` containing `ids` repeated `repeats` times, i.e. the
|
||||
/// shape the append paths produce when a peer's whole journal is
|
||||
/// concatenated onto this one over and over.
|
||||
|
||||
@@ -0,0 +1,290 @@
|
||||
//! Set-union merge of `.ecj` deletion journals for EC shard copy / index
|
||||
//! recovery. Mirrors Go's `weed/storage/erasure_coding/ecj_merge.go`.
|
||||
//!
|
||||
//! An EC volume's deletion journal (`<vid>.ecj`) is a *set* of deleted needle
|
||||
//! ids stored as 8-byte big-endian records. Shard copy and index recovery fold
|
||||
//! a peer's journal into the local one; they must append only the ids the
|
||||
//! local journal lacks, or every `ec_balance` round trip doubles the file.
|
||||
//!
|
||||
//! The journal is only ever appended to, never replaced: a mounted `EcVolume`
|
||||
//! holds it open, and a rename would leave that handle writing to an unlinked
|
||||
//! inode, losing every later delete at the next mount. A mounted volume merges
|
||||
//! through [`EcVolume::merge_journal`](super::ec_volume::EcVolume::merge_journal);
|
||||
//! [`append_ecj_ids`] is for a journal no volume has open.
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::fs::{self, OpenOptions};
|
||||
use std::io::{self, Read, Seek, SeekFrom, Write};
|
||||
use std::path::Path;
|
||||
|
||||
use crate::storage::types::{NEEDLE_ID_SIZE, NeedleId};
|
||||
use crate::storage::volume::fsync_dir;
|
||||
use crate::storage::volume_open::open_volume_file;
|
||||
|
||||
/// Bytes per read when scanning a journal; a multiple of `NEEDLE_ID_SIZE`.
|
||||
const ECJ_READ_CHUNK_BYTES: usize = 1 << 20;
|
||||
|
||||
/// Decodes `.ecj` records from a byte stream split at arbitrary boundaries (a
|
||||
/// CopyFile stream, chunked reads), collecting the distinct ids. Memory
|
||||
/// follows the number of distinct ids, not the journal's length. A trailing
|
||||
/// partial record is never decoded.
|
||||
#[derive(Default)]
|
||||
pub(crate) struct EcjIdDecoder {
|
||||
ids: HashSet<NeedleId>,
|
||||
partial: [u8; NEEDLE_ID_SIZE],
|
||||
pending: usize,
|
||||
}
|
||||
|
||||
impl EcjIdDecoder {
|
||||
pub(crate) fn push(&mut self, mut bytes: &[u8]) {
|
||||
if self.pending > 0 {
|
||||
let n = (NEEDLE_ID_SIZE - self.pending).min(bytes.len());
|
||||
self.partial[self.pending..self.pending + n].copy_from_slice(&bytes[..n]);
|
||||
self.pending += n;
|
||||
bytes = &bytes[n..];
|
||||
if self.pending < NEEDLE_ID_SIZE {
|
||||
return;
|
||||
}
|
||||
self.ids.insert(NeedleId::from_bytes(&self.partial));
|
||||
self.pending = 0;
|
||||
}
|
||||
let mut records = bytes.chunks_exact(NEEDLE_ID_SIZE);
|
||||
for record in &mut records {
|
||||
self.ids.insert(NeedleId::from_bytes(record));
|
||||
}
|
||||
let rest = records.remainder();
|
||||
self.partial[..rest.len()].copy_from_slice(rest);
|
||||
self.pending = rest.len();
|
||||
}
|
||||
|
||||
pub(crate) fn into_ids(self) -> HashSet<NeedleId> {
|
||||
self.ids
|
||||
}
|
||||
}
|
||||
|
||||
/// Read the distinct ids of the journal at `path` in bounded chunks. A missing
|
||||
/// file reads as empty. Also returns the whole-record length read; a torn
|
||||
/// trailing partial record is excluded from it.
|
||||
pub(crate) fn read_ecj_ids(path: &str) -> io::Result<(HashSet<NeedleId>, u64)> {
|
||||
let mut file = match fs::File::open(path) {
|
||||
Ok(f) => f,
|
||||
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok((HashSet::new(), 0)),
|
||||
Err(e) => return Err(e),
|
||||
};
|
||||
let len = file.metadata()?.len();
|
||||
let size = len - len % NEEDLE_ID_SIZE as u64;
|
||||
let mut decoder = EcjIdDecoder::default();
|
||||
let mut buf = vec![0u8; (ECJ_READ_CHUNK_BYTES as u64).min(size) as usize];
|
||||
let mut off = 0u64;
|
||||
while off < size {
|
||||
let want = (buf.len() as u64).min(size - off) as usize;
|
||||
file.read_exact(&mut buf[..want])?;
|
||||
decoder.push(&buf[..want]);
|
||||
off += want as u64;
|
||||
}
|
||||
Ok((decoder.into_ids(), size))
|
||||
}
|
||||
|
||||
/// The ids of `incoming` that `has` does not report, ascending, so a merge
|
||||
/// appends deterministic output.
|
||||
pub(crate) fn ecj_delta(
|
||||
incoming: &HashSet<NeedleId>,
|
||||
has: impl Fn(&NeedleId) -> bool,
|
||||
) -> Vec<NeedleId> {
|
||||
let mut delta: Vec<NeedleId> = incoming.iter().copied().filter(|id| !has(id)).collect();
|
||||
delta.sort_unstable();
|
||||
delta
|
||||
}
|
||||
|
||||
pub(crate) fn encode_ecj_ids(ids: &[NeedleId]) -> Vec<u8> {
|
||||
let mut buf = vec![0u8; ids.len() * NEEDLE_ID_SIZE];
|
||||
for (record, id) in buf.chunks_exact_mut(NEEDLE_ID_SIZE).zip(ids) {
|
||||
id.to_bytes(record);
|
||||
}
|
||||
buf
|
||||
}
|
||||
|
||||
/// Append to the journal at `path` the ids of `incoming` that `local` lacks,
|
||||
/// in one write and one fsync, returning how many were added. `local` and
|
||||
/// `size` come from [`read_ecj_ids`] on the same path: if the journal's
|
||||
/// whole-record length is no longer `size`, returns `Ok(None)` so the caller
|
||||
/// re-reads. A torn tail past `size` is truncated first so the new records
|
||||
/// stay aligned.
|
||||
pub(crate) fn append_ecj_ids(
|
||||
path: &str,
|
||||
local: &HashSet<NeedleId>,
|
||||
incoming: &HashSet<NeedleId>,
|
||||
size: u64,
|
||||
) -> io::Result<Option<usize>> {
|
||||
let delta = ecj_delta(incoming, |id| local.contains(id));
|
||||
if delta.is_empty() {
|
||||
return Ok(Some(0));
|
||||
}
|
||||
let created = !Path::new(path).exists();
|
||||
let mut file = open_volume_file(OpenOptions::new().read(true).write(true).create(true), path)?;
|
||||
let len = file.metadata()?.len();
|
||||
if len - len % NEEDLE_ID_SIZE as u64 != size {
|
||||
return Ok(None);
|
||||
}
|
||||
if len != size {
|
||||
file.set_len(size)?;
|
||||
}
|
||||
let appended = file
|
||||
.seek(SeekFrom::Start(size))
|
||||
.and_then(|_| file.write_all(&encode_ecj_ids(&delta)))
|
||||
.and_then(|_| file.sync_all());
|
||||
if let Err(e) = appended {
|
||||
let _ = file.set_len(size);
|
||||
return Err(e);
|
||||
}
|
||||
if created {
|
||||
fsync_dir(path)?;
|
||||
}
|
||||
Ok(Some(delta.len()))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn ids(v: &[u64]) -> HashSet<NeedleId> {
|
||||
v.iter().map(|&id| NeedleId(id)).collect()
|
||||
}
|
||||
|
||||
fn bytes(v: &[u64]) -> Vec<u8> {
|
||||
encode_ecj_ids(&v.iter().map(|&id| NeedleId(id)).collect::<Vec<_>>())
|
||||
}
|
||||
|
||||
fn records(path: &str) -> Vec<u64> {
|
||||
let data = fs::read(path).expect("read ecj");
|
||||
assert_eq!(data.len() % NEEDLE_ID_SIZE, 0, "journal must stay aligned");
|
||||
data.chunks_exact(NEEDLE_ID_SIZE)
|
||||
.map(|c| NeedleId::from_bytes(c).0)
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Run the unmounted merge the way the server does: read, then append.
|
||||
fn merge_file(path: &str, incoming: &HashSet<NeedleId>) -> usize {
|
||||
let (local, size) = read_ecj_ids(path).expect("read");
|
||||
append_ecj_ids(path, &local, incoming, size)
|
||||
.expect("append")
|
||||
.expect("journal unchanged")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn decoder_handles_records_split_across_chunks() {
|
||||
let mut stream = bytes(&[1, 2, 3, 2, 0x0102030405060708]);
|
||||
stream.extend_from_slice(&[9, 9, 9]);
|
||||
for chunk in 1..=stream.len() {
|
||||
let mut d = EcjIdDecoder::default();
|
||||
for piece in stream.chunks(chunk) {
|
||||
d.push(piece);
|
||||
}
|
||||
assert_eq!(
|
||||
d.into_ids(),
|
||||
ids(&[1, 2, 3, 0x0102030405060708]),
|
||||
"chunk {}",
|
||||
chunk
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_ecj_ids_dedups_and_ignores_torn_tail() {
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let missing = dir.path().join("missing.ecj");
|
||||
let (got, size) = read_ecj_ids(missing.to_str().unwrap()).expect("read");
|
||||
assert!(got.is_empty());
|
||||
assert_eq!(size, 0);
|
||||
|
||||
let torn = dir.path().join("torn.ecj");
|
||||
let mut data = bytes(&[1, 2, 1]);
|
||||
data.extend_from_slice(&[7, 7, 7]);
|
||||
fs::write(&torn, data).unwrap();
|
||||
let (got, size) = read_ecj_ids(torn.to_str().unwrap()).expect("read");
|
||||
assert_eq!(got, ids(&[1, 2]));
|
||||
assert_eq!(size, 3 * NEEDLE_ID_SIZE as u64);
|
||||
|
||||
// A bloated journal repeating a few ids across several read chunks
|
||||
// keeps only the distinct ids.
|
||||
let bloated = dir.path().join("bloated.ecj");
|
||||
let data: Vec<u8> = (0..300_000u64).flat_map(|i| bytes(&[i % 3])).collect();
|
||||
fs::write(&bloated, &data).unwrap();
|
||||
let (got, size) = read_ecj_ids(bloated.to_str().unwrap()).expect("read");
|
||||
assert_eq!(got, ids(&[0, 1, 2]));
|
||||
assert_eq!(size, data.len() as u64);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn appends_only_missing_ids() {
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let path = dir.path().join("vol.ecj");
|
||||
let path = path.to_str().unwrap();
|
||||
fs::write(path, bytes(&[1, 2, 3])).unwrap();
|
||||
assert_eq!(merge_file(path, &ids(&[3, 4])), 1);
|
||||
assert_eq!(records(path), vec![1, 2, 3, 4]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn round_trip_stays_constant() {
|
||||
// A->B->A->B 20 times: the old append path doubled the journal each trip.
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let a = dir.path().join("a.ecj");
|
||||
let b = dir.path().join("b.ecj");
|
||||
let (a, b) = (a.to_str().unwrap(), b.to_str().unwrap());
|
||||
fs::write(a, bytes(&[1, 2])).unwrap();
|
||||
fs::write(b, bytes(&[2, 3])).unwrap();
|
||||
for i in 0..20 {
|
||||
let (src, dst) = if i % 2 == 0 { (a, b) } else { (b, a) };
|
||||
let (incoming, _) = read_ecj_ids(src).expect("read");
|
||||
merge_file(dst, &incoming);
|
||||
}
|
||||
for path in [a, b] {
|
||||
let mut got = records(path);
|
||||
got.sort_unstable();
|
||||
assert_eq!(got, vec![1, 2, 3]);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn nothing_new_leaves_journal_alone() {
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let missing = dir.path().join("missing.ecj");
|
||||
assert_eq!(merge_file(missing.to_str().unwrap(), &ids(&[])), 0);
|
||||
assert!(
|
||||
!missing.exists(),
|
||||
"an empty merge must not create a journal"
|
||||
);
|
||||
|
||||
let path = dir.path().join("vol.ecj");
|
||||
let path = path.to_str().unwrap();
|
||||
fs::write(path, bytes(&[1, 2])).unwrap();
|
||||
assert_eq!(merge_file(path, &ids(&[2, 1])), 0);
|
||||
assert_eq!(records(path), vec![1, 2]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn repairs_torn_tail() {
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let path = dir.path().join("vol.ecj");
|
||||
let path = path.to_str().unwrap();
|
||||
let mut data = bytes(&[1, 2]);
|
||||
data.extend_from_slice(&[9, 9, 9]);
|
||||
fs::write(path, data).unwrap();
|
||||
assert_eq!(merge_file(path, &ids(&[3])), 1);
|
||||
assert_eq!(records(path), vec![1, 2, 3]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_changed_journal() {
|
||||
let dir = tempfile::tempdir().expect("tempdir");
|
||||
let path = dir.path().join("vol.ecj");
|
||||
let path = path.to_str().unwrap();
|
||||
fs::write(path, bytes(&[1])).unwrap();
|
||||
let (local, size) = read_ecj_ids(path).expect("read");
|
||||
fs::write(path, bytes(&[1, 5])).unwrap();
|
||||
let outcome = append_ecj_ids(path, &local, &ids(&[2]), size).expect("append");
|
||||
assert_eq!(outcome, None);
|
||||
assert_eq!(records(path), vec![1, 5]);
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,7 @@ pub mod ec_encoder;
|
||||
pub mod ec_locate;
|
||||
pub mod ec_shard;
|
||||
pub mod ec_volume;
|
||||
pub mod ecj_merge;
|
||||
|
||||
pub use ec_shard::{
|
||||
DATA_SHARDS_COUNT, EcVolumeShard, MAX_SHARD_COUNT, MIN_TOTAL_DISKS, PARITY_SHARDS_COUNT,
|
||||
|
||||
@@ -6,6 +6,7 @@ pub(crate) mod io_error;
|
||||
pub mod needle;
|
||||
pub mod needle_map;
|
||||
pub mod store;
|
||||
pub mod store_ec_journal;
|
||||
pub mod store_ec_mirror;
|
||||
pub mod store_ec_reconcile;
|
||||
pub mod super_block;
|
||||
|
||||
@@ -0,0 +1,369 @@
|
||||
//! Merging a peer's `.ecj` deletion ids into a local EC journal. Mirrors Go's
|
||||
//! `Store.MergeEcJournal` (`weed/storage/store_ec_journal.go`).
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::io;
|
||||
use std::path::Path;
|
||||
use std::sync::RwLock;
|
||||
|
||||
use crate::storage::erasure_coding::ecj_merge::{append_ecj_ids, read_ecj_ids};
|
||||
use crate::storage::store::Store;
|
||||
use crate::storage::types::{NeedleId, VolumeId};
|
||||
|
||||
/// How often an unmounted merge re-reads a journal that changed under it. Only
|
||||
/// a mount-delete-unmount or a concurrent merge between the read and the
|
||||
/// append changes it, so one retry is nearly always enough.
|
||||
const ECJ_MERGE_ATTEMPTS: usize = 5;
|
||||
|
||||
/// Fold a peer's deletion `ids` into the local journal of EC volume `vid` on
|
||||
/// the receiving disk, the one whose data directory is `data_dir`; `ecj_path`
|
||||
/// is that journal's path in the disk's index directory. Appends only the ids
|
||||
/// the journal lacks and returns how many it added. Blocking: call it from
|
||||
/// `spawn_blocking`.
|
||||
///
|
||||
/// A mounted volume owns its journal: the merge goes through its open handle
|
||||
/// and in-memory set. That is the receiving disk's own runtime for `vid`,
|
||||
/// wherever its journal lives (it may sit in the data dir rather than
|
||||
/// `ecj_path`'s index dir), else a sibling runtime journaling into `ecj_path`
|
||||
/// itself: disks sharing one index directory, or reconciliation mounting `vid`
|
||||
/// on a disk that journals into another's (#9212). Otherwise `ecj_path` is
|
||||
/// appended to under the store write lock, which mounts take, so no mount can
|
||||
/// open it mid-append. The read that computes the delta runs outside the lock,
|
||||
/// and a journal that changed in between — including by a concurrent merge —
|
||||
/// is re-read.
|
||||
pub fn merge_ec_journal(
|
||||
store: &RwLock<Store>,
|
||||
vid: VolumeId,
|
||||
data_dir: &str,
|
||||
ecj_path: &str,
|
||||
ids: &HashSet<NeedleId>,
|
||||
) -> io::Result<usize> {
|
||||
merge_ec_journal_with(store, vid, data_dir, ecj_path, ids, read_ecj_ids)
|
||||
}
|
||||
|
||||
/// `merge_ec_journal` with the unlocked journal read injected, so a test can
|
||||
/// mount the volume between that read and the append.
|
||||
fn merge_ec_journal_with(
|
||||
store: &RwLock<Store>,
|
||||
vid: VolumeId,
|
||||
data_dir: &str,
|
||||
ecj_path: &str,
|
||||
ids: &HashSet<NeedleId>,
|
||||
mut read: impl FnMut(&str) -> io::Result<(HashSet<NeedleId>, u64)>,
|
||||
) -> io::Result<usize> {
|
||||
for _ in 0..ECJ_MERGE_ATTEMPTS {
|
||||
{
|
||||
let mut store = store
|
||||
.write()
|
||||
.map_err(|_| io::Error::other("store lock poisoned"))?;
|
||||
if let Some(primary) = mounted_ec_journal(&store, vid, data_dir, ecj_path)? {
|
||||
let ecv = store.locations[primary]
|
||||
.find_ec_volume_mut(vid)
|
||||
.expect("mounted journal runtime");
|
||||
let journal_path = ecv.ecj_file_name();
|
||||
let added = ecv.merge_journal(ids)?;
|
||||
// Publish to the holders of the file the merge wrote to — the
|
||||
// picked runtime's journal may live outside ecj_path, and a
|
||||
// holder of a different file must not claim ids it lacks.
|
||||
publish_to_journal_siblings(&store, primary, vid, &journal_path, ids);
|
||||
return Ok(added);
|
||||
}
|
||||
}
|
||||
let (local, size) = read(ecj_path)?;
|
||||
let mut store = store
|
||||
.write()
|
||||
.map_err(|_| io::Error::other("store lock poisoned"))?;
|
||||
if let Some(primary) = mounted_ec_journal(&store, vid, data_dir, ecj_path)? {
|
||||
// Mounted since the read: its handle owns the journal now.
|
||||
let ecv = store.locations[primary]
|
||||
.find_ec_volume_mut(vid)
|
||||
.expect("mounted journal runtime");
|
||||
let journal_path = ecv.ecj_file_name();
|
||||
let added = ecv.merge_journal(ids)?;
|
||||
publish_to_journal_siblings(&store, primary, vid, &journal_path, ids);
|
||||
return Ok(added);
|
||||
}
|
||||
if let Some(added) = append_ecj_ids(ecj_path, &local, ids, size)? {
|
||||
return Ok(added);
|
||||
}
|
||||
}
|
||||
Err(io::Error::other(format!(
|
||||
"ec volume {}: journal {} kept changing during merge",
|
||||
vid.0, ecj_path
|
||||
)))
|
||||
}
|
||||
|
||||
/// The disk index of the runtime holding `ecj_path` open, if any: the disk at
|
||||
/// `data_dir`'s own, else the first sibling journaling into it.
|
||||
fn mounted_ec_journal(
|
||||
store: &Store,
|
||||
vid: VolumeId,
|
||||
data_dir: &str,
|
||||
ecj_path: &str,
|
||||
) -> io::Result<Option<usize>> {
|
||||
let owner = store
|
||||
.locations
|
||||
.iter()
|
||||
.position(|loc| Path::new(&loc.directory) == Path::new(data_dir))
|
||||
.ok_or_else(|| {
|
||||
io::Error::other(format!(
|
||||
"ec volume {}: no disk at {} owns journal {}",
|
||||
vid.0, data_dir, ecj_path
|
||||
))
|
||||
})?;
|
||||
let runtime = if store.locations[owner].has_ec_volume(vid) {
|
||||
Some(owner)
|
||||
} else {
|
||||
store.locations.iter().position(|loc| {
|
||||
loc.find_ec_volume(vid)
|
||||
.is_some_and(|ecv| Path::new(&ecv.ecj_file_name()) == Path::new(ecj_path))
|
||||
})
|
||||
};
|
||||
Ok(runtime)
|
||||
}
|
||||
|
||||
/// Publishes merged ids into every other runtime journaling into `ecj_path`.
|
||||
fn publish_to_journal_siblings(
|
||||
store: &Store,
|
||||
primary: usize,
|
||||
vid: VolumeId,
|
||||
ecj_path: &str,
|
||||
ids: &HashSet<NeedleId>,
|
||||
) {
|
||||
for (i, loc) in store.locations.iter().enumerate() {
|
||||
if i == primary {
|
||||
continue;
|
||||
}
|
||||
if let Some(ecv) = loc.find_ec_volume(vid) {
|
||||
if Path::new(&ecv.ecj_file_name()) == Path::new(ecj_path) {
|
||||
ecv.publish_merged_ids(ids);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::config::MinFreeSpace;
|
||||
use crate::storage::needle_map::NeedleMapKind;
|
||||
use crate::storage::types::DiskType;
|
||||
use crate::storage::volume::{VifEcShardConfig, VifVolumeInfo};
|
||||
use tempfile::TempDir;
|
||||
|
||||
const COLLECTION: &str = "c";
|
||||
const VID: VolumeId = VolumeId(9);
|
||||
|
||||
/// A store with one disk per entry of `data`, all sharing `idx` when given,
|
||||
/// else each indexing into its own data dir.
|
||||
fn make_store(tmp: &TempDir, data: &[&str], idx: Option<&str>) -> RwLock<Store> {
|
||||
let mut store = Store::new(NeedleMapKind::InMemory);
|
||||
for d in data {
|
||||
let dir = tmp.path().join(d).to_string_lossy().into_owned();
|
||||
let idx_dir = idx
|
||||
.map(|i| tmp.path().join(i).to_string_lossy().into_owned())
|
||||
.unwrap_or_else(|| dir.clone());
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
std::fs::create_dir_all(&idx_dir).unwrap();
|
||||
store
|
||||
.add_location(
|
||||
&dir,
|
||||
&idx_dir,
|
||||
100,
|
||||
DiskType::HardDrive,
|
||||
MinFreeSpace::Percent(0.0),
|
||||
Vec::new(),
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
RwLock::new(store)
|
||||
}
|
||||
|
||||
fn dir(tmp: &TempDir, d: &str) -> String {
|
||||
tmp.path().join(d).to_string_lossy().into_owned()
|
||||
}
|
||||
|
||||
fn records(ids: &[u64]) -> Vec<u8> {
|
||||
let ids: Vec<NeedleId> = ids.iter().copied().map(NeedleId).collect();
|
||||
crate::storage::erasure_coding::ecj_merge::encode_ecj_ids(&ids)
|
||||
}
|
||||
|
||||
fn id_set(ids: &[u64]) -> HashSet<NeedleId> {
|
||||
ids.iter().copied().map(NeedleId).collect()
|
||||
}
|
||||
|
||||
/// Shard 0 of `VID` and its `.vif` in `data_dir`.
|
||||
fn write_shard0(data_dir: &str) {
|
||||
let base = format!("{}/{}_{}", data_dir, COLLECTION, VID.0);
|
||||
std::fs::write(format!("{}.ec00", base), b"shard data nonempty").unwrap();
|
||||
let vif = VifVolumeInfo {
|
||||
version: 3,
|
||||
ec_shard_config: Some(VifEcShardConfig {
|
||||
data_shards: 10,
|
||||
parity_shards: 4,
|
||||
..Default::default()
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
std::fs::write(
|
||||
format!("{}.vif", base),
|
||||
serde_json::to_string(&vif).unwrap(),
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// `VID`'s `.ecx` and a `.ecj` holding `deleted` in `dir`; returns the
|
||||
/// journal path.
|
||||
fn write_index(dir: &str, deleted: &[u64]) -> String {
|
||||
let base = format!("{}/{}_{}", dir, COLLECTION, VID.0);
|
||||
std::fs::write(format!("{}.ecx", base), vec![0u8; 16]).unwrap();
|
||||
let ecj = format!("{}.ecj", base);
|
||||
std::fs::write(&ecj, records(deleted)).unwrap();
|
||||
ecj
|
||||
}
|
||||
|
||||
fn deleted_on(store: &RwLock<Store>, disk: usize, id: u64) -> bool {
|
||||
store.read().unwrap().locations[disk]
|
||||
.find_ec_volume(VID)
|
||||
.expect("mounted")
|
||||
.is_needle_deleted(NeedleId(id))
|
||||
}
|
||||
|
||||
/// Disks sharing one index directory all hold the same journal path.
|
||||
/// Copying shards onto a disk that has not mounted `vid` must still reach
|
||||
/// the sibling runtime holding that journal open.
|
||||
#[test]
|
||||
fn shared_index_dir_reaches_sibling_mount() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let store = make_store(&tmp, &["d0", "d1"], Some("idx"));
|
||||
write_shard0(&dir(&tmp, "d0"));
|
||||
let ecj = write_index(&dir(&tmp, "idx"), &[1]);
|
||||
store.write().unwrap().locations[0]
|
||||
.mount_ec_shards(VID, COLLECTION, &[0], "")
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
store.read().unwrap().locations[0]
|
||||
.find_ec_volume(VID)
|
||||
.unwrap()
|
||||
.ecj_file_name(),
|
||||
ecj
|
||||
);
|
||||
|
||||
let added =
|
||||
merge_ec_journal(&store, VID, &dir(&tmp, "d1"), &ecj, &id_set(&[1, 2])).unwrap();
|
||||
assert_eq!(added, 1);
|
||||
assert!(
|
||||
deleted_on(&store, 0, 2),
|
||||
"the mounted sibling must see id 2"
|
||||
);
|
||||
assert_eq!(std::fs::read(&ecj).unwrap(), records(&[1, 2]));
|
||||
}
|
||||
|
||||
/// Every runtime holding the journal open must see merged ids in memory.
|
||||
#[test]
|
||||
fn shared_journal_reaches_every_holder() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let store = make_store(&tmp, &["d0", "d1"], Some("idx"));
|
||||
write_shard0(&dir(&tmp, "d0"));
|
||||
write_shard0(&dir(&tmp, "d1"));
|
||||
let ecj = write_index(&dir(&tmp, "idx"), &[1]);
|
||||
for i in 0..2 {
|
||||
store.write().unwrap().locations[i]
|
||||
.mount_ec_shards(VID, COLLECTION, &[0], "")
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
let added =
|
||||
merge_ec_journal(&store, VID, &dir(&tmp, "d1"), &ecj, &id_set(&[1, 2])).unwrap();
|
||||
assert_eq!(added, 1);
|
||||
assert!(
|
||||
deleted_on(&store, 0, 2) && deleted_on(&store, 1, 2),
|
||||
"every journal holder must see the merged id"
|
||||
);
|
||||
assert_eq!(std::fs::read(&ecj).unwrap(), records(&[1, 2]));
|
||||
}
|
||||
|
||||
/// The picked runtime may journal to a different file than the copied
|
||||
/// one — its index lives in its data directory while a sibling's lives
|
||||
/// in the index directory. The ids must be published only to holders of
|
||||
/// the file they were written to.
|
||||
#[test]
|
||||
fn publishes_to_actual_journal_holders() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let store = make_store(&tmp, &["d0", "d1"], Some("idx"));
|
||||
write_shard0(&dir(&tmp, "d0"));
|
||||
write_shard0(&dir(&tmp, "d1"));
|
||||
let data_ecj = write_index(&dir(&tmp, "d0"), &[1]);
|
||||
let idx_ecj = write_index(&dir(&tmp, "idx"), &[1]);
|
||||
for i in 0..2 {
|
||||
store.write().unwrap().locations[i]
|
||||
.mount_ec_shards(VID, COLLECTION, &[0], "")
|
||||
.unwrap();
|
||||
}
|
||||
assert_eq!(
|
||||
store.read().unwrap().locations[0]
|
||||
.find_ec_volume(VID)
|
||||
.unwrap()
|
||||
.ecj_file_name(),
|
||||
data_ecj
|
||||
);
|
||||
|
||||
let added =
|
||||
merge_ec_journal(&store, VID, &dir(&tmp, "d0"), &idx_ecj, &id_set(&[1, 2])).unwrap();
|
||||
assert_eq!(added, 1);
|
||||
assert_eq!(std::fs::read(&data_ecj).unwrap(), records(&[1, 2]));
|
||||
assert_eq!(std::fs::read(&idx_ecj).unwrap(), records(&[1]));
|
||||
assert!(deleted_on(&store, 0, 2));
|
||||
assert!(
|
||||
!deleted_on(&store, 1, 2),
|
||||
"a different journal's holder must not claim the merged id"
|
||||
);
|
||||
}
|
||||
|
||||
/// A sibling disk can mount `vid` from the receiving disk's index (#9212)
|
||||
/// while the merge reads the journal unlocked. The merge must go through
|
||||
/// that mount and report what it added.
|
||||
#[test]
|
||||
fn mount_during_read_is_merged_through_and_counted() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let store = make_store(&tmp, &["d0", "d1"], None);
|
||||
let owner = dir(&tmp, "d0");
|
||||
let ecj = write_index(&owner, &[1]);
|
||||
write_shard0(&dir(&tmp, "d1"));
|
||||
|
||||
let mut mounted = false;
|
||||
let read = |path: &str| {
|
||||
let read = read_ecj_ids(path);
|
||||
if !mounted {
|
||||
store.write().unwrap().locations[1]
|
||||
.mount_ec_shards_with_idx_dir(VID, COLLECTION, &[0], &owner, "")
|
||||
.unwrap();
|
||||
mounted = true;
|
||||
}
|
||||
read
|
||||
};
|
||||
let added =
|
||||
merge_ec_journal_with(&store, VID, &owner, &ecj, &id_set(&[1, 2, 3]), read).unwrap();
|
||||
assert_eq!(added, 2, "the ids merged through the new mount are counted");
|
||||
assert!(deleted_on(&store, 1, 2) && deleted_on(&store, 1, 3));
|
||||
assert_eq!(std::fs::read(&ecj).unwrap(), records(&[1, 2, 3]));
|
||||
}
|
||||
|
||||
/// An unmounted journal is merged on disk and a repeat adds nothing; a
|
||||
/// data dir that is no disk is refused.
|
||||
#[test]
|
||||
fn unmounted_journal_is_idempotent() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let store = make_store(&tmp, &["d0"], Some("idx"));
|
||||
let ecj = write_index(&dir(&tmp, "idx"), &[1, 2]);
|
||||
let d0 = dir(&tmp, "d0");
|
||||
for want in [2, 0, 0] {
|
||||
let added = merge_ec_journal(&store, VID, &d0, &ecj, &id_set(&[2, 3, 4])).unwrap();
|
||||
assert_eq!(added, want);
|
||||
}
|
||||
assert_eq!(std::fs::read(&ecj).unwrap(), records(&[1, 2, 3, 4]));
|
||||
assert!(
|
||||
merge_ec_journal(&store, VID, &dir(&tmp, "elsewhere"), &ecj, &id_set(&[1])).is_err()
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -5,8 +5,8 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"math"
|
||||
"io/fs"
|
||||
"math"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
@@ -437,8 +437,9 @@ func (vs *VolumeServer) VolumeEcShardsCopy(ctx context.Context, req *volume_serv
|
||||
}
|
||||
|
||||
if req.CopyEcjFile {
|
||||
// copy ecj file
|
||||
if _, err := vs.doCopyFileWithThrottler(client, true, req.Collection, req.VolumeId, math.MaxUint32, math.MaxInt64, indexBaseFileName, ".ecj", true, true, nil, throttler); err != nil {
|
||||
// The journal is a *set* of ids: merge the source's into the
|
||||
// local one as a union, never append it whole.
|
||||
if err := vs.copyEcjAndMerge(client, req.Collection, req.VolumeId, location.Directory, indexBaseFileName, throttler); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
@@ -469,6 +470,61 @@ func (vs *VolumeServer) VolumeEcShardsCopy(ctx context.Context, req *volume_serv
|
||||
return &volume_server_pb.VolumeEcShardsCopyResponse{}, nil
|
||||
}
|
||||
|
||||
// copyEcjAndMerge folds the source peer's .ecj into the local journal of vid
|
||||
// as a set union: only ids the local journal lacks are appended, so a
|
||||
// shard bounced between servers cannot grow it. The source journal streams
|
||||
// straight into memory — no staging file — and a source without one is not an
|
||||
// error. destDir is the receiving disk's data directory and destBase its index
|
||||
// base name.
|
||||
func (vs *VolumeServer) copyEcjAndMerge(client volume_server_pb.VolumeServerClient, collection string, vid uint32, destDir, destBase string, throttler *util.WriteThrottler) error {
|
||||
stream, err := client.CopyFile(context.Background(), &volume_server_pb.CopyFileRequest{
|
||||
VolumeId: vid,
|
||||
Ext: ".ecj",
|
||||
CompactionRevision: math.MaxUint32,
|
||||
StopOffset: math.MaxInt64,
|
||||
Collection: collection,
|
||||
IsEcVolume: true,
|
||||
IgnoreSourceFileNotFound: true,
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("volume %d: start copying .ecj: %w", vid, err)
|
||||
}
|
||||
ids, found, err := receiveEcjIds(stream, throttler)
|
||||
if err != nil {
|
||||
return fmt.Errorf("volume %d: copy .ecj: %w", vid, err)
|
||||
}
|
||||
if !found {
|
||||
return nil
|
||||
}
|
||||
if _, err := vs.store.MergeEcJournal(needle.VolumeId(vid), destDir, destBase+".ecj", ids); err != nil {
|
||||
return fmt.Errorf("volume %d: merge .ecj: %w", vid, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// receiveEcjIds decodes a CopyFile stream of an .ecj into its distinct ids.
|
||||
// found is false only when the source has no journal, which the source
|
||||
// signals with neither a modified time nor any bytes; an empty journal still
|
||||
// carries its modified time.
|
||||
func receiveEcjIds(stream volume_server_pb.VolumeServer_CopyFileClient, throttler *util.WriteThrottler) (ids map[types.NeedleId]struct{}, found bool, err error) {
|
||||
decoder := erasure_coding.NewEcjIdDecoder()
|
||||
for {
|
||||
resp, recvErr := stream.Recv()
|
||||
if recvErr == io.EOF {
|
||||
break
|
||||
}
|
||||
if recvErr != nil {
|
||||
return nil, false, recvErr
|
||||
}
|
||||
if resp.ModifiedTsNs != 0 || len(resp.FileContent) > 0 {
|
||||
found = true
|
||||
}
|
||||
decoder.Write(resp.FileContent)
|
||||
throttler.MaybeSlowdown(int64(len(resp.FileContent)))
|
||||
}
|
||||
return decoder.Ids(), found, nil
|
||||
}
|
||||
|
||||
// VolumeEcShardsDelete local delete the .ecx and some ec data slices if not needed
|
||||
// the shard should not be mounted before calling this.
|
||||
// Allowed in maintenance mode: like VolumeDelete it only removes data, and
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func ecjStreamBytes(ids ...types.NeedleId) []byte {
|
||||
b := make([]byte, len(ids)*types.NeedleIdSize)
|
||||
for i, id := range ids {
|
||||
types.NeedleIdToBytes(b[i*types.NeedleIdSize:], id)
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
func TestReceiveEcjIds(t *testing.T) {
|
||||
journal := ecjStreamBytes(1, 2, 3, 2)
|
||||
tests := []struct {
|
||||
name string
|
||||
responses []*volume_server_pb.CopyFileResponse
|
||||
wantIds []types.NeedleId
|
||||
wantFound bool
|
||||
}{
|
||||
{
|
||||
name: "source has no journal",
|
||||
responses: nil,
|
||||
},
|
||||
{
|
||||
name: "empty journal still carries its modified time",
|
||||
responses: []*volume_server_pb.CopyFileResponse{{ModifiedTsNs: 42}},
|
||||
wantFound: true,
|
||||
},
|
||||
{
|
||||
// A source that cannot read its journal's mtime reports 0; its
|
||||
// bytes still count (the Rust server sends unwrap_or(0)).
|
||||
name: "bytes without a modified time",
|
||||
responses: []*volume_server_pb.CopyFileResponse{{FileContent: journal}},
|
||||
wantIds: []types.NeedleId{1, 2, 3},
|
||||
wantFound: true,
|
||||
},
|
||||
{
|
||||
name: "records split across chunks",
|
||||
responses: []*volume_server_pb.CopyFileResponse{
|
||||
{FileContent: journal[:5], ModifiedTsNs: 7},
|
||||
{FileContent: journal[5:19]},
|
||||
{FileContent: journal[19:]},
|
||||
},
|
||||
wantIds: []types.NeedleId{1, 2, 3},
|
||||
wantFound: true,
|
||||
},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
stream := &fakeCopyFileStream{responses: tt.responses}
|
||||
ids, found, err := receiveEcjIds(stream, util.NewWriteThrottler(0))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, tt.wantFound, found)
|
||||
got := make([]types.NeedleId, 0, len(ids))
|
||||
for id := range ids {
|
||||
got = append(got, id)
|
||||
}
|
||||
assert.ElementsMatch(t, tt.wantIds, got)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestReceiveEcjIds_StreamError(t *testing.T) {
|
||||
stream := &fakeCopyFileStream{
|
||||
responses: []*volume_server_pb.CopyFileResponse{{FileContent: ecjStreamBytes(1), ModifiedTsNs: 1}},
|
||||
finalErr: errors.New("peer went away"),
|
||||
}
|
||||
_, _, err := receiveEcjIds(stream, util.NewWriteThrottler(0))
|
||||
assert.Error(t, err)
|
||||
}
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
// ecRecoveryLookupTimeout bounds the master LookupEcVolume call so a slow or
|
||||
@@ -119,7 +120,6 @@ func (vs *VolumeServer) fetchEcIndexFromPeers(peers []pb.ServerAddress, m storag
|
||||
idxBaseFileName := storage.VolumeFileName(m.IdxDir, m.Collection, int(m.VolumeId))
|
||||
dataBaseFileName := storage.VolumeFileName(m.DataDir, m.Collection, int(m.VolumeId))
|
||||
ecxPath := idxBaseFileName + ".ecx"
|
||||
ecjPath := idxBaseFileName + ".ecj"
|
||||
|
||||
removePartial := func(path string) {
|
||||
if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
|
||||
@@ -145,12 +145,12 @@ func (vs *VolumeServer) fetchEcIndexFromPeers(peers []pb.ServerAddress, m storag
|
||||
}
|
||||
// .ecj is the source peer's deletion journal; .vif carries EC params
|
||||
// and EncodeTsNs. Both are best-effort: a missing .ecj is recreated at
|
||||
// mount and a missing .vif falls back to default EC parameters. A failed
|
||||
// .ecj copy is an in-place append, so drop the partial file; the .vif
|
||||
// copy stages and renames, leaving nothing to clean up.
|
||||
if _, err := vs.doCopyFile(client, true, m.Collection, uint32(m.VolumeId), math.MaxUint32, math.MaxInt64, idxBaseFileName, ".ecj", true, true, nil); err != nil {
|
||||
// mount and a missing .vif falls back to default EC parameters. The
|
||||
// journal is a *set*: merge it as a union so a bounced volume
|
||||
// cannot double it. The merge only appends whole records, and the
|
||||
// .vif copy stages and renames, so a failure leaves nothing to clean.
|
||||
if err := vs.copyEcjAndMerge(client, m.Collection, uint32(m.VolumeId), m.DataDir, idxBaseFileName, util.NewWriteThrottler(vs.maintenanceBytePerSecond)); err != nil {
|
||||
glog.Warningf("ec volume %d: copy .ecj from %s: %v", m.VolumeId, peer, err)
|
||||
removePartial(ecjPath)
|
||||
}
|
||||
if _, err := vs.doCopyFile(client, true, m.Collection, uint32(m.VolumeId), math.MaxUint32, math.MaxInt64, dataBaseFileName, ".vif", false, true, nil); err != nil {
|
||||
glog.Warningf("ec volume %d: copy .vif from %s: %v", m.VolumeId, peer, err)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"net"
|
||||
"os"
|
||||
@@ -109,7 +110,7 @@ func TestFetchEcIndexFromPeers_CopiesIndexOverGrpc(t *testing.T) {
|
||||
if err := os.WriteFile(srcBase+".ecx", ecxBytes, 0o644); err != nil {
|
||||
t.Fatalf("write source .ecx: %v", err)
|
||||
}
|
||||
if err := os.WriteFile(srcBase+".ecj", []byte("journal"), 0o644); err != nil {
|
||||
if err := os.WriteFile(srcBase+".ecj", ecjStreamBytes(3, 5, 3), 0o644); err != nil {
|
||||
t.Fatalf("write source .ecj: %v", err)
|
||||
}
|
||||
// A real .vif: the receiver MOUNTS the volume from the copied files, and a
|
||||
@@ -160,8 +161,11 @@ func TestFetchEcIndexFromPeers_CopiesIndexOverGrpc(t *testing.T) {
|
||||
if len(got) != len(ecxBytes) {
|
||||
t.Errorf("copied .ecx size = %d, want %d", len(got), len(ecxBytes))
|
||||
}
|
||||
if _, err := os.Stat(dstBase + ".ecj"); err != nil {
|
||||
// The journal is merged as a set: the duplicate record collapses.
|
||||
if got, err := os.ReadFile(dstBase + ".ecj"); err != nil {
|
||||
t.Errorf("copied .ecj missing: %v", err)
|
||||
} else if want := ecjStreamBytes(3, 5); !bytes.Equal(got, want) {
|
||||
t.Errorf("copied .ecj = %x, want %x", got, want)
|
||||
}
|
||||
if _, err := os.Stat(dstBase + ".vif"); err != nil {
|
||||
t.Errorf("copied .vif missing: %v", err)
|
||||
|
||||
@@ -74,7 +74,21 @@ func (ev *EcVolume) DeleteNeedleFromEcx(needleId types.NeedleId) (err error) {
|
||||
|
||||
b := make([]byte, types.NeedleIdSize)
|
||||
types.NeedleIdToBytes(b, needleId)
|
||||
if err := ev.appendJournalLocked(b); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Publish into the in-memory set only after the journal is durable.
|
||||
ev.markNeedleDeletedInMemory(needleId)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// appendJournalLocked appends whole records to .ecj and syncs them. A partial
|
||||
// write is truncated back to the known-good size so the on-disk journal and
|
||||
// deletedNeedles cannot drift. Callers hold ecjFileAccessLock and have checked
|
||||
// that ecjFile is open.
|
||||
func (ev *EcVolume) appendJournalLocked(b []byte) error {
|
||||
prevEcjSize := ev.ecjFileSize
|
||||
if _, seekErr := ev.ecjFile.Seek(0, io.SeekEnd); seekErr != nil {
|
||||
return fmt.Errorf("seek ecj: %w", seekErr)
|
||||
@@ -93,10 +107,6 @@ func (ev *EcVolume) DeleteNeedleFromEcx(needleId types.NeedleId) (err error) {
|
||||
return fmt.Errorf("sync ecj: %w", syncErr)
|
||||
}
|
||||
ev.ecjFileSize += int64(n)
|
||||
|
||||
// Publish into the in-memory set only after the journal is durable.
|
||||
ev.markNeedleDeletedInMemory(needleId)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,335 @@
|
||||
package erasure_coding
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"slices"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
// An EC volume's deletion journal (<vid>.ecj) is a *set* of deleted needle ids
|
||||
// stored as 8-byte big-endian records. Shard copy and index recovery fold a
|
||||
// peer's journal into the local one; they must append only the ids the local
|
||||
// journal lacks, or every ec_balance round trip doubles the file.
|
||||
//
|
||||
// The journal is only ever appended to, never replaced: a mounted EcVolume
|
||||
// holds it open, and a rename would leave that handle writing to an unlinked
|
||||
// inode, losing every later delete at the next mount.
|
||||
|
||||
// ErrEcjChanged reports that a journal's length moved between the read that
|
||||
// computed a merge delta and the append that would publish it.
|
||||
var ErrEcjChanged = errors.New("ec journal changed during merge")
|
||||
|
||||
// EcjIdDecoder decodes .ecj records from a byte stream split at arbitrary
|
||||
// boundaries (a CopyFile stream, chunked reads), collecting the distinct ids.
|
||||
// Memory follows the number of distinct ids, not the journal's length. A
|
||||
// trailing partial record is never decoded.
|
||||
type EcjIdDecoder struct {
|
||||
ids map[types.NeedleId]struct{}
|
||||
partial [types.NeedleIdSize]byte
|
||||
pending int
|
||||
}
|
||||
|
||||
func NewEcjIdDecoder() *EcjIdDecoder {
|
||||
return &EcjIdDecoder{ids: make(map[types.NeedleId]struct{})}
|
||||
}
|
||||
|
||||
func (d *EcjIdDecoder) Write(p []byte) {
|
||||
if d.pending > 0 {
|
||||
n := copy(d.partial[d.pending:], p)
|
||||
d.pending += n
|
||||
p = p[n:]
|
||||
if d.pending < types.NeedleIdSize {
|
||||
return
|
||||
}
|
||||
d.ids[types.BytesToNeedleId(d.partial[:])] = struct{}{}
|
||||
d.pending = 0
|
||||
}
|
||||
whole := len(p) - len(p)%types.NeedleIdSize
|
||||
for i := 0; i < whole; i += types.NeedleIdSize {
|
||||
d.ids[types.BytesToNeedleId(p[i:i+types.NeedleIdSize])] = struct{}{}
|
||||
}
|
||||
d.pending = copy(d.partial[:], p[whole:])
|
||||
}
|
||||
|
||||
func (d *EcjIdDecoder) Ids() map[types.NeedleId]struct{} {
|
||||
return d.ids
|
||||
}
|
||||
|
||||
// ReadEcjIds reads the distinct ids of the journal at path in bounded chunks.
|
||||
// A missing file reads as empty. size is the whole-record length read; a torn
|
||||
// trailing partial record is excluded from it.
|
||||
func ReadEcjIds(path string) (ids map[types.NeedleId]struct{}, size int64, err error) {
|
||||
f, err := os.Open(path)
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
return make(map[types.NeedleId]struct{}), 0, nil
|
||||
}
|
||||
return nil, 0, err
|
||||
}
|
||||
defer f.Close()
|
||||
fi, err := f.Stat()
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
size = fi.Size() - fi.Size()%int64(types.NeedleIdSize)
|
||||
d := NewEcjIdDecoder()
|
||||
buf := make([]byte, min(int64(ecjLoadChunkBytes), size))
|
||||
for off := int64(0); off < size; {
|
||||
want := min(int64(len(buf)), size-off)
|
||||
if _, err := f.ReadAt(buf[:want], off); err != nil {
|
||||
return nil, 0, fmt.Errorf("read %s at %d: %w", path, off, err)
|
||||
}
|
||||
d.Write(buf[:want])
|
||||
off += want
|
||||
}
|
||||
return d.Ids(), size, nil
|
||||
}
|
||||
|
||||
// ecjDelta returns the ids of incoming that has does not report, ascending,
|
||||
// so a merge appends deterministic output.
|
||||
func ecjDelta(incoming map[types.NeedleId]struct{}, has func(types.NeedleId) bool) []types.NeedleId {
|
||||
var delta []types.NeedleId
|
||||
for id := range incoming {
|
||||
if !has(id) {
|
||||
delta = append(delta, id)
|
||||
}
|
||||
}
|
||||
slices.Sort(delta)
|
||||
return delta
|
||||
}
|
||||
|
||||
func encodeEcjIds(ids []types.NeedleId) []byte {
|
||||
b := make([]byte, len(ids)*types.NeedleIdSize)
|
||||
for i, id := range ids {
|
||||
types.NeedleIdToBytes(b[i*types.NeedleIdSize:], id)
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
// AppendEcjIds appends to the journal at path the ids of incoming that local
|
||||
// lacks, in one write and one fsync, rolling the write back if the fsync
|
||||
// fails. It is WriteEcjIds followed by EcjAppend.Sync; see WriteEcjIds for the
|
||||
// contract.
|
||||
func AppendEcjIds(path string, local, incoming map[types.NeedleId]struct{}, size int64) (added int, err error) {
|
||||
a, err := WriteEcjIds(path, local, incoming, size)
|
||||
if a == nil || err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if err := a.Sync(); err != nil {
|
||||
a.Rollback()
|
||||
return 0, err
|
||||
}
|
||||
a.Close()
|
||||
return a.Added, nil
|
||||
}
|
||||
|
||||
// EcjAppend is a journal append WriteEcjIds wrote but did not sync. The owner
|
||||
// must end it with Close (after a successful Sync) or Rollback.
|
||||
type EcjAppend struct {
|
||||
Added int
|
||||
f *os.File
|
||||
path string
|
||||
size int64 // whole-record length before the append
|
||||
delta []types.NeedleId
|
||||
created bool
|
||||
}
|
||||
|
||||
// WriteEcjIds appends to the journal at path the ids of incoming that local
|
||||
// lacks, in one write, without syncing; it returns nil when there is nothing
|
||||
// to add. It is for a journal no EcVolume has open; a mounted volume merges
|
||||
// through EcVolume.MergeJournal instead. local and size come from ReadEcjIds
|
||||
// on the same path: if the journal's whole-record length is no longer size,
|
||||
// it returns ErrEcjChanged so the caller re-reads. A torn tail past size is
|
||||
// truncated first so the new records stay aligned.
|
||||
func WriteEcjIds(path string, local, incoming map[types.NeedleId]struct{}, size int64) (*EcjAppend, error) {
|
||||
delta := ecjDelta(incoming, func(id types.NeedleId) bool {
|
||||
_, ok := local[id]
|
||||
return ok
|
||||
})
|
||||
if len(delta) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
_, statErr := os.Stat(path)
|
||||
created := os.IsNotExist(statErr)
|
||||
f, err := backend.OpenVolumeFile(path, os.O_RDWR|os.O_CREATE)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open %s: %w", path, err)
|
||||
}
|
||||
a := &EcjAppend{Added: len(delta), f: f, path: path, size: size, delta: delta, created: created}
|
||||
fi, err := f.Stat()
|
||||
if err != nil {
|
||||
a.Close()
|
||||
return nil, fmt.Errorf("stat %s: %w", path, err)
|
||||
}
|
||||
if fi.Size()-fi.Size()%int64(types.NeedleIdSize) != size {
|
||||
a.Close()
|
||||
return nil, ErrEcjChanged
|
||||
}
|
||||
if fi.Size() != size {
|
||||
if err := f.Truncate(size); err != nil {
|
||||
a.Close()
|
||||
return nil, fmt.Errorf("truncate torn tail of %s: %w", path, err)
|
||||
}
|
||||
}
|
||||
if _, err := f.WriteAt(encodeEcjIds(delta), size); err != nil {
|
||||
_ = f.Truncate(size)
|
||||
a.Close()
|
||||
return nil, fmt.Errorf("append %s: %w", path, err)
|
||||
}
|
||||
return a, nil
|
||||
}
|
||||
|
||||
// Sync makes the append durable, including the journal's directory entry if
|
||||
// the append created it.
|
||||
func (a *EcjAppend) Sync() error {
|
||||
if err := a.f.Sync(); err != nil {
|
||||
return fmt.Errorf("sync %s: %w", a.path, err)
|
||||
}
|
||||
if a.created {
|
||||
if err := util.FsyncDir(filepath.Dir(a.path)); err != nil {
|
||||
return fmt.Errorf("fsync dir for %s: %w", a.path, err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Rollback removes the append, for a journal no EcVolume has open, and closes
|
||||
// it. See RollbackUnsynced.
|
||||
func (a *EcjAppend) Rollback() {
|
||||
a.RollbackUnsynced(nil)
|
||||
a.Close()
|
||||
}
|
||||
|
||||
// RollbackUnsynced removes an append whose fsync failed and reports whether it
|
||||
// did. holders are the mounted volumes that opened the journal since the write
|
||||
// (none can predate it: the write happens only while none is mounted); each
|
||||
// loaded the append's records, and they leave its in-memory set too, so disk
|
||||
// and memory agree the merge did not happen and a retried merge appends and
|
||||
// syncs them again. The caller must exclude new mounts and other unmounted
|
||||
// merges into the path.
|
||||
//
|
||||
// Invariant: only the journal's actual length decides, never a holder's cached
|
||||
// ecjFileSize, which another holder's appends leave stale. With every holder's
|
||||
// ecjFileAccessLock held nothing can append, so if the length is still this
|
||||
// append's end, its records are the tail and removing them loses nothing. If
|
||||
// anything follows them, it changes nothing and returns false; the caller then
|
||||
// keeps the records and rewrites and syncs them (Rewrite, Sync), or, if that
|
||||
// fails too, takes the ids back out of the holders' sets with Unpublish.
|
||||
func (a *EcjAppend) RollbackUnsynced(holders []*EcVolume) bool {
|
||||
for _, ev := range holders {
|
||||
ev.ecjFileAccessLock.Lock()
|
||||
defer ev.ecjFileAccessLock.Unlock()
|
||||
}
|
||||
if fi, err := a.f.Stat(); err != nil || fi.Size() != a.end() {
|
||||
return false
|
||||
}
|
||||
if err := a.f.Truncate(a.size); err != nil {
|
||||
return false
|
||||
}
|
||||
for _, ev := range holders {
|
||||
if ev.ecjFile != nil { // closed: it serves nothing and journals nothing
|
||||
ev.ecjFileSize = a.size
|
||||
}
|
||||
}
|
||||
a.forget(holders)
|
||||
return true
|
||||
}
|
||||
|
||||
// Unpublish takes the append's ids back out of the holders' deleted sets, for
|
||||
// records that could neither be removed nor made durable. The records stay in
|
||||
// the file, where they are legitimate deletes, but as in DeleteNeedleFromEcx a
|
||||
// set only claims ids whose record is known durable: a retried merge then
|
||||
// sees them missing and appends and syncs them again.
|
||||
func (a *EcjAppend) Unpublish(holders []*EcVolume) {
|
||||
for _, ev := range holders {
|
||||
ev.ecjFileAccessLock.Lock()
|
||||
defer ev.ecjFileAccessLock.Unlock()
|
||||
}
|
||||
a.forget(holders)
|
||||
}
|
||||
|
||||
// forget removes the append's ids from the holders' deleted sets. Each holder
|
||||
// loaded exactly the ids read before the append, which exclude the delta, plus
|
||||
// the delta; it cannot have journaled a delta id itself, as it already held
|
||||
// them all. So every delta id in a holder's set came from this append. Callers
|
||||
// hold every holder's ecjFileAccessLock.
|
||||
func (a *EcjAppend) forget(holders []*EcVolume) {
|
||||
for _, ev := range holders {
|
||||
ev.deletedNeedlesLock.Lock()
|
||||
for _, id := range a.delta {
|
||||
delete(ev.deletedNeedles, id)
|
||||
}
|
||||
ev.deletedNeedlesLock.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
// Rewrite writes the append's records again in place, for records that later
|
||||
// ones now follow and so cannot be removed; a Sync after it makes them
|
||||
// durable. A bare second fsync would prove nothing: after a failed writeback
|
||||
// the kernel may have dropped the pages or marked them clean and still report
|
||||
// the next fsync as clean. Rewriting the same bytes dirties them again.
|
||||
// Writing outside the locks is safe: appends land past these records, and
|
||||
// every other truncate (a holder's failed append, a mount's torn-tail repair)
|
||||
// lands at or past their end, since each holder loaded at least that much.
|
||||
func (a *EcjAppend) Rewrite() error {
|
||||
if fi, err := a.f.Stat(); err != nil {
|
||||
return fmt.Errorf("stat %s: %w", a.path, err)
|
||||
} else if fi.Size() < a.end() {
|
||||
return fmt.Errorf("%s shrank below the merged records", a.path)
|
||||
}
|
||||
if _, err := a.f.WriteAt(encodeEcjIds(a.delta), a.size); err != nil {
|
||||
return fmt.Errorf("rewrite %s: %w", a.path, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (a *EcjAppend) end() int64 {
|
||||
return a.size + int64(len(a.delta)*types.NeedleIdSize)
|
||||
}
|
||||
|
||||
// Close releases the journal, keeping whatever the append wrote.
|
||||
func (a *EcjAppend) Close() {
|
||||
_ = a.f.Close()
|
||||
}
|
||||
|
||||
// MergeJournal folds a peer's deletion journal into this mounted volume. Under
|
||||
// ecjFileAccessLock it appends only the ids not already deleted here, in one
|
||||
// write and one fsync, then publishes them into the in-memory set — the same
|
||||
// commit order as DeleteNeedleFromEcx, which it serializes with. The live
|
||||
// handle is appended to in place, so no delete can land in an orphaned file.
|
||||
func (ev *EcVolume) MergeJournal(ids map[types.NeedleId]struct{}) (added int, err error) {
|
||||
ev.ecjFileAccessLock.Lock()
|
||||
defer ev.ecjFileAccessLock.Unlock()
|
||||
if ev.ecjFile == nil {
|
||||
return 0, fmt.Errorf("ec volume %d closed", ev.VolumeId)
|
||||
}
|
||||
delta := ecjDelta(ids, ev.IsNeedleDeleted)
|
||||
if len(delta) == 0 {
|
||||
return 0, nil
|
||||
}
|
||||
if err := ev.appendJournalLocked(encodeEcjIds(delta)); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
ev.deletedNeedlesLock.Lock()
|
||||
for _, id := range delta {
|
||||
ev.deletedNeedles[id] = struct{}{}
|
||||
}
|
||||
ev.deletedNeedlesLock.Unlock()
|
||||
return len(delta), nil
|
||||
}
|
||||
|
||||
// PublishMergedIds marks ids deleted in memory only, for runtimes sharing a
|
||||
// journal another runtime already appended them to.
|
||||
func (ev *EcVolume) PublishMergedIds(ids map[types.NeedleId]struct{}) {
|
||||
ev.deletedNeedlesLock.Lock()
|
||||
for id := range ids {
|
||||
ev.deletedNeedles[id] = struct{}{}
|
||||
}
|
||||
ev.deletedNeedlesLock.Unlock()
|
||||
}
|
||||
@@ -0,0 +1,262 @@
|
||||
package erasure_coding_test
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
erasure_coding "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func idSet(ids ...types.NeedleId) map[types.NeedleId]struct{} {
|
||||
s := make(map[types.NeedleId]struct{}, len(ids))
|
||||
for _, id := range ids {
|
||||
s[id] = struct{}{}
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
func readEcjRecords(t *testing.T, path string) []types.NeedleId {
|
||||
t.Helper()
|
||||
data, err := os.ReadFile(path)
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, len(data)%types.NeedleIdSize, "journal must stay aligned")
|
||||
var ids []types.NeedleId
|
||||
for i := 0; i+types.NeedleIdSize <= len(data); i += types.NeedleIdSize {
|
||||
ids = append(ids, types.BytesToNeedleId(data[i:i+types.NeedleIdSize]))
|
||||
}
|
||||
return ids
|
||||
}
|
||||
|
||||
// mergeFile runs the unmounted merge the way the store does: read, then append.
|
||||
func mergeFile(t *testing.T, path string, incoming map[types.NeedleId]struct{}) int {
|
||||
t.Helper()
|
||||
local, size, err := erasure_coding.ReadEcjIds(path)
|
||||
require.NoError(t, err)
|
||||
added, err := erasure_coding.AppendEcjIds(path, local, incoming, size)
|
||||
require.NoError(t, err)
|
||||
return added
|
||||
}
|
||||
|
||||
// A CopyFile stream splits records at arbitrary byte boundaries.
|
||||
func TestEcjIdDecoder_RecordsSplitAcrossChunks(t *testing.T) {
|
||||
stream := append(ecjBytes(1, 2, 3, 2, 0x0102030405060708), 9, 9, 9)
|
||||
for chunk := 1; chunk <= len(stream); chunk++ {
|
||||
d := erasure_coding.NewEcjIdDecoder()
|
||||
for off := 0; off < len(stream); off += chunk {
|
||||
d.Write(stream[off:min(off+chunk, len(stream))])
|
||||
}
|
||||
assert.Equal(t, idSet(1, 2, 3, 0x0102030405060708), d.Ids(), "chunk size %d", chunk)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadEcjIds(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
|
||||
ids, size, err := erasure_coding.ReadEcjIds(filepath.Join(dir, "missing.ecj"))
|
||||
require.NoError(t, err)
|
||||
assert.Empty(t, ids)
|
||||
assert.Zero(t, size)
|
||||
|
||||
torn := filepath.Join(dir, "torn.ecj")
|
||||
require.NoError(t, os.WriteFile(torn, append(ecjBytes(1, 2, 1), 7, 7, 7), 0644))
|
||||
ids, size, err = erasure_coding.ReadEcjIds(torn)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, idSet(1, 2), ids)
|
||||
assert.Equal(t, int64(3*types.NeedleIdSize), size, "size covers whole records only")
|
||||
|
||||
// A bloated journal repeating a few ids across several read chunks keeps
|
||||
// only the distinct ids.
|
||||
bloated := filepath.Join(dir, "bloated.ecj")
|
||||
var data []byte
|
||||
for i := 0; i < 300_000; i++ {
|
||||
data = append(data, ecjBytes(types.NeedleId(i%3))...)
|
||||
}
|
||||
require.NoError(t, os.WriteFile(bloated, data, 0644))
|
||||
ids, size, err = erasure_coding.ReadEcjIds(bloated)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, idSet(0, 1, 2), ids)
|
||||
assert.Equal(t, int64(len(data)), size)
|
||||
}
|
||||
|
||||
// Local {1,2,3} + source {3,4} => {1,2,3,4}: only the missing id is appended.
|
||||
func TestAppendEcjIds_AppendsOnlyMissingIds(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "vol.ecj")
|
||||
require.NoError(t, os.WriteFile(path, ecjBytes(1, 2, 3), 0644))
|
||||
|
||||
assert.Equal(t, 1, mergeFile(t, path, idSet(3, 4)))
|
||||
assert.Equal(t, []types.NeedleId{1, 2, 3, 4}, readEcjRecords(t, path))
|
||||
}
|
||||
|
||||
// Copying the same shard A->B->A->B keeps the journal size constant; the old
|
||||
// append path doubled it on every trip.
|
||||
func TestAppendEcjIds_RoundTripStaysConstant(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
a := filepath.Join(dir, "a.ecj")
|
||||
b := filepath.Join(dir, "b.ecj")
|
||||
require.NoError(t, os.WriteFile(a, ecjBytes(1, 2), 0644))
|
||||
require.NoError(t, os.WriteFile(b, ecjBytes(2, 3), 0644))
|
||||
for i := 0; i < 20; i++ {
|
||||
src, dst := a, b
|
||||
if i%2 == 1 {
|
||||
src, dst = b, a
|
||||
}
|
||||
ids, _, err := erasure_coding.ReadEcjIds(src)
|
||||
require.NoError(t, err)
|
||||
mergeFile(t, dst, ids)
|
||||
}
|
||||
for _, path := range []string{a, b} {
|
||||
assert.ElementsMatch(t, []types.NeedleId{1, 2, 3}, readEcjRecords(t, path))
|
||||
}
|
||||
}
|
||||
|
||||
func TestAppendEcjIds_NothingNewLeavesJournalAlone(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
missing := filepath.Join(dir, "missing.ecj")
|
||||
assert.Zero(t, mergeFile(t, missing, idSet()))
|
||||
assert.NoFileExists(t, missing, "an empty merge must not create a journal")
|
||||
|
||||
path := filepath.Join(dir, "vol.ecj")
|
||||
require.NoError(t, os.WriteFile(path, ecjBytes(1, 2), 0644))
|
||||
assert.Zero(t, mergeFile(t, path, idSet(2, 1)))
|
||||
assert.Equal(t, []types.NeedleId{1, 2}, readEcjRecords(t, path))
|
||||
}
|
||||
|
||||
func TestAppendEcjIds_RepairsTornTail(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "vol.ecj")
|
||||
require.NoError(t, os.WriteFile(path, append(ecjBytes(1, 2), 9, 9, 9), 0644))
|
||||
|
||||
assert.Equal(t, 1, mergeFile(t, path, idSet(3)))
|
||||
assert.Equal(t, []types.NeedleId{1, 2, 3}, readEcjRecords(t, path))
|
||||
}
|
||||
|
||||
// A journal that grew after the read must not receive a delta computed
|
||||
// against its old contents.
|
||||
func TestAppendEcjIds_RejectsChangedJournal(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "vol.ecj")
|
||||
require.NoError(t, os.WriteFile(path, ecjBytes(1), 0644))
|
||||
local, size, err := erasure_coding.ReadEcjIds(path)
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, os.WriteFile(path, ecjBytes(1, 5), 0644))
|
||||
|
||||
_, err = erasure_coding.AppendEcjIds(path, local, idSet(2), size)
|
||||
assert.ErrorIs(t, err, erasure_coding.ErrEcjChanged)
|
||||
assert.Equal(t, []types.NeedleId{1, 5}, readEcjRecords(t, path))
|
||||
}
|
||||
|
||||
// Rolling back an unsynced append removes exactly its records, but leaves the
|
||||
// journal alone once something has been appended after them.
|
||||
func TestWriteEcjIds_Rollback(t *testing.T) {
|
||||
for _, appendedAfter := range []bool{false, true} {
|
||||
path := filepath.Join(t.TempDir(), "vol.ecj")
|
||||
require.NoError(t, os.WriteFile(path, ecjBytes(1), 0644))
|
||||
local, size, err := erasure_coding.ReadEcjIds(path)
|
||||
require.NoError(t, err)
|
||||
|
||||
a, err := erasure_coding.WriteEcjIds(path, local, idSet(1, 2), size)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 1, a.Added)
|
||||
want := []types.NeedleId{1}
|
||||
if appendedAfter {
|
||||
f, err := os.OpenFile(path, os.O_WRONLY|os.O_APPEND, 0644)
|
||||
require.NoError(t, err)
|
||||
_, err = f.Write(ecjBytes(3))
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, f.Close())
|
||||
want = []types.NeedleId{1, 2, 3}
|
||||
}
|
||||
a.Rollback()
|
||||
assert.Equal(t, want, readEcjRecords(t, path), "appended after: %v", appendedAfter)
|
||||
}
|
||||
}
|
||||
|
||||
// A mounted volume merges through its own handle: the journal keeps its inode,
|
||||
// the in-memory set follows, and later deletes land in the same live file.
|
||||
func TestMergeJournal_MountedVolumeAppendsInPlace(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
var ecx []byte
|
||||
for id := types.NeedleId(1); id <= 5; id++ {
|
||||
ecx = append(ecx, makeNeedleMapEntry(id, types.ToOffset(int64(id)*8), types.Size(100))...)
|
||||
}
|
||||
ev, base := mountEcVolume(t, dir, ecx, ecjBytes(1))
|
||||
before, err := os.Stat(base + ".ecj")
|
||||
require.NoError(t, err)
|
||||
|
||||
added, err := ev.MergeJournal(idSet(1, 2, 3))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 2, added)
|
||||
for _, id := range []types.NeedleId{1, 2, 3} {
|
||||
assert.True(t, ev.IsNeedleDeleted(id), "id %d", id)
|
||||
}
|
||||
_, deleteCount := ev.FileAndDeleteCount()
|
||||
assert.Equal(t, uint64(3), deleteCount)
|
||||
|
||||
added, err = ev.MergeJournal(idSet(2, 3))
|
||||
require.NoError(t, err)
|
||||
assert.Zero(t, added, "a repeated merge appends nothing")
|
||||
|
||||
require.NoError(t, ev.DeleteNeedleFromEcx(4))
|
||||
after, err := os.Stat(base + ".ecj")
|
||||
require.NoError(t, err)
|
||||
assert.True(t, os.SameFile(before, after), "the journal must never be replaced under the open handle")
|
||||
assert.Equal(t, []types.NeedleId{1, 2, 3, 4}, readEcjRecords(t, base+".ecj"))
|
||||
ev.Close()
|
||||
|
||||
ev, _ = mountEcVolume(t, dir, ecx, nil)
|
||||
defer ev.Close()
|
||||
for _, id := range []types.NeedleId{1, 2, 3, 4} {
|
||||
assert.True(t, ev.IsNeedleDeleted(id), "id %d survives remount", id)
|
||||
}
|
||||
}
|
||||
|
||||
// Deletes racing a merge must all reach the journal, and each id only once.
|
||||
func TestMergeJournal_ConcurrentDeletesPersist(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
const n = 200
|
||||
var ecx []byte
|
||||
for id := types.NeedleId(1); id <= 2*n; id++ {
|
||||
ecx = append(ecx, makeNeedleMapEntry(id, types.ToOffset(int64(id)*8), types.Size(100))...)
|
||||
}
|
||||
ev, base := mountEcVolume(t, dir, ecx, nil)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
for id := types.NeedleId(1); id <= n; id++ {
|
||||
assert.NoError(t, ev.DeleteNeedleFromEcx(id))
|
||||
}
|
||||
}()
|
||||
for id := types.NeedleId(n/2 + 1); id <= 2*n; id += 10 {
|
||||
batch := idSet()
|
||||
for k := id; k < id+10 && k <= 2*n; k++ {
|
||||
batch[k] = struct{}{}
|
||||
}
|
||||
_, err := ev.MergeJournal(batch)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
wg.Wait()
|
||||
ev.Close()
|
||||
|
||||
records := readEcjRecords(t, base+".ecj")
|
||||
assert.Len(t, records, 2*n, "every id journaled exactly once")
|
||||
ids, _, err := erasure_coding.ReadEcjIds(base + ".ecj")
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, ids, 2*n)
|
||||
}
|
||||
|
||||
// A merge that loses a race with unmount must fail cleanly and must not bring
|
||||
// a destroyed journal back.
|
||||
func TestMergeJournal_ClosedVolume(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
ev, base := mountEcVolume(t, dir, makeNeedleMapEntry(1, types.ToOffset(8), types.Size(100)), nil)
|
||||
ev.Destroy()
|
||||
|
||||
_, err := ev.MergeJournal(idSet(1))
|
||||
assert.Error(t, err)
|
||||
assert.NoFileExists(t, base+".ecj")
|
||||
}
|
||||
@@ -0,0 +1,246 @@
|
||||
package storage
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
)
|
||||
|
||||
// ecjMergeAttempts bounds how often an unmounted merge re-reads a journal
|
||||
// that changed under it. Only a mount-delete-unmount between the read and the
|
||||
// append changes it, so one retry is nearly always enough.
|
||||
const ecjMergeAttempts = 5
|
||||
|
||||
// ecjMergeLocks serializes merges into the same journal path, so concurrent
|
||||
// copies of one volume (a balance racing a rebuild) cannot both append the
|
||||
// same delta.
|
||||
var ecjMergeLocks = struct {
|
||||
sync.Mutex
|
||||
byPath map[string]*ecjPathLock
|
||||
}{byPath: make(map[string]*ecjPathLock)}
|
||||
|
||||
type ecjPathLock struct {
|
||||
sync.Mutex
|
||||
refs int
|
||||
}
|
||||
|
||||
func lockEcjPath(path string) (unlock func()) {
|
||||
ecjMergeLocks.Lock()
|
||||
l := ecjMergeLocks.byPath[path]
|
||||
if l == nil {
|
||||
l = &ecjPathLock{}
|
||||
ecjMergeLocks.byPath[path] = l
|
||||
}
|
||||
l.refs++
|
||||
ecjMergeLocks.Unlock()
|
||||
|
||||
l.Lock()
|
||||
return func() {
|
||||
l.Unlock()
|
||||
ecjMergeLocks.Lock()
|
||||
if l.refs--; l.refs == 0 {
|
||||
delete(ecjMergeLocks.byPath, path)
|
||||
}
|
||||
ecjMergeLocks.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
// ecjMergeIO is the disk I/O an unmounted merge does outside every lock,
|
||||
// injectable so tests can act between the steps.
|
||||
type ecjMergeIO struct {
|
||||
read func(path string) (ids map[types.NeedleId]struct{}, size int64, err error)
|
||||
sync func(*erasure_coding.EcjAppend) error
|
||||
}
|
||||
|
||||
var defaultEcjMergeIO = ecjMergeIO{
|
||||
read: erasure_coding.ReadEcjIds,
|
||||
sync: (*erasure_coding.EcjAppend).Sync,
|
||||
}
|
||||
|
||||
// MergeEcJournal folds a peer's deletion ids into the local journal of EC
|
||||
// volume vid on the receiving disk, the one whose data directory is dataDir;
|
||||
// ecjPath is that journal's path in the disk's index directory. It appends
|
||||
// only the ids the journal lacks and returns how many it added.
|
||||
//
|
||||
// A mounted volume owns its journal: the merge goes through its open handle
|
||||
// and in-memory set. That is the receiving disk's own runtime for vid,
|
||||
// wherever its journal lives (it may sit in the data dir rather than
|
||||
// ecjPath's index dir), else a sibling runtime journaling into ecjPath itself:
|
||||
// disks sharing one index directory, or reconciliation mounting vid on a disk
|
||||
// that journals into another's (#9212). Otherwise ecjPath is written while
|
||||
// every disk's EC lock is held, so no mount anywhere can open it mid-append,
|
||||
// and synced after they are released.
|
||||
func (s *Store) MergeEcJournal(vid needle.VolumeId, dataDir, ecjPath string, ids map[types.NeedleId]struct{}) (int, error) {
|
||||
return s.mergeEcJournal(vid, dataDir, ecjPath, ids, defaultEcjMergeIO)
|
||||
}
|
||||
|
||||
func (s *Store) mergeEcJournal(vid needle.VolumeId, dataDir, ecjPath string, ids map[types.NeedleId]struct{}, mio ecjMergeIO) (int, error) {
|
||||
unlock := lockEcjPath(ecjPath)
|
||||
defer unlock()
|
||||
|
||||
var owner *DiskLocation
|
||||
for _, loc := range s.Locations {
|
||||
if filepath.Clean(loc.Directory) == filepath.Clean(dataDir) {
|
||||
owner = loc
|
||||
break
|
||||
}
|
||||
}
|
||||
if owner == nil {
|
||||
return 0, fmt.Errorf("ec volume %d: no disk at %s owns journal %s", vid, dataDir, ecjPath)
|
||||
}
|
||||
for attempt := 0; attempt < ecjMergeAttempts; attempt++ {
|
||||
if added, merged, err := s.mergeIntoMountedEcJournal(owner, vid, ecjPath, ids); merged {
|
||||
return added, err
|
||||
}
|
||||
// Read outside the locks: a bloated journal can take a while, and a
|
||||
// queued mount would otherwise stall every EC read on its disk.
|
||||
local, size, err := mio.read(ecjPath)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("read %s: %w", ecjPath, err)
|
||||
}
|
||||
added, mounted, err := s.appendUnmountedEcJournal(owner, vid, ecjPath, local, ids, size, mio.sync)
|
||||
if mounted || errors.Is(err, erasure_coding.ErrEcjChanged) {
|
||||
continue
|
||||
}
|
||||
return added, err
|
||||
}
|
||||
return 0, fmt.Errorf("ec volume %d: journal %s kept changing during merge", vid, ecjPath)
|
||||
}
|
||||
|
||||
// rLockEcVolumes read-locks every disk's EC volume map in location order. A
|
||||
// mount registers its EcVolume, and reads its journal, under its own disk's
|
||||
// write lock, so holding all of them excludes a mount on any disk. No path
|
||||
// holds two disks' EC locks at once, so the fixed order cannot deadlock.
|
||||
// Holders must not wait on disk I/O beyond a page-cache write: a queued mount
|
||||
// on any disk stalls that disk's EC reads until they let go.
|
||||
func (s *Store) rLockEcVolumes() (unlock func()) {
|
||||
for _, loc := range s.Locations {
|
||||
loc.ecVolumesLock.RLock()
|
||||
}
|
||||
return func() {
|
||||
for _, loc := range s.Locations {
|
||||
loc.ecVolumesLock.RUnlock()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// mountedEcJournal returns the runtime a merge into ecjPath on owner must go
|
||||
// through and the disk it is mounted on, or nil when there is none: owner's
|
||||
// own runtime for vid, else the first sibling's whose journal is ecjPath.
|
||||
// Callers hold rLockEcVolumes.
|
||||
func (s *Store) mountedEcJournal(owner *DiskLocation, vid needle.VolumeId, ecjPath string) (*erasure_coding.EcVolume, *DiskLocation) {
|
||||
if ev, found := owner.ecVolumes[vid]; found {
|
||||
return ev, owner
|
||||
}
|
||||
for _, loc := range s.Locations {
|
||||
if ev, found := loc.ecVolumes[vid]; found && filepath.Clean(ev.FileName(".ecj")) == filepath.Clean(ecjPath) {
|
||||
return ev, loc
|
||||
}
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// mergeIntoMountedEcJournal merges ids through the runtime mountedEcJournal
|
||||
// picks and publishes them to the other runtimes sharing that journal.
|
||||
// merged reports whether there was one. Only the picked runtime's disk stays
|
||||
// locked across the merge's fsync, which keeps it mounted.
|
||||
func (s *Store) mergeIntoMountedEcJournal(owner *DiskLocation, vid needle.VolumeId, ecjPath string, ids map[types.NeedleId]struct{}) (added int, merged bool, err error) {
|
||||
unlock := s.rLockEcVolumes()
|
||||
ev, mountedOn := s.mountedEcJournal(owner, vid, ecjPath)
|
||||
if ev == nil {
|
||||
unlock()
|
||||
return 0, false, nil
|
||||
}
|
||||
for _, loc := range s.Locations {
|
||||
if loc != mountedOn {
|
||||
loc.ecVolumesLock.RUnlock()
|
||||
}
|
||||
}
|
||||
journalPath := ev.FileName(".ecj")
|
||||
added, err = ev.MergeJournal(ids)
|
||||
mountedOn.ecVolumesLock.RUnlock()
|
||||
if err != nil {
|
||||
return added, true, err
|
||||
}
|
||||
// Publish to the holders of the file the merge wrote to — the picked
|
||||
// runtime's journal may live outside ecjPath, and a holder of a different
|
||||
// file must not claim ids that file lacks.
|
||||
s.withEcJournalHolders(vid, journalPath, func(holders []*erasure_coding.EcVolume) {
|
||||
for _, h := range holders {
|
||||
if h != ev {
|
||||
h.PublishMergedIds(ids)
|
||||
}
|
||||
}
|
||||
})
|
||||
return added, true, nil
|
||||
}
|
||||
|
||||
// appendUnmountedEcJournal writes the missing ids under every disk's EC lock,
|
||||
// after confirming no runtime has mounted the journal since it was read, and
|
||||
// syncs them once the locks are released: a slow fsync must not hold off
|
||||
// mounts, and the EC reads queued behind them, on every disk. A mount after
|
||||
// the write reads the new records like any others. mounted reports that a
|
||||
// runtime appeared before the write; the caller then merges through it.
|
||||
func (s *Store) appendUnmountedEcJournal(owner *DiskLocation, vid needle.VolumeId, ecjPath string, local, ids map[types.NeedleId]struct{}, size int64, sync func(*erasure_coding.EcjAppend) error) (added int, mounted bool, err error) {
|
||||
pending, mounted, err := s.writeUnmountedEcJournal(owner, vid, ecjPath, local, ids, size)
|
||||
if pending == nil {
|
||||
return 0, mounted, err
|
||||
}
|
||||
defer pending.Close()
|
||||
if err := sync(pending); err != nil {
|
||||
var rolledBack bool
|
||||
s.withEcJournalHolders(vid, ecjPath, func(holders []*erasure_coding.EcVolume) {
|
||||
rolledBack = pending.RollbackUnsynced(holders)
|
||||
})
|
||||
if rolledBack {
|
||||
return 0, false, err
|
||||
}
|
||||
// Later records follow these, so they stay: make them durable,
|
||||
// still outside the locks.
|
||||
resyncErr := pending.Rewrite()
|
||||
if resyncErr == nil {
|
||||
resyncErr = sync(pending)
|
||||
}
|
||||
if resyncErr != nil {
|
||||
// Neither removable nor durable: no deleted set may claim them,
|
||||
// so a retried merge writes and syncs them again.
|
||||
s.withEcJournalHolders(vid, ecjPath, pending.Unpublish)
|
||||
glog.Errorf("ec volume %d: merged records in %s may not be durable: %v; resync: %v", vid, ecjPath, err, resyncErr)
|
||||
return 0, false, fmt.Errorf("%w; resync: %v", err, resyncErr)
|
||||
}
|
||||
}
|
||||
return pending.Added, false, nil
|
||||
}
|
||||
|
||||
func (s *Store) writeUnmountedEcJournal(owner *DiskLocation, vid needle.VolumeId, ecjPath string, local, ids map[types.NeedleId]struct{}, size int64) (pending *erasure_coding.EcjAppend, mounted bool, err error) {
|
||||
unlock := s.rLockEcVolumes()
|
||||
defer unlock()
|
||||
if ev, _ := s.mountedEcJournal(owner, vid, ecjPath); ev != nil {
|
||||
return nil, true, nil
|
||||
}
|
||||
pending, err = erasure_coding.WriteEcjIds(ecjPath, local, ids, size)
|
||||
return pending, false, err
|
||||
}
|
||||
|
||||
// withEcJournalHolders runs fn with every runtime that has mounted the journal
|
||||
// at ecjPath, for undoing an append whose fsync failed (see
|
||||
// EcjAppend.RollbackUnsynced and Unpublish). The disk locks stop new mounts
|
||||
// and the caller's path lock stops other unmounted merges; fn runs no fsync,
|
||||
// so they are held only for in-memory work and a truncate.
|
||||
func (s *Store) withEcJournalHolders(vid needle.VolumeId, ecjPath string, fn func([]*erasure_coding.EcVolume)) {
|
||||
unlock := s.rLockEcVolumes()
|
||||
defer unlock()
|
||||
var holders []*erasure_coding.EcVolume
|
||||
for _, loc := range s.Locations {
|
||||
if ev, found := loc.ecVolumes[vid]; found && filepath.Clean(ev.FileName(".ecj")) == filepath.Clean(ecjPath) {
|
||||
holders = append(holders, ev)
|
||||
}
|
||||
}
|
||||
fn(holders)
|
||||
}
|
||||
@@ -0,0 +1,419 @@
|
||||
package storage
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/stats"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/volume_info"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func ecjIdSet(ids ...types.NeedleId) map[types.NeedleId]struct{} {
|
||||
s := make(map[types.NeedleId]struct{}, len(ids))
|
||||
for _, id := range ids {
|
||||
s[id] = struct{}{}
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
func ecjRecords(ids ...types.NeedleId) []byte {
|
||||
b := make([]byte, len(ids)*types.NeedleIdSize)
|
||||
for i, id := range ids {
|
||||
types.NeedleIdToBytes(b[i*types.NeedleIdSize:], id)
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
func startEcJournalStore(t *testing.T, dataDir, idxDir string) *Store {
|
||||
t.Helper()
|
||||
return startEcJournalStoreDisks(t, idxDir, dataDir)
|
||||
}
|
||||
|
||||
// startEcJournalStoreDisks starts a store with one disk per data dir. An
|
||||
// empty idxDir gives each disk its data dir as index dir; otherwise all disks
|
||||
// share idxDir.
|
||||
func startEcJournalStoreDisks(t *testing.T, idxDir string, dataDirs ...string) *Store {
|
||||
t.Helper()
|
||||
dirs := append([]string{idxDir}, dataDirs...)
|
||||
if idxDir == "" {
|
||||
dirs = dataDirs
|
||||
}
|
||||
for _, d := range dirs {
|
||||
require.NoError(t, os.MkdirAll(d, 0o755))
|
||||
}
|
||||
n := len(dataDirs)
|
||||
maxCounts := make([]int32, n)
|
||||
minFree := make([]util.MinFreeSpace, n)
|
||||
diskTypes := make([]types.DiskType, n)
|
||||
for i := range dataDirs {
|
||||
maxCounts[i] = 100
|
||||
diskTypes[i] = types.HardDriveType
|
||||
}
|
||||
store := NewStore(nil, "localhost", 8080, 18080, "http://localhost:8080", "store-id",
|
||||
dataDirs,
|
||||
maxCounts,
|
||||
minFree,
|
||||
idxDir,
|
||||
NeedleMapInMemory,
|
||||
diskTypes,
|
||||
nil,
|
||||
3,
|
||||
stats.DefaultDiskIOProbeConfig(),
|
||||
)
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case <-store.NewEcShardsChan:
|
||||
case <-store.NewVolumesChan:
|
||||
case <-store.DeletedVolumesChan:
|
||||
case <-store.DeletedEcShardsChan:
|
||||
case <-store.StateUpdateChan:
|
||||
case <-done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
t.Cleanup(func() {
|
||||
store.Close()
|
||||
close(done)
|
||||
})
|
||||
return store
|
||||
}
|
||||
|
||||
// A mounted volume whose index was found in its data directory keeps its live
|
||||
// journal there, while a shard copy targets the index directory. The merge
|
||||
// must reach the live journal and the in-memory set, not a sibling file.
|
||||
func TestMergeEcJournal_MountedVolumeJournalInDataDir(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
dataDir := filepath.Join(tempDir, "data")
|
||||
idxDir := filepath.Join(tempDir, "idx")
|
||||
store := startEcJournalStore(t, dataDir, idxDir)
|
||||
|
||||
const collection = "c"
|
||||
vid := needle.VolumeId(9)
|
||||
writeEcShard0(t, dataDir, collection, vid)
|
||||
dataBase := writeEcIndex(t, dataDir, collection, vid, 1)
|
||||
require.NoError(t, store.MountEcShards(collection, vid, 0, ""))
|
||||
ev, found := store.Locations[0].FindEcVolume(vid)
|
||||
require.True(t, found)
|
||||
require.Equal(t, dataBase+".ecj", ev.FileName(".ecj"))
|
||||
|
||||
idxJournal := erasure_coding.EcShardFileName(collection, idxDir, int(vid)) + ".ecj"
|
||||
added, err := store.MergeEcJournal(vid, dataDir, idxJournal, ecjIdSet(1, 2))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, added)
|
||||
assert.True(t, ev.IsNeedleDeleted(2))
|
||||
assert.Equal(t, ecjRecords(1, 2), mustReadFile(t, dataBase+".ecj"))
|
||||
assert.NoFileExists(t, idxJournal)
|
||||
}
|
||||
|
||||
// An unmounted journal is merged on disk, and merging the same peer journal
|
||||
// again adds nothing.
|
||||
func TestMergeEcJournal_UnmountedVolumeIsIdempotent(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
dataDir := filepath.Join(tempDir, "data")
|
||||
idxDir := filepath.Join(tempDir, "idx")
|
||||
store := startEcJournalStore(t, dataDir, idxDir)
|
||||
|
||||
vid := needle.VolumeId(9)
|
||||
journal := erasure_coding.EcShardFileName("", idxDir, int(vid)) + ".ecj"
|
||||
require.NoError(t, os.WriteFile(journal, ecjRecords(1, 2), 0o644))
|
||||
|
||||
for i, want := range []int{2, 0, 0} {
|
||||
added, err := store.MergeEcJournal(vid, dataDir, journal, ecjIdSet(2, 3, 4))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, want, added, "merge %d", i)
|
||||
}
|
||||
assert.Equal(t, ecjRecords(1, 2, 3, 4), mustReadFile(t, journal))
|
||||
|
||||
_, err := store.MergeEcJournal(vid, filepath.Join(tempDir, "elsewhere"), journal, ecjIdSet(1))
|
||||
assert.Error(t, err, "a data dir that is no disk owns no journal")
|
||||
}
|
||||
|
||||
// Disks sharing one index directory all hold the same journal path. Copying
|
||||
// shards onto a disk that has not mounted vid must still reach the sibling
|
||||
// runtime that holds that journal open, or the sibling keeps serving needles
|
||||
// the peer deleted until it remounts.
|
||||
func TestMergeEcJournal_SharedIndexDirReachesSiblingMount(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
disk0 := filepath.Join(tempDir, "d0")
|
||||
disk1 := filepath.Join(tempDir, "d1")
|
||||
idxDir := filepath.Join(tempDir, "idx")
|
||||
store := startEcJournalStoreDisks(t, idxDir, disk0, disk1)
|
||||
|
||||
const collection = "c"
|
||||
vid := needle.VolumeId(9)
|
||||
writeEcShard0(t, disk0, collection, vid)
|
||||
idxBase := writeEcIndex(t, idxDir, collection, vid, 1)
|
||||
ev, err := store.Locations[0].LoadEcShard(collection, vid, 0)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, idxBase+".ecj", ev.FileName(".ecj"))
|
||||
|
||||
// The copy lands on disk1, whose index dir is the shared one.
|
||||
added, err := store.MergeEcJournal(vid, disk1, idxBase+".ecj", ecjIdSet(1, 2))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, added)
|
||||
assert.True(t, ev.IsNeedleDeleted(2), "the mounted sibling must see the merged id")
|
||||
assert.Equal(t, ecjRecords(1, 2), mustReadFile(t, idxBase+".ecj"))
|
||||
}
|
||||
|
||||
// Every runtime holding the journal open must see merged ids in memory.
|
||||
func TestMergeEcJournal_SharedJournalReachesEveryHolder(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
disk0 := filepath.Join(tempDir, "d0")
|
||||
disk1 := filepath.Join(tempDir, "d1")
|
||||
idxDir := filepath.Join(tempDir, "idx")
|
||||
store := startEcJournalStoreDisks(t, idxDir, disk0, disk1)
|
||||
|
||||
const collection = "c"
|
||||
vid := needle.VolumeId(9)
|
||||
writeEcShard0(t, disk0, collection, vid)
|
||||
writeEcShard0(t, disk1, collection, vid)
|
||||
idxBase := writeEcIndex(t, idxDir, collection, vid, 1)
|
||||
ev0, err := store.Locations[0].LoadEcShard(collection, vid, 0)
|
||||
require.NoError(t, err)
|
||||
ev1, err := store.Locations[1].LoadEcShard(collection, vid, 0)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, idxBase+".ecj", ev0.FileName(".ecj"))
|
||||
require.Equal(t, idxBase+".ecj", ev1.FileName(".ecj"))
|
||||
|
||||
added, err := store.MergeEcJournal(vid, disk0, idxBase+".ecj", ecjIdSet(1, 2))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, added)
|
||||
assert.True(t, ev0.IsNeedleDeleted(2))
|
||||
assert.True(t, ev1.IsNeedleDeleted(2), "every journal holder must see the merged id")
|
||||
}
|
||||
|
||||
// The picked runtime may journal to a different file than the copied one —
|
||||
// its index lives in its data directory while a sibling's lives in the index
|
||||
// directory. The ids must be published only to holders of the file they were
|
||||
// written to: the sibling's deleted set must not claim records its journal
|
||||
// lacks, or they come back after its remount.
|
||||
func TestMergeEcJournal_PublishesToActualJournalHolders(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
disk0 := filepath.Join(tempDir, "d0")
|
||||
disk1 := filepath.Join(tempDir, "d1")
|
||||
idxDir := filepath.Join(tempDir, "idx")
|
||||
store := startEcJournalStoreDisks(t, idxDir, disk0, disk1)
|
||||
|
||||
const collection = "c"
|
||||
vid := needle.VolumeId(9)
|
||||
writeEcShard0(t, disk0, collection, vid)
|
||||
writeEcShard0(t, disk1, collection, vid)
|
||||
dataBase := writeEcIndex(t, disk0, collection, vid, 1)
|
||||
idxBase := writeEcIndex(t, idxDir, collection, vid, 1)
|
||||
|
||||
ev0, err := store.Locations[0].LoadEcShard(collection, vid, 0)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, dataBase+".ecj", ev0.FileName(".ecj"))
|
||||
ev1, err := store.Locations[1].LoadEcShard(collection, vid, 0)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, idxBase+".ecj", ev1.FileName(".ecj"))
|
||||
|
||||
// The copy targets the index-dir journal, but the receiving disk's runtime
|
||||
// owns the data-dir one and the merge goes through it.
|
||||
added, err := store.MergeEcJournal(vid, disk0, idxBase+".ecj", ecjIdSet(1, 2))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, added)
|
||||
assert.Equal(t, ecjRecords(1, 2), mustReadFile(t, dataBase+".ecj"))
|
||||
assert.Equal(t, ecjRecords(1), mustReadFile(t, idxBase+".ecj"))
|
||||
assert.True(t, ev0.IsNeedleDeleted(2))
|
||||
assert.False(t, ev1.IsNeedleDeleted(2), "a different journal's holder must not claim the merged id")
|
||||
}
|
||||
|
||||
// A sibling disk can mount vid from the receiving disk's index (#9212) while
|
||||
// the merge reads the journal unlocked. The append must notice that mount and
|
||||
// go through it instead of writing behind its open handle.
|
||||
func TestMergeEcJournal_SiblingMountDuringReadIsMergedThrough(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
disk0 := filepath.Join(tempDir, "d0")
|
||||
disk1 := filepath.Join(tempDir, "d1")
|
||||
store := startEcJournalStoreDisks(t, "", disk0, disk1)
|
||||
|
||||
const collection = "c"
|
||||
vid := needle.VolumeId(9)
|
||||
ownerBase := writeEcIndex(t, disk0, collection, vid, 1)
|
||||
writeEcShard0(t, disk1, collection, vid)
|
||||
|
||||
var sibling *erasure_coding.EcVolume
|
||||
read := func(path string) (map[types.NeedleId]struct{}, int64, error) {
|
||||
ids, size, err := erasure_coding.ReadEcjIds(path)
|
||||
if sibling == nil {
|
||||
var mountErr error
|
||||
sibling, mountErr = store.Locations[1].loadEcShardWithIdxDir(collection, vid, 0, disk0)
|
||||
require.NoError(t, mountErr)
|
||||
require.Equal(t, ownerBase+".ecj", sibling.FileName(".ecj"))
|
||||
}
|
||||
return ids, size, err
|
||||
}
|
||||
mio := defaultEcjMergeIO
|
||||
mio.read = read
|
||||
added, err := store.mergeEcJournal(vid, disk0, ownerBase+".ecj", ecjIdSet(1, 2), mio)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, added)
|
||||
require.NotNil(t, sibling)
|
||||
assert.True(t, sibling.IsNeedleDeleted(2), "the sibling mounted mid-merge must see the merged id")
|
||||
assert.Equal(t, ecjRecords(1, 2), mustReadFile(t, ownerBase+".ecj"))
|
||||
}
|
||||
|
||||
// The unmounted append's fsync runs with no disk's EC lock held, so a slow
|
||||
// sync cannot hold off mounts, or the EC reads queued behind them, on every
|
||||
// disk.
|
||||
func TestMergeEcJournal_UnmountedSyncHoldsNoDiskLock(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
disk0 := filepath.Join(tempDir, "d0")
|
||||
disk1 := filepath.Join(tempDir, "d1")
|
||||
store := startEcJournalStoreDisks(t, "", disk0, disk1)
|
||||
|
||||
vid := needle.VolumeId(9)
|
||||
journal := erasure_coding.EcShardFileName("", disk0, int(vid)) + ".ecj"
|
||||
require.NoError(t, os.WriteFile(journal, ecjRecords(1), 0o644))
|
||||
|
||||
synced := false
|
||||
mio := defaultEcjMergeIO
|
||||
mio.sync = func(a *erasure_coding.EcjAppend) error {
|
||||
for i, loc := range store.Locations {
|
||||
if assert.True(t, loc.ecVolumesLock.TryLock(), "disk %d EC lock held across fsync", i) {
|
||||
loc.ecVolumesLock.Unlock()
|
||||
}
|
||||
}
|
||||
synced = true
|
||||
return a.Sync()
|
||||
}
|
||||
added, err := store.mergeEcJournal(vid, disk0, journal, ecjIdSet(1, 2), mio)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, added)
|
||||
assert.True(t, synced)
|
||||
assert.Equal(t, ecjRecords(1, 2), mustReadFile(t, journal))
|
||||
}
|
||||
|
||||
// A failed fsync removes the append when its records are still the journal's
|
||||
// tail, also through every runtime that mounted the journal since the write:
|
||||
// their deleted sets must not keep ids the journal may never have persisted,
|
||||
// and a retried merge appends them again. Once any runtime has journaled
|
||||
// after them, removing them would lose that delete, even if the runtime doing
|
||||
// the rollback cached the older length; they stay and are rewritten and
|
||||
// synced instead. If that sync fails too, the records stay but leave every
|
||||
// deleted set, so a retried merge appends and syncs them again rather than
|
||||
// trusting records never shown durable.
|
||||
func TestMergeEcJournal_FailedSync(t *testing.T) {
|
||||
for _, tt := range []struct {
|
||||
name string
|
||||
mounts int // runtimes on disks 1.. mounting the journal mid-sync
|
||||
journalVia int // 1-based runtime that journals id 3 mid-sync, 0 for none
|
||||
resyncFails bool // the fsync after the rewrite fails too
|
||||
wantJournal []types.NeedleId
|
||||
}{
|
||||
{name: "unmounted", wantJournal: []types.NeedleId{1}},
|
||||
{name: "mounted since the write", mounts: 1, wantJournal: []types.NeedleId{1}},
|
||||
{name: "two mounted since the write", mounts: 2, wantJournal: []types.NeedleId{1}},
|
||||
{name: "mounted and journaled since", mounts: 1, journalVia: 1, wantJournal: []types.NeedleId{1, 2, 3}},
|
||||
{name: "another runtime journaled since", mounts: 2, journalVia: 2, wantJournal: []types.NeedleId{1, 2, 3}},
|
||||
{name: "resync fails", mounts: 2, journalVia: 2, resyncFails: true, wantJournal: []types.NeedleId{1, 2, 3}},
|
||||
} {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
disks := []string{filepath.Join(tempDir, "d0"), filepath.Join(tempDir, "d1"), filepath.Join(tempDir, "d2")}
|
||||
store := startEcJournalStoreDisks(t, "", disks...)
|
||||
|
||||
const collection = "c"
|
||||
vid := needle.VolumeId(9)
|
||||
journal := writeEcIndex(t, disks[0], collection, vid, 1) + ".ecj"
|
||||
for _, d := range disks[1:] {
|
||||
writeEcShard0(t, d, collection, vid)
|
||||
}
|
||||
|
||||
var holders []*erasure_coding.EcVolume
|
||||
syncs := 0
|
||||
mio := defaultEcjMergeIO
|
||||
mio.sync = func(a *erasure_coding.EcjAppend) error {
|
||||
if syncs++; syncs > 1 {
|
||||
if tt.resyncFails {
|
||||
return errors.New("injected resync failure")
|
||||
}
|
||||
return a.Sync()
|
||||
}
|
||||
for i := 1; i <= tt.mounts; i++ {
|
||||
ev, err := store.Locations[i].loadEcShardWithIdxDir(collection, vid, 0, disks[0])
|
||||
require.NoError(t, err)
|
||||
require.True(t, ev.IsNeedleDeleted(2), "the mount loads the unsynced record")
|
||||
holders = append(holders, ev)
|
||||
}
|
||||
if tt.journalVia > 0 {
|
||||
_, err := holders[tt.journalVia-1].MergeJournal(ecjIdSet(3))
|
||||
require.NoError(t, err)
|
||||
}
|
||||
return errors.New("injected fsync failure")
|
||||
}
|
||||
added, err := store.mergeEcJournal(vid, disks[0], journal, ecjIdSet(1, 2), mio)
|
||||
assert.Equal(t, ecjRecords(tt.wantJournal...), mustReadFile(t, journal))
|
||||
if tt.journalVia > 0 {
|
||||
assert.True(t, holders[tt.journalVia-1].IsNeedleDeleted(3), "a later delete is never lost")
|
||||
}
|
||||
if tt.journalVia > 0 && !tt.resyncFails {
|
||||
require.NoError(t, err, "records followed by a delete are resynced")
|
||||
assert.Equal(t, 1, added)
|
||||
for _, ev := range holders {
|
||||
assert.True(t, ev.IsNeedleDeleted(2))
|
||||
}
|
||||
return
|
||||
}
|
||||
require.Error(t, err)
|
||||
for _, ev := range holders {
|
||||
assert.False(t, ev.IsNeedleDeleted(2), "no deleted set claims a record not shown durable")
|
||||
}
|
||||
|
||||
// A retried merge appends and syncs the id again.
|
||||
added, err = store.MergeEcJournal(vid, disks[0], journal, ecjIdSet(1, 2))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, added)
|
||||
assert.Equal(t, ecjRecords(append(tt.wantJournal, 2)...), mustReadFile(t, journal))
|
||||
if len(holders) > 0 {
|
||||
assert.True(t, holders[0].IsNeedleDeleted(2))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// writeEcShard0 writes shard 0 of vid and its .vif into dataDir.
|
||||
func writeEcShard0(t *testing.T, dataDir, collection string, vid needle.VolumeId) {
|
||||
t.Helper()
|
||||
const datSize int64 = 1024 * 1024
|
||||
dataBase := erasure_coding.EcShardFileName(collection, dataDir, int(vid))
|
||||
f, err := os.Create(dataBase + erasure_coding.ToExt(0))
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, f.Truncate(calculateExpectedShardSize(datSize, 10)))
|
||||
require.NoError(t, f.Close())
|
||||
require.NoError(t, volume_info.SaveVolumeInfo(dataBase+".vif", &volume_server_pb.VolumeInfo{
|
||||
Version: uint32(needle.Version3),
|
||||
DatFileSize: datSize,
|
||||
EcShardConfig: &volume_server_pb.EcShardConfig{DataShards: 10, ParityShards: 4},
|
||||
}))
|
||||
}
|
||||
|
||||
// writeEcIndex writes vid's .ecx and a .ecj holding deleted into dir and
|
||||
// returns their base name.
|
||||
func writeEcIndex(t *testing.T, dir, collection string, vid needle.VolumeId, deleted ...types.NeedleId) string {
|
||||
t.Helper()
|
||||
base := erasure_coding.EcShardFileName(collection, dir, int(vid))
|
||||
require.NoError(t, os.WriteFile(base+".ecx", make([]byte, types.NeedleMapEntrySize), 0o644))
|
||||
require.NoError(t, os.WriteFile(base+".ecj", ecjRecords(deleted...), 0o644))
|
||||
return base
|
||||
}
|
||||
|
||||
func mustReadFile(t *testing.T, path string) []byte {
|
||||
t.Helper()
|
||||
b, err := os.ReadFile(path)
|
||||
require.NoError(t, err)
|
||||
return b
|
||||
}
|
||||
Reference in new issue
Block a user