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...)