orchard controller run: introduce --experimental-ping-interval (#316)

* orchard controller run: introduce --experimental-ping-interval

* Ensure that --experimental-ping-interval is always larger than 5s
This commit is contained in:
Nikolay Edigaryev
2025-05-15 21:14:17 +04:00
committed by GitHub
parent d52aa91927
commit a37a8914cd
6 changed files with 25 additions and 3 deletions
+14
View File
@@ -15,6 +15,7 @@ import (
"os" "os"
"path/filepath" "path/filepath"
"strconv" "strconv"
"time"
) )
var ErrRunFailed = errors.New("failed to run controller") var ErrRunFailed = errors.New("failed to run controller")
@@ -26,6 +27,7 @@ var noTLS bool
var sshNoClientAuth bool var sshNoClientAuth bool
var experimentalRPCV2 bool var experimentalRPCV2 bool
var noExperimentalRPCV2 bool var noExperimentalRPCV2 bool
var experimentalPingInterval time.Duration
func newRunCommand() *cobra.Command { func newRunCommand() *cobra.Command {
cmd := &cobra.Command{ cmd := &cobra.Command{
@@ -65,6 +67,10 @@ func newRunCommand() *cobra.Command {
_ = cmd.PersistentFlags().MarkHidden("experimental-rpc-v2") _ = cmd.PersistentFlags().MarkHidden("experimental-rpc-v2")
cmd.PersistentFlags().BoolVar(&noExperimentalRPCV2, "no-experimental-rpc-v2", false, cmd.PersistentFlags().BoolVar(&noExperimentalRPCV2, "no-experimental-rpc-v2", false,
"disable experimental RPC v2 (https://github.com/cirruslabs/orchard/issues/235)") "disable experimental RPC v2 (https://github.com/cirruslabs/orchard/issues/235)")
cmd.PersistentFlags().DurationVar(&experimentalPingInterval, "experimental-ping-interval", 0,
"interval between WebSocket PING's sent by the controller to workers and clients, "+
"useful when facing intermediate load balancers/proxies that have timeouts "+
"smaller than the controller's default 30 second interval")
return cmd return cmd
} }
@@ -156,6 +162,14 @@ func runController(cmd *cobra.Command, args []string) (err error) {
controllerOpts = append(controllerOpts, controller.WithExperimentalRPCV2()) controllerOpts = append(controllerOpts, controller.WithExperimentalRPCV2())
} }
if experimentalPingInterval != 0 {
if experimentalPingInterval < 5*time.Second {
return fmt.Errorf("--experimental-ping-interval's value cannot be less than 5 seconds")
}
controllerOpts = append(controllerOpts, controller.WithPingInterval(experimentalPingInterval))
}
controllerInstance, err := controller.New(controllerOpts...) controllerInstance, err := controller.New(controllerOpts...)
if err != nil { if err != nil {
return err return err
+1 -1
View File
@@ -43,7 +43,7 @@ func (controller *Controller) rpcPortForward(ctx *gin.Context) responder.Respond
case <-proxyCtx.Done(): case <-proxyCtx.Done():
// Do not close the WebSocket connection as it should be already closed by our rendezvous party // Do not close the WebSocket connection as it should be already closed by our rendezvous party
return responder.Empty() return responder.Empty()
case <-time.After(30 * time.Second): case <-time.After(controller.pingInterval):
pingCtx, pingCtxCancel := context.WithTimeout(ctx, 5*time.Second) pingCtx, pingCtxCancel := context.WithTimeout(ctx, 5*time.Second)
if err := wsConn.Ping(pingCtx); err != nil { if err := wsConn.Ping(pingCtx); err != nil {
+1 -1
View File
@@ -78,7 +78,7 @@ func (controller *Controller) rpcWatch(ctx *gin.Context) responder.Responder {
return controller.wsError(wsConn, websocket.StatusInternalError, "watch RPC", return controller.wsError(wsConn, websocket.StatusInternalError, "watch RPC",
"failure to write the watch instruction", err) "failure to write the watch instruction", err)
} }
case <-time.After(30 * time.Second): case <-time.After(controller.pingInterval):
pingCtx, pingCtxCancel := context.WithTimeout(ctx, 5*time.Second) pingCtx, pingCtxCancel := context.WithTimeout(ctx, 5*time.Second)
if err := wsConn.Ping(pingCtx); err != nil { if err := wsConn.Ping(pingCtx); err != nil {
+1 -1
View File
@@ -142,7 +142,7 @@ func (controller *Controller) portForward(
} }
return responder.Empty() return responder.Empty()
case <-time.After(30 * time.Second): case <-time.After(controller.pingInterval):
pingCtx, pingCtxCancel := context.WithTimeout(ctx, 5*time.Second) pingCtx, pingCtxCancel := context.WithTimeout(ctx, 5*time.Second)
if err := wsConn.Ping(pingCtx); err != nil { if err := wsConn.Ping(pingCtx); err != nil {
+2
View File
@@ -60,6 +60,7 @@ type Controller struct {
workerOfflineTimeout time.Duration workerOfflineTimeout time.Duration
maxWorkersPerLicense uint maxWorkersPerLicense uint
experimentalRPCV2 bool experimentalRPCV2 bool
pingInterval time.Duration
sshListenAddr string sshListenAddr string
sshSigner ssh.Signer sshSigner ssh.Signer
@@ -75,6 +76,7 @@ func New(opts ...Option) (*Controller, error) {
ipRendezvous: rendezvous.New[rendezvous.ResultWithErrorMessage[string]](), ipRendezvous: rendezvous.New[rendezvous.ResultWithErrorMessage[string]](),
workerOfflineTimeout: 3 * time.Minute, workerOfflineTimeout: 3 * time.Minute,
maxWorkersPerLicense: maxWorkersPerDefaultLicense, maxWorkersPerLicense: maxWorkersPerDefaultLicense,
pingInterval: 30 * time.Second,
} }
// Apply options // Apply options
+6
View File
@@ -59,6 +59,12 @@ func WithExperimentalRPCV2() Option {
} }
} }
func WithPingInterval(pingInterval time.Duration) Option {
return func(controller *Controller) {
controller.pingInterval = pingInterval
}
}
func WithLogger(logger *zap.Logger) Option { func WithLogger(logger *zap.Logger) Option {
return func(controller *Controller) { return func(controller *Controller) {
controller.logger = logger.Sugar() controller.logger = logger.Sugar()