mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
merge
This commit is contained in:
commit
b7d9f91d2c
2 changed files with 49 additions and 43 deletions
83
core/etcd.go
83
core/etcd.go
|
|
@ -3,7 +3,6 @@ package core
|
|||
import (
|
||||
"github.com/davecgh/go-spew/spew"
|
||||
"github.com/coreos/go-etcd/etcd"
|
||||
"encoding/gob"
|
||||
"pilosa/db"
|
||||
"log"
|
||||
"strings"
|
||||
|
|
@ -12,22 +11,47 @@ import (
|
|||
"errors"
|
||||
)
|
||||
|
||||
func (service *Service) SetupEtcd() {
|
||||
gob.Register(db.Location{})
|
||||
service.Etcd = etcd.NewClient(nil)
|
||||
type MetaWatcher struct {
|
||||
service *Service
|
||||
namespace string
|
||||
}
|
||||
|
||||
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)...)
|
||||
func (self *MetaWatcher) Run() {
|
||||
log.Println(self.namespace + "/db")
|
||||
resp, err := self.service.Etcd.Get(self.namespace + "/db", false, true)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
return nodes
|
||||
for _, node := range flatten(resp.Node) {
|
||||
err := self.handlenode(node)
|
||||
if err != nil {
|
||||
spew.Dump(node)
|
||||
log.Println(err)
|
||||
}
|
||||
}
|
||||
receiver := make(chan *etcd.Response)
|
||||
stop := make(chan bool)
|
||||
go func() {
|
||||
// TODO: error check and restart watcher
|
||||
_, _ = self.service.Etcd.Watch(self.namespace + "/db", 0, true, receiver, stop)
|
||||
}()
|
||||
go func() {
|
||||
for resp = range receiver {
|
||||
switch resp.Action {
|
||||
case "set":
|
||||
self.handlenode(resp.Node)
|
||||
}
|
||||
// TODO: handle deletes
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func handlenode(node *etcd.Node, namespace string, cluster *db.Cluster) error {
|
||||
key := node.Key[len(namespace)+1:]
|
||||
func NewMetaWatcher(service *Service, namespace string) *MetaWatcher {
|
||||
return &MetaWatcher{service, namespace}
|
||||
}
|
||||
|
||||
func (self *MetaWatcher) handlenode(node *etcd.Node) error {
|
||||
key := node.Key[len(self.namespace)+1:]
|
||||
bits := strings.Split(key, "/")
|
||||
var database *db.Database
|
||||
var frame *db.Frame
|
||||
|
|
@ -42,7 +66,7 @@ func handlenode(node *etcd.Node, namespace string, cluster *db.Cluster) error {
|
|||
return nil
|
||||
}
|
||||
if len(bits) > 1 {
|
||||
database = cluster.GetOrCreateDatabase(bits[1])
|
||||
database = self.service.Cluster.GetOrCreateDatabase(bits[1])
|
||||
}
|
||||
if len(bits) > 2 {
|
||||
if bits[2] != "frame" {
|
||||
|
|
@ -91,32 +115,11 @@ func handlenode(node *etcd.Node, namespace string, cluster *db.Cluster) error {
|
|||
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)
|
||||
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)...)
|
||||
}
|
||||
for _, node := range flatten(resp.Node) {
|
||||
err := handlenode(node, namespace, cluster)
|
||||
if err != nil {
|
||||
spew.Dump(node)
|
||||
log.Println(err)
|
||||
}
|
||||
}
|
||||
receiver := make(chan *etcd.Response)
|
||||
stop := make(chan bool)
|
||||
go func() {
|
||||
_, _ = service.Etcd.Watch(namespace + "/db", 0, true, receiver, stop)
|
||||
}()
|
||||
go func() {
|
||||
for resp = range receiver {
|
||||
switch resp.Action {
|
||||
case "set":
|
||||
handlenode(resp.Node, namespace, cluster)
|
||||
}
|
||||
}
|
||||
}()
|
||||
return nodes
|
||||
}
|
||||
|
|
|
|||
|
|
@ -50,8 +50,9 @@ type Service struct {
|
|||
Listener net.Listener
|
||||
ConnectionRegisterChannel chan *PersistentConnection
|
||||
Stats *Stats
|
||||
//Cluster query.Cluster
|
||||
Cluster *db.Cluster
|
||||
Cruncher *Cruncher
|
||||
MetaWatcher *MetaWatcher
|
||||
}
|
||||
|
||||
func NewService(tcp, http *db.Location) *Service {
|
||||
|
|
@ -62,6 +63,9 @@ func NewService(tcp, http *db.Location) *Service {
|
|||
service.Inbox = make(chan *db.Message)
|
||||
service.Stats = new(Stats)
|
||||
service.Cruncher = new(Cruncher)
|
||||
service.Etcd = etcd.NewClient(nil)
|
||||
service.Cluster = db.NewCluster()
|
||||
service.MetaWatcher = &MetaWatcher{service, "/pilosa/0"}
|
||||
return service
|
||||
}
|
||||
|
||||
|
|
@ -300,7 +304,6 @@ func (service *Service) NewListener() chan *db.Message {
|
|||
|
||||
func (service *Service) Run() {
|
||||
log.Println("Running service...")
|
||||
service.SetupEtcd()
|
||||
//go r.SyncEtcd()
|
||||
//go service.WatchEtcd()
|
||||
//go service.HandleConnections()
|
||||
|
|
@ -308,7 +311,7 @@ func (service *Service) Run() {
|
|||
//go service.Serve()
|
||||
//go service.HandleInbox()
|
||||
//go service.ServeHTTP()
|
||||
go service.MetaWatcher()
|
||||
go service.MetaWatcher.Run()
|
||||
go service.Cruncher.Run(config.GetInt("port_tcp"))
|
||||
|
||||
sigterm, sighup := service.GetSignals()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue