基于serf的数据库同步模块库
zhangzengfei
2023-05-15 2863a050be2530afc452e48aae8b4be9b3965ebd
config.go
@@ -17,6 +17,7 @@
package syncdb
import (
   "context"
   "fmt"
   "net"
   "os"
@@ -46,7 +47,7 @@
   BroadcastInterval  = 5
   MaxQueryRespSize   = 50 * 1024 * 1024
   MaxQuerySize       = 50 * 1024 * 1024
   MaxUserEventSize   = 9 * 1024
   MaxUserEventSize   = 9 * 1024 * 10
   ReplayOnJoinDefault = false
   SnapshotPathDefault = "./serfSnapShot"
   MaxEventBufferCount = 2048
@@ -64,6 +65,32 @@
      Mode:        ModeSingle,
      Config:      agentConf,
      ClusterPort: DefaultClusterPort,
      Ctx:        context.Background(),
   }
}
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
      }
   }
}
@@ -79,6 +106,7 @@
   // port to communicate between cluster members
   ClusterPort int       `yaml:"cluster_port"`
   RPCPort     int       `yaml:"-"`
   Ctx       context.Context
}
// readConfigFile reads configuration from config file
@@ -135,7 +163,7 @@
   serfConf.MemberlistConfig.BindPort = bindPort
   serfConf.NodeName = c.NodeName
   serfConf.Tags = map[string]string{TagKeyRPCPort: strconv.Itoa(c.RPCPort)}
   serfConf.Tags = map[string]string{TagKeyRPCPort: strconv.Itoa(c.RPCPort), "role": "slave"}
   if c.ClusterID != "" {
      serfConf.Tags[tagKeyClusterID] = c.ClusterID