mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
volume server: group-commit fsync writes in the write queue (#11543)
* volume: split the write path into reusable steps do_write_request ran its pre-append checks, the append, the sync rollback, the index publish and the post-write bookkeeping inline, so a batched write could only reuse it one needle at a time. Pull the steps out (check_writable, prepare_write, undo_unsynced_append, publish_write, finish_write) and the store's volume lookup plus disk-space check (writable_volume_mut). do_write_request composes them in the same order with the same early returns; no behaviour change. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume server: group-commit fsync writes in the write queue The write queue holds one store lock for a batch of up to 128 needles but wrote them one at a time, so every fsync needle paid its own .dat sync and its own .idx sync: 2N syncs per batch. Add Volume::write_needles_grouped, after Go's processBatch. A volume's entries are split into runs of distinct needle ids (a repeated id starts a new run, so its dedup and cookie checks see the earlier write). A run with a durable entry appends everything with append_at_ns chained through a local, syncs the .dat once, and only then publishes the entries and syncs the .idx once. A failed .dat sync truncates the .dat back to the run start (marking the volume unavailable if that fails), leaves last_append_at_ns and last_modified untouched, and fails every entry of the run. Runs with no durable entry go through the unchanged per-needle path. Store::write_volume_needles is the queue's entry point; the handlers' non-queue path is unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume server: fail closed when a failed append's rollback fails append_needle discarded the truncate-back result, so a partial write that could not be rolled back left unindexed bytes on the .dat while the volume stayed writable; the next append would bury them mid-file, past the load-time tail check. Route the rollback through undo_unsynced_append, which marks the volume unavailable when the truncate fails, so nothing more is appended over an unverified tail. * volume server: keep a grouped run's I/O error streak from later appends A synced run stages every append before any entry finishes, so the success reset in finish_write ran after the failed appends queued behind the last write to land and erased their media-error streak. Sent one at a time, those errors would have counted and quarantined the volume. Skip the reset when an append after the last landed write added to the streak. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume server: replay a grouped run's I/O error streak in queue order Skipping the run's success reset whenever an append after the last landed write failed kept the errors from before that write as well, so a run like [EIO, EIO, landed, EIO] reached the quarantine count that the same writes one at a time (one error) do not. Mark the streak where each entry is staged and record the run's success at the last landed write's mark: errors before it are cleared, the ones after it still count. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume server: replay a run's I/O error streak in one atomic step Reads record their outcomes on the tracker without the volume's write lock, so record_success_at's separate load and store could drop an error a read counted in between, or restore errors a read had just cleared. Keep the count and the clear counter in one atomic word and apply the replay with a single fetch_update. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Co-authored-by: Chris Lu <chris.lu@gmail.com>
This commit is contained in:
4 files changed
+837
-68
No files matched your search
@@ -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<VolumeServerState>, batch: Vec<WriteRequest>) {
|
||||
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.
|
||||
|
||||
@@ -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<Option<String>>,
|
||||
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<u64>) -> Option<u64> {
|
||||
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<String>, 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();
|
||||
|
||||
@@ -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<Result<(u64, Size, bool), VolumeError>> {
|
||||
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.
|
||||
|
||||
@@ -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<NeedleId>,
|
||||
#[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<DataFileAccessControl>,
|
||||
|
||||
@@ -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<Result<(u64, Size, bool), VolumeError>> {
|
||||
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<Result<(u64, Size, bool), VolumeError>> {
|
||||
// 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<Option<u64>, 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<bool, VolumeError> {
|
||||
// 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<bool>) {
|
||||
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();
|
||||
|
||||
Reference in new issue
Block a user