You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

159 lines
3.8 KiB

5 years ago
5 years ago
3 years ago
3 years ago
3 years ago
11 years ago
3 years ago
3 years ago
  1. package weed_server
  2. import (
  3. "encoding/json"
  4. "math/rand"
  5. "os"
  6. "path"
  7. "time"
  8. "google.golang.org/grpc"
  9. "github.com/chrislusf/seaweedfs/weed/pb"
  10. "github.com/chrislusf/raft"
  11. "github.com/chrislusf/seaweedfs/weed/glog"
  12. "github.com/chrislusf/seaweedfs/weed/topology"
  13. )
  14. type RaftServerOption struct {
  15. GrpcDialOption grpc.DialOption
  16. Peers []pb.ServerAddress
  17. ServerAddr pb.ServerAddress
  18. DataDir string
  19. Topo *topology.Topology
  20. RaftResumeState bool
  21. HeartbeatInterval time.Duration
  22. ElectionTimeout time.Duration
  23. }
  24. type RaftServer struct {
  25. peers map[string]pb.ServerAddress // initial peers to join with
  26. raftServer raft.Server
  27. dataDir string
  28. serverAddr pb.ServerAddress
  29. topo *topology.Topology
  30. *raft.GrpcServer
  31. }
  32. type StateMachine struct {
  33. raft.StateMachine
  34. topo *topology.Topology
  35. }
  36. func (s StateMachine) Save() ([]byte, error) {
  37. state := topology.MaxVolumeIdCommand{
  38. MaxVolumeId: s.topo.GetMaxVolumeId(),
  39. }
  40. glog.V(1).Infof("Save raft state %+v", state)
  41. return json.Marshal(state)
  42. }
  43. func (s StateMachine) Recovery(data []byte) error {
  44. state := topology.MaxVolumeIdCommand{}
  45. err := json.Unmarshal(data, &state)
  46. if err != nil {
  47. return err
  48. }
  49. glog.V(1).Infof("Recovery raft state %+v", state)
  50. s.topo.UpAdjustMaxVolumeId(state.MaxVolumeId)
  51. return nil
  52. }
  53. func NewRaftServer(option *RaftServerOption) (*RaftServer, error) {
  54. peers := make(map[string]pb.ServerAddress)
  55. for _, peer := range option.Peers {
  56. peers[peer.String()] = peer
  57. }
  58. s := &RaftServer{
  59. peers: peers,
  60. serverAddr: option.ServerAddr,
  61. dataDir: option.DataDir,
  62. topo: option.Topo,
  63. }
  64. if glog.V(4) {
  65. raft.SetLogLevel(2)
  66. }
  67. raft.RegisterCommand(&topology.MaxVolumeIdCommand{})
  68. var err error
  69. transporter := raft.NewGrpcTransporter(option.GrpcDialOption)
  70. glog.V(0).Infof("Starting RaftServer with %v", option.ServerAddr)
  71. // always clear previous log to avoid server is promotable
  72. os.RemoveAll(path.Join(s.dataDir, "log"))
  73. if !option.RaftResumeState {
  74. // always clear previous metadata
  75. os.RemoveAll(path.Join(s.dataDir, "conf"))
  76. os.RemoveAll(path.Join(s.dataDir, "snapshot"))
  77. }
  78. if err := os.MkdirAll(path.Join(s.dataDir, "snapshot"), 0600); err != nil {
  79. return nil, err
  80. }
  81. stateMachine := StateMachine{topo: option.Topo}
  82. s.raftServer, err = raft.NewServer(string(s.serverAddr), s.dataDir, transporter, stateMachine, option.Topo, "")
  83. if err != nil {
  84. glog.V(0).Infoln(err)
  85. return nil, err
  86. }
  87. heartbeatInterval := time.Duration(float64(option.HeartbeatInterval) * (rand.Float64()*0.25 + 1))
  88. s.raftServer.SetHeartbeatInterval(heartbeatInterval)
  89. s.raftServer.SetElectionTimeout(option.ElectionTimeout)
  90. if err := s.raftServer.LoadSnapshot(); err != nil {
  91. return nil, err
  92. }
  93. if err := s.raftServer.Start(); err != nil {
  94. return nil, err
  95. }
  96. for name, peer := range s.peers {
  97. if err := s.raftServer.AddPeer(name, peer.ToGrpcAddress()); err != nil {
  98. return nil, err
  99. }
  100. }
  101. // Remove deleted peers
  102. for existsPeerName := range s.raftServer.Peers() {
  103. if existingPeer, found := s.peers[existsPeerName]; !found {
  104. if err := s.raftServer.RemovePeer(existsPeerName); err != nil {
  105. glog.V(0).Infoln(err)
  106. return nil, err
  107. } else {
  108. glog.V(0).Infof("removing old peer: %s", existingPeer)
  109. }
  110. }
  111. }
  112. s.GrpcServer = raft.NewGrpcServer(s.raftServer)
  113. glog.V(0).Infof("current cluster leader: %v", s.raftServer.Leader())
  114. return s, nil
  115. }
  116. func (s *RaftServer) Peers() (members []string) {
  117. peers := s.raftServer.Peers()
  118. for _, p := range peers {
  119. members = append(members, p.Name)
  120. }
  121. return
  122. }
  123. func (s *RaftServer) DoJoinCommand() {
  124. glog.V(0).Infoln("Initializing new cluster")
  125. if _, err := s.raftServer.Do(&raft.DefaultJoinCommand{
  126. Name: s.raftServer.Name(),
  127. ConnectionString: s.serverAddr.ToGrpcAddress(),
  128. }); err != nil {
  129. glog.Errorf("fail to send join command: %v", err)
  130. }
  131. }