From 74f47e32eeaa705325253f47a139de5b7e771c6d Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 11 Aug 2014 19:47:25 +0000 Subject: [PATCH] range pass1 --- core/query.go | 23 ++++++++++++++++ executor/executor.go | 4 +-- query/planner.go | 63 +++++++++++++++++++++++++++++++++----------- 3 files changed, 73 insertions(+), 17 deletions(-) diff --git a/core/query.go b/core/query.go index cf645d9a3..e075cb50b 100644 --- a/core/query.go +++ b/core/query.go @@ -272,3 +272,26 @@ func (self *Service) SetQueryStepHandler(msg *db.Message) { result_message := db.Message{Data: query.SetQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} self.Transport.Send(&result_message, qs.Destination.ProcessId) } + +func (self *Service) RangeQueryStepHandler(msg *db.Message) { + qs := msg.Data.(query.RangeQueryStep) + //spew.Dump("RANDE QUERYSTEP") + + bh, err := self.Index.Range(qs.Location.FragmentId, qs.Bitmap.Id, qs.Start, qs.End) + if err != nil { + spew.Dump(err) + } + + var result interface{} + if qs.LocIsDest() { + result = bh + } else { + bm, err := self.Index.GetBytes(qs.Location.FragmentId, bh) + if err != nil { + spew.Dump(err) + } + result = bm + } + result_message := db.Message{Data: query.RangeQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + self.Transport.Send(&result_message, qs.Destination.ProcessId) +} diff --git a/executor/executor.go b/executor/executor.go index db8f5275f..d383a932b 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -44,10 +44,10 @@ func (self *Executor) NewJob(job *db.Message) { self.service.GetQueryStepHandler(job) case query.SetQueryStep: self.service.SetQueryStepHandler(job) + case query.RangeQueryStep: + self.service.RangeQueryStepHandler(job) // case query.MaskQueryStep: // self.service.MaskQueryStepHandler(job) - // case query.RangeQueryStep: - // self.service.RangeQueryStepHandler(job) default: fmt.Println("unknown") } diff --git a/query/planner.go b/query/planner.go index df441f48d..b5bb77834 100644 --- a/query/planner.go +++ b/query/planner.go @@ -390,6 +390,26 @@ type QueryTree interface { getLocation(d *db.Database) (*db.Location, error) } +func validateRange(Args map[string]interface{}) error { + _, ok := Args["id"] + if !ok { + return errors.New("missing bitmap id") + } + _, ok = Args["frame"] + if !ok { + return errors.New("missing frame ") + } + _, ok = Args["start"] + if !ok { + return errors.New("missing start time") + } + _, ok = Args["end"] + if !ok { + return errors.New("missing end time") + } + return nil +} + // Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { var tree QueryTree @@ -426,6 +446,19 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) if query.Operation == "get" { tree = &GetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), 0}, slice} return tree, nil + } else if query.Operation == "range" { + err := validateRange(query.Args) + if err != nil { + return nil, err + } + tree = &RangeQueryTree{ + &db.Bitmap{ + query.Args["id"].(uint64), + query.Args["frame"].(string), 0}, + slice, + query.Args["start"].(time.Time), + query.Args["end"].(time.Time)} + return tree, nil } else if query.Operation == "count" { subquery, err := self.buildTree(&query.Subqueries[0], slice) if err != nil { @@ -581,15 +614,15 @@ func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Loca step := MaskQueryStep{&BaseQueryStep{id, "mask", loc, location}, mask.start, mask.end} plan := QueryPlan{step} return &plan, nil - } else if rang, ok := qt.(*RangeQueryTree); ok { - loc, err := rang.getLocation(self.Database) - if err != nil { - return nil, err - } - step := RangeQueryStep{&BaseQueryStep{id, "range", loc, location}, rang.bitmap, rang.start, rang.end} - plan := QueryPlan{step} - return &plan, nil */ + } else if rang, ok := qt.(*RangeQueryTree); ok { + loc, err := rang.getLocation(self.Database) + if err != nil { + return nil, err + } + step := RangeQueryStep{&BaseQueryStep{id, "range", loc, location}, rang.bitmap, rang.start, rang.end} + plan := QueryPlan{step} + return &plan, nil } else if set, ok := qt.(*SetQueryTree); ok { loc, err := set.getLocation(self.Database) if err != nil { @@ -673,8 +706,8 @@ func (qt *MaskQueryTree) getLocation(d *db.Database) (*db.Location, error) { /////////////////////////////////////////////////////////////////////////////////////////////////// type RangeQueryStep struct { *BaseQueryStep - bitmap *db.Bitmap - start, end time.Time + Bitmap *db.Bitmap + Start, End time.Time } type RangeQueryResult struct { @@ -683,17 +716,17 @@ type RangeQueryResult struct { // QueryTree for Mask queries type RangeQueryTree struct { + bitmap *db.Bitmap + slice int + start, end time.Time } -// Uses consistent hashing function to select node containing data for GET operation -/* -func (self *RangeQueryTree) getLocation(d *db.Database) (*db.Location, error) { +func (qt *RangeQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create) - fragment, err := d.GetFragmentForBitmap(slice, self.bitmap) + fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { log.Println("GetFragmenForBitmapFailed GetQueryTree", slice) return nil, err } return fragment.GetLocation(), nil } -*/