mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 12:27:52 +00:00
unpushed WIP on cluster-tests
This commit is contained in:
parent
05211bdf70
commit
8db7a78a93
7 changed files with 113 additions and 14 deletions
|
|
@ -205,8 +205,18 @@ func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*mem
|
|||
conf.Name = api.Node().ID
|
||||
conf.BindAddr = api.Node().URI.Host
|
||||
conf.BindPort = port
|
||||
if cfg.AdvertisePort != "" {
|
||||
port, err = strconv.Atoi(cfg.Port)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("convert advertise port: %s", err)
|
||||
}
|
||||
}
|
||||
conf.AdvertisePort = port
|
||||
conf.AdvertiseAddr = hostToIP(api.Node().URI.Host)
|
||||
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
|
||||
|
|
@ -447,10 +457,20 @@ func newTransport(conf *memberlist.Config) (*memberlist.NetTransport, error) {
|
|||
|
||||
// Config holds toml-friendly memberlist configuration.
|
||||
type Config struct {
|
||||
// Host is the host gossip will bind to. If left blank it will be set to the
|
||||
// host from Pilosa.
|
||||
Host string `toml:"host"`
|
||||
// Port indicates the port to which pilosa should bind for internal state sharing.
|
||||
Port string `toml:"port"`
|
||||
Seeds []string `toml:"seeds"`
|
||||
Key string `toml:"key"`
|
||||
Port string `toml:"port"`
|
||||
// AdvertiseHost is the hostname or IP other nodes should use to connect to
|
||||
// this host. If left blank, the value for Host will be used. This is useful
|
||||
// in some proxy and NAT scenarios.
|
||||
AdvertiseHost string `toml:"advertise-host`
|
||||
// AdvertisePort is the port other nodes will use to connect to this one.
|
||||
// Behaves like AdvertiseHost.
|
||||
AdvertisePort string `toml:"advertise-port"`
|
||||
Seeds []string `toml:"seeds"`
|
||||
Key string `toml:"key"`
|
||||
// StreamTimeout is the timeout for establishing a stream connection with
|
||||
// a remote node for a full state sync, and for stream read and write
|
||||
// operations. Maps to memberlist TCPTimeout.
|
||||
|
|
|
|||
16
server.go
16
server.go
|
|
@ -62,6 +62,7 @@ type Server struct { // nolint: maligned
|
|||
|
||||
nodeID string
|
||||
uri URI
|
||||
advertiseURI URI
|
||||
antiEntropyInterval time.Duration
|
||||
metricInterval time.Duration
|
||||
diagnosticInterval time.Duration
|
||||
|
|
@ -73,7 +74,7 @@ type Server struct { // nolint: maligned
|
|||
dataDir string
|
||||
}
|
||||
|
||||
// TODO: have this return an interface for Holder instead of concrete object?
|
||||
// TODO (2.0): have this return an interface for Holder instead of concrete object?
|
||||
func (s *Server) Holder() *Holder {
|
||||
return s.holder
|
||||
}
|
||||
|
|
@ -160,6 +161,13 @@ func OptServerInternalClient(c InternalClient) ServerOption {
|
|||
}
|
||||
}
|
||||
|
||||
func OptServerAdvertiseURI(u *URI) ServerOption {
|
||||
return func(s *Server) error {
|
||||
s.advertiseURI = *u
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// DEPRECATED
|
||||
func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption {
|
||||
return func(s *Server) error {
|
||||
|
|
@ -296,7 +304,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
|
|||
// Set Cluster Node.
|
||||
node := &Node{
|
||||
ID: s.nodeID,
|
||||
URI: s.uri,
|
||||
URI: s.advertiseURI,
|
||||
IsCoordinator: s.cluster.Coordinator == s.nodeID,
|
||||
}
|
||||
s.cluster.Node = node
|
||||
|
|
@ -581,7 +589,7 @@ func (s *Server) SendSync(m Message) error {
|
|||
node := node
|
||||
s.logger.Printf("SendSync to: %s", node.URI)
|
||||
// Don't forward the message to ourselves.
|
||||
if s.uri == node.URI {
|
||||
if s.advertiseURI == node.URI {
|
||||
continue
|
||||
}
|
||||
|
||||
|
|
@ -676,7 +684,7 @@ func (s *Server) monitorDiagnostics() {
|
|||
|
||||
s.diagnostics.Logger = s.logger
|
||||
s.diagnostics.SetVersion(Version)
|
||||
s.diagnostics.Set("Host", s.uri.Host)
|
||||
s.diagnostics.Set("Host", s.advertiseURI.Host)
|
||||
s.diagnostics.Set("Cluster", strings.Join(s.cluster.nodeIDs(), ","))
|
||||
s.diagnostics.Set("NumNodes", len(s.cluster.nodes))
|
||||
s.diagnostics.Set("NumCPU", runtime.NumCPU())
|
||||
|
|
|
|||
|
|
@ -36,9 +36,15 @@ type Config struct {
|
|||
// DataDir is the directory where Pilosa stores both indexed data and
|
||||
// running state such as cluster topology information.
|
||||
DataDir string `toml:"data-dir"`
|
||||
|
||||
// Bind is the host:port on which Pilosa will listen.
|
||||
Bind string `toml:"bind"`
|
||||
|
||||
// Advertise is the host:port that this node will report as its address to
|
||||
// others. If left blank (the default), this will be set to the bind address
|
||||
// once it is listening.
|
||||
Advertise string `toml:"advertise"`
|
||||
|
||||
// MaxWritesPerRequest limits the number of mutating commands that can be in
|
||||
// a single request to the server. This includes Set, Clear,
|
||||
// SetRowAttrs & SetColumnAttrs.
|
||||
|
|
|
|||
|
|
@ -248,6 +248,11 @@ func (m *Command) SetupServer() error {
|
|||
uri.SetPort(uint16(m.ln.Addr().(*net.TCPAddr).Port))
|
||||
}
|
||||
|
||||
advertURI, err := pilosa.NewURIFromAddressWithDefault(m.Config.Advertise, uri)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "processing avertise address '%s'", m.Config.Advertise)
|
||||
}
|
||||
|
||||
c := http.GetHTTPClient(TLSConfig)
|
||||
|
||||
// Primary store configuration is handled automatically now.
|
||||
|
|
@ -276,6 +281,7 @@ func (m *Command) SetupServer() error {
|
|||
pilosa.OptServerGCNotifier(gcnotify.NewActiveGCNotifier()),
|
||||
pilosa.OptServerStatsClient(statsClient),
|
||||
pilosa.OptServerURI(uri),
|
||||
pilosa.OptServerAdvertiseURI(advertURI),
|
||||
pilosa.OptServerInternalClient(http.NewInternalClientFromURI(uri, c)),
|
||||
pilosa.OptServerPrimaryTranslateStoreFunc(http.NewTranslateStore),
|
||||
pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts),
|
||||
|
|
|
|||
|
|
@ -405,7 +405,6 @@ func TestClusteringNodesReplica1(t *testing.T) {
|
|||
// Create new main with the same config.
|
||||
config := cluster[2].Command.Config
|
||||
config.Translation.MapSize = 100000
|
||||
// config.Bind = cluster[2].API.Node().URI.HostPort()
|
||||
|
||||
// this isn't necessary, but makes the test run way faster
|
||||
config.Gossip.Port = strconv.Itoa(int(cluster[2].Command.GossipTransport().URI.Port))
|
||||
|
|
@ -413,6 +412,8 @@ func TestClusteringNodesReplica1(t *testing.T) {
|
|||
cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr)
|
||||
cluster[2].Command.Config = config
|
||||
|
||||
time.Sleep(time.Second * 40)
|
||||
|
||||
// Run new program.
|
||||
if err := cluster[2].Start(); err != nil {
|
||||
t.Fatalf("restarting node 2: %v", err)
|
||||
|
|
|
|||
|
|
@ -19,11 +19,13 @@ import (
|
|||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"math/rand"
|
||||
gohttp "net/http"
|
||||
"os"
|
||||
"path"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
|
|
@ -31,8 +33,37 @@ import (
|
|||
"github.com/pilosa/pilosa/http"
|
||||
"github.com/pilosa/pilosa/server"
|
||||
"github.com/pkg/errors"
|
||||
|
||||
"github.com/Shopify/toxiproxy"
|
||||
tox "github.com/Shopify/toxiproxy/client"
|
||||
)
|
||||
|
||||
type portAllocator struct {
|
||||
port *uint32
|
||||
}
|
||||
|
||||
var ports *portAllocator
|
||||
|
||||
// Next generates a new port and returns it
|
||||
func (n *portAllocator) Next() (nextPort uint32) {
|
||||
return atomic.AddUint32(n.port, 1)
|
||||
}
|
||||
|
||||
var proxy string
|
||||
|
||||
func init() {
|
||||
r := rand.New(rand.NewSource(int64(time.Now().Nanosecond())))
|
||||
p := uint32(r.Uint32()%40000 + 20000)
|
||||
ports = &portAllocator{
|
||||
port: &p,
|
||||
}
|
||||
|
||||
tport := strconv.Itoa(int(ports.Next()))
|
||||
proxy = "localhost:" + tport
|
||||
go toxiproxy.NewServer().Listen("localhost", tport)
|
||||
|
||||
}
|
||||
|
||||
////////////////////////////////////////////////////////////////////////////////////
|
||||
// Command represents a test wrapper for server.Command.
|
||||
type Command struct {
|
||||
|
|
@ -232,6 +263,8 @@ func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu
|
|||
|
||||
// newCluster creates a new cluster
|
||||
func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
||||
tclient := tox.NewClient(proxy)
|
||||
|
||||
if size == 0 {
|
||||
return nil, errors.New("cluster must contain at least one node")
|
||||
}
|
||||
|
|
@ -245,8 +278,24 @@ func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
|
|||
if len(opts) > 0 {
|
||||
commandOpts = opts[i%len(opts)]
|
||||
}
|
||||
aport := strconv.Itoa(int(ports.Next()))
|
||||
name := "node" + strconv.Itoa(i)
|
||||
m := NewCommandNode(i == 0, commandOpts...)
|
||||
err := ioutil.WriteFile(path.Join(m.Config.DataDir, ".id"), []byte("node"+strconv.Itoa(i)), 0600)
|
||||
m.Config.Bind = "localhost:" + strconv.Itoa(int(ports.Next()))
|
||||
m.Config.Advertise = "localhost:" + aport
|
||||
p, err := tclient.CreateProxy(name+aport, m.Config.Advertise, m.Config.Bind)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "setting up toxiproxy")
|
||||
}
|
||||
m.Config.Gossip.Port = strconv.Itoa(int(ports.Next()))
|
||||
aport = strconv.Itoa(int(ports.Next()))
|
||||
m.Config.Gossip.AdvertisePort = aport
|
||||
p, err = tclient.CreateProxy(name+"-gossip"+aport, "localhost:"+m.Config.Gossip.AdvertisePort, "localhost:"+m.Config.Gossip.Port)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "setting up toxiproxy for gossip")
|
||||
}
|
||||
fmt.Println(p)
|
||||
err = ioutil.WriteFile(path.Join(m.Config.DataDir, ".id"), []byte(name), 0600)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "writing node id")
|
||||
}
|
||||
|
|
|
|||
17
uri.go
17
uri.go
|
|
@ -82,6 +82,10 @@ func NewURIFromAddress(address string) (*URI, error) {
|
|||
return parseAddress(address)
|
||||
}
|
||||
|
||||
func NewURIFromAddressWithDefault(address string, base *URI) (*URI, error) {
|
||||
return parseAddressWithDefault(address, base)
|
||||
}
|
||||
|
||||
// setScheme sets the scheme of this URI.
|
||||
func (u *URI) setScheme(scheme string) error {
|
||||
m := schemeRegexp.FindStringSubmatch(scheme)
|
||||
|
|
@ -154,20 +158,20 @@ func (u URI) Type() string {
|
|||
return "URI"
|
||||
}
|
||||
|
||||
func parseAddress(address string) (uri *URI, err error) {
|
||||
func parseAddressWithDefault(address string, def *URI) (uri *URI, err error) {
|
||||
m := addressRegexp.FindStringSubmatch(address)
|
||||
if m == nil {
|
||||
return nil, errors.New("invalid address")
|
||||
}
|
||||
scheme := "http"
|
||||
scheme := def.Scheme
|
||||
if m[2] != "" {
|
||||
scheme = m[2]
|
||||
}
|
||||
host := "localhost"
|
||||
host := def.Host
|
||||
if m[3] != "" {
|
||||
host = m[3]
|
||||
}
|
||||
var port = 10101
|
||||
var port = int(def.Port)
|
||||
if m[5] != "" {
|
||||
port, err = strconv.Atoi(m[5])
|
||||
if err != nil {
|
||||
|
|
@ -185,6 +189,11 @@ func parseAddress(address string) (uri *URI, err error) {
|
|||
return uri, nil
|
||||
}
|
||||
|
||||
func parseAddress(address string) (uri *URI, err error) {
|
||||
u, err := parseAddressWithDefault(address, defaultURI())
|
||||
return u, err
|
||||
}
|
||||
|
||||
// MarshalJSON marshals URI into a JSON-encoded byte slice.
|
||||
func (u *URI) MarshalJSON() ([]byte, error) {
|
||||
var output struct {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue