From 76bd3cc222018ca66748a7da874a0bfb451842c8 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 13 Dec 2013 13:23:10 -0600 Subject: [PATCH] Messaging and interface changes. --- core/etcd.go | 129 ++++++++++++++++- core/http.go | 32 ++--- core/interfaces.go | 18 +++ core/service.go | 274 ++---------------------------------- db/db.go | 6 - transport/http.go | 53 +++++++ transport/inproc.go | 1 + transport/tcp.go | 2 + transport/transport.go | 1 + transport/transport_test.go | 18 +++ 10 files changed, 241 insertions(+), 293 deletions(-) create mode 100644 core/interfaces.go create mode 100644 transport/http.go create mode 100644 transport/inproc.go create mode 100644 transport/tcp.go create mode 100644 transport/transport.go create mode 100644 transport/transport_test.go diff --git a/core/etcd.go b/core/etcd.go index d4d6c8491..934e4a037 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -11,12 +11,12 @@ import ( "errors" ) -type MetaWatcher struct { +type TopologyMapper struct { service *Service namespace string } -func (self *MetaWatcher) Run() { +func (self *TopologyMapper) Run() { log.Println(self.namespace + "/db") resp, err := self.service.Etcd.Get(self.namespace + "/db", false, true) if err != nil { @@ -33,6 +33,7 @@ func (self *MetaWatcher) Run() { stop := make(chan bool) go func() { // TODO: error check and restart watcher + // TODO: use modindex to make sure watch catches everything _, _ = self.service.Etcd.Watch(self.namespace + "/db", 0, true, receiver, stop) }() go func() { @@ -46,11 +47,11 @@ func (self *MetaWatcher) Run() { }() } -func NewMetaWatcher(service *Service, namespace string) *MetaWatcher { - return &MetaWatcher{service, namespace} +func NewTopologyMapper(service *Service, namespace string) *TopologyMapper { + return &TopologyMapper{service, namespace} } -func (self *MetaWatcher) handlenode(node *etcd.Node) error { +func (self *TopologyMapper) handlenode(node *etcd.Node) error { key := node.Key[len(self.namespace)+1:] bits := strings.Split(key, "/") var database *db.Database @@ -123,3 +124,121 @@ func flatten(node *etcd.Node) []*etcd.Node { } return nodes } + +type Node struct { + id *uuid.UUID + ip string + port_tcp int + port_http int +} + +type ProcessMap struct { + nodes []*db.Process +} + +type ProcessMapper struct { + etcd *etcd.Client + nodes []Node + receiver chan *etcd.Response + commands chan *ProcessMapperCommand +} + +func NewProcessMapper(service *Service) *ProcessMapper { + return &ProcessMapper{} +} + +type ProcessMapperCommand struct { + key string +} + +func getKey(input string) string { + bits := strings.Split(input, "/") + return bits[len(bits)-1] +} + +func(self *ProcessMapper) getnode(u *uuid.UUID) *Node { + return new(Node) +} + +func (self *ProcessMapper) Run() { + var modindex uint64 + response, err := self.etcd.Get("nodes", false, true) + if err != nil { + log.Fatal(err) + } + //modindex = response.ModifiedIndex + log.Println(modindex) + nodes := make([]Node, 0) + spew.Dump(response) + for _, noderef := range response.Node.Nodes { + nodestring := getKey(noderef.Key) + u, err := uuid.ParseHex(nodestring) + if err != nil { + log.Fatal("Not a valid UUID: ", nodestring) + } + node := Node{id: u} + for _, prop := range noderef.Nodes { + switch getKey(prop.Key) { + case "port_tcp": + node.port_tcp, _ = strconv.Atoi(prop.Value) + case "port_http": + node.port_http, _ = strconv.Atoi(prop.Value) + case "ip": + node.ip = prop.Value + } + } + nodes = append(nodes, node) + } + + + self.nodes = nodes + spew.Dump(self.nodes) + go func() { + stop := make(chan bool) + _, err := self.etcd.Watch("nodes/", 0, true, self.receiver, stop) + if err != nil { + log.Fatal(err) + } + }() + go func() { + for { + select { + case cmd := <-self.commands: + spew.Dump(cmd) + case response := <-self.receiver: + switch response.Action { + case "set": + bits := strings.Split(response.Node.Key, "/") + if len(bits) != 4 { + log.Fatal("bug in etcd sync or etcd data") + } + //router, err := db.NewLocation(response.Node.Value) + u, err := uuid.ParseHex(bits[2]) + node := self.getnode(u) + if err != nil { + log.Fatal(err) + } + switch bits[3] { + case "port_tcp": + node.port_tcp, _ = strconv.Atoi(response.Node.Value) + case "port_http": + node.port_http, _ = strconv.Atoi(response.Node.Value) + case "ip": + node.ip = response.Node.Value + } + spew.Dump(node) + case "delete": + spew.Dump("delete", response) + + default: + spew.Dump("unhandled", response) + } + } + //_, err = self.etcd.Get("nodes", false, false) + //if err != nil { + // log.Fatal(err) + //} + //spew.Dump(nodes) + } + }() +} diff --git a/core/http.go b/core/http.go index cb4f84be6..fdb2e893c 100644 --- a/core/http.go +++ b/core/http.go @@ -4,7 +4,6 @@ import ( "net/http" "encoding/json" "log" - "strconv" "io/ioutil" "pilosa/db" "pilosa/query" @@ -21,7 +20,7 @@ func (service *Service) HandleMessage(w http.ResponseWriter, r *http.Request) { http.Error(w, "Invalid JSON", http.StatusBadRequest) return } - service.Inbox <- &message + //service.Inbox <- &message } func (service *Service) HandleQuery(w http.ResponseWriter, r *http.Request) { @@ -45,7 +44,8 @@ func (service *Service) HandleStats(w http.ResponseWriter, r *http.Request) { return } encoder := json.NewEncoder(w) - stats := service.GetStats() + //stats := service.GetStats() + stats := "" err := encoder.Encode(stats) if err != nil { log.Fatal("Error encoding stats") @@ -53,18 +53,18 @@ func (service *Service) HandleStats(w http.ResponseWriter, r *http.Request) { } func (service *Service) HandleListen(w http.ResponseWriter, r *http.Request) { - listener := service.NewListener() - encoder := json.NewEncoder(w) - for { - select { - case message := <-listener: - err := encoder.Encode(message) - if err != nil { - log.Println("Error sending message") - return - } - } - } + //listener := service.NewListener() + //encoder := json.NewEncoder(w) + //for { + // select { + // case message := <-listener: + // err := encoder.Encode(message) + // if err != nil { + // log.Println("Error sending message") + // return + // } + // } + //} } func (service *Service) ServeHTTP() { @@ -73,5 +73,5 @@ func (service *Service) ServeHTTP() { http.HandleFunc("/query", service.HandleQuery) http.HandleFunc("/stats", service.HandleStats) http.HandleFunc("/listen", service.HandleListen) - http.ListenAndServe(":" + strconv.Itoa(service.HttpLocation.Port), nil) + //http.ListenAndServe(":" + strconv.Itoa(service.HttpLocation.Port), nil) } diff --git a/core/interfaces.go b/core/interfaces.go new file mode 100644 index 000000000..00935754c --- /dev/null +++ b/core/interfaces.go @@ -0,0 +1,18 @@ +package core + +import ( + "pilosa/db" +) + +type Transporter interface { + Init() error + Close() + Send(*db.Message) + Receive() (*db.Message) +} + +type Dispatcher interface { + Init() error + Close() + Run() +} diff --git a/core/service.go b/core/service.go index ae43e9bce..f3233c879 100644 --- a/core/service.go +++ b/core/service.go @@ -6,15 +6,9 @@ import ( "os" "os/signal" "syscall" - "time" - "sync" "net" "encoding/gob" "pilosa/db" - //"net" - //"flag" - //"encoding/gob" - //"io" "pilosa/config" ) @@ -34,38 +28,25 @@ type Connection struct { type Service struct { Stopper - Handler func(db.Message) - Mailbox chan db.Message - Inbox chan *db.Message Port string PortHttp string Etcd *etcd.Client - Location *db.Location - HttpLocation *db.Location - NodeMap db.NodeMap - NodeMapMutex sync.RWMutex - Outbox chan *db.Envelope - ConnectionMap map[db.Location]*Connection - ConnectionMapMutex sync.RWMutex - Listener net.Listener - ConnectionRegisterChannel chan *PersistentConnection - Stats *Stats Cluster *db.Cluster Cruncher *Cruncher - MetaWatcher *MetaWatcher + TopologyMapper *TopologyMapper + ProcessMapper *ProcessMapper + ProcessMap *ProcessMap + Transport *Transporter + Dispatcher *Dispatcher } func NewService(tcp, http *db.Location) *Service { service := new(Service) - service.Location = tcp - service.HttpLocation = http - service.Outbox = make(chan *db.Envelope) - service.Inbox = make(chan *db.Message) - service.Stats = new(Stats) service.Cruncher = new(Cruncher) service.Etcd = etcd.NewClient(nil) service.Cluster = db.NewCluster() - service.MetaWatcher = &MetaWatcher{service, "/pilosa/0"} + service.TopologyMapper = &TopologyMapper{service, "/pilosa/0"} + service.ProcessMapper = NewProcessMapper(service) return service } @@ -76,242 +57,12 @@ func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) { signal.Notify(termChan, syscall.SIGINT, syscall.SIGTERM) return termChan, hupChan } - -func (service *Service) SendMessage(message *db.Message) error { - if message.Destination == *service.Location { - return service.DeliverMessage(message) - } - router := service.GetRouterLocation(message.Destination) - var dest *db.Location - if router == *service.Location { - dest = &message.Destination - } else { - dest = &router - } - return service.DoSendMessage(message, dest) -} - -func (service *Service) DeliverMessage(message *db.Message) error { - log.Println("do handle of", message) - switch message.Key { - case "ping": - location, ok := message.Data.(db.Location) - if !ok { - log.Println("Invalid ping request!", message) - return nil - } - service.SendMessage(&db.Message{"pong", nil, location}) - } - return nil -} - -func (service *Service) DoSendMessage(message *db.Message, destination *db.Location) error { - log.Println("send message", message, "to", destination) - service.Outbox <- &db.Envelope{message, destination} - return nil -} - -type PersistentConnection struct { - Outbox chan *db.Message - Inbox chan *db.Message - Encoder *gob.Encoder - Decoder *gob.Decoder - Location db.Location - Service *Service - Connection *net.Conn - Identified bool -} - -func NewPersistentConnection(service *Service, location *db.Location) *PersistentConnection { - conn := &PersistentConnection{} - conn.Service = service - conn.Identified = true - conn.Outbox = make(chan *db.Message) - conn.Inbox = make(chan *db.Message) - if location != nil { - conn.Location = *location - } - return conn -} - -func (conn *PersistentConnection) TryConnect() error { - locationString := conn.Location.ToString() - log.Println("Dialing", locationString) - connection, err := net.Dial("tcp", locationString) - if err != nil { - log.Println("Error in dial!") - return err - } - conn.Connection = &connection - conn.Encoder = gob.NewEncoder(connection) - conn.Decoder = gob.NewDecoder(connection) - return nil -} - -func (conn *PersistentConnection) Manage() { - log.Println("Starting manager") - var connectionError chan error - var err error - - go func() { - var message *db.Message - var err error - for { - err = conn.Decoder.Decode(&message) - if err != nil { - connectionError <- err - return - } - conn.Inbox <- message - } - }() - - var message *db.Message - for { - select { - case message = <-conn.Outbox: - log.Println("sending", message) - err = conn.Encoder.Encode(message) - if err != nil { -// if e, ok := err.(*net.OpError); ok { -// if e.Err == syscall.EPIPE { -// // Client disconnected - - log.Println("error sending", err) - go func() { conn.Outbox <- message }() // Resend failed message - if !conn.Identified { - log.Println("Stopping because connection not identified") - return - } - conn.Connect() - } - case message = <-conn.Inbox: - log.Println("receiving", message) - if message.Key == "identify" { - log.Println("Registering connection") - conn.Location = message.Data.(db.Location) - go conn.Service.RegisterConnection(conn) - } else { - conn.Service.Inbox <- message - } - case err = <-connectionError: - log.Println("error receiving", err) - if !conn.Identified { - log.Println("Stopping because connection not identified") - return - } - conn.Connect() - } - } -} - -func (conn *PersistentConnection) Connect() { - log.Println("Connnecting to", conn.Location) - err := conn.TryConnect() - if err != nil { - log.Println(err) - log.Println("Connection failed! Waiting 1 second...") - time.Sleep(time.Second) - // Infinite recursion if node never comes up. Not sure if this is a problem. - conn.Connect() - } else { - log.Println("Send identify message") - conn.Encoder.Encode(db.Message{"identify", *conn.Service.Location, conn.Location}) - //go func() { conn.Outbox <- &db.Message{"identify", *conn.Service.Location, conn.Location} }() - } - log.Println("Connect ending") -} - -type ConnectionMapping map[db.Location]*PersistentConnection - -func (service *Service) RegisterConnection(conn *PersistentConnection) { - service.ConnectionRegisterChannel <- conn -} -func (service *Service) HandleConnections() { - connections := ConnectionMapping{} - - log.Println("Handling connections...") - for { - select { - case envelope := <-service.Outbox: - conn, ok := connections[*envelope.Location] - if !ok { - conn = NewPersistentConnection(service, envelope.Location) - connections[*envelope.Location] = conn - conn.Identified = true - go func () { - conn.Connect() - conn.Manage() - }() - } - // Spawning new goroutine so it doesn't block the main event loop while it's sending. - // This emulates an infinitely buffered channel. - go func() { conn.Outbox <- envelope.Message }() - case conn := <-service.ConnectionRegisterChannel: - connections[conn.Location] = conn - } - } -} - -func (service *Service) GetRouterLocation(node db.Location) db.Location { - service.NodeMapMutex.RLock() - defer service.NodeMapMutex.RUnlock() - location, ok := service.NodeMap[node] - if !ok { - location, ok = service.NodeMap[*service.Location] - if !ok { - log.Fatal("Cannot find router!!!") - } - } - return location -} - -func (service *Service) SetupNetwork() { - log.Println("Setup network") - var err error - locationString := service.Location.ToString() - service.Listener, err = net.Listen("tcp", locationString) - if err != nil { - log.Fatal(err) - } -} -func (service *Service) Serve() { - for { - conn, err := service.Listener.Accept() - if err != nil { - log.Fatal(err) - } - con := NewPersistentConnection(service, nil) - con.Connection = &conn - con.Encoder = gob.NewEncoder(*con.Connection) - con.Decoder = gob.NewDecoder(*con.Connection) - go con.Manage() - } -} - -func (service *Service) GetStats() *Stats { - return service.Stats -} - -func (service *Service) NewListener() chan *db.Message { - ch := make(chan *db.Message) - return ch -} - - //////////////////////////////////////////////// func (service *Service) Run() { log.Println("Running service...") - //go r.SyncEtcd() - //go service.WatchEtcd() - //go service.HandleConnections() - //service.SetupNetwork() - //go service.Serve() - //go service.HandleInbox() - //go service.ServeHTTP() - go service.MetaWatcher.Run() + go service.TopologyMapper.Run() go service.Cruncher.Run(config.GetInt("port_tcp")) sigterm, sighup := service.GetSignals() @@ -327,12 +78,3 @@ func (service *Service) Run() { } } } - -func (service *Service) HandleInbox() { - for { - select { - case message := <-service.Inbox: - log.Println("process", message) - } - } -} diff --git a/db/db.go b/db/db.go index 112794942..02f0d0e58 100644 --- a/db/db.go +++ b/db/db.go @@ -5,9 +5,3 @@ type Message struct { Data interface{} `json:data` Destination Location } - -type Envelope struct { - Message *Message - Location *Location -} - diff --git a/transport/http.go b/transport/http.go new file mode 100644 index 000000000..913031b86 --- /dev/null +++ b/transport/http.go @@ -0,0 +1,53 @@ +package transport + +import ( + "pilosa/db" + "log" +) + +type HttpTransport struct { + port int + outbox chan *db.Message + done chan int +} + +func (trans *HttpTransport) Init() error { + log.Println("Bind to port", trans.port) + trans.done = make(chan int) + go trans.Loop() + return nil +} + +func (trans *HttpTransport) Loop() { + var message *db.Message + for { + select { + case message = <-trans.outbox: + log.Println(message) + case <-trans.done: + return + } + } +} + +func (trans *HttpTransport) Close() { + log.Println("Closing HTTP transport.") + trans.done <- 1 +} + +func (trans *HttpTransport) Send(node string, message *db.Message) error { + log.Println("Send", message, "to", node) + trans.outbox <- message + return nil +} + +func (trans *HttpTransport) Receive() (*db.Message, error) { + return &db.Message{}, nil +} + +func NewHttpTransport(port int) *HttpTransport { + trans := new(HttpTransport) + trans.port = port + trans.outbox = make(chan *db.Message, 10) + return trans +} diff --git a/transport/inproc.go b/transport/inproc.go new file mode 100644 index 000000000..d11d0be55 --- /dev/null +++ b/transport/inproc.go @@ -0,0 +1 @@ +package transport diff --git a/transport/tcp.go b/transport/tcp.go new file mode 100644 index 000000000..a5cfc8343 --- /dev/null +++ b/transport/tcp.go @@ -0,0 +1,2 @@ +package transport + diff --git a/transport/transport.go b/transport/transport.go new file mode 100644 index 000000000..d11d0be55 --- /dev/null +++ b/transport/transport.go @@ -0,0 +1 @@ +package transport diff --git a/transport/transport_test.go b/transport/transport_test.go new file mode 100644 index 000000000..966d3f6a9 --- /dev/null +++ b/transport/transport_test.go @@ -0,0 +1,18 @@ +package transport + +import ( + "testing" + "pilosa/db" + . "github.com/smartystreets/goconvey/convey" +) + +func TestHttpTransport(t *testing.T) { + Convey("Test HTTP transport", t, func() { + var com Transporter + com = NewHttpTransport(9009) + com.Init() + com.Send("derp", &db.Message{"derp", 42, db.Location{}}) + com.Close() + So(1, ShouldEqual, 1) + }) +}