[Volume] Keep DAT and index state consistent after async batch Sync failure (#11425)

* fix 11400

* persist failed-recovery quarantine and harden rollback

- record the unavailable state in a .unavailable marker, fsync it, and
  re-arm it on load so a restart cannot serve an unverified pair
- quarantine the volume so heartbeats stop advertising it
- block MarkVolumeWritable while unavailable, rechecked under noWriteLock
- fail every request of a failed batch, not only the succeeded ones
- restore the needle map and truncate .dat on inline fsync rollback failure
- add truncateIndex for the sorted-file needle map
- mirror the fail-closed semantics in the Rust volume server

* volume: erase rolled-back mappings instead of leaving tombstones

A rolled-back batch or failed inline write used Delete() to undo a
needle that did not exist beforehand, leaving a tombstoned map entry
whose stale offset makes the next write to that needle fail reading a
header that no longer exists. Add removeMapping/restoreMapping to the
mappers so recovery erases entries that were absent before the batch
and reinstates the exact prior offset/size for ones that were,
including tombstones. The index row still goes through Delete so a
replay forgets the needle.

* volume: gate bulk readers on unavailable and fsync the marker's dir

- fsync_dir(&self.dir) synced the volume dir's parent, not the dir
  holding .unavailable; pass the marker path so the create survives
  a host crash
- export UnavailableError and check it in ReadAllNeedles,
  VolumeTailSender, VolumeIncrementalCopy, and IncrementalBackup so
  replica-sync paths cannot stream or append data from an unverified
  .dat/.idx pair; mirror on the Rust side via read_dat_slice,
  read_all_needles, dat_scan_plan, and the incremental-copy handler

* volume: drop issue references from comments near touched code

* volume: stop active scans when the volume becomes unavailable

The stream entry-point checks ran once per RPC, so a volume quarantined
by a failed recovery mid-scan kept serving data. Recheck availability
per needle/chunk on the detached read paths: tail scan and heartbeat,
read-all, incremental copy, incremental backup writes, and the Rust
StreamingBody chunk reads. Rust incremental copy also rejects a
quarantined volume before sync_to_disk touches the backend.

---------

Co-authored-by: Chris Lu <chris.lu@gmail.com>
This commit is contained in:
ssshr-66
2026-09-24 06:57:44 +08:00
committed by GitHub
co-authored by Chris Lu
parent d848b8ed00
commit b9ad62fc16
21 changed files with 1181 additions and 90 deletions
+41 -3
View File
@@ -1353,18 +1353,25 @@ impl VolumeServer for VolumeGrpcService {
let req = request.into_inner();
let vid = VolumeId(req.volume_id);
// Sync to disk first
// A quarantined volume is rejected before the sync does disk I/O for it.
{
let mut store = self.state.store.write().unwrap();
if let Some((_, v)) = store.find_volume_mut(vid) {
let _ = v.sync_to_disk();
let Some((_, v)) = store.find_volume_mut(vid) else {
return Err(Status::not_found(format!("not found volume id {}", vid)));
};
if let Some(e) = v.unavailable_error() {
return Err(Status::unavailable(e.to_string()));
}
let _ = v.sync_to_disk();
}
let store = self.state.store.read().unwrap();
let (_, v) = store
.find_volume(vid)
.ok_or_else(|| Status::not_found(format!("not found volume id {}", vid)))?;
if let Some(e) = v.unavailable_error() {
return Err(Status::unavailable(e.to_string()));
}
let dat_size = v.dat_file_size().unwrap_or(0);
let super_block_size = v.super_block.block_size() as u64;
@@ -1438,6 +1445,7 @@ impl VolumeServer for VolumeGrpcService {
.map_err(|e| Status::internal(format!("open {}: {}", path, e)))?;
DatReader::Local(file)
};
let state = self.state.clone();
drop(store);
let total = dat_size - start_offset;
@@ -1445,7 +1453,20 @@ impl VolumeServer for VolumeGrpcService {
Result<volume_server_pb::VolumeIncrementalCopyResponse, Status>,
>(8);
// The reader handle is detached from the volume, so each chunk
// re-checks that the volume has not been quarantined mid-stream.
tokio::task::spawn_blocking(move || {
macro_rules! bail_if_unavailable {
() => {{
let store = state.store.read().unwrap();
if let Some((_, v)) = store.find_volume(vid)
&& let Some(e) = v.unavailable_error()
{
let _ = tx.blocking_send(Err(Status::unavailable(e.to_string())));
return;
}
}};
}
let buffer_size = 2 * 1024 * 1024u64; // 2MB chunks
let mut bytes_to_read = total;
let mut offset = start_offset;
@@ -1460,6 +1481,7 @@ impl VolumeServer for VolumeGrpcService {
return;
}
while bytes_to_read > 0 {
bail_if_unavailable!();
let chunk = std::cmp::min(bytes_to_read, buffer_size) as usize;
let mut buf = vec![0u8; chunk];
match reader.read(&mut buf) {
@@ -1486,6 +1508,7 @@ impl VolumeServer for VolumeGrpcService {
// handle. No store lock is held while the (potentially slow)
// S3 fetch runs, so it never blocks store writers.
while bytes_to_read > 0 {
bail_if_unavailable!();
let chunk = std::cmp::min(bytes_to_read, buffer_size) as usize;
match remote.read_slice(offset, chunk) {
Ok(buf) if buf.is_empty() => break,
@@ -5904,7 +5927,19 @@ fn tail_pass(
let mut last_processed_ns = last_timestamp_ns;
let mut sent_any = false;
let mut client_gone = false;
let mut unavailable: Option<String> = None;
let scanned = plan.scan(|needle| {
// The plan is detached from the volume, so a failed recovery marking
// it unavailable mid-scan would otherwise keep streaming.
{
let store = state.store.read().unwrap();
if let Some((_, vol)) = store.find_volume(vid)
&& let Some(e) = vol.unavailable_error()
{
unavailable = Some(e.to_string());
return ControlFlow::Break(());
}
}
// Notice a receiver that hung up between sends too, so a pass over
// needles it already has does not read on for nobody.
if tx.is_closed() {
@@ -5938,6 +5973,9 @@ fn tail_pass(
ControlFlow::Continue(())
});
if let Some(reason) = unavailable {
return TailPass::Failed(format!("volume {} is unavailable: {}", vid, reason));
}
if client_gone {
return TailPass::ClientGone;
}
+3
View File
@@ -211,6 +211,9 @@ impl http_body::Body for StreamingBody {
let relookup_result = {
let store = self.server_state.store.read().unwrap();
if let Some((_, vol)) = store.find_volume(self.volume_id) {
if let Some(e) = vol.unavailable_error() {
return std::task::Poll::Ready(Some(Err(std::io::Error::other(e))));
}
if vol.super_block.compaction_revision != self.compaction_revision {
// Compaction occurred — re-lookup the needle's data offset
Some(vol.re_lookup_needle_data_offset(self.needle_id))
+218 -20
View File
@@ -63,6 +63,9 @@ pub enum VolumeError {
#[error("volume is read-only")]
ReadOnly,
#[error("volume is unavailable: {0}")]
Unavailable(String),
#[error("volume size limit exceeded: current {current}, limit {limit}")]
SizeLimitExceeded { current: u64, limit: u64 },
@@ -667,6 +670,8 @@ pub struct Volume {
fail_fsync_for_test: bool,
#[cfg(test)]
fail_idx_sync_for_test: bool,
#[cfg(test)]
fail_truncate_for_test: bool,
needle_map_kind: NeedleMapKind,
data_file_access_control: Arc<DataFileAccessControl>,
@@ -675,6 +680,11 @@ pub struct Volume {
no_write_or_delete: bool,
no_write_can_delete: bool,
/// Set when a failed recovery leaves the .dat/index pair unverified: all
/// I/O is refused and a `.unavailable` marker keeps the volume quarantined
/// across restarts. Mirrors Go's ioUnavailable.
io_unavailable: Option<String>,
/// Shared flag from the parent DiskLocation indicating low disk space.
/// Matches Go's `v.location.isDiskSpaceLow` checked in `IsReadOnly()`.
pub location_disk_space_low: Arc<AtomicBool>,
@@ -757,6 +767,8 @@ impl Volume {
fail_fsync_for_test: false,
#[cfg(test)]
fail_idx_sync_for_test: false,
#[cfg(test)]
fail_truncate_for_test: false,
nm: None,
needle_map_kind,
data_file_access_control: Arc::new(DataFileAccessControl::default()),
@@ -767,6 +779,7 @@ impl Volume {
},
no_write_or_delete: false,
no_write_can_delete: false,
io_unavailable: None,
location_disk_space_low: Arc::new(AtomicBool::new(false)),
last_modified_ts_seconds: 0,
last_append_at_ns: 0,
@@ -798,12 +811,15 @@ impl Volume {
fail_fsync_for_test: false,
#[cfg(test)]
fail_idx_sync_for_test: false,
#[cfg(test)]
fail_truncate_for_test: false,
nm: None,
needle_map_kind: NeedleMapKind::InMemory,
data_file_access_control: Arc::new(DataFileAccessControl::default()),
super_block: SuperBlock::default(),
no_write_or_delete: false,
no_write_can_delete: false,
io_unavailable: None,
location_disk_space_low: Arc::new(AtomicBool::new(false)),
last_modified_ts_seconds: 0,
last_append_at_ns: 0,
@@ -1020,7 +1036,6 @@ impl Volume {
// — no extra disk I/O. A violation marks the volume read-only
// so vacuum doesn't silently drop reachable data based on a
// corrupt .idx left over from a crashed batched write.
// See issue #8928.
if let Some(ref nm) = self.nm
&& let Ok(dat_size) = self.current_dat_file_size()
{
@@ -1038,6 +1053,8 @@ impl Volume {
}
}
self.restore_unavailable();
// Match Go: if no .vif file existed, create one with version and bytes_offset
if !has_volume_info_file {
self.volume_info.version = self.super_block.version.0 as u32;
@@ -1338,6 +1355,9 @@ impl Volume {
/// remote-only tiered volumes whose `.dat` is no longer present locally.
pub fn read_dat_slice(&self, offset: u64, size: usize) -> Result<Vec<u8>, VolumeError> {
let _guard = self.data_file_access_control.read_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
let dat_size = self.current_dat_file_size()?;
if size == 0 || offset >= dat_size {
return Ok(Vec::new());
@@ -1425,6 +1445,9 @@ impl Volume {
read_option: &mut ReadOption,
) -> Result<i32, VolumeError> {
let _guard = self.data_file_access_control.read_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
let nm = self.nm_or_not_found()?;
let nv = nm.get(n.id)?.ok_or(VolumeError::NotFound)?;
@@ -1487,6 +1510,9 @@ impl Volume {
size: Size,
) -> Result<(), VolumeError> {
let _guard = self.data_file_access_control.read_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
let mut read_option = ReadOption::default();
self.read_needle_data_at_unlocked(n, offset, size, &mut read_option)
}
@@ -1540,6 +1566,9 @@ impl Volume {
/// Read raw needle blob at a specific offset.
pub fn read_needle_blob(&self, offset: i64, size: Size) -> Result<Vec<u8>, VolumeError> {
let _guard = self.data_file_access_control.read_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
self.read_needle_blob_unlocked(offset, size)
}
@@ -1571,6 +1600,9 @@ impl Volume {
size: Size,
) -> Result<(), VolumeError> {
let _guard = self.data_file_access_control.read_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
self.read_needle_meta_at_unlocked(n, offset, size)
}
@@ -1685,6 +1717,9 @@ impl Volume {
read_deleted: bool,
) -> Result<NeedleStreamInfo, VolumeError> {
let _guard = self.data_file_access_control.read_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
let nm = self.nm_or_not_found()?;
let nv = nm.get(n.id)?.ok_or(VolumeError::NotFound)?;
@@ -1829,6 +1864,9 @@ impl Volume {
fsync: bool,
) -> Result<(u64, Size, bool), VolumeError> {
let _guard = self.data_file_access_control.write_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
if self.is_read_only() {
return Err(VolumeError::ReadOnly);
}
@@ -1886,6 +1924,18 @@ impl Volume {
Ok(())
}
/// Take an unflushed append back off the .dat after its sync failed.
fn truncate_dat(&self, len: u64) -> io::Result<()> {
#[cfg(test)]
if self.fail_truncate_for_test {
return Err(io::Error::other("injected truncate failure"));
}
match self.dat_file.as_ref() {
Some(dat_file) => dat_file.set_len(len),
None => Ok(()),
}
}
fn do_write_request(
&mut self,
n: &mut Needle,
@@ -1951,22 +2001,14 @@ impl Volume {
// 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));
let truncated = match self.dat_file.as_ref() {
Some(dat_file) => dat_file.set_len(offset),
None => Ok(()),
};
if let Err(te) = truncated {
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 stop taking writes instead.
self.no_write_or_delete = true;
tracing::error!(
"volume {}: failed to truncate back to {} after a failed fsync, \
marking read only: {}",
self.id.0,
offset,
te
);
// 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));
}
@@ -2166,6 +2208,9 @@ impl Volume {
/// Delete a needle from the volume.
pub fn delete_needle(&mut self, n: &mut Needle) -> Result<Size, VolumeError> {
let _guard = self.data_file_access_control.write_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
if self.no_write_or_delete {
return Err(VolumeError::ReadOnly);
}
@@ -2248,9 +2293,74 @@ impl Volume {
pub fn is_read_only(&self) -> bool {
self.no_write_or_delete
|| self.no_write_can_delete
|| self.io_unavailable.is_some()
|| self.location_disk_space_low.load(Ordering::Relaxed)
}
/// The reason the volume refuses all I/O, when a failed recovery left the
/// .dat/index pair unverified. Mirrors Go's unavailableError.
pub fn unavailable_error(&self) -> Option<VolumeError> {
self.io_unavailable
.as_ref()
.map(|reason| VolumeError::Unavailable(reason.clone()))
}
/// Fail closed after a recovery could not return the volume to a verified
/// state: refuse all I/O, leave heartbeats out, and record the state so a
/// reload stays unavailable until an operator verifies the volume.
fn mark_io_unavailable(&mut self, reason: String) {
self.no_write_or_delete = true;
self.io_unavailable = Some(reason.clone());
self.mark_io_quarantined();
if let Err(e) = self.persist_unavailable(&reason) {
warn!(
volume_id = self.id.0,
error = %e,
"failed to persist unavailable marker"
);
}
if let Err(e) = self.set_read_only_persist(false, true) {
warn!(
volume_id = self.id.0,
error = %e,
"failed to persist unavailable state"
);
}
error!(
volume_id = self.id.0,
"volume entered unavailable state after failed recovery: {}", reason
);
}
fn persist_unavailable(&self, reason: &str) -> io::Result<()> {
let marker = self.file_name(".unavailable");
let mut f = OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.open(&marker)?;
f.write_all(reason.as_bytes())?;
f.sync_all()?;
drop(f);
fsync_dir(&marker)
}
/// Re-arm the state a `.unavailable` marker recorded. The marker is
/// deleted manually once the .dat/index pair is verified.
fn restore_unavailable(&mut self) {
let Ok(reason) = fs::read_to_string(self.file_name(".unavailable")) else {
return;
};
self.no_write_or_delete = true;
self.io_unavailable = Some(reason.trim().to_string());
self.mark_io_quarantined();
warn!(
volume_id = self.id.0,
"volume is unavailable: {}",
reason.trim()
);
}
pub fn is_no_write_or_delete(&self) -> bool {
self.no_write_or_delete
}
@@ -2327,6 +2437,9 @@ impl Volume {
/// Read all live needles from the volume (for ReadAllNeedles streaming RPC).
pub fn read_all_needles(&self) -> Result<Vec<Needle>, VolumeError> {
let _guard = self.data_file_access_control.read_lock();
if let Some(e) = self.unavailable_error() {
return Err(e);
}
let nm = self.nm_or_not_found()?;
let version = self.version();
let dat_size = self.current_dat_file_size()? as i64;
@@ -2410,7 +2523,7 @@ impl Volume {
let version = self.version();
// The deeper-than-tail structural check (every (offset + actual size)
// fits inside .dat — issue #8928) is now handled in load() via the
// fits inside .dat) is now handled in load() via the
// needle map's max_needle_end accumulator, so we don't pay for a
// second linear scan of the .idx here.
@@ -2894,6 +3007,9 @@ impl Volume {
/// dropping its store guard. See `DatScanPlan` for why the offset, the
/// handle and the end bound must come from the same guard.
pub(crate) fn dat_scan_plan(&self, from_offset: u64) -> Result<DatScanPlan, VolumeError> {
if let Some(e) = self.unavailable_error() {
return Err(e);
}
let source = if self.dat_file.is_some() {
// A fresh open, not `try_clone`: a duplicated handle shares the
// file position, and on Windows `read_exact_at` goes through
@@ -2991,6 +3107,9 @@ impl Volume {
/// surviving until the next restart, then vanishing. Re-attach a writer
/// here so writes persist again.
pub fn set_writable(&mut self) -> Result<(), VolumeError> {
if let Some(e) = self.unavailable_error() {
return Err(e);
}
let was_no_write_or_delete = self.no_write_or_delete;
let was_no_write_can_delete = self.no_write_can_delete;
self.no_write_or_delete = false;
@@ -3787,7 +3906,6 @@ impl Volume {
// surfaces as a generic Io rather than UnexpectedEof) must
// abort so an operator notices, rather than silently
// compacting away data that might come back on retry.
// See issue #8928.
if !is_skippable_needle_read_error(&e) {
return Err(VolumeError::Io(io::Error::other(format!(
"cannot hydrate needle from file: {}",
@@ -4420,6 +4538,11 @@ impl Volume {
self.fail_idx_sync_for_test = fail;
}
#[cfg(test)]
pub(crate) fn fail_next_truncate_for_test(&mut self, fail: bool) {
self.fail_truncate_for_test = fail;
}
#[cfg(test)]
pub(crate) fn set_last_modified_ts_for_test(&mut self, ts_seconds: u64) {
self.last_modified_ts_seconds = ts_seconds;
@@ -4541,7 +4664,16 @@ pub(crate) fn fsync_dir(path: &str) -> io::Result<()> {
pub(crate) fn remove_volume_files(base: &str, keep_vif: bool) {
for ext in &[
".dat", ".idx", ".vif", ".sdx", ".cpd", ".cpx", ".cpc", ".note", ".rdb",
".dat",
".idx",
".vif",
".sdx",
".cpd",
".cpx",
".cpc",
".note",
".rdb",
".unavailable",
] {
if *ext == ".vif" && keep_vif {
continue;
@@ -5553,6 +5685,72 @@ mod tests {
);
}
/// A failed fsync whose .dat rollback cannot complete leaves the tail
/// unverified, so the volume fails closed: reads and writes are refused,
/// the state is persisted, and a reload stays quarantined.
#[test]
fn test_failed_rollback_marks_volume_unavailable() {
let tmp = TempDir::new().unwrap();
let dir = tmp.path().to_str().unwrap();
let mut v = make_test_volume(dir);
v.fail_next_fsync_for_test(true);
v.fail_next_truncate_for_test(true);
let mut n = Needle {
id: NeedleId(1),
cookie: Cookie(0xaa),
data: b"never-landed".to_vec(),
data_size: 12,
..Needle::default()
};
v.write_needle(&mut n, true, true).unwrap_err();
v.fail_next_fsync_for_test(false);
v.fail_next_truncate_for_test(false);
assert!(v.is_read_only());
assert!(v.unavailable_error().is_some());
assert!(
v.should_quarantine(),
"an unavailable volume must stay out of heartbeats"
);
assert!(Path::new(&v.file_name(".unavailable")).exists());
let mut read_n = Needle {
id: NeedleId(1),
..Needle::default()
};
assert!(matches!(
v.read_needle(&mut read_n),
Err(VolumeError::Unavailable(_))
));
let mut later = Needle {
id: NeedleId(2),
cookie: Cookie(0xbb),
data: b"should-be-refused".to_vec(),
data_size: 17,
..Needle::default()
};
assert!(matches!(
v.write_needle(&mut later, true, true),
Err(VolumeError::Unavailable(_))
));
assert!(
v.set_writable().is_err(),
"a manual writable transition must not bypass quarantine"
);
assert!(v.unavailable_error().is_some());
let reloaded = reload_volume(dir);
assert!(reloaded.unavailable_error().is_some());
assert!(reloaded.is_read_only());
assert!(reloaded.should_quarantine());
let mut read_n = Needle {
id: NeedleId(1),
..Needle::default()
};
assert!(reloaded.read_needle(&mut read_n).is_err());
}
/// Same for a needle the failed write introduced: it never becomes visible,
/// rather than being published and then tombstoned back out.
#[test]
@@ -7058,7 +7256,7 @@ mod tests {
}
/// Vacuum compaction must tolerate an .idx entry whose offset points past
/// the end of the .dat file (the failure mode in issue #8928). The bad
/// the end of the .dat file. The bad
/// entry is silently dropped from the resulting .cpx; healthy needles
/// survive untouched.
#[test]
@@ -7125,7 +7323,7 @@ mod tests {
/// The needle map's max_needle_end accumulator must let volume.load
/// detect an .idx whose entries point past the end of the .dat — the
/// deeper-than-tail corruption shape from issue #8928 that the existing
/// deeper-than-tail corruption shape that the existing
/// last-10-entries scan cannot see. The check is populated by the load
/// walk and read in volume.load() to flip the volume read-only.
#[test]
+10 -4
View File
@@ -6,7 +6,7 @@ import (
"io"
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
"github.com/seaweedfs/seaweedfs/weed/storage"
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
)
@@ -16,6 +16,9 @@ func (vs *VolumeServer) VolumeIncrementalCopy(req *volume_server_pb.VolumeIncrem
if v == nil {
return fmt.Errorf("not found volume id %d", req.VolumeId)
}
if err := v.UnavailableError(); err != nil {
return err
}
stopOffset, _, _ := v.FileStat()
foundOffset, isLastOne, err := v.BinarySearchByAppendAtNs(req.SinceNs)
@@ -30,7 +33,7 @@ func (vs *VolumeServer) VolumeIncrementalCopy(req *volume_server_pb.VolumeIncrem
startOffset := foundOffset.ToActualOffset()
buf := make([]byte, 1024*1024*2)
return sendFileContent(v.DataBackend, buf, startOffset, int64(stopOffset), stream)
return sendFileContent(v, buf, startOffset, int64(stopOffset), stream)
}
@@ -47,10 +50,13 @@ func (vs *VolumeServer) VolumeSyncStatus(ctx context.Context, req *volume_server
}
func sendFileContent(datBackend backend.BackendStorageFile, buf []byte, startOffset, stopOffset int64, stream volume_server_pb.VolumeServer_VolumeIncrementalCopyServer) error {
func sendFileContent(v *storage.Volume, buf []byte, startOffset, stopOffset int64, stream volume_server_pb.VolumeServer_VolumeIncrementalCopyServer) error {
var blockSizeLimit = int64(len(buf))
for i := int64(0); i < stopOffset-startOffset; i += blockSizeLimit {
n, readErr := datBackend.ReadAt(buf, startOffset+i)
if err := v.UnavailableError(); err != nil {
return err
}
n, readErr := v.DataBackend.ReadAt(buf, startOffset+i)
if readErr == nil || readErr == io.EOF {
resp := &volume_server_pb.VolumeIncrementalCopyResponse{}
resp.FileContent = buf[:int64(n)]
+3
View File
@@ -29,6 +29,9 @@ func (vs *VolumeServer) streamReadOneVolume(vid needle.VolumeId, stream volume_s
if v == nil {
return fmt.Errorf("not found volume id %d", vid)
}
if err := v.UnavailableError(); err != nil {
return err
}
scanner := &storage.VolumeFileScanner4ReadAll{
Stream: stream,
+14
View File
@@ -20,6 +20,9 @@ func (vs *VolumeServer) VolumeTailSender(req *volume_server_pb.VolumeTailSenderR
if v == nil {
return fmt.Errorf("not found volume id %d", req.VolumeId)
}
if err := v.UnavailableError(); err != nil {
return err
}
defer glog.V(1).Infof("tailing volume %d finished", v.Id)
@@ -27,6 +30,9 @@ func (vs *VolumeServer) VolumeTailSender(req *volume_server_pb.VolumeTailSenderR
drainingSeconds := req.IdleTimeoutSeconds
for {
if err := v.UnavailableError(); err != nil {
return err
}
lastProcessedTimestampNs, err := sendNeedlesSince(stream, v, lastTimestampNs)
if err != nil {
glog.Infof("sendNeedlesSince: %v", err)
@@ -64,6 +70,9 @@ func sendNeedlesSince(stream volume_server_pb.VolumeServer_VolumeTailSenderServe
// log.Printf("reading ts %d offset %d isLast %v", lastTimestampNs, foundOffset, isLastOne)
if isLastOne {
if err := v.UnavailableError(); err != nil {
return 0, err
}
// need to heart beat to the client to ensure the connection health
sendErr := stream.Send(&volume_server_pb.VolumeTailSenderResponse{IsLastChunk: true, Version: uint32(v.Version())})
return lastTimestampNs, sendErr
@@ -72,6 +81,7 @@ func sendNeedlesSince(stream volume_server_pb.VolumeServer_VolumeTailSenderServe
scanner := &VolumeFileScanner4Tailing{
stream: stream,
version: uint32(v.Version()),
v: v,
}
err = storage.ScanVolumeFileFrom(v.Version(), v.DataBackend, foundOffset.ToActualOffset(), scanner)
@@ -112,6 +122,7 @@ type VolumeFileScanner4Tailing struct {
stream volume_server_pb.VolumeServer_VolumeTailSenderServer
lastProcessedTimestampNs uint64
version uint32
v *storage.Volume
}
func (scanner *VolumeFileScanner4Tailing) VisitSuperBlock(superBlock super_block.SuperBlock) error {
@@ -123,6 +134,9 @@ func (scanner *VolumeFileScanner4Tailing) ReadNeedleBody() bool {
}
func (scanner *VolumeFileScanner4Tailing) VisitNeedle(n *needle.Needle, offset int64, needleHeader, needleBody []byte) error {
if err := scanner.v.UnavailableError(); err != nil {
return err
}
isLastChunk := false
// need to send body by chunks
+28
View File
@@ -37,6 +37,23 @@ type NeedleMapper interface {
ReadIndexEntry(n int64) (key NeedleId, offset Offset, size Size, err error)
}
type batchIndexRollbacker interface {
truncateIndex(offset int64) error
}
// batchMapRollbacker restores the in-memory/durable needle mapping without
// touching the index file: a rolled-back batch rewrites the index wholesale
// via truncateIndex, so replay-correcting entries are not needed here.
type batchMapRollbacker interface {
removeMapping(key NeedleId) error
restoreMapping(key NeedleId, offset Offset, size Size) error
}
type batchMetricRollbacker interface {
snapshotBatchMetrics() batchMapMetricSnapshot
restoreBatchMetrics(snapshot batchMapMetricSnapshot)
}
type baseNeedleMapper struct {
mapMetric
@@ -75,6 +92,17 @@ func (nm *baseNeedleMapper) Sync() error {
return nm.indexFile.Sync()
}
func (nm *baseNeedleMapper) truncateIndex(offset int64) error {
nm.indexFileAccessLock.Lock()
defer nm.indexFileAccessLock.Unlock()
if err := nm.indexFile.Truncate(offset); err != nil {
return err
}
nm.indexFileOffset = offset
return nil
}
func (nm *baseNeedleMapper) ReadIndexEntry(n int64) (key NeedleId, offset Offset, size Size, err error) {
bytes := make([]byte, NeedleMapEntrySize)
var readCount int
+31
View File
@@ -190,6 +190,23 @@ func (cs *CompactMapSegment) delete(key types.NeedleId) types.Size {
return types.Size(0)
}
// remove erases a map entry entirely, returning whether it existed.
func (cs *CompactMapSegment) remove(key types.NeedleId) bool {
i, found := cs.bsearchKey(key)
if !found {
return false
}
copy(cs.list[i:], cs.list[i+1:])
cs.list = cs.list[:len(cs.list)-1]
if len(cs.list) == 0 {
cs.firstKey, cs.lastKey = MaxCompactKey, 0
} else {
cs.firstKey = cs.list[0].key
cs.lastKey = cs.list[len(cs.list)-1].key
}
return true
}
func NewCompactMap() *CompactMap {
return &CompactMap{
segments: map[Chunk]*CompactMapSegment{},
@@ -273,6 +290,20 @@ func (cm *CompactMap) Delete(key types.NeedleId) types.Size {
return cs.delete(key)
}
// Remove erases a map entry entirely, returning whether it existed. Unlike
// Delete it leaves no tombstoned entry behind.
func (cm *CompactMap) Remove(key types.NeedleId) bool {
cm.Lock()
defer cm.Unlock()
chunk := Chunk(key / SegmentChunkSize)
cs, ok := cm.segments[chunk]
if !ok {
return false
}
return cs.remove(key)
}
// AscendingVisit runs a function on all entries, in ascending key order. Returns any errors hit while visiting.
func (cm *CompactMap) AscendingVisit(visit func(NeedleValue) error) error {
cm.RLock()
@@ -7,6 +7,7 @@ import (
type NeedleValueMap interface {
Set(key NeedleId, offset Offset, size Size) (oldOffset Offset, oldSize Size)
Delete(key NeedleId) Size
Remove(key NeedleId) bool
Get(key NeedleId) (*NeedleValue, bool)
AscendingVisit(visit func(NeedleValue) error) error
}
+28
View File
@@ -189,6 +189,34 @@ func (m *LevelDbNeedleMap) Put(key NeedleId, offset Offset, size Size) error {
return levelDbWrite(m.db, key, offset, size, watermark != 0, watermark)
}
func (m *LevelDbNeedleMap) truncateIndex(offset int64) error {
if err := m.baseNeedleMapper.truncateIndex(offset); err != nil {
return err
}
m.recordCount = uint64(offset / NeedleMapEntrySize)
return nil
}
func (m *LevelDbNeedleMap) removeMapping(key NeedleId) error {
if m.ldbTimeout > 0 {
if err := m.ensureLdbLoaded(); err != nil {
return err
}
defer m.ldbAccessLock.RUnlock()
}
return levelDbDelete(m.db, key)
}
func (m *LevelDbNeedleMap) restoreMapping(key NeedleId, offset Offset, size Size) error {
if m.ldbTimeout > 0 {
if err := m.ensureLdbLoaded(); err != nil {
return err
}
defer m.ldbAccessLock.RUnlock()
}
return levelDbWrite(m.db, key, offset, size, false, 0)
}
func getWatermark(db *leveldb.DB) uint64 {
data, err := db.Get(watermarkKey, nil)
if err != nil || len(data) != 8 {
+8
View File
@@ -71,6 +71,14 @@ func (nm *NeedleMap) Delete(key NeedleId, offset Offset) error {
nm.logDelete(deletedBytes)
return nm.appendToIndexFile(key, offset, TombstoneFileSize)
}
func (nm *NeedleMap) removeMapping(key NeedleId) error {
nm.m.Remove(NeedleId(key))
return nil
}
func (nm *NeedleMap) restoreMapping(key NeedleId, offset Offset, size Size) error {
nm.m.Set(NeedleId(key), offset, size)
return nil
}
func (nm *NeedleMap) Close() {
if nm.indexFile == nil {
return
+29
View File
@@ -26,6 +26,35 @@ type mapMetric struct {
MaximumNeedleEnd int64 `json:"MaxNeedleEnd"`
}
type batchMapMetricSnapshot struct {
deletionCounter uint32
fileCounter uint32
deletionByteCounter uint64
fileByteCounter uint64
maximumFileKey uint64
maximumNeedleEnd int64
}
func (mm *mapMetric) snapshotBatchMetrics() batchMapMetricSnapshot {
return batchMapMetricSnapshot{
deletionCounter: atomic.LoadUint32(&mm.DeletionCounter),
fileCounter: atomic.LoadUint32(&mm.FileCounter),
deletionByteCounter: atomic.LoadUint64(&mm.DeletionByteCounter),
fileByteCounter: atomic.LoadUint64(&mm.FileByteCounter),
maximumFileKey: atomic.LoadUint64(&mm.MaximumFileKey),
maximumNeedleEnd: atomic.LoadInt64(&mm.MaximumNeedleEnd),
}
}
func (mm *mapMetric) restoreBatchMetrics(snapshot batchMapMetricSnapshot) {
atomic.StoreUint32(&mm.DeletionCounter, snapshot.deletionCounter)
atomic.StoreUint32(&mm.FileCounter, snapshot.fileCounter)
atomic.StoreUint64(&mm.DeletionByteCounter, snapshot.deletionByteCounter)
atomic.StoreUint64(&mm.FileByteCounter, snapshot.fileByteCounter)
atomic.StoreUint64(&mm.MaximumFileKey, snapshot.maximumFileKey)
atomic.StoreInt64(&mm.MaximumNeedleEnd, snapshot.maximumNeedleEnd)
}
func (mm *mapMetric) logDelete(deletedByteCount Size) {
if mm == nil {
return
+27
View File
@@ -101,6 +101,14 @@ func (m *SortedFileNeedleMap) Put(key NeedleId, offset Offset, size Size) error
return fmt.Errorf("needle map %s.sdx is read only: %w", m.baseFileName, os.ErrInvalid)
}
func (m *SortedFileNeedleMap) removeMapping(key NeedleId) error {
return fmt.Errorf("needle map %s.sdx is read only: %w", m.baseFileName, os.ErrInvalid)
}
func (m *SortedFileNeedleMap) restoreMapping(key NeedleId, offset Offset, size Size) error {
return fmt.Errorf("needle map %s.sdx is read only: %w", m.baseFileName, os.ErrInvalid)
}
func (m *SortedFileNeedleMap) Delete(key NeedleId, offset Offset) error {
f, err := pooledIndexFiles.borrow(m.dbFileName, true)
@@ -180,6 +188,25 @@ func (m *SortedFileNeedleMap) Sync() error {
return nil
}
// truncateIndex drops .idx tombstones a failed batch appended. The .sdx is
// untouched: a delete already marked there cannot be unmarked, so callers
// treat the batch as unrecoverable before reaching this.
func (m *SortedFileNeedleMap) truncateIndex(offset int64) error {
f, err := pooledIndexFiles.borrow(m.indexFileName, true)
if err != nil {
return err
}
defer pooledIndexFiles.release(f)
m.indexFileAccessLock.Lock()
defer m.indexFileAccessLock.Unlock()
if err := f.file.Truncate(offset); err != nil {
return err
}
m.indexFileOffset = offset
return nil
}
func (m *SortedFileNeedleMap) ReadIndexEntry(n int64) (key NeedleId, offset Offset, size Size, err error) {
var f *pooledFile
if f, err = pooledIndexFiles.borrow(m.indexFileName, false); err != nil {
+7
View File
@@ -935,6 +935,9 @@ func (s *Store) MarkVolumeWritable(i needle.VolumeId) error {
if v == nil {
return fmt.Errorf("volume %d not found", i)
}
if err := v.UnavailableError(); err != nil {
return fmt.Errorf("volume %d cannot be marked writable: %w", i, err)
}
// If the volume booted with .vif ReadOnly=true, .idx is opened O_RDONLY
// and v.nm is a SortedFileNeedleMap that rejects Put. Swap to writable
// form before flipping the flag so the next write doesn't race past a
@@ -943,6 +946,10 @@ func (s *Store) MarkVolumeWritable(i needle.VolumeId) error {
return fmt.Errorf("volume %d reopen idx for write: %v", i, err)
}
v.noWriteLock.Lock()
if v.ioUnavailable {
v.noWriteLock.Unlock()
return fmt.Errorf("volume %d cannot be marked writable: %w", i, v.UnavailableError())
}
prevNoWriteOrDelete := v.noWriteOrDelete
prevNoWriteCanDelete := v.noWriteCanDelete
v.noWriteOrDelete = false
+81 -1
View File
@@ -6,6 +6,7 @@ import (
"os"
"path"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
@@ -18,6 +19,7 @@ import (
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
"github.com/seaweedfs/seaweedfs/weed/storage/types"
"github.com/seaweedfs/seaweedfs/weed/storage/volume_info"
"github.com/seaweedfs/seaweedfs/weed/util"
"github.com/seaweedfs/seaweedfs/weed/glog"
)
@@ -33,6 +35,8 @@ type Volume struct {
needleMapKind NeedleMapKind
noWriteOrDelete bool // if readonly, either noWriteOrDelete or noWriteCanDelete
noWriteCanDelete bool // if readonly, either noWriteOrDelete or noWriteCanDelete
ioUnavailable bool
ioUnavailableError string
noWriteLock sync.RWMutex
hasRemoteFile atomic.Bool // if the volume is tiered: data lives in a remote backend
MemoryMapMaxSizeMb uint32
@@ -444,7 +448,7 @@ func (v *Volume) ToVolumeInformationMessage(into *master_pb.VolumeInformationMes
// disk-path operation can ever succeed. Skip remote-tiered volumes, whose .dat
// legitimately lives in cloud storage. Only a present .dat is cached for 30s; a
// missing one is re-checked every heartbeat so the volume stays suppressed until
// the file returns. See github.com/seaweedfs/seaweedfs/issues/10004
// the file returns.
if fileCount > 0 && !v.HasRemoteFile() {
const diskCheckIntervalNs = 30 * int64(time.Second)
now := time.Now().UnixNano()
@@ -500,6 +504,9 @@ func (v *Volume) IsReadOnly() bool {
func (v *Volume) ReadOnlyReasons() (readOnly, noWriteOrDelete, noWriteCanDelete, diskSpaceLow bool) {
v.noWriteLock.RLock()
noWriteOrDelete, noWriteCanDelete = v.noWriteOrDelete, v.noWriteCanDelete
if v.ioUnavailable {
noWriteOrDelete = true
}
v.noWriteLock.RUnlock()
// The location is attached when the volume joins a disk location, which is
// after NewVolume hands it back.
@@ -507,6 +514,79 @@ func (v *Volume) ReadOnlyReasons() (readOnly, noWriteOrDelete, noWriteCanDelete,
return noWriteOrDelete || noWriteCanDelete || diskSpaceLow, noWriteOrDelete, noWriteCanDelete, diskSpaceLow
}
var errVolumeUnavailable = errors.New("volume unavailable")
func (v *Volume) UnavailableError() error {
v.noWriteLock.RLock()
unavailable := v.ioUnavailable
reason := v.ioUnavailableError
v.noWriteLock.RUnlock()
if !unavailable {
return nil
}
if reason == "" {
return fmt.Errorf("volume %d is unavailable: %w", v.Id, errVolumeUnavailable)
}
return fmt.Errorf("volume %d is unavailable: %s: %w", v.Id, reason, errVolumeUnavailable)
}
func (v *Volume) markIoUnavailable(err error) {
v.noWriteLock.Lock()
v.noWriteOrDelete = true
v.ioUnavailable = true
v.ioUnavailableError = err.Error()
v.noWriteLock.Unlock()
v.markIoQuarantined()
if persistErr := v.persistUnavailable(err.Error()); persistErr != nil {
glog.Warningf("volume %d: failed to persist unavailable marker: %v", v.Id, persistErr)
}
if v.volumeInfo != nil {
if persistErr := v.PersistReadOnly(true, false); persistErr != nil {
glog.Warningf("volume %d: failed to persist unavailable state: %v", v.Id, persistErr)
}
}
glog.Errorf("volume %d entered unavailable state after failed recovery: %v", v.Id, err)
}
// persistUnavailable records the failed-recovery state so a reload keeps the
// volume unavailable instead of serving an unverified .dat/index pair.
func (v *Volume) persistUnavailable(reason string) error {
marker := v.FileName(".unavailable")
f, err := os.OpenFile(marker, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0644)
if err != nil {
return err
}
if _, err := f.WriteString(reason); err != nil {
f.Close()
return err
}
if err := f.Sync(); err != nil {
f.Close()
return err
}
if err := f.Close(); err != nil {
return err
}
return util.FsyncDir(v.dir)
}
// restoreUnavailable re-arms the in-memory state the .unavailable marker
// recorded. The marker is deleted manually once the volume pair is verified.
func (v *Volume) restoreUnavailable() {
reason, err := os.ReadFile(v.FileName(".unavailable"))
if err != nil {
return
}
v.noWriteLock.Lock()
v.noWriteOrDelete = true
v.ioUnavailable = true
v.ioUnavailableError = strings.TrimSpace(string(reason))
v.noWriteLock.Unlock()
v.markIoQuarantined()
glog.Warningf("volume %d is unavailable: %s", v.Id, v.ioUnavailableError)
}
func (v *Volume) PersistReadOnly(readOnly bool, canDelete bool) error {
v.volumeInfoRWLock.Lock()
defer v.volumeInfoRWLock.Unlock()
+7
View File
@@ -68,6 +68,10 @@ update needle map when receiving new .dat bytes. But seems not necessary now.)
func (v *Volume) IncrementalBackup(volumeServer pb.ServerAddress, grpcDialOption grpc.DialOption) error {
if err := v.UnavailableError(); err != nil {
return err
}
startFromOffset, _, _ := v.FileStat()
appendAtNs, err := v.findLastAppendAtNs()
if err != nil {
@@ -96,6 +100,9 @@ func (v *Volume) IncrementalBackup(volumeServer pb.ServerAddress, grpcDialOption
}
}
if err := v.UnavailableError(); err != nil {
return err
}
n, writeErr := v.DataBackend.WriteAt(resp.FileContent, writeOffset)
if writeErr != nil {
return writeErr
+5 -3
View File
@@ -297,8 +297,8 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
}
// The post-load structural check below uses the in-memory needle map
// to verify that no .idx entry references bytes past the end of .dat
// (issue #8928). The check piggybacks on MaxNeedleEnd, which the load
// to verify that no .idx entry references bytes past the end of .dat.
// The check piggybacks on MaxNeedleEnd, which the load
// walks below populate without a second linear scan.
// Loaders can return a typed-nil pointer with err set; assigning that
@@ -386,7 +386,7 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
// MaximumNeedleEnd, so this is just a numeric comparison — no extra
// disk I/O. A violation marks the volume read-only so a corrupt
// .idx left over from a crashed batched write does not silently
// power vacuum to drop reachable data. See issue #8928. err == nil
// power vacuum to drop reachable data. err == nil
// guards against a partial-walk MaximumNeedleEnd.
if err == nil && !v.HasRemoteFile() && v.nm != nil && v.DataBackend != nil {
if datSize, _, statErr := v.DataBackend.GetStat(); statErr == nil && datSize > 0 {
@@ -399,6 +399,8 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
}
}
v.restoreUnavailable()
if !hasVolumeInfoFile {
v.volumeInfo.Version = uint32(v.SuperBlock.Version)
v.volumeInfo.BytesOffset = uint32(types.OffsetSize)
+23
View File
@@ -22,6 +22,10 @@ func (v *Volume) readNeedle(n *needle.Needle, readOption *ReadOption, onReadSize
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
if err := v.UnavailableError(); err != nil {
return 0, err
}
if v.nm == nil {
glog.V(0).Infof("volume %d: needle map not loaded; read returns not-found", v.Id)
return -1, ErrorNotFound
@@ -89,6 +93,11 @@ func (v *Volume) readNeedle(n *needle.Needle, readOption *ReadOption, onReadSize
func (v *Volume) readNeedleMetaAt(n *needle.Needle, offset int64, size int32) (err error) {
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
if err := v.UnavailableError(); err != nil {
return err
}
// read deleted needle meta data
if size < 0 {
size = 0
@@ -114,6 +123,12 @@ func (v *Volume) readNeedleDataInto(n *needle.Needle, readOption *ReadOption, wr
if readOption.HasSlowRead {
v.dataFileAccessLock.RLock()
}
if err := v.UnavailableError(); err != nil {
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
}
return err
}
if v.nm == nil {
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
@@ -156,6 +171,10 @@ func (v *Volume) readNeedleDataInto(n *needle.Needle, readOption *ReadOption, wr
if readOption.HasSlowRead {
v.dataFileAccessLock.RLock()
if err := v.UnavailableError(); err != nil {
v.dataFileAccessLock.RUnlock()
return err
}
}
// possibly re-read needle offset if volume is compacted
if readOption.VolumeRevision != v.SuperBlock.CompactionRevision {
@@ -243,6 +262,10 @@ func (v *Volume) ReadNeedleBlob(offset int64, size Size) ([]byte, error) {
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
if err := v.UnavailableError(); err != nil {
return nil, err
}
blob, err := needle.ReadNeedleBlob(v.DataBackend, offset, size, v.Version())
v.checkReadWriteError(err)
return blob, err
+3
View File
@@ -21,6 +21,9 @@ func (scanner *VolumeFileScanner4ReadAll) ReadNeedleBody() bool {
func (scanner *VolumeFileScanner4ReadAll) VisitNeedle(n *needle.Needle, offset int64, needleHeader, needleBody []byte) error {
if err := scanner.V.UnavailableError(); err != nil {
return err
}
nv, ok := scanner.V.nm.Get(n.Id)
if !ok {
return nil
+208 -55
View File
@@ -16,6 +16,101 @@ var ErrorNotFound = errors.New("not found")
var ErrorDeleted = errors.New("already deleted")
var ErrorSizeMismatch = errors.New("size mismatch")
type batchNeedleSnapshot struct {
id NeedleId
found bool
offset Offset
size Size
changed bool
}
func newBatchNeedleSnapshot(v *Volume, id NeedleId) *batchNeedleSnapshot {
snapshot := &batchNeedleSnapshot{id: id}
if element, found := v.nm.Get(id); found && element != nil {
snapshot.found = true
snapshot.offset = element.Offset
snapshot.size = element.Size
}
return snapshot
}
func (s *batchNeedleSnapshot) observeCurrent(v *Volume) {
element, found := v.nm.Get(s.id)
if found != s.found {
s.changed = true
return
}
if found && (element == nil || element.Offset != s.offset || element.Size != s.size) {
s.changed = true
}
}
func (v *Volume) restoreBatchNeedle(snapshot *batchNeedleSnapshot) error {
mapRollbacker, ok := v.nm.(batchMapRollbacker)
if !ok {
return fmt.Errorf("mapper %T cannot restore its mappings", v.nm)
}
if snapshot.found {
return mapRollbacker.restoreMapping(snapshot.id, snapshot.offset, snapshot.size)
}
// A plain Delete would leave a tombstoned entry whose stale offset makes
// the next write to this needle fail reading a header that no longer exists.
return mapRollbacker.removeMapping(snapshot.id)
}
func (v *Volume) rollbackBatch(end int64, indexEnd int64, snapshots []*batchNeedleSnapshot,
metricRollbacker batchMetricRollbacker, metrics batchMapMetricSnapshot) error {
var recoveryErrors []error
for _, snapshot := range snapshots {
if !snapshot.changed {
continue
}
if err := v.restoreBatchNeedle(snapshot); err != nil {
recoveryErrors = append(recoveryErrors,
fmt.Errorf("restore needle %d: %w", snapshot.id, err))
}
}
if len(recoveryErrors) == 0 {
rollbacker, ok := v.nm.(batchIndexRollbacker)
if !ok {
recoveryErrors = append(recoveryErrors,
fmt.Errorf("mapper %T cannot truncate its index", v.nm))
} else if err := rollbacker.truncateIndex(indexEnd); err != nil {
recoveryErrors = append(recoveryErrors,
fmt.Errorf("truncate index to %d: %w", indexEnd, err))
}
}
if len(recoveryErrors) == 0 {
if err := v.nm.Sync(); err != nil {
recoveryErrors = append(recoveryErrors,
fmt.Errorf("sync recovered index: %w", err))
}
}
if len(recoveryErrors) == 0 {
if err := v.DataBackend.Truncate(end); err != nil {
recoveryErrors = append(recoveryErrors,
fmt.Errorf("truncate %s to %d: %w", v.DataBackend.Name(), end, err))
} else if err := v.DataBackend.Sync(); err != nil {
recoveryErrors = append(recoveryErrors,
fmt.Errorf("sync truncated %s: %w", v.DataBackend.Name(), err))
}
}
if len(recoveryErrors) == 0 {
if metricRollbacker == nil {
recoveryErrors = append(recoveryErrors,
fmt.Errorf("mapper %T cannot restore its metrics", v.nm))
} else {
metricRollbacker.restoreBatchMetrics(metrics)
}
}
return errors.Join(recoveryErrors...)
}
// isFileUnchanged checks whether this needle to write is same as last one.
// It requires serialized access in the same volume.
func (v *Volume) isFileUnchanged(n *needle.Needle) bool {
@@ -128,6 +223,8 @@ func removeVolumeFiles(filename string, keepVif bool) {
deleteAndLog("rdb")
// marker for damaged or incomplete volume
deleteAndLog("note")
// marker for a volume whose failed-batch recovery could not be verified
deleteAndLog("unavailable")
}
// asyncRequestAppend queues a request for the batch worker, starting it on the
@@ -147,6 +244,10 @@ func (v *Volume) syncWrite(n *needle.Needle, checkCookie bool, fsync bool) (offs
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
if err := v.UnavailableError(); err != nil {
return 0, 0, false, err
}
// A caller can still hold the volume after it was closed or destroyed, which
// leaves both of these nil. Refuse the write rather than dereference them.
if v.nm == nil || v.DataBackend == nil {
@@ -184,22 +285,35 @@ func (v *Volume) syncWrite(n *needle.Needle, checkCookie bool, fsync bool) (offs
// data we can vouch for, so they come back off the .dat and the needle map goes
// back to what it pointed at before, rather than at an offset past the new end.
func (v *Volume) rollbackUnflushedWrite(n *needle.Needle, offset uint64, end int64, priorOffset Offset, priorSize Size, hasPrior bool) {
var recoveryErr error
if te := v.DataBackend.Truncate(end); te != nil {
glog.V(0).Infof("Failed to truncate %s back to %d with error: %v", v.DataBackend.Name(), end, te)
recoveryErr = fmt.Errorf("truncate %s back to %d: %w", v.DataBackend.Name(), end, te)
}
current, found := v.nm.Get(n.Id)
if !found || current.Offset.ToActualOffset() != int64(offset) {
// doWriteRequest kept an existing mapping at a higher offset
return
if found && current.Offset.ToActualOffset() == int64(offset) {
var err error
if hasPrior {
err = v.nm.Put(n.Id, priorOffset, priorSize)
} else {
// The tombstone must reach .idx so a replay forgets the needle, but
// the negated entry it leaves in memory points at truncated bytes
// and would fail the next write, so erase the mapping as well.
err = v.nm.Delete(n.Id, ToOffset(int64(offset)))
if err == nil {
if rb, ok := v.nm.(batchMapRollbacker); ok {
err = rb.removeMapping(n.Id)
} else {
err = fmt.Errorf("mapper %T cannot remove mapping", v.nm)
}
}
}
if err != nil {
recoveryErr = errors.Join(recoveryErr,
fmt.Errorf("roll back the index of needle %d in volume %d: %w", n.Id, v.Id, err))
}
}
var err error
if hasPrior {
err = v.nm.Put(n.Id, priorOffset, priorSize)
} else {
err = v.nm.Delete(n.Id, ToOffset(int64(offset)))
}
if err != nil {
glog.V(0).Infof("Failed to roll back the index of needle %d in volume %d: %v", n.Id, v.Id, err)
if recoveryErr != nil {
v.markIoUnavailable(recoveryErr)
}
}
@@ -209,6 +323,10 @@ func (v *Volume) rollbackUnflushedWrite(n *needle.Needle, offset uint64, end int
// paths only return once the .dat is on disk.
func (v *Volume) writeNeedle2(n *needle.Needle, checkCookie bool, fsync bool, isStopping bool) (offset uint64, size Size, isUnchanged bool, err error) {
// glog.V(4).Infof("writing needle %s", needle.NewFileIdFromNeedle(v.Id, n).String())
if err := v.UnavailableError(); err != nil {
return 0, 0, false, err
}
if n.Ttl == needle.EMPTY_TTL && v.Ttl != needle.EMPTY_TTL {
n.SetHasTtl()
n.Ttl = v.Ttl
@@ -288,6 +406,10 @@ func (v *Volume) syncDelete(n *needle.Needle) (Size, error) {
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
if err := v.UnavailableError(); err != nil {
return 0, err
}
if v.nm == nil {
return 0, nil
}
@@ -340,6 +462,75 @@ func (v *Volume) doDeleteRequest(n *needle.Needle) (Size, error) {
return 0, nil
}
func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) {
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
end, e := int64(0), error(nil)
if unavailableErr := v.UnavailableError(); unavailableErr != nil {
e = unavailableErr
} else if v.nm == nil || v.DataBackend == nil {
e = fmt.Errorf("volume %d is closed", v.Id)
} else {
end, _, e = v.DataBackend.GetStat()
}
if e != nil {
for i := 0; i < len(currentRequests); i++ {
currentRequests[i].Complete(0, 0, false,
fmt.Errorf("cannot read current volume position: %v", e))
}
return
}
batchSnapshots := make(map[NeedleId]*batchNeedleSnapshot, len(currentRequests))
orderedSnapshots := make([]*batchNeedleSnapshot, 0, len(currentRequests))
metricRollbacker, _ := v.nm.(batchMetricRollbacker)
var batchMetrics batchMapMetricSnapshot
if metricRollbacker != nil {
batchMetrics = metricRollbacker.snapshotBatchMetrics()
}
indexEnd := int64(v.nm.IndexFileSize())
batchLastAppendAtNs := v.lastAppendAtNs
batchLastModifiedTsSeconds := v.lastModifiedTsSeconds
for i := 0; i < len(currentRequests); i++ {
needleID := currentRequests[i].N.Id
snapshot, found := batchSnapshots[needleID]
if !found {
snapshot = newBatchNeedleSnapshot(v, needleID)
batchSnapshots[needleID] = snapshot
orderedSnapshots = append(orderedSnapshots, snapshot)
}
if currentRequests[i].IsWriteRequest {
offset, size, isUnchanged, err := v.doWriteRequest(currentRequests[i].N, true)
currentRequests[i].UpdateResult(offset, uint64(size), isUnchanged, err)
} else {
size, err := v.doDeleteRequest(currentRequests[i].N)
currentRequests[i].UpdateResult(0, uint64(size), false, err)
}
snapshot.observeCurrent(v)
}
// if sync error the batch is not durable; restore it before another
// operation observes the volume
if syncErr := v.DataBackend.Sync(); syncErr != nil {
v.checkReadWriteError(syncErr)
v.lastAppendAtNs = batchLastAppendAtNs
v.lastModifiedTsSeconds = batchLastModifiedTsSeconds
batchErr := syncErr
if recoveryErr := v.rollbackBatch(end, indexEnd, orderedSnapshots, metricRollbacker, batchMetrics); recoveryErr != nil {
batchErr = errors.Join(syncErr, recoveryErr)
v.markIoUnavailable(batchErr)
}
for i := 0; i < len(currentRequests); i++ {
currentRequests[i].UpdateResult(0, 0, false, batchErr)
}
}
for i := 0; i < len(currentRequests); i++ {
currentRequests[i].Submit()
}
}
// startWorker returns the volume's batch-write channel, creating it and its
// goroutine on first use, and nil once stopWorker has run.
func (v *Volume) startWorker() chan *needle.AsyncRequest {
@@ -385,49 +576,7 @@ func (v *Volume) startWorker() chan *needle.AsyncRequest {
if len(currentRequests) == 0 {
continue
}
v.dataFileAccessLock.Lock()
end, e := int64(0), error(nil)
if v.nm == nil || v.DataBackend == nil {
e = fmt.Errorf("volume %d is closed", v.Id)
} else {
end, _, e = v.DataBackend.GetStat()
}
if e != nil {
for i := 0; i < len(currentRequests); i++ {
currentRequests[i].Complete(0, 0, false,
fmt.Errorf("cannot read current volume position: %v", e))
}
v.dataFileAccessLock.Unlock()
continue
}
for i := 0; i < len(currentRequests); i++ {
if currentRequests[i].IsWriteRequest {
offset, size, isUnchanged, err := v.doWriteRequest(currentRequests[i].N, true)
currentRequests[i].UpdateResult(offset, uint64(size), isUnchanged, err)
} else {
size, err := v.doDeleteRequest(currentRequests[i].N)
currentRequests[i].UpdateResult(0, uint64(size), false, err)
}
}
// if sync error, data is not reliable, we should mark the completed request as fail and rollback
if err := v.DataBackend.Sync(); err != nil {
// todo: this may generate dirty data or cause data inconsistent, may be weed need to panic?
if te := v.DataBackend.Truncate(end); te != nil {
glog.V(0).Infof("Failed to truncate %s back to %d with error: %v", v.DataBackend.Name(), end, te)
}
for i := 0; i < len(currentRequests); i++ {
if currentRequests[i].IsSucceed() {
currentRequests[i].UpdateResult(0, 0, false, err)
}
}
}
for i := 0; i < len(currentRequests); i++ {
currentRequests[i].Submit()
}
v.dataFileAccessLock.Unlock()
v.processBatch(currentRequests)
}
}()
return requests
@@ -453,6 +602,10 @@ func (v *Volume) WriteNeedleBlob(needleId NeedleId, needleBlob []byte, size Size
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
if err := v.UnavailableError(); err != nil {
return err
}
// nm.Put on a read-only volume fails only after the blob is appended to .dat.
if v.IsReadOnly() {
return fmt.Errorf("volume %d is read only", v.Id)
+406 -4
View File
@@ -14,22 +14,42 @@ import (
// countingBackend counts Sync calls and can be made to fail them.
type countingBackend struct {
backend.BackendStorageFile
syncCount int
syncErr error
syncCount int
syncErr error
syncErrOnce bool
truncateErr error
truncateCount int
}
func (b *countingBackend) Sync() error {
b.syncCount++
if b.syncErr != nil {
return b.syncErr
err := b.syncErr
if b.syncErrOnce {
b.syncErr = nil
b.syncErrOnce = false
}
return err
}
return b.BackendStorageFile.Sync()
}
func (b *countingBackend) Truncate(off int64) error {
b.truncateCount++
if b.truncateErr != nil {
return b.truncateErr
}
return b.BackendStorageFile.Truncate(off)
}
func newCountingVolume(t *testing.T) (*Volume, *countingBackend) {
return newCountingVolumeWithKind(t, NeedleMapInMemory)
}
func newCountingVolumeWithKind(t *testing.T, needleMapKind NeedleMapKind) (*Volume, *countingBackend) {
t.Helper()
dir := t.TempDir()
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0)
v, err := NewVolume(dir, dir, "", 1, needleMapKind, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0)
require.NoError(t, err)
t.Cleanup(v.Close)
counting := &countingBackend{BackendStorageFile: v.DataBackend}
@@ -37,6 +57,17 @@ func newCountingVolume(t *testing.T) (*Volume, *countingBackend) {
return v, counting
}
func reopenCountingVolume(t *testing.T, v *Volume) *Volume {
t.Helper()
dir, dirIdx, id, needleMapKind := v.dir, v.dirIdx, v.Id, v.needleMapKind
v.Close()
reloaded, err := NewVolume(dir, dirIdx, "", id, needleMapKind,
&super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0)
require.NoError(t, err)
t.Cleanup(reloaded.Close)
return reloaded
}
// A durable write reaching a stopping server used to silently drop its fsync,
// so the ack promised durability the .dat did not have. It now flushes inline
// instead of queueing on the batch worker that is winding down.
@@ -118,6 +149,377 @@ func TestWriteNeedle2DropsIndexOfUnflushedNewNeedle(t *testing.T) {
require.Error(t, err, "reading the rolled-back needle should fail cleanly, not read past the end")
}
// The async durable path must roll back the mapper as well as the data tail
// when its batch fsync fails. The current worker only truncates the data file,
// so this test exposes the dangling mapping.
func TestWriteNeedle2RollsBackBatchFsyncFailure(t *testing.T) {
v, counting := newCountingVolume(t)
initialFileCount := v.nm.FileCount()
initialDeletedCount := v.nm.DeletedCount()
initialContentSize := v.nm.ContentSize()
initialDeletedSize := v.nm.DeletedSize()
initialMaxFileKey := v.nm.MaxFileKey()
before, _, err := v.DataBackend.GetStat()
require.NoError(t, err)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
fresh := fixedNeedle(8, "batch-never-landed")
_, _, _, err = v.writeNeedle2(fresh, true, true, false)
require.Error(t, err, "a batch whose fsync failed must not be acknowledged")
after, _, err := v.DataBackend.GetStat()
require.NoError(t, err)
require.Equal(t, before, after, "the failed batch append should be truncated")
if entry, found := v.nm.Get(fresh.Id); found {
require.True(t, entry.Size.IsDeleted(), "the failed batch must not leave a live mapping")
}
require.Equal(t, initialFileCount, v.nm.FileCount())
require.Equal(t, initialDeletedCount, v.nm.DeletedCount())
require.Equal(t, initialContentSize, v.nm.ContentSize())
require.Equal(t, initialDeletedSize, v.nm.DeletedSize())
require.Equal(t, initialMaxFileKey, v.nm.MaxFileKey())
}
func TestWriteNeedle2RollsBackBatchFsyncFailureAndReloadsNewNeedle(t *testing.T) {
v, counting := newCountingVolume(t)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
fresh := fixedNeedle(9, "batch-never-landed")
_, _, _, err := v.writeNeedle2(fresh, true, true, false)
require.Error(t, err)
reloaded := reopenCountingVolume(t, v)
if entry, found := reloaded.nm.Get(fresh.Id); found {
require.True(t, entry.Size.IsDeleted(), "the failed batch must not reload a live mapping")
}
readBack := new(needle.Needle)
readBack.Id = fresh.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.Error(t, err, "the failed batch must not be readable after reload")
require.False(t, reloaded.IsReadOnly(), "a completed rollback should keep the volume healthy")
}
func TestWriteNeedle2RollsBackBatchFsyncFailureAndReloadsOverwrite(t *testing.T) {
v, counting := newCountingVolume(t)
kept := fixedNeedle(10, "first-copy")
_, _, _, err := v.writeNeedle2(kept, true, true, true)
require.NoError(t, err)
keptEntry, found := v.nm.Get(kept.Id)
require.True(t, found)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
_, _, _, err = v.writeNeedle2(fixedNeedle(10, "second-copy"), true, true, false)
require.Error(t, err)
now, found := v.nm.Get(kept.Id)
require.True(t, found)
require.Equal(t, keptEntry.Offset, now.Offset)
require.Equal(t, keptEntry.Size, now.Size)
reloaded := reopenCountingVolume(t, v)
readBack := new(needle.Needle)
readBack.Id = kept.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("first-copy"), readBack.Data)
}
func TestBatchDeleteRollsBackOnFsyncFailure(t *testing.T) {
v, counting := newCountingVolume(t)
kept := fixedNeedle(11, "delete-me")
_, _, _, err := v.writeNeedle2(kept, true, true, true)
require.NoError(t, err)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
deleteRequest := needle.NewAsyncRequest(&needle.Needle{Id: kept.Id}, false)
deleteRequest.ActualSize = needle.GetActualSize(0, v.Version())
require.True(t, v.asyncRequestAppend(deleteRequest))
_, _, _, err = deleteRequest.WaitComplete()
require.Error(t, err)
readBack := new(needle.Needle)
readBack.Id = kept.Id
_, err = v.readNeedle(readBack, nil, nil)
require.NoError(t, err, "a failed batch delete must restore the live mapping")
require.Equal(t, []byte("delete-me"), readBack.Data)
reloaded := reopenCountingVolume(t, v)
readBack = new(needle.Needle)
readBack.Id = kept.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("delete-me"), readBack.Data)
}
func TestBatchSyncFailureEntersFailClosedWhenTruncateFails(t *testing.T) {
v, counting := newCountingVolume(t)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
counting.truncateErr = errors.New("truncate failed")
fresh := fixedNeedle(12, "cannot-be-recovered")
_, _, _, err := v.writeNeedle2(fresh, true, true, false)
require.Error(t, err)
require.True(t, v.IsReadOnly(), "a failed recovery must make the volume unavailable")
require.NotNil(t, v.UnavailableError())
readBack := new(needle.Needle)
readBack.Id = fresh.Id
_, err = v.readNeedle(readBack, nil, nil)
require.Error(t, err, "an unavailable volume must reject reads")
_, _, _, err = v.writeNeedle2(fixedNeedle(13, "after-failure"), true, true, false)
require.Error(t, err, "an unavailable volume must reject later writes")
require.Equal(t, 1, counting.truncateCount)
}
func TestUnavailableVolumeCannotBeMarkedWritable(t *testing.T) {
dir := t.TempDir()
store := newIdxSplitStore(t, dir, dir)
const vid = needle.VolumeId(19)
require.NoError(t, store.AddVolume(vid, "", NeedleMapInMemory, "000", "", 0,
needle.GetCurrentVersion(), 0, types.HardDriveType, 0))
v := store.findVolume(vid)
require.NotNil(t, v)
v.markIoUnavailable(errors.New("batch recovery failed"))
require.Error(t, store.MarkVolumeWritable(vid), "manual writable transition must not bypass quarantine")
require.NotNil(t, v.UnavailableError())
}
func TestMixedBatchSyncFailureRollsBackAsOneUnit(t *testing.T) {
v, counting := newCountingVolume(t)
overwritten := fixedNeedle(14, "keep-overwrite")
deleted := fixedNeedle(15, "keep-delete")
_, _, _, err := v.writeNeedle2(overwritten, true, true, true)
require.NoError(t, err)
_, _, _, err = v.writeNeedle2(deleted, true, true, true)
require.NoError(t, err)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
replacement := fixedNeedle(14, "failed-overwrite")
fresh := fixedNeedle(16, "failed-new")
deleteRequest := needle.NewAsyncRequest(&needle.Needle{Id: deleted.Id}, false)
requests := []*needle.AsyncRequest{
needle.NewAsyncRequest(replacement, true),
deleteRequest,
needle.NewAsyncRequest(fresh, true),
}
v.processBatch(requests)
for _, request := range requests {
_, _, _, err = request.WaitComplete()
require.Error(t, err, "every request in a failed batch must fail")
}
readBack := new(needle.Needle)
readBack.Id = overwritten.Id
_, err = v.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("keep-overwrite"), readBack.Data)
readBack = new(needle.Needle)
readBack.Id = deleted.Id
_, err = v.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("keep-delete"), readBack.Data)
readBack = new(needle.Needle)
readBack.Id = fresh.Id
_, err = v.readNeedle(readBack, nil, nil)
require.Error(t, err)
reloaded := reopenCountingVolume(t, v)
readBack = new(needle.Needle)
readBack.Id = overwritten.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("keep-overwrite"), readBack.Data)
readBack = new(needle.Needle)
readBack.Id = deleted.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("keep-delete"), readBack.Data)
readBack = new(needle.Needle)
readBack.Id = fresh.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.Error(t, err)
}
func TestWriteNeedle2EntersFailClosedWhenInlineRollbackFails(t *testing.T) {
v, counting := newCountingVolume(t)
counting.syncErr = errors.New("inline fsync failed")
counting.syncErrOnce = true
counting.truncateErr = errors.New("truncate failed")
_, _, _, err := v.writeNeedle2(fixedNeedle(20, "cannot-be-recovered"), true, true, true)
require.Error(t, err)
require.NotNil(t, v.UnavailableError())
readBack := new(needle.Needle)
readBack.Id = 20
_, err = v.readNeedle(readBack, nil, nil)
require.Error(t, err)
}
func TestUnavailableVolumeStaysUnavailableAfterReload(t *testing.T) {
v, _ := newCountingVolume(t)
v.markIoUnavailable(errors.New("batch recovery failed"))
reloaded := reopenCountingVolume(t, v)
require.NotNil(t, reloaded.UnavailableError(), "the unavailable state must survive a reload")
require.True(t, reloaded.IsReadOnly())
_, _, quarantined := reloaded.getIoErrorState()
require.True(t, quarantined, "a reloaded unavailable volume must stay out of heartbeats")
readBack := new(needle.Needle)
readBack.Id = 1
_, err := reloaded.readNeedle(readBack, nil, nil)
require.Error(t, err)
_, _, _, err = reloaded.writeNeedle2(fixedNeedle(1, "after-reload"), true, true, false)
require.Error(t, err)
}
func TestBatchFsyncRollbackWithLevelDbMapper(t *testing.T) {
v, counting := newCountingVolumeWithKind(t, NeedleMapLevelDb)
kept := fixedNeedle(17, "leveldb-old")
_, _, _, err := v.writeNeedle2(kept, true, true, true)
require.NoError(t, err)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
_, _, _, err = v.writeNeedle2(fixedNeedle(17, "leveldb-new"), true, true, false)
require.Error(t, err)
readBack := new(needle.Needle)
readBack.Id = kept.Id
_, err = v.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("leveldb-old"), readBack.Data)
reloaded := reopenCountingVolume(t, v)
readBack = new(needle.Needle)
readBack.Id = kept.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("leveldb-old"), readBack.Data)
}
func TestBatchFsyncRollbackRestoresOriginalValueAfterRepeatedNeedleId(t *testing.T) {
v, counting := newCountingVolume(t)
kept := fixedNeedle(18, "original-value")
_, _, _, err := v.writeNeedle2(kept, true, true, true)
require.NoError(t, err)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
firstReplacement := fixedNeedle(18, "first-failed-value")
secondReplacement := fixedNeedle(18, "second-failed-value")
requests := []*needle.AsyncRequest{
needle.NewAsyncRequest(firstReplacement, true),
needle.NewAsyncRequest(secondReplacement, true),
}
v.processBatch(requests)
for _, request := range requests {
_, _, _, err = request.WaitComplete()
require.Error(t, err)
}
readBack := new(needle.Needle)
readBack.Id = kept.Id
_, err = v.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("original-value"), readBack.Data)
reloaded := reopenCountingVolume(t, v)
readBack = new(needle.Needle)
readBack.Id = kept.Id
_, err = reloaded.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("original-value"), readBack.Data)
}
// Rolling back a first-time write must not leave a tombstoned map entry: the
// stale offset would make the next write to that needle fail reading a header
// that no longer exists.
func TestBatchRollbackLeavesNoPhantomMapping(t *testing.T) {
v, counting := newCountingVolume(t)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
fresh := fixedNeedle(24, "batch-never-landed")
_, _, _, err := v.writeNeedle2(fresh, true, true, false)
require.Error(t, err)
_, _, _, err = v.writeNeedle2(fresh, true, true, false)
require.NoError(t, err, "rewriting a rolled-back needle must succeed")
readBack := new(needle.Needle)
readBack.Id = fresh.Id
_, err = v.readNeedle(readBack, nil, nil)
require.NoError(t, err)
require.Equal(t, []byte("batch-never-landed"), readBack.Data)
}
// A request that already failed on its own must still surface the batch
// failure: keeping its earlier error would hide that nothing was persisted,
// or that recovery itself failed.
func TestFailedBatchMarksEveryRequestFailed(t *testing.T) {
v, counting := newCountingVolume(t)
kept := fixedNeedle(22, "original")
_, _, _, err := v.writeNeedle2(kept, true, true, true)
require.NoError(t, err)
counting.syncErr = errors.New("batch fsync failed")
counting.syncErrOnce = true
badCookie := fixedNeedle(22, "wrong-cookie")
badCookie.Cookie = kept.Cookie + 1
requests := []*needle.AsyncRequest{
needle.NewAsyncRequest(badCookie, true),
needle.NewAsyncRequest(fixedNeedle(23, "fresh"), true),
}
v.processBatch(requests)
for _, request := range requests {
_, _, _, err = request.WaitComplete()
require.ErrorContains(t, err, "batch fsync failed")
}
}
// The quarantine is what CollectHeartbeat keys on: an unavailable volume must
// not be announced to the master at all.
func TestUnavailableVolumeIsSkippedInHeartbeat(t *testing.T) {
store := newTestStore(t, 1)
v := mountTestVolume(t, store.Locations[0], 1, "pics")
fillTestVolume(t, v)
v.markIoUnavailable(errors.New("batch recovery failed"))
heartbeat := store.CollectHeartbeat()
for _, m := range heartbeat.Volumes {
require.NotEqual(t, uint32(1), m.Id, "a failed-recovery volume must not be announced")
}
}
// The pre-stop drain exists so writes already assigned to this server land.
// Refusing them once stopping would turn every rolling restart into client
// write failures for the length of the drain.