From 10b0f2b8adf2a585fe0b1ad700a76bfa2beeda68 Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Sat, 3 Oct 2026 16:04:30 +0300 Subject: [PATCH] volume server: refuse the rest of a grouped run after a durable index failure (#11576) * volume server: refuse the rest of a grouped run after a durable index failure A durable write whose needle-map put fails stops the volume taking writes (#10825): sent on its own, the next write then fails read only before it appends. The grouped run from #11543 appends and syncs every entry before publishing any, then kept publishing the entries after the failed one and acked them once the shared .idx sync went through. When the failed put tore its .idx row, the rows appended after it land off alignment, so the next load parses them as garbage and the acked writes are gone. Once a durable entry fails to publish, refuse every later entry of the run with ReadOnly, as the per-needle path does. The entries before it stay acked; their rows go down with the run's one .idx sync. The refused records stay on the .dat unindexed, as the failed one does on its own. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: refuse a grouped entry staged as a cookie mismatch too After a durable entry in a grouped run fails to index, the entries after it are refused as they would be on their own. On its own an entry meets check_writable before its cookie check, so one staged as a cookie mismatch now gets the refusal too, instead of keeping its staging error. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: trim a torn .idx row back so the next stays aligned A failed write_index_entry can leave half a row in the .idx. With the writer appending at the tail, every row written after it lands off alignment and the next load parses them as garbage, so a write acked behind a torn row does not come back. Trim the file back to idx_file_offset on a failed append, in both needle maps, and cover it with a test that writes past a torn row and reloads. * volume server: refuse queued Go writes once a durable index update fails processBatch kept writing after a failed nm.Put, and the single-write path checked IsReadOnly only outside the volume lock. A durable write whose index update fails now marks the volume noWriteOrDelete, and each queued request is checked before it appends, so the ones after a failed durable entry are refused the way a lone write is. Deletes get the same noWriteOrDelete refusal a lone delete gets. * volume server: refuse appends while a torn .idx row cannot be trimmed When trimming back a half-written .idx row itself fails, the next append would land after the torn bytes and every later row would parse off alignment on load. Latch the map as torn and refuse appends until the trim succeeds, on both CompactNeedleMap and RedbNeedleMap; the same latch covers an orphan row that could not be trimmed after a failed redb commit. The .idx writer is now opened with write+append access so truncate_to (set_len) works on Windows, where an append-only handle cannot trim. * volume server: write .idx rows at idx_file_offset, not via append mode Rust's OpenOptions on Windows strips FILE_WRITE_DATA whenever append is set so the handle stays strictly append-only, which makes set_len fail - the torn-row trim could never succeed there. Open the .idx writer with plain write access and seek to idx_file_offset before each row, the same positioned-write model the Go server uses. --------- Co-authored-by: Claude Opus 5.5 (1M context) Co-authored-by: Chris Lu --- seaweed-volume/src/storage/needle_map.rs | 110 ++++++++-- seaweed-volume/src/storage/volume.rs | 244 ++++++++++++++++++++++- weed/storage/store_duplicate_vid_test.go | 2 +- weed/storage/volume_write.go | 43 +++- weed/storage/volume_write_fsync_test.go | 55 +++++ 5 files changed, 426 insertions(+), 28 deletions(-) diff --git a/seaweed-volume/src/storage/needle_map.rs b/seaweed-volume/src/storage/needle_map.rs index e413a5333..9e683f0c9 100644 --- a/seaweed-volume/src/storage/needle_map.rs +++ b/seaweed-volume/src/storage/needle_map.rs @@ -195,7 +195,12 @@ impl NeedleMapKind { // ============================================================================ /// Trait for appending to an index file. -pub trait IdxFileWriter: Write + Send + Sync { +/// +/// The file is opened without append mode and each row is written at the +/// current `idx_file_offset` — the same positioned-write model the Go +/// server uses — because an append-mode handle cannot truncate on Windows, +/// where the std library keeps it strictly append-only. +pub trait IdxFileWriter: Write + Seek + Send + Sync { fn sync_all(&self) -> io::Result<()>; /// Truncate the file to `len` bytes. Used to remove an orphan .idx row /// left by a failed redb commit so `idx_file_offset` stays a contiguous @@ -224,6 +229,9 @@ pub struct CompactNeedleMap { metric: NeedleMapMetric, idx_file: Option>, idx_file_offset: u64, + /// The file holds bytes past `idx_file_offset` that must be trimmed + /// before another row can land aligned. + idx_torn: bool, } impl Default for CompactNeedleMap { @@ -240,6 +248,7 @@ impl CompactNeedleMap { metric: NeedleMapMetric::default(), idx_file: None, idx_file_offset: 0, + idx_torn: false, } } @@ -264,6 +273,7 @@ impl CompactNeedleMap { pub fn set_idx_file(&mut self, file: Box, offset: u64) { self.idx_file = Some(file); self.idx_file_offset = offset; + self.idx_torn = false; } /// True when an .idx file writer is attached. A read-only load leaves @@ -278,8 +288,8 @@ impl CompactNeedleMap { /// Insert or update an entry. Appends to .idx file if present. pub fn put(&mut self, key: NeedleId, offset: Offset, size: Size) -> io::Result<()> { // Persist to idx file BEFORE mutating in-memory state for crash consistency - if let Some(ref mut idx_file) = self.idx_file { - idx::write_index_entry(idx_file, key, offset, size)?; + self.append_to_index_file(key, offset, size)?; + if self.idx_file.is_some() { self.idx_file_offset += NEEDLE_MAP_ENTRY_SIZE as u64; } @@ -289,6 +299,41 @@ impl CompactNeedleMap { Ok(()) } + /// Write one row to the .idx file at `idx_file_offset`. A row left + /// half-written by a failed write is trimmed back to the offset so the + /// next row still lands aligned; while the trim keeps failing no row is + /// written at all, or it would sit off alignment and parse as garbage on + /// load. The offset itself is advanced by the caller once the row counts. + fn append_to_index_file( + &mut self, + key: NeedleId, + offset: Offset, + size: Size, + ) -> io::Result<()> { + let Some(idx_file) = self.idx_file.as_mut() else { + return Ok(()); + }; + if self.idx_torn { + match idx_file.truncate_to(self.idx_file_offset) { + Ok(()) => self.idx_torn = false, + Err(e) => { + return Err(io::Error::other(format!( + "index file still holds a torn row: {e}" + ))); + } + } + } + idx_file.seek(io::SeekFrom::Start(self.idx_file_offset))?; + if let Err(e) = idx::write_index_entry(idx_file, key, offset, size) { + if let Err(te) = idx_file.truncate_to(self.idx_file_offset) { + self.idx_torn = true; + tracing::warn!("failed to trim torn .idx row: {}", te); + } + return Err(e); + } + Ok(()) + } + /// Look up a needle. pub fn get(&self, key: NeedleId) -> Option { self.map.get(key) @@ -313,8 +358,8 @@ impl CompactNeedleMap { } // Always write tombstone to idx file (matching Go) - if let Some(ref mut idx_file) = self.idx_file { - idx::write_index_entry(idx_file, key, offset, TOMBSTONE_FILE_SIZE)?; + self.append_to_index_file(key, offset, TOMBSTONE_FILE_SIZE)?; + if self.idx_file.is_some() { self.idx_file_offset += NEEDLE_MAP_ENTRY_SIZE as u64; } @@ -464,6 +509,9 @@ pub struct RedbNeedleMap { metric: NeedleMapMetric, idx_file: Option>, idx_file_offset: u64, + /// The file holds bytes past `idx_file_offset` that must be trimmed + /// before another row can land aligned. + idx_torn: bool, /// Puts/deletes since the last durable checkpoint. writes_since_checkpoint: u32, } @@ -570,6 +618,7 @@ impl RedbNeedleMap { metric: NeedleMapMetric::default(), idx_file: None, idx_file_offset: 0, + idx_torn: false, writes_since_checkpoint: 0, }) } @@ -667,6 +716,7 @@ impl RedbNeedleMap { metric: NeedleMapMetric::default(), idx_file: None, idx_file_offset: 0, + idx_torn: false, writes_since_checkpoint: 0, }; @@ -858,6 +908,7 @@ impl RedbNeedleMap { pub fn set_idx_file(&mut self, file: Box, offset: u64) { self.idx_file = Some(file); self.idx_file_offset = offset; + self.idx_torn = false; } /// True when an .idx file writer is attached. See CompactNeedleMap. @@ -874,9 +925,7 @@ impl RedbNeedleMap { // commit leaves an orphan row in .idx that redb doesn't reflect, and // advancing the offset here would let a later checkpoint record it as // reflected, making the reload skip it permanently. - if let Some(ref mut idx_file) = self.idx_file { - idx::write_index_entry(idx_file, key, offset, size)?; - } + self.append_to_index_file(key, offset, size)?; let key_u64: u64 = key.into(); let packed = pack_needle_value(&NeedleValue { offset, size }); @@ -947,6 +996,41 @@ impl RedbNeedleMap { Ok(()) } + /// Write one row to the .idx file at `idx_file_offset`. A row left + /// half-written by a failed write is trimmed back to the offset so the + /// next row still lands aligned; while the trim keeps failing no row is + /// written at all, or it would sit off alignment and parse as garbage on + /// load. The offset itself is advanced by the caller once the row counts. + fn append_to_index_file( + &mut self, + key: NeedleId, + offset: Offset, + size: Size, + ) -> io::Result<()> { + let Some(idx_file) = self.idx_file.as_mut() else { + return Ok(()); + }; + if self.idx_torn { + match idx_file.truncate_to(self.idx_file_offset) { + Ok(()) => self.idx_torn = false, + Err(e) => { + return Err(io::Error::other(format!( + "index file still holds a torn row: {e}" + ))); + } + } + } + idx_file.seek(io::SeekFrom::Start(self.idx_file_offset))?; + if let Err(e) = idx::write_index_entry(idx_file, key, offset, size) { + if let Err(te) = idx_file.truncate_to(self.idx_file_offset) { + self.idx_torn = true; + tracing::warn!("failed to trim torn .idx row: {}", te); + } + return Err(e); + } + Ok(()) + } + /// Look up a needle. A redb failure is an ERROR, not an absent needle: /// answering "not found" would turn a database problem into a read miss /// and let a delete report success without recording a tombstone. @@ -993,9 +1077,7 @@ impl RedbNeedleMap { return Ok(None); }; - if let Some(ref mut idx_file) = self.idx_file { - idx::write_index_entry(idx_file, key, offset, TOMBSTONE_FILE_SIZE)?; - } + self.append_to_index_file(key, offset, TOMBSTONE_FILE_SIZE)?; let deleted_nv = NeedleValue { offset: old.offset, @@ -1085,10 +1167,13 @@ impl RedbNeedleMap { /// a failed redb commit. Without this the next successful write appends /// after the orphan, `idx_file_offset` advances past it, and a later /// checkpoint records an offset that makes the reload skip the orphan. + /// When the trim fails the file is latched torn so no later row lands + /// after bytes the map does not reflect. fn truncate_idx_to_offset(&mut self) { if let Some(ref mut idx_file) = self.idx_file && let Err(e) = idx_file.truncate_to(self.idx_file_offset) { + self.idx_torn = true; tracing::warn!("failed to truncate orphan .idx row: {}", e); } } @@ -1139,6 +1224,7 @@ impl RedbNeedleMap { self.db = reopened.db; self.metric = reopened.metric; self.idx_file_offset = actual_idx_size; + self.idx_torn = reopened.idx_torn; // The reopen replayed all rows since the last durable checkpoint // non-durably; start the counter fresh. self.writes_since_checkpoint = 0; @@ -1691,7 +1777,7 @@ mod tests { ) .unwrap(); let writer = std::fs::OpenOptions::new() - .append(true) + .write(true) .open(&idx_path) .unwrap(); nm.set_idx_file(Box::new(writer), idx_size); diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index bfe04df4c..ca7b76c71 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -1730,8 +1730,10 @@ impl Volume { let mut idx_reader = io::BufReader::new(&idx_file); let mut nm = CompactNeedleMap::load_from_idx(&mut idx_reader, self.version())?; - // Re-open for append-only writes - let write_file = OpenOptions::new().append(true).open(idx_path)?; + // Re-open for positioned writes: rows are written at + // idx_file_offset so a torn tail can be trimmed (an append-mode + // handle cannot set_len on Windows). + let write_file = OpenOptions::new().write(true).open(idx_path)?; nm.set_idx_file(Box::new(write_file), idx_size); self.nm = Some(NeedleMap::InMemory(nm)); } @@ -1788,8 +1790,10 @@ impl Volume { cache_bytes, )?; - // Re-open for append-only writes - let write_file = OpenOptions::new().append(true).open(idx_path)?; + // Re-open for positioned writes: rows are written at + // idx_file_offset so a torn tail can be trimmed (an append-mode + // handle cannot set_len on Windows). + let write_file = OpenOptions::new().write(true).open(idx_path)?; nm.set_idx_file(Box::new(write_file), idx_size); self.nm = Some(NeedleMap::Redb(nm)); } @@ -2270,8 +2274,10 @@ impl Volume { /// sync and one .idx sync per run of distinct needle ids that holds a /// durable write, instead of two per durable needle. Nothing in such a /// run is published before its sync, so a failed sync takes the whole - /// run back off the .dat and fails every entry. A repeated id starts a - /// new run, so its dedup and cookie checks see the earlier write. + /// run back off the .dat and fails every entry. A durable entry that + /// fails to index stops the volume taking writes, and the entries after + /// it in the run are refused read only. A repeated id starts a new run, + /// so its dedup and cookie checks see the earlier write. pub fn write_needles_grouped( &mut self, writes: &mut [(Needle, bool)], @@ -2334,11 +2340,21 @@ impl Volume { } self.last_append_at_ns = last_append_at_ns; + // A durable entry that fails to publish stops the volume taking + // writes, so the entries after it are refused the way a lone + // write would be. + let mut refused = false; for ((n, fsync), r) in run.iter().zip(staged.iter_mut()) { - if let Ok(Some(offset)) = *r + if refused { + *r = Err(self + .check_writable() + .err() + .unwrap_or(VolumeError::ReadOnly(self.id))); + } else if let Ok(Some(offset)) = *r && let Err(e) = self.publish_write(n, offset, *fsync) { *r = Err(e); + refused = *fsync; } } @@ -3635,9 +3651,12 @@ impl Volume { .unwrap_or(false); if needs_idx_writer { let idx_path = self.file_name(".idx"); + // Positioned writes: an append-mode handle cannot set_len on + // Windows. let write_file = OpenOptions::new() - .append(true) + .write(true) .create(true) + .truncate(false) .open(&idx_path)?; let idx_size = trim_torn_idx_tail(&write_file, &idx_path)?; if let Some(ref mut nm) = self.nm { @@ -6492,6 +6511,215 @@ mod tests { assert!(matches!(results[0], Err(VolumeError::Unavailable(_)))); } + /// An .idx writer that tears its `tear_at`-th row: half of the row + /// reaches the file, then the write fails. + struct TornIdxWriter { + file: File, + writes: usize, + tear_at: usize, + fail_truncates: bool, + } + + impl Write for TornIdxWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.writes += 1; + if self.writes == self.tear_at { + self.file.write_all(&buf[..buf.len() / 2])?; + return Err(io::Error::other("injected torn .idx write")); + } + self.file.write(buf) + } + + fn flush(&mut self) -> io::Result<()> { + self.file.flush() + } + } + + impl Seek for TornIdxWriter { + fn seek(&mut self, pos: SeekFrom) -> io::Result { + self.file.seek(pos) + } + } + + impl crate::storage::needle_map::IdxFileWriter for TornIdxWriter { + fn sync_all(&self) -> io::Result<()> { + self.file.sync_all() + } + + fn truncate_to(&mut self, len: u64) -> io::Result<()> { + if self.fail_truncates { + return Err(io::Error::other("injected trim failure")); + } + self.file.set_len(len) + } + } + + /// A durable entry whose index update fails stops the volume taking + /// writes, as it does when sent on its own, so the entries after it in + /// the run are refused instead of indexed behind a row that may be + /// torn. Entries before it stay acked. + #[test] + fn test_grouped_failed_durable_index_refuses_rest_of_run() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + let mut kept = batch_needle(4, 0xdd, b"kept"); + v.write_needle(&mut kept, true, true).unwrap(); + let mut other = batch_needle(5, 0xee, b"other"); + v.write_needle(&mut other, true, true).unwrap(); + + let nm = v.nm.as_mut().unwrap(); + let idx_len = nm.index_file_size(); + let file = OpenOptions::new() + .write(true) + .open(format!("{dir}/1.idx")) + .unwrap(); + nm.set_idx_file( + Box::new(TornIdxWriter { + file, + writes: 0, + tear_at: 2, + fail_truncates: false, + }), + idx_len, + ); + + let mut writes = vec![ + (batch_needle(1, 0xaa, b"before"), true), + (batch_needle(2, 0xbb, b"torn"), true), + (batch_needle(3, 0xcc, b"after"), true), + (batch_needle(4, 0xdd, b"kept"), true), + (batch_needle(5, 0xef, b"wrong cookie"), true), + ]; + let results = v.write_needles_grouped(&mut writes); + // Two setup writes, then the run's one sync of each file. + assert_eq!(v.sync_counts_for_test(), (3, 3)); + assert!(v.is_read_only()); + drop(v); + + let reopened = match Volume::new( + dir, + dir, + VolumeId(1), + NeedleMapKind::InMemory, + &VolumeSpec::default(), + ) { + Ok(v) => v, + Err(e) => panic!("the volume does not reload: {e}"), + }; + for ((n, _), r) in writes.iter().zip(&results) { + if r.is_ok() { + let mut got = Needle { + id: n.id, + ..Needle::default() + }; + let read = reopened.read_needle(&mut got); + assert!( + read.is_ok() && got.data == n.data, + "acked write {} lost on reload: {read:?}", + n.id.0 + ); + } + } + + assert!(matches!(results[0], Ok((_, _, false))), "{results:?}"); + assert!(matches!(results[1], Err(VolumeError::Io(_))), "{results:?}"); + for r in &results[2..] { + assert!( + matches!(r, Err(VolumeError::ReadOnly(VolumeId(1)))), + "refused as it would be sent on its own: {results:?}" + ); + } + } + + /// A torn .idx row is trimmed back, so the row a later write appends + /// still lands aligned and survives a reload. + #[test] + fn test_failed_index_write_is_trimmed() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let nm = v.nm.as_mut().unwrap(); + let idx_len = nm.index_file_size(); + let file = OpenOptions::new() + .write(true) + .open(format!("{dir}/1.idx")) + .unwrap(); + nm.set_idx_file( + Box::new(TornIdxWriter { + file, + writes: 0, + tear_at: 1, + fail_truncates: false, + }), + idx_len, + ); + + let mut writes = vec![ + (batch_needle(1, 0xaa, b"torn"), false), + (batch_needle(2, 0xbb, b"after"), true), + ]; + let results = v.write_needles_grouped(&mut writes); + assert!(matches!(results[0], Err(VolumeError::Io(_))), "{results:?}"); + assert!(matches!(results[1], Ok((_, _, false))), "{results:?}"); + drop(v); + + let reopened = match Volume::new( + dir, + dir, + VolumeId(1), + NeedleMapKind::InMemory, + &VolumeSpec::default(), + ) { + Ok(v) => v, + Err(e) => panic!("the volume does not reload: {e}"), + }; + let mut got = Needle { + id: NeedleId(2), + ..Needle::default() + }; + reopened.read_needle(&mut got).unwrap(); + assert_eq!(got.data, b"after"); + } + + /// When the trim of a torn .idx row itself fails, no later row is + /// appended after the torn bytes. + #[test] + fn test_untrimmed_torn_row_refuses_later_appends() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + + let nm = v.nm.as_mut().unwrap(); + let idx_len = nm.index_file_size(); + let file = OpenOptions::new() + .write(true) + .open(format!("{dir}/1.idx")) + .unwrap(); + nm.set_idx_file( + Box::new(TornIdxWriter { + file, + writes: 0, + tear_at: 1, + fail_truncates: true, + }), + idx_len, + ); + + let mut writes = vec![ + (batch_needle(1, 0xaa, b"torn"), false), + (batch_needle(2, 0xbb, b"after"), false), + ]; + let results = v.write_needles_grouped(&mut writes); + assert!(matches!(results[0], Err(VolumeError::Io(_))), "{results:?}"); + assert!(results[1].is_err(), "{results:?}"); + + // The second row was never appended behind the torn bytes. + let size = std::fs::metadata(format!("{dir}/1.idx")).unwrap().len(); + assert_eq!(size, idx_len + 8); + } + /// The I/O error streak after writing four durable needles, the ones in /// `failing` with a media error on their append, first one at a time and /// then as one grouped run. Also returns which grouped writes landed. diff --git a/weed/storage/store_duplicate_vid_test.go b/weed/storage/store_duplicate_vid_test.go index 49169439c..740d2bab5 100644 --- a/weed/storage/store_duplicate_vid_test.go +++ b/weed/storage/store_duplicate_vid_test.go @@ -117,7 +117,7 @@ func TestDeleteVolumeGuardWaitsForInFlightCopyWrite(t *testing.T) { // The write wins; the delete must see the new needle and refuse, leaving // every copy intact. n := &needle.Needle{Id: types.Uint64ToNeedleId(1), Data: []byte("x")} - _, _, _, err := copy1.doWriteRequest(n, false) + _, _, _, err := copy1.doWriteRequest(n, false, false) require.NoError(t, err) copy1.dataFileAccessLock.Unlock() diff --git a/weed/storage/volume_write.go b/weed/storage/volume_write.go index 95f84f6e8..2b03cb21c 100644 --- a/weed/storage/volume_write.go +++ b/weed/storage/volume_write.go @@ -266,6 +266,9 @@ func (v *Volume) syncWrite(n *needle.Needle, checkCookie bool, fsync bool) (offs if err := v.UnavailableError(); err != nil { return 0, 0, false, err } + if v.IsReadOnly() { + return 0, 0, false, fmt.Errorf("volume %d is read only", v.Id) + } // 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. @@ -274,7 +277,7 @@ func (v *Volume) syncWrite(n *needle.Needle, checkCookie bool, fsync bool) (offs } if !fsync { - return v.doWriteRequest(n, checkCookie) + return v.doWriteRequest(n, checkCookie, fsync) } end, _, statErr := v.DataBackend.GetStat() @@ -286,7 +289,7 @@ func (v *Volume) syncWrite(n *needle.Needle, checkCookie bool, fsync bool) (offs priorOffset, priorSize, hasPrior = nv.Offset, nv.Size, true } - offset, size, isUnchanged, err = v.doWriteRequest(n, checkCookie) + offset, size, isUnchanged, err = v.doWriteRequest(n, checkCookie, fsync) if err != nil { return } @@ -367,7 +370,7 @@ func (v *Volume) writeNeedle2(n *needle.Needle, checkCookie bool, fsync bool, is } } -func (v *Volume) doWriteRequest(n *needle.Needle, checkCookie bool) (offset uint64, size Size, isUnchanged bool, err error) { +func (v *Volume) doWriteRequest(n *needle.Needle, checkCookie bool, fsync bool) (offset uint64, size Size, isUnchanged bool, err error) { // glog.V(4).Infof("writing needle %s", needle.NewFileIdFromNeedle(v.Id, n).String()) if v.isFileUnchanged(n) { size = Size(n.DataSize) @@ -412,6 +415,14 @@ func (v *Volume) doWriteRequest(n *needle.Needle, checkCookie bool) (offset uint if err = v.nm.Put(n.Id, ToOffset(int64(offset)), n.Size); err != nil { err = fmt.Errorf("index needle %d of volume %d at offset %d: %w", n.Id, v.Id, offset, err) glog.V(0).Info(err) + if fsync { + // The record is down but nothing indexes it. Stop taking + // writes rather than append past it, same as the Rust + // volume server. + v.noWriteLock.Lock() + v.noWriteOrDelete = true + v.noWriteLock.Unlock() + } } } if v.lastModifiedTsSeconds < n.LastModified { @@ -428,6 +439,9 @@ func (v *Volume) syncDelete(n *needle.Needle) (Size, error) { if err := v.UnavailableError(); err != nil { return 0, err } + if _, noWriteOrDelete, _, _ := v.ReadOnlyReasons(); noWriteOrDelete { + return 0, fmt.Errorf("volume %d is read only", v.Id) + } if v.nm == nil { return 0, nil @@ -519,12 +533,27 @@ func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) { batchSnapshots[needleID] = snapshot orderedSnapshots = append(orderedSnapshots, snapshot) } + // Every queued request meets the same refusal it would get sent on + // its own: once a durable index update fails the volume stops taking + // writes, and a lone write is turned away before it appends. if currentRequests[i].IsWriteRequest { - offset, size, isUnchanged, err := v.doWriteRequest(currentRequests[i].N, true) - currentRequests[i].UpdateResult(offset, uint64(size), isUnchanged, err) + if unavailableErr := v.UnavailableError(); unavailableErr != nil { + currentRequests[i].UpdateResult(0, 0, false, unavailableErr) + } else if v.IsReadOnly() { + currentRequests[i].UpdateResult(0, 0, false, fmt.Errorf("volume %d is read only", v.Id)) + } else { + offset, size, isUnchanged, err := v.doWriteRequest(currentRequests[i].N, true, 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 unavailableErr := v.UnavailableError(); unavailableErr != nil { + currentRequests[i].UpdateResult(0, 0, false, unavailableErr) + } else if _, noWriteOrDelete, _, _ := v.ReadOnlyReasons(); noWriteOrDelete { + currentRequests[i].UpdateResult(0, 0, false, fmt.Errorf("volume %d is read only", v.Id)) + } else { + size, err := v.doDeleteRequest(currentRequests[i].N) + currentRequests[i].UpdateResult(0, uint64(size), false, err) + } } snapshot.observeCurrent(v) } diff --git a/weed/storage/volume_write_fsync_test.go b/weed/storage/volume_write_fsync_test.go index f842923c7..8bf5662a0 100644 --- a/weed/storage/volume_write_fsync_test.go +++ b/weed/storage/volume_write_fsync_test.go @@ -514,6 +514,61 @@ func TestFailedBatchMarksEveryRequestFailed(t *testing.T) { } } +// failingPutMapper fails the failAt-th Put. +type failingPutMapper struct { + NeedleMapper + puts int + failAt int + err error +} + +func (m *failingPutMapper) Put(key types.NeedleId, offset types.Offset, size types.Size) error { + m.puts++ + if m.puts == m.failAt { + return m.err + } + return m.NeedleMapper.Put(key, offset, size) +} + +// A durable write whose index update fails stops the volume taking writes: +// the requests after it in the same batch are refused, as a write sent on +// its own would be, instead of indexed behind it. Entries before it stay +// acked. +func TestProcessBatchRefusesWritesAfterDurableIndexFailure(t *testing.T) { + v, _ := newCountingVolume(t) + + kept := fixedNeedle(4, "kept") + _, _, _, err := v.writeNeedle2(kept, true, true, true) + require.NoError(t, err) + + v.nm = &failingPutMapper{ + NeedleMapper: v.nm, + failAt: 2, + err: errors.New("index write failed"), + } + + wrongCookie := fixedNeedle(4, "kept") + wrongCookie.Cookie = kept.Cookie + 1 + requests := []*needle.AsyncRequest{ + needle.NewAsyncRequest(fixedNeedle(1, "before"), true), + needle.NewAsyncRequest(fixedNeedle(2, "torn"), true), + needle.NewAsyncRequest(fixedNeedle(3, "after"), true), + needle.NewAsyncRequest(fixedNeedle(4, "kept"), true), + needle.NewAsyncRequest(wrongCookie, true), + } + v.processBatch(requests) + + _, _, _, err = requests[0].WaitComplete() + require.NoError(t, err) + _, _, _, err = requests[1].WaitComplete() + require.Error(t, err) + for _, request := range requests[2:] { + _, _, _, err = request.WaitComplete() + require.ErrorContains(t, err, "read only", "refused as it would be sent on its own") + } + require.True(t, v.IsReadOnly()) +} + // The quarantine is what CollectHeartbeat keys on: an unavailable volume must // not be announced to the master at all. func TestUnavailableVolumeIsSkippedInHeartbeat(t *testing.T) {