Messaging and interface changes.

This commit is contained in:
Cody Soyland 2013-12-13 13:23:10 -06:00
parent 7ec666d912
commit 76bd3cc222
10 changed files with 241 additions and 293 deletions

View file

@ -11,12 +11,12 @@ import (
"errors"
)
type MetaWatcher struct {
type TopologyMapper struct {
service *Service
namespace string
}
func (self *MetaWatcher) Run() {
func (self *TopologyMapper) Run() {
log.Println(self.namespace + "/db")
resp, err := self.service.Etcd.Get(self.namespace + "/db", false, true)
if err != nil {
@ -33,6 +33,7 @@ func (self *MetaWatcher) Run() {
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(self.namespace + "/db", 0, true, receiver, stop)
}()
go func() {
@ -46,11 +47,11 @@ func (self *MetaWatcher) Run() {
}()
}
func NewMetaWatcher(service *Service, namespace string) *MetaWatcher {
return &MetaWatcher{service, namespace}
func NewTopologyMapper(service *Service, namespace string) *TopologyMapper {
return &TopologyMapper{service, namespace}
}
func (self *MetaWatcher) handlenode(node *etcd.Node) error {
func (self *TopologyMapper) handlenode(node *etcd.Node) error {
key := node.Key[len(self.namespace)+1:]
bits := strings.Split(key, "/")
var database *db.Database
@ -123,3 +124,121 @@ func flatten(node *etcd.Node) []*etcd.Node {
}
return nodes
}
type Node struct {
id *uuid.UUID
ip string
port_tcp int
port_http int
}
type ProcessMap struct {
nodes []*db.Process
}
type ProcessMapper struct {
etcd *etcd.Client
nodes []Node
receiver chan *etcd.Response
commands chan *ProcessMapperCommand
}
func NewProcessMapper(service *Service) *ProcessMapper {
return &ProcessMapper{}
}
type ProcessMapperCommand struct {
key string
}
func getKey(input string) string {
bits := strings.Split(input, "/")
return bits[len(bits)-1]
}
func(self *ProcessMapper) getnode(u *uuid.UUID) *Node {
return new(Node)
}
func (self *ProcessMapper) Run() {
var modindex uint64
response, err := self.etcd.Get("nodes", false, true)
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)
if err != nil {
log.Fatal("Not a valid UUID: ", nodestring)
}
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
}
}
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)
default:
spew.Dump("unhandled", response)
}
}
//_, err = self.etcd.Get("nodes", false, false)
//if err != nil {
// log.Fatal(err)
//}
//spew.Dump(nodes)
}
}()
}

View file

@ -4,7 +4,6 @@ import (
"net/http"
"encoding/json"
"log"
"strconv"
"io/ioutil"
"pilosa/db"
"pilosa/query"
@ -21,7 +20,7 @@ func (service *Service) HandleMessage(w http.ResponseWriter, r *http.Request) {
http.Error(w, "Invalid JSON", http.StatusBadRequest)
return
}
service.Inbox <- &message
//service.Inbox <- &message
}
func (service *Service) HandleQuery(w http.ResponseWriter, r *http.Request) {
@ -45,7 +44,8 @@ func (service *Service) HandleStats(w http.ResponseWriter, r *http.Request) {
return
}
encoder := json.NewEncoder(w)
stats := service.GetStats()
//stats := service.GetStats()
stats := ""
err := encoder.Encode(stats)
if err != nil {
log.Fatal("Error encoding stats")
@ -53,18 +53,18 @@ func (service *Service) HandleStats(w http.ResponseWriter, r *http.Request) {
}
func (service *Service) HandleListen(w http.ResponseWriter, r *http.Request) {
listener := service.NewListener()
encoder := json.NewEncoder(w)
for {
select {
case message := <-listener:
err := encoder.Encode(message)
if err != nil {
log.Println("Error sending message")
return
}
}
}
//listener := service.NewListener()
//encoder := json.NewEncoder(w)
//for {
// select {
// case message := <-listener:
// err := encoder.Encode(message)
// if err != nil {
// log.Println("Error sending message")
// return
// }
// }
//}
}
func (service *Service) ServeHTTP() {
@ -73,5 +73,5 @@ func (service *Service) ServeHTTP() {
http.HandleFunc("/query", service.HandleQuery)
http.HandleFunc("/stats", service.HandleStats)
http.HandleFunc("/listen", service.HandleListen)
http.ListenAndServe(":" + strconv.Itoa(service.HttpLocation.Port), nil)
//http.ListenAndServe(":" + strconv.Itoa(service.HttpLocation.Port), nil)
}

18
core/interfaces.go Normal file
View file

@ -0,0 +1,18 @@
package core
import (
"pilosa/db"
)
type Transporter interface {
Init() error
Close()
Send(*db.Message)
Receive() (*db.Message)
}
type Dispatcher interface {
Init() error
Close()
Run()
}

View file

@ -6,15 +6,9 @@ import (
"os"
"os/signal"
"syscall"
"time"
"sync"
"net"
"encoding/gob"
"pilosa/db"
//"net"
//"flag"
//"encoding/gob"
//"io"
"pilosa/config"
)
@ -34,38 +28,25 @@ type Connection struct {
type Service struct {
Stopper
Handler func(db.Message)
Mailbox chan db.Message
Inbox chan *db.Message
Port string
PortHttp string
Etcd *etcd.Client
Location *db.Location
HttpLocation *db.Location
NodeMap db.NodeMap
NodeMapMutex sync.RWMutex
Outbox chan *db.Envelope
ConnectionMap map[db.Location]*Connection
ConnectionMapMutex sync.RWMutex
Listener net.Listener
ConnectionRegisterChannel chan *PersistentConnection
Stats *Stats
Cluster *db.Cluster
Cruncher *Cruncher
MetaWatcher *MetaWatcher
TopologyMapper *TopologyMapper
ProcessMapper *ProcessMapper
ProcessMap *ProcessMap
Transport *Transporter
Dispatcher *Dispatcher
}
func NewService(tcp, http *db.Location) *Service {
service := new(Service)
service.Location = tcp
service.HttpLocation = http
service.Outbox = make(chan *db.Envelope)
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"}
service.TopologyMapper = &TopologyMapper{service, "/pilosa/0"}
service.ProcessMapper = NewProcessMapper(service)
return service
}
@ -76,242 +57,12 @@ func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) {
signal.Notify(termChan, syscall.SIGINT, syscall.SIGTERM)
return termChan, hupChan
}
func (service *Service) SendMessage(message *db.Message) error {
if message.Destination == *service.Location {
return service.DeliverMessage(message)
}
router := service.GetRouterLocation(message.Destination)
var dest *db.Location
if router == *service.Location {
dest = &message.Destination
} else {
dest = &router
}
return service.DoSendMessage(message, dest)
}
func (service *Service) DeliverMessage(message *db.Message) error {
log.Println("do handle of", message)
switch message.Key {
case "ping":
location, ok := message.Data.(db.Location)
if !ok {
log.Println("Invalid ping request!", message)
return nil
}
service.SendMessage(&db.Message{"pong", nil, location})
}
return nil
}
func (service *Service) DoSendMessage(message *db.Message, destination *db.Location) error {
log.Println("send message", message, "to", destination)
service.Outbox <- &db.Envelope{message, destination}
return nil
}
type PersistentConnection struct {
Outbox chan *db.Message
Inbox chan *db.Message
Encoder *gob.Encoder
Decoder *gob.Decoder
Location db.Location
Service *Service
Connection *net.Conn
Identified bool
}
func NewPersistentConnection(service *Service, location *db.Location) *PersistentConnection {
conn := &PersistentConnection{}
conn.Service = service
conn.Identified = true
conn.Outbox = make(chan *db.Message)
conn.Inbox = make(chan *db.Message)
if location != nil {
conn.Location = *location
}
return conn
}
func (conn *PersistentConnection) TryConnect() error {
locationString := conn.Location.ToString()
log.Println("Dialing", locationString)
connection, err := net.Dial("tcp", locationString)
if err != nil {
log.Println("Error in dial!")
return err
}
conn.Connection = &connection
conn.Encoder = gob.NewEncoder(connection)
conn.Decoder = gob.NewDecoder(connection)
return nil
}
func (conn *PersistentConnection) Manage() {
log.Println("Starting manager")
var connectionError chan error
var err error
go func() {
var message *db.Message
var err error
for {
err = conn.Decoder.Decode(&message)
if err != nil {
connectionError <- err
return
}
conn.Inbox <- message
}
}()
var message *db.Message
for {
select {
case message = <-conn.Outbox:
log.Println("sending", message)
err = conn.Encoder.Encode(message)
if err != nil {
// if e, ok := err.(*net.OpError); ok {
// if e.Err == syscall.EPIPE {
// // Client disconnected
log.Println("error sending", err)
go func() { conn.Outbox <- message }() // Resend failed message
if !conn.Identified {
log.Println("Stopping because connection not identified")
return
}
conn.Connect()
}
case message = <-conn.Inbox:
log.Println("receiving", message)
if message.Key == "identify" {
log.Println("Registering connection")
conn.Location = message.Data.(db.Location)
go conn.Service.RegisterConnection(conn)
} else {
conn.Service.Inbox <- message
}
case err = <-connectionError:
log.Println("error receiving", err)
if !conn.Identified {
log.Println("Stopping because connection not identified")
return
}
conn.Connect()
}
}
}
func (conn *PersistentConnection) Connect() {
log.Println("Connnecting to", conn.Location)
err := conn.TryConnect()
if err != nil {
log.Println(err)
log.Println("Connection failed! Waiting 1 second...")
time.Sleep(time.Second)
// Infinite recursion if node never comes up. Not sure if this is a problem.
conn.Connect()
} else {
log.Println("Send identify message")
conn.Encoder.Encode(db.Message{"identify", *conn.Service.Location, conn.Location})
//go func() { conn.Outbox <- &db.Message{"identify", *conn.Service.Location, conn.Location} }()
}
log.Println("Connect ending")
}
type ConnectionMapping map[db.Location]*PersistentConnection
func (service *Service) RegisterConnection(conn *PersistentConnection) {
service.ConnectionRegisterChannel <- conn
}
func (service *Service) HandleConnections() {
connections := ConnectionMapping{}
log.Println("Handling connections...")
for {
select {
case envelope := <-service.Outbox:
conn, ok := connections[*envelope.Location]
if !ok {
conn = NewPersistentConnection(service, envelope.Location)
connections[*envelope.Location] = conn
conn.Identified = true
go func () {
conn.Connect()
conn.Manage()
}()
}
// Spawning new goroutine so it doesn't block the main event loop while it's sending.
// This emulates an infinitely buffered channel.
go func() { conn.Outbox <- envelope.Message }()
case conn := <-service.ConnectionRegisterChannel:
connections[conn.Location] = conn
}
}
}
func (service *Service) GetRouterLocation(node db.Location) db.Location {
service.NodeMapMutex.RLock()
defer service.NodeMapMutex.RUnlock()
location, ok := service.NodeMap[node]
if !ok {
location, ok = service.NodeMap[*service.Location]
if !ok {
log.Fatal("Cannot find router!!!")
}
}
return location
}
func (service *Service) SetupNetwork() {
log.Println("Setup network")
var err error
locationString := service.Location.ToString()
service.Listener, err = net.Listen("tcp", locationString)
if err != nil {
log.Fatal(err)
}
}
func (service *Service) Serve() {
for {
conn, err := service.Listener.Accept()
if err != nil {
log.Fatal(err)
}
con := NewPersistentConnection(service, nil)
con.Connection = &conn
con.Encoder = gob.NewEncoder(*con.Connection)
con.Decoder = gob.NewDecoder(*con.Connection)
go con.Manage()
}
}
func (service *Service) GetStats() *Stats {
return service.Stats
}
func (service *Service) NewListener() chan *db.Message {
ch := make(chan *db.Message)
return ch
}
////////////////////////////////////////////////
func (service *Service) Run() {
log.Println("Running service...")
//go r.SyncEtcd()
//go service.WatchEtcd()
//go service.HandleConnections()
//service.SetupNetwork()
//go service.Serve()
//go service.HandleInbox()
//go service.ServeHTTP()
go service.MetaWatcher.Run()
go service.TopologyMapper.Run()
go service.Cruncher.Run(config.GetInt("port_tcp"))
sigterm, sighup := service.GetSignals()
@ -327,12 +78,3 @@ func (service *Service) Run() {
}
}
}
func (service *Service) HandleInbox() {
for {
select {
case message := <-service.Inbox:
log.Println("process", message)
}
}
}

View file

@ -5,9 +5,3 @@ type Message struct {
Data interface{} `json:data`
Destination Location
}
type Envelope struct {
Message *Message
Location *Location
}

53
transport/http.go Normal file
View file

@ -0,0 +1,53 @@
package transport
import (
"pilosa/db"
"log"
)
type HttpTransport struct {
port int
outbox chan *db.Message
done chan int
}
func (trans *HttpTransport) Init() error {
log.Println("Bind to port", trans.port)
trans.done = make(chan int)
go trans.Loop()
return nil
}
func (trans *HttpTransport) Loop() {
var message *db.Message
for {
select {
case message = <-trans.outbox:
log.Println(message)
case <-trans.done:
return
}
}
}
func (trans *HttpTransport) Close() {
log.Println("Closing HTTP transport.")
trans.done <- 1
}
func (trans *HttpTransport) Send(node string, message *db.Message) error {
log.Println("Send", message, "to", node)
trans.outbox <- message
return nil
}
func (trans *HttpTransport) Receive() (*db.Message, error) {
return &db.Message{}, nil
}
func NewHttpTransport(port int) *HttpTransport {
trans := new(HttpTransport)
trans.port = port
trans.outbox = make(chan *db.Message, 10)
return trans
}

1
transport/inproc.go Normal file
View file

@ -0,0 +1 @@
package transport

2
transport/tcp.go Normal file
View file

@ -0,0 +1,2 @@
package transport

1
transport/transport.go Normal file
View file

@ -0,0 +1 @@
package transport

View file

@ -0,0 +1,18 @@
package transport
import (
"testing"
"pilosa/db"
. "github.com/smartystreets/goconvey/convey"
)
func TestHttpTransport(t *testing.T) {
Convey("Test HTTP transport", t, func() {
var com Transporter
com = NewHttpTransport(9009)
com.Init()
com.Send("derp", &db.Message{"derp", 42, db.Location{}})
com.Close()
So(1, ShouldEqual, 1)
})
}