From 73631de9e8ccc3cfc88adffdfbfeeeb73a92024d Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 20 Aug 2015 13:46:39 -0600 Subject: [PATCH] vet --- core/batch.go | 2 +- core/load.go | 4 +- core/stopper.go | 4 +- db/db.go | 2 +- dispatch/dispatch.go | 4 +- executor/executor.go | 76 +++++++++++++++++++++++++------- index/benchmark_autogen_test.go | 2 +- index/fragment_container_test.go | 6 +-- query/lexer.go | 2 - query/planner.go | 11 ++--- query/planner_test.go | 18 ++++---- transport/tcp.go | 4 +- 12 files changed, 89 insertions(+), 46 deletions(-) diff --git a/core/batch.go b/core/batch.go index 618895150..753891d6b 100644 --- a/core/batch.go +++ b/core/batch.go @@ -56,7 +56,7 @@ func (b *Batcher) Batch(database_name, frame, compressed_bitmap string, bitmap_i oslice := database.GetOrCreateSlice(slice) //need to find processid and fragment id for that slice - fragment, err := database.GetFragmentForBitmap(oslice, &db.Bitmap{bitmap_id, frame, filter}) + fragment, err := database.GetFragmentForBitmap(oslice, &db.Bitmap{Id: bitmap_id, FrameType: frame, Filter: filter}) if err == nil { id := util.RandomUUID() batch := db.Message{Data: BatchRequest{Id: &id, Source: &b.ID, Fragment_id: fragment.GetId(), Bitmap_id: bitmap_id, Compressed_bitmap: compressed_bitmap}} diff --git a/core/load.go b/core/load.go index de07dcdf2..234e04ae4 100644 --- a/core/load.go +++ b/core/load.go @@ -18,7 +18,7 @@ func copy_raw(src [32]uint64) index.BlockArray { for k, v := range src { o[k] = v } - return index.BlockArray{o} + return index.BlockArray{Block: o} } func sendBitmap(batcher *Batcher, bitmap index.IBitmap, db string, frame string, bitmap_id, filter uint64, slice int, finish chan error) { if slice < 0 { @@ -73,7 +73,7 @@ func FromApiString(batcher *Batcher, db string, frame string, api_string string, last_slice = slice } o := copy_raw(raw.Block) - chunk := &index.Chunk{raw.Key, o} + chunk := &index.Chunk{Key: raw.Key, Value: o} bitmap.AddChunk(chunk) } diff --git a/core/stopper.go b/core/stopper.go index 3e3c9fda8..5569d2e21 100644 --- a/core/stopper.go +++ b/core/stopper.go @@ -15,9 +15,9 @@ func (stopper *Stopper) Stop() { var o chan int stopper.Mutex.RLock() for _, i = range stopper.TermChans { - go func() { + go func(i chan int) { i <- 1 - }() + }(i) } for _, o = range stopper.DoneChans { <-o diff --git a/db/db.go b/db/db.go index cde9fa215..158d8cbfa 100644 --- a/db/db.go +++ b/db/db.go @@ -7,7 +7,7 @@ import ( ) type Message struct { - Data interface{} `json:data` + Data interface{} `json:"data"` } type Envelope struct { diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 7ea06632a..af2bff717 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -71,7 +71,7 @@ func (d *Dispatch) Run() { } else { result, _ = d.Index.ClearBit(v.Fragment_id, v.Bitmap_id, v.Profile_id) } - bundle := core.SBResult{v.Bitmap_id, v.Frame, v.Filter, v.Profile_id, result} + bundle := core.SBResult{Bitmap_id: v.Bitmap_id, Frame: v.Frame, Filter: v.Filter, Profile_id: v.Profile_id, Result: result} results = append(results, bundle) } response := db.Message{Data: core.BitsResponse{Id: &data.QueryId, Items: results}} @@ -113,7 +113,7 @@ func (d *Dispatch) topFillHandler(m *db.Message) { d.Transport.Send(&db.Message{ Data: query.FillResult{ - &query.BaseQueryResult{ + BaseQueryResult: &query.BaseQueryResult{ Id: &topfill.QueryId, Data: topn, }, diff --git a/executor/executor.go b/executor/executor.go index 0fbbb4603..a10470ea7 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -114,7 +114,11 @@ func (self *Executor) CountQueryStepHandler(msg *db.Message) { spew.Dump(err) } //spew.Dump("SLICE COUNT", count) - result_message := db.Message{Data: query.CountQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: count}}} + result_message := db.Message{ + Data: query.CountQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: count}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -148,7 +152,11 @@ func (self *Executor) TopNQueryStepHandler(msg *db.Message) { topnPackage = TopNPackage{*qs.Location.ProcessId, qs.Location.FragmentId, topn, bh} } - result_message := db.Message{Data: query.TopNQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: topnPackage}}} + result_message := db.Message{ + Data: query.TopNQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: topnPackage}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -185,7 +193,11 @@ func (self *Executor) UnionQueryStepHandler(msg *db.Message) { } result = bm } - result_message := db.Message{Data: query.UnionQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.UnionQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -221,7 +233,11 @@ func (self *Executor) IntersectQueryStepHandler(msg *db.Message) { } result = bm } - result_message := db.Message{Data: query.IntersectQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.IntersectQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -257,7 +273,11 @@ func (self *Executor) DifferenceQueryStepHandler(msg *db.Message) { } result = bm } - result_message := db.Message{Data: query.DifferenceQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.DifferenceQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -348,7 +368,7 @@ func (self *Executor) CatQueryStepHandler(msg *db.Message) { continue //shouldn't be getting 0 keys or values anyway } rank := new(index.Rank) - rank.Pair = &index.Pair{k, v} + rank.Pair = &index.Pair{Key: k, Count: v} rank_list = append(rank_list, rank) } sort.Sort(rank_list) // kinda seems like this copy is wasteful..i'll ponder @@ -365,7 +385,11 @@ func (self *Executor) CatQueryStepHandler(msg *db.Message) { } else { result = "NONE" } - result_message := db.Message{Data: query.CatQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.CatQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -435,7 +459,11 @@ func (self *Executor) GetQueryStepHandler(msg *db.Message) { } result = bm } - result_message := db.Message{Data: query.GetQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.GetQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -444,7 +472,11 @@ func (self *Executor) SetQueryStepHandler(msg *db.Message) { qs := msg.Data.(query.SetQueryStep) result, _ := self.Index.SetBit(qs.Location.FragmentId, qs.Bitmap.Id, qs.ProfileId, qs.Bitmap.Filter) - result_message := db.Message{Data: query.SetQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.SetQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -453,7 +485,11 @@ func (self *Executor) ClearQueryStepHandler(msg *db.Message) { 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}}} + result_message := db.Message{ + Data: query.ClearQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -476,7 +512,11 @@ func (self *Executor) RangeQueryStepHandler(msg *db.Message) { } result = bm } - result_message := db.Message{Data: query.RangeQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.RangeQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -505,7 +545,7 @@ func (self *Executor) StashQueryStepHandler(msg *db.Message) { //result.Handles = append(result.Handles, val) case []byte: bh, _ := self.Index.FromBytes(qs.Location.FragmentId, val) - item := query.CacheItem{qs.Location.FragmentId, bh} + item := query.CacheItem{FragmentId: qs.Location.FragmentId, Handle: bh} result.Stash = append(result.Stash, item) case query.Stash: result.Stash = append(result.Stash, val.Stash...) @@ -513,7 +553,11 @@ func (self *Executor) StashQueryStepHandler(msg *db.Message) { log.Warn("UNEXCPECTED MESSAGE", value) } } - result_message := db.Message{Data: query.StashQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + result_message := db.Message{ + Data: query.StashQueryResult{ + BaseQueryResult: &query.BaseQueryResult{Id: qs.Id, Data: result}, + }, + } self.Transport.Send(&result_message, qs.Destination.ProcessId) } @@ -530,7 +574,7 @@ func (self *Executor) runQuery(database *db.Database, qry *query.Query) error { } process_id := process.Id() fragment_id := util.SUUID(0) - destination := db.Location{&process_id, fragment_id} + destination := db.Location{ProcessId: &process_id, FragmentId: fragment_id} query_plan, err := query.QueryPlanForQuery(database, qry, &destination) if err != nil { @@ -671,7 +715,7 @@ func newtask(p util.GUID) *Task { 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 = index.FillArgs{Frag_id: frag, Handle: handle, Bitmaps: make([]uint64, 0, 0)} } fa.Bitmaps = append(fa.Bitmaps, bitmap_id) t.f[frag] = fa @@ -729,7 +773,7 @@ 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 + return &db.Location{ProcessId: &self.DestProcessId, FragmentId: 0} //this message is a broadcast to many fragments so i'm choosing fragmentzero } func (self *Executor) Run() { diff --git a/index/benchmark_autogen_test.go b/index/benchmark_autogen_test.go index b3b96d6be..5e86e2110 100644 --- a/index/benchmark_autogen_test.go +++ b/index/benchmark_autogen_test.go @@ -43,7 +43,7 @@ func benchmarkDifferentCombinations(b *testing.B, op string, b1, b2 int, s1, s2 b.ResetTimer() for i := 0; i < b.N; i++ { if f(m1, m2) == nil { - b.Fatal("Problem with %s benchmark at i =", op, i) + b.Fatalf("Problem with %s benchmark at i = %d", op, i) } } } diff --git a/index/fragment_container_test.go b/index/fragment_container_test.go index 577ebe173..bd068e3eb 100644 --- a/index/fragment_container_test.go +++ b/index/fragment_container_test.go @@ -64,7 +64,7 @@ func TestFragmentContainer_Count(t *testing.T) { } else if n, err := fc.Count(util.SUUID(1), bh); err != nil { t.Fatal(err) } else if n != 1 { - t.Fatal("unexpected count: %d", n) + t.Fatalf("unexpected count: %d", n) } } @@ -108,7 +108,7 @@ func TestFragmentContainer_Intersect(t *testing.T) { if result, err := fc.Intersect(1, []index.BitmapHandle{fc.MustGet(1, 1234), fc.MustGet(1, 4321)}); err != nil { t.Fatal(err) } else if n := fc.MustCount(1, result); n != 0 { - t.Fatal("unexpected intersect bit count: %d", n) + t.Fatalf("unexpected intersect bit count: %d", n) } } @@ -123,7 +123,7 @@ func TestFragmentContainer_Difference(t *testing.T) { if result, err := fc.Difference(1, []index.BitmapHandle{fc.MustGet(1, 1234), fc.MustGet(1, 4321)}); err != nil { t.Fatal(err) } else if n := fc.MustCount(1, result); n != 1 { - t.Fatalf("unexpected difference bit count: %d", err) + t.Fatalf("unexpected difference bit count: %s", err) } } diff --git a/query/lexer.go b/query/lexer.go index 7ac91396d..1c7ce0ea4 100644 --- a/query/lexer.go +++ b/query/lexer.go @@ -226,7 +226,6 @@ func stateArgs(lexer *Lexer) statefn { default: return stateError(errors.New("Expecting arguments!")) } - return nil } func stateKeyword(lexer *Lexer) statefn { @@ -273,7 +272,6 @@ func stateValue(lexer *Lexer) statefn { default: return stateError(errors.New("Unexpected character!")) } - return nil } func stateRP(lexer *Lexer) statefn { diff --git a/query/planner.go b/query/planner.go index fefbe967c..4a0326616 100644 --- a/query/planner.go +++ b/query/planner.go @@ -466,12 +466,12 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) var tree QueryTree // handle SET operation regardless of the slice if query.Operation == "set" { - tree = &SetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), query.Args["filter"].(uint64)}, query.Args["profile_id"].(uint64)} + tree = &SetQueryTree{&db.Bitmap{Id: query.Args["id"].(uint64), FrameType: query.Args["frame"].(string), Filter: query.Args["filter"].(uint64)}, query.Args["profile_id"].(uint64)} return tree, nil } if query.Operation == "clear" { - tree = &ClearQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), query.Args["filter"].(uint64)}, query.Args["profile_id"].(uint64)} + tree = &ClearQueryTree{&db.Bitmap{Id: query.Args["id"].(uint64), FrameType: query.Args["frame"].(string), Filter: query.Args["filter"].(uint64)}, query.Args["profile_id"].(uint64)} return tree, nil } @@ -509,7 +509,7 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) } } else { if query.Operation == "get" { - tree = &GetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), 0}, slice} + tree = &GetQueryTree{&db.Bitmap{Id: query.Args["id"].(uint64), FrameType: query.Args["frame"].(string), Filter: 0}, slice} return tree, nil } else if query.Operation == "range" { err := validateRange(query.Args) @@ -518,8 +518,9 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) } tree = &RangeQueryTree{ &db.Bitmap{ - query.Args["id"].(uint64), - query.Args["frame"].(string), 0}, + Id: query.Args["id"].(uint64), + FrameType: query.Args["frame"].(string), + Filter: 0}, slice, query.Args["start"].(time.Time), query.Args["end"].(time.Time)} diff --git a/query/planner_test.go b/query/planner_test.go index 5f0a04f2b..d5f268468 100644 --- a/query/planner_test.go +++ b/query/planner_test.go @@ -34,7 +34,7 @@ func TestQueryPlanner_Plan_Get(t *testing.T) { t.Fatalf("unexpected step(0) operation: %s", step.Operation) } else if step.Slice != 0 { t.Fatalf("unexpected step(0) slice: %d", step.Slice) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{10, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 10, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(0) bitmap: %s", spew.Sprint(step.Bitmap)) } @@ -42,7 +42,7 @@ func TestQueryPlanner_Plan_Get(t *testing.T) { t.Fatalf("unexpected step(1) operation: %s", step.Operation) } else if step.Slice != 1 { t.Fatalf("unexpected step(1) slice: %d", step.Slice) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{10, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 10, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(1) bitmap: %s", spew.Sprint(step.Bitmap)) } @@ -80,7 +80,7 @@ func TestQueryPlanner_Plan_Set(t *testing.T) { t.Fatalf("unexpected step(0) operation: %s", step.Operation) } else if step.ProfileId != 100 { t.Fatalf("unexpected step(0) profile id: %d", step.ProfileId) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{10, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 10, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(0) bitmap: %s", spew.Sprint(step.Bitmap)) } } @@ -107,7 +107,7 @@ func TestQueryPlanner_Plan_TopN(t *testing.T) { if step := (*plan)[0].(query.GetQueryStep); step.Operation != "get" { t.Fatalf("unexpected step(0) operation: %s", step.Operation) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{10, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 10, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(0) bitmap: %s", spew.Sprint(step.Bitmap)) } @@ -123,7 +123,7 @@ func TestQueryPlanner_Plan_TopN(t *testing.T) { if step := (*plan)[2].(query.GetQueryStep); step.Operation != "get" { t.Fatalf("unexpected step(2) operation: %s", step.Operation) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{10, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 10, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(2) bitmap: %s", spew.Sprint(step.Bitmap)) } @@ -196,7 +196,7 @@ func TestQueryPlanner_Plan_Union(t *testing.T) { t.Fatalf("unexpected step(0) operation: %s", step.Operation) } else if step.Slice != 0 { t.Fatalf("unexpected step(0) slice: %d", step.Slice) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{10, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 10, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(0) bitmap: %s", spew.Sprint(step.Bitmap)) } @@ -205,7 +205,7 @@ func TestQueryPlanner_Plan_Union(t *testing.T) { t.Fatalf("unexpected step(1) operation: %s", step.Operation) } else if step.Slice != 0 { t.Fatalf("unexpected step(1) slice: %d", step.Slice) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{20, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 20, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(1) bitmap: %s", spew.Sprint(step.Bitmap)) } @@ -224,7 +224,7 @@ func TestQueryPlanner_Plan_Union(t *testing.T) { t.Fatalf("unexpected step(3) operation: %s", step.Operation) } else if step.Slice != 1 { t.Fatalf("unexpected step(3) slice: %d", step.Slice) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{10, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 10, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(3) bitmap: %s", spew.Sprint(step.Bitmap)) } @@ -233,7 +233,7 @@ func TestQueryPlanner_Plan_Union(t *testing.T) { t.Fatalf("unexpected step(4) operation: %s", step.Operation) } else if step.Slice != 1 { t.Fatalf("unexpected step(4) slice: %d", step.Slice) - } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{20, "default", 0}) { + } else if !reflect.DeepEqual(step.Bitmap, &db.Bitmap{Id: 20, FrameType: "default", Filter: 0}) { t.Fatalf("unexpected step(4) bitmap: %s", spew.Sprint(step.Bitmap)) } diff --git a/transport/tcp.go b/transport/tcp.go index 31415b2b6..5421ae771 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -52,7 +52,7 @@ BeginManageConnection: } self.conn = &conn go func() { - self.outbox <- &db.Message{self.transport.ID.String()} + self.outbox <- &db.Message{Data: self.transport.ID.String()} }() } encoder := gob.NewEncoder(*self.conn) @@ -170,7 +170,7 @@ func (self *TcpTransport) Close() { func (self *TcpTransport) Send(message *db.Message, host *util.GUID) { log.Trace("TcpTransport.Send", message, host) - envelope := db.Envelope{message, host} + envelope := db.Envelope{Message: message, Host: host} notify.Post("outbox", &envelope) self.outbox <- envelope }