diff --git a/core/core_test.go b/core/core_test.go new file mode 100644 index 000000000..d48ca3fea --- /dev/null +++ b/core/core_test.go @@ -0,0 +1,14 @@ +package core + +import ( + "testing" + + . "github.com/smartystreets/goconvey/convey" +) + +func TestCompile(t *testing.T) { + Convey("has asm", t, func() { + + So(canCompile(), ShouldEqual, true) + }) +} diff --git a/core/query.go b/core/query.go index 09eb4b732..cf645d9a3 100644 --- a/core/query.go +++ b/core/query.go @@ -4,6 +4,7 @@ import ( "pilosa/db" "pilosa/index" "pilosa/query" + "pilosa/util" "sort" "github.com/davecgh/go-spew/spew" @@ -30,39 +31,6 @@ func (self *Service) CountQueryStepHandler(msg *db.Message) { self.Transport.Send(&result_message, qs.Destination.ProcessId) } -func (self *Service) TopNQueryStepHandler(msg *db.Message) { - //spew.Dump("TOP-N QUERYSTEP") - qs := msg.Data.(query.TopNQueryStep) - //spew.Dump(qs) - //need categories in qs I just added so it would compile - var categoryleaves []uint64 - input := qs.Input - value, _ := self.Hold.Get(input, 10) - var bh index.BitmapHandle - switch val := value.(type) { - case index.BitmapHandle: - bh = val - case []byte: - bh, _ = self.Index.FromBytes(qs.Location.FragmentId, val) - } - - categoryleaves = qs.Filters - - /* - spew.Dump(bh, qs.N*2, categoryleaves) - result_message := db.Message{Data: "foobar"} - self.Transport.Send(&result_message, qs.Destination.ProcessId) - */ - - topn, err := self.Index.TopN(qs.Location.FragmentId, bh, qs.N*2, categoryleaves) - if err != nil { - spew.Dump(err) - } - result_message := db.Message{Data: query.TopNQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: topn}}} - self.Transport.Send(&result_message, qs.Destination.ProcessId) - -} - func (self *Service) UnionQueryStepHandler(msg *db.Message) { //spew.Dump("UNION QUERYSTEP") qs := msg.Data.(query.UnionQueryStep) @@ -174,10 +142,28 @@ func (self *Service) CatQueryStepHandler(msg *db.Message) { var handles []index.BitmapHandle return_type := "bitmap-handles" var sum uint64 - merge_map := map[uint64]uint64{} + merge_map := make(map[uint64]uint64) + slice_map := make(map[uint64]map[util.SUUID]struct{}) + all_slice := make(map[util.SUUID]struct { + process util.GUID + handle index.BitmapHandle + }) + // either create a list of bitmap handles to cat (i.e. union), or sum the integer values + part := make(chan interface{}) + num_parts := len(qs.Inputs) + for _, input := range qs.Inputs { - value, _ := self.Hold.Get(input, 10) + go func(id *util.GUID, part chan interface{}) { + value, _ := self.Hold.Get(id, 10) + part <- value + }(input, part) + } + + //for _, input := range qs.Inputs { + for i := 0; i < num_parts; i++ { + value := <-part + switch val := value.(type) { case index.BitmapHandle: handles = append(handles, val) @@ -188,15 +174,32 @@ func (self *Service) CatQueryStepHandler(msg *db.Message) { //spew.Dump(val) return_type = "sum" sum += val - case []index.Pair: - //spew.Dump(val) + case TopNPackage: return_type = "pair-list" - for _, pair := range val { + var e struct{} + for _, pair := range val.Pairs { + //merge_map[pair.Key] += pair.Count merge_map[pair.Key] += pair.Count + mm, ok := slice_map[pair.Key] + if !ok { + mm = make(map[util.SUUID]struct{}) + slice_map[pair.Key] = mm + } + mm[val.FragmentId] = e } + all_slice[val.FragmentId] = struct { + process util.GUID + handle index.BitmapHandle + }{val.ProcessId, val.HBitmap} } } + tasks := BuildTask(merge_map, slice_map, all_slice) + FetchMissing(tasks, self) + for k, v := range GatherResults(tasks, self) { + merge_map[k] += v + } + // either return the sum, or return the compressed bitmap resulting from the cat (union) var result interface{} if return_type == "sum" { diff --git a/core/topn.go b/core/topn.go new file mode 100644 index 000000000..675733125 --- /dev/null +++ b/core/topn.go @@ -0,0 +1,191 @@ +package core + +import ( + "log" + "pilosa/db" + "pilosa/index" + "pilosa/query" + "pilosa/util" + + "github.com/davecgh/go-spew/spew" +) + +type Task struct { + processid util.GUID + f map[util.SUUID]index.FillArgs + hold_id util.GUID +} + +type TopFill struct { + Args []index.FillArgs + ReturnProcessId util.GUID + QueryId util.GUID + DestProcessId util.GUID +} + +//portable query step +func (self *TopFill) GetId() *util.GUID { + return &self.QueryId +} +func (self *TopFill) GetLocation() *db.Location { + return &db.Location{&self.DestProcessId, 0} //this message is a broadcast to many fragments so i'm choosing fragmentzero +} + +// + +type hole struct { + process util.GUID + handle index.BitmapHandle + fragment util.SUUID +} + +func missing(fids map[util.SUUID]struct{}, all map[util.SUUID]struct { + process util.GUID + handle index.BitmapHandle +}) []hole { + results := make([]hole, 0, 0) + + for k, v := range all { + _, ok := fids[k] + if ok { + results = append(results, hole{v.process, v.handle, k}) + } + } + return results +} + +func newtask(p util.GUID) *Task { + result := new(Task) + result.processid = p + result.f = make(map[util.SUUID]index.FillArgs) + result.hold_id = util.RandomUUID() + return result +} + +func (t *Task) Add(frag util.SUUID, bitmap_id uint64, handle index.BitmapHandle) { + fa, ok := t.f[frag] + if ok { + fa = index.FillArgs{frag, handle, make([]uint64, 0, 0)} + } + fa.Bitmaps = append(fa.Bitmaps, bitmap_id) + t.f[frag] = fa +} + +func BuildTask(merge_map map[uint64]uint64, + slice_map map[uint64]map[util.SUUID]struct{}, + total_fragments map[util.SUUID]struct { + process util.GUID + handle index.BitmapHandle + }) map[util.GUID]*Task { + + tasks := make(map[util.GUID]*Task) + for bitmap_id, _ := range merge_map { //for all brands + //for fragment_id, reported_fragments := range slice_map[bitmap_id] { //find missing fragments + reporting_fragments := slice_map[bitmap_id] + //id slice ==> SUUID,BitmapHandle + for _, p := range missing(reporting_fragments, total_fragments) { + task, ok := tasks[p.process] + if ok { + task = newtask(p.process) + tasks[p.process] = task + } + task.Add(p.fragment, bitmap_id, p.handle) + + } + //} + } + + return tasks +} +func (self *Service) TopFillHandler(msg *db.Message) { //in order for this to get executed it needs to be a portable query step + topfill := msg.Data.(TopFill) + topn, err := self.Index.TopFillBatch(topfill.Args) + if err != nil { + log.Println("TopFileHandler:", err) + } + sendmsg := new(db.Message) + sendmsg.Data = query.BaseQueryResult{Id: &topfill.QueryId, Data: topn} + self.Transport.Send(sendmsg, &topfill.ReturnProcessId) +} + +func SendRequest(process_id util.GUID, t *Task, service *Service) { + args := make([]index.FillArgs, len(t.f), len(t.f)) + for _, v := range t.f { + args = append(args, v) + } + msg := new(db.Message) + p, _ := service.GetProcess() + msg.Data = TopFill{args, p.Id(), t.hold_id, process_id} + service.Transport.Send(msg, &process_id) +} + +func FetchMissing(tasks map[util.GUID]*Task, service *Service) { + for k, v := range tasks { + go SendRequest(k, v, service) + } +} + +func GatherResults(tasks map[util.GUID]*Task, service *Service) map[uint64]uint64 { + results := make(map[uint64]uint64) + answers := make(chan []index.Pair) + for _, task := range tasks { + go func(id util.GUID) { + value, _ := service.Hold.Get(&id, 10) //eiher need to be the frame process or the handler process? + answers <- value.([]index.Pair) + }(task.hold_id) + } + for i := 0; i < len(tasks); i++ { + batch := <-answers + for _, pair := range batch { + results[pair.Key] += pair.Count + } + } + close(answers) + return results +} + +type TopNPackage struct { + ProcessId util.GUID + FragmentId util.SUUID + Pairs []index.Pair + HBitmap index.BitmapHandle +} + +func (self *Service) TopNQueryStepHandler(msg *db.Message) { + //spew.Dump("TOP-N QUERYSTEP") + qs := msg.Data.(query.TopNQueryStep) + //spew.Dump(qs) + //need categories in qs I just added so it would compile + var categoryleaves []uint64 + input := qs.Input + value, _ := self.Hold.Get(input, 10) + var bh index.BitmapHandle + switch val := value.(type) { + case index.BitmapHandle: + bh = val + case []byte: + bh, _ = self.Index.FromBytes(qs.Location.FragmentId, val) + } + + categoryleaves = qs.Filters + + /* + spew.Dump(bh, qs.N*2, categoryleaves) + result_message := db.Message{Data: "foobar"} + self.Transport.Send(&result_message, qs.Destination.ProcessId) + */ + + topn, err := self.Index.TopN(qs.Location.FragmentId, bh, qs.N*2, categoryleaves) + topnPackage := TopNPackage{*qs.Location.ProcessId, qs.Location.FragmentId, topn, bh} + + if err != nil { + spew.Dump(err) + } + result_message := db.Message{Data: query.TopNQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: topnPackage}}} + self.Transport.Send(&result_message, qs.Destination.ProcessId) + +} + +func canCompile() bool { + return true +} diff --git a/db/topology_test.go b/db/topology_test.go index 637b7a2f4..2367d1ea1 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -23,8 +23,8 @@ func TestTopology(t *testing.T) { spew.Dump(fragment_id) database.GetOrCreateFragment(frame, slice, fragment_id) - spew.Dump(database) - spew.Dump("DONE") + // spew.Dump(database) + // spew.Dump("DONE") }) } diff --git a/executor/executor.go b/executor/executor.go index a24923778..db8f5275f 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -44,6 +44,10 @@ func (self *Executor) NewJob(job *db.Message) { self.service.GetQueryStepHandler(job) case query.SetQueryStep: self.service.SetQueryStepHandler(job) + // case query.MaskQueryStep: + // self.service.MaskQueryStepHandler(job) + // case query.RangeQueryStep: + // self.service.RangeQueryStepHandler(job) default: fmt.Println("unknown") } @@ -99,7 +103,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"} + reserved_functions := stringSlice{"get", "set", "union", "intersect", "difference", "count", "top-n", "mask", "range"} tokens, err := query.Lex(pql) if err != nil { return nil, err diff --git a/index/bitmap.go b/index/bitmap.go index b4210f353..f68fd162e 100644 --- a/index/bitmap.go +++ b/index/bitmap.go @@ -67,7 +67,6 @@ func popcount(i uint64)uint64{ //x:= uint64(val) return uint64(val) } -*/ func popcount(x uint64) (n uint64) { // bit population count, see // http://graphics.stanford.edu/~seander/bithacks.html#CountBitsSetParallel @@ -78,6 +77,7 @@ func popcount(x uint64) (n uint64) { x *= 0x0101010101010101 return uint64(x >> 56) } +*/ type BlockArray struct { Block [32]uint64 @@ -89,6 +89,7 @@ func (s *BlockArray) bitcount() uint64 { sum += popcount(b) } return sum + // return popcntSlice(s.Block) } func BlockArray_union(a *BlockArray, b *BlockArray) BlockArray { var o = BlockArray{} diff --git a/index/commands.go b/index/commands.go index 71f6a1ede..7169c8d7c 100644 --- a/index/commands.go +++ b/index/commands.go @@ -291,3 +291,23 @@ func (self *CmdMask) Execute(f *Fragment) Calculation { } return f.AllocHandle(result) } + +type CmdTopFill struct { + *Responder + args FillArgs +} + +func NewTopFill(a FillArgs) *CmdTopFill { + return &CmdTopFill{NewResponder("TopFill"), a} +} + +func (self *CmdTopFill) Execute(f *Fragment) Calculation { + result := make([]Pair, len(self.args.Bitmaps)) + for _, v := range self.args.Bitmaps { + a := f.NewHandle(v) + res := f.intersect([]BitmapHandle{self.args.Handle, a}) + bm, _ := f.getBitmap(res) + result = append(result, Pair{v, BitCount(bm)}) + } + return result +} diff --git a/index/fragment_container.go b/index/fragment_container.go index a61f1a549..27dc3b040 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -39,6 +39,11 @@ func NewFragmentContainer() *FragmentContainer { } type BitmapHandle uint64 +type FillArgs struct { + Frag_id util.SUUID + Handle BitmapHandle + Bitmaps []uint64 +} func init() { var vh BitmapHandle @@ -103,7 +108,6 @@ func (self *FragmentContainer) Intersect(frag_id util.SUUID, bh []BitmapHandle) } return 0, errors.New("Invalid Bitmap Handle") } - func (self *FragmentContainer) Union(frag_id util.SUUID, bh []BitmapHandle) (BitmapHandle, error) { if fragment, found := self.GetFragment(frag_id); found { request := NewUnion(bh) @@ -169,6 +173,35 @@ func (self *FragmentContainer) TopN(frag_id util.SUUID, bh BitmapHandle, n int, } return nil, errors.New(fmt.Sprintf("Fragment not found:%s", util.SUUID_to_Hex(frag_id))) } +func (self *FragmentContainer) TopFillBatch(args []FillArgs) ([]Pair, error) { + //should probaly make this concurrent but then all hell breaks lose + results := make(map[uint64]uint64) + for _, v := range args { + items, _ := self.TopFillFragment(v) + if len(args) == 1 { + return items, nil + } + for _, pair := range items { + results[pair.Key] += pair.Count + } + } + ret_val := make([]Pair, len(results)) + for k, v := range results { + ret_val = append(ret_val, Pair{k, v}) + } + return ret_val, nil + +} +func (self *FragmentContainer) TopFillFragment(arg FillArgs) ([]Pair, error) { + if fragment, found := self.GetFragment(arg.Frag_id); found { + request := NewTopFill(arg) + fragment.requestChan <- request + result := request.Response() + util.SendTimer("fragmant_container_TopFillFragment", result.exec_time.Nanoseconds()) + return result.answer.([]Pair), nil + } + return nil, errors.New("Invalid Bitmap Handle") +} func (self *FragmentContainer) GetList(frag_id util.SUUID, bitmap_id []uint64) ([]BitmapHandle, error) { if fragment, found := self.GetFragment(frag_id); found { diff --git a/index/popcnt.go b/index/popcnt.go new file mode 100644 index 000000000..a1f25c8a8 --- /dev/null +++ b/index/popcnt.go @@ -0,0 +1,53 @@ +package index + +// bit population count, take from +// https://code.google.com/p/go/issues/detail?id=4988#c11 +// credit: https://code.google.com/u/arnehormann/ +func popcount(x uint64) (n uint64) { + x -= (x >> 1) & 0x5555555555555555 + x = (x>>2)&0x3333333333333333 + x&0x3333333333333333 + x += x >> 4 + x &= 0x0f0f0f0f0f0f0f0f + x *= 0x0101010101010101 + return x >> 56 +} + +func popcntSliceGo(s []uint64) uint64 { + cnt := uint64(0) + for _, x := range s { + cnt += popcount(x) + } + return cnt +} + +func popcntMaskSliceGo(s, m []uint64) uint64 { + cnt := uint64(0) + for i := range s { + cnt += popcount(s[i] &^ m[i]) + } + return cnt +} + +func popcntAndSliceGo(s, m []uint64) uint64 { + cnt := uint64(0) + for i := range s { + cnt += popcount(s[i] & m[i]) + } + return cnt +} + +func popcntOrSliceGo(s, m []uint64) uint64 { + cnt := uint64(0) + for i := range s { + cnt += popcount(s[i] | m[i]) + } + return cnt +} + +func popcntXorSliceGo(s, m []uint64) uint64 { + cnt := uint64(0) + for i := range s { + cnt += popcount(s[i] ^ m[i]) + } + return cnt +} diff --git a/index/popcnt_amd64.s b/index/popcnt_amd64.s new file mode 100644 index 000000000..121ce7bfa --- /dev/null +++ b/index/popcnt_amd64.s @@ -0,0 +1,102 @@ +TEXT ·hasAsm(SB),4,$0 +MOVQ $1, AX +CPUID +SHRQ $23, CX +ANDQ $1, CX +MOVB CX, ret+0(FP) +RET + + +#define POPCNTQ_DX_DX BYTE $0xf3; BYTE $0x48; BYTE $0x0f; BYTE $0xb8; BYTE $0xd2 + +TEXT ·popcntSliceAsm(SB),4,$0-32 +XORQ AX, AX +MOVQ s+0(FP), SI +MOVQ s+8(FP), CX +TESTQ CX, CX +JZ popcntSliceEnd +popcntSliceLoop: +BYTE $0xf3; BYTE $0x48; BYTE $0x0f; BYTE $0xb8; BYTE $0x16 // POPCNTQ (SI), DX +ADDQ DX, AX +ADDQ $8, SI +LOOP popcntSliceLoop +popcntSliceEnd: +MOVQ AX, ret+24(FP) +RET + +TEXT ·popcntMaskSliceAsm(SB),4,$0-56 +XORQ AX, AX +MOVQ s+0(FP), SI +MOVQ s+8(FP), CX +TESTQ CX, CX +JZ popcntMaskSliceEnd +MOVQ m+24(FP), DI +popcntMaskSliceLoop: +MOVQ (DI), DX +NOTQ DX +ANDQ (SI), DX +POPCNTQ_DX_DX +ADDQ DX, AX +ADDQ $8, SI +ADDQ $8, DI +LOOP popcntMaskSliceLoop +popcntMaskSliceEnd: +MOVQ AX, ret+48(FP) +RET + +TEXT ·popcntAndSliceAsm(SB),4,$0-56 +XORQ AX, AX +MOVQ s+0(FP), SI +MOVQ s+8(FP), CX +TESTQ CX, CX +JZ popcntAndSliceEnd +MOVQ m+24(FP), DI +popcntAndSliceLoop: +MOVQ (DI), DX +ANDQ (SI), DX +POPCNTQ_DX_DX +ADDQ DX, AX +ADDQ $8, SI +ADDQ $8, DI +LOOP popcntAndSliceLoop +popcntAndSliceEnd: +MOVQ AX, ret+48(FP) +RET + +TEXT ·popcntOrSliceAsm(SB),4,$0-56 +XORQ AX, AX +MOVQ s+0(FP), SI +MOVQ s+8(FP), CX +TESTQ CX, CX +JZ popcntOrSliceEnd +MOVQ m+24(FP), DI +popcntOrSliceLoop: +MOVQ (DI), DX +ORQ (SI), DX +POPCNTQ_DX_DX +ADDQ DX, AX +ADDQ $8, SI +ADDQ $8, DI +LOOP popcntOrSliceLoop +popcntOrSliceEnd: +MOVQ AX, ret+48(FP) +RET + +TEXT ·popcntXorSliceAsm(SB),4,$0-56 +XORQ AX, AX +MOVQ s+0(FP), SI +MOVQ s+8(FP), CX +TESTQ CX, CX +JZ popcntXorSliceEnd +MOVQ m+24(FP), DI +popcntXorSliceLoop: +MOVQ (DI), DX +XORQ (SI), DX +POPCNTQ_DX_DX +ADDQ DX, AX +ADDQ $8, SI +ADDQ $8, DI +LOOP popcntXorSliceLoop +popcntXorSliceEnd: +MOVQ AX, ret+48(FP) +RET diff --git a/index/popcnt_asm.go b/index/popcnt_asm.go new file mode 100644 index 000000000..6af8ad422 --- /dev/null +++ b/index/popcnt_asm.go @@ -0,0 +1,64 @@ +// +build amd64 + +package index + +//go:noescape + +func hasAsm() bool + +var useAsm = hasAsm() + +//go:noescape + +func popcntSliceAsm(s []uint64) uint64 + +//go:noescape + +func popcntMaskSliceAsm(s, m []uint64) uint64 + +//go:noescape + +func popcntAndSliceAsm(s, m []uint64) uint64 + +//go:noescape + +func popcntOrSliceAsm(s, m []uint64) uint64 + +//go:noescape + +func popcntXorSliceAsm(s, m []uint64) uint64 + +func popcntSlice(s []uint64) uint64 { + if useAsm { + return popcntSliceAsm(s) + } + return popcntSliceGo(s) +} + +func popcntMaskSlice(s, m []uint64) uint64 { + if useAsm { + return popcntMaskSliceAsm(s, m) + } + return popcntMaskSliceGo(s, m) +} + +func popcntAndSlice(s, m []uint64) uint64 { + if useAsm { + return popcntAndSliceAsm(s, m) + } + return popcntAndSliceGo(s, m) +} + +func popcntOrSlice(s, m []uint64) uint64 { + if useAsm { + return popcntOrSliceAsm(s, m) + } + return popcntOrSliceGo(s, m) +} + +func popcntXorSlice(s, m []uint64) uint64 { + if useAsm { + return popcntXorSliceAsm(s, m) + } + return popcntXorSliceGo(s, m) +} diff --git a/index/popcnt_generic.go b/index/popcnt_generic.go new file mode 100644 index 000000000..d79741573 --- /dev/null +++ b/index/popcnt_generic.go @@ -0,0 +1,23 @@ +// +build !amd64 + +package index + +func popcntSlice(s []uint64) uint64 { + return popcntSliceGo(s) +} + +func popcntMaskSlice(s, m []uint64) uint64 { + return popcntMaskSliceGo(s, m) +} + +func popcntAndSlice(s, m []uint64) uint64 { + return popcntAndSliceGo(s, m) +} + +func popcntOrSlice(s, m []uint64) uint64 { + return popcntOrSliceGo(s, m) +} + +func popcntXorSlice(s, m []uint64) uint64 { + return popcntXorSliceGo(s, m) +} diff --git a/index/timeframe_test.go b/index/timeframe_test.go index cbbfe5cd2..fb8814967 100644 --- a/index/timeframe_test.go +++ b/index/timeframe_test.go @@ -2,6 +2,7 @@ package index import ( "fmt" + "log" "testing" "time" @@ -9,6 +10,30 @@ import ( . "github.com/smartystreets/goconvey/convey" ) +func getTime(id uint64, s string) { + const shortForm = "2006-01-02 15:04" + t1, _ := time.Parse(shortForm, s) + + for i, v := range GetTimeIds(uint64(id), t1, YMDH) { + log.Println(i, v, s) + } + log.Println() +} +func TestDemo(t *testing.T) { + Convey("Test ID", t, func() { + getTime(uint64(666), "2014-01-01 00:00") + getTime(uint64(666), "2014-02-01 00:00") + getTime(uint64(666), "2014-03-01 00:00") + getTime(uint64(666), "2014-04-01 00:00") + getTime(uint64(666), "2014-04-02 00:00") + getTime(uint64(666), "2014-04-03 00:00") + getTime(uint64(666), "2014-04-04 00:00") + getTime(uint64(666), "2014-04-05 00:00") + getTime(uint64(666), "2014-05-01 00:00") + getTime(uint64(666), "2014-06-01 00:00") + So(1, ShouldEqual, 1) + }) +} func TestTimeFrame(t *testing.T) { //print get_YMD_id(2014,3,28,1234) diff --git a/query/parser.go b/query/parser.go index 907927e12..5abd6be66 100644 --- a/query/parser.go +++ b/query/parser.go @@ -6,6 +6,7 @@ import ( "log" "pilosa/util" "strconv" + "time" "github.com/davecgh/go-spew/spew" ) @@ -81,6 +82,7 @@ func (self *QueryParser) Parse() (query *Query, err error) { if token.Type != TYPE_LP { return nil, fmt.Errorf("Expected '(', found token %v.", token) } + const shortForm = "2006-01-02 15:04" ArgLoop: for { @@ -99,6 +101,46 @@ ArgLoop: query.Subqueries = append(query.Subqueries, *subquery) case TYPE_VALUE: switch query.Operation { + case "range": + 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: + t, err := time.Parse(shortForm, token.Text) + if err != nil { + return nil, fmt.Errorf("Expecting integer DateTime (%v)", err) + } + query.Args["start"] = t + case 3: + t, err := time.Parse(shortForm, token.Text) + if err != nil { + return nil, fmt.Errorf("Expecting integer DateTime (%v)", err) + } + query.Args["end"] = t + } + case "mask": + switch len(query.Args) { + case 0: + query.Args["frame"] = token.Text + case 1: + i, err := strconv.ParseUint(token.Text, 10, 64) + if err != nil { + return nil, fmt.Errorf("Expecting integer id! (%v)", err) + } + query.Args["start"] = i + case 2: + i, err := strconv.ParseUint(token.Text, 10, 64) + if err != nil { + return nil, fmt.Errorf("Expecting integer id! (%v)", err) + } + query.Args["end"] = i + } case "get": switch len(query.Args) { case 0: @@ -227,6 +269,12 @@ ArgLoop: if query.Operation == "difference" { return nil, fmt.Errorf("No Args Given") } + if query.Operation == "range" { + return nil, fmt.Errorf("No Args Given") + } + if query.Operation == "mask" { + return nil, fmt.Errorf("No Args Given") + } } return query, nil } diff --git a/query/planner.go b/query/planner.go index 4ef76306a..5d2584f4d 100644 --- a/query/planner.go +++ b/query/planner.go @@ -8,6 +8,7 @@ import ( "math/rand" "pilosa/db" "pilosa/util" + "time" "github.com/davecgh/go-spew/spew" ) @@ -146,7 +147,7 @@ func (qt *TopNQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice := d.GetOrCreateSlice(qt.Slice) fragment, err := d.GetFragmentForFrameSlice(frame, slice) if err != nil { - log.Println("GetFragmentForFrameSliceFailed GetQueryTree", frame, slice) + log.Println("GetFragmentForFrameSliceFailed TopNQueryTree", frame, slice) return nil, err } qt.location = fragment.GetLocation() @@ -389,7 +390,7 @@ type QueryTree interface { } // Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices -func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { +func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { var tree QueryTree // handle SET operation regardless of the slice @@ -407,13 +408,13 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { n = n_.(int) } tree = &CatQueryTree{N: n} - numSlices, err := qp.Database.NumSlices() + numSlices, err := self.Database.NumSlices() if err != nil { return nil, err } for slice := 0; slice < numSlices; slice++ { //for slice := 0; slice < 3; slice++ { - subtree, err := qp.buildTree(query, slice) + subtree, err := self.buildTree(query, slice) if err != nil { return nil, err } @@ -425,7 +426,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { tree = &GetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), 0}, slice} return tree, nil } else if query.Operation == "count" { - subquery, err := qp.buildTree(&query.Subqueries[0], slice) + subquery, err := self.buildTree(&query.Subqueries[0], slice) if err != nil { return nil, err } @@ -447,7 +448,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { filters = filters_.([]uint64) } - subquery, err := qp.buildTree(&query.Subqueries[0], slice) + subquery, err := self.buildTree(&query.Subqueries[0], slice) if err != nil { return nil, err } @@ -456,7 +457,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { subqueries := make([]QueryTree, len(query.Subqueries)) var err error for i, query := range query.Subqueries { - subqueries[i], err = qp.buildTree(&query, slice) + subqueries[i], err = self.buildTree(&query, slice) if err != nil { return nil, err } @@ -466,7 +467,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { subqueries := make([]QueryTree, len(query.Subqueries)) var err error for i, query := range query.Subqueries { - subqueries[i], err = qp.buildTree(&query, slice) + subqueries[i], err = self.buildTree(&query, slice) if err != nil { return nil, err } @@ -476,7 +477,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { subqueries := make([]QueryTree, len(query.Subqueries)) var err error for i, query := range query.Subqueries { - subqueries[i], err = qp.buildTree(&query, slice) + subqueries[i], err = self.buildTree(&query, slice) if err != nil { return nil, err } @@ -492,11 +493,11 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { } // Produces flattened QueryPlan from QueryTree input -func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Location) (*QueryPlan, error) { +func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Location) (*QueryPlan, error) { plan := QueryPlan{} if cat, ok := qt.(*CatQueryTree); ok { inputs := make([]*util.GUID, len(cat.subqueries)) - loc, err := cat.getLocation(qp.Database) + loc, err := cat.getLocation(self.Database) if err != nil { return nil, err } @@ -504,7 +505,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati for index, subq := range cat.subqueries { sub_id := util.RandomUUID() step.Inputs[index] = &sub_id - subq_steps, err := qp.flatten(subq, &sub_id, loc) + subq_steps, err := self.flatten(subq, &sub_id, loc) if err != nil { return nil, err } @@ -513,7 +514,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati plan = append(plan, step) } else if union, ok := qt.(*UnionQueryTree); ok { inputs := make([]*util.GUID, len(union.subqueries)) - loc, err := union.getLocation(qp.Database) + loc, err := union.getLocation(self.Database) if err != nil { return nil, err } @@ -521,7 +522,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati for index, subq := range union.subqueries { sub_id := util.RandomUUID() step.Inputs[index] = &sub_id - subq_steps, err := qp.flatten(subq, &sub_id, loc) + subq_steps, err := self.flatten(subq, &sub_id, loc) if err != nil { return nil, err } @@ -530,7 +531,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati plan = append(plan, step) } else if intersect, ok := qt.(*IntersectQueryTree); ok { inputs := make([]*util.GUID, len(intersect.subqueries)) - loc, err := intersect.getLocation(qp.Database) + loc, err := intersect.getLocation(self.Database) if err != nil { return nil, err } @@ -538,7 +539,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati for index, subq := range intersect.subqueries { sub_id := util.RandomUUID() step.Inputs[index] = &sub_id - subq_steps, err := qp.flatten(subq, &sub_id, loc) + subq_steps, err := self.flatten(subq, &sub_id, loc) if err != nil { return nil, err } @@ -547,7 +548,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati plan = append(plan, step) } else if difference, ok := qt.(*DifferenceQueryTree); ok { inputs := make([]*util.GUID, len(difference.subqueries)) - loc, err := difference.getLocation(qp.Database) + loc, err := difference.getLocation(self.Database) if err != nil { return nil, err } @@ -555,7 +556,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati for index, subq := range difference.subqueries { sub_id := util.RandomUUID() step.Inputs[index] = &sub_id - subq_steps, err := qp.flatten(subq, &sub_id, loc) + subq_steps, err := self.flatten(subq, &sub_id, loc) if err != nil { return nil, err } @@ -563,15 +564,33 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati } plan = append(plan, step) } else if get, ok := qt.(*GetQueryTree); ok { - loc, err := get.getLocation(qp.Database) + loc, err := get.getLocation(self.Database) if err != nil { return nil, err } step := GetQueryStep{&BaseQueryStep{id, "get", loc, location}, get.bitmap, get.slice} plan := QueryPlan{step} return &plan, nil + /* + } else if mask, ok := qt.(*MaskQueryTree); ok { + loc, err := mask.getLocation(self.Database) + if err != nil { + return nil, err + } + 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 set, ok := qt.(*SetQueryTree); ok { - loc, err := set.getLocation(qp.Database) + loc, err := set.getLocation(self.Database) if err != nil { return nil, err } @@ -580,12 +599,12 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati return &plan, nil } else if cnt, ok := qt.(*CountQueryTree); ok { sub_id := util.RandomUUID() - loc, err := cnt.getLocation(qp.Database) + loc, err := cnt.getLocation(self.Database) if err != nil { return nil, err } step := &CountQueryStep{&BaseQueryStep{id, "count", loc, location}, &sub_id} - subq_steps, err := qp.flatten(cnt.subquery, &sub_id, loc) + subq_steps, err := self.flatten(cnt.subquery, &sub_id, loc) if err != nil { return nil, err } @@ -593,12 +612,12 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati plan = append(plan, step) } else if topn, ok := qt.(*TopNQueryTree); ok { sub_id := util.RandomUUID() - loc, err := topn.getLocation(qp.Database) + loc, err := topn.getLocation(self.Database) if err != nil { return nil, err } step := &TopNQueryStep{&BaseQueryStep{id, "top-n", loc, location}, &sub_id, topn.Filters, topn.N, topn.Frame} - subq_steps, err := qp.flatten(topn.subquery, &sub_id, loc) + subq_steps, err := self.flatten(topn.subquery, &sub_id, loc) if err != nil { return nil, err } @@ -609,10 +628,71 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati } // Transforms Query into QueryTree and flattens to QueryPlan object -func (qp *QueryPlanner) Plan(query *Query, id *util.GUID, destination *db.Location) (*QueryPlan, error) { - queryTree, err := qp.buildTree(query, -1) +func (self *QueryPlanner) Plan(query *Query, id *util.GUID, destination *db.Location) (*QueryPlan, error) { + queryTree, err := self.buildTree(query, -1) if err != nil { return nil, err } - return qp.flatten(queryTree, query.Id, destination) + return self.flatten(queryTree, query.Id, destination) } + +/////////////////////////////////////////////////////////////////////////////////////////////////// +//MASK +/////////////////////////////////////////////////////////////////////////////////////////////////// +type MaskQueryStep struct { + *BaseQueryStep + start, end uint64 +} + +type MaskQueryResult struct { + *BaseQueryResult +} + +// QueryTree for Mask queries +type MaskQueryTree struct { + start, end uint64 + bitmap *db.Bitmap +} + +// Uses consistent hashing function to select node containing data for GET operation +/* +func (qt *MaskQueryTree) 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, qt.bitmap) + if err != nil { + log.Println("GetFragmenForBitmapFailed GetQueryTree", slice) + return nil, err + } + return fragment.GetLocation(), nil +} +*/ + +/////////////////////////////////////////////////////////////////////////////////////////////////// +//Range +/////////////////////////////////////////////////////////////////////////////////////////////////// +type RangeQueryStep struct { + *BaseQueryStep + bitmap *db.Bitmap + start, end time.Time +} + +type RangeQueryResult struct { + *BaseQueryResult +} + +// QueryTree for Mask queries +type RangeQueryTree struct { +} + +// Uses consistent hashing function to select node containing data for GET operation +/* +func (self *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) + if err != nil { + log.Println("GetFragmenForBitmapFailed GetQueryTree", slice) + return nil, err + } + return fragment.GetLocation(), nil +} +*/ diff --git a/util/id.go b/util/id.go index 88b69fa0c..a76908198 100644 --- a/util/id.go +++ b/util/id.go @@ -108,3 +108,23 @@ func ParseGUID(input string) (GUID, error) { } return u, nil } + +func In(val int, list []int) bool { + for _, v := range list { + if val == v { + return true + } + } + return false +} + +func Difference(a, b []int) []int { + results := make([]int, 0, len(a)) + for _, v := range a { + if !In(v, b) { + results = append(results, v) + } + + } + return results +}