From fc6b0ea0e89ee2b79ddbb9bcbe59148f0af9c3b5 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 21 Feb 2018 15:10:52 -0600 Subject: [PATCH 1/3] Add support for comma-separated list of gossip seeds for redundancy. --- docs/configuration.md | 2 +- gossip/gossip.go | 16 +++++++---- server/cluster_test.go | 61 ++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 73 insertions(+), 6 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index 0798d1381..6c20e6ec5 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -111,7 +111,7 @@ Any flag that has a value that is a comma separated list on the command line bec #### Gossip Seed -* Description: When using the gossip [Cluster Type]({{< ref "#cluster-type" >}}), this specifies which internal host should be used to initialize membership in the cluster. Typcially this can be the address of any available host in the cluster. For example, when starting a three-node cluster made up of `node0`, `node1`, and `node2`, the `gossip-seed` for all three nodes can be configured to be the address of `node0`. +* Description: When using the gossip [Cluster Type]({{< ref "#cluster-type" >}}), this specifies which internal host(s) should be used to initialize membership in the cluster. Typcially this can be the address of any available host in the cluster. You may enter multiple seeds by separating them with a comma. For example, when starting a three-node cluster made up of `node0`, `node1`, and `node2`, the `gossip-seed` for all three nodes can be configured to be the address of `node0`. * Flag: `--gossip.seed="localhost:11101"` * Env: `PILOSA_GOSSIP_SEED="localhost:11101"` * Config: diff --git a/gossip/gossip.go b/gossip/gossip.go index acfad66a1..a61a5b9b7 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -92,13 +92,19 @@ func (g *GossipMemberSet) Open(n *pilosa.Node) error { RetransmitMult: 3, } - uri, err := pilosa.NewURIFromAddress(g.config.gossipSeed) - if err != nil { - return fmt.Errorf("new uri from address: %s", err) + parts := strings.Split(g.config.gossipSeed, ",") + var uris = make([]*pilosa.URI, len(parts)) + for i, addr := range parts { + uris[i], err = pilosa.NewURIFromAddress(addr) + if err != nil { + return fmt.Errorf("new uri from address: %s", err) + } } - // attach to gossip seed node - nodes := []*pilosa.Node{&pilosa.Node{URI: *uri}} //TODO: support a list of seeds + var nodes = make([]*pilosa.Node, len(uris)) + for i, uri := range uris { + nodes[i] = &pilosa.Node{URI: *uri} + } g.mu.RLock() err = g.joinWithRetry(pilosa.URIs(pilosa.Nodes(nodes).URIs()).HostPortStrings()) diff --git a/server/cluster_test.go b/server/cluster_test.go index 23284bb00..82c54bb96 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -435,3 +435,64 @@ func TestClusterResize_AddNode(t *testing.T) { } }) } + +// Ensure that redundant gossip seeds are used +func TestCluster_GossipMembership(t *testing.T) { + t.Run("Node0Down", func(t *testing.T) { + // Configure node0 + m0 := test.NewMainWithCluster() + defer m0.Close() + + seed, coord, err := m0.RunWithTransport("localhost", 0, "", pilosa.URI{}) + if err != nil { + t.Fatal(err) + } + + // Configure node1 + m1 := test.NewMainWithCluster() + defer m1.Close() + + var eg errgroup.Group + eg.Go(func() error { + // Pass invalid seed as first in list + _, _, err = m1.RunWithTransport("localhost", 0, "http://localhost:8765,"+seed, coord) + if err != nil { + return err + } + return nil + }) + + // Configure node2 + m2 := test.NewMainWithCluster() + defer m2.Close() + + eg.Go(func() error { + // Pass invalid seed as last in list + _, _, err = m2.RunWithTransport("localhost", 0, seed+",http://localhost:8765", coord) + if err != nil { + return err + } + return nil + }) + + if err := eg.Wait(); err != nil { + t.Fatal(err) + } + + // Give the cluster time to settle. + time.Sleep(1 * time.Second) + + if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) + } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + } else if m2.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node2 cluster state: %s", m2.Server.Cluster.State()) + } + + numNodes := len(m0.Server.Cluster.Status().Nodes) + if numNodes != 3 { + t.Fatalf("Expected 3 nodes, got %d", numNodes) + } + }) +} From a23fe868391779123b9d5e2d7f61e16a138842bc Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Thu, 22 Feb 2018 11:40:28 -0600 Subject: [PATCH 2/3] Convert gossip seed string to slice --- config.go | 9 +-------- ctl/server.go | 4 +--- ctl/server_test.go | 3 --- docs/configuration.md | 12 ++++++------ gossip/gossip.go | 21 ++++++++++----------- server/cluster_test.go | 30 +++++++++++++++--------------- server/server.go | 3 --- test/pilosa.go | 15 ++++++++------- 8 files changed, 41 insertions(+), 56 deletions(-) diff --git a/config.go b/config.go index 05f132db2..234334f56 100644 --- a/config.go +++ b/config.go @@ -127,10 +127,6 @@ type TLSConfig struct { type Config struct { DataDir string `toml:"data-dir"` Bind string `toml:"bind"` - // GossipPort DEPRECATED - GossipPort string `toml:"gossip-port"` - // GossipSeed DEPRECATED - GossipSeed string `toml:"gossip-seed"` // Limits the number of mutating commands that can be in a single request to // the server. This includes SetBit, ClearBit, SetRowAttrs & SetColumnAttrs. @@ -151,7 +147,7 @@ type Config struct { Gossip struct { Port string `toml:"port"` - Seed string `toml:"seed"` + Seeds []string `toml:"seeds"` Key string `toml:"key"` StreamTimeout Duration `toml:"stream-timeout"` SuspicionMult int `toml:"suspicion-mult"` @@ -193,9 +189,6 @@ func NewConfig() *Config { c.Cluster.LongQueryTime = Duration(time.Minute) // Gossip config. - // c.Gossip.Port = "" - // c.Gossip.Seed = "" - // c.Gossip.Key = "" c.Gossip.StreamTimeout = Duration(DefaultGossipStreamTimeout) c.Gossip.SuspicionMult = DefaultGossipSuspicionMult c.Gossip.PushPullInterval = Duration(DefaultGossipPushPullInterval) diff --git a/ctl/server.go b/ctl/server.go index 3c5b1a058..18f4f6f93 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -26,8 +26,6 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags := cmd.Flags() flags.StringVarP(&srv.Config.DataDir, "data-dir", "d", srv.Config.DataDir, "Directory to store pilosa data files.") flags.StringVarP(&srv.Config.Bind, "bind", "b", srv.Config.Bind, "Default URI on which pilosa should listen.") - flags.StringVarP(&srv.Config.GossipPort, "gossip-port", "", "", "(DEPRECATED) Port to which pilosa should bind for internal state sharing.") - flags.StringVarP(&srv.Config.GossipSeed, "gossip-seed", "", "", "(DEPRECATED) Host with which to seed the gossip membership.") flags.IntVarP(&srv.Config.MaxWritesPerRequest, "max-writes-per-request", "", srv.Config.MaxWritesPerRequest, "Number of write commands per request.") flags.StringVar(&srv.Config.LogPath, "log-path", srv.Config.LogPath, "Log path") @@ -43,7 +41,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { // Gossip flags.StringVarP(&srv.Config.Gossip.Port, "gossip.port", "", srv.Config.Gossip.Port, "Port to which pilosa should bind for internal state sharing.") - flags.StringVarP(&srv.Config.Gossip.Seed, "gossip.seed", "", srv.Config.Gossip.Seed, "Host with which to seed the gossip membership.") + flags.StringSliceVarP(&srv.Config.Gossip.Seeds, "gossip.seeds", "", srv.Config.Gossip.Seeds, "Host with which to seed the gossip membership.") flags.StringVarP(&srv.Config.Gossip.Key, "gossip.key", "", srv.Config.Gossip.Key, "The path to file of the encryption key for gossip. The contents of the file should be either 16, 24, or 32 bytes to select AES-128, AES-192, or AES-256.") flags.DurationVarP((*time.Duration)(&srv.Config.Gossip.StreamTimeout), "gossip.stream-timeout", "", (time.Duration)(srv.Config.Gossip.StreamTimeout), "Timeout for establishing a stream connection with a remote node for a full state sync.") flags.IntVarP(&srv.Config.Gossip.SuspicionMult, "gossip.suspicion-mult", "", srv.Config.Gossip.SuspicionMult, "Multiplier for determining the time an inaccessible node is considered suspect before declaring it dead.") diff --git a/ctl/server_test.go b/ctl/server_test.go index dc31426ac..9866b497d 100644 --- a/ctl/server_test.go +++ b/ctl/server_test.go @@ -28,9 +28,6 @@ func TestBuildServerFlags(t *testing.T) { stdin, stdout, stderr := GetIO(buf) Server := server.NewCommand(stdin, stdout, stderr) BuildServerFlags(cm, Server) - if cm.Flags().Lookup("gossip-port").Name == "" { - t.Fatal("gossip-port flag is required") - } if cm.Flags().Lookup("data-dir").Name == "" { t.Fatal("data-dir flag is required") } diff --git a/docs/configuration.md b/docs/configuration.md index 6c20e6ec5..dfaa01d64 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -109,16 +109,16 @@ Any flag that has a value that is a comma separated list on the command line bec port = 11101 ``` -#### Gossip Seed +#### Gossip Seeds -* Description: When using the gossip [Cluster Type]({{< ref "#cluster-type" >}}), this specifies which internal host(s) should be used to initialize membership in the cluster. Typcially this can be the address of any available host in the cluster. You may enter multiple seeds by separating them with a comma. For example, when starting a three-node cluster made up of `node0`, `node1`, and `node2`, the `gossip-seed` for all three nodes can be configured to be the address of `node0`. -* Flag: `--gossip.seed="localhost:11101"` -* Env: `PILOSA_GOSSIP_SEED="localhost:11101"` +* Description: This specifies which internal host(s) should be used to initialize membership in the cluster. Typcially this can be the address of any available host in the cluster. For example, when starting a three-node cluster made up of `node0`, `node1`, and `node2`, the `gossip.seeds` for all three nodes can be configured to be the address of `node0`. Multiple seeds should be comma-separated in the flag and env forms. +* Flag: `--gossip.seeds="localhost:11101"` +* Env: `PILOSA_GOSSIP_SEEDS="localhost:11101"` * Config: ```toml [gossip] - seed = "localhost:11101" + seeds = ["localhost:11101"] ``` #### Gossip Key @@ -134,7 +134,7 @@ Any flag that has a value that is a comma separated list on the command line bec #### Cluster Hosts -* Description: List of hosts in the cluster. Multiple hosts should be comma separated in the flag and env forms. +* Description: List of hosts in the cluster. Multiple hosts should be comma-separated in the flag and env forms. * Flag: `--cluster.hosts="localhost:10101"` * Env: `PILOSA_CLUSTER_HOSTS="localhost:10101"` * Config: diff --git a/gossip/gossip.go b/gossip/gossip.go index a61a5b9b7..4b79c41ac 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -62,9 +62,9 @@ func (g *GossipMemberSet) Start(h pilosa.BroadcastHandler) error { return nil } -// Seed returns the gossipSeed determined by the config. -func (g *GossipMemberSet) Seed() string { - return g.config.gossipSeed +// Seeds returns the gossipSeeds determined by the config. +func (g *GossipMemberSet) Seeds() []string { + return g.config.gossipSeeds } // Open implements the MemberSet interface to start network activity. @@ -92,9 +92,8 @@ func (g *GossipMemberSet) Open(n *pilosa.Node) error { RetransmitMult: 3, } - parts := strings.Split(g.config.gossipSeed, ",") - var uris = make([]*pilosa.URI, len(parts)) - for i, addr := range parts { + var uris = make([]*pilosa.URI, len(g.config.gossipSeeds)) + for i, addr := range g.config.gossipSeeds { uris[i], err = pilosa.NewURIFromAddress(addr) if err != nil { return fmt.Errorf("new uri from address: %s", err) @@ -148,7 +147,7 @@ func (g *GossipMemberSet) logger() *log.Logger { //////////////////////////////////////////////////////////////// type gossipConfig struct { - gossipSeed string + gossipSeeds []string memberlistConfig *memberlist.Config } @@ -197,14 +196,14 @@ func NewGossipMemberSetWithTransport(name string, cfg *pilosa.Config, transport g.config = &gossipConfig{ memberlistConfig: conf, - gossipSeed: cfg.Gossip.Seed, + gossipSeeds: cfg.Gossip.Seeds, } g.statusHandler = server - // If no gossipSeed is provided, use local host:port. - if cfg.Gossip.Seed == "" { - g.config.gossipSeed = fmt.Sprintf("%s:%d", host, port) + // If no gossipSeeds is provided, use local host:port. + if len(cfg.Gossip.Seeds) == 0 { + g.config.gossipSeeds = []string{fmt.Sprintf("%s:%d", host, port)} } return g, nil diff --git a/server/cluster_test.go b/server/cluster_test.go index 82c54bb96..bfc3640b9 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -47,7 +47,7 @@ func TestMain_SendReceiveMessage(t *testing.T) { // get the host portion of addr to use for binding m0.Config.Gossip.Port = "0" - m0.Config.Gossip.Seed = "" + m0.Config.Gossip.Seeds = []string{} m0.Server.Cluster.Coordinator = m0.Server.URI m0.Server.Cluster.Topology = &pilosa.Topology{NodeIDs: []string{m0.Server.NodeID, m1.Server.NodeID}} @@ -75,7 +75,7 @@ func TestMain_SendReceiveMessage(t *testing.T) { // get the host portion of addr to use for binding m1.Config.Gossip.Port = "0" - m1.Config.Gossip.Seed = gossipMemberSet0.Seed() + m1.Config.Gossip.Seeds = gossipMemberSet0.Seeds() m1.Server.Cluster.Coordinator = m0.Server.URI m1.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver() @@ -222,7 +222,7 @@ func TestClusterResize_EmptyNodes(t *testing.T) { gossipHost := "localhost" gossipPort := 0 - seed, coord, err := m0.RunWithTransport(gossipHost, gossipPort, "", pilosa.URI{}) + seed, coord, err := m0.RunWithTransport(gossipHost, gossipPort, []string{}, pilosa.URI{}) if err != nil { t.Fatal(err) } @@ -231,7 +231,7 @@ func TestClusterResize_EmptyNodes(t *testing.T) { m1 := test.NewMainWithCluster() defer m1.Close() - seed, coord, err = m1.RunWithTransport(gossipHost, gossipPort, seed, coord) + seed, coord, err = m1.RunWithTransport(gossipHost, gossipPort, []string{seed}, coord) if err != nil { t.Fatal(err) } @@ -250,7 +250,7 @@ func TestClusterResize_AddNode(t *testing.T) { m0 := test.NewMainWithCluster() defer m0.Close() - seed, coord, err := m0.RunWithTransport("localhost", 0, "", pilosa.URI{}) + seed, coord, err := m0.RunWithTransport("localhost", 0, []string{}, pilosa.URI{}) if err != nil { t.Fatal(err) } @@ -261,7 +261,7 @@ func TestClusterResize_AddNode(t *testing.T) { var eg errgroup.Group eg.Go(func() error { - _, _, err = m1.RunWithTransport("localhost", 0, seed, coord) + _, _, err = m1.RunWithTransport("localhost", 0, []string{seed}, coord) if err != nil { return err } @@ -284,7 +284,7 @@ func TestClusterResize_AddNode(t *testing.T) { m0 := test.NewMainWithCluster() defer m0.Close() - seed, coord, err := m0.RunWithTransport("localhost", 0, "", pilosa.URI{}) + seed, coord, err := m0.RunWithTransport("localhost", 0, []string{}, pilosa.URI{}) if err != nil { t.Fatal(err) } @@ -305,7 +305,7 @@ func TestClusterResize_AddNode(t *testing.T) { var eg errgroup.Group eg.Go(func() error { - _, _, err = m1.RunWithTransport("localhost", 0, seed, coord) + _, _, err = m1.RunWithTransport("localhost", 0, []string{seed}, coord) if err != nil { return err } @@ -330,7 +330,7 @@ func TestClusterResize_AddNode(t *testing.T) { m0 := test.NewMainWithCluster() defer m0.Close() - seed, coord, err := m0.RunWithTransport("localhost", 0, "", pilosa.URI{}) + seed, coord, err := m0.RunWithTransport("localhost", 0, []string{}, pilosa.URI{}) if err != nil { t.Fatal(err) } @@ -360,7 +360,7 @@ func TestClusterResize_AddNode(t *testing.T) { var eg errgroup.Group eg.Go(func() error { - _, _, err = m1.RunWithTransport("localhost", 0, seed, coord) + _, _, err = m1.RunWithTransport("localhost", 0, []string{seed}, coord) if err != nil { return err } @@ -385,7 +385,7 @@ func TestClusterResize_AddNode(t *testing.T) { m0 := test.NewMainWithCluster() defer m0.Close() - seed, coord, err := m0.RunWithTransport("localhost", 0, "", pilosa.URI{}) + seed, coord, err := m0.RunWithTransport("localhost", 0, []string{}, pilosa.URI{}) if err != nil { t.Fatal(err) } @@ -415,7 +415,7 @@ func TestClusterResize_AddNode(t *testing.T) { var eg errgroup.Group eg.Go(func() error { - _, _, err = m1.RunWithTransport("localhost", 0, seed, coord) + _, _, err = m1.RunWithTransport("localhost", 0, []string{seed}, coord) if err != nil { return err } @@ -443,7 +443,7 @@ func TestCluster_GossipMembership(t *testing.T) { m0 := test.NewMainWithCluster() defer m0.Close() - seed, coord, err := m0.RunWithTransport("localhost", 0, "", pilosa.URI{}) + seed, coord, err := m0.RunWithTransport("localhost", 0, []string{}, pilosa.URI{}) if err != nil { t.Fatal(err) } @@ -455,7 +455,7 @@ func TestCluster_GossipMembership(t *testing.T) { var eg errgroup.Group eg.Go(func() error { // Pass invalid seed as first in list - _, _, err = m1.RunWithTransport("localhost", 0, "http://localhost:8765,"+seed, coord) + _, _, err = m1.RunWithTransport("localhost", 0, []string{"http://localhost:8765", seed}, coord) if err != nil { return err } @@ -468,7 +468,7 @@ func TestCluster_GossipMembership(t *testing.T) { eg.Go(func() error { // Pass invalid seed as last in list - _, _, err = m2.RunWithTransport("localhost", 0, seed+",http://localhost:8765", coord) + _, _, err = m2.RunWithTransport("localhost", 0, []string{seed, "http://localhost:8765"}, coord) if err != nil { return err } diff --git a/server/server.go b/server/server.go index 19a08f982..7b623461a 100644 --- a/server/server.go +++ b/server/server.go @@ -237,11 +237,8 @@ func (m *Command) SetupNetworking() error { // Set internal port (string). gossipPortStr := pilosa.DefaultGossipPort - // Config.GossipPort is deprecated, so Config.Gossip.Port has priority if m.Config.Gossip.Port != "" { gossipPortStr = m.Config.Gossip.Port - } else if m.Config.GossipPort != "" { - gossipPortStr = m.Config.GossipPort } gossipPort, err := strconv.Atoi(gossipPortStr) diff --git a/test/pilosa.go b/test/pilosa.go index 3d741af76..de03b1c28 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -93,13 +93,13 @@ func runMainWithCluster(size int) ([]*Main, error) { gossipHost := "localhost" gossipPort := 0 var err error - var gossipSeed string + var gossipSeeds = make([]string, size) var coordinator pilosa.URI for i := 0; i < size; i++ { m := NewMainWithCluster() - gossipSeed, coordinator, err = m.RunWithTransport(gossipHost, gossipPort, gossipSeed, coordinator) + gossipSeeds[i], coordinator, err = m.RunWithTransport(gossipHost, gossipPort, gossipSeeds[:i], coordinator) if err != nil { return nil, errors.Wrap(err, "RunWithTransport") } @@ -146,7 +146,7 @@ func (m *Main) Reopen() error { } // RunWithTransport runs Main and returns the dynamically allocated gossip port. -func (m *Main) RunWithTransport(host string, bindPort int, joinSeed string, coordinator pilosa.URI) (seed string, coord pilosa.URI, err error) { +func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string, coordinator pilosa.URI) (seed string, coord pilosa.URI, err error) { defer close(m.Started) /* @@ -182,12 +182,13 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeed string, coor } m.GossipTransport = transport - if joinSeed != "" { - m.Config.Gossip.Seed = joinSeed + if len(joinSeeds) != 0 { + m.Config.Gossip.Seeds = joinSeeds } else { - m.Config.Gossip.Seed = transport.URI.String() + m.Config.Gossip.Seeds = []string{transport.URI.String()} } - seed = m.Config.Gossip.Seed + + seed = transport.URI.String() // SetupNetworking err = m.SetupNetworking() From 3bba35236f9037f71f7edd56d57bd447737d904f Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 23 Feb 2018 11:11:45 -0600 Subject: [PATCH 3/3] Add back example settings --- config.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/config.go b/config.go index 234334f56..4fa06fc51 100644 --- a/config.go +++ b/config.go @@ -189,6 +189,9 @@ func NewConfig() *Config { c.Cluster.LongQueryTime = Duration(time.Minute) // Gossip config. + // c.Gossip.Port = "" + // c.Gossip.Seeds = []string{} + // c.Gossip.Key = "" c.Gossip.StreamTimeout = Duration(DefaultGossipStreamTimeout) c.Gossip.SuspicionMult = DefaultGossipSuspicionMult c.Gossip.PushPullInterval = Duration(DefaultGossipPushPullInterval)