ehanced transport

This commit is contained in:
Todd Gruben 2014-12-20 18:15:32 +00:00
parent ba274ea998
commit 1e7cbc88f2

View file

@ -1,6 +1,7 @@
package transport
import (
"bytes"
"encoding/gob"
"fmt"
"log"
@ -9,6 +10,7 @@ import (
"pilosa/core"
"pilosa/db"
. "pilosa/util"
"sync"
"time"
@ -19,8 +21,11 @@ type connection struct {
transport *TcpTransport
inbox chan *db.Message
outbox chan *db.Message
conn *net.Conn
conn net.Conn
process *GUID
id int
terminate bool
exit chan int
}
type newconnection struct {
@ -32,70 +37,100 @@ func init() {
gob.Register(GUID{})
}
func newConnection(transport *TcpTransport, conn net.Conn, proc *GUID) *connection {
p := new(connection)
p.transport = transport
p.outbox = make(chan *db.Message, 100)
p.inbox = make(chan *db.Message, 100)
p.conn = conn
p.process = proc
p.exit = make(chan int)
return p
}
func (self *connection) manage() {
BeginManageConnection:
for {
if self.conn == nil {
process, err := self.transport.service.ProcessMap.GetProcess(self.process)
if err != nil {
log.Println("transport/tcp: error getting process, retrying in 2 seconds... ", self.process, err)
time.Sleep(2 * time.Second)
continue
}
host_string := fmt.Sprintf("%s:%d", process.Host(), process.PortTcp())
conn, err := net.Dial("tcp", host_string)
if err != nil {
log.Println("transport/tcp: error dialing: ", host_string, " Retrying in 2 seconds...")
time.Sleep(2 * time.Second)
continue
}
self.conn = &conn
go func() {
self.outbox <- &db.Message{self.transport.service.Id}
}()
self.exit = make(chan int)
self.serviceConnection()
close(self.exit)
time.Sleep(2 * time.Second)
self.conn = nil
if self.terminate {
break
}
encoder := gob.NewEncoder(*self.conn)
decoder := gob.NewDecoder(*self.conn)
var exit = make(chan int)
}
}
func (self *connection) Shutdown() {
self.terminate = true
if self.conn != nil {
self.conn.Close()
}
}
func (self *connection) serviceConnection() {
var host_string string
if self.conn == nil {
process, err := self.transport.service.ProcessMap.GetProcess(self.process)
if err != nil {
log.Println("transport/tcp: error getting process, retrying in 2 seconds... ", self.process, err)
return
}
host_string = fmt.Sprintf("%s:%d", process.Host(), process.PortTcp())
log.Println("Connecting:", host_string)
conn, err := net.Dial("tcp", host_string)
if err != nil {
log.Println("transport/tcp: error dialing: ", host_string, " Retrying in 2 seconds...")
return
}
self.conn = conn
go func() {
for {
var mess *db.Message
err := decoder.Decode(&mess)
if err != nil {
log.Println("transport/tcp: error decoding message: ", err.Error())
exit <- 1
return
}
self.inbox <- mess
}
//register on server
self.outbox <- &db.Message{self.transport.service.Id}
}()
}
encoder := gob.NewEncoder(self.conn)
decoder := gob.NewDecoder(self.conn)
var wg sync.WaitGroup
wg.Add(1)
go func() {
for {
select {
case message := <-self.outbox:
err := encoder.Encode(message)
if err != nil {
log.Println(err.Error())
return
}
case message := <-self.inbox:
identifier, ok := message.Data.(GUID)
if ok {
// message is connection registration; bypass inbox and register
self.process = &identifier
self.transport.reg <- &newconnection{&identifier, self}
} else {
self.transport.inbox <- message
}
case <-exit:
if self.process != nil {
self.conn = nil
continue BeginManageConnection
} else {
return
}
var mess *db.Message
err := decoder.Decode(&mess)
if err != nil {
log.Println("transport/tcp: error decoding message: ", host_string, err.Error())
self.exit <- 1
wg.Done()
return
}
self.inbox <- mess
}
}()
breakout := true
for breakout {
select {
case message := <-self.outbox:
err := encoder.Encode(message)
if err != nil {
log.Println("Sending to ", host_string, err.Error())
breakout = false
}
case message := <-self.inbox:
identifier, ok := message.Data.(GUID)
if ok {
// message is connection registration; bypass inbox and register
self.process = &identifier
self.transport.reg <- &newconnection{&identifier, self}
} else {
self.transport.inbox <- message
}
case <-self.exit:
breakout = false
}
}
self.conn.Close()
wg.Wait()
}
type TcpTransport struct {
@ -105,6 +140,22 @@ type TcpTransport struct {
outbox chan db.Envelope
connections map[GUID]*connection
reg chan *newconnection
enc *gob.Encoder
dec *gob.Decoder
}
func NewTcpTransport(service *core.Service) *TcpTransport {
p := new(TcpTransport)
var network bytes.Buffer // Stand-in for a network connection
p.service = service
p.port = config.GetInt("port_tcp")
p.inbox = make(chan *db.Message, 100)
p.outbox = make(chan db.Envelope, 100)
p.connections = make(map[GUID]*connection)
p.reg = make(chan *newconnection)
p.enc = gob.NewEncoder(&network) // Will write to network.
p.dec = gob.NewDecoder(&network) // Will read from network.
return p
}
func (self *TcpTransport) Run() {
@ -112,15 +163,19 @@ func (self *TcpTransport) Run() {
go self.listen()
for {
select {
case env := <-self.outbox:
case env := <-self.outbox: //transport outbox
con, ok := self.connections[*(env.Host)]
if !ok {
con = &connection{self, make(chan *db.Message, 100), make(chan *db.Message, 100), nil, env.Host}
con = newConnection(self, nil, env.Host)
go con.manage()
self.connections[*env.Host] = con
}
con.outbox <- env.Message
case nc := <-self.reg:
before, present := self.connections[*nc.id]
if present {
before.Shutdown()
}
self.connections[*nc.id] = nc.connection
}
}
@ -139,23 +194,37 @@ func (self *TcpTransport) listen() {
time.Sleep(2 * time.Second)
continue
}
go self.manage(&conn)
go self.manage(conn)
}
}
func (self *TcpTransport) manage(conn *net.Conn) {
con := &connection{self, make(chan *db.Message, 1024), make(chan *db.Message, 1024), conn, nil}
func (self *TcpTransport) manage(c net.Conn) {
con := newConnection(self, c, nil)
con.manage()
}
func (self *TcpTransport) Close() {
log.Println("Shutting down TCP transport")
}
func (self *TcpTransport) adjust(in *db.Message) *db.Message {
err := self.enc.Encode(in)
var out db.Message
err = self.dec.Decode(&out)
if err != nil {
log.Println(err)
}
return &out
}
func (self *TcpTransport) Send(message *db.Message, host *GUID) {
envelope := db.Envelope{message, host}
notify.Post("outbox", &envelope)
self.outbox <- envelope
if Equal(host, self.service.Id) {
self.inbox <- self.adjust(message)
} else {
envelope := db.Envelope{message, host}
notify.Post("outbox", &envelope)
self.outbox <- envelope
}
}
func (self *TcpTransport) Receive() *db.Message {
@ -167,7 +236,3 @@ func (self *TcpTransport) Receive() *db.Message {
func (self *TcpTransport) Push(message *db.Message) {
self.inbox <- message
}
func NewTcpTransport(service *core.Service) *TcpTransport {
return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message, 100), make(chan db.Envelope, 100), make(map[GUID]*connection), make(chan *newconnection)}
}