diff --git a/e2e/go.mod b/e2e/go.mod index ed358a56..0fc80dc2 100644 --- a/e2e/go.mod +++ b/e2e/go.mod @@ -250,9 +250,9 @@ require ( github.com/pion/dtls/v3 v3.1.4 // indirect github.com/pion/logging v0.2.4 // indirect github.com/pion/randutil v0.1.0 // indirect - github.com/pion/stun/v3 v3.0.0 // indirect + github.com/pion/stun/v3 v3.1.5 // indirect github.com/pion/transport/v3 v3.0.7 // indirect - github.com/pion/transport/v4 v4.0.1 // indirect + github.com/pion/transport/v4 v4.0.2 // indirect github.com/pion/turn/v4 v4.0.0 // indirect github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c // indirect github.com/pkg/errors v0.9.1 // indirect diff --git a/e2e/go.sum b/e2e/go.sum index e873f8fb..9ed8bee5 100644 --- a/e2e/go.sum +++ b/e2e/go.sum @@ -857,12 +857,12 @@ github.com/pion/logging v0.2.4 h1:tTew+7cmQ+Mc1pTBLKH2puKsOvhm32dROumOZ655zB8= github.com/pion/logging v0.2.4/go.mod h1:DffhXTKYdNZU+KtJ5pyQDjvOAh/GsNSyv1lbkFbe3so= github.com/pion/randutil v0.1.0 h1:CFG1UdESneORglEsnimhUjf33Rwjubwj6xfiOXBa3mA= github.com/pion/randutil v0.1.0/go.mod h1:XcJrSMMbbMRhASFVOlj/5hQial/Y8oH/HVo7TBZq+j8= -github.com/pion/stun/v3 v3.0.0 h1:4h1gwhWLWuZWOJIJR9s2ferRO+W3zA/b6ijOI6mKzUw= -github.com/pion/stun/v3 v3.0.0/go.mod h1:HvCN8txt8mwi4FBvS3EmDghW6aQJ24T+y+1TKjB5jyU= +github.com/pion/stun/v3 v3.1.5 h1:Y1FHlhaI6+4UoC5i/zQf4F7JvdZtB24/05oyy/GF1x8= +github.com/pion/stun/v3 v3.1.5/go.mod h1:zRUghXSQU32Lx5orJsz3uYMkIihweXb3mu5gIns02fs= github.com/pion/transport/v3 v3.0.7 h1:iRbMH05BzSNwhILHoBoAPxoB9xQgOaJk+591KC9P1o0= github.com/pion/transport/v3 v3.0.7/go.mod h1:YleKiTZ4vqNxVwh77Z0zytYi7rXHl7j6uPLGhhz9rwo= -github.com/pion/transport/v4 v4.0.1 h1:sdROELU6BZ63Ab7FrOLn13M6YdJLY20wldXW2Cu2k8o= -github.com/pion/transport/v4 v4.0.1/go.mod h1:nEuEA4AD5lPdcIegQDpVLgNoDGreqM/YqmEx3ovP4jM= +github.com/pion/transport/v4 v4.0.2 h1:ifYlPqNwsy6aKQ9y8yzxXlHae5431ZrH2avkD/Rn6Tk= +github.com/pion/transport/v4 v4.0.2/go.mod h1:06hFI+jCFcok2X2MekVufNZ/uzNZXivGBPfviSVcjgM= github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c h1:+mdjkGKdHQG3305AYmdv1U2eRNDiU2ErMBj1gwrq8eQ= github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c/go.mod h1:7rwL4CYBLnjLxUqIJNnCWiEdr3bn6IUYi15bNlnbCCU= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= diff --git a/go.mod b/go.mod index e5acef0f..4dcfc337 100644 --- a/go.mod +++ b/go.mod @@ -181,9 +181,9 @@ require ( github.com/pingcap/log v1.1.1-0.20241212030209-7e3ff8601a2a // indirect github.com/pingcap/tidb/pkg/parser v0.0.0-20250421232622-526b2c79173d // indirect github.com/pion/randutil v0.1.0 // indirect - github.com/pion/stun/v3 v3.0.0 // indirect + github.com/pion/stun/v3 v3.1.5 // indirect github.com/pion/transport/v3 v3.0.7 // indirect - github.com/pion/transport/v4 v4.0.1 // indirect + github.com/pion/transport/v4 v4.0.2 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/rivo/uniseg v0.2.0 // indirect github.com/shopspring/decimal v1.4.0 // indirect diff --git a/go.sum b/go.sum index 993a0108..c69d3a3a 100644 --- a/go.sum +++ b/go.sum @@ -570,12 +570,12 @@ github.com/pion/logging v0.2.4 h1:tTew+7cmQ+Mc1pTBLKH2puKsOvhm32dROumOZ655zB8= github.com/pion/logging v0.2.4/go.mod h1:DffhXTKYdNZU+KtJ5pyQDjvOAh/GsNSyv1lbkFbe3so= github.com/pion/randutil v0.1.0 h1:CFG1UdESneORglEsnimhUjf33Rwjubwj6xfiOXBa3mA= github.com/pion/randutil v0.1.0/go.mod h1:XcJrSMMbbMRhASFVOlj/5hQial/Y8oH/HVo7TBZq+j8= -github.com/pion/stun/v3 v3.0.0 h1:4h1gwhWLWuZWOJIJR9s2ferRO+W3zA/b6ijOI6mKzUw= -github.com/pion/stun/v3 v3.0.0/go.mod h1:HvCN8txt8mwi4FBvS3EmDghW6aQJ24T+y+1TKjB5jyU= +github.com/pion/stun/v3 v3.1.5 h1:Y1FHlhaI6+4UoC5i/zQf4F7JvdZtB24/05oyy/GF1x8= +github.com/pion/stun/v3 v3.1.5/go.mod h1:zRUghXSQU32Lx5orJsz3uYMkIihweXb3mu5gIns02fs= github.com/pion/transport/v3 v3.0.7 h1:iRbMH05BzSNwhILHoBoAPxoB9xQgOaJk+591KC9P1o0= github.com/pion/transport/v3 v3.0.7/go.mod h1:YleKiTZ4vqNxVwh77Z0zytYi7rXHl7j6uPLGhhz9rwo= -github.com/pion/transport/v4 v4.0.1 h1:sdROELU6BZ63Ab7FrOLn13M6YdJLY20wldXW2Cu2k8o= -github.com/pion/transport/v4 v4.0.1/go.mod h1:nEuEA4AD5lPdcIegQDpVLgNoDGreqM/YqmEx3ovP4jM= +github.com/pion/transport/v4 v4.0.2 h1:ifYlPqNwsy6aKQ9y8yzxXlHae5431ZrH2avkD/Rn6Tk= +github.com/pion/transport/v4 v4.0.2/go.mod h1:06hFI+jCFcok2X2MekVufNZ/uzNZXivGBPfviSVcjgM= github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c h1:+mdjkGKdHQG3305AYmdv1U2eRNDiU2ErMBj1gwrq8eQ= github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c/go.mod h1:7rwL4CYBLnjLxUqIJNnCWiEdr3bn6IUYi15bNlnbCCU= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= diff --git a/packages/api/api.go b/packages/api/api.go index 9bce7e93..1707cfd2 100644 --- a/packages/api/api.go +++ b/packages/api/api.go @@ -1,6 +1,7 @@ package api import ( + "context" "encoding/base64" "errors" "fmt" @@ -45,6 +46,7 @@ const ( operationCallExchangeRelayCertV1 = "CallExchangeRelayCertV1" operationCallGatewayHeartBeatV1 = "CallGatewayHeartBeatV1" operationCallGatewayHeartBeatV2 = "CallGatewayHeartBeatV2" + operationCallGatewayMetricsReportV2 = "CallGatewayMetricsReportV2" operationCallBootstrapInstance = "CallBootstrapInstance" operationCallRegisterInstanceRelay = "CallRegisterInstanceRelay" operationCallRegisterOrgRelay = "CallRegisterOrgRelay" @@ -847,6 +849,25 @@ func CallGatewayHeartBeatV2(httpClient *resty.Client, request GatewayHeartbeatRe return nil } +func CallGatewayMetricsReportV2(ctx context.Context, httpClient *resty.Client, request GatewayMetricsReportRequest) error { + response, err := httpClient. + R(). + SetContext(ctx). + SetHeader("User-Agent", USER_AGENT). + SetBody(request). + Post(fmt.Sprintf("%v/v2/gateways/metrics", config.INFISICAL_URL)) + + if err != nil { + return NewGenericRequestError(operationCallGatewayMetricsReportV2, err) + } + + if response.IsError() { + return NewAPIErrorWithResponse(operationCallGatewayMetricsReportV2, response, nil) + } + + return nil +} + func CallOrgRelayHeartBeat(httpClient *resty.Client, request RelayHeartbeatRequest) error { response, err := httpClient. R(). diff --git a/packages/api/model.go b/packages/api/model.go index ec9aeddb..3188f73c 100644 --- a/packages/api/model.go +++ b/packages/api/model.go @@ -1084,6 +1084,10 @@ type GatewayHeartbeatRequest struct { Capabilities map[string]any `json:"capabilities,omitempty"` } +type GatewayMetricsReportRequest struct { + ActiveChannels int64 `json:"activeChannels"` +} + type RelayLoginRequest struct { Method string `json:"method"` Token string `json:"token,omitempty"` diff --git a/packages/gateway-v2/gateway.go b/packages/gateway-v2/gateway.go index 9fdae907..8eadba7f 100644 --- a/packages/gateway-v2/gateway.go +++ b/packages/gateway-v2/gateway.go @@ -57,6 +57,16 @@ const ( const heartbeatInterval = 3 * time.Minute +// Reported over the same direct HTTP client as the heartbeat. The relay connection cannot carry this: +// the gateway only accepts channels on it, so it runs platform-to-gateway and has no reverse path. +// Kept off the heartbeat itself because that handler probes back through the relay and writes to the +// platform's database, which is far too costly at the cadence pool selection needs. +const metricsReportInterval = 10 * time.Second + +// Must stay below the interval so a stalled endpoint cannot hold the reporting loop past the next +// tick. Without a bound the loop would also stop noticing cancellation and shutdown would hang. +const metricsReportTimeout = 5 * time.Second + const GATEWAY_ROUTING_INFO_OID = "1.3.6.1.4.1.12345.100.1" const GATEWAY_ACTOR_OID = "1.3.6.1.4.1.12345.100.2" const PAM_INFO_OID = "1.3.6.1.4.1.12345.100.3" @@ -144,6 +154,11 @@ type Gateway struct { mongoProxies map[string]*mongoProxyEntry mongoProxiesMu sync.Mutex pkcs11Module Pkcs11Module + + // Every channel reaching this gateway is accepted in handleIncomingChannel, whoever opened it, + // so counting there is the only place that cannot be bypassed by a new caller. The platform uses + // it to pick the least loaded member of a gateway pool. + activeChannels atomic.Int64 } // mongoProxyEntry holds a session-level MongoDB proxy with a ready signal. @@ -384,6 +399,42 @@ func (g *Gateway) reapIdleSessions() { } } +// reportLoad publishes the gateway's active channel count so the platform can route new work to the +// least loaded member of a pool. Failures are not surfaced: a missed report only costs accuracy, and +// the platform falls back to its own view when a gateway stops reporting. +func (g *Gateway) sendMetricsReport(ctx context.Context, count int64) error { + reqCtx, cancel := context.WithTimeout(ctx, metricsReportTimeout) + defer cancel() + return api.CallGatewayMetricsReportV2(reqCtx, g.httpClient, api.GatewayMetricsReportRequest{ActiveChannels: count}) +} + +func (g *Gateway) reportMetrics(ctx context.Context) { + go func() { + ticker := time.NewTicker(metricsReportInterval) + defer ticker.Stop() + + var last int64 = -1 + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + count := g.activeChannels.Load() + // Still republish an unchanged count so the platform can tell a quiet gateway from + // one that has stopped reporting. + if err := g.sendMetricsReport(ctx, count); err != nil { + log.Debug().Msgf("Load report failed: %v", err) + continue + } + if count != last { + log.Debug().Msgf("Reported %d active channels", count) + last = count + } + } + } + }() +} + func (g *Gateway) registerHeartBeat(ctx context.Context, errCh chan error) { sendHeartbeat := func() error { capabilities := map[string]any{} @@ -575,6 +626,7 @@ func (g *Gateway) startHeartbeatOnce(ctx context.Context, errCh chan error) { defer g.heartbeatMu.Unlock() if !g.heartbeatStarted { g.registerHeartBeat(ctx, errCh) + g.reportMetrics(ctx) g.heartbeatStarted = true } } @@ -884,6 +936,9 @@ func (g *Gateway) handleIncomingChannel(newChannel ssh.NewChannel) { } defer channel.Close() + g.activeChannels.Add(1) + defer g.activeChannels.Add(-1) + go ssh.DiscardRequests(requests) // Create mTLS server configuration