Files
seaweedfs/weed/storage/blockvol/flusher.go
T
pingqiuandClaude Opus 4.6 ebe95b6e2e fix: flusher OOM on multi-block writes + testrunner enhancements
Bug: flusher.go:336 allocated make([]byte, entryLen) per dirty block
instead of per unique WAL entry. A 4MB WriteLBA creates 1024 dirty map
entries (one per 4KB block), all sharing the same WAL offset. The flusher
read the full 4MB WAL entry 1024 times into separate buffers:
1024 × 4MB = 4GB per 4MB write → OOM on mkfs.ext4.

Root cause: flusher assumed 1:1 dirty-block-to-WAL-entry mapping.
WriteLBA supports multi-block writes but the flusher never deduplicated
shared WAL offsets.

Fix: deduplicate WAL reads by WalOffset in flushOnceLocked(). Multiple
dirty blocks from the same WAL entry share one read buffer and one
DecodeWALEntry call. Memory: O(WAL_entries × size) not O(blocks × size).
For a 4MB write: 4GB → 4MB.

Verified on hardware (m01/M02 25Gbps RoCE):
- Before: mkfs.ext4 → VS RSS 100MB→25GB → OOM killed
- After: mkfs.ext4 → VS RSS 129MB stable, mkfs succeeds
- pgbench TPC-B c=4: 1,248 TPS (RF=1, previously blocked by OOM)

Tests added:
- flusher_test.go: flush_multiblock_shared_wal_read (16 blocks share
  one WAL offset, flush dedup verified)
- flusher_test.go: flush_multiblock_data_correct (3 mixed multi-block
  writes, all data correct after flush)
- test/component/large_write_test.go: 7 component tests (single 4MB,
  sequential mkfs sim, concurrent, mixed sizes, production volume,
  flusher throughput 30s sustained)
- iscsi/large_write_mem_test.go: 2 iSCSI session memory tests (4MB
  R2T flow, slow device)

Testrunner enhancements (same commit — all tested on hardware):
- discover_primary action: maps primary IP → topology node name,
  supports alt_ips for multi-NIC (RoCE + management)
- NodeSpec.AltIPs field for multi-NIC node identification
- 5 new YAML scenarios: ec3, ec5, degraded sync_all/best_effort, pgbench
- All 13 hardware-verified scenarios PASS

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-04-02 14:24:10 -07:00

536 lines
15 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package blockvol
import (
"encoding/binary"
"fmt"
"log"
"os"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol/batchio"
)
// Flusher copies WAL entries to the extent region and frees WAL space.
// It runs as a background goroutine and can also be triggered manually.
type Flusher struct {
fd *os.File
super *Superblock
superMu *sync.Mutex // serializes superblock writes (shared with group commit)
wal *WALWriter
retentionFloorFn func() (uint64, bool) // CP13-6
evaluateRetentionBudgetsFn func() // CP13-6
dirtyMap *DirtyMap
walOffset uint64 // absolute file offset of WAL region
walSize uint64
blockSize uint32
extentStart uint64 // absolute file offset of extent region
mu sync.Mutex
checkpointLSN uint64 // last flushed LSN
checkpointTail uint64 // WAL physical tail after last flush
// flushMu serializes FlushOnce calls and is acquired by CreateSnapshot
// to pause the flusher while the snapshot is being set up.
// Lock order: flushMu -> snapMu (flushMu acquired first).
flushMu sync.Mutex
// snapMu protects the snapshots slice. Acquired under flushMu in
// FlushOnce (RLock) and under flushMu in PauseAndFlush callers.
snapMu sync.RWMutex
snapshots []*activeSnapshot
bio batchio.BatchIO // batch I/O backend (default: standard sequential)
logger *log.Logger
lastErr bool // true if last FlushOnce returned error
metrics *EngineMetrics
interval time.Duration
notifyCh chan struct{}
stopCh chan struct{}
done chan struct{}
stopOnce sync.Once
}
// FlusherConfig configures the flusher.
type FlusherConfig struct {
FD *os.File
Super *Superblock
SuperMu *sync.Mutex // serializes superblock writes (shared with group commit)
WAL *WALWriter
DirtyMap *DirtyMap
Interval time.Duration // default 100ms
Logger *log.Logger // optional; defaults to log.Default()
Metrics *EngineMetrics // optional; if nil, no metrics recorded
BatchIO batchio.BatchIO // optional; defaults to batchio.NewStandard()
// CP13-6: replica-aware WAL retention.
RetentionFloorFn func() (floorLSN uint64, hasFloor bool) // nil = no replica hold
EvaluateRetentionBudgetsFn func() // nil = no budget evaluation
}
// NewFlusher creates a flusher. Call Run() in a goroutine.
func NewFlusher(cfg FlusherConfig) *Flusher {
if cfg.Interval == 0 {
cfg.Interval = 100 * time.Millisecond
}
if cfg.Logger == nil {
cfg.Logger = log.Default()
}
if cfg.BatchIO == nil {
cfg.BatchIO = batchio.NewStandard()
}
return &Flusher{
fd: cfg.FD,
super: cfg.Super,
superMu: cfg.SuperMu,
wal: cfg.WAL,
dirtyMap: cfg.DirtyMap,
walOffset: cfg.Super.WALOffset,
walSize: cfg.Super.WALSize,
blockSize: cfg.Super.BlockSize,
extentStart: cfg.Super.WALOffset + cfg.Super.WALSize,
bio: cfg.BatchIO,
logger: cfg.Logger,
metrics: cfg.Metrics,
checkpointLSN: cfg.Super.WALCheckpointLSN,
checkpointTail: 0,
interval: cfg.Interval,
notifyCh: make(chan struct{}, 1),
stopCh: make(chan struct{}),
done: make(chan struct{}),
retentionFloorFn: cfg.RetentionFloorFn,
evaluateRetentionBudgetsFn: cfg.EvaluateRetentionBudgetsFn,
}
}
// Run is the flusher main loop. Call in a goroutine.
func (f *Flusher) Run() {
defer close(f.done)
ticker := time.NewTicker(f.interval)
defer ticker.Stop()
for {
select {
case <-f.stopCh:
return
case <-ticker.C:
if err := f.FlushOnce(); err != nil {
if !f.lastErr {
f.logger.Printf("flusher error: %v", err)
}
f.lastErr = true
if f.metrics != nil {
f.metrics.RecordFlusherError()
}
} else {
f.lastErr = false
}
case <-f.notifyCh:
if err := f.FlushOnce(); err != nil {
if !f.lastErr {
f.logger.Printf("flusher error: %v", err)
}
f.lastErr = true
if f.metrics != nil {
f.metrics.RecordFlusherError()
}
} else {
f.lastErr = false
}
}
}
}
// Notify wakes up the flusher for an immediate flush cycle.
func (f *Flusher) Notify() {
select {
case f.notifyCh <- struct{}{}:
default:
}
}
// NotifyUrgent wakes the flusher for an urgent flush (WAL pressure).
// Phase 3 MVP: delegates to Notify(). Future: may use a priority channel.
func (f *Flusher) NotifyUrgent() {
f.Notify()
}
// Stop shuts down the flusher. Safe to call multiple times.
func (f *Flusher) Stop() {
f.stopOnce.Do(func() {
close(f.stopCh)
})
<-f.done
}
// AddSnapshot adds a snapshot to the flusher's active list.
func (f *Flusher) AddSnapshot(snap *activeSnapshot) {
f.snapMu.Lock()
f.snapshots = append(f.snapshots, snap)
f.snapMu.Unlock()
}
// RemoveSnapshot removes a snapshot from the flusher's active list by ID.
func (f *Flusher) RemoveSnapshot(id uint32) {
f.snapMu.Lock()
for i, s := range f.snapshots {
if s.id == id {
f.snapshots = append(f.snapshots[:i], f.snapshots[i+1:]...)
break
}
}
f.snapMu.Unlock()
}
// HasActiveSnapshots returns true if there are active snapshots needing CoW.
func (f *Flusher) HasActiveSnapshots() bool {
f.snapMu.RLock()
n := len(f.snapshots)
f.snapMu.RUnlock()
return n > 0
}
// PauseAndFlush acquires flushMu (pausing the flusher), then runs FlushOnce.
// The caller must call Resume() when done.
func (f *Flusher) PauseAndFlush() error {
f.flushMu.Lock()
return f.flushOnceLocked()
}
// Pause acquires flushMu, pausing the flusher without flushing.
// The caller must call Resume() when done.
// Used by TruncateToLSN to prevent the flusher from flushing ahead
// entries while the dirty map and WAL are being cleared.
func (f *Flusher) Pause() {
f.flushMu.Lock()
}
// Resume releases flushMu, allowing the flusher to resume.
func (f *Flusher) Resume() {
f.flushMu.Unlock()
}
// FlushOnce performs a single flush cycle: scan dirty map, CoW for active
// snapshots, copy data to extent region, fsync, update checkpoint, advance WAL tail.
func (f *Flusher) FlushOnce() error {
f.flushMu.Lock()
defer f.flushMu.Unlock()
return f.flushOnceLocked()
}
// flushOnceLocked is the inner FlushOnce. Caller must hold flushMu.
func (f *Flusher) flushOnceLocked() error {
flushStart := time.Now()
entries := f.dirtyMap.Snapshot()
if len(entries) == 0 {
return nil
}
// --- Phase 1: CoW for active snapshots ---
f.snapMu.RLock()
snaps := make([]*activeSnapshot, len(f.snapshots))
copy(snaps, f.snapshots)
f.snapMu.RUnlock()
if len(snaps) > 0 {
cowDirty := false
for _, e := range entries {
for _, snap := range snaps {
if !snap.bitmap.Get(e.Lba) {
// Read old data from extent (pre-modification state).
oldData := make([]byte, f.blockSize)
extentOff := int64(f.extentStart + e.Lba*uint64(f.blockSize))
if _, err := f.fd.ReadAt(oldData, extentOff); err != nil {
return fmt.Errorf("flusher: CoW read extent LBA %d: %w", e.Lba, err)
}
// Write old data to delta file.
deltaOff := int64(snap.dataOffset + e.Lba*uint64(f.blockSize))
if _, err := snap.fd.WriteAt(oldData, deltaOff); err != nil {
return fmt.Errorf("flusher: CoW write delta LBA %d snap %d: %w", e.Lba, snap.id, err)
}
snap.bitmap.Set(e.Lba)
snap.dirty = true
cowDirty = true
}
}
}
if cowDirty {
// Crash safety: delta data -> fsync -> bitmap persist -> fsync -> extent write.
for _, snap := range snaps {
if snap.dirty {
if err := snap.fd.Sync(); err != nil {
return fmt.Errorf("flusher: fsync delta snap %d: %w", snap.id, err)
}
if err := snap.bitmap.WriteTo(snap.fd, SnapHeaderSize); err != nil {
return fmt.Errorf("flusher: persist bitmap snap %d: %w", snap.id, err)
}
if err := snap.fd.Sync(); err != nil {
return fmt.Errorf("flusher: fsync bitmap snap %d: %w", snap.id, err)
}
snap.dirty = false
}
}
}
}
// --- Phase 2: Extent writes via BatchIO ---
var maxLSN uint64
var maxWALEnd uint64
// Step 2a: Batch-read WAL headers.
headerOps := make([]batchio.Op, len(entries))
for i, e := range entries {
headerOps[i] = batchio.Op{
Buf: make([]byte, walEntryHeaderSize),
Offset: int64(f.walOffset + e.WalOffset),
}
}
if err := f.bio.PreadBatch(f.fd, headerOps); err != nil {
return fmt.Errorf("flusher: batch read WAL headers: %w", err)
}
// Step 2b: Identify entries needing full WAL read, batch-read them.
type pendingEntry struct {
idx int // index into entries
entryType uint8
entryLen int
}
var pending []pendingEntry
for i, e := range entries {
hdr := headerOps[i].Buf
entryLSN := binary.LittleEndian.Uint64(hdr[0:8])
if entryLSN != e.Lsn {
continue // stale — WAL slot reused
}
entryType := hdr[16]
if entryType == EntryTypeWrite {
dataLen := parseLength(hdr)
if dataLen > 0 {
pending = append(pending, pendingEntry{
idx: i,
entryType: entryType,
entryLen: walEntryHeaderSize + int(dataLen),
})
}
} else if entryType == EntryTypeTrim {
pending = append(pending, pendingEntry{
idx: i,
entryType: entryType,
entryLen: walEntryHeaderSize,
})
}
if e.Lsn > maxLSN {
maxLSN = e.Lsn
}
}
// Batch-read full WAL entries for write ops.
// Deduplicate by WalOffset: multiple dirty blocks from the same multi-block
// write share one WAL entry. Without dedup, a 4MB write (1024 blocks) would
// allocate 1024 × 4MB buffers = 4GB instead of 1 × 4MB.
type walReadInfo struct {
buf []byte
entry *WALEntry // decoded once, shared by all blocks in this WAL entry
}
walReadByOffset := make(map[uint64]*walReadInfo, len(pending))
var walReadOps []batchio.Op
for _, p := range pending {
if p.entryType != EntryTypeWrite {
continue
}
off := entries[p.idx].WalOffset
if _, ok := walReadByOffset[off]; !ok {
buf := make([]byte, p.entryLen)
walReadByOffset[off] = &walReadInfo{buf: buf}
walReadOps = append(walReadOps, batchio.Op{
Buf: buf,
Offset: int64(f.walOffset + off),
})
}
}
if len(walReadOps) > 0 {
if err := f.bio.PreadBatch(f.fd, walReadOps); err != nil {
return fmt.Errorf("flusher: batch read WAL entries: %w", err)
}
// Decode each WAL entry once.
for _, ri := range walReadByOffset {
decoded, err := DecodeWALEntry(ri.buf)
if err == nil {
ri.entry = &decoded
}
}
}
// Step 2c: Build extent write ops from decoded (shared) WAL entries.
var extentWriteOps []batchio.Op
for _, p := range pending {
e := entries[p.idx]
if p.entryType == EntryTypeWrite {
ri := walReadByOffset[e.WalOffset]
if ri == nil || ri.entry == nil {
continue
}
entry := ri.entry
if e.Lba < entry.LBA {
continue
}
blockIdx := e.Lba - entry.LBA
dataStart := blockIdx * uint64(f.blockSize)
if dataStart+uint64(f.blockSize) <= uint64(len(entry.Data)) {
extentOff := int64(f.extentStart + e.Lba*uint64(f.blockSize))
blockData := entry.Data[dataStart : dataStart+uint64(f.blockSize)]
extentWriteOps = append(extentWriteOps, batchio.Op{
Buf: blockData,
Offset: extentOff,
})
}
walEnd := e.WalOffset + uint64(p.entryLen)
if walEnd > maxWALEnd {
maxWALEnd = walEnd
}
} else if p.entryType == EntryTypeTrim {
zeroBlock := make([]byte, f.blockSize)
extentOff := int64(f.extentStart + e.Lba*uint64(f.blockSize))
extentWriteOps = append(extentWriteOps, batchio.Op{
Buf: zeroBlock,
Offset: extentOff,
})
walEnd := e.WalOffset + uint64(walEntryHeaderSize)
if walEnd > maxWALEnd {
maxWALEnd = walEnd
}
}
}
// Step 2d: Batch-write extents + fsync.
if len(extentWriteOps) > 0 {
if err := f.bio.PwriteBatch(f.fd, extentWriteOps); err != nil {
return fmt.Errorf("flusher: batch write extents: %w", err)
}
}
if err := f.bio.Fsync(f.fd); err != nil {
return fmt.Errorf("flusher: fsync extent: %w", err)
}
// Remove flushed entries from dirty map.
f.mu.Lock()
for _, e := range entries {
_, currentLSN, _, ok := f.dirtyMap.Get(e.Lba)
if ok && currentLSN == e.Lsn {
f.dirtyMap.Delete(e.Lba)
}
}
// CP13-6: replica-aware WAL retention.
// Evaluate retention budgets first (may escalate stale replicas).
if f.evaluateRetentionBudgetsFn != nil {
f.evaluateRetentionBudgetsFn()
}
// Check retention floor — hold WAL for recoverable replicas.
effectiveLSN := maxLSN
effectiveWALEnd := maxWALEnd
if f.retentionFloorFn != nil {
floorLSN, hasFloor := f.retentionFloorFn()
if hasFloor && maxLSN > floorLSN {
// Hold: don't advance checkpoint/tail past the floor.
// Extent writes are done (reads work), but WAL space is kept.
effectiveLSN = floorLSN
effectiveWALEnd = 0
}
}
f.checkpointLSN = effectiveLSN
f.checkpointTail = effectiveWALEnd
f.mu.Unlock()
// Advance WAL tail to free space (only if not held by retention).
if effectiveWALEnd > 0 {
f.wal.AdvanceTail(effectiveWALEnd)
}
// Update superblock checkpoint.
f.updateSuperblockCheckpoint(effectiveLSN, f.wal.Tail())
// Record metrics.
if f.metrics != nil {
bytesWritten := uint64(len(entries)) * uint64(f.blockSize)
f.metrics.RecordFlusherFlush(bytesWritten, time.Since(flushStart))
}
return nil
}
// updateSuperblockCheckpoint writes the updated checkpoint to disk.
// This is the primary place where WALHead is persisted to the superblock.
// Recovery's extended scan handles the gap between checkpoints.
func (f *Flusher) updateSuperblockCheckpoint(checkpointLSN uint64, walTail uint64) error {
f.superMu.Lock()
defer f.superMu.Unlock()
f.super.WALCheckpointLSN = checkpointLSN
f.super.WALHead = f.wal.LogicalHead()
f.super.WALTail = f.wal.LogicalTail()
if _, err := f.fd.Seek(0, 0); err != nil {
return fmt.Errorf("flusher: seek to superblock: %w", err)
}
if _, err := f.super.WriteTo(f.fd); err != nil {
return fmt.Errorf("flusher: write superblock: %w", err)
}
return f.fd.Sync()
}
// CheckpointLSN returns the last flushed LSN.
func (f *Flusher) CheckpointLSN() uint64 {
f.mu.Lock()
defer f.mu.Unlock()
return f.checkpointLSN
}
// SetCheckpointLSN updates the flusher's internal checkpoint state.
// Used after rebuild to sync flusher with the rebuilt superblock state.
func (f *Flusher) SetCheckpointLSN(lsn uint64) {
f.mu.Lock()
f.checkpointLSN = lsn
f.mu.Unlock()
}
// RetentionFloorFn returns the current retention floor function.
func (f *Flusher) RetentionFloorFn() func() (uint64, bool) {
return f.retentionFloorFn
}
// SetRetentionFloorFn replaces the retention floor function.
// Used by V2 bridge to chain additional retention holds.
func (f *Flusher) SetRetentionFloorFn(fn func() (uint64, bool)) {
f.retentionFloorFn = fn
}
// CloseBatchIO releases the batch I/O backend resources (e.g. io_uring ring).
// Must be called after Stop() and the final FlushOnce().
func (f *Flusher) CloseBatchIO() error {
if f.bio != nil {
return f.bio.Close()
}
return nil
}
// SetFD replaces the file descriptor used for extent writes. Test-only.
func (f *Flusher) SetFD(fd *os.File) {
f.mu.Lock()
defer f.mu.Unlock()
f.fd = fd
}
// parseLength extracts the Length field from a WAL entry header buffer.
func parseLength(headerBuf []byte) uint32 {
// Length at LSN(8)+Epoch(8)+Type(1)+Flags(1)+LBA(8) = 26
return binary.LittleEndian.Uint32(headerBuf[26:])
}