diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index b525b5e5b..74b269d89 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -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, +) -> std::io::Result<(std::collections::HashSet, 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, + vid: VolumeId, + data_dir: String, + ecj_path: String, + ids: std::collections::HashSet, +) -> std::io::Result { + 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> { diff --git a/seaweed-volume/src/server/store_ec.rs b/seaweed-volume/src/server/store_ec.rs index 6d8b3ca12..d54c8063f 100644 --- a/seaweed-volume/src/server/store_ec.rs +++ b/seaweed-volume/src/server/store_ec.rs @@ -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, 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 diff --git a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs index de26e0ba9..2544887a7 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs @@ -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) -> io::Result { + 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) { + 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 = [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::>() + ); + 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. diff --git a/seaweed-volume/src/storage/erasure_coding/ecj_merge.rs b/seaweed-volume/src/storage/erasure_coding/ecj_merge.rs new file mode 100644 index 000000000..d9a316847 --- /dev/null +++ b/seaweed-volume/src/storage/erasure_coding/ecj_merge.rs @@ -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 (`.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, + 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 { + 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, 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, + has: impl Fn(&NeedleId) -> bool, +) -> Vec { + let mut delta: Vec = incoming.iter().copied().filter(|id| !has(id)).collect(); + delta.sort_unstable(); + delta +} + +pub(crate) fn encode_ecj_ids(ids: &[NeedleId]) -> Vec { + 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, + incoming: &HashSet, + size: u64, +) -> io::Result> { + 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 { + v.iter().map(|&id| NeedleId(id)).collect() + } + + fn bytes(v: &[u64]) -> Vec { + encode_ecj_ids(&v.iter().map(|&id| NeedleId(id)).collect::>()) + } + + fn records(path: &str) -> Vec { + 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) -> 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 = (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]); + } +} diff --git a/seaweed-volume/src/storage/erasure_coding/mod.rs b/seaweed-volume/src/storage/erasure_coding/mod.rs index c6cf4d8f3..678b83961 100644 --- a/seaweed-volume/src/storage/erasure_coding/mod.rs +++ b/seaweed-volume/src/storage/erasure_coding/mod.rs @@ -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, diff --git a/seaweed-volume/src/storage/mod.rs b/seaweed-volume/src/storage/mod.rs index 38220cc62..cd7bd1df0 100644 --- a/seaweed-volume/src/storage/mod.rs +++ b/seaweed-volume/src/storage/mod.rs @@ -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; diff --git a/seaweed-volume/src/storage/store_ec_journal.rs b/seaweed-volume/src/storage/store_ec_journal.rs new file mode 100644 index 000000000..d31cd135f --- /dev/null +++ b/seaweed-volume/src/storage/store_ec_journal.rs @@ -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, + vid: VolumeId, + data_dir: &str, + ecj_path: &str, + ids: &HashSet, +) -> io::Result { + 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, + vid: VolumeId, + data_dir: &str, + ecj_path: &str, + ids: &HashSet, + mut read: impl FnMut(&str) -> io::Result<(HashSet, u64)>, +) -> io::Result { + 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> { + 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, +) { + 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 { + 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 { + let ids: Vec = ids.iter().copied().map(NeedleId).collect(); + crate::storage::erasure_coding::ecj_merge::encode_ecj_ids(&ids) + } + + fn id_set(ids: &[u64]) -> HashSet { + 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, 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() + ); + } +} diff --git a/weed/server/volume_grpc_erasure_coding.go b/weed/server/volume_grpc_erasure_coding.go index 011b18ecd..aef7162a9 100644 --- a/weed/server/volume_grpc_erasure_coding.go +++ b/weed/server/volume_grpc_erasure_coding.go @@ -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 diff --git a/weed/server/volume_grpc_erasure_coding_ecj_test.go b/weed/server/volume_grpc_erasure_coding_ecj_test.go new file mode 100644 index 000000000..a4a61508f --- /dev/null +++ b/weed/server/volume_grpc_erasure_coding_ecj_test.go @@ -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) +} diff --git a/weed/server/volume_grpc_erasure_coding_recover.go b/weed/server/volume_grpc_erasure_coding_recover.go index c05def22f..3e28b4ae8 100644 --- a/weed/server/volume_grpc_erasure_coding_recover.go +++ b/weed/server/volume_grpc_erasure_coding_recover.go @@ -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) diff --git a/weed/server/volume_grpc_erasure_coding_recover_test.go b/weed/server/volume_grpc_erasure_coding_recover_test.go index fc0d0d1ab..36689cd74 100644 --- a/weed/server/volume_grpc_erasure_coding_recover_test.go +++ b/weed/server/volume_grpc_erasure_coding_recover_test.go @@ -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) diff --git a/weed/storage/erasure_coding/ec_volume_delete.go b/weed/storage/erasure_coding/ec_volume_delete.go index 1bf4b9c2a..44929935f 100644 --- a/weed/storage/erasure_coding/ec_volume_delete.go +++ b/weed/storage/erasure_coding/ec_volume_delete.go @@ -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 } diff --git a/weed/storage/erasure_coding/ecj_merge.go b/weed/storage/erasure_coding/ecj_merge.go new file mode 100644 index 000000000..0cc11517a --- /dev/null +++ b/weed/storage/erasure_coding/ecj_merge.go @@ -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 (.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() +} diff --git a/weed/storage/erasure_coding/ecj_merge_test.go b/weed/storage/erasure_coding/ecj_merge_test.go new file mode 100644 index 000000000..d845762c0 --- /dev/null +++ b/weed/storage/erasure_coding/ecj_merge_test.go @@ -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") +} diff --git a/weed/storage/store_ec_journal.go b/weed/storage/store_ec_journal.go new file mode 100644 index 000000000..c7e579252 --- /dev/null +++ b/weed/storage/store_ec_journal.go @@ -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) +} diff --git a/weed/storage/store_ec_journal_test.go b/weed/storage/store_ec_journal_test.go new file mode 100644 index 000000000..a40624099 --- /dev/null +++ b/weed/storage/store_ec_journal_test.go @@ -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 +}