Files
Felis/internal/offsite/offsite_test.go
T

581 lines
19 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
// stall holds keys whose upload never finishes: Put waits out its context.
stall map[string]bool
// started lists every Put in the order it began.
started []string
// modified is what List reports as each key's modification time.
modified map[string]time.Time
// getErr fails Get for a key the way an unreachable bucket would.
getErr map[string]error
}
func newMemBucket() *memBucket {
return &memBucket{objs: map[string][]byte{}, putErr: map[string]error{}}
}
func (b *memBucket) Put(ctx context.Context, key string, r io.Reader, size int64) error {
b.mu.Lock()
b.started = append(b.started, key)
b.mu.Unlock()
if b.stall[key] {
<-ctx.Done()
return ctx.Err()
}
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()
if err := b.getErr[key]; err != nil {
return nil, err
}
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)), Modified: b.modified[k]})
}
}
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.status == "expired") && r.expires.Before(now) && !r.offsite.IsZero() {
out = append(out, r.Ref)
}
}
return out, nil
}
func (c *fakeCatalog) NewestOffsite(context.Context) (time.Time, error) {
var at time.Time
for _, r := range c.rows {
if r.offsite.After(at) {
at = r.offsite
}
}
return at, nil
}
func (c *fakeCatalog) KeptRefs(_ context.Context, now time.Time) ([]string, error) {
var out []string
for _, r := range c.rows {
if r.status == "present" || !r.expires.Before(now) {
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)
// newSyncer syncs into a bucket that already records its key, as every bucket
// does after its first run.
func newSyncer(t *testing.T, cat *fakeCatalog) (*Syncer, *memBucket) {
t.Helper()
b := newMemBucket()
key := testKey(t)
b.objs[keyMark] = []byte(KeyID(key) + "\n")
return &Syncer{
Bucket: b, Catalog: cat, Key: key,
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, or once its owner deleted it (status expired, expires_at pulled to
// the delete); 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)},
{WorldBackup: WorldBackup{ID: "dropped", Ref: "/a/dropped.tar.gz"}, status: "expired", expires: now.Add(-time.Minute), 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/dropped.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)
}
sort.Strings(b.removed)
if res.WorldsExpired != 2 || len(b.removed) != 2 || b.removed[0] != "worlds/dropped.tar.gz.fenc" || b.removed[1] != "worlds/old.tar.gz.fenc" {
t.Fatalf("removed %v (expired %d), want the expired and the owner-deleted 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 " + keyMark
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)
}
}
}
// TestSyncBigArchiveDoesNotStallThePass: an archive the uplink cannot send in
// time fails at its own deadline and stays pending, while the database bundle
// (sent before any world) and the archives after it still reach the bucket.
// One deadline for the whole pass used to go to the big archive, every hour,
// with the bundles queued behind it.
func TestSyncBigArchiveDoesNotStallThePass(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)
s.UploadGrace = 50 * time.Millisecond
writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", 100)
writeFile(t, s.ArchiveDir, "beta-2.tar.gz", 100)
writeFile(t, s.DBDir, "felis-db-20260924T030000Z-daily.tar", 50)
b.stall = map[string]bool{"worlds/alpha-1.tar.gz.fenc": true}
type ran struct {
res Result
err error
}
done := make(chan ran, 1)
go func() {
res, err := s.Run(context.Background())
done <- ran{res, err}
}()
var r ran
select {
case r = <-done:
case <-time.After(10 * time.Second):
t.Fatal("the pass is still waiting on the archive that cannot be sent")
}
if r.err == nil || !strings.Contains(r.err.Error(), "alpha-1.tar.gz") || !strings.Contains(r.err.Error(), "the next run tries again") {
t.Fatalf("run error = %v, want the stalled archive named as retried next run", r.err)
}
if r.res.DBUploaded != 1 || r.res.WorldsUploaded != 1 || r.res.WorldsPending != 1 {
t.Fatalf("result = %+v, want the bundle and beta copied, alpha pending", r.res)
}
if !cat.rows[0].offsite.IsZero() || cat.rows[1].offsite.IsZero() {
t.Fatalf("offsite_at alpha=%v beta=%v, want only beta recorded", cat.rows[0].offsite, cat.rows[1].offsite)
}
if len(b.started) == 0 || !strings.HasPrefix(b.started[0], dbDir) {
t.Fatalf("uploads began in the order %v, want the database bundle first", b.started)
}
}
// TestUploadBudgetFitsTheUplink: by default a 10 GiB archive may take well over
// what a 20 Mbit/s uplink needs to send it, and the smallest object still gets
// the grace.
func TestUploadBudgetFitsTheUplink(t *testing.T) {
s := &Syncer{}
const big = 10 << 30
need := time.Duration(float64(big) / (20e6 / 8) * float64(time.Second))
if got := s.uploadBudget(big); got < 3*need {
t.Errorf("budget for 10 GiB = %s, want at least %s (3× a 20 Mbit/s uplink)", got, 3*need)
}
if got := s.uploadBudget(0); got < defaultUploadGrace {
t.Errorf("budget for an empty object = %s, want at least %s", got, defaultUploadGrace)
}
}