From 85201f3acb8b24f9d7255a87f318db28fd98f19e Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 21 Oct 2013 12:01:09 -0500 Subject: [PATCH] Move things around. --- core/etcd.go | 87 ++++++++++++++++++++++++++++++++ core/http.go | 23 +++++++++ core/service.go | 130 ------------------------------------------------ core/stopper.go | 37 ++++++++++++++ 4 files changed, 147 insertions(+), 130 deletions(-) create mode 100644 core/etcd.go create mode 100644 core/http.go create mode 100644 core/stopper.go diff --git a/core/etcd.go b/core/etcd.go new file mode 100644 index 000000000..b20294eb8 --- /dev/null +++ b/core/etcd.go @@ -0,0 +1,87 @@ +package core + +import ( + "github.com/coreos/go-etcd/etcd" + "log" + "strings" + "time" + "encoding/gob" +) + +func (service *Service) SetupEtcd() { + gob.Register(Location{}) + service.Etcd = etcd.NewClient(nil) + service.NodeMapMutex.Lock() + defer service.NodeMapMutex.Unlock() + service.NodeMap = NodeMap{} + + nodes, err := service.Etcd.Get("nodes") + if err != nil { + log.Fatal(err) + } + for _, node := range nodes { + nodestring := strings.Split(node.Key, "/")[2] + location, err := NewLocation(nodestring) + if err != nil { + log.Fatal(err) + } + routerlocation, err := 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 := NewLocation(nodestring) + if err != nil { + log.Fatal(err) + } + router, err := 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 := 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(2*time.Second) + log.Println("done!") + done <- 1 + } + } +} diff --git a/core/http.go b/core/http.go new file mode 100644 index 000000000..efa180d1f --- /dev/null +++ b/core/http.go @@ -0,0 +1,23 @@ +package core + +import ( + "net/http" + "encoding/json" +) + +func (service *Service) ServeHTTP() { + http.HandleFunc("/message", func(w http.ResponseWriter, r *http.Request) { + if r.Method != "POST" { + http.Error(w, "Only POST allowed", http.StatusMethodNotAllowed) + return + } + var message Message + decoder := json.NewDecoder(r.Body) + if decoder.Decode(&message) != nil { + http.Error(w, "Invalid JSON", http.StatusBadRequest) + return + } + service.Inbox <- &message + }) + http.ListenAndServe(string(service.HttpLocation.Port), nil) +} diff --git a/core/service.go b/core/service.go index 00abb1568..a2b2ba1c5 100644 --- a/core/service.go +++ b/core/service.go @@ -7,15 +7,12 @@ import ( "os/signal" "syscall" "time" - "strings" "sync" "net" "encoding/gob" //"net" //"flag" //"encoding/gob" - "net/http" - "encoding/json" //"io" ) type Message struct { @@ -35,38 +32,6 @@ type Connection struct { Decoder *gob.Decoder } -type Stopper struct { - TermChans []chan int - DoneChans []chan int - Mutex sync.RWMutex -} - -func (stopper *Stopper) Stop() { - var i chan int - var o chan int - stopper.Mutex.RLock() - for _, i = range stopper.TermChans { - go func() { - i <- 1 - }() - } - for _, o = range stopper.DoneChans { - <-o - } - stopper.Mutex.RUnlock() - return -} - -func (stopper *Stopper) GetExitChannels() (chan int, chan int) { - termchan := make(chan int, 1) - donechan := make(chan int, 1) - stopper.Mutex.Lock() - stopper.TermChans = append(stopper.TermChans, termchan) - stopper.DoneChans = append(stopper.DoneChans, donechan) - stopper.Mutex.Unlock() - return termchan, donechan -} - type Service struct { Stopper Handler func(Message) @@ -280,84 +245,6 @@ func (service *Service) HandleConnections() { } } -func (service *Service) SetupEtcd() { - gob.Register(Location{}) - service.Etcd = etcd.NewClient(nil) - service.NodeMapMutex.Lock() - defer service.NodeMapMutex.Unlock() - service.NodeMap = NodeMap{} - - nodes, err := service.Etcd.Get("nodes") - if err != nil { - log.Fatal(err) - } - for _, node := range nodes { - nodestring := strings.Split(node.Key, "/")[2] - location, err := NewLocation(nodestring) - if err != nil { - log.Fatal(err) - } - routerlocation, err := 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 := NewLocation(nodestring) - if err != nil { - log.Fatal(err) - } - router, err := 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 := 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(2*time.Second) - log.Println("done!") - done <- 1 - } - } -} - func (service *Service) GetRouterLocation(node Location) Location { service.NodeMapMutex.RLock() defer service.NodeMapMutex.RUnlock() @@ -393,20 +280,3 @@ func (service *Service) Serve() { go con.Manage() } } - -func (service *Service) ServeHTTP() { - http.HandleFunc("/message", func(w http.ResponseWriter, r *http.Request) { - if r.Method != "POST" { - http.Error(w, "Only POST allowed", http.StatusMethodNotAllowed) - return - } - var message Message - decoder := json.NewDecoder(r.Body) - if decoder.Decode(&message) != nil { - http.Error(w, "Invalid JSON", http.StatusBadRequest) - return - } - service.Inbox <- &message - }) - http.ListenAndServe(string(service.HttpLocation.Port), nil) -} diff --git a/core/stopper.go b/core/stopper.go new file mode 100644 index 000000000..271f605a4 --- /dev/null +++ b/core/stopper.go @@ -0,0 +1,37 @@ +package core + +import ( + "sync" +) + +type Stopper struct { + TermChans []chan int + DoneChans []chan int + Mutex sync.RWMutex +} + +func (stopper *Stopper) Stop() { + var i chan int + var o chan int + stopper.Mutex.RLock() + for _, i = range stopper.TermChans { + go func() { + i <- 1 + }() + } + for _, o = range stopper.DoneChans { + <-o + } + stopper.Mutex.RUnlock() + return +} + +func (stopper *Stopper) GetExitChannels() (chan int, chan int) { + termchan := make(chan int, 1) + donechan := make(chan int, 1) + stopper.Mutex.Lock() + stopper.TermChans = append(stopper.TermChans, termchan) + stopper.DoneChans = append(stopper.DoneChans, donechan) + stopper.Mutex.Unlock() + return termchan, donechan +}