This is an automated email from the ASF dual-hosted git repository.

zenlin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/servicecomb-service-center.git


The following commit(s) were added to refs/heads/master by this push:
     new 5a3cd64  Simplify serf module code and optimize dependencies
5a3cd64 is described below

commit 5a3cd6464a62ab734100b9c7d35169bffa6aff6e
Author: chinxââ <[email protected]>
AuthorDate: Fri Mar 27 09:35:16 2020 +0800

    Simplify serf module code and optimize dependencies
---
 syncer/serf/agent.go      | 204 ------------------------------------
 syncer/serf/agent_test.go |  74 -------------
 syncer/serf/config.go     | 115 --------------------
 syncer/serf/handler.go    | 174 +++++++++++++++++++++++++++++++
 syncer/serf/option.go     |  80 ++++++++++++++
 syncer/serf/query.go      |  52 ++++++++++
 syncer/serf/serf.go       | 259 ++++++++++++++++++++++++++++++++++++++++++++++
 syncer/serf/serf_test.go  | 168 ++++++++++++++++++++++++++++++
 syncer/server/convert.go  |  43 +++++---
 syncer/server/handler.go  |  49 +++------
 syncer/server/server.go   |  52 +++++-----
 11 files changed, 803 insertions(+), 467 deletions(-)

diff --git a/syncer/serf/agent.go b/syncer/serf/agent.go
deleted file mode 100644
index 59a9d0a..0000000
--- a/syncer/serf/agent.go
+++ /dev/null
@@ -1,204 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements.  See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package serf
-
-import (
-       "context"
-       "errors"
-       "time"
-
-       "github.com/apache/servicecomb-service-center/pkg/log"
-       "github.com/hashicorp/serf/cmd/serf/command/agent"
-       "github.com/hashicorp/serf/serf"
-)
-
-// Agent warps the serf agent
-type Agent struct {
-       *agent.Agent
-       conf    *Config
-       readyCh chan struct{}
-       stopCh  chan struct{}
-}
-
-// Create create serf agent with config
-func Create(conf *Config) (*Agent, error) {
-       // config cover to serf config
-       serfConf, err := conf.convertToSerf()
-       if err != nil {
-               return nil, err
-       }
-
-       // create serf agent with serf config
-       serfAgent, err := agent.Create(conf.Config, serfConf, nil)
-       if err != nil {
-               return nil, err
-       }
-       return &Agent{
-               Agent:   serfAgent,
-               conf:    conf,
-               readyCh: make(chan struct{}),
-               stopCh:  make(chan struct{}),
-       }, nil
-}
-
-// Start agent
-func (a *Agent) Start(ctx context.Context) {
-       err := a.Agent.Start()
-       if err == nil {
-               a.RegisterEventHandler(a)
-               err = a.retryJoin(ctx)
-       }
-
-       if err != nil {
-               log.Errorf(err, "start serf agent failed")
-               close(a.stopCh)
-       }
-}
-
-// HandleEvent Handles serf.EventMemberJoin events,
-// which will wait for members to join until the number of group members is 
equal to "groupExpect"
-// when the startup mode is "ModeCluster",
-// used for logical grouping of serf nodes
-func (a *Agent) HandleEvent(event serf.Event) {
-       if event.EventType() != serf.EventMemberJoin {
-               return
-       }
-
-       if a.conf.Mode == ModeCluster {
-               if len(a.GroupMembers(a.conf.ClusterName)) < groupExpect {
-                       return
-               }
-       }
-       a.DeregisterEventHandler(a)
-       close(a.readyCh)
-}
-
-// Ready Returns a channel that will be closed when serf is ready
-func (a *Agent) Ready() <-chan struct{} {
-       return a.readyCh
-}
-
-// Error Returns a channel that will be closed when serf is stopped
-func (a *Agent) Stopped() <-chan struct{} {
-       return a.stopCh
-}
-
-// Stop serf agent
-func (a *Agent) Stop() {
-       a.Leave()
-       a.Shutdown()
-}
-
-// LocalMember returns the Member information for the local node
-func (a *Agent) LocalMember() *serf.Member {
-       serfAgent := a.Agent.Serf()
-       if serfAgent != nil {
-               member := serfAgent.LocalMember()
-               return &member
-       }
-       return nil
-}
-
-// GroupMembers returns a point-in-time snapshot of the members of by groupName
-func (a *Agent) GroupMembers(groupName string) (members []serf.Member) {
-       serfAgent := a.Agent.Serf()
-       if serfAgent != nil {
-               for _, member := range serfAgent.Members() {
-                       log.Debugf("member = %s, groupName = %s", member.Name, 
member.Tags[tagKeyClusterName])
-                       if member.Tags[tagKeyClusterName] == groupName {
-                               members = append(members, member)
-                       }
-               }
-       }
-       return
-}
-
-// Member get member information with node
-func (a *Agent) Member(node string) *serf.Member {
-       serfAgent := a.Agent.Serf()
-       if serfAgent != nil {
-               ms := serfAgent.Members()
-               for _, m := range ms {
-                       if m.Name == node {
-                               return &m
-                       }
-               }
-       }
-       return nil
-}
-
-// SerfConfig get serf config
-func (a *Agent) SerfConfig() *serf.Config {
-       return a.Agent.SerfConfig()
-}
-
-// Join serf clusters through one or more members
-func (a *Agent) Join(addrs []string, replay bool) (n int, err error) {
-       return a.Agent.Join(addrs, replay)
-}
-
-// UserEvent sends a UserEvent on Serf
-func (a *Agent) UserEvent(name string, payload []byte, coalesce bool) error {
-       return a.Agent.UserEvent(name, payload, coalesce)
-}
-
-// Query sends a Query on Serf
-func (a *Agent) Query(name string, payload []byte, params *serf.QueryParam) 
(*serf.QueryResponse, error) {
-       return a.Agent.Query(name, payload, params)
-}
-
-func (a *Agent) retryJoin(ctx context.Context) (err error) {
-       if len(a.conf.RetryJoin) == 0 {
-               log.Infof("retry join mumber %d", len(a.conf.RetryJoin))
-               return nil
-       }
-
-       // Count of attempts
-       attempt := 0
-       ticker := time.NewTicker(a.conf.RetryInterval)
-       for {
-               log.Infof("serf: Joining cluster...(replay: %v)", 
a.conf.ReplayOnJoin)
-               var n int
-
-               // Try to join the specified serf nodes
-               n, err = a.Join(a.conf.RetryJoin, a.conf.ReplayOnJoin)
-               if err == nil {
-                       log.Infof("serf: Join completed. Synced with %d initial 
agents", n)
-                       break
-               }
-               attempt++
-
-               // If RetryMaxAttempts is greater than 0, agent will exit
-               // and throw an error when the number of attempts exceeds 
RetryMaxAttempts,
-               // else agent will try to join other nodes until successful 
always
-               if a.conf.RetryMaxAttempts > 0 && attempt > 
a.conf.RetryMaxAttempts {
-                       err = errors.New("serf: maximum retry join attempts 
made, exiting")
-                       log.Errorf(err, err.Error())
-                       break
-               }
-               select {
-               case <-ctx.Done():
-                       err = ctx.Err()
-                       goto done
-               // Waiting for ticker to trigger
-               case <-ticker.C:
-               }
-       }
-done:
-       ticker.Stop()
-       return
-}
diff --git a/syncer/serf/agent_test.go b/syncer/serf/agent_test.go
deleted file mode 100644
index af68f8f..0000000
--- a/syncer/serf/agent_test.go
+++ /dev/null
@@ -1,74 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements.  See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package serf
-
-import (
-       "context"
-       "testing"
-       "time"
-
-       "github.com/hashicorp/serf/serf"
-)
-
-func TestAgent(t *testing.T) {
-       conf := DefaultConfig()
-       agent, err := Create(conf)
-       if err != nil {
-               t.Errorf("create agent failed, error: %s", err)
-       }
-       agent.Start(context.Background())
-       <- agent.readyCh
-       go func() {
-               agent.ShutdownCh()
-       }()
-       time.Sleep(time.Second)
-
-       err = agent.UserEvent("test", []byte("test"), true)
-       if err != nil {
-               t.Errorf("send user event failed, error: %s", err)
-       }
-
-       _, err = agent.Query("test", []byte("test"), &serf.QueryParam{})
-       if err != nil {
-               t.Errorf("query for other node failed, error: %s", err)
-       }
-       agent.LocalMember()
-
-       agent.Member("testnode")
-
-       agent.SerfConfig()
-
-       _, err = agent.Join([]string{"127.0.0.1:9999"}, true)
-       if err != nil {
-               t.Logf("join to other node failed, error: %s", err)
-       }
-
-       err = agent.Leave()
-       if err != nil {
-               t.Errorf("angent leave failed, error: %s", err)
-       }
-
-       err = agent.ForceLeave("testnode")
-       if err != nil {
-               t.Errorf("angent force leave failed, error: %s", err)
-       }
-
-       err = agent.Shutdown()
-       if err != nil {
-               t.Errorf("angent shutdown failed, error: %s", err)
-       }
-}
diff --git a/syncer/serf/config.go b/syncer/serf/config.go
deleted file mode 100644
index e954862..0000000
--- a/syncer/serf/config.go
+++ /dev/null
@@ -1,115 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements.  See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package serf
-
-import (
-       "fmt"
-       "strconv"
-
-       "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"
-)
-
-const (
-       DefaultBindPort    = 30190
-       DefaultRPCPort     = 30191
-       DefaultClusterPort = 30192
-       ModeSingle         = "single"
-       ModeCluster        = "cluster"
-       retryMaxAttempts   = 3
-       groupExpect        = 3
-       tagKeyClusterName  = "syncer-cluster-name"
-       TagKeyClusterPort  = "syncer-cluster-port"
-       TagKeyRPCPort      = "syncer-rpc-port"
-       TagKeyTLSEnabled   = "syncer-tls-enabled"
-)
-
-// DefaultConfig default config
-func DefaultConfig() *Config {
-       agentConf := agent.DefaultConfig()
-       agentConf.BindAddr = fmt.Sprintf("0.0.0.0:%d", DefaultBindPort)
-       agentConf.RPCAddr = fmt.Sprintf("0.0.0.0:%d", DefaultRPCPort)
-       return &Config{
-               Mode:        ModeSingle,
-               Config:      agentConf,
-               ClusterPort: DefaultClusterPort,
-       }
-}
-
-// Config struct
-type Config struct {
-       // config from serf agent
-       *agent.Config
-       Mode string `json:"mode"`
-
-       // name to group members into cluster
-       ClusterName string `json:"cluster_name"`
-
-       // port to communicate between cluster members
-       ClusterPort int  `yaml:"cluster_port"`
-       RPCPort     int  `yaml:"-"`
-       TLSEnabled  bool `json:"-"`
-}
-
-// readConfigFile reads configuration from config file
-func (c *Config) readConfigFile(filepath string) error {
-       if filepath != "" {
-               // todo:
-       }
-       return nil
-}
-
-// convertToSerf convert Config to serf.Config
-func (c *Config) convertToSerf() (*serf.Config, error) {
-       serfConf := serf.DefaultConfig()
-
-       bindIP, bindPort, err := utils.SplitHostPort(c.BindAddr, 
DefaultBindPort)
-       if err != nil {
-               return nil, fmt.Errorf("invalid bind address: %s", err)
-       }
-
-       switch c.Profile {
-       case "lan":
-               serfConf.MemberlistConfig = memberlist.DefaultLANConfig()
-       case "wan":
-               serfConf.MemberlistConfig = memberlist.DefaultWANConfig()
-       case "local":
-               serfConf.MemberlistConfig = memberlist.DefaultLocalConfig()
-       default:
-               serfConf.MemberlistConfig = memberlist.DefaultLANConfig()
-       }
-
-       serfConf.MemberlistConfig.BindAddr = bindIP
-       serfConf.MemberlistConfig.BindPort = bindPort
-       serfConf.NodeName = c.NodeName
-       serfConf.Tags = map[string]string{
-               TagKeyRPCPort:    strconv.Itoa(c.RPCPort),
-               TagKeyTLSEnabled: strconv.FormatBool(c.TLSEnabled),
-       }
-
-       if c.ClusterName != "" {
-               serfConf.Tags[tagKeyClusterName] = c.ClusterName
-               serfConf.Tags[TagKeyClusterPort] = strconv.Itoa(c.ClusterPort)
-       }
-
-       if c.Mode == ModeCluster && c.RetryMaxAttempts <= 0 {
-               c.RetryMaxAttempts = retryMaxAttempts
-       }
-       return serfConf, nil
-}
diff --git a/syncer/serf/handler.go b/syncer/serf/handler.go
new file mode 100644
index 0000000..a73a6ac
--- /dev/null
+++ b/syncer/serf/handler.go
@@ -0,0 +1,174 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package serf
+
+import (
+       "fmt"
+       "sync"
+
+       "github.com/hashicorp/serf/serf"
+)
+
+// EventHandler interface
+type EventHandler interface {
+       Handle(event serf.Event) bool
+       String() string
+}
+
+type eventHandler struct {
+       filter  EventFilter
+       handler HandleFunc
+}
+
+// NewEventHandler returns serf event handler
+func NewEventHandler(filter EventFilter, handler HandleFunc) EventHandler {
+       return &eventHandler{
+               filter:  filter,
+               handler: handler,
+       }
+}
+
+// Handle invoke event handler
+func (h *eventHandler) Handle(event serf.Event) bool {
+       return h.filter.Invoke(event, h.handler)
+}
+
+// String returns event handler string
+func (h *eventHandler) String() string {
+       return "handler: filter = " + h.filter.String()
+}
+
+type onceHandler struct {
+       once    sync.Once
+       readyCh chan struct{}
+       EventHandler
+}
+
+func onceEventHandler(handler EventHandler) *onceHandler {
+       return &onceHandler{
+               EventHandler: handler,
+               readyCh:      make(chan struct{}),
+       }
+}
+
+// Handle invoke once event handler
+func (h *onceHandler) Handle(event serf.Event) bool {
+       done := h.EventHandler.Handle(event)
+       if done {
+               h.once.Do(func() {
+                       close(h.readyCh)
+               })
+       }
+       return done
+}
+
+// Ready Returns a channel that will be closed when once handler is invoked
+func (h *onceHandler) Ready() <-chan struct{} {
+       return h.readyCh
+}
+
+// EventFilter interface
+type EventFilter interface {
+       Invoke(event serf.Event, handler HandleFunc) bool
+       String() string
+}
+
+// UserEventFilter event filter of serf user event
+func UserEventFilter(name string) EventFilter {
+       return userFilter{name: name}
+}
+
+// QueryFilter event filter of serf query
+func QueryFilter(name string) EventFilter {
+       return queryFilter{name: name}
+}
+
+// MemberJoinFilter event filter of member join
+func MemberJoinFilter() EventFilter {
+       return memberFilter{kind: serf.EventMemberJoin}
+}
+
+// MemberLeaveFilter event filter of member leave
+func MemberLeaveFilter() EventFilter {
+       return memberFilter{kind: serf.EventMemberLeave}
+}
+
+// MemberFailedFilter event filter of member failed
+func MemberFailedFilter() EventFilter {
+       return memberFilter{kind: serf.EventMemberFailed}
+}
+
+// MemberUpdateFilter event filter of member update
+func MemberUpdateFilter() EventFilter {
+       return memberFilter{kind: serf.EventMemberUpdate}
+}
+
+// MemberReapFilter event filter of member reap
+func MemberReapFilter() EventFilter {
+       return memberFilter{kind: serf.EventMemberReap}
+}
+
+type userFilter struct {
+       name string
+}
+
+// Invoke user filter handler
+func (f userFilter) Invoke(event serf.Event, handler HandleFunc) bool {
+       if event.EventType() != serf.EventUser {
+               return false
+       }
+       user, ok := event.(serf.UserEvent)
+       return ok && f.name == user.Name && handler(user.Payload)
+}
+
+// String returns user filter string
+func (f userFilter) String() string {
+       return fmt.Sprintf("event kind = %s, name = %s", 
serf.EventUser.String(), f.name)
+}
+
+type queryFilter struct {
+       name string
+}
+
+// Invoke query filter handler
+func (f queryFilter) Invoke(event serf.Event, handler HandleFunc) bool {
+       if event.EventType() != serf.EventQuery {
+               return false
+       }
+       query, ok := event.(*serf.Query)
+       return ok && f.name == query.Name && handler(query.Payload)
+}
+
+// String returns query filter string
+func (f queryFilter) String() string {
+       return fmt.Sprintf("event kind = %s, name = %s", 
serf.EventQuery.String(), f.name)
+}
+
+type memberFilter struct {
+       kind serf.EventType
+}
+
+// Invoke member filter handler
+func (f memberFilter) Invoke(event serf.Event, handler HandleFunc) bool {
+       return event.EventType() == f.kind && handler()
+}
+
+// String returns member filter string
+func (f memberFilter) String() string {
+       return fmt.Sprintf("event kind = %s", f.kind.String())
+}
diff --git a/syncer/serf/option.go b/syncer/serf/option.go
new file mode 100644
index 0000000..102fa2b
--- /dev/null
+++ b/syncer/serf/option.go
@@ -0,0 +1,80 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package serf
+
+import (
+       "io"
+
+       "github.com/hashicorp/serf/serf"
+)
+
+// Option func
+type Option func(*serf.Config)
+
+// WithNode returns node option
+func WithNode(nodeName string) Option {
+       return func(c *serf.Config) { c.NodeName = nodeName }
+}
+
+// WithTags returns tags option
+func WithTags(tags map[string]string) Option {
+       return func(c *serf.Config) { c.Tags = tags }
+}
+
+// WithAddTag returns add tag option
+func WithAddTag(key, val string) Option {
+       return func(c *serf.Config) { c.Tags[key] = val }
+}
+
+// WithBindAddr returns bind addr option
+func WithBindAddr(addr string) Option {
+       return func(c *serf.Config) { c.MemberlistConfig.BindAddr = addr }
+}
+
+// WithBindPort returns bind port option
+func WithBindPort(port int) Option {
+       return func(c *serf.Config) { c.MemberlistConfig.BindPort = port }
+}
+
+// WithAdvertiseAddr returns advertise addr option
+func WithAdvertiseAddr(addr string) Option {
+       return func(c *serf.Config) { c.MemberlistConfig.AdvertiseAddr = addr }
+}
+
+// WithAdvertisePort returns advertise port option
+func WithAdvertisePort(port int) Option {
+       return func(c *serf.Config) { c.MemberlistConfig.AdvertisePort = port }
+}
+
+// WithEnableCompression returns enable compression option
+func WithEnableCompression(enable bool) Option {
+       return func(c *serf.Config) { c.MemberlistConfig.EnableCompression = 
enable }
+}
+
+// WithSecretKey returns secret key option
+func WithSecretKey(secretKey []byte) Option {
+       return func(c *serf.Config) { c.MemberlistConfig.SecretKey = secretKey }
+}
+
+// WithLogOutput returns log output option
+func WithLogOutput(logOutput io.Writer) Option {
+       return func(c *serf.Config) {
+               c.LogOutput = logOutput
+               c.MemberlistConfig.LogOutput = logOutput
+       }
+}
diff --git a/syncer/serf/query.go b/syncer/serf/query.go
new file mode 100644
index 0000000..da87937
--- /dev/null
+++ b/syncer/serf/query.go
@@ -0,0 +1,52 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package serf
+
+import (
+       "time"
+
+       "github.com/hashicorp/serf/serf"
+)
+
+// QueryOption func
+type QueryOption func(*serf.QueryParam)
+
+// WithFilterNodes returns node filter query option
+func WithFilterNodes(nodes ...string) QueryOption {
+       return func(p *serf.QueryParam) { p.FilterNodes = nodes }
+}
+
+// WithFilterTags returns tags filter query option
+func WithFilterTags(tags map[string]string) QueryOption {
+       return func(p *serf.QueryParam) { p.FilterTags = tags }
+}
+
+// WithRequestAck returns request ack query option
+func WithRequestAck(ack bool) QueryOption {
+       return func(p *serf.QueryParam) { p.RequestAck = ack }
+}
+
+// WithRelayFactor returns relay factor query option
+func WithRelayFactor(num uint8) QueryOption {
+       return func(p *serf.QueryParam) { p.RelayFactor = num }
+}
+
+// WithTimeout returns timeout query option
+func WithTimeout(timeout time.Duration) QueryOption {
+       return func(p *serf.QueryParam) { p.Timeout = timeout }
+}
diff --git a/syncer/serf/serf.go b/syncer/serf/serf.go
new file mode 100644
index 0000000..ee8a407
--- /dev/null
+++ b/syncer/serf/serf.go
@@ -0,0 +1,259 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package serf
+
+import (
+       "context"
+       "sync"
+       "time"
+
+       "github.com/apache/servicecomb-service-center/pkg/log"
+       "github.com/apache/servicecomb-service-center/syncer/pkg/utils"
+       "github.com/hashicorp/serf/serf"
+       "github.com/pkg/errors"
+)
+
+// HandleFunc handle user event
+type HandleFunc func(data ...[]byte) bool
+
+//// QueryHandler handle query event
+//type QueryHandler func(data []byte) []byte
+
+// CallbackFunc callback handler for query event
+type CallbackFunc func(from string, data []byte)
+
+// Server serf server
+type Server struct {
+       conf    *serf.Config
+       serf    *serf.Serf
+       running *utils.AtomicBool
+
+       eventCh chan serf.Event
+       readyCh chan struct{}
+       stopCh  chan struct{}
+
+       handlerMap *sync.Map
+}
+
+// NewServer new serf server with options
+func NewServer(opts ...Option) *Server {
+       conf := serf.DefaultConfig()
+       conf.Tags = map[string]string{}
+       for _, opt := range opts {
+               opt(conf)
+       }
+
+       eventCh := make(chan serf.Event, 64)
+       conf.EventCh = eventCh
+
+       return &Server{
+               conf:       conf,
+               running:    utils.NewAtomicBool(false),
+               eventCh:    eventCh,
+               readyCh:    make(chan struct{}),
+               stopCh:     make(chan struct{}),
+               handlerMap: &sync.Map{},
+       }
+}
+
+// Start serf server
+func (s *Server) Start(ctx context.Context) {
+       s.running.DoToReverse(false, func() {
+               sf, err := serf.Create(s.conf)
+               if err != nil {
+                       log.Error("serf: start server failed", err)
+                       close(s.stopCh)
+                       return
+               }
+               s.serf = sf
+               close(s.readyCh)
+               go s.waitEvent(ctx)
+       })
+}
+
+// Stop serf server
+func (s *Server) Stop() {
+       s.running.DoToReverse(true, func() {
+               if s.serf != nil {
+                       log.Info("serf: begin shutdown")
+                       if err := s.serf.Shutdown(); err != nil {
+                               log.Error("serf: shutdown failed", err)
+                       }
+                       close(s.stopCh)
+               }
+
+               log.Info("serf: shutdown complete")
+       })
+}
+
+// Ready Returns a channel that will be closed when serf is ready
+func (s *Server) Ready() <-chan struct{} {
+       return s.readyCh
+}
+
+// Stopped Returns a channel that will be closed when serf is stopped
+func (s *Server) Stopped() <-chan struct{} {
+       return s.stopCh
+}
+
+// OnceEventHandler Add serf event handler, the handler will be automatically 
deleted after executing it once
+func (s *Server) OnceEventHandler(handler EventHandler) {
+       onceHandler := onceEventHandler(handler)
+       s.AddEventHandler(onceHandler)
+       go s.waitOnceEventHandler(onceHandler)
+}
+
+// AddEventHandler Add serf event handler
+func (s *Server) AddEventHandler(handler EventHandler) {
+       _, ok := s.handlerMap.Load(handler)
+       if ok {
+               log.Warn("serf: event handle is already exits, " + 
handler.String())
+       }
+       s.handlerMap.Store(handler, struct{}{})
+}
+
+// RemoveEventHandler remove serf event handler
+func (s *Server) RemoveEventHandler(handler EventHandler) {
+       _, ok := s.handlerMap.Load(handler)
+       if !ok {
+               log.Warn("serf: event handle is notfound, " + handler.String())
+               return
+       }
+       s.handlerMap.Delete(handler)
+}
+
+func (s *Server) waitOnceEventHandler(once *onceHandler) {
+       <-once.Ready()
+       s.RemoveEventHandler(once)
+}
+
+// UserEvent send user event
+func (s *Server) UserEvent(name string, payload []byte) error {
+       err := s.serf.UserEvent(name, payload, true)
+       if err != nil {
+               err = errors.Wrapf(err, "serf: send user event '%s' failed", 
name)
+       }
+       return err
+}
+
+// Query send query
+func (s *Server) Query(name string, payload []byte, callback CallbackFunc, 
opts ...QueryOption) error {
+       param := s.serf.DefaultQueryParams()
+       for _, opt := range opts {
+               opt(param)
+       }
+
+       resp, err := s.serf.Query(name, payload, param)
+       if err != nil {
+               err = errors.Wrapf(err, "serf: send query '%s' failed", name)
+               return err
+       }
+       go s.responseCallback(resp, callback)
+       return nil
+}
+
+// Join asks the Serf instance to join. See the Serf.Join function.
+func (s *Server) Join(addrs []string) (n int, err error) {
+       log.Infof("serf: join to: %v replay : %v", addrs)
+       n, err = s.serf.Join(addrs, true)
+       if n > 0 {
+               log.Infof("serf: joined: %d nodes", n)
+       }
+       if err != nil {
+               log.Warnf("serf: error joining: %v", err)
+       }
+       return
+}
+
+// MembersByTags Returns members matching the tags
+func (s *Server) MembersByTags(tags map[string]string) (members []serf.Member) 
{
+       if s.serf == nil {
+               return
+       }
+
+next:
+       for _, member := range s.serf.Members() {
+               for key, val := range tags {
+                       if member.Tags[key] != val {
+                               continue next
+                       }
+               }
+               members = append(members, member)
+       }
+       return
+}
+
+// LocalMember returns the Member information for the local node
+func (s *Server) LocalMember() *serf.Member {
+       if s.serf != nil {
+               member := s.serf.LocalMember()
+               return &member
+       }
+       return nil
+}
+
+// Member get member information with node
+func (s *Server) Member(node string) *serf.Member {
+       if s.serf != nil {
+               ms := s.serf.Members()
+               for _, m := range ms {
+                       if m.Name == node {
+                               return &m
+                       }
+               }
+       }
+       return nil
+}
+
+func (s *Server) responseCallback(resp *serf.QueryResponse, callback 
CallbackFunc) {
+       hourglass := time.After(resp.Deadline().Sub(time.Now()))
+       for {
+               select {
+               case a := <-resp.AckCh():
+                       log.Infof("query response ack: %s", a)
+               case r := <-resp.ResponseCh():
+                       log.Infof("query response: from %s, content %s", 
r.From, string(r.Payload))
+                       callback(r.From, r.Payload)
+               case <-hourglass:
+                       log.Info("query response timeout")
+                       return
+               }
+       }
+}
+
+func (s *Server) waitEvent(ctx context.Context) {
+       for {
+               select {
+               case e := <-s.eventCh:
+                       s.handlerMap.Range(func(key, value interface{}) bool {
+                               if handler, ok := key.(EventHandler); ok {
+                                       handler.Handle(e)
+                               }
+                               return true
+                       })
+               case <-s.serf.ShutdownCh():
+                       log.Warn("serf: server stopped, exited")
+                       s.Stop()
+                       return
+               case <-ctx.Done():
+                       log.Warn("serf: cancel server by context")
+                       s.Stop()
+                       return
+               }
+       }
+}
diff --git a/syncer/serf/serf_test.go b/syncer/serf/serf_test.go
new file mode 100644
index 0000000..b350222
--- /dev/null
+++ b/syncer/serf/serf_test.go
@@ -0,0 +1,168 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package serf
+
+import (
+       "context"
+       "errors"
+       "os"
+       "testing"
+       "time"
+
+       "github.com/hashicorp/serf/serf"
+       "github.com/stretchr/testify/assert"
+)
+
+func TestSerfServer(t *testing.T) {
+       svr := defaultServer()
+
+       ctx, cancel := context.WithCancel(context.Background())
+       err := startServer(ctx, svr)
+       assert.Nil(t, err)
+
+       err = svr.UserEvent("test-event", []byte("test-data"))
+       assert.Nil(t, err)
+
+       _, err = svr.Join([]string{"127.0.0.1:35151"})
+       assert.Nil(t, err)
+
+       list := svr.MembersByTags(map[string]string{"test-key": "test-value"})
+       assert.Equal(t, 0, len(list))
+
+       self := svr.LocalMember()
+       assert.NotNil(t, self)
+
+       m := svr.Member("syncer-test")
+       assert.NotNil(t, m)
+
+       cancel()
+       svr.Stop()
+}
+
+func TestServerFailed(t *testing.T) {
+       svr := NewServer(
+               WithNode("syncer-test"),
+               WithTags(map[string]string{"test-key": "test-value"}),
+               WithAddTag("added-key", "added-value"),
+               WithBindAddr("127.0.0.1"),
+               WithBindPort(35151),
+               WithAdvertiseAddr(""),
+               WithAdvertisePort(0),
+               WithEnableCompression(true),
+               WithSecretKey([]byte("123456")),
+               WithLogOutput(os.Stdout),
+       )
+       err := startServer(context.Background(), svr)
+       assert.NotNil(t, err)
+
+       svr.Stop()
+}
+
+func TestServerEventHandler(t *testing.T) {
+       svr := defaultServer()
+       startServer(context.Background(), svr)
+       svr.OnceEventHandler(NewEventHandler(MemberJoinFilter(), func(data 
...[]byte) bool {
+               t.Log("Once event form member join triggered")
+               return true
+       }))
+
+       // wait for trigger
+       <-time.After(time.Second)
+
+       handler := NewEventHandler(MemberJoinFilter(), func(data ...[]byte) 
bool {
+               return false
+       })
+
+       svr.RemoveEventHandler(handler)
+
+       svr.AddEventHandler(handler)
+
+       svr.AddEventHandler(handler)
+
+       svr.RemoveEventHandler(handler)
+
+       svr.Stop()
+}
+
+func TestUserQuery(t *testing.T) {
+       svr := defaultServer()
+       startServer(context.Background(), svr)
+       err := svr.Query("test-query", []byte("test-data"), func(from string, 
data []byte) {},
+               WithFilterNodes("syncer-test"),
+               WithFilterTags(map[string]string{"test-key": "test-value"}),
+               WithRequestAck(false),
+               WithRelayFactor(0),
+               WithTimeout(time.Second),
+       )
+       assert.Nil(t, err)
+
+       svr.Stop()
+}
+
+func TestEventHandler(t *testing.T) {
+       filter := UserEventFilter("test-event")
+       ok := filter.Invoke(serf.UserEvent{Name: "test-event", Payload: 
[]byte("test-data")}, func(data ...[]byte) bool {
+               return true
+       })
+       assert.True(t, ok)
+
+       ok = filter.Invoke(&serf.Query{Name: "test-event", Payload: 
[]byte("test-data")}, func(data ...[]byte) bool {
+               return true
+       })
+       assert.False(t, ok)
+
+       t.Log(filter.String())
+
+       filter = QueryFilter("test-event")
+       ok = filter.Invoke(&serf.Query{Name: "test-event", Payload: 
[]byte("test-data")}, func(data ...[]byte) bool {
+               return true
+       })
+       assert.True(t, ok)
+
+       ok = filter.Invoke(serf.UserEvent{Name: "test-event", Payload: 
[]byte("test-data")}, func(data ...[]byte) bool {
+               return true
+       })
+
+       assert.False(t, ok)
+       t.Log(filter.String())
+
+       filter = MemberLeaveFilter()
+       filter = MemberFailedFilter()
+       filter = MemberUpdateFilter()
+       filter = MemberReapFilter()
+}
+
+func defaultServer() *Server {
+       return NewServer(
+               WithNode("syncer-test"),
+               WithBindAddr("127.0.0.1"),
+               WithBindPort(35151),
+       )
+}
+
+func startServer(ctx context.Context, svr *Server) (err error) {
+       svr.Start(ctx)
+       select {
+       case <-svr.Ready():
+       case <-svr.Stopped():
+               err = errors.New("start serf server failed")
+       case <-time.After(time.Second * 3):
+               err = errors.New("start serf server timeout")
+       }
+       return
+}
diff --git a/syncer/server/convert.go b/syncer/server/convert.go
index d9d2ccc..0ac5c8c 100644
--- a/syncer/server/convert.go
+++ b/syncer/server/convert.go
@@ -19,8 +19,8 @@ package server
 
 import (
        "crypto/tls"
+       "strconv"
        "strings"
-       "time"
 
        "github.com/apache/servicecomb-service-center/pkg/tlsutil"
        "github.com/apache/servicecomb-service-center/syncer/config"
@@ -31,21 +31,34 @@ import (
        "github.com/apache/servicecomb-service-center/syncer/task"
 )
 
-func convertSerfConfig(c *config.Config) *serf.Config {
-       conf := serf.DefaultConfig()
-       conf.NodeName = c.Node
-       conf.ClusterName = c.Cluster
-       conf.Mode = c.Mode
-       conf.TLSEnabled = c.Listener.TLSMount.Enabled
-       conf.BindAddr = c.Listener.BindAddr
-       _, conf.ClusterPort, _ = utils.ResolveAddr(c.Listener.PeerAddr)
-       _, conf.RPCPort, _ = utils.ResolveAddr(c.Listener.RPCAddr)
-       if c.Join.Enabled {
-               conf.RetryJoin = strings.Split(c.Join.Address, ",")
-               conf.RetryInterval, _ = time.ParseDuration(c.Join.RetryInterval)
-               conf.RetryMaxAttempts = c.Join.RetryMax
+const (
+       tagKeyClusterName = "syncer-cluster-name"
+       tagKeyClusterPort = "syncer-cluster-port"
+       tagKeyRPCPort     = "syncer-rpc-port"
+       tagKeyTLSEnabled  = "syncer-tls-enabled"
+
+       groupExpect = 3
+)
+
+func convertSerfOptions(c *config.Config) []serf.Option {
+       bindHost, bindPort, _ := utils.ResolveAddr(c.Listener.BindAddr)
+       _, rpcPort, _ := utils.ResolveAddr(c.Listener.RPCAddr)
+       opts := []serf.Option{
+               serf.WithNode(c.Node),
+               serf.WithBindAddr(bindHost),
+               serf.WithBindPort(bindPort),
+               serf.WithAddTag(tagKeyRPCPort, strconv.Itoa(rpcPort)),
+               serf.WithAddTag(tagKeyTLSEnabled, 
strconv.FormatBool(c.Listener.TLSMount.Enabled)),
        }
-       return conf
+
+       if c.Cluster != "" {
+               _, peerPort, _ := utils.ResolveAddr(c.Listener.PeerAddr)
+               opts = append(opts,
+                       serf.WithAddTag(tagKeyClusterName, c.Cluster),
+                       serf.WithAddTag(tagKeyClusterPort, 
strconv.Itoa(peerPort)),
+               )
+       }
+       return opts
 }
 
 func convertEtcdOptions(c *config.Config) []etcd.Option {
diff --git a/syncer/server/handler.go b/syncer/server/handler.go
index e7c349e..f00d8ba 100644
--- a/syncer/server/handler.go
+++ b/syncer/server/handler.go
@@ -14,6 +14,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package server
 
 import (
@@ -27,8 +28,6 @@ import (
        "github.com/apache/servicecomb-service-center/pkg/util"
        "github.com/apache/servicecomb-service-center/syncer/grpc"
        pb "github.com/apache/servicecomb-service-center/syncer/proto"
-       myserf "github.com/apache/servicecomb-service-center/syncer/serf"
-       "github.com/hashicorp/serf/serf"
 )
 
 const (
@@ -36,7 +35,7 @@ const (
 )
 
 // tickHandler Timed task handler
-func (s *Server) tickHandler(ctx context.Context) {
+func (s *Server) tickHandler() {
        log.Debugf("is leader: %v", s.etcd.IsLeader())
        if !s.etcd.IsLeader() {
                return
@@ -46,57 +45,41 @@ func (s *Server) tickHandler(ctx context.Context) {
        s.servicecenter.FlushData()
 
        // sends a UserEvent on Serf, the event will be broadcast between 
members
-       err := s.agent.UserEvent(EventDiscovered, 
util.StringToBytesWithNoCopy(s.conf.Cluster), true)
+       err := s.serf.UserEvent(EventDiscovered, 
util.StringToBytesWithNoCopy(s.conf.Cluster))
        if err != nil {
                log.Errorf(err, "Syncer send user event failed")
        }
 }
 
-// GetData Sync Data to GRPC
+// Discovery discovery sync data from servicecenter
 func (s *Server) Discovery() *pb.SyncData {
        return s.servicecenter.Discovery()
 }
 
-// HandleEvent Handles serf.EventUser/serf.EventQuery,
-// used for message passing and processing between serf nodes
-func (s *Server) HandleEvent(event serf.Event) {
-       log.Debugf("is leader: %v", s.etcd.IsLeader())
-       if !s.etcd.IsLeader() {
-               return
-       }
-       switch event.EventType() {
-       case serf.EventUser:
-               s.userEvent(event.(serf.UserEvent))
-       case serf.EventQuery:
-               s.queryEvent(event.(*serf.Query))
-       default:
-               log.Infof("serf event = %s", event)
-       }
-}
-
 // userEvent Handles "EventUser" notification events, no response required
-func (s *Server) userEvent(event serf.UserEvent) {
+func (s *Server) userEvent(data ...[]byte) (success bool) {
        log.Debug("Receive serf user event")
-       clusterName := util.BytesToStringWithNoCopy(event.Payload)
+       clusterName := util.BytesToStringWithNoCopy(data[0])
 
        // Excludes notifications from self, as the gossip protocol inevitably 
has redundant notifications
        if s.conf.Cluster == clusterName {
                return
        }
 
+       tags := map[string]string{tagKeyClusterName: clusterName}
        // Get member information and get synchronized data from it
-       members := s.agent.GroupMembers(clusterName)
-       if members == nil || len(members) == 0 {
+       members := s.serf.MembersByTags(tags)
+       if len(members) == 0 {
                log.Warnf("serf member = %s is not found", clusterName)
                return
        }
 
        // todo: grpc supports multi-address polling
        // Get dta from remote member
-       endpoint := fmt.Sprintf("%s:%s", members[0].Addr, 
members[0].Tags[myserf.TagKeyRPCPort])
+       endpoint := fmt.Sprintf("%s:%s", members[0].Addr, 
members[0].Tags[tagKeyRPCPort])
        log.Debugf("Going to pull data from %s %s", members[0].Name, endpoint)
 
-       enabled, err := 
strconv.ParseBool(members[0].Tags[myserf.TagKeyTLSEnabled])
+       enabled, err := strconv.ParseBool(members[0].Tags[tagKeyTLSEnabled])
        if err != nil {
                log.Warnf("get tls enabled failed, err = %s", err)
        }
@@ -111,16 +94,12 @@ func (s *Server) userEvent(event serf.UserEvent) {
                }
        }
 
-       data, err := grpc.Pull(context.Background(), endpoint, tlsConfig)
+       syncData, err := grpc.Pull(context.Background(), endpoint, tlsConfig)
        if err != nil {
                log.Errorf(err, "Pull other serf instances failed, node name is 
'%s'", members[0].Name)
                return
        }
        // Registry instances to servicecenter and update storage of it
-       s.servicecenter.Registry(clusterName, data)
-}
-
-// queryEvent Handles "EventQuery" query events and respond if conditions are 
met
-func (s *Server) queryEvent(query *serf.Query) {
-       // todo: Get instances requested
+       s.servicecenter.Registry(clusterName, syncData)
+       return true
 }
diff --git a/syncer/server/server.go b/syncer/server/server.go
index 5e81b19..df0e23f 100644
--- a/syncer/server/server.go
+++ b/syncer/server/server.go
@@ -76,7 +76,7 @@ type Server struct {
        etcd *etcd.Server
 
        // Wraps the serf agent
-       agent *serf.Agent
+       serf *serf.Server
 
        // Wraps the grpc server
        grpc *grpc.Server
@@ -107,12 +107,7 @@ func (s *Server) Run(ctx context.Context) {
        // Start system signal listening, wait for user interrupt program
        gopool.Go(syssig.Run)
 
-       err = s.startModuleServer(s.agent)
-       if err != nil {
-               return
-       }
-
-       err = s.configureCluster()
+       err = s.startModuleServer(s.serf)
        if err != nil {
                return
        }
@@ -129,11 +124,7 @@ func (s *Server) Run(ctx context.Context) {
 
        s.servicecenter.SetStorageEngine(s.etcd.Storage())
 
-       s.agent.RegisterEventHandler(s)
-
-       s.task.Handle(func() {
-               s.tickHandler(ctx)
-       })
+       s.task.Handle(s.tickHandler)
 
        s.task.Run(ctx)
 
@@ -147,11 +138,9 @@ func (s *Server) Run(ctx context.Context) {
 
 // Stop Syncer Server
 func (s *Server) Stop() {
-       if s.agent != nil {
-               // removes the serf eventHandler
-               s.agent.DeregisterEventHandler(s)
+       if s.serf != nil {
                //stop serf agent
-               s.agent.Stop()
+               s.serf.Stop()
        }
 
        if s.grpc != nil {
@@ -191,11 +180,8 @@ func (s *Server) initialization() (err error) {
                return
        }
 
-       s.agent, err = serf.Create(convertSerfConfig(s.conf))
-       if err != nil {
-               log.Errorf(err, "Create serf failed, %s", err)
-               return
-       }
+       s.serf = serf.NewServer(convertSerfOptions(s.conf)...)
+       s.serf.OnceEventHandler(serf.NewEventHandler(serf.MemberJoinFilter(), 
s.waitClusterMembers))
 
        s.etcd, err = etcd.NewServer(convertEtcdOptions(s.conf)...)
        if err != nil {
@@ -236,16 +222,34 @@ func (s *Server) initPlugin() {
        plugins.LoadPlugins()
 }
 
+func (s *Server) waitClusterMembers(data ...[]byte) bool {
+       if s.conf.Mode == config.ModeCluster {
+               tags := map[string]string{tagKeyClusterName: s.conf.Cluster}
+               if len(s.serf.MembersByTags(tags)) < groupExpect {
+                       return false
+               }
+               err := s.configureCluster()
+               if err != nil {
+                       log.Error("configure cluster failed", err)
+                       s.Stop()
+                       return false
+               }
+       }
+       
s.serf.AddEventHandler(serf.NewEventHandler(serf.UserEventFilter(EventDiscovered),
 s.userEvent))
+       return true
+}
+
 // configureCluster Configuring the cluster by serf group member information
 func (s *Server) configureCluster() error {
        // get local member of serf
-       self := s.agent.LocalMember()
+       self := s.serf.LocalMember()
        _, peerPort, _ := utils.SplitAddress(s.conf.Listener.PeerAddr)
        ops := []etcd.Option{etcd.WithPeerAddr(self.Addr.String() + ":" + 
strconv.Itoa(peerPort))}
 
        // group members from serf as initial cluster members
-       for _, member := range s.agent.GroupMembers(s.conf.Cluster) {
-               ops = append(ops, etcd.WithAddPeers(member.Name, 
member.Addr.String()+":"+member.Tags[serf.TagKeyClusterPort]))
+       tags := map[string]string{tagKeyClusterName: s.conf.Cluster}
+       for _, member := range s.serf.MembersByTags(tags) {
+               ops = append(ops, etcd.WithAddPeers(member.Name, 
member.Addr.String()+":"+member.Tags[tagKeyClusterPort]))
        }
 
        return s.etcd.AddOptions(ops...)

Reply via email to