From d4dd91973c74f8a7e82d018dcdf36461496a658d Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Tue, 11 Feb 2014 18:28:55 -0600 Subject: [PATCH 1/8] 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)} } From 1e932aca9a1cdaf573f243401fcb31d75674ec20 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 12 Feb 2014 18:02:32 -0600 Subject: [PATCH 2/8] Fix race condition, add flags. --- commands/pilosa-nexter/nexter.go | 46 ++++++++++++++++++++------------ 1 file changed, 29 insertions(+), 17 deletions(-) diff --git a/commands/pilosa-nexter/nexter.go b/commands/pilosa-nexter/nexter.go index 31df74c2c..d5f735dbc 100644 --- a/commands/pilosa-nexter/nexter.go +++ b/commands/pilosa-nexter/nexter.go @@ -2,15 +2,24 @@ package main import ( "encoding/json" - "github.com/coreos/go-etcd/etcd" + "flag" "log" "net/http" "strconv" "strings" "time" + + "github.com/coreos/go-etcd/etcd" ) -const blocksize = 64 +var port uint +var blocksize uint64 + +func init() { + flag.UintVar(&port, "port", 9000, "Port to run HTTP server on") + flag.Uint64Var(&blocksize, "blocksize", 64, "Block size") + flag.Parse() +} type Req interface{} @@ -36,30 +45,32 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { for { node, err := client.Get(path, false, false) if err != nil { - ee, ok := err.(etcd.EtcdError) + ee, ok := err.(*etcd.EtcdError) if ok && ee.ErrorCode == 100 { // node does not exist - /* start = 0 - end = blocksize - 1 - client.Set(path, "0", 0) // TODO: error check - _ , err := client.CompareAndSwap(path, "0", 0, "0", 0) - - */ - //_,_ := c.put(path, "0", 0, options) - client.RawCreate(path, "0", 0) + _, err := client.Create(path, "0", 0) + if err != nil { + ee, ok := err.(*etcd.EtcdError) + // Catch race condition where another node did the same client.Create() + if ok && ee.ErrorCode == 105 { // Node has been created + continue + } else { + log.Fatal(err) + } + } continue } else { log.Fatal(err) } } else { // No error, get start of series from etcd node start, err = strconv.ParseUint(node.Node.Value, 10, 0) - end = start + blocksize if err != nil { log.Fatal(err) } + end = start + blocksize } for { newval, err := client.CompareAndSwap(path, strconv.FormatUint(end, 10), 0, strconv.FormatUint(start, 10), 0) - if err != nil { + if err == nil { break } else { log.Println("Error with CompareAndSet! Trying again in 1 second...") @@ -84,7 +95,7 @@ func (self *Nexter) loop() { select { case req := <-self.reqchan: switch req.(type) { - case IncReq: + case *IncReq: countreq := req.(*IncReq) counter, ok := counters[countreq.id] if !ok { @@ -93,12 +104,12 @@ func (self *Nexter) loop() { go self.countloop(counter, countreq.id, client) } go func() { countreq.ret <- <-counter }() - case DelReq: + case *DelReq: delreq := req.(*DelReq) delete(counters, delreq.id) path := "nexter/" + strconv.Itoa(delreq.id) client.Delete(path, true) - + default: } case <-self.done: break @@ -145,5 +156,6 @@ func main() { dec := json.NewEncoder(w) dec.Encode(num) }) - http.ListenAndServe(":9000", nil) + port_string := strconv.FormatUint(uint64(port), 10) + http.ListenAndServe(":"+port_string, nil) } From dec92588b58e5043a28da211999648d4a5ed4926 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 12 Feb 2014 18:43:43 -0600 Subject: [PATCH 3/8] Add listen streamer that works with curl -N --- core/http.go | 47 +++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 47 insertions(+) diff --git a/core/http.go b/core/http.go index ea3da8ef4..c0a0b9bdf 100644 --- a/core/http.go +++ b/core/http.go @@ -35,6 +35,7 @@ func (self *WebService) Run() { mux.HandleFunc("/info", self.HandleInfo) mux.HandleFunc("/processes", self.HandleProcesses) mux.HandleFunc("/listen/ws", self.HandleListenWS) + mux.HandleFunc("/listen/stream", self.HandleListenStream) mux.HandleFunc("/listen", self.HandleListen) mux.HandleFunc("/test", self.HandleTest) mux.HandleFunc("/version", self.HandleVersion) @@ -342,3 +343,49 @@ func drain(ch chan interface{}) { } } } + +func (self *WebService) HandleListenStream(w http.ResponseWriter, r *http.Request) { + defer func() { + err := recover() + spew.Dump(err) + }() + 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) + + writer := json.NewEncoder(w) + err := writer.Encode(map[string]interface{}{ + "type": typ.String(), + "dump": spew.Sdump(data), + "host": host, + }) + w.(http.Flusher).Flush() + if err != nil { + log.Println("stopping") + notify.Stop("inbox", inbox) + notify.Stop("outbox", outbox) + drain(inbox) + drain(outbox) + return + } + } +} From 1574341dbabddf31584a1ff9ac230a7009b66f00 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Thu, 13 Feb 2014 14:54:39 -0600 Subject: [PATCH 4/8] Add etcd node flag to nexter --- commands/pilosa-nexter/nexter.go | 16 +++++++++++++++- 1 file changed, 15 insertions(+), 1 deletion(-) diff --git a/commands/pilosa-nexter/nexter.go b/commands/pilosa-nexter/nexter.go index d5f735dbc..a7a6d3d77 100644 --- a/commands/pilosa-nexter/nexter.go +++ b/commands/pilosa-nexter/nexter.go @@ -3,6 +3,7 @@ package main import ( "encoding/json" "flag" + "fmt" "log" "net/http" "strconv" @@ -12,12 +13,25 @@ import ( "github.com/coreos/go-etcd/etcd" ) +type stringslice []string + +func (self *stringslice) String() string { + return fmt.Sprintf("%d", *self) +} + +func (self *stringslice) Set(value string) error { + *self = append(*self, value) + return nil +} + var port uint var blocksize uint64 +var etcd_nodes stringslice func init() { flag.UintVar(&port, "port", 9000, "Port to run HTTP server on") flag.Uint64Var(&blocksize, "blocksize", 64, "Block size") + flag.Var(&etcd_nodes, "etcd", "Etcd server") flag.Parse() } @@ -90,7 +104,7 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { func (self *Nexter) loop() { counters := make(map[int]chan uint64) - client := etcd.NewClient(nil) + client := etcd.NewClient(etcd_nodes) for { select { case req := <-self.reqchan: From 9c9ba2a8a8d709bc51f90780ef7b294f8e3ec3a6 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 14 Feb 2014 11:29:09 -0600 Subject: [PATCH 5/8] Simplify logic, fix potential nil pointer dereference --- commands/pilosa-nexter/nexter.go | 18 +++++------------- 1 file changed, 5 insertions(+), 13 deletions(-) diff --git a/commands/pilosa-nexter/nexter.go b/commands/pilosa-nexter/nexter.go index a7a6d3d77..22e360756 100644 --- a/commands/pilosa-nexter/nexter.go +++ b/commands/pilosa-nexter/nexter.go @@ -82,19 +82,11 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { } end = start + blocksize } - for { - newval, err := client.CompareAndSwap(path, strconv.FormatUint(end, 10), 0, strconv.FormatUint(start, 10), 0) - if err == nil { - break - } else { - log.Println("Error with CompareAndSet! Trying again in 1 second...") - time.Sleep(time.Second) - start, err = strconv.ParseUint(newval.Node.Value, 10, 0) - if err != nil { - log.Fatal(err) - } - end = start + blocksize - } + _, err = client.CompareAndSwap(path, strconv.FormatUint(end, 10), 0, strconv.FormatUint(start, 10), 0) + if err != nil { + log.Println("Error with CompareAndSet! Trying again in 1 second...") + time.Sleep(time.Second) + continue } for c := start; c < end; c += 1 { ch <- c From 146dc92acce3f8a656d4c7a14613d9f72cac4214 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 14 Feb 2014 11:48:19 -0600 Subject: [PATCH 6/8] Remove delay --- commands/pilosa-nexter/nexter.go | 3 --- 1 file changed, 3 deletions(-) diff --git a/commands/pilosa-nexter/nexter.go b/commands/pilosa-nexter/nexter.go index 22e360756..f2821c40e 100644 --- a/commands/pilosa-nexter/nexter.go +++ b/commands/pilosa-nexter/nexter.go @@ -8,7 +8,6 @@ import ( "net/http" "strconv" "strings" - "time" "github.com/coreos/go-etcd/etcd" ) @@ -84,8 +83,6 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { } _, err = client.CompareAndSwap(path, strconv.FormatUint(end, 10), 0, strconv.FormatUint(start, 10), 0) if err != nil { - log.Println("Error with CompareAndSet! Trying again in 1 second...") - time.Sleep(time.Second) continue } for c := start; c < end; c += 1 { From 4d23b5af03d1014e913b5787d8c96623fd5475f8 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 17 Feb 2014 16:59:47 -0600 Subject: [PATCH 7/8] Refactor --- core/http.go | 80 ++++++++++++++++++---------------------------------- 1 file changed, 27 insertions(+), 53 deletions(-) diff --git a/core/http.go b/core/http.go index c0a0b9bdf..7c50e0d53 100644 --- a/core/http.go +++ b/core/http.go @@ -282,19 +282,7 @@ func (self *WebService) HandleListen(w http.ResponseWriter, r *http.Request) { `)) } -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 - } +func (self *WebService) streamer(writer func(map[string]interface{}) error) { inbox := make(chan interface{}, 10) outbox := make(chan interface{}, 10) notify.Start("inbox", inbox) @@ -318,7 +306,7 @@ func (self *WebService) HandleListenWS(w http.ResponseWriter, r *http.Request) { typ := reflect.TypeOf(data) - err := ws.WriteJSON(map[string]interface{}{ + err := writer(map[string]interface{}{ "type": typ.String(), "dump": spew.Sdump(data), "host": host, @@ -334,6 +322,24 @@ func (self *WebService) HandleListenWS(w http.ResponseWriter, r *http.Request) { } } +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 + } + self.streamer(func(data map[string]interface{}) error { + return ws.WriteJSON(data) + }) +} + func drain(ch chan interface{}) { for { switch { @@ -349,43 +355,11 @@ func (self *WebService) HandleListenStream(w http.ResponseWriter, r *http.Reques err := recover() spew.Dump(err) }() - 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) - - writer := json.NewEncoder(w) - err := writer.Encode(map[string]interface{}{ - "type": typ.String(), - "dump": spew.Sdump(data), - "host": host, - }) - w.(http.Flusher).Flush() - if err != nil { - log.Println("stopping") - notify.Stop("inbox", inbox) - notify.Stop("outbox", outbox) - drain(inbox) - drain(outbox) - return - } - } + writer := json.NewEncoder(w) + flusher := w.(http.Flusher) + self.streamer(func(data map[string]interface{}) error { + err := writer.Encode(data) + flusher.Flush() + return err + }) } From a288b8760f800a654634365d565e167cd0838f69 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 17 Feb 2014 19:05:48 -0600 Subject: [PATCH 8/8] Add status screen --- core/http.go | 38 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 38 insertions(+) diff --git a/core/http.go b/core/http.go index 7c50e0d53..41c4cd547 100644 --- a/core/http.go +++ b/core/http.go @@ -32,6 +32,7 @@ func (self *WebService) Run() { mux.HandleFunc("/message", self.HandleMessage) mux.HandleFunc("/query", self.HandleQuery) mux.HandleFunc("/stats", self.HandleStats) + mux.HandleFunc("/status", self.HandleStatus) mux.HandleFunc("/info", self.HandleInfo) mux.HandleFunc("/processes", self.HandleProcesses) mux.HandleFunc("/listen/ws", self.HandleListenWS) @@ -363,3 +364,40 @@ func (self *WebService) HandleListenStream(w http.ResponseWriter, r *http.Reques return err }) } + +func (self *WebService) HandleStatus(w http.ResponseWriter, r *http.Request) { + w.Write([]byte(` + Pilosa - status + + + + + `)) +}