From ae4839e005597cc5b1170fff0554c15bff2a5679 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Fri, 31 Jul 2026 01:16:02 -0700 Subject: [PATCH] mount: keep a sealed chunk alive until its own upload finishes (#10504) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Sealing a logic chunk index that already held a sealed chunk dropped the old chunk's only reference and freed its page chunk. That chunk's upload may not have started reading it yet — Execute() returns as soon as the job is handed to a goroutine — so mem.Free could hand a live 2 MiB mem chunk back to the slot pool, the next NewMemChunk would overwrite it, and the in-flight upload shipped whatever bytes were there. Under fio randwrite the volume server rejected those needles with "Content-MD5 did not match md5 of file data" and the FUSE write failed with EIO. Give the sealed chunk a second reference for its upload, dropped only by the upload itself, and let the upload unindex itself only while it still owns the index — the unconditional delete could evict a newer sealed chunk and hide its dirty pages from readers. --- weed/mount/page_writer/upload_pipeline.go | 20 ++- .../upload_pipeline_reseal_test.go | 121 ++++++++++++++++++ 2 files changed, 134 insertions(+), 7 deletions(-) create mode 100644 weed/mount/page_writer/upload_pipeline_reseal_test.go diff --git a/weed/mount/page_writer/upload_pipeline.go b/weed/mount/page_writer/upload_pipeline.go index 275a174eb..1af97fa25 100644 --- a/weed/mount/page_writer/upload_pipeline.go +++ b/weed/mount/page_writer/upload_pipeline.go @@ -221,8 +221,11 @@ func (up *UploadPipeline) moveToSealed(memChunk PageChunk, logicChunkIndex Logic oldMemChunk.FreeReference(fmt.Sprintf("%s replace chunk %d", up.filepath, logicChunkIndex)) } sealedChunk := &SealedChunk{ - chunk: memChunk, - referenceCounter: 1, // default 1 is for uploading process + chunk: memChunk, + // One for the sealedChunks slot, one for the upload below, which + // drops its own. Execute() returns before SaveContent reads the + // chunk, so a later seal must not free the buffer under it. + referenceCounter: 2, accountant: up.accountant, chunkSize: up.ChunkSize, } @@ -252,8 +255,12 @@ func (up *UploadPipeline) moveToSealed(memChunk PageChunk, logicChunkIndex Logic up.readerCountCond.Wait() } - // then remove from sealed chunks - delete(up.sealedChunks, logicChunkIndex) + // then remove from sealed chunks, unless a later seal took over the + // index — deleting that one would hide its dirty pages from readers + if up.sealedChunks[logicChunkIndex] == sealedChunk { + delete(up.sealedChunks, logicChunkIndex) + sealedChunk.FreeReference(fmt.Sprintf("%s unindex chunk %d", up.filepath, logicChunkIndex)) + } sealedChunk.FreeReference(fmt.Sprintf("%s finished uploading chunk %d", up.filepath, logicChunkIndex)) }) @@ -388,9 +395,8 @@ func (up *UploadPipeline) Shutdown() { delete(up.writableChunks, logicChunkIndex) } for logicChunkIndex, sealedChunk := range up.sealedChunks { - // FreeReference releases the accountant slot on the refcount-zero - // transition; a racing async uploader will call FreeReference again - // and be a no-op, so there is no double-release. + // only the slot reference; an in-flight upload frees the chunk itself sealedChunk.FreeReference(fmt.Sprintf("%s uploadpipeline shutdown chunk %d", up.filepath, logicChunkIndex)) + delete(up.sealedChunks, logicChunkIndex) } } diff --git a/weed/mount/page_writer/upload_pipeline_reseal_test.go b/weed/mount/page_writer/upload_pipeline_reseal_test.go new file mode 100644 index 000000000..41d78ef0b --- /dev/null +++ b/weed/mount/page_writer/upload_pipeline_reseal_test.go @@ -0,0 +1,121 @@ +package page_writer + +import ( + "sync/atomic" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/util" +) + +// gatedChunk parks in SaveContent until the test opens its gate: the state +// moveToSealed leaves an upload in, submitted but not yet reading its buffer. +type gatedChunk struct { + gate chan struct{} + saved chan struct{} + freedBefore bool // FreeResource had already run when SaveContent began + freeCount atomic.Int32 +} + +func newGatedChunk() *gatedChunk { + return &gatedChunk{ + gate: make(chan struct{}), + saved: make(chan struct{}), + } +} + +func (c *gatedChunk) FreeResource() { c.freeCount.Add(1) } +func (c *gatedChunk) WriteDataAt([]byte, int64, int64) int { return 0 } +func (c *gatedChunk) ReadDataAt([]byte, int64, int64) int64 { return 0 } +func (c *gatedChunk) IsComplete() bool { return false } +func (c *gatedChunk) IsContiguouslyWritten() bool { return true } +func (c *gatedChunk) ActivityScore() int64 { return 0 } +func (c *gatedChunk) WrittenSize() int64 { return 0 } +func (c *gatedChunk) LastWriteTsNs() int64 { return 0 } +func (c *gatedChunk) SaveContent(saveFn SaveToStorageFunc) { + <-c.gate + c.freedBefore = c.freeCount.Load() > 0 + close(c.saved) +} + +// Sealing the same index twice must not free the first chunk: freeing it +// returns a live buffer to the slot pool and the in-flight upload then ships +// whatever the next NewMemChunk writes there. +func TestMoveToSealed_ReplaceKeepsUploadReference(t *testing.T) { + up := NewUploadPipeline(util.NewLimitedConcurrentExecutor(2), 2*1024*1024, nil, 16, "", nil) + + first, second := newGatedChunk(), newGatedChunk() + + up.chunksLock.Lock() + up.moveToSealed(first, 0) + up.moveToSealed(second, 0) // takes over index 0 while first is uploading + up.chunksLock.Unlock() + + // both seals are done, so any premature free has already happened + close(first.gate) + select { + case <-first.saved: + case <-time.After(5 * time.Second): + t.Fatal("first chunk's SaveContent never ran") + } + if first.freedBefore { + t.Error("first chunk was freed before its upload read it") + } + + // a finished upload must leave the newer chunk indexed for readers + waitForFree(t, first) + up.chunksLock.Lock() + indexed := up.sealedChunks[0] + up.chunksLock.Unlock() + if indexed == nil || indexed.chunk != second { + t.Errorf("index 0 should still hold the second chunk, got %v", indexed) + } + + close(second.gate) + waitForFree(t, second) + if got := first.freeCount.Load(); got != 1 { + t.Errorf("first chunk freed %d times, want 1", got) + } + if got := second.freeCount.Load(); got != 1 { + t.Errorf("second chunk freed %d times, want 1", got) + } +} + +// Shutdown drops only the slot reference; the in-flight upload keeps the +// chunk alive until it is done, then frees it exactly once. +func TestShutdown_LeavesInFlightUploadItsReference(t *testing.T) { + up := NewUploadPipeline(util.NewLimitedConcurrentExecutor(2), 2*1024*1024, nil, 16, t.TempDir(), nil) + + chunk := newGatedChunk() + up.chunksLock.Lock() + up.moveToSealed(chunk, 0) + up.chunksLock.Unlock() + + up.Shutdown() + + close(chunk.gate) + select { + case <-chunk.saved: + case <-time.After(5 * time.Second): + t.Fatal("SaveContent never ran") + } + if chunk.freedBefore { + t.Error("Shutdown freed the chunk before its upload read it") + } + waitForFree(t, chunk) + if got := chunk.freeCount.Load(); got != 1 { + t.Errorf("chunk freed %d times, want 1", got) + } +} + +func waitForFree(t *testing.T, c *gatedChunk) { + t.Helper() + <-c.saved + deadline := time.Now().Add(5 * time.Second) + for c.freeCount.Load() == 0 { + if time.Now().After(deadline) { + t.Fatal("chunk was never freed after its upload finished") + } + time.Sleep(time.Millisecond) + } +}