From 3f2ff1a0e3e1f0a3d9f8679797548511c0efa9e7 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 16 Oct 2013 17:34:43 -0500 Subject: [PATCH] Initial work on query planner, parser, and connection logic. --- commands/pilosa-cruncher/cruncher.go | 25 ++ commands/pilosa-router/router.go | 25 ++ core/service.go | 413 +++++++++++++++++++++++++++ core/topology.go | 134 +++++++++ cruncher/cruncher.go | 53 ++++ query/parser.go | 61 ++++ query/planner.go | 141 +++++++++ query/planner_example.go | 67 +++++ query/planner_test.go | 10 + router/router.go | 82 ++++++ 10 files changed, 1011 insertions(+) create mode 100644 commands/pilosa-cruncher/cruncher.go create mode 100644 commands/pilosa-router/router.go create mode 100644 core/service.go create mode 100644 core/topology.go create mode 100644 cruncher/cruncher.go create mode 100644 query/parser.go create mode 100644 query/planner.go create mode 100644 query/planner_example.go create mode 100644 query/planner_test.go create mode 100644 router/router.go diff --git a/commands/pilosa-cruncher/cruncher.go b/commands/pilosa-cruncher/cruncher.go new file mode 100644 index 000000000..0336de333 --- /dev/null +++ b/commands/pilosa-cruncher/cruncher.go @@ -0,0 +1,25 @@ +package main + +import ( + "pilosa/cruncher" + "pilosa/core" + "flag" + "log" +) + +var locationString string + +func init() { + flag.StringVar(&locationString, "l", "127.0.0.1:1300", "ip:port to listen on") + flag.Parse() +} + +func main() { + location, err := core.NewLocation(locationString) + if err != nil { + log.Fatal("Location not valid:", locationString) + } + + cruncher := cruncher.NewCruncher(location) + cruncher.Run() +} diff --git a/commands/pilosa-router/router.go b/commands/pilosa-router/router.go new file mode 100644 index 000000000..e56e7cde6 --- /dev/null +++ b/commands/pilosa-router/router.go @@ -0,0 +1,25 @@ +package main + +import ( + "pilosa/router" + "pilosa/core" + "flag" + "log" +) + +var locationString string + +func init() { + flag.StringVar(&locationString, "l", "127.0.0.1:1200", "ip:port to listen on") + flag.Parse() +} + +func main() { + location, err := core.NewLocation(locationString) + if err != nil { + log.Fatal("Location not valid:", locationString) + } + + router := router.NewRouter(location) + router.Run() +} diff --git a/core/service.go b/core/service.go new file mode 100644 index 000000000..b10cdbc8a --- /dev/null +++ b/core/service.go @@ -0,0 +1,413 @@ +package core + +import ( + "github.com/coreos/go-etcd/etcd" + "github.com/coreos/etcd/store" + "log" + "os" + "os/signal" + "syscall" + "time" + "strings" + "sync" + "net" + "encoding/gob" + //"net" + //"flag" + //"encoding/gob" + //"net/http" + //"encoding/json" + //"io" +) +type Message struct { + Key string `json:key` + Data interface{} `json:data` + Destination Location +} + +type Envelope struct { + Message *Message + Location *Location +} + +type Connection struct { + Conn net.Conn + Encoder *gob.Encoder + Decoder *gob.Decoder +} + +type Stopper struct { + TermChans []chan int + DoneChans []chan int + Mutex sync.RWMutex +} + +func (stopper *Stopper) Stop() { + var i chan int + var o chan int + stopper.Mutex.RLock() + for _, i = range stopper.TermChans { + go func() { + i <- 1 + }() + } + for _, o = range stopper.DoneChans { + <-o + } + stopper.Mutex.RUnlock() + return +} + +func (stopper *Stopper) GetExitChannels() (chan int, chan int) { + termchan := make(chan int, 1) + donechan := make(chan int, 1) + stopper.Mutex.Lock() + stopper.TermChans = append(stopper.TermChans, termchan) + stopper.DoneChans = append(stopper.DoneChans, donechan) + stopper.Mutex.Unlock() + return termchan, donechan +} + +type Service struct { + Stopper + Handler func(Message) + Mailbox chan Message + Inbox chan *Message + Port string + PortHttp string + Etcd *etcd.Client + Location *Location + NodeMap NodeMap + NodeMapMutex sync.RWMutex + Outbox chan *Envelope + ConnectionMap map[Location]*Connection + ConnectionMapMutex sync.RWMutex + Listener net.Listener + ConnectionRegisterChannel chan *PersistentConnection + //Cluster query.Cluster +} + +func NewService(location *Location) *Service { + service := new(Service) + service.Location = location + service.Outbox = make(chan *Envelope) + service.Inbox = make(chan *Message) + return service +} + +func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) { + hupChan := make(chan os.Signal, 1) + termChan := make(chan os.Signal, 1) + signal.Notify(hupChan, syscall.SIGHUP) + signal.Notify(termChan, syscall.SIGINT, syscall.SIGTERM) + return termChan, hupChan +} + +func (service *Service) SendMessage(message *Message) error { + if message.Destination == *service.Location { + return service.DeliverMessage(message) + } + router := service.GetRouterLocation(message.Destination) + var dest *Location + if router == *service.Location { + dest = &message.Destination + } else { + dest = &router + } + return service.DoSendMessage(message, dest) +} + +func (service *Service) DeliverMessage(message *Message) error { + log.Println("do handle of", message) + switch message.Key { + case "ping": + location, ok := message.Data.(Location) + if !ok { + log.Println("Invalid ping request!", message) + return nil + } + service.SendMessage(&Message{"pong", nil, location}) + } + return nil +} + +func (service *Service) DoSendMessage(message *Message, destination *Location) error { + log.Println("send message", message, "to", destination) + service.Outbox <- &Envelope{message, destination} + return nil +} + +type PersistentConnection struct { + Outbox chan *Message + Inbox chan *Message + Encoder *gob.Encoder + Decoder *gob.Decoder + Location Location + Service *Service + Connection *net.Conn + Identified bool +} + +func NewPersistentConnection(service *Service, location *Location) *PersistentConnection { + conn := &PersistentConnection{} + conn.Service = service + conn.Identified = true + conn.Outbox = make(chan *Message) + conn.Inbox = make(chan *Message) + if location != nil { + conn.Location = *location + } + return conn +} + +func (conn *PersistentConnection) TryConnect() error { + locationString := conn.Location.ToString() + log.Println("Dialing", locationString) + connection, err := net.Dial("tcp", locationString) + if err != nil { + log.Println("Error in dial!") + return err + } + conn.Connection = &connection + conn.Encoder = gob.NewEncoder(connection) + conn.Decoder = gob.NewDecoder(connection) + return nil +} + +func (conn *PersistentConnection) Manage() { + log.Println("Starting manager") + var connectionError chan error + var err error + + go func() { + var message *Message + var err error + for { + err = conn.Decoder.Decode(&message) + if err != nil { + connectionError <- err + return + } + conn.Inbox <- message + } + }() + + var message *Message + for { + select { + case message = <-conn.Outbox: + log.Println("sending", message) + err = conn.Encoder.Encode(message) + if err != nil { +// if e, ok := err.(*net.OpError); ok { +// if e.Err == syscall.EPIPE { +// // Client disconnected + + log.Println("error sending", err) + go func() { conn.Outbox <- message }() // Resend failed message + if !conn.Identified { + log.Println("Stopping because connection not identified") + return + } + conn.Connect() + } + case message = <-conn.Inbox: + log.Println("receiving", message) + if message.Key == "identify" { + log.Println("Registering connection") + conn.Location = message.Data.(Location) + go conn.Service.RegisterConnection(conn) + } else { + conn.Service.Inbox <- message + } + case err = <-connectionError: + log.Println("error receiving", err) + if !conn.Identified { + log.Println("Stopping because connection not identified") + return + } + conn.Connect() + } + } +} + +func (conn *PersistentConnection) Connect() { + log.Println("Connnecting to", conn.Location) + err := conn.TryConnect() + if err != nil { + log.Println(err) + log.Println("Connection failed! Waiting 1 second...") + time.Sleep(time.Second) + // Infinite recursion if node never comes up. Not sure if this is a problem. + conn.Connect() + } else { + log.Println("Send identify message") + conn.Encoder.Encode(Message{"identify", *conn.Service.Location, conn.Location}) + //go func() { conn.Outbox <- &Message{"identify", *conn.Service.Location, conn.Location} }() + } + log.Println("Connect ending") +} + +type ConnectionMapping map[Location]*PersistentConnection + +func (service *Service) RegisterConnection(conn *PersistentConnection) { + service.ConnectionRegisterChannel <- conn +} +func (service *Service) HandleConnections() { + connections := ConnectionMapping{} + + log.Println("Handling connections...") + for { + select { + case envelope := <-service.Outbox: + conn, ok := connections[*envelope.Location] + if !ok { + conn = NewPersistentConnection(service, envelope.Location) + connections[*envelope.Location] = conn + conn.Identified = true + go func () { + conn.Connect() + conn.Manage() + }() + } + // Spawning new goroutine so it doesn't block the main event loop while it's sending. + // This emulates an infinitely buffered channel. + go func() { conn.Outbox <- envelope.Message }() + case conn := <-service.ConnectionRegisterChannel: + connections[conn.Location] = conn + } + } +} + +func (service *Service) SetupEtcd() { + gob.Register(Location{}) + service.Etcd = etcd.NewClient() + service.NodeMapMutex.Lock() + defer service.NodeMapMutex.Unlock() + service.NodeMap = NodeMap{} + + nodes, err := service.Etcd.Get("nodes") + if err != nil { + log.Fatal(err) + } + for _, node := range nodes { + nodestring := strings.Split(node.Key, "/")[2] + location, err := NewLocation(nodestring) + if err != nil { + log.Fatal(err) + } + routerlocation, err := NewLocation(node.Value) + if err != nil { + log.Fatal(err) + } + service.NodeMap[*location] = *routerlocation + } + log.Println(service.NodeMap) +} + +func (service *Service) WatchEtcd() { + var receiver = make(chan *store.Response) + var stop chan bool + go func () { + _, err := service.Etcd.Watch("nodes/", 0, receiver, stop) + if err != nil { + log.Fatal(err) + } + }() + + exit, done := service.GetExitChannels() + + for { + select { + case response := <-receiver: + switch response.Action { + case "SET": + nodestring := strings.Split(response.Key, "/")[2] + node, err := NewLocation(nodestring) + if err != nil { + log.Fatal(err) + } + router, err := NewLocation(response.Value) + if err != nil { + log.Fatal(err) + } + service.NodeMapMutex.Lock() + service.NodeMap[*node] = *router + service.NodeMapMutex.Unlock() + case "DELETE": + nodestring := strings.Split(response.Key, "/")[2] + node, err := NewLocation(nodestring) + if err != nil { + log.Fatal(err) + } + service.NodeMapMutex.Lock() + delete(service.NodeMap, *node) + service.NodeMapMutex.Unlock() + default: + log.Println("unhandled etcd message", response) + } + //log.Println(response.Action, response.Key, response.Value) + log.Println(service.NodeMap) + case <-exit: + log.Println("cleaning up watchetcd service thing.") + time.Sleep(2*time.Second) + log.Println("done!") + done <- 1 + } + } +} + +func (service *Service) GetRouterLocation(node Location) Location { + service.NodeMapMutex.RLock() + defer service.NodeMapMutex.RUnlock() + location, ok := service.NodeMap[node] + if !ok { + location, ok = service.NodeMap[*service.Location] + if !ok { + log.Fatal("Cannot find router!!!") + } + } + return location +} + +func (service *Service) SetupNetwork() { + log.Println("Setup network") + var err error + locationString := service.Location.ToString() + service.Listener, err = net.Listen("tcp", locationString) + if err != nil { + log.Fatal(err) + } +} +func (service *Service) Serve() { + for { + conn, err := service.Listener.Accept() + if err != nil { + log.Fatal(err) + } + con := NewPersistentConnection(service, nil) + con.Connection = &conn + con.Encoder = gob.NewEncoder(*con.Connection) + con.Decoder = gob.NewDecoder(*con.Connection) + go con.Manage() + } +} + +//func (app *Application) serveHTTP() { +// http.HandleFunc("/message", func(w http.ResponseWriter, r *http.Request) { +// if r.Method != "POST" { +// http.Error(w, "Only POST allowed", http.StatusMethodNotAllowed) +// return +// } +// var message Message +// decoder := json.NewDecoder(r.Body) +// if decoder.Decode(&message) != nil { +// http.Error(w, "Invalid JSON", http.StatusBadRequest) +// return +// } +// app.Mailbox <- message +// }) +// http.ListenAndServe(app.PortHttp, nil) +//} +// + diff --git a/core/topology.go b/core/topology.go new file mode 100644 index 000000000..36f87c517 --- /dev/null +++ b/core/topology.go @@ -0,0 +1,134 @@ +package core + +import ( + "github.com/stathat/consistent" + "log" + "fmt" + "errors" + "strings" + "strconv" +) + +var FrameDoesNotExistError = errors.New("Frame does not exist.") + +type Location struct { + Ip string + Port int +} + +// Create a Location struct given a string in form "0.0.0.0:0" +func NewLocation(location_string string) (*Location, error) { + splitstring := strings.Split(location_string, ":") + if len(splitstring) != 2 { + return nil, errors.New("Location string must be in form 0.0.0.0:0") + } + ip := splitstring[0] + port, err := strconv.Atoi(splitstring[1]) + if err != nil{ + return nil, errors.New("Port is not a number!") + } + return &Location{ip, port}, nil +} + +func (location *Location) ToString() string { + return fmt.Sprintf("%s:%d", location.Ip, location.Port) +} + +// Map of node location to their router +type NodeMap map[Location]Location + +// A fragment is a collection of bitmaps within a slice. The fragment contains a reference to the responsible node for that fragment. The node is in the form ip:port +type Fragment struct { + Node string +} + +// A slice is the vertical combination of every fragment. It contains the hashring used to delegate bitmaps to fragments +type Slice struct { + Fragments []Fragment + Hashring *consistent.Consistent +} + +// A frame is a collection of slices in a given category (brands, demographics, etc), specific to a database +type Frame struct { + Name string + Slices []*Slice +} + +// Add a slice to a frame with given Node addresses +func (f *Frame) AddSlice(addrs ...string) *Slice { + slice := Slice{} + slice.Hashring = consistent.New() + slice.Hashring.NumberOfReplicas = 200 + sliceIndex := len(f.Slices) + for index, addr := range addrs { + slice.Fragments = append(slice.Fragments, Fragment{addr}) + slice.Hashring.Add(fmt.Sprintf("%d %d", sliceIndex, index)) + } + f.Slices = append(f.Slices, &slice) + return &slice +} + +// Represents the entire cluster, and a reference to the Node this instance is running on +type Cluster struct { + Databases map[string]*Database + Self string +} + +// Add a database to a cluster +func (c *Cluster) AddDatabase(name string) *Database { + database := Database{Name: name} + if c.Databases == nil { + c.Databases = make(map[string]*Database) + } + c.Databases[name] = &database + return &database +} + +// A database is a collection of all the frames within a given profile space +type Database struct { + Name string + Frames []*Frame +} + +// Count the number of slices in a database +func (d *Database) NumSlices() (int, error) { + if len(d.Frames) < 1 { + return 0, errors.New("Database is empty") + } + return len(d.Frames[0].Slices), nil +} + +// Add a frame to a database +func (d *Database) AddFrame(name string) *Frame { + frame := Frame{Name: name} + d.Frames = append(d.Frames, &frame) + return &frame +} + +// Get a frame from a database +func (d *Database) GetFrame(name string) (*Frame, error) { + for _, frame := range d.Frames { + if frame.Name == name { + return frame, nil + } + } + return nil, FrameDoesNotExistError +} + +// For debugging, prints cluster information +func (c *Cluster) Describe() { + for _, database := range c.Databases { + log.Println("frames", database.Frames) + for _, frame := range database.Frames { + log.Println(frame.Name, database.Name) + for _, slice := range frame.Slices { + log.Println(" ", slice) + } + } + } +} + +type Bitmap struct { + FrameType string + Id int +} diff --git a/cruncher/cruncher.go b/cruncher/cruncher.go new file mode 100644 index 000000000..3217f0388 --- /dev/null +++ b/cruncher/cruncher.go @@ -0,0 +1,53 @@ +package cruncher + +import ( + "pilosa/core" + "log" +) + +type Cruncher struct { + core.Service +} + +func (c *Cruncher) Init() { + log.Println("Initializing cruncher...") +} + +func (c *Cruncher) Run() { + log.Println("Running cruncher...") + c.SetupEtcd() + //go r.SyncEtcd() + go c.WatchEtcd() + go c.HandleConnections() + c.SetupNetwork() + go c.Serve() + go c.HandleInbox() + + sigterm, sighup := c.GetSignals() + for { + select { + case <- sighup: + log.Println("SIGHUP! Reloading configuration...") + // TODO: reload configuration + case <- sigterm: + log.Println("SIGTERM! Cleaning up...") + c.Stop() + return + } + } +} + +func NewCruncher(location *core.Location) *Cruncher { + service := core.NewService(location) + cruncher := Cruncher{*service} + return &cruncher +} + +func (c *Cruncher) HandleInbox() { + for { + select { + case message := <-c.Inbox: + log.Println("process", message) + } + } +} diff --git a/query/parser.go b/query/parser.go new file mode 100644 index 000000000..67f3c0a14 --- /dev/null +++ b/query/parser.go @@ -0,0 +1,61 @@ +package query + +import ( + "encoding/json" + "errors" + "pilosa/core" +) + +var InvalidQueryError = errors.New("Invalid query format.") + +type QueryParser struct { + QueryString string +} + +func (q *QueryParser) Walk(data interface{}) (*Query, error) { + query := new(Query) + + slice, ok := data.([]interface{}) + if !ok { + return nil, InvalidQueryError + } + operation, ok := slice[0].(string) + + if !ok { + return nil, InvalidQueryError + } + if operation == "union" || operation == "intersect" { + query.Operation = operation + inputs := slice[1:] + query.Inputs = make([]QueryInput, len(inputs)) + for idx, input := range inputs { + subquery, err := q.Walk(input) + if err != nil { + return nil, err + } + query.Inputs[idx] = subquery + } + } else if operation == "bitmap" { + query.Operation = "get" + frame, ok := slice[1].(string) + if !ok { + return nil, InvalidQueryError + } + id, ok := slice[2].(float64) + if !ok { + return nil, InvalidQueryError + } + id_int := int(id) + query.Inputs = []QueryInput{core.Bitmap{frame, id_int}} + } + + return query, nil +} + +func (q *QueryParser) Parse() (*Query, error) { + var data interface{} + if err := json.Unmarshal([]byte(q.QueryString), &data); err != nil { + return nil, err + } + return q.Walk(data) +} diff --git a/query/planner.go b/query/planner.go new file mode 100644 index 000000000..b7f96994e --- /dev/null +++ b/query/planner.go @@ -0,0 +1,141 @@ +package query + +import ( + "github.com/nu7hatch/gouuid" + "strconv" + "fmt" + "math/rand" + "pilosa/core" +) + +// A single step in the query plan. +type QueryStep struct { + id uuid.UUID + operation string + inputs []QueryInput + location string + destination string +} + +func (q QueryStep) String() string { + return fmt.Sprintf("%s %s %s, LOC: %s, DEST: %s", q.operation, q.id.String(), q.inputs, q.location, q.destination) +} + +type QueryInput interface {} + +// Represents a parsed query. Inputs can be Query or Bitmap objects +type Query struct{ + Operation string + Inputs []QueryInput // Maybe Bitmap and Query objects should have different fields to avoid using interface{} +} + +// This is the output of the query planner. Contains a list of steps which can be performed in parallel +type QueryPlan []QueryStep + +type QueryPlanner struct { + Cluster *core.Cluster + Database *core.Database +} + +type QueryTree interface { + getLocation(d *core.Database) string +} + +// QueryTree for UNION, INTER, and CAT queries +type CompositeQueryTree struct { + operation string + subqueries []QueryTree + location string +} + +// Randomly select location from subqueries (so subqueries roll up into composite queries while minimizing inter-node data traffic) +func (qt *CompositeQueryTree) getLocation(d *core.Database) string { + if qt.location == "" { + subqueryLength := len(qt.subqueries) + if subqueryLength > 1 { + locationIndex := rand.Intn(subqueryLength) + subquery := qt.subqueries[locationIndex] + qt.location = subquery.getLocation(d) + } + } + return qt.location +} + +// QueryTree for GET queries +type GetQueryTree struct { + bitmap core.Bitmap + slice int +} + +// Uses consistent hashing function to select node containing data for GET operation +func (qt *GetQueryTree) getLocation(d *core.Database) string { + frame, err := d.GetFrame(qt.bitmap.FrameType) + if err != nil { + panic(err) + } + slice := frame.Slices[qt.slice] + hashString, err := slice.Hashring.Get(strconv.Itoa(qt.bitmap.Id)) + + var sliceIndex int + var fragIndex int + fmt.Sscan(hashString, &fragIndex, &sliceIndex) + fragment := slice.Fragments[fragIndex] + + return fmt.Sprintf(fragment.Node) +} + +// Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices +func (qp *QueryPlanner) buildTree(query *Query, slice int) QueryTree { + var tree QueryTree + if slice == -1 { + tree = &CompositeQueryTree{operation:"cat"} + numSlices, err := qp.Database.NumSlices() + if err != nil { + panic(err) + } + for slice := 0; slice < numSlices; slice++ { + subtree := qp.buildTree(query, slice) + composite := tree.(*CompositeQueryTree) + composite.subqueries = append(composite.subqueries, subtree) + } + } else { + if query.Operation == "get" { + tree = &GetQueryTree{query.Inputs[0].(core.Bitmap), slice} + return tree + } else { + subqueries := make([]QueryTree, len(query.Inputs)) + for i, input := range query.Inputs { + subqueries[i] = qp.buildTree(input.(*Query), slice) + } + tree = &CompositeQueryTree{operation:query.Operation, subqueries:subqueries} + } + } + return tree +} + +// Produces flattened QueryPlan from QueryTree input +func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location string) *QueryPlan { + plan := QueryPlan{} + if composite, ok := qt.(*CompositeQueryTree); ok { + inputs := make([]QueryInput, len(composite.subqueries)) + step := QueryStep{*id, composite.operation, inputs, composite.getLocation(qp.Database), location} + for index, subq := range composite.subqueries { + sub_id, _ := uuid.NewV4() + step.inputs[index] = []QueryInput{"wait", sub_id} + subq_steps := qp.flatten(subq, sub_id, composite.getLocation(qp.Database)) + plan = append(plan, *subq_steps...) + } + plan = append(plan, step) + } else if get, ok := qt.(*GetQueryTree); ok { + step := QueryStep{*id, "get", []QueryInput{get.bitmap, get.slice}, get.getLocation(qp.Database), location} + plan := QueryPlan{step} + return &plan + } + return &plan +} + +// Transforms Query into QueryTree and flattens to QueryPlan object +func (qp *QueryPlanner) Plan(query *Query, id *uuid.UUID, destination string, slice int) *QueryPlan { + queryTree := qp.buildTree(query, -1) + return qp.flatten(queryTree, id, destination) +} diff --git a/query/planner_example.go b/query/planner_example.go new file mode 100644 index 000000000..7e4ceb12b --- /dev/null +++ b/query/planner_example.go @@ -0,0 +1,67 @@ +package query +//package main +// +//import ( +// "pilosa/query" +// "pilosa/core" +// "math/rand" +// "github.com/nu7hatch/gouuid" +// "time" +// "log" +//) +// +//func main() { +// rand.Seed(time.Now().UnixNano()) +// cluster := core.Cluster{Self:"192.168.1.100:1201"} +// database := cluster.AddDatabase("property49") +// frame := database.AddFrame("general") +// frame.AddSlice("192.168.1.100:1201", "192.168.1.100:1202", "192.168.1.100:1203") +// frame.AddSlice("192.168.1.101:1201", "192.168.1.101:1202", "192.168.1.101:1203") +// frame.AddSlice("192.168.1.102:1201", "192.168.1.102:1202", "192.168.1.102:1203") +// frame2 := database.AddFrame("brands") +// frame2.AddSlice("192.168.1.200:1201", "192.168.1.200:1202", "192.168.1.200:1203") +// frame2.AddSlice("192.168.1.201:1201", "192.168.1.201:1202", "192.168.1.201:1203") +// frame2.AddSlice("192.168.1.202:1201", "192.168.1.202:1202", "192.168.1.202:1203") +// +// //cluster.describe() +// +// //query := Query{"get", []QueryInput{Bitmap{"general", 10}}} +// +// //query := Query{"union", []QueryInput{ +// // &Query{"get", []QueryInput{Bitmap{"general", 20}}}, +// // &Query{"get", []QueryInput{Bitmap{"brands", 30}}}, +// //}} +// +// //queryString = "union(bitmap(general, 33), bitmap(brands, 44))" +// //queryString := `["union", ["intersect", ["bitmap", "general", 33], ["bitmap", "brands", 44]], ["bitmap", "general", 55]]` +// //query := `count(union(bitmap(""), bitmap(brands, 55)))` +// //query := `setbit(bitmap(cats, 33), 75000)` +// queryString := `["union", ["intersect", ["bitmap", "general", 33], ["bitmap", "brands", 44]], ["bitmap", "general", 55]]` +// queryParser := query.QueryParser{queryString} +// queryParsed, err := queryParser.Parse() +// log.Println(queryParsed) +// if err != nil { +// log.Println("ERROR!", err) +// } +// +//// query := qp.Query{"inter", []qp.QueryInput{ +//// &qp.Query{"union", []qp.QueryInput{ +//// &qp.Query{"get", []qp.QueryInput{core.Bitmap{"general", 20}}}, +//// &qp.Query{"get", []qp.QueryInput{core.Bitmap{"brands", 30}}}, +//// }}, +//// &qp.Query{"get", []qp.QueryInput{core.Bitmap{"general", 10}}}, +//// }} +//// +// planner := query.QueryPlanner{&cluster, database} +// +// id, _ := uuid.NewV4() +// dest := "1.2.3.4:1234" +// log.Println("Query id", id, "Dest", dest) +// +// queryplan := planner.Plan(queryParsed, id, dest, -1) +// +// for _, i := range *queryplan { +// log.Println(i) +// } +// //fmt.Println(queryplan) +//} diff --git a/query/planner_test.go b/query/planner_test.go new file mode 100644 index 000000000..db6e6a623 --- /dev/null +++ b/query/planner_test.go @@ -0,0 +1,10 @@ +package query + +import ( + "testing" + "log" +) + +func TestQueryPlanner(t *testing.T) { + log.Println("query planner test") +} diff --git a/router/router.go b/router/router.go new file mode 100644 index 000000000..dc0dcaf7c --- /dev/null +++ b/router/router.go @@ -0,0 +1,82 @@ +package router + +import ( + "log" + "time" + //"strings" + "pilosa/core" +) + +type Router struct { + core.Service +} + +func (r *Router) Init() { + log.Println("Initializing router...") +} + +func (r *Router) Run() { + log.Println("Running router...") + r.SetupEtcd() + //go r.SyncEtcd() + go r.WatchEtcd() + go r.HandleConnections() + + go func() { + for { + r.SendMessage(&core.Message{"ping", core.Location{"127.0.0.1", 1200}, core.Location{"127.0.0.1", 1300}}) + time.Sleep(2*time.Second) + //log.Println(r.GetRouterLocation(core.Location{"127.0.0.1", 1200})) + } + }() + + sigterm, sighup := r.GetSignals() + for { + select { + case <- sighup: + log.Println("SIGHUP! Reloading configuration...") + // TODO: reload configuration + case <- sigterm: + log.Println("SIGTERM! Cleaning up...") + r.Stop() + return + } + } +} + +func (r *Router) SetupEtcd() { + r.Service.SetupEtcd() + //routers, err := r.Etcd.Get("topology") + //if err != nil { + // log.Fatal(err) + //} + //for _, router := range routers { + // routerstring := strings.Split(router.Key, "/")[2] + // routerlocation, err := core.NewLocation(routerstring) + // if err != nil { + // log.Fatal(err) + // } + // nodes, err := r.Etcd.Get("topology/" + routerstring) + // if err != nil { + // log.Fatal(err) + // } + // for _, node := range nodes { + // nodestring := strings.Split(node.Key, "/")[3] + // nodelocation, err := core.NewLocation(nodestring) + // if err != nil { + // log.Fatal(err) + // } + // r.NodeMap[*nodelocation] = *routerlocation + // } + //} +} + +func (r *Router) HandleMessage(m *core.Message) { + log.Println(m) +} + +func NewRouter(location *core.Location) *Router { + service := core.NewService(location) + router := Router{*service} + return &router +}