mirror of
https://github.com/cirruslabs/orchard.git
synced 2026-09-30 03:51:43 +02:00
Fix websocket.(*Conn).timeoutLoop goroutine leak (#329)
This commit is contained in:
@@ -17,6 +17,8 @@ import (
|
|||||||
"go.uber.org/zap/zapcore"
|
"go.uber.org/zap/zapcore"
|
||||||
"gopkg.in/natefinch/lumberjack.v2"
|
"gopkg.in/natefinch/lumberjack.v2"
|
||||||
"io"
|
"io"
|
||||||
|
"net/http"
|
||||||
|
_ "net/http/pprof"
|
||||||
"os"
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
)
|
)
|
||||||
@@ -37,6 +39,7 @@ var noPKI bool
|
|||||||
var defaultCPU uint64
|
var defaultCPU uint64
|
||||||
var defaultMemory uint64
|
var defaultMemory uint64
|
||||||
var username string
|
var username string
|
||||||
|
var addressPprof string
|
||||||
var debug bool
|
var debug bool
|
||||||
|
|
||||||
func newRunCommand() *cobra.Command {
|
func newRunCommand() *cobra.Command {
|
||||||
@@ -71,6 +74,8 @@ func newRunCommand() *cobra.Command {
|
|||||||
"(\"Local Network\" permission workaround: requires starting \"orchard worker run\" as \"root\", "+
|
"(\"Local Network\" permission workaround: requires starting \"orchard worker run\" as \"root\", "+
|
||||||
"the privileges will be then dropped to the specified user after starting the \"orchard localnetworkhelper\" "+
|
"the privileges will be then dropped to the specified user after starting the \"orchard localnetworkhelper\" "+
|
||||||
"helper process)")
|
"helper process)")
|
||||||
|
cmd.Flags().StringVar(&addressPprof, "listen-pprof", "",
|
||||||
|
"start pprof HTTP server on localhost:6060 for diagnostic purposes (e.g. \"localhost:6060\")")
|
||||||
cmd.Flags().BoolVar(&debug, "debug", false, "enable debug logging")
|
cmd.Flags().BoolVar(&debug, "debug", false, "enable debug logging")
|
||||||
|
|
||||||
return cmd
|
return cmd
|
||||||
@@ -154,6 +159,14 @@ func runWorker(cmd *cobra.Command, args []string) (err error) {
|
|||||||
}()
|
}()
|
||||||
workerOpts = append(workerOpts, worker.WithLogger(logger))
|
workerOpts = append(workerOpts, worker.WithLogger(logger))
|
||||||
|
|
||||||
|
if addressPprof != "" {
|
||||||
|
go func() {
|
||||||
|
if err := http.ListenAndServe(addressPprof, nil); err != nil {
|
||||||
|
logger.Sugar().Errorf("pprof server failed: %v", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
|
||||||
workerInstance, err := worker.New(
|
workerInstance, err := worker.New(
|
||||||
controllerClient,
|
controllerClient,
|
||||||
workerOpts...,
|
workerOpts...,
|
||||||
|
|||||||
@@ -27,6 +27,13 @@ func (controller *Controller) rpcPortForward(ctx *gin.Context) responder.Respond
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return responder.Error(err)
|
return responder.Error(err)
|
||||||
}
|
}
|
||||||
|
defer func() {
|
||||||
|
// Ensure that we always close the accepted WebSocket connection,
|
||||||
|
// otherwise resource leak is possible[1]
|
||||||
|
//
|
||||||
|
// [1]: https://github.com/coder/websocket/issues/445#issuecomment-2053792044
|
||||||
|
_ = wsConn.CloseNow()
|
||||||
|
}()
|
||||||
|
|
||||||
// Respond with the established connection
|
// Respond with the established connection
|
||||||
proxyCtx, err := controller.connRendezvous.Respond(session, rendezvous.ResultWithErrorMessage[net.Conn]{
|
proxyCtx, err := controller.connRendezvous.Respond(session, rendezvous.ResultWithErrorMessage[net.Conn]{
|
||||||
|
|||||||
@@ -37,6 +37,13 @@ func (controller *Controller) rpcWatch(ctx *gin.Context) responder.Responder {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return responder.Error(err)
|
return responder.Error(err)
|
||||||
}
|
}
|
||||||
|
defer func() {
|
||||||
|
// Ensure that we always close the accepted WebSocket connection,
|
||||||
|
// otherwise resource leak is possible[1]
|
||||||
|
//
|
||||||
|
// [1]: https://github.com/coder/websocket/issues/445#issuecomment-2053792044
|
||||||
|
_ = wsConn.CloseNow()
|
||||||
|
}()
|
||||||
|
|
||||||
// Ensure that pongs will be processed by reading
|
// Ensure that pongs will be processed by reading
|
||||||
// from the connection in the background
|
// from the connection in the background
|
||||||
|
|||||||
@@ -102,6 +102,13 @@ func (controller *Controller) portForward(
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return responder.Error(err)
|
return responder.Error(err)
|
||||||
}
|
}
|
||||||
|
defer func() {
|
||||||
|
// Ensure that we always close the accepted WebSocket connection,
|
||||||
|
// otherwise resource leak is possible[1]
|
||||||
|
//
|
||||||
|
// [1]: https://github.com/coder/websocket/issues/445#issuecomment-2053792044
|
||||||
|
_ = wsConn.CloseNow()
|
||||||
|
}()
|
||||||
|
|
||||||
expectedMsgType := websocket.MessageBinary
|
expectedMsgType := websocket.MessageBinary
|
||||||
|
|
||||||
|
|||||||
@@ -55,6 +55,13 @@ func (worker *Worker) handlePortForwardV2(ctx context.Context, portForward *v1.P
|
|||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
defer func() {
|
||||||
|
// Ensure that we always close the accepted WebSocket connection,
|
||||||
|
// otherwise resource leak is possible[1]
|
||||||
|
//
|
||||||
|
// [1]: https://github.com/coder/websocket/issues/445#issuecomment-2053792044
|
||||||
|
_ = netConn.Close()
|
||||||
|
}()
|
||||||
|
|
||||||
// Proxy bytes if the connection was established without errors
|
// Proxy bytes if the connection was established without errors
|
||||||
if errorMessage == "" {
|
if errorMessage == "" {
|
||||||
|
|||||||
Reference in New Issue
Block a user