From 6f23834b87f50614e2ba3e9fbe73598f2c1315f5 Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Fri, 18 Sep 2026 23:16:43 +0200 Subject: [PATCH] fix(local): allow out-of-order chunk writes in Put (#3157) The local driver rejected any chunk whose offset exceeded the current file size, which assumed sequential arrival. With chunk_concurrency > 1 later chunks routinely land before earlier ones finish, so parallel uploads to local storage failed with "size of unfinished uploaded chunks is not as expected" while single-chunk concurrency worked. Positioned writes are safe with per-request file descriptors; the assembled file size is now verified in CompleteUpload so a missing or lost range still fails loudly instead of silently corrupting the blob. Authored By: TDvorak Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/filemanager/driver/local/local.go | 32 +++++++++++------ pkg/filemanager/driver/local/local_test.go | 41 ++++++++++++++++++++++ 2 files changed, 63 insertions(+), 10 deletions(-) diff --git a/pkg/filemanager/driver/local/local.go b/pkg/filemanager/driver/local/local.go index 127eaa6a..757f1926 100644 --- a/pkg/filemanager/driver/local/local.go +++ b/pkg/filemanager/driver/local/local.go @@ -152,16 +152,10 @@ func (handler *Driver) Put(ctx context.Context, file *fs.UploadRequest) error { } defer out.Close() - stat, err := out.Stat() - if err != nil { - handler.l.Warning("Failed to read file info: %s", err) - return err - } - - if stat.Size() < file.Offset { - return errors.New("size of unfinished uploaded chunks is not as expected") - } - + // Chunks may arrive out of order under concurrent uploads, so the current + // file size is not a valid precondition for the chunk offset. Positioned + // writes are safe here; a missing range is caught by the assembled-size + // check in CompleteUpload. if _, err := out.Seek(file.Offset, io.SeekStart); err != nil { return fmt.Errorf("failed to seek to desired offset %d: %s", file.Offset, err) } @@ -279,6 +273,24 @@ func (handler *Driver) CancelToken(ctx context.Context, uploadSession *fs.Upload } func (handler *Driver) CompleteUpload(ctx context.Context, session *fs.UploadSession) error { + // Concurrent chunks are written at their own offsets and may land out of + // order; verify the assembled file matches the declared size before + // reporting success so a missing or lost range fails loudly instead of + // leaving a corrupted blob. + if session.Props != nil && session.Props.SavePath != "" { + stat, err := os.Stat(handler.LocalPath(ctx, session.Props.SavePath)) + if err != nil { + return fmt.Errorf("failed to stat uploaded file: %w", err) + } + if stat.Size() != session.Props.Size { + return serializer.NewError( + serializer.CodeUploadFailed, + fmt.Sprintf("uploaded file size mismatch: expected %d, got %d", session.Props.Size, stat.Size()), + nil, + ) + } + } + if session.Callback == "" { return nil } diff --git a/pkg/filemanager/driver/local/local_test.go b/pkg/filemanager/driver/local/local_test.go index 165be0e0..999d78a6 100644 --- a/pkg/filemanager/driver/local/local_test.go +++ b/pkg/filemanager/driver/local/local_test.go @@ -2,10 +2,13 @@ package local import ( "context" + "io" "os" "path/filepath" + "strings" "testing" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" "github.com/cloudreve/Cloudreve/v4/pkg/logging" "github.com/cloudreve/Cloudreve/v4/pkg/util" "github.com/stretchr/testify/require" @@ -37,3 +40,41 @@ func TestDeletePrunesEmptyAncestors(t *testing.T) { require.DirExists(t, "uploads") require.FileExists(t, filepath.Join("uploads", "keep.txt")) } + +func TestPutOutOfOrderChunks(t *testing.T) { + old := util.UseWorkingDir + util.UseWorkingDir = true + t.Cleanup(func() { util.UseWorkingDir = old }) + t.Chdir(t.TempDir()) + + d := New(nil, logging.NewConsoleLogger(logging.LevelDebug), nil) + ctx := context.Background() + savePath := filepath.Join("uploads", "blob.bin") + + props := &fs.UploadProps{SavePath: savePath, Size: 12} + put := func(offset int64, data string) { + require.NoError(t, d.Put(ctx, &fs.UploadRequest{ + Props: props, + Mode: fs.ModeOverwrite, + Offset: offset, + File: io.NopCloser(strings.NewReader(data)), + })) + } + + // Later chunk lands first — must not be rejected by the offset check. + put(8, "IJXL") + put(4, "EFGH") + put(0, "ABCD") + + content, err := os.ReadFile(savePath) + require.NoError(t, err) + require.Equal(t, "ABCDEFGHIJXL", string(content)) + + // Completion verifies the assembled size. + require.NoError(t, d.CompleteUpload(ctx, &fs.UploadSession{Props: props})) + + // A lost chunk fails loudly at completion instead of silently corrupting. + require.NoError(t, os.WriteFile(savePath, []byte("short"), 0644)) + err = d.CompleteUpload(ctx, &fs.UploadSession{Props: props}) + require.Error(t, err) +}