Files
vaultik/internal/storage/file.go
T
clawbot ea72697992
check / check (push) Successful in 5m53s
check / check (pull_request) Successful in 4m53s
List a missing file:// destination directory as an error (closes #220)
The file backend listed a destination directory that does not exist as
an empty store. With the volume unplugged, snapshot list reported every
local snapshot as missing from the store, snapshot remove said it had
removed metadata it never reached, and prune dropped every local
snapshot record. List and ListStream now fail when the destination
directory is missing, so those commands take their existing path for a
store that cannot be listed. A missing prefix under an existing
directory is still an empty listing, and a first backup still creates
the directory.

Three tests listed a file:// destination nothing had created; they now
create it.

Model: opus-5-5
2026-10-06 08:46:17 +02:00

344 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. 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 and is overwritten when the same key is written again.
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
}