diff --git a/db/constants.go b/db/constants.go new file mode 100644 index 000000000..d2fbdd0eb --- /dev/null +++ b/db/constants.go @@ -0,0 +1,3 @@ +package db + +const SLICE_WIDTH = 65536 diff --git a/db/topology.go b/db/topology.go index 64eba14bc..0e815342f 100644 --- a/db/topology.go +++ b/db/topology.go @@ -10,6 +10,8 @@ import ( ) var FrameDoesNotExistError = errors.New("Frame does not exist.") +var SliceDoesNotExistError = errors.New("Slice does not exist.") +var FrameSliceIntersectDoesNotExistError = errors.New("FrameSliceIntersect does not exist.") type Location struct { Ip string @@ -39,32 +41,36 @@ type NodeMap map[Location]Location // A fragment is a collection of bitmaps within a slice. The fragment contains a reference to the responsible node for that fragment. The node is in the form ip:port type Fragment struct { - Node string + location *Location + id int } // A slice is the vertical combination of every fragment. It contains the hashring used to delegate bitmaps to fragments type Slice struct { - Fragments []Fragment - Hashring *consistent.Consistent + id int } // A frame is a collection of slices in a given category (brands, demographics, etc), specific to a database type Frame struct { Name string - Slices []*Slice } -// Add a slice to a frame with given Node addresses -func (f *Frame) AddSlice(addrs ...string) *Slice { - slice := Slice{} - slice.Hashring = consistent.New() - slice.Hashring.NumberOfReplicas = 200 - sliceIndex := len(f.Slices) - for index, addr := range addrs { - slice.Fragments = append(slice.Fragments, Fragment{addr}) - slice.Hashring.Add(fmt.Sprintf("%d %d", sliceIndex, index)) - } - f.Slices = append(f.Slices, &slice) +type FrameSliceIntersect struct { + slice *Slice + frame *Frame + Fragments []Fragment + Hashring *consistent.Consistent +} + +// Add a slice to a database +func (d *Database) AddSlice() *Slice { + slice_id := len(d.slices) + slice := Slice{id: slice_id} + d.slices = append(d.slices, &slice) + // add intersections + for _, frame := range d.frames { + d.AddFrameSliceIntersect(frame, &slice) + } return &slice } @@ -87,27 +93,63 @@ func (c *Cluster) AddDatabase(name string) *Database { // A database is a collection of all the frames within a given profile space type Database struct { Name string - Frames []*Frame + frames []*Frame + slices []*Slice + FrameSliceIntersects []*FrameSliceIntersect } // Count the number of slices in a database func (d *Database) NumSlices() (int, error) { - if len(d.Frames) < 1 { + if len(d.slices) < 1 { return 0, errors.New("Database is empty") } - return len(d.Frames[0].Slices), nil + return len(d.slices), nil } // Add a frame to a database func (d *Database) AddFrame(name string) *Frame { frame := Frame{Name: name} - d.Frames = append(d.Frames, &frame) + d.frames = append(d.frames, &frame) + // add intersections + for _, slice := range d.slices { + d.AddFrameSliceIntersect(&frame, slice) + } return &frame } +func (d *Database) AddFragment(frame *Frame, slice *Slice, location *Location, fragment_id int) *Fragment { + + frameslice, _ := d.GetFrameSliceIntersect(frame, slice) + fragment := Fragment{location: location, id: fragment_id} + frameslice.Fragments = append(frameslice.Fragments, fragment) + + frameslice.Hashring.Add(fmt.Sprintf("%d", fragment_id)) + + return &fragment +} + +func (d *Database) AddFrameSliceIntersect(frame *Frame, slice *Slice) *FrameSliceIntersect { + frameslice := FrameSliceIntersect{frame: frame, slice: slice} + d.FrameSliceIntersects = append(d.FrameSliceIntersects, &frameslice) + + frameslice.Hashring = consistent.New() + frameslice.Hashring.NumberOfReplicas = 16 + + return &frameslice +} + +func (d *Database) GetFrameSliceIntersect(frame *Frame, slice *Slice) (*FrameSliceIntersect, error) { + for _, frameslice := range d.FrameSliceIntersects { + if frameslice.frame == frame && frameslice.slice == slice { + return frameslice, nil + } + } + return nil, FrameSliceIntersectDoesNotExistError +} + // Get a frame from a database func (d *Database) GetFrame(name string) (*Frame, error) { - for _, frame := range d.Frames { + for _, frame := range d.frames { if frame.Name == name { return frame, nil } @@ -115,20 +157,38 @@ func (d *Database) GetFrame(name string) (*Frame, error) { return nil, FrameDoesNotExistError } -// For debugging, prints cluster information -func (c *Cluster) Describe() { - for _, database := range c.Databases { - log.Println("frames", database.Frames) - for _, frame := range database.Frames { - log.Println(frame.Name, database.Name) - for _, slice := range frame.Slices { - log.Println(" ", slice) - } +// Get a slice from a database +func (d *Database) GetSlice(id int) (*Slice, error) { + for _, slice := range d.slices { + if slice.id == id { + return slice, nil } } + return nil, SliceDoesNotExistError } -type Bitmap struct { - FrameType string - Id int +// Get a slice from a database +func (d *Database) GetSliceForProfile(profile_id int) (*Slice, error) { + log.Println("GetSliceForProfile") + log.Println("profile_id:",profile_id) + slice_id := profile_id / SLICE_WIDTH + return d.GetSlice(slice_id) +} + + +type Bitmap struct { + Id int + FrameType string +} + +func (d *Database) TestSetBit(bitmap Bitmap, profile_id int) { + log.Println("TestSetBit") + slice, _ := d.GetSliceForProfile(profile_id) + log.Println("slice:",slice) + frame, _ := d.GetFrame(bitmap.FrameType) + fsi, _ := d.GetFrameSliceIntersect(frame, slice) + fragment,erry := fsi.Hashring.Get(fmt.Sprintf("%d", bitmap.Id)) + + log.Println("fragment:",fragment) + log.Println("error:",erry) } diff --git a/db/topology_test.go b/db/topology_test.go new file mode 100644 index 000000000..d26efdd0b --- /dev/null +++ b/db/topology_test.go @@ -0,0 +1,84 @@ +package db + +import ( + "testing" + "log" + . "github.com/smartystreets/goconvey/convey" +) + +func TestTopology(t *testing.T) { + Convey("Basic DB structures", t, func() { + log.Println("topology test") + + /* + cluster := Cluster{Self:"192.168.1.100:1201"} + database := cluster.AddDatabase("property49") + frame := database.AddFrame("general") + frame.AddSlice("192.168.1.100:1201", "192.168.1.100:1202", "192.168.1.100:1203") + + log.Println(cluster) + log.Println(database) + log.Println(frame) + */ + + cluster := Cluster{Self:"192.168.1.100:1201"} + database := cluster.AddDatabase("property49") + database.AddFrame("general") + //database.AddFrame("brands") + database.AddSlice() + database.AddSlice() + //database.AddSlice() + + log.Println(database) + log.Println("----------------------------------") + for _, fsi := range database.FrameSliceIntersects { + log.Println(fsi.frame,fsi.slice) + } + num_slices, _ := database.NumSlices() + log.Println(num_slices) + + + frame, _ := database.GetFrame("general") + slice, _ := database.GetSlice(0) + + loc1, _ := NewLocation("192.168.1.100:8001") + /* + loc2, _ := NewLocation("192.168.1.100:8002") + loc3, _ := NewLocation("192.168.1.100:8003") + loc4, _ := NewLocation("192.168.1.100:8004") + loc5, _ := NewLocation("192.168.1.100:8005") + */ + + database.AddFragment(frame, slice, loc1, 0) + database.AddFragment(frame, slice, loc1, 1) + database.AddFragment(frame, slice, loc1, 2) + database.AddFragment(frame, slice, loc1, 3) + database.AddFragment(frame, slice, loc1, 4) + + + fsi, _ := database.GetFrameSliceIntersect(frame, slice) + log.Println(fsi) + //log.Println(fsi.Hashring) + + x,_ := fsi.Hashring.Get("able") + log.Println(x) + + /* + slicer, errer := database.GetSliceForProfile(131072) + log.Println(slicer) + log.Println(errer) + database.AddSlice() + slicer, errer = database.GetSliceForProfile(131072) + log.Println(slicer) + log.Println(errer) + */ + + + bitmap := Bitmap{Id: 555, FrameType: "general"} + log.Println("bitmap:",bitmap) + + database.TestSetBit(bitmap, 65535) + + + }) +} diff --git a/query/lexer.go b/query/lexer.go new file mode 100644 index 000000000..68762b6ac --- /dev/null +++ b/query/lexer.go @@ -0,0 +1,186 @@ +package query + +import ( + "errors" + "strings" + "log" + "unicode" + "unicode/utf8" + //"github.com/davecgh/go-spew/spew" +) + +const ( + TYPE_FUNC = iota + TYPE_LP = iota + TYPE_RP = iota + TYPE_ID = iota + TYPE_COMMA = iota +) + +type Token struct { + Text string + Type int +} + +type statefn func(lexer *Lexer) statefn + +type Lexer struct { + text string // the string being scanned. + pos int // current position in the input. + width int // width of last rune read from input. + start int // start position of this item. + state int // current state of lexer NEEDED??? + ch chan Token // channel of scanned items (Tokens). +} + + +func (lexer *Lexer) emit(typ int) { + lexer.ch <- Token{lexer.text[lexer.start:lexer.pos], typ} + lexer.start = lexer.pos +} + +func (lexer *Lexer) acceptUntil(chars string) error { + for { + if strings.HasPrefix(lexer.text[lexer.pos:], chars) { + return nil + } + // if we receive a reserved character that we are not expecting, throw a parse error + lexer.pos += 1 + if lexer.pos > len(lexer.text) { + return errors.New("Parse error, expecting " + string(chars)) + } + } +} + +func (lexer *Lexer) acceptRun(valid string) { + for strings.IndexRune(valid, lexer.next()) >= 0 { + } + lexer.backup() +} + +// next returns the next rune in the input. +func (lexer *Lexer) next() (runey rune) { + if lexer.pos >= len(lexer.text) { + lexer.width = 0 + return 0 + } + runey, lexer.width = utf8.DecodeRuneInString(lexer.text[lexer.pos:]) + lexer.pos += lexer.width + return runey +} + +// ignore skips over the pending input before this point. +func (lexer *Lexer) ignore() { + lexer.start = lexer.pos +} + +// backup steps back one rune. +// Can be called only once per call of next. +func (lexer *Lexer) backup() { + lexer.pos -= lexer.width +} + +// peek returns but does not consume +// the next rune in the input. +func (lexer *Lexer) peek() rune { + for { + next_rune := lexer.next() + // ignore spaces + if next_rune != rune(' ') { + lexer.backup() + return next_rune + } + lexer.ignore() + } +} + +func stateFunc(lexer *Lexer) statefn { + err := lexer.acceptUntil("(") + if err != nil { + log.Fatal(err) + } + lexer.emit(TYPE_FUNC) + return stateLP +} + +func stateLP(lexer *Lexer) statefn { + lexer.pos += 1 + lexer.emit(TYPE_LP) + // handle multiple LPs + if lexer.peek() == rune('(') { + return stateLP + } + return stateArgs +} + +func stateArgs(lexer *Lexer) statefn { + if unicode.IsNumber(lexer.peek()) { + return stateID + } else { + return stateFunc + } +} + +func stateID(lexer *Lexer) statefn { + + digits := "0123456789" + lexer.acceptRun(digits) + lexer.emit(TYPE_ID) + + // if next is comma + peeked := lexer.peek() + if peeked == rune(',') { + return stateComma + } else if peeked == rune(')') { + return stateRP + } else { + return stateID + } +} + +func stateRP(lexer *Lexer) statefn { + lexer.pos += 1 + lexer.emit(TYPE_RP) + + peeked := lexer.peek() + if peeked == rune(',') { + return stateComma + } else if peeked == rune(')') { + return stateRP + } else { + return stateEOF + } +} + +func stateComma(lexer *Lexer) statefn { + lexer.pos += 1 + lexer.emit(TYPE_COMMA) + return stateArgs +} + +func stateEOF(lexer *Lexer) statefn { + close(lexer.ch) + return nil +} + +func (lexer *Lexer) Lex() []Token{ + tokens := make([]Token, 0) + state := stateFunc + go func () { + for { + state = state(lexer) + if state == nil { + return + } + } + }() + for t := range lexer.ch { + tokens = append(tokens, t) + } + return tokens +} + +func Lex(input string) []Token { + lexer := Lexer{input, 0, 0, 0, TYPE_FUNC, make(chan Token)} + return lexer.Lex() +} diff --git a/query/lexer_test.go b/query/lexer_test.go new file mode 100644 index 000000000..f0c632384 --- /dev/null +++ b/query/lexer_test.go @@ -0,0 +1,104 @@ +package query + +import ( + "testing" + . "github.com/smartystreets/goconvey/convey" +) + +func TestLexer(t *testing.T) { + Convey("Basic lexical analysis", t, func() { + tokens := Lex("get(10)") + So(len(tokens), ShouldEqual, 4) + So(tokens[0].Text, ShouldEqual, "get") + So(tokens[0].Type, ShouldEqual, TYPE_FUNC) + So(tokens[1].Text, ShouldEqual, "(") + So(tokens[1].Type, ShouldEqual, TYPE_LP) + So(tokens[2].Text, ShouldEqual, "10") + So(tokens[2].Type, ShouldEqual, TYPE_ID) + So(tokens[3].Text, ShouldEqual, ")") + So(tokens[3].Type, ShouldEqual, TYPE_RP) + + tokens2 := Lex("intersect(get(10), get(11), get(12))") + So(len(tokens2), ShouldEqual, 17) + So(tokens2[0].Text, ShouldEqual, "intersect") + So(tokens2[0].Type, ShouldEqual, TYPE_FUNC) + So(tokens2[1].Text, ShouldEqual, "(") + So(tokens2[1].Type, ShouldEqual, TYPE_LP) + So(tokens2[2].Text, ShouldEqual, "get") + So(tokens2[2].Type, ShouldEqual, TYPE_FUNC) + So(tokens2[3].Text, ShouldEqual, "(") + So(tokens2[3].Type, ShouldEqual, TYPE_LP) + So(tokens2[4].Text, ShouldEqual, "10") + So(tokens2[4].Type, ShouldEqual, TYPE_ID) + So(tokens2[5].Text, ShouldEqual, ")") + So(tokens2[5].Type, ShouldEqual, TYPE_RP) + So(tokens2[6].Text, ShouldEqual, ",") + So(tokens2[6].Type, ShouldEqual, TYPE_COMMA) + So(tokens2[7].Text, ShouldEqual, "get") + So(tokens2[7].Type, ShouldEqual, TYPE_FUNC) + So(tokens2[8].Text, ShouldEqual, "(") + So(tokens2[8].Type, ShouldEqual, TYPE_LP) + So(tokens2[9].Text, ShouldEqual, "11") + So(tokens2[9].Type, ShouldEqual, TYPE_ID) + So(tokens2[10].Text, ShouldEqual, ")") + So(tokens2[10].Type, ShouldEqual, TYPE_RP) + So(tokens2[11].Text, ShouldEqual, ",") + So(tokens2[11].Type, ShouldEqual, TYPE_COMMA) + So(tokens2[12].Text, ShouldEqual, "get") + So(tokens2[12].Type, ShouldEqual, TYPE_FUNC) + So(tokens2[13].Text, ShouldEqual, "(") + So(tokens2[13].Type, ShouldEqual, TYPE_LP) + So(tokens2[14].Text, ShouldEqual, "12") + So(tokens2[14].Type, ShouldEqual, TYPE_ID) + So(tokens2[15].Text, ShouldEqual, ")") + So(tokens2[15].Type, ShouldEqual, TYPE_RP) + So(tokens2[16].Text, ShouldEqual, ")") + So(tokens2[16].Type, ShouldEqual, TYPE_RP) + + tokens3 := Lex("intersect(get(10), get(11), concat(12,14))") + So(len(tokens3), ShouldEqual, 19) + So(tokens3[0].Text, ShouldEqual, "intersect") + So(tokens3[0].Type, ShouldEqual, TYPE_FUNC) + So(tokens3[1].Text, ShouldEqual, "(") + So(tokens3[1].Type, ShouldEqual, TYPE_LP) + So(tokens3[2].Text, ShouldEqual, "get") + So(tokens3[2].Type, ShouldEqual, TYPE_FUNC) + So(tokens3[3].Text, ShouldEqual, "(") + So(tokens3[3].Type, ShouldEqual, TYPE_LP) + So(tokens3[4].Text, ShouldEqual, "10") + So(tokens3[4].Type, ShouldEqual, TYPE_ID) + So(tokens3[5].Text, ShouldEqual, ")") + So(tokens3[5].Type, ShouldEqual, TYPE_RP) + So(tokens3[6].Text, ShouldEqual, ",") + So(tokens3[6].Type, ShouldEqual, TYPE_COMMA) + So(tokens3[7].Text, ShouldEqual, "get") + So(tokens3[7].Type, ShouldEqual, TYPE_FUNC) + So(tokens3[8].Text, ShouldEqual, "(") + So(tokens3[8].Type, ShouldEqual, TYPE_LP) + So(tokens3[9].Text, ShouldEqual, "11") + So(tokens3[9].Type, ShouldEqual, TYPE_ID) + So(tokens3[10].Text, ShouldEqual, ")") + So(tokens3[10].Type, ShouldEqual, TYPE_RP) + So(tokens3[11].Text, ShouldEqual, ",") + So(tokens3[11].Type, ShouldEqual, TYPE_COMMA) + So(tokens3[12].Text, ShouldEqual, "concat") + So(tokens3[12].Type, ShouldEqual, TYPE_FUNC) + So(tokens3[13].Text, ShouldEqual, "(") + So(tokens3[13].Type, ShouldEqual, TYPE_LP) + So(tokens3[14].Text, ShouldEqual, "12") + So(tokens3[14].Type, ShouldEqual, TYPE_ID) + So(tokens3[15].Text, ShouldEqual, ",") + So(tokens3[15].Type, ShouldEqual, TYPE_COMMA) + So(tokens3[16].Text, ShouldEqual, "14") + So(tokens3[16].Type, ShouldEqual, TYPE_ID) + So(tokens3[17].Text, ShouldEqual, ")") + So(tokens3[17].Type, ShouldEqual, TYPE_RP) + So(tokens3[18].Text, ShouldEqual, ")") + So(tokens3[18].Type, ShouldEqual, TYPE_RP) + + tokens4 := Lex("concat(1,2,345,890)") + So(len(tokens4), ShouldEqual, 10) + So(tokens4[6].Text, ShouldEqual, "345") + So(tokens4[6].Type, ShouldEqual, TYPE_ID) + }) +} diff --git a/query/parser.go b/query/parser.go index 130b48bba..c94b8cc26 100644 --- a/query/parser.go +++ b/query/parser.go @@ -1,105 +1,12 @@ package query import ( - "encoding/json" + "encoding/json" "errors" - "log" "pilosa/db" - "github.com/davecgh/go-spew/spew" + //"github.com/davecgh/go-spew/spew" ) -const ( - TYPE_FUNC = iota - TYPE_LP = iota - TYPE_RP = iota - TYPE_ID = iota -) - -type Token struct { - Text string - Type int -} - -type statefn func(lexer *Lexer) statefn - -type Lexer struct { - text string - pos int - start int - state int - ch chan Token -} - -func (lexer *Lexer) emit(typ int) { - lexer.ch <- Token{lexer.text[lexer.start:lexer.pos], typ} - lexer.start = lexer.pos -} - -func (lexer *Lexer) accept(char uint8) error { - for { - if lexer.text[lexer.pos] == char { - return nil - } - lexer.pos += 1 - if lexer.pos > len(lexer.text) { - return errors.New("Parse error, expecting " + string(char)) - } - } -} - -func stateFunc(lexer *Lexer) statefn { - err := lexer.accept('(') - if err != nil { - log.Fatal(err) - } - lexer.emit(TYPE_FUNC) - return stateLP -} - -func stateLP(lexer *Lexer) statefn { - lexer.pos += 1 - lexer.emit(TYPE_LP) - return stateID -} - -func stateID(lexer *Lexer) statefn { - err := lexer.accept(')') - if err != nil { - log.Fatal(err) - } - lexer.emit(TYPE_ID) - return stateRP -} - -func stateRP(lexer *Lexer) statefn { - lexer.pos += 1 - lexer.emit(TYPE_RP) - close(lexer.ch) - return nil -} - -func (lexer *Lexer) Lex() []Token{ - tokens := make([]Token, 0) - state := stateFunc - go func () { - for { - state = state(lexer) - if state == nil { - return - } - } - }() - for t := range lexer.ch { - spew.Dump(t) - tokens = append(tokens, t) - } - return tokens -} - -func Lex(input string) []Token { - lexer := Lexer{input, 0, 0, TYPE_FUNC, make(chan Token)} - return lexer.Lex() -} var InvalidQueryError = errors.New("Invalid query format.") @@ -141,7 +48,7 @@ func (q *QueryParser) Walk(data interface{}) (*Query, error) { return nil, InvalidQueryError } id_int := int(id) - query.Inputs = []QueryInput{db.Bitmap{frame, id_int}} + query.Inputs = []QueryInput{db.Bitmap{id_int, frame}} } return query, nil diff --git a/query/parser_test.go b/query/parser_test.go deleted file mode 100644 index 772aab41c..000000000 --- a/query/parser_test.go +++ /dev/null @@ -1,21 +0,0 @@ -package query - -import ( - "testing" - . "github.com/smartystreets/goconvey/convey" -) - -func TestParser(t *testing.T) { - Convey("Basic parsing", t, func() { - tokens := Lex("get(10)") - So(len(tokens), ShouldEqual, 4) - So(tokens[0].Text, ShouldEqual, "get") - So(tokens[0].Type, ShouldEqual, TYPE_FUNC) - So(tokens[1].Text, ShouldEqual, "(") - So(tokens[1].Type, ShouldEqual, TYPE_LP) - So(tokens[2].Text, ShouldEqual, "10") - So(tokens[2].Type, ShouldEqual, TYPE_ID) - So(tokens[3].Text, ShouldEqual, ")") - So(tokens[3].Type, ShouldEqual, TYPE_RP) - }) -} diff --git a/query/planner.go b/query/planner.go index b70261468..8551e0438 100644 --- a/query/planner.go +++ b/query/planner.go @@ -2,7 +2,7 @@ package query import ( "github.com/nu7hatch/gouuid" - "strconv" + //"strconv" "fmt" "math/rand" "pilosa/db" @@ -69,6 +69,7 @@ type GetQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *GetQueryTree) getLocation(d *db.Database) string { + /* frame, err := d.GetFrame(qt.bitmap.FrameType) if err != nil { panic(err) @@ -82,6 +83,8 @@ func (qt *GetQueryTree) getLocation(d *db.Database) string { fragment := slice.Fragments[fragIndex] return fmt.Sprintf(fragment.Node) + */ + return "Nothing yet" } // Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices diff --git a/query/planner_test.go b/query/planner_test.go index db6e6a623..f63aa4160 100644 --- a/query/planner_test.go +++ b/query/planner_test.go @@ -3,8 +3,11 @@ package query import ( "testing" "log" + . "github.com/smartystreets/goconvey/convey" ) func TestQueryPlanner(t *testing.T) { - log.Println("query planner test") + Convey("Basic query plan", t, func() { + log.Println("query planner test") + }) }