基于serf的数据库同步模块库
liuxiaolong
2021-06-07 aea48291949e54554f5c484cca32a1d3119d8e76
config.go
@@ -17,14 +17,17 @@
package syncdb
import (
   "context"
   "fmt"
   "net"
   "os"
   "strconv"
   "strings"
   //"github.com/apache/servicecomb-service-center/syncer/pkg/utils"
   "github.com/hashicorp/memberlist"
   "github.com/hashicorp/serf/cmd/serf/command/agent"
   "github.com/hashicorp/serf/serf"
   "basic.com/valib/serf.git/cmd/serf/command/agent"
   "basic.com/valib/serf.git/serf"
)
const (
@@ -46,10 +49,8 @@
   MaxQuerySize       = 50 * 1024 * 1024
   MaxUserEventSize   = 9 * 1024
   ReplayOnJoinDefault = false
   SnapshotPathDefault = "/opt/vasystem/serfSnapShot"
   SnapshotPathDefault = "./serfSnapShot"
   MaxEventBufferCount = 2048
   TcpTransportPort = 30194 //tcp传输大数据量接口
)
// DefaultConfig default config
@@ -67,18 +68,44 @@
   }
}
func (c *Config) MergeConf(s *Config) {
   if s != nil {
      if s.Ctx != nil {
         c.Ctx = s.Ctx
      } else {
         c.Ctx = context.Background()
      }
      c.BindAddr = s.BindAddr
      c.RPCAddr = s.RPCAddr
      c.RPCPort = s.RPCPort
      //serf快照地址
      if s.SnapshotPath != "" {
         c.SnapshotPath = s.SnapshotPath
      }
      if s.EncryptKey != "" {
         //报文加密的key
         c.EncryptKey = s.EncryptKey
      }
      if s.RPCAuthKey != "" {
         //RPC认证的key
         c.RPCAuthKey = s.RPCAuthKey
      }
   }
}
// Config struct
type Config struct {
   // config from serf agent
   *agent.Config
   Mode string `json:"mode"`
   Mode       string       `json:"mode"`
   // name to group members into cluster
   ClusterID string `json:"cluster_name"`
   ClusterID    string       `json:"cluster_name"`
   // port to communicate between cluster members
   ClusterPort int `yaml:"cluster_port"`
   RPCPort     int `yaml:"-"`
   ClusterPort int       `yaml:"cluster_port"`
   RPCPort     int       `yaml:"-"`
   Ctx       context.Context
}
// readConfigFile reads configuration from config file
@@ -89,8 +116,30 @@
   return nil
}
func isFileRightful(filePath string) bool {
   if filePath != "" {
      _, err := os.Stat(filePath)
      if err != nil && os.IsNotExist(err) {
         pos := strings.LastIndex(filePath, "/")
         if pos != -1 {
            filePath = filePath[0:pos]
         }
         _, err = os.Stat(filePath)
         if err == nil || !os.IsNotExist(err) {
            return true
         } else {
            return false
         }
      } else {
         return false
      }
   }
   return false
}
// convertToSerf convert Config to serf.Config
func (c *Config) convertToSerf() (*serf.Config, error) {
func (c *Config) convertToSerf(snapshotPath string) (*serf.Config, error) {
   serfConf := serf.DefaultConfig()
   bindIP, bindPort, err := SplitHostPort(c.BindAddr, DefaultBindPort)
@@ -123,7 +172,12 @@
   if c.Mode == ModeCluster && c.RetryMaxAttempts <= 0 {
      c.RetryMaxAttempts = retryMaxAttempts
   }
   c.SnapshotPath = SnapshotPathDefault
   if isFileRightful(snapshotPath) {
      c.SnapshotPath = snapshotPath
   }
   c.ReplayOnJoin = ReplayOnJoinDefault
   serfConf.QueryResponseSizeLimit = c.QueryResponseSizeLimit