mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-10 04:47:53 +00:00
resolved merge conflicts
This commit is contained in:
commit
6968507296
9 changed files with 194 additions and 24 deletions
32
core/http.go
32
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)
|
||||
|
|
|
|||
41
core/ping.go
Normal file
41
core/ping.go
Normal file
|
|
@ -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
|
||||
}
|
||||
14
db/db.go
14
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{})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
||||
/*
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
24
hold/hold.go
24
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)
|
||||
}()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue