Files
vaultik/internal/storage/file.go
T
sneak 9974434a3a
check / check (push) Waiting to run
Write rclone uploads under a temporary name and move them into place (closes #266)
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. The next
backup found the key, skipped the upload and recorded a snapshot that
could not be restored.

On every remote with a server-side move, an object is now written under
a name ending in `.partial` and moved onto its key; listings skip such
names. Rclone's own copy also requires the remote's 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, as the README
now says.

Model: opus-5-5
2026-10-07 17:38:34 +00:00

345 lines
8.4 KiB
Go

package storage
import (
"context"
"fmt"
"io"
"os"
"path/filepath"
"strings"
"github.com/spf13/afero"
)
// FileStorer implements Storer using the local filesystem.
// It mirrors the S3 path structure for consistency.
type FileStorer struct {
fs afero.Fs
basePath string
}
// NewFileStorer creates a new filesystem storage backend.
//
// Construction is intentionally cheap and does not touch the filesystem.
// The basePath is recorded; the directory is created lazily on first
// write. Write operations (Put/PutWithProgress) call MkdirAll for the
// per-blob parent directory, which also covers basePath on first use.
// Get and Stat report a key under a missing basePath as ErrNotFound.
// List and ListStream fail on a missing basePath, because listing it as
// an empty store would make `prune` drop every local snapshot record.
//
// Uses the real OS filesystem by default; call SetFilesystem to
// override for testing.
func NewFileStorer(basePath string) (*FileStorer, error) {
return &FileStorer{
fs: afero.NewOsFs(),
basePath: basePath,
}, nil
}
// SetFilesystem overrides the filesystem for testing.
func (f *FileStorer) SetFilesystem(fs afero.Fs) {
f.fs = fs
}
// storageDirPerm is the mode used for directories created under the
// storage base path.
const storageDirPerm = 0o755
// tempSuffix marks a partially written object. writeAtomic streams into a
// temp file carrying this suffix and only renames it onto the real key once
// the whole object is on disk, so an interrupted write can never leave a
// truncated object at the key a later run would Stat and trust as a complete
// blob. The rclone backend's upload does the same on remotes with a
// server-side move. List and ListStream skip these files, so a leftover from
// an interrupted write is never listed or trusted as a blob; it is otherwise
// harmless.
const tempSuffix = ".partial"
// Put stores data at the specified key.
func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error {
return f.writeAtomic(key, data, nil)
}
// PutWithProgress stores data with progress reporting.
func (f *FileStorer) PutWithProgress(
_ context.Context, key string, data io.Reader,
_ int64, progress ProgressCallback,
) error {
return f.writeAtomic(key, data, progress)
}
// Get retrieves data from the specified key.
func (f *FileStorer) Get(_ context.Context, key string) (io.ReadCloser, error) {
path := f.fullPath(key)
file, err := f.fs.Open(path)
if err != nil {
if os.IsNotExist(err) {
return nil, ErrNotFound
}
return nil, fmt.Errorf("opening file: %w", err)
}
return file, nil
}
// Stat returns metadata about an object without retrieving its contents.
func (f *FileStorer) Stat(_ context.Context, key string) (*ObjectInfo, error) {
path := f.fullPath(key)
info, err := f.fs.Stat(path)
if err != nil {
if os.IsNotExist(err) {
return nil, ErrNotFound
}
return nil, fmt.Errorf("stat file: %w", err)
}
return &ObjectInfo{
Key: key,
Size: info.Size(),
}, nil
}
// Delete removes an object.
func (f *FileStorer) Delete(_ context.Context, key string) error {
path := f.fullPath(key)
err := f.fs.Remove(path)
if os.IsNotExist(err) {
return nil // Match S3 behavior: no error if doesn't exist
}
if err != nil {
return fmt.Errorf("removing file: %w", err)
}
return nil
}
// List returns all keys with the given prefix. It fails when the
// destination directory is missing; a missing prefix under it is an
// empty listing.
func (f *FileStorer) List(ctx context.Context, prefix string) ([]string, error) {
var keys []string
_, err := f.fs.Stat(f.basePath)
if err != nil {
return nil, fmt.Errorf("checking destination directory: %w", err)
}
basePath := f.fullPath(prefix)
// Check if the prefix exists
exists, err := afero.Exists(f.fs, basePath)
if err != nil {
return nil, fmt.Errorf("checking path: %w", err)
}
if !exists {
return keys, nil // Empty list for non-existent prefix
}
err = afero.Walk(f.fs, basePath, func(path string, info os.FileInfo, err error) error {
if err != nil {
return err
}
// Check context cancellation
select {
case <-ctx.Done():
return ctx.Err()
default:
}
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
// Convert back to key (relative path from basePath)
relPath, err := filepath.Rel(f.basePath, path)
if err != nil {
return fmt.Errorf("computing relative path: %w", err)
}
// Normalize path separators to forward slashes for consistency
relPath = strings.ReplaceAll(relPath, string(filepath.Separator), "/")
keys = append(keys, relPath)
}
return nil
})
if err != nil {
return nil, fmt.Errorf("walking directory: %w", err)
}
return keys, nil
}
// ListStream returns a channel of ObjectInfo for large result sets. Like
// List, it sends an error when the destination directory is missing.
func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan ObjectInfo {
ch := make(chan ObjectInfo)
go func() {
defer close(ch)
_, err := f.fs.Stat(f.basePath)
if err != nil {
ch <- ObjectInfo{Err: fmt.Errorf("checking destination directory: %w", err)}
return
}
basePath := f.fullPath(prefix)
// Check if the prefix exists
exists, err := afero.Exists(f.fs, basePath)
if err != nil {
ch <- ObjectInfo{Err: fmt.Errorf("checking path: %w", err)}
return
}
if !exists {
return // Empty channel for non-existent prefix
}
_ = afero.Walk(f.fs, basePath, func(path string, info os.FileInfo, err error) error {
// Check context cancellation
select {
case <-ctx.Done():
ch <- ObjectInfo{Err: ctx.Err()}
return ctx.Err()
default:
}
if err != nil {
ch <- ObjectInfo{Err: err}
return nil //nolint:nilerr // continue walking despite errors
}
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
relPath, err := filepath.Rel(f.basePath, path)
if err != nil {
ch <- ObjectInfo{Err: fmt.Errorf("computing relative path: %w", err)}
return nil
}
// Normalize path separators
relPath = strings.ReplaceAll(relPath, string(filepath.Separator), "/")
ch <- ObjectInfo{
Key: relPath,
Size: info.Size(),
}
}
return nil
})
}()
return ch
}
// Info returns human-readable storage location information.
func (f *FileStorer) Info() Info {
return Info{
Type: schemeFile,
Location: f.basePath,
}
}
// writeAtomic streams data into a temp file in the destination directory,
// fsyncs it, and renames it onto the final key. The key therefore appears
// only once the whole object has been durably written; a failure part-way
// leaves a temp file (removed here on the failing path) rather than a
// truncated object at the key.
func (f *FileStorer) writeAtomic(
key string, data io.Reader, progress ProgressCallback,
) error {
path := f.fullPath(key)
dir := filepath.Dir(path)
err := f.fs.MkdirAll(dir, storageDirPerm)
if err != nil {
return fmt.Errorf("creating directories: %w", err)
}
tmp, err := afero.TempFile(f.fs, dir, filepath.Base(path)+"-*"+tempSuffix)
if err != nil {
return fmt.Errorf("creating temp file: %w", err)
}
tmpPath := tmp.Name()
// Remove the temp file unless the rename below claims it. On the success
// path renamed is true, so the deferred Close and Remove are harmless
// no-ops on a name that no longer exists.
renamed := false
defer func() {
_ = tmp.Close()
if !renamed {
_ = f.fs.Remove(tmpPath)
}
}()
var w io.Writer = tmp
if progress != nil {
w = &progressWriter{writer: tmp, callback: progress}
}
_, err = io.Copy(w, data)
if err != nil {
return fmt.Errorf("writing file: %w", err)
}
err = tmp.Sync()
if err != nil {
return fmt.Errorf("syncing temp file: %w", err)
}
err = tmp.Close()
if err != nil {
return fmt.Errorf("closing temp file: %w", err)
}
err = f.fs.Rename(tmpPath, path)
if err != nil {
return fmt.Errorf("renaming temp file: %w", err)
}
renamed = true
return nil
}
// fullPath returns the full filesystem path for a key.
func (f *FileStorer) fullPath(key string) string {
return filepath.Join(f.basePath, key)
}
// progressWriter wraps an io.Writer to track write progress.
type progressWriter struct {
writer io.Writer
written int64
callback ProgressCallback
}
func (pw *progressWriter) Write(p []byte) (int, error) {
n, err := pw.writer.Write(p)
if n > 0 {
pw.written += int64(n)
if pw.callback != nil {
callbackErr := pw.callback(pw.written)
if callbackErr != nil {
return n, callbackErr
}
}
}
return n, err
}