Skip to content

Commit 9663554

Browse files
committed
feat: implement Backup functionality for datastore
- Added Backup method to Store interface and DefaultStore implementation to stream a Badger backup of the datastore. - Introduced BackupRequest and BackupResponse messages in the state_rpc.proto file to handle backup requests and responses. - Implemented backup streaming logic in StoreServer, including metadata handling for current and target heights. - Created a backupStreamWriter to manage chunked writing of backup data. - Updated client tests to validate the Backup functionality. - Enhanced mock store to support Backup method for testing. - Added unit tests for Backup functionality in the store package.
1 parent 8abdd83 commit 9663554

19 files changed

Lines changed: 724 additions & 42 deletions

pkg/rpc/client/client.go

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@ package client
22

33
import (
44
"context"
5+
"fmt"
6+
"io"
57
"net/http"
68

79
"connectrpc.com/connect"
@@ -92,6 +94,54 @@ func (c *Client) GetMetadata(ctx context.Context, key string) ([]byte, error) {
9294
return resp.Msg.Value, nil
9395
}
9496

97+
// Backup streams a datastore backup into the provided writer and returns the final metadata emitted by the server.
98+
// The writer is not closed by this method.
99+
func (c *Client) Backup(ctx context.Context, params *pb.BackupRequest, dst io.Writer) (*pb.BackupMetadata, error) {
100+
if dst == nil {
101+
return nil, fmt.Errorf("backup destination writer cannot be nil")
102+
}
103+
104+
if params == nil {
105+
params = &pb.BackupRequest{}
106+
}
107+
108+
stream, err := c.storeClient.Backup(ctx, connect.NewRequest(params))
109+
if err != nil {
110+
return nil, err
111+
}
112+
defer stream.Close() // Best effort; ignore close error to preserve primary result.
113+
114+
var lastMetadata *pb.BackupMetadata
115+
for stream.Receive() {
116+
msg := stream.Msg()
117+
if metadata := msg.GetMetadata(); metadata != nil {
118+
lastMetadata = metadata
119+
continue
120+
}
121+
122+
if chunk := msg.GetChunk(); chunk != nil {
123+
if _, err := dst.Write(chunk); err != nil {
124+
_ = stream.Close()
125+
return lastMetadata, fmt.Errorf("failed to write backup chunk: %w", err)
126+
}
127+
}
128+
}
129+
130+
if err := stream.Err(); err != nil {
131+
return lastMetadata, err
132+
}
133+
134+
if lastMetadata == nil {
135+
return nil, fmt.Errorf("backup stream completed without metadata")
136+
}
137+
138+
if !lastMetadata.GetCompleted() {
139+
return lastMetadata, fmt.Errorf("backup stream ended without completion metadata")
140+
}
141+
142+
return lastMetadata, nil
143+
}
144+
95145
// GetPeerInfo returns information about the connected peers
96146
func (c *Client) GetPeerInfo(ctx context.Context) ([]*pb.PeerInfo, error) {
97147
req := connect.NewRequest(&emptypb.Empty{})

pkg/rpc/client/client_test.go

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,9 @@
11
package client
22

33
import (
4+
"bytes"
45
"context"
6+
"io"
57
"net/http"
68
"net/http/httptest"
79
"testing"
@@ -20,6 +22,7 @@ import (
2022
"github.com/evstack/ev-node/pkg/rpc/server"
2123
"github.com/evstack/ev-node/test/mocks"
2224
"github.com/evstack/ev-node/types"
25+
pb "github.com/evstack/ev-node/types/pb/evnode/v1"
2326
rpc "github.com/evstack/ev-node/types/pb/evnode/v1/v1connect"
2427
)
2528

@@ -173,6 +176,34 @@ func TestClientGetBlockByHash(t *testing.T) {
173176
mockStore.AssertExpectations(t)
174177
}
175178

179+
func TestClientBackup(t *testing.T) {
180+
mockStore := mocks.NewMockStore(t)
181+
mockP2P := mocks.NewMockP2PRPC(t)
182+
183+
mockStore.On("Height", mock.Anything).Return(uint64(15), nil)
184+
mockStore.On("Backup", mock.Anything, mock.Anything, mock.Anything).Run(func(args mock.Arguments) {
185+
writer := args.Get(1).(io.Writer)
186+
_, _ = writer.Write([]byte("chunk-1"))
187+
_, _ = writer.Write([]byte("chunk-2"))
188+
}).Return(uint64(42), nil)
189+
190+
testServer, client := setupTestServer(t, mockStore, mockP2P)
191+
defer testServer.Close()
192+
193+
var buf bytes.Buffer
194+
metadata, err := client.Backup(context.Background(), &pb.BackupRequest{TargetHeight: 10}, &buf)
195+
require.NoError(t, err)
196+
require.NotNil(t, metadata)
197+
require.Equal(t, "chunk-1chunk-2", buf.String())
198+
require.True(t, metadata.GetCompleted())
199+
require.Equal(t, uint64(15), metadata.GetCurrentHeight())
200+
require.Equal(t, uint64(10), metadata.GetTargetHeight())
201+
require.Equal(t, uint64(0), metadata.GetSinceVersion())
202+
require.Equal(t, uint64(42), metadata.GetLastVersion())
203+
204+
mockStore.AssertExpectations(t)
205+
}
206+
176207
func TestClientGetPeerInfo(t *testing.T) {
177208
// Create mocks
178209
mockStore := mocks.NewMockStore(t)

pkg/rpc/server/server.go

Lines changed: 161 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package server
33
import (
44
"context"
55
"fmt"
6+
"io"
67

78
"net/http"
89
"time"
@@ -188,6 +189,166 @@ func (s *StoreServer) GetMetadata(
188189
}), nil
189190
}
190191

192+
// Backup streams a Badger backup of the datastore so it can be persisted externally.
193+
func (s *StoreServer) Backup(
194+
ctx context.Context,
195+
req *connect.Request[pb.BackupRequest],
196+
stream *connect.ServerStream[pb.BackupResponse],
197+
) error {
198+
since := req.Msg.GetSinceVersion()
199+
targetHeight := req.Msg.GetTargetHeight()
200+
201+
currentHeight, err := s.store.Height(ctx)
202+
if err != nil {
203+
return connect.NewError(connect.CodeInternal, fmt.Errorf("failed to get current height: %w", err))
204+
}
205+
206+
if targetHeight != 0 && targetHeight > currentHeight {
207+
return connect.NewError(
208+
connect.CodeFailedPrecondition,
209+
fmt.Errorf("requested target height %d exceeds current height %d", targetHeight, currentHeight),
210+
)
211+
}
212+
213+
initialMetadata := &pb.BackupMetadata{
214+
CurrentHeight: currentHeight,
215+
TargetHeight: targetHeight,
216+
SinceVersion: since,
217+
Completed: false,
218+
LastVersion: 0,
219+
}
220+
221+
if err := stream.Send(&pb.BackupResponse{
222+
Response: &pb.BackupResponse_Metadata{
223+
Metadata: initialMetadata,
224+
},
225+
}); err != nil {
226+
return err
227+
}
228+
229+
writer := newBackupStreamWriter(stream, defaultBackupChunkSize)
230+
version, err := s.store.Backup(ctx, writer, since)
231+
if err != nil {
232+
var connectErr *connect.Error
233+
if errors.As(err, &connectErr) {
234+
return connectErr
235+
}
236+
if errors.Is(err, context.Canceled) {
237+
return connect.NewError(connect.CodeCanceled, err)
238+
}
239+
if errors.Is(err, context.DeadlineExceeded) {
240+
return connect.NewError(connect.CodeDeadlineExceeded, err)
241+
}
242+
return connect.NewError(connect.CodeInternal, fmt.Errorf("failed to execute backup: %w", err))
243+
}
244+
245+
if err := writer.Flush(); err != nil {
246+
var connectErr *connect.Error
247+
if errors.As(err, &connectErr) {
248+
return connectErr
249+
}
250+
if errors.Is(err, context.Canceled) {
251+
return connect.NewError(connect.CodeCanceled, err)
252+
}
253+
if errors.Is(err, context.DeadlineExceeded) {
254+
return connect.NewError(connect.CodeDeadlineExceeded, err)
255+
}
256+
return connect.NewError(connect.CodeInternal, fmt.Errorf("failed to flush backup stream: %w", err))
257+
}
258+
259+
completedMetadata := &pb.BackupMetadata{
260+
CurrentHeight: currentHeight,
261+
TargetHeight: targetHeight,
262+
SinceVersion: since,
263+
LastVersion: version,
264+
Completed: true,
265+
}
266+
267+
if err := stream.Send(&pb.BackupResponse{
268+
Response: &pb.BackupResponse_Metadata{
269+
Metadata: completedMetadata,
270+
},
271+
}); err != nil {
272+
return err
273+
}
274+
275+
return nil
276+
}
277+
278+
const defaultBackupChunkSize = 128 * 1024
279+
280+
var _ io.Writer = (*backupStreamWriter)(nil)
281+
282+
type backupStreamWriter struct {
283+
stream *connect.ServerStream[pb.BackupResponse]
284+
buf []byte
285+
chunkSize int
286+
}
287+
288+
func newBackupStreamWriter(stream *connect.ServerStream[pb.BackupResponse], chunkSize int) *backupStreamWriter {
289+
if chunkSize <= 0 {
290+
chunkSize = defaultBackupChunkSize
291+
}
292+
293+
return &backupStreamWriter{
294+
stream: stream,
295+
buf: make([]byte, 0, chunkSize),
296+
chunkSize: chunkSize,
297+
}
298+
}
299+
300+
func (w *backupStreamWriter) Write(p []byte) (int, error) {
301+
written := 0
302+
for len(p) > 0 {
303+
space := w.chunkSize - len(w.buf)
304+
if space == 0 {
305+
if err := w.flush(); err != nil {
306+
return written, err
307+
}
308+
space = w.chunkSize - len(w.buf)
309+
}
310+
311+
if space > len(p) {
312+
space = len(p)
313+
}
314+
315+
w.buf = append(w.buf, p[:space]...)
316+
p = p[space:]
317+
written += space
318+
319+
if len(w.buf) == w.chunkSize {
320+
if err := w.flush(); err != nil {
321+
return written, err
322+
}
323+
}
324+
}
325+
return written, nil
326+
}
327+
328+
func (w *backupStreamWriter) Flush() error {
329+
return w.flush()
330+
}
331+
332+
func (w *backupStreamWriter) flush() error {
333+
if len(w.buf) == 0 {
334+
return nil
335+
}
336+
337+
chunk := make([]byte, len(w.buf))
338+
copy(chunk, w.buf)
339+
340+
if err := w.stream.Send(&pb.BackupResponse{
341+
Response: &pb.BackupResponse_Chunk{
342+
Chunk: chunk,
343+
},
344+
}); err != nil {
345+
return err
346+
}
347+
348+
w.buf = w.buf[:0]
349+
return nil
350+
}
351+
191352
type ConfigServer struct {
192353
config config.Config
193354
signer []byte

pkg/store/backup.go

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
package store
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"io"
7+
8+
badger4 "github.com/ipfs/go-ds-badger4"
9+
)
10+
11+
// Backup streams the underlying Badger datastore snapshot into the provided writer.
12+
// The returned uint64 corresponds to the last version contained in the backup stream,
13+
// which can be re-used to generate incremental backups via the since parameter.
14+
func (s *DefaultStore) Backup(ctx context.Context, writer io.Writer, since uint64) (uint64, error) {
15+
if err := ctx.Err(); err != nil {
16+
return 0, err
17+
}
18+
19+
// Try to leverage a native backup implementation if the underlying datastore exposes one.
20+
type backupable interface {
21+
Backup(io.Writer, uint64) (uint64, error)
22+
}
23+
if dsBackup, ok := s.db.(backupable); ok {
24+
version, err := dsBackup.Backup(writer, since)
25+
if err != nil {
26+
return 0, fmt.Errorf("datastore backup failed: %w", err)
27+
}
28+
return version, nil
29+
}
30+
31+
// Default Badger datastore used across ev-node.
32+
badgerDatastore, ok := s.db.(*badger4.Datastore)
33+
if !ok {
34+
return 0, fmt.Errorf("backup is not supported by the configured datastore")
35+
}
36+
37+
// `badger.DB.Backup` internally orchestrates a consistent snapshot without pausing writes.
38+
version, err := badgerDatastore.DB.Backup(writer, since)
39+
if err != nil {
40+
return 0, fmt.Errorf("badger backup failed: %w", err)
41+
}
42+
return version, nil
43+
}

pkg/store/store_test.go

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package store
22

33
import (
4+
"bytes"
45
"context"
56
"encoding/binary"
67
"errors"
@@ -1127,3 +1128,24 @@ func TestRollbackDAIncludedHeightGetMetadataError(t *testing.T) {
11271128
require.Contains(err.Error(), "failed to get DA included height")
11281129
require.Contains(err.Error(), "metadata retrieval failed")
11291130
}
1131+
1132+
func TestDefaultStoreBackup(t *testing.T) {
1133+
t.Parallel()
1134+
1135+
ctx := context.Background()
1136+
kv, err := NewDefaultInMemoryKVStore()
1137+
require.NoError(t, err)
1138+
1139+
s := New(kv)
1140+
t.Cleanup(func() {
1141+
require.NoError(t, s.Close())
1142+
})
1143+
1144+
require.NoError(t, s.SetMetadata(ctx, "backup-test", []byte("value")))
1145+
1146+
var buf bytes.Buffer
1147+
version, err := s.Backup(ctx, &buf, 0)
1148+
require.NoError(t, err)
1149+
require.NotZero(t, buf.Len())
1150+
require.NotZero(t, version)
1151+
}

pkg/store/types.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package store
22

33
import (
44
"context"
5+
"io"
56

67
"github.com/evstack/ev-node/types"
78
)
@@ -50,6 +51,10 @@ type Store interface {
5051
// Aggregator is used to determine if the rollback is performed on the aggregator node.
5152
Rollback(ctx context.Context, height uint64, aggregator bool) error
5253

54+
// Backup writes a consistent backup stream to writer. The returned version can be used
55+
// as the starting point for incremental backups.
56+
Backup(ctx context.Context, writer io.Writer, since uint64) (uint64, error)
57+
5358
// Close safely closes underlying data storage, to ensure that data is actually saved.
5459
Close() error
5560
}

0 commit comments

Comments
 (0)