mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 03:47:51 +00:00
Add port wrapper POC
Signed-off-by: Antonio Navarro Perez <antnavper@gmail.com>
This commit is contained in:
parent
61edff3eee
commit
d20b831084
12 changed files with 261 additions and 252 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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() {
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue