Write rclone uploads under a temporary name and move them into place (closes #266)
check / check (push) Waiting to run
check / check (push) Waiting to run
The rclone backend wrote each object straight to its key, so killing an upload to a local or sftp remote left a truncated object there that the next backup trusted. On every remote with a server-side move, an object is now written under a name ending in `.partial` and moved onto its key with rclone's operations.Move, which removes an object already at the key first; drive, dropbox and others will not move onto one. Listings skip `.partial` names. Rclone's own copy also requires the PartialUploads flag; this does not, because hdfs shows a file while it is written without setting it. Remotes without a move are written in place. Model: opus-5-5
This commit is contained in:
+64
-15
@@ -3,6 +3,7 @@ package storage
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -68,14 +69,7 @@ func (r *RcloneStorer) Put(ctx context.Context, key string, data io.Reader) erro
|
||||
return fmt.Errorf("reading data: %w", err)
|
||||
}
|
||||
|
||||
// Upload the object
|
||||
_, err = operations.Rcat(ctx, r.fsys, key,
|
||||
io.NopCloser(bytes.NewReader(buf)), time.Now(), nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
return r.upload(ctx, key, bytes.NewReader(buf))
|
||||
}
|
||||
|
||||
// PutWithProgress stores data with progress reporting.
|
||||
@@ -89,13 +83,7 @@ func (r *RcloneStorer) PutWithProgress(
|
||||
callback: progress,
|
||||
}
|
||||
|
||||
// Upload the object
|
||||
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(pr), time.Now(), nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
return r.upload(ctx, key, pr)
|
||||
}
|
||||
|
||||
// Get retrieves data from the specified key.
|
||||
@@ -173,6 +161,10 @@ func (r *RcloneStorer) List(ctx context.Context, prefix string) ([]string, error
|
||||
|
||||
err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) {
|
||||
key := obj.Remote()
|
||||
if strings.HasSuffix(key, tempSuffix) {
|
||||
return
|
||||
}
|
||||
|
||||
if prefix == "" || strings.HasPrefix(key, prefix) {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
@@ -202,6 +194,10 @@ func (r *RcloneStorer) ListStream(
|
||||
}
|
||||
|
||||
key := obj.Remote()
|
||||
if strings.HasSuffix(key, tempSuffix) {
|
||||
return
|
||||
}
|
||||
|
||||
if prefix == "" || strings.HasPrefix(key, prefix) {
|
||||
ch <- ObjectInfo{
|
||||
Key: key,
|
||||
@@ -230,6 +226,59 @@ func (r *RcloneStorer) Info() Info {
|
||||
}
|
||||
}
|
||||
|
||||
// upload writes data to key. Where the remote has a server-side move, it
|
||||
// writes under a temporary name ending in tempSuffix and moves the object
|
||||
// onto key once it is complete, so a killed upload cannot leave a truncated
|
||||
// object at key; a remote without one is written in place. List and
|
||||
// ListStream skip a temporary object left behind.
|
||||
//
|
||||
// rclone's own copy does this only where the remote also sets
|
||||
// PartialUploads. That flag is not checked here: hdfs, for one, shows a
|
||||
// file while it is written without setting it.
|
||||
func (r *RcloneStorer) upload(ctx context.Context, key string, data io.Reader) error {
|
||||
if r.fsys.Features().Move == nil {
|
||||
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(data), time.Now(), nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
tempKey := key + "-" + rand.Text() + tempSuffix
|
||||
|
||||
obj, err := operations.Rcat(ctx, r.fsys, tempKey, io.NopCloser(data), time.Now(), nil)
|
||||
if err != nil {
|
||||
// Rcat returns the object it wrote when the written data fails its check.
|
||||
if obj != nil {
|
||||
_ = obj.Remove(ctx)
|
||||
}
|
||||
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
// On drive, dropbox, onedrive and others the remote's own move does not
|
||||
// replace an object already at key. operations.Move removes the object
|
||||
// it is given first, and copies where the remote refuses the move.
|
||||
existing, err := r.fsys.NewObject(ctx, key)
|
||||
if errors.Is(err, fs.ErrorObjectNotFound) {
|
||||
existing = nil
|
||||
} else if err != nil {
|
||||
_ = obj.Remove(ctx)
|
||||
|
||||
return fmt.Errorf("looking up existing object: %w", err)
|
||||
}
|
||||
|
||||
_, err = operations.Move(ctx, r.fsys, existing, key, obj)
|
||||
if err != nil {
|
||||
_ = obj.Remove(ctx)
|
||||
|
||||
return fmt.Errorf("moving object into place: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// progressReader wraps an io.Reader to track read progress.
|
||||
type progressReader struct {
|
||||
reader io.Reader
|
||||
|
||||
Reference in New Issue
Block a user