From d4dd91973c74f8a7e82d018dcdf36461496a658d Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Tue, 11 Feb 2014 18:28:55 -0600 Subject: [PATCH] Add websocket listener for messaging layer. --- core/http.go | 104 +++++++++++++++++++++++++++++++++++++++++------ db/db.go | 16 ++------ deps.json | 10 +++++ transport/tcp.go | 26 ++++++------ 4 files changed, 118 insertions(+), 38 deletions(-) diff --git a/core/http.go b/core/http.go index 820330746..ea3da8ef4 100644 --- a/core/http.go +++ b/core/http.go @@ -7,10 +7,13 @@ import ( "net/http" "pilosa/config" "pilosa/db" + "reflect" "runtime" "strconv" + notify "github.com/bitly/go-notify" "github.com/davecgh/go-spew/spew" + "github.com/gorilla/websocket" "tux21b.org/v1/gocql/uuid" ) @@ -31,6 +34,7 @@ func (self *WebService) Run() { mux.HandleFunc("/stats", self.HandleStats) mux.HandleFunc("/info", self.HandleInfo) mux.HandleFunc("/processes", self.HandleProcesses) + mux.HandleFunc("/listen/ws", self.HandleListenWS) mux.HandleFunc("/listen", self.HandleListen) mux.HandleFunc("/test", self.HandleTest) mux.HandleFunc("/version", self.HandleVersion) @@ -249,16 +253,92 @@ func (self *WebService) HandlePing(w http.ResponseWriter, r *http.Request) { } func (self *WebService) 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 - // } - // } - //} + w.Write([]byte(` + Pilosa - streaming client + + + + + `)) +} + +func (self *WebService) HandleListenWS(w http.ResponseWriter, r *http.Request) { + defer func() { + err := recover() + spew.Dump(err) + }() + ws, err := websocket.Upgrade(w, r, nil, 1024, 1024) + if _, ok := err.(websocket.HandshakeError); ok { + http.Error(w, "Not a websocket handshake", 400) + return + } else if err != nil { + log.Println(err) + return + } + inbox := make(chan interface{}, 10) + outbox := make(chan interface{}, 10) + notify.Start("inbox", inbox) + notify.Start("outbox", outbox) + var obj interface{} + var data interface{} + var inmessage *db.Message + var outmessage *db.Envelope + var host string + for { + select { + case obj = <-outbox: + outmessage = obj.(*db.Envelope) + data = outmessage.Message.Data + host = outmessage.Host.String() + case obj = <-inbox: + inmessage = obj.(*db.Message) + data = inmessage.Data + host = "" + } + + typ := reflect.TypeOf(data) + + err := ws.WriteJSON(map[string]interface{}{ + "type": typ.String(), + "dump": spew.Sdump(data), + "host": host, + }) + if err != nil { + log.Println("stopping") + notify.Stop("inbox", inbox) + notify.Stop("outbox", outbox) + drain(inbox) + drain(outbox) + return + } + } +} + +func drain(ch chan interface{}) { + for { + switch { + case <-ch: + default: + return + } + } } diff --git a/db/db.go b/db/db.go index 402cf4018..57139433a 100644 --- a/db/db.go +++ b/db/db.go @@ -10,21 +10,11 @@ type Message struct { Data interface{} `json:data` } -/* -import "pilosa/core" - -type Message interface { - Handle(*core.Service) +type Envelope struct { + Message *Message + Host *uuid.UUID } - -type Message struct { - Key string `json:key` - Data interface{} `json:data` - Destination Location -} -*/ - type HoldResult interface { ResultId() *uuid.UUID ResultData() interface{} diff --git a/deps.json b/deps.json index b3fad2633..db940d04e 100644 --- a/deps.json +++ b/deps.json @@ -63,5 +63,15 @@ "repo": "github.com/davecgh/go-spew/spew", "version": "e762b3d1320b76030bd7f6cc2bfc3d9acce874c0", "type": "git" + }, + "notify": { + "repo": "github.com/bitly/go-notify", + "version": "0a148b8111d688ba7550fc7119fe0d5d8e650838", + "type": "git" + }, + "websocket": { + "repo": "github.com/gorilla/websocket", + "version": "92334662baa9cbebc2e6e68b8d56bc1233f85a4c", + "type": "git" } } diff --git a/transport/tcp.go b/transport/tcp.go index aa1632ad5..6e7032cc9 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -10,14 +10,10 @@ import ( "pilosa/db" "time" + notify "github.com/bitly/go-notify" "tux21b.org/v1/gocql/uuid" ) -type envelope struct { - message *db.Message - host *uuid.UUID -} - type connection struct { transport *TcpTransport inbox chan *db.Message @@ -106,7 +102,7 @@ type TcpTransport struct { service *core.Service port int inbox chan *db.Message - outbox chan envelope + outbox chan db.Envelope connections map[uuid.UUID]*connection reg chan *newconnection } @@ -117,13 +113,13 @@ func (self *TcpTransport) Run() { for { select { case env := <-self.outbox: - con, ok := self.connections[*(env.host)] + con, ok := self.connections[*(env.Host)] if !ok { - con = &connection{self, make(chan *db.Message, 100), make(chan *db.Message, 100), nil, env.host} + con = &connection{self, make(chan *db.Message, 100), make(chan *db.Message, 100), nil, env.Host} go con.manage() - self.connections[*env.host] = con + self.connections[*env.Host] = con } - con.outbox <- env.message + con.outbox <- env.Message case nc := <-self.reg: self.connections[*nc.id] = nc.connection } @@ -157,11 +153,15 @@ func (self *TcpTransport) Close() { } func (self *TcpTransport) Send(message *db.Message, host *uuid.UUID) { - self.outbox <- envelope{message, host} + envelope := db.Envelope{message, host} + notify.Post("outbox", &envelope) + self.outbox <- envelope } func (self *TcpTransport) Receive() *db.Message { - return <-self.inbox + message := <-self.inbox + notify.Post("inbox", message) + return message } func (self *TcpTransport) Push(message *db.Message) { @@ -169,5 +169,5 @@ func (self *TcpTransport) Push(message *db.Message) { } func NewTcpTransport(service *core.Service) *TcpTransport { - return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message, 100), make(chan envelope, 100), make(map[uuid.UUID]*connection), make(chan *newconnection)} + return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message, 100), make(chan db.Envelope, 100), make(map[uuid.UUID]*connection), make(chan *newconnection)} }