From 94a68fa9b9f1cd178cd3f481af1f250dc504fe45 Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Thu, 24 Sep 2026 02:08:26 +0300 Subject: [PATCH] volume: walk_index_file keeps row alignment across short reads (#11412) * volume: walk_index_file keeps row alignment across short reads walk_index_file issued one Read::read per batch and decoded whatever came back. Read::read may legally return a short count that is not a multiple of the 17-byte entry size (FUSE and network filesystems, a BufReader whose capacity is not a multiple of 17). The split entry at the end of the batch was dropped with no carry and the next read started mid-entry, so every later row was decoded from misaligned bytes and fed to the index as a garbage key/offset/size. This function backs every in-memory index load. Go's WalkIndexFile is immune because it reads through io.ReaderAt, which returns a full buffer or an error. Fill the batch buffer until it is full or the reader reports EOF, retrying ErrorKind::Interrupted, and only then decode whole entries. Reads stay batched at ROWS_TO_READ entries. EOF semantics are unchanged and match Go: on io.EOF Go decodes the whole entries in the final buffer, ignores a trailing partial entry and returns nil. A torn final entry is still skipped without an error here. Co-Authored-By: Claude Fable 5.1 * volume: trim walk_index_file comments The batch-fill loop and the ShortReader test helper each carried a paragraph where a sentence suffices. --------- Co-authored-by: Claude Fable 5.1 Co-authored-by: Chris Lu --- seaweed-volume/src/storage/idx/mod.rs | 136 ++++++++++++++++++++++++-- 1 file changed, 130 insertions(+), 6 deletions(-) diff --git a/seaweed-volume/src/storage/idx/mod.rs b/seaweed-volume/src/storage/idx/mod.rs index 0a7182c3a..a56a91d81 100644 --- a/seaweed-volume/src/storage/idx/mod.rs +++ b/seaweed-volume/src/storage/idx/mod.rs @@ -21,12 +21,26 @@ where let mut buf = vec![0u8; NEEDLE_MAP_ENTRY_SIZE * ROWS_TO_READ]; loop { - let count = match reader.read(&mut buf) { - Ok(0) => return Ok(()), - Ok(n) => n, - Err(ref e) if e.kind() == io::ErrorKind::UnexpectedEof => return Ok(()), - Err(e) => return Err(e), - }; + // Fill the batch before decoding: `read` may return a count that is + // not a multiple of the entry size, and a split entry would misalign + // every later row. Go is immune: `ReadAt` fills or errors. + let mut count = 0; + let mut eof = false; + while count < buf.len() { + match reader.read(&mut buf[count..]) { + Ok(0) => { + eof = true; + break; + } + Ok(n) => count += n, + Err(ref e) if e.kind() == io::ErrorKind::Interrupted => continue, + Err(ref e) if e.kind() == io::ErrorKind::UnexpectedEof => { + eof = true; + break; + } + Err(e) => return Err(e), + } + } let mut i = 0; while i + NEEDLE_MAP_ENTRY_SIZE <= count { @@ -34,6 +48,11 @@ where f(key, offset, size)?; i += NEEDLE_MAP_ENTRY_SIZE; } + + // A trailing partial entry at EOF is ignored, as Go does on `io.EOF`. + if eof { + return Ok(()); + } } } @@ -177,6 +196,111 @@ mod tests { data } + /// Reader that hands back at most `chunk` bytes per `read`. 7 is coprime + /// with the 17-byte entry size, so nearly every read ends mid-entry. With + /// `interrupts`, every other call fails with `ErrorKind::Interrupted`. + struct ShortReader { + inner: Cursor>, + chunk: usize, + interrupts: bool, + interrupt_next: bool, + } + + impl ShortReader { + fn new(data: Vec, interrupts: bool) -> Self { + ShortReader { + inner: Cursor::new(data), + chunk: 7, + interrupts, + interrupt_next: false, + } + } + } + + impl Read for ShortReader { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + if self.interrupt_next { + self.interrupt_next = false; + return Err(io::Error::from(io::ErrorKind::Interrupted)); + } + self.interrupt_next = self.interrupts; + let n = buf.len().min(self.chunk); + self.inner.read(&mut buf[..n]) + } + } + + impl Seek for ShortReader { + fn seek(&mut self, pos: SeekFrom) -> io::Result { + self.inner.seek(pos) + } + } + + fn walk_all(reader: &mut R, start_from: u64) -> Vec<(NeedleId, i64, Size)> { + let mut collected = Vec::new(); + walk_index_file(reader, start_from, |key, offset, size| { + collected.push((key, offset.to_actual_offset(), size)); + Ok(()) + }) + .unwrap(); + collected + } + + /// More than one ROWS_TO_READ batch, so the walk crosses a buffer refill. + fn many_entries() -> Vec<(NeedleId, Offset, Size)> { + (0..(ROWS_TO_READ as u64 * 2 + 37)) + .map(|i| { + ( + NeedleId(i * 7 + 1), + Offset::from_actual_offset(i as i64 * 128), + Size(i as i32 + 1), + ) + }) + .collect() + } + + #[test] + fn test_walk_index_file_short_reads_keep_alignment() { + let data = idx_bytes(&many_entries()); + let expected = walk_all(&mut Cursor::new(data.clone()), 0); + assert_eq!(expected.len(), ROWS_TO_READ * 2 + 37); + + let mut short = ShortReader::new(data, false); + assert_eq!(walk_all(&mut short, 0), expected); + } + + #[test] + fn test_walk_index_file_retries_interrupted_reads() { + let data = idx_bytes(&many_entries()); + let expected = walk_all(&mut Cursor::new(data.clone()), 0); + + let mut short = ShortReader::new(data, true); + assert_eq!(walk_all(&mut short, 0), expected); + } + + #[test] + fn test_walk_index_file_short_reads_start_from() { + let data = idx_bytes(&many_entries()); + let expected = walk_all(&mut Cursor::new(data.clone()), 0); + + let start = ROWS_TO_READ as u64 + 5; + let mut short = ShortReader::new(data, false); + assert_eq!(walk_all(&mut short, start), expected[start as usize..]); + } + + #[test] + fn test_walk_index_file_ignores_trailing_partial_entry() { + // A torn final entry is dropped without an error, as Go does on io.EOF. + let entries = many_entries(); + let mut data = idx_bytes(&entries); + data.extend_from_slice(&[0xAB; NEEDLE_MAP_ENTRY_SIZE - 1]); + + let expected = walk_all(&mut Cursor::new(idx_bytes(&entries)), 0); + assert_eq!(walk_all(&mut Cursor::new(data.clone()), 0), expected); + + let mut short = ShortReader::new(data, false); + assert_eq!(walk_all(&mut short, 0), expected); + } + #[test] fn test_check_index_file_clean() { let data = idx_bytes(&[