diff --git a/cluster_internal_test.go b/cluster_internal_test.go index 01c877b5e..27b60deff 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -434,15 +434,19 @@ func TestCluster_ContainsShards(t *testing.T) { } func TestCluster_Nodes(t *testing.T) { - uri0 := NewTestURIFromHostPort("node0", getport()) - uri1 := NewTestURIFromHostPort("node1", getport()) - uri2 := NewTestURIFromHostPort("node2", getport()) - uri3 := NewTestURIFromHostPort("node3", getport()) + const urisCount = 4 + var uris []pnet.URI + 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) - node0 := &topology.Node{ID: "node0", URI: uri0} - node1 := &topology.Node{ID: "node1", URI: uri1} - node2 := &topology.Node{ID: "node2", URI: uri2} - node3 := &topology.Node{ID: "node3", URI: uri3} + node0 := &topology.Node{ID: "node0", URI: uris[0]} + node1 := &topology.Node{ID: "node1", URI: uris[1]} + node2 := &topology.Node{ID: "node2", URI: uris[2]} + node3 := &topology.Node{ID: "node3", URI: uris[3]} nodes := []*topology.Node{node0, node1, node2} @@ -456,15 +460,15 @@ func TestCluster_Nodes(t *testing.T) { t.Run("Filter", func(t *testing.T) { actual := topology.Nodes(topology.Nodes(nodes).Filter(nodes[1])).URIs() - expected := []pnet.URI{uri0, uri2} + expected := []pnet.URI{uris[0], uris[2]} if !reflect.DeepEqual(actual, expected) { t.Errorf("expected: %v, but got: %v", expected, actual) } }) t.Run("FilterURI", func(t *testing.T) { - actual := topology.Nodes(topology.Nodes(nodes).FilterURI(uri1)).URIs() - expected := []pnet.URI{uri0, uri2} + actual := topology.Nodes(topology.Nodes(nodes).FilterURI(uris[1])).URIs() + expected := []pnet.URI{uris[0], uris[2]} if !reflect.DeepEqual(actual, expected) { t.Errorf("expected: %v, but got: %v", expected, actual) } @@ -484,7 +488,7 @@ func TestCluster_Nodes(t *testing.T) { t.Run("Clone", func(t *testing.T) { clone := topology.Nodes(nodes).Clone() actual := topology.Nodes(clone).URIs() - expected := []pnet.URI{uri0, uri1, uri2} + expected := []pnet.URI{uris[0], uris[1], uris[2]} if !reflect.DeepEqual(actual, expected) { t.Errorf("expected: %v, but got: %v", expected, actual) } @@ -547,11 +551,17 @@ func TestCluster_PreviousNode(t *testing.T) { // NEXT: move this test to internal and unexport IsCoordinator func TestCluster_Coordinator(t *testing.T) { - uri1 := NewTestURIFromHostPort("node1", getport()) - uri2 := NewTestURIFromHostPort("node2", getport()) + const urisCount = 2 + var uris []pnet.URI + 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) - node1 := &topology.Node{ID: "node1", URI: uri1} - node2 := &topology.Node{ID: "node2", URI: uri2} + node1 := &topology.Node{ID: "node1", URI: uris[0]} + node2 := &topology.Node{ID: "node2", URI: uris[1]} c1 := *newCluster() c1.Node = node1 @@ -569,22 +579,22 @@ func TestCluster_Coordinator(t *testing.T) { }) } -func getport() uint16 { - return uint16(port.GlobalPortMap.MustGetPort()) -} - func TestCluster_Topology(t *testing.T) { c1 := NewTestCluster(t, 1) // automatically creates Node{ID: "node0"} - uri0 := NewTestURIFromHostPort("host0", getport()) - uri1 := NewTestURIFromHostPort("host1", getport()) - uri2 := NewTestURIFromHostPort("host2", getport()) - invalid := NewTestURIFromHostPort("invalid", getport()) + const urisCount = 4 + var uris []pnet.URI + 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) - node0 := &topology.Node{ID: "node0", URI: uri0} - node1 := &topology.Node{ID: "node1", URI: uri1} - node2 := &topology.Node{ID: "node2", URI: uri2} - nodeinvalid := &topology.Node{ID: "nodeinvalid", URI: invalid} + node0 := &topology.Node{ID: "node0", URI: uris[0]} + node1 := &topology.Node{ID: "node1", URI: uris[1]} + node2 := &topology.Node{ID: "node2", URI: uris[2]} + nodeinvalid := &topology.Node{ID: "nodeinvalid", URI: uris[3]} t.Run("AddNode", func(t *testing.T) { err := c1.addNode(node1) diff --git a/http/handler_test.go b/http/handler_test.go index 370dfb92d..3a3f66be9 100644 --- a/http/handler_test.go +++ b/http/handler_test.go @@ -35,11 +35,18 @@ func TestHandlerOptions(t *testing.T) { if err == nil { t.Fatalf("expected error making handler without options, got nil") } - ln, err := net.Listen("tcp", fmt.Sprintf(":%d", port.MustGetPort())) - if err != nil { - t.Fatal(err) - } - _, err = http.NewHandler(http.OptHandlerListener(ln)) + + var ln net.Listener + var err error + err = port.GetPort(func(p int) error { + ln, err = net.Listen("tcp", port.ColonZeroString(p) + if err != nil { + t.Fatal(err) + } + + return err + }, 10) + if err == nil { t.Fatalf("expected error making handler without options, got nil") } diff --git a/main_test.go b/main_test.go index 88443b2f6..467f8be53 100644 --- a/main_test.go +++ b/main_test.go @@ -19,14 +19,13 @@ import ( "net/http" "testing" + _ "net/http/pprof" + "github.com/pilosa/pilosa/v2/test/port" "github.com/pilosa/pilosa/v2/testhook" - _ "net/http/pprof" ) func TestMain(m *testing.M) { - port.RaiseUlimitNofiles() - port := port.MustGetPort() fmt.Printf("pilosa/ TestMain: online stack-traces: curl http://localhost:%v/debug/pprof/goroutine?debug=2\n", port) go func() { diff --git a/pg/server_test.go b/pg/server_test.go index c5310a40a..41983887a 100644 --- a/pg/server_test.go +++ b/pg/server_test.go @@ -110,7 +110,15 @@ func TestPQConnect(t *testing.T) { StartupTimeout: time.Second, Logger: logger.NopLogger, } - addr, shutdown, err := pgtest.ServeTCP(port.ColonZeroString(), server) + + 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) + if err != nil { t.Fatalf("starting postgres server: %v", err) } @@ -141,7 +149,14 @@ func TestPQConnectSSL(t *testing.T) { StartupTimeout: time.Second, Logger: logger.NopLogger, } - addr, shutdown, err := pgtest.ServeTLS(port.ColonZeroString(), server) + + 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) if err != nil { t.Fatalf("starting postgres server: %v", err) } @@ -205,7 +220,14 @@ func TestPSQLQuery(t *testing.T) { StartupTimeout: time.Second, Logger: logger.NopLogger, } - addr, shutdown, err := pgtest.ServeTCP(port.ColonZeroString(), server) + + 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) if err != nil { t.Fatalf("starting postgres server: %v", err) } @@ -266,7 +288,14 @@ func TestPSQLQuery(t *testing.T) { Logger: logger.NopLogger, CancellationManager: pg.NewLocalCancellationManager(rand.Reader), } - addr, shutdown, err := pgtest.ServeTCP(port.ColonZeroString(), server) + + 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) if err != nil { t.Fatalf("starting postgres server: %v", err) } diff --git a/server/cluster_test.go b/server/cluster_test.go index 1d1125d76..eab551ba4 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -182,7 +182,12 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) + m1.Config.Gossip.Seeds = []string{seed} err := m1.Start() if err != nil { @@ -231,7 +236,11 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) m1.Config.Gossip.Seeds = []string{seed} err := m1.Start() if err != nil { @@ -280,7 +289,10 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) m1.Config.Gossip.Seeds = []string{seed} err := m1.Start() if err != nil { @@ -335,7 +347,10 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) m1.Config.Gossip.Seeds = []string{seed} err := m1.Start() if err != nil { @@ -384,7 +399,10 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) m1.Config.Gossip.Seeds = []string{seed} err := m1.Start() if err != nil { @@ -437,7 +455,10 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) m1.Config.Gossip.Seeds = []string{seed} err := m1.Start() if err != nil { @@ -496,7 +517,10 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) m1.Config.Gossip.Seeds = []string{seed} errc := make(chan error, 1) go func() { @@ -552,7 +576,10 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { // Configure node1 m1 := test.NewCommandNode(t, false) - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) m1.Config.Gossip.Seeds = []string{seed} errc := make(chan error, 1) go func() { @@ -591,7 +618,10 @@ func TestCluster_GossipMembership(t *testing.T) { m1 := test.NewCommandNode(t, false) defer m1.Close() eg.Go(func() error { - m1.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) // Pass invalid seed as first in list m1.Config.Gossip.Seeds = []string{"http://localhost:8765", seed} err := m1.Start() @@ -605,7 +635,10 @@ func TestCluster_GossipMembership(t *testing.T) { m2 := test.NewCommandNode(t, false) defer m2.Close() eg.Go(func() error { - m2.Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + port.GetPort(func(p int) error { + m1.Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) // Pass invalid seed as first in list m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"} err := m2.Start() diff --git a/server/handler_test.go b/server/handler_test.go index c1a8cf0dc..3f6b78160 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -1399,7 +1399,12 @@ func TestCluster_TranslateStore(t *testing.T) { pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderWithLockerFunc(nil, &sync.Mutex{})), ), ) - cluster.GetNode(0).Config.Gossip.Port = fmt.Sprintf("%d", port.MustGetPort()) + + port.GetPort(func(p int) error { + cluster.GetNode(0).Config.Gossip.Port = fmt.Sprintf("%d", p) + return nil + }, 10) + err := cluster.GetNode(0).Start() if err != nil { t.Fatalf("starting node 0: %v", err) diff --git a/server/server.go b/server/server.go index 36e1c4f12..a28be44d1 100644 --- a/server/server.go +++ b/server/server.go @@ -515,7 +515,10 @@ 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) - gossipPort = port.MustGetPort() + port.GetPort(func(p int) error { + gossipPort = p + return nil + }, 10) m.Config.Gossip.Port = fmt.Sprintf(":%d", gossipPort) m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) } diff --git a/test/cluster.go b/test/cluster.go index 122360a60..36b3420e5 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -257,7 +257,11 @@ func (c *Cluster) Start() error { if err != nil { return errors.Wrap(err, "processing bind address") } - cc.Config.Gossip.Port = fmt.Sprint(port.GlobalPortMap.MustGetPort()) // 63965 given out here. gossip port. + + port.GetPort(func(p int) error { + cc.Config.Gossip.Port = fmt.Sprint(p) // 63965 given out here. gossip port. + return nil + }, 10) gossipHost := uri.Host gossipPort := cc.Config.Gossip.Port diff --git a/test/disco.go b/test/disco.go index 955bd8921..14730b3e3 100644 --- a/test/disco.go +++ b/test/disco.go @@ -30,20 +30,29 @@ func GenDisCoConfig(clusterSize int) []*server.Config { clusterURLs := make([]string, clusterSize) for i := range cfgs { name := fmt.Sprintf("server%d", i) - lClientURL := fmt.Sprintf("http://localhost:%d", port.GlobalPortMap.MustGetPort()) - lPeerURL := fmt.Sprintf("http://localhost:%d", port.GlobalPortMap.MustGetPort()) - cfgs[i] = &server.Config{ - BindGRPC: port.ColonZeroString(), - DisCo: etcd.Options{ - Name: name, - Dir: "", - ClusterName: "bartholemuuuuu", - LClientURL: lClientURL, - AClientURL: lClientURL, - LPeerURL: lPeerURL, - APeerURL: lPeerURL, - }, - } + + var lClientURL, lPeerURL string + port.GetPorts(func(ports []int) error { + lClientURL = fmt.Sprintf("http://localhost:%d", ports[0]) + lPeerURL = fmt.Sprintf("http://localhost:%d", ports[1]) + + cfgs[i] = &server.Config{ + BindGRPC: port.ColonZeroString(ports[2]), + DisCo: etcd.Options{ + Name: name, + Dir: "", + ClusterName: "bartholemuuuuu", + LClientURL: lClientURL, + AClientURL: lClientURL, + LPeerURL: lPeerURL, + APeerURL: lPeerURL, + }, + } + + return nil + + }, 3, 10) + 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) } diff --git a/test/pilosa.go b/test/pilosa.go index 0642a835f..f1a9f1924 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -71,12 +71,18 @@ func newCommand(tb testing.TB, opts ...server.CommandOption) *Command { m.Command = server.NewCommand(bytes.NewReader(nil), ioutil.Discard, ioutil.Discard, opts...) m.Config.DataDir = path defaultConf := server.NewConfig() - if m.Config.Bind == defaultConf.Bind { - m.Config.Bind = fmt.Sprintf("http://localhost:%d", port.GlobalPortMap.MustGetPort()) - } - if m.Config.BindGRPC == defaultConf.BindGRPC { - m.Config.BindGRPC = fmt.Sprintf("http://localhost:%d", port.GlobalPortMap.MustGetPort()) - } + + port.GetPorts(func(ports []int) error { + if m.Config.Bind == defaultConf.Bind { + m.Config.Bind = fmt.Sprintf("http://localhost:%d", ports[0]) + } + if m.Config.BindGRPC == defaultConf.BindGRPC { + m.Config.BindGRPC = fmt.Sprintf("http://localhost:%d", ports[1]) + } + + return nil + }, 2, 10) + m.Config.Translation.MapSize = 140000 m.Config.WorkerPoolSize = 2 diff --git a/test/port/port_mapper.go b/test/port/port_mapper.go index a20e868a0..37295bc60 100644 --- a/test/port/port_mapper.go +++ b/test/port/port_mapper.go @@ -16,173 +16,74 @@ package port import ( "fmt" + "log" "net" - "sync" "syscall" ) -const BlockOfPortsSize = 2000 +// lsn, err := net.Listen("tcp", ":0") +// if err != nil { +// panic(err) +// } +// // must be available to UDP too! +// addr := lsn.Addr() +// port := addr.(*net.TCPAddr).Port +// udpConn, err := net.ListenUDP("udp4", &net.UDPAddr{ +// IP: net.IP{}, // listen on all non-multicast addresses... +// Port: port, +// }) +// if err != nil { +// fmt.Printf("UDP port %v was available on tcp but not udp: %v\n", port, err) +// lsn.Close() +// } else { +// _ = udpConn.Close() +// if lsn == nil { +// panic("lsn should never be nil") +// } +// pm.availPorts[i] = lsn +// i++ +// //println("------ bulk reservation: port mapping reserves port ", lsn.Addr().(*net.TCPAddr).Port) +// } -// GlobalPortMap avoids many races and port conflicts when setting -// up ports for test clusters. Used for tests only. -var GlobalPortMap *globalPortMapper -var GlobalPortMapMu sync.Mutex - -func init() { - RaiseUlimitNofiles() - GlobalPortMap = NewGlobalPortMapper(BlockOfPortsSize) +func ColonZeroString(port int) string { + return fmt.Sprintf(":%d", port) } -func MustGetPort() int { - port := GlobalPortMap.MustGetPort() - return port -} -func ColonZeroString() string { - return fmt.Sprintf(":%d", MustGetPort()) +func GetPort(wrapper func(int) error, retries int) error { + f := func(ports []int) error { return wrapper(ports[0]) } + return GetPorts(f, 1, retries) } -// globalPortMapper maintains a pool of available ports by -// holding them open until GetPort() is called. -type globalPortMapper struct { - numPorts int - availPorts []net.Listener -} - -// newGlobalPortMapper initalizes a globalPortMapper with n ports. -func NewGlobalPortMapper(n int) (pm *globalPortMapper) { - GlobalPortMapMu.Lock() - defer GlobalPortMapMu.Unlock() - - pm = &globalPortMapper{ - numPorts: n, - } - pm.allocateAtTop() - return -} - -var _ = (&globalPortMapper{}).allocate // happy linter - -func (pm *globalPortMapper) allocate() { - println("888888 allocate ports called") - pm.availPorts = make([]net.Listener, pm.numPorts) - i := 0 - for i < pm.numPorts { - lsn, err := net.Listen("tcp", ":0") - if err != nil { - panic(err) - } - // must be available to UDP too! - addr := lsn.Addr() - port := addr.(*net.TCPAddr).Port - udpConn, err := net.ListenUDP("udp4", &net.UDPAddr{ - IP: net.IP{}, // listen on all non-multicast addresses... - Port: port, - }) - if err != nil { - fmt.Printf("UDP port %v was available on tcp but not udp: %v\n", port, err) - lsn.Close() - } else { - _ = udpConn.Close() - if lsn == nil { - panic("lsn should never be nil") +func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error { + for i := 0; i < retries; i++ { + // get all requested ports + listeners := make([]net.Listener, 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) } - pm.availPorts[i] = lsn - i++ - //println("------ bulk reservation: port mapping reserves port ", lsn.Addr().(*net.TCPAddr).Port) + + ports[i] = l.Addr().(*net.TCPAddr).Port + listeners[i] = l } - } -} - -func (pm *globalPortMapper) allocateAtTop() { - println("888888 allocateAtTop ports called") - pm.availPorts = make([]net.Listener, pm.numPorts) - i := 0 - - for next := 65000; i < pm.numPorts && next > 1000; next-- { - lsn, err := net.Listen("tcp", fmt.Sprintf(":%d", next)) - if err != nil { - //fmt.Printf("next=%v, err = %v\n", next+1, err) + 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 == syscall.EADDRINUSE { + log.Println("[port_mapper] address already in use error calling the wrapper", err) + // only retry on addres already in use error continue } - //println("next was avail: ", next+1) - // must be available to UDP too! - addr := lsn.Addr() - port := addr.(*net.TCPAddr).Port - udpConn, err := net.ListenUDP("udp4", &net.UDPAddr{ - IP: net.IP{}, // listen on all non-multicast addresses... - Port: port, - }) - if err != nil { - fmt.Printf("UDP port %v was available on tcp but not udp: %v\n", port, err) - lsn.Close() - } else { - _ = udpConn.Close() - if lsn == nil { - panic("lsn should never be nil") - } - pm.availPorts[i] = lsn - i++ - //println("------ bulk reservation: port mapping reserves port ", lsn.Addr().(*net.TCPAddr).Port) - } + return err } -} - -func (pm *globalPortMapper) GetPort() (port int, err error) { - GlobalPortMapMu.Lock() - defer GlobalPortMapMu.Unlock() - - i := len(pm.availPorts) - if i < 1 { - panic(fmt.Sprintf("ran out of ports, allocate more up front for these tests. had BlockOfPortsSize=%v", BlockOfPortsSize)) - } - - lsn := pm.availPorts[i-1] - addr := lsn.Addr() - port = addr.(*net.TCPAddr).Port - - println("port mapping gives out port ", port) - lsn.Close() - pm.availPorts = pm.availPorts[:i-1] - - // verify that it IS usable again - lsn, err = net.Listen("tcp", fmt.Sprintf(":%d", port)) - if err != nil { - panic(err) - } - lsn.Close() - - return port, nil -} - -func (pm *globalPortMapper) MustGetPort() int { - port, err := pm.GetPort() - if err != nil { - panic(err) - } - //fmt.Printf("port %v allocated at stack:\n'%v'", port, string(debug.Stack())) - return port -} - -// RaiseUlimitNofiles raises the number of open file handles -// to at least 3000. This allows us to reserve 2000 open -// ports for the etcd tests that need to know their ports -// up front and not have them re-used quickly (since -// a socket might be still in TIME_WAIT closing state if -// the server closes first). -func RaiseUlimitNofiles() { - var rLimit syscall.Rlimit - err := syscall.Getrlimit(syscall.RLIMIT_NOFILE, &rLimit) - if err != nil { - panic(fmt.Sprintf("Error Getting Rlimit '%v'", err)) - } - - if rLimit.Cur < 6000 { - rLimit.Cur = 6000 - err = syscall.Setrlimit(syscall.RLIMIT_NOFILE, &rLimit) - if err != nil { - fmt.Println("Error Setting Rlimit ", err) - } - } - fmt.Printf("RaiseUlimitNofiles is now %v\n", rLimit.Cur) + + return nil } diff --git a/test/port/port_mapper_test.go b/test/port/port_mapper_test.go index ceaf7a02b..cbd9616b0 100644 --- a/test/port/port_mapper_test.go +++ b/test/port/port_mapper_test.go @@ -16,6 +16,7 @@ package port_test import ( "fmt" + "log" "net" "testing" @@ -23,36 +24,38 @@ import ( ) func TestPortsAreUnique(t *testing.T) { - - local := port.NewGlobalPortMapper(port.BlockOfPortsSize) - - oracle := make(map[int]bool) - - for i := 0; i < port.BlockOfPortsSize; i++ { - port := local.MustGetPort() - if oracle[port] { - panic(fmt.Sprintf("port %v was already issued!", port)) + portmap := make(map[int]struct{}) + port.GetPorts(func(ports []int) error { + for _, p := range ports { + log.Println("PORTTT", p) + if _, exists := portmap[p]; exists { + panic(fmt.Sprintf("port %v was already issued!", p)) + } + portmap[p] = struct{}{} } - oracle[port] = true - } + + return nil + }, 2000, 3) } func TestPortsAreUsable(t *testing.T) { + portmap := make(map[int]struct{}) + port.GetPorts(func(ports []int) error { + for _, p := range ports { + if _, exists := portmap[p]; exists { + panic(fmt.Sprintf("port %v was already issued!", p)) + } - local := port.NewGlobalPortMapper(port.BlockOfPortsSize) + lsn, err := net.Listen("tcp", fmt.Sprintf(":%v", p)) + if err != nil { + panic(err) + } - oracle := make(map[int]bool) - - for i := 0; i < port.BlockOfPortsSize; i++ { - port := local.MustGetPort() - if oracle[port] { - panic(fmt.Sprintf("port %v was already issued!", port)) + portmap[p] = struct{}{} + lsn.Close() } - lsn, err := net.Listen("tcp", fmt.Sprintf(":%v", port)) - if err != nil { - panic(err) - } - oracle[port] = true - lsn.Close() - } + + return nil + }, 2000, 3) + }