diff --git a/gossip/gossip.go b/gossip/gossip.go index 7085710b3..e413410d8 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -18,6 +18,8 @@ import ( "fmt" "io" "log" + "os" + "strings" "time" "golang.org/x/sync/errgroup" @@ -59,6 +61,11 @@ func (g *GossipNodeSet) Start(h pilosa.BroadcastHandler) error { return nil } +// Seed returns the gossipSeed determined by the config. +func (g *GossipNodeSet) Seed() string { + return g.config.gossipSeed +} + // Open implements the NodeSet interface to start network activity. func (g *GossipNodeSet) Open() error { if g.handler == nil { @@ -122,28 +129,109 @@ type gossipConfig struct { memberlistConfig *memberlist.Config } +// newTransport returns a NetTransport based on the memberlist configuration. +// It will dynamically bind to a port if conf.BindPort is 0. +// This is useful for test cases where specifiying a port is not reasonable. +func newTransport(conf *memberlist.Config) (*memberlist.NetTransport, error) { + if conf.LogOutput != nil && conf.Logger != nil { + return nil, fmt.Errorf("Cannot specify both LogOutput and Logger. Please choose a single log configuration setting.") + } + + logDest := conf.LogOutput + if logDest == nil { + logDest = os.Stderr + } + + logger := conf.Logger + if logger == nil { + logger = log.New(logDest, "", log.LstdFlags) + } + + nc := &memberlist.NetTransportConfig{ + BindAddrs: []string{conf.BindAddr}, + BindPort: conf.BindPort, + Logger: logger, + } + + // See comment below for details about the retry in here. + makeNetRetry := func(limit int) (*memberlist.NetTransport, error) { + var err error + for try := 0; try < limit; try++ { + var nt *memberlist.NetTransport + if nt, err = memberlist.NewNetTransport(nc); err == nil { + return nt, nil + } + if strings.Contains(err.Error(), "address already in use") { + logger.Printf("[DEBUG] Got bind error: %v", err) + continue + } + } + + return nil, fmt.Errorf("failed to obtain an address: %v", err) + } + + // The dynamic bind port operation is inherently racy because + // even though we are using the kernel to find a port for us, we + // are attempting to bind multiple protocols (and potentially + // multiple addresses) with the same port number. We build in a + // few retries here since this often gets transient errors in + // busy unit tests. + limit := 1 + if conf.BindPort == 0 { + limit = 10 + } + + nt, err := makeNetRetry(limit) + if err != nil { + return nil, fmt.Errorf("Could not set up network transport: %v", err) + } + if conf.BindPort == 0 { + port := nt.GetAutoBindPort() + conf.BindPort = port + conf.AdvertisePort = port + logger.Printf("[DEBUG] Using dynamic bind port %d", port) + } + + return nt, nil +} + // NewGossipNodeSet returns a new instance of GossipNodeSet. -func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed string, server *pilosa.Server, secretKey []byte) *GossipNodeSet { +func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed string, server *pilosa.Server, secretKey []byte) (*GossipNodeSet, error) { g := &GossipNodeSet{ LogOutput: server.LogOutput, } + conf := memberlist.DefaultLocalConfig() + conf.BindPort = gossipPort + conf.AdvertisePort = gossipPort + //TODO: pull memberlist config from pilosa.cfg file g.config = &gossipConfig{ - memberlistConfig: memberlist.DefaultLocalConfig(), + memberlistConfig: conf, gossipSeed: gossipSeed, } + g.config.memberlistConfig.Name = name g.config.memberlistConfig.BindAddr = gossipHost - g.config.memberlistConfig.BindPort = gossipPort g.config.memberlistConfig.AdvertiseAddr = pilosa.HostToIP(gossipHost) - g.config.memberlistConfig.AdvertisePort = gossipPort g.config.memberlistConfig.Delegate = g g.config.memberlistConfig.SecretKey = secretKey g.statusHandler = server - return g + // set up the transport + transport, err := newTransport(g.config.memberlistConfig) + if err != nil { + return nil, err + } + g.config.memberlistConfig.Transport = transport + + // If no gossipSeed is provided, use local host:port. + if gossipSeed == "" { + g.config.gossipSeed = fmt.Sprintf("%s:%d", gossipHost, g.config.memberlistConfig.BindPort) + } + + return g, nil } // SendSync implementation of the Broadcaster interface. diff --git a/server/server.go b/server/server.go index a5f17f6d7..93b0f184f 100644 --- a/server/server.go +++ b/server/server.go @@ -221,7 +221,10 @@ func (m *Command) SetupServer() error { // get the host portion of addr to use for binding gossipHost := uri.Host() - gossipNodeSet := gossip.NewGossipNodeSet(uri.HostPort(), gossipHost, gossipPort, gossipSeed, m.Server, gossipKey) + gossipNodeSet, err := gossip.NewGossipNodeSet(uri.HostPort(), gossipHost, gossipPort, gossipSeed, m.Server, gossipKey) + if err != nil { + return err + } m.Server.Cluster.NodeSet = gossipNodeSet m.Server.Broadcaster = gossipNodeSet m.Server.BroadcastReceiver = gossipNodeSet diff --git a/server/server_test.go b/server/server_test.go index 28886bc67..ca8ac7f09 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -22,19 +22,19 @@ import ( "io" "io/ioutil" "math/rand" - "net" "net/http" "os" "reflect" "runtime" "sort" - "strconv" "strings" "testing" "testing/quick" + "time" "github.com/BurntSushi/toml" "github.com/pilosa/pilosa" + "github.com/pilosa/pilosa/gossip" "github.com/pilosa/pilosa/server" "github.com/pilosa/pilosa/test" ) @@ -384,21 +384,15 @@ func TestCountOpenFiles(t *testing.T) { } } -/* TODO: Fix this test. See #951. // Ensure program can send/receive broadcast messages. func TestMain_SendReceiveMessage(t *testing.T) { + m0 := MustRunMain() defer m0.Close() m1 := MustRunMain() defer m1.Close() - // Get available ports for internal messaging - freePorts, err := availablePorts(2) - if err != nil { - t.Fatal(err) - } - // Update cluster config m0.Server.Cluster.Nodes = []*pilosa.Node{ {Host: m0.Server.URI.HostPort()}, @@ -409,17 +403,14 @@ func TestMain_SendReceiveMessage(t *testing.T) { // Configure node0 // get the host portion of addr to use for binding - gossipHost, _, err := net.SplitHostPort(m0.Server.URI.HostPort()) - if err != nil { - gossipHost = m0.Server.URI.HostPort() - } - gossipPort, err := strconv.Atoi(freePorts[0]) + gossipHost := m0.Server.URI.Host() + gossipPort := 0 + gossipSeed := "" + + gossipNodeSet0, err := gossip.NewGossipNodeSet(m0.Server.URI.HostPort(), gossipHost, gossipPort, gossipSeed, m0.Server, nil) if err != nil { t.Fatal(err) } - gossipSeed := gossipHost + ":" + freePorts[0] - - gossipNodeSet0 := gossip.NewGossipNodeSet(m0.Server.URI.HostPort(), gossipHost, gossipPort, gossipSeed, m0.Server, nil) m0.Server.Cluster.NodeSet = gossipNodeSet0 m0.Server.Broadcaster = gossipNodeSet0 m0.Server.Handler.Broadcaster = m0.Server.Broadcaster @@ -437,16 +428,14 @@ func TestMain_SendReceiveMessage(t *testing.T) { // Configure node1 // get the host portion of addr to use for binding - gossipHost, _, err = net.SplitHostPort(m1.Server.URI.HostPort()) - if err != nil { - gossipHost = m1.Server.URI.HostPort() - } - gossipPort, err = strconv.Atoi(freePorts[1]) + gossipHost = m1.Server.URI.Host() + gossipPort = 0 + gossipSeed = gossipNodeSet0.Seed() + + gossipNodeSet1, err := gossip.NewGossipNodeSet(m1.Server.URI.HostPort(), gossipHost, gossipPort, gossipSeed, m1.Server, nil) if err != nil { t.Fatal(err) } - - gossipNodeSet1 := gossip.NewGossipNodeSet(m1.Server.URI.HostPort(), gossipHost, gossipPort, gossipSeed, m1.Server, nil) m1.Server.Cluster.NodeSet = gossipNodeSet1 m1.Server.Broadcaster = gossipNodeSet1 m1.Server.Handler.Broadcaster = m1.Server.Broadcaster @@ -566,37 +555,6 @@ func TestMain_SendReceiveMessage(t *testing.T) { t.Fatal("frame not found") } } -*/ - -// availablePorts returns a slice of ports that can be used for testing. -func availablePorts(cnt int) ([]string, error) { - rtn := []string{} - - for i := 0; i < cnt; i++ { - port, err := getPort() - if err != nil { - return nil, err - } - rtn = append(rtn, strconv.Itoa(port)) - } - return rtn, nil -} - -// Ask the kernel for a free open port that is ready to use -func getPort() (int, error) { - addr, err := net.ResolveTCPAddr("tcp", "localhost:0") - if err != nil { - return 0, err - } - - l, err := net.ListenTCP("tcp", addr) - if err != nil { - return 0, err - } - defer l.Close() - - return l.Addr().(*net.TCPAddr).Port, nil -} // Main represents a test wrapper for main.Main. type Main struct {