Files
Felis/internal/offsite/offsite_test.go
T

474 lines
15 KiB
Go

package offsite
import (
"bytes"
"context"
"crypto/rand"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"testing"
"time"
)
func testKey(t *testing.T) []byte {
t.Helper()
s, err := NewKey()
if err != nil {
t.Fatal(err)
}
k, err := ParseKey(s)
if err != nil {
t.Fatal(err)
}
return k
}
func seal(t *testing.T, plain, key []byte) []byte {
t.Helper()
var buf bytes.Buffer
if err := Encrypt(&buf, bytes.NewReader(plain), key); err != nil {
t.Fatal(err)
}
return buf.Bytes()
}
func TestEncryptRoundTrip(t *testing.T) {
key := testKey(t)
for _, n := range []int{0, 1, segmentSize - 1, segmentSize, segmentSize + 1, 3*segmentSize + 5} {
plain := make([]byte, n)
rand.Read(plain)
sealed := seal(t, plain, key)
if int64(len(sealed)) != SealedSize(int64(n)) {
t.Errorf("n=%d: sealed %d bytes, SealedSize says %d", n, len(sealed), SealedSize(int64(n)))
}
var out bytes.Buffer
if err := Decrypt(&out, bytes.NewReader(sealed), key); err != nil {
t.Fatalf("n=%d: decrypt: %v", n, err)
}
if !bytes.Equal(out.Bytes(), plain) {
t.Fatalf("n=%d: round trip changed the data", n)
}
}
}
// TestDecryptRejectsTampering: a wrong key, a flipped bit, a missing tail at a
// segment boundary and swapped segments all fail instead of yielding data.
func TestDecryptRejectsTampering(t *testing.T) {
key := testKey(t)
plain := make([]byte, 3*segmentSize)
rand.Read(plain)
sealed := seal(t, plain, key)
seg := segmentSize + 16
flipped := bytes.Clone(sealed)
flipped[headerSize+10] ^= 1
cut := sealed[:headerSize+2*seg]
swapped := bytes.Clone(sealed)
copy(swapped[headerSize:], sealed[headerSize+seg:headerSize+2*seg])
copy(swapped[headerSize+seg:], sealed[headerSize:headerSize+seg])
cases := map[string]struct {
data []byte
key []byte
}{
"wrong key": {sealed, testKey(t)},
"flipped bit": {flipped, key},
"cut short": {cut, key},
"swapped": {swapped, key},
"header only": {sealed[:headerSize], key},
"not ours": {[]byte("PK\x03\x04 a zip file, not a sealed object"), key},
}
for name, c := range cases {
if err := Decrypt(io.Discard, bytes.NewReader(c.data), c.key); err == nil {
t.Errorf("%s: decrypted without error", name)
}
}
if err := Decrypt(io.Discard, bytes.NewReader(sealed), testKey(t)); !errors.Is(err, ErrAuth) {
t.Errorf("wrong key error = %v, want ErrAuth", err)
}
}
func TestParseKey(t *testing.T) {
s, _ := NewKey()
k, err := ParseKey(" " + s + "\n")
if err != nil || len(k) != KeySize {
t.Fatalf("ParseKey(NewKey()) = %d bytes, %v", len(k), err)
}
for _, bad := range []string{"", "short", strings.Repeat("A", 40)} {
if _, err := ParseKey(bad); err == nil {
t.Errorf("ParseKey(%q) accepted", bad)
}
}
if KeyID(k) == KeyID(testKey(t)) || len(KeyID(k)) != 16 {
t.Errorf("KeyID = %q, want 16 hex chars that differ per key", KeyID(k))
}
}
// ---- fakes ----------------------------------------------------------------
type memBucket struct {
mu sync.Mutex
objs map[string][]byte
putErr map[string]error
puts int
removed []string
}
func newMemBucket() *memBucket {
return &memBucket{objs: map[string][]byte{}, putErr: map[string]error{}}
}
func (b *memBucket) Put(_ context.Context, key string, r io.Reader, size int64) error {
if err := b.putErr[key]; err != nil {
io.Copy(io.Discard, r)
return err
}
data, err := io.ReadAll(r)
if err != nil {
return err
}
if int64(len(data)) != size {
return fmt.Errorf("put %s: got %d bytes, declared %d", key, len(data), size)
}
b.mu.Lock()
defer b.mu.Unlock()
b.objs[key] = data
b.puts++
return nil
}
func (b *memBucket) Get(_ context.Context, key string) (io.ReadCloser, error) {
b.mu.Lock()
defer b.mu.Unlock()
data, ok := b.objs[key]
if !ok {
return nil, fmt.Errorf("%w: %s", ErrNotFound, key)
}
return io.NopCloser(bytes.NewReader(data)), nil
}
func (b *memBucket) List(_ context.Context, prefix string) ([]Object, error) {
b.mu.Lock()
defer b.mu.Unlock()
var out []Object
for k, v := range b.objs {
if strings.HasPrefix(k, prefix) {
out = append(out, Object{Key: k, Size: int64(len(v))})
}
}
sort.Slice(out, func(i, j int) bool { return out[i].Key < out[j].Key })
return out, nil
}
func (b *memBucket) Remove(_ context.Context, key string) error {
b.mu.Lock()
defer b.mu.Unlock()
delete(b.objs, key)
b.removed = append(b.removed, key)
return nil
}
type row struct {
WorldBackup
status string
expires time.Time
offsite time.Time
}
type fakeCatalog struct{ rows []*row }
func (c *fakeCatalog) PendingWorlds(context.Context) ([]WorldBackup, error) {
var out []WorldBackup
for _, r := range c.rows {
if r.status == "present" && r.offsite.IsZero() {
out = append(out, r.WorldBackup)
}
}
return out, nil
}
func (c *fakeCatalog) PresentWorlds(context.Context) ([]WorldBackup, error) {
var out []WorldBackup
for _, r := range c.rows {
if r.status == "present" {
out = append(out, r.WorldBackup)
}
}
return out, nil
}
func (c *fakeCatalog) MarkOffsite(_ context.Context, id string, at time.Time) error {
for _, r := range c.rows {
if r.ID == id {
r.offsite = at
}
}
return nil
}
func (c *fakeCatalog) ExpiredRefs(_ context.Context, now time.Time) ([]string, error) {
var out []string
for _, r := range c.rows {
if r.status == "deleted" && r.expires.Before(now) && !r.offsite.IsZero() {
out = append(out, r.Ref)
}
}
return out, nil
}
func writeFile(t *testing.T, dir, name string, size int) []byte {
t.Helper()
data := make([]byte, size)
rand.Read(data)
if err := os.WriteFile(filepath.Join(dir, name), data, 0o600); err != nil {
t.Fatal(err)
}
return data
}
var now = time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC)
func newSyncer(t *testing.T, cat *fakeCatalog) (*Syncer, *memBucket) {
t.Helper()
b := newMemBucket()
return &Syncer{
Bucket: b, Catalog: cat, Key: testKey(t),
ArchiveDir: t.TempDir(), DBDir: t.TempDir(), DBKeep: 2,
Now: func() time.Time { return now },
}, b
}
// TestSyncCopiesWorldsAndRecordsThem: each pending archive is encrypted into
// the bucket and marked; a second run sends nothing again.
func TestSyncCopiesWorldsAndRecordsThem(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Server: "alpha", Ref: "/var/lib/felis/archives/alpha-1.tar.gz"}, status: "present"},
{WorldBackup: WorldBackup{ID: "b2", Server: "beta", Ref: "/var/lib/felis/archives/beta-2.tar.gz"}, status: "present"},
}}
s, b := newSyncer(t, cat)
alpha := writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", 3*segmentSize+7)
writeFile(t, s.ArchiveDir, "beta-2.tar.gz", 10)
res, err := s.Run(context.Background())
if err != nil {
t.Fatalf("run: %v (%+v)", err, res)
}
if res.WorldsUploaded != 2 || res.WorldsPending != 0 || res.RemoteWorlds != 2 {
t.Fatalf("result = %+v, want 2 uploaded, 0 pending, 2 remote", res)
}
for _, r := range cat.rows {
if !r.offsite.Equal(now) {
t.Errorf("row %s offsite_at = %v, want %v", r.ID, r.offsite, now)
}
}
var out bytes.Buffer
if err := Decrypt(&out, bytes.NewReader(b.objs["worlds/alpha-1.tar.gz.fenc"]), s.Key); err != nil || !bytes.Equal(out.Bytes(), alpha) {
t.Fatalf("stored alpha does not decrypt to the archive (err %v)", err)
}
puts := b.puts
if res, err := s.Run(context.Background()); err != nil || res.WorldsUploaded != 0 || b.puts != puts {
t.Fatalf("second run uploaded again: %+v, puts %d→%d, err %v", res, puts, b.puts, err)
}
}
// TestSyncResumesWithoutResending: an object stored by a run that died before
// recording it is recorded, not uploaded twice; a partial one is replaced.
func TestSyncResumesWithoutResending(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Server: "alpha", Ref: "/a/alpha-1.tar.gz"}, status: "present"},
{WorldBackup: WorldBackup{ID: "b2", Server: "beta", Ref: "/a/beta-2.tar.gz"}, status: "present"},
}}
s, b := newSyncer(t, cat)
alpha := writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", 1000)
writeFile(t, s.ArchiveDir, "beta-2.tar.gz", 1000)
b.objs["worlds/alpha-1.tar.gz.fenc"] = seal(t, alpha, s.Key)
b.objs["worlds/beta-2.tar.gz.fenc"] = []byte("partial")
res, err := s.Run(context.Background())
if err != nil {
t.Fatal(err)
}
if res.WorldsUploaded != 1 || b.puts != 1 {
t.Fatalf("uploaded %d (puts %d), want only the partial beta resent", res.WorldsUploaded, b.puts)
}
if cat.rows[0].offsite.IsZero() || cat.rows[1].offsite.IsZero() {
t.Fatal("both rows should be recorded as copied")
}
}
// TestSyncFailureLeavesRowPending: a failed upload is reported, the row stays
// unmarked (so the reaper keeps the world), and the rest still go.
func TestSyncFailureLeavesRowPending(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Server: "alpha", Ref: "/a/alpha-1.tar.gz"}, status: "present"},
{WorldBackup: WorldBackup{ID: "b2", Server: "beta", Ref: "/a/beta-2.tar.gz"}, status: "present"},
{WorldBackup: WorldBackup{ID: "b3", Server: "gamma", Ref: "/a/gamma-3.tar.gz"}, status: "present"},
}}
s, b := newSyncer(t, cat)
writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", 100)
writeFile(t, s.ArchiveDir, "beta-2.tar.gz", 100)
b.putErr["worlds/alpha-1.tar.gz.fenc"] = errors.New("503 slow down")
res, err := s.Run(context.Background())
if err == nil {
t.Fatal("run reported success despite a failed upload")
}
if !cat.rows[0].offsite.IsZero() {
t.Error("the failed archive was recorded as copied")
}
if cat.rows[1].offsite.IsZero() {
t.Error("the archive after the failure was not copied")
}
if res.WorldsPending != 1 || len(res.WorldsMissing) != 1 || !strings.Contains(res.WorldsMissing[0], "gamma") {
t.Fatalf("result = %+v, want alpha pending, gamma missing", res)
}
}
// TestSyncExpiresOnlyPastRetention: a remote archive goes once its row has
// expired; one evicted early from the local disk stays until then, and an
// object with no row at all is left alone.
func TestSyncExpiresOnlyPastRetention(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "old", Ref: "/a/old.tar.gz"}, status: "deleted", expires: now.Add(-time.Hour), offsite: now.Add(-100 * 24 * time.Hour)},
{WorldBackup: WorldBackup{ID: "evicted", Ref: "/a/evicted.tar.gz"}, status: "deleted", expires: now.Add(30 * 24 * time.Hour), offsite: now.Add(-24 * time.Hour)},
}}
s, b := newSyncer(t, cat)
for _, k := range []string{"worlds/old.tar.gz.fenc", "worlds/evicted.tar.gz.fenc", "worlds/unknown.tar.gz.fenc"} {
b.objs[k] = []byte("x")
}
res, err := s.Run(context.Background())
if err != nil {
t.Fatal(err)
}
if res.WorldsExpired != 1 || len(b.removed) != 1 || b.removed[0] != "worlds/old.tar.gz.fenc" {
t.Fatalf("removed %v (expired %d), want only the expired archive", b.removed, res.WorldsExpired)
}
if res.RemoteWorlds != 2 {
t.Fatalf("remote worlds = %d, want 2 left", res.RemoteWorlds)
}
}
// TestSyncDBKeepsNewest: only the newest DBKeep bundles are sent, and older
// remote bundles are pruned down to DBKeep.
func TestSyncDBKeepsNewest(t *testing.T) {
s, b := newSyncer(t, &fakeCatalog{})
for _, name := range []string{
"felis-db-20260920T030000Z-daily.tar",
"felis-db-20260922T030000Z-daily.tar",
"felis-db-20260923T030000Z-manual.tar",
"felis-db-20260924T030000Z-daily.tar",
} {
writeFile(t, s.DBDir, name, 50)
}
b.objs["db/felis-db-20260901T030000Z-daily.tar.fenc"] = []byte("old")
b.objs["db/notes.txt"] = []byte("not ours")
res, err := s.Run(context.Background())
if err != nil {
t.Fatal(err)
}
if res.DBUploaded != 2 || res.RemoteDB != 2 || res.NewestDB != "felis-db-20260924T030000Z-daily.tar" {
t.Fatalf("result = %+v, want the newest 2 uploaded", res)
}
var keys []string
for k := range b.objs {
keys = append(keys, k)
}
sort.Strings(keys)
want := "db/felis-db-20260923T030000Z-manual.tar.fenc db/felis-db-20260924T030000Z-daily.tar.fenc db/notes.txt"
if strings.Join(keys, " ") != want {
t.Fatalf("bucket = %v, want %s", keys, want)
}
listed, err := ListDB(context.Background(), b)
if err != nil || len(listed) != 2 || listed[0].Key != "felis-db-20260924T030000Z-daily.tar" {
t.Fatalf("ListDB = %+v, %v", listed, err)
}
}
// TestSyncDBRanksWithTheBucket: an old local bundle (kept here by its label's
// own retention) that newer bundles in the bucket outrank is not sent, so a
// pass does not upload what it then prunes, and the next pass the same again.
func TestSyncDBRanksWithTheBucket(t *testing.T) {
s, b := newSyncer(t, &fakeCatalog{})
writeFile(t, s.DBDir, "felis-db-20260910T030000Z-pre-migrate.tar", 50)
writeFile(t, s.DBDir, "felis-db-20260924T030000Z-daily.tar", 50)
b.objs["db/felis-db-20260923T030000Z-daily.tar.fenc"] = []byte("gone here, kept there")
for pass := 1; pass <= 2; pass++ {
res, err := s.Run(context.Background())
if err != nil {
t.Fatal(err)
}
want := 0
if pass == 1 {
want = 1
}
if res.DBUploaded != want || res.DBPruned != 0 || res.RemoteDB != 2 {
t.Fatalf("pass %d: result = %+v, want %d uploaded, none pruned, 2 held", pass, res, want)
}
}
if _, ok := b.objs["db/felis-db-20260910T030000Z-pre-migrate.tar.fenc"]; ok {
t.Fatal("the outranked bundle was sent")
}
}
// TestFetchWorldsRestoresVolume: after a rebuild, every present archive the
// volume lacks comes back byte for byte; ones already there are left alone.
func TestFetchWorldsRestoresVolume(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Server: "alpha", Ref: "/a/alpha-1.tar.gz"}, status: "present"},
{WorldBackup: WorldBackup{ID: "b2", Server: "beta", Ref: "/a/beta-2.tar.gz"}, status: "present"},
{WorldBackup: WorldBackup{ID: "b3", Server: "gamma", Ref: "/a/gamma-3.tar.gz"}, status: "present"},
}}
s, b := newSyncer(t, cat)
alpha := writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", segmentSize+3)
writeFile(t, s.ArchiveDir, "beta-2.tar.gz", 20)
if res, err := s.Run(context.Background()); err != nil || len(res.WorldsMissing) != 1 {
t.Fatalf("run = %+v, %v; want gamma reported missing", res, err)
}
fresh := t.TempDir()
os.WriteFile(filepath.Join(fresh, "beta-2.tar.gz"), []byte("kept"), 0o600)
res, err := FetchWorlds(context.Background(), b, cat, s.Key, fresh, nil)
if err != nil {
t.Fatal(err)
}
if strings.Join(res.Fetched, ",") != "alpha-1.tar.gz" || len(res.Missing) != 1 {
t.Fatalf("fetch = %+v, want alpha fetched, gamma missing", res)
}
got, _ := os.ReadFile(filepath.Join(fresh, "alpha-1.tar.gz"))
if !bytes.Equal(got, alpha) {
t.Fatal("fetched alpha differs from the original")
}
if kept, _ := os.ReadFile(filepath.Join(fresh, "beta-2.tar.gz")); string(kept) != "kept" {
t.Fatal("fetch overwrote an archive already on the volume")
}
// A wrong key leaves no file behind, partial or whole.
other := t.TempDir()
if err := FetchObject(context.Background(), b, testKey(t), "worlds/alpha-1.tar.gz.fenc", filepath.Join(other, "alpha-1.tar.gz"), 0o600); !errors.Is(err, ErrAuth) {
t.Fatalf("wrong key fetch = %v, want ErrAuth", err)
}
if entries, _ := os.ReadDir(other); len(entries) != 0 {
t.Fatalf("wrong key left files: %v", entries)
}
}
func TestWorldKeyRejectsOddRefs(t *testing.T) {
if k, ok := WorldKey("/var/lib/felis/archives/alpha-1.tar.gz"); !ok || k != "worlds/alpha-1.tar.gz.fenc" {
t.Fatalf("WorldKey = %q, %v", k, ok)
}
for _, ref := range []string{"", "/", "..", "/a/.."} {
if _, ok := WorldKey(ref); ok {
t.Errorf("WorldKey(%q) accepted", ref)
}
}
}