feat(offsite): 世界归档与数据库备份加密同步到异地 S3,reaper 确认异地副本后才删除世界

This commit is contained in:
Lemon-miaow committed 2026-09-24 19:25:16 +08:00
1 parent fa8db4c7e8
commit 47890ca913
27 files changed
+2928 -51

No files matched your search

+160
View File
@@ -0,0 +1,160 @@
package offsite
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"strings"
"time"
"github.com/minio/minio-go/v7"
"github.com/minio/minio-go/v7/pkg/credentials"
)
// Object is one object in the bucket, its key relative to the configured
// prefix.
type Object struct {
Key string
Size int64
Modified time.Time
}
// Bucket is the object store the sync writes to. S3 is the production one; the
// tests use an in-memory map.
type Bucket interface {
// Put stores exactly size bytes read from r under key.
Put(ctx context.Context, key string, r io.Reader, size int64) error
// Get opens key. A missing key is ErrNotFound.
Get(ctx context.Context, key string) (io.ReadCloser, error)
// List returns every object whose key starts with prefix.
List(ctx context.Context, prefix string) ([]Object, error)
Remove(ctx context.Context, key string) error
}
// ErrNotFound is a key the bucket does not hold.
var ErrNotFound = errors.New("offsite: no such object")
// S3Config locates an S3-compatible bucket. Endpoint takes an http:// or
// https:// scheme; a bare host means TLS.
type S3Config struct {
Endpoint string
Region string
Bucket string
Prefix string
AccessKey string
SecretKey string
}
// S3 is a Bucket on an S3-compatible store (AWS S3, Backblaze B2, Cloudflare
// R2, Wasabi, MinIO...). Every key is placed under Prefix.
type S3 struct {
client *minio.Client
bucket string
prefix string
}
// partSize bounds what one multipart upload holds in memory.
const partSize = 16 << 20
// NewS3 builds the client. It does not touch the network; Check does.
func NewS3(c S3Config) (*S3, error) {
host, secure, err := splitEndpoint(c.Endpoint)
if err != nil {
return nil, err
}
if c.Bucket == "" {
return nil, errors.New("offsite: no bucket configured")
}
if c.AccessKey == "" || c.SecretKey == "" {
return nil, errors.New("offsite: the access key or the secret key is empty")
}
cl, err := minio.New(host, &minio.Options{
Creds: credentials.NewStaticV4(c.AccessKey, c.SecretKey, ""),
Secure: secure,
Region: c.Region,
})
if err != nil {
return nil, fmt.Errorf("offsite: %w", err)
}
return &S3{client: cl, bucket: c.Bucket, prefix: cleanPrefix(c.Prefix)}, nil
}
func splitEndpoint(ep string) (host string, secure bool, err error) {
ep = strings.TrimSpace(ep)
switch {
case ep == "":
return "", false, errors.New("offsite: no endpoint configured")
case strings.HasPrefix(ep, "https://"):
return strings.Trim(strings.TrimPrefix(ep, "https://"), "/"), true, nil
case strings.HasPrefix(ep, "http://"):
return strings.Trim(strings.TrimPrefix(ep, "http://"), "/"), false, nil
default:
return strings.Trim(ep, "/"), true, nil
}
}
func cleanPrefix(p string) string {
p = strings.Trim(p, "/")
if p == "" {
return ""
}
return p + "/"
}
// Check proves the bucket is reachable and the credentials may use it.
func (s *S3) Check(ctx context.Context) error {
ok, err := s.client.BucketExists(ctx, s.bucket)
if err != nil {
return fmt.Errorf("offsite: cannot reach bucket %q (endpoint unreachable or credentials rejected): %w", s.bucket, err)
}
if !ok {
return fmt.Errorf("offsite: bucket %q does not exist; create it first", s.bucket)
}
return nil
}
func (s *S3) Put(ctx context.Context, key string, r io.Reader, size int64) error {
_, err := s.client.PutObject(ctx, s.bucket, s.prefix+key, r, size, minio.PutObjectOptions{
ContentType: "application/octet-stream",
PartSize: partSize,
})
return err
}
func (s *S3) Get(ctx context.Context, key string) (io.ReadCloser, error) {
obj, err := s.client.GetObject(ctx, s.bucket, s.prefix+key, minio.GetObjectOptions{})
if err != nil {
return nil, err
}
// GetObject is lazy; Stat surfaces a missing key before the first read.
if _, err := obj.Stat(); err != nil {
obj.Close()
if isNotFound(err) {
return nil, fmt.Errorf("%w: %s", ErrNotFound, key)
}
return nil, err
}
return obj, nil
}
func (s *S3) List(ctx context.Context, prefix string) ([]Object, error) {
var out []Object
for o := range s.client.ListObjects(ctx, s.bucket, minio.ListObjectsOptions{Prefix: s.prefix + prefix, Recursive: true}) {
if o.Err != nil {
return nil, o.Err
}
out = append(out, Object{Key: strings.TrimPrefix(o.Key, s.prefix), Size: o.Size, Modified: o.LastModified})
}
return out, nil
}
func (s *S3) Remove(ctx context.Context, key string) error {
return s.client.RemoveObject(ctx, s.bucket, s.prefix+key, minio.RemoveObjectOptions{})
}
func isNotFound(err error) bool {
resp := minio.ToErrorResponse(err)
return resp.Code == "NoSuchKey" || resp.StatusCode == http.StatusNotFound
}
+195
View File
@@ -0,0 +1,195 @@
package offsite
import (
"bufio"
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/binary"
"encoding/hex"
"errors"
"fmt"
"io"
"strings"
)
// An object in the bucket is its file encrypted with AES-256-GCM in fixed
// segments, so a multi-gigabyte world streams through in constant memory:
//
// magic "FELISOS1" | 7-byte random nonce prefix | segment 0 | segment 1 | ...
//
// Segment i seals up to segmentSize bytes under the nonce
// prefix || uint32(i) || last, where last is 1 on the final segment only. The
// counter stops segments being reordered or dropped, and the last flag stops
// the object being cut short at a segment boundary: either change fails to
// authenticate.
const (
magic = "FELISOS1"
prefixSize = 7
segmentSize = 64 << 10
headerSize = len(magic) + prefixSize
// KeySize is the key length: AES-256.
KeySize = 32
)
// ErrAuth is a segment that does not authenticate: the wrong key, or an object
// damaged in the bucket or on the way.
var ErrAuth = errors.New("offsite: object does not decrypt with this key (wrong key, or the object is damaged)")
// NewKey returns a fresh random key in the text form ParseKey reads.
func NewKey() (string, error) {
k := make([]byte, KeySize)
if _, err := rand.Read(k); err != nil {
return "", err
}
return base64.StdEncoding.EncodeToString(k), nil
}
// ParseKey decodes a key written by NewKey (standard base64 of 32 bytes).
func ParseKey(s string) ([]byte, error) {
k, err := base64.StdEncoding.DecodeString(strings.TrimSpace(s))
if err != nil || len(k) != KeySize {
return nil, fmt.Errorf("offsite: the key must be %d random bytes in base64 (generate one with `felis offsite keygen`)", KeySize)
}
return k, nil
}
// KeyID names a key without revealing it, so status output and logs can say
// which key a bucket's objects were written with.
func KeyID(key []byte) string {
h := sha256.Sum256(append([]byte("felis-offsite-key-id\x00"), key...))
return hex.EncodeToString(h[:8])
}
// SealedSize is the size of the object Encrypt makes from n plaintext bytes.
// Uploads need it up front: an S3 upload of unknown length buffers far more.
func SealedSize(n int64) int64 {
segs := (n + segmentSize - 1) / segmentSize
if segs == 0 {
segs = 1
}
return int64(headerSize) + n + segs*16
}
func newAEAD(key []byte) (cipher.AEAD, error) {
if len(key) != KeySize {
return nil, fmt.Errorf("offsite: key is %d bytes, want %d", len(key), KeySize)
}
block, err := aes.NewCipher(key)
if err != nil {
return nil, err
}
return cipher.NewGCM(block)
}
func segmentNonce(prefix []byte, i uint32, last bool) []byte {
n := make([]byte, 12)
copy(n, prefix)
binary.BigEndian.PutUint32(n[prefixSize:], i)
if last {
n[11] = 1
}
return n
}
// Encrypt writes src to dst in the segmented format.
func Encrypt(dst io.Writer, src io.Reader, key []byte) error {
aead, err := newAEAD(key)
if err != nil {
return err
}
prefix := make([]byte, prefixSize)
if _, err := rand.Read(prefix); err != nil {
return err
}
if _, err := io.WriteString(dst, magic); err != nil {
return err
}
if _, err := dst.Write(prefix); err != nil {
return err
}
in := bufio.NewReaderSize(src, segmentSize+1)
buf := make([]byte, segmentSize)
out := make([]byte, 0, segmentSize+aead.Overhead())
for i := uint32(0); ; i++ {
n, err := io.ReadFull(in, buf)
switch {
case err == io.EOF || err == io.ErrUnexpectedEOF:
err = nil
case err != nil:
return err
}
last := n < segmentSize
if !last {
if _, perr := in.Peek(1); perr == io.EOF {
last = true
} else if perr != nil {
return perr
}
}
out = aead.Seal(out[:0], segmentNonce(prefix, i, last), buf[:n], []byte(magic))
if _, err := dst.Write(out); err != nil {
return err
}
if last {
return nil
}
if i == ^uint32(0) {
return errors.New("offsite: file too large to encrypt")
}
}
}
// Decrypt reverses Encrypt. Nothing unauthenticated is written: each segment
// is checked before its plaintext reaches dst, and a missing tail is an error.
// A caller writing to a file still has to discard it on error, since earlier
// segments were already written.
func Decrypt(dst io.Writer, src io.Reader, key []byte) error {
aead, err := newAEAD(key)
if err != nil {
return err
}
hdr := make([]byte, headerSize)
if _, err := io.ReadFull(src, hdr); err != nil {
return fmt.Errorf("offsite: object too short for its header: %w", err)
}
if string(hdr[:len(magic)]) != magic {
return errors.New("offsite: not a Felis off-site object (bad magic)")
}
prefix := hdr[len(magic):]
sealed := segmentSize + aead.Overhead()
in := bufio.NewReaderSize(src, sealed+1)
buf := make([]byte, sealed)
var plain []byte
for i := uint32(0); ; i++ {
n, err := io.ReadFull(in, buf)
switch {
case err == io.EOF:
return fmt.Errorf("offsite: object is cut short after %d segments", i)
case err == io.ErrUnexpectedEOF:
err = nil
case err != nil:
return err
}
last := n < sealed
if !last {
if _, perr := in.Peek(1); perr == io.EOF {
last = true
} else if perr != nil {
return perr
}
}
plain, err = aead.Open(plain[:0], segmentNonce(prefix, i, last), buf[:n], []byte(magic))
if err != nil {
return ErrAuth
}
if _, err := dst.Write(plain); err != nil {
return err
}
if last {
return nil
}
}
}
+446
View File
@@ -0,0 +1,446 @@
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)
}
}
// 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)
}
}
}
+73
View File
@@ -0,0 +1,73 @@
package offsite
import (
"context"
"database/sql"
"time"
)
// PGCatalog is the Catalog on the felis database.
type PGCatalog struct{ DB *sql.DB }
func (c PGCatalog) PendingWorlds(ctx context.Context) ([]WorldBackup, error) {
return c.query(ctx, `SELECT id, server_name, backup_ref, created_at FROM world_backups
WHERE status = 'present' AND offsite_at IS NULL ORDER BY created_at`)
}
func (c PGCatalog) PresentWorlds(ctx context.Context) ([]WorldBackup, error) {
return c.query(ctx, `SELECT id, server_name, backup_ref, created_at FROM world_backups
WHERE status = 'present' ORDER BY created_at`)
}
func (c PGCatalog) query(ctx context.Context, q string, args ...any) ([]WorldBackup, error) {
rows, err := c.DB.QueryContext(ctx, q, args...)
if err != nil {
return nil, err
}
defer rows.Close()
var out []WorldBackup
for rows.Next() {
var w WorldBackup
if err := rows.Scan(&w.ID, &w.Server, &w.Ref, &w.Created); err != nil {
return nil, err
}
out = append(out, w)
}
return out, rows.Err()
}
func (c PGCatalog) MarkOffsite(ctx context.Context, id string, at time.Time) error {
_, err := c.DB.ExecContext(ctx,
`UPDATE world_backups SET offsite_at = $2 WHERE id = $1 AND offsite_at IS NULL`, id, at)
return err
}
func (c PGCatalog) ExpiredRefs(ctx context.Context, now time.Time) ([]string, error) {
rows, err := c.DB.QueryContext(ctx, `SELECT backup_ref FROM world_backups
WHERE status = 'deleted' AND expires_at < $1 AND offsite_at IS NOT NULL`, now)
if err != nil {
return nil, err
}
defer rows.Close()
var out []string
for rows.Next() {
var ref string
if err := rows.Scan(&ref); err != nil {
return nil, err
}
out = append(out, ref)
}
return out, rows.Err()
}
// PendingCount is how many present archives wait for their copy, and since
// when the oldest has waited.
func (c PGCatalog) PendingCount(ctx context.Context) (int, time.Time, error) {
var (
n int
oldest sql.NullTime
)
err := c.DB.QueryRowContext(ctx, `SELECT count(*), min(created_at) FROM world_backups
WHERE status = 'present' AND offsite_at IS NULL`).Scan(&n, &oldest)
return n, oldest.Time, err
}
+64
View File
@@ -0,0 +1,64 @@
package offsite
import (
"encoding/json"
"os"
"path/filepath"
"time"
)
// DefaultStatusFile is where `felis offsite sync` records its last run; the
// watchdog and `felis offsite status` read it.
const DefaultStatusFile = "/var/lib/felis/offsite/status.json"
// StaleAfter is how long the sync may go without a clean run before the
// watchdog mails the owners. The timer runs hourly, so this rides out a
// provider's bad morning, and still leaves a day's database bundle uncopied
// for at most half a day.
const StaleAfter = 12 * time.Hour
// Status is the record of the last run.
type Status struct {
LastAttempt time.Time `json:"last_attempt"`
// LastSuccess is the last run in which every step succeeded.
LastSuccess time.Time `json:"last_success,omitempty"`
LastError string `json:"last_error,omitempty"`
Endpoint string `json:"endpoint"`
Bucket string `json:"bucket"`
Prefix string `json:"prefix,omitempty"`
KeyID string `json:"key_id"`
Result Result `json:"result"`
}
// ReadStatus reads the status file. A missing file is (nil, nil): no sync has
// run yet.
func ReadStatus(p string) (*Status, error) {
raw, err := os.ReadFile(p)
if os.IsNotExist(err) {
return nil, nil
}
if err != nil {
return nil, err
}
var st Status
if err := json.Unmarshal(raw, &st); err != nil {
return nil, err
}
return &st, nil
}
// WriteStatus replaces the status file atomically.
func WriteStatus(p string, st Status) error {
if err := os.MkdirAll(filepath.Dir(p), 0o700); err != nil {
return err
}
raw, err := json.MarshalIndent(st, "", " ")
if err != nil {
return err
}
tmp := p + ".tmp"
if err := os.WriteFile(tmp, append(raw, '\n'), 0o600); err != nil {
return err
}
return os.Rename(tmp, p)
}
+464
View File
@@ -0,0 +1,464 @@
// Package offsite keeps a second copy of what a lost node would take with it:
// every world archive (world_backups) and the newest control-plane database
// bundles (internal/dbbackup), encrypted, in an S3-compatible bucket off the
// machine. `felis offsite sync` runs it from felis-offsite.timer on the host,
// which is where both the archive volume and the bundle directory live.
//
// The database records the copy: world_backups.offsite_at is set once an
// archive's object is in the bucket, and with [offsite] configured the reaper
// deletes an idle world only after that (internal/reaper). Remote world
// objects go when their row has expired, so the bucket keeps each archive for
// the same retention the panel promises, including archives evicted early from
// the local disk to make room. Remote bundles are pruned to the newest DBKeep.
package offsite
import (
"context"
"errors"
"fmt"
"io"
"os"
"path"
"path/filepath"
"sort"
"strings"
"syscall"
"time"
"felis.lolicon.best/internal/dbbackup"
)
// Object key layout under the configured prefix.
const (
worldsDir = "worlds/"
dbDir = "db/"
objExt = ".fenc"
)
// WorldKey is the object key of the world archive stored at ref (the
// world_backups.backup_ref, an in-pod path whose last element is the file name
// on the archive volume).
func WorldKey(ref string) (string, bool) {
name := path.Base(ref)
if !safeName(name) {
return "", false
}
return worldsDir + name + objExt, true
}
// DBKey is the object key of a database bundle.
func DBKey(bundle string) string { return dbDir + bundle + objExt }
func safeName(name string) bool {
return name != "" && name != "." && name != ".." && name != "/" && !strings.ContainsAny(name, `/\`)
}
// WorldBackup is one world_backups row the sync works on.
type WorldBackup struct {
ID string
Server string
Ref string
Created time.Time
}
// Catalog is the world_backups view the sync needs; PGCatalog in production.
type Catalog interface {
// PendingWorlds lists present archives without an off-site copy yet,
// oldest first.
PendingWorlds(ctx context.Context) ([]WorldBackup, error)
// MarkOffsite records that the archive of row id is in the bucket.
MarkOffsite(ctx context.Context, id string, at time.Time) error
// ExpiredRefs lists the backup_ref of every row past its retention
// (deleted and expires_at < now) whose archive was copied off-site.
ExpiredRefs(ctx context.Context, now time.Time) ([]string, error)
// PresentWorlds lists every present archive, for a restore of the volume.
PresentWorlds(ctx context.Context) ([]WorldBackup, error)
}
// Syncer copies what is missing from the bucket and prunes what has expired.
type Syncer struct {
Bucket Bucket
Catalog Catalog
Key []byte
// ArchiveDir is the host directory of the world archive volume; "" when
// there is none yet (nothing has been archived on this install).
ArchiveDir string
// DBDir holds the database bundles; DBKeep is how many of the newest the
// bucket keeps.
DBDir string
DBKeep int
Now func() time.Time
Log io.Writer
}
// Result is what one Run did and found.
type Result struct {
WorldsUploaded int `json:"worlds_uploaded"`
BytesUploaded int64 `json:"bytes_uploaded"`
// WorldsPending are present archives still without an off-site copy
// after this run; the missing ones below are counted there alone.
WorldsPending int `json:"worlds_pending"`
// WorldsMissing are present rows whose archive file is not on the volume:
// nothing to copy, and nothing a restore could use.
WorldsMissing []string `json:"worlds_missing,omitempty"`
WorldsExpired int `json:"worlds_expired"`
RemoteWorlds int `json:"remote_worlds"`
RemoteBytes int64 `json:"remote_bytes"`
DBUploaded int `json:"db_uploaded"`
DBPruned int `json:"db_pruned"`
RemoteDB int `json:"remote_db"`
NewestDB string `json:"newest_db,omitempty"`
Errors []string `json:"errors,omitempty"`
}
func (s *Syncer) now() time.Time {
if s.Now != nil {
return s.Now()
}
return time.Now()
}
func (s *Syncer) logf(format string, args ...any) {
if s.Log != nil {
fmt.Fprintf(s.Log, "felis offsite: "+format+"\n", args...)
}
}
// Run does one pass: world archives, then database bundles, then expiry. A
// failure on one item is recorded and the pass carries on; the returned error
// is non-nil when anything failed.
func (s *Syncer) Run(ctx context.Context) (Result, error) {
var res Result
fail := func(format string, args ...any) {
msg := fmt.Sprintf(format, args...)
res.Errors = append(res.Errors, msg)
s.logf("%s", msg)
}
remoteWorlds, err := s.listSizes(ctx, worldsDir)
if err != nil {
return res, fmt.Errorf("list %s in the bucket: %w", worldsDir, err)
}
s.syncWorlds(ctx, remoteWorlds, &res, fail)
s.syncDB(ctx, &res, fail)
s.expireWorlds(ctx, remoteWorlds, &res, fail)
for _, size := range remoteWorlds {
res.RemoteWorlds++
res.RemoteBytes += size
}
if len(res.Errors) > 0 {
return res, fmt.Errorf("%d of this run's steps failed; first: %s", len(res.Errors), res.Errors[0])
}
return res, nil
}
func (s *Syncer) listSizes(ctx context.Context, prefix string) (map[string]int64, error) {
objs, err := s.Bucket.List(ctx, prefix)
if err != nil {
return nil, err
}
out := make(map[string]int64, len(objs))
for _, o := range objs {
out[o.Key] = o.Size
}
return out, nil
}
func (s *Syncer) syncWorlds(ctx context.Context, remote map[string]int64, res *Result, fail func(string, ...any)) {
pending, err := s.Catalog.PendingWorlds(ctx)
if err != nil {
fail("list world archives waiting for a copy: %v", err)
return
}
if s.ArchiveDir == "" {
res.WorldsPending = len(pending)
if len(pending) > 0 {
fail("%d world archives wait for a copy, but there is no archive volume to read them from", len(pending))
}
return
}
for _, w := range pending {
if ctx.Err() != nil {
fail("stopped: %v", ctx.Err())
res.WorldsPending++
continue
}
key, ok := WorldKey(w.Ref)
if !ok {
fail("world archive %s of %s has an unusable path %q", w.ID, w.Server, w.Ref)
res.WorldsPending++
continue
}
local := filepath.Join(s.ArchiveDir, path.Base(w.Ref))
fi, err := os.Stat(local)
if errors.Is(err, os.ErrNotExist) {
res.WorldsMissing = append(res.WorldsMissing, fmt.Sprintf("%s (%s, backup %s)", path.Base(w.Ref), w.Server, w.ID))
continue
}
if err != nil {
fail("world archive %s: %v", local, err)
res.WorldsPending++
continue
}
want := SealedSize(fi.Size())
// An earlier run may have stored the object and died before recording
// it; a complete object is only recorded, not sent again.
if size, ok := remote[key]; !ok || size != want {
if err := s.putFile(ctx, key, local, fi.Size()); err != nil {
fail("upload %s (%s): %v", path.Base(w.Ref), w.Server, err)
res.WorldsPending++
continue
}
remote[key] = want
res.WorldsUploaded++
res.BytesUploaded += fi.Size()
s.logf("copied world archive %s (%s, %s)", path.Base(w.Ref), w.Server, HumanBytes(fi.Size()))
}
if err := s.Catalog.MarkOffsite(ctx, w.ID, s.now()); err != nil {
fail("record the copy of %s: %v", path.Base(w.Ref), err)
res.WorldsPending++
}
}
}
func (s *Syncer) syncDB(ctx context.Context, res *Result, fail func(string, ...any)) {
if s.DBDir == "" {
return
}
keep := s.DBKeep
if keep < 1 {
keep = 1
}
local, err := dbbackup.List(s.DBDir)
if err != nil {
fail("list database bundles in %s: %v", s.DBDir, err)
return
}
remote, err := s.listSizes(ctx, dbDir)
if err != nil {
fail("list %s in the bucket: %v", dbDir, err)
return
}
// Only the newest keep bundles are worth sending: older ones would be
// pruned again at the end of this very pass.
if len(local) > keep {
local = local[:keep]
}
for _, b := range local {
key := DBKey(b.Name)
want := SealedSize(b.Size)
if size, ok := remote[key]; ok && size == want {
continue
}
if err := s.putFile(ctx, key, b.Path, b.Size); err != nil {
fail("upload database bundle %s: %v", b.Name, err)
continue
}
remote[key] = want
res.DBUploaded++
res.BytesUploaded += b.Size
s.logf("copied database bundle %s (%s)", b.Name, HumanBytes(b.Size))
}
var names []string
for key := range remote {
name := strings.TrimSuffix(strings.TrimPrefix(key, dbDir), objExt)
if _, _, ok := dbbackup.ParseBundleName(name); ok && strings.HasSuffix(key, objExt) {
names = append(names, name)
}
}
// Bundle names start with their UTC stamp, so newest sorts last.
sort.Sort(sort.Reverse(sort.StringSlice(names)))
for i, name := range names {
if i < keep {
continue
}
if err := s.Bucket.Remove(ctx, DBKey(name)); err != nil {
fail("prune database bundle %s: %v", name, err)
continue
}
res.DBPruned++
}
res.RemoteDB = min(len(names), keep)
if len(names) > 0 {
res.NewestDB = names[0]
}
}
func (s *Syncer) expireWorlds(ctx context.Context, remote map[string]int64, res *Result, fail func(string, ...any)) {
refs, err := s.Catalog.ExpiredRefs(ctx, s.now())
if err != nil {
fail("list expired world archives: %v", err)
return
}
for _, ref := range refs {
key, ok := WorldKey(ref)
if !ok {
continue
}
if _, ok := remote[key]; !ok {
continue
}
if err := s.Bucket.Remove(ctx, key); err != nil {
fail("remove expired %s: %v", key, err)
continue
}
delete(remote, key)
res.WorldsExpired++
}
}
// putFile encrypts the file at p into key. The sealed size is known in
// advance, so the upload streams: nothing larger than one part is buffered.
func (s *Syncer) putFile(ctx context.Context, key, p string, size int64) error {
f, err := os.Open(p)
if err != nil {
return err
}
defer f.Close()
pr, pw := io.Pipe()
go func() {
// A file that changed size under us would not match the declared
// length; LimitReader keeps the stream to the size we announced and
// the length check below catches a short one.
err := Encrypt(pw, io.LimitReader(f, size), s.Key)
pw.CloseWithError(err)
}()
err = s.Bucket.Put(ctx, key, pr, SealedSize(size))
pr.CloseWithError(errors.New("upload finished"))
return err
}
// FetchResult is what Fetch did.
type FetchResult struct {
Fetched []string
Present int
Missing []string // present rows with no object in the bucket
Failures []string
}
// FetchWorlds downloads every present archive the volume lacks: the volume
// half of a rebuild, after the database came back from a bundle.
func FetchWorlds(ctx context.Context, b Bucket, cat Catalog, key []byte, archiveDir string, log io.Writer) (FetchResult, error) {
var res FetchResult
worlds, err := cat.PresentWorlds(ctx)
if err != nil {
return res, err
}
res.Present = len(worlds)
for _, w := range worlds {
objKey, ok := WorldKey(w.Ref)
if !ok {
res.Failures = append(res.Failures, fmt.Sprintf("%s: unusable path %q", w.ID, w.Ref))
continue
}
dst := filepath.Join(archiveDir, path.Base(w.Ref))
if _, err := os.Stat(dst); err == nil {
continue
}
// World-readable like the archives the backup Jobs write: the restore
// Job reads them as its own non-root user.
switch err := FetchObject(ctx, b, key, objKey, dst, 0o644); {
case errors.Is(err, ErrNotFound):
res.Missing = append(res.Missing, fmt.Sprintf("%s (%s)", path.Base(w.Ref), w.Server))
case err != nil:
res.Failures = append(res.Failures, fmt.Sprintf("%s: %v", path.Base(w.Ref), err))
default:
alignOwner(dst, archiveDir)
res.Fetched = append(res.Fetched, path.Base(w.Ref))
if log != nil {
fmt.Fprintf(log, "felis offsite: restored %s (%s)\n", path.Base(w.Ref), w.Server)
}
}
}
if len(res.Failures) > 0 {
return res, fmt.Errorf("%d archives failed; first: %s", len(res.Failures), res.Failures[0])
}
return res, nil
}
// alignOwner gives a restored archive the archive volume's owner, which is the
// user the backup Jobs write as, so the volume looks as they left it. Only root
// can; for anyone else the file stays theirs, readable all the same.
func alignOwner(file, dir string) {
if os.Geteuid() != 0 {
return
}
fi, err := os.Stat(dir)
if err != nil {
return
}
if st, ok := fi.Sys().(*syscall.Stat_t); ok {
_ = os.Lchown(file, int(st.Uid), int(st.Gid))
}
}
// FetchObject downloads and decrypts key into dst, created with mode. The file
// appears under its name only once it has decrypted completely; until then it
// is a hidden .partial next to it.
func FetchObject(ctx context.Context, b Bucket, key []byte, objKey, dst string, mode os.FileMode) error {
rc, err := b.Get(ctx, objKey)
if err != nil {
return err
}
defer rc.Close()
tmp := filepath.Join(filepath.Dir(dst), "."+filepath.Base(dst)+".partial")
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600)
if err != nil {
return err
}
if err := f.Chmod(mode); err != nil {
f.Close()
os.Remove(tmp)
return err
}
if err := Decrypt(f, rc, key); err != nil {
f.Close()
os.Remove(tmp)
return err
}
if err := f.Sync(); err != nil {
f.Close()
os.Remove(tmp)
return err
}
if err := f.Close(); err != nil {
os.Remove(tmp)
return err
}
return os.Rename(tmp, dst)
}
// ListDB returns the database bundles in the bucket, newest first.
func ListDB(ctx context.Context, b Bucket) ([]Object, error) {
objs, err := b.List(ctx, dbDir)
if err != nil {
return nil, err
}
var out []Object
for _, o := range objs {
name := strings.TrimSuffix(strings.TrimPrefix(o.Key, dbDir), objExt)
if _, _, ok := dbbackup.ParseBundleName(name); !ok || !strings.HasSuffix(o.Key, objExt) {
continue
}
o.Key = name
out = append(out, o)
}
sort.Slice(out, func(i, j int) bool { return out[i].Key > out[j].Key })
return out, nil
}
// HumanBytes formats n in binary units.
func HumanBytes(n int64) string {
const unit = 1024
if n < unit {
return fmt.Sprintf("%d B", n)
}
div, exp := int64(unit), 0
for m := n / unit; m >= unit; m /= unit {
div *= unit
exp++
}
return fmt.Sprintf("%.1f %ciB", float64(n)/float64(div), "KMGTPE"[exp])
}