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) {