Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3543f45ba8 | ||
|
|
c355ef4d25 | ||
|
|
5927e1aa3d |
@@ -25,6 +25,16 @@ release" is exactly the contradiction
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
|
- 2026-09-21: Stopped `prune` from reporting a failed row count as 0
|
||||||
|
([issue #96](https://git.eeqj.de/sneak/vaultik/issues/96)). The seven
|
||||||
|
`getTableCount` reads in `PruneDatabase` discarded their error, so a
|
||||||
|
query that could not run became a plausible `0` and the before/after
|
||||||
|
delta computed from it looked like real work. Each read now logs at
|
||||||
|
warn on failure and renders as `unknown`, never `0`, so an empty table
|
||||||
|
is distinguishable from one that could not be queried. The counts have
|
||||||
|
no `--json` representation — under `--json` the summary is suppressed
|
||||||
|
entirely — so nothing there can show a false `0`.
|
||||||
|
|
||||||
- 2026-09-21: Made the s3 storage backend report a missing object as
|
- 2026-09-21: Made the s3 storage backend report a missing object as
|
||||||
`storage.ErrNotFound`, like the `file` and `rclone` backends and as the
|
`storage.ErrNotFound`, like the `file` and `rclone` backends and as the
|
||||||
`Storer` interface documents. `S3Storer.Get` and `Stat` returned the raw
|
`Storer` interface documents. `S3Storer.Get` and `Stat` returned the raw
|
||||||
@@ -33,7 +43,6 @@ release" is exactly the contradiction
|
|||||||
helper (reused by `HeadObject`) and a test that a missing key maps to
|
helper (reused by `HeadObject`) and a test that a missing key maps to
|
||||||
`ErrNotFound`
|
`ErrNotFound`
|
||||||
([issue #129](https://git.eeqj.de/sneak/vaultik/issues/129)).
|
([issue #129](https://git.eeqj.de/sneak/vaultik/issues/129)).
|
||||||
|
|
||||||
- 2026-09-21: Fixed `verify --deep` reporting healthy snapshots as
|
- 2026-09-21: Fixed `verify --deep` reporting healthy snapshots as
|
||||||
corrupt. Its final blob-integrity check hashed the encrypted
|
corrupt. Its final blob-integrity check hashed the encrypted
|
||||||
downloaded bytes with a single SHA256 and compared that to the blob
|
downloaded bytes with a single SHA256 and compared that to the blob
|
||||||
|
|||||||
@@ -0,0 +1,198 @@
|
|||||||
|
package storage_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"io"
|
||||||
|
"reflect"
|
||||||
|
"sort"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"sneak.berlin/go/vaultik/internal/storage"
|
||||||
|
)
|
||||||
|
|
||||||
|
// runStorerConformance is the shared Storer contract. Every backend that
|
||||||
|
// can run in-process is expected to pass it: TestFileStorer runs it against
|
||||||
|
// file://, TestS3Storer against s3://. A new backend inherits this coverage
|
||||||
|
// by passing its own constructor, so the contract is defined once.
|
||||||
|
//
|
||||||
|
// It exercises the public Storer interface: round-trip, stat, list with
|
||||||
|
// prefix filtering, overwrite, delete, delete-of-missing, and not-found on
|
||||||
|
// Get and Stat. Each section takes its own fresh backend instance, so the
|
||||||
|
// order of sections never matters and no section sees another's objects.
|
||||||
|
func runStorerConformance(t *testing.T, newStorer func(*testing.T) storage.Storer) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
conformanceRoundTrip(t, newStorer(t))
|
||||||
|
conformanceOverwrite(t, newStorer(t))
|
||||||
|
conformanceList(t, newStorer(t))
|
||||||
|
conformanceDelete(t, newStorer(t))
|
||||||
|
conformanceNotFound(t, newStorer(t))
|
||||||
|
}
|
||||||
|
|
||||||
|
// conformanceRoundTrip stores a nested key, then reads it back and stats it.
|
||||||
|
func conformanceRoundTrip(t *testing.T, s storage.Storer) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
key := "blobs/aa/bb/object.bin"
|
||||||
|
want := []byte("round-trip payload")
|
||||||
|
|
||||||
|
err := s.Put(ctx, key, bytes.NewReader(want))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Put: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
got := getBytes(t, s, key)
|
||||||
|
if !bytes.Equal(got, want) {
|
||||||
|
t.Errorf("Get returned %q, want %q", got, want)
|
||||||
|
}
|
||||||
|
|
||||||
|
info, err := s.Stat(ctx, key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Stat: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if info.Key != key {
|
||||||
|
t.Errorf("Stat key = %q, want %q", info.Key, key)
|
||||||
|
}
|
||||||
|
|
||||||
|
if info.Size != int64(len(want)) {
|
||||||
|
t.Errorf("Stat size = %d, want %d", info.Size, len(want))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// conformanceOverwrite checks that a second Put replaces the first.
|
||||||
|
func conformanceOverwrite(t *testing.T, s storage.Storer) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
key := "meta/snapshot.json"
|
||||||
|
|
||||||
|
err := s.Put(ctx, key, bytes.NewReader([]byte("first")))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("first Put: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
want := []byte("second and longer payload")
|
||||||
|
|
||||||
|
err = s.Put(ctx, key, bytes.NewReader(want))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("second Put: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
got := getBytes(t, s, key)
|
||||||
|
if !bytes.Equal(got, want) {
|
||||||
|
t.Errorf("after overwrite Get returned %q, want %q", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// conformanceList checks prefix filtering and the empty result for a
|
||||||
|
// prefix that matches nothing.
|
||||||
|
func conformanceList(t *testing.T, s storage.Storer) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
keys := []string{"blobs/aa/one", "blobs/bb/two", "meta/three"}
|
||||||
|
|
||||||
|
for _, k := range keys {
|
||||||
|
err := s.Put(ctx, k, bytes.NewReader([]byte("data")))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Put %q: %v", k, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if got := listSorted(t, s, ""); !reflect.DeepEqual(got, keys) {
|
||||||
|
t.Errorf("List(\"\") = %v, want %v", got, keys)
|
||||||
|
}
|
||||||
|
|
||||||
|
wantBlobs := []string{"blobs/aa/one", "blobs/bb/two"}
|
||||||
|
if got := listSorted(t, s, "blobs/"); !reflect.DeepEqual(got, wantBlobs) {
|
||||||
|
t.Errorf("List(\"blobs/\") = %v, want %v", got, wantBlobs)
|
||||||
|
}
|
||||||
|
|
||||||
|
if got := listSorted(t, s, "absent/"); len(got) != 0 {
|
||||||
|
t.Errorf("List(\"absent/\") = %v, want empty", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// conformanceDelete checks that Delete removes an object and that deleting
|
||||||
|
// a missing key is not an error.
|
||||||
|
func conformanceDelete(t *testing.T, s storage.Storer) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
key := "blobs/cc/gone.bin"
|
||||||
|
|
||||||
|
err := s.Put(ctx, key, bytes.NewReader([]byte("temporary")))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Put: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = s.Delete(ctx, key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Delete: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = s.Get(ctx, key)
|
||||||
|
if !errors.Is(err, storage.ErrNotFound) {
|
||||||
|
t.Errorf("Get after Delete error = %v, want ErrNotFound", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = s.Delete(ctx, key)
|
||||||
|
if err != nil {
|
||||||
|
t.Errorf("Delete of missing key = %v, want nil", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// conformanceNotFound checks Get and Stat on an absent key.
|
||||||
|
func conformanceNotFound(t *testing.T, s storage.Storer) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
key := "never/written"
|
||||||
|
|
||||||
|
_, err := s.Get(ctx, key)
|
||||||
|
if !errors.Is(err, storage.ErrNotFound) {
|
||||||
|
t.Errorf("Get error = %v, want ErrNotFound", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = s.Stat(ctx, key)
|
||||||
|
if !errors.Is(err, storage.ErrNotFound) {
|
||||||
|
t.Errorf("Stat error = %v, want ErrNotFound", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// getBytes reads a key fully and closes the reader.
|
||||||
|
func getBytes(t *testing.T, s storage.Storer, key string) []byte {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
rc, err := s.Get(context.Background(), key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Get %q: %v", key, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
defer func() { _ = rc.Close() }()
|
||||||
|
|
||||||
|
data, err := io.ReadAll(rc)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("read %q: %v", key, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return data
|
||||||
|
}
|
||||||
|
|
||||||
|
// listSorted returns the keys under a prefix in a stable order.
|
||||||
|
func listSorted(t *testing.T, s storage.Storer, prefix string) []string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
keys, err := s.List(context.Background(), prefix)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("List %q: %v", prefix, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
sort.Strings(keys)
|
||||||
|
|
||||||
|
return keys
|
||||||
|
}
|
||||||
+79
-54
@@ -46,31 +46,18 @@ func (f *FileStorer) SetFilesystem(fs afero.Fs) {
|
|||||||
// storage base path.
|
// storage base path.
|
||||||
const storageDirPerm = 0o755
|
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.
|
// Put stores data at the specified key.
|
||||||
func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error {
|
func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error {
|
||||||
path := f.fullPath(key)
|
return f.writeAtomic(key, data, nil)
|
||||||
|
|
||||||
// Create parent directories
|
|
||||||
dir := filepath.Dir(path)
|
|
||||||
|
|
||||||
err := f.fs.MkdirAll(dir, storageDirPerm)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("creating directories: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
file, err := f.fs.Create(path)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("creating file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
defer func() { _ = file.Close() }()
|
|
||||||
|
|
||||||
_, err = io.Copy(file, data)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("writing file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// PutWithProgress stores data with progress reporting.
|
// PutWithProgress stores data with progress reporting.
|
||||||
@@ -78,35 +65,7 @@ func (f *FileStorer) PutWithProgress(
|
|||||||
_ context.Context, key string, data io.Reader,
|
_ context.Context, key string, data io.Reader,
|
||||||
_ int64, progress ProgressCallback,
|
_ int64, progress ProgressCallback,
|
||||||
) error {
|
) error {
|
||||||
path := f.fullPath(key)
|
return f.writeAtomic(key, data, progress)
|
||||||
|
|
||||||
// Create parent directories
|
|
||||||
dir := filepath.Dir(path)
|
|
||||||
|
|
||||||
err := f.fs.MkdirAll(dir, storageDirPerm)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("creating directories: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
file, err := f.fs.Create(path)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("creating file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
defer func() { _ = file.Close() }()
|
|
||||||
|
|
||||||
// Wrap with progress tracking
|
|
||||||
pw := &progressWriter{
|
|
||||||
writer: file,
|
|
||||||
callback: progress,
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = io.Copy(pw, data)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("writing file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get retrieves data from the specified key.
|
// Get retrieves data from the specified key.
|
||||||
@@ -188,7 +147,7 @@ func (f *FileStorer) List(ctx context.Context, prefix string) ([]string, error)
|
|||||||
default:
|
default:
|
||||||
}
|
}
|
||||||
|
|
||||||
if !info.IsDir() {
|
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
|
||||||
// Convert back to key (relative path from basePath)
|
// Convert back to key (relative path from basePath)
|
||||||
relPath, err := filepath.Rel(f.basePath, path)
|
relPath, err := filepath.Rel(f.basePath, path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -245,7 +204,7 @@ func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan Objec
|
|||||||
return nil //nolint:nilerr // continue walking despite errors
|
return nil //nolint:nilerr // continue walking despite errors
|
||||||
}
|
}
|
||||||
|
|
||||||
if !info.IsDir() {
|
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
|
||||||
relPath, err := filepath.Rel(f.basePath, path)
|
relPath, err := filepath.Rel(f.basePath, path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
ch <- ObjectInfo{Err: fmt.Errorf("computing relative path: %w", err)}
|
ch <- ObjectInfo{Err: fmt.Errorf("computing relative path: %w", err)}
|
||||||
@@ -275,6 +234,72 @@ func (f *FileStorer) Info() Info {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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.
|
// fullPath returns the full filesystem path for a key.
|
||||||
func (f *FileStorer) fullPath(key string) string {
|
func (f *FileStorer) fullPath(key string) string {
|
||||||
return filepath.Join(f.basePath, key)
|
return filepath.Join(f.basePath, key)
|
||||||
|
|||||||
@@ -0,0 +1,119 @@
|
|||||||
|
package storage_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"sneak.berlin/go/vaultik/internal/storage"
|
||||||
|
)
|
||||||
|
|
||||||
|
// errStreamInterrupted stands in for an upload cut off mid-stream.
|
||||||
|
var errStreamInterrupted = errors.New("connection reset mid-upload")
|
||||||
|
|
||||||
|
// failingReader yields its data once, then fails.
|
||||||
|
type failingReader struct {
|
||||||
|
data []byte
|
||||||
|
done bool
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *failingReader) Read(p []byte) (int, error) {
|
||||||
|
if r.done {
|
||||||
|
return 0, errStreamInterrupted
|
||||||
|
}
|
||||||
|
|
||||||
|
n := copy(p, r.data)
|
||||||
|
r.done = true
|
||||||
|
|
||||||
|
return n, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestFileStorer_InterruptedWriteLeavesNoTrustedObject checks that a write
|
||||||
|
// cut off mid-stream leaves nothing at the destination key, so a later run
|
||||||
|
// cannot Stat a truncated object and trust it as a complete blob.
|
||||||
|
func TestFileStorer_InterruptedWriteLeavesNoTrustedObject(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
f, err := storage.NewFileStorer(t.TempDir())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("NewFileStorer: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
key := "blobs/aa/bb/aabbccddeeff"
|
||||||
|
|
||||||
|
err = f.PutWithProgress(ctx, key, &failingReader{data: []byte("partial")}, 4096, nil)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected the interrupted write to fail, got nil")
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = f.Stat(ctx, key)
|
||||||
|
if !errors.Is(err, storage.ErrNotFound) {
|
||||||
|
t.Fatalf("expected key absent after interrupted write, got Stat err %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
keys, err := f.List(ctx, "blobs/")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("List: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(keys) != 0 {
|
||||||
|
t.Fatalf("expected no keys listed after interrupted write, got %v", keys)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestFileStorer_ListSkipsPartialFiles checks that a leftover temp file (the
|
||||||
|
// storage layer names them with a ".partial" suffix) is never surfaced as a
|
||||||
|
// key by List or ListStream.
|
||||||
|
func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
base := t.TempDir()
|
||||||
|
|
||||||
|
f, err := storage.NewFileStorer(base)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("NewFileStorer: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
realKey := "blobs/aa/bb/aabbccddeeff"
|
||||||
|
|
||||||
|
err = f.Put(ctx, realKey, strings.NewReader("blob-bytes"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Put: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A stray temp file, as an interrupted write would leave behind.
|
||||||
|
leftover := filepath.Join(base, "blobs/aa/bb/aabbccddeeff-123456.partial")
|
||||||
|
|
||||||
|
err = os.WriteFile(leftover, []byte("half"), 0o600)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("writing leftover temp file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
keys, err := f.List(ctx, "blobs/")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("List: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(keys) != 1 || keys[0] != realKey {
|
||||||
|
t.Fatalf("List should return only the real key, got %v", keys)
|
||||||
|
}
|
||||||
|
|
||||||
|
var streamed []string
|
||||||
|
|
||||||
|
for obj := range f.ListStream(ctx, "blobs/") {
|
||||||
|
if obj.Err != nil {
|
||||||
|
t.Fatalf("ListStream: %v", obj.Err)
|
||||||
|
}
|
||||||
|
|
||||||
|
streamed = append(streamed, obj.Key)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(streamed) != 1 || streamed[0] != realKey {
|
||||||
|
t.Fatalf("ListStream should return only the real key, got %v", streamed)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,12 +1,6 @@
|
|||||||
package storage_test
|
package storage_test
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"io"
|
|
||||||
"reflect"
|
|
||||||
"sort"
|
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"sneak.berlin/go/vaultik/internal/storage"
|
"sneak.berlin/go/vaultik/internal/storage"
|
||||||
@@ -26,189 +20,8 @@ func newFileStorer(t *testing.T) storage.Storer {
|
|||||||
return s
|
return s
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestFileStorer runs the Storer contract against the file:// backend.
|
// TestFileStorer runs the shared Storer contract against the file:// backend.
|
||||||
// The conformance helper is backend-agnostic, so a new backend inherits
|
|
||||||
// this coverage by passing its own constructor.
|
|
||||||
func TestFileStorer(t *testing.T) {
|
func TestFileStorer(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
runStorerConformance(t, newFileStorer)
|
runStorerConformance(t, newFileStorer)
|
||||||
}
|
}
|
||||||
|
|
||||||
// runStorerConformance exercises the public Storer contract: round-trip,
|
|
||||||
// stat, list, overwrite, delete, and not-found behaviour. Each section
|
|
||||||
// uses its own backend instance so ordering never matters.
|
|
||||||
func runStorerConformance(t *testing.T, newStorer func(*testing.T) storage.Storer) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
conformanceRoundTrip(t, newStorer(t))
|
|
||||||
conformanceOverwrite(t, newStorer(t))
|
|
||||||
conformanceList(t, newStorer(t))
|
|
||||||
conformanceDelete(t, newStorer(t))
|
|
||||||
conformanceNotFound(t, newStorer(t))
|
|
||||||
}
|
|
||||||
|
|
||||||
// conformanceRoundTrip stores a nested key, then reads it back and stats it.
|
|
||||||
func conformanceRoundTrip(t *testing.T, s storage.Storer) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
key := "blobs/aa/bb/object.bin"
|
|
||||||
want := []byte("round-trip payload")
|
|
||||||
|
|
||||||
err := s.Put(ctx, key, bytes.NewReader(want))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Put: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
got := getBytes(t, s, key)
|
|
||||||
if !bytes.Equal(got, want) {
|
|
||||||
t.Errorf("Get returned %q, want %q", got, want)
|
|
||||||
}
|
|
||||||
|
|
||||||
info, err := s.Stat(ctx, key)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Stat: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if info.Key != key {
|
|
||||||
t.Errorf("Stat key = %q, want %q", info.Key, key)
|
|
||||||
}
|
|
||||||
|
|
||||||
if info.Size != int64(len(want)) {
|
|
||||||
t.Errorf("Stat size = %d, want %d", info.Size, len(want))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// conformanceOverwrite checks that a second Put replaces the first.
|
|
||||||
func conformanceOverwrite(t *testing.T, s storage.Storer) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
key := "meta/snapshot.json"
|
|
||||||
|
|
||||||
err := s.Put(ctx, key, bytes.NewReader([]byte("first")))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("first Put: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
want := []byte("second and longer payload")
|
|
||||||
|
|
||||||
err = s.Put(ctx, key, bytes.NewReader(want))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("second Put: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
got := getBytes(t, s, key)
|
|
||||||
if !bytes.Equal(got, want) {
|
|
||||||
t.Errorf("after overwrite Get returned %q, want %q", got, want)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// conformanceList checks prefix filtering and the empty result for a
|
|
||||||
// prefix that matches nothing.
|
|
||||||
func conformanceList(t *testing.T, s storage.Storer) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
keys := []string{"blobs/aa/one", "blobs/bb/two", "meta/three"}
|
|
||||||
|
|
||||||
for _, k := range keys {
|
|
||||||
err := s.Put(ctx, k, bytes.NewReader([]byte("data")))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Put %q: %v", k, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if got := listSorted(t, s, ""); !reflect.DeepEqual(got, keys) {
|
|
||||||
t.Errorf("List(\"\") = %v, want %v", got, keys)
|
|
||||||
}
|
|
||||||
|
|
||||||
wantBlobs := []string{"blobs/aa/one", "blobs/bb/two"}
|
|
||||||
if got := listSorted(t, s, "blobs/"); !reflect.DeepEqual(got, wantBlobs) {
|
|
||||||
t.Errorf("List(\"blobs/\") = %v, want %v", got, wantBlobs)
|
|
||||||
}
|
|
||||||
|
|
||||||
if got := listSorted(t, s, "absent/"); len(got) != 0 {
|
|
||||||
t.Errorf("List(\"absent/\") = %v, want empty", got)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// conformanceDelete checks that Delete removes an object and that deleting
|
|
||||||
// a missing key is not an error.
|
|
||||||
func conformanceDelete(t *testing.T, s storage.Storer) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
key := "blobs/cc/gone.bin"
|
|
||||||
|
|
||||||
err := s.Put(ctx, key, bytes.NewReader([]byte("temporary")))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Put: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = s.Delete(ctx, key)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Delete: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = s.Get(ctx, key)
|
|
||||||
if !errors.Is(err, storage.ErrNotFound) {
|
|
||||||
t.Errorf("Get after Delete error = %v, want ErrNotFound", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = s.Delete(ctx, key)
|
|
||||||
if err != nil {
|
|
||||||
t.Errorf("Delete of missing key = %v, want nil", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// conformanceNotFound checks Get and Stat on an absent key.
|
|
||||||
func conformanceNotFound(t *testing.T, s storage.Storer) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
key := "never/written"
|
|
||||||
|
|
||||||
_, err := s.Get(ctx, key)
|
|
||||||
if !errors.Is(err, storage.ErrNotFound) {
|
|
||||||
t.Errorf("Get error = %v, want ErrNotFound", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = s.Stat(ctx, key)
|
|
||||||
if !errors.Is(err, storage.ErrNotFound) {
|
|
||||||
t.Errorf("Stat error = %v, want ErrNotFound", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// getBytes reads a key fully and closes the reader.
|
|
||||||
func getBytes(t *testing.T, s storage.Storer, key string) []byte {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
rc, err := s.Get(context.Background(), key)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Get %q: %v", key, err)
|
|
||||||
}
|
|
||||||
|
|
||||||
defer func() { _ = rc.Close() }()
|
|
||||||
|
|
||||||
data, err := io.ReadAll(rc)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("read %q: %v", key, err)
|
|
||||||
}
|
|
||||||
|
|
||||||
return data
|
|
||||||
}
|
|
||||||
|
|
||||||
// listSorted returns the keys under a prefix in a stable order.
|
|
||||||
func listSorted(t *testing.T, s storage.Storer, prefix string) []string {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
keys, err := s.List(context.Background(), prefix)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("List %q: %v", prefix, err)
|
|
||||||
}
|
|
||||||
|
|
||||||
sort.Strings(keys)
|
|
||||||
|
|
||||||
return keys
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -0,0 +1,58 @@
|
|||||||
|
package storage_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"sneak.berlin/go/vaultik/internal/storage"
|
||||||
|
)
|
||||||
|
|
||||||
|
// The rclone backend is a thin adapter over the rclone library: it turns a
|
||||||
|
// (remote, path) pair into rclone's "remote:path" string, hands it to
|
||||||
|
// rclone, and maps rclone's own results back to the Storer interface. What
|
||||||
|
// can be tested in-process, without a configured remote or network, is that
|
||||||
|
// adapter layer — how the arguments are shaped and how construction errors
|
||||||
|
// are reported. The data-plane operations (Put/Get/List/Delete) are rclone's
|
||||||
|
// own, exercised against a real provider (drive, s3-via-rclone, ...), which
|
||||||
|
// needs a configured remote with credentials and network access and so is
|
||||||
|
// out of reach of a unit test. The shared Storer conformance suite therefore
|
||||||
|
// runs against the in-process file and s3 backends; the rclone backend
|
||||||
|
// inherits that contract once a remote is configured.
|
||||||
|
//
|
||||||
|
// These tests use rclone's ":local:" on-the-fly backend, which addresses the
|
||||||
|
// local filesystem directly without any configured remote, so construction
|
||||||
|
// runs entirely in-process.
|
||||||
|
|
||||||
|
// TestNewRcloneStorerConstruction checks that a valid remote constructs a
|
||||||
|
// backend and that Info() reports the shaped "remote:path" location.
|
||||||
|
//
|
||||||
|
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||||
|
func TestNewRcloneStorerConstruction(t *testing.T) {
|
||||||
|
dir := t.TempDir()
|
||||||
|
|
||||||
|
s, err := storage.NewRcloneStorer(context.Background(), ":local", dir)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("NewRcloneStorer: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Info().Location is the "remote:path" string the adapter builds from
|
||||||
|
// its two arguments, so asserting it confirms the argument shaping.
|
||||||
|
want := ":local:" + dir
|
||||||
|
if got := s.Info().Location; got != want {
|
||||||
|
t.Errorf("Info().Location = %q, want %q", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestNewRcloneStorerUnknownRemote checks that a remote that is not in the
|
||||||
|
// rclone config fails construction with the ErrRemoteNotFound sentinel,
|
||||||
|
// rather than silently returning a backend pointed nowhere.
|
||||||
|
//
|
||||||
|
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||||
|
func TestNewRcloneStorerUnknownRemote(t *testing.T) {
|
||||||
|
_, err := storage.NewRcloneStorer(
|
||||||
|
context.Background(), "vaultik-no-such-remote", "path")
|
||||||
|
if !errors.Is(err, storage.ErrRemoteNotFound) {
|
||||||
|
t.Errorf("NewRcloneStorer error = %v, want ErrRemoteNotFound", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
+36
-14
@@ -13,18 +13,23 @@ import (
|
|||||||
"sneak.berlin/go/vaultik/internal/storage"
|
"sneak.berlin/go/vaultik/internal/storage"
|
||||||
)
|
)
|
||||||
|
|
||||||
// TestS3StorerMissingKeyMapsToErrNotFound verifies that the s3 backend reports
|
// s3TestBucket is the bucket created for each in-process S3 server.
|
||||||
// a missing object as storage.ErrNotFound, matching the file and rclone
|
const s3TestBucket = "test-bucket"
|
||||||
// backends and the Storer contract. Without the mapping, Get and Stat leak the
|
|
||||||
// raw SDK error and errors.Is(err, storage.ErrNotFound) is false.
|
// newS3Storer builds an s3:// backend backed by a fresh in-process
|
||||||
|
// S3 server. It reuses the same in-memory S3 harness (gofakes3 + s3mem
|
||||||
|
// over httptest) that internal/s3 and the not-found regression test use,
|
||||||
|
// so no new mock or dependency is introduced. Each call gets its own
|
||||||
|
// server, bucket, and client, so the conformance suite's per-section
|
||||||
|
// instances stay isolated.
|
||||||
//
|
//
|
||||||
//nolint:paralleltest // shares an in-process S3 server via t.Cleanup
|
//nolint:ireturn // conformance runs against the Storer interface by design
|
||||||
func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
func newS3Storer(t *testing.T) storage.Storer {
|
||||||
const bucket = "test-bucket"
|
t.Helper()
|
||||||
|
|
||||||
backend := s3mem.New()
|
backend := s3mem.New()
|
||||||
|
|
||||||
err := backend.CreateBucket(bucket)
|
err := backend.CreateBucket(s3TestBucket)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("create bucket: %v", err)
|
t.Fatalf("create bucket: %v", err)
|
||||||
}
|
}
|
||||||
@@ -32,11 +37,9 @@ func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
|||||||
srv := httptest.NewServer(gofakes3.New(backend).Server())
|
srv := httptest.NewServer(gofakes3.New(backend).Server())
|
||||||
t.Cleanup(srv.Close)
|
t.Cleanup(srv.Close)
|
||||||
|
|
||||||
ctx := context.Background()
|
client, err := s3.NewClient(context.Background(), s3.Config{
|
||||||
|
|
||||||
client, err := s3.NewClient(ctx, s3.Config{
|
|
||||||
Endpoint: srv.URL,
|
Endpoint: srv.URL,
|
||||||
Bucket: bucket,
|
Bucket: s3TestBucket,
|
||||||
AccessKeyID: "test",
|
AccessKeyID: "test",
|
||||||
SecretAccessKey: "test",
|
SecretAccessKey: "test",
|
||||||
Region: "us-east-1",
|
Region: "us-east-1",
|
||||||
@@ -45,9 +48,28 @@ func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
|||||||
t.Fatalf("new client: %v", err)
|
t.Fatalf("new client: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
storer := storage.NewS3Storer(client)
|
return storage.NewS3Storer(client)
|
||||||
|
}
|
||||||
|
|
||||||
_, err = storer.Get(ctx, "does-not-exist")
|
// TestS3Storer runs the shared Storer contract against the s3:// backend,
|
||||||
|
// so it is held to the same round-trip, list, delete, and not-found
|
||||||
|
// behaviour as the file:// backend.
|
||||||
|
func TestS3Storer(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
runStorerConformance(t, newS3Storer)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestS3StorerMissingKeyMapsToErrNotFound pins the specific contract that a
|
||||||
|
// missing object surfaces as storage.ErrNotFound rather than the raw AWS SDK
|
||||||
|
// error. Without the mapping, errors.Is(err, storage.ErrNotFound) is false on
|
||||||
|
// s3 and callers would branch differently per backend.
|
||||||
|
func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
storer := newS3Storer(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
_, err := storer.Get(ctx, "does-not-exist")
|
||||||
if !errors.Is(err, storage.ErrNotFound) {
|
if !errors.Is(err, storage.ErrNotFound) {
|
||||||
t.Errorf("Get on missing key: got %v, want ErrNotFound", err)
|
t.Errorf("Get on missing key: got %v, want ErrNotFound", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,79 @@
|
|||||||
|
package vaultik //nolint:testpackage // exercises unexported count helpers
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"sneak.berlin/go/vaultik/internal/database"
|
||||||
|
"sneak.berlin/go/vaultik/internal/log"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestTableCountForReportSurfacesReadFailure is the regression guard for
|
||||||
|
// the discarded-error bug: getTableCount for a table its query cannot
|
||||||
|
// resolve must not silently become 0. A count that could not be read is
|
||||||
|
// reported as unknown, which a reader can tell apart from an empty table.
|
||||||
|
//
|
||||||
|
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||||
|
func TestTableCountForReportSurfacesReadFailure(t *testing.T) {
|
||||||
|
log.Initialize(log.Config{})
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
db, err := database.New(ctx, ":memory:")
|
||||||
|
require.NoError(t, err)
|
||||||
|
t.Cleanup(func() { _ = db.Close() })
|
||||||
|
|
||||||
|
v := &Vaultik{DB: db}
|
||||||
|
v.SetContext(ctx)
|
||||||
|
|
||||||
|
// A table present in the schema reads as a real count.
|
||||||
|
blobs := v.tableCountForReport("blobs")
|
||||||
|
require.NotNil(t, blobs, "an existing table must read as a real count")
|
||||||
|
assert.Equal(t, int64(0), *blobs)
|
||||||
|
|
||||||
|
// A syntactically valid name the sanitizer accepts but whose table
|
||||||
|
// the query cannot resolve is the exact shape #96 describes: a
|
||||||
|
// would-be loud failure that used to be discarded into a 0.
|
||||||
|
_, err = v.getTableCount("snapshots_missing")
|
||||||
|
require.Error(t, err, "a query against a nonexistent table must fail")
|
||||||
|
|
||||||
|
missing := v.tableCountForReport("snapshots_missing")
|
||||||
|
assert.Nil(t, missing, "a failed read is unknown, not a count")
|
||||||
|
|
||||||
|
// The rendered count for a failed read must say unknown, never 0.
|
||||||
|
assert.Equal(t, countUnknown, countText(missing))
|
||||||
|
assert.NotEqual(t, "0", countText(missing))
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestCountTextDistinguishesEmptyFromUnknown pins the distinction the
|
||||||
|
// output has to preserve: 0 means the table was empty, "unknown" means
|
||||||
|
// the count could not be read.
|
||||||
|
func TestCountTextDistinguishesEmptyFromUnknown(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
zero := int64(0)
|
||||||
|
seven := int64(7)
|
||||||
|
|
||||||
|
assert.Equal(t, "0", countText(&zero))
|
||||||
|
assert.Equal(t, "7", countText(&seven))
|
||||||
|
assert.Equal(t, countUnknown, countText(nil))
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestCountDiffUnknownWhenEitherSideUnknown checks that a delta computed
|
||||||
|
// from an unreadable count is itself unknown rather than a plausible
|
||||||
|
// number.
|
||||||
|
func TestCountDiffUnknownWhenEitherSideUnknown(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
before := int64(10)
|
||||||
|
after := int64(3)
|
||||||
|
|
||||||
|
require.NotNil(t, countDiff(&before, &after))
|
||||||
|
assert.Equal(t, int64(7), *countDiff(&before, &after))
|
||||||
|
|
||||||
|
assert.Nil(t, countDiff(nil, &after), "unknown before yields unknown delta")
|
||||||
|
assert.Nil(t, countDiff(&before, nil), "unknown after yields unknown delta")
|
||||||
|
assert.Nil(t, countDiff(nil, nil))
|
||||||
|
}
|
||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"regexp"
|
"regexp"
|
||||||
"sort"
|
"sort"
|
||||||
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -1540,12 +1541,17 @@ func (v *Vaultik) outputRemoveJSON(result *RemoveResult) error {
|
|||||||
return encoder.Encode(result)
|
return encoder.Encode(result)
|
||||||
}
|
}
|
||||||
|
|
||||||
// PruneResult contains statistics about the prune operation
|
// PruneResult contains statistics about the prune operation.
|
||||||
|
// SnapshotsDeleted counts snapshots actually deleted. FilesDeleted,
|
||||||
|
// ChunksDeleted, and BlobsDeleted are derived from before/after row
|
||||||
|
// counts of the local index; each is nil when a count could not be read,
|
||||||
|
// so an unreadable count is reported as unknown rather than silently
|
||||||
|
// as 0.
|
||||||
type PruneResult struct {
|
type PruneResult struct {
|
||||||
SnapshotsDeleted int64
|
SnapshotsDeleted int64
|
||||||
FilesDeleted int64
|
FilesDeleted *int64
|
||||||
ChunksDeleted int64
|
ChunksDeleted *int64
|
||||||
BlobsDeleted int64
|
BlobsDeleted *int64
|
||||||
}
|
}
|
||||||
|
|
||||||
// PruneDatabase removes incomplete snapshots and orphaned files, chunks,
|
// PruneDatabase removes incomplete snapshots and orphaned files, chunks,
|
||||||
@@ -1560,7 +1566,7 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
|||||||
result := &PruneResult{}
|
result := &PruneResult{}
|
||||||
|
|
||||||
// Snapshot counts before deletion of incompletes.
|
// Snapshot counts before deletion of incompletes.
|
||||||
snapshotCountBefore, _ := v.getTableCount("snapshots")
|
snapshotCountBefore := v.tableCountForReport("snapshots")
|
||||||
|
|
||||||
// First, delete any incomplete snapshots
|
// First, delete any incomplete snapshots
|
||||||
incompleteSnapshots, err := v.Repositories.Snapshots.GetIncompleteSnapshots(v.ctx)
|
incompleteSnapshots, err := v.Repositories.Snapshots.GetIncompleteSnapshots(v.ctx)
|
||||||
@@ -1575,9 +1581,9 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Get counts before cleanup for reporting
|
// Get counts before cleanup for reporting
|
||||||
fileCountBefore, _ := v.getTableCount("files")
|
fileCountBefore := v.tableCountForReport("files")
|
||||||
chunkCountBefore, _ := v.getTableCount("chunks")
|
chunkCountBefore := v.tableCountForReport("chunks")
|
||||||
blobCountBefore, _ := v.getTableCount("blobs")
|
blobCountBefore := v.tableCountForReport("blobs")
|
||||||
|
|
||||||
// Run the cleanup
|
// Run the cleanup
|
||||||
err = v.SnapshotManager.CleanupOrphanedData(v.ctx)
|
err = v.SnapshotManager.CleanupOrphanedData(v.ctx)
|
||||||
@@ -1586,36 +1592,83 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Get counts after cleanup
|
// Get counts after cleanup
|
||||||
fileCountAfter, _ := v.getTableCount("files")
|
fileCountAfter := v.tableCountForReport("files")
|
||||||
chunkCountAfter, _ := v.getTableCount("chunks")
|
chunkCountAfter := v.tableCountForReport("chunks")
|
||||||
blobCountAfter, _ := v.getTableCount("blobs")
|
blobCountAfter := v.tableCountForReport("blobs")
|
||||||
|
|
||||||
result.FilesDeleted = fileCountBefore - fileCountAfter
|
result.FilesDeleted = countDiff(fileCountBefore, fileCountAfter)
|
||||||
result.ChunksDeleted = chunkCountBefore - chunkCountAfter
|
result.ChunksDeleted = countDiff(chunkCountBefore, chunkCountAfter)
|
||||||
result.BlobsDeleted = blobCountBefore - blobCountAfter
|
result.BlobsDeleted = countDiff(blobCountBefore, blobCountAfter)
|
||||||
|
|
||||||
log.Info("Local database prune complete",
|
log.Info("Local database prune complete",
|
||||||
"incomplete_snapshots", result.SnapshotsDeleted,
|
"incomplete_snapshots", result.SnapshotsDeleted,
|
||||||
"orphaned_files", result.FilesDeleted,
|
"orphaned_files", countText(result.FilesDeleted),
|
||||||
"orphaned_chunks", result.ChunksDeleted,
|
"orphaned_chunks", countText(result.ChunksDeleted),
|
||||||
"orphaned_blobs", result.BlobsDeleted,
|
"orphaned_blobs", countText(result.BlobsDeleted),
|
||||||
)
|
)
|
||||||
|
|
||||||
snapshotCountAfter := snapshotCountBefore - result.SnapshotsDeleted
|
// Snapshots remaining after removing the incomplete ones; unknown if
|
||||||
|
// the pre-prune snapshot count could not be read.
|
||||||
|
snapshotsRemain := countDiff(snapshotCountBefore, &result.SnapshotsDeleted)
|
||||||
|
|
||||||
v.UI.Completef("Pruned local index database.")
|
v.UI.Completef("Pruned local index database.")
|
||||||
v.UI.Detailf("Incomplete snapshots: %d removed (%d remain).",
|
v.UI.Detailf("Incomplete snapshots: %s removed (%s remain).",
|
||||||
result.SnapshotsDeleted, snapshotCountAfter)
|
countText(&result.SnapshotsDeleted), countText(snapshotsRemain))
|
||||||
v.UI.Detailf("Orphaned files: %d removed (%d remain).",
|
v.UI.Detailf("Orphaned files: %s removed (%s remain).",
|
||||||
result.FilesDeleted, fileCountAfter)
|
countText(result.FilesDeleted), countText(fileCountAfter))
|
||||||
v.UI.Detailf("Orphaned chunks: %d removed (%d remain).",
|
v.UI.Detailf("Orphaned chunks: %s removed (%s remain).",
|
||||||
result.ChunksDeleted, chunkCountAfter)
|
countText(result.ChunksDeleted), countText(chunkCountAfter))
|
||||||
v.UI.Detailf("Orphaned blobs: %d removed (%d remain).",
|
v.UI.Detailf("Orphaned blobs: %s removed (%s remain).",
|
||||||
result.BlobsDeleted, blobCountAfter)
|
countText(result.BlobsDeleted), countText(blobCountAfter))
|
||||||
|
|
||||||
return result, nil
|
return result, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// countUnknown is what a count reads as when its query could not be run,
|
||||||
|
// distinct from "0", which means the table really was empty.
|
||||||
|
const countUnknown = "unknown"
|
||||||
|
|
||||||
|
// tableCountForReport returns the row count of a table for the prune
|
||||||
|
// summary, or nil if the count could not be read. A read failure is
|
||||||
|
// logged at warn — visible even under --json, which routes warnings to
|
||||||
|
// stderr — and then rendered as unknown rather than silently becoming 0,
|
||||||
|
// so a broken query is a visible failure instead of a plausible wrong
|
||||||
|
// number.
|
||||||
|
func (v *Vaultik) tableCountForReport(tableName string) *int64 {
|
||||||
|
count, err := v.getTableCount(tableName)
|
||||||
|
if err != nil {
|
||||||
|
log.Warn("could not read table row count for prune summary",
|
||||||
|
"table", tableName, "error", err)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
return &count
|
||||||
|
}
|
||||||
|
|
||||||
|
// countDiff returns before-after, or nil if either count is unknown so
|
||||||
|
// that an unreadable count does not collapse into a plausible delta.
|
||||||
|
func countDiff(before, after *int64) *int64 {
|
||||||
|
if before == nil || after == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
diff := *before - *after
|
||||||
|
|
||||||
|
return &diff
|
||||||
|
}
|
||||||
|
|
||||||
|
// countText renders a count that may be unknown: nil (the read failed)
|
||||||
|
// becomes "unknown", never "0", so a reader can tell an empty table from
|
||||||
|
// one that could not be queried.
|
||||||
|
func countText(count *int64) string {
|
||||||
|
if count == nil {
|
||||||
|
return countUnknown
|
||||||
|
}
|
||||||
|
|
||||||
|
return strconv.FormatInt(*count, 10)
|
||||||
|
}
|
||||||
|
|
||||||
// validTableNameRe matches table names containing only lowercase
|
// validTableNameRe matches table names containing only lowercase
|
||||||
// alphanumeric characters and underscores.
|
// alphanumeric characters and underscores.
|
||||||
var validTableNameRe = regexp.MustCompile(`^[a-z0-9_]+$`)
|
var validTableNameRe = regexp.MustCompile(`^[a-z0-9_]+$`)
|
||||||
|
|||||||
Reference in New Issue
Block a user