Skip to content
Merged
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
10 changes: 5 additions & 5 deletions cmds/core-service/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ func createAuxServer(ctx context.Context, locality string, publicEndpoint string
return nil, stacktrace.NewError("Public endpoint not set")
}

auxStore, err := auxs.Init(ctx, logger, true, locality)
auxStore, err := auxs.Init(ctx, logger, true)
if err != nil {
return nil, err
}
Expand All @@ -122,7 +122,7 @@ func createAuxServer(ctx context.Context, locality string, publicEndpoint string

func createRIDServers(ctx context.Context, locality string, logger *zap.Logger) (*rid_v1.Server, *rid_v2.Server, error) {

ridStore, err := rids.Init(ctx, logger, true, locality)
ridStore, err := rids.Init(ctx, logger, true)
if err != nil {
return nil, nil, err
}
Expand Down Expand Up @@ -151,9 +151,9 @@ func createRIDServers(ctx context.Context, locality string, logger *zap.Logger)
}, nil
}

func createSCDServer(ctx context.Context, logger *zap.Logger, locality string) (*scd.Server, error) {
func createSCDServer(ctx context.Context, logger *zap.Logger) (*scd.Server, error) {

scdStore, err := scds.Init(ctx, logger, true, locality)
scdStore, err := scds.Init(ctx, logger, true)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -350,7 +350,7 @@ func RunHTTPServer(ctx context.Context, ctxCanceler func(), address, locality st

// Initialize strategic conflict detection
if *enableSCD {
scdV1Server, err = createSCDServer(ctx, logger, locality)
scdV1Server, err = createSCDServer(ctx, logger)
if err != nil {
return stacktrace.Propagate(err, "Failed to create strategic conflict detection server")
}
Expand Down
4 changes: 2 additions & 2 deletions cmds/db-manager/cleanup/evict.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,12 +64,12 @@ func evict(cmd *cobra.Command, _ []string) error {

logger := logging.WithValuesFromContext(ctx, logging.Logger)

scdStore, err := scds.Init(ctx, logger, false, *locality)
scdStore, err := scds.Init(ctx, logger, false)
if err != nil {
return err
}

ridStore, err := rids.Init(ctx, logger, false, *locality)
ridStore, err := rids.Init(ctx, logger, false)
if err != nil {
return err
}
Expand Down
4 changes: 2 additions & 2 deletions pkg/aux_/store/raftstore/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ type repo struct {
*memstore.Store[repos.Repository]
}

func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.Store[repos.Repository], error) {
func Init(ctx context.Context, logger *zap.Logger) (*raftstore.Store[repos.Repository], error) {
params, err := auxraftparams.GetConnectParameters()
if err != nil {
return nil, stacktrace.Propagate(err, "failed to get aux raft parameters")
Expand All @@ -39,7 +39,7 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.
}

r := &repo{Store: memStore}
store, err := raftstore.Init(ctx, logger.With(zap.String("service", "aux_")), locality, params, r, nil)
store, err := raftstore.Init(ctx, logger.With(zap.String("service", "aux_")), params, r, nil)
if err != nil {
return nil, stacktrace.Propagate(err, "failed to initialize aux raftstore")
}
Expand Down
4 changes: 2 additions & 2 deletions pkg/aux_/store/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,12 @@ import (
type Store = dssstore.Store[repos.Repository]

// Init selects and initializes the aux store backend.
func Init(ctx context.Context, logger *zap.Logger, withCheckCron bool, locality string) (Store, error) {
func Init(ctx context.Context, logger *zap.Logger, withCheckCron bool) (Store, error) {
switch storeType := params.GetStoreParameters().StoreType; storeType {
case params.SQLStoreType:
return auxsqlstore.Init(ctx, logger, withCheckCron)
case params.RaftStoreType:
return auxraftstore.Init(ctx, logger, locality)
return auxraftstore.Init(ctx, logger)
case params.MemStoreType:
return auxmemstore.Init(ctx, logger)
default:
Expand Down
16 changes: 7 additions & 9 deletions pkg/raftstore/consensus/consensus.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,8 @@ import (
type Consensus struct {
logger *zap.Logger

locality string
nodeID uint64
node raft.Node
nodeID uint64
node raft.Node

transport *rafthttp.Transport
server *http.Server
Expand All @@ -47,7 +46,7 @@ type Consensus struct {
appliedIndex uint64
}

func NewConsensus(ctx context.Context, logger *zap.Logger, locality string, connectParams params.ConnectParameters, provider snapshotProvider, commitC chan<- EntryCommit) (*Consensus, error) {
func NewConsensus(ctx context.Context, logger *zap.Logger, connectParams params.ConnectParameters, provider snapshotProvider, commitC chan<- EntryCommit) (*Consensus, error) {
storage, old, err := newStorage(ctx, logger.With(zap.String("component", "storage")), connectParams.DataDir, connectParams.NodeID, provider, connectParams.SnapshotCatchupEntries)
if err != nil {
return nil, stacktrace.Propagate(err, "failed to initialize storage")
Expand Down Expand Up @@ -76,11 +75,10 @@ func NewConsensus(ctx context.Context, logger *zap.Logger, locality string, conn
consensus := &Consensus{
logger: logging.WithValuesFromContext(ctx, logger),

nodeID: connectParams.NodeID,
node: node,
locality: locality,
storage: storage,
commitC: commitC,
nodeID: connectParams.NodeID,
node: node,
storage: storage,
commitC: commitC,

shutdownTimeout: 2 * connectParams.ElectionInterval(),
serverErrC: make(chan error, 1),
Expand Down
2 changes: 0 additions & 2 deletions pkg/raftstore/consensus/proposal.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ type EntryCommit struct {

type Proposal struct {
ID string `json:"id"`
Locality string `json:"locality"`
NodeID uint64 `json:"node_id"`
Timestamp time.Time `json:"timestamp"`
RequestType string `json:"request_type"`
Expand All @@ -33,7 +32,6 @@ func (c *Consensus) newProposal(ctx context.Context, requestType string, value [

return Proposal{
ID: uuid.NewString(),
Locality: c.locality,
NodeID: c.nodeID,
Timestamp: timestamp.UTC(),
RequestType: requestType,
Expand Down
4 changes: 2 additions & 2 deletions pkg/raftstore/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ type Store[R any] struct {
done chan struct{}
}

func Init[R any](ctx context.Context, logger *zap.Logger, locality string, params raftparams.ConnectParameters, r RaftRepo[R], registry map[string]store.OperationHandler[R]) (*Store[R], error) {
func Init[R any](ctx context.Context, logger *zap.Logger, params raftparams.ConnectParameters, r RaftRepo[R], registry map[string]store.OperationHandler[R]) (*Store[R], error) {
ctx, cancel := context.WithCancel(ctx)

store := &Store[R]{
Expand All @@ -50,7 +50,7 @@ func Init[R any](ctx context.Context, logger *zap.Logger, locality string, param
store.processCommits(ctx, commitC)
}()

consensusInstance, err := consensus.NewConsensus(ctx, logger, locality, params, r.GetSnapshot, commitC)
consensusInstance, err := consensus.NewConsensus(ctx, logger, params, r.GetSnapshot, commitC)
if err != nil {
return nil, stacktrace.Propagate(err, "failed to initialize consensus")
}
Expand Down
4 changes: 2 additions & 2 deletions pkg/rid/store/raftstore/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ type repo struct {
*memstore.Store[repos.Repository]
}

func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.Store[repos.Repository], error) {
func Init(ctx context.Context, logger *zap.Logger) (*raftstore.Store[repos.Repository], error) {
params, err := ridraftparams.GetConnectParameters()
if err != nil {
return nil, stacktrace.Propagate(err, "failed to get rid raft parameters")
Expand All @@ -36,7 +36,7 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.
}

r := &repo{Store: memStore}
store, err := raftstore.Init(ctx, logger.With(zap.String("service", "rid")), locality, params, r, operations.Registry)
store, err := raftstore.Init(ctx, logger.With(zap.String("service", "rid")), params, r, operations.Registry)
if err != nil {
return nil, stacktrace.Propagate(err, "failed to initialize rid raftstore")
}
Expand Down
4 changes: 2 additions & 2 deletions pkg/rid/store/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,13 @@ import (
type Store = dssstore.Store[repos.Repository]

// Init selects and initializes the rid store backend.
func Init(ctx context.Context, logger *zap.Logger, withCheckCron bool, locality string) (Store, error) {
func Init(ctx context.Context, logger *zap.Logger, withCheckCron bool) (Store, error) {
storeType := params.GetStoreParameters().StoreType
switch storeType {
case params.SQLStoreType:
return ridsqlstore.Init(ctx, logger, withCheckCron)
case params.RaftStoreType:
return ridraftstore.Init(ctx, logger, locality)
return ridraftstore.Init(ctx, logger)
case params.MemStoreType:
return ridmemstore.Init(ctx, logger)
default:
Expand Down
4 changes: 2 additions & 2 deletions pkg/scd/store/raftstore/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ type repo struct {
*memstore.Store[repos.Repository]
}

func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.Store[repos.Repository], error) {
func Init(ctx context.Context, logger *zap.Logger) (*raftstore.Store[repos.Repository], error) {
params, err := scdraftparams.GetConnectParameters()
if err != nil {
return nil, stacktrace.Propagate(err, "failed to get scd raft parameters")
Expand All @@ -36,7 +36,7 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.
}

r := &repo{Store: memStore}
store, err := raftstore.Init(ctx, logger.With(zap.String("service", "scd")), locality, params, r, operations.Registry)
store, err := raftstore.Init(ctx, logger.With(zap.String("service", "scd")), params, r, operations.Registry)
if err != nil {
return nil, stacktrace.Propagate(err, "failed to initialize scd raftstore")
}
Expand Down
4 changes: 2 additions & 2 deletions pkg/scd/store/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,13 @@ import (
type Store = dssstore.Store[repos.Repository]

// Init selects and initializes the scd store backend.
func Init(ctx context.Context, logger *zap.Logger, withCheckCron bool, locality string) (Store, error) {
func Init(ctx context.Context, logger *zap.Logger, withCheckCron bool) (Store, error) {
storeType := params.GetStoreParameters().StoreType
switch storeType {
case params.SQLStoreType:
return scdsqlstore.Init(ctx, logger, withCheckCron)
case params.RaftStoreType:
return scdraftstore.Init(ctx, logger, locality)
return scdraftstore.Init(ctx, logger)
case params.MemStoreType:
return scdmemstore.Init(ctx, logger)
default:
Expand Down
Loading