From e72d147c3d4bcd60238ef902d6e609c1493ceef3 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 6 Dec 2013 15:54:33 -0600 Subject: [PATCH] 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 -}