diff --git a/seaweed-volume/src/server/write_queue.rs b/seaweed-volume/src/server/write_queue.rs index 528007325..e27ed486a 100644 --- a/seaweed-volume/src/server/write_queue.rs +++ b/seaweed-volume/src/server/write_queue.rs @@ -3,8 +3,8 @@ //! Instead of each upload handler directly calling `write_needle`, writes are //! submitted to a queue. A background worker drains the queue in batches (up to //! 128 entries), groups them by volume ID, and processes them together under a -//! single store lock. Requests that asked for `fsync` are flushed by -//! `write_needle` itself, one flush per durable write. +//! single store lock. Durable writes to a volume share their .dat and .idx +//! flushes (see `Volume::write_needles_grouped`). use std::sync::Arc; @@ -159,8 +159,12 @@ fn process_batch(state: Arc, batch: Vec) { let mut store = state.store.write().unwrap(); for (vid, entries) in groups { - for (mut needle, fsync, response_tx) in entries { - let result = store.write_volume_needle(vid, &mut needle, fsync); + let (mut writes, senders): (Vec<_>, Vec<_>) = entries + .into_iter() + .map(|(needle, fsync, response_tx)| ((needle, fsync), response_tx)) + .unzip(); + let results = store.write_volume_needles(vid, &mut writes); + for (response_tx, result) in senders.into_iter().zip(results) { // Send result back; ignore error if receiver dropped. let _ = response_tx.send(result); } @@ -312,6 +316,63 @@ mod tests { } } + /// The queue hands a volume's batch to the grouped path, so ten durable + /// writes cost one .dat sync and one .idx sync, not ten of each. + #[test] + fn test_process_batch_group_commits_fsync_writes() { + use crate::config::MinFreeSpace; + use crate::storage::types::{DiskType, NeedleId}; + use crate::storage::volume::VolumeSpec; + + let tmp = tempfile::TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let state = make_test_state(); + { + let mut store = state.store.write().unwrap(); + store + .add_location( + dir, + dir, + 10, + DiskType::HardDrive, + MinFreeSpace::Percent(1.0), + Vec::new(), + ) + .unwrap(); + store + .add_volume(VolumeId(1), DiskType::HardDrive, &VolumeSpec::default()) + .unwrap(); + } + + let mut receivers = Vec::new(); + let batch = (1..=10u64) + .map(|id| { + let (response_tx, response_rx) = oneshot::channel(); + receivers.push(response_rx); + WriteRequest { + volume_id: VolumeId(1), + needle: Needle { + id: NeedleId(id), + cookie: 0x1111.into(), + data: vec![id as u8; 8], + data_size: 8, + ..Needle::default() + }, + fsync: true, + response_tx, + } + }) + .collect(); + process_batch(state.clone(), batch); + + for mut rx in receivers { + assert!(matches!(rx.try_recv().unwrap(), Ok((_, _, false)))); + } + let store = state.store.read().unwrap(); + let (_, vol) = store.find_volume(VolumeId(1)).unwrap(); + assert_eq!(vol.sync_counts_for_test(), (1, 1)); + } + #[tokio::test] async fn test_write_queue_dropped_sender() { // When the queue is dropped, subsequent submits should fail gracefully. diff --git a/seaweed-volume/src/storage/io_error.rs b/seaweed-volume/src/storage/io_error.rs index d1cf03c88..a07341930 100644 --- a/seaweed-volume/src/storage/io_error.rs +++ b/seaweed-volume/src/storage/io_error.rs @@ -3,7 +3,7 @@ use std::io; use std::sync::Mutex; -use std::sync::atomic::{AtomicBool, AtomicI32, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; /// Consecutive storage-media errors allowed before the volume is quarantined. pub(crate) const IO_ERROR_TOLERANCE: i32 = 3; @@ -34,10 +34,30 @@ pub(crate) fn is_storage_io_error(e: &io::Error) -> bool { #[derive(Default)] pub(crate) struct IoErrorTracker { last: Mutex>, - count: AtomicI32, + /// The consecutive error count in the low 32 bits and, in the high 32, + /// how many times it has been cleared. They share one word so that + /// `record_success_at` updates both in one step: reads record their + /// outcomes here without the volume's write lock. + streak: AtomicU64, quarantined: AtomicBool, } +const STREAK_COUNT_BITS: u64 = 0xffff_ffff; + +fn streak_count(streak: u64) -> i32 { + (streak & STREAK_COUNT_BITS) as i32 +} + +/// `streak` with its count cleared and one more clear on record. +fn streak_cleared(streak: u64) -> u64 { + (streak >> 32).wrapping_add(1) << 32 +} + +/// A point in the error streak, taken where a write landed whose success +/// is only recorded later. See `IoErrorTracker::record_success_at`. +#[derive(Clone, Copy)] +pub(crate) struct StreakMark(u64); + impl IoErrorTracker { /// `Some(e)` records a failure, `None` a success. Only storage-media /// failures count; every other outcome clears the count and last error. @@ -45,14 +65,22 @@ impl IoErrorTracker { if let Some(e) = err && is_storage_io_error(e) { - self.count.fetch_add(1, Ordering::Relaxed); + self.streak.fetch_add(1, Ordering::Relaxed); if let Ok(mut guard) = self.last.lock() { *guard = Some(e.to_string()); } crate::metrics::STORAGE_IO_ERROR_COUNTER.inc(); return; } - self.count.store(0, Ordering::Relaxed); + self.clear_count(); + self.clear_last(); + } + + fn clear_count(&self) { + self.update_streak(|streak| Some(streak_cleared(streak))); + } + + fn clear_last(&self) { if let Ok(mut guard) = self.last.lock() && guard.is_some() { @@ -60,17 +88,58 @@ impl IoErrorTracker { } } + /// Apply `f` to the streak atomically; `None` leaves it as it is. + /// Returns the streak `f` produced, if any. + fn update_streak(&self, mut f: impl FnMut(u64) -> Option) -> Option { + let mut updated = None; + let _ = self + .streak + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |streak| { + updated = f(streak); + updated + }); + updated + } + + pub(crate) fn mark(&self) -> StreakMark { + StreakMark(self.streak.load(Ordering::Relaxed)) + } + + /// Record a success as if it had come at `mark`: the errors counted + /// before the mark are cleared and the ones counted since still stand, + /// as they would had each outcome been recorded in order. A streak + /// cleared since the mark is left as it is. + pub(crate) fn record_success_at(&self, mark: StreakMark) { + let before = streak_count(mark.0); + let updated = self.update_streak(|streak| { + if streak >> 32 != mark.0 >> 32 { + return None; + } + if streak_count(streak) <= before { + return Some(streak_cleared(streak)); + } + Some(streak - before as u64) + }); + // The last error stays when one counted since the mark is left. + if updated.is_some_and(|streak| streak_count(streak) == 0) { + self.clear_last(); + } + } + /// The last recorded error, the consecutive count, and the quarantine flag. pub(crate) fn get_io_error_state(&self) -> (Option, i32, bool) { let err = self.last.lock().ok().and_then(|g| g.clone()); - let count = self.count.load(Ordering::Relaxed); + let count = self.count(); let quarantined = self.quarantined.load(Ordering::Relaxed); (err, count, quarantined) } + fn count(&self) -> i32 { + streak_count(self.streak.load(Ordering::Relaxed)) + } + pub(crate) fn should_quarantine(&self) -> bool { - self.quarantined.load(Ordering::Relaxed) - || self.count.load(Ordering::Relaxed) >= IO_ERROR_TOLERANCE + self.quarantined.load(Ordering::Relaxed) || self.count() >= IO_ERROR_TOLERANCE } pub(crate) fn mark_io_quarantined(&self) { @@ -78,7 +147,7 @@ impl IoErrorTracker { } pub(crate) fn reset_io_error_state(&self) { - self.count.store(0, Ordering::Relaxed); + self.clear_count(); self.quarantined.store(false, Ordering::Relaxed); if let Ok(mut guard) = self.last.lock() { *guard = None; @@ -91,9 +160,11 @@ impl IoErrorTracker { *guard = err.map(|value| value.to_string()); } if err.is_some() { - self.count.store(IO_ERROR_TOLERANCE, Ordering::Relaxed); + self.update_streak(|streak| { + Some((streak & !STREAK_COUNT_BITS) | IO_ERROR_TOLERANCE as u64) + }); } else { - self.count.store(0, Ordering::Relaxed); + self.clear_count(); } } } @@ -171,6 +242,87 @@ mod tests { assert!(tracker.should_quarantine()); } + #[test] + fn success_at_a_mark_keeps_only_the_errors_after_it() { + let tracker = IoErrorTracker::default(); + tracker.check_read_write_error(Some(&media_error())); + tracker.check_read_write_error(Some(&media_error())); + let mark = tracker.mark(); + tracker.check_read_write_error(Some(&media_error())); + tracker.record_success_at(mark); + + assert_eq!( + tracker.get_io_error_state(), + (Some(media_error().to_string()), 1, false) + ); + } + + #[test] + fn success_at_a_mark_with_nothing_after_it_clears_the_streak() { + let tracker = IoErrorTracker::default(); + tracker.check_read_write_error(Some(&media_error())); + let mark = tracker.mark(); + tracker.record_success_at(mark); + + assert_eq!(tracker.get_io_error_state(), (None, 0, false)); + } + + #[test] + fn success_at_a_mark_leaves_a_streak_cleared_since() { + let tracker = IoErrorTracker::default(); + tracker.check_read_write_error(Some(&media_error())); + tracker.check_read_write_error(Some(&media_error())); + let mark = tracker.mark(); + tracker.check_read_write_error(None); + tracker.check_read_write_error(Some(&media_error())); + tracker.record_success_at(mark); + + assert_eq!( + tracker.get_io_error_state(), + (Some(media_error().to_string()), 1, false) + ); + } + + /// Reads update the tracker without the volume's write lock, so a + /// success replayed at a mark must not lose the errors they record + /// while it runs. + #[test] + fn success_at_a_mark_keeps_concurrent_errors() { + use std::sync::{Arc, Barrier}; + const READERS: i32 = 4; + const ERRORS: i32 = 200; + + for _ in 0..500 { + let tracker = Arc::new(IoErrorTracker::default()); + tracker.check_read_write_error(Some(&media_error())); + tracker.check_read_write_error(Some(&media_error())); + let mark = tracker.mark(); + let start = Arc::new(Barrier::new(READERS as usize + 1)); + let readers: Vec<_> = (0..READERS) + .map(|_| { + let (tracker, start) = (tracker.clone(), start.clone()); + std::thread::spawn(move || { + start.wait(); + for _ in 0..ERRORS { + tracker.check_read_write_error(Some(&media_error())); + } + }) + }) + .collect(); + start.wait(); + while tracker.get_io_error_state().1 < 2 + READERS * ERRORS / 2 { + std::hint::spin_loop(); + } + tracker.record_success_at(mark); + for reader in readers { + reader.join().unwrap(); + } + // The two errors before the mark are cleared; every error the + // readers recorded after it stands. + assert_eq!(tracker.get_io_error_state().1, READERS * ERRORS); + } + } + #[test] fn reset_io_error_state_lifts_the_quarantine() { let tracker = IoErrorTracker::default(); diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index 8be0fa6b8..e94612e65 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -18,7 +18,7 @@ use crate::storage::needle::needle::{Needle, get_actual_size}; use crate::storage::needle_map::NeedleMapKind; use crate::storage::super_block::{ReplicaPlacement, SUPER_BLOCK_SIZE}; use crate::storage::types::*; -use crate::storage::volume::{CompactionJob, VifVolumeInfo, VolumeError, VolumeSpec}; +use crate::storage::volume::{CompactionJob, VifVolumeInfo, Volume, VolumeError, VolumeSpec}; /// Top-level storage manager containing all disk locations and their volumes. pub struct Store { @@ -745,6 +745,30 @@ impl Store { n: &mut Needle, fsync: bool, ) -> Result<(u64, Size, bool), VolumeError> { + self.writable_volume_mut(vid)?.write_needle(n, true, fsync) + } + + /// Write a batch of needles to one volume, sharing the syncs of its + /// durable writes. See `Volume::write_needles_grouped`. + pub fn write_volume_needles( + &mut self, + vid: VolumeId, + writes: &mut [(Needle, bool)], + ) -> Vec> { + match self.writable_volume_mut(vid) { + Ok(vol) => vol.write_needles_grouped(writes), + // The lookup fails only with NotFound or the disk-space ReadOnly. + Err(e) => writes + .iter() + .map(|_| match e { + VolumeError::ReadOnly => Err(VolumeError::ReadOnly), + _ => Err(VolumeError::NotFound), + }) + .collect(), + } + } + + fn writable_volume_mut(&mut self, vid: VolumeId) -> Result<&mut Volume, VolumeError> { // Check disk space on the location containing this volume. // We do this before the mutable borrow to avoid borrow conflicts. let loc_idx = self @@ -759,7 +783,7 @@ impl Store { } let (_, vol) = self.find_volume_mut(vid).ok_or(VolumeError::NotFound)?; - vol.write_needle(n, true, fsync) + Ok(vol) } /// Delete a needle from a volume. diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index 0f8322a8d..7072cb454 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -9,6 +9,7 @@ //! Matches Go's storage/volume.go, volume_loading.go, volume_read.go, //! volume_write.go, volume_super_block.go. +use std::collections::HashSet; use std::fs::{self, File, OpenOptions}; use std::io::{self, Read, Seek, SeekFrom, Write}; use std::ops::ControlFlow; @@ -22,7 +23,7 @@ use tracing::{debug, error, info, warn}; use crate::storage::idx; use crate::storage::io::read_exact_at; -use crate::storage::io_error::IoErrorTracker; +use crate::storage::io_error::{IoErrorTracker, StreakMark}; use crate::storage::needle::needle::{self, Needle, NeedleError, get_actual_size}; use crate::storage::needle_map::sorted_file::SortedFileNeedleMap; use crate::storage::needle_map::{CompactNeedleMap, NeedleMap, NeedleMapKind, RedbNeedleMap}; @@ -1146,6 +1147,13 @@ pub struct Volume { fail_idx_sync_for_test: bool, #[cfg(test)] fail_truncate_for_test: bool, + /// Needle ids whose .dat append fails with a media error. + #[cfg(test)] + fail_append_for_test: HashSet, + #[cfg(test)] + dat_syncs_for_test: std::sync::atomic::AtomicUsize, + #[cfg(test)] + idx_syncs_for_test: std::sync::atomic::AtomicUsize, needle_map_kind: NeedleMapKind, data_file_access_control: Arc, @@ -1244,6 +1252,12 @@ impl Volume { fail_idx_sync_for_test: false, #[cfg(test)] fail_truncate_for_test: false, + #[cfg(test)] + fail_append_for_test: HashSet::new(), + #[cfg(test)] + dat_syncs_for_test: Default::default(), + #[cfg(test)] + idx_syncs_for_test: Default::default(), nm: None, needle_map_kind, data_file_access_control: Arc::new(DataFileAccessControl::default()), @@ -1288,6 +1302,12 @@ impl Volume { fail_idx_sync_for_test: false, #[cfg(test)] fail_truncate_for_test: false, + #[cfg(test)] + fail_append_for_test: HashSet::new(), + #[cfg(test)] + dat_syncs_for_test: Default::default(), + #[cfg(test)] + idx_syncs_for_test: Default::default(), nm: None, needle_map_kind: NeedleMapKind::InMemory, data_file_access_control: Arc::new(DataFileAccessControl::default()), @@ -2267,20 +2287,152 @@ impl Volume { fsync: bool, ) -> Result<(u64, Size, bool), VolumeError> { let _guard = self.data_file_access_control.write_lock(); + self.check_writable()?; + self.do_write_request(n, check_cookie, fsync) + } + + /// Write a batch of needles the way Go's processBatch does: one .dat + /// sync and one .idx sync per run of distinct needle ids that holds a + /// durable write, instead of two per durable needle. Nothing in such a + /// run is published before its sync, so a failed sync takes the whole + /// run back off the .dat and fails every entry. A repeated id starts a + /// new run, so its dedup and cookie checks see the earlier write. + pub fn write_needles_grouped( + &mut self, + writes: &mut [(Needle, bool)], + ) -> Vec> { + let _guard = self.data_file_access_control.write_lock(); + let mut results = Vec::with_capacity(writes.len()); + let mut rest = writes; + while !rest.is_empty() { + let mut ids = HashSet::new(); + let len = rest.iter().take_while(|(n, _)| ids.insert(n.id)).count(); + let (run, tail) = std::mem::take(&mut rest).split_at_mut(len); + if run.iter().any(|(_, fsync)| *fsync) { + results.extend(self.write_synced_run(run)); + } else { + for (n, _) in run.iter_mut() { + results.push( + self.check_writable() + .and_then(|()| self.do_write_request(n, true, false)), + ); + } + } + rest = tail; + } + results + } + + fn write_synced_run( + &mut self, + run: &mut [(Needle, bool)], + ) -> Vec> { + // Per entry: Some(offset) once appended, None when it dedups. + let mut staged = Vec::with_capacity(run.len()); + // Per entry: the I/O error streak once it is staged, where a write + // sent on its own would have recorded its success. + let mut marks = Vec::with_capacity(run.len()); + let mut last_append_at_ns = self.last_append_at_ns; + let mut run_start = None; + let mut sync = false; + for (n, fsync) in run.iter_mut() { + let r = self.append_unpublished(n, &mut last_append_at_ns); + marks.push(self.io_errors.mark()); + if let Ok(Some(offset)) = r { + run_start.get_or_insert(offset); + } + // Only a durable entry that got this far needs the sync. + sync |= *fsync && r.is_ok(); + staged.push(r); + } + + if sync && let Err(e) = self.flush_dat() { + self.check_read_write_error(Some(&e)); + if let Some(start) = run_start { + self.undo_unsynced_append(start); + } + let e = VolumeError::Io(e); + return staged + .into_iter() + .map(|r| r.and_then(|_| Err(run_error(&e)))) + .collect(); + } + self.last_append_at_ns = last_append_at_ns; + + for ((n, fsync), r) in run.iter().zip(staged.iter_mut()) { + if let Ok(Some(offset)) = *r + && let Err(e) = self.publish_write(n, offset, *fsync) + { + *r = Err(e); + } + } + + if sync && let Err(e) = self.flush_idx() { + return staged + .into_iter() + .map(|r| r.and_then(|_| Err(run_error(&e)))) + .collect(); + } + + let written = run + .iter() + .zip(&staged) + .filter(|(_, r)| matches!(r, Ok(Some(_)))) + .map(|((n, _), _)| n.last_modified) + .max(); + let last = staged.iter().rposition(|r| matches!(r, Ok(Some(_)))); + if let (Some(last_modified), Some(last)) = (written, last) { + // Sent one at a time, the last write to land would have cleared + // the streak of the appends before it, and the ones after it + // would have failed on top of that. + self.finish_write(last_modified, sync, marks[last]); + } + + run.iter() + .zip(staged) + .map(|((n, _), r)| { + let size = Size(n.data_size as i32); + r.map(|staged| match staged { + Some(offset) => (offset, size, false), + None => (0, size, true), + }) + }) + .collect() + } + + /// The checks and the append for one entry of a synced run, leaving the + /// publish to the caller. Returns the offset, or None when it dedups. + fn append_unpublished( + &mut self, + n: &mut Needle, + last_append_at_ns: &mut u64, + ) -> Result, VolumeError> { + self.check_writable()?; + if self.prepare_write(n, true)? { + return Ok(None); + } + n.append_at_ns = get_append_at_ns(*last_append_at_ns); + let (offset, _, _) = self.append_needle(n)?; + *last_append_at_ns = n.append_at_ns; + Ok(Some(offset)) + } + + fn check_writable(&self) -> Result<(), VolumeError> { if let Some(e) = self.unavailable_error() { return Err(e); } if self.is_read_only() { return Err(VolumeError::ReadOnly); } - - self.do_write_request(n, check_cookie, fsync) + Ok(()) } /// Flush the .dat, the first half of a durable write. The .idx is flushed /// separately by flush_idx once the row is published; the two are split so /// nothing is indexed before the bytes it points at are down. fn flush_dat(&self) -> io::Result<()> { + #[cfg(test)] + self.dat_syncs_for_test.fetch_add(1, Ordering::Relaxed); #[cfg(test)] if self.fail_fsync_for_test { return Err(io::Error::other("injected fsync failure")); @@ -2303,6 +2455,8 @@ impl Volume { /// record whose index may not survive, and the master routes writes /// elsewhere once the volume heartbeats read only. fn flush_idx(&mut self) -> Result<(), VolumeError> { + #[cfg(test)] + self.idx_syncs_for_test.fetch_add(1, Ordering::Relaxed); #[cfg(test)] if self.fail_idx_sync_for_test { let e = io::Error::other("injected idx sync failure"); @@ -2345,6 +2499,54 @@ impl Volume { check_cookie: bool, fsync: bool, ) -> Result<(u64, Size, bool), VolumeError> { + if self.prepare_write(n, check_cookie)? { + // Nothing to append, but the write this matched may have been + // non-durable, and the caller is asking for the content to be on + // disk. Its .idx row can be sitting in the page cache too, so both + // files get flushed exactly as they would for a fresh append. + if fsync { + self.flush_dat().map_err(|e| { + self.check_read_write_error(Some(&e)); + VolumeError::Io(e) + })?; + self.flush_idx()?; + } + return Ok((0, Size(n.data_size as i32), true)); + } + + // Update append timestamp + n.append_at_ns = get_append_at_ns(self.last_append_at_ns); + + // Append to .dat file + let (offset, _body_size, _actual_size) = self.append_needle(n)?; + + // Nothing is published until the bytes are down: an index entry for an + // unflushed append would resolve past the end of the file after a crash, + // and undoing it afterwards would double-count the volume's metrics. + if fsync && let Err(e) = self.flush_dat() { + self.check_read_write_error(Some(&e)); + self.undo_unsynced_append(offset); + return Err(VolumeError::Io(e)); + } + + self.last_append_at_ns = n.append_at_ns; + + self.publish_write(n, offset, fsync)?; + + if fsync { + self.flush_idx()?; + } + + self.finish_write(n.last_modified, fsync, self.io_errors.mark()); + + // Return Size(n.DataSize) as the logical size, matching Go's doWriteRequest + Ok((offset, Size(n.data_size as i32), false)) + } + + /// The checks a write passes before anything is appended: TTL + /// inheritance, checksum, dedup and cookie. Returns true when the needle + /// matches the stored copy and there is nothing to append. + fn prepare_write(&self, n: &mut Needle, check_cookie: bool) -> Result { // TTL inheritance from volume (matching Go's writeNeedle2) { use crate::storage::needle::ttl::TTL; @@ -2363,18 +2565,7 @@ impl Volume { // Dedup check (matches Go: n.DataSize = oldNeedle.DataSize on dedup) if let Some(old_data_size) = self.is_file_unchanged(n) { n.data_size = old_data_size; - // Nothing to append, but the write this matched may have been - // non-durable, and the caller is asking for the content to be on - // disk. Its .idx row can be sitting in the page cache too, so both - // files get flushed exactly as they would for a fresh append. - if fsync { - self.flush_dat().map_err(|e| { - self.check_read_write_error(Some(&e)); - VolumeError::Io(e) - })?; - self.flush_idx()?; - } - return Ok((0, Size(n.data_size as i32), true)); + return Ok(true); } // Cookie validation for existing needle (matches Go: check whenever nm.Get returns ok) @@ -2392,32 +2583,24 @@ impl Volume { return Err(VolumeError::CookieMismatch(n.cookie.0)); } } + Ok(false) + } - // Update append timestamp - n.append_at_ns = get_append_at_ns(self.last_append_at_ns); - - // Append to .dat file - let (offset, _body_size, _actual_size) = self.append_needle(n)?; - - // Nothing is published until the bytes are down: an index entry for an - // unflushed append would resolve past the end of the file after a crash, - // and undoing it afterwards would double-count the volume's metrics. - if fsync && let Err(e) = self.flush_dat() { - self.check_read_write_error(Some(&e)); - if let Err(te) = self.truncate_dat(offset) { - // The rejected record is still on the end. A later append - // would bury it mid-file, where the .dat tail check cannot - // see it, so the volume fails closed instead. - self.mark_io_unavailable(format!( - "failed to truncate back to {} after a failed fsync: {}", - offset, te - )); - } - return Err(VolumeError::Io(e)); + /// Take an append whose sync failed back off the .dat. + fn undo_unsynced_append(&mut self, offset: u64) { + if let Err(te) = self.truncate_dat(offset) { + // The rejected record is still on the end. A later append + // would bury it mid-file, where the .dat tail check cannot + // see it, so the volume fails closed instead. + self.mark_io_unavailable(format!( + "failed to truncate back to {} after a failed fsync: {}", + offset, te + )); } + } - self.last_append_at_ns = n.append_at_ns; - + /// Index an appended needle. `fsync` means its record is already down. + fn publish_write(&mut self, n: &Needle, offset: u64, fsync: bool) -> Result<(), VolumeError> { // Update needle map (uses n.size = full body size, matching Go's nm.Put) let prior = match self.nm.as_ref() { Some(nm) => nm.get(n.id), @@ -2465,26 +2648,25 @@ impl Volume { return Err(VolumeError::Io(e)); } } + Ok(()) + } - if fsync { - self.flush_idx()?; + /// The bookkeeping after a write is fully down. `landed_at` is the point + /// in the I/O error streak where the write landed; errors counted after + /// it, by later appends of the same run, are not cleared. + fn finish_write(&mut self, last_modified: u64, idx_synced: bool, landed_at: StreakMark) { + if self.last_modified_ts_seconds < last_modified { + self.last_modified_ts_seconds = last_modified; } - if self.last_modified_ts_seconds < n.last_modified { - self.last_modified_ts_seconds = n.last_modified; - } - - let checkpoint_ok = self.maybe_checkpoint_index(fsync); + let checkpoint_ok = self.maybe_checkpoint_index(idx_synced); // Clear the EIO streak only after the full write (data + flush + // index + checkpoint) succeeds, so a failed fsync or checkpoint // does not get its EIO erased by the success reset. if checkpoint_ok { - self.check_read_write_error(None); + self.io_errors.record_success_at(landed_at); } - - // Return Size(n.DataSize) as the logical size, matching Go's doWriteRequest - Ok((offset, Size(n.data_size as i32), false)) } /// Take the index checkpoint the needle map asked for, data first: the @@ -2596,9 +2778,15 @@ impl Volume { }); } - if let Err(e) = dat_file.write_all(&bytes) { - // Truncate back to pre-write position on error (matching Go) - let _ = dat_file.set_len(offset); + let written = dat_file.write_all(&bytes); + #[cfg(test)] + let written = if self.fail_append_for_test.contains(&n.id) { + Err(media_error_for_test()) + } else { + written + }; + if let Err(e) = written { + self.undo_unsynced_append(offset); self.check_read_write_error(Some(&e)); return Err(VolumeError::Io(e)); } @@ -4830,6 +5018,20 @@ impl Volume { self.fail_truncate_for_test = fail; } + #[cfg(test)] + pub(crate) fn fail_append_for_test(&mut self, ids: &[NeedleId]) { + self.fail_append_for_test = ids.iter().copied().collect(); + } + + /// (.dat syncs, .idx syncs) attempted since the volume was opened. + #[cfg(test)] + pub(crate) fn sync_counts_for_test(&self) -> (usize, usize) { + ( + self.dat_syncs_for_test.load(Ordering::Relaxed), + self.idx_syncs_for_test.load(Ordering::Relaxed), + ) + } + #[cfg(test)] pub(crate) fn set_last_modified_ts_for_test(&mut self, ts_seconds: u64) { self.last_modified_ts_seconds = ts_seconds; @@ -4847,6 +5049,14 @@ impl Volume { // Helpers // ============================================================================ +/// A copy of a run-wide failure for each entry it fails. +fn run_error(e: &VolumeError) -> VolumeError { + match e { + VolumeError::Io(e) => VolumeError::Io(io::Error::new(e.kind(), e.to_string())), + e => VolumeError::Io(io::Error::other(e.to_string())), + } +} + /// Generate volume file base name: dir/collection_id or dir/id /// Byte offset just past the needle's on-disk record. Mirrors Go's /// needleDiskEnd. @@ -5096,6 +5306,25 @@ fn preallocate_file(file: &File, size: u64) { // Tests // ============================================================================ +/// An OS error the platform reports for failing storage media, which is +/// what counts toward the I/O error streak. +#[cfg(test)] +fn media_error_for_test() -> io::Error { + #[cfg(unix)] + { + io::Error::from_raw_os_error(libc::EIO) + } + #[cfg(windows)] + { + const ERROR_IO_DEVICE: i32 = 1117; + io::Error::from_raw_os_error(ERROR_IO_DEVICE) + } + #[cfg(not(any(unix, windows)))] + { + io::Error::other("injected media error") + } +} + #[cfg(test)] mod tests { use super::*; @@ -6074,6 +6303,309 @@ mod tests { assert!(v.read_needle(&mut read_n).is_err()); } + fn batch_needle(id: u64, cookie: u32, data: &[u8]) -> Needle { + Needle { + id: NeedleId(id), + cookie: Cookie(cookie), + data: data.to_vec(), + data_size: data.len() as u32, + ..Needle::default() + } + } + + fn dat_len(v: &Volume) -> u64 { + std::fs::metadata(v.file_name(".dat")).unwrap().len() + } + + /// A batch of durable writes shares one .dat sync and one .idx sync, + /// where writing them one at a time syncs both files per needle. + #[test] + fn test_grouped_fsync_writes_share_one_sync() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let mut writes: Vec<_> = (1..=10u64) + .map(|id| { + ( + batch_needle(id, 0xaa, format!("body-{id}").as_bytes()), + true, + ) + }) + .collect(); + let results = v.write_needles_grouped(&mut writes); + assert!(results.iter().all(|r| matches!(r, Ok((_, _, false))))); + assert_eq!(v.sync_counts_for_test(), (1, 1)); + for id in 1..=10u64 { + let mut n = Needle { + id: NeedleId(id), + ..Needle::default() + }; + v.read_needle(&mut n).unwrap(); + assert_eq!(n.data, format!("body-{id}").as_bytes()); + } + + // A batch with no durable write syncs nothing, as before. + let mut writes: Vec<_> = (11..=20u64) + .map(|id| (batch_needle(id, 0xaa, b"lazy"), false)) + .collect(); + let results = v.write_needles_grouped(&mut writes); + assert!(results.iter().all(|r| r.is_ok())); + assert_eq!(v.sync_counts_for_test(), (1, 1)); + + // The per-needle path pays both syncs for every durable write. + let tmp2 = TempDir::new().unwrap(); + let mut one_by_one = make_test_volume(tmp2.path().to_str().unwrap()); + for id in 1..=10u64 { + let mut n = batch_needle(id, 0xaa, format!("body-{id}").as_bytes()); + one_by_one.write_needle(&mut n, true, true).unwrap(); + } + assert_eq!(one_by_one.sync_counts_for_test(), (10, 10)); + } + + /// A failed shared sync fails every entry of the run, durable or not, + /// and leaves the volume as it was before the run: the .dat back at the + /// run start, the clocks where they were, nothing published. + #[test] + fn test_grouped_failed_sync_rolls_back_the_run() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let mut kept = batch_needle(1, 0xaa, b"first-copy"); + v.write_needle(&mut kept, true, true).unwrap(); + let prior = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap(); + let dat_len_before = dat_len(&v); + let last_append_before = v.last_append_at_ns; + let last_modified_before = v.last_modified_ts_seconds; + let file_count_before = v.file_count(); + + let mut writes = vec![ + (batch_needle(1, 0xaa, b"second-copy"), false), + (batch_needle(2, 0xbb, b"durable"), true), + (batch_needle(3, 0xcc, b"lazy"), false), + ]; + for (n, _) in writes.iter_mut() { + n.last_modified = last_modified_before + 1000; + } + v.fail_next_fsync_for_test(true); + let results = v.write_needles_grouped(&mut writes); + v.fail_next_fsync_for_test(false); + + assert!( + results.iter().all(|r| matches!(r, Err(VolumeError::Io(_)))), + "every entry of the run shares its failed sync: {results:?}" + ); + assert_eq!(dat_len(&v), dat_len_before, "the run is off the .dat"); + assert_eq!(v.last_append_at_ns, last_append_before); + assert_eq!(v.last_modified_ts_seconds, last_modified_before); + let now = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap(); + assert_eq!((now.offset, now.size), (prior.offset, prior.size)); + assert!(v.nm.as_ref().unwrap().get(NeedleId(2)).unwrap().is_none()); + assert!(v.nm.as_ref().unwrap().get(NeedleId(3)).unwrap().is_none()); + assert_eq!(v.file_count(), file_count_before); + assert!( + !v.is_read_only(), + "a rolled-back run keeps the volume writable" + ); + + let mut read_n = Needle { + id: NeedleId(1), + ..Needle::default() + }; + v.read_needle(&mut read_n).unwrap(); + assert_eq!(read_n.data, b"first-copy"); + + // The same run goes through once the disk recovers. + let results = v.write_needles_grouped(&mut writes); + assert!(results.iter().all(|r| r.is_ok())); + assert!(v.last_append_at_ns > last_append_before); + } + + /// A repeated id starts a new run, so the second write sees the first + /// one published: a different cookie is refused and the same content + /// dedups, exactly as two sequential writes would. + #[test] + fn test_grouped_repeated_id_behaves_as_sequential_writes() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let mut writes = vec![ + (batch_needle(1, 0xaa, b"one"), true), + (batch_needle(1, 0xbb, b"imposter"), true), + (batch_needle(2, 0xcc, b"two"), true), + (batch_needle(2, 0xcc, b"two"), true), + ]; + let results = v.write_needles_grouped(&mut writes); + + assert!(matches!(results[0], Ok((_, _, false)))); + assert!(matches!(results[1], Err(VolumeError::CookieMismatch(0xbb)))); + assert!(matches!(results[2], Ok((_, _, false)))); + assert!( + matches!(results[3], Ok((0, _, true))), + "the same content dedups against the write just before it" + ); + // Runs [1], [1, 2], [2]: one pair of syncs each. + assert_eq!(v.sync_counts_for_test(), (3, 3)); + + let mut read_n = Needle { + id: NeedleId(1), + ..Needle::default() + }; + v.read_needle(&mut read_n).unwrap(); + assert_eq!(read_n.data, b"one"); + } + + /// A failed shared .idx sync fails every published entry and stops the + /// volume taking writes, as it does for a single durable write. + #[test] + fn test_grouped_failed_idx_sync_quarantines_volume() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let mut writes = vec![ + (batch_needle(1, 0xaa, b"durable"), true), + (batch_needle(2, 0xbb, b"lazy"), false), + ]; + v.fail_next_idx_sync_for_test(true); + let results = v.write_needles_grouped(&mut writes); + v.fail_next_idx_sync_for_test(false); + + assert!(results.iter().all(|r| matches!(r, Err(VolumeError::Io(_))))); + assert_eq!(v.sync_counts_for_test(), (1, 1)); + assert!(v.is_read_only()); + } + + /// A failed shared sync whose truncate also fails leaves an unverified + /// tail, so the volume fails closed exactly as a single write does. + #[test] + fn test_grouped_failed_rollback_marks_volume_unavailable() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let mut writes = vec![ + (batch_needle(1, 0xaa, b"never"), true), + (batch_needle(2, 0xbb, b"landed"), false), + ]; + v.fail_next_fsync_for_test(true); + v.fail_next_truncate_for_test(true); + let results = v.write_needles_grouped(&mut writes); + v.fail_next_fsync_for_test(false); + v.fail_next_truncate_for_test(false); + + assert!(results.iter().all(|r| r.is_err())); + assert!(v.unavailable_error().is_some()); + assert!(v.is_read_only()); + assert!(v.should_quarantine()); + assert!(v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().is_none()); + + let mut later = vec![(batch_needle(3, 0xcc, b"refused"), true)]; + let results = v.write_needles_grouped(&mut later); + assert!(matches!(results[0], Err(VolumeError::Unavailable(_)))); + } + + /// The I/O error streak after writing four durable needles, the ones in + /// `failing` with a media error on their append, first one at a time and + /// then as one grouped run. Also returns which grouped writes landed. + #[cfg(any(unix, windows))] + fn io_error_streaks(failing: &[u64]) -> (i32, i32, Vec) { + let writes = || -> Vec<_> { + (1..=4u64) + .map(|id| { + ( + batch_needle(id, 0xaa, format!("body-{id}").as_bytes()), + true, + ) + }) + .collect() + }; + let failing: Vec<_> = failing.iter().map(|&id| NeedleId(id)).collect(); + + let tmp = TempDir::new().unwrap(); + let mut one_by_one = make_test_volume(tmp.path().to_str().unwrap()); + one_by_one.fail_append_for_test(&failing); + for (mut n, fsync) in writes() { + let _ = one_by_one.write_needle(&mut n, true, fsync); + } + + let tmp2 = TempDir::new().unwrap(); + let mut v = make_test_volume(tmp2.path().to_str().unwrap()); + v.fail_append_for_test(&failing); + let mut grouped = writes(); + let results = v.write_needles_grouped(&mut grouped); + assert!(!v.is_read_only(), "the failed appends were truncated back"); + + let landed = results + .iter() + .map(|r| matches!(r, Ok((_, _, false)))) + .collect(); + ( + one_by_one.get_io_error_state().1, + v.get_io_error_state().1, + landed, + ) + } + + /// A run counts I/O errors as the same writes sent one at a time would: + /// a success early in the run must not wipe out the streak that the + /// failed appends after it built up, or the volume escapes quarantine. + #[cfg(any(unix, windows))] + #[test] + fn test_grouped_run_keeps_the_io_error_streak_of_later_appends() { + use crate::storage::io_error::IO_ERROR_TOLERANCE; + + let (one_by_one, grouped, landed) = io_error_streaks(&[2, 3, 4]); + + assert_eq!(landed, [true, false, false, false]); + assert_eq!(one_by_one, IO_ERROR_TOLERANCE); + assert_eq!(grouped, IO_ERROR_TOLERANCE); + } + + /// The other half: a write that lands clears the errors of the appends + /// queued before it, even when another append after it fails, or the + /// run reaches a quarantine that the same writes one at a time do not. + #[cfg(any(unix, windows))] + #[test] + fn test_grouped_run_clears_the_io_error_streak_of_earlier_appends() { + let (one_by_one, grouped, landed) = io_error_streaks(&[1, 2, 4]); + + assert_eq!(landed, [false, false, true, false]); + assert_eq!(one_by_one, 1); + assert_eq!(grouped, 1); + } + + /// An append whose partial bytes cannot be truncated back leaves the .dat + /// tail unverified, so the volume fails closed rather than let a later + /// append bury them mid-file. + #[test] + fn test_failed_append_rollback_marks_volume_unavailable() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let mut first = batch_needle(1, 0xaa, b"landed"); + v.write_needle(&mut first, true, true).unwrap(); + + // A read-only .dat handle fails the append and the truncate-back alike. + v.dat_file = Some(File::open(v.file_name(".dat")).unwrap()); + let mut n = batch_needle(2, 0xbb, b"never-lands"); + v.write_needle(&mut n, true, false).unwrap_err(); + + assert!(v.unavailable_error().is_some()); + assert!(v.is_read_only()); + assert!(v.should_quarantine()); + assert!(v.nm.as_ref().unwrap().get(NeedleId(2)).unwrap().is_none()); + + let mut later = batch_needle(3, 0xcc, b"refused"); + assert!(matches!( + v.write_needle(&mut later, true, true), + Err(VolumeError::Unavailable(_)) + )); + } + #[test] fn test_volume_write_dedup() { let tmp = TempDir::new().unwrap();