Merge pull request #1338 from jaten-molecula/etcd-listner

Etcd listener
This commit is contained in:
jaten-molecula 2021-01-18 14:40:32 -06:00 committed by GitHub
commit 3928b6854b
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
9 changed files with 172 additions and 140 deletions

View file

@ -19,6 +19,7 @@ import (
"context"
"fmt"
"log"
"net"
"path"
"strings"
"time"
@ -45,6 +46,9 @@ type Options struct {
ClusterURL string `toml:"cluster-url"`
ClusterName string `toml:"cluster-name"`
HeartbeatTTL int64 `toml:"heartbeat-ttl"`
LPeerSocket []*net.TCPListener
LClientSocket []*net.TCPListener
}
var (
@ -124,6 +128,19 @@ func parseOptions(opt Options) *embed.Config {
cfg.LPUrls = types.MustNewURLs([]string{opt.LPeerURL})
cfg.APUrls = types.MustNewURLs([]string{opt.APeerURL})
lps := make([]*net.TCPListener, len(opt.LPeerSocket))
copy(lps, opt.LPeerSocket)
cfg.LPeerSocket = lps
lcs := make([]*net.TCPListener, len(opt.LPeerSocket))
copy(lcs, opt.LClientSocket)
cfg.LClientSocket = lcs
cfg.Logger = "zap"
cfg.ZapLoggerBuilder = func(*embed.Config) error {
return nil
}
if opt.InitCluster != "" {
cfg.InitialCluster = opt.InitCluster
cfg.ClusterState = embed.ClusterStateFlagNew

2
go.mod
View file

@ -2,7 +2,7 @@ module github.com/pilosa/pilosa/v2
replace github.com/hashicorp/memberlist => github.com/pilosa/memberlist v0.1.4-0.20190415211605-f6512523c021
replace go.etcd.io/etcd => github.com/molecula/etcd v0.0.0-20210108232729-18e95f2f5b93
replace go.etcd.io/etcd => github.com/molecula/etcd v0.0.0-20210115113447-5d28bda617d2
require (
github.com/CAFxX/gcnotifier v0.0.0-20190112062741-224a280d589d

7
go.sum
View file

@ -247,8 +247,8 @@ github.com/modern-go/reflect2 v1.0.1 h1:9f412s+6RmYXLWZSEzVVgPGK7C2PphHj5RJrvfx9
github.com/modern-go/reflect2 v1.0.1/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0=
github.com/molecula/apophenia v0.0.0-20190827192002-68b7a14a478b h1:cZADDaNYM7xn/nklO3g198JerGQjadFuA0ofxBJgK0Y=
github.com/molecula/apophenia v0.0.0-20190827192002-68b7a14a478b/go.mod h1:uXd1BiH7xLmgkhVmspdJLENv6uGWrTL/MQX2TN7Yz9s=
github.com/molecula/etcd v0.0.0-20210108232729-18e95f2f5b93 h1:9a+hOGmPrcJEfpK07rzeA0D+F99a+2iha5PfDXGLrbE=
github.com/molecula/etcd v0.0.0-20210108232729-18e95f2f5b93/go.mod h1:yVHk9ub3CSBatqGNg7GRmsnfLWtoW60w4eDYfh7vHDg=
github.com/molecula/etcd v0.0.0-20210115113447-5d28bda617d2 h1:pkzCVLSrFQGVQv3raVGJw6aJCJdIZC/z59tUsSU1Zws=
github.com/molecula/etcd v0.0.0-20210115113447-5d28bda617d2/go.mod h1:1X1h4BZ44WjM0LJof1gKKLap1OA4RsicGCDRtACTkLI=
github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223 h1:F9x/1yl3T2AeKLr2AMdilSD8+f9bvMnNN8VS5iDtovc=
github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U=
github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e h1:fD57ERR4JtEqsWbfPhv4DMiApHyliiK5xCTNVSPiaAs=
@ -367,8 +367,6 @@ github.com/zeebo/blake3 v0.0.4/go.mod h1:YOZo8A49yNqM0X/Y+JmDUZshJWLt1laHsNSn5ny
github.com/zeebo/pcg v0.0.0-20181207190024-3cdc6b625a05 h1:4pW5fMvVkrgkMXdvIsVRRTs69DWYA8uNNQsu1stfVKU=
github.com/zeebo/pcg v0.0.0-20181207190024-3cdc6b625a05/go.mod h1:Gr+78ptB0MwXxm//LBaEvBiaXY7hXJ6KGe2V32X2F6E=
go.etcd.io/bbolt v1.3.2/go.mod h1:IbVyRI1SCnLcuJnV2u8VeU0CEYM7e686BmAb1XKL+uU=
go.etcd.io/bbolt v1.3.3 h1:MUGmc65QhB3pIlaQ5bB4LwqSj6GIonVJXpZiaKNyaKk=
go.etcd.io/bbolt v1.3.3/go.mod h1:IbVyRI1SCnLcuJnV2u8VeU0CEYM7e686BmAb1XKL+uU=
go.etcd.io/bbolt v1.3.5 h1:XAzx9gjCb0Rxj7EoqcClPD1d5ZBxZJk0jbuoPHenBt0=
go.etcd.io/bbolt v1.3.5/go.mod h1:G5EMThwa9y8QZGBClrRx5EY+Yw9kAhnjy3bSjsnlVTQ=
go.opencensus.io v0.21.0/go.mod h1:mSImk1erAIZhrmZN+AvHh14ztQfjbGwt4TtuofqLduU=
@ -453,7 +451,6 @@ golang.org/x/sys v0.0.0-20190502145724-3ef323f4f1fd/go.mod h1:h1NjWce9XRLGQEsW7w
golang.org/x/sys v0.0.0-20190507160741-ecd444e8653b/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20190606165138-5da285871e9c/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20190624142023-c5567b49c5d0/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20190826190057-c7b8b68b1456/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191005200804-aed5e4c7ecf9/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191220142924-d4481acd189f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=

View file

@ -18,6 +18,7 @@ import (
"context"
"encoding/json"
"fmt"
"net"
"net/http"
"os"
"reflect"
@ -185,8 +186,8 @@ func TestClusterResize_AddNode(t *testing.T) {
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -242,8 +243,8 @@ func TestClusterResize_AddNode(t *testing.T) {
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -298,8 +299,8 @@ func TestClusterResize_AddNode(t *testing.T) {
m1 := test.NewCommandNode(t, false)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -360,8 +361,8 @@ func TestClusterResize_AddNode(t *testing.T) {
m1 := test.NewCommandNode(t, false)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -416,8 +417,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -474,8 +475,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -538,8 +539,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -600,8 +601,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t, false)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetPorts(func(ports []int) error {
portsCfg := test.GenPortsConfig(test.NewPorts(ports))
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
m1.Config.Gossip.Port = portsCfg[0].Gossip.Port
m1.Config.DisCo = portsCfg[0].DisCo
@ -661,10 +662,12 @@ func TestCluster_GossipMembership(t *testing.T) {
eg.Go(func() error {
// Pass invalid seed as first in list
m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"}
if err := port.GetPort(func(p int) error {
err := port.GetPort(func(p int) error {
m2.Config.Gossip.Port = fmt.Sprintf("%d", p)
return m2.Start()
}, 10); err != nil {
}, 10)
if err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m2.Close()

View file

@ -604,6 +604,10 @@ func (m *Command) Close() error {
}
}
// prevent the closed sockets from being re-injected into etcd.
m.Config.DisCo.LPeerSocket = nil
m.Config.DisCo.LClientSocket = nil
err := eg.Wait()
_ = testhook.Closed(pilosa.NewAuditor(), m, nil)
return errors.Wrap(err, "closing everything")

View file

@ -19,6 +19,7 @@ import (
"fmt"
"io/ioutil"
"math"
"net"
"path"
"strconv"
"strings"
@ -269,40 +270,52 @@ 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.GetPorts(func(ports []int) error {
portsCfg := GenPortsConfig(NewPorts(ports))
err := port.GetListeners(
var gossipSeeds []string
for i, cc := range c.Nodes {
i := i
// get the bind uri to use as the host portion of the gossip seed.
uri, err := pilosa.AddressWithDefaults(cc.Config.Bind)
if err != nil {
return errors.Wrap(err, "processing bind address")
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)
var gossipSeeds []string
for i, cc := range c.Nodes {
i := i
// get the bind uri to use as the host portion of the gossip seed.
uri, err := pilosa.AddressWithDefaults(cc.Config.Bind)
if err != nil {
return errors.Wrap(err, "processing bind address")
}
cc.Config.Gossip.Port = portsCfg[i].Gossip.Port
gossipHost := uri.Host
gossipPort := cc.Config.Gossip.Port
gossipSeeds = append(gossipSeeds, fmt.Sprintf("%s:%s", gossipHost, gossipPort))
}
cc.Config.Gossip.Port = portsCfg[i].Gossip.Port
gossipHost := uri.Host
gossipPort := cc.Config.Gossip.Port
for i, cc := range c.Nodes {
cc := cc
cc.Config.DisCo = portsCfg[i].DisCo
cc.Config.BindGRPC = portsCfg[i].BindGRPC
gossipSeeds = append(gossipSeeds, fmt.Sprintf("%s:%s", gossipHost, gossipPort))
}
eg.Go(func() error {
fmt.Printf("DISCO CONFIG: %+v\n", cc.Config.DisCo)
cc.Config.Gossip.Seeds = gossipSeeds
for i, cc := range c.Nodes {
cc := cc
cc.Config.DisCo = portsCfg[i].DisCo
cc.Config.BindGRPC = portsCfg[i].BindGRPC
return cc.Start()
})
}
eg.Go(func() error {
fmt.Printf("DISCO CONFIG: %+v\n", cc.Config.DisCo)
cc.Config.Gossip.Seeds = gossipSeeds
return eg.Wait()
}, 4*len(c.Nodes), 10)
return cc.Start()
})
}
return eg.Wait()
}, 4*len(c.Nodes), 10)
if err != nil {
return err
}

View file

@ -17,18 +17,33 @@ package test
import (
"fmt"
"io/ioutil"
"net"
"strings"
"time"
"github.com/pilosa/pilosa/v2/etcd"
"github.com/pilosa/pilosa/v2/gossip"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test/port"
)
type Ports struct {
Client, Peer int
Grpc, Gossip int //TODO remove
LsnC *net.TCPListener
PortC int
LsnP *net.TCPListener
PortP int
Grpc int
Gossip int //TODO remove
}
func (ports *Ports) Close() error {
err := ports.LsnC.Close()
err2 := ports.LsnP.Close()
if err != nil {
return err
}
return err2
}
//GenPortsConfig creates specific configuration for etcd.
@ -38,9 +53,11 @@ func GenPortsConfig(ports []Ports) []*server.Config {
for i := range cfgs {
name := fmt.Sprintf("server%d", i)
var lClientURL, lPeerURL string
lClientURL = fmt.Sprintf("http://localhost:%d", ports[i].Client)
lPeerURL = fmt.Sprintf("http://localhost:%d", ports[i].Peer)
lsnC, portC := ports[i].LsnC, ports[i].PortC
lClientURL := fmt.Sprintf("http://localhost:%d", portC)
lsnP, portP := ports[i].LsnP, ports[i].PortP
lPeerURL := fmt.Sprintf("http://localhost:%d", portP)
discoDir := ""
if d, err := ioutil.TempDir("/tmp", "disco."); err == nil {
discoDir = d
@ -50,22 +67,24 @@ func GenPortsConfig(ports []Ports) []*server.Config {
Gossip: gossip.Config{
Port: fmt.Sprint(ports[i].Gossip),
},
BindGRPC: port.ColonZeroString(ports[i].Grpc),
BindGRPC: fmt.Sprintf(":%d", ports[i].Grpc),
DisCo: etcd.Options{
Name: name,
Dir: discoDir,
ClusterName: "bartholemuuuuu",
LClientURL: lClientURL,
AClientURL: lClientURL,
LPeerURL: lPeerURL,
APeerURL: lPeerURL,
HeartbeatTTL: 5 * int64(time.Second),
Name: name,
Dir: discoDir,
ClusterName: "bartholemuuuuu",
LClientURL: lClientURL,
AClientURL: lClientURL,
LPeerURL: lPeerURL,
APeerURL: lPeerURL,
HeartbeatTTL: 5 * int64(time.Second),
LPeerSocket: []*net.TCPListener{lsnP},
LClientSocket: []*net.TCPListener{lsnC},
},
}
clusterURLs[i] = fmt.Sprintf("%s=%s", name, lPeerURL)
fmt.Printf("\ndebug test/disco.go: on i=%v, GenPortsConfig Gossip: %v, DisCo.Client: %v, DisCo.Peer: %v, BindGRPC: %v\n",
i, ports[i].Gossip, ports[i].Client, ports[i].Peer, ports[i].Grpc)
i, ports[i].Gossip, portC, portP, ports[i].Grpc)
}
for i := range cfgs {
cfgs[i].DisCo.InitCluster = strings.Join(clusterURLs, ",")
@ -74,15 +93,29 @@ func GenPortsConfig(ports []Ports) []*server.Config {
return cfgs
}
func NewPorts(ports []int) []Ports {
func NewPorts(lsn []*net.TCPListener) []Ports {
var out []Ports
for i := 0; i < len(ports); i = i + 4 {
n := len(lsn)
ports := make([]int, n)
for i := 0; i < n; i++ {
ports[i] = lsn[i].Addr().(*net.TCPAddr).Port
}
for i := 0; i < n; i = i + 4 {
out = append(out, Ports{
Client: ports[i],
Peer: ports[i+1],
LsnC: lsn[i],
PortC: ports[i],
LsnP: lsn[i+1],
PortP: ports[i+1],
Grpc: ports[i+2],
Gossip: ports[i+3],
})
// make Grpc and Gossip ports available to
// be rebound.
lsn[i+2].Close()
lsn[i+3].Close()
}
return out

View file

@ -27,14 +27,16 @@ func ColonZeroString(port int) string {
}
func GetPort(wrapper func(int) error, retries int) error {
f := func(ports []int) error { return wrapper(ports[0]) }
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.Listener, requestedPorts)
listeners := make([]*net.TCPListener, requestedPorts)
ports := make([]int, requestedPorts)
for i := 0; i < requestedPorts; i++ {
l, err := net.Listen("tcp", ":0")
@ -44,7 +46,7 @@ func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error {
}
ports[i] = l.Addr().(*net.TCPAddr).Port
listeners[i] = l
listeners[i] = l.(*net.TCPListener)
}
for _, l := range listeners {
if err := l.Close(); err != nil {
@ -64,3 +66,32 @@ func GetPorts(wrapper func([]int) error, requestedPorts, retries int) error {
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
}

View file

@ -1,66 +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_test
import (
"fmt"
"net"
"testing"
"github.com/pilosa/pilosa/v2/test/port"
)
func TestPortsAreUnique(t *testing.T) {
t.Skip("do we use this anymore?")
portmap := make(map[int]struct{})
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))
}
portmap[p] = struct{}{}
}
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{})
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))
}
lsn, err := net.Listen("tcp", fmt.Sprintf(":%v", p))
if err != nil {
panic(err)
}
portmap[p] = struct{}{}
lsn.Close()
}
return nil
}, 2000, 3)
if err != nil {
t.Fatal(err)
}
}