mirror of
https://github.com/rclone/rclone.git
synced 2026-10-09 06:28:15 -04:00
webdav: retry Nextcloud chunk uploads individually - fixes #10005
Nextcloud chunk uploads currently return a retryable error when a chunk PUT fails. The outer low-level retry then restarts the whole upload, removing the temporary upload directory and sending previously accepted chunks again. Retry each chunk in place using the configured low-level retry budget, rewinding its repeatable body before each attempt. Once that budget is exhausted, prevent the outer retry from restarting the entire object.
This commit is contained in:
1 parent
f01af78486
commit
039b69adb8
2 files changed
+105
-1
No files matched your search
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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, `<d:multistatus xmlns:d="DAV:"><d:response><d:href>/file</d:href><d:propstat><d:prop><d:getcontentlength>10</d:getcontentlength><d:resourcetype/></d:prop><d:status>HTTP/1.1 200 OK</d:status></d:propstat></d:response></d:multistatus>`)
|
||||
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"}
|
||||
|
||||
Reference in new issue
Block a user