diff --git a/backend/webdav/chunking.go b/backend/webdav/chunking.go index ce29507c9..a871353c3 100644 --- a/backend/webdav/chunking.go +++ b/backend/webdav/chunking.go @@ -17,6 +17,7 @@ import ( "time" "github.com/rclone/rclone/fs" + "github.com/rclone/rclone/fs/fserrors" "github.com/rclone/rclone/lib/readers" "github.com/rclone/rclone/lib/rest" ) @@ -133,7 +134,26 @@ func (o *Object) uploadChunks(ctx context.Context, in0 io.Reader, size int64, pa return io.NopCloser(in), nil } - err := partObj.updateSimple(ctx, in, getBody, partObj.remote, contentLength, "application/x-www-form-urlencoded", nil, o.fs.chunksUploadURL, options...) + // Retry this chunk in place so a transient failure does not make the caller + // restart the whole upload and resend chunks that were already accepted. + maxTries := max(fs.GetConfig(ctx).LowLevelRetries, 1) + var err error + for try := 1; try <= maxTries; try++ { + retryBody, bodyErr := getBody() + if bodyErr != nil { + return fmt.Errorf("rewinding chunk for retry failed: %w", bodyErr) + } + err = partObj.updateSimple(ctx, retryBody, getBody, partObj.remote, contentLength, "application/x-www-form-urlencoded", nil, o.fs.chunksUploadURL, options...) + _ = retryBody.Close() + if err == nil || (!fserrors.IsRetryError(err) && !fserrors.ShouldRetry(err)) { + break + } + if try == maxTries { + // The chunk retry budget is exhausted. Do not let the outer + // low-level retry restart the upload from its first chunk. + err = fserrors.NoLowLevelRetryError(errors.New(err.Error())) + } + } if err != nil { return fmt.Errorf("uploading chunk failed: %w", err) } diff --git a/backend/webdav/webdav_internal_test.go b/backend/webdav/webdav_internal_test.go index f5de11d8b..87f223120 100644 --- a/backend/webdav/webdav_internal_test.go +++ b/backend/webdav/webdav_internal_test.go @@ -1,6 +1,7 @@ package webdav_test import ( + "bytes" "context" "fmt" "io" @@ -9,6 +10,7 @@ import ( "strings" "sync/atomic" "testing" + "time" "github.com/rclone/rclone/backend/local" "github.com/rclone/rclone/backend/webdav" @@ -16,11 +18,93 @@ import ( "github.com/rclone/rclone/fs/config/configfile" "github.com/rclone/rclone/fs/config/configmap" "github.com/rclone/rclone/fs/config/obscure" + "github.com/rclone/rclone/fs/object" "github.com/rclone/rclone/fs/operations" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) +func TestChunkedUploadRetriesOnlyFailedChunk(t *testing.T) { + tests := []struct { + name string + failChunk1 bool + wantError bool + wantChunk2 int32 + }{ + {name: "retry failed chunk", wantChunk2: 1}, + {name: "stop after chunk retry budget", failChunk1: true, wantError: true}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + var chunk0, chunk1, chunk2 atomic.Int32 + handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch { + case r.Method == http.MethodPut && strings.Contains(r.URL.Path, "/dav/uploads/alice/"): + switch { + case strings.HasSuffix(r.URL.Path, "000000000000000-000000000000003"): + chunk0.Add(1) + case strings.HasSuffix(r.URL.Path, "000000000000004-000000000000007"): + if test.failChunk1 || chunk1.Add(1) == 1 { + if test.failChunk1 { + chunk1.Add(1) + } + w.WriteHeader(http.StatusServiceUnavailable) + return + } + case strings.HasSuffix(r.URL.Path, "000000000000008-000000000000009"): + chunk2.Add(1) + default: + t.Errorf("unexpected chunk path %q", r.URL.Path) + } + w.WriteHeader(http.StatusCreated) + case r.Method == "MKCOL": + w.WriteHeader(http.StatusCreated) + case r.Method == "MOVE": + w.WriteHeader(http.StatusCreated) + case r.Method == "DELETE": + w.WriteHeader(http.StatusNotFound) + case r.Method == "PROPFIND": + w.WriteHeader(http.StatusMultiStatus) + _, _ = io.WriteString(w, `/file10HTTP/1.1 200 OK`) + case r.Method == "PATCH": + w.Header().Set("OC-Checksum", "SHA1:0000000000000000000000000000000000000000") + w.WriteHeader(http.StatusOK) + default: + t.Errorf("unexpected request %s %s", r.Method, r.URL.Path) + w.WriteHeader(http.StatusNotFound) + } + }) + ts := httptest.NewServer(handler) + defer ts.Close() + configfile.Install() + ctx, ci := fs.AddConfig(context.Background()) + ci.LowLevelRetries = 2 + m := configmap.Simple{ + "type": "webdav", + "url": ts.URL + "/remote.php/dav/files/alice", + "vendor": "nextcloud", + "nextcloud_chunk_size": "4B", + } + f, err := webdav.NewFs(ctx, remoteName, "", m) + require.NoError(t, err) + src := object.NewStaticObjectInfo("file", time.Now(), 10, true, nil, f) + body := []byte("abcdefghij") + err = operations.Retry(ctx, src, ci.LowLevelRetries, func() error { + _, err := f.Put(ctx, bytes.NewReader(body), src) + return err + }) + if test.wantError { + require.Error(t, err) + } else { + require.NoError(t, err) + } + assert.EqualValues(t, 1, chunk0.Load(), "successful first chunk should not be uploaded again") + assert.EqualValues(t, 2, chunk1.Load(), "only the failed chunk should be retried") + assert.EqualValues(t, test.wantChunk2, chunk2.Load(), "later chunks should not upload after failure") + }) + } +} + var ( remoteName = "TestWebDAV" headers = []string{"X-Potato", "sausage", "X-Rhubarb", "cucumber"}