Finish implementing port wrapper

This commit is contained in:
Travis 2021-01-13 21:29:58 -06:00
parent d20b831084
commit 27614c42f7
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
15 changed files with 111 additions and 64 deletions

View file

@ -436,12 +436,14 @@ func TestCluster_ContainsShards(t *testing.T) {
func TestCluster_Nodes(t *testing.T) {
const urisCount = 4
var uris []pnet.URI
port.GetPorts(func(ports []int) error {
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)
}, urisCount, 10); err != nil {
t.Fatalf("getting ports: %v", err)
}
node0 := &topology.Node{ID: "node0", URI: uris[0]}
node1 := &topology.Node{ID: "node1", URI: uris[1]}
@ -553,12 +555,14 @@ func TestCluster_PreviousNode(t *testing.T) {
func TestCluster_Coordinator(t *testing.T) {
const urisCount = 2
var uris []pnet.URI
port.GetPorts(func(ports []int) error {
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)
}, urisCount, 10); err != nil {
t.Fatalf("getting ports: %v", err)
}
node1 := &topology.Node{ID: "node1", URI: uris[0]}
node2 := &topology.Node{ID: "node2", URI: uris[1]}
@ -584,12 +588,14 @@ func TestCluster_Topology(t *testing.T) {
const urisCount = 4
var uris []pnet.URI
port.GetPorts(func(ports []int) error {
if err := port.GetPorts(func(ports []int) error {
for i := 0; i < urisCount; i++ {
uris = append(uris, NewTestURIFromHostPort(fmt.Sprintf("host%d", i), uint16(ports[i])))
}
return nil
}, urisCount, 10)
}, urisCount, 10); err != nil {
t.Fatalf("getting ports: %v", err)
}
node0 := &topology.Node{ID: "node0", URI: uris[0]}
node1 := &topology.Node{ID: "node1", URI: uris[1]}

View file

@ -23,7 +23,6 @@ import (
"github.com/pilosa/pilosa/v2/cmd"
_ "github.com/pilosa/pilosa/v2/test"
"github.com/pilosa/pilosa/v2/test/port"
"github.com/pilosa/pilosa/v2/toml"
"github.com/pkg/errors"
)
@ -37,7 +36,7 @@ func TestServerHelp(t *testing.T) {
}
func nextPort() string {
return fmt.Sprintf(`"localhost:%d"`, port.GlobalPortMap.MustGetPort())
return fmt.Sprintf(`"localhost:%d"`, 0)
}
var _ = nextPort // happy linter

View file

@ -16,7 +16,6 @@ package http_test
import (
"encoding/json"
"fmt"
"net"
"testing"
@ -37,9 +36,8 @@ func TestHandlerOptions(t *testing.T) {
}
var ln net.Listener
var err error
err = port.GetPort(func(p int) error {
ln, err = net.Listen("tcp", port.ColonZeroString(p)
ln, err = net.Listen("tcp", port.ColonZeroString(p))
if err != nil {
t.Fatal(err)
}
@ -47,6 +45,7 @@ func TestHandlerOptions(t *testing.T) {
return err
}, 10)
_, err = http.NewHandler(http.OptHandlerListener(ln))
if err == nil {
t.Fatalf("expected error making handler without options, got nil")
}

View file

@ -26,10 +26,14 @@ import (
)
func TestMain(m *testing.M) {
port := port.MustGetPort()
fmt.Printf("pilosa/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
go func() {
_ = http.ListenAndServe(fmt.Sprintf("127.0.0.1:%v", port), nil)
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)
if err != nil {
panic(err)
}
}()
testhook.RunTestsWithHooks(m)
}

View file

@ -154,7 +154,7 @@ func TestPQConnectSSL(t *testing.T) {
var shutdown pgtest.ShutdownFunc
var err error
err = port.GetPort(func(p int) error {
addr, shutdown, err = pgtest.ServeTCP(port.ColonZeroString(p), server)
addr, shutdown, err = pgtest.ServeTLS(port.ColonZeroString(p), server)
return err
}, 10)
if err != nil {

View file

@ -23,11 +23,12 @@ import (
"testing"
"time"
_ "net/http/pprof"
"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"
_ "net/http/pprof"
)
func TestDB_Open(t *testing.T) {
@ -350,17 +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) {
port := port.MustGetPort()
fmt.Printf("rbf/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
go func() {
_ = http.ListenAndServe(fmt.Sprintf("127.0.0.1:%v", port), nil)
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)
if err != nil {
panic(err)
}
}()
os.Exit(m.Run())
}
/*func getAvailPort() int {
l, _ := net.Listen("tcp", ":0")
r := l.Addr()
l.Close()
return r.(*net.TCPAddr).Port
}*/

View file

@ -183,10 +183,12 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
@ -237,10 +239,12 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
if err != nil {
@ -289,10 +293,12 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
if err != nil {
@ -347,10 +353,12 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
if err != nil {
@ -399,10 +407,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
if err != nil {
@ -455,10 +465,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
err := m1.Start()
if err != nil {
@ -517,10 +529,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
errc := make(chan error, 1)
go func() {
@ -576,10 +590,12 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
m1.Config.Gossip.Seeds = []string{seed}
errc := make(chan error, 1)
go func() {
@ -618,10 +634,12 @@ func TestCluster_GossipMembership(t *testing.T) {
m1 := test.NewCommandNode(t, false)
defer m1.Close()
eg.Go(func() error {
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
// Pass invalid seed as first in list
m1.Config.Gossip.Seeds = []string{"http://localhost:8765", seed}
err := m1.Start()
@ -635,10 +653,12 @@ func TestCluster_GossipMembership(t *testing.T) {
m2 := test.NewCommandNode(t, false)
defer m2.Close()
eg.Go(func() error {
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting gossip port: %v", err)
}
// Pass invalid seed as first in list
m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"}
err := m2.Start()

View file

@ -1400,10 +1400,12 @@ func TestCluster_TranslateStore(t *testing.T) {
),
)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
cluster.GetNode(0).Config.Gossip.Port = fmt.Sprintf("%d", p)
return nil
}, 10)
}, 10); err != nil {
t.Fatalf("getting port: %v", err)
}
err := cluster.GetNode(0).Start()
if err != nil {

View file

@ -515,10 +515,12 @@ func (m *Command) setupNetworking() error {
// new port. See also the gossip config in gossip/gossip.go.
// TODO: Maybe make that more configurable here.
m.logger.Printf("ephemeral port %d already occupied, switching to :0 (%v)", gossipPort, err)
port.GetPort(func(p int) error {
if err := port.GetPort(func(p int) error {
gossipPort = p
return nil
}, 10)
}, 10); err != nil {
return errors.Wrap(err, "getting port")
}
m.Config.Gossip.Port = fmt.Sprintf(":%d", gossipPort)
m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger())
}

View file

@ -1228,10 +1228,14 @@ Set("h", adec=100.22)
}
func TestMain(m *testing.M) {
port := port.MustGetPort()
fmt.Printf("server/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port)
go func() {
_ = nethttp.ListenAndServe(fmt.Sprintf("127.0.0.1:%v", port), nil)
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)
if err != nil {
panic(err)
}
}()
os.Exit(m.Run())
}

View file

@ -258,10 +258,12 @@ func (c *Cluster) Start() error {
return errors.Wrap(err, "processing bind address")
}
port.GetPort(func(p int) error {
cc.Config.Gossip.Port = fmt.Sprint(p) // 63965 given out here. gossip port.
if err := port.GetPort(func(p int) error {
cc.Config.Gossip.Port = fmt.Sprint(p)
return nil
}, 10)
}, 10); err != nil {
return errors.Wrap(err, "getting gossip port")
}
gossipHost := uri.Host
gossipPort := cc.Config.Gossip.Port

View file

@ -32,7 +32,7 @@ func GenDisCoConfig(clusterSize int) []*server.Config {
name := fmt.Sprintf("server%d", i)
var lClientURL, lPeerURL string
port.GetPorts(func(ports []int) error {
err := port.GetPorts(func(ports []int) error {
lClientURL = fmt.Sprintf("http://localhost:%d", ports[0])
lPeerURL = fmt.Sprintf("http://localhost:%d", ports[1])
@ -50,8 +50,10 @@ func GenDisCoConfig(clusterSize int) []*server.Config {
}
return nil
}, 3, 10)
if err != nil {
panic(err)
}
clusterURLs[i] = fmt.Sprintf("%s=%s", name, lPeerURL)
fmt.Printf("\ndebug test/disco.go: on i=%v, GenDisCoConfig BindGRPC: %v\n", i, cfgs[i].BindGRPC)

View file

@ -72,7 +72,7 @@ func newCommand(tb testing.TB, opts ...server.CommandOption) *Command {
m.Config.DataDir = path
defaultConf := server.NewConfig()
port.GetPorts(func(ports []int) error {
if err := port.GetPorts(func(ports []int) error {
if m.Config.Bind == defaultConf.Bind {
m.Config.Bind = fmt.Sprintf("http://localhost:%d", ports[0])
}
@ -81,7 +81,9 @@ func newCommand(tb testing.TB, opts ...server.CommandOption) *Command {
}
return nil
}, 2, 10)
}, 2, 10); err != nil {
panic(err)
}
m.Config.Translation.MapSize = 140000
m.Config.WorkerPoolSize = 2

View file

@ -78,7 +78,7 @@ func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error {
err := wrapper(ports)
if err == syscall.EADDRINUSE {
log.Println("[port_mapper] address already in use error calling the wrapper", err)
// only retry on addres already in use error
// only retry on address already in use error
continue
}

View file

@ -24,8 +24,9 @@ import (
)
func TestPortsAreUnique(t *testing.T) {
t.Skip("do we use this anymore?")
portmap := make(map[int]struct{})
port.GetPorts(func(ports []int) error {
err := port.GetPorts(func(ports []int) error {
for _, p := range ports {
log.Println("PORTTT", p)
if _, exists := portmap[p]; exists {
@ -36,11 +37,15 @@ func TestPortsAreUnique(t *testing.T) {
return nil
}, 2000, 3)
if err != nil {
t.Fatal(err)
}
}
func TestPortsAreUsable(t *testing.T) {
t.Skip("do we use this anymore?")
portmap := make(map[int]struct{})
port.GetPorts(func(ports []int) error {
err := port.GetPorts(func(ports []int) error {
for _, p := range ports {
if _, exists := portmap[p]; exists {
panic(fmt.Sprintf("port %v was already issued!", p))
@ -57,5 +62,7 @@ func TestPortsAreUsable(t *testing.T) {
return nil
}, 2000, 3)
if err != nil {
t.Fatal(err)
}
}