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 <info@tdvorak.dev>

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

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
pull/3582/head
Tomas Dvorak 2 weeks ago
parent 93e11533c0
commit 6f23834b87

@ -152,16 +152,10 @@ func (handler *Driver) Put(ctx context.Context, file *fs.UploadRequest) error {
} }
defer out.Close() defer out.Close()
stat, err := out.Stat() // Chunks may arrive out of order under concurrent uploads, so the current
if err != nil { // file size is not a valid precondition for the chunk offset. Positioned
handler.l.Warning("Failed to read file info: %s", err) // writes are safe here; a missing range is caught by the assembled-size
return err // check in CompleteUpload.
}
if stat.Size() < file.Offset {
return errors.New("size of unfinished uploaded chunks is not as expected")
}
if _, err := out.Seek(file.Offset, io.SeekStart); err != nil { if _, err := out.Seek(file.Offset, io.SeekStart); err != nil {
return fmt.Errorf("failed to seek to desired offset %d: %s", file.Offset, err) 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 { 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 == "" { if session.Callback == "" {
return nil return nil
} }

@ -2,10 +2,13 @@ package local
import ( import (
"context" "context"
"io"
"os" "os"
"path/filepath" "path/filepath"
"strings"
"testing" "testing"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs"
"github.com/cloudreve/Cloudreve/v4/pkg/logging" "github.com/cloudreve/Cloudreve/v4/pkg/logging"
"github.com/cloudreve/Cloudreve/v4/pkg/util" "github.com/cloudreve/Cloudreve/v4/pkg/util"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
@ -37,3 +40,41 @@ func TestDeletePrunesEmptyAncestors(t *testing.T) {
require.DirExists(t, "uploads") require.DirExists(t, "uploads")
require.FileExists(t, filepath.Join("uploads", "keep.txt")) 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)
}

Loading…
Cancel
Save