Add support for Node.InternalHost.

This was required for the HTTPBroadcaster to run locally on different InternalPorts.
This commit is contained in:
Travis 2017-04-19 09:09:21 -05:00
parent a11daecfcb
commit b4a69be751
7 changed files with 24 additions and 34 deletions

View file

@ -51,13 +51,11 @@ type nopBroadcaster struct{}
// SendSync A no-op implemenetation of Broadcaster SendSync method.
func (c *nopBroadcaster) SendSync(pb proto.Message) error {
fmt.Println("NOPBroadcaster: SendSync") // TODO remove or log properly?
return nil
}
// SendAsync A no-op implemenetation of Broadcaster SendAsync method.
func (c *nopBroadcaster) SendAsync(pb proto.Message) error {
fmt.Println("NOPBroadcaster: SendAsync") // TODO remove or log properly?
return nil
}

View file

@ -19,7 +19,8 @@ const (
// Node represents a node in the cluster.
type Node struct {
Host string `json:"host"`
Host string `json:"host"`
InternalHost string `json:"internal_host"`
}
// Nodes represents a list of nodes.

View file

@ -77,6 +77,7 @@ on the configured port.`,
flags.StringVarP(&Server.Config.Host, "bind", "b", ":10101", "Default URI on which pilosa should listen.")
flags.IntVarP(&Server.Config.Cluster.ReplicaN, "cluster.replicas", "", 1, "Number of hosts each piece of data should be stored on.")
flags.StringSliceVarP(&Server.Config.Cluster.Nodes, "cluster.hosts", "", []string{}, "Comma separated list of hosts in cluster.")
flags.StringSliceVarP(&Server.Config.Cluster.InternalNodes, "cluster.internal-hosts", "", []string{}, "Comma separated list of hosts in cluster used for internal communication.")
flags.DurationVarP((*time.Duration)(&Server.Config.Cluster.PollingInterval), "cluster.poll-interval", "", time.Minute, "Polling interval for cluster.") // TODO what actually is this?
flags.StringVarP(&Server.Config.Plugins.Path, "plugins.path", "", "", "Path to plugin directory.")
flags.StringVar(&Server.Config.LogPath, "log-path", "", "Log path")

View file

@ -4,10 +4,10 @@ import "time"
const (
// DefaultHost is the default hostname and port to use.
DefaultHost = "localhost"
DefaultPort = "10101"
DefaultClusterType = "static"
DefaultGossipPort = "14000"
DefaultHost = "localhost"
DefaultPort = "10101"
DefaultClusterType = "static"
DefaultInternalPort = "14000"
)
// Config represents the configuration for the command.
@ -19,6 +19,7 @@ type Config struct {
ReplicaN int `toml:"replicas"`
Type string `toml:"type"`
Nodes []string `toml:"hosts"`
InternalNodes []string `toml:"internal-hosts"`
PollingInterval Duration `toml:"polling-interval"`
InternalPort string `toml:"internal-port"`
GossipSeed string `toml:"gossip-seed"`
@ -44,20 +45,11 @@ func NewConfig() *Config {
c.Cluster.Type = DefaultClusterType
c.Cluster.PollingInterval = Duration(DefaultPollingInterval)
c.Cluster.Nodes = []string{}
c.Cluster.InternalNodes = []string{}
c.AntiEntropy.Interval = Duration(DefaultAntiEntropyInterval)
return c
}
// NewConfigForHosts returns a Config object with Config.Cluster.Nodes already
// set up.
func NewConfigForHosts(hosts []string) *Config {
conf := NewConfig()
for _, hostport := range hosts {
conf.Cluster.Nodes = append(conf.Cluster.Nodes, hostport)
}
return conf
}
// Duration is a TOML wrapper type for time.Duration.
type Duration time.Duration

View file

@ -763,7 +763,7 @@ func TestHandler_Fragment_Nodes(t *testing.T) {
h.ServeHTTP(w, r)
if w.Code != http.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code)
} else if w.Body.String() != `[{"host":"host1"},{"host":"host2"}]`+"\n" {
} else if w.Body.String() != `[{"host":"host1","internal_host":""},{"host":"host2","internal_host":""}]`+"\n" {
t.Fatalf("unexpected body: %q", w.Body.String())
}
}

View file

@ -11,8 +11,6 @@ import (
"golang.org/x/sync/errgroup"
"net"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa"
)
@ -25,7 +23,7 @@ type HTTPBroadcaster struct {
// NewHTTPBroadcaster returns a new instance of HTTPBroadcaster.
func NewHTTPBroadcaster(s *pilosa.Server, internalPort string) *HTTPBroadcaster {
return &HTTPBroadcaster{server: s}
return &HTTPBroadcaster{server: s, internalPort: internalPort}
}
// SendSync sends a protobuf message to all nodes simultaneously.
@ -77,16 +75,11 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *pilosa.Node, msg []byte) error {
var client *http.Client
client = http.DefaultClient
host, _, err := net.SplitHostPort(node.Host)
// Create HTTP request.
req, err := http.NewRequest("POST", (&url.URL{
Scheme: "http",
Host: host + ":" + h.internalPort,
Host: node.InternalHost,
}).String(), bytes.NewReader(msg))
if err != nil {
return err
}
// Require protobuf encoding.
req.Header.Set("Content-Type", "application/x-protobuf")

View file

@ -96,6 +96,9 @@ func (m *Command) SetupServer() error {
for _, hostport := range m.Config.Cluster.Nodes {
cluster.Nodes = append(cluster.Nodes, &pilosa.Node{Host: hostport})
}
for i, internalhostport := range m.Config.Cluster.InternalNodes {
cluster.Nodes[i].InternalHost = internalhostport
}
m.Server.Cluster = cluster
// Setup logging output.
@ -120,21 +123,23 @@ func (m *Command) SetupServer() error {
return err
}
// Set internal port (string).
internalPortStr := pilosa.DefaultInternalPort
if m.Config.Cluster.InternalPort != "" {
internalPortStr = m.Config.Cluster.InternalPort
}
switch m.Config.Cluster.Type {
case "http":
m.Server.Broadcaster = httpbroadcast.NewHTTPBroadcaster(m.Server, m.Config.Cluster.InternalPort)
m.Server.BroadcastReceiver = httpbroadcast.NewHTTPBroadcastReceiver(m.Config.Cluster.InternalPort, m.Stderr)
m.Server.Broadcaster = httpbroadcast.NewHTTPBroadcaster(m.Server, internalPortStr)
m.Server.BroadcastReceiver = httpbroadcast.NewHTTPBroadcastReceiver(internalPortStr, m.Stderr)
m.Server.Cluster.NodeSet = httpbroadcast.NewHTTPNodeSet()
err := m.Server.Cluster.NodeSet.(*httpbroadcast.HTTPNodeSet).Join(m.Server.Cluster.Nodes)
if err != nil {
return err
}
case "gossip":
gossipPortStr := pilosa.DefaultGossipPort
if m.Config.Cluster.InternalPort != "" {
gossipPortStr = m.Config.Cluster.InternalPort
}
gossipPort, err := strconv.Atoi(gossipPortStr)
gossipPort, err := strconv.Atoi(internalPortStr)
if err != nil {
return err
}