From c1edbb9382edfc7f58aa43a5a0086dfc84cdc6e6 Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 10 Jan 2014 13:24:53 -0600 Subject: [PATCH 1/4] preparation for macro support. clean up the RunQuery logic --- core/http.go | 9 +++- core/query.go | 7 ++- dispatch/dispatch.go | 3 ++ executor/executor.go | 105 +++++++++++++++++++++++++++++++++++------- interfaces/core.go | 2 +- query/parser.go | 4 ++ query/planner.go | 8 +++- query/planner_test.go | 9 ++-- query/query.go | 30 ++++++++++++ 9 files changed, 154 insertions(+), 23 deletions(-) diff --git a/core/http.go b/core/http.go index a76a6fb73..5aafaa786 100644 --- a/core/http.go +++ b/core/http.go @@ -68,7 +68,14 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) { // TODO: we need to get the database name from the query string (for now, hard-coded) database_name := "main" pql := string(body) - self.service.Executor.RunQuery(database_name, pql) + results := self.service.Executor.RunPQL(database_name, pql) + + encoder := json.NewEncoder(w) + err = encoder.Encode(results) + if err != nil { + log.Fatal("Error encoding stats") + } + } func (self *WebService) HandleStats(w http.ResponseWriter, r *http.Request) { diff --git a/core/query.go b/core/query.go index 4a1a37655..412242fda 100644 --- a/core/query.go +++ b/core/query.go @@ -160,5 +160,10 @@ func (self *Service) GetQueryStepHandler(msg *db.Message) { func (self *Service) SetQueryStepHandler(msg *db.Message) { spew.Dump("SET QUERYSTEP") qs := msg.Data.(query.SetQueryStep) - self.Index.SetBit(qs.Location.FragmentId, qs.Bitmap.Id, qs.ProfileId) + result, err := self.Index.SetBit(qs.Location.FragmentId, qs.Bitmap.Id, qs.ProfileId) + spew.Dump("result:", result) + spew.Dump("err:", err) + + result_message := db.Message{Data: query.SetQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + self.Transport.Send(&result_message, qs.Destination.ProcessId) } diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 3b8754419..7a3facd27 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -5,6 +5,8 @@ import ( "pilosa/core" "pilosa/db" "pilosa/query" + + "github.com/davecgh/go-spew/spew" ) type Dispatch struct { @@ -29,6 +31,7 @@ func (self *Dispatch) Run() { pong := db.Message{Data: core.PongRequest{Id: data.Id}} self.service.Transport.Send(&pong, data.Source) case db.HoldResult: + spew.Dump("HOLD-SET", data.ResultId()) self.service.Hold.Set(data.ResultId(), data.ResultData(), 30) case query.PortableQueryStep: go self.service.Executor.NewJob(message) diff --git a/executor/executor.go b/executor/executor.go index 7c85fb069..e208447f4 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -2,12 +2,13 @@ package executor import ( "fmt" + "io/ioutil" "log" + "pilosa/config" "pilosa/core" "pilosa/db" "pilosa/query" "pilosa/util" - "reflect" "tux21b.org/v1/gocql/uuid" "github.com/davecgh/go-spew/spew" @@ -46,8 +47,35 @@ func (self *Executor) NewJob(job *db.Message) { } } -func (self *Executor) RunQuery(database_name string, pql string) { - database := self.service.Cluster.GetOrCreateDatabase(database_name) +type stringSlice []string + +func (slice stringSlice) pos(value string) int { + for p, v := range slice { + if v == value { + return p + } + } + return -1 +} + +func (self *Executor) RunQueryTest(database_name string, pql string) string { + return pql +} + +type queryItem struct { + label string + pql string +} +type queryItemCopy struct { + id *uuid.UUID + label string +} +type QueryItemResult struct { + Label string + Result interface{} +} + +func (self *Executor) runQuery(database *db.Database, qry *query.Query) { process, err := self.service.GetProcess() if err != nil { spew.Dump(err) @@ -56,33 +84,78 @@ func (self *Executor) RunQuery(database_name string, pql string) { fragment_id := util.SUUID(0) destination := db.Location{&process_id, fragment_id} - query_plan := query.QueryPlanForPQL(database, pql, &destination) - //spew.Dump(query_plan) - + spew.Dump("QUERY.ID:", qry.Id) + query_plan := query.QueryPlanForQuery(database, qry, &destination) // loop over the query steps and send to Transport - var last_id *uuid.UUID for _, qs := range *query_plan { msg := new(db.Message) msg.Data = qs switch step := qs.(type) { case query.PortableQueryStep: self.service.Transport.Send(msg, step.GetLocation().ProcessId) - if reflect.TypeOf(step) != reflect.TypeOf(query.SetQueryStep{}) { - last_id = step.GetId() - } } } +} - // add an entry to my execute map[key] that is waiting for the final result - if last_id != nil { - final, err := self.service.Hold.Get(last_id, 10) +func (self *Executor) RunPQL(database_name string, pql string) interface{} { + database := self.service.Cluster.GetOrCreateDatabase(database_name) + + // see if the outer query function is a custom query + reserved_functions := stringSlice{"get", "set", "union", "intersect", "count"} + tokens := query.Lex(pql) + outer_token := tokens[0].Text + + if reserved_functions.pos(outer_token) != -1 { + + qry := query.QueryForTokens(tokens) + go self.runQuery(database, qry) + + var final interface{} + final, err := self.service.Hold.Get(qry.Id, 10) if err != nil { spew.Dump(err) } - spew.Dump("*******************************************************") - spew.Dump("FINAL", final) - spew.Dump("*******************************************************") + return final + + } else { + macros_dir := config.Get("macros").(string) + macros_file := macros_dir + "/" + outer_token + ".js" + file_data, err := ioutil.ReadFile(macros_file) + if err != nil { + spew.Dump(err) + } + spew.Dump(file_data) + + // CUSTOM QUERY LIST /////////////// + var query_list []queryItem + query_list = append(query_list, queryItem{"set1", "set(20, 1)"}) + query_list = append(query_list, queryItem{"set2", "set(20, 1)"}) + query_list = append(query_list, queryItem{"set3", "set(20, 2)"}) + query_list = append(query_list, queryItem{"set4", "set(20, 3)"}) + query_list = append(query_list, queryItem{"set5", "set(20, 4)"}) + query_list = append(query_list, queryItem{"count", "count(get(20))"}) + var query_list_copy []*queryItemCopy + //////////////////////////////////// + + for _, qi := range query_list { + qry := query.QueryForPQL(qi.pql) + go self.runQuery(database, qry) + query_list_copy = append(query_list_copy, &queryItemCopy{qry.Id, qi.label}) + } + + var final_result []*QueryItemResult + for _, qlc := range query_list_copy { + final, err := self.service.Hold.Get(qlc.id, 10) + if err != nil { + spew.Dump(err) + } + spew.Dump(final) + final_result = append(final_result, &QueryItemResult{qlc.label, final}) + } + + return final_result } + } func (self *Executor) Run() { diff --git a/interfaces/core.go b/interfaces/core.go index fd4aafda1..77d82aa1d 100644 --- a/interfaces/core.go +++ b/interfaces/core.go @@ -25,5 +25,5 @@ type Executorer interface { Close() Run() NewJob(*db.Message) - RunQuery(string, string) + RunPQL(string, string) interface{} } diff --git a/query/parser.go b/query/parser.go index 43f63b673..940b8fd2d 100644 --- a/query/parser.go +++ b/query/parser.go @@ -4,6 +4,8 @@ import ( "errors" "pilosa/db" "strconv" + + "tux21b.org/v1/gocql/uuid" ) var InvalidQueryError = errors.New("Invalid query format.") @@ -75,6 +77,8 @@ func (qp *QueryParser) walk(tokens []Token) (*Query, error) { } q := new(Query) + id := uuid.RandomUUID() + q.Id = &id q.Operation = tokens[0].Text // scan from open to close paren diff --git a/query/planner.go b/query/planner.go index d66a1383d..b0ee427ff 100644 --- a/query/planner.go +++ b/query/planner.go @@ -219,6 +219,10 @@ type SetQueryStep struct { ProfileId uint64 } +type SetQueryResult struct { + *BaseQueryResult +} + // QueryTree for SET queries type SetQueryTree struct { bitmap *db.Bitmap @@ -238,6 +242,7 @@ func (qt *SetQueryTree) getLocation(d *db.Database) *db.Location { /////////////////////////////////////////////////////////////////////////////////////////////////// func init() { + gob.Register(SetQueryResult{}) gob.Register(GetQueryResult{}) gob.Register(CatQueryResult{}) gob.Register(UnionQueryResult{}) @@ -401,5 +406,6 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Locati // Transforms Query into QueryTree and flattens to QueryPlan object func (qp *QueryPlanner) Plan(query *Query, id *uuid.UUID, destination *db.Location) *QueryPlan { queryTree := qp.buildTree(query, -1) - return qp.flatten(queryTree, id, destination) + //return qp.flatten(queryTree, id, destination) // TODO: remove the "id" parameter, since we are using the query.Id as the value + return qp.flatten(queryTree, query.Id, destination) } diff --git a/query/planner_test.go b/query/planner_test.go index 3bb714890..97a67ba8e 100644 --- a/query/planner_test.go +++ b/query/planner_test.go @@ -13,16 +13,19 @@ import ( func TestQueryPlanner(t *testing.T) { Convey("Basic query plan", t, func() { + id1 := uuid.RandomUUID() bm1 := db.Bitmap{10, "general"} inputs1 := []QueryInput{&bm1} - query1 := Query{"get", inputs1, 0} + query1 := Query{&id1, "get", inputs1, 0} + id2 := uuid.RandomUUID() bm2 := db.Bitmap{20, "general"} inputs2 := []QueryInput{&bm2} - query2 := Query{"get", inputs2, 0} + query2 := Query{&id2, "get", inputs2, 0} + id3 := uuid.RandomUUID() inputs := []QueryInput{&query1, &query2} - query := Query{"union", inputs, 0} + query := Query{&id3, "union", inputs, 0} /* bm1 := db.Bitmap{10, "general"} diff --git a/query/query.go b/query/query.go index cf1ab9f0d..a82711ffb 100644 --- a/query/query.go +++ b/query/query.go @@ -2,6 +2,7 @@ package query import ( "pilosa/db" + "tux21b.org/v1/gocql/uuid" ) @@ -12,6 +13,7 @@ type QueryResults struct { } type Query struct { + Id *uuid.UUID Operation string Inputs []QueryInput //"strconv" // Represents a parsed query. Inputs can be Query or Bitmap objects @@ -21,13 +23,41 @@ type Query struct { func QueryPlanForPQL(database *db.Database, pql string, destination *db.Location) *QueryPlan { tokens := Lex(pql) + return QueryPlanForTokens(database, tokens, destination) +} + +func QueryForPQL(pql string) *Query { + tokens := Lex(pql) + return QueryForTokens(tokens) +} + +func QueryForTokens(tokens []Token) *Query { query_parser := QueryParser{} query, err := query_parser.Parse(tokens) if err != nil { panic(err) } + return query +} + +func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.Location) *QueryPlan { + query := QueryForTokens(tokens) + return QueryPlanForQuery(database, query, destination) + /* + //spew.Dump(query) + query_planner := QueryPlanner{Database: database} + id := uuid.RandomUUID() + query_plan := query_planner.Plan(query, &id, destination) + //spew.Dump(query_plan) + return query_plan + */ +} + +func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Location) *QueryPlan { + //spew.Dump(query) query_planner := QueryPlanner{Database: database} id := uuid.RandomUUID() query_plan := query_planner.Plan(query, &id, destination) + //spew.Dump(query_plan) return query_plan } From e3b417539a0a272a5b6b3c1bc23e7e6e6302b664 Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 10 Jan 2014 14:32:21 -0600 Subject: [PATCH 2/4] removed some of the slice range copy --- executor/executor.go | 41 +++++++++++++++++------------------------ 1 file changed, 17 insertions(+), 24 deletions(-) diff --git a/executor/executor.go b/executor/executor.go index e208447f4..8b10a755f 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -62,17 +62,10 @@ func (self *Executor) RunQueryTest(database_name string, pql string) string { return pql } -type queryItem struct { - label string - pql string -} -type queryItemCopy struct { +type queryListItem struct { id *uuid.UUID label string -} -type QueryItemResult struct { - Label string - Result interface{} + pql string } func (self *Executor) runQuery(database *db.Database, qry *query.Query) { @@ -127,30 +120,30 @@ func (self *Executor) RunPQL(database_name string, pql string) interface{} { spew.Dump(file_data) // CUSTOM QUERY LIST /////////////// - var query_list []queryItem - query_list = append(query_list, queryItem{"set1", "set(20, 1)"}) - query_list = append(query_list, queryItem{"set2", "set(20, 1)"}) - query_list = append(query_list, queryItem{"set3", "set(20, 2)"}) - query_list = append(query_list, queryItem{"set4", "set(20, 3)"}) - query_list = append(query_list, queryItem{"set5", "set(20, 4)"}) - query_list = append(query_list, queryItem{"count", "count(get(20))"}) - var query_list_copy []*queryItemCopy + var query_list []queryListItem + query_list = append(query_list, queryListItem{label: "set1", pql: "set(20, 1)"}) + query_list = append(query_list, queryListItem{label: "set2", pql: "set(20, 1)"}) + query_list = append(query_list, queryListItem{label: "set3", pql: "set(20, 2)"}) + query_list = append(query_list, queryListItem{label: "set4", pql: "set(20, 3)"}) + query_list = append(query_list, queryListItem{label: "set5", pql: "set(20, 4)"}) + query_list = append(query_list, queryListItem{label: "count", pql: "count(get(20))"}) //////////////////////////////////// - for _, qi := range query_list { - qry := query.QueryForPQL(qi.pql) + for i, _ := range query_list { + qry := query.QueryForPQL(query_list[i].pql) go self.runQuery(database, qry) - query_list_copy = append(query_list_copy, &queryItemCopy{qry.Id, qi.label}) + query_list[i].id = qry.Id } - var final_result []*QueryItemResult - for _, qlc := range query_list_copy { - final, err := self.service.Hold.Get(qlc.id, 10) + // technically, this is blocking on the hold for each query, but it may be ok, since we need them all for the final anyway. + final_result := make(map[string]interface{}) + for i, _ := range query_list { + final, err := self.service.Hold.Get(query_list[i].id, 10) if err != nil { spew.Dump(err) } spew.Dump(final) - final_result = append(final_result, &QueryItemResult{qlc.label, final}) + final_result[query_list[i].label] = final } return final_result From 66f3aeb21c56580257f4ceb0d828ddb4fb4ccf6f Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 10 Jan 2014 17:51:19 -0600 Subject: [PATCH 3/4] add support for macros, with Otto javascript to insert filters into macro --- executor/executor.go | 33 ++++++----------------------- executor/utils.go | 50 ++++++++++++++++++++++++++++++++++++++++++++ query/query.go | 31 ++++++++++++++++++--------- 3 files changed, 77 insertions(+), 37 deletions(-) create mode 100644 executor/utils.go diff --git a/executor/executor.go b/executor/executor.go index 8b10a755f..33ca5baab 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -2,14 +2,12 @@ package executor import ( "fmt" - "io/ioutil" "log" "pilosa/config" "pilosa/core" "pilosa/db" "pilosa/query" "pilosa/util" - "tux21b.org/v1/gocql/uuid" "github.com/davecgh/go-spew/spew" ) @@ -62,12 +60,6 @@ func (self *Executor) RunQueryTest(database_name string, pql string) string { return pql } -type queryListItem struct { - id *uuid.UUID - label string - pql string -} - func (self *Executor) runQuery(database *db.Database, qry *query.Query) { process, err := self.service.GetProcess() if err != nil { @@ -113,37 +105,24 @@ func (self *Executor) RunPQL(database_name string, pql string) interface{} { } else { macros_dir := config.Get("macros").(string) macros_file := macros_dir + "/" + outer_token + ".js" - file_data, err := ioutil.ReadFile(macros_file) - if err != nil { - spew.Dump(err) - } - spew.Dump(file_data) - - // CUSTOM QUERY LIST /////////////// - var query_list []queryListItem - query_list = append(query_list, queryListItem{label: "set1", pql: "set(20, 1)"}) - query_list = append(query_list, queryListItem{label: "set2", pql: "set(20, 1)"}) - query_list = append(query_list, queryListItem{label: "set3", pql: "set(20, 2)"}) - query_list = append(query_list, queryListItem{label: "set4", pql: "set(20, 3)"}) - query_list = append(query_list, queryListItem{label: "set5", pql: "set(20, 4)"}) - query_list = append(query_list, queryListItem{label: "count", pql: "count(get(20))"}) - //////////////////////////////////// + filter := query.TokensToString(tokens) + query_list := GetMacro(macros_file, filter).(query.QueryList) for i, _ := range query_list { - qry := query.QueryForPQL(query_list[i].pql) + qry := query.QueryForPQL(query_list[i].PQL) go self.runQuery(database, qry) - query_list[i].id = qry.Id + query_list[i].Id = qry.Id } // technically, this is blocking on the hold for each query, but it may be ok, since we need them all for the final anyway. final_result := make(map[string]interface{}) for i, _ := range query_list { - final, err := self.service.Hold.Get(query_list[i].id, 10) + final, err := self.service.Hold.Get(query_list[i].Id, 10) if err != nil { spew.Dump(err) } spew.Dump(final) - final_result[query_list[i].label] = final + final_result[query_list[i].Label] = final } return final_result diff --git a/executor/utils.go b/executor/utils.go new file mode 100644 index 000000000..b4ee7cb6e --- /dev/null +++ b/executor/utils.go @@ -0,0 +1,50 @@ +package executor + +import ( + "io/ioutil" + "pilosa/query" + + "github.com/davecgh/go-spew/spew" + "github.com/robertkrimen/otto" +) + +func GetMacro(file_name string, filter string) interface{} { + + file_data, err := ioutil.ReadFile(file_name) + if err != nil { + spew.Dump(err) + } + s := string(file_data[:]) + + js := "query_list = (function (filter){" + s + "})('" + filter + "');" + + Otto := otto.New() + Otto.Run(js) + query_objects, err := Otto.Get("query_list") + + query_list_interface, err := query_objects.Export() + if err != nil { + spew.Dump(err) + } + + var query_list query.QueryList + + // ql is []interface{} + switch ql := query_list_interface.(type) { + case []interface{}: + spew.Dump("INTERFACE", ql) + // q is map[string]interface{} + for i, _ := range ql { + q := ql[i].(map[string]interface{}) + spew.Dump(ql[i]) + spew.Dump(q) + spew.Dump(q["label"].(string)) + spew.Dump(q["pql"].(string)) + query_list = append(query_list, query.QueryListItem{Label: q["label"].(string), PQL: q["pql"].(string)}) + } + default: + spew.Dump("DEFAULT") + } + + return query_list +} diff --git a/query/query.go b/query/query.go index a82711ffb..cba9e5a69 100644 --- a/query/query.go +++ b/query/query.go @@ -2,6 +2,8 @@ package query import ( "pilosa/db" + "strings" + "github.com/davecgh/go-spew/spew" "tux21b.org/v1/gocql/uuid" ) @@ -12,6 +14,14 @@ type QueryResults struct { Data interface{} } +type QueryList []QueryListItem + +type QueryListItem struct { + Id *uuid.UUID + Label string + PQL string +} + type Query struct { Id *uuid.UUID Operation string @@ -43,21 +53,22 @@ func QueryForTokens(tokens []Token) *Query { func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.Location) *QueryPlan { query := QueryForTokens(tokens) return QueryPlanForQuery(database, query, destination) - /* - //spew.Dump(query) - query_planner := QueryPlanner{Database: database} - id := uuid.RandomUUID() - query_plan := query_planner.Plan(query, &id, destination) - //spew.Dump(query_plan) - return query_plan - */ } func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Location) *QueryPlan { - //spew.Dump(query) query_planner := QueryPlanner{Database: database} id := uuid.RandomUUID() query_plan := query_planner.Plan(query, &id, destination) - //spew.Dump(query_plan) return query_plan } + +func TokensToString(tokens []Token) string { + var str []string + for i, _ := range tokens { + spew.Dump(tokens[i].Text) + str = append(str, tokens[i].Text) + } + // for now, we're just using this function to pull the filter out of the outer function "outerfunc(filter)" + str = str[2 : len(str)-1] + return strings.Join(str, "") +} From 52f1b3661e2033b638106db235f6db55e8e35db5 Mon Sep 17 00:00:00 2001 From: travisturner Date: Mon, 13 Jan 2014 11:16:23 -0600 Subject: [PATCH 4/4] rename QueryList to PqlList for clarity --- executor/executor.go | 4 ++-- executor/utils.go | 4 ++-- query/query.go | 4 ++-- 3 files changed, 6 insertions(+), 6 deletions(-) diff --git a/executor/executor.go b/executor/executor.go index 33ca5baab..f421e7d0c 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -103,10 +103,10 @@ func (self *Executor) RunPQL(database_name string, pql string) interface{} { return final } else { - macros_dir := config.Get("macros").(string) + macros_dir := config.GetString("macros") macros_file := macros_dir + "/" + outer_token + ".js" filter := query.TokensToString(tokens) - query_list := GetMacro(macros_file, filter).(query.QueryList) + query_list := GetMacro(macros_file, filter).(query.PqlList) for i, _ := range query_list { qry := query.QueryForPQL(query_list[i].PQL) diff --git a/executor/utils.go b/executor/utils.go index b4ee7cb6e..d6a989793 100644 --- a/executor/utils.go +++ b/executor/utils.go @@ -27,7 +27,7 @@ func GetMacro(file_name string, filter string) interface{} { spew.Dump(err) } - var query_list query.QueryList + var query_list query.PqlList // ql is []interface{} switch ql := query_list_interface.(type) { @@ -40,7 +40,7 @@ func GetMacro(file_name string, filter string) interface{} { spew.Dump(q) spew.Dump(q["label"].(string)) spew.Dump(q["pql"].(string)) - query_list = append(query_list, query.QueryListItem{Label: q["label"].(string), PQL: q["pql"].(string)}) + query_list = append(query_list, query.PqlListItem{Label: q["label"].(string), PQL: q["pql"].(string)}) } default: spew.Dump("DEFAULT") diff --git a/query/query.go b/query/query.go index cba9e5a69..3fd1b4a89 100644 --- a/query/query.go +++ b/query/query.go @@ -14,9 +14,9 @@ type QueryResults struct { Data interface{} } -type QueryList []QueryListItem +type PqlList []PqlListItem -type QueryListItem struct { +type PqlListItem struct { Id *uuid.UUID Label string PQL string