From 622b3b2351cd5dc8a2c0647d869e512759f03870 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 3 Jan 2014 17:27:19 -0600 Subject: [PATCH 1/3] Add timeouts. --- hold/hold.go | 24 +++++++++++++++++++----- hold/hold_test.go | 8 ++++---- 2 files changed, 23 insertions(+), 9 deletions(-) diff --git a/hold/hold.go b/hold/hold.go index dd10c7070..5819e1134 100644 --- a/hold/hold.go +++ b/hold/hold.go @@ -1,6 +1,11 @@ package hold -import "github.com/nu7hatch/gouuid" +import ( + "errors" + "time" + + "github.com/nu7hatch/gouuid" +) type holdchan chan interface{} type gethold struct { @@ -31,15 +36,24 @@ func (self *Holder) GetChan(id *uuid.UUID) holdchan { return <-reply } -func (self *Holder) Get(id *uuid.UUID) interface{} { +func (self *Holder) Get(id *uuid.UUID, timeout int) (interface{}, error) { ch := self.GetChan(id) - return <-ch + select { + case val := <-ch: + return val, nil + case <-time.After(time.Duration(timeout) * time.Second): + self.DelChan(id) + return nil, errors.New("Timeout getting from holder") + } } -func (self *Holder) Set(id *uuid.UUID, value interface{}) { +func (self *Holder) Set(id *uuid.UUID, value interface{}, timeout int) { ch := self.GetChan(id) go func() { - ch <- value + select { + case ch <- value: + case <-time.After(time.Duration(timeout) * time.Second): + } self.DelChan(id) }() } diff --git a/hold/hold_test.go b/hold/hold_test.go index a57edcc8f..93c9e9623 100644 --- a/hold/hold_test.go +++ b/hold/hold_test.go @@ -11,17 +11,17 @@ import ( func TestHoldChan(t *testing.T) { Convey("set then get", t, func() { id, _ := uuid.NewV4() - Hold.Set(id, "derp") - derp := Hold.Get(id) + Hold.Set(id, "derp", 10) + derp, _ := Hold.Get(id, 10) So(derp, ShouldEqual, "derp") }) Convey("get then set", t, func() { id, _ := uuid.NewV4() go func() { time.Sleep(time.Second / 10) - Hold.Set(id, "derpsy") + Hold.Set(id, "derpsy", 10) }() - derp := Hold.Get(id) + derp, _ := Hold.Get(id, 10) So(derp, ShouldEqual, "derpsy") }) } From 11d964339cb78c1d21652a9ba3f4507f972f8e6f Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 3 Jan 2014 18:14:04 -0600 Subject: [PATCH 2/3] Basic ping working. --- core/http.go | 31 +++++++++++++++++++ db/db.go | 21 +++++++++++++ dispatch/dispatch.go | 13 +++++--- interfaces/core.go | 4 ++- transport/tcp.go | 72 ++++++++++++++++++++++++++++++++++++++++---- 5 files changed, 130 insertions(+), 11 deletions(-) diff --git a/core/http.go b/core/http.go index 4a7e2bfdf..d40e4f8a8 100644 --- a/core/http.go +++ b/core/http.go @@ -7,10 +7,12 @@ import ( "net/http" "pilosa/config" "pilosa/db" + "pilosa/hold" "pilosa/query" "strconv" "github.com/davecgh/go-spew/spew" + "github.com/nu7hatch/gouuid" ) type WebService struct { @@ -31,6 +33,7 @@ func (self *WebService) Run() { mux.HandleFunc("/info", self.HandleInfo) mux.HandleFunc("/processes", self.HandleProcesses) mux.HandleFunc("/listen", self.HandleListen) + mux.HandleFunc("/ping", self.HandlePing) s := &http.Server{ Addr: ":" + port_string, Handler: mux, @@ -111,6 +114,34 @@ func (self *WebService) HandleProcesses(w http.ResponseWriter, r *http.Request) } } +func (self *WebService) HandlePing(w http.ResponseWriter, r *http.Request) { + if r.Method != "GET" { + http.Error(w, "Only GET allowed", http.StatusMethodNotAllowed) + return + } + err := r.ParseForm() + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + process_string := r.Form.Get("process") + process_id, err := uuid.ParseHex(process_string) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + _, err = self.service.ProcessMap.GetProcess(process_id) + if err != nil { + http.Error(w, err.Error(), http.StatusNotFound) + return + } + id, _ := uuid.NewV4() + ping := db.Message{Data: db.PingRequest{Id: id, Source: self.service.Id}} + self.service.Transport.Send(&ping, process_id) + data, err := hold.Hold.Get(id, 60) + spew.Fdump(w, data, err) +} + func (self *WebService) HandleListen(w http.ResponseWriter, r *http.Request) { //listener := service.NewListener() //encoder := json.NewEncoder(w) diff --git a/db/db.go b/db/db.go index f49216985..ff7d36279 100644 --- a/db/db.go +++ b/db/db.go @@ -1,7 +1,28 @@ package db +import ( + "encoding/gob" + + "github.com/nu7hatch/gouuid" +) + type Message struct { Key string `json:key` Data interface{} `json:data` Destination Location } + +type PingRequest struct { + Id *uuid.UUID + Source *uuid.UUID +} + +type PongRequest struct { + Id *uuid.UUID +} + +func init() { + gob.Register(Message{}) + gob.Register(PingRequest{}) + gob.Register(PongRequest{}) +} diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 74b03d747..800860312 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -3,8 +3,8 @@ package dispatch import ( "log" "pilosa/core" - - "github.com/davecgh/go-spew/spew" + "pilosa/db" + "pilosa/hold" ) type Dispatch struct { @@ -25,8 +25,13 @@ func (self *Dispatch) Run() { for { message := self.service.Transport.Receive() log.Println("Processing ", message) - spew.Dump(message.Key) - spew.Dump(message.Data) + switch data := message.Data.(type) { + case db.PingRequest: + pong := db.Message{Data: db.PongRequest{Id: data.Id}} + self.service.Transport.Send(&pong, data.Source) + case db.PongRequest: + hold.Hold.Set(data.Id, 1, 10) + } /* path := message.Data.(string) diff --git a/interfaces/core.go b/interfaces/core.go index d48ad9c43..5dc535c16 100644 --- a/interfaces/core.go +++ b/interfaces/core.go @@ -3,12 +3,14 @@ package interfaces import ( "pilosa/db" "pilosa/query" + + "github.com/nu7hatch/gouuid" ) type Transporter interface { Run() Close() - Send(*db.Message) + Send(*db.Message, *uuid.UUID) Receive() *db.Message } diff --git a/transport/tcp.go b/transport/tcp.go index 34e73d5ad..5dfd16613 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -1,28 +1,88 @@ package transport import ( + "encoding/gob" + "fmt" "log" + "net" "pilosa/config" "pilosa/core" "pilosa/db" + "time" + + "github.com/nu7hatch/gouuid" ) +type qmessage struct { + message *db.Message + host *uuid.UUID +} + type TcpTransport struct { - port int - inbox chan *db.Message - stop chan bool + service *core.Service + port int + inbox chan *db.Message + outbox chan qmessage + stop chan bool } func (self *TcpTransport) Run() { log.Println("Initializing TCP transport") + go self.listen() + for { + select { + case qmessage := <-self.outbox: + // todo: persistent connections + process, err := self.service.ProcessMap.GetProcess(qmessage.host) + if err != nil { + log.Println("transport/tcp", err) + return + } + host_string := fmt.Sprintf("%s:%d", process.Host(), process.PortTcp()) + conn, err := net.Dial("tcp", host_string) + encoder := gob.NewEncoder(conn) + err = encoder.Encode(qmessage.message) + if err != nil { + log.Println(err.Error()) + return + } + time.Sleep(1 * time.Second) + } + } +} + +func (self *TcpTransport) listen() { + port_string := fmt.Sprintf(":%d", self.port) + l, e := net.Listen("tcp", port_string) + if e != nil { + log.Fatal("Cannot bind to port!", self.port) + } + for { + conn, err := l.Accept() + if err != nil { + log.Fatal("Cannot accept message!") + } + go self.manage(&conn) + } +} + +func (self *TcpTransport) manage(conn *net.Conn) { + decoder := gob.NewDecoder(*conn) + var mess *db.Message + err := decoder.Decode(&mess) + if err != nil { + log.Println("tcp/transport", err.Error()) + return + } + self.inbox <- mess } func (self *TcpTransport) Close() { log.Println("Shutting down TCP transport") } -func (self *TcpTransport) Send(message *db.Message) { - log.Println("Send", message) +func (self *TcpTransport) Send(message *db.Message, host *uuid.UUID) { + self.outbox <- qmessage{message, host} } func (self *TcpTransport) Receive() *db.Message { @@ -30,5 +90,5 @@ func (self *TcpTransport) Receive() *db.Message { } func NewTcpTransport(service *core.Service) *TcpTransport { - return &TcpTransport{config.GetInt("port_tcp"), make(chan *db.Message), nil} + return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message), make(chan qmessage), nil} } From 0a50f18eca78abb9f8b654bb388e714575845e8a Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 6 Jan 2014 16:20:31 -0600 Subject: [PATCH 3/3] Use HoldResult interface to route responses from dispatch. Added standalone ping function. --- core/http.go | 13 +++++++------ core/ping.go | 42 ++++++++++++++++++++++++++++++++++++++++++ db/db.go | 11 ++--------- dispatch/dispatch.go | 14 +++++++++----- 4 files changed, 60 insertions(+), 20 deletions(-) create mode 100644 core/ping.go diff --git a/core/http.go b/core/http.go index d40e4f8a8..a6862128a 100644 --- a/core/http.go +++ b/core/http.go @@ -7,7 +7,6 @@ import ( "net/http" "pilosa/config" "pilosa/db" - "pilosa/hold" "pilosa/query" "strconv" @@ -135,11 +134,13 @@ func (self *WebService) HandlePing(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusNotFound) return } - id, _ := uuid.NewV4() - ping := db.Message{Data: db.PingRequest{Id: id, Source: self.service.Id}} - self.service.Transport.Send(&ping, process_id) - data, err := hold.Hold.Get(id, 60) - spew.Fdump(w, data, err) + duration, err := self.service.Ping(process_id) + if err != nil { + spew.Fdump(w, err) + return + } + encoder := json.NewEncoder(w) + encoder.Encode(map[string]float64{"duration": duration.Seconds()}) } func (self *WebService) HandleListen(w http.ResponseWriter, r *http.Request) { diff --git a/core/ping.go b/core/ping.go new file mode 100644 index 000000000..164453264 --- /dev/null +++ b/core/ping.go @@ -0,0 +1,42 @@ +package core + +import ( + "encoding/gob" + "pilosa/db" + "pilosa/hold" + "time" + + "github.com/nu7hatch/gouuid" +) + +type PingRequest struct { + Id *uuid.UUID + Source *uuid.UUID +} + +type PongRequest struct { + Id *uuid.UUID +} + +func (self PongRequest) ResultId() *uuid.UUID { + return self.Id +} + +func init() { + gob.Register(PingRequest{}) + gob.Register(PongRequest{}) +} + +func (self *Service) Ping(process_id *uuid.UUID) (*time.Duration, error) { + id, _ := uuid.NewV4() + ping := db.Message{Data: PingRequest{Id: id, Source: self.Id}} + start := time.Now() + self.Transport.Send(&ping, process_id) + _, err := hold.Hold.Get(id, 60) + if err != nil { + return nil, err + } + end := time.Now() + dur := end.Sub(start) + return &dur, nil +} diff --git a/db/db.go b/db/db.go index ff7d36279..46f93992d 100644 --- a/db/db.go +++ b/db/db.go @@ -12,17 +12,10 @@ type Message struct { Destination Location } -type PingRequest struct { - Id *uuid.UUID - Source *uuid.UUID -} - -type PongRequest struct { - Id *uuid.UUID +type HoldResult interface { + ResultId() *uuid.UUID } func init() { gob.Register(Message{}) - gob.Register(PingRequest{}) - gob.Register(PongRequest{}) } diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 800860312..f79062f4f 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -5,6 +5,8 @@ import ( "pilosa/core" "pilosa/db" "pilosa/hold" + + "github.com/davecgh/go-spew/spew" ) type Dispatch struct { @@ -24,13 +26,15 @@ func (self *Dispatch) Run() { log.Println("Dispatch Run...") for { message := self.service.Transport.Receive() - log.Println("Processing ", message) + spew.Dump("Processing ", message) switch data := message.Data.(type) { - case db.PingRequest: - pong := db.Message{Data: db.PongRequest{Id: data.Id}} + case core.PingRequest: + pong := db.Message{Data: core.PongRequest{Id: data.Id}} self.service.Transport.Send(&pong, data.Source) - case db.PongRequest: - hold.Hold.Set(data.Id, 1, 10) + case db.HoldResult: + hold.Hold.Set(data.ResultId(), data, 30) + default: + log.Println("Unprocessed message", data) } /*