Skip to content

Commit 76cf2b8

Browse files
committed
[raft/aux] Add aux Raftstore
1 parent 7ea03cd commit 76cf2b8

7 files changed

Lines changed: 131 additions & 22 deletions

File tree

pkg/aux_/store/raftstore/dss.go

Lines changed: 50 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -2,25 +2,67 @@ package raftstore
22

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

67
auxmodels "github.com/interuss/dss/pkg/aux_/models"
8+
auxraftparams "github.com/interuss/dss/pkg/aux_/store/raftstore/params"
79
dsserr "github.com/interuss/dss/pkg/errors"
10+
"github.com/interuss/dss/pkg/timestamp"
811
"github.com/interuss/stacktrace"
912
)
1013

11-
// SaveOwnMetadata returns nil instead of dsserr.NotImplemented because it is needed to allow the server to startup.
12-
func (r *repo) SaveOwnMetadata(_ context.Context, locality string, publicEndpoint string) error {
13-
return nil
14+
type saveOwnMetadataPayload struct {
15+
Locality string
16+
PublicEndpoint string
1417
}
1518

16-
func (r *repo) GetDSSMetadata(_ context.Context) ([]*auxmodels.DSSMetadata, error) {
17-
return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "GetDSSMetadata not implemented for raftstore")
19+
func (r *repo) SaveOwnMetadata(ctx context.Context, locality string, publicEndpoint string) error {
20+
if locality == "" {
21+
return stacktrace.NewErrorWithCode(dsserr.BadRequest, "Locality not set")
22+
}
23+
if publicEndpoint == "" {
24+
return stacktrace.NewErrorWithCode(dsserr.BadRequest, "Public endpoint not set")
25+
}
26+
27+
_, err := r.consensus.ProposeValue(ctx, saveOwnMetadata, saveOwnMetadataPayload{
28+
Locality: locality,
29+
PublicEndpoint: publicEndpoint,
30+
}, false)
31+
return err
32+
}
33+
34+
func (r *repo) GetDSSMetadata(ctx context.Context) ([]*auxmodels.DSSMetadata, error) {
35+
result, err := r.consensus.ProposeValue(ctx, getDSSMetadata, nil, true)
36+
if err != nil {
37+
return nil, stacktrace.Propagate(err, "failed to propose %s", getDSSMetadata)
38+
}
39+
if result == nil {
40+
return nil, nil
41+
}
42+
return result.([]*auxmodels.DSSMetadata), nil
1843
}
1944

20-
func (r *repo) RecordHeartbeat(_ context.Context, heartbeat auxmodels.Heartbeat) error {
21-
return stacktrace.NewErrorWithCode(dsserr.NotImplemented, "RecordHeartbeat not implemented for raftstore")
45+
func (r *repo) RecordHeartbeat(ctx context.Context, heartbeat auxmodels.Heartbeat) error {
46+
if heartbeat.Locality == "" {
47+
return stacktrace.NewErrorWithCode(dsserr.BadRequest, "Locality not set")
48+
}
49+
if heartbeat.Source == "" {
50+
return stacktrace.NewErrorWithCode(dsserr.BadRequest, "Source not set")
51+
}
52+
53+
if heartbeat.Timestamp == nil {
54+
now := timestamp.NowFromContext(ctx)
55+
heartbeat.Timestamp = &now
56+
}
57+
58+
if heartbeat.NextHeartbeatExpectedBefore != nil && heartbeat.NextHeartbeatExpectedBefore.Before(*heartbeat.Timestamp) {
59+
return stacktrace.NewErrorWithCode(dsserr.BadRequest, "Cannot expect the timestamp of the next heartbeat before the timestamp of the new heartbeat")
60+
}
61+
62+
_, err := r.consensus.ProposeValue(ctx, recordHeartbeat, heartbeat, false)
63+
return err
2264
}
2365

2466
func (r *repo) GetDSSAirspaceRepresentationID(_ context.Context) (string, error) {
25-
return "", stacktrace.NewErrorWithCode(dsserr.NotImplemented, "GetDSSAirspaceRepresentationID not implemented for raftstore")
67+
return strconv.Itoa(int(auxraftparams.GetClusterID())), nil
2668
}

pkg/aux_/store/raftstore/params/params.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,3 +27,7 @@ func GetConnectParameters() (raftparams.ConnectParameters, error) {
2727
p.Peers = peers
2828
return p, nil
2929
}
30+
31+
func GetClusterID() uint64 {
32+
return storeParams.ClusterID
33+
}

pkg/aux_/store/raftstore/store.go

Lines changed: 66 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2,39 +2,96 @@ package raftstore
22

33
import (
44
"context"
5+
"encoding/json"
56

7+
auxmodels "github.com/interuss/dss/pkg/aux_/models"
68
"github.com/interuss/dss/pkg/aux_/repos"
9+
auxmemstore "github.com/interuss/dss/pkg/aux_/store/memstore"
710
auxraftparams "github.com/interuss/dss/pkg/aux_/store/raftstore/params"
8-
dsserr "github.com/interuss/dss/pkg/errors"
11+
"github.com/interuss/dss/pkg/memstore"
912
"github.com/interuss/dss/pkg/raftstore"
1013
"github.com/interuss/dss/pkg/raftstore/consensus"
1114
"github.com/interuss/stacktrace"
1215
"go.uber.org/zap"
1316
)
1417

18+
const storeID = "aux_"
19+
20+
const (
21+
saveOwnMetadata raftstore.RequestType = "saveOwnMetadata"
22+
getDSSMetadata raftstore.RequestType = "getDSSMetadata"
23+
recordHeartbeat raftstore.RequestType = "recordHeartbeat"
24+
)
25+
1526
// repo is a full implementation of aux_.repos.Repository for Raft-based storage.
16-
type repo struct{}
27+
type repo struct {
28+
consensus *consensus.Consensus
29+
memStore *memstore.Store[repos.Repository]
30+
}
1731

1832
func Init(ctx context.Context, logger *zap.Logger) (*raftstore.Store[repos.Repository], error) {
1933
params, err := auxraftparams.GetConnectParameters()
2034
if err != nil {
2135
return nil, stacktrace.Propagate(err, "failed to get aux raft parameters")
2236
}
23-
return raftstore.Init(ctx, logger, params, "aux", &repo{})
37+
38+
memStore, err := auxmemstore.Init(ctx, logger)
39+
if err != nil {
40+
return nil, stacktrace.Propagate(err, "failed to initialize aux memstore")
41+
}
42+
43+
r := &repo{memStore: memStore}
44+
store, err := raftstore.Init(ctx, logger, params, r)
45+
if err != nil {
46+
return nil, stacktrace.Propagate(err, "failed to initialize aux raftstore")
47+
}
48+
49+
r.consensus = store.Consensus
50+
51+
return store, nil
2452
}
2553

2654
func (r *repo) GetRepo() repos.Repository { return r }
2755

28-
func (r *repo) IsReadOnly(_ raftstore.RequestType) bool { return false }
56+
func (r *repo) IsReadOnly(requestType raftstore.RequestType) bool {
57+
return requestType == getDSSMetadata
58+
}
2959

3060
func (r *repo) GetSnapshot() ([]byte, error) {
31-
return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "not implemented yet")
61+
return r.memStore.GetSnapshot()
3262
}
3363

34-
func (r *repo) RestoreFromSnapshot([]byte) error {
35-
return stacktrace.NewErrorWithCode(dsserr.NotImplemented, "not implemented yet")
64+
func (r *repo) RestoreFromSnapshot(data []byte) error {
65+
return r.memStore.RestoreFromSnapshot(data)
3666
}
3767

38-
func (r *repo) Apply(_ context.Context, _ consensus.Proposal) (any, error) {
39-
return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "not implemented yet")
68+
func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, error) {
69+
mem, err := r.memStore.Interact(ctx)
70+
if err != nil {
71+
return nil, stacktrace.Propagate(err, "failed to obtain aux memstore repository")
72+
}
73+
74+
switch raftstore.RequestType(proposal.RequestType) {
75+
case saveOwnMetadata:
76+
var p saveOwnMetadataPayload
77+
if err := json.Unmarshal(proposal.Value, &p); err != nil {
78+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", saveOwnMetadata)
79+
}
80+
81+
return nil, mem.SaveOwnMetadata(ctx, p.Locality, p.PublicEndpoint)
82+
83+
case getDSSMetadata:
84+
return mem.GetDSSMetadata(ctx)
85+
86+
case recordHeartbeat:
87+
var hb auxmodels.Heartbeat
88+
if err := json.Unmarshal(proposal.Value, &hb); err != nil {
89+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", recordHeartbeat)
90+
}
91+
92+
return nil, mem.RecordHeartbeat(ctx, hb)
93+
94+
default:
95+
return nil, stacktrace.NewError("unknown request type: %q", proposal.RequestType)
96+
}
4097
}

pkg/memstore/store.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,14 @@ func (s *Store[R]) Restore(cp any) error {
7474
return s.memRepo.Restore(cp)
7575
}
7676

77+
func (s *Store[R]) GetSnapshot() ([]byte, error) {
78+
return s.memRepo.GetSnapshot()
79+
}
80+
81+
func (s *Store[R]) RestoreFromSnapshot(data []byte) error {
82+
return s.memRepo.RestoreFromSnapshot(data)
83+
}
84+
7785
func (s *Store[R]) Close() error {
7886
return nil
7987
}

pkg/raftstore/store.go

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,18 +25,16 @@ type RaftRepo[R any] interface {
2525
type Store[R any] struct {
2626
logger *zap.Logger
2727

28-
name string
2928
raftRepo RaftRepo[R]
3029
cancel context.CancelFunc
3130

3231
Consensus *consensus.Consensus
3332
}
3433

35-
func Init[R any](ctx context.Context, logger *zap.Logger, params raftparams.ConnectParameters, name string, r RaftRepo[R]) (*Store[R], error) {
34+
func Init[R any](ctx context.Context, logger *zap.Logger, params raftparams.ConnectParameters, r RaftRepo[R]) (*Store[R], error) {
3635
ctx, cancel := context.WithCancel(ctx)
3736

3837
store := &Store[R]{
39-
name: name,
4038
raftRepo: r,
4139
logger: logging.WithValuesFromContext(ctx, logger),
4240
cancel: cancel,

pkg/rid/store/raftstore/store.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ func Init(ctx context.Context, logger *zap.Logger) (*raftstore.Store[repos.Repos
2020
if err != nil {
2121
return nil, stacktrace.Propagate(err, "failed to get rid raft parameters")
2222
}
23-
return raftstore.Init(ctx, logger, params, "rid", &repo{})
23+
return raftstore.Init(ctx, logger, params, &repo{})
2424
}
2525

2626
func (r *repo) GetRepo() repos.Repository { return r }

pkg/scd/store/raftstore/store.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ func Init(ctx context.Context, logger *zap.Logger) (*raftstore.Store[repos.Repos
2020
if err != nil {
2121
return nil, stacktrace.Propagate(err, "failed to get scd raft parameters")
2222
}
23-
return raftstore.Init(ctx, logger, params, "scd", &repo{})
23+
return raftstore.Init(ctx, logger, params, &repo{})
2424
}
2525

2626
func (r *repo) GetRepo() repos.Repository { return r }

0 commit comments

Comments
 (0)