From bc8f7f702f10b89cfb37b3d2072b1e857614127a Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 17 Apr 2017 17:09:52 -0500 Subject: [PATCH] add BroadcastReceiver and collapse Gossip structs into one GossipNodeSet and GossipBroadcaster had a lot of cross dependency - I made them the same object and had it implement all three interfaces (NodeSet, Broadcaster, BroadcastReceiver). I implmented an HTTPBroadcastReceiver which runs as as separate server. BroadcastReceiver is a separate entity on the pilosa.Server object and is started separately from broadcaster and nodeset. In the case of GossipNodeSet, the broadcast receiver must be started before Open()ing the Nodeset, because it doesn't actually start listening until Open() is called, but it needs the handler set up before then. server/server.go was heavily refactored to configure the new structs and interfaces on the pilosa.Server object. --- broadcast.go | 83 +++++++++++++++++++++++++++++++++ gossip.go | 87 ++++++++++++++--------------------- handler.go | 30 ------------ messenger.go | 12 +++-- server.go | 18 +++++--- server/server.go | 117 ++++++++++++++++++++--------------------------- 6 files changed, 186 insertions(+), 161 deletions(-) create mode 100644 broadcast.go 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) {