Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
89e609f063 |
@@ -236,8 +236,7 @@ func (r *RcloneStorer) Info() Info {
|
|||||||
// PartialUploads. That flag is not checked here: hdfs, for one, shows a
|
// PartialUploads. That flag is not checked here: hdfs, for one, shows a
|
||||||
// file while it is written without setting it.
|
// file while it is written without setting it.
|
||||||
func (r *RcloneStorer) upload(ctx context.Context, key string, data io.Reader) error {
|
func (r *RcloneStorer) upload(ctx context.Context, key string, data io.Reader) error {
|
||||||
features := r.fsys.Features()
|
if r.fsys.Features().Move == nil {
|
||||||
if features.Move == nil {
|
|
||||||
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(data), time.Now(), nil)
|
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(data), time.Now(), nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("uploading object: %w", err)
|
return fmt.Errorf("uploading object: %w", err)
|
||||||
@@ -258,7 +257,19 @@ func (r *RcloneStorer) upload(ctx context.Context, key string, data io.Reader) e
|
|||||||
return fmt.Errorf("uploading object: %w", err)
|
return fmt.Errorf("uploading object: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
_, err = features.Move(ctx, obj, key)
|
// 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 {
|
if err != nil {
|
||||||
_ = obj.Remove(ctx)
|
_ = obj.Remove(ctx)
|
||||||
|
|
||||||
|
|||||||
@@ -4,12 +4,14 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"io"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/rclone/rclone/fs"
|
"github.com/rclone/rclone/fs"
|
||||||
|
"github.com/rclone/rclone/fs/config/configmap"
|
||||||
"sneak.berlin/go/vaultik/internal/storage"
|
"sneak.berlin/go/vaultik/internal/storage"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -162,6 +164,85 @@ func TestRcloneStorerListSkipsPartialFiles(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
var errNameConflict = errors.New("an object with this name already exists")
|
||||||
|
|
||||||
|
// moveRefusesExisting is rclone's local backend with a server-side move
|
||||||
|
// that, like dropbox's or onedrive's, refuses to move onto an existing
|
||||||
|
// object.
|
||||||
|
type moveRefusesExisting struct {
|
||||||
|
fs.Fs
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *moveRefusesExisting) Features() *fs.Features {
|
||||||
|
features := *f.Fs.Features()
|
||||||
|
features.Move = f.move
|
||||||
|
|
||||||
|
return &features
|
||||||
|
}
|
||||||
|
|
||||||
|
//nolint:ireturn // the signature is rclone's
|
||||||
|
func (f *moveRefusesExisting) move(
|
||||||
|
ctx context.Context, src fs.Object, remote string,
|
||||||
|
) (fs.Object, error) {
|
||||||
|
_, err := f.NewObject(ctx, remote)
|
||||||
|
if err == nil {
|
||||||
|
return nil, errNameConflict
|
||||||
|
}
|
||||||
|
|
||||||
|
return f.Fs.Features().Move(ctx, src, remote)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRcloneStorerOverwritesWhereMoveRefusesExisting checks that writing a
|
||||||
|
// key twice replaces the object on a remote whose server-side move will not
|
||||||
|
// replace an existing object.
|
||||||
|
//
|
||||||
|
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||||
|
func TestRcloneStorerOverwritesWhereMoveRefusesExisting(t *testing.T) {
|
||||||
|
fs.Register(&fs.RegInfo{
|
||||||
|
Name: "moverefusesexisting",
|
||||||
|
NewFs: func(
|
||||||
|
ctx context.Context, _, root string, _ configmap.Mapper,
|
||||||
|
) (fs.Fs, error) {
|
||||||
|
local, err := fs.NewFs(ctx, ":local:"+root)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return &moveRefusesExisting{Fs: local}, nil
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
s, err := storage.NewRcloneStorer(ctx, ":moverefusesexisting", t.TempDir())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("NewRcloneStorer: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, content := range []string{"first", "second"} {
|
||||||
|
err = s.Put(ctx, testBlobKey, strings.NewReader(content))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Put %q: %v", content, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
rc, err := s.Get(ctx, testBlobKey)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Get: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
defer func() { _ = rc.Close() }()
|
||||||
|
|
||||||
|
got, err := io.ReadAll(rc)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("reading object: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if string(got) != "second" {
|
||||||
|
t.Errorf("object = %q, want %q", got, "second")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// TestNewRcloneStorerUnknownRemote checks that a remote that is not in the
|
// TestNewRcloneStorerUnknownRemote checks that a remote that is not in the
|
||||||
// rclone config fails construction with the ErrRemoteNotFound sentinel,
|
// rclone config fails construction with the ErrRemoteNotFound sentinel,
|
||||||
// rather than silently returning a backend pointed nowhere.
|
// rather than silently returning a backend pointed nowhere.
|
||||||
|
|||||||
Reference in New Issue
Block a user