bangbang theory

This commit is contained in:
Kuba Podgórski 2021-01-14 16:25:44 +01:00
parent 1a5ab4b155
commit 1f6ed4feca
2 changed files with 33 additions and 50 deletions

View file

@ -248,41 +248,32 @@ func (c *Cluster) Start() error {
err := port.GetPorts(func(ports []int) error {
portsCfg := GenPortsConfig(NewPorts(ports))
// seedCh is a channel of host:port values to use
// as gossip seeds during startup.
seedCh := make(chan string, len(c.Nodes))
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.DisCo = portsCfg[i].DisCo
cc.Config.BindGRPC = portsCfg[i].BindGRPC
eg.Go(func() error {
// 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")
}
fmt.Printf("DISCO CONFIG: %+v\n", cc.Config.DisCo)
cc.Config.Gossip.Seeds = gossipSeeds
cc.Config.Gossip.Port = portsCfg[i].Gossip.Port
gossipHost := uri.Host
gossipPort := cc.Config.Gossip.Port
if gossipPort == "0" || gossipPort == "" {
panic("gossipPort not allowed to be 0!")
}
println("gossipPort is ", gossipPort)
// the first node doesn't need to wait for a seed.
if i > 0 {
x := <-seedCh
cc.Config.Gossip.Seeds = []string{x}
}
seedCh <- fmt.Sprintf("%s:%s", gossipHost, gossipPort)
if err := cc.Start(); err != nil {
return errors.Wrapf(err, "starting server %d", i)
}
return nil
return cc.Start()
})
// fixes race on gossip: time.Sleep(time.Second)
}
@ -421,20 +412,9 @@ func newCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (*Clust
// MustRunCluster creates and starts a new cluster. The opts parameter
// is slightly magical; see MustNewCluster.
func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) *Cluster {
var tries int = 5
var cluster *Cluster
var err error
for i := 0; i < tries; i++ {
if i > 0 {
fmt.Printf("--- try starting cluster again: %d\n", i)
}
cluster = MustNewCluster(tb, size, opts...)
if err = cluster.Start(); err == nil {
break
}
}
cluster := MustNewCluster(tb, size, opts...)
err := cluster.Start()
if err != nil {
tb.Fatalf("run cluster: %v", err)
}

View file

@ -17,6 +17,7 @@ package test
import (
"fmt"
"strings"
"time"
"github.com/pilosa/pilosa/v2/etcd"
"github.com/pilosa/pilosa/v2/gossip"
@ -46,18 +47,20 @@ func GenPortsConfig(ports []Ports) []*server.Config {
},
BindGRPC: port.ColonZeroString(ports[i].Grpc),
DisCo: etcd.Options{
Name: name,
Dir: "",
ClusterName: "bartholemuuuuu",
LClientURL: lClientURL,
AClientURL: lClientURL,
LPeerURL: lPeerURL,
APeerURL: lPeerURL,
Name: name,
Dir: "",
ClusterName: "bartholemuuuuu",
LClientURL: lClientURL,
AClientURL: lClientURL,
LPeerURL: lPeerURL,
APeerURL: lPeerURL,
HeartbeatTTL: 5 * int64(time.Second),
},
}
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)
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)
}
for i := range cfgs {
cfgs[i].DisCo.InitCluster = strings.Join(clusterURLs, ",")