From ddb1b0cc09220fc0622d8b997cceade75fd92274 Mon Sep 17 00:00:00 2001 From: Mariem Baccari Date: Wed, 19 Aug 2026 13:51:12 +0200 Subject: [PATCH] [raftstore] Embed memstore and use checkpoint --- pkg/aux_/store/raftstore/store.go | 19 +++++-------------- pkg/raftstore/store.go | 17 ++++++++--------- pkg/rid/store/raftstore/store.go | 15 +++------------ pkg/scd/store/raftstore/store.go | 15 +++------------ 4 files changed, 19 insertions(+), 47 deletions(-) diff --git a/pkg/aux_/store/raftstore/store.go b/pkg/aux_/store/raftstore/store.go index 8f8a4270f..7b55231ed 100644 --- a/pkg/aux_/store/raftstore/store.go +++ b/pkg/aux_/store/raftstore/store.go @@ -24,8 +24,7 @@ const ( // repo is a full implementation of aux_.repos.Repository for Raft-based storage. type repo struct { consensus *consensus.Consensus - memStore *memstore.Store[repos.Repository] - memRepo repos.Repository + *memstore.Store[repos.Repository] } func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.Store[repos.Repository], error) { @@ -39,7 +38,7 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore. return nil, stacktrace.Propagate(err, "failed to initialize aux memstore") } - r := &repo{memStore: memStore, memRepo: memStore.GetRepo()} + r := &repo{Store: memStore} store, err := raftstore.Init(ctx, logger.With(zap.String("service", "aux_")), locality, params, r, nil) if err != nil { return nil, stacktrace.Propagate(err, "failed to initialize aux raftstore") @@ -52,14 +51,6 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore. func (r *repo) GetRepo() repos.Repository { return r } -func (r *repo) GetSnapshot() ([]byte, error) { - return r.memStore.GetSnapshot() -} - -func (r *repo) RestoreFromSnapshot(data []byte) error { - return r.memStore.RestoreFromSnapshot(data) -} - func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, error) { switch proposal.RequestType { case saveOwnMetadata: @@ -68,10 +59,10 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", saveOwnMetadata) } - return nil, r.memRepo.SaveOwnMetadata(ctx, payload.Locality, payload.PublicEndpoint) + return nil, r.Store.GetRepo().SaveOwnMetadata(ctx, payload.Locality, payload.PublicEndpoint) case getDSSMetadata: - return r.memRepo.GetDSSMetadata(ctx) + return r.Store.GetRepo().GetDSSMetadata(ctx) case recordHeartbeat: var heartbeat auxmodels.Heartbeat @@ -79,7 +70,7 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", recordHeartbeat) } - return nil, r.memRepo.RecordHeartbeat(ctx, heartbeat) + return nil, r.Store.GetRepo().RecordHeartbeat(ctx, heartbeat) default: return nil, stacktrace.NewError("unknown request type: %q", proposal.RequestType) diff --git a/pkg/raftstore/store.go b/pkg/raftstore/store.go index b3c6919aa..436f4bf8f 100644 --- a/pkg/raftstore/store.go +++ b/pkg/raftstore/store.go @@ -5,6 +5,7 @@ import ( "github.com/interuss/dss/pkg/locality" "github.com/interuss/dss/pkg/logging" + "github.com/interuss/dss/pkg/memstore" "github.com/interuss/dss/pkg/raftstore/consensus" raftparams "github.com/interuss/dss/pkg/raftstore/params" "github.com/interuss/dss/pkg/random" @@ -15,19 +16,12 @@ import ( ) type RaftRepo[R any] interface { - GetRepo() R + memstore.MemRepo[R] + // Apply is called on every committed entry. The proposal must be applied atomically. // The any return mirrors store.OperationHandler.Execute: different requests yield different // concrete result types. Callers recover the type via store.TransactWithResult. Apply(ctx context.Context, proposal consensus.Proposal) (any, error) - - // GetSnapshot returns a serialized view of current state, suitable - // for restoring via RestoreFromSnapshot. - GetSnapshot() ([]byte, error) - - // RestoreFromSnapshot replaces all state with the snapshot in data. - // data is always the output of a prior GetSnapshot. - RestoreFromSnapshot(data []byte) error } type Store[R any] struct { @@ -119,7 +113,12 @@ func (s *Store[R]) processCommits(ctx context.Context, commitCh <-chan consensus proposalCtx := timestamp.NewContext(ctx, commit.Prop.Timestamp) proposalCtx = locality.NewContext(proposalCtx, commit.Prop.Locality) proposalCtx = random.NewContext(proposalCtx, commit.Prop.Seed) + s.raftRepo.Checkpoint() result, err := s.raftRepo.Apply(proposalCtx, commit.Prop) + if err != nil { + s.logger.Warn("failed to apply proposal, rolling back", zap.String("proposal_id", commit.Prop.ID), zap.String("proposal_type", string(commit.Prop.RequestType)), zap.Error(err)) + s.raftRepo.Restore() + } commit.Done <- consensus.ProposalResult{Result: result, Error: err} } } diff --git a/pkg/rid/store/raftstore/store.go b/pkg/rid/store/raftstore/store.go index 00b4bf0c5..cfbf6dd51 100644 --- a/pkg/rid/store/raftstore/store.go +++ b/pkg/rid/store/raftstore/store.go @@ -17,8 +17,7 @@ import ( // repo is a full implementation of rid.repos.Repository for Raft-based storage. type repo struct { consensus *consensus.Consensus - memStore *memstore.Store[repos.Repository] - memRepo repos.Repository + *memstore.Store[repos.Repository] } func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.Store[repos.Repository], error) { @@ -32,7 +31,7 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore. return nil, stacktrace.Propagate(err, "failed to initialize rid memstore") } - r := &repo{memStore: memStore, memRepo: memStore.GetRepo()} + r := &repo{Store: memStore} store, err := raftstore.Init(ctx, logger.With(zap.String("service", "rid")), locality, params, r, operations.Registry) if err != nil { return nil, stacktrace.Propagate(err, "failed to initialize rid raftstore") @@ -45,14 +44,6 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore. func (r *repo) GetRepo() repos.Repository { return r } -func (r *repo) GetSnapshot() ([]byte, error) { - return r.memStore.GetSnapshot() -} - -func (r *repo) RestoreFromSnapshot(data []byte) error { - return r.memStore.RestoreFromSnapshot(data) -} - func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, error) { switch proposal.RequestType { @@ -67,6 +58,6 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err return nil, stacktrace.Propagate(err, "failed to decode %s payload", proposal.RequestType) } - return handler.Execute(ctx, r.memRepo, request) + return handler.Execute(ctx, r.Store.GetRepo(), request) } } diff --git a/pkg/scd/store/raftstore/store.go b/pkg/scd/store/raftstore/store.go index baa1a11ae..7718c7d17 100644 --- a/pkg/scd/store/raftstore/store.go +++ b/pkg/scd/store/raftstore/store.go @@ -17,8 +17,7 @@ import ( // repo is a full implementation of scd.repos.Repository for Raft-based storage. type repo struct { consensus *consensus.Consensus - memStore *memstore.Store[repos.Repository] - memRepo repos.Repository + *memstore.Store[repos.Repository] } func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.Store[repos.Repository], error) { @@ -32,7 +31,7 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore. return nil, stacktrace.Propagate(err, "failed to initialize scd memstore") } - r := &repo{memStore: memStore, memRepo: memStore.GetRepo()} + r := &repo{Store: memStore} store, err := raftstore.Init(ctx, logger.With(zap.String("service", "scd")), locality, params, r, operations.Registry) if err != nil { return nil, stacktrace.Propagate(err, "failed to initialize scd raftstore") @@ -45,14 +44,6 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore. func (r *repo) GetRepo() repos.Repository { return r } -func (r *repo) GetSnapshot() ([]byte, error) { - return r.memStore.GetSnapshot() -} - -func (r *repo) RestoreFromSnapshot(data []byte) error { - return r.memStore.RestoreFromSnapshot(data) -} - func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, error) { switch proposal.RequestType { @@ -67,6 +58,6 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err return nil, stacktrace.Propagate(err, "failed to decode %s payload", proposal.RequestType) } - return handler.Execute(ctx, r.memRepo, request) + return handler.Execute(ctx, r.Store.GetRepo(), request) } }