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.

155 lines
3.7 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 map[string]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. s := &RaftServer{
  55. peers: option.Peers,
  56. serverAddr: option.ServerAddr,
  57. dataDir: option.DataDir,
  58. topo: option.Topo,
  59. }
  60. if glog.V(4) {
  61. raft.SetLogLevel(2)
  62. }
  63. raft.RegisterCommand(&topology.MaxVolumeIdCommand{})
  64. var err error
  65. transporter := raft.NewGrpcTransporter(option.GrpcDialOption)
  66. glog.V(0).Infof("Starting RaftServer with %v", option.ServerAddr)
  67. // always clear previous log to avoid server is promotable
  68. os.RemoveAll(path.Join(s.dataDir, "log"))
  69. if !option.RaftResumeState {
  70. // always clear previous metadata
  71. os.RemoveAll(path.Join(s.dataDir, "conf"))
  72. os.RemoveAll(path.Join(s.dataDir, "snapshot"))
  73. }
  74. if err := os.MkdirAll(path.Join(s.dataDir, "snapshot"), 0600); err != nil {
  75. return nil, err
  76. }
  77. stateMachine := StateMachine{topo: option.Topo}
  78. s.raftServer, err = raft.NewServer(string(s.serverAddr), s.dataDir, transporter, stateMachine, option.Topo, "")
  79. if err != nil {
  80. glog.V(0).Infoln(err)
  81. return nil, err
  82. }
  83. heartbeatInterval := time.Duration(float64(option.HeartbeatInterval) * (rand.Float64()*0.25 + 1))
  84. s.raftServer.SetHeartbeatInterval(heartbeatInterval)
  85. s.raftServer.SetElectionTimeout(option.ElectionTimeout)
  86. if err := s.raftServer.LoadSnapshot(); err != nil {
  87. return nil, err
  88. }
  89. if err := s.raftServer.Start(); err != nil {
  90. return nil, err
  91. }
  92. for name, peer := range s.peers {
  93. if err := s.raftServer.AddPeer(name, peer.ToGrpcAddress()); err != nil {
  94. return nil, err
  95. }
  96. }
  97. // Remove deleted peers
  98. for existsPeerName := range s.raftServer.Peers() {
  99. if existingPeer, found := s.peers[existsPeerName]; !found {
  100. if err := s.raftServer.RemovePeer(existsPeerName); err != nil {
  101. glog.V(0).Infoln(err)
  102. return nil, err
  103. } else {
  104. glog.V(0).Infof("removing old peer: %s", existingPeer)
  105. }
  106. }
  107. }
  108. s.GrpcServer = raft.NewGrpcServer(s.raftServer)
  109. glog.V(0).Infof("current cluster leader: %v", s.raftServer.Leader())
  110. return s, nil
  111. }
  112. func (s *RaftServer) Peers() (members []string) {
  113. peers := s.raftServer.Peers()
  114. for _, p := range peers {
  115. members = append(members, p.Name)
  116. }
  117. return
  118. }
  119. func (s *RaftServer) DoJoinCommand() {
  120. glog.V(0).Infoln("Initializing new cluster")
  121. if _, err := s.raftServer.Do(&raft.DefaultJoinCommand{
  122. Name: s.raftServer.Name(),
  123. ConnectionString: s.serverAddr.ToGrpcAddress(),
  124. }); err != nil {
  125. glog.Errorf("fail to send join command: %v", err)
  126. }
  127. }