Merge pull request #1133 from codysoyland/redundant-gossip-seeds

Add support for lists of gossip seeds for redundancy
This commit is contained in:
Cody Soyland 2018-02-23 12:03:16 -06:00 • committed by GitHub
commit 66cc7a35c0
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
8 changed files with 108 additions and 53 deletions

View file

@ -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"`
@ -194,7 +190,7 @@ func NewConfig() *Config {
// Gossip config.
// c.Gossip.Port = ""
// c.Gossip.Seed = ""
// c.Gossip.Seeds = []string{}
// c.Gossip.Key = ""
c.Gossip.StreamTimeout = Duration(DefaultGossipStreamTimeout)
c.Gossip.SuspicionMult = DefaultGossipSuspicionMult

View file

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

View file

@ -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")
}

View file

@ -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 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`.
* 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:

View file

@ -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,13 +92,18 @@ 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)
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)
}
}
// 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())
@ -142,7 +147,7 @@ func (g *GossipMemberSet) logger() *log.Logger {
////////////////////////////////////////////////////////////////
type gossipConfig struct {
gossipSeed string
gossipSeeds []string
memberlistConfig *memberlist.Config
}
@ -191,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

View file

@ -50,7 +50,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}}
@ -78,7 +78,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()
@ -225,7 +225,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)
}
@ -234,7 +234,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)
}
@ -253,7 +253,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)
}
@ -264,7 +264,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
}
@ -287,7 +287,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)
}
@ -308,7 +308,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
}
@ -333,7 +333,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)
}
@ -363,7 +363,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
}
@ -388,7 +388,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)
}
@ -418,7 +418,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
}
@ -439,6 +439,67 @@ 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, []string{}, 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, []string{"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, []string{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)
}
})
}
func TestClusterResize_RemoveNode(t *testing.T) {
cluster := test.MustRunMainWithCluster(t, 3)
m0 := cluster[0]

View file

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

View file

@ -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()