remove headerPool to simplify the code

This commit is contained in:
gzdaijie
2020-10-01 19:24:57 +08:00
parent 2763266fef
commit 3792afdf1f
8 changed files with 23 additions and 46 deletions
-5
View File
@@ -2,7 +2,6 @@ package codec
import (
"io"
"sync"
)
type Header struct {
@@ -11,10 +10,6 @@ type Header struct {
Error string
}
var HeaderPool = sync.Pool{
New: func() interface{} { return &Header{} },
}
type Codec interface {
io.Closer
ReadHeader(*Header) error
+3 -5
View File
@@ -89,15 +89,14 @@ type request struct {
}
func (server *Server) readRequestHeader(cc codec.Codec) (*codec.Header, error) {
h, _ := codec.HeaderPool.Get().(*codec.Header)
if err := cc.ReadHeader(h); err != nil {
codec.HeaderPool.Put(h)
var h codec.Header
if err := cc.ReadHeader(&h); err != nil {
if err != io.EOF && err != io.ErrUnexpectedEOF {
log.Println("rpc server: read header error:", err)
}
return nil, err
}
return h, nil
return &h, nil
}
func (server *Server) readRequest(cc codec.Codec) (*request, error) {
@@ -121,7 +120,6 @@ func (server *Server) sendResponse(cc codec.Codec, h *codec.Header, body interfa
if err := cc.Write(h, body); err != nil {
log.Println("rpc server: write response error:", err)
}
codec.HeaderPool.Put(h) // recycle Header object
}
func (server *Server) handleRequest(cc codec.Codec, req *request, sending *sync.Mutex, wg *sync.WaitGroup) {
+7 -8
View File
@@ -34,7 +34,8 @@ func (call *Call) done() {
// multiple goroutines simultaneously.
type Client struct {
cc codec.Codec
sending sync.Mutex // protect sending a complete request
sending sync.Mutex // protect following
header codec.Header
mu sync.Mutex // protect following
seq uint64
pending map[uint64]*Call
@@ -101,14 +102,12 @@ func (client *Client) send(call *Call) {
}
// prepare request header
h, _ := codec.HeaderPool.Get().(*codec.Header)
h.ServiceMethod = call.ServiceMethod
h.Seq = seq
h.Error = ""
defer codec.HeaderPool.Put(h)
client.header.ServiceMethod = call.ServiceMethod
client.header.Seq = seq
client.header.Error = ""
// encode and send the request
if err := client.cc.Write(h, call.Args); err != nil {
if err := client.cc.Write(&client.header, call.Args); err != nil {
call := client.removeCall(seq)
// call may be nil, it usually means that Write partially failed,
// client has received the response and handled
@@ -120,9 +119,9 @@ func (client *Client) send(call *Call) {
}
func (client *Client) receive() {
var h codec.Header
var err error
for err == nil {
var h codec.Header
if err = client.cc.ReadHeader(&h); err != nil {
break
}
-5
View File
@@ -2,7 +2,6 @@ package codec
import (
"io"
"sync"
)
type Header struct {
@@ -11,10 +10,6 @@ type Header struct {
Error string
}
var HeaderPool = sync.Pool{
New: func() interface{} { return &Header{} },
}
type Codec interface {
io.Closer
ReadHeader(*Header) error
+3 -5
View File
@@ -89,15 +89,14 @@ type request struct {
}
func (server *Server) readRequestHeader(cc codec.Codec) (*codec.Header, error) {
h, _ := codec.HeaderPool.Get().(*codec.Header)
if err := cc.ReadHeader(h); err != nil {
codec.HeaderPool.Put(h)
var h codec.Header
if err := cc.ReadHeader(&h); err != nil {
if err != io.EOF && err != io.ErrUnexpectedEOF {
log.Println("rpc server: read header error:", err)
}
return nil, err
}
return h, nil
return &h, nil
}
func (server *Server) readRequest(cc codec.Codec) (*request, error) {
@@ -121,7 +120,6 @@ func (server *Server) sendResponse(cc codec.Codec, h *codec.Header, body interfa
if err := cc.Write(h, body); err != nil {
log.Println("rpc server: write response error:", err)
}
codec.HeaderPool.Put(h) // recycle Header object
}
func (server *Server) handleRequest(cc codec.Codec, req *request, sending *sync.Mutex, wg *sync.WaitGroup) {
+7 -8
View File
@@ -34,7 +34,8 @@ func (call *Call) done() {
// multiple goroutines simultaneously.
type Client struct {
cc codec.Codec
sending sync.Mutex // protect sending a complete request
sending sync.Mutex // protect following
header codec.Header
mu sync.Mutex // protect following
seq uint64
pending map[uint64]*Call
@@ -101,14 +102,12 @@ func (client *Client) send(call *Call) {
}
// prepare request header
h, _ := codec.HeaderPool.Get().(*codec.Header)
h.ServiceMethod = call.ServiceMethod
h.Seq = seq
h.Error = ""
defer codec.HeaderPool.Put(h)
client.header.ServiceMethod = call.ServiceMethod
client.header.Seq = seq
client.header.Error = ""
// encode and send the request
if err := client.cc.Write(h, call.Args); err != nil {
if err := client.cc.Write(&client.header, call.Args); err != nil {
call := client.removeCall(seq)
// call may be nil, it usually means that Write partially failed,
// client has received the response and handled
@@ -120,9 +119,9 @@ func (client *Client) send(call *Call) {
}
func (client *Client) receive() {
var h codec.Header
var err error
for err == nil {
var h codec.Header
if err = client.cc.ReadHeader(&h); err != nil {
break
}
-5
View File
@@ -2,7 +2,6 @@ package codec
import (
"io"
"sync"
)
type Header struct {
@@ -11,10 +10,6 @@ type Header struct {
Error string
}
var HeaderPool = sync.Pool{
New: func() interface{} { return &Header{} },
}
type Codec interface {
io.Closer
ReadHeader(*Header) error
+3 -5
View File
@@ -94,15 +94,14 @@ type request struct {
}
func (server *Server) readRequestHeader(cc codec.Codec) (*codec.Header, error) {
h, _ := codec.HeaderPool.Get().(*codec.Header)
if err := cc.ReadHeader(h); err != nil {
codec.HeaderPool.Put(h)
var h codec.Header
if err := cc.ReadHeader(&h); err != nil {
if err != io.EOF && err != io.ErrUnexpectedEOF {
log.Println("rpc server: read header error:", err)
}
return nil, err
}
return h, nil
return &h, nil
}
func (server *Server) findService(serviceMethod string) (svc *service, mtype *methodType, err error) {
@@ -156,7 +155,6 @@ func (server *Server) sendResponse(cc codec.Codec, h *codec.Header, body interfa
if err := cc.Write(h, body); err != nil {
log.Println("rpc server: write response error:", err)
}
codec.HeaderPool.Put(h) // recycle Header object
}
func (server *Server) handleRequest(cc codec.Codec, req *request, sending *sync.Mutex, wg *sync.WaitGroup) {