diff --git a/core/http.go b/core/http.go index bcacb6c20..a76a6fb73 100644 --- a/core/http.go +++ b/core/http.go @@ -10,6 +10,7 @@ import ( "strconv" "github.com/davecgh/go-spew/spew" + "tux21b.org/v1/gocql/uuid" ) type WebService struct { @@ -31,6 +32,7 @@ func (self *WebService) Run() { mux.HandleFunc("/processes", self.HandleProcesses) mux.HandleFunc("/listen", self.HandleListen) mux.HandleFunc("/test", self.HandleTest) + mux.HandleFunc("/ping", self.HandlePing) s := &http.Server{ Addr: ":" + port_string, Handler: mux, @@ -120,6 +122,36 @@ 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.ParseUUID(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 + } + 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) { //listener := service.NewListener() //encoder := json.NewEncoder(w) diff --git a/core/ping.go b/core/ping.go new file mode 100644 index 000000000..0c0a9e1bb --- /dev/null +++ b/core/ping.go @@ -0,0 +1,41 @@ +package core + +import ( + "encoding/gob" + "pilosa/db" + "time" + + "tux21b.org/v1/gocql/uuid" +) + +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.RandomUUID() + ping := db.Message{Data: PingRequest{Id: &id, Source: self.Id}} + start := time.Now() + self.Transport.Send(&ping, process_id) + _, err := self.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 0433ace2e..434c583c0 100644 --- a/db/db.go +++ b/db/db.go @@ -1,5 +1,11 @@ package db +import ( + "encoding/gob" + + "tux21b.org/v1/gocql/uuid" +) + type Message struct { Data interface{} `json:data` } @@ -18,3 +24,11 @@ type Message struct { Destination Location } */ + +type HoldResult interface { + ResultId() *uuid.UUID +} + +func init() { + gob.Register(Message{}) +} diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 12eba7f04..ea32b3b7b 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -4,6 +4,7 @@ import ( "fmt" "log" "pilosa/core" + "pilosa/db" "pilosa/query" "github.com/davecgh/go-spew/spew" @@ -26,15 +27,18 @@ func (self *Dispatch) Run() { log.Println("Dispatch Run...") for { message := self.service.Transport.Receive() - log.Println("Processing ", message) - spew.Dump(message.Data) - - switch message.Data.(type) { + spew.Dump("Processing ", message) + switch data := message.Data.(type) { + case core.PingRequest: + pong := db.Message{Data: core.PongRequest{Id: data.Id}} + self.service.Transport.Send(&pong, data.Source) + case db.HoldResult: + self.service.Hold.Set(data.ResultId(), data, 30) case query.CatQueryStep, query.GetQueryStep, query.SetQueryStep: fmt.Println("CAT/GET/SET QUERYSTEP") go self.service.Executor.NewJob(message) default: - fmt.Println("unknown") + log.Println("Unprocessed message", data) } /* diff --git a/executor/executor.go b/executor/executor.go index 9ded6e7a2..e633b025e 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -36,7 +36,8 @@ func (self *Executor) NewJob(job *db.Message) { spew.Dump(qs) for _, input := range qs.Inputs { spew.Dump(input) - bh := self.service.Hold.Get(input).(index.BitmapHandle) + bhi, err := self.service.Hold.Get(input, 10) + bh := bhi.(index.BitmapHandle) // TODO: git rid of this count, need the cat to do a sum() or a true cat() count, err := self.service.Index.Count(qs.Location.FragmentId, bh) if err != nil { @@ -57,7 +58,7 @@ func (self *Executor) NewJob(job *db.Message) { } //spew.Dump("COUNT", count) // push results to the map - self.service.Hold.Set(qs.Id, bh) + self.service.Hold.Set(qs.Id, bh, 10) case query.SetQueryStep: qs := job.Data.(query.SetQueryStep) diff --git a/hold/hold.go b/hold/hold.go index 5b0dbe55f..17affd846 100644 --- a/hold/hold.go +++ b/hold/hold.go @@ -1,6 +1,11 @@ package hold -import "tux21b.org/v1/gocql/uuid" +import ( + "errors" + "time" + + "tux21b.org/v1/gocql/uuid" +) 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 e3348604b..a9e057108 100644 --- a/hold/hold_test.go +++ b/hold/hold_test.go @@ -15,17 +15,17 @@ func TestHoldChan(t *testing.T) { Convey("set then get", t, func() { id := uuid.RandomUUID() - 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.RandomUUID() 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") }) } diff --git a/interfaces/core.go b/interfaces/core.go index e916f1c53..fd4aafda1 100644 --- a/interfaces/core.go +++ b/interfaces/core.go @@ -1,11 +1,15 @@ package interfaces -import "pilosa/db" +import ( + "pilosa/db" + + "tux21b.org/v1/gocql/uuid" +) type Transporter interface { Run() Close() - Send(*db.Message) + Send(*db.Message, *uuid.UUID) Receive() *db.Message Push(*db.Message) } diff --git a/transport/tcp.go b/transport/tcp.go index e29707937..3b9c0b409 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" + + "tux21b.org/v1/gocql/uuid" ) +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 { @@ -34,5 +94,5 @@ func (self *TcpTransport) Push(message *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} }