From c6b3b2431232c38bd569d4ab7f1aaa4152ae2df6 Mon Sep 17 00:00:00 2001 From: = Date: Mon, 17 Feb 2025 13:53:07 +0530 Subject: [PATCH] feat: gateway in cli first version\ --- cli/go.mod | 17 ++- cli/go.sum | 34 ++++-- cli/packages/api/api.go | 39 +++++++ cli/packages/api/model.go | 19 ++++ cli/packages/cmd/gateway.go | 61 ++++++++++ cli/packages/gateway/connection.go | 109 ++++++++++++++++++ cli/packages/gateway/gateway.go | 174 +++++++++++++++++++++++++++++ 7 files changed, 438 insertions(+), 15 deletions(-) create mode 100644 cli/packages/cmd/gateway.go create mode 100644 cli/packages/gateway/connection.go create mode 100644 cli/packages/gateway/gateway.go diff --git a/cli/go.mod b/cli/go.mod index b96772c91..7339ec68f 100644 --- a/cli/go.mod +++ b/cli/go.mod @@ -18,14 +18,16 @@ require ( github.com/muesli/reflow v0.3.0 github.com/muesli/roff v0.1.0 github.com/petar-dambovaliev/aho-corasick v0.0.0-20211021192214-5ab2d9280aa9 + github.com/pion/logging v0.2.3 + github.com/pion/turn/v4 v4.0.0 github.com/posthog/posthog-go v0.0.0-20221221115252-24dfed35d71a github.com/rs/cors v1.11.0 github.com/rs/zerolog v1.26.1 github.com/spf13/cobra v1.6.1 github.com/spf13/viper v1.8.1 github.com/stretchr/testify v1.9.0 - golang.org/x/crypto v0.31.0 - golang.org/x/term v0.27.0 + golang.org/x/crypto v0.33.0 + golang.org/x/term v0.29.0 gopkg.in/yaml.v2 v2.4.0 ) @@ -81,6 +83,10 @@ require ( github.com/muesli/termenv v0.15.2 // indirect github.com/oklog/ulid v1.3.1 // indirect github.com/pelletier/go-toml v1.9.3 // indirect + github.com/pion/dtls/v3 v3.0.4 // indirect + github.com/pion/randutil v0.1.0 // indirect + github.com/pion/stun/v3 v3.0.0 // indirect + github.com/pion/transport/v3 v3.0.7 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/pmezard/go-difflib v1.0.0 // indirect github.com/rivo/uniseg v0.2.0 // indirect @@ -88,6 +94,7 @@ require ( github.com/spf13/cast v1.3.1 // indirect github.com/spf13/jwalterweatherman v1.1.0 // indirect github.com/subosito/gotenv v1.2.0 // indirect + github.com/wlynxg/anet v0.0.5 // indirect github.com/xtgo/uuid v0.0.0-20140804021211-a0b114877d4c // indirect go.mongodb.org/mongo-driver v1.10.0 // indirect go.opencensus.io v0.24.0 // indirect @@ -98,9 +105,9 @@ require ( go.opentelemetry.io/otel/trace v1.24.0 // indirect golang.org/x/net v0.33.0 // indirect golang.org/x/oauth2 v0.21.0 // indirect - golang.org/x/sync v0.10.0 // indirect - golang.org/x/sys v0.28.0 // indirect - golang.org/x/text v0.21.0 // indirect + golang.org/x/sync v0.11.0 // indirect + golang.org/x/sys v0.30.0 // indirect + golang.org/x/text v0.22.0 // indirect golang.org/x/time v0.6.0 // indirect google.golang.org/api v0.188.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20240701130421-f6361c86f094 // indirect diff --git a/cli/go.sum b/cli/go.sum index 67a3d3275..20215d776 100644 --- a/cli/go.sum +++ b/cli/go.sum @@ -347,6 +347,18 @@ github.com/pelletier/go-toml v1.9.3 h1:zeC5b1GviRUyKYd6OJPvBU/mcVDVoL1OhT17FCt5d github.com/pelletier/go-toml v1.9.3/go.mod h1:u1nR/EPcESfeI/szUZKdtJ0xRNbUoANCkoOuaOx1Y+c= github.com/petar-dambovaliev/aho-corasick v0.0.0-20211021192214-5ab2d9280aa9 h1:lL+y4Xv20pVlCGyLzNHRC0I0rIHhIL1lTvHizoS/dU8= github.com/petar-dambovaliev/aho-corasick v0.0.0-20211021192214-5ab2d9280aa9/go.mod h1:EHPiTAKtiFmrMldLUNswFwfZ2eJIYBHktdaUTZxYWRw= +github.com/pion/dtls/v3 v3.0.4 h1:44CZekewMzfrn9pmGrj5BNnTMDCFwr+6sLH+cCuLM7U= +github.com/pion/dtls/v3 v3.0.4/go.mod h1:R373CsjxWqNPf6MEkfdy3aSe9niZvL/JaKlGeFphtMg= +github.com/pion/logging v0.2.3 h1:gHuf0zpoh1GW67Nr6Gj4cv5Z9ZscU7g/EaoC/Ke/igI= +github.com/pion/logging v0.2.3/go.mod h1:z8YfknkquMe1csOrxK5kc+5/ZPAzMxbKLX5aXpbpC90= +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/transport/v3 v3.0.7 h1:iRbMH05BzSNwhILHoBoAPxoB9xQgOaJk+591KC9P1o0= +github.com/pion/transport/v3 v3.0.7/go.mod h1:YleKiTZ4vqNxVwh77Z0zytYi7rXHl7j6uPLGhhz9rwo= +github.com/pion/turn/v4 v4.0.0 h1:qxplo3Rxa9Yg1xXDxxH8xaqcyGUtbHYw4QSCvmFWvhM= +github.com/pion/turn/v4 v4.0.0/go.mod h1:MuPDkm15nYSklKpN8vWJ9W2M0PlyQZqYt1McGuxG7mA= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= @@ -410,6 +422,8 @@ github.com/subosito/gotenv v1.2.0/go.mod h1:N0PQaV/YGNqwC0u51sEeR/aUtSLEXKX9iv69 github.com/tidwall/pretty v1.0.0 h1:HsD+QiTn7sK6flMKIvNmpqz1qrpP3Ps6jOKIKMooyg4= github.com/tidwall/pretty v1.0.0/go.mod h1:XNkn88O1ChpSDQmQeStsy+sBenx6DDtFZJxhVysOjyk= github.com/urfave/cli v1.22.5/go.mod h1:Gos4lmkARVdJ6EkW0WaNv/tZAAMe9V7XWyB60NtXRu0= +github.com/wlynxg/anet v0.0.5 h1:J3VJGi1gvo0JwZ/P1/Yc/8p63SoW98B5dHkYDmpgvvU= +github.com/wlynxg/anet v0.0.5/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA= github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= github.com/xdg-go/scram v1.1.1/go.mod h1:RaEWvsqvNKKvBPvcKeFjrG2cJqOkHTiyTpzz23ni57g= github.com/xdg-go/stringprep v1.0.3/go.mod h1:W3f5j4i+9rC0kuIEJL0ky1VpHXQU3ocBgklLGvcBnW8= @@ -458,8 +472,8 @@ golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8U golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.0.0-20211215165025-cf75a172585e/go.mod h1:P+XmwS30IXTQdn5tA2iutPOUgjI07+tq3H3K9MVA1s8= golang.org/x/crypto v0.0.0-20220622213112-05595931fe9d/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= -golang.org/x/crypto v0.31.0 h1:ihbySMvVjLAeSH1IbfcRTkD/iNscyz8rGzjF/E5hV6U= -golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk= +golang.org/x/crypto v0.33.0 h1:IOBPskki6Lysi0lo9qQvbxiQ+FvsCC/YWOecCHAixus= +golang.org/x/crypto v0.33.0/go.mod h1:bVdXmD7IV/4GdElGPozy6U7lWdRXA4qyRVGJV57uQ5M= golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= golang.org/x/exp v0.0.0-20190306152737-a1d7652674e8/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= golang.org/x/exp v0.0.0-20190510132918-efd6b22b2522/go.mod h1:ZjyILWgesfNpC6sMxTJOJm9Kp84zZh5NQWvqDGG3Qr8= @@ -560,8 +574,8 @@ golang.org/x/sync v0.0.0-20200625203802-6e8e738ad208/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.10.0 h1:3NQrjDixjgGwUOCaF8w2+VYHv0Ve/vGYSbdkTa98gmQ= -golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/sync v0.11.0 h1:GGz8+XQP4FvTTrjZPzNKTMFtSXH80RAzG+5ghFPgK9w= +golang.org/x/sync v0.11.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= golang.org/x/sys v0.0.0-20180823144017-11551d06cbcc/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20181026203630-95b1ffbd15a5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -610,11 +624,11 @@ golang.org/x/sys v0.0.0-20210809222454-d867a43fc93e/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.0.0-20220310020820-b874c991c1a5/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.28.0 h1:Fksou7UEQUWlKvIdsqzJmUmCX3cZuD2+P3XyyzwMhlA= -golang.org/x/sys v0.28.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc= +golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= -golang.org/x/term v0.27.0 h1:WP60Sv1nlK1T6SupCHbXzSaN0b9wUmsPoRS9b61A23Q= -golang.org/x/term v0.27.0/go.mod h1:iMsnZpn0cago0GOrHO2+Y7u7JPn5AylBrcoWkElMTSM= +golang.org/x/term v0.29.0 h1:L6pJp37ocefwRRtYPKSWOWzOtWSxVajvz2ldH/xi3iU= +golang.org/x/term v0.29.0/go.mod h1:6bl4lRlvVuDgSf3179VpIxBF0o10JUpXWOnI7nErv7s= golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.1-0.20180807135948-17ff2d5776d2/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= @@ -624,8 +638,8 @@ golang.org/x/text v0.3.4/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.5/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= -golang.org/x/text v0.21.0 h1:zyQAAkrwaneQ066sspRyJaG9VNi/YJ1NfzcGB3hZ/qo= -golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ= +golang.org/x/text v0.22.0 h1:bofq7m3/HAFvbF51jz3Q9wLg3jkvSPuiZu/pD1XwgtM= +golang.org/x/text v0.22.0/go.mod h1:YRoo4H8PVmsu+E3Ou7cqLVH8oXWIHVoX0jqUWALQhfY= golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= golang.org/x/time v0.0.0-20190308202827-9d24e82272b4/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= golang.org/x/time v0.0.0-20191024005414-555d28b269f0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= diff --git a/cli/packages/api/api.go b/cli/packages/api/api.go index 3d050cd97..7b13c6e02 100644 --- a/cli/packages/api/api.go +++ b/cli/packages/api/api.go @@ -544,3 +544,42 @@ func CallUpdateRawSecretsV3(httpClient *resty.Client, request UpdateRawSecretByN return nil } + +func CallRegisterGatewayIdentityV1(httpClient *resty.Client) (*GetRelayCredentialsResponseV1, error) { + var resBody GetRelayCredentialsResponseV1 + response, err := httpClient. + R(). + SetResult(&resBody). + SetHeader("User-Agent", USER_AGENT). + Post(fmt.Sprintf("%v/v1/gateways/register-identity", config.INFISICAL_URL)) + + if err != nil { + return nil, fmt.Errorf("CallRegisterGatewayIdentityV1: Unable to complete api request [err=%w]", err) + } + + if response.IsError() { + return nil, fmt.Errorf("CallRegisterGatewayIdentityV1: Unsuccessful response [%v %v] [status-code=%v] [response=%v]", response.Request.Method, response.Request.URL, response.StatusCode(), response.String()) + } + + return &resBody, nil +} + +func CallExchangeRelayCertV1(httpClient *resty.Client, request ExchangeRelayCertRequestV1) (*ExchangeRelayCertResponseV1, error) { + var resBody ExchangeRelayCertResponseV1 + response, err := httpClient. + R(). + SetResult(&resBody). + SetBody(request). + SetHeader("User-Agent", USER_AGENT). + Post(fmt.Sprintf("%v/v1/gateways/exchange-cert", config.INFISICAL_URL)) + + if err != nil { + return nil, fmt.Errorf("CallExchangeRelayCertV1: Unable to complete api request [err=%w]", err) + } + + if response.IsError() { + return nil, fmt.Errorf("CallExchangeRelayCertV1: Unsuccessful response [%v %v] [status-code=%v] [response=%v]", response.Request.Method, response.Request.URL, response.StatusCode(), response.String()) + } + + return &resBody, nil +} diff --git a/cli/packages/api/model.go b/cli/packages/api/model.go index 73a04200d..72dbbc97b 100644 --- a/cli/packages/api/model.go +++ b/cli/packages/api/model.go @@ -629,3 +629,22 @@ type GetRawSecretV3ByNameResponse struct { } `json:"secret"` ETag string } + +type GetRelayCredentialsResponseV1 struct { + TurnServerUsername string `json:"turnServerUsername"` + TurnServerPassword string `json:"turnServerPassword"` + TurnServerRealm string `json:"turnServerRealm"` + TurnServerAddress string `json:"turnServerAddress"` + InfisicalStaticIp string `json:"infisicalStaticIp"` +} + +type ExchangeRelayCertRequestV1 struct { + RelayAddress string `json:"relayAddress"` +} + +type ExchangeRelayCertResponseV1 struct { + SerialNumber string `json:"serialNumber"` + PrivateKey string `json:"privateKey"` + Certificate string `json:"certificate"` + CertificateChain string `json:"certificateChain"` +} diff --git a/cli/packages/cmd/gateway.go b/cli/packages/cmd/gateway.go new file mode 100644 index 000000000..945340b22 --- /dev/null +++ b/cli/packages/cmd/gateway.go @@ -0,0 +1,61 @@ +package cmd + +import ( + // "fmt" + + // "github.com/Infisical/infisical-merge/packages/api" + // "github.com/Infisical/infisical-merge/packages/models" + "fmt" + + "github.com/Infisical/infisical-merge/packages/gateway" + "github.com/Infisical/infisical-merge/packages/util" + // "github.com/Infisical/infisical-merge/packages/visualize" + // "github.com/rs/zerolog/log" + + // "github.com/go-resty/resty/v2" + "github.com/posthog/posthog-go" + "github.com/spf13/cobra" +) + +var gatewayCmd = &cobra.Command{ + Example: `infisical gateway`, + Short: "Used to infisical gateway", + Use: "gateway", + DisableFlagsInUseLine: true, + Args: cobra.NoArgs, + Run: func(cmd *cobra.Command, args []string) { + token, err := util.GetInfisicalToken(cmd) + if err != nil { + util.HandleError(err, "Unable to parse flag") + } + + if token == nil { + util.HandleError(fmt.Errorf("Token not found")) + } + + gatewayInstance, err := gateway.NewGateway(token.Token) + if err != nil { + util.HandleError(err) + } + + if err = gatewayInstance.ConnectWithRelay(); err != nil { + util.HandleError(err) + } + + if err := gatewayInstance.Listen(); err != nil { + util.HandleError(err) + } + + Telemetry.CaptureEvent("cli-command:gateway", posthog.NewProperties().Set("version", util.CLI_VERSION)) + }, +} + +func init() { + gatewayCmd.SetHelpFunc(func(command *cobra.Command, strings []string) { + command.Flags().MarkHidden("domain") + command.Parent().HelpFunc()(command, strings) + }) + gatewayCmd.Flags().String("token", "", "Connect with Infisical using machine identity access token") + + rootCmd.AddCommand(gatewayCmd) +} diff --git a/cli/packages/gateway/connection.go b/cli/packages/gateway/connection.go new file mode 100644 index 000000000..dfde595cb --- /dev/null +++ b/cli/packages/gateway/connection.go @@ -0,0 +1,109 @@ +package gateway + +import ( + "bufio" + "bytes" + "fmt" + "io" + "net" + "sync" + + "github.com/rs/zerolog/log" +) + +func handleConnection(conn net.Conn) { + defer conn.Close() + log.Info().Msgf("New connection from: %s", conn.RemoteAddr().String()) + + // Use buffered reader for better handling of fragmented data + reader := bufio.NewReader(conn) + for { + msg, err := reader.ReadBytes('\n') + if err != nil { + log.Error().Msgf("Error reading command: %s", err) + return + } + + cmd := bytes.ToUpper(bytes.TrimSpace(bytes.Split(msg, []byte(" "))[0])) + args := bytes.TrimSpace(bytes.TrimPrefix(msg, cmd)) + + switch string(cmd) { + case "FORWARD-TCP": + proxyAddress := string(bytes.Split(args, []byte(" "))[0]) + fmt.Println(proxyAddress) + destTarget, err := net.Dial("tcp", proxyAddress) + fmt.Println(err) + if err != nil { + log.Error().Msgf("Failed to connect to target: %v", err) + return + } + defer destTarget.Close() + + // Handle buffered data + buffered := reader.Buffered() + if buffered > 0 { + bufferedData := make([]byte, buffered) + _, err := reader.Read(bufferedData) + if err != nil { + log.Error().Msgf("Error reading buffered data: %v", err) + return + } + + if _, err = destTarget.Write(bufferedData); err != nil { + log.Error().Msgf("Error writing buffered data: %v", err) + return + } + } + + CopyData(conn, destTarget) + break + case "PING": + conn.Write([]byte("PONG\n")) + default: + log.Error().Msgf("Unknown command: %s", string(cmd)) + break + } + } +} + +type CloseWrite interface { + CloseWrite() error +} + +func CopyData(src, dst net.Conn) { + // Create a WaitGroup to wait for both copy operations + var wg sync.WaitGroup + wg.Add(2) + + // Start copying in both directions + go func() { + defer wg.Done() + if _, err := io.Copy(dst, src); err != nil { + log.Error().Msgf("Error copying postgres->client: %v", err) + } + + if e, ok := dst.(CloseWrite); ok { + log.Print("Closing dst") + e.CloseWrite() + } else { + + log.Print("Not closed") + } + }() + + go func() { + defer wg.Done() + if _, err := io.Copy(src, dst); err != nil { + log.Error().Msgf("Error copying client->postgres: %v", err) + } + if e, ok := src.(CloseWrite); ok { + log.Print("Closing src") + e.CloseWrite() + } else { + log.Print("Not closed") + } + }() + + // Wait for both copies to complete + wg.Wait() +} diff --git a/cli/packages/gateway/gateway.go b/cli/packages/gateway/gateway.go new file mode 100644 index 000000000..484737b79 --- /dev/null +++ b/cli/packages/gateway/gateway.go @@ -0,0 +1,174 @@ +package gateway + +import ( + "crypto/tls" + "crypto/x509" + "fmt" + "net" + "strings" + "time" + + "github.com/Infisical/infisical-merge/packages/api" + "github.com/go-resty/resty/v2" + "github.com/pion/logging" + "github.com/pion/turn/v4" + "github.com/rs/zerolog/log" +) + +type GatewayConfig struct { + TurnServerUsername string + TurnServerPassword string + TurnServerAddress string + InfisicalStaticIp string + SerialNumber string + PrivateKey string + Certificate string + CertificateChain string +} + +type Gateway struct { + httpClient *resty.Client + config *GatewayConfig + client *turn.Client +} + +func NewGateway(identityToken string) (Gateway, error) { + httpClient := resty.New() + httpClient.SetAuthToken(identityToken) + + return Gateway{ + httpClient: httpClient, + config: &GatewayConfig{}, + }, nil +} + +func (g *Gateway) ConnectWithRelay() error { + relayDetails, err := api.CallRegisterGatewayIdentityV1(g.httpClient) + if err != nil { + return err + } + + // Dial TURN Server + conn, err := net.Dial("tcp", relayDetails.TurnServerAddress) + if err != nil { + return fmt.Errorf("Failed to connect with relay server: %w", err) + } + + if tcpConn, ok := conn.(*net.TCPConn); ok { + tcpConn.SetKeepAlive(true) + tcpConn.SetKeepAlivePeriod(10 * time.Second) + tcpConn.SetNoDelay(true) + } + + // Start a new TURN Client and wrap our net.Conn in a STUNConn + // This allows us to simulate datagram based communication over a net.Conn + cfg := &turn.ClientConfig{ + STUNServerAddr: relayDetails.TurnServerAddress, + TURNServerAddr: relayDetails.TurnServerAddress, + Conn: turn.NewSTUNConn(conn), + Username: relayDetails.TurnServerUsername, + Password: relayDetails.TurnServerPassword, + Realm: relayDetails.TurnServerRealm, + LoggerFactory: logging.NewDefaultLoggerFactory(), + } + + client, err := turn.NewClient(cfg) + if err != nil { + return fmt.Errorf("Failed to create relay client: %w", err) + } + + err = client.Listen() + if err != nil { + return fmt.Errorf("Failed to listen to relay server: %w", err) + } + + g.config = &GatewayConfig{ + TurnServerUsername: relayDetails.TurnServerUsername, + TurnServerPassword: relayDetails.TurnServerPassword, + TurnServerAddress: relayDetails.TurnServerAddress, + InfisicalStaticIp: relayDetails.InfisicalStaticIp, + } + // if port not specific allow all port + if !strings.Contains(relayDetails.InfisicalStaticIp, ":") { + g.config.InfisicalStaticIp = g.config.InfisicalStaticIp + ":0" + } + + g.client = client + return nil +} + +func (g *Gateway) Listen() error { + defer g.client.Close() + // Allocate a relay socket on the TURN server. On success, it + // will return a net.PacketConn which represents the remote + // socket. + relayNonTlsConn, err := g.client.AllocateTCP() + if err != nil { + return fmt.Errorf("Failed to allocate relay connection: %w", err) + } + defer func() { + if closeErr := relayNonTlsConn.Close(); closeErr != nil { + log.Error().Msgf("Failed to close connection: %s", closeErr) + } + }() + + peerAddr, err := net.ResolveTCPAddr("tcp", g.config.InfisicalStaticIp) + if err != nil { + return fmt.Errorf("Failed to parse infisical static ip: %w", err) + } + gatewayCert, err := api.CallExchangeRelayCertV1(g.httpClient, api.ExchangeRelayCertRequestV1{ + RelayAddress: relayNonTlsConn.Addr().String(), + }) + if err != nil { + return err + } + + g.config.SerialNumber = gatewayCert.SerialNumber + g.config.PrivateKey = gatewayCert.PrivateKey + g.config.Certificate = gatewayCert.Certificate + g.config.CertificateChain = gatewayCert.CertificateChain + + go func() { + err := relayNonTlsConn.CreatePermissions(peerAddr) + if err != nil { + log.Error().Msgf("Failed to refresh permission: %s", err) + } + log.Printf("Created permission for incoming connections") + ticker := time.NewTicker(2 * time.Minute) // Refresh before 5-min expiry + for range ticker.C { + err := relayNonTlsConn.CreatePermissions(peerAddr) + if err != nil { + log.Error().Msgf("Failed to refresh permission: %s", err) + } + } + }() + + cert, err := tls.X509KeyPair([]byte(gatewayCert.Certificate), []byte(gatewayCert.PrivateKey)) + if err != nil { + return fmt.Errorf("failed to parse cert: %s", err) + } + + fmt.Println(relayNonTlsConn.Addr().String()) + caCertPool := x509.NewCertPool() + caCertPool.AppendCertsFromPEM([]byte(gatewayCert.CertificateChain)) + + relayConn := tls.NewListener(relayNonTlsConn, &tls.Config{ + Certificates: []tls.Certificate{cert}, + MinVersion: tls.VersionTLS12, + ClientCAs: caCertPool, + ClientAuth: tls.RequireAndVerifyClientCert, + }) + + for { + log.Info().Msg("Connector started successfully") + // Accept new relay connection + conn, err := relayConn.Accept() + if err != nil { + log.Error().Msgf("Failed to accept connection: %v", err) + continue + } + + // Handle the connection in a goroutine + go handleConnection(conn) + } +}