diff --git a/weed/filer/filer.go b/weed/filer/filer.go index ef9b18f8b..ad3a186a0 100644 --- a/weed/filer/filer.go +++ b/weed/filer/filer.go @@ -727,6 +727,8 @@ func (f *Filer) Shutdown() { f.EmptyFolderCleaner.Stop() } f.LocalMetaLogBuffer.ShutdownLogBuffer() + // The final metadata-log flush still needs the store to append its entry. + f.LocalMetaLogBuffer.WaitForShutdown() f.Store.Shutdown() } diff --git a/weed/filer/filer_shutdown_test.go b/weed/filer/filer_shutdown_test.go new file mode 100644 index 000000000..c01350d45 --- /dev/null +++ b/weed/filer/filer_shutdown_test.go @@ -0,0 +1,64 @@ +package filer + +import ( + "testing" + "testing/synctest" + "time" + + "github.com/seaweedfs/seaweedfs/weed/util/log_buffer" +) + +type shutdownStore struct { + VirtualFilerStore + closed chan struct{} +} + +func (s *shutdownStore) Shutdown() { + close(s.closed) +} + +func TestShutdownKeepsStoreOpenUntilMetadataIsFlushed(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + store := &shutdownStore{closed: make(chan struct{})} + releaseFlush := make(chan struct{}) + flushed := false + lb := log_buffer.NewLogBuffer("filer shutdown", time.Hour, + func(_ *log_buffer.LogBuffer, _, _ time.Time, _ []byte, _, _ int64) { + <-releaseFlush + select { + case <-store.closed: + t.Error("metadata store closed before the pending log write") + default: + flushed = true + } + }, nil, nil) + f := &Filer{Store: store, LocalMetaLogBuffer: lb, deletionQuit: make(chan struct{})} + 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) + }() + synctest.Wait() + select { + case <-store.closed: + t.Error("metadata store closed while the log write was blocked") + default: + } + + close(releaseFlush) + <-done + lb.WaitForShutdown() + if !flushed { + t.Error("shutdown returned without persisting the pending metadata") + } + select { + case <-store.closed: + default: + t.Error("metadata store was not closed after the log write finished") + } + }) +} diff --git a/weed/filer/leveldb/object_size_metric_test.go b/weed/filer/leveldb/object_size_metric_test.go index 4e4d2c525..3400fdb84 100644 --- a/weed/filer/leveldb/object_size_metric_test.go +++ b/weed/filer/leveldb/object_size_metric_test.go @@ -34,7 +34,7 @@ func TestCreateEntryRecordsObjectSize(t *testing.T) { } testFiler.SetStore(store) defer testFiler.Shutdown() - ctx := context.Background() + ctx := filer.WithSuppressedMetadataEvents(context.Background()) file := func(path string, size uint64) *filer.Entry { return &filer.Entry{ diff --git a/weed/util/log_buffer/flush_budget_test.go b/weed/util/log_buffer/flush_budget_test.go index 89050b244..06a97b03e 100644 --- a/weed/util/log_buffer/flush_budget_test.go +++ b/weed/util/log_buffer/flush_budget_test.go @@ -149,9 +149,7 @@ func TestQueueFlushChargesTheSlabNotTheWindow(t *testing.T) { t.Skipf("allocator returned an exact fit (%d bytes), nothing to distinguish", cap(data)) } - if !lb.queueFlush(&dataToFlush{data: data}) { - t.Fatal("queueFlush dropped the window") - } + lb.queueFlush(&dataToFlush{data: data}) lb.flushBudget.mu.Lock() queued := lb.flushBudget.queued @@ -162,31 +160,6 @@ func TestQueueFlushChargesTheSlabNotTheWindow(t *testing.T) { } } -// The window is already sealed by the time queueFlush runs, so dropping it -// loses records the caller was told were accepted. A shutdown racing the -// hand-off must not cost data while the queue still has room for it. -func TestQueueFlushKeepsSealedWindowWhenQueueHasRoom(t *testing.T) { - stall := make(chan struct{}) - var flushed atomic.Int64 - lb := NewLogBuffer("shutdown-race", time.Hour, func(_ *LogBuffer, _, _ time.Time, _ []byte, _, _ int64) { - <-stall - flushed.Add(1) - }, nil, func() {}) - defer func() { close(stall) }() - - // Simulate a shutdown landing between the seal and the hand-off. - close(lb.shutdownCh) - lb.isStopping.Store(true) - - // Stay inside the channel's capacity so room is guaranteed for every one of - // them; loopFlush is parked in the stalled flushFn and drains nothing. - for i := 0; i < flushQueueDepth; i++ { - if !lb.queueFlush(&dataToFlush{data: mem.Allocate(1024), seq: uint64(i)}) { - t.Fatalf("window %d was dropped even though the queue had room", i) - } - } -} - // The point of the budget: a stalled flush must stop producers rather than let // them keep handing over copies. Asserted on how many producers get through, // which is independent of the counter reserve maintains. diff --git a/weed/util/log_buffer/log_buffer.go b/weed/util/log_buffer/log_buffer.go index 78235a5cc..2a5e9023f 100644 --- a/weed/util/log_buffer/log_buffer.go +++ b/weed/util/log_buffer/log_buffer.go @@ -48,6 +48,8 @@ const flushQueueBudget = flushQueueDepth * BufferSize var ( // ErrBufferCorrupted indicates the log buffer contains corrupted data ErrBufferCorrupted = fmt.Errorf("log buffer is corrupted") + // ErrBufferStopped indicates that shutdown has closed write admission. + ErrBufferStopped = fmt.Errorf("log buffer is stopping") ) type dataToFlush struct { @@ -198,7 +200,8 @@ type LogBuffer struct { isAllFlushed bool flushChan chan *dataToFlush flushBudget *flushBudget - flushSeq uint64 // seal counter, assigned under the write lock + 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 @@ -443,6 +446,9 @@ func (logBuffer *LogBuffer) AddLogEntryToBuffer(logEntry *filer_pb.LogEntry) err logBuffer.flushBudget.waitForRoom(len(logEntry.Data)) } + if !logBuffer.beginWrite() { + return ErrBufferStopped + } var toFlush *dataToFlush var marshalErr error logBuffer.Lock() @@ -451,6 +457,7 @@ func (logBuffer *LogBuffer) AddLogEntryToBuffer(logEntry *filer_pb.LogEntry) err if toFlush != nil { logBuffer.queueFlush(toFlush) } + logBuffer.writeMu.Unlock() // Only notify if there was no error if marshalErr == nil { if logBuffer.notifyFn != nil { @@ -570,6 +577,9 @@ func (logBuffer *LogBuffer) AddDataToBuffer(partitionKey, data []byte, processin Key: partitionKey, } + if !logBuffer.beginWrite() { + return ErrBufferStopped + } var toFlush *dataToFlush var marshalErr error logBuffer.Lock() @@ -578,6 +588,7 @@ func (logBuffer *LogBuffer) AddDataToBuffer(partitionKey, data []byte, processin if toFlush != nil { logBuffer.queueFlush(toFlush) } + logBuffer.writeMu.Unlock() // Only notify if there was no error if marshalErr == nil { if logBuffer.notifyFn != nil { @@ -676,25 +687,22 @@ func (logBuffer *LogBuffer) IsStopping() bool { return logBuffer.isStopping.Load() } -// ForceFlush immediately flushes the current buffer content and WAITS for completion +// ForceFlush queues the current buffer content, then waits up to 5 seconds for completion // This is useful for critical topics that need immediate persistence -// CRITICAL: This function is now SYNCHRONOUS - it blocks until the flush completes +// Queueing itself has no timeout. func (logBuffer *LogBuffer) ForceFlush() { - if logBuffer.isStopping.Load() { - return // Don't flush if we're shutting down + if !logBuffer.beginWrite() { + return } - logBuffer.Lock() toFlush := logBuffer.copyToFlushWithCallback() logBuffer.Unlock() + if toFlush != nil { + logBuffer.queueFlush(toFlush) + } + logBuffer.writeMu.Unlock() if toFlush != nil { - // The live buffer was already sealed and reset by copyToFlushWithCallback, - // so dropping toFlush on a timeout would lose it. Block until queued, - // bailing out only on shutdown. - if !logBuffer.queueFlush(toFlush) { - return - } select { case <-toFlush.done: // Flush completed @@ -704,6 +712,17 @@ func (logBuffer *LogBuffer) ForceFlush() { } } +// beginWrite serializes admission and batch handoff with shutdown. +// Callers hold writeMu until any sealed batch has been queued. +func (logBuffer *LogBuffer) beginWrite() bool { + logBuffer.writeMu.Lock() + if logBuffer.isStopping.Load() { + logBuffer.writeMu.Unlock() + return false + } + return true +} + // ShutdownLogBuffer flushes the buffer and stops the log buffer func (logBuffer *LogBuffer) ShutdownLogBuffer() { isAlreadyStopped := logBuffer.isStopping.Swap(true) @@ -714,51 +733,39 @@ func (logBuffer *LogBuffer) ShutdownLogBuffer() { // notice IsStopping() and exit promptly, even on an idle buffer where no // flush notification would otherwise fire. close(logBuffer.shutdownCh) - // Let go of the flush budget before sealing the last window, so a producer - // parked on it wakes up and the hand-off below cannot wait on a reservation. + // Wake oversized writers waiting for room; they will observe shutdown + // before appending. Already admitted writes must finish their handoff. logBuffer.flushBudget.close() + logBuffer.writeMu.Lock() + defer logBuffer.writeMu.Unlock() logBuffer.Lock() toFlush := logBuffer.copyToFlush() logBuffer.Unlock() if toFlush != nil { - toFlush.budget = logBuffer.flushBudget.reserve(toFlush.seq, cap(toFlush.data)) - logBuffer.flushChan <- toFlush + logBuffer.queueFlush(toFlush) } - // nil is the shutdown sentinel: loopFlush drains everything queued before - // it and exits. The channel is never closed, so a sender racing shutdown - // can never panic on a closed channel. + // Every accepted batch is now queued, and no producer can append after + // this sentinel. loopFlush drains the queue before exiting. logBuffer.flushChan <- nil } +// WaitForShutdown waits for pending flushes and background loops to finish. +// Call ShutdownLogBuffer first, after stopping producers. +func (logBuffer *LogBuffer) WaitForShutdown() { + logBuffer.loopsDone.Wait() +} + // IsAllFlushed returns true if all data in the buffer has been flushed, after calling ShutdownLogBuffer(). func (logBuffer *LogBuffer) IsAllFlushed() bool { return logBuffer.isAllFlushed } -// queueFlush hands a sealed window to loopFlush, reserving its bytes first so -// a producer waits for the queue to drain rather than adding another copy to -// it. Reports false when the buffer shut down before the hand-off. -func (logBuffer *LogBuffer) queueFlush(d *dataToFlush) bool { - // Charge the slab, not the window: mem.Allocate rounds up to a size class, - // so the bytes actually held are cap(data), and charging len would let the - // queue hold up to twice the ceiling. +// queueFlush hands a sealed window to loopFlush. The caller holds writeMu +// from before sealing until this handoff completes, so shutdown cannot pass it. +func (logBuffer *LogBuffer) queueFlush(d *dataToFlush) { + // Charge the pooled slab, whose capacity may exceed the window length. d.budget = logBuffer.flushBudget.reserve(d.seq, cap(d.data)) - // The window is already sealed, so dropping it here loses records the - // caller was told were accepted. Take any room in the queue first, and only - // fall back to the shutdown escape when there is none. - select { - case logBuffer.flushChan <- d: - return true - default: - } - select { - case logBuffer.flushChan <- d: - return true - case <-logBuffer.shutdownCh: - // shutting down; loopFlush may be gone, do not park forever - logBuffer.flushBudget.release(d.budget) - return false - } + logBuffer.flushChan <- d } func (logBuffer *LogBuffer) loopFlush() { @@ -813,12 +820,16 @@ func (logBuffer *LogBuffer) loopInterval() { case <-ticker.C: } + if !logBuffer.beginWrite() { + return + } logBuffer.Lock() toFlush := logBuffer.copyToFlush() logBuffer.Unlock() if toFlush != nil { logBuffer.queueFlush(toFlush) } + logBuffer.writeMu.Unlock() } } diff --git a/weed/util/log_buffer/shutdown_test.go b/weed/util/log_buffer/shutdown_test.go new file mode 100644 index 000000000..7ce70b0fa --- /dev/null +++ b/weed/util/log_buffer/shutdown_test.go @@ -0,0 +1,112 @@ +package log_buffer + +import ( + "bytes" + "errors" + "fmt" + "reflect" + "testing" + "testing/synctest" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +func TestShutdownDrainsAcceptedRecords(t *testing.T) { + for _, seal := range []string{"append", "interval", "force"} { + t.Run(seal, func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + releaseFlush := make(chan struct{}) + var persisted bytes.Buffer + lb := NewLogBuffer("shutdown", time.Hour, func(_ *LogBuffer, _, _ time.Time, data []byte, _, _ int64) { + <-releaseFlush + persisted.Write(data) + }, nil, nil) + + // Fill the queue behind a blocked persistence call, leaving one + // record in the current window. Each append seals its predecessor. + var want []string + appendRecord := func() { + value := fmt.Sprint(len(want)) + ts := time.Now().Add(time.Duration(len(want)) * 2 * time.Hour).UnixNano() + if err := lb.AddDataToBuffer(nil, []byte(value), ts); err != nil { + t.Error(err) + } + want = append(want, value) + } + for range flushQueueDepth + 2 { + appendRecord() + synctest.Wait() + } + switch seal { + case "append": + go appendRecord() + case "interval": + time.Sleep(time.Hour) + case "force": + go lb.ForceFlush() + } + synctest.Wait() + + done := make(chan struct{}) + go func() { + lb.ShutdownLogBuffer() + lb.WaitForShutdown() + close(done) + }() + <-lb.shutdownCh + select { + case <-done: + t.Error("shutdown completed while persistence was blocked") + default: + } + + close(releaseFlush) + <-done + synctest.Wait() + var got []string + for persisted.Len() > 0 { + size := util.BytesToUint32(persisted.Next(4)) + var entry filer_pb.LogEntry + if err := entry.UnmarshalVT(persisted.Next(int(size))); err != nil { + t.Fatal(err) + } + got = append(got, string(entry.Data)) + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("persisted records = %v, want %v", got, want) + } + }) + }) + } +} + +func TestShutdownRejectsNewRecords(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + lb := NewLogBuffer("stopped", time.Hour, nil, nil, nil) + lb.ShutdownLogBuffer() + for _, appendRecord := range []func() error{ + func() error { return lb.AddDataToBuffer(nil, []byte("late data"), 0) }, + func() error { return lb.AddLogEntryToBuffer(&filer_pb.LogEntry{Data: []byte("late entry")}) }, + } { + if err := appendRecord(); !errors.Is(err, ErrBufferStopped) { + t.Errorf("append after shutdown = %v, want ErrBufferStopped", err) + } + } + lb.ForceFlush() + lb.WaitForShutdown() + if lb.pos != 0 { + t.Errorf("stopped buffer retained %d bytes", lb.pos) + } + }) +} + +func TestWaitForShutdownWithEmptyBuffer(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + lb := NewLogBuffer("empty", time.Hour, nil, nil, nil) + lb.ShutdownLogBuffer() + lb.WaitForShutdown() + lb.WaitForShutdown() + }) +}