Retry transient uploads and continue backups
All checks were successful
Build / Test and build (push) Successful in 5m53s

This commit is contained in:
2026-08-15 15:32:43 +08:00
parent 977fcefc2f
commit c602698ce9
7 changed files with 175 additions and 28 deletions

View File

@@ -25,16 +25,17 @@ const (
)
type Client struct {
httpClient *http.Client
apiBase string
oauthBase string
uploadBase string
userAgent string
root string
uploadParts int
partSize int64
downloadMu sync.Mutex
downloadURL map[int64]cachedDownloadURL
httpClient *http.Client
apiBase string
oauthBase string
uploadBase string
userAgent string
root string
uploadParts int
partSize int64
uploadRetryDelay func(int) time.Duration
downloadMu sync.Mutex
downloadURL map[int64]cachedDownloadURL
mu sync.Mutex
clientID string
@@ -75,8 +76,9 @@ func New(cfg *config.Config, opts ...Option) *Client {
apiBase: defaultAPIBase, oauthBase: defaultOAuthBase, uploadBase: defaultUploadBase,
userAgent: cfg.UserAgent, root: cfg.Root,
uploadParts: cfg.UploadParts, partSize: cfg.PartSize,
downloadURL: make(map[int64]cachedDownloadURL),
clientID: cfg.ClientID, secret: cfg.ClientSecret,
uploadRetryDelay: defaultUploadRetryDelay,
downloadURL: make(map[int64]cachedDownloadURL),
clientID: cfg.ClientID, secret: cfg.ClientSecret,
accessToken: cfg.AccessToken, refresh: cfg.RefreshToken, expiresAt: cfg.ExpiresAt,
}
for _, opt := range opts {

View File

@@ -12,6 +12,7 @@ import (
"path"
"reflect"
"strconv"
"strings"
"sync"
"testing"
"time"
@@ -205,6 +206,60 @@ func TestUploadMultipartFlow(t *testing.T) {
}
}
func TestUploadPartRetriesGatewayTimeout(t *testing.T) {
file, err := os.CreateTemp(t.TempDir(), "retry-upload-*")
if err != nil {
t.Fatal(err)
}
defer file.Close()
content := []byte("retry payload")
if _, err := file.Write(content); err != nil {
t.Fatal(err)
}
attempts := 0
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/rest/2.0/pcs/superfile2" {
http.NotFound(w, r)
return
}
attempts++
_, _ = io.Copy(io.Discard, r.Body)
if attempts < 4 {
http.Error(w, "gateway timeout", http.StatusGatewayTimeout)
return
}
fmt.Fprint(w, `{"md5":"ok"}`)
}))
defer server.Close()
client := New(testConfig("token"), WithEndpoints(server.URL, server.URL, server.URL))
client.uploadRetryDelay = func(int) time.Duration { return 0 }
if err := client.uploadPartWithRetry(context.Background(), file, "/apps/bdrclone/retry.bin", "upload-1", 0, 0, int64(len(content))); err != nil {
t.Fatal(err)
}
if attempts != 4 {
t.Fatalf("attempts = %d, want 4", attempts)
}
failureAttempts := 0
failureServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
failureAttempts++
_, _ = io.Copy(io.Discard, r.Body)
http.Error(w, "gateway timeout", http.StatusGatewayTimeout)
}))
defer failureServer.Close()
failureClient := New(testConfig("token"), WithEndpoints(failureServer.URL, failureServer.URL, failureServer.URL))
failureClient.uploadRetryDelay = func(int) time.Duration { return 0 }
err = failureClient.uploadPartWithRetry(context.Background(), file, "/apps/bdrclone/retry.bin", "upload-2", 0, 0, int64(len(content)))
if err == nil || !strings.Contains(err.Error(), "after 6 attempts") {
t.Fatalf("error = %v", err)
}
if failureAttempts != maxUploadAttempts {
t.Fatalf("attempts = %d, want %d", failureAttempts, maxUploadAttempts)
}
}
func TestMkdirAllCreatesMissingRemoteDirectories(t *testing.T) {
directories := map[string]bool{"/apps/bdrclone": true}
var created []string

View File

@@ -21,10 +21,11 @@ import (
)
const (
defaultPartSize = int64(4 << 20)
vipPartSize = int64(16 << 20)
svipPartSize = int64(32 << 20)
maxPartCount = 2048
defaultPartSize = int64(4 << 20)
vipPartSize = int64(16 << 20)
svipPartSize = int64(32 << 20)
maxPartCount = 2048
maxUploadAttempts = 6
)
type UploadProgress func(uploaded, total int64)
@@ -247,7 +248,7 @@ func (c *Client) uploadPartsParallel(ctx context.Context, file *os.File, remote,
func (c *Client) uploadPartWithRetry(ctx context.Context, file *os.File, remote, uploadID string, part int, offset, size int64) error {
var lastErr error
for attempt := 0; attempt < 3; attempt++ {
for attempt := 0; attempt < maxUploadAttempts; attempt++ {
if err := ctx.Err(); err != nil {
return err
}
@@ -256,13 +257,18 @@ func (c *Client) uploadPartWithRetry(ctx context.Context, file *os.File, remote,
if lastErr == nil {
return nil
}
if attempt < 2 {
if err := sleepContext(ctx, time.Duration(1<<attempt)*time.Second); err != nil {
if attempt < maxUploadAttempts-1 {
if err := sleepContext(ctx, c.uploadRetryDelay(attempt)); err != nil {
return err
}
}
}
return lastErr
return fmt.Errorf("upload failed after %d attempts: %w", maxUploadAttempts, lastErr)
}
func defaultUploadRetryDelay(attempt int) time.Duration {
delay := time.Second * time.Duration(1<<attempt)
return min(delay, 30*time.Second)
}
func (c *Client) uploadPart(ctx context.Context, section *io.SectionReader, filename, remote, uploadID string, part int) error {