This commit is contained in:
Ben Johnson 2015-08-20 13:46:39 -06:00
parent 6fdd32d5ce
commit 73631de9e8
12 changed files with 89 additions and 46 deletions

View file

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

View file

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

View file

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

View file

@ -7,7 +7,7 @@ import (
)
type Message struct {
Data interface{} `json:data`
Data interface{} `json:"data"`
}
type Envelope struct {

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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