From b4a69be751f20f8e5c076e2b58af0088eb051699 Mon Sep 17 00:00:00 2001 From: Travis Date: Wed, 19 Apr 2017 09:09:21 -0500 Subject: [PATCH] Add support for Node.InternalHost. This was required for the HTTPBroadcaster to run locally on different InternalPorts. --- broadcast.go | 2 -- cluster.go | 3 ++- cmd/server.go | 1 + config.go | 20 ++++++-------------- handler_test.go | 2 +- httpbroadcast/messenger.go | 11 ++--------- server/server.go | 19 ++++++++++++------- 7 files changed, 24 insertions(+), 34 deletions(-) diff --git a/broadcast.go b/broadcast.go index 6eaf834cc..a3f3cb6c9 100644 --- a/broadcast.go +++ b/broadcast.go @@ -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 } diff --git a/cluster.go b/cluster.go index d8f4401e3..6290f0050 100644 --- a/cluster.go +++ b/cluster.go @@ -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. diff --git a/cmd/server.go b/cmd/server.go index 51d92824a..d4e72a411 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -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") diff --git a/config.go b/config.go index 22750d72d..b577f6281 100644 --- a/config.go +++ b/config.go @@ -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 diff --git a/handler_test.go b/handler_test.go index 51f48fef7..2611da50e 100644 --- a/handler_test.go +++ b/handler_test.go @@ -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()) } } diff --git a/httpbroadcast/messenger.go b/httpbroadcast/messenger.go index 1545e3328..ebd4b2e6a 100644 --- a/httpbroadcast/messenger.go +++ b/httpbroadcast/messenger.go @@ -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") diff --git a/server/server.go b/server/server.go index 1bf3fe5ce..574941c42 100644 --- a/server/server.go +++ b/server/server.go @@ -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 }