From 2df5b961f541943a01df4c97d8ac4a3453299c1a Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Thu, 19 Dec 2013 16:33:25 -0600 Subject: [PATCH] Updated ProcessMapper, flatten() fix, etc --- core/etcd.go | 234 +++++++++++++++++++++++++++++++----------------- core/service.go | 5 +- db/topology.go | 58 +++++++++++- 3 files changed, 209 insertions(+), 88 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index d600e94e7..3ac5615c3 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -2,13 +2,16 @@ package core import ( "errors" - "github.com/coreos/go-etcd/etcd" - "github.com/davecgh/go-spew/spew" - "github.com/nu7hatch/gouuid" "log" + "pilosa/config" "pilosa/db" "strconv" "strings" + "sync" + + "github.com/coreos/go-etcd/etcd" + "github.com/davecgh/go-spew/spew" + "github.com/nu7hatch/gouuid" ) type TopologyMapper struct { @@ -117,10 +120,9 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { } 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)...) + nodes := []*etcd.Node{node} + for i := 0; i < len(node.Nodes); i++ { + nodes = append(nodes, flatten(&node.Nodes[i])...) } return nodes } @@ -133,18 +135,85 @@ type Node struct { } type ProcessMap struct { - nodes []*db.Process + nodes map[uuid.UUID]*db.Process + mutex sync.Mutex +} + +func NewProcessMap() *ProcessMap { + p := ProcessMap{} + p.nodes = make(map[uuid.UUID]*db.Process) + return &p +} + +func (self *ProcessMap) AddProcess(process *db.Process) { + self.mutex.Lock() + defer self.mutex.Unlock() + self.nodes[process.Id()] = process +} + +func (self *ProcessMap) GetProcess(id *uuid.UUID) (*db.Process, error) { + self.mutex.Lock() + defer self.mutex.Unlock() + process, ok := self.nodes[*id] + if !ok { + return nil, errors.New("No such process") + } + return process, nil +} + +func (self *ProcessMap) GetOrAddProcess(id *uuid.UUID) *db.Process { + process, err := self.GetProcess(id) + if err != nil { + process = db.NewProcess(id) + self.AddProcess(process) + } + return process +} + +func (self *ProcessMap) GetHost(id *uuid.UUID) (string, error) { + self.mutex.Lock() + defer self.mutex.Unlock() + process, ok := self.nodes[*id] + if !ok { + return "", errors.New("Process does not exist") + } + return process.Host(), nil +} + +func (self *ProcessMap) GetPortTcp(id *uuid.UUID) (int, error) { + self.mutex.Lock() + defer self.mutex.Unlock() + process, ok := self.nodes[*id] + if !ok { + return 0, errors.New("Process does not exist") + } + return process.PortTcp(), nil +} + +func (self *ProcessMap) GetPortHttp(id *uuid.UUID) (int, error) { + self.mutex.Lock() + defer self.mutex.Unlock() + process, ok := self.nodes[*id] + if !ok { + return 0, errors.New("Process does not exist") + } + return process.PortHttp(), nil } type ProcessMapper struct { - etcd *etcd.Client - nodes []Node - receiver chan *etcd.Response - commands chan *ProcessMapperCommand + service *Service + receiver chan etcd.Response + commands chan ProcessMapperCommand + namespace string } -func NewProcessMapper(service *Service) *ProcessMapper { - return &ProcessMapper{} +func NewProcessMapper(service *Service, namespace string) *ProcessMapper { + return &ProcessMapper{ + service: service, + receiver: make(chan etcd.Response), + commands: make(chan ProcessMapperCommand), + namespace: namespace, + } } type ProcessMapperCommand struct { @@ -160,85 +229,84 @@ func (self *ProcessMapper) getnode(u *uuid.UUID) *Node { return new(Node) } -func (self *ProcessMapper) Run() { - return - var modindex uint64 - response, err := self.etcd.Get("nodes", false, true) +func crash_on_error(err error) { if err != nil { log.Fatal(err) } - //modindex = response.ModifiedIndex - log.Println(modindex) - nodes := make([]Node, 0) - spew.Dump(response) - for _, noderef := range response.Node.Nodes { - nodestring := getKey(noderef.Key) - u, err := uuid.ParseHex(nodestring) +} + +func (self *ProcessMapper) handlenode(node *etcd.Node) error { + var err error + var process *db.Process + key := node.Key[len(self.namespace)+1:] + bits := strings.Split(key, "/") + if len(bits) <= 1 || bits[0] != "process" { + return nil + } + if len(bits) >= 2 { + id_string := bits[1] + id, err := uuid.ParseHex(id_string) if err != nil { - log.Fatal("Not a valid UUID: ", nodestring) + return errors.New("Invalid UUID: " + id_string) } - node := Node{id: u} - for _, prop := range noderef.Nodes { - switch getKey(prop.Key) { - case "port_tcp": - node.port_tcp, _ = strconv.Atoi(prop.Value) - case "port_http": - node.port_http, _ = strconv.Atoi(prop.Value) - case "ip": - node.ip = prop.Value - } + process = self.service.ProcessMap.GetOrAddProcess(id) + } + if len(bits) >= 3 { + switch bits[2] { + case "port_tcp": + port_tcp, _ := strconv.Atoi(node.Value) + process.SetPortTcp(port_tcp) + case "port_http": + port_http, _ := strconv.Atoi(node.Value) + process.SetPortHttp(port_http) + case "host": + host := node.Value + process.SetHost(host) } - nodes = append(nodes, node) } - self.nodes = nodes - spew.Dump(self.nodes) - go func() { - stop := make(chan bool) - _, err := self.etcd.Watch("nodes/", 0, true, self.receiver, stop) - if err != nil { - log.Fatal(err) - } - }() - go func() { - for { - select { - case cmd := <-self.commands: - spew.Dump(cmd) - case response := <-self.receiver: - switch response.Action { - case "set": - bits := strings.Split(response.Node.Key, "/") - if len(bits) != 4 { - log.Fatal("bug in etcd sync or etcd data") - } - //router, err := db.NewLocation(response.Node.Value) - u, err := uuid.ParseHex(bits[2]) - node := self.getnode(u) - if err != nil { - log.Fatal(err) - } - switch bits[3] { - case "port_tcp": - node.port_tcp, _ = strconv.Atoi(response.Node.Value) - case "port_http": - node.port_http, _ = strconv.Atoi(response.Node.Value) - case "ip": - node.ip = response.Node.Value - } - spew.Dump(node) - case "delete": - spew.Dump("delete", response) + return err +} - default: - spew.Dump("unhandled", response) - } +func (self *ProcessMapper) Run() { + id_string := self.service.Id.String() + path := self.namespace + "/process" + self_path := path + "/" + id_string + + log.Println("Writing configuration to etcd...") + + var err error + _, err = self.service.Etcd.Set(self_path+"/port_tcp", strconv.Itoa(config.GetInt("port_tcp")), 0) + crash_on_error(err) + _, err = self.service.Etcd.Set(self_path+"/port_http", strconv.Itoa(config.GetInt("port_http")), 0) + crash_on_error(err) + _, err = self.service.Etcd.Set(self_path+"/host", config.GetString("host"), 0) + crash_on_error(err) + + response, err := self.service.Etcd.Get(path, false, true) + for _, node := range flatten(response.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 + // TODO: use modindex to make sure watch catches everything + _, _ = self.service.Etcd.Watch(path, 0, true, receiver, stop) + }() + + go func() { + for response = range receiver { + switch response.Action { + case "set": + self.handlenode(response.Node) } - //_, err = self.etcd.Get("nodes", false, false) - //if err != nil { - // log.Fatal(err) - //} - //spew.Dump(nodes) + // TODO: handle deletes } }() } diff --git a/core/service.go b/core/service.go index 79e8d7b6a..570fbb446 100644 --- a/core/service.go +++ b/core/service.go @@ -30,8 +30,9 @@ func NewService() *Service { service.init_id() service.Etcd = etcd.NewClient(nil) service.Cluster = db.NewCluster() - service.TopologyMapper = &TopologyMapper{service, "/pilosa/0"} - service.ProcessMapper = NewProcessMapper(service) + service.TopologyMapper = NewTopologyMapper(service, "/pilosa/0") + service.ProcessMapper = NewProcessMapper(service, "/pilosa/0") + service.ProcessMap = NewProcessMap() return service } diff --git a/db/topology.go b/db/topology.go index a6efa87f9..cbcd80692 100644 --- a/db/topology.go +++ b/db/topology.go @@ -5,11 +5,11 @@ import ( //"github.com/davecgh/go-spew/spew" "errors" "fmt" - "github.com/nu7hatch/gouuid" "log" "strconv" "strings" "sync" + "github.com/nu7hatch/gouuid" ) var FrameDoesNotExistError = errors.New("Frame does not exist.") @@ -23,11 +23,63 @@ type Location struct { } type Process struct { - id *uuid.UUID + id *uuid.UUID + host string + port_tcp int + port_http int + mutex sync.Mutex } func NewProcess(id *uuid.UUID) *Process { - return &Process{id} + return &Process{id: id} +} + +func (self *Process) Id() uuid.UUID { + self.mutex.Lock() + defer self.mutex.Unlock() + return *self.id +} + +func (self *Process) SetId(id uuid.UUID) { + self.mutex.Lock() + defer self.mutex.Unlock() + self.id = &id +} + +func (self *Process) Host() string { + self.mutex.Lock() + defer self.mutex.Unlock() + return self.host +} + +func (self *Process) SetHost(host string) { + self.mutex.Lock() + defer self.mutex.Unlock() + self.host = host +} + +func (self *Process) PortTcp() int { + self.mutex.Lock() + defer self.mutex.Unlock() + return self.port_tcp +} + +func (self *Process) SetPortTcp(port int) { + self.mutex.Lock() + defer self.mutex.Unlock() + self.port_tcp = port +} + +func (self *Process) PortHttp() int { + self.mutex.Lock() + defer self.mutex.Unlock() + return self.port_http +} + +func (self *Process) SetPortHttp(port int) { + self.mutex.Lock() + defer self.mutex.Unlock() + self.port_http = port } // Create a Location struct given a string in form "0.0.0.0:0"