From f4834a3a0532a88c680ccb08848659ee83a9b32e Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 17 Jan 2014 16:54:36 -0600 Subject: [PATCH 1/4] fix hold test. --- hold/hold_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hold/hold_test.go b/hold/hold_test.go index a9e057108..0ca3ce154 100644 --- a/hold/hold_test.go +++ b/hold/hold_test.go @@ -22,9 +22,9 @@ func TestHoldChan(t *testing.T) { Convey("get then set", t, func() { id := uuid.RandomUUID() go func() { - time.Sleep(time.Second / 10) Hold.Set(&id, "derpsy", 10) }() + time.Sleep(time.Second / 10) derp, _ := Hold.Get(&id, 10) So(derp, ShouldEqual, "derpsy") }) From ad96f3f7209aca63cb6317cf787bd04cc1a83eeb Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 17 Jan 2014 16:55:32 -0600 Subject: [PATCH 2/4] persistent connection support. --- transport/tcp.go | 137 ++++++++++++++++++++++++++++++++++++----------- 1 file changed, 106 insertions(+), 31 deletions(-) diff --git a/transport/tcp.go b/transport/tcp.go index c8f242571..0be715a01 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -8,21 +8,108 @@ import ( "pilosa/config" "pilosa/core" "pilosa/db" + "github.com/davecgh/go-spew/spew" + "time" "tux21b.org/v1/gocql/uuid" ) -type qmessage struct { +type envelope struct { message *db.Message host *uuid.UUID } +type connection struct { + transport *TcpTransport + inbox chan *db.Message + outbox chan *db.Message + conn *net.Conn + process *uuid.UUID +} + +type newconnection struct { + id *uuid.UUID + connection *connection +} + +func init() { + gob.Register(uuid.UUID{}) +} + +func (self *connection) manage() { +BeginManageConnection: + for { + log.Println("manage", self) + if self.conn == nil { + process, err := self.transport.service.ProcessMap.GetProcess(self.process) + if err != nil { + log.Fatal("transport/tcp: ", err) + return + } + 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} + }() + } + encoder := gob.NewEncoder(*self.conn) + decoder := gob.NewDecoder(*self.conn) + var exit = make(chan int) + go func() { + for { + var mess *db.Message + err := decoder.Decode(&mess) + if err != nil { + exit <- 1 + return + } + self.inbox <- mess + } + }() + 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.(uuid.UUID) + if ok { + spew.Dump("ok, sending to reg chan", identifier) + // message is connection registration; bypass inbox and register + self.process = &identifier + self.transport.reg <- &newconnection{&identifier, self} + } else { + self.transport.inbox <- message + } + case <-exit: + log.Println("decoder stopped!") + if self.process != nil { + self.conn = nil + continue BeginManageConnection + } else { + return + } + } + } + } +} + type TcpTransport struct { - service *core.Service - port int - inbox chan *db.Message - outbox chan qmessage - stop chan bool + service *core.Service + port int + inbox chan *db.Message + outbox chan envelope + connections map[uuid.UUID]*connection + reg chan *newconnection } func (self *TcpTransport) Run() { @@ -30,22 +117,16 @@ func (self *TcpTransport) Run() { go self.listen() for { select { - case qmessage := <-self.outbox: - // todo: persistent connections - process, err := self.service.ProcessMap.GetProcess(qmessage.host) - if err != nil { - log.Println("transport/tcp", err) - return + case env := <-self.outbox: + con, ok := self.connections[*(env.host)] + if !ok { + con = &connection{self, make(chan *db.Message), make(chan *db.Message), nil, env.host} + go con.manage() + self.connections[*env.host] = con } - host_string := fmt.Sprintf("%s:%d", process.Host(), process.PortTcp()) - conn, err := net.Dial("tcp", host_string) - encoder := gob.NewEncoder(conn) - err = encoder.Encode(qmessage.message) - if err != nil { - log.Println(err.Error()) - return - } - //time.Sleep(1 * time.Second) + con.outbox <- env.message + case nc := <-self.reg: + self.connections[*nc.id] = nc.connection } } } @@ -66,14 +147,8 @@ func (self *TcpTransport) listen() { } func (self *TcpTransport) manage(conn *net.Conn) { - decoder := gob.NewDecoder(*conn) - var mess *db.Message - err := decoder.Decode(&mess) - if err != nil { - log.Println("tcp/transport", err.Error()) - return - } - self.inbox <- mess + con := &connection{self, make(chan *db.Message), make(chan *db.Message), conn, nil} + con.manage() } func (self *TcpTransport) Close() { @@ -81,7 +156,7 @@ func (self *TcpTransport) Close() { } func (self *TcpTransport) Send(message *db.Message, host *uuid.UUID) { - self.outbox <- qmessage{message, host} + self.outbox <- envelope{message, host} } func (self *TcpTransport) Receive() *db.Message { @@ -93,5 +168,5 @@ func (self *TcpTransport) Push(message *db.Message) { } func NewTcpTransport(service *core.Service) *TcpTransport { - return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message), make(chan qmessage), nil} + return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message, 100), make(chan envelope, 100), make(map[uuid.UUID]*connection), make(chan *newconnection)} } From a9c1927b8dab5ee96d99832bca671ac27f3704ea Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 17 Jan 2014 19:09:25 -0600 Subject: [PATCH 3/4] Better logging and error handling. Removal of fatal errors. --- transport/tcp.go | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/transport/tcp.go b/transport/tcp.go index 0be715a01..1f5e91956 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -8,7 +8,6 @@ import ( "pilosa/config" "pilosa/core" "pilosa/db" - "github.com/davecgh/go-spew/spew" "time" "tux21b.org/v1/gocql/uuid" @@ -43,8 +42,9 @@ BeginManageConnection: if self.conn == nil { process, err := self.transport.service.ProcessMap.GetProcess(self.process) if err != nil { - log.Fatal("transport/tcp: ", err) - return + 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) @@ -66,6 +66,7 @@ BeginManageConnection: var mess *db.Message err := decoder.Decode(&mess) if err != nil { + log.Println("transport/tcp: error decoding message: ", err.Error()) exit <- 1 return } @@ -83,7 +84,6 @@ BeginManageConnection: case message := <-self.inbox: identifier, ok := message.Data.(uuid.UUID) if ok { - spew.Dump("ok, sending to reg chan", identifier) // message is connection registration; bypass inbox and register self.process = &identifier self.transport.reg <- &newconnection{&identifier, self} @@ -91,7 +91,6 @@ BeginManageConnection: self.transport.inbox <- message } case <-exit: - log.Println("decoder stopped!") if self.process != nil { self.conn = nil continue BeginManageConnection @@ -135,12 +134,14 @@ func (self *TcpTransport) listen() { port_string := fmt.Sprintf(":%d", self.port) l, e := net.Listen("tcp", port_string) if e != nil { - log.Fatal("Cannot bind to port!", self.port) + log.Fatal("Cannot bind to port! ", self.port) } for { conn, err := l.Accept() if err != nil { - log.Fatal("Cannot accept message!") + log.Println("Error accepting, trying again in 2 sec... ", err) + time.Sleep(2 * time.Second) + continue } go self.manage(&conn) } From d6040199091c33388b70864d8812005ddace237f Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 17 Jan 2014 19:10:23 -0600 Subject: [PATCH 4/4] Use buffered channels for new connections. --- transport/tcp.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/transport/tcp.go b/transport/tcp.go index 1f5e91956..aa1632ad5 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -119,7 +119,7 @@ func (self *TcpTransport) Run() { case env := <-self.outbox: con, ok := self.connections[*(env.host)] if !ok { - con = &connection{self, make(chan *db.Message), make(chan *db.Message), nil, env.host} + con = &connection{self, make(chan *db.Message, 100), make(chan *db.Message, 100), nil, env.host} go con.manage() self.connections[*env.host] = con } @@ -148,7 +148,7 @@ func (self *TcpTransport) listen() { } func (self *TcpTransport) manage(conn *net.Conn) { - con := &connection{self, make(chan *db.Message), make(chan *db.Message), conn, nil} + con := &connection{self, make(chan *db.Message, 100), make(chan *db.Message, 100), conn, nil} con.manage() }