mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-20 06:54:24 +00:00
fix(volume): stop ScanVolumeFileFrom at a header it cannot advance past (#11398)
fix(volume): stop scans at a header they cannot advance past A corrupt .dat header with a very negative size gives a record length (NeedleHeaderSize + NeedleBodyLength) of zero or less: v3 sizes -43..-36 and v2 sizes -35..-28 give exactly zero, and smaller sizes give a negative length. ScanVolumeFileFrom advanced by that length, so it re-read the same header forever or stepped back into the record before it. weed fix, weed export, weed compact, incremental weed backup and the tail sender behind volume.move and volume.merge could hang on such a volume, and weed compact could also finish with a .cpx that had dropped every needle after the header. Return an error wrapping needle.ErrorCorrupted instead. The check runs after the visitor has seen the record, so the rebuild scanner still stops quietly with io.EOF. Smaller negative sizes whose record length is positive are still stepped over, preserving the salvage behavior compaction relies on. Mirror the guard into the Rust volume scans: DatScanPlan::scan and read_all_needles fail on a non-positive record length, as does scan_dat_head, so a corrupt header cannot stall a tail pass or leave the repair scan walking stale offsets.
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -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:?}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user