ttrpc: refactor channel to take a conn
Signed-off-by: Stephen J Day <stephen.day@docker.com>
This commit is contained in:
@@ -5,6 +5,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
"io"
|
"io"
|
||||||
|
"net"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"github.com/pkg/errors"
|
"github.com/pkg/errors"
|
||||||
@@ -60,16 +61,18 @@ func writeMessageHeader(w io.Writer, p []byte, mh messageHeader) error {
|
|||||||
var buffers sync.Pool
|
var buffers sync.Pool
|
||||||
|
|
||||||
type channel struct {
|
type channel struct {
|
||||||
|
conn net.Conn
|
||||||
bw *bufio.Writer
|
bw *bufio.Writer
|
||||||
br *bufio.Reader
|
br *bufio.Reader
|
||||||
hrbuf [messageHeaderLength]byte // avoid alloc when reading header
|
hrbuf [messageHeaderLength]byte // avoid alloc when reading header
|
||||||
hwbuf [messageHeaderLength]byte
|
hwbuf [messageHeaderLength]byte
|
||||||
}
|
}
|
||||||
|
|
||||||
func newChannel(w io.Writer, r io.Reader) *channel {
|
func newChannel(conn net.Conn) *channel {
|
||||||
return &channel{
|
return &channel{
|
||||||
bw: bufio.NewWriter(w),
|
conn: conn,
|
||||||
br: bufio.NewReader(r),
|
bw: bufio.NewWriter(conn),
|
||||||
|
br: bufio.NewReader(conn),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,10 +1,10 @@
|
|||||||
package ttrpc
|
package ttrpc
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bufio"
|
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"io"
|
"io"
|
||||||
|
"net"
|
||||||
"reflect"
|
"reflect"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
@@ -16,28 +16,29 @@ import (
|
|||||||
func TestReadWriteMessage(t *testing.T) {
|
func TestReadWriteMessage(t *testing.T) {
|
||||||
var (
|
var (
|
||||||
ctx = context.Background()
|
ctx = context.Background()
|
||||||
buffer bytes.Buffer
|
w, r = net.Pipe()
|
||||||
w = bufio.NewWriter(&buffer)
|
ch = newChannel(w)
|
||||||
ch = newChannel(w, nil)
|
rch = newChannel(r)
|
||||||
messages = [][]byte{
|
messages = [][]byte{
|
||||||
[]byte("hello"),
|
[]byte("hello"),
|
||||||
[]byte("this is a test"),
|
[]byte("this is a test"),
|
||||||
[]byte("of message framing"),
|
[]byte("of message framing"),
|
||||||
}
|
}
|
||||||
)
|
|
||||||
|
|
||||||
for i, msg := range messages {
|
|
||||||
if err := ch.send(ctx, uint32(i), 1, msg); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
var (
|
|
||||||
received [][]byte
|
received [][]byte
|
||||||
r = bufio.NewReader(bytes.NewReader(buffer.Bytes()))
|
errs = make(chan error, 1)
|
||||||
rch = newChannel(nil, r)
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
for i, msg := range messages {
|
||||||
|
if err := ch.send(ctx, uint32(i), 1, msg); err != nil {
|
||||||
|
errs <- err
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
w.Close()
|
||||||
|
}()
|
||||||
|
|
||||||
for {
|
for {
|
||||||
_, p, err := rch.recv(ctx)
|
_, p, err := rch.recv(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -48,31 +49,44 @@ func TestReadWriteMessage(t *testing.T) {
|
|||||||
break
|
break
|
||||||
}
|
}
|
||||||
received = append(received, p)
|
received = append(received, p)
|
||||||
|
|
||||||
|
// make sure we don't have send errors
|
||||||
|
select {
|
||||||
|
case err := <-errs:
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
default:
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if !reflect.DeepEqual(received, messages) {
|
if !reflect.DeepEqual(received, messages) {
|
||||||
t.Fatalf("didn't received expected set of messages: %v != %v", received, messages)
|
t.Fatalf("didn't received expected set of messages: %v != %v", received, messages)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case err := <-errs:
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
default:
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestMessageOversize(t *testing.T) {
|
func TestMessageOversize(t *testing.T) {
|
||||||
var (
|
var (
|
||||||
ctx = context.Background()
|
ctx = context.Background()
|
||||||
buffer bytes.Buffer
|
w, r = net.Pipe()
|
||||||
w = bufio.NewWriter(&buffer)
|
wch, rch = newChannel(w), newChannel(r)
|
||||||
ch = newChannel(w, nil)
|
msg = bytes.Repeat([]byte("a message of massive length"), 512<<10)
|
||||||
msg = bytes.Repeat([]byte("a message of massive length"), 512<<10)
|
errs = make(chan error, 1)
|
||||||
)
|
)
|
||||||
|
|
||||||
if err := ch.send(ctx, 1, 1, msg); err != nil {
|
go func() {
|
||||||
t.Fatal(err)
|
if err := wch.send(ctx, 1, 1, msg); err != nil {
|
||||||
}
|
errs <- err
|
||||||
|
}
|
||||||
// now, read it off the channel with a small buffer
|
}()
|
||||||
var (
|
|
||||||
r = bufio.NewReader(bytes.NewReader(buffer.Bytes()))
|
|
||||||
rch = newChannel(nil, r)
|
|
||||||
)
|
|
||||||
|
|
||||||
_, _, err := rch.recv(ctx)
|
_, _, err := rch.recv(ctx)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
@@ -87,4 +101,12 @@ func TestMessageOversize(t *testing.T) {
|
|||||||
if status.Code() != codes.ResourceExhausted {
|
if status.Code() != codes.ResourceExhausted {
|
||||||
t.Fatalf("expected grpc status code: %v != %v", status.Code(), codes.ResourceExhausted)
|
t.Fatalf("expected grpc status code: %v != %v", status.Code(), codes.ResourceExhausted)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case err := <-errs:
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
default:
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -35,7 +35,7 @@ func NewClient(conn net.Conn) *Client {
|
|||||||
c := &Client{
|
c := &Client{
|
||||||
codec: codec{},
|
codec: codec{},
|
||||||
conn: conn,
|
conn: conn,
|
||||||
channel: newChannel(conn, conn),
|
channel: newChannel(conn),
|
||||||
calls: make(chan *callRequest),
|
calls: make(chan *callRequest),
|
||||||
closed: make(chan struct{}),
|
closed: make(chan struct{}),
|
||||||
done: make(chan struct{}),
|
done: make(chan struct{}),
|
||||||
|
|||||||
@@ -281,7 +281,7 @@ func (c *serverConn) run(sctx context.Context) {
|
|||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
ch = newChannel(c.conn, c.conn)
|
ch = newChannel(c.conn)
|
||||||
ctx, cancel = context.WithCancel(sctx)
|
ctx, cancel = context.WithCancel(sctx)
|
||||||
active int
|
active int
|
||||||
state connState = connStateIdle
|
state connState = connStateIdle
|
||||||
|
|||||||
Reference in New Issue
Block a user