Merge pull request #1432 from seebs/discoPorts

possible fix for CI deadlocks after port_mapper messages in CI
This commit is contained in:
seebs 2021-02-24 12:03:14 -06:00 committed by GitHub
commit bffdee1e8e
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
17 changed files with 272 additions and 275 deletions

View file

@ -27,7 +27,6 @@ import (
"github.com/davecgh/go-spew/spew"
pnet "github.com/pilosa/pilosa/v2/net"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/test/port"
"github.com/pilosa/pilosa/v2/testhook"
"github.com/pilosa/pilosa/v2/topology"
)
@ -498,13 +497,9 @@ func TestCluster_ContainsShards(t *testing.T) {
func TestCluster_Nodes(t *testing.T) {
const urisCount = 4
var uris []pnet.URI
if err := port.GetPorts(func(ports []int) error {
for i := 0; i < urisCount; i++ {
uris = append(uris, NewTestURIFromHostPort(fmt.Sprintf("node%d", i), uint16(ports[i])))
}
return nil
}, urisCount, 10); err != nil {
t.Fatalf("getting ports: %v", err)
arbitraryPorts := []int{17384, 17385, 17386, 17387}
for i := 0; i < urisCount; i++ {
uris = append(uris, NewTestURIFromHostPort(fmt.Sprintf("node%d", i), uint16(arbitraryPorts[i])))
}
node0 := &topology.Node{ID: "node0", URI: uris[0]}

View file

@ -26,8 +26,10 @@ import (
"strings"
"time"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/disco"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/testhook"
"github.com/pilosa/pilosa/v2/topology"
"github.com/pkg/errors"
"go.etcd.io/etcd/clientv3"
@ -100,6 +102,7 @@ func NewEtcd(opt Options, replicas int) *Etcd {
// Close implements io.Closer
func (e *Etcd) Close() error {
_ = testhook.Closed(pilosa.NewAuditor(), e, nil)
if e.e != nil {
if e.resizeCancel != nil {
e.resizeCancel()
@ -153,10 +156,11 @@ func parseOptions(opt Options) *embed.Config {
if opt.ClusterURL != "" {
cfg.ClusterState = embed.ClusterStateFlagExisting
cli, err := clientv3.NewFromURL(opt.ClusterURL)
t, err := clientv3.NewFromURL(opt.ClusterURL)
if err != nil {
panic(err)
}
cli := &hookedClient{Client: t}
defer cli.Close()
log.Println("Cluster Members:")
@ -183,6 +187,7 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) {
if err != nil {
return state, errors.Wrap(err, "starting etcd")
}
_ = testhook.Opened(pilosa.NewAuditor(), e, nil)
e.e = etcd
select {
@ -205,12 +210,11 @@ func (e *Etcd) startHeartbeat() error {
}
defer cli.Close()
heartbeatID, heartbeatFunc, err := e.leaseKeepAlive(e.options.HeartbeatTTL)
heartbeatID, ctx, heartbeatCancel, err := e.leaseKeepAlive(context.Background(), e.options.HeartbeatTTL)
if err != nil {
return errors.Wrap(err, "startHeartbeat: creates a new hearbeat")
}
ctx, heartbeatCancel := context.WithCancel(context.Background())
key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.ClusterStateStarting
if e.e.Config().ClusterState == embed.ClusterStateFlagExisting {
value = disco.ClusterStateResizing
@ -222,7 +226,6 @@ func (e *Etcd) startHeartbeat() error {
}
e.heartbeatID, e.heartbeatCancel = heartbeatID, heartbeatCancel
go heartbeatFunc(ctx, time.Second)
return nil
}
@ -237,7 +240,7 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e
return e.nodeState(ctx, cli, peerID)
}
func (e *Etcd) nodeState(ctx context.Context, cli *clientv3.Client, peerID string) (disco.NodeState, error) {
func (e *Etcd) nodeState(ctx context.Context, cli *hookedClient, peerID string) (disco.NodeState, error) {
resp, err := cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly())
if err != nil {
return disco.NodeStateUnknown, err
@ -387,12 +390,11 @@ func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) {
}
defer cli.Close()
resizeID, resizeFunc, err := e.leaseKeepAlive(e.options.HeartbeatTTL)
resizeID, ctx, resizeCancel, err := e.leaseKeepAlive(ctx, e.options.HeartbeatTTL)
if err != nil {
return nil, errors.Wrap(err, "Resize: creates a new hearbeat")
}
ctx, resizeCancel := context.WithCancel(ctx)
// Check if key exists - maybe we are still resizing
key := path.Join(resizePrefix, e.e.Server.ID().String())
txnResp, err := cli.Txn(ctx).
@ -410,7 +412,6 @@ func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) {
}
e.resizeCancel = resizeCancel
go resizeFunc(ctx, time.Second)
return func(value []byte) error {
log.Println("Update progress:", key, string(value))
@ -811,37 +812,45 @@ func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) error {
return err
}
func (e *Etcd) leaseKeepAlive(ttl int64) (clientv3.LeaseID, func(context.Context, time.Duration), error) {
// leaseKeepAlive creates a lease with the given ttl (treated as a time.Duration),
// then refreshes it periodically, and cancels it when done. it yields the lease ID,
// and also a context and cancelfunc that can be used to abort the heartbeat.
func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID, context.Context, context.CancelFunc, error) {
cli, err := e.client()
if err != nil {
return 0, nil, errors.Wrap(err, "leaseKeepAlive: creates a new client")
return 0, nil, nil, errors.Wrap(err, "leaseKeepAlive: creates a new client")
}
defer cli.Close()
leaseResp, err := cli.Grant(context.TODO(), ttl)
ctx, cancelFunc := context.WithCancel(ctx)
leaseResp, err := cli.Grant(ctx, ttl)
if err != nil {
return 0, nil, errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %v)", ttl)
cancelFunc()
return 0, nil, nil, errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %v)", ttl)
}
keepaliveFunc := func(ctx context.Context, tick time.Duration) {
keepaliveFunc := func(tick time.Duration) {
ticker := time.NewTicker(tick)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
log.Printf("leaseKeepAlive: %v\n", ctx.Err())
// Because of the load balancer, this can take ridiculously
// long times to run if the cluster's already down when we get
// here, resulting in massive piles of excess goroutines.
revoker, cancel := context.WithTimeout(context.Background(), time.Duration(ttl))
defer cancel()
if cli, err := e.client(); err != nil {
log.Printf("leaseKeepAlive: creates a new client: %v\n", err)
} else {
if _, err := cli.Revoke(context.Background(), leaseResp.ID); err != nil {
log.Printf("leaseKeepAlive: revokes the lease (ID: %x): %v\n", leaseResp.ID, err)
if _, err := cli.Revoke(revoker, leaseResp.ID); err != nil {
log.Printf("leaseKeepAlive: revokes the lease (ID: %x): %#v\n", leaseResp.ID, err)
}
cli.Close()
}
return
case <-ticker.C:
if cli, err := e.client(); err != nil {
log.Printf("leaseKeepAlive: creates a new client: %v\n", err)
@ -854,11 +863,26 @@ func (e *Etcd) leaseKeepAlive(ttl int64) (clientv3.LeaseID, func(context.Context
}
}
}
go keepaliveFunc(1 * time.Second)
return leaseResp.ID, keepaliveFunc, nil
return leaseResp.ID, ctx, cancelFunc, nil
}
func (e *Etcd) client() (*clientv3.Client, error) {
type hookedClient struct {
*clientv3.Client
}
func (h *hookedClient) Close() {
// The hook open/closed test here is disabled because there's a
// slight delay before the client actually gets closed in
// some cases, which is long enough to frequently be caught
// if there was a client in the last test run, even though it'd
// be fine a few seconds later.
// _ = testhook.Closed(pilosa.NewAuditor(), h.Client, nil)
h.Client.Close()
}
func (e *Etcd) client() (*hookedClient, error) {
urls := e.e.Server.Cluster().ClientURLs()
cli, err := clientv3.NewFromURLs(urls)
@ -866,10 +890,12 @@ func (e *Etcd) client() (*clientv3.Client, error) {
return nil, errors.Wrapf(err, "creates a new etcd client from URLs (%v)", urls)
}
return cli, nil
// Temporarily disabled, see comment in Close above.
// _ = testhook.Opened(pilosa.NewAuditor(), cli, nil)
return &hookedClient{Client: cli}, nil
}
func memberList(cli *clientv3.Client) (ids []uint64, names []string, urls []string) {
func memberList(cli *hookedClient) (ids []uint64, names []string, urls []string) {
ml, err := cli.MemberList(context.TODO())
if err != nil {
panic(err)
@ -885,7 +911,7 @@ func memberList(cli *clientv3.Client) (ids []uint64, names []string, urls []stri
return
}
func memberAdd(cli *clientv3.Client, peerURL string) (id uint64, name string) {
func memberAdd(cli *hookedClient, peerURL string) (id uint64, name string) {
ma, err := cli.MemberAdd(context.TODO(), []string{peerURL})
if err != nil {
return 0, ""
@ -905,7 +931,7 @@ func (e *Etcd) Shards(ctx context.Context, index, field string) (*roaring.Bitmap
return e.shards(ctx, cli, index, field)
}
func (e *Etcd) shards(ctx context.Context, cli *clientv3.Client, index, field string) (*roaring.Bitmap, error) {
func (e *Etcd) shards(ctx context.Context, cli *hookedClient, index, field string) (*roaring.Bitmap, error) {
key := path.Join(shardPrefix, index, field)
// Get the current shards for the field.
@ -948,7 +974,7 @@ func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roari
// }
// Create a session to acquire a lock.
sess, _ := concurrency.NewSession(cli)
sess, _ := concurrency.NewSession(cli.Client)
defer sess.Close()
muKey := path.Join(lockPrefix, index, field)
@ -1013,7 +1039,7 @@ func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64)
// write shards to etcd.
// Create a session to acquire a lock.
sess, _ := concurrency.NewSession(cli)
sess, _ := concurrency.NewSession(cli.Client)
defer sess.Close()
muKey := path.Join(lockPrefix, index, field)
@ -1082,7 +1108,7 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6
// write shards to etcd.
// Create a session to acquire a lock.
sess, _ := concurrency.NewSession(cli)
sess, _ := concurrency.NewSession(cli.Client)
defer sess.Close()
muKey := path.Join(lockPrefix, index, field)

View file

@ -6848,20 +6848,20 @@ func TestMissingKeyRegression(t *testing.T) {
// queries (HTTP, GRPC, Postgres), etc.).
func TestVariousQueries(t *testing.T) {
for _, clusterSize := range []int{1, 3, 4, 7} {
clusterSize := clusterSize
t.Run(fmt.Sprintf("%d-node", clusterSize), func(t *testing.T) {
t.Parallel()
c := test.MustRunCluster(t, clusterSize)
defer c.Close()
variousQueries(t, clusterSize)
variousQueriesOnTimeFields(t, clusterSize)
variousQueries(t, c)
variousQueriesOnTimeFields(t, c)
})
}
}
// tests for abbreviating time values in queries
func variousQueriesOnTimeFields(t *testing.T, clusterSize int) {
c := test.MustRunCluster(t, clusterSize)
defer c.Close()
func variousQueriesOnTimeFields(t *testing.T, c *test.Cluster) {
ts := func(t time.Time) int64 {
return t.Unix() * 1e+9
}
@ -6984,10 +6984,7 @@ func variousQueriesOnTimeFields(t *testing.T, clusterSize int) {
}
}
func variousQueries(t *testing.T, clusterSize int) {
c := test.MustRunCluster(t, clusterSize)
defer c.Close()
func variousQueries(t *testing.T, c *test.Cluster) {
// Create and populate "likenums" similar to "likes", but without keys on the field.
c.CreateField(t, "users", pilosa.IndexOptions{Keys: true, TrackExistence: true}, "likenums")
c.ImportIDKey(t, "users", "likenums", []test.KeyID{

View file

@ -1596,15 +1596,15 @@ func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, option
requiredDepth = v
}
// Increase bit depth if required.
f.mu.Lock()
bitDepth := bsig.BitDepth
if requiredDepth > bitDepth {
f.mu.Lock()
bsig.BitDepth = requiredDepth
f.options.BitDepth = requiredDepth
f.mu.Unlock()
} else {
requiredDepth = bitDepth
}
f.mu.Unlock()
// Import into each fragment.
for key, data := range dataByFragment {

View file

@ -22,7 +22,6 @@ import (
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/test"
"github.com/pilosa/pilosa/v2/test/port"
)
func TestHandlerOptions(t *testing.T) {
@ -35,15 +34,10 @@ func TestHandlerOptions(t *testing.T) {
t.Fatalf("expected error making handler without options, got nil")
}
var ln net.Listener
err = port.GetPort(func(p int) error {
ln, err = net.Listen("tcp", port.ColonZeroString(p))
if err != nil {
t.Fatal(err)
}
return err
}, 10)
ln, err := net.Listen("tcp", ":0")
if err != nil {
t.Fatalf("creating listener: %v", err)
}
_, err = http.NewHandler(http.OptHandlerListener(ln, ln.Addr().String()))
if err == nil {

View file

@ -16,24 +16,28 @@ package pilosa_test
import (
"fmt"
"net"
"net/http"
"testing"
_ "net/http/pprof"
"github.com/pilosa/pilosa/v2/test/port"
"github.com/pilosa/pilosa/v2/testhook"
)
func TestMain(m *testing.M) {
l, err := net.Listen("tcp", ":0")
if err != nil {
panic(err)
}
port := l.Addr().(*net.TCPAddr).Port
fmt.Printf("pilosa/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
go func() {
err := port.GetPort(func(port int) error {
fmt.Printf("pilosa/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
return http.ListenAndServe(fmt.Sprintf("127.0.0.1:%v", port), nil)
}, 10)
err := http.Serve(l, nil)
if err != nil {
panic(err)
}
}()
testhook.RunTestsWithHooks(m)
}

View file

@ -37,13 +37,8 @@ func (f ShutdownFunc) Finish(tb testing.TB, name string) {
}
}
// ServeTCP creates a TCP listener and serves postgres wire protocol on it.
func ServeTCP(addr string, server *pg.Server) (net.Addr, ShutdownFunc, error) {
listener, err := net.Listen("tcp", addr)
if err != nil {
return nil, nil, errors.Wrap(err, "listening on TCP")
}
// ServeListener serves postgres wire protocol on a listener.
func ServeListener(listener net.Listener, server *pg.Server) (net.Addr, ShutdownFunc, error) {
laddr := listener.Addr()
ctx, cancel := context.WithCancel(context.Background())
@ -58,6 +53,37 @@ func ServeTCP(addr string, server *pg.Server) (net.Addr, ShutdownFunc, error) {
nil
}
// ServeTCP creates a TCP listener and serves postgres wire protocol on it.
func ServeTCP(addr string, server *pg.Server) (net.Addr, ShutdownFunc, error) {
listener, err := net.Listen("tcp", addr)
if err != nil {
return nil, nil, errors.Wrap(err, "listening on TCP")
}
return ServeListener(listener, server)
}
// ServeTLSListener sets up TLS on the server and invokes ServeListener.
func ServeTLSListener(listener net.Listener, server *pg.Server) (net.Addr, ShutdownFunc, error) {
err := SetupTLS(server)
if err != nil {
return nil, nil, errors.Wrap(err, "server TLS setup failed")
}
var tries int = 5
var netAddr net.Addr
var shutdown ShutdownFunc
for i := 0; i < tries; i++ {
if i > 0 {
fmt.Printf("--- try serving TLS again: %d\n", i)
}
if netAddr, shutdown, err = ServeListener(listener, server); err == nil {
break
}
}
return netAddr, shutdown, err
}
// ServeTLS sets up TLS on the server and invokes ServeTCP.
func ServeTLS(addr string, server *pg.Server) (net.Addr, ShutdownFunc, error) {
err := SetupTLS(server)

View file

@ -31,7 +31,6 @@ import (
"github.com/pilosa/pilosa/v2/logger"
"github.com/pilosa/pilosa/v2/pg"
"github.com/pilosa/pilosa/v2/pg/pgtest"
"github.com/pilosa/pilosa/v2/test/port"
)
// TestStartupTimeout tests that an incoming connection that does nothing times out and gets closed.
@ -111,14 +110,7 @@ func TestPQConnect(t *testing.T) {
Logger: logger.NopLogger,
}
var addr net.Addr
var shutdown pgtest.ShutdownFunc
var err error
err = port.GetPort(func(p int) error {
addr, shutdown, err = pgtest.ServeTCP(port.ColonZeroString(p), server)
return err
}, 10)
addr, shutdown, err := pgtest.ServeTCP(":0", server)
if err != nil {
t.Fatalf("starting postgres server: %v", err)
}
@ -150,13 +142,7 @@ func TestPQConnectSSL(t *testing.T) {
Logger: logger.NopLogger,
}
var addr net.Addr
var shutdown pgtest.ShutdownFunc
var err error
err = port.GetPort(func(p int) error {
addr, shutdown, err = pgtest.ServeTLS(port.ColonZeroString(p), server)
return err
}, 10)
addr, shutdown, err := pgtest.ServeTLS(":0", server)
if err != nil {
t.Fatalf("starting postgres server: %v", err)
}
@ -221,13 +207,7 @@ func TestPSQLQuery(t *testing.T) {
Logger: logger.NopLogger,
}
var addr net.Addr
var shutdown pgtest.ShutdownFunc
var err error
err = port.GetPort(func(p int) error {
addr, shutdown, err = pgtest.ServeTCP(port.ColonZeroString(p), server)
return err
}, 10)
addr, shutdown, err := pgtest.ServeTCP(":0", server)
if err != nil {
t.Fatalf("starting postgres server: %v", err)
}
@ -289,13 +269,7 @@ func TestPSQLQuery(t *testing.T) {
CancellationManager: pg.NewLocalCancellationManager(rand.Reader),
}
var addr net.Addr
var shutdown pgtest.ShutdownFunc
var err error
err = port.GetPort(func(p int) error {
addr, shutdown, err = pgtest.ServeTCP(port.ColonZeroString(p), server)
return err
}, 10)
addr, shutdown, err := pgtest.ServeTCP(":0", server)
if err != nil {
t.Fatalf("starting postgres server: %v", err)
}

View file

@ -18,6 +18,7 @@ import (
"context"
"fmt"
"math/rand"
"net"
"net/http"
"os"
"testing"
@ -27,7 +28,6 @@ import (
"github.com/pilosa/pilosa/v2/rbf"
rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg"
"github.com/pilosa/pilosa/v2/test/port"
"golang.org/x/sync/errgroup"
)
@ -351,11 +351,14 @@ func TestDB_MultiTx(t *testing.T) {
// better diagnosis of deadlocks/hung situations versus just really slow "Quick" tests.
func TestMain(m *testing.M) {
l, err := net.Listen("tcp", ":0")
if err != nil {
panic(err)
}
port := l.Addr().(*net.TCPAddr).Port
fmt.Printf("rbf/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
go func() {
err := port.GetPort(func(port int) error {
fmt.Printf("rbf/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
return http.ListenAndServe(fmt.Sprintf("127.0.0.1:%v", port), nil)
}, 10)
err := http.Serve(l, nil)
if err != nil {
panic(err)
}

View file

@ -29,7 +29,6 @@ import (
"github.com/pilosa/pilosa/v2/disco"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
"github.com/pilosa/pilosa/v2/test/port"
)
// Ensure program can send/receive broadcast messages.
@ -185,19 +184,27 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Etcd = portsCfg[0].Etcd
m1.Config.Name = portsCfg[0].Name
m1.Config.Cluster.Name = portsCfg[0].Cluster.Name
m1.Config.BindGRPC = portsCfg[0].BindGRPC
return m1.Start()
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
lsns := make([]*net.TCPListener, 3)
for i := range lsns {
l, err := net.Listen("tcp", ":0")
if err != nil {
t.Fatal(err)
}
lsns[i] = l.(*net.TCPListener)
}
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Etcd = portsCfg[0].Etcd
m1.Config.Name = portsCfg[0].Name
m1.Config.Cluster.Name = portsCfg[0].Cluster.Name
m1.Config.BindGRPC = portsCfg[0].BindGRPC
m1.Config.GRPCListener = portsCfg[0].GRPCListener
err := m1.Start()
if err != nil {
t.Fatal(err)
}
defer m1.Close()
state0, err0 := m0.API.State()

View file

@ -66,6 +66,12 @@ type Config struct {
// BindGRPC is the host:port on which Pilosa will bind for gRPC.
BindGRPC string `toml:"bind-grpc"`
// GRPCListener is an already-bound listener to use for gRPC.
// This is for use by test infrastructure, where it's useful to
// be able to dynamically generate the bindings by actually binding
// to :0, and avoid "address already in use" errors.
GRPCListener *net.TCPListener
// Advertise is the address advertised by the server to other nodes
// in the cluster. It should be reachable by all other nodes and should
// route to an interface that Bind is listening on.

View file

@ -516,7 +516,7 @@ func TestHandler_Endpoints(t *testing.T) {
for _, nodeUsage := range nodeUsages {
numIndexes := len(nodeUsage.Disk.IndexUsage)
if nodeUsage.Disk.TotalUse < 75000 || nodeUsage.Disk.TotalUse > 500000 {
if nodeUsage.Disk.TotalUse < 75000 || nodeUsage.Disk.TotalUse > 600000 {
// Usage measurements are not consistent between machines, or
// over time, as features and implementations change, so checking
// for a range of sizes may be most useful way to test the details of this.

View file

@ -296,16 +296,18 @@ func (m *Command) SetupServer() error {
if err != nil {
return errors.Wrap(err, "processing bind grpc address")
}
// create gRPC listener
m.grpcLn, err = net.Listen("tcp", grpcURI.HostPort())
if err != nil {
return errors.Wrap(err, "creating grpc listener")
}
// If grpc port is 0, get auto-allocated port from listener
if grpcURI.Port == 0 {
grpcURI.SetPort(uint16(m.grpcLn.Addr().(*net.TCPAddr).Port))
if m.Config.GRPCListener == nil {
// create gRPC listener
m.grpcLn, err = net.Listen("tcp", grpcURI.HostPort())
if err != nil {
return errors.Wrap(err, "creating grpc listener")
}
// If grpc port is 0, get auto-allocated port from listener
if grpcURI.Port == 0 {
grpcURI.SetPort(uint16(m.grpcLn.Addr().(*net.TCPAddr).Port))
}
} else {
m.grpcLn = m.Config.GRPCListener
}
// Setup TLS

View file

@ -22,8 +22,8 @@ import (
"fmt"
"io/ioutil"
"math/rand"
"net"
nethttp "net/http"
"os"
"reflect"
"sort"
"strings"
@ -37,7 +37,7 @@ import (
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
"github.com/pilosa/pilosa/v2/test/port"
"github.com/pilosa/pilosa/v2/testhook"
"github.com/pkg/errors"
"golang.org/x/sync/errgroup"
)
@ -1190,16 +1190,19 @@ Set("h", adec=100.22)
}
func TestMain(m *testing.M) {
l, err := net.Listen("tcp", ":0")
if err != nil {
panic(err)
}
port := l.Addr().(*net.TCPAddr).Port
fmt.Printf("server/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
go func() {
err := port.GetPort(func(port int) error {
fmt.Printf("server/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
return nethttp.ListenAndServe(fmt.Sprintf("127.0.0.1:%v", port), nil)
}, 10)
err := nethttp.Serve(l, nil)
if err != nil {
panic(err)
}
}()
os.Exit(m.Run())
testhook.RunTestsWithHooks(m)
}
// TestClusterCreatedAtRace is a regression test for an issue where

View file

@ -18,7 +18,6 @@ import (
"context"
"fmt"
"math"
"net"
"sort"
"strings"
"testing"
@ -30,7 +29,6 @@ import (
"github.com/pilosa/pilosa/v2/proto"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/storage"
"github.com/pilosa/pilosa/v2/test/port"
"github.com/pkg/errors"
"golang.org/x/sync/errgroup"
)
@ -45,6 +43,7 @@ func (*ModHasher) Name() string { return "mod" }
// Cluster represents a Pilosa cluster (multiple Command instances)
type Cluster struct {
Nodes []*Command
tb testing.TB
}
// Query executes an API.Query through one of the cluster's node's API. It fails
@ -82,6 +81,7 @@ func (c *Cluster) QueryGRPC(t testing.TB, index, query string) *proto.TableRespo
if err != nil {
t.Fatalf("getting GRPC client: %v", err)
}
defer grpcClient.Close()
tableResp, err := grpcClient.QueryUnary(context.Background(), index, query)
if err != nil {
@ -376,40 +376,21 @@ func (c *Cluster) CreateField(t testing.TB, index string, iopts pilosa.IndexOpti
// Start runs a Cluster
func (c *Cluster) Start() error {
var eg errgroup.Group
err := port.GetListeners(
func(lsns []*net.TCPListener) (err0 error) {
sliceOfPorts := NewPorts(lsns)
defer func() {
if err0 != nil {
// going to retry. Close the still open listeners
for _, ports := range sliceOfPorts {
_ = ports.Close()
}
}
}()
portsCfg := GenPortsConfig(sliceOfPorts)
for i, cc := range c.Nodes {
cc := cc
cc.Config.Etcd = portsCfg[i].Etcd
cc.Config.Name = portsCfg[i].Name
cc.Config.Cluster.Name = portsCfg[i].Cluster.Name
cc.Config.BindGRPC = portsCfg[i].BindGRPC
eg.Go(func() error {
return cc.Start()
})
}
return eg.Wait()
}, 3*len(c.Nodes), 10)
err := GetPortsGenConfigs(c.tb, c.Nodes)
if err != nil {
return err
return errors.Wrap(err, "configuring cluster ports")
}
var eg errgroup.Group
for _, cc := range c.Nodes {
cc := cc
eg.Go(func() error {
return cc.Start()
})
}
err = eg.Wait()
if err != nil {
return errors.Wrap(err, "starting cluster")
}
return c.GetNode(0).AwaitState(disco.ClusterStateNormal, 30*time.Second)
}
@ -488,13 +469,14 @@ func newCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (*Clust
return nil, errors.New("Slice of CommandOptions must be of length 0, 1, or equal to the number of cluster nodes")
}
cluster := &Cluster{Nodes: make([]*Command, size)}
cluster := &Cluster{Nodes: make([]*Command, size), tb: tb}
for i := 0; i < size; i++ {
var commandOpts []server.CommandOption
if len(opts) > 0 {
commandOpts = opts[i%len(opts)]
}
m := NewCommandNode(tb, commandOpts...)
m.Config.ImportWorkerPoolSize = 2
cluster.Nodes[i] = m
}

View file

@ -19,10 +19,13 @@ import (
"io/ioutil"
"net"
"strings"
"testing"
"time"
"github.com/pilosa/pilosa/v2/etcd"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/testhook"
"github.com/pkg/errors"
)
type Ports struct {
@ -32,16 +35,90 @@ type Ports struct {
LsnP *net.TCPListener
PortP int
LsnG *net.TCPListener
Grpc int
}
func (ports *Ports) Close() error {
err := ports.LsnC.Close()
err2 := ports.LsnP.Close()
err3 := ports.LsnG.Close()
if err != nil {
return err
}
return err2
if err2 != nil {
return err2
}
return err3
}
// listenerPortURL builds a TCP listener and corresponding http://localhost:%d
// URL, and returns those.
func listenerWithURL() (listener *net.TCPListener, url string, err error) {
l, err := net.Listen("tcp", ":0")
if err != nil {
return listener, url, err
}
listener = l.(*net.TCPListener)
port := listener.Addr().(*net.TCPAddr).Port
url = fmt.Sprintf("http://localhost:%d", port)
return listener, url, err
}
// GetPortsGenConfigs creates listener ports, and updates the configurations
// of servers to match these created ports, including cross-references
// like updating the InitCluster values in the Etcd configs.
func GetPortsGenConfigs(tb testing.TB, nodes []*Command) error {
peerUrls := make([]string, len(nodes))
for i := range nodes {
if nodes[i].Config == nil {
nodes[i].Config = &server.Config{}
}
config := nodes[i].Config
name := fmt.Sprintf("server%d", i)
clusterName := fmt.Sprintf("cluster-%s", tb.Name())
discoDir, err := testhook.TempDir(tb, "disco.")
if err != nil {
return errors.Wrap(err, "creating temp directory")
}
clientListener, clientURL, err := listenerWithURL()
if err != nil {
return errors.Wrap(err, "creating client listener")
}
peerListener, peerURL, err := listenerWithURL()
if err != nil {
return errors.Wrap(err, "creating peer listener")
}
grpcListener, grpcUrl, err := listenerWithURL()
if err != nil {
return errors.Wrap(err, "creating gRPC listener")
}
// for grpc, we don't want the http part...
colon := strings.LastIndexByte(grpcUrl, ':')
if colon != -1 {
grpcUrl = grpcUrl[colon:]
}
config.Name = name
config.Cluster.Name = clusterName
config.BindGRPC = grpcUrl
config.GRPCListener = grpcListener
config.Etcd = etcd.Options{
Dir: discoDir,
LClientURL: clientURL,
AClientURL: clientURL,
LPeerURL: peerURL,
APeerURL: peerURL,
HeartbeatTTL: 5 * int64(time.Second),
LPeerSocket: []*net.TCPListener{peerListener},
LClientSocket: []*net.TCPListener{clientListener},
}
peerUrls[i] = fmt.Sprintf("%s=%s", name, peerURL)
}
allPeerUrls := strings.Join(peerUrls, ",")
for i := range nodes {
nodes[i].Config.Etcd.InitCluster = allPeerUrls
}
return nil
}
//GenPortsConfig creates specific configuration for etcd.
@ -63,8 +140,9 @@ func GenPortsConfig(ports []Ports) []*server.Config {
}
cfgs[i] = &server.Config{
Name: name,
BindGRPC: fmt.Sprintf(":%d", ports[i].Grpc),
Name: name,
BindGRPC: fmt.Sprintf(":%d", ports[i].Grpc),
GRPCListener: ports[i].LsnG,
Etcd: etcd.Options{
Dir: discoDir,
LClientURL: lClientURL,
@ -102,12 +180,9 @@ func NewPorts(lsn []*net.TCPListener) []Ports {
PortC: ports[i],
LsnP: lsn[i+1],
PortP: ports[i+1],
Grpc: ports[i+2],
Grpc: ports[i+2],
LsnG: lsn[i+2],
})
// make Grpc port available to
// be rebound.
lsn[i+2].Close()
}
return out

View file

@ -1,97 +0,0 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed 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 port
import (
"fmt"
"log"
"net"
"strings"
"syscall"
)
func ColonZeroString(port int) string {
return fmt.Sprintf(":%d", port)
}
func GetPort(wrapper func(int) error, retries int) error {
f := func(ports []int) error {
return wrapper(ports[0])
}
return GetPorts(f, 1, retries)
}
func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error {
for i := 0; i < retries; i++ {
// get all requested ports
listeners := make([]*net.TCPListener, requestedPorts)
ports := make([]int, requestedPorts)
for i := 0; i < requestedPorts; i++ {
l, err := net.Listen("tcp", ":0")
if err != nil {
log.Println("[port_mapper] error getting a free port", err)
return GetPorts(wrapper, requestedPorts, retries-1)
}
ports[i] = l.Addr().(*net.TCPAddr).Port
listeners[i] = l.(*net.TCPListener)
}
for _, l := range listeners {
if err := l.Close(); err != nil {
log.Println("[port_mapper] error closing the listener", err)
}
}
// send to wrapper and check output error
err := wrapper(ports)
if (err != nil) && (err == syscall.EADDRINUSE || strings.Contains(err.Error(), "address already in use")) {
log.Printf("[port_mapper: %+v] address already in use error calling the wrapper: %v\n", ports, err)
// only retry on address already in use error
continue
}
return err
}
return nil
}
func GetListeners(wrapper func([]*net.TCPListener) error, requestedPorts, retries int) error {
for i := 0; i < retries; i++ {
// get all requested ports
listeners := make([]*net.TCPListener, requestedPorts)
ports := make([]int, requestedPorts)
for i := 0; i < requestedPorts; i++ {
l, err := net.Listen("tcp", ":0")
if err != nil {
log.Println("[port_mapper] error getting a free port", err)
return GetListeners(wrapper, requestedPorts, retries-1)
}
ports[i] = l.Addr().(*net.TCPAddr).Port
listeners[i] = l.(*net.TCPListener)
}
// send to wrapper and check output error
err := wrapper(listeners)
if (err != nil) && (err == syscall.EADDRINUSE || strings.Contains(err.Error(), "address already in use")) {
log.Printf("[port_mapper: %+v] address already in use error calling the wrapper: %v\n", ports, err)
// only retry on address already in use error
continue
}
return err
}
return nil
}