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
344 lines
8.4 KiB
Go
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
|
|
}
|