filer: bound metadata log flush retries during shutdown (#11596)

* filer: bound metadata log flush retries during shutdown

On SIGTERM the filer could hang forever in Shutdown: the final
LocalMetaLogBuffer flush retries appendToFile indefinitely, and with
the master already down each AssignVolume attempt just kept failing.
WaitForShutdown never returned, the interrupt hook never reached
os.Exit, and the half-dead filer kept its ports bound.

Thread a context through appendToFile/assignAndUpload and switch to
a 15s-bounded context once the filer is stopping: the flush abandons
with a log line instead of retrying forever. Normal operation keeps
the unbounded retry so no metadata is dropped while running.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: share one shutdown deadline across all pending meta log flushes

Review feedback on the per-flush timeout: a deadline armed at flush start
could already be expired when shutdown arrived, the first append of a flush
still ran unbounded, each queued window got a fresh budget (16 windows *
15s), and an abandoned window still advanced the flushed watermark as if it
had landed.

Rework to a single shared flush context on the Filer, cancelled once by
Shutdown via AfterFunc. Every append - the in-flight one and every queued
window - observes the same deadline, so the whole drain is bounded at 15s.
A flush that gives up reports its dropped bytes through the new
LogBuffer.NoteFlushDropped, and loopFlush then skips the offset/timestamp
advance and subscriber notifications for that window.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: bound append attempts and commit uploaded pieces detached

Store operations check then drop request cancellation, so a stalled
backend could still hold flushFn past the shutdown deadline; run each
append attempt on its own goroutine and give up on it at the deadline.
Once a piece is uploaded, commit its entry on a detached context so the
expired deadline cannot strand the chunk.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: keep shutdown flushes synchronous

Detaching the append attempt let flushFn return while the goroutine
still held the pooled flush buffer and could commit after the metadata
store closed; abandonment is only safe for the cancelable assign/upload
phase, which the shared flush context already bounds.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Chris LuandDevin authored and GitHub committed 2026-10-06 09:51:02 +08:00
1 parent 90eb4ec091
commit 076fd24186
5 files changed
+115 -16

No files matched your search

+14
View File
@@ -110,6 +110,12 @@ type Filer struct {
// overwriting the unread ledger with a partial set.
deletionLedgerBlocked atomic.Bool
deletionLedgerFlush chan struct{}
// flushCtx is the context metadata-log flushes append under. It stays
// live for the process lifetime and is cancelled once during Shutdown,
// bounding every retry - in-flight and queued alike - to one deadline.
// Nil for Filer literals built without NewFiler.
flushCtx context.Context
flushCancel context.CancelFunc
}
func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerHost pb.ServerAddress, filerGroup string, collection string, replication string, dataCenter string, maxFilenameLength uint32, notifyFn func()) *Filer {
@@ -128,6 +134,7 @@ func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerH
persistedLogCache: newPersistedLogCache(persistedLogCacheMaxBytes),
remoteTombstones: newRemoteDeletionTombstones(),
}
f.flushCtx, f.flushCancel = context.WithCancel(context.Background())
if f.UniqueFilerId < 0 {
f.UniqueFilerId = -f.UniqueFilerId
}
@@ -808,6 +815,13 @@ func (f *Filer) Shutdown() {
if f.EmptyFolderCleaner != nil {
f.EmptyFolderCleaner.Stop()
}
// Bound the remaining flush retries: with the cluster already gone each
// append keeps failing, and an unbounded retry would hold the shutdown
// drain open indefinitely. One shared deadline covers every in-flight and
// queued window; the flush loop drops what it cannot write.
if f.flushCancel != nil {
time.AfterFunc(shutdownMetadataLogFlushBudget, f.flushCancel)
}
f.LocalMetaLogBuffer.ShutdownLogBuffer()
// The final metadata-log flush still needs the store to append its entry.
f.LocalMetaLogBuffer.WaitForShutdown()
+20 -2
View File
@@ -247,6 +247,11 @@ func (f *Filer) logMetaEvent(ctx context.Context, event *filer_pb.SubscribeMetad
// in the rejection and volumeFileSizeLimit picks the real limit up from there.
const metadataLogUploadLimit = log_buffer.BufferSize
// shutdownMetadataLogFlushBudget bounds how long metadata-log flushes may keep
// retrying once the filer is shutting down; Shutdown cancels the shared flush
// context when it expires.
const shutdownMetadataLogFlushBudget = 15 * time.Second
var fileSizeLimitPattern = regexp.MustCompile(`file over the limited (\d+) bytes`)
// volumeFileSizeLimit reads the byte limit back out of a volume server's size
@@ -278,17 +283,30 @@ func (f *Filer) logFlushFunc(logBuffer *log_buffer.LogBuffer, startTime, stopTim
// One piece at a time, each retried on its own so a partial success is not
// replayed, and the piece size follows the limit the volume servers report.
// While the filer keeps running the retry is unbounded so no metadata is
// dropped; once Shutdown arms the flush deadline the shared context cancels
// and the flush drops what is left instead of holding shutdown open.
limit := metadataLogUploadLimit
ctx := f.flushCtx
if ctx == nil {
ctx = context.Background()
}
for len(buf) > 0 {
piece := nextLogPiece(buf, limit)
if err := f.appendToFile(targetFile, piece); err != nil {
if err := f.appendToFile(ctx, targetFile, piece); err != nil {
glog.V(0).Infof("metadata log write failed %s: %v", targetFile, err)
if reported := volumeFileSizeLimit(err); reported > 0 && reported < limit {
glog.V(0).Infof("metadata log upload limit lowered to %d bytes", reported)
limit = reported
continue
}
time.Sleep(737 * time.Millisecond)
select {
case <-ctx.Done():
logBuffer.NoteFlushDropped(len(buf))
glog.V(0).Infof("metadata log flush abandoned during shutdown, %d bytes left for %s: %v", len(buf), targetFile, ctx.Err())
return
case <-time.After(737 * time.Millisecond):
}
continue
}
buf = buf[len(piece):]
+11 -7
View File
@@ -11,16 +11,20 @@ import (
"github.com/seaweedfs/seaweedfs/weed/util"
)
func (f *Filer) appendToFile(targetFile string, data []byte) error {
func (f *Filer) appendToFile(ctx context.Context, targetFile string, data []byte) error {
assignResult, uploadResult, err2 := f.assignAndUpload(targetFile, data)
assignResult, uploadResult, err2 := f.assignAndUpload(ctx, targetFile, data)
if err2 != nil {
return err2
}
// The piece is already uploaded; commit it on a detached context so an
// expired shutdown deadline does not strand the chunk.
ctx = context.WithoutCancel(ctx)
// find out existing entry
fullpath := util.FullPath(targetFile)
entry, err := f.FindEntry(context.Background(), fullpath)
entry, err := f.FindEntry(ctx, fullpath)
var offset int64 = 0
if err == filer_pb.ErrNotFound {
entry = &Entry{
@@ -43,7 +47,7 @@ func (f *Filer) appendToFile(targetFile string, data []byte) error {
entry.Chunks = append(entry.GetChunks(), uploadResult.ToPbFileChunk(assignResult.Fid, offset, time.Now().UnixNano()))
// update the entry
err = f.CreateEntry(context.Background(), entry, nil, false, false, nil, false, f.MaxFilenameLength)
err = f.CreateEntry(ctx, entry, nil, false, false, nil, false, f.MaxFilenameLength)
return err
}
@@ -70,7 +74,7 @@ func (f *Filer) metaLogReplicationFor(ruleReplication string) string {
return util.Nvl(f.metaLogTargetReplication, f.metaLogReplication, ruleReplication)
}
func (f *Filer) assignAndUpload(targetFile string, data []byte) (*operation.AssignResult, *operation.UploadResult, error) {
func (f *Filer) assignAndUpload(ctx context.Context, targetFile string, data []byte) (*operation.AssignResult, *operation.UploadResult, error) {
// assign a volume location
diskType, rule := f.resolveMetadataLogAssignDiskType(targetFile)
assignRequest := &operation.VolumeAssignRequest{
@@ -82,7 +86,7 @@ func (f *Filer) assignAndUpload(targetFile string, data []byte) (*operation.Assi
ExpectedDataSize: uint64(len(data)),
}
assignResult, err := operation.Assign(context.Background(), f.GetMaster, f.GrpcDialOption, assignRequest)
assignResult, err := operation.Assign(ctx, f.GetMaster, f.GrpcDialOption, assignRequest)
if err != nil {
return nil, nil, fmt.Errorf("AssignVolume: %w", err)
}
@@ -107,7 +111,7 @@ func (f *Filer) assignAndUpload(targetFile string, data []byte) (*operation.Assi
return nil, nil, fmt.Errorf("upload data %s: %v", targetUrl, err)
}
uploadResult, err := uploader.UploadData(context.Background(), data, uploadOption)
uploadResult, err := uploader.UploadData(ctx, data, uploadOption)
if err != nil {
return nil, nil, fmt.Errorf("upload data %s: %v", targetUrl, err)
}
+43
View File
@@ -1,6 +1,8 @@
package filer
import (
"context"
"sync/atomic"
"testing"
"testing/synctest"
"time"
@@ -62,3 +64,44 @@ func TestShutdownKeepsStoreOpenUntilMetadataIsFlushed(t *testing.T) {
}
})
}
func TestShutdownBoundsBlockedMetadataFlush(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
store := &shutdownStore{closed: make(chan struct{})}
f := &Filer{Store: store, deletionQuit: make(chan struct{})}
f.flushCtx, f.flushCancel = context.WithCancel(context.Background())
var flushStarted atomic.Bool
lb := log_buffer.NewLogBuffer("blocked flush", time.Hour,
func(lb *log_buffer.LogBuffer, _, _ time.Time, buf []byte, _, _ int64) {
flushStarted.Store(true)
// A dead cluster makes append retries hang here; the shutdown
// deadline cancels the shared flush context to unblock it.
<-f.flushCtx.Done()
lb.NoteFlushDropped(len(buf))
}, nil, nil)
f.LocalMetaLogBuffer = lb
if err := lb.AddDataToBuffer(nil, []byte("last metadata event"), 0); err != nil {
t.Fatal(err)
}
done := make(chan struct{})
go func() {
f.Shutdown()
close(done)
}()
<-done
if !flushStarted.Load() {
t.Error("shutdown finished without running the pending flush")
}
if ts := lb.GetLastFlushTsNs(); ts != 0 {
t.Errorf("dropped flush advanced the flushed watermark to %d", ts)
}
select {
case <-store.closed:
default:
t.Error("metadata store was not closed after the bounded wait")
}
})
}
+27 -7
View File
@@ -195,13 +195,17 @@ type LogBuffer struct {
// Notified only when a flush lands, for readers that cannot act on an append
flushSubscribers map[string]*subscription
isStopping *atomic.Bool
shutdownCh chan struct{} // closed by ShutdownLogBuffer to wake blocked subscribers
loopsDone sync.WaitGroup // loopFlush and loopInterval signal exit
isAllFlushed bool
flushChan chan *dataToFlush
flushBudget *flushBudget
flushSeq uint64 // seal counter, assigned under the write lock
writeMu sync.Mutex // serializes sealing and enqueueing, including shutdown
shutdownCh chan struct{} // closed by ShutdownLogBuffer to wake blocked subscribers
// flushDroppedBytes counts bytes the flushFn gave up on during shutdown.
// loopFlush reads it after each window so an abandoned flush does not
// advance the flushed watermark or wake subscribers as if it had landed.
flushDroppedBytes atomic.Int64
loopsDone sync.WaitGroup // loopFlush and loopInterval signal exit
isAllFlushed bool
flushChan chan *dataToFlush
flushBudget *flushBudget
flushSeq uint64 // seal counter, assigned under the write lock
writeMu sync.Mutex // serializes sealing and enqueueing, including shutdown
// Offset range tracking for Kafka integration
hasOffsets bool
// Disk chunk cache for historical data reads
@@ -755,6 +759,14 @@ func (logBuffer *LogBuffer) WaitForShutdown() {
logBuffer.loopsDone.Wait()
}
// NoteFlushDropped records that the flush function abandoned n bytes instead
// of writing them (a shutdown-bounded retry giving up). loopFlush then skips
// the watermark advance and subscriber notifications for that window, so no
// reader is told the bytes were flushed when they were not.
func (logBuffer *LogBuffer) NoteFlushDropped(n int) {
logBuffer.flushDroppedBytes.Add(int64(n))
}
// IsAllFlushed returns true if all data in the buffer has been flushed, after calling ShutdownLogBuffer().
func (logBuffer *LogBuffer) IsAllFlushed() bool {
return logBuffer.isAllFlushed
@@ -775,8 +787,16 @@ func (logBuffer *LogBuffer) loopFlush() {
break // shutdown sentinel
}
logBuffer.flushFn(logBuffer, d.startTime, d.stopTime, d.data, d.minOffset, d.maxOffset)
dropped := logBuffer.flushDroppedBytes.Swap(0)
d.releaseMemory()
logBuffer.flushBudget.release(d.budget)
if dropped > 0 {
glog.Warningf("log buffer %s: shutdown flush dropped %d bytes", logBuffer.name, dropped)
if d.done != nil {
close(d.done)
}
continue
}
// local logbuffer is different from aggregate logbuffer here
if d.maxOffset >= 0 {
logBuffer.lastFlushedOffset.Store(d.maxOffset)