fix: Race when writing to a closed pipe (#1916)

This commit is contained in:
Kyle Carberry
2022-06-01 07:59:03 -05:00
committed by GitHub
parent 1c5d94ed5b
commit 1fa50a9da1
2 changed files with 13 additions and 17 deletions
+9 -1
View File
@@ -31,7 +31,7 @@ func Serve(ctx context.Context, server proto.DRPCProvisionerServer, options *Ser
if options.Listener == nil {
config := yamux.DefaultConfig()
config.LogOutput = io.Discard
stdio, err := yamux.Server(readWriteCloser{
stdio, err := yamux.Server(&readWriteCloser{
ReadCloser: os.Stdin,
Writer: os.Stdout,
}, config)
@@ -54,6 +54,9 @@ func Serve(ctx context.Context, server proto.DRPCProvisionerServer, options *Ser
// short-lived processes that can be executed concurrently.
err = srv.Serve(ctx, options.Listener)
if err != nil {
if errors.Is(err, io.EOF) {
return nil
}
if errors.Is(err, context.Canceled) {
return nil
}
@@ -67,3 +70,8 @@ func Serve(ctx context.Context, server proto.DRPCProvisionerServer, options *Ser
}
return nil
}
type readWriteCloser struct {
io.ReadCloser
io.Writer
}
+4 -16
View File
@@ -3,6 +3,7 @@ package provisionersdk
import (
"context"
"io"
"net"
"github.com/hashicorp/yamux"
"storj.io/drpc"
@@ -17,22 +18,14 @@ const (
// TransportPipe creates an in-memory pipe for dRPC transport.
func TransportPipe() (*yamux.Session, *yamux.Session) {
clientReader, clientWriter := io.Pipe()
serverReader, serverWriter := io.Pipe()
c1, c2 := net.Pipe()
yamuxConfig := yamux.DefaultConfig()
yamuxConfig.LogOutput = io.Discard
client, err := yamux.Client(&readWriteCloser{
ReadCloser: clientReader,
Writer: serverWriter,
}, yamuxConfig)
client, err := yamux.Client(c1, yamuxConfig)
if err != nil {
panic(err)
}
server, err := yamux.Server(&readWriteCloser{
ReadCloser: serverReader,
Writer: clientWriter,
}, yamuxConfig)
server, err := yamux.Server(c2, yamuxConfig)
if err != nil {
panic(err)
}
@@ -44,11 +37,6 @@ func Conn(session *yamux.Session) drpc.Conn {
return &multiplexedDRPC{session}
}
type readWriteCloser struct {
io.ReadCloser
io.Writer
}
// Allows concurrent requests on a single dRPC connection.
// Required for calling functions concurrently.
type multiplexedDRPC struct {