Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
240 changes: 204 additions & 36 deletions adapter/admin_backup.go

Large diffs are not rendered by default.

181 changes: 173 additions & 8 deletions adapter/admin_backup_test.go
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
package adapter

import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
stderrors "errors"
"sync"
"sync/atomic"
Expand All @@ -11,6 +13,7 @@ import (

logicalbackup "github.com/bootjp/elastickv/internal/backup"
"github.com/bootjp/elastickv/internal/raftengine"
"github.com/bootjp/elastickv/internal/s3keys"
"github.com/bootjp/elastickv/kv"
pb "github.com/bootjp/elastickv/proto"
kvstore "github.com/bootjp/elastickv/store"
Expand Down Expand Up @@ -123,6 +126,7 @@ type backupTestStore struct {
captureErr error
validateErr error
valueKeys [][]byte
values map[string][]byte
}

func (s *backupTestStore) ValidateBackupSnapshotAt(context.Context, kv.BackupRouteSnapshot, uint64, int) error {
Expand Down Expand Up @@ -185,10 +189,54 @@ func (s *backupTestStore) newBackupScannerAtSnapshot(ts uint64, keyFilter kv.Bac
}
}
s.valueKeys = append(s.valueKeys, append([]byte(nil), key...))
pairs = append(pairs, &kvstore.KVPair{Key: append([]byte(nil), key...), Value: []byte("value")})
value := s.values[string(key)]
if value == nil {
value = backupTestDefaultValueForKey(key)
}
pairs = append(pairs, &kvstore.KVPair{Key: append([]byte(nil), key...), Value: append([]byte(nil), value...)})
}
onExhaust := s.onExhaust
delay := s.scanDelay
closeErr := s.pairCloseErr
return &backupPairScanner{pairs: pairs, err: filterErr, closeErr: closeErr}
return &backupPairScanner{pairs: pairs, err: filterErr, onExhaust: onExhaust, delay: delay, closeErr: closeErr}
}

func backupTestDefaultValueForKey(key []byte) []byte {
scope, scoped, err := logicalbackup.ScopeForKey(key)
if err != nil || !scoped {
return []byte("value")
}
if scope.Adapter != "dynamodb" {
return []byte("value")
}
if bytes.HasPrefix(key, []byte(logicalbackup.DDBTableMetaPrefix)) {
value, err := encodeStoredDynamoTableSchema(&dynamoTableSchema{
TableName: scope.Name,
AttributeDefinitions: map[string]string{"id": "S"},
PrimaryKey: dynamoKeySchema{HashKey: "id"},
Generation: 1,
})
if err != nil {
panic(err)
}
return value
}
if bytes.HasPrefix(key, []byte(logicalbackup.DDBItemPrefix)) {
id := "value"
value, err := encodeStoredDynamoItem(map[string]attributeValue{"id": {S: &id}})
if err != nil {
panic(err)
}
return value
}
return []byte("value")
}

func backupTestJSON(t *testing.T, v any) []byte {
t.Helper()
out, err := json.Marshal(v)
require.NoError(t, err)
return out
}

type backupSliceScanner struct {
Expand Down Expand Up @@ -229,10 +277,13 @@ func (s *backupSliceScanner) Next(ctx context.Context) ([]byte, bool, error) {
func (s *backupSliceScanner) Close() error { return s.closeErr }

type backupPairScanner struct {
pairs []*kvstore.KVPair
index int
err error
closeErr error
pairs []*kvstore.KVPair
index int
err error
onExhaust func()
once sync.Once
delay time.Duration
closeErr error
}

func (s *backupPairScanner) Next(ctx context.Context) (*kvstore.KVPair, bool, error) {
Expand All @@ -244,7 +295,21 @@ func (s *backupPairScanner) Next(ctx context.Context) (*kvstore.KVPair, bool, er
s.err = nil
return nil, false, err
}
if s.delay > 0 && s.index < len(s.pairs) {
timer := time.NewTimer(s.delay)
select {
case <-ctx.Done():
timer.Stop()
return nil, false, ctx.Err()
case <-timer.C:
}
}
if s.index >= len(s.pairs) {
s.once.Do(func() {
if s.onExhaust != nil {
s.onExhaust()
}
})
return nil, false, nil
}
pair := s.pairs[s.index]
Expand Down Expand Up @@ -350,6 +415,103 @@ func TestBeginBackupLifecycleAndBaselineAtPinnedTimestamp(t *testing.T) {
}
}

func TestBeginBackupExpectedKeysUseRetainedCounts(t *testing.T) {
t.Parallel()
const table = "orders"
activeID := "active"
staleID := "stale"
schema, err := encodeStoredDynamoTableSchema(&dynamoTableSchema{
TableName: table,
AttributeDefinitions: map[string]string{"id": "S"},
PrimaryKey: dynamoKeySchema{HashKey: "id"},
Generation: 2,
})
require.NoError(t, err)
activeItem, err := encodeStoredDynamoItem(map[string]attributeValue{"id": {S: &activeID}})
require.NoError(t, err)
staleItem, err := encodeStoredDynamoItem(map[string]attributeValue{"id": {S: &staleID}})
require.NoError(t, err)
schemaKey := logicalbackup.EncodeDDBTableMetaKey(table)
activeKey := logicalbackup.EncodeDDBItemKey(table, 2, activeID, "")
staleKey := logicalbackup.EncodeDDBItemKey(table, 1, staleID, "")
store := &backupTestStore{
keys: [][]byte{staleKey, activeKey, schemaKey},
values: map[string][]byte{
string(schemaKey): schema,
string(activeKey): activeItem,
string(staleKey): staleItem,
},
}
group := &backupTestGroup{status: raftengine.Status{AppliedIndex: 100}, every: 10_000}
proposer := newBackupTestProposer()
srv := newBackupControlTestServer(t, store, map[uint64]*backupTestGroup{1: group}, map[uint64]*backupTestProposer{1: proposer}, nil)

begin, err := srv.BeginBackup(context.Background(), &pb.BeginBackupRequest{})
require.NoError(t, err)
require.Len(t, begin.GetExpectedKeys(), 1)
require.Equal(t, "dynamodb", begin.GetExpectedKeys()[0].GetAdapter())
require.Equal(t, table, begin.GetExpectedKeys()[0].GetScope())
require.EqualValues(t, 2, begin.GetExpectedKeys()[0].GetKeyCount())
require.Equal(t, [][]byte{schemaKey}, store.valueKeys)
}

func TestBeginBackupExpectedKeysAvoidMaterializingBlobValues(t *testing.T) {
t.Parallel()
const (
bucket = "photos"
object = "large.bin"
uploadID = "upload-1"
)
bucketKey := s3keys.BucketMetaKey(bucket)
manifestKey := s3keys.ObjectManifestKey(bucket, 1, object)
blobKey := s3keys.BlobKey(bucket, 1, object, uploadID, 1, 0)
store := &backupTestStore{
keys: [][]byte{bucketKey, manifestKey, blobKey},
values: map[string][]byte{
string(bucketKey): backupTestJSON(t, map[string]any{
"bucket_name": bucket,
"generation": 1,
}),
string(manifestKey): backupTestJSON(t, map[string]any{
"upload_id": uploadID,
"parts": []map[string]any{{
"part_no": 1,
"chunk_count": 1,
}},
}),
string(blobKey): []byte("large-blob-body"),
},
}
group := &backupTestGroup{status: raftengine.Status{AppliedIndex: 100}, every: 10_000}
proposer := newBackupTestProposer()
srv := newBackupControlTestServer(t, store, map[uint64]*backupTestGroup{1: group}, map[uint64]*backupTestProposer{1: proposer}, nil)

begin, err := srv.BeginBackup(context.Background(), &pb.BeginBackupRequest{Adapters: []string{"s3"}})
require.NoError(t, err)
require.Len(t, begin.GetExpectedKeys(), 1)
require.Equal(t, "s3", begin.GetExpectedKeys()[0].GetAdapter())
require.Equal(t, bucket, begin.GetExpectedKeys()[0].GetScope())
require.EqualValues(t, 3, begin.GetExpectedKeys()[0].GetKeyCount())
require.Equal(t, [][]byte{bucketKey, manifestKey}, store.valueKeys)
}

func TestBeginBackupBaselineSkipsUnselectedAdapterValues(t *testing.T) {
t.Parallel()
schemaKey := logicalbackup.EncodeDDBTableMetaKey("orders")
store := &backupTestStore{
keys: [][]byte{schemaKey},
values: map[string][]byte{string(schemaKey): []byte("future-format")},
}
group := &backupTestGroup{status: raftengine.Status{AppliedIndex: 100}, every: 10_000}
proposer := newBackupTestProposer()
srv := newBackupControlTestServer(t, store, map[uint64]*backupTestGroup{1: group}, map[uint64]*backupTestProposer{1: proposer}, nil)

begin, err := srv.BeginBackup(context.Background(), &pb.BeginBackupRequest{Adapters: []string{"redis"}})
require.NoError(t, err)
require.Empty(t, begin.GetExpectedKeys())
require.Empty(t, store.valueKeys)
}

func TestSnapshotBackupGroupsExcludesReservedTSOGroup(t *testing.T) {
t.Parallel()
groups := map[uint64]*backupTestGroup{
Expand Down Expand Up @@ -486,6 +648,9 @@ func TestStreamBackupUsesPinTimestampAndScopeFilter(t *testing.T) {
srv := newBackupControlTestServer(t, store, map[uint64]*backupTestGroup{1: group}, map[uint64]*backupTestProposer{1: proposer}, nil)
begin, err := srv.BeginBackup(context.Background(), &pb.BeginBackupRequest{})
require.NoError(t, err)
store.mu.Lock()
store.valueKeys = nil
store.mu.Unlock()

stream := &backupTestStream{ctx: context.Background()}
err = srv.StreamBackup(&pb.StreamBackupRequest{
Expand Down Expand Up @@ -644,7 +809,7 @@ func TestListBackupScopesReportsScannerCloseError(t *testing.T) {
require.NoError(t, err)

store.mu.Lock()
store.keyCloseErr = stderrors.New("close failed")
store.pairCloseErr = stderrors.New("close failed")
store.mu.Unlock()
_, err = srv.ListAdaptersAndScopes(context.Background(), &pb.ListAdaptersAndScopesRequest{PinToken: begin.GetPinToken()})
require.Equal(t, codes.FailedPrecondition, status.Code(err))
Expand Down Expand Up @@ -884,7 +1049,7 @@ func TestBackupTokenDeadlineRotatesAndFailsClosed(t *testing.T) {
_, err = srv.ListAdaptersAndScopes(context.Background(), &pb.ListAdaptersAndScopesRequest{PinToken: begin.GetPinToken()})
require.Equal(t, codes.FailedPrecondition, status.Code(err))
err = srv.StreamBackup(&pb.StreamBackupRequest{PinToken: begin.GetPinToken()}, &backupTestStream{ctx: context.Background()})
require.Equal(t, codes.FailedPrecondition, status.Code(err))
require.NoError(t, err, "streams use the renewed server-side session deadline")

_, err = srv.RenewBackup(context.Background(), &pb.RenewBackupRequest{PinToken: renewed.GetPinToken()})
require.NoError(t, err)
Expand Down
Loading
Loading