Add persistent stores for cluster raft data

This commit is contained in:
Shrenuj Bansal
2022-09-29 12:56:22 +00:00
parent f89a682d98
commit b8060cd8f7
5 changed files with 65 additions and 44 deletions
+11 -8
View File
@@ -1,7 +1,9 @@
package consensus
import (
"github.com/filecoin-project/lotus/node/repo"
"io/ioutil"
"path/filepath"
"time"
hraft "github.com/hashicorp/raft"
@@ -18,7 +20,7 @@ var envConfigKey = "cluster_raft"
// Configuration defaults
var (
DefaultDataSubFolder = "raft"
DefaultDataSubFolder = "raft-cluster"
DefaultWaitForLeaderTimeout = 15 * time.Second
DefaultCommitRetries = 1
DefaultNetworkTimeout = 100 * time.Second
@@ -350,13 +352,14 @@ func ValidateConfig(cfg *ClusterRaftConfig) error {
// return cfg.applyJSONConfig(jcfg)
//}
//
//// GetDataFolder returns the Raft data folder that we are using.
//func (cfg *Config) GetDataFolder() string {
// if cfg.DataFolder == "" {
// return filepath.Join(cfg.BaseDir, DefaultDataSubFolder)
// }
// return cfg.DataFolder
//}
// GetDataFolder returns the Raft data folder that we are using.
func (cfg *ClusterRaftConfig) GetDataFolder(repo repo.LockedRepo) string {
if cfg.DataFolder == "" {
return filepath.Join(repo.Path() + DefaultDataSubFolder)
}
return filepath.Join(repo.Path() + cfg.DataFolder)
}
//
//// ToDisplayJSON returns JSON config as a string.
//func (cfg *Config) ToDisplayJSON() ([]byte, error) {
+8 -4
View File
@@ -6,6 +6,7 @@ import (
"context"
"errors"
"fmt"
"github.com/filecoin-project/lotus/node/repo"
"sort"
"time"
@@ -83,6 +84,7 @@ type Consensus struct {
readyCh chan struct{}
peerSet []peer.ID
repo repo.LockedRepo
//shutdownLock sync.RWMutex
//shutdown bool
@@ -96,7 +98,7 @@ type Consensus struct {
//
// The staging parameter controls if the Raft peer should start in
// staging mode (used when joining a new Raft peerset with other peers).
func NewConsensus(host host.Host, cfg *ClusterRaftConfig, mpool *messagepool.MessagePool, staging bool) (*Consensus, error) {
func NewConsensus(host host.Host, cfg *ClusterRaftConfig, mpool *messagepool.MessagePool, repo repo.LockedRepo, staging bool) (*Consensus, error) {
err := ValidateConfig(cfg)
if err != nil {
return nil, err
@@ -109,7 +111,7 @@ func NewConsensus(host host.Host, cfg *ClusterRaftConfig, mpool *messagepool.Mes
consensus := libp2praft.NewOpLog(state, &ConsensusOp{})
raft, err := newRaftWrapper(host, cfg, consensus.FSM(), staging)
raft, err := newRaftWrapper(host, cfg, consensus.FSM(), repo, staging)
if err != nil {
logger.Error("error creating raft: ", err)
cancel()
@@ -130,6 +132,7 @@ func NewConsensus(host host.Host, cfg *ClusterRaftConfig, mpool *messagepool.Mes
peerSet: cfg.InitPeerset,
rpcReady: make(chan struct{}, 1),
readyCh: make(chan struct{}, 1),
repo: repo,
}
go cc.finishBootstrap()
@@ -141,10 +144,11 @@ func NewConsensusWithRPCClient(staging bool) func(host host.Host,
cfg *ClusterRaftConfig,
rpcClient *rpc.Client,
mpool *messagepool.MessagePool,
repo repo.LockedRepo,
) (*Consensus, error) {
return func(host host.Host, cfg *ClusterRaftConfig, rpcClient *rpc.Client, mpool *messagepool.MessagePool) (*Consensus, error) {
cc, err := NewConsensus(host, cfg, mpool, staging)
return func(host host.Host, cfg *ClusterRaftConfig, rpcClient *rpc.Client, mpool *messagepool.MessagePool, repo repo.LockedRepo) (*Consensus, error) {
cc, err := NewConsensus(host, cfg, mpool, repo, staging)
if err != nil {
return nil, err
}
+40 -30
View File
@@ -4,16 +4,23 @@ import (
"context"
"errors"
"fmt"
"github.com/filecoin-project/lotus/node/repo"
"github.com/ipfs/go-log/v2"
"go.uber.org/zap"
"io"
"os"
"path/filepath"
"time"
hraft "github.com/hashicorp/raft"
raftboltdb "github.com/hashicorp/raft-boltdb"
p2praft "github.com/libp2p/go-libp2p-raft"
host "github.com/libp2p/go-libp2p/core/host"
peer "github.com/libp2p/go-libp2p/core/peer"
)
var raftLogger = log.Logger("raft-cluster")
// ErrWaitingForSelf is returned when we are waiting for ourselves to depart
// the peer set, which won't happen
var errWaitingForSelf = errors.New("waiting for ourselves to depart")
@@ -49,8 +56,9 @@ type raftWrapper struct {
snapshotStore hraft.SnapshotStore
logStore hraft.LogStore
stableStore hraft.StableStore
//boltdb *raftboltdb.BoltStore
staging bool
boltdb *raftboltdb.BoltStore
repo repo.LockedRepo
staging bool
}
// newRaftWrapper creates a Raft instance and initializes
@@ -60,6 +68,7 @@ func newRaftWrapper(
host host.Host,
cfg *ClusterRaftConfig,
fsm hraft.FSM,
repo repo.LockedRepo,
staging bool,
) (*raftWrapper, error) {
@@ -67,18 +76,19 @@ func newRaftWrapper(
raftW.config = cfg
raftW.host = host
raftW.staging = staging
raftW.repo = repo
// Set correct LocalID
cfg.RaftConfig.LocalID = hraft.ServerID(peer.Encode(host.ID()))
//df := cfg.GetDataFolder()
//err := makeDataFolder(df)
//if err != nil {
// return nil, err
//}
df := cfg.GetDataFolder(repo)
err := makeDataFolder(df)
if err != nil {
return nil, err
}
raftW.makeServerConfig()
err := raftW.makeTransport()
err = raftW.makeTransport()
if err != nil {
return nil, err
}
@@ -124,13 +134,13 @@ func (rw *raftWrapper) makeTransport() (err error) {
}
func (rw *raftWrapper) makeStores() error {
//logger.Debug("creating BoltDB store")
//df := rw.config.GetDataFolder()
//store, err := raftboltdb.NewBoltStore(filepath.Join(df, "raft.db"))
//if err != nil {
// return err
//}
store := hraft.NewInmemStore()
logger.Debug("creating BoltDB store")
df := rw.config.GetDataFolder(rw.repo)
store, err := raftboltdb.NewBoltStore(filepath.Join(df, "raft.db"))
if err != nil {
return err
}
//store := hraft.NewInmemStore()
// wraps the store in a LogCache to improve performance.
// See consul/agent/consul/server.go
cacheStore, err := hraft.NewLogCache(RaftLogCacheSize, store)
@@ -138,22 +148,22 @@ func (rw *raftWrapper) makeStores() error {
return err
}
//logger.Debug("creating raft snapshot store")
//snapstore, err := hraft.NewFileSnapshotStoreWithLogger(
// df,
// RaftMaxSnapshots,
// raftStdLogger,
//)
logger.Debug("creating raft snapshot store")
snapstore, err := hraft.NewFileSnapshotStoreWithLogger(
df,
RaftMaxSnapshots,
zap.NewStdLog(log.Logger("raft-snapshot").SugaredLogger.Desugar()),
)
if err != nil {
return err
}
snapstore := hraft.NewInmemSnapshotStore()
//if err != nil {
// return err
//}
//snapstore := hraft.NewInmemSnapshotStore()
rw.logStore = cacheStore
rw.stableStore = store
rw.snapshotStore = snapstore
//rw.boltdb = store
rw.boltdb = store
return nil
}
@@ -420,10 +430,10 @@ func (rw *raftWrapper) Shutdown(ctx context.Context) error {
errMsgs += "could not shutdown raft: " + err.Error() + ".\n"
}
//err = rw.boltdb.Close() // important!
//if err != nil {
// errMsgs += "could not close boltdb: " + err.Error()
//}
err = rw.boltdb.Close() // important!
if err != nil {
errMsgs += "could not close boltdb: " + err.Error()
}
if errMsgs != "" {
return errors.New(errMsgs)