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 +}