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. Reads (Get/Stat/List) tolerate a missing basePath — a missing // or unmounted destination during `snapshot list` should NOT block the // command, it should degrade to "no remote snapshots reachable" with a // warning. Write operations (Put/PutWithProgress) call MkdirAll for the // per-blob parent directory, which also covers basePath on first use. // // 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. func (f *FileStorer) List(ctx context.Context, prefix string) ([]string, error) { var keys []string basePath := f.fullPath(prefix) // Check if base path 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. func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan ObjectInfo { ch := make(chan ObjectInfo) go func() { defer close(ch) basePath := f.fullPath(prefix) // Check if base path 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 }