mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-10 15:01:03 +00:00
Add beginnings of etcd topology sync.
This commit is contained in:
parent
63f3311602
commit
e72d147c3d
7 changed files with 183 additions and 179 deletions
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
213
core/etcd.go
213
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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -6,7 +6,7 @@
|
|||
},
|
||||
"etcd": {
|
||||
"repo": "github.com/coreos/go-etcd/etcd",
|
||||
"version": "8a4461a676eb65fb74f10da1f8198cc9f67da366",
|
||||
"version": "8a4461a",
|
||||
"type": "git"
|
||||
},
|
||||
"goconvey": {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue