Fix: grpc transport concurrent write
This commit is contained in:
parent
dfe601b377
commit
ff31722d77
2 changed files with 13 additions and 42 deletions
|
@ -5,6 +5,7 @@ package gun
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bufio"
|
"bufio"
|
||||||
|
"bytes"
|
||||||
"crypto/tls"
|
"crypto/tls"
|
||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
"errors"
|
"errors"
|
||||||
|
@ -30,6 +31,7 @@ var (
|
||||||
"content-type": []string{"application/grpc"},
|
"content-type": []string{"application/grpc"},
|
||||||
"user-agent": []string{"grpc-go/1.36.0"},
|
"user-agent": []string{"grpc-go/1.36.0"},
|
||||||
}
|
}
|
||||||
|
bufferPool = sync.Pool{New: func() interface{} { return &bytes.Buffer{} }}
|
||||||
)
|
)
|
||||||
|
|
||||||
type DialFn = func(network, addr string) (net.Conn, error)
|
type DialFn = func(network, addr string) (net.Conn, error)
|
||||||
|
@ -119,13 +121,20 @@ func (g *Conn) Read(b []byte) (n int, err error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (g *Conn) Write(b []byte) (n int, err error) {
|
func (g *Conn) Write(b []byte) (n int, err error) {
|
||||||
protobufHeader := appendUleb128([]byte{0x0A}, uint64(len(b)))
|
protobufHeader := [binary.MaxVarintLen64 + 1]byte{0x0A}
|
||||||
|
varuintSize := binary.PutUvarint(protobufHeader[1:], uint64(len(b)))
|
||||||
grpcHeader := make([]byte, 5)
|
grpcHeader := make([]byte, 5)
|
||||||
grpcPayloadLen := uint32(len(protobufHeader) + len(b))
|
grpcPayloadLen := uint32(varuintSize + 1 + len(b))
|
||||||
binary.BigEndian.PutUint32(grpcHeader[1:5], grpcPayloadLen)
|
binary.BigEndian.PutUint32(grpcHeader[1:5], grpcPayloadLen)
|
||||||
|
|
||||||
buffers := net.Buffers{grpcHeader, protobufHeader, b}
|
buf := bufferPool.Get().(*bytes.Buffer)
|
||||||
_, err = buffers.WriteTo(g.writer)
|
defer bufferPool.Put(buf)
|
||||||
|
defer buf.Reset()
|
||||||
|
buf.Write(grpcHeader)
|
||||||
|
buf.Write(protobufHeader[:varuintSize+1])
|
||||||
|
buf.Write(b)
|
||||||
|
|
||||||
|
_, err = g.writer.Write(buf.Bytes())
|
||||||
if err == io.ErrClosedPipe && g.err != nil {
|
if err == io.ErrClosedPipe && g.err != nil {
|
||||||
err = g.err
|
err = g.err
|
||||||
}
|
}
|
||||||
|
|
|
@ -1,38 +0,0 @@
|
||||||
// Copy from: https://github.com/Equim-chan/leb128
|
|
||||||
// License: BSD-3-Clause
|
|
||||||
|
|
||||||
package gun
|
|
||||||
|
|
||||||
var sevenbits = [...]byte{
|
|
||||||
0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, 0x0d, 0x0e, 0x0f,
|
|
||||||
0x10, 0x11, 0x12, 0x13, 0x14, 0x15, 0x16, 0x17, 0x18, 0x19, 0x1a, 0x1b, 0x1c, 0x1d, 0x1e, 0x1f,
|
|
||||||
0x20, 0x21, 0x22, 0x23, 0x24, 0x25, 0x26, 0x27, 0x28, 0x29, 0x2a, 0x2b, 0x2c, 0x2d, 0x2e, 0x2f,
|
|
||||||
0x30, 0x31, 0x32, 0x33, 0x34, 0x35, 0x36, 0x37, 0x38, 0x39, 0x3a, 0x3b, 0x3c, 0x3d, 0x3e, 0x3f,
|
|
||||||
0x40, 0x41, 0x42, 0x43, 0x44, 0x45, 0x46, 0x47, 0x48, 0x49, 0x4a, 0x4b, 0x4c, 0x4d, 0x4e, 0x4f,
|
|
||||||
0x50, 0x51, 0x52, 0x53, 0x54, 0x55, 0x56, 0x57, 0x58, 0x59, 0x5a, 0x5b, 0x5c, 0x5d, 0x5e, 0x5f,
|
|
||||||
0x60, 0x61, 0x62, 0x63, 0x64, 0x65, 0x66, 0x67, 0x68, 0x69, 0x6a, 0x6b, 0x6c, 0x6d, 0x6e, 0x6f,
|
|
||||||
0x70, 0x71, 0x72, 0x73, 0x74, 0x75, 0x76, 0x77, 0x78, 0x79, 0x7a, 0x7b, 0x7c, 0x7d, 0x7e, 0x7f,
|
|
||||||
}
|
|
||||||
|
|
||||||
func appendUleb128(b []byte, v uint64) []byte {
|
|
||||||
// If it's less than or equal to 7-bit
|
|
||||||
if v < 0x80 {
|
|
||||||
return append(b, sevenbits[v])
|
|
||||||
}
|
|
||||||
|
|
||||||
for {
|
|
||||||
c := uint8(v & 0x7f)
|
|
||||||
v >>= 7
|
|
||||||
|
|
||||||
if v != 0 {
|
|
||||||
c |= 0x80
|
|
||||||
}
|
|
||||||
|
|
||||||
b = append(b, c)
|
|
||||||
if c&0x80 == 0 {
|
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return b
|
|
||||||
}
|
|
Loading…
Reference in a new issue