From f6c712caf62cccbac3a0f00088582f94a17cf190 Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 6 Dec 2013 07:54:08 -0600 Subject: [PATCH 01/14] Lexer updates. DB topology updates --- db/constants.go | 3 + db/topology.go | 124 ++++++++++++++++++++-------- db/topology_test.go | 84 +++++++++++++++++++ query/lexer.go | 186 ++++++++++++++++++++++++++++++++++++++++++ query/lexer_test.go | 104 +++++++++++++++++++++++ query/parser.go | 99 +--------------------- query/parser_test.go | 21 ----- query/planner.go | 5 +- query/planner_test.go | 5 +- 9 files changed, 480 insertions(+), 151 deletions(-) create mode 100644 db/constants.go create mode 100644 db/topology_test.go create mode 100644 query/lexer.go create mode 100644 query/lexer_test.go delete mode 100644 query/parser_test.go 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") + }) } From 63f33116028875c25dfca5622372deda8cacce55 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 6 Dec 2013 11:01:22 -0600 Subject: [PATCH 02/14] Fix nexter --- commands/pilosa-nexter/nexter.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/commands/pilosa-nexter/nexter.go b/commands/pilosa-nexter/nexter.go index 52abc33bb..17364eeb5 100644 --- a/commands/pilosa-nexter/nexter.go +++ b/commands/pilosa-nexter/nexter.go @@ -51,7 +51,7 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { log.Fatal(err) } } else { // No error, get start of series from etcd node - start, err = strconv.ParseUint(node.Node.Value, 10, 0) + start, err = strconv.ParseUint(node.Value, 10, 0) end = start + blocksize if err != nil { log.Fatal(err) @@ -64,7 +64,7 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { } else { log.Println("Error with CompareAndSet! Trying again in 1 second...") time.Sleep(time.Second) - start, err = strconv.ParseUint(newval.Node.Value, 10, 0) + start, err = strconv.ParseUint(newval.Value, 10, 0) if err != nil { log.Fatal(err) } From e72d147c3d4bcd60238ef902d6e609c1493ceef3 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 6 Dec 2013 15:54:33 -0600 Subject: [PATCH 03/14] Add beginnings of etcd topology sync. --- commands/pilosa-nexter/nexter.go | 4 +- core/etcd.go | 213 +++++++++++++++++++++---------- cruncher/cruncher.go | 13 +- db/topology.go | 34 ++++- db/topology_test.go | 2 +- deps.json | 2 +- router/router.go | 94 -------------- 7 files changed, 183 insertions(+), 179 deletions(-) delete mode 100644 router/router.go diff --git a/commands/pilosa-nexter/nexter.go b/commands/pilosa-nexter/nexter.go index 17364eeb5..52abc33bb 100644 --- a/commands/pilosa-nexter/nexter.go +++ b/commands/pilosa-nexter/nexter.go @@ -51,7 +51,7 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { log.Fatal(err) } } else { // No error, get start of series from etcd node - start, err = strconv.ParseUint(node.Value, 10, 0) + start, err = strconv.ParseUint(node.Node.Value, 10, 0) end = start + blocksize if err != nil { log.Fatal(err) @@ -64,7 +64,7 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { } else { log.Println("Error with CompareAndSet! Trying again in 1 second...") time.Sleep(time.Second) - start, err = strconv.ParseUint(newval.Value, 10, 0) + start, err = strconv.ParseUint(newval.Node.Value, 10, 0) if err != nil { log.Fatal(err) } diff --git a/core/etcd.go b/core/etcd.go index bb19b91d4..66679beef 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -2,87 +2,162 @@ package core import ( "github.com/coreos/go-etcd/etcd" - "log" - "strings" - "time" "encoding/gob" "pilosa/db" + "github.com/davecgh/go-spew/spew" + "github.com/nu7hatch/gouuid" + "strconv" + "log" ) func (service *Service) SetupEtcd() { gob.Register(db.Location{}) service.Etcd = etcd.NewClient(nil) - service.NodeMapMutex.Lock() - defer service.NodeMapMutex.Unlock() - service.NodeMap = db.NodeMap{} + //service.NodeMapMutex.Lock() + //defer service.NodeMapMutex.Unlock() + //service.NodeMap = db.NodeMap{} - nodes, err := service.Etcd.Get("nodes", false) + //nodes, err := service.Etcd.Get("nodes", false) + //if err != nil { + // log.Fatal(err) + //} + //for _, node := range nodes.Kvs { + // nodestring := strings.Split(node.Key, "/")[2] + // location, err := db.NewLocation(nodestring) + // if err != nil { + // log.Fatal(err) + // } + // routerlocation, err := db.NewLocation(node.Value) + // if err != nil { + // log.Fatal(err) + // } + // service.NodeMap[*location] = *routerlocation + //} + //log.Println(service.NodeMap) +} + +//func (service *Service) WatchEtcd() { +// var receiver = make(chan *etcd.Response) +// var stop chan bool +// go func () { +// _, err := service.Etcd.Watch("nodes/", 0, receiver, stop) +// if err != nil { +// log.Fatal(err) +// } +// }() +// +// exit, done := service.GetExitChannels() +// +// for { +// select { +// case response := <-receiver: +// switch response.Action { +// case "SET": +// nodestring := strings.Split(response.Key, "/")[2] +// node, err := db.NewLocation(nodestring) +// if err != nil { +// log.Fatal(err) +// } +// router, err := db.NewLocation(response.Value) +// if err != nil { +// log.Fatal(err) +// } +// service.NodeMapMutex.Lock() +// service.NodeMap[*node] = *router +// service.NodeMapMutex.Unlock() +// case "DELETE": +// nodestring := strings.Split(response.Key, "/")[2] +// node, err := db.NewLocation(nodestring) +// if err != nil { +// log.Fatal(err) +// } +// service.NodeMapMutex.Lock() +// delete(service.NodeMap, *node) +// service.NodeMapMutex.Unlock() +// default: +// log.Println("unhandled etcd message", response) +// } +// //log.Println(response.Action, response.Key, response.Value) +// log.Println(service.NodeMap) +// case <-exit: +// log.Println("cleaning up watchetcd service thing.") +// time.Sleep(time.Second/2) +// log.Println("done!") +// done <- 1 +// } +// } +//} + + +func (service *Service) MetaWatcher() { + namespace := "/pilosa/0" + log.Println(namespace + "/db") + resp, err := service.Etcd.Get(namespace + "/db", false, true) if err != nil { log.Fatal(err) } - for _, node := range nodes.Kvs { - nodestring := strings.Split(node.Key, "/")[2] - location, err := db.NewLocation(nodestring) - if err != nil { - log.Fatal(err) - } - routerlocation, err := db.NewLocation(node.Value) - if err != nil { - log.Fatal(err) - } - service.NodeMap[*location] = *routerlocation - } - log.Println(service.NodeMap) -} + cluster := db.NewCluster() -func (service *Service) WatchEtcd() { - var receiver = make(chan *etcd.Response) - var stop chan bool - go func () { - _, err := service.Etcd.Watch("nodes/", 0, receiver, stop) - if err != nil { - log.Fatal(err) + for _, database_ref := range resp.Node.Nodes { + database_name := database_ref.Key[len(namespace)+4:] + database := cluster.AddDatabase(database_name) + for _, database_attr_ref := range database_ref.Nodes { + key := database_attr_ref.Key[len(database_ref.Key)+1:] + if key == "frame" { + for _, frame_ref := range database_attr_ref.Nodes { + frame_name := frame_ref.Key[len(database_attr_ref.Key)+1:] + frame := database.AddFrame(frame_name) + for _, frame_attr_ref := range frame_ref.Nodes { + key = frame_attr_ref.Key[len(frame_ref.Key)+1:] + if key == "slice" { + for _, slice_ref := range frame_attr_ref.Nodes { + slice_name := slice_ref.Key[len(frame_attr_ref.Key)+1:] + slice_id, err := strconv.Atoi(slice_name) + if err != nil { + log.Fatal(err) + } + slice := database.AddSlice(slice_id) + for _, slice_attr_ref := range slice_ref.Nodes { + key = slice_attr_ref.Key[len(slice_ref.Key)+1:] + if key == "fragment" { + for _, fragment_ref := range slice_attr_ref.Nodes { + fragment_name := fragment_ref.Key[len(slice_attr_ref.Key)+1:] + fragment_id, err := strconv.Atoi(fragment_name) + if err != nil { + log.Fatal(err) + } + for _, fragment_attr_ref := range fragment_ref.Nodes { + key = fragment_attr_ref.Key[len(fragment_ref.Key)+1:] + if key == "node" { + uuid, err := uuid.ParseHex(fragment_attr_ref.Value) + if err != nil { + log.Fatal(err) + } + process := db.NewProcess(uuid) + database.AddFragment(frame, slice, process, fragment_id) + } + } + } + } + } + } + } + } + } + } + } + } + database, _ := cluster.GetDatabase("main") + database.TestSetBit(db.Bitmap{0, "general"}, 0) + + receiver := make(chan *etcd.Response) + stop := make(chan bool) + go func() { + _, _ = service.Etcd.Watch(namespace + "/db", 0, true, receiver, stop) + }() + go func() { + for x := range receiver { + spew.Dump(x) } }() - - exit, done := service.GetExitChannels() - - for { - select { - case response := <-receiver: - switch response.Action { - case "SET": - nodestring := strings.Split(response.Key, "/")[2] - node, err := db.NewLocation(nodestring) - if err != nil { - log.Fatal(err) - } - router, err := db.NewLocation(response.Value) - if err != nil { - log.Fatal(err) - } - service.NodeMapMutex.Lock() - service.NodeMap[*node] = *router - service.NodeMapMutex.Unlock() - case "DELETE": - nodestring := strings.Split(response.Key, "/")[2] - node, err := db.NewLocation(nodestring) - if err != nil { - log.Fatal(err) - } - service.NodeMapMutex.Lock() - delete(service.NodeMap, *node) - service.NodeMapMutex.Unlock() - default: - log.Println("unhandled etcd message", response) - } - //log.Println(response.Action, response.Key, response.Value) - log.Println(service.NodeMap) - case <-exit: - log.Println("cleaning up watchetcd service thing.") - time.Sleep(time.Second/2) - log.Println("done!") - done <- 1 - } - } } diff --git a/cruncher/cruncher.go b/cruncher/cruncher.go index 309a22a48..980699f63 100644 --- a/cruncher/cruncher.go +++ b/cruncher/cruncher.go @@ -18,12 +18,13 @@ func (c *Cruncher) Run() { log.Println("Running cruncher...") c.SetupEtcd() //go r.SyncEtcd() - go c.WatchEtcd() - go c.HandleConnections() - c.SetupNetwork() - go c.Serve() - go c.HandleInbox() - go c.ServeHTTP() + //go c.WatchEtcd() + //go c.HandleConnections() + //c.SetupNetwork() + //go c.Serve() + //go c.HandleInbox() + //go c.ServeHTTP() + go c.MetaWatcher() sigterm, sighup := c.GetSignals() for { diff --git a/db/topology.go b/db/topology.go index 0e815342f..0f343e251 100644 --- a/db/topology.go +++ b/db/topology.go @@ -2,6 +2,7 @@ package db import ( "github.com/stathat/consistent" + "github.com/nu7hatch/gouuid" "log" "fmt" "errors" @@ -18,6 +19,14 @@ type Location struct { Port int } +type Process struct { + id *uuid.UUID +} + +func NewProcess(id *uuid.UUID) *Process { + return &Process{id} +} + // Create a Location struct given a string in form "0.0.0.0:0" func NewLocation(location_string string) (*Location, error) { splitstring := strings.Split(location_string, ":") @@ -41,7 +50,7 @@ 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 { - location *Location + process *Process id int } @@ -63,8 +72,7 @@ type FrameSliceIntersect struct { } // Add a slice to a database -func (d *Database) AddSlice() *Slice { - slice_id := len(d.slices) +func (d *Database) AddSlice(slice_id int) *Slice { slice := Slice{id: slice_id} d.slices = append(d.slices, &slice) // add intersections @@ -77,7 +85,12 @@ func (d *Database) AddSlice() *Slice { // Represents the entire cluster, and a reference to the Node this instance is running on type Cluster struct { Databases map[string]*Database - Self string +} + +func NewCluster() *Cluster { + cluster := Cluster{} + cluster.Databases = make(map[string]*Database) + return &cluster } // Add a database to a cluster @@ -90,6 +103,15 @@ func (c *Cluster) AddDatabase(name string) *Database { return &database } +func (c *Cluster) GetDatabase(name string) (*Database, error) { + value, ok := c.Databases[name] + if !ok { + return nil, errors.New("The database does not exist!!!!!!!!!!!!!!!!!!!!!!!! -HO") + } else { + return value, nil + } +} + // A database is a collection of all the frames within a given profile space type Database struct { Name string @@ -117,10 +139,10 @@ func (d *Database) AddFrame(name string) *Frame { return &frame } -func (d *Database) AddFragment(frame *Frame, slice *Slice, location *Location, fragment_id int) *Fragment { +func (d *Database) AddFragment(frame *Frame, slice *Slice, process *Process, fragment_id int) *Fragment { frameslice, _ := d.GetFrameSliceIntersect(frame, slice) - fragment := Fragment{location: location, id: fragment_id} + fragment := Fragment{process: process, id: fragment_id} frameslice.Fragments = append(frameslice.Fragments, fragment) frameslice.Hashring.Add(fmt.Sprintf("%d", fragment_id)) diff --git a/db/topology_test.go b/db/topology_test.go index d26efdd0b..8de4132be 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -21,7 +21,7 @@ func TestTopology(t *testing.T) { log.Println(frame) */ - cluster := Cluster{Self:"192.168.1.100:1201"} + cluster := NewCluster() database := cluster.AddDatabase("property49") database.AddFrame("general") //database.AddFrame("brands") diff --git a/deps.json b/deps.json index 8df70badd..fc18b8318 100644 --- a/deps.json +++ b/deps.json @@ -6,7 +6,7 @@ }, "etcd": { "repo": "github.com/coreos/go-etcd/etcd", - "version": "8a4461a676eb65fb74f10da1f8198cc9f67da366", + "version": "8a4461a", "type": "git" }, "goconvey": { diff --git a/router/router.go b/router/router.go deleted file mode 100644 index 5dc700c97..000000000 --- a/router/router.go +++ /dev/null @@ -1,94 +0,0 @@ -package router - -import ( - "log" - //"time" - //"strings" - "pilosa/core" - "pilosa/db" -) - -type Router struct { - core.Service -} - -func (r *Router) Init() { - log.Println("Initializing router...") -} - -func (r *Router) Run() { - log.Println("Running router...") - r.SetupEtcd() - //go r.SyncEtcd() - go r.WatchEtcd() - go r.HandleConnections() - go r.HandleInbox() - go r.ServeHTTP() - - //go func() { - // for { - // r.SendMessage(&core.Message{"ping", core.Location{"127.0.0.1", 1200}, core.Location{"127.0.0.1", 1300}}) - // time.Sleep(2*time.Second) - // //log.Println(r.GetRouterLocation(core.Location{"127.0.0.1", 1200})) - // } - //}() - - sigterm, sighup := r.GetSignals() - for { - select { - case <- sighup: - log.Println("SIGHUP! Reloading configuration...") - // TODO: reload configuration - case <- sigterm: - log.Println("SIGTERM! Cleaning up...") - r.Stop() - return - } - } -} - -func (r *Router) HandleInbox() { - for { - select { - case message := <-r.Inbox: - log.Println("process", message) - } - } -} - -func (r *Router) SetupEtcd() { - r.Service.SetupEtcd() - //routers, err := r.Etcd.Get("topology") - //if err != nil { - // log.Fatal(err) - //} - //for _, router := range routers { - // routerstring := strings.Split(router.Key, "/")[2] - // routerlocation, err := core.NewLocation(routerstring) - // if err != nil { - // log.Fatal(err) - // } - // nodes, err := r.Etcd.Get("topology/" + routerstring) - // if err != nil { - // log.Fatal(err) - // } - // for _, node := range nodes { - // nodestring := strings.Split(node.Key, "/")[3] - // nodelocation, err := core.NewLocation(nodestring) - // if err != nil { - // log.Fatal(err) - // } - // r.NodeMap[*nodelocation] = *routerlocation - // } - //} -} - -func (r *Router) HandleMessage(m *db.Message) { - log.Println(m) -} - -func NewRouter(tcp, http *db.Location) *Router { - service := core.NewService(tcp, http) - router := Router{*service} - return &router -} From f33cbdf168dcccbc87d1dbc14310165871ba602c Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 6 Dec 2013 16:13:25 -0600 Subject: [PATCH 04/14] AddFragment method --- db/topology.go | 14 ++++++++++++-- db/topology_test.go | 8 +++++--- 2 files changed, 17 insertions(+), 5 deletions(-) diff --git a/db/topology.go b/db/topology.go index 0f343e251..b3b016c8b 100644 --- a/db/topology.go +++ b/db/topology.go @@ -12,6 +12,7 @@ import ( var FrameDoesNotExistError = errors.New("Frame does not exist.") var SliceDoesNotExistError = errors.New("Slice does not exist.") +var FragmentDoesNotExistError = errors.New("Fragment does not exist.") var FrameSliceIntersectDoesNotExistError = errors.New("FrameSliceIntersect does not exist.") type Location struct { @@ -67,10 +68,19 @@ type Frame struct { type FrameSliceIntersect struct { slice *Slice frame *Frame - Fragments []Fragment + Fragments []*Fragment Hashring *consistent.Consistent } +func (fsi *FrameSliceIntersect) GetFragment(fragment_id int) (*Fragment, error) { + for _, fragment := range fsi.Fragments { + if fragment.id == fragment_id { + return fragment, nil + } + } + return nil, FragmentDoesNotExistError +} + // Add a slice to a database func (d *Database) AddSlice(slice_id int) *Slice { slice := Slice{id: slice_id} @@ -143,7 +153,7 @@ func (d *Database) AddFragment(frame *Frame, slice *Slice, process *Process, fra frameslice, _ := d.GetFrameSliceIntersect(frame, slice) fragment := Fragment{process: process, id: fragment_id} - frameslice.Fragments = append(frameslice.Fragments, fragment) + frameslice.Fragments = append(frameslice.Fragments, &fragment) frameslice.Hashring.Add(fmt.Sprintf("%d", fragment_id)) diff --git a/db/topology_test.go b/db/topology_test.go index 8de4132be..1feb7f795 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -25,8 +25,8 @@ func TestTopology(t *testing.T) { database := cluster.AddDatabase("property49") database.AddFrame("general") //database.AddFrame("brands") - database.AddSlice() - database.AddSlice() + database.AddSlice(0) + database.AddSlice(1) //database.AddSlice() log.Println(database) @@ -41,7 +41,7 @@ func TestTopology(t *testing.T) { frame, _ := database.GetFrame("general") slice, _ := database.GetSlice(0) - loc1, _ := NewLocation("192.168.1.100:8001") + //loc1, _ := NewLocation("192.168.1.100:8001") /* loc2, _ := NewLocation("192.168.1.100:8002") loc3, _ := NewLocation("192.168.1.100:8003") @@ -49,11 +49,13 @@ func TestTopology(t *testing.T) { 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) From a3f93a617b386efacba42ab988309f9f5d541c75 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 6 Dec 2013 16:38:29 -0600 Subject: [PATCH 05/14] add things. --- core/etcd.go | 3 ++- db/topology.go | 12 ++++++++---- 2 files changed, 10 insertions(+), 5 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 66679beef..95eaeb093 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -148,7 +148,8 @@ func (service *Service) MetaWatcher() { } } database, _ := cluster.GetDatabase("main") - database.TestSetBit(db.Bitmap{0, "general"}, 0) + spew.Dump(database) + database.TestSetBit(db.Bitmap{1200, "general"}, 1) receiver := make(chan *etcd.Response) stop := make(chan bool) diff --git a/db/topology.go b/db/topology.go index b3b016c8b..bb1521eda 100644 --- a/db/topology.go +++ b/db/topology.go @@ -2,6 +2,7 @@ package db import ( "github.com/stathat/consistent" + "github.com/davecgh/go-spew/spew" "github.com/nu7hatch/gouuid" "log" "fmt" @@ -219,8 +220,11 @@ func (d *Database) TestSetBit(bitmap Bitmap, profile_id int) { 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) + frag_id_s,err := fsi.Hashring.Get(fmt.Sprintf("%d", bitmap.Id)) + frag_id, err := strconv.Atoi(frag_id_s) + fragment, err := fsi.GetFragment(frag_id) + if err != nil { + log.Fatal(err) + } + spew.Dump(fragment) } From 40e7474d82052d70cf82cfed270ddb9ea9fef1c5 Mon Sep 17 00:00:00 2001 From: travisturner Date: Mon, 9 Dec 2013 13:57:40 -0600 Subject: [PATCH 06/14] lowercase some attributes --- db/topology.go | 42 +++++++++++++++++++++--------------------- db/topology_test.go | 19 ++++++++----------- 2 files changed, 29 insertions(+), 32 deletions(-) diff --git a/db/topology.go b/db/topology.go index bb1521eda..15c61561b 100644 --- a/db/topology.go +++ b/db/topology.go @@ -63,18 +63,18 @@ type Slice struct { // A frame is a collection of slices in a given category (brands, demographics, etc), specific to a database type Frame struct { - Name string + name string } type FrameSliceIntersect struct { slice *Slice frame *Frame - Fragments []*Fragment - Hashring *consistent.Consistent + fragments []*Fragment + hashring *consistent.Consistent } func (fsi *FrameSliceIntersect) GetFragment(fragment_id int) (*Fragment, error) { - for _, fragment := range fsi.Fragments { + for _, fragment := range fsi.fragments { if fragment.id == fragment_id { return fragment, nil } @@ -95,29 +95,29 @@ func (d *Database) AddSlice(slice_id int) *Slice { // Represents the entire cluster, and a reference to the Node this instance is running on type Cluster struct { - Databases map[string]*Database + databases map[string]*Database } func NewCluster() *Cluster { cluster := Cluster{} - cluster.Databases = make(map[string]*Database) + cluster.databases = make(map[string]*Database) return &cluster } // Add a database to a cluster func (c *Cluster) AddDatabase(name string) *Database { database := Database{Name: name} - if c.Databases == nil { - c.Databases = make(map[string]*Database) + if c.databases == nil { + c.databases = make(map[string]*Database) } - c.Databases[name] = &database + c.databases[name] = &database return &database } func (c *Cluster) GetDatabase(name string) (*Database, error) { - value, ok := c.Databases[name] + value, ok := c.databases[name] if !ok { - return nil, errors.New("The database does not exist!!!!!!!!!!!!!!!!!!!!!!!! -HO") + return nil, errors.New("The database does not exist!") } else { return value, nil } @@ -128,7 +128,7 @@ type Database struct { Name string frames []*Frame slices []*Slice - FrameSliceIntersects []*FrameSliceIntersect + frame_slice_intersects []*FrameSliceIntersect } // Count the number of slices in a database @@ -141,7 +141,7 @@ func (d *Database) NumSlices() (int, error) { // Add a frame to a database func (d *Database) AddFrame(name string) *Frame { - frame := Frame{Name: name} + frame := Frame{name: name} d.frames = append(d.frames, &frame) // add intersections for _, slice := range d.slices { @@ -154,25 +154,25 @@ func (d *Database) AddFragment(frame *Frame, slice *Slice, process *Process, fra frameslice, _ := d.GetFrameSliceIntersect(frame, slice) fragment := Fragment{process: process, id: fragment_id} - frameslice.Fragments = append(frameslice.Fragments, &fragment) + frameslice.fragments = append(frameslice.fragments, &fragment) - frameslice.Hashring.Add(fmt.Sprintf("%d", fragment_id)) + 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) + d.frame_slice_intersects = append(d.frame_slice_intersects, &frameslice) - frameslice.Hashring = consistent.New() - frameslice.Hashring.NumberOfReplicas = 16 + 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 { + for _, frameslice := range d.frame_slice_intersects { if frameslice.frame == frame && frameslice.slice == slice { return frameslice, nil } @@ -183,7 +183,7 @@ func (d *Database) GetFrameSliceIntersect(frame *Frame, slice *Slice) (*FrameSli // Get a frame from a database func (d *Database) GetFrame(name string) (*Frame, error) { for _, frame := range d.frames { - if frame.Name == name { + if frame.name == name { return frame, nil } } @@ -220,7 +220,7 @@ func (d *Database) TestSetBit(bitmap Bitmap, profile_id int) { log.Println("slice:",slice) frame, _ := d.GetFrame(bitmap.FrameType) fsi, _ := d.GetFrameSliceIntersect(frame, slice) - frag_id_s,err := fsi.Hashring.Get(fmt.Sprintf("%d", bitmap.Id)) + frag_id_s,err := fsi.hashring.Get(fmt.Sprintf("%d", bitmap.Id)) frag_id, err := strconv.Atoi(frag_id_s) fragment, err := fsi.GetFragment(frag_id) if err != nil { diff --git a/db/topology_test.go b/db/topology_test.go index 1feb7f795..aa7489739 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -31,7 +31,7 @@ func TestTopology(t *testing.T) { log.Println(database) log.Println("----------------------------------") - for _, fsi := range database.FrameSliceIntersects { + for _, fsi := range database.frame_slice_intersects { log.Println(fsi.frame,fsi.slice) } num_slices, _ := database.NumSlices() @@ -40,6 +40,8 @@ func TestTopology(t *testing.T) { frame, _ := database.GetFrame("general") slice, _ := database.GetSlice(0) + log.Println(frame) + log.Println(slice) //loc1, _ := NewLocation("192.168.1.100:8001") /* @@ -58,13 +60,10 @@ func TestTopology(t *testing.T) { */ - fsi, _ := database.GetFrameSliceIntersect(frame, slice) - log.Println(fsi) + //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) @@ -75,12 +74,10 @@ func TestTopology(t *testing.T) { log.Println(errer) */ + //bitmap := Bitmap{Id: 555, FrameType: "general"} + //log.Println("bitmap:",bitmap) - bitmap := Bitmap{Id: 555, FrameType: "general"} - log.Println("bitmap:",bitmap) - - database.TestSetBit(bitmap, 65535) - + //database.TestSetBit(bitmap, 65535) }) } From fc2ea849c36b524820796b8881647192d8ed1caf Mon Sep 17 00:00:00 2001 From: travisturner Date: Mon, 9 Dec 2013 15:21:17 -0600 Subject: [PATCH 07/14] remove Cruncher() struct. Use Service() instead --- commands/pilosa-cruncher/cruncher.go | 6 +-- core/service.go | 39 +++++++++++++++++++ cruncher/cruncher.go | 56 ---------------------------- 3 files changed, 42 insertions(+), 59 deletions(-) delete mode 100644 cruncher/cruncher.go diff --git a/commands/pilosa-cruncher/cruncher.go b/commands/pilosa-cruncher/cruncher.go index a94dd842f..0efe20441 100644 --- a/commands/pilosa-cruncher/cruncher.go +++ b/commands/pilosa-cruncher/cruncher.go @@ -1,7 +1,7 @@ package main import ( - "pilosa/cruncher" + "pilosa/core" "pilosa/db" "flag" "log" @@ -26,6 +26,6 @@ func main() { log.Fatal("Location not valid:", httpLoc) } - cruncher := cruncher.NewCruncher(tcp, http) - cruncher.Run() + service := core.NewService(tcp, http) + service.Run() } diff --git a/core/service.go b/core/service.go index c8914255e..871022d8b 100644 --- a/core/service.go +++ b/core/service.go @@ -290,3 +290,42 @@ func (service *Service) NewListener() chan *db.Message { ch := make(chan *db.Message) return ch } + + +//////////////////////////////////////////////// + + +func (service *Service) Run() { + log.Println("Running service...") + service.SetupEtcd() + //go r.SyncEtcd() + //go service.WatchEtcd() + //go service.HandleConnections() + //service.SetupNetwork() + //go service.Serve() + //go service.HandleInbox() + //go service.ServeHTTP() + go service.MetaWatcher() + + sigterm, sighup := service.GetSignals() + for { + select { + case <- sighup: + log.Println("SIGHUP! Reloading configuration...") + // TODO: reload configuration + case <- sigterm: + log.Println("SIGTERM! Cleaning up...") + service.Stop() + return + } + } +} + +func (service *Service) HandleInbox() { + for { + select { + case message := <-service.Inbox: + log.Println("process", message) + } + } +} diff --git a/cruncher/cruncher.go b/cruncher/cruncher.go deleted file mode 100644 index 980699f63..000000000 --- a/cruncher/cruncher.go +++ /dev/null @@ -1,56 +0,0 @@ -package cruncher - -import ( - "pilosa/core" - "pilosa/db" - "log" -) - -type Cruncher struct { - core.Service -} - -func (c *Cruncher) Init() { - log.Println("Initializing cruncher...") -} - -func (c *Cruncher) Run() { - log.Println("Running cruncher...") - c.SetupEtcd() - //go r.SyncEtcd() - //go c.WatchEtcd() - //go c.HandleConnections() - //c.SetupNetwork() - //go c.Serve() - //go c.HandleInbox() - //go c.ServeHTTP() - go c.MetaWatcher() - - sigterm, sighup := c.GetSignals() - for { - select { - case <- sighup: - log.Println("SIGHUP! Reloading configuration...") - // TODO: reload configuration - case <- sigterm: - log.Println("SIGTERM! Cleaning up...") - c.Stop() - return - } - } -} - -func NewCruncher(tcp, http *db.Location) *Cruncher { - service := core.NewService(tcp, http) - cruncher := Cruncher{*service} - return &cruncher -} - -func (c *Cruncher) HandleInbox() { - for { - select { - case message := <-c.Inbox: - log.Println("process", message) - } - } -} From 0cf2ddb08678277446f88067004c009bad235350 Mon Sep 17 00:00:00 2001 From: travisturner Date: Tue, 10 Dec 2013 14:26:59 -0600 Subject: [PATCH 08/14] first pass at GetOrCreate methods in db.topology --- core/etcd.go | 8 +- db/topology.go | 304 +++++++++++++++++++++++++++++++------------- db/topology_test.go | 25 +++- 3 files changed, 246 insertions(+), 91 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 95eaeb093..795ce4b1f 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -121,11 +121,13 @@ func (service *Service) MetaWatcher() { key = slice_attr_ref.Key[len(slice_ref.Key)+1:] if key == "fragment" { for _, fragment_ref := range slice_attr_ref.Nodes { + /* fragment_name := fragment_ref.Key[len(slice_attr_ref.Key)+1:] fragment_id, err := strconv.Atoi(fragment_name) if err != nil { log.Fatal(err) } + */ for _, fragment_attr_ref := range fragment_ref.Nodes { key = fragment_attr_ref.Key[len(fragment_ref.Key)+1:] if key == "node" { @@ -133,8 +135,8 @@ func (service *Service) MetaWatcher() { if err != nil { log.Fatal(err) } - process := db.NewProcess(uuid) - database.AddFragment(frame, slice, process, fragment_id) + process := uuid + database.AddFragment(frame, slice, process) } } } @@ -149,7 +151,7 @@ func (service *Service) MetaWatcher() { } database, _ := cluster.GetDatabase("main") spew.Dump(database) - database.TestSetBit(db.Bitmap{1200, "general"}, 1) + database.OldGetFragment(db.Bitmap{1200, "general"}, 1) receiver := make(chan *etcd.Response) stop := make(chan bool) diff --git a/db/topology.go b/db/topology.go index 15c61561b..226b48404 100644 --- a/db/topology.go +++ b/db/topology.go @@ -2,13 +2,14 @@ package db import ( "github.com/stathat/consistent" - "github.com/davecgh/go-spew/spew" + //"github.com/davecgh/go-spew/spew" "github.com/nu7hatch/gouuid" "log" "fmt" "errors" "strings" "strconv" + "sync" ) var FrameDoesNotExistError = errors.New("Frame does not exist.") @@ -50,52 +51,15 @@ func (location *Location) ToString() string { // Map of node location to their router 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 { - process *Process - id int -} -// A slice is the vertical combination of every fragment. It contains the hashring used to delegate bitmaps to fragments -type Slice struct { - 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 -} -type FrameSliceIntersect struct { - slice *Slice - frame *Frame - fragments []*Fragment - hashring *consistent.Consistent -} - -func (fsi *FrameSliceIntersect) GetFragment(fragment_id int) (*Fragment, error) { - for _, fragment := range fsi.fragments { - if fragment.id == fragment_id { - return fragment, nil - } - } - return nil, FragmentDoesNotExistError -} - -// Add a slice to a database -func (d *Database) AddSlice(slice_id int) *Slice { - slice := Slice{id: slice_id} - d.slices = append(d.slices, &slice) - // add intersections - for _, frame := range d.frames { - d.AddFrameSliceIntersect(frame, &slice) - } - return &slice -} +/////////// CLUSTERS //////////////////////////////////////////////////////////////////// // Represents the entire cluster, and a reference to the Node this instance is running on type Cluster struct { databases map[string]*Database + mutex sync.Mutex } func NewCluster() *Cluster { @@ -104,8 +68,22 @@ func NewCluster() *Cluster { return &cluster } + +/////////// DATABASES //////////////////////////////////////////////////////////////////// + +// A database is a collection of all the frames within a given profile space +type Database struct { + Name string + frames []*Frame + slices []*Slice + frame_slice_intersects []*FrameSliceIntersect + mutex sync.Mutex +} + // Add a database to a cluster func (c *Cluster) AddDatabase(name string) *Database { + c.mutex.Lock() + defer c.mutex.Unlock() database := Database{Name: name} if c.databases == nil { c.databases = make(map[string]*Database) @@ -115,6 +93,8 @@ func (c *Cluster) AddDatabase(name string) *Database { } func (c *Cluster) GetDatabase(name string) (*Database, error) { + c.mutex.Lock() + defer c.mutex.Unlock() value, ok := c.databases[name] if !ok { return nil, errors.New("The database does not exist!") @@ -123,12 +103,14 @@ func (c *Cluster) GetDatabase(name string) (*Database, error) { } } -// A database is a collection of all the frames within a given profile space -type Database struct { - Name string - frames []*Frame - slices []*Slice - frame_slice_intersects []*FrameSliceIntersect +func (c *Cluster) GetOrCreateDatabase(name string) *Database { + c.mutex.Lock() + defer c.mutex.Unlock() + database, err := c.GetDatabase(name) + if err == nil { + return database + } + return c.AddDatabase(name) } // Count the number of slices in a database @@ -139,8 +121,30 @@ func (d *Database) NumSlices() (int, error) { return len(d.slices), nil } +///////// FRAMES //////////////////////////////////////////////////////////////////// + +// A frame is a collection of slices in a given category +// (brands, demographics, etc), specific to a database +type Frame struct { + name string +} + +// Get a frame from a database +func (d *Database) GetFrame(name string) (*Frame, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + for _, frame := range d.frames { + if frame.name == name { + return frame, nil + } + } + return nil, FrameDoesNotExistError +} + // Add a frame to a database func (d *Database) AddFrame(name string) *Frame { + d.mutex.Lock() + defer d.mutex.Unlock() frame := Frame{name: name} d.frames = append(d.frames, &frame) // add intersections @@ -150,24 +154,74 @@ func (d *Database) AddFrame(name string) *Frame { return &frame } -func (d *Database) AddFragment(frame *Frame, slice *Slice, process *Process, fragment_id int) *Fragment { +func (d *Database) GetOrCreateFrame(name string) *Frame { + d.mutex.Lock() + defer d.mutex.Unlock() + frame, err := d.GetFrame(name) + if err == nil { + return frame + } + return d.AddFrame(name) +} - frameslice, _ := d.GetFrameSliceIntersect(frame, slice) - fragment := Fragment{process: process, id: fragment_id} - frameslice.fragments = append(frameslice.fragments, &fragment) - frameslice.hashring.Add(fmt.Sprintf("%d", fragment_id)) +///////// SLICES ///////////////////////////////////////////////////////////////////////// - return &fragment +// A slice is the vertical combination of every fragment. +type Slice struct { + id int +} + +// Get a slice from a database +func (d *Database) GetSlice(slice_id int) (*Slice, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + for _, slice := range d.slices { + if slice.id == slice_id { + return slice, nil + } + } + return nil, SliceDoesNotExistError +} + +// Add a slice to a database +func (d *Database) AddSlice(slice_id int) *Slice { + d.mutex.Lock() + defer d.mutex.Unlock() + slice := Slice{id: slice_id} + d.slices = append(d.slices, &slice) + // add intersections + for _, frame := range d.frames { + d.AddFrameSliceIntersect(frame, &slice) + } + return &slice +} + +func (d *Database) GetOrCreateSlice(slice_id int) *Slice { + d.mutex.Lock() + defer d.mutex.Unlock() + slice, err := d.GetSlice(slice_id) + if err == nil { + return slice + } + return d.AddSlice(slice_id) +} + + +///////// FRAME-SLICE INTERSECT ////////////////////////////////////////////////////////////// + +type FrameSliceIntersect struct { + frame *Frame + slice *Slice + fragments []*Fragment + hashring *consistent.Consistent } func (d *Database) AddFrameSliceIntersect(frame *Frame, slice *Slice) *FrameSliceIntersect { frameslice := FrameSliceIntersect{frame: frame, slice: slice} d.frame_slice_intersects = append(d.frame_slice_intersects, &frameslice) - frameslice.hashring = consistent.New() frameslice.hashring.NumberOfReplicas = 16 - return &frameslice } @@ -180,30 +234,125 @@ func (d *Database) GetFrameSliceIntersect(frame *Frame, slice *Slice) (*FrameSli return nil, FrameSliceIntersectDoesNotExistError } -// Get a frame from a database -func (d *Database) GetFrame(name string) (*Frame, error) { - for _, frame := range d.frames { - if frame.name == name { - return frame, nil +func (fsi *FrameSliceIntersect) GetFragment(fragment_id *uuid.UUID) (*Fragment, error) { + for _, fragment := range fsi.fragments { + if fragment.id == fragment_id { + return fragment, nil } } - return nil, FrameDoesNotExistError + return nil, FragmentDoesNotExistError } -// 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 +func (fsi *FrameSliceIntersect) AddFragment(fragment *Fragment) { + fsi.fragments = append(fsi.fragments, fragment) + fsi.hashring.Add(fragment.id.String()) } + + +///////// FRAGMENTS //////////////////////////////////////////////////////////////////////// + +// 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 { + id *uuid.UUID + process *Process +} + +// rename this one +func (d *Database) OldGetFragment(bitmap Bitmap, profile_id int) (*Fragment, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + slice, _ := d.GetSliceForProfile(profile_id) + frame, _ := d.GetFrame(bitmap.FrameType) + fsi, err := d.GetFrameSliceIntersect(frame, slice) + frag_id_s, err := fsi.hashring.Get(fmt.Sprintf("%d", bitmap.Id)) + frag_id, err := uuid.ParseHex(frag_id_s) + if err != nil { + log.Fatal(err) + } + return fsi.GetFragment(frag_id) +} + + +/* +// NOT IMPLEMENTED +// this would loop through all frame_slice_intersect[], then all fragmments to find a match +func (d *Database) GetFragment(fragment_id *uuid.UUID) *Fragment { +} +*/ +func (d *Database) GetFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) (*Fragment, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + fsi, err := d.GetFrameSliceIntersect(frame, slice) + if err != nil { + log.Fatal(err) + } + return fsi.GetFragment(fragment_id) +} + +func (d *Database) AddFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) *Fragment { + d.mutex.Lock() + defer d.mutex.Unlock() + fsi, err := d.GetFrameSliceIntersect(frame, slice) + if err != nil { + log.Fatal(err) + } + fragment := Fragment{id: fragment_id} + fsi.AddFragment(&fragment) + return &fragment +} + +/* +func (d *Database) AllocateFragment(frame *Frame, slice *Slice) *Fragment { + // from ETCD, randomly get a process that has available_fragments > 0 + // atomically decrement available_fragments (as long as it's not 0) + // if it IS 0, try until we find a process with available capacity + + * + process, err := GetAvailableProcess() + if err != nil { + log.Fatal(err) + } + * + process_id, _ := uuid.NewV4() + process := NewProcess(process_id) + return nil + //return d.AddFragment(&frame, &slice, process) +} +*/ + +/* +func (d *Database) AddFragmentByProcess(frame *Frame, slice *Slice, process *Process) *Fragment { + d.mutex.Lock() + defer d.mutex.Unlock() + frameslice, _ := d.GetFrameSliceIntersect(frame, slice) + fragment_id, _ := uuid.NewV4() + fragment := Fragment{id: fragment_id, process: process} + frameslice.fragments = append(frameslice.fragments, &fragment) + frameslice.hashring.Add(fragment.id.String()) + return &fragment +} +*/ + +func (d *Database) GetOrCreateFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) *Fragment { + d.mutex.Lock() + defer d.mutex.Unlock() + fragment, err := d.GetFragment(frame, slice, fragment_id) + if err == nil { + return fragment + } + return d.AddFragment(frame, slice, fragment_id) +} + +func (f *Fragment) SetProcess(process *Process) { + f.process = process +} + +/////////////////////////////////////////////////////////////////////////////////////////////// + + // 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) } @@ -213,18 +362,3 @@ 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) - frag_id_s,err := fsi.hashring.Get(fmt.Sprintf("%d", bitmap.Id)) - frag_id, err := strconv.Atoi(frag_id_s) - fragment, err := fsi.GetFragment(frag_id) - if err != nil { - log.Fatal(err) - } - spew.Dump(fragment) -} diff --git a/db/topology_test.go b/db/topology_test.go index aa7489739..f5ac55c00 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -3,6 +3,8 @@ package db import ( "testing" "log" + "github.com/nu7hatch/gouuid" + "github.com/davecgh/go-spew/spew" . "github.com/smartystreets/goconvey/convey" ) @@ -59,6 +61,16 @@ func TestTopology(t *testing.T) { database.AddFragment(frame, slice, loc1, 4) */ + uuid, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") + database.AddFragment(frame, slice, uuid) + /* + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + */ + //fsi, _ := database.GetFrameSliceIntersect(frame, slice) //log.Println(fsi) @@ -74,10 +86,17 @@ func TestTopology(t *testing.T) { log.Println(errer) */ - //bitmap := Bitmap{Id: 555, FrameType: "general"} - //log.Println("bitmap:",bitmap) + bitmap := Bitmap{Id: 555, FrameType: "general"} + log.Println("bitmap:",bitmap) - //database.TestSetBit(bitmap, 65535) + /* + profile_id := 65535 + fragment, _ := database.OldGetFragment(bitmap, profile_id) + + spew.Dump("FRAGMENT") + spew.Dump(fragment) + */ + spew.Dump("DONE") }) } From 7f52ff759ea90fbbca073442511fa9eb094d9746 Mon Sep 17 00:00:00 2001 From: travisturner Date: Tue, 10 Dec 2013 14:54:06 -0600 Subject: [PATCH 09/14] made add* and get* methods private --- core/etcd.go | 4 ++-- db/topology.go | 34 ++++++++++++++-------------------- db/topology_test.go | 34 +++++++++++++++++++--------------- 3 files changed, 35 insertions(+), 37 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 795ce4b1f..76bda70ea 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -106,7 +106,7 @@ func (service *Service) MetaWatcher() { if key == "frame" { for _, frame_ref := range database_attr_ref.Nodes { frame_name := frame_ref.Key[len(database_attr_ref.Key)+1:] - frame := database.AddFrame(frame_name) + frame := database.GetOrCreateFrame(frame_name) for _, frame_attr_ref := range frame_ref.Nodes { key = frame_attr_ref.Key[len(frame_ref.Key)+1:] if key == "slice" { @@ -116,7 +116,7 @@ func (service *Service) MetaWatcher() { if err != nil { log.Fatal(err) } - slice := database.AddSlice(slice_id) + slice := database.GetOrCreateSlice(slice_id) for _, slice_attr_ref := range slice_ref.Nodes { key = slice_attr_ref.Key[len(slice_ref.Key)+1:] if key == "fragment" { diff --git a/db/topology.go b/db/topology.go index 226b48404..6c8856bc3 100644 --- a/db/topology.go +++ b/db/topology.go @@ -2,7 +2,7 @@ package db import ( "github.com/stathat/consistent" - //"github.com/davecgh/go-spew/spew" + "github.com/davecgh/go-spew/spew" "github.com/nu7hatch/gouuid" "log" "fmt" @@ -130,9 +130,7 @@ type Frame struct { } // Get a frame from a database -func (d *Database) GetFrame(name string) (*Frame, error) { - d.mutex.Lock() - defer d.mutex.Unlock() +func (d *Database) getFrame(name string) (*Frame, error) { for _, frame := range d.frames { if frame.name == name { return frame, nil @@ -142,9 +140,7 @@ func (d *Database) GetFrame(name string) (*Frame, error) { } // Add a frame to a database -func (d *Database) AddFrame(name string) *Frame { - d.mutex.Lock() - defer d.mutex.Unlock() +func (d *Database) addFrame(name string) *Frame { frame := Frame{name: name} d.frames = append(d.frames, &frame) // add intersections @@ -157,11 +153,11 @@ func (d *Database) AddFrame(name string) *Frame { func (d *Database) GetOrCreateFrame(name string) *Frame { d.mutex.Lock() defer d.mutex.Unlock() - frame, err := d.GetFrame(name) + frame, err := d.getFrame(name) if err == nil { return frame } - return d.AddFrame(name) + return d.addFrame(name) } @@ -173,9 +169,7 @@ type Slice struct { } // Get a slice from a database -func (d *Database) GetSlice(slice_id int) (*Slice, error) { - d.mutex.Lock() - defer d.mutex.Unlock() +func (d *Database) getSlice(slice_id int) (*Slice, error) { for _, slice := range d.slices { if slice.id == slice_id { return slice, nil @@ -185,9 +179,7 @@ func (d *Database) GetSlice(slice_id int) (*Slice, error) { } // Add a slice to a database -func (d *Database) AddSlice(slice_id int) *Slice { - d.mutex.Lock() - defer d.mutex.Unlock() +func (d *Database) addSlice(slice_id int) *Slice { slice := Slice{id: slice_id} d.slices = append(d.slices, &slice) // add intersections @@ -200,11 +192,11 @@ func (d *Database) AddSlice(slice_id int) *Slice { func (d *Database) GetOrCreateSlice(slice_id int) *Slice { d.mutex.Lock() defer d.mutex.Unlock() - slice, err := d.GetSlice(slice_id) + slice, err := d.getSlice(slice_id) if err == nil { return slice } - return d.AddSlice(slice_id) + return d.addSlice(slice_id) } @@ -246,6 +238,8 @@ func (fsi *FrameSliceIntersect) GetFragment(fragment_id *uuid.UUID) (*Fragment, func (fsi *FrameSliceIntersect) AddFragment(fragment *Fragment) { fsi.fragments = append(fsi.fragments, fragment) fsi.hashring.Add(fragment.id.String()) + spew.Dump("DUMPY") + spew.Dump(fragment.id.String()) } @@ -263,7 +257,7 @@ func (d *Database) OldGetFragment(bitmap Bitmap, profile_id int) (*Fragment, err d.mutex.Lock() defer d.mutex.Unlock() slice, _ := d.GetSliceForProfile(profile_id) - frame, _ := d.GetFrame(bitmap.FrameType) + frame, _ := d.getFrame(bitmap.FrameType) fsi, err := d.GetFrameSliceIntersect(frame, slice) frag_id_s, err := fsi.hashring.Get(fmt.Sprintf("%d", bitmap.Id)) frag_id, err := uuid.ParseHex(frag_id_s) @@ -277,7 +271,7 @@ func (d *Database) OldGetFragment(bitmap Bitmap, profile_id int) (*Fragment, err /* // NOT IMPLEMENTED // this would loop through all frame_slice_intersect[], then all fragmments to find a match -func (d *Database) GetFragment(fragment_id *uuid.UUID) *Fragment { +func (d *Database) GetFragmentById(fragment_id *uuid.UUID) *Fragment { } */ func (d *Database) GetFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) (*Fragment, error) { @@ -354,7 +348,7 @@ func (f *Fragment) SetProcess(process *Process) { // Get a slice from a database func (d *Database) GetSliceForProfile(profile_id int) (*Slice, error) { slice_id := profile_id / SLICE_WIDTH - return d.GetSlice(slice_id) + return d.getSlice(slice_id) } diff --git a/db/topology_test.go b/db/topology_test.go index f5ac55c00..ed4c13ec9 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -24,13 +24,18 @@ func TestTopology(t *testing.T) { */ cluster := NewCluster() - database := cluster.AddDatabase("property49") - database.AddFrame("general") - //database.AddFrame("brands") - database.AddSlice(0) - database.AddSlice(1) - //database.AddSlice() + database := cluster.AddDatabase("main") + frame := database.GetOrCreateFrame("general") + slice := database.GetOrCreateSlice(0) + + fragment_id, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") + spew.Dump(fragment_id) + database.AddFragment(frame, slice, fragment_id) + + + spew.Dump(database) + /* log.Println(database) log.Println("----------------------------------") for _, fsi := range database.frame_slice_intersects { @@ -38,12 +43,10 @@ func TestTopology(t *testing.T) { } num_slices, _ := database.NumSlices() log.Println(num_slices) + */ - - frame, _ := database.GetFrame("general") - slice, _ := database.GetSlice(0) - log.Println(frame) - log.Println(slice) + //frame, _ := database.GetFrame("general") + //slice, _ := database.GetSlice(0) //loc1, _ := NewLocation("192.168.1.100:8001") /* @@ -61,8 +64,9 @@ func TestTopology(t *testing.T) { database.AddFragment(frame, slice, loc1, 4) */ - uuid, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") - database.AddFragment(frame, slice, uuid) + //uuid, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") + //spew.Dump(uuid) + ////database.AddFragment(frame, slice, uuid) /* database.AddFragment(frame, slice, process) database.AddFragment(frame, slice, process) @@ -86,8 +90,8 @@ func TestTopology(t *testing.T) { log.Println(errer) */ - bitmap := Bitmap{Id: 555, FrameType: "general"} - log.Println("bitmap:",bitmap) + ////bitmap := Bitmap{Id: 555, FrameType: "general"} + ////log.Println("bitmap:",bitmap) /* profile_id := 65535 From ec676186b37edb75bfa9c02f8aaabe1a83646bd7 Mon Sep 17 00:00:00 2001 From: travisturner Date: Tue, 10 Dec 2013 15:01:58 -0600 Subject: [PATCH 10/14] made add* and get* methods private (part 2) --- core/etcd.go | 6 +++--- db/topology.go | 24 ++++++++---------------- db/topology_test.go | 4 ++-- 3 files changed, 13 insertions(+), 21 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 76bda70ea..287d140f7 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -100,7 +100,7 @@ func (service *Service) MetaWatcher() { for _, database_ref := range resp.Node.Nodes { database_name := database_ref.Key[len(namespace)+4:] - database := cluster.AddDatabase(database_name) + database := cluster.GetOrCreateDatabase(database_name) for _, database_attr_ref := range database_ref.Nodes { key := database_attr_ref.Key[len(database_ref.Key)+1:] if key == "frame" { @@ -136,7 +136,7 @@ func (service *Service) MetaWatcher() { log.Fatal(err) } process := uuid - database.AddFragment(frame, slice, process) + database.GetOrCreateFragment(frame, slice, process) } } } @@ -149,7 +149,7 @@ func (service *Service) MetaWatcher() { } } } - database, _ := cluster.GetDatabase("main") + database := cluster.GetOrCreateDatabase("main") spew.Dump(database) database.OldGetFragment(db.Bitmap{1200, "general"}, 1) diff --git a/db/topology.go b/db/topology.go index 6c8856bc3..38de0ab1d 100644 --- a/db/topology.go +++ b/db/topology.go @@ -81,9 +81,7 @@ type Database struct { } // Add a database to a cluster -func (c *Cluster) AddDatabase(name string) *Database { - c.mutex.Lock() - defer c.mutex.Unlock() +func (c *Cluster) addDatabase(name string) *Database { database := Database{Name: name} if c.databases == nil { c.databases = make(map[string]*Database) @@ -92,9 +90,7 @@ func (c *Cluster) AddDatabase(name string) *Database { return &database } -func (c *Cluster) GetDatabase(name string) (*Database, error) { - c.mutex.Lock() - defer c.mutex.Unlock() +func (c *Cluster) getDatabase(name string) (*Database, error) { value, ok := c.databases[name] if !ok { return nil, errors.New("The database does not exist!") @@ -106,11 +102,11 @@ func (c *Cluster) GetDatabase(name string) (*Database, error) { func (c *Cluster) GetOrCreateDatabase(name string) *Database { c.mutex.Lock() defer c.mutex.Unlock() - database, err := c.GetDatabase(name) + database, err := c.getDatabase(name) if err == nil { return database } - return c.AddDatabase(name) + return c.addDatabase(name) } // Count the number of slices in a database @@ -274,9 +270,7 @@ func (d *Database) OldGetFragment(bitmap Bitmap, profile_id int) (*Fragment, err func (d *Database) GetFragmentById(fragment_id *uuid.UUID) *Fragment { } */ -func (d *Database) GetFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) (*Fragment, error) { - d.mutex.Lock() - defer d.mutex.Unlock() +func (d *Database) getFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) (*Fragment, error) { fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { log.Fatal(err) @@ -284,9 +278,7 @@ func (d *Database) GetFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUI return fsi.GetFragment(fragment_id) } -func (d *Database) AddFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) *Fragment { - d.mutex.Lock() - defer d.mutex.Unlock() +func (d *Database) addFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) *Fragment { fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { log.Fatal(err) @@ -331,11 +323,11 @@ func (d *Database) AddFragmentByProcess(frame *Frame, slice *Slice, process *Pro func (d *Database) GetOrCreateFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) *Fragment { d.mutex.Lock() defer d.mutex.Unlock() - fragment, err := d.GetFragment(frame, slice, fragment_id) + fragment, err := d.getFragment(frame, slice, fragment_id) if err == nil { return fragment } - return d.AddFragment(frame, slice, fragment_id) + return d.addFragment(frame, slice, fragment_id) } func (f *Fragment) SetProcess(process *Process) { diff --git a/db/topology_test.go b/db/topology_test.go index ed4c13ec9..532e3c61b 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -24,14 +24,14 @@ func TestTopology(t *testing.T) { */ cluster := NewCluster() - database := cluster.AddDatabase("main") + database := cluster.GetOrCreateDatabase("main") frame := database.GetOrCreateFrame("general") slice := database.GetOrCreateSlice(0) fragment_id, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") spew.Dump(fragment_id) - database.AddFragment(frame, slice, fragment_id) + database.GetOrCreateFragment(frame, slice, fragment_id) spew.Dump(database) From cd6dc19ea162d196b7b6485ee7c4b84c4c8cb7b9 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Tue, 10 Dec 2013 18:10:22 -0600 Subject: [PATCH 11/14] New less ugly etcd stuff! --- core/etcd.go | 211 ++++++++++++++++++++------------------------------- 1 file changed, 82 insertions(+), 129 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 287d140f7..143b09083 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -1,158 +1,111 @@ package core import ( + "github.com/davecgh/go-spew/spew" "github.com/coreos/go-etcd/etcd" "encoding/gob" "pilosa/db" - "github.com/davecgh/go-spew/spew" + "log" + "strings" "github.com/nu7hatch/gouuid" "strconv" - "log" + "errors" ) func (service *Service) SetupEtcd() { gob.Register(db.Location{}) service.Etcd = etcd.NewClient(nil) - //service.NodeMapMutex.Lock() - //defer service.NodeMapMutex.Unlock() - //service.NodeMap = db.NodeMap{} - - //nodes, err := service.Etcd.Get("nodes", false) - //if err != nil { - // log.Fatal(err) - //} - //for _, node := range nodes.Kvs { - // nodestring := strings.Split(node.Key, "/")[2] - // location, err := db.NewLocation(nodestring) - // if err != nil { - // log.Fatal(err) - // } - // routerlocation, err := db.NewLocation(node.Value) - // if err != nil { - // log.Fatal(err) - // } - // service.NodeMap[*location] = *routerlocation - //} - //log.Println(service.NodeMap) } -//func (service *Service) WatchEtcd() { -// var receiver = make(chan *etcd.Response) -// var stop chan bool -// go func () { -// _, err := service.Etcd.Watch("nodes/", 0, receiver, stop) -// if err != nil { -// log.Fatal(err) -// } -// }() -// -// exit, done := service.GetExitChannels() -// -// for { -// select { -// case response := <-receiver: -// switch response.Action { -// case "SET": -// nodestring := strings.Split(response.Key, "/")[2] -// node, err := db.NewLocation(nodestring) -// if err != nil { -// log.Fatal(err) -// } -// router, err := db.NewLocation(response.Value) -// if err != nil { -// log.Fatal(err) -// } -// service.NodeMapMutex.Lock() -// service.NodeMap[*node] = *router -// service.NodeMapMutex.Unlock() -// case "DELETE": -// nodestring := strings.Split(response.Key, "/")[2] -// node, err := db.NewLocation(nodestring) -// if err != nil { -// log.Fatal(err) -// } -// service.NodeMapMutex.Lock() -// delete(service.NodeMap, *node) -// service.NodeMapMutex.Unlock() -// default: -// log.Println("unhandled etcd message", response) -// } -// //log.Println(response.Action, response.Key, response.Value) -// log.Println(service.NodeMap) -// case <-exit: -// log.Println("cleaning up watchetcd service thing.") -// time.Sleep(time.Second/2) -// log.Println("done!") -// done <- 1 -// } -// } -//} +func flatten(node *etcd.Node) []*etcd.Node { + nodes := make([]*etcd.Node, 0) + nodes = append(nodes, node) + for _, node := range node.Nodes { + nodes = append(nodes, flatten(&node)...) + } + return nodes +} +func handlenode(node *etcd.Node, namespace string, cluster *db.Cluster) error { + key := node.Key[len(namespace)+1:] + bits := strings.Split(key, "/") + var database *db.Database + var frame *db.Frame + var fragment *db.Fragment + var fragment_uuid *uuid.UUID + var slice *db.Slice + var process_uuid *uuid.UUID + var process *db.Process + var err error + + if len(bits) <= 1 || bits[0] != "db" { + return nil + } + if len(bits) > 1 { + database = cluster.GetOrCreateDatabase(bits[1]) + } + if len(bits) > 2 { + if bits[2] != "frame" { + return errors.New("no frame") + } + } + if len(bits) > 3 { + frame = database.GetOrCreateFrame(bits[3]) + } + if len(bits) > 4 { + if bits[4] != "slice" { + return errors.New("no slice") + } + } + if len(bits) > 5 { + slice_int, err := strconv.Atoi(bits[5]) + if err != nil { + return err + } + slice = database.GetOrCreateSlice(slice_int) + } + if len(bits) > 6 { + if bits[6] != "fragment" { + return errors.New("no fragment") + } + } + if len(bits) > 7 { + fragment_uuid, err = uuid.ParseHex(bits[7]) + if err != nil { + return err + } + fragment = database.GetOrCreateFragment(frame, slice, fragment_uuid) + } + + if len(bits) > 8 { + if bits[8] != "process" { + return errors.New("no process") + } + process_uuid, err = uuid.ParseHex(node.Value) + if err != nil { + return err + } + process = db.NewProcess(process_uuid) + fragment.SetProcess(process) + } + return err +} func (service *Service) MetaWatcher() { namespace := "/pilosa/0" log.Println(namespace + "/db") + cluster := db.NewCluster() resp, err := service.Etcd.Get(namespace + "/db", false, true) if err != nil { log.Fatal(err) } - cluster := db.NewCluster() - - for _, database_ref := range resp.Node.Nodes { - database_name := database_ref.Key[len(namespace)+4:] - database := cluster.GetOrCreateDatabase(database_name) - for _, database_attr_ref := range database_ref.Nodes { - key := database_attr_ref.Key[len(database_ref.Key)+1:] - if key == "frame" { - for _, frame_ref := range database_attr_ref.Nodes { - frame_name := frame_ref.Key[len(database_attr_ref.Key)+1:] - frame := database.GetOrCreateFrame(frame_name) - for _, frame_attr_ref := range frame_ref.Nodes { - key = frame_attr_ref.Key[len(frame_ref.Key)+1:] - if key == "slice" { - for _, slice_ref := range frame_attr_ref.Nodes { - slice_name := slice_ref.Key[len(frame_attr_ref.Key)+1:] - slice_id, err := strconv.Atoi(slice_name) - if err != nil { - log.Fatal(err) - } - slice := database.GetOrCreateSlice(slice_id) - for _, slice_attr_ref := range slice_ref.Nodes { - key = slice_attr_ref.Key[len(slice_ref.Key)+1:] - if key == "fragment" { - for _, fragment_ref := range slice_attr_ref.Nodes { - /* - fragment_name := fragment_ref.Key[len(slice_attr_ref.Key)+1:] - fragment_id, err := strconv.Atoi(fragment_name) - if err != nil { - log.Fatal(err) - } - */ - for _, fragment_attr_ref := range fragment_ref.Nodes { - key = fragment_attr_ref.Key[len(fragment_ref.Key)+1:] - if key == "node" { - uuid, err := uuid.ParseHex(fragment_attr_ref.Value) - if err != nil { - log.Fatal(err) - } - process := uuid - database.GetOrCreateFragment(frame, slice, process) - } - } - } - } - } - } - } - } - } - } + for _, node := range flatten(resp.Node) { + err := handlenode(node, namespace, cluster) + if err != nil { + spew.Dump(node) + log.Println(err) } } - database := cluster.GetOrCreateDatabase("main") - spew.Dump(database) - database.OldGetFragment(db.Bitmap{1200, "general"}, 1) - receiver := make(chan *etcd.Response) stop := make(chan bool) go func() { From 04aadebaa52f89d80d2b270e64bd7d383b3497b1 Mon Sep 17 00:00:00 2001 From: travisturner Date: Wed, 11 Dec 2013 10:13:51 -0600 Subject: [PATCH 12/14] stub out Cruncher.Run --- core/cruncher.go | 13 +++++++++++++ core/cruncher_test.go | 30 ++++++++++++++++++++++++++++++ core/service.go | 3 +++ 3 files changed, 46 insertions(+) create mode 100644 core/cruncher.go create mode 100644 core/cruncher_test.go diff --git a/core/cruncher.go b/core/cruncher.go new file mode 100644 index 000000000..1e2af8684 --- /dev/null +++ b/core/cruncher.go @@ -0,0 +1,13 @@ +package core + +import ( + "github.com/davecgh/go-spew/spew" +) + + +type Cruncher struct { +} + +func (cruncher *Cruncher) Run() { + spew.Dump("Cruncher.Run") +} diff --git a/core/cruncher_test.go b/core/cruncher_test.go new file mode 100644 index 000000000..9d0b42df3 --- /dev/null +++ b/core/cruncher_test.go @@ -0,0 +1,30 @@ +package core + +import ( + "testing" + //"github.com/nu7hatch/gouuid" + "github.com/davecgh/go-spew/spew" + . "github.com/smartystreets/goconvey/convey" +) + +func TestCruncher(t *testing.T) { + Convey("Basic Cruncher Tests", t, func() { + spew.Dump("cruncher test") + + /* + cluster := NewCluster() + database := cluster.GetOrCreateDatabase("main") + + frame := database.GetOrCreateFrame("general") + slice := database.GetOrCreateSlice(0) + + fragment_id, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") + spew.Dump(fragment_id) + database.GetOrCreateFragment(frame, slice, fragment_id) + + spew.Dump(database) + spew.Dump("DONE") + */ + + }) +} diff --git a/core/service.go b/core/service.go index 871022d8b..90b62b3de 100644 --- a/core/service.go +++ b/core/service.go @@ -50,6 +50,7 @@ type Service struct { ConnectionRegisterChannel chan *PersistentConnection Stats *Stats //Cluster query.Cluster + Cruncher *Cruncher } func NewService(tcp, http *db.Location) *Service { @@ -59,6 +60,7 @@ func NewService(tcp, http *db.Location) *Service { service.Outbox = make(chan *db.Envelope) service.Inbox = make(chan *db.Message) service.Stats = new(Stats) + service.Cruncher = new(Cruncher) return service } @@ -306,6 +308,7 @@ func (service *Service) Run() { //go service.HandleInbox() //go service.ServeHTTP() go service.MetaWatcher() + go service.Cruncher.Run() sigterm, sighup := service.GetSignals() for { From 7d31f487322cd17af4aaebebf8f849c7a0367a34 Mon Sep 17 00:00:00 2001 From: travisturner Date: Wed, 11 Dec 2013 10:14:29 -0600 Subject: [PATCH 13/14] remove some debugging lines --- db/topology.go | 4 +-- db/topology_test.go | 77 --------------------------------------------- 2 files changed, 1 insertion(+), 80 deletions(-) diff --git a/db/topology.go b/db/topology.go index 38de0ab1d..b4332bc7d 100644 --- a/db/topology.go +++ b/db/topology.go @@ -2,7 +2,7 @@ package db import ( "github.com/stathat/consistent" - "github.com/davecgh/go-spew/spew" + //"github.com/davecgh/go-spew/spew" "github.com/nu7hatch/gouuid" "log" "fmt" @@ -234,8 +234,6 @@ func (fsi *FrameSliceIntersect) GetFragment(fragment_id *uuid.UUID) (*Fragment, func (fsi *FrameSliceIntersect) AddFragment(fragment *Fragment) { fsi.fragments = append(fsi.fragments, fragment) fsi.hashring.Add(fragment.id.String()) - spew.Dump("DUMPY") - spew.Dump(fragment.id.String()) } diff --git a/db/topology_test.go b/db/topology_test.go index 532e3c61b..0666ef0e7 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -12,17 +12,6 @@ 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 := NewCluster() database := cluster.GetOrCreateDatabase("main") @@ -33,73 +22,7 @@ func TestTopology(t *testing.T) { spew.Dump(fragment_id) database.GetOrCreateFragment(frame, slice, fragment_id) - spew.Dump(database) - /* - log.Println(database) - log.Println("----------------------------------") - for _, fsi := range database.frame_slice_intersects { - 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) - */ - - //uuid, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") - //spew.Dump(uuid) - ////database.AddFragment(frame, slice, uuid) - /* - database.AddFragment(frame, slice, process) - database.AddFragment(frame, slice, process) - database.AddFragment(frame, slice, process) - database.AddFragment(frame, slice, process) - database.AddFragment(frame, slice, process) - */ - - - //fsi, _ := database.GetFrameSliceIntersect(frame, slice) - //log.Println(fsi) - //log.Println(fsi.Hashring) - - /* - 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) - - /* - profile_id := 65535 - fragment, _ := database.OldGetFragment(bitmap, profile_id) - - spew.Dump("FRAGMENT") - spew.Dump(fragment) - */ spew.Dump("DONE") }) From f7b5550315b491f609897a6733df5b714d7e6752 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 11 Dec 2013 10:45:37 -0600 Subject: [PATCH 14/14] Update local metadata from etcd watcher. --- core/etcd.go | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 143b09083..054ffe96a 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -112,8 +112,11 @@ func (service *Service) MetaWatcher() { _, _ = service.Etcd.Watch(namespace + "/db", 0, true, receiver, stop) }() go func() { - for x := range receiver { - spew.Dump(x) + for resp = range receiver { + switch resp.Action { + case "set": + handlenode(resp.Node, namespace, cluster) + } } }() }