From 3aa2882ed9394d18b9080fb142f0b23c11227061 Mon Sep 17 00:00:00 2001 From: Mariem Baccari Date: Tue, 6 Oct 2026 14:37:57 +0200 Subject: [PATCH] [raft] Remove locality from proposals --- cmds/core-service/main.go | 10 +++++----- cmds/db-manager/cleanup/evict.go | 4 ++-- pkg/aux_/store/raftstore/store.go | 4 ++-- pkg/aux_/store/store.go | 4 ++-- pkg/raftstore/consensus/consensus.go | 16 +++++++--------- pkg/raftstore/consensus/proposal.go | 2 -- pkg/raftstore/store.go | 4 ++-- pkg/rid/store/raftstore/store.go | 4 ++-- pkg/rid/store/store.go | 4 ++-- pkg/scd/store/raftstore/store.go | 4 ++-- pkg/scd/store/store.go | 4 ++-- 11 files changed, 28 insertions(+), 32 deletions(-) diff --git a/cmds/core-service/main.go b/cmds/core-service/main.go index ac11141e7..af395a434 100644 --- a/cmds/core-service/main.go +++ b/cmds/core-service/main.go @@ -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 } @@ -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 } @@ -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 } @@ -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") } diff --git a/cmds/db-manager/cleanup/evict.go b/cmds/db-manager/cleanup/evict.go index 1e7737500..4bb446454 100644 --- a/cmds/db-manager/cleanup/evict.go +++ b/cmds/db-manager/cleanup/evict.go @@ -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 } diff --git a/pkg/aux_/store/raftstore/store.go b/pkg/aux_/store/raftstore/store.go index c05025757..0530b1f54 100644 --- a/pkg/aux_/store/raftstore/store.go +++ b/pkg/aux_/store/raftstore/store.go @@ -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") @@ -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") } diff --git a/pkg/aux_/store/store.go b/pkg/aux_/store/store.go index 93148ff80..2532842a7 100644 --- a/pkg/aux_/store/store.go +++ b/pkg/aux_/store/store.go @@ -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: diff --git a/pkg/raftstore/consensus/consensus.go b/pkg/raftstore/consensus/consensus.go index 49fa06776..2c981eb0c 100644 --- a/pkg/raftstore/consensus/consensus.go +++ b/pkg/raftstore/consensus/consensus.go @@ -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 @@ -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") @@ -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), diff --git a/pkg/raftstore/consensus/proposal.go b/pkg/raftstore/consensus/proposal.go index 17e39056c..a7433d54b 100644 --- a/pkg/raftstore/consensus/proposal.go +++ b/pkg/raftstore/consensus/proposal.go @@ -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"` @@ -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, diff --git a/pkg/raftstore/store.go b/pkg/raftstore/store.go index 380401fc4..1de45b76d 100644 --- a/pkg/raftstore/store.go +++ b/pkg/raftstore/store.go @@ -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]{ @@ -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") } diff --git a/pkg/rid/store/raftstore/store.go b/pkg/rid/store/raftstore/store.go index 5272f88ca..ebbe4c40d 100644 --- a/pkg/rid/store/raftstore/store.go +++ b/pkg/rid/store/raftstore/store.go @@ -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") @@ -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") } diff --git a/pkg/rid/store/store.go b/pkg/rid/store/store.go index d9a1235c2..7f2dc6c20 100644 --- a/pkg/rid/store/store.go +++ b/pkg/rid/store/store.go @@ -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: diff --git a/pkg/scd/store/raftstore/store.go b/pkg/scd/store/raftstore/store.go index 18b962bd3..44abc5ceb 100644 --- a/pkg/scd/store/raftstore/store.go +++ b/pkg/scd/store/raftstore/store.go @@ -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") @@ -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") } diff --git a/pkg/scd/store/store.go b/pkg/scd/store/store.go index 6a6b70665..0cdc48b93 100644 --- a/pkg/scd/store/store.go +++ b/pkg/scd/store/store.go @@ -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: