diff --git a/weed/filer/filer.go b/weed/filer/filer.go index b2c09fcd1..a78d90691 100644 --- a/weed/filer/filer.go +++ b/weed/filer/filer.go @@ -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() diff --git a/weed/filer/filer_notify.go b/weed/filer/filer_notify.go index e75bc632b..6101715eb 100644 --- a/weed/filer/filer_notify.go +++ b/weed/filer/filer_notify.go @@ -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):] diff --git a/weed/filer/filer_notify_append.go b/weed/filer/filer_notify_append.go index c23c33e7f..24f4e733b 100644 --- a/weed/filer/filer_notify_append.go +++ b/weed/filer/filer_notify_append.go @@ -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) } diff --git a/weed/filer/filer_shutdown_test.go b/weed/filer/filer_shutdown_test.go index c01350d45..8383086d0 100644 --- a/weed/filer/filer_shutdown_test.go +++ b/weed/filer/filer_shutdown_test.go @@ -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") + } + }) +} diff --git a/weed/util/log_buffer/log_buffer.go b/weed/util/log_buffer/log_buffer.go index 2a5e9023f..6d71105d6 100644 --- a/weed/util/log_buffer/log_buffer.go +++ b/weed/util/log_buffer/log_buffer.go @@ -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)