refactor setbit/clearbit

This commit is contained in:
Todd Gruben 2015-03-04 20:41:37 +00:00
parent 1dffa07fd6
commit 2524ff2a51
7 changed files with 114 additions and 37 deletions

View file

@ -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 {

View file

@ -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")

View file

@ -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

View file

@ -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

View file

@ -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:

View file

@ -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)

View file

@ -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
}