remove the rest of the gossip code (except config)

This commit is contained in:
Travis 2021-02-04 10:55:55 -06:00
parent 3e90a88c34
commit a4b37273ea
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
11 changed files with 16 additions and 800 deletions

1
api.go
View file

@ -214,7 +214,6 @@ func (api *API) CreateIndex(ctx context.Context, indexName string, options Index
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
if !snap.IsPrimaryFieldTranslationNode(api.Node().ID) {
fmt.Println("--- DEBUG: forward to coordinator")
if err := api.server.defaultClient.CreateIndex(ctx, indexName, options); err != nil {
return nil, errors.Wrap(err, "forwarding CreateIndex to coordinator")
}

View file

@ -105,7 +105,7 @@ type cluster struct { // nolint: maligned
sharder disco.Sharder
// Required for cluster Resize.
Static bool // Static is primarily used for testing in a non-gossip environment.
Static bool // Static is primarily used for testing.
holder *Holder
broadcaster broadcaster

View file

@ -201,8 +201,6 @@ func TestServerConfig_DeprecateLongQueryTime(t *testing.T) {
bind = ` + nextPort() + `
bind-grpc = ` + nextPort() + `
data-dir = "` + actualDataDir + `"
[gossip]
port = "14321"
`,
validation: func() error {
v := validator{}
@ -218,8 +216,6 @@ func TestServerConfig_DeprecateLongQueryTime(t *testing.T) {
cfgFileContent: `
bind = ` + nextPort() + `
bind-grpc = ` + nextPort() + `
[gossip]
port = "14321"
`,
validation: func() error {
v := validator{}
@ -235,8 +231,6 @@ func TestServerConfig_DeprecateLongQueryTime(t *testing.T) {
cfgFileContent: `
bind = ` + nextPort() + `
bind-grpc = ` + nextPort() + `
[gossip]
port = "14321"
`,
validation: func() error {
v := validator{}

View file

@ -15,556 +15,9 @@
package gossip
import (
"bytes"
"context"
"fmt"
"io"
"io/ioutil"
"log"
"net"
"os"
"strconv"
"strings"
"sync"
"time"
"github.com/hashicorp/memberlist"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/logger"
pnet "github.com/pilosa/pilosa/v2/net"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/toml"
"github.com/pilosa/pilosa/v2/topology"
"github.com/pkg/errors"
)
// Ensure GossipMemberSet implements interfaces.
var _ memberlist.Delegate = &memberSet{}
// memberSet represents a gossip implementation of MemberSet using memberlist.
type memberSet struct {
mu sync.RWMutex
memberlist *memberlist.Memberlist
broadcasts *memberlist.TransmitLimitedQueue
papi *pilosa.API
config *config
Logger logger.Logger
// stdLogger is only used when passed into memberlist library things that take a std library logger rather than an interface.
stdLogger *log.Logger
// logOutput is similar to stdLogger in that it's passed to memberlist things which can't take a pilosa Logger.
logOutput io.Writer
transport *Transport
eventReceiver *eventReceiver
}
// Open implements the MemberSet interface to start network activity.
func (g *memberSet) Open() (err error) {
g.mu.Lock()
defer g.mu.Unlock()
g.memberlist, err = memberlist.Create(g.config.memberlistConfig)
if err != nil {
return errors.Wrap(err, "creating memberlist")
}
g.broadcasts = &memberlist.TransmitLimitedQueue{
NumNodes: func() int {
g.mu.RLock()
defer g.mu.RUnlock()
return g.memberlist.NumMembers()
},
RetransmitMult: 3,
}
var uris = make([]*pnet.URI, len(g.config.gossipSeeds))
for i, addr := range g.config.gossipSeeds {
uris[i], err = pnet.NewURIFromAddress(addr)
if err != nil {
return fmt.Errorf("new uri from address: %s", err)
}
}
var nodes = make([]*topology.Node, len(uris))
for i, uri := range uris {
nodes[i] = &topology.Node{URI: *uri}
}
err = g.joinWithRetry(pnet.URIs(topology.Nodes(nodes).URIs()).HostPortStrings())
if err != nil {
return errors.Wrap(err, "joinWithRetry")
}
return nil
}
// Close attempts to gracefully leave the cluster, and finally calls shutdown
// after (at most) a timeout period.
func (g *memberSet) Close() error {
defer g.eventReceiver.Close()
leaveErr := g.memberlist.Leave(5 * time.Second)
shutdownErr := g.memberlist.Shutdown()
if leaveErr != nil || shutdownErr != nil {
return fmt.Errorf("leaving: '%v', shutting down: '%v'", leaveErr, shutdownErr)
}
return nil
}
// joinWithRetry wraps the standard memberlist Join function in a retry.
func (g *memberSet) joinWithRetry(hosts []string) error {
err := retry(60, 2*time.Second, func() error {
_, err := g.memberlist.Join(hosts)
return err
})
return err
}
// retry periodically retries function fn a specified number of attempts.
func retry(attempts int, sleep time.Duration, fn func() error) (err error) { // nolint: unparam
for i := 0; ; i++ {
err = fn()
if err == nil {
return
}
if i >= (attempts - 1) {
break
}
time.Sleep(sleep)
log.Println("retrying after error:", err)
}
return fmt.Errorf("after %d attempts, last error: %s", attempts, err)
}
////////////////////////////////////////////////////////////////
type config struct {
gossipSeeds []string
memberlistConfig *memberlist.Config
}
// memberSetOption describes a functional option for GossipMemberSet.
type memberSetOption func(*memberSet) error
// WithTransport is a functional option for providing a transport to NewMemberSet.
func WithTransport(transport *Transport) memberSetOption {
return func(g *memberSet) error {
g.transport = transport
return nil
}
}
// WithLogger is a functional option for providing a Go logger to NewMemberSet.
// If the memberSet's transport is nil, this logger will be used when creating
// one. If WithLogOutput is not used, this logger will be passed to memberlist
// for it to use internally. This logger is not used for logging by code in this
// (gossip) package - for that, use the WithPilosaLogger option.
func WithLogger(logger *log.Logger) memberSetOption {
return func(g *memberSet) error {
g.stdLogger = logger
return nil
}
}
// WithLogOutput allows one to pass a Writer which will in turn be passed to
// memberlist for use in logging.
func WithLogOutput(o io.Writer) memberSetOption {
return func(g *memberSet) error {
g.logOutput = o
return nil
}
}
// WithPilosaLogger allows one to configure a memberSet with a logger of their
// choice which satisfies the pilosa logger interface.
func WithPilosaLogger(l logger.Logger) memberSetOption {
return func(g *memberSet) error {
g.Logger = l
return nil
}
}
// NewMemberSet returns a new instance of GossipMemberSet based on options. The
// logging options which can be passed to NewMemberSet are complicated for
// historical reasons - please pass WithPilosaLogger, and either WithLogOutput
// or WithLogger. If you pass WithLogOutput, be sure to also pass in a Transport
// using WithTransport.
func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*memberSet, error) {
host := api.Node().URI.Host
g := &memberSet{
papi: api,
Logger: logger.NopLogger,
}
// options
for _, opt := range options {
if err := opt(g); err != nil {
return nil, errors.Wrap(err, "executing option")
}
}
ger := newEventReceiver(g.Logger, api)
g.eventReceiver = ger
if g.transport == nil {
port, err := strconv.Atoi(cfg.Port)
if err != nil {
return nil, fmt.Errorf("convert port: %s", err)
}
if g.stdLogger == nil {
if g.logOutput != nil {
g.stdLogger = logger.NewStandardLogger(g.logOutput).Logger()
} else {
g.stdLogger = log.New(os.Stderr, "", log.LstdFlags)
}
}
// Set up the transport.
transport, err := NewTransport(host, port, g.stdLogger)
if err != nil {
return nil, fmt.Errorf("new tranport: %s", err)
}
g.transport = transport
}
port := g.transport.net.GetAutoBindPort()
var gossipKey []byte
var err error
if cfg.Key != "" {
gossipKey, err = ioutil.ReadFile(cfg.Key)
if err != nil {
return nil, fmt.Errorf("reading gossip key: %s", err)
}
}
////////////////////
// memberlist config
conf := memberlist.DefaultWANConfig()
conf.Transport = g.transport.net
conf.Name = api.Node().ID
conf.BindAddr = api.Node().URI.Host
conf.BindPort = port
// AdvertisePort
if cfg.AdvertisePort != "" {
if p, err := strconv.Atoi(cfg.Port); err != nil {
return nil, fmt.Errorf("convert advertise port: %s", err)
} else {
conf.AdvertisePort = p
}
} else {
conf.AdvertisePort = port
}
// AdvertiseHost
if cfg.AdvertiseHost != "" {
conf.AdvertiseAddr = cfg.AdvertiseHost
} else {
conf.AdvertiseAddr = hostToIP(api.Node().URI.Host)
}
//
conf.TCPTimeout = time.Duration(cfg.StreamTimeout)
conf.SuspicionMult = cfg.SuspicionMult
conf.PushPullInterval = time.Duration(cfg.PushPullInterval)
conf.ProbeTimeout = time.Duration(cfg.ProbeTimeout)
conf.ProbeInterval = time.Duration(cfg.ProbeInterval)
conf.GossipNodes = cfg.Nodes
conf.GossipInterval = time.Duration(cfg.Interval)
conf.GossipToTheDeadTime = time.Duration(cfg.ToTheDeadTime)
//
conf.Delegate = g
conf.SecretKey = gossipKey
conf.Events = ger
if g.logOutput != nil {
conf.LogOutput = g.logOutput
} else {
conf.Logger = g.stdLogger
}
g.config = &config{
memberlistConfig: conf,
gossipSeeds: cfg.Seeds,
}
return g, nil
}
// NodeMeta implementation of the memberlist.Delegate interface.
func (g *memberSet) NodeMeta(limit int) []byte {
buf, err := g.papi.Serializer.Marshal(g.papi.Node())
if err != nil {
g.Logger.Printf("marshal message error: %s", err)
return []byte{}
}
return buf
}
// NotifyMsg implementation of the memberlist.Delegate interface
// called when a user-data message is received.
func (g *memberSet) NotifyMsg(b []byte) {
err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(b))
if err != nil {
g.Logger.Printf("cluster message error: %s", err)
}
}
// GetBroadcasts implementation of the memberlist.Delegate interface
// called when user data messages can be broadcast.
func (g *memberSet) GetBroadcasts(overhead, limit int) [][]byte {
return g.broadcasts.GetBroadcasts(overhead, limit)
}
// LocalState implementation of the memberlist.Delegate interface
// sends this Node's state data.
func (g *memberSet) LocalState(join bool) []byte {
schema, err := g.papi.Schema(context.Background())
if err != nil {
// just panic, this code will be removed soon
panic(err)
}
m := &pilosa.NodeStatus{
Node: g.papi.Node(),
Schema: &pilosa.Schema{Indexes: schema},
}
for _, idx := range m.Schema.Indexes {
is := &pilosa.IndexStatus{Name: idx.Name, CreatedAt: idx.CreatedAt}
for _, f := range idx.Fields {
availableShards := roaring.NewBitmap()
if field, _ := g.papi.Field(context.Background(), idx.Name, f.Name); field != nil {
availableShards = field.AvailableShards(false)
}
fs := &pilosa.FieldStatus{
Name: f.Name,
CreatedAt: f.CreatedAt,
AvailableShards: availableShards,
}
is.Fields = append(is.Fields, fs)
}
m.Indexes = append(m.Indexes, is)
}
// Marshal nodestate data to bytes.
buf, err := pilosa.MarshalInternalMessage(m, g.papi.Serializer)
if err != nil {
g.Logger.Printf("error marshalling nodestate data, err=%s", err)
return []byte{}
}
return buf
}
// MergeRemoteState implementation of the memberlist.Delegate interface
// receive and process the remote side's LocalState.
func (g *memberSet) MergeRemoteState(buf []byte, join bool) {
err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(buf))
if err != nil {
g.Logger.Printf("merge state error: %s", err)
}
}
// eventReceiver is used to enable an application to receive
// events about joins and leaves over a channel.
//
// Care must be taken that events are processed in a timely manner from
// the channel, since this delegate will block until an event can be sent.
type eventReceiver struct {
ch chan memberlist.NodeEvent
closed chan struct{}
papi *pilosa.API
logger logger.Logger
}
// newEventReceiver returns a new instance of GossipEventReceiver.
func newEventReceiver(logger logger.Logger, papi *pilosa.API) *eventReceiver {
ger := &eventReceiver{
ch: make(chan memberlist.NodeEvent, 1),
closed: make(chan struct{}),
logger: logger,
papi: papi,
}
go ger.listen()
return ger
}
func (g *eventReceiver) NotifyJoin(n *memberlist.Node) {
// copy node to avoid data race
n2 := *n
n2.Meta = make([]byte, len(n.Meta))
copy(n2.Meta, n.Meta)
select {
case g.ch <- memberlist.NodeEvent{Event: memberlist.NodeJoin, Node: &n2}:
case <-g.closed:
}
}
func (g *eventReceiver) NotifyLeave(n *memberlist.Node) {
// copy node to avoid data race
n2 := *n
n2.Meta = make([]byte, len(n.Meta))
copy(n2.Meta, n.Meta)
select {
case g.ch <- memberlist.NodeEvent{Event: memberlist.NodeLeave, Node: &n2}:
case <-g.closed:
}
}
func (g *eventReceiver) NotifyUpdate(n *memberlist.Node) {
// copy node to avoid data race
n2 := *n
n2.Meta = make([]byte, len(n.Meta))
copy(n2.Meta, n.Meta)
select {
case g.ch <- memberlist.NodeEvent{Event: memberlist.NodeUpdate, Node: &n2}:
case <-g.closed:
}
}
func (g *eventReceiver) Close() {
// TODO workaround to make tests pass. We are going to delete this code anyways.
select {
case <-g.closed:
return
default:
close(g.closed)
}
}
func (g *eventReceiver) listen() {
var nodeEventType pilosa.NodeEventType
for {
var e memberlist.NodeEvent
select {
case <-g.closed:
return
case e = <-g.ch:
}
switch e.Event {
case memberlist.NodeJoin:
nodeEventType = pilosa.NodeJoin
case memberlist.NodeLeave:
nodeEventType = pilosa.NodeLeave
case memberlist.NodeUpdate:
nodeEventType = pilosa.NodeUpdate
default:
continue
}
// Get the node from the event.Node meta data.
var n topology.Node
if err := g.papi.Serializer.Unmarshal(e.Node.Meta, &n); err != nil {
panic("failed to unmarshal event node meta into node")
}
ne := &pilosa.NodeEvent{
Event: nodeEventType,
Node: &n,
}
buf, err := pilosa.MarshalInternalMessage(ne, g.papi.Serializer)
if err != nil {
panic(err)
}
if err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(buf)); err != nil {
g.logger.Printf("receive event error: %s", err)
}
}
}
// Transport is a gossip transport for binding to a port.
type Transport struct {
//memberlist.Transport
net *memberlist.NetTransport
URI *pnet.URI
}
// NewTransport returns a NetTransport based on the given host and port.
// It will dynamically bind to a port if port is 0.
// This is useful for test cases where specifying a port is not reasonable.
//func NewTransport(host string, port int) (*memberlist.NetTransport, error) {
func NewTransport(host string, port int, logger *log.Logger) (*Transport, error) {
// memberlist config
conf := memberlist.DefaultWANConfig()
conf.BindAddr = host
conf.BindPort = port
conf.AdvertisePort = port
conf.Logger = logger
net, err := newTransport(conf)
if err != nil {
return nil, fmt.Errorf("new transport: %s", err)
}
uri, err := pnet.NewURIFromHostPort(host, uint16(net.GetAutoBindPort()))
if err != nil {
return nil, fmt.Errorf("new uri from host port: %s", err)
}
return &Transport{
net: net,
URI: uri,
}, nil
}
// newTransport returns a NetTransport based on the memberlist configuration.
// It will dynamically bind to a port if conf.BindPort is 0.
func newTransport(conf *memberlist.Config) (*memberlist.NetTransport, error) {
nc := &memberlist.NetTransportConfig{
BindAddrs: []string{conf.BindAddr},
BindPort: conf.BindPort,
Logger: conf.Logger,
}
if conf.BindPort == 0 {
panic("TODO: remove this. problem: gossip conf.BindPort was 0!")
}
// See comment below for details about the retry in here.
makeNetRetry := func(limit int) (*memberlist.NetTransport, error) {
var err error
for try := 0; try < limit; try++ {
var nt *memberlist.NetTransport
if nt, err = memberlist.NewNetTransport(nc); err == nil {
return nt, nil
}
if strings.Contains(err.Error(), "address already in use") {
conf.Logger.Printf("[DEBUG] Got bind error: %v", err)
continue
}
}
return nil, fmt.Errorf("failed to obtain an address: %v", err)
}
// The dynamic bind port operation is inherently racy because
// even though we are using the kernel to find a port for us, we
// are attempting to bind multiple protocols (and potentially
// multiple addresses) with the same port number. We build in a
// few retries here since this often gets transient errors in
// busy unit tests.
limit := 1
if conf.BindPort == 0 {
limit = 10
}
nt, err := makeNetRetry(limit)
if err != nil {
return nil, errors.Wrap(err, "could not set up network transport")
}
return nt, nil
}
// Config holds toml-friendly memberlist configuration.
type Config struct {
// Port indicates the port to which pilosa should bind for internal state sharing.
@ -638,21 +91,3 @@ type Config struct {
Nodes int `toml:"nodes"`
ToTheDeadTime toml.Duration `toml:"to-the-dead-time"`
}
// hostToIP converts host to an IP4 address based on net.LookupIP().
func hostToIP(host string) string {
// if host is not an IP addr, check net.LookupIP()
if net.ParseIP(host) == nil {
hosts, err := net.LookupIP(host)
if err != nil {
return host
}
for _, h := range hosts {
// this restricts pilosa to IP4
if h.To4() != nil {
return h.String()
}
}
}
return host
}

View file

@ -29,7 +29,6 @@ import (
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
"github.com/pilosa/pilosa/v2/test/port"
"golang.org/x/sync/errgroup"
)
// Ensure program can send/receive broadcast messages.
@ -173,8 +172,6 @@ func TestClusterResize_AddNode(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -188,19 +185,16 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.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()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m1.Close()
@ -219,8 +213,6 @@ func TestClusterResize_AddNode(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -249,19 +241,16 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.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()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m1.Close()
@ -283,8 +272,6 @@ func TestClusterResize_AddNode(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -309,19 +296,17 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.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()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m1.Close()
@ -345,8 +330,6 @@ func TestClusterResize_AddNode(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -375,19 +358,17 @@ func TestClusterResize_AddNode(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.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()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
@ -416,8 +397,6 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -436,17 +415,15 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.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()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m1.Close()
@ -468,8 +445,6 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -498,17 +473,15 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.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()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
errc := make(chan error, 1)
@ -536,8 +509,6 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -566,11 +537,9 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.Etcd = portsCfg[0].Etcd
m1.Config.Name = portsCfg[0].Name
m1.Config.Cluster.Name = portsCfg[0].Cluster.Name
@ -582,7 +551,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
errc <- err
}()
return m1.Start()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m1.Close()
@ -604,8 +573,6 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
// Create a client for each node.
client0 := m0.Client()
@ -632,11 +599,9 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
// Configure node1
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
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.Etcd = portsCfg[0].Etcd
m1.Config.Name = portsCfg[0].Name
m1.Config.Cluster.Name = portsCfg[0].Cluster.Name
@ -648,7 +613,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
errc <- err
}()
return m1.Start()
}, 4, 10); err != nil {
}, 3, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
@ -664,74 +629,6 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
})
}
// Ensure that redundant gossip seeds are used
func TestCluster_GossipMembership(t *testing.T) {
t.Skip("skipping gossip test")
t.Run("Node0Down", func(t *testing.T) {
// Configure node0
m0 := test.MustRunCluster(t, 1).GetNode(0)
defer m0.Close()
seed := ""
var eg errgroup.Group
// Configure node1
m1 := test.NewCommandNode(t)
defer m1.Close()
eg.Go(func() error {
// Pass invalid seed as first in list
m1.Config.Gossip.Seeds = []string{"http://localhost:8765", seed}
if err := port.GetPort(func(p int) error {
m1.Config.Gossip.Port = fmt.Sprintf("%d", p)
return m1.Start()
}, 10); err != nil {
t.Fatalf("starting second main: %v", err)
}
return nil
})
// Configure node1
m2 := test.NewCommandNode(t)
defer m2.Close()
eg.Go(func() error {
// Pass invalid seed as first in list
m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"}
err := port.GetPort(func(p int) error {
m2.Config.Gossip.Port = fmt.Sprintf("%d", p)
return m2.Start()
}, 10)
if err != nil {
t.Fatalf("starting second main: %v", err)
}
defer m2.Close()
return nil
})
if err := eg.Wait(); err != nil {
t.Fatal(err)
}
state0, err0 := m0.API.State()
state1, err1 := m1.API.State()
state2, err2 := m2.API.State()
if err0 != nil || !test.CheckClusterState(m0, string(pilosa.ClusterStateNormal), 1000) {
t.Fatalf("unexpected node0 cluster state: %s, error: %v", state0, err0)
} else if err1 != nil || !test.CheckClusterState(m1, string(pilosa.ClusterStateNormal), 1000) {
t.Fatalf("unexpected node1 cluster state: %s, error: %v", state1, err1)
} else if err2 != nil || !test.CheckClusterState(m2, string(pilosa.ClusterStateNormal), 1000) {
t.Fatalf("unexpected node2 cluster state: %s, error: %v", state2, err2)
}
numNodes := len(m0.API.Hosts(context.Background()))
if numNodes != 3 {
t.Fatalf("Expected 3 nodes, got %d", numNodes)
}
})
}
func TestClusterResize_RemoveNode(t *testing.T) {
cluster := test.MustRunCluster(t, 3)
defer cluster.Close()

View file

@ -40,7 +40,6 @@ import (
pb "github.com/pilosa/pilosa/v2/proto"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
"github.com/pilosa/pilosa/v2/test/port"
)
func TestHandler_PostSchemaCluster(t *testing.T) {
@ -1405,10 +1404,7 @@ func TestCluster_TranslateStore(t *testing.T) {
),
)
if err := port.GetPort(func(p int) error {
cluster.GetIdleNode(0).Config.Gossip.Port = fmt.Sprintf("%d", p)
return cluster.GetIdleNode(0).Start()
}, 10); err != nil {
if err := cluster.GetIdleNode(0).Start(); err != nil {
t.Fatalf("starting node 0: %v", err)
}
defer cluster.GetIdleNode(0).Close()

View file

@ -20,7 +20,6 @@
package server
import (
"bytes"
"context"
"crypto/tls"
"io"
@ -47,7 +46,6 @@ import (
petcd "github.com/pilosa/pilosa/v2/etcd"
"github.com/pilosa/pilosa/v2/gcnotify"
"github.com/pilosa/pilosa/v2/gopsutil"
"github.com/pilosa/pilosa/v2/gossip"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/logger"
pnet "github.com/pilosa/pilosa/v2/net"
@ -72,10 +70,6 @@ type Command struct {
// Configuration.
Config *Config
// Gossip transport
gossipTransport *gossip.Transport
gossipMemberSet io.Closer
// Standard input/output
*pilosa.CmdIO
@ -84,7 +78,6 @@ type Command struct {
// done will be closed when Command.Close() is called
done chan struct{}
// Passed to the Gossip implementation.
logOutput io.Writer
logger loggerLogger
@ -233,11 +226,6 @@ func (m *Command) UpAndDown() (err error) {
return errors.Wrap(err, "setting up server")
}
// SetupNetworking (so we'll have profiling)
err = m.setupNetworking()
if err != nil {
return errors.Wrap(err, "setting up networking")
}
go func() {
err := m.Handler.Serve()
if err != nil {
@ -469,35 +457,6 @@ func (m *Command) SetupServer() error {
return errors.Wrap(err, "new handler")
}
// setupNetworking sets up internode communication based on the configuration.
func (m *Command) setupNetworking() error {
gossipPort, err := strconv.Atoi(m.Config.Gossip.Port)
if err != nil {
return errors.Wrap(err, "parsing port")
}
// get the host portion of addr to use for binding
gossipHost := m.listenURI.Host
m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger())
if err != nil {
return errors.Wrap(err, "getting transport")
}
gossipMemberSet, err := gossip.NewMemberSet(
m.Config.Gossip,
m.API,
gossip.WithLogOutput(&filteredWriter{logOutput: m.logOutput, v: m.Config.Verbose}),
gossip.WithPilosaLogger(m.logger),
gossip.WithTransport(m.gossipTransport),
)
if err != nil {
return errors.Wrap(err, "getting memberset")
}
m.gossipMemberSet = gossipMemberSet
return errors.Wrap(gossipMemberSet.Open(), "opening gossip memberset")
}
// setupLogger sets up the logger based on the configuration.
func (m *Command) setupLogger() error {
var f *logger.FileWriter
@ -539,13 +498,6 @@ func (m *Command) setupLogger() error {
return nil
}
// GossipTransport allows a caller to return the gossip transport created when
// setting up the GossipMemberSet. This is useful if one needs to determine the
// allocated ephemeral port programmatically. (usually used in tests)
func (m *Command) GossipTransport() *gossip.Transport {
return m.gossipTransport
}
// Close shuts down the server.
func (m *Command) Close() error {
select {
@ -558,9 +510,6 @@ func (m *Command) Close() error {
eg.Go(m.Server.Close)
eg.Go(m.API.Close)
eg.Go(m.pgserver.Close)
if m.gossipMemberSet != nil {
eg.Go(m.gossipMemberSet.Close)
}
if closer, ok := m.logOutput.(io.Closer); ok {
// If closer is os.Stdout or os.Stderr, don't close it.
if closer != os.Stdout && closer != os.Stderr {
@ -617,27 +566,6 @@ func getListener(uri pnet.URI, tlsconf *tls.Config) (ln net.Listener, err error)
return ln, nil
}
type filteredWriter struct {
v bool
logOutput io.Writer
}
// Write forwards the write to logOutput if verbose is true, or it doesn't
// contain [DEBUG] or [INFO]. This implementation isn't technically correct
// since Write could be called with only part of a log line, but I don't think
// that actually happens, so until it becomes a problem, I don't think it's
// worth dealing with the extra complexity. (jaffee)
func (f *filteredWriter) Write(p []byte) (n int, err error) {
if bytes.Contains(p, []byte("[DEBUG]")) || bytes.Contains(p, []byte("[INFO]")) {
if f.v {
return f.logOutput.Write(p)
}
} else {
return f.logOutput.Write(p)
}
return len(p), nil
}
// ParseConfig parses s into a Config.
func ParseConfig(s string) (Config, error) {
var c Config

View file

@ -26,7 +26,6 @@ import (
"os"
"reflect"
"sort"
"strconv"
"strings"
"testing"
"time"
@ -977,8 +976,6 @@ func TestClusterQueriesAfterRestart(t *testing.T) {
config := cmd1.Command.Config
config.Bind = cmd1.API.Node().URI.HostPort()
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cmd1.Command.GossipTransport().URI.Port))
cmd1.Command = server.NewCommand(cmd1.Stdin, cmd1.Stdout, cmd1.Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cmd1.Command.Config = config
err = cmd1.Start()

View file

@ -401,22 +401,6 @@ func (c *Cluster) Start() error {
}()
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))
}
for i, cc := range c.Nodes {
cc := cc
cc.Config.Etcd = portsCfg[i].Etcd
@ -425,14 +409,12 @@ func (c *Cluster) Start() error {
cc.Config.BindGRPC = portsCfg[i].BindGRPC
eg.Go(func() error {
cc.Config.Gossip.Seeds = gossipSeeds
return cc.Start()
})
}
return eg.Wait()
}, 4*len(c.Nodes), 10)
}, 3*len(c.Nodes), 10)
if err != nil {
return err

View file

@ -22,7 +22,6 @@ import (
"time"
"github.com/pilosa/pilosa/v2/etcd"
"github.com/pilosa/pilosa/v2/gossip"
"github.com/pilosa/pilosa/v2/server"
)
@ -33,8 +32,7 @@ type Ports struct {
LsnP *net.TCPListener
PortP int
Grpc int
Gossip int //TODO remove
Grpc int
}
func (ports *Ports) Close() error {
@ -65,10 +63,7 @@ func GenPortsConfig(ports []Ports) []*server.Config {
}
cfgs[i] = &server.Config{
Name: name,
Gossip: gossip.Config{
Port: fmt.Sprint(ports[i].Gossip),
},
Name: name,
BindGRPC: fmt.Sprintf(":%d", ports[i].Grpc),
Etcd: etcd.Options{
Dir: discoDir,
@ -101,20 +96,18 @@ func NewPorts(lsn []*net.TCPListener) []Ports {
ports[i] = lsn[i].Addr().(*net.TCPAddr).Port
}
for i := 0; i < n; i = i + 4 {
for i := 0; i < n; i = i + 3 {
out = append(out, Ports{
LsnC: lsn[i],
PortC: ports[i],
LsnP: lsn[i+1],
PortP: ports[i+1],
Grpc: ports[i+2],
Gossip: ports[i+3],
Grpc: ports[i+2],
})
// make Grpc and Gossip ports available to
// make Grpc port available to
// be rebound.
lsn[i+2].Close()
lsn[i+3].Close()
}
return out

View file

@ -263,17 +263,12 @@ func TestTranslation_Reset(t *testing.T) {
if err := node0.SoftOpen(); err != nil {
t.Fatal(err)
}
gossipSeeds := []string{}
node1.Config.Gossip.Seeds = gossipSeeds
if err := node1.SoftOpen(); err != nil {
t.Fatal(err)
}
node2.Config.Gossip.Seeds = gossipSeeds
if err := node2.SoftOpen(); err != nil {
t.Fatal(err)
}
node3.Config.Gossip.Seeds = gossipSeeds
if err := node3.SoftOpen(); err != nil {
t.Fatal(err)
}