Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
91e4bf92ab |
@@ -22,12 +22,14 @@ the tag exists and is exercised; what is left is merging `next` to
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 2026-10-08: Made `remote nuke` delete the `.partial` files that
|
||||
uploads cut off part-way leave on the destination store
|
||||
([issue #281](https://git.eeqj.de/sneak/vaultik/issues/281)). The
|
||||
file and rclone listings skip such a file, so the command left it in
|
||||
place and still reported the store empty. It now removes them under
|
||||
`metadata/` and `blobs/` as its last step.
|
||||
- 2026-10-08: Made a local index error while recording a directory or
|
||||
symlink stop a backup under `--skip-errors`
|
||||
([issue #284](https://git.eeqj.de/sneak/vaultik/issues/284)). Phase 2
|
||||
only records such an entry in the local index, since it has no data
|
||||
to open or read, but an error doing so was skipped like an unreadable
|
||||
file. The snapshot then completed without the entry, and a restore
|
||||
did not recreate it. The error now stops the backup, as it does
|
||||
without the flag.
|
||||
|
||||
- 2026-10-08: Counted a file that a backup could not store as failed
|
||||
([issue #280](https://git.eeqj.de/sneak/vaultik/issues/280)). A file
|
||||
|
||||
@@ -1348,11 +1348,25 @@ func (s *Scanner) processPhase(
|
||||
return s.finalizeProcessPhase(ctx, result)
|
||||
}
|
||||
|
||||
// processFileWithErrorHandling wraps processFileStreaming with error recovery for
|
||||
// deleted files and skip-errors mode. Returns (skipped, error).
|
||||
// processFileWithErrorHandling records a directory or symlink, or wraps
|
||||
// processFileStreaming for a regular file with error recovery for deleted
|
||||
// files and skip-errors mode. Returns (skipped, error).
|
||||
func (s *Scanner) processFileWithErrorHandling(
|
||||
ctx context.Context, fileToProcess *FileToProcess, result *ScanResult,
|
||||
) (bool, error) {
|
||||
// A directory or symlink has no data to open or read; it is only
|
||||
// recorded in the local index. An error recording it stops the run
|
||||
// even under --skip-errors, or the snapshot would complete without it.
|
||||
mode := os.FileMode(fileToProcess.File.Mode)
|
||||
if mode&os.ModeSymlink != 0 || mode.IsDir() {
|
||||
err := s.recordNonRegularFile(ctx, fileToProcess)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("processing file %s: %w", fileToProcess.Path, err)
|
||||
}
|
||||
|
||||
return false, nil
|
||||
}
|
||||
|
||||
err := s.processFileStreaming(ctx, fileToProcess, result)
|
||||
if err != nil {
|
||||
// A packer/database/encryption/upload failure means the chunk's data
|
||||
@@ -1392,12 +1406,7 @@ func (s *Scanner) processFileWithErrorHandling(
|
||||
|
||||
// countFailedFile counts a file that phase 2 could not store as failed
|
||||
// and takes its size back out of BytesScanned, where phase 1 put it.
|
||||
// Phase 1 counts no directories, so a directory is not counted here.
|
||||
func countFailedFile(fileToProcess *FileToProcess, result *ScanResult) {
|
||||
if fileToProcess.FileInfo.IsDir() {
|
||||
return
|
||||
}
|
||||
|
||||
result.FilesFailed++
|
||||
result.BytesScanned -= fileToProcess.FileInfo.Size()
|
||||
}
|
||||
@@ -1770,16 +1779,11 @@ func (e *packerError) Error() string { return e.err.Error() }
|
||||
|
||||
func (e *packerError) Unwrap() error { return e.err }
|
||||
|
||||
// processFileStreaming processes a file by streaming chunks directly to the packer
|
||||
// processFileStreaming processes a regular file by streaming chunks directly
|
||||
// to the packer
|
||||
func (s *Scanner) processFileStreaming(
|
||||
ctx context.Context, fileToProcess *FileToProcess, result *ScanResult,
|
||||
) error {
|
||||
// Symlinks and directories have no data to chunk — just record them in the DB.
|
||||
mode := os.FileMode(fileToProcess.File.Mode)
|
||||
if mode&os.ModeSymlink != 0 || mode.IsDir() {
|
||||
return s.recordNonRegularFile(ctx, fileToProcess)
|
||||
}
|
||||
|
||||
file, err := s.fs.Open(fileToProcess.Path)
|
||||
if err != nil {
|
||||
return fmt.Errorf("opening file: %w", wrapPermissionError(fileToProcess.Path, err))
|
||||
|
||||
@@ -171,11 +171,6 @@ func (f *Storer) ListStream(
|
||||
return f.inner.ListStream(ctx, prefix)
|
||||
}
|
||||
|
||||
// DeletePartialUploads delegates unchanged.
|
||||
func (f *Storer) DeletePartialUploads(ctx context.Context, prefix string) error {
|
||||
return f.inner.DeletePartialUploads(ctx, prefix)
|
||||
}
|
||||
|
||||
// Info delegates unchanged.
|
||||
func (f *Storer) Info() storage.Info {
|
||||
return f.inner.Info()
|
||||
|
||||
@@ -242,42 +242,6 @@ func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan Objec
|
||||
return ch
|
||||
}
|
||||
|
||||
// DeletePartialUploads removes every file under prefix whose name ends in
|
||||
// tempSuffix. A missing prefix has none to remove.
|
||||
func (f *FileStorer) DeletePartialUploads(ctx context.Context, prefix string) error {
|
||||
basePath := f.fullPath(prefix)
|
||||
|
||||
exists, err := afero.Exists(f.fs, basePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("checking path: %w", err)
|
||||
}
|
||||
|
||||
if !exists {
|
||||
return nil
|
||||
}
|
||||
|
||||
err = afero.Walk(f.fs, basePath, func(path string, info os.FileInfo, err error) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
if info.IsDir() || !strings.HasSuffix(info.Name(), tempSuffix) {
|
||||
return nil
|
||||
}
|
||||
|
||||
return f.fs.Remove(path)
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("walking directory: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Info returns human-readable storage location information.
|
||||
func (f *FileStorer) Info() Info {
|
||||
return Info{
|
||||
|
||||
@@ -120,45 +120,3 @@ func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
|
||||
t.Fatalf("ListStream should return only the real key, got %v", streamed)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFileStorer_DeletePartialUploads checks that a leftover temp file is
|
||||
// removed and the object at the real key is kept.
|
||||
func TestFileStorer_DeletePartialUploads(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
base := t.TempDir()
|
||||
|
||||
f, err := storage.NewFileStorer(base)
|
||||
if err != nil {
|
||||
t.Fatalf("NewFileStorer: %v", err)
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
err = f.Put(ctx, testBlobKey, strings.NewReader("blob-bytes"))
|
||||
if err != nil {
|
||||
t.Fatalf("Put: %v", err)
|
||||
}
|
||||
|
||||
leftover := filepath.Join(base, testBlobKey+"-123456.partial")
|
||||
|
||||
err = os.WriteFile(leftover, []byte("half"), 0o600)
|
||||
if err != nil {
|
||||
t.Fatalf("writing leftover temp file: %v", err)
|
||||
}
|
||||
|
||||
err = f.DeletePartialUploads(ctx, "blobs/")
|
||||
if err != nil {
|
||||
t.Fatalf("DeletePartialUploads: %v", err)
|
||||
}
|
||||
|
||||
_, err = os.Stat(leftover)
|
||||
if !os.IsNotExist(err) {
|
||||
t.Errorf("leftover temp file was not removed: %v", err)
|
||||
}
|
||||
|
||||
_, err = f.Stat(ctx, testBlobKey)
|
||||
if err != nil {
|
||||
t.Errorf("Stat of the real key: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -213,31 +213,6 @@ func (r *RcloneStorer) ListStream(
|
||||
return ch
|
||||
}
|
||||
|
||||
// DeletePartialUploads removes every object under prefix whose name ends
|
||||
// in tempSuffix.
|
||||
func (r *RcloneStorer) DeletePartialUploads(ctx context.Context, prefix string) error {
|
||||
var partial []fs.Object
|
||||
|
||||
err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) {
|
||||
key := obj.Remote()
|
||||
if strings.HasPrefix(key, prefix) && strings.HasSuffix(key, tempSuffix) {
|
||||
partial = append(partial, obj)
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("listing objects: %w", err)
|
||||
}
|
||||
|
||||
for _, obj := range partial {
|
||||
err = obj.Remove(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("removing object: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Info returns human-readable storage location information.
|
||||
func (r *RcloneStorer) Info() Info {
|
||||
location := r.remote
|
||||
|
||||
@@ -198,47 +198,6 @@ func TestRcloneStorerListSkipsPartialFiles(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestRcloneStorerDeletePartialUploads checks that a temporary file left
|
||||
// by a killed upload is removed and the object at the real key is kept.
|
||||
//
|
||||
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||
func TestRcloneStorerDeletePartialUploads(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
ctx := context.Background()
|
||||
|
||||
s, err := storage.NewRcloneStorer(ctx, ":local", dir)
|
||||
if err != nil {
|
||||
t.Fatalf("NewRcloneStorer: %v", err)
|
||||
}
|
||||
|
||||
err = s.Put(ctx, testBlobKey, strings.NewReader("blob-bytes"))
|
||||
if err != nil {
|
||||
t.Fatalf("Put: %v", err)
|
||||
}
|
||||
|
||||
leftover := filepath.Join(dir, testBlobKey+"-123456.partial")
|
||||
|
||||
err = os.WriteFile(leftover, []byte("half"), 0o600)
|
||||
if err != nil {
|
||||
t.Fatalf("writing leftover temp file: %v", err)
|
||||
}
|
||||
|
||||
err = s.DeletePartialUploads(ctx, "blobs/")
|
||||
if err != nil {
|
||||
t.Fatalf("DeletePartialUploads: %v", err)
|
||||
}
|
||||
|
||||
_, err = os.Stat(leftover)
|
||||
if !os.IsNotExist(err) {
|
||||
t.Errorf("leftover temp file was not removed: %v", err)
|
||||
}
|
||||
|
||||
_, err = s.Stat(ctx, testBlobKey)
|
||||
if err != nil {
|
||||
t.Errorf("Stat of the real key: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// newRcloneStorerOnWrappedLocal registers name as rclone's local backend
|
||||
// wrapped by wrap, and builds an rclone backend on it rooted at a fresh
|
||||
// temp directory. wrap changes the features the local backend reports, so
|
||||
|
||||
@@ -99,12 +99,6 @@ func (s *S3Storer) ListStream(ctx context.Context, prefix string) <-chan ObjectI
|
||||
return ch
|
||||
}
|
||||
|
||||
// DeletePartialUploads has nothing to remove: S3 shows an object only once
|
||||
// its upload has completed, so an upload cut off part-way leaves none.
|
||||
func (s *S3Storer) DeletePartialUploads(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Info returns human-readable storage location information.
|
||||
func (s *S3Storer) Info() Info {
|
||||
return Info{
|
||||
|
||||
@@ -71,12 +71,6 @@ type Storer interface {
|
||||
// If an error occurs during listing, the final item will have Err set.
|
||||
ListStream(ctx context.Context, prefix string) <-chan ObjectInfo
|
||||
|
||||
// DeletePartialUploads removes every object under prefix that an
|
||||
// upload cut off part-way left under a temporary name ending in
|
||||
// `.partial`. The file and rclone backends' List and ListStream skip
|
||||
// such an object.
|
||||
DeletePartialUploads(ctx context.Context, prefix string) error
|
||||
|
||||
// Info returns human-readable storage location information.
|
||||
Info() Info
|
||||
}
|
||||
|
||||
@@ -257,10 +257,6 @@ func (s *stubLister) List(_ context.Context, _ string) ([]string, error) {
|
||||
return nil, errStubUnused
|
||||
}
|
||||
|
||||
func (s *stubLister) DeletePartialUploads(_ context.Context, _ string) error {
|
||||
return errStubUnused
|
||||
}
|
||||
|
||||
func (s *stubLister) Info() storage.Info {
|
||||
return storage.Info{}
|
||||
}
|
||||
|
||||
@@ -155,12 +155,6 @@ func (m *MockStorer) ListStream(
|
||||
return ch
|
||||
}
|
||||
|
||||
// DeletePartialUploads has nothing to remove: Put stores each object
|
||||
// under its key at once.
|
||||
func (m *MockStorer) DeletePartialUploads(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *MockStorer) Info() storage.Info {
|
||||
return storage.Info{
|
||||
Type: "mock",
|
||||
|
||||
@@ -1,64 +0,0 @@
|
||||
package vaultik_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io/fs"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"sneak.berlin/go/vaultik/internal/snapshot"
|
||||
)
|
||||
|
||||
// TestNukeRemoteLeavesNoFiles checks that remote nuke leaves no file
|
||||
// under a file:// destination, including the `.partial` files uploads
|
||||
// killed part-way leave next to a snapshot's metadata and next to where a
|
||||
// blob would have been.
|
||||
func TestNukeRemoteLeavesNoFiles(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
storeDir := filepath.Join(t.TempDir(), "store")
|
||||
v, repos, _ := backUpToFileDestination(ctx, t, storeDir)
|
||||
|
||||
snapshots, err := repos.Snapshots.ListRecent(ctx, listRecentTestLimit)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, snapshots, 1)
|
||||
|
||||
snapshotKey := snapshot.RemoteSnapshotKey(snapshots[0].ID.String())
|
||||
hash := testBlobHashA
|
||||
leftovers := []string{
|
||||
filepath.Join(storeDir, "metadata", snapshotKey,
|
||||
"db.zst.age-123456.partial"),
|
||||
filepath.Join(storeDir, "blobs", hash[:2], hash[2:4],
|
||||
hash+"-123456.partial"),
|
||||
}
|
||||
|
||||
for _, leftover := range leftovers {
|
||||
require.NoError(t, os.MkdirAll(filepath.Dir(leftover), 0o750))
|
||||
require.NoError(t, os.WriteFile(leftover, []byte("half an upload"), 0o600))
|
||||
}
|
||||
|
||||
require.NoError(t, v.NukeRemote(true))
|
||||
|
||||
var files []string
|
||||
|
||||
err = filepath.WalkDir(storeDir,
|
||||
func(path string, entry fs.DirEntry, err error) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if !entry.IsDir() {
|
||||
files = append(files, path)
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
require.NoError(t, err)
|
||||
assert.Empty(t, files)
|
||||
}
|
||||
@@ -24,9 +24,8 @@ var errNukeRequiresForce = errors.New(
|
||||
const metadataDirName = "metadata"
|
||||
|
||||
// NukeRemote deletes every snapshot's metadata and every blob from remote
|
||||
// storage, along with any object an upload cut off part-way left under a
|
||||
// temporary `.partial` name. After this returns successfully the bucket
|
||||
// prefix is empty and the next backup starts from scratch.
|
||||
// storage. After this returns successfully the bucket prefix is empty and
|
||||
// the next backup starts from scratch.
|
||||
//
|
||||
// Refuses to run unless force is true. The caller is responsible for
|
||||
// confirming with the user.
|
||||
@@ -49,15 +48,6 @@ func (v *Vaultik) NukeRemote(force bool) error {
|
||||
return fmt.Errorf("pruning blobs: %w", err)
|
||||
}
|
||||
|
||||
// The file and rclone listings skip `.partial` objects, so the two
|
||||
// steps above never delete them.
|
||||
for _, prefix := range []string{"metadata/", "blobs/"} {
|
||||
err = v.Storage.DeletePartialUploads(v.ctx, prefix)
|
||||
if err != nil {
|
||||
return fmt.Errorf("deleting partial uploads: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
v.UI.Completef("Backup destination store is now empty.")
|
||||
|
||||
return nil
|
||||
|
||||
@@ -125,12 +125,6 @@ func (s *testStorer) ListStream(
|
||||
return ch
|
||||
}
|
||||
|
||||
// DeletePartialUploads has nothing to remove: Put stores each object
|
||||
// under its key at once.
|
||||
func (s *testStorer) DeletePartialUploads(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *testStorer) Info() storage.Info {
|
||||
return storage.Info{
|
||||
Type: testLabel,
|
||||
|
||||
@@ -0,0 +1,75 @@
|
||||
package vaultik_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/spf13/afero"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"sneak.berlin/go/vaultik/internal/config"
|
||||
"sneak.berlin/go/vaultik/internal/database"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
"sneak.berlin/go/vaultik/internal/vaultik"
|
||||
)
|
||||
|
||||
// --skip-errors skips only a file that cannot be opened or read. A local
|
||||
// index error while a directory is recorded stops the backup, and the
|
||||
// snapshot is not recorded as complete. See
|
||||
// https://git.eeqj.de/sneak/vaultik/issues/284.
|
||||
func TestSkipErrorsBackupStopsOnIndexErrorRecordingDirectory(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
// The scan walks the source path with symlinks resolved, so dirPath
|
||||
// must be spelled the same way to match the row the scan inserts.
|
||||
tempDir, err := filepath.EvalSymlinks(t.TempDir())
|
||||
require.NoError(t, err)
|
||||
|
||||
srcDir := filepath.Join(tempDir, "src")
|
||||
dirPath := filepath.Join(srcDir, "dir")
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
require.NoError(t, fs.MkdirAll(dirPath, 0o755))
|
||||
require.NoError(t, afero.WriteFile(fs,
|
||||
filepath.Join(dirPath, "file.txt"), []byte("file content"), 0o644))
|
||||
|
||||
cfg := faultTestConfig()
|
||||
cfg.IndexPath = filepath.Join(tempDir, "index.sqlite")
|
||||
cfg.Snapshots = map[string]config.SnapshotConfig{
|
||||
"tree": {Paths: []string{srcDir}},
|
||||
}
|
||||
|
||||
store, err := storage.NewFileStorer(filepath.Join(tempDir, "remote"))
|
||||
require.NoError(t, err)
|
||||
|
||||
db, err := database.New(ctx, cfg.IndexPath)
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = db.Close() })
|
||||
|
||||
// The local index refuses the directory's files row.
|
||||
_, err = db.Conn().ExecContext(ctx, fmt.Sprintf(`
|
||||
CREATE TRIGGER refuse_directory BEFORE INSERT ON files
|
||||
WHEN NEW.path = '%s'
|
||||
BEGIN SELECT RAISE(ABORT, 'simulated index error'); END`, dirPath))
|
||||
require.NoError(t, err)
|
||||
|
||||
repos := database.NewRepositories(db)
|
||||
v := newBackupVaultik(ctx, cfg, store, repos, db, fs)
|
||||
|
||||
err = v.CreateSnapshot(&vaultik.SnapshotCreateOptions{
|
||||
SkipErrors: true,
|
||||
Snapshots: []string{"tree"},
|
||||
})
|
||||
require.ErrorContains(t, err, "simulated index error")
|
||||
|
||||
snapshots, err := repos.Snapshots.ListRecent(ctx, listRecentTestLimit)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, snapshots, 1)
|
||||
assert.Nil(t, snapshots[0].CompletedAt)
|
||||
}
|
||||
Reference in New Issue
Block a user