Resume test: TestMain_SendReceiveMessage

Create a custom memberlist NetTransport (which will bind to an available
port when port = 0 in the configuration). This allows us to bind
to dynamic ports in tests while at the same time determining a valid
seed for the cluster.
This commit is contained in:
Travis Turner 2017-12-04 14:56:21 -06:00
parent 29aca37d83
commit 9b54259cd7
No known key found for this signature in database
GPG key ID: 7F08008DFD9314C9
3 changed files with 110 additions and 61 deletions

View file

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

View file

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

View file

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