From a0ecea3fa7576311dcf4f3b25fe90a5736d917ab Mon Sep 17 00:00:00 2001 From: Felix Lange Date: Thu, 24 Jan 2019 13:42:25 +0100 Subject: [PATCH] rpc: fix two concurrency issues in client Read operations could be dropped after reconnect because the new reader goroutine would be started before the previous one exited. Start it later. For short-lived connections performing a single call, the result would sometimes be dropped because we closed the connection before waiting for handler shutdown. Reverse shutdown order to fix this problem. --- rpc/client.go | 4 ++-- rpc/server_test.go | 41 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 43 insertions(+), 2 deletions(-) diff --git a/rpc/client.go b/rpc/client.go index 6ca5b3c6cc..1b3c3ada0c 100644 --- a/rpc/client.go +++ b/rpc/client.go @@ -117,8 +117,8 @@ func (c *Client) newClientConn(conn ServerCodec) *clientConn { } func (cc *clientConn) close(err error, inflightReq *requestOp) { - cc.codec.Close() cc.handler.close(err, inflightReq) + cc.codec.Close() } type readOp struct { @@ -558,7 +558,6 @@ func (c *Client) dispatch(codec ServerCodec) { // Reconnect: case newcodec := <-c.reconnected: log.Debug("RPC client reconnected", "reading", reading, "conn", newcodec.RemoteAddr()) - go c.read(newcodec) if reading { // Wait for the previous read loop to exit. This is a rare case which // happens if this loop isn't notified in time after the connection breaks. @@ -568,6 +567,7 @@ func (c *Client) dispatch(codec ServerCodec) { conn.close(errClientReconnected, lastOp) c.drainRead() } + go c.read(newcodec) reading = true conn = c.newClientConn(newcodec) // Re-register the in-flight request on the new handler diff --git a/rpc/server_test.go b/rpc/server_test.go index 13d268248c..39099546bb 100644 --- a/rpc/server_test.go +++ b/rpc/server_test.go @@ -18,6 +18,7 @@ package rpc import ( "bufio" + "bytes" "io" "io/ioutil" "net" @@ -109,3 +110,43 @@ func runTestScript(t *testing.T, file string) { } } } + +// This test checks that responses are delivered for very short-lived connections that +// only carry a single request. +func TestServerShortLivedConn(t *testing.T) { + server := newTestServer() + defer server.Stop() + + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal("can't listen:", err) + } + defer listener.Close() + go server.ServeListener(listener) + + var ( + request = `{"jsonrpc":"2.0","id":1,"method":"rpc_modules"}` + "\n" + wantResp = `{"jsonrpc":"2.0","id":1,"result":{"nftest":"1.0","rpc":"1.0","test":"1.0"}}` + "\n" + deadline = time.Now().Add(10 * time.Second) + ) + for i := 0; i < 20; i++ { + conn, err := net.Dial("tcp", listener.Addr().String()) + if err != nil { + t.Fatal("can't dial:", err) + } + defer conn.Close() + conn.SetDeadline(deadline) + // Write the request, then half-close the connection so the server stops reading. + conn.Write([]byte(request)) + conn.(*net.TCPConn).CloseWrite() + // Now try to get the response. + buf := make([]byte, 2000) + n, err := conn.Read(buf) + if err != nil { + t.Fatal("read error:", err) + } + if !bytes.Equal(buf[:n], []byte(wantResp)) { + t.Fatalf("wrong response: %s", buf[:n]) + } + } +}