From 80fd3635d29af32359ca4eb7d3036ff3ca551788 Mon Sep 17 00:00:00 2001 From: Mohd Quamar Tyagi <104281681+Tyagiquamar@users.noreply.github.com> Date: Sat, 26 Sep 2026 17:11:59 +0530 Subject: [PATCH] volume: skip TTL last-write scan when it cannot fit its budget (#11472) * volume: skip TTL last-write scan when it cannot fit its budget * Update weed/storage/volume_checking.go Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --------- Co-authored-by: Chris Lu Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- weed/storage/volume_checking.go | 57 ++++++++++++++++++++++ weed/storage/volume_ttl_expiry_test.go | 63 +++++++++++++++++++++++++ weed/storage/volume_write_fsync_test.go | 10 +++- 3 files changed, 129 insertions(+), 1 deletion(-) diff --git a/weed/storage/volume_checking.go b/weed/storage/volume_checking.go index 89e17a61d..7d007f927 100644 --- a/weed/storage/volume_checking.go +++ b/weed/storage/volume_checking.go @@ -300,6 +300,9 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui } scanEveryWrite := v.SuperBlock.CompactionRevision > 0 entryBudget := vacuumedLastWriteScanEntries + if scanEveryWrite && !affordableVacuumedScan(indexFile, indexSize, v.Id, v.FileName(".dat")) { + return 0, nil + } var lastWriteAppendAtNs uint64 block := make([]byte, types.NeedleMapEntrySize*idx.RowsToRead) for end := indexSize; end > 0; { @@ -340,6 +343,60 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui return lastWriteAppendAtNs, nil } +// affordableVacuumedScan reports whether a vacuumed volume's recovery scan fits +// its budget. A vacuumed volume must take the maximum over every write it +// indexes, so a volume holding more live needles than vacuumedLastWriteScanEntries +// would spend up to two random .dat reads per needle and then give up, keeping +// the .dat mtime anyway. Counting the live .idx entries first is one short +// sequential pass over a small file, and it skips the scan before any .dat I/O +// when the scan provably cannot succeed. +func affordableVacuumedScan(indexFile *os.File, indexSize int64, volumeId needle.VolumeId, datFileName string) bool { + liveEntries, err := countLiveIndexEntries(indexFile, indexSize, vacuumedLastWriteScanEntries) + if err != nil { + glog.Warningf("volume %d count live entries in %s: %v", volumeId, indexFile.Name(), err) + return true + } + if liveEntries >= vacuumedLastWriteScanEntries { + glog.V(0).Infof("volume %d: more than %d needles to scan for its last write, keeping the %s mtime", + volumeId, vacuumedLastWriteScanEntries, datFileName) + return false + } + return true +} + +// countLiveIndexEntries counts the .idx entries describing a live needle, +// stopping early once the count passes limit. Tombstones and zero offsets are +// skipped exactly as findLastWriteAppendAtNs skips them. +func countLiveIndexEntries(indexFile *os.File, indexSize int64, limit int) (int, error) { + count := 0 + block := make([]byte, types.NeedleMapEntrySize*idx.RowsToRead) + for start := int64(0); start < indexSize; { + end := start + int64(len(block)) + if end > indexSize { + end = indexSize + } + entries := block[:end-start] + readCount, err := indexFile.ReadAt(entries, start) + if err == io.EOF && readCount == len(entries) { + err = nil + } + if err != nil { + return 0, fmt.Errorf("read %s at %d: %v", indexFile.Name(), start, err) + } + for i := 0; i+types.NeedleMapEntrySize <= len(entries); i += types.NeedleMapEntrySize { + _, offset, size := idx.IdxFileEntry(entries[i : i+types.NeedleMapEntrySize]) + if offset.IsZero() || size.IsDeleted() { + continue + } + if count++; count > limit { + return count, nil + } + } + start = end + } + return count, nil +} + // findNeedleOffset returns the .dat offset holding the needle an .idx entry // describes, or -1 when no needle there matches it. A .dat past // MaxPossibleVolumeSize wraps the 4-byte offsets in its .idx, so the needle can diff --git a/weed/storage/volume_ttl_expiry_test.go b/weed/storage/volume_ttl_expiry_test.go index 3fd689c99..3dd18dd5d 100644 --- a/weed/storage/volume_ttl_expiry_test.go +++ b/weed/storage/volume_ttl_expiry_test.go @@ -1,6 +1,7 @@ package storage import ( + "os" "testing" "time" @@ -199,6 +200,68 @@ func TestVolumeTtlClockDeclinesUnaffordableScan(t *testing.T) { } } +// TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads is the regression test +// for https://github.com/seaweedfs/seaweedfs/issues/11469: a vacuumed volume +// holding more live needles than the scan budget must decline recovery before +// doing any .dat I/O, instead of spending up to two random reads per needle +// and then keeping the mtime anyway. +func TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads(t *testing.T) { + dir := t.TempDir() + ttl, err := needle.ReadTTL("5m") + if err != nil { + t.Fatalf("read ttl: %v", err) + } + + v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0) + if err != nil { + t.Fatalf("volume creation: %v", err) + } + defer v.Close() + + for i := 1; i <= 5; i++ { + if _, _, _, err := v.writeNeedle2(newRandomNeedle(uint64(i)), true, false, false); err != nil { + t.Fatalf("write needle %d: %v", i, err) + } + } + if err := v.CompactByIndex(nil); err != nil { + t.Fatalf("compact: %v", err) + } + if err := v.CommitCompact(); err != nil { + t.Fatalf("commit compact: %v", err) + } + if v.SuperBlock.CompactionRevision == 0 { + t.Fatal("vacuum must bump CompactionRevision for this test to exercise the vacuumed path") + } + + indexFile, err := os.Open(v.FileName(".idx")) + if err != nil { + t.Fatalf("open .idx: %v", err) + } + defer indexFile.Close() + indexStat, err := indexFile.Stat() + if err != nil { + t.Fatalf("stat .idx: %v", err) + } + + var reads countingBackend + reads.BackendStorageFile = v.DataBackend + v.DataBackend = &reads + + defer func(budget int) { vacuumedLastWriteScanEntries = budget }(vacuumedLastWriteScanEntries) + vacuumedLastWriteScanEntries = 2 + + appendAtNs, err := findLastWriteAppendAtNs(v, indexFile, indexStat.Size()) + if err != nil { + t.Fatalf("recover last write: %v", err) + } + if appendAtNs != 0 { + t.Errorf("over-budget scan must decline with 0, got %d", appendAtNs) + } + if got := reads.readCount; got != 0 { + t.Errorf("over-budget scan did %d .dat reads before declining, want 0", got) + } +} + // TestVolumeExpireAtSecCountsFromLastWrite guards the destroy time an EC volume // is reclaimed on (erasure_coding.EcVolume.IsTimeToDestroy). It was recomputed // as now+TTL on every .vif write, so a read-only mark, a tier upload or an EC diff --git a/weed/storage/volume_write_fsync_test.go b/weed/storage/volume_write_fsync_test.go index 71c7884e3..f842923c7 100644 --- a/weed/storage/volume_write_fsync_test.go +++ b/weed/storage/volume_write_fsync_test.go @@ -11,7 +11,9 @@ import ( "github.com/stretchr/testify/require" ) -// countingBackend counts Sync calls and can be made to fail them. +// countingBackend counts Sync calls and can be made to fail them. It also +// counts random reads, so recovery paths can prove they declined before +// touching the .dat. type countingBackend struct { backend.BackendStorageFile syncCount int @@ -19,6 +21,12 @@ type countingBackend struct { syncErrOnce bool truncateErr error truncateCount int + readCount int +} + +func (b *countingBackend) ReadAt(p []byte, off int64) (int, error) { + b.readCount++ + return b.BackendStorageFile.ReadAt(p, off) } func (b *countingBackend) Sync() error {