From 4bef6c3d46309c2c5c379cbacdfca6268f0b448b Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 11 Dec 2013 11:26:55 -0600 Subject: [PATCH] move things around, use new MetaWatcher struct --- core/etcd.go | 83 +++++++++++++++++++++++++------------------------ core/service.go | 9 ++++-- 2 files changed, 49 insertions(+), 43 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 054ffe96a..d4d6c8491 100644 --- a/core/etcd.go +++ b/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 } diff --git a/core/service.go b/core/service.go index 0c7c5cc10..ae43e9bce 100644 --- a/core/service.go +++ b/core/service.go @@ -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()