diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index ed82c2769..e5b9fc30b 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -497,10 +497,11 @@ impl DatScanPlan { /// `ScanVolumeFileFrom` feeds a `VolumeFileScanner`: only the record being /// visited is in memory. Reads are positional and never touch the /// `Volume`. The pass ends `Ok` at the end of the data, at a record that - /// does not fit before `end`, at a corrupt header, or when `visit` breaks. - /// It fails on any read error, including a short read below `end`: every - /// byte below `end` existed when the plan was taken, so a short read means - /// the file was truncated under the plan and the pass is incomplete. + /// does not fit before `end`, or when `visit` breaks. A corrupt header + /// fails the pass: ending early would let a truncated tail look complete. + /// It also fails on any read error, including a short read below `end`: + /// every byte below `end` existed when the plan was taken, so a short read + /// means the file was truncated under the plan and the pass is incomplete. pub(crate) fn scan( &self, mut visit: impl FnMut(RawNeedle<'_>) -> ControlFlow<()>, @@ -516,11 +517,18 @@ impl DatScanPlan { if size.0 == 0 && id.is_empty() { break; } - // A negative size is a corrupt header: body_length would size the - // buffer from a negative length or walk the scan from a wrong - // offset. Go's scanners stop here by returning io.EOF. + // A negative size is a corrupt header: the record length it + // implies cannot advance the scan past it, and stopping quietly + // would let a truncated tail look complete. Go's scan fails here. if size.0 < 0 { - break; + return Err(VolumeError::Io(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "corrupt needle header at offset {offset}: size {}, record length {}", + size.0, + get_actual_size(size, self.version) + ), + ))); } // Nothing past `end` was a complete record when the plan was @@ -2356,8 +2364,18 @@ impl Volume { break; } - let body_length = needle::needle_body_length(size, version); - let total_size = NEEDLE_HEADER_SIZE as i64 + body_length; + let total_size = get_actual_size(size, version); + // A corrupt header can make the record length zero or negative; + // the scan cannot advance past it. + if total_size <= 0 { + return Err(VolumeError::Io(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "corrupt needle header at offset {offset}: size {}, record length {total_size}", + size.0 + ), + ))); + } if size.is_deleted() || size.0 <= 0 { offset += total_size; @@ -4609,7 +4627,8 @@ pub fn scan_volume_file( break; // end of valid data } // A negative size is a corrupt header, and body_length would advance the - // walk backwards from it. Go's scanners stop here by returning io.EOF. + // walk backwards from it. The index rebuild this feeds salvages the + // records before it, like Go's rebuild scanner stopping on io.EOF. if size.0 < 0 { break; } @@ -5207,12 +5226,11 @@ mod tests { #[test] fn dat_scan_plan_ends_the_pass_at_a_corrupt_header() { - // A negative size is a corrupt header: sizing a buffer from it - // overflows, and a small one walks the scan from a wrong offset. A - // size running past the end bound cannot be a complete record either, - // and one near i32::MAX overflows the padding arithmetic. Both end - // the pass after the records before them, without reading or - // allocating the bogus body. + // A negative size is a corrupt header: the record length it implies + // cannot advance the scan past it, so the pass fails. A size running + // past the end bound, like one near i32::MAX, cannot be a complete + // record and ends the pass after the records before it, without + // reading or allocating the bogus body. for bad_size in [-1000, i32::MAX] { let tmp = TempDir::new().unwrap(); let dir = tmp.path().to_str().unwrap(); @@ -5234,8 +5252,21 @@ mod tests { .write_all(&bad) .unwrap(); - let records = scan_all(&v.dat_scan_plan(sb_size).unwrap()); - assert_eq!(records.len(), 1, "size {}: only the good record", bad_size); + let mut visited = 0; + let result = v.dat_scan_plan(sb_size).unwrap().scan(|_| { + visited += 1; + ControlFlow::Continue(()) + }); + if bad_size < 0 { + let err = result.unwrap_err(); + assert!( + matches!(&err, VolumeError::Io(e) if e.kind() == io::ErrorKind::InvalidData), + "size {bad_size}: expected a corrupt-data error, got {err:?}" + ); + } else { + result.unwrap(); + } + assert_eq!(visited, 1, "size {bad_size}: only the good record"); } } @@ -5276,6 +5307,62 @@ mod tests { ); } + fn append_corrupt_header(v: &Volume, size: i32) { + let mut dat = OpenOptions::new() + .append(true) + .open(v.file_name(".dat")) + .unwrap(); + let mut corrupt = [0u8; NEEDLE_HEADER_SIZE]; + NeedleId(99).to_bytes(&mut corrupt[4..12]); + Size(size).to_bytes(&mut corrupt[12..16]); + dat.write_all(&corrupt).unwrap(); + } + + #[test] + fn dat_scan_plan_fails_at_a_header_it_cannot_advance_past() { + // A size so negative the record length is zero or less cannot be + // advanced past; the pass must fail rather than end early and let a + // truncated tail look complete. + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + write_test_needle(&mut v, 1, b"good"); + append_corrupt_header(&v, -100); + + let plan = v.dat_scan_plan(v.super_block.block_size() as u64).unwrap(); + let mut visited = 0; + let err = plan + .scan(|_| { + visited += 1; + ControlFlow::Continue(()) + }) + .unwrap_err(); + + assert_eq!( + visited, 1, + "the record before the corrupt header is visited" + ); + assert!( + matches!(&err, VolumeError::Io(e) if e.kind() == io::ErrorKind::InvalidData), + "expected a corrupt-data error, got {err:?}" + ); + } + + #[test] + fn read_all_needles_fails_at_a_header_it_cannot_advance_past() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + write_test_needle(&mut v, 1, b"good"); + append_corrupt_header(&v, -100); + + let err = v.read_all_needles().unwrap_err(); + assert!( + matches!(&err, VolumeError::Io(e) if e.kind() == io::ErrorKind::InvalidData), + "expected a corrupt-data error, got {err:?}" + ); + } + #[test] fn test_volume_write_read() { let tmp = TempDir::new().unwrap(); diff --git a/seaweed-volume/src/storage/volume_idx_repair.rs b/seaweed-volume/src/storage/volume_idx_repair.rs index e67a0c881..e731269a6 100644 --- a/seaweed-volume/src/storage/volume_idx_repair.rs +++ b/seaweed-volume/src/storage/volume_idx_repair.rs @@ -122,7 +122,19 @@ impl Volume { } else { found.remove(&id); } - offset += NEEDLE_HEADER_SIZE as i64 + needle_body_length(size, version); + let record_size = NEEDLE_HEADER_SIZE as i64 + needle_body_length(size, version); + // A corrupt header can make the record length zero or negative; + // the scan cannot advance past it. + if record_size <= 0 { + return Err(VolumeError::Io(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "corrupt needle header at offset {offset}: size {}, record length {record_size}", + size.0 + ), + ))); + } + offset += record_size; } Ok((found, order)) @@ -426,4 +438,29 @@ mod tests { "needle 1 should have been recovered" ); } + + #[test] + fn test_scan_dat_head_fails_at_a_header_it_cannot_advance_past() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + write_test_volume(dir, 2); + + let mut dat = OpenOptions::new() + .append(true) + .open(format!("{}/1.dat", dir)) + .unwrap(); + let mut corrupt = [0u8; NEEDLE_HEADER_SIZE]; + NeedleId(99).to_bytes(&mut corrupt[4..12]); + Size(-100).to_bytes(&mut corrupt[12..16]); + dat.write_all(&corrupt).unwrap(); + drop(dat); + + let v = open_volume(dir); + let first = v.super_block.block_size() as i64; + let err = v.scan_dat_head(v.version(), first, 10).unwrap_err(); + assert!( + matches!(&err, VolumeError::Io(e) if e.kind() == io::ErrorKind::InvalidData), + "expected a corrupt-data error, got {err:?}" + ); + } } diff --git a/weed/storage/volume_read.go b/weed/storage/volume_read.go index fa96693de..4976f7ff4 100644 --- a/weed/storage/volume_read.go +++ b/weed/storage/volume_read.go @@ -299,7 +299,13 @@ func ScanVolumeFileFrom(version needle.Version, datBackend backend.BackendStorag glog.V(0).Infof("visit needle error: %v", err) return fmt.Errorf("visit needle error: %w", err) } - offset += NeedleHeaderSize + rest + // A corrupt header can carry a size so negative that the record length + // is zero or less; the scan cannot advance past it. + recordSize := NeedleHeaderSize + rest + if recordSize <= 0 { + return fmt.Errorf("%s: needle header at offset %d has size %d, record length %d: %w", datBackend.Name(), offset, n.Size, recordSize, needle.ErrorCorrupted) + } + offset += recordSize glog.V(4).Infof("==> new entry offset %d", offset) if n, nh, rest, err = needle.ReadNeedleHeader(datBackend, version, offset); err != nil { if err == io.EOF { diff --git a/weed/storage/volume_read_test.go b/weed/storage/volume_read_test.go index a31287598..9cc256c2d 100644 --- a/weed/storage/volume_read_test.go +++ b/weed/storage/volume_read_test.go @@ -2,8 +2,15 @@ package storage import ( "bytes" + "errors" + "fmt" + "os" + "path/filepath" + "reflect" "testing" + "time" + "github.com/seaweedfs/seaweedfs/weed/storage/backend" "github.com/seaweedfs/seaweedfs/weed/storage/needle" "github.com/seaweedfs/seaweedfs/weed/storage/super_block" "github.com/seaweedfs/seaweedfs/weed/storage/types" @@ -126,3 +133,112 @@ func TestReadNeedMetaWithDeletesThenWrites(t *testing.T) { expectedLastUpdateTime += 2000 } } + +// scanRecorder records visited offsets and fails once a scan visits more +// records than the file holds. +type scanRecorder struct { + readBody bool + maxVisits int + offsets []int64 +} + +func (s *scanRecorder) VisitSuperBlock(super_block.SuperBlock) error { return nil } + +func (s *scanRecorder) ReadNeedleBody() bool { return s.readBody } + +func (s *scanRecorder) VisitNeedle(_ *needle.Needle, offset int64, _, _ []byte) error { + s.offsets = append(s.offsets, offset) + if len(s.offsets) > s.maxVisits { + return fmt.Errorf("visited %d records in a file of %d, offsets %v", len(s.offsets), s.maxVisits, s.offsets) + } + return nil +} + +var errDidNotReturn = errors.New("did not return in time") + +// runWithTimeout returns errDidNotReturn if fn does not finish within d. +func runWithTimeout(d time.Duration, fn func() error) error { + done := make(chan error, 1) + go func() { done <- fn() }() + select { + case err := <-done: + return err + case <-time.After(d): + return errDidNotReturn + } +} + +// A corrupt .dat header can make a record's length zero or negative; the scan +// must stop there with ErrorCorrupted rather than re-read the header or step +// back into the previous record. A negative size whose record length stays +// positive is stepped over as before. +func TestScanVolumeFileFrom_StopsAtRecordThatCannotAdvance(t *testing.T) { + cases := []struct { + version needle.Version + size types.Size + stops bool + }{ + {needle.Version3, -1, false}, // record length 32 + {needle.Version3, -36, true}, // record length 0 + {needle.Version3, -43, true}, // record length 0 + {needle.Version3, -44, true}, // record length -8 + {needle.Version3, -4096, true}, // record length -4056 + {needle.Version2, -1, false}, // record length 24 + {needle.Version2, -28, true}, // record length 0 + {needle.Version2, -35, true}, // record length 0 + {needle.Version2, -36, true}, // record length -8 + } + for _, tc := range cases { + recordLen := needle.GetActualSize(tc.size, tc.version) + if (recordLen <= 0) != tc.stops { + t.Fatalf("v%d size %d: record length %d, case expects stops=%v", tc.version, tc.size, recordLen, tc.stops) + } + for _, readBody := range []bool{false, true} { + t.Run(fmt.Sprintf("v%d/size%d/readBody=%v", tc.version, tc.size, readBody), func(t *testing.T) { + f, err := os.Create(filepath.Join(t.TempDir(), "1.dat")) + if err != nil { + t.Fatalf("create dat: %v", err) + } + dat := backend.NewDiskFile(f) + defer dat.Close() + + first, _, _, err := newRandomNeedle(1).Append(dat, tc.version) + if err != nil { + t.Fatalf("append needle 1: %v", err) + } + corruptAt, _, err := dat.GetStat() + if err != nil { + t.Fatalf("stat dat: %v", err) + } + raw := make([]byte, max(recordLen, types.NeedleHeaderSize)) + types.NeedleIdToBytes(raw[types.CookieSize:types.CookieSize+types.NeedleIdSize], types.Uint64ToNeedleId(99)) + types.SizeToBytes(raw[types.CookieSize+types.NeedleIdSize:types.NeedleHeaderSize], tc.size) + if _, err := dat.WriteAt(raw, corruptAt); err != nil { + t.Fatalf("append corrupt record: %v", err) + } + second, _, _, err := newRandomNeedle(2).Append(dat, tc.version) + if err != nil { + t.Fatalf("append needle 2: %v", err) + } + + scanner := &scanRecorder{readBody: readBody, maxVisits: 3} + err = runWithTimeout(10*time.Second, func() error { + return ScanVolumeFileFrom(tc.version, dat, 0, scanner) + }) + + want := []int64{int64(first), corruptAt, int64(second)} + if tc.stops { + want = want[:2] + if !errors.Is(err, needle.ErrorCorrupted) { + t.Errorf("scan error = %v, want one wrapping ErrorCorrupted", err) + } + } else if err != nil { + t.Errorf("scan error = %v, want nil", err) + } + if !reflect.DeepEqual(scanner.offsets, want) { + t.Errorf("visited offsets %v, want %v", scanner.offsets, want) + } + }) + } + } +} diff --git a/weed/storage/volume_vacuum_test.go b/weed/storage/volume_vacuum_test.go index 6026b3a91..ec2a574c9 100644 --- a/weed/storage/volume_vacuum_test.go +++ b/weed/storage/volume_vacuum_test.go @@ -1,6 +1,9 @@ package storage import ( + "bytes" + "errors" + "fmt" "math/rand" "os" "path/filepath" @@ -470,6 +473,74 @@ func TestCompactByVolumeData_SkipsNegativeSizeHeader(t *testing.T) { } } +// A corrupt header the scan cannot advance past must fail compaction and +// leave the original volume untouched. +func TestCompactByVolumeData_FailsOnHeaderThatCannotAdvance(t *testing.T) { + for _, size := range []types.Size{-36, -100} { + t.Run(fmt.Sprintf("size%d", size), func(t *testing.T) { + dir := t.TempDir() + + v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0) + if err != nil { + t.Fatalf("volume creation: %v", err) + } + stuck := false + defer func() { + if !stuck { + v.Close() // blocks until a running compaction ends + } + }() + + if _, _, _, err := v.writeNeedle2(newRandomNeedle(1), true, false, false); err != nil { + t.Fatalf("write needle 1: %v", err) + } + datSize, _, err := v.DataBackend.GetStat() + if err != nil { + t.Fatalf("stat .dat: %v", err) + } + corrupt := make([]byte, types.NeedleHeaderSize) + types.NeedleIdToBytes(corrupt[types.CookieSize:types.CookieSize+types.NeedleIdSize], types.Uint64ToNeedleId(99)) + types.SizeToBytes(corrupt[types.CookieSize+types.NeedleIdSize:], size) + if _, err := v.DataBackend.WriteAt(corrupt, datSize); err != nil { + t.Fatalf("append corrupt header: %v", err) + } + if _, _, _, err := v.writeNeedle2(newRandomNeedle(2), true, false, false); err != nil { + t.Fatalf("write needle 2: %v", err) + } + datBefore, err := os.ReadFile(v.FileName(".dat")) + if err != nil { + t.Fatalf("read .dat: %v", err) + } + + err = runWithTimeout(10*time.Second, func() error { return v.CompactByVolumeData(nil) }) + if errors.Is(err, errDidNotReturn) { + stuck = true + t.Fatalf("CompactByVolumeData: %v", err) + } + if !errors.Is(err, needle.ErrorCorrupted) { + t.Fatalf("CompactByVolumeData error = %v, want one wrapping ErrorCorrupted", err) + } + + if _, err := os.Stat(v.FileName(".cpx")); !os.IsNotExist(err) { + t.Errorf("failed compaction left a compacted index: %v", err) + } + datAfter, err := os.ReadFile(v.FileName(".dat")) + if err != nil { + t.Fatalf("read .dat: %v", err) + } + if !bytes.Equal(datBefore, datAfter) { + t.Errorf(".dat changed: %d bytes before, %d after", len(datBefore), len(datAfter)) + } + for _, id := range []uint64{1, 2} { + n := &needle.Needle{Id: types.Uint64ToNeedleId(id)} + if _, err := v.readNeedle(n, &ReadOption{}, nil); err != nil { + t.Errorf("read needle %d from the original volume: %v", id, err) + } + } + }) + } +} + // TestExceedsExpectedCompactedSize guards the copy-phase integrity check // against regressing into double-subtracting skipped bytes: expectedLiveBytes // already excludes needles dropped as unreadable (they return before being