diff --git a/core/http.go b/core/http.go index 79cc60fc2..9afd98fe8 100644 --- a/core/http.go +++ b/core/http.go @@ -515,10 +515,7 @@ func (self *WebService) HandleBit(w http.ResponseWriter, r *http.Request, ToSet http.Error(w, "Request To large", http.StatusBadRequest) return } - result := false - remoteSetBit := NewRemoteSetBit(self.service) - var database *db.Database - olddbs := "" + //remoteSetBit := NewRemoteSetBit(self.service) for _, obj := range args { if obj["profile_id"] == nil { http.Error(w, "Missing Profile", http.StatusBadRequest) @@ -546,38 +543,17 @@ func (self *WebService) HandleBit(w http.ResponseWriter, r *http.Request, ToSet } t = float64(obj["filter"].(float64)) filter := uint64(t) - if dbs != olddbs { - database = self.service.Cluster.GetOrCreateDatabase(dbs) - olddbs = dbs - } - frag, err := database.GetFragmentFromProfile(frame, profile_id) - if err != nil { - //no fragment - if ToSet { - self.service.TopologyMapper.MakeFragments(dbs, db.GetSlice(profile_id)) - time.Sleep(2 * time.Second) - } - break - } - isLocal := util.Equal(frag.GetProcessId(), self.service.Id) for bitmap_id := range bitmaps(frame, obj) { - if isLocal { - // The Local Route - if ToSet { - result, _ = self.service.Index.SetBit(frag.GetId(), bitmap_id, profile_id, filter) - } else { - result, _ = self.service.Index.ClearBit(frag.GetId(), bitmap_id, profile_id) - } - bundle := SBResult{bitmap_id, frame, filter, profile_id, result} - results = append(results, bundle) + var pql string + if ToSet { + pql = fmt.Sprintf("set(%d, %s, %d, %d)", bitmap_id, frame, filter, profile_id) } else { - - remoteSetBit.Add(frag, bitmap_id, profile_id, filter, frame, ToSet) + pql = fmt.Sprintf("clear(%d, %s, %d, %d)", bitmap_id, frame, filter, profile_id) } - - //result, err := self.service.Executor.RunPQL(db, pql) - //pql := fmt.Sprintf("set(%d, %s, %d, %d)", bitmap_id, frame, filter, profile_id) + result, err := self.service.Executor.RunPQL(dbs, pql) + bundle := SBResult{bitmap_id, frame, filter, profile_id, result} + results = append(results, bundle) if err != nil { log.Warn("Error running set_bit", dbs, frame, profile_id, ToSet) @@ -587,8 +563,8 @@ func (self *WebService) HandleBit(w http.ResponseWriter, r *http.Request, ToSet } } - remoteSetBit.Request() - results = remoteSetBit.MergeResults(results) + // remoteSetBit.Request() + // results = remoteSetBit.MergeResults(results) encoder := json.NewEncoder(w) err = encoder.Encode(results) if err != nil { diff --git a/core/query.go b/core/query.go index 4c0ab5039..329fa6f92 100644 --- a/core/query.go +++ b/core/query.go @@ -321,6 +321,15 @@ func (self *Service) SetQueryStepHandler(msg *db.Message) { self.Transport.Send(&result_message, qs.Destination.ProcessId) } +func (self *Service) ClearQueryStepHandler(msg *db.Message) { + //spew.Dump("SET QUERYSTEP") + qs := msg.Data.(query.ClearQueryStep) + result, _ := self.Index.ClearBit(qs.Location.FragmentId, qs.Bitmap.Id, qs.ProfileId) + + result_message := db.Message{Data: query.ClearQueryResult{&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") diff --git a/executor/executor.go b/executor/executor.go index 699415e4b..770aeb5b9 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -44,6 +44,8 @@ func (self *Executor) NewJob(job *db.Message) { self.service.GetQueryStepHandler(job) case query.SetQueryStep: self.service.SetQueryStepHandler(job) + case query.ClearQueryStep: + self.service.ClearQueryStepHandler(job) case query.RangeQueryStep: self.service.RangeQueryStepHandler(job) case query.StashQueryStep: @@ -115,7 +117,7 @@ func (self *Executor) RunPQL(database_name string, pql string) (interface{}, err database := self.service.Cluster.GetOrCreateDatabase(database_name) // see if the outer query function is a custom query - reserved_functions := stringSlice{"get", "set", "union", "intersect", "difference", "count", "top-n", "mask", "range", "stash", "recall"} + reserved_functions := stringSlice{"get", "set", "clear", "union", "intersect", "difference", "count", "top-n", "mask", "range", "stash", "recall"} tokens, err := query.Lex(pql) if err != nil { return nil, err diff --git a/index/fragment_container.go b/index/fragment_container.go index 692f8ca16..e2f09f618 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -259,6 +259,7 @@ func (self *FragmentContainer) FromBytes(frag_id util.SUUID, bytes []byte) (Bitm } func (self *FragmentContainer) SetBit(frag_id util.SUUID, bitmap_id uint64, pos uint64, category uint64) (bool, error) { + log.Trace("SetBit", frag_id, bitmap_id, pos, category) if fragment, found := self.GetFragment(frag_id); found { request := NewSetBit(bitmap_id, pos, category) fragment.requestChan <- request @@ -271,6 +272,7 @@ func (self *FragmentContainer) SetBit(frag_id util.SUUID, bitmap_id uint64, pos } func (self *FragmentContainer) ClearBit(frag_id util.SUUID, bitmap_id uint64, pos uint64) (bool, error) { + log.Trace("ClearBit", frag_id, bitmap_id, pos) if fragment, found := self.GetFragment(frag_id); found { request := NewClearBit(bitmap_id, pos) fragment.requestChan <- request diff --git a/query/parser.go b/query/parser.go index 16af7a979..75b7a93c2 100644 --- a/query/parser.go +++ b/query/parser.go @@ -159,6 +159,31 @@ ArgLoop: default: return nil, fmt.Errorf("Unexpected argument! (%v)", token) } + case "clear": + switch len(query.Args) { + case 0: + i, err := strconv.ParseUint(token.Text, 10, 64) + if err != nil { + return nil, fmt.Errorf("Expecting integer id! (%v)", err) + } + query.Args["id"] = i + case 1: + query.Args["frame"] = token.Text + case 2: + i, err := strconv.ParseUint(token.Text, 10, 64) + if err != nil { + return nil, fmt.Errorf("Expecting integer id! (%v)", err) + } + query.Args["filter"] = i + case 3: + i, err := strconv.ParseUint(token.Text, 10, 64) + if err != nil { + return nil, fmt.Errorf("Expecting integer id! (%v)", err) + } + query.Args["profile_id"] = i + default: + return nil, fmt.Errorf("Unexpected argument! (%v)", token) + } case "set": switch len(query.Args) { case 0: @@ -184,7 +209,6 @@ ArgLoop: default: return nil, fmt.Errorf("Unexpected argument! (%v)", token) } - case "top-n": switch len(query.Args) { case 0: diff --git a/query/parser_test.go b/query/parser_test.go index dfa48d53f..212741d58 100644 --- a/query/parser_test.go +++ b/query/parser_test.go @@ -18,6 +18,16 @@ func TestQueryParser(t *testing.T) { So(query.Operation, ShouldEqual, "get") So(query.Args, ShouldResemble, map[string]interface{}{"id": uint64(10), "frame": "general"}) }) + Convey("Basic parse - clear()", t, func() { + tokens, err := Lex("clear(10, general, 0, 20)") + So(err, ShouldBeNil) + + query, err := Parse(tokens) + So(err, ShouldBeNil) + spew.Dump(query) + So(query.Operation, ShouldEqual, "clear") + So(query.Args, ShouldResemble, map[string]interface{}{"id": uint64(10), "frame": "general", "filter": uint64(0), "profile_id": uint64(20)}) + }) Convey("Basic parse - set()", t, func() { tokens, err := Lex("set(10, general, 0, 20)") So(err, ShouldBeNil) diff --git a/query/planner.go b/query/planner.go index c2abfcf34..d81c9f425 100644 --- a/query/planner.go +++ b/query/planner.go @@ -398,6 +398,7 @@ func (qt *SetQueryTree) getLocation(d *db.Database) (*db.Location, error) { func init() { gob.Register(BaseQueryResult{}) gob.Register(SetQueryResult{}) + gob.Register(ClearQueryResult{}) gob.Register(GetQueryResult{}) gob.Register(RangeQueryResult{}) gob.Register(CatQueryResult{}) @@ -412,6 +413,7 @@ func init() { gob.Register(Stash{}) gob.Register(SetQueryStep{}) + gob.Register(ClearQueryStep{}) gob.Register(GetQueryStep{}) gob.Register(RangeQueryStep{}) gob.Register(CatQueryStep{}) @@ -462,14 +464,24 @@ func validateRange(Args map[string]interface{}) error { func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { log.Trace("QueryPlanner.buildTree", query, slice) var tree QueryTree - + spew.Dump(query) // handle SET operation regardless of the slice if query.Operation == "set" { + println("SET") + log.Warn("SET") tree = &SetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), query.Args["filter"].(uint64)}, query.Args["profile_id"].(uint64)} return tree, nil } + if query.Operation == "clear" { + + println("CLEAR") + log.Warn("CLEAR") + tree = &ClearQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), query.Args["filter"].(uint64)}, query.Args["profile_id"].(uint64)} + return tree, nil + } + if query.Operation == "recall" { tree = &RecallQueryTree{query.Args["stash"].(Stash)} @@ -733,6 +745,14 @@ func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Loca step := SetQueryStep{&BaseQueryStep{id, "set", loc, location}, set.bitmap, set.profile_id} plan := QueryPlan{step} return &plan, nil + } else if clear, ok := qt.(*ClearQueryTree); ok { + loc, err := clear.getLocation(self.Database) + if err != nil { + return nil, err + } + step := ClearQueryStep{&BaseQueryStep{id, "clear", loc, location}, clear.bitmap, clear.profile_id} + plan := QueryPlan{step} + return &plan, nil } else if cnt, ok := qt.(*CountQueryTree); ok { sub_id := util.RandomUUID() loc, err := cnt.getLocation(self.Database) @@ -924,3 +944,37 @@ func (rqt *RecallQueryTree) getLocation(d *db.Database) (*db.Location, error) { type RecallQueryResult struct { *BaseQueryResult } + +type ClearQueryStep struct { + *BaseQueryStep + Bitmap *db.Bitmap + ProfileId uint64 +} + +type ClearQueryResult struct { + *BaseQueryResult +} + +type ClearQueryTree struct { + bitmap *db.Bitmap + profile_id uint64 +} + +// Uses consistent hashing function to select node containing data for GET operation +func (qt *ClearQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("ClearQueryTree.getLocation", d) + // check here for supported frames + if !d.IsValidFrame(qt.bitmap.FrameType) { + return nil, NewInvalidFrame(d.Name, qt.bitmap.FrameType) + } + slice, err := d.GetSliceForProfile(qt.profile_id) + if err != nil { + return nil, NewFragmentNotFound(d.Name, qt.bitmap.FrameType, db.GetSlice(qt.profile_id)) + } + fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) + if err != nil { + log.Warn("NOT FOUND:", slice, qt.bitmap) + return nil, NewFragmentNotFound(d.Name, qt.bitmap.FrameType, db.GetSlice(qt.profile_id)) + } + return fragment.GetLocation(), nil +}