diff --git a/broadcast.go b/broadcast.go new file mode 100644 index 000000000..feaf7f55f --- /dev/null +++ b/broadcast.go @@ -0,0 +1,83 @@ +package pilosa + +import ( + "fmt" + "io" + "io/ioutil" + "net/http" + + "github.com/gogo/protobuf/proto" +) + +// BroadcastHandler is the interface for the pilosa object which knows how to +// handle broadcast messages. (Hint: this is implemented by pilosa.Server) +type BroadcastHandler interface { + ReceiveMessage(pb proto.Message) error +} + +// BroadcastReceiver is the interface for the object which will listen for and +// decode broadcast messages before passing them to pilosa to handle. The +// implementation of this could be an http server which listens for messages, +// gets the protobuf payload, and then passes it to +// BroadcastHandler.ReceiveMessage. +type BroadcastReceiver interface { + // Start starts listening for broadcast messages - it should return + // immediately, spawning a goroutine if necessary. + Start(BroadcastHandler) error +} + +type nopBroadcastReceiver struct{} + +func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil } + +var NopBroadcastReceiver = &nopBroadcastReceiver{} + +type HTTPBroadcastReceiver struct { + port string + handler BroadcastHandler + logOutput io.Writer +} + +func NewHTTPBroadcastReceiver(port string, logOutput io.Writer) *HTTPBroadcastReceiver { + return &HTTPBroadcastReceiver{ + port: port, + logOutput: logOutput, + } +} + +func (rec *HTTPBroadcastReceiver) Start(b BroadcastHandler) error { + rec.handler = b + go func() { + err := http.ListenAndServe(":"+rec.port, rec) + if err != nil { + fmt.Fprintf(rec.logOutput, "Error listening on %v for HTTPBroadcastReceiver: %v\n", ":"+rec.port, err) + } + }() + return nil +} + +func (rec *HTTPBroadcastReceiver) ServeHTTP(w http.ResponseWriter, r *http.Request) { + if r.Header.Get("Content-Type") != "application/x-protobuf" { + http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType) + return + } + + // Read entire body. + body, err := ioutil.ReadAll(r.Body) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + // Unmarshal message to specific proto type. + m, err := UnmarshalMessage(body) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + if err := rec.handler.ReceiveMessage(m); err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } +} diff --git a/gossip.go b/gossip.go index da9c847f5..1224cf69b 100644 --- a/gossip.go +++ b/gossip.go @@ -1,6 +1,7 @@ package pilosa import ( + "fmt" "io" "log" "os" @@ -13,9 +14,15 @@ import ( ) // GossipNodeSet represents a gossip implementation of NodeSet using memberlist +// GossipNodeSet also represents a gossip implementation of pilosa.Broadcaster // GossipNodeSet also represents an implementation of memberlist.Delegate type GossipNodeSet struct { memberlist *memberlist.Memberlist + handler BroadcastHandler + + broadcasts *memberlist.TransmitLimitedQueue + + server *Server config *GossipConfig @@ -23,10 +30,6 @@ type GossipNodeSet struct { LogOutput io.Writer } -func (g *GossipNodeSet) AttachBroadcaster(mb *GossipBroadcaster) { - g.config.memberlistConfig.Delegate = mb -} - func (g *GossipNodeSet) Nodes() []*Node { a := make([]*Node, 0, g.memberlist.NumMembers()) for _, n := range g.memberlist.Members() { @@ -35,7 +38,15 @@ func (g *GossipNodeSet) Nodes() []*Node { return a } +func (g *GossipNodeSet) Start(h BroadcastHandler) error { + g.handler = h + return nil +} + func (g *GossipNodeSet) Open() error { + if g.handler == nil { + return fmt.Errorf("opening GossipNodeSet: you must call Start(pilosa.BroadcastHandler) before calling Open()") + } ml, err := memberlist.Create(g.config.memberlistConfig) if err != nil { return err @@ -48,6 +59,12 @@ func (g *GossipNodeSet) Open() error { if err != nil { return err } + g.broadcasts = &memberlist.TransmitLimitedQueue{ + NumNodes: func() int { + return ml.NumMembers() + }, + RetransmitMult: 3, + } return nil } @@ -64,7 +81,7 @@ type GossipConfig struct { } // NewGossipNodeSet returns a new instance of GossipNodeSet. -func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed string) *GossipNodeSet { +func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed string, s *Server) *GossipNodeSet { g := &GossipNodeSet{ LogOutput: os.Stderr, } @@ -79,25 +96,15 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed g.config.memberlistConfig.BindPort = gossipPort g.config.memberlistConfig.AdvertiseAddr = gossipHost g.config.memberlistConfig.AdvertisePort = gossipPort + g.config.memberlistConfig.Delegate = g + + g.server = s return g } -//////////////////////////////////////////////////////////////// - -// GossipBroadcaster represents a gossip implementation of pilosa.Broadcaster -// GossipBroadcaster also represents an implementation of memberlist.Delegate -type GossipBroadcaster struct { - broadcasts *memberlist.TransmitLimitedQueue - - server *Server - - // The writer for any logging. - LogOutput io.Writer -} - // SendSync implementation of the Broadcaster interface -func (g *GossipBroadcaster) SendSync(pb proto.Message) error { +func (g *GossipNodeSet) SendSync(pb proto.Message) error { msg, err := MarshalMessage(pb) if err != nil { return err @@ -125,7 +132,7 @@ func (g *GossipBroadcaster) SendSync(pb proto.Message) error { } // SendAsync implementation of the Broadcaster interface -func (g *GossipBroadcaster) SendAsync(pb proto.Message) error { +func (g *GossipNodeSet) SendAsync(pb proto.Message) error { msg, err := MarshalMessage(pb) if err != nil { return err @@ -139,19 +146,19 @@ func (g *GossipBroadcaster) SendAsync(pb proto.Message) error { return nil } -func (g *GossipBroadcaster) Receive(pb proto.Message) error { - if err := g.server.ReceiveMessage(pb); err != nil { +func (g *GossipNodeSet) Receive(pb proto.Message) error { + if err := g.handler.ReceiveMessage(pb); err != nil { return err } return nil } // implementation of the memberlist.Delegate interface -func (g *GossipBroadcaster) NodeMeta(limit int) []byte { +func (g *GossipNodeSet) NodeMeta(limit int) []byte { return []byte{} } -func (g *GossipBroadcaster) NotifyMsg(b []byte) { +func (g *GossipNodeSet) NotifyMsg(b []byte) { m, err := UnmarshalMessage(b) if err != nil { g.logger().Printf("unmarshal message error: %s", err) @@ -163,11 +170,11 @@ func (g *GossipBroadcaster) NotifyMsg(b []byte) { } } -func (g *GossipBroadcaster) GetBroadcasts(overhead, limit int) [][]byte { +func (g *GossipNodeSet) GetBroadcasts(overhead, limit int) [][]byte { return g.broadcasts.GetBroadcasts(overhead, limit) } -func (g *GossipBroadcaster) LocalState(join bool) []byte { +func (g *GossipNodeSet) LocalState(join bool) []byte { pb, err := g.server.LocalState() if err != nil { g.logger().Printf("error getting local state, err=%s", err) @@ -183,7 +190,7 @@ func (g *GossipBroadcaster) LocalState(join bool) []byte { return buf } -func (g *GossipBroadcaster) MergeRemoteState(buf []byte, join bool) { +func (g *GossipNodeSet) MergeRemoteState(buf []byte, join bool) { // Unmarshal nodestate data. var pb internal.NodeState if err := proto.Unmarshal(buf, &pb); err != nil { @@ -196,32 +203,6 @@ func (g *GossipBroadcaster) MergeRemoteState(buf []byte, join bool) { } } -// logger returns a logger for the GossipBroadcaster -func (g *GossipBroadcaster) logger() *log.Logger { - return log.New(g.LogOutput, "", log.LstdFlags) -} - -//////////////////////////////////////////////////////////////// - -// NewGossipBroadcaster returns a new instance of GossipBroadcaster. -func NewGossipBroadcaster(s *Server) *GossipBroadcaster { - g := &GossipBroadcaster{ - LogOutput: os.Stderr, - server: s, - } - - g.broadcasts = &memberlist.TransmitLimitedQueue{ - NumNodes: func() int { - return g.server.Cluster.NodeSet.(*GossipNodeSet).memberlist.NumMembers() - }, - RetransmitMult: 3, - } - - return g -} - -//////////////////////////////////////////////////////////////// - // broadcast represents an implementation of memberlist.Broadcast type broadcast struct { msg []byte diff --git a/handler.go b/handler.go index 83f969d22..94b99bbe0 100644 --- a/handler.go +++ b/handler.go @@ -192,36 +192,6 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { } } -// handlePostMessage handles /message requests. -func (h *Handler) handlePostMessage(w http.ResponseWriter, r *http.Request) { - // Verify that request is only communicating over protobufs. - if r.Header.Get("Content-Type") != "application/x-protobuf" { - http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType) - return - } - - // Read entire body. - body, err := ioutil.ReadAll(r.Body) - if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } - - // Unmarshal message to specific proto type. - m, err := UnmarshalMessage(body) - if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } - - if err := h.Server.ReceiveMessage(m); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } - - return -} - func (h *Handler) handleGetSliceMax(w http.ResponseWriter, r *http.Request) { var ms map[string]uint64 if inverse, _ := strconv.ParseBool(r.URL.Query().Get("inverse")); inverse { diff --git a/messenger.go b/messenger.go index 0549e35fc..f0b1e76a6 100644 --- a/messenger.go +++ b/messenger.go @@ -11,6 +11,8 @@ import ( "golang.org/x/sync/errgroup" + "net" + "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" ) @@ -46,11 +48,12 @@ func (c *nopBroadcaster) SendAsync(pb proto.Message) error { // HTTPBroadcaster represents a NodeSet that broadcasts messages over HTTP. type HTTPBroadcaster struct { - server *Server + server *Server + internalPort string } // NewHTTPBroadcaster returns a new instance of HTTPBroadcaster. -func NewHTTPBroadcaster(s *Server) *HTTPBroadcaster { +func NewHTTPBroadcaster(s *Server, internalPort string) *HTTPBroadcaster { return &HTTPBroadcaster{server: s} } @@ -103,11 +106,12 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *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: node.Host, - Path: "/message", + Host: host + ":" + h.internalPort, }).String(), bytes.NewReader(msg)) if err != nil { return err diff --git a/server.go b/server.go index 2f0ace9ce..111e94ee7 100644 --- a/server.go +++ b/server.go @@ -33,9 +33,10 @@ type Server struct { closing chan struct{} // Data storage and HTTP interface. - Index *Index - Handler *Handler - Broadcaster Broadcaster + Index *Index + Handler *Handler + Broadcaster Broadcaster + BroadcastReceiver BroadcastReceiver // Cluster configuration. // Host is replaced with actual host after opening if port is ":0". @@ -54,9 +55,10 @@ func NewServer() *Server { s := &Server{ closing: make(chan struct{}), - Index: NewIndex(), - Handler: NewHandler(), - Broadcaster: NopBroadcaster, + Index: NewIndex(), + Handler: NewHandler(), + Broadcaster: NopBroadcaster, + BroadcastReceiver: NopBroadcastReceiver, AntiEntropyInterval: DefaultAntiEntropyInterval, PollingInterval: DefaultPollingInterval, @@ -100,6 +102,10 @@ func (s *Server) Open() error { return err } + if err := s.BroadcastReceiver.Start(s); err != nil { + return err + } + // Open NodeSet communication if err := s.Cluster.NodeSet.Open(); err != nil { return err diff --git a/server/server.go b/server/server.go index 251bf77a5..9d00eabfb 100644 --- a/server/server.go +++ b/server/server.go @@ -73,6 +73,29 @@ func (m *Command) Run(args ...string) (err error) { m.Config.DataDir = filepath.Join(HomeDir, strings.TrimPrefix(m.Config.DataDir, prefix)) } + // SetupServer + err = m.SetupServer() + if err != nil { + return err + } + + // Initialize server. + if err = m.Server.Open(); err != nil { + return fmt.Errorf("server.Open: %v", err) + } + fmt.Fprintf(m.Stderr, "Listening as http://%s\n", m.Server.Host) + return nil +} + +func (m *Command) SetupServer() error { + cluster := pilosa.NewCluster() + cluster.ReplicaN = m.Config.Cluster.ReplicaN + + for _, hostport := range m.Config.Cluster.Nodes { + cluster.Nodes = append(cluster.Nodes, &pilosa.Node{Host: hostport}) + } + m.Server.Cluster = cluster + // Setup logging output. if m.Config.LogPath == "" { m.Server.LogOutput = m.Stderr @@ -90,93 +113,51 @@ func (m *Command) Run(args ...string) (err error) { m.Server.Index.Stats = pilosa.NewExpvarStatsClient() // Build cluster from config file. + var err error m.Server.Host, err = normalizeHost(m.Config.Host) if err != nil { return err } - m.Server.Broadcaster = PilosaBroadcaster(m.Config, m.Server) - m.Server.Cluster = PilosaCluster(m.Config) - // Associate objects to the Broadcaster based on config. - AssociateBroadcaster(m.Server, m.Config) - - // Set configuration options. - m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval) - - // Initialize server. - if err = m.Server.Open(); err != nil { - return fmt.Errorf("server.Open: %v", err) - } - fmt.Fprintf(m.Stderr, "Listening as http://%s\n", m.Server.Host) - return nil -} - -// PilosaBroadcaster returns a new instance of Broadcaster based on the config. -func PilosaBroadcaster(c *pilosa.Config, server *pilosa.Server) (broadcaster pilosa.Broadcaster) { - switch c.Cluster.BroadcasterType { + switch m.Config.Cluster.BroadcasterType { // TODO change name to something that encompasses broadcasting, receiving broadcasts, and tracking cluster membership case "http": - broadcaster = pilosa.NewHTTPBroadcaster(server) + port := strconv.Itoa(m.Config.Cluster.Gossip.Port) + m.Server.Broadcaster = pilosa.NewHTTPBroadcaster(m.Server, port) + m.Server.BroadcastReceiver = pilosa.NewHTTPBroadcastReceiver(port, m.Stderr) + m.Server.Cluster.NodeSet = pilosa.NewHTTPNodeSet() + m.Server.Cluster.NodeSet.(*pilosa.HTTPNodeSet).Join(m.Server.Cluster.Nodes) case "gossip": - broadcaster = pilosa.NewGossipBroadcaster(server) - case "static": - broadcaster = pilosa.NopBroadcaster - } - return broadcaster -} - -// PilosaCluster returns a new instance of Cluster based on the config. -func PilosaCluster(c *pilosa.Config) *pilosa.Cluster { - cluster := pilosa.NewCluster() - cluster.ReplicaN = c.Cluster.ReplicaN - - for _, hostport := range c.Cluster.Nodes { - cluster.Nodes = append(cluster.Nodes, &pilosa.Node{Host: hostport}) - } - - // Setup a Broadcast (over HTTP) or Gossip NodeSet based on config. - switch c.Cluster.BroadcasterType { - case "http": - cluster.NodeSet = pilosa.NewHTTPNodeSet() - cluster.NodeSet.(*pilosa.HTTPNodeSet).Join(cluster.Nodes) - case "gossip": - gport, err := strconv.Atoi(pilosa.DefaultGossipPort) + gossipPort, err := strconv.Atoi(pilosa.DefaultGossipPort) if err != nil { panic(err) // Atoi on a compile-time constant should never fail. } - gossipPort := gport gossipSeed := pilosa.DefaultHost - if c.Cluster.Gossip.Port != 0 { - gossipPort = c.Cluster.Gossip.Port + if m.Config.Cluster.Gossip.Port != 0 { + gossipPort = m.Config.Cluster.Gossip.Port } - if c.Cluster.Gossip.Seed != "" { - gossipSeed = c.Cluster.Gossip.Seed + if m.Config.Cluster.Gossip.Seed != "" { + gossipSeed = m.Config.Cluster.Gossip.Seed } // get the host portion of addr to use for binding - gossipHost, _, err := net.SplitHostPort(c.Host) + gossipHost, _, err := net.SplitHostPort(m.Config.Host) if err != nil { - gossipHost = c.Host + gossipHost = m.Config.Host } - cluster.NodeSet = pilosa.NewGossipNodeSet(c.Host, gossipHost, gossipPort, gossipSeed) - case "static": - cluster.NodeSet = pilosa.NewStaticNodeSet() + gossipNodeSet := pilosa.NewGossipNodeSet(m.Config.Host, gossipHost, gossipPort, gossipSeed, m.Server) + m.Server.Cluster.NodeSet = gossipNodeSet + m.Server.Broadcaster = gossipNodeSet + m.Server.BroadcastReceiver = gossipNodeSet + case "static", "": + m.Server.Broadcaster = pilosa.NopBroadcaster + m.Server.Cluster.NodeSet = pilosa.NewStaticNodeSet() + m.Server.BroadcastReceiver = pilosa.NopBroadcastReceiver default: - cluster.NodeSet = pilosa.NewStaticNodeSet() + return fmt.Errorf("'%v' is not a supported value for broadcaster type.", m.Config.Cluster.BroadcasterType) } - return cluster -} - -// AssociateBroadcaster allows an implementation to associate objects to the Broadcaster -// after cluster configuration. -func AssociateBroadcaster(s *pilosa.Server, c *pilosa.Config) { - switch c.Cluster.BroadcasterType { - case "http": - // nop - case "gossip": - s.Cluster.NodeSet.(*pilosa.GossipNodeSet).AttachBroadcaster(s.Broadcaster.(*pilosa.GossipBroadcaster)) - case "static": - // nop - } + // Set configuration options. + m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval) + return nil } func normalizeHost(host string) (string, error) {