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
303 lines
7.4 KiB
Go
303 lines
7.4 KiB
Go
package storage
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/rand"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/rclone/rclone/fs"
|
|
"github.com/rclone/rclone/fs/config/configfile"
|
|
"github.com/rclone/rclone/fs/operations"
|
|
|
|
// Import all rclone backends
|
|
_ "github.com/rclone/rclone/backend/all"
|
|
)
|
|
|
|
// ErrRemoteNotFound is returned when an rclone remote is not configured.
|
|
var ErrRemoteNotFound = errors.New("rclone remote not found in config")
|
|
|
|
// RcloneStorer implements Storer using rclone's filesystem abstraction.
|
|
// This allows vaultik to use any of rclone's 70+ supported storage providers.
|
|
type RcloneStorer struct {
|
|
fsys fs.Fs // rclone filesystem
|
|
remote string // remote name (for Info())
|
|
path string // path within remote (for Info())
|
|
}
|
|
|
|
// NewRcloneStorer creates a new rclone storage backend.
|
|
// The remote parameter is the rclone remote name (as configured via `rclone config`).
|
|
// The path parameter is the path within the remote.
|
|
func NewRcloneStorer(ctx context.Context, remote, path string) (*RcloneStorer, error) {
|
|
// Install the default config file handler
|
|
configfile.Install()
|
|
|
|
// Build the rclone path string (e.g., "myremote:path/to/backups")
|
|
rclonePath := remote + ":"
|
|
if path != "" {
|
|
rclonePath += path
|
|
}
|
|
|
|
// Create the rclone filesystem
|
|
fsys, err := fs.NewFs(ctx, rclonePath)
|
|
if err != nil {
|
|
// Check for remote not found error
|
|
if strings.Contains(err.Error(), "didn't find section in config file") ||
|
|
strings.Contains(err.Error(), "failed to find remote") {
|
|
return nil, fmt.Errorf("%w: %s", ErrRemoteNotFound, remote)
|
|
}
|
|
|
|
return nil, fmt.Errorf("creating rclone filesystem: %w", err)
|
|
}
|
|
|
|
return &RcloneStorer{
|
|
fsys: fsys,
|
|
remote: remote,
|
|
path: path,
|
|
}, nil
|
|
}
|
|
|
|
// Put stores data at the specified key.
|
|
func (r *RcloneStorer) Put(ctx context.Context, key string, data io.Reader) error {
|
|
// Read all data into memory to get size (required by rclone)
|
|
buf, err := io.ReadAll(data)
|
|
if err != nil {
|
|
return fmt.Errorf("reading data: %w", err)
|
|
}
|
|
|
|
return r.upload(ctx, key, bytes.NewReader(buf))
|
|
}
|
|
|
|
// PutWithProgress stores data with progress reporting.
|
|
func (r *RcloneStorer) PutWithProgress(
|
|
ctx context.Context, key string, data io.Reader,
|
|
_ int64, progress ProgressCallback,
|
|
) error {
|
|
// Wrap reader with progress tracking
|
|
pr := &progressReader{
|
|
reader: data,
|
|
callback: progress,
|
|
}
|
|
|
|
return r.upload(ctx, key, pr)
|
|
}
|
|
|
|
// Get retrieves data from the specified key.
|
|
func (r *RcloneStorer) Get(ctx context.Context, key string) (io.ReadCloser, error) {
|
|
// Get the object
|
|
obj, err := r.fsys.NewObject(ctx, key)
|
|
if err != nil {
|
|
if errors.Is(err, fs.ErrorObjectNotFound) {
|
|
return nil, ErrNotFound
|
|
}
|
|
|
|
if errors.Is(err, fs.ErrorDirNotFound) {
|
|
return nil, ErrNotFound
|
|
}
|
|
|
|
return nil, fmt.Errorf("getting object: %w", err)
|
|
}
|
|
|
|
// Open the object for reading
|
|
reader, err := obj.Open(ctx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("opening object: %w", err)
|
|
}
|
|
|
|
return reader, nil
|
|
}
|
|
|
|
// Stat returns metadata about an object without retrieving its contents.
|
|
func (r *RcloneStorer) Stat(ctx context.Context, key string) (*ObjectInfo, error) {
|
|
obj, err := r.fsys.NewObject(ctx, key)
|
|
if err != nil {
|
|
if errors.Is(err, fs.ErrorObjectNotFound) {
|
|
return nil, ErrNotFound
|
|
}
|
|
|
|
if errors.Is(err, fs.ErrorDirNotFound) {
|
|
return nil, ErrNotFound
|
|
}
|
|
|
|
return nil, fmt.Errorf("getting object: %w", err)
|
|
}
|
|
|
|
return &ObjectInfo{
|
|
Key: key,
|
|
Size: obj.Size(),
|
|
}, nil
|
|
}
|
|
|
|
// Delete removes an object.
|
|
func (r *RcloneStorer) Delete(ctx context.Context, key string) error {
|
|
obj, err := r.fsys.NewObject(ctx, key)
|
|
if err != nil {
|
|
if errors.Is(err, fs.ErrorObjectNotFound) {
|
|
return nil // Match S3 behavior: no error if doesn't exist
|
|
}
|
|
|
|
if errors.Is(err, fs.ErrorDirNotFound) {
|
|
return nil
|
|
}
|
|
|
|
return fmt.Errorf("getting object: %w", err)
|
|
}
|
|
|
|
err = obj.Remove(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("removing object: %w", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// List returns all keys with the given prefix.
|
|
func (r *RcloneStorer) List(ctx context.Context, prefix string) ([]string, error) {
|
|
var keys []string
|
|
|
|
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)
|
|
}
|
|
})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("listing objects: %w", err)
|
|
}
|
|
|
|
return keys, nil
|
|
}
|
|
|
|
// ListStream returns a channel of ObjectInfo for large result sets.
|
|
func (r *RcloneStorer) ListStream(
|
|
ctx context.Context, prefix string,
|
|
) <-chan ObjectInfo {
|
|
ch := make(chan ObjectInfo)
|
|
|
|
go func() {
|
|
defer close(ch)
|
|
|
|
err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) {
|
|
// Check context cancellation
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
}
|
|
|
|
key := obj.Remote()
|
|
if strings.HasSuffix(key, tempSuffix) {
|
|
return
|
|
}
|
|
|
|
if prefix == "" || strings.HasPrefix(key, prefix) {
|
|
ch <- ObjectInfo{
|
|
Key: key,
|
|
Size: obj.Size(),
|
|
}
|
|
}
|
|
})
|
|
if err != nil {
|
|
ch <- ObjectInfo{Err: fmt.Errorf("listing objects: %w", err)}
|
|
}
|
|
}()
|
|
|
|
return ch
|
|
}
|
|
|
|
// Info returns human-readable storage location information.
|
|
func (r *RcloneStorer) Info() Info {
|
|
location := r.remote
|
|
if r.path != "" {
|
|
location += ":" + r.path
|
|
}
|
|
|
|
return Info{
|
|
Type: schemeRclone,
|
|
Location: location,
|
|
}
|
|
}
|
|
|
|
// 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
|
|
read int64
|
|
callback ProgressCallback
|
|
}
|
|
|
|
func (pr *progressReader) Read(p []byte) (int, error) {
|
|
n, err := pr.reader.Read(p)
|
|
if n > 0 {
|
|
pr.read += int64(n)
|
|
if pr.callback != nil {
|
|
callbackErr := pr.callback(pr.read)
|
|
if callbackErr != nil {
|
|
return n, callbackErr
|
|
}
|
|
}
|
|
}
|
|
|
|
return n, err
|
|
}
|