From 7ec666d912cabd2d35b9c9ad2e00210cc8467004 Mon Sep 17 00:00:00 2001 From: travisturner Date: Thu, 12 Dec 2013 08:50:25 -0600 Subject: [PATCH 1/5] move cruncher out of core and into its own package --- {core => cruncher}/cruncher.go | 4 ++-- {core => cruncher}/cruncher_test.go | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) rename {core => cruncher}/cruncher.go (93%) rename {core => cruncher}/cruncher_test.go (97%) diff --git a/core/cruncher.go b/cruncher/cruncher.go similarity index 93% rename from core/cruncher.go rename to cruncher/cruncher.go index fceda00a9..ceef6a8ce 100644 --- a/core/cruncher.go +++ b/cruncher/cruncher.go @@ -1,4 +1,4 @@ -package core +package cruncher import ( "github.com/davecgh/go-spew/spew" @@ -7,7 +7,7 @@ import ( type Cruncher struct { - close_chan chan bool + close_chan chan bool } func (cruncher *Cruncher) Run(port int) { diff --git a/core/cruncher_test.go b/cruncher/cruncher_test.go similarity index 97% rename from core/cruncher_test.go rename to cruncher/cruncher_test.go index 9d0b42df3..b38ac3b6e 100644 --- a/core/cruncher_test.go +++ b/cruncher/cruncher_test.go @@ -1,4 +1,4 @@ -package core +package cruncher import ( "testing" From 76bd3cc222018ca66748a7da874a0bfb451842c8 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 13 Dec 2013 13:23:10 -0600 Subject: [PATCH 2/5] 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) + }) +} From 9a18d9139948f481c7563ba54f1097ecffa19295 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 13 Dec 2013 13:35:12 -0600 Subject: [PATCH 3/5] Move things around. --- commands/pilosa-cruncher/cruncher.go | 28 ++++------------------------ core/service.go | 6 +----- cruncher/cruncher.go | 27 +++++++++++++++------------ 3 files changed, 20 insertions(+), 41 deletions(-) diff --git a/commands/pilosa-cruncher/cruncher.go b/commands/pilosa-cruncher/cruncher.go index 0efe20441..865431d5b 100644 --- a/commands/pilosa-cruncher/cruncher.go +++ b/commands/pilosa-cruncher/cruncher.go @@ -1,31 +1,11 @@ package main import ( - "pilosa/core" - "pilosa/db" - "flag" - "log" + "pilosa/cruncher" + "pilosa/config" ) -var tcpLoc string -var httpLoc string - -func init() { - flag.StringVar(&tcpLoc, "l", "127.0.0.1:1300", "ip:port to listen on (tcp)") - flag.StringVar(&httpLoc, "h", "127.0.0.1:1400", "ip:port to listen on (http)") - flag.Parse() -} - func main() { - tcp, err := db.NewLocation(tcpLoc) - if err != nil { - log.Fatal("Location not valid:", tcpLoc) - } - http, err := db.NewLocation(httpLoc) - if err != nil { - log.Fatal("Location not valid:", httpLoc) - } - - service := core.NewService(tcp, http) - service.Run() + cruncher := cruncher.NewCruncher() + cruncher.Run(config.GetInt("port_tcp")) } diff --git a/core/service.go b/core/service.go index f3233c879..0c2e461f4 100644 --- a/core/service.go +++ b/core/service.go @@ -9,7 +9,6 @@ import ( "net" "encoding/gob" "pilosa/db" - "pilosa/config" ) type Stats struct { @@ -32,7 +31,6 @@ type Service struct { PortHttp string Etcd *etcd.Client Cluster *db.Cluster - Cruncher *Cruncher TopologyMapper *TopologyMapper ProcessMapper *ProcessMapper ProcessMap *ProcessMap @@ -40,9 +38,8 @@ type Service struct { Dispatcher *Dispatcher } -func NewService(tcp, http *db.Location) *Service { +func NewService() *Service { service := new(Service) - service.Cruncher = new(Cruncher) service.Etcd = etcd.NewClient(nil) service.Cluster = db.NewCluster() service.TopologyMapper = &TopologyMapper{service, "/pilosa/0"} @@ -63,7 +60,6 @@ func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) { func (service *Service) Run() { log.Println("Running service...") go service.TopologyMapper.Run() - go service.Cruncher.Run(config.GetInt("port_tcp")) sigterm, sighup := service.GetSignals() for { diff --git a/cruncher/cruncher.go b/cruncher/cruncher.go index ceef6a8ce..51aee8ecb 100644 --- a/cruncher/cruncher.go +++ b/cruncher/cruncher.go @@ -3,24 +3,27 @@ package cruncher import ( "github.com/davecgh/go-spew/spew" "pilosa/index" + "pilosa/core" ) type Cruncher struct { + core.Service close_chan chan bool } func (cruncher *Cruncher) Run(port int) { - spew.Dump("Cruncher.Run") - spew.Dump(port) - web_api:= index.NewFragmentContainer() -// web_api.AddFragment("general", "25", 0, "AAA-BBB-CCC") -// web_api.AddFragment("general", "25", 1, "AAA-BBB-CCC") -// web_api.AddFragment("general", "25", 2, "AAA-BBB-CCC") - -started:= make(chan bool) -go web_api.RunServer(port , cruncher.close_chan ,started ) -<-started -//server is listening and going - + spew.Dump("Cruncher.Run") + spew.Dump(port) + web_api:= index.NewFragmentContainer() + started:= make(chan bool) + go web_api.RunServer(port , cruncher.close_chan ,started ) + cruncher.Service.Run() + <-started +} + +func NewCruncher() *Cruncher { + service := core.NewService() + cruncher := Cruncher{*service, make(chan bool)} + return &cruncher } From c842a04e029613beb2bd128d1c506514bb836828 Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 13 Dec 2013 13:38:23 -0600 Subject: [PATCH 4/5] comment out transport test for now --- transport/transport_test.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/transport/transport_test.go b/transport/transport_test.go index 966d3f6a9..d68bf216a 100644 --- a/transport/transport_test.go +++ b/transport/transport_test.go @@ -2,17 +2,19 @@ package transport import ( "testing" - "pilosa/db" + //"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) + */ }) } From 37167299cbcb6ec12836a7d19ef0d7a5e285eebf Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 13 Dec 2013 14:06:40 -0600 Subject: [PATCH 5/5] Move things around, lay out framework for transport/dispatch logic --- core/etcd.go | 1 + core/service.go | 26 +++----------------- cruncher/cruncher.go | 4 +++ dispatch/dispatch.go | 30 +++++++++++++++++++++++ core/interfaces.go => interfaces/core.go | 2 +- transport/tcp.go | 31 ++++++++++++++++++++++++ 6 files changed, 71 insertions(+), 23 deletions(-) create mode 100644 dispatch/dispatch.go rename core/interfaces.go => interfaces/core.go (90%) diff --git a/core/etcd.go b/core/etcd.go index 934e4a037..38c35eff3 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -161,6 +161,7 @@ func(self *ProcessMapper) getnode(u *uuid.UUID) *Node { } func (self *ProcessMapper) Run() { + return var modindex uint64 response, err := self.etcd.Get("nodes", false, true) if err != nil { diff --git a/core/service.go b/core/service.go index 0c2e461f4..a298a2ecf 100644 --- a/core/service.go +++ b/core/service.go @@ -6,36 +6,19 @@ import ( "os" "os/signal" "syscall" - "net" - "encoding/gob" "pilosa/db" + "pilosa/interfaces" ) -type Stats struct { - MessageInCount int - MessageOutCount int - MessageProcessedCount int - Uptime int - MemoryUsed int -} - -type Connection struct { - Conn net.Conn - Encoder *gob.Encoder - Decoder *gob.Decoder -} - type Service struct { Stopper - Port string - PortHttp string Etcd *etcd.Client Cluster *db.Cluster TopologyMapper *TopologyMapper ProcessMapper *ProcessMapper ProcessMap *ProcessMap - Transport *Transporter - Dispatcher *Dispatcher + Transport interfaces.Transporter + Dispatch interfaces.Dispatcher } func NewService() *Service { @@ -54,12 +37,11 @@ func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) { signal.Notify(termChan, syscall.SIGINT, syscall.SIGTERM) return termChan, hupChan } -//////////////////////////////////////////////// - func (service *Service) Run() { log.Println("Running service...") go service.TopologyMapper.Run() + go service.ProcessMapper.Run() sigterm, sighup := service.GetSignals() for { diff --git a/cruncher/cruncher.go b/cruncher/cruncher.go index 51aee8ecb..2cec4f4c2 100644 --- a/cruncher/cruncher.go +++ b/cruncher/cruncher.go @@ -4,6 +4,8 @@ import ( "github.com/davecgh/go-spew/spew" "pilosa/index" "pilosa/core" + "pilosa/transport" + "pilosa/dispatch" ) @@ -25,5 +27,7 @@ func (cruncher *Cruncher) Run(port int) { func NewCruncher() *Cruncher { service := core.NewService() cruncher := Cruncher{*service, make(chan bool)} + cruncher.Transport = transport.NewTcpTransport(service) + cruncher.Dispatch = dispatch.NewCruncherDispatch(service) return &cruncher } diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go new file mode 100644 index 000000000..8201fba2f --- /dev/null +++ b/dispatch/dispatch.go @@ -0,0 +1,30 @@ +package dispatch + +import ( + "pilosa/core" + "log" +) + +type CruncherDispatch struct { + service *core.Service +} + +func (self *CruncherDispatch) Init() error { + log.Println("Starting Dispatcher") + return nil +} + +func (self *CruncherDispatch) Close() { + log.Println("Shutting down Dispatcher") +} + +func (self *CruncherDispatch) Run() { + for { + message := self.service.Transport.Receive() + log.Println("Processing ", message) + } +} + +func NewCruncherDispatch(service *core.Service) *CruncherDispatch { + return &CruncherDispatch{service} +} diff --git a/core/interfaces.go b/interfaces/core.go similarity index 90% rename from core/interfaces.go rename to interfaces/core.go index 00935754c..c6acfeb07 100644 --- a/core/interfaces.go +++ b/interfaces/core.go @@ -1,4 +1,4 @@ -package core +package interfaces import ( "pilosa/db" diff --git a/transport/tcp.go b/transport/tcp.go index a5cfc8343..1f80fae04 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -1,2 +1,33 @@ package transport +import ( + "pilosa/core" + "pilosa/db" + "log" +) + +type TcpTransport struct { + port int + inbox chan *db.Message +} + +func (self *TcpTransport) Init() error { + log.Println("Initializing TCP transport") + return nil +} + +func (self *TcpTransport) Close() { + log.Println("Shutting down TCP transport") +} + +func (self *TcpTransport) Send(message *db.Message) { + log.Println("Send", message) +} + +func (self *TcpTransport) Receive() *db.Message { + return <-self.inbox +} + +func NewTcpTransport(service *core.Service) *TcpTransport { + return &TcpTransport{12000, make(chan *db.Message)} //TODO: make port configurable +}