From b9ad62fc16255545f2963b504c569011a7d1f3f5 Mon Sep 17 00:00:00 2001 From: ssshr-66 Date: Thu, 24 Sep 2026 06:57:44 +0800 Subject: [PATCH] [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 --- seaweed-volume/src/server/grpc_server.rs | 44 ++- seaweed-volume/src/server/handlers.rs | 3 + seaweed-volume/src/storage/volume.rs | 238 +++++++++++- weed/server/volume_grpc_copy_incremental.go | 14 +- weed/server/volume_grpc_read_all.go | 3 + weed/server/volume_grpc_tail.go | 14 + weed/storage/needle_map.go | 28 ++ weed/storage/needle_map/compact_map.go | 31 ++ weed/storage/needle_map/needle_value_map.go | 1 + weed/storage/needle_map_leveldb.go | 28 ++ weed/storage/needle_map_memory.go | 8 + weed/storage/needle_map_metric.go | 29 ++ weed/storage/needle_map_sorted_file.go | 27 ++ weed/storage/store.go | 7 + weed/storage/volume.go | 82 +++- weed/storage/volume_backup.go | 7 + weed/storage/volume_loading.go | 8 +- weed/storage/volume_read.go | 23 ++ weed/storage/volume_read_all.go | 3 + weed/storage/volume_write.go | 263 ++++++++++--- weed/storage/volume_write_fsync_test.go | 410 +++++++++++++++++++- 21 files changed, 1181 insertions(+), 90 deletions(-) diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index de651c611..4f1fedb21 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -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, >(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 = 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; } diff --git a/seaweed-volume/src/server/handlers.rs b/seaweed-volume/src/server/handlers.rs index 80367055b..a204b2dc2 100644 --- a/seaweed-volume/src/server/handlers.rs +++ b/seaweed-volume/src/server/handlers.rs @@ -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)) diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index 0b7f4fd62..df9941e1e 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -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, @@ -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, + /// 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, @@ -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, 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 { 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, 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 { 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 { 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 { + 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, 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 { + 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] diff --git a/weed/server/volume_grpc_copy_incremental.go b/weed/server/volume_grpc_copy_incremental.go index 449adb933..335b5c63e 100644 --- a/weed/server/volume_grpc_copy_incremental.go +++ b/weed/server/volume_grpc_copy_incremental.go @@ -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)] diff --git a/weed/server/volume_grpc_read_all.go b/weed/server/volume_grpc_read_all.go index 43334d7ad..5a40eeb39 100644 --- a/weed/server/volume_grpc_read_all.go +++ b/weed/server/volume_grpc_read_all.go @@ -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, diff --git a/weed/server/volume_grpc_tail.go b/weed/server/volume_grpc_tail.go index 120675b7f..2c42d314e 100644 --- a/weed/server/volume_grpc_tail.go +++ b/weed/server/volume_grpc_tail.go @@ -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 diff --git a/weed/storage/needle_map.go b/weed/storage/needle_map.go index 3311e931c..e13903e18 100644 --- a/weed/storage/needle_map.go +++ b/weed/storage/needle_map.go @@ -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 diff --git a/weed/storage/needle_map/compact_map.go b/weed/storage/needle_map/compact_map.go index df3c86d25..7b7169064 100644 --- a/weed/storage/needle_map/compact_map.go +++ b/weed/storage/needle_map/compact_map.go @@ -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() diff --git a/weed/storage/needle_map/needle_value_map.go b/weed/storage/needle_map/needle_value_map.go index 067d6358c..a2a258b38 100644 --- a/weed/storage/needle_map/needle_value_map.go +++ b/weed/storage/needle_map/needle_value_map.go @@ -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 } diff --git a/weed/storage/needle_map_leveldb.go b/weed/storage/needle_map_leveldb.go index f012f443d..19a535e3d 100644 --- a/weed/storage/needle_map_leveldb.go +++ b/weed/storage/needle_map_leveldb.go @@ -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 { diff --git a/weed/storage/needle_map_memory.go b/weed/storage/needle_map_memory.go index 4abb8891d..e8c21864e 100644 --- a/weed/storage/needle_map_memory.go +++ b/weed/storage/needle_map_memory.go @@ -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 diff --git a/weed/storage/needle_map_metric.go b/weed/storage/needle_map_metric.go index e1de97c6e..c8b412610 100644 --- a/weed/storage/needle_map_metric.go +++ b/weed/storage/needle_map_metric.go @@ -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 diff --git a/weed/storage/needle_map_sorted_file.go b/weed/storage/needle_map_sorted_file.go index ebf7d78fe..e5de53925 100644 --- a/weed/storage/needle_map_sorted_file.go +++ b/weed/storage/needle_map_sorted_file.go @@ -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 { diff --git a/weed/storage/store.go b/weed/storage/store.go index a8983bb36..329208483 100644 --- a/weed/storage/store.go +++ b/weed/storage/store.go @@ -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 diff --git a/weed/storage/volume.go b/weed/storage/volume.go index 1f5b582bd..129d69b89 100644 --- a/weed/storage/volume.go +++ b/weed/storage/volume.go @@ -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() diff --git a/weed/storage/volume_backup.go b/weed/storage/volume_backup.go index d56b49878..80377b34c 100644 --- a/weed/storage/volume_backup.go +++ b/weed/storage/volume_backup.go @@ -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 diff --git a/weed/storage/volume_loading.go b/weed/storage/volume_loading.go index 5a120f371..797e5c3ef 100644 --- a/weed/storage/volume_loading.go +++ b/weed/storage/volume_loading.go @@ -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) diff --git a/weed/storage/volume_read.go b/weed/storage/volume_read.go index 4976f7ff4..91292c4ed 100644 --- a/weed/storage/volume_read.go +++ b/weed/storage/volume_read.go @@ -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 diff --git a/weed/storage/volume_read_all.go b/weed/storage/volume_read_all.go index 7f58abfd7..74e181fa2 100644 --- a/weed/storage/volume_read_all.go +++ b/weed/storage/volume_read_all.go @@ -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 diff --git a/weed/storage/volume_write.go b/weed/storage/volume_write.go index a4964832c..13b26a1b4 100644 --- a/weed/storage/volume_write.go +++ b/weed/storage/volume_write.go @@ -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) diff --git a/weed/storage/volume_write_fsync_test.go b/weed/storage/volume_write_fsync_test.go index ab51bcff0..71c7884e3 100644 --- a/weed/storage/volume_write_fsync_test.go +++ b/weed/storage/volume_write_fsync_test.go @@ -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.