From 323aa46d1c049d25c25b55c24674ec6c96163825 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Wed, 6 May 2020 08:59:58 +0000 Subject: [PATCH 1/2] fix (pipes): check websocket errors inside CopyToWebsocket() Previously we were treating EOF on the reader as no-error, meaning that operations like Kubernetes Describe would retry endlessly when finished. --- app/multitenant/consul_pipe_router.go | 4 ++-- app/pipes.go | 6 ++--- common/xfer/pipes.go | 32 +++++++++++++++------------ probe/appclient/app_client.go | 6 ++--- probe/docker/controls_test.go | 10 ++++----- 5 files changed, 29 insertions(+), 29 deletions(-) diff --git a/app/multitenant/consul_pipe_router.go b/app/multitenant/consul_pipe_router.go index 88f6954b1..7546c1a8a 100644 --- a/app/multitenant/consul_pipe_router.go +++ b/app/multitenant/consul_pipe_router.go @@ -256,7 +256,7 @@ func (pr *consulPipeRouter) privateAPI() { defer conn.Close() end, _ := pipe.Ends() - if err := pipe.CopyToWebsocket(end, conn); err != nil && !xfer.IsExpectedWSCloseError(err) { + if _, err := pipe.CopyToWebsocket(end, conn); err != nil { log.Errorf("%s: Server bridge connection; Error copying pipe to websocket: %v", key, err) } }) @@ -446,7 +446,7 @@ func (bc *bridgeConnection) loop() { bc.conn = conn bc.mtx.Unlock() - if err := bc.pipe.CopyToWebsocket(end, conn); err != nil && !xfer.IsExpectedWSCloseError(err) { + if _, err := bc.pipe.CopyToWebsocket(end, conn); err != nil { log.Errorf("%s: Client bridge connection; Error copying pipe to websocket: %v", bc.key, err) } conn.Close() diff --git a/app/pipes.go b/app/pipes.go index c88fc0549..994ebd27a 100644 --- a/app/pipes.go +++ b/app/pipes.go @@ -67,13 +67,11 @@ func handlePipeWs(pr PipeRouter, end End) CtxHandlerFunc { } defer conn.Close() - if err := pipe.CopyToWebsocket(endIO, conn); err != nil { + if _, err := pipe.CopyToWebsocket(endIO, conn); err != nil { if span := opentracing.SpanFromContext(ctx); span != nil { span.LogKV("error", err.Error()) } - if !xfer.IsExpectedWSCloseError(err) { - log.Errorf("Error copying to pipe %s (%d) websocket: %v", id, end, err) - } + log.Errorf("Error copying to pipe %s (%d) websocket: %v", id, end, err) } } } diff --git a/common/xfer/pipes.go b/common/xfer/pipes.go index 0a60af507..ffd3a8837 100644 --- a/common/xfer/pipes.go +++ b/common/xfer/pipes.go @@ -11,7 +11,7 @@ import ( // to the UI. type Pipe interface { Ends() (io.ReadWriter, io.ReadWriter) - CopyToWebsocket(io.ReadWriter, Websocket) error + CopyToWebsocket(io.ReadWriter, Websocket) (bool, error) Close() error Closed() bool @@ -99,27 +99,26 @@ func (p *pipe) OnClose(f func()) { } // CopyToWebsocket copies pipe data to/from a websocket. It blocks. -func (p *pipe) CopyToWebsocket(end io.ReadWriter, conn Websocket) error { +// Returns bool 'done' and an error, masked if websocket closed in an expected manner. +func (p *pipe) CopyToWebsocket(end io.ReadWriter, conn Websocket) (bool, error) { p.mtx.Lock() if p.closed { p.mtx.Unlock() - return nil + return true, nil } p.wg.Add(1) p.mtx.Unlock() defer p.wg.Done() - // The goroutines below both post their errors to the channel, but if you close() - // the pipe before any errors then the pipe may not get read from. Therefore it - // needs up to 2 slots free. - errors := make(chan error, 2) + endError := make(chan error, 1) + connError := make(chan error, 1) // Read-from-UI loop go func() { for { _, buf, err := conn.ReadMessage() // TODO type should be binary message if err != nil { - errors <- err + connError <- err return } @@ -128,7 +127,7 @@ func (p *pipe) CopyToWebsocket(end io.ReadWriter, conn Websocket) error { } if _, err := end.Write(buf); err != nil { - errors <- err + endError <- err return } } @@ -140,7 +139,7 @@ func (p *pipe) CopyToWebsocket(end io.ReadWriter, conn Websocket) error { for { n, err := end.Read(buf) if err != nil { - errors <- err + endError <- err return } @@ -149,7 +148,7 @@ func (p *pipe) CopyToWebsocket(end io.ReadWriter, conn Websocket) error { } if err := conn.WriteMessage(websocket.BinaryMessage, buf[:n]); err != nil { - errors <- err + connError <- err return } } @@ -158,9 +157,14 @@ func (p *pipe) CopyToWebsocket(end io.ReadWriter, conn Websocket) error { // block until one of the goroutines exits // this convoluted mechanism is to ensure we only close the websocket once. select { - case err := <-errors: - return err + case err := <-endError: + return false, err + case err := <-connError: + if IsExpectedWSCloseError(err) { + return false, nil + } + return false, err case <-p.quit: - return nil + return true, nil } } diff --git a/probe/appclient/app_client.go b/probe/appclient/app_client.go index 3c04e4e40..04bbe7e4b 100644 --- a/probe/appclient/app_client.go +++ b/probe/appclient/app_client.go @@ -373,10 +373,8 @@ func (c *appClient) pipeConnection(id string, pipe xfer.Pipe) (bool, error) { defer c.closeConn(id) _, remote := pipe.Ends() - if err := pipe.CopyToWebsocket(remote, conn); err != nil && !xfer.IsExpectedWSCloseError(err) { - return false, err - } - return false, nil + done, err := pipe.CopyToWebsocket(remote, conn) + return done, err } func (c *appClient) PipeConnection(id string, pipe xfer.Pipe) { diff --git a/probe/docker/controls_test.go b/probe/docker/controls_test.go index afb157056..0af1a0c45 100644 --- a/probe/docker/controls_test.go +++ b/probe/docker/controls_test.go @@ -46,11 +46,11 @@ func TestControls(t *testing.T) { type mockPipe struct{} -func (mockPipe) Ends() (io.ReadWriter, io.ReadWriter) { return nil, nil } -func (mockPipe) CopyToWebsocket(io.ReadWriter, xfer.Websocket) error { return nil } -func (mockPipe) Close() error { return nil } -func (mockPipe) Closed() bool { return false } -func (mockPipe) OnClose(func()) {} +func (mockPipe) Ends() (io.ReadWriter, io.ReadWriter) { return nil, nil } +func (mockPipe) CopyToWebsocket(io.ReadWriter, xfer.Websocket) (bool, error) { return true, nil } +func (mockPipe) Close() error { return nil } +func (mockPipe) Closed() bool { return false } +func (mockPipe) OnClose(func()) {} func TestPipes(t *testing.T) { oldNewPipe := controls.NewPipe From 24197bc71ee0f335330874499e9589873e6de8c8 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Wed, 6 May 2020 09:05:04 +0000 Subject: [PATCH 2/2] fix (pipes): silence EOF error This comes back when an operation like Kubernetes Describe finishes; hide it so we don't print as an error. --- probe/appclient/app_client.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/probe/appclient/app_client.go b/probe/appclient/app_client.go index 04bbe7e4b..4b01220b8 100644 --- a/probe/appclient/app_client.go +++ b/probe/appclient/app_client.go @@ -374,6 +374,9 @@ func (c *appClient) pipeConnection(id string, pipe xfer.Pipe) (bool, error) { _, remote := pipe.Ends() done, err := pipe.CopyToWebsocket(remote, conn) + if err == io.EOF { + return true, nil + } return done, err }