diff --git a/core/http.go b/core/http.go index 7d6c1e8c2..24c0b3f3a 100644 --- a/core/http.go +++ b/core/http.go @@ -6,6 +6,7 @@ import ( "encoding/base64" "encoding/json" "fmt" + "io/ioutil" "log" "net/http" "net/http/httputil" @@ -287,6 +288,7 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) { http.Error(w, "Provide a valid query string (pql)", http.StatusNotFound) return } + _, bits := r.Form["bits"] results, err := self.service.Executor.RunPQL(database_name, pql) switch r := results.(type) { //a hack to handle empty sets @@ -295,6 +297,15 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) { results = []int{} } + case []byte: //can i figure out the type of the compressed string here? + if bits { + reader, _ := gzip.NewReader(bytes.NewReader(r)) + b, _ := ioutil.ReadAll(reader) + result := index.NewBitmap() + result.FromBytes(b) + results = result.Bits() + } + } if results == nil { diff --git a/core/query.go b/core/query.go index c51d15862..c88345721 100644 --- a/core/query.go +++ b/core/query.go @@ -11,6 +11,14 @@ import ( "github.com/davecgh/go-spew/spew" ) +/* +func (self *Service) RecallQueryStepHandler(msg *db.Message) { + qs := msg.Data.(query.RecallQueryStep) + + result_message := db.Message{Data: query.RecallQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: count}}} + self.Transport.Send(&result_message, qs.Destination.ProcessId) +} +*/ func (self *Service) CountQueryStepHandler(msg *db.Message) { //spew.Dump("COUNT QUERYSTEP") qs := msg.Data.(query.CountQueryStep) @@ -150,7 +158,7 @@ func (self *Service) StashQueryStepHandler(msg *db.Message) { }(input, part) } //just collect all the handles and return them - result := query.Stash{make([]query.CacheItem, 0)} + result := query.NewStash() //query.Stash{make([]query.CacheItem, 0), false} for i := 0; i < num_parts; i++ { value := <-part @@ -176,7 +184,6 @@ func (self *Service) StashQueryStepHandler(msg *db.Message) { func (self *Service) CatQueryStepHandler(msg *db.Message) { qs := msg.Data.(query.CatQueryStep) - spew.Dump("CAT QUERYSTEP") var handles []index.BitmapHandle return_type := "bitmap-handles" var sum uint64 diff --git a/db/topology.go b/db/topology.go index 824ee550a..e5916d5ea 100644 --- a/db/topology.go +++ b/db/topology.go @@ -299,6 +299,15 @@ func (self *FrameSliceIntersect) GetFragments() []*Fragment { return self.fragments } +func (d *Database) GetFragment(fragment_id util.SUUID) (*Fragment, error) { + for _, fsi := range d.frame_slice_intersects { + f, err := fsi.GetFragment(fragment_id) + if err == nil { + return f, nil + } + } + return nil, FragmentDoesNotExistError +} func (self *FrameSliceIntersect) GetFragment(fragment_id util.SUUID) (*Fragment, error) { for _, fragment := range self.fragments { if fragment.id == fragment_id { diff --git a/executor/executor.go b/executor/executor.go index a6f839db9..ae1db3673 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -48,8 +48,10 @@ func (self *Executor) NewJob(job *db.Message) { self.service.RangeQueryStepHandler(job) case query.StashQueryStep: self.service.StashQueryStepHandler(job) - // case query.MaskQueryStep: - // self.service.MaskQueryStepHandler(job) + //case query.RecallQueryStep: + // self.service.RecallQueryStepHandler(job) + // case query.MaskQueryStep: + // self.service.MaskQueryStepHandler(job) default: fmt.Println("unknown") } @@ -127,7 +129,7 @@ func (self *Executor) RunPQL(database_name string, pql string) (interface{}, err } return final, nil - } else { + } else { //want to refactor this down to just RunPlugin(tokens) plugins_dir := config.GetString("plugins") plugins_file := plugins_dir + "/" + outer_token + ".js" filter := query.TokensToString(tokens) diff --git a/query/planner.go b/query/planner.go index 85923075c..35d7d9674 100644 --- a/query/planner.go +++ b/query/planner.go @@ -436,6 +436,12 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) return tree, nil } + if query.Operation == "recall" { + + tree = &RecallQueryTree{query.Args["stash"].(Stash)} + return tree, nil + } + // handle the remaining operations, taking slice into consideration //I'm kinda thinking for stash it needs a differnt handler..i guess for now I'll see if i can use the cathandler if slice == -1 { @@ -581,6 +587,15 @@ func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Loca plan = append(plan, *subq_steps...) } plan = append(plan, step) + /*} else if recall, ok := qt.(*RecallQueryTree); ok { + //need to go through the stash and fetch + //from self.Database i can fetch the locations + step := RecallQueryStep{&BaseQueryStep{id, "recall", loc, location}, recall.Stash} + + plan := QueryPlan{step} + return &plan, nil + */ + } else if stash, ok := qt.(*StashQueryTree); ok { inputs := make([]*util.GUID, len(stash.subqueries)) loc, err := stash.getLocation(self.Database) @@ -801,6 +816,10 @@ type Stash struct { incomplete bool } +func NewStash() Stash { + return Stash{make([]CacheItem, 0), false} +} + func (st *Stash) Add(i util.SUUID) { item := CacheItem{i, 0} st.Stash = append(st.Stash, item) @@ -846,3 +865,19 @@ func (qt *StashQueryTree) getLocation(d *db.Database) (*db.Location, error) { } return qt.location, err } + +type RecallQueryStep struct { + *BaseQueryStep + Stash Stash +} +type RecallQueryTree struct { + Stash Stash +} + +func (rqt *RecallQueryTree) getLocation(d *db.Database) (*db.Location, error) { + return nil, nil +} + +type RecallQueryResult struct { + *BaseQueryResult +}