diff --git a/README.md b/README.md index ee7d3ed..042732a 100644 --- a/README.md +++ b/README.md @@ -40,4 +40,6 @@ This will start Orchard Controller and a single Orchard Worker on your local mac You can interact with the newly created cluster using the `orchard` CLI or programmatically, through the built-in REST API server. +Controller-worker communication uses WebSocket-based RPC v2. Controllers that only support gRPC-based RPC v1 must be upgraded before starting a current worker. When upgrading an existing controller, first ensure its workers support RPC v2. The `--experimental-rpc-v2` and `--no-experimental-rpc-v2` flags have been removed. + Please check out the [official documentation](https://tart.run/orchard/quick-start/) for more information and/or feel free to use [issues](https://github.com/cirruslabs/orchard/issues) for the remaining questions. diff --git a/buf.gen.yaml b/buf.gen.yaml index e79d27e..54c1cf1 100644 --- a/buf.gen.yaml +++ b/buf.gen.yaml @@ -3,6 +3,3 @@ plugins: - plugin: go out: rpc opt: paths=source_relative - - plugin: go-grpc - out: rpc - opt: paths=source_relative diff --git a/go.mod b/go.mod index f6f1b91..7060a45 100644 --- a/go.mod +++ b/go.mod @@ -17,14 +17,12 @@ require ( github.com/gin-gonic/gin v1.12.0 github.com/go-openapi/runtime v0.29.3 github.com/gofrs/flock v0.13.0 - github.com/golang/protobuf v1.5.4 github.com/google/go-cmp v0.7.0 github.com/google/uuid v1.6.0 github.com/gosuri/uitable v0.0.4 github.com/hashicorp/go-multierror v1.1.1 github.com/hashicorp/go-version v1.8.0 github.com/manifoldco/promptui v0.9.0 - github.com/mitchellh/go-grpc-net-conn v0.0.0-20200427190222-eb030e4876f0 github.com/pkg/errors v0.9.1 github.com/pterm/pterm v0.12.83 github.com/puzpuzpuz/xsync/v4 v4.4.0 @@ -104,6 +102,7 @@ require ( github.com/gogo/protobuf v1.3.2 // indirect github.com/golang/glog v1.2.5 // indirect github.com/golang/groupcache v0.0.0-20210331224755-41bb18bfe9da // indirect + github.com/golang/protobuf v1.5.4 // indirect github.com/golang/snappy v0.0.4 // indirect github.com/google/flatbuffers v1.12.1 // indirect github.com/gookit/color v1.6.0 // indirect diff --git a/go.sum b/go.sum index 0edf211..b0eedd7 100644 --- a/go.sum +++ b/go.sum @@ -33,7 +33,6 @@ github.com/bytedance/sonic/loader v0.5.0 h1:gXH3KVnatgY7loH5/TkeVyXPfESoqSBSBEiD github.com/bytedance/sonic/loader v0.5.0/go.mod h1:AR4NYCk5DdzZizZ5djGqQ92eEhCCcdf5x77udYiSJRo= github.com/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1xcsSM= github.com/cenkalti/backoff/v5 v5.0.3/go.mod h1:rkhZdG3JZukswDf7f0cwqPNk4K0sa+F97BxZthm/crw= -github.com/census-instrumentation/opencensus-proto v0.2.1/go.mod h1:f6KPmirojxKA12rnyqOA5BBL4O983OfeGPqjHWSTneU= github.com/cespare/xxhash v1.1.0 h1:a6HrQnmkObjyL+Gs60czilIUGqrzKutQD6XZog3p+ko= github.com/cespare/xxhash v1.1.0/go.mod h1:XrSqR1VqqWfGrhpAt58auRo0WTKS1nRRg3ghfAqPWnc= github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= @@ -52,7 +51,6 @@ github.com/clipperhouse/uax29/v2 v2.7.0 h1:+gs4oBZ2gPfVrKPthwbMzWZDaAFPGYK72F0NJ github.com/clipperhouse/uax29/v2 v2.7.0/go.mod h1:EFJ2TJMRUaplDxHKj1qAEhCtQPW2tJSwu5BF98AuoVM= github.com/cloudwego/base64x v0.1.6 h1:t11wG9AECkCDk5fMSoxmufanudBtJ+/HemLstXDLI2M= github.com/cloudwego/base64x v0.1.6/go.mod h1:OFcloc187FXDaYHvrNIjxSe8ncn0OOM8gEHfghB2IPU= -github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGXZJjfX53e64911xZQV5JYwmTeXPW+k8Sc= github.com/coder/websocket v1.8.14 h1:9L0p0iKiNOibykf283eHkKUHHrpG7f65OE3BhhO7v9g= github.com/coder/websocket v1.8.14/go.mod h1:NX3SzP+inril6yawo5CQXx8+fk145lPDC6pumgx0mVg= github.com/containerd/console v1.0.3/go.mod h1:7LqA/THxQ86k76b8c/EMSiaJ3h1eZkMkXar0TQ1gf3U= @@ -80,9 +78,6 @@ github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkp github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= github.com/ebitengine/purego v0.10.0 h1:QIw4xfpWT6GWTzaW5XEKy3HXoqrJGx1ijYHzTF0/ISU= github.com/ebitengine/purego v0.10.0/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= -github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= -github.com/envoyproxy/go-control-plane v0.9.4/go.mod h1:6rpuAdCZL397s3pYoYcLgu1mIlRU8Am5FuJP05cCM98= -github.com/envoyproxy/protoc-gen-validate v0.1.0/go.mod h1:iSmxcyjqTsJpI2R4NaDN7+kN2VEUnK/pcBlmesArF7c= github.com/fatih/color v1.18.0 h1:S8gINlzdQ840/4pfAwic/ZE0djQEH3wM94VfqLTZcOM= github.com/fatih/color v1.18.0/go.mod h1:4FelSpRwEGDpQ12mAdzqdOukCy4u8WUtOY6lkT/6HfU= github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo= @@ -170,8 +165,6 @@ github.com/golang/groupcache v0.0.0-20210331224755-41bb18bfe9da/go.mod h1:cIg4er github.com/golang/mock v1.1.1/go.mod h1:oTYuIxOrZwtPieC+H1uAHpcLFnEyAGVDL/k47Jfbm0A= github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.3.3/go.mod h1:vzj43D7+SQXF/4pzW/hwtAqwc6iTitCiVSaWz5lYuqw= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= github.com/golang/snappy v0.0.3/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= @@ -179,7 +172,6 @@ github.com/golang/snappy v0.0.4 h1:yAGX7huGHXlcLOEtBnF4w7FQwA26wojNCwOYAEhLjQM= github.com/golang/snappy v0.0.4/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= github.com/google/flatbuffers v1.12.1 h1:MVlul7pQNoDzWRLTw5imwYsl+usrS1TXG2H4jg6ImGw= github.com/google/flatbuffers v1.12.1/go.mod h1:1AeVuKshWv4vARoZatz6mlQ0JxURH0Kv5+zNeJKJCa8= -github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M= github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.5.4/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.6/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= @@ -245,8 +237,6 @@ github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D github.com/mattn/go-runewidth v0.0.13/go.mod h1:Jdepj2loyihRzMpdS35Xk/zdY8IAYHsh153qUoGf23w= github.com/mattn/go-runewidth v0.0.20 h1:WcT52H91ZUAwy8+HUkdM3THM6gXqXuLJi9O3rjcQQaQ= github.com/mattn/go-runewidth v0.0.20/go.mod h1:XBkDxAl56ILZc9knddidhrOlY5R/pDhgLpndooCuJAs= -github.com/mitchellh/go-grpc-net-conn v0.0.0-20200427190222-eb030e4876f0 h1:oZuel4h7224ILBLg2SlTxdaMYXDyqcVfL4Cg1PJQHZs= -github.com/mitchellh/go-grpc-net-conn v0.0.0-20200427190222-eb030e4876f0/go.mod h1:ZCzL0JMR6qfm7VrDC8HGwVtPA8D2Ijc/edUSBw58x94= github.com/mitchellh/go-homedir v1.1.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0= github.com/mitchellh/mapstructure v1.1.2/go.mod h1:FVVH3fgwuzCH5S8UJGiWEs2h04kUh9fWfEaFds41c1Y= github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= @@ -267,7 +257,6 @@ github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRI github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 h1:o4JXh1EVt9k/+g42oCprj/FisM4qX9L3sZB3upGN2ZU= github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55/go.mod h1:OmDBASR4679mdNQnz2pUhc2G8CO2JrUAVFDRBDP/hJE= -github.com/prometheus/client_model v0.0.0-20190812154241-14fe0d1b01d4/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA= github.com/pterm/pterm v0.12.27/go.mod h1:PhQ89w4i95rhgE+xedAoqous6K9X+r6aSOI2eFF7DZI= github.com/pterm/pterm v0.12.29/go.mod h1:WI3qxgvoQFFGKGjGnJR849gU0TsEOvKn5Q8LlY1U7lg= github.com/pterm/pterm v0.12.30/go.mod h1:MOqLIyMOgmTDz9yorcYbcw+HsgoZo3BQfg2wtl3HEFE= @@ -320,7 +309,6 @@ github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/ github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= -github.com/stretchr/testify v1.5.1/go.mod h1:5W2xD1RspED5o8YsWQXVCued0rvSQ+mT+I5cxcmMvtA= github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= @@ -470,7 +458,6 @@ golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGm golang.org/x/tools v0.0.0-20190114222345-bf090417da8b/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.0.0-20190226205152-f727befe758c/go.mod h1:9Yl7xja0Znq3iFh3HoIrodX9oNMXvdceNzlUR8zjMvY= golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= -golang.org/x/tools v0.0.0-20190524140312-2c0ae7006135/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA= @@ -486,14 +473,10 @@ google.golang.org/appengine v1.1.0/go.mod h1:EbEs0AVv82hx2wNQdGPgUI5lhzA/G0D9Ywl google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= google.golang.org/genproto v0.0.0-20180817151627-c66870c02cf8/go.mod h1:JiN7NxoALGmiZfu7CAH4rXhgtRTLTxftemlI0sWmxmc= google.golang.org/genproto v0.0.0-20190425155659-357c62f0e4bb/go.mod h1:VzzqZJRnGkLBvHegQrXjBqPurQTc5/KpmUdxsrq26oE= -google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc= google.golang.org/genproto v0.0.0-20230410155749-daa745c078e1 h1:KpwkzHKEF7B9Zxg18WzOa7djJ+Ha5DzthMyZYQfEn2A= google.golang.org/genproto v0.0.0-20230410155749-daa745c078e1/go.mod h1:nKE/iIaLqn2bQwXBg8f1g2Ylh6r5MN5CmZvuzZCgsCU= google.golang.org/grpc v1.19.0/go.mod h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c= google.golang.org/grpc v1.20.1/go.mod h1:10oTOabMzJvdu6/UiuZezV6QK5dSlG84ov/aaiqXj38= -google.golang.org/grpc v1.23.0/go.mod h1:Y5yQAOtifL1yxbo5wqy6BxZv8vAUGQwXBOALyacEbxg= -google.golang.org/grpc v1.25.1/go.mod h1:c3i+UQWmh7LiEpx4sFZnkU36qjEYZ0imhYfXVyQciAY= -google.golang.org/grpc v1.28.1/go.mod h1:rpkK4SK4GF4Ach/+MFLZUBavHOvF2JJB5uozKKal+60= google.golang.org/grpc v1.79.3 h1:sybAEdRIEtvcD68Gx7dmnwjZKlyfuc61Dyo9pGXXkKE= google.golang.org/grpc v1.79.3/go.mod h1:KmT0Kjez+0dde/v2j9vzwoAScgEPx/Bw1CYChhHLrHQ= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= @@ -512,6 +495,5 @@ gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/go.mod h1:K4uyk7z7BCEPqu6E+C gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= honnef.co/go/tools v0.0.0-20190102054323-c2f93a96b099/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= -honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= howett.net/plist v1.0.1 h1:37GdZ8tP09Q35o9ych3ehygcsL+HqKSwzctveSlarvM= howett.net/plist v1.0.1/go.mod h1:lqaXoTrLY4hg8tnEzNru53gicrbv7rrk+2xJA/7hw9g= diff --git a/internal/command/controller/run.go b/internal/command/controller/run.go index 8c37222..8600d1f 100644 --- a/internal/command/controller/run.go +++ b/internal/command/controller/run.go @@ -31,8 +31,6 @@ var debug bool var noTLS bool var sshNoClientAuth bool var insecureAllowHostDirs bool -var experimentalRPCV2 bool -var noExperimentalRPCV2 bool var experimentalPingInterval time.Duration var experimentalDisableDBCompression bool var workerOfflineTimeout time.Duration @@ -77,11 +75,6 @@ func newRunCommand() *cobra.Command { "thus only authenticating on the target worker/VM's SSH server") cmd.Flags().BoolVar(&insecureAllowHostDirs, "insecure-allow-host-dirs", false, "allow unsafe path-based local host directory sharing") - cmd.Flags().BoolVar(&experimentalRPCV2, "experimental-rpc-v2", false, - "enable experimental RPC v2 (https://github.com/cirruslabs/orchard/issues/235)") - _ = cmd.Flags().MarkHidden("experimental-rpc-v2") - cmd.Flags().BoolVar(&noExperimentalRPCV2, "no-experimental-rpc-v2", false, - "disable experimental RPC v2 (https://github.com/cirruslabs/orchard/issues/235)") cmd.Flags().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 "+ @@ -186,9 +179,7 @@ func runController(cmd *cobra.Command, args []string) (err error) { Certificates: []tls.Certificate{ controllerCert, }, - // Since gRPC clients started enforcing ALPN at some point, we need to advertise it - // - // See https://github.com/grpc/grpc-go/issues/7922 for more details. + // Advertise the HTTP protocols supported by the controller. NextProtos: []string{"http/1.1", "h2"}, })) } @@ -202,18 +193,6 @@ func runController(cmd *cobra.Command, args []string) (err error) { controllerOpts = append(controllerOpts, controller.WithSSHServer(addressSSH, signer, sshNoClientAuth)) } - if experimentalRPCV2 && noExperimentalRPCV2 { - return fmt.Errorf("--experimental-rpc-v2 and --no-experimental-rpc-v2 flags are mutually exclusive") - } - - if experimentalRPCV2 { - logger.Warn("--experimental-rpc-v2 flag is deprecated: experimental RPC v2 is now enabled by default") - } - - if !noExperimentalRPCV2 { - 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") diff --git a/internal/command/dev/dev.go b/internal/command/dev/dev.go index 632181f..e3d04c8 100644 --- a/internal/command/dev/dev.go +++ b/internal/command/dev/dev.go @@ -31,7 +31,6 @@ var ErrFailed = errors.New("failed to run development controller and worker") var devDataDirPath string var apiPrefix string var stringToStringResources map[string]string -var experimentalRPCV2 bool var addressPprof string var synthetic bool var workers int @@ -51,8 +50,6 @@ func NewCommand() *cobra.Command { "behind an HTTP proxy together with other services") command.Flags().StringToStringVar(&stringToStringResources, "resources", map[string]string{}, "resources that the development worker will provide") - command.Flags().BoolVar(&experimentalRPCV2, "experimental-rpc-v2", false, - "enable experimental RPC v2 (https://github.com/cirruslabs/orchard/issues/235)") command.Flags().StringVar(&addressPprof, "listen-pprof", "", "start pprof HTTP server on localhost:6060 for diagnostic purposes (e.g. \"localhost:6060\")") command.Flags().BoolVar(&synthetic, "synthetic", false, @@ -104,10 +101,6 @@ func runDev(cmd *cobra.Command, args []string) error { additionalControllerOpts = append(additionalControllerOpts, controller.WithAPIPrefix(apiPrefix)) } - if experimentalRPCV2 { - additionalControllerOpts = append(additionalControllerOpts, controller.WithExperimentalRPCV2()) - } - if insecureAllowHostDirs { additionalControllerOpts = append(additionalControllerOpts, controller.WithInsecureAllowHostDirs()) } diff --git a/internal/controller/api.go b/internal/controller/api.go index dd15dd8..821461d 100644 --- a/internal/controller/api.go +++ b/internal/controller/api.go @@ -1,7 +1,6 @@ package controller import ( - "context" "crypto/subtle" "errors" "net/http" @@ -12,7 +11,6 @@ import ( storepkg "github.com/cirruslabs/orchard/internal/controller/store" "github.com/cirruslabs/orchard/internal/responder" v1pkg "github.com/cirruslabs/orchard/pkg/resource/v1" - "github.com/cirruslabs/orchard/rpc" "github.com/deckarep/golang-set/v2" ginzap "github.com/gin-contrib/zap" "github.com/gin-gonic/gin" @@ -20,7 +18,6 @@ import ( "github.com/samber/lo" "go.opentelemetry.io/contrib/instrumentation/github.com/gin-gonic/gin/otelgin" "go.uber.org/zap" - "google.golang.org/grpc/metadata" ) const ctxServiceAccountKey = "service-account" @@ -311,28 +308,6 @@ func (controller *Controller) authorizeBase( NewErrorResponse("%s: %s", hint, strings.Join(humanizedRoles, ", "))) } -func (controller *Controller) authorizeGRPC(ctx context.Context, scopes ...v1pkg.ServiceAccountRole) bool { - if controller.insecureAuthDisabled { - return true - } - - name := metadata.ValueFromIncomingContext(ctx, rpc.MetadataServiceAccountNameKey) - if len(name) != 1 { - return false - } - token := metadata.ValueFromIncomingContext(ctx, rpc.MetadataServiceAccountTokenKey) - if len(token) != 1 { - return false - } - - serviceAccount, err := controller.fetchServiceAccount(name[0], token[0]) - if err != nil { - return false - } - - return mapset.NewSet[v1pkg.ServiceAccountRole](serviceAccount.Roles...).Contains(scopes...) -} - type storeTransactionFunc func(operation func(txn storepkg.Transaction) error) error func (controller *Controller) storeView(view func(txn storepkg.Transaction) responder.Responder) responder.Responder { diff --git a/internal/controller/api_controller.go b/internal/controller/api_controller.go index b67ca2a..af6da5c 100644 --- a/internal/controller/api_controller.go +++ b/internal/controller/api_controller.go @@ -17,14 +17,10 @@ func (controller *Controller) controllerInfo(ctx *gin.Context) responder.Respond } capabilities := []v1pkg.ControllerCapability{ - v1pkg.ControllerCapabilityRPCV1, + v1pkg.ControllerCapabilityRPCV2, v1pkg.ControllerCapabilityVMStateEndpoint, } - if controller.experimentalRPCV2 { - capabilities = append(capabilities, v1pkg.ControllerCapabilityRPCV2) - } - return responder.JSON(http.StatusOK, &v1pkg.ControllerInfo{ Version: version.Version, Commit: version.Commit, diff --git a/internal/controller/controller.go b/internal/controller/controller.go index 9c55807..0344af2 100644 --- a/internal/controller/controller.go +++ b/internal/controller/controller.go @@ -20,7 +20,6 @@ import ( "github.com/cirruslabs/orchard/internal/netconstants" "github.com/cirruslabs/orchard/internal/opentelemetry" v1 "github.com/cirruslabs/orchard/pkg/resource/v1" - "github.com/cirruslabs/orchard/rpc" "github.com/samber/lo" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/metric" @@ -29,8 +28,6 @@ import ( "golang.org/x/net/http2" "golang.org/x/net/http2/h2c" "golang.org/x/sync/singleflight" - "google.golang.org/grpc" - "google.golang.org/grpc/keepalive" ) var ( @@ -50,7 +47,6 @@ type Controller struct { scheduler *scheduler.Scheduler store storepkg.Store logger *zap.SugaredLogger - grpcServer *grpc.Server workerNotifier *notifier.Notifier connRendezvous *rendezvous.Rendezvous[rendezvous.ResultWithErrorMessage[net.Conn]] ipRendezvous *rendezvous.Rendezvous[rendezvous.ResultWithErrorMessage[string]] @@ -58,7 +54,6 @@ type Controller struct { workerOfflineTimeout time.Duration execSessionRetentionTTL time.Duration execSSHConnectionKeepaliveInterval time.Duration - experimentalRPCV2 bool disableDBCompression bool pingInterval time.Duration synthetic bool @@ -71,8 +66,6 @@ type Controller struct { execSSHClients *execSSHClientPool single singleflight.Group - - rpc.UnimplementedControllerServer } func New(opts ...Option) (*Controller, error) { @@ -148,23 +141,8 @@ func New(opts ...Option) (*Controller, error) { apiServer := controller.initAPI() - controller.grpcServer = grpc.NewServer( - grpc.KeepaliveParams(keepalive.ServerParameters{ - Time: 30 * time.Second, - }), - ) - rpc.RegisterControllerServer(controller.grpcServer, controller) - - handler := http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) { - if request.Header.Get("Content-Type") == "application/grpc" { - controller.grpcServer.ServeHTTP(writer, request) - } else { - apiServer.ServeHTTP(writer, request) - } - }) - controller.httpServer = &http.Server{ - Handler: h2c.NewHandler(handler, &http2.Server{}), + Handler: h2c.NewHandler(apiServer, &http2.Server{}), ReadHeaderTimeout: 60 * time.Second, } diff --git a/internal/controller/option.go b/internal/controller/option.go index e29b99a..b8854d0 100644 --- a/internal/controller/option.go +++ b/internal/controller/option.go @@ -78,12 +78,6 @@ func WithExecSSHConnectionKeepaliveInterval(execSSHConnectionKeepaliveInterval t } } -func WithExperimentalRPCV2() Option { - return func(controller *Controller) { - controller.experimentalRPCV2 = true - } -} - func WithDisableDBCompression() Option { return func(controller *Controller) { controller.disableDBCompression = true diff --git a/internal/controller/rpc.go b/internal/controller/rpc.go deleted file mode 100644 index d3edc33..0000000 --- a/internal/controller/rpc.go +++ /dev/null @@ -1,104 +0,0 @@ -package controller - -import ( - "context" - "github.com/cirruslabs/orchard/internal/controller/rendezvous" - v1pkg "github.com/cirruslabs/orchard/pkg/resource/v1" - "github.com/cirruslabs/orchard/rpc" - "google.golang.org/grpc/metadata" - "google.golang.org/protobuf/types/known/emptypb" - "net" - - //nolint:staticcheck // https://github.com/mitchellh/go-grpc-net-conn/pull/1 - "github.com/golang/protobuf/proto" - grpc_net_conn "github.com/mitchellh/go-grpc-net-conn" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" -) - -func (controller *Controller) Watch(_ *emptypb.Empty, stream rpc.Controller_WatchServer) error { - if !controller.authorizeGRPC(stream.Context(), v1pkg.ServiceAccountRoleComputeWrite) { - return status.Errorf(codes.Unauthenticated, "auth failed") - } - - workerMetadataValue := metadata.ValueFromIncomingContext(stream.Context(), rpc.MetadataWorkerNameKey) - if len(workerMetadataValue) == 0 { - return status.Errorf(codes.InvalidArgument, "no worker ident in metadata") - } - - worker := workerMetadataValue[0] - workerCh, cancel := controller.workerNotifier.Register(stream.Context(), worker) - defer cancel() - - for { - select { - case msg := <-workerCh: - if err := stream.Send(msg); err != nil { - return err - } - case <-stream.Context().Done(): - return stream.Context().Err() - } - } -} - -func (controller *Controller) PortForward(stream rpc.Controller_PortForwardServer) error { - if !controller.authorizeGRPC(stream.Context(), v1pkg.ServiceAccountRoleComputeWrite) { - return status.Errorf(codes.Unauthenticated, "auth failed") - } - - sessionMetadataValue := metadata.ValueFromIncomingContext(stream.Context(), rpc.MetadataWorkerPortForwardingSessionKey) - if len(sessionMetadataValue) == 0 { - return status.Errorf(codes.InvalidArgument, "no session in metadata") - } - - conn := &grpc_net_conn.Conn{ - Stream: stream, - Request: &rpc.PortForwardData{}, - Response: &rpc.PortForwardData{}, - Encode: grpc_net_conn.SimpleEncoder(func(message proto.Message) *[]byte { - return &message.(*rpc.PortForwardData).Data - }), - Decode: grpc_net_conn.SimpleDecoder(func(message proto.Message) *[]byte { - return &message.(*rpc.PortForwardData).Data - }), - } - - // make connection rendezvous aware of the connection - proxyCtx, err := controller.connRendezvous.Respond(sessionMetadataValue[0], - rendezvous.ResultWithErrorMessage[net.Conn]{ - Result: conn, - }, - ) - if err != nil { - return err - } - - select { - case <-proxyCtx.Done(): - return proxyCtx.Err() - case <-stream.Context().Done(): - return stream.Context().Err() - } -} - -func (controller *Controller) ResolveIP(ctx context.Context, request *rpc.ResolveIPResult) (*emptypb.Empty, error) { - if !controller.authorizeGRPC(ctx, v1pkg.ServiceAccountRoleComputeWrite) { - return nil, status.Errorf(codes.Unauthenticated, "auth failed") - } - - sessionMetadataValue := metadata.ValueFromIncomingContext(ctx, rpc.MetadataWorkerPortForwardingSessionKey) - if len(sessionMetadataValue) == 0 { - return nil, status.Errorf(codes.InvalidArgument, "no session in metadata") - } - - // Respond with the resolved IP address - _, err := controller.ipRendezvous.Respond(sessionMetadataValue[0], rendezvous.ResultWithErrorMessage[string]{ - Result: request.Ip, - }) - if err != nil { - return nil, err - } - - return &emptypb.Empty{}, nil -} diff --git a/internal/tests/devcontroller/devcontroller.go b/internal/tests/devcontroller/devcontroller.go index eb00031..22f0f2b 100644 --- a/internal/tests/devcontroller/devcontroller.go +++ b/internal/tests/devcontroller/devcontroller.go @@ -29,9 +29,6 @@ func StartIntegrationTestEnvironmentWithAdditionalOpts( ) (*client.Client, *controller.Controller, *worker.Worker) { t.Setenv("ORCHARD_HOME", t.TempDir()) - // Enable experimental RPC v2 by default in tests - additionalControllerOpts = append(additionalControllerOpts, controller.WithExperimentalRPCV2()) - devController, devWorker, err := dev.CreateDevControllerAndWorker(t.TempDir(), ":0", nil, additionalControllerOpts, additionalWorkerOpts) require.NoError(t, err) diff --git a/internal/worker/rpc.go b/internal/worker/rpc.go deleted file mode 100644 index f3a87ba..0000000 --- a/internal/worker/rpc.go +++ /dev/null @@ -1,190 +0,0 @@ -package worker - -import ( - "context" - "fmt" - "time" - - "github.com/cirruslabs/orchard/internal/proxy" - "github.com/cirruslabs/orchard/internal/worker/vmmanager" - "github.com/cirruslabs/orchard/rpc" - "google.golang.org/grpc/keepalive" - "google.golang.org/protobuf/types/known/emptypb" - - "net" - - //nolint:staticcheck // https://github.com/mitchellh/go-grpc-net-conn/pull/1 - "github.com/golang/protobuf/proto" - grpc_net_conn "github.com/mitchellh/go-grpc-net-conn" - "google.golang.org/grpc" - "google.golang.org/grpc/metadata" - - "github.com/samber/lo" -) - -func (worker *Worker) watchRPC(ctx context.Context, operationCtx context.Context, onEstablished func()) error { - worker.logger.Infof("connecting to %s over gRPC", worker.client.GRPCTarget()) - - conn, err := grpc.NewClient(worker.client.GRPCTarget(), - grpc.WithTransportCredentials(worker.client.GRPCTransportCredentials()), - grpc.WithKeepaliveParams(keepalive.ClientParameters{ - Time: 30 * time.Second, - }), - ) - if err != nil { - return err - } - - worker.logger.Infof("gRPC connection established, starting gRPC stream with the controller") - - client := rpc.NewControllerClient(conn) - - ctxWithMetadata := metadata.NewOutgoingContext(ctx, worker.grpcMetadata()) - operationCtxWithMetadata := metadata.NewOutgoingContext(operationCtx, worker.grpcMetadata()) - - stream, err := client.Watch(ctxWithMetadata, &emptypb.Empty{}) - if err != nil { - return err - } - onEstablished() - - worker.logger.Infof("running gRPC stream with the controller") - - for { - watchFromController, err := stream.Recv() - if err != nil { - return err - } - - switch action := watchFromController.Action.(type) { - case *rpc.WatchInstruction_PortForwardAction: - go worker.handlePortForward(operationCtxWithMetadata, client, action.PortForwardAction) - case *rpc.WatchInstruction_SyncVmsAction: - worker.requestVMSyncing() - case *rpc.WatchInstruction_ResolveIpAction: - go worker.handleGetIP(operationCtxWithMetadata, client, action.ResolveIpAction) - } - } -} - -func (worker *Worker) handlePortForward( - ctx context.Context, - client rpc.ControllerClient, - portForwardAction *rpc.WatchInstruction_PortForward, -) { - subCtx, cancel := context.WithCancel(ctx) - defer cancel() - - grpcMetadata := metadata.Join( - worker.grpcMetadata(), - metadata.Pairs(rpc.MetadataWorkerPortForwardingSessionKey, portForwardAction.Session), - ) - ctxWithMetadata := metadata.NewOutgoingContext(subCtx, grpcMetadata) - stream, err := client.PortForward(ctxWithMetadata) - if err != nil { - worker.logger.Warnf("port forwarding failed: failed to call PortForward() RPC method: %v", err) - - return - } - - var host string - - if portForwardAction.VmUid == "" { - // Port-forwarding request to a worker - host = "localhost" - } else { - // Port-forwarding request to a VM, find that VM - vm, ok := lo.Find(worker.vmm.List(), func(item vmmanager.VM) bool { - return item.Resource().UID == portForwardAction.VmUid - }) - if !ok { - worker.logger.Warnf("port forwarding failed: failed to get the VM: %v", err) - - return - } - - // Obtain VM's IP address - host, err = vm.IP(ctx) - if err != nil { - worker.logger.Warnf("port forwarding failed: failed to get VM's IP: %v", err) - - return - } - } - - // Connect to the VM's port - var vmConn net.Conn - - if worker.dialer != nil { - vmConn, err = worker.dialer.DialContext(ctx, "tcp", - fmt.Sprintf("%s:%d", host, portForwardAction.Port)) - } else { - dialer := net.Dialer{} - - vmConn, err = dialer.DialContext(ctx, "tcp", - fmt.Sprintf("%s:%d", host, portForwardAction.Port)) - } - if err != nil { - worker.logger.Warnf("port forwarding failed: failed to connect to the VM: %v", err) - - return - } - - // Proxy bytes - grpcConn := &grpc_net_conn.Conn{ - Stream: stream, - Request: &rpc.PortForwardData{}, - Response: &rpc.PortForwardData{}, - Encode: grpc_net_conn.SimpleEncoder(func(message proto.Message) *[]byte { - return &message.(*rpc.PortForwardData).Data - }), - Decode: grpc_net_conn.SimpleDecoder(func(message proto.Message) *[]byte { - return &message.(*rpc.PortForwardData).Data - }), - } - - _ = proxy.Connections(vmConn, grpcConn) -} - -func (worker *Worker) handleGetIP( - ctx context.Context, - client rpc.ControllerClient, - resolveIP *rpc.WatchInstruction_ResolveIP, -) { - grpcMetadata := metadata.Join( - worker.grpcMetadata(), - metadata.Pairs(rpc.MetadataWorkerPortForwardingSessionKey, resolveIP.Session), - ) - ctxWithMetadata := metadata.NewOutgoingContext(ctx, grpcMetadata) - - // Find the desired VM - vm, ok := lo.Find(worker.vmm.List(), func(item vmmanager.VM) bool { - return item.Resource().UID == resolveIP.VmUid - }) - if !ok { - worker.logger.Warnf("failed to resolve IP for the VM with UID %q: VM not found", - resolveIP.VmUid) - - return - } - - // Obtain VM's IP address - ip, err := vm.IP(ctx) - if err != nil { - worker.logger.Warnf("failed to resolve IP for the VM with UID %q: \"tart ip\" failed: %v", - resolveIP.VmUid, err) - - return - } - - _, err = client.ResolveIP(ctxWithMetadata, &rpc.ResolveIPResult{ - Session: resolveIP.Session, - Ip: ip, - }) - if err != nil { - worker.logger.Warnf("failed to resolve IP for the VM with UID %q: "+ - "failed to call back to the controller: %v", resolveIP.VmUid, err) - - return - } -} diff --git a/internal/worker/worker.go b/internal/worker/worker.go index fda971e..f786588 100644 --- a/internal/worker/worker.go +++ b/internal/worker/worker.go @@ -20,7 +20,6 @@ import ( "github.com/cirruslabs/orchard/internal/worker/vmmanager" "github.com/cirruslabs/orchard/pkg/client" v1 "github.com/cirruslabs/orchard/pkg/resource/v1" - "github.com/cirruslabs/orchard/rpc" mapset "github.com/deckarep/golang-set/v2" "github.com/dustin/go-humanize" "github.com/hashicorp/go-multierror" @@ -31,7 +30,6 @@ import ( "go.opentelemetry.io/otel/metric" "go.uber.org/zap" "golang.org/x/sync/errgroup" - "google.golang.org/grpc/metadata" ) const ( @@ -247,7 +245,7 @@ func (worker *Worker) runNewSession(ctx context.Context, onWatchHealthy func()) } group, sessionCtx := errgroup.WithContext(subCtx) - worker.superviseRPCWatch(sessionCtx, ctx, group, info, onWatchHealthy) + worker.superviseRPCWatch(sessionCtx, ctx, group, onWatchHealthy) // Sync on-disk VMs if err := worker.syncOnDiskVMsWithInventory(sessionCtx, vmInfos); err != nil { @@ -329,32 +327,22 @@ func (worker *Worker) superviseRPCWatch( sessionCtx context.Context, operationCtx context.Context, group *errgroup.Group, - info v1.ControllerInfo, onWatchHealthy func(), ) { - watchRPC := worker.watchRPC - rpcVersion := "v1" - - if info.Capabilities.Has(v1.ControllerCapabilityRPCV2) { - worker.logger.Infof("using WebSocket-based v2 RPC") - watchRPC = worker.watchRPCV2 - rpcVersion = "v2" - } else { - worker.logger.Infof("using gRPC-based v1 RPC") - } + worker.logger.Infof("using WebSocket-based v2 RPC") watchEstablished := make(chan struct{}) group.Go(func() error { - if err := watchRPC(sessionCtx, operationCtx, func() { close(watchEstablished) }); err != nil { + if err := worker.watchRPCV2(sessionCtx, operationCtx, func() { close(watchEstablished) }); err != nil { if sessionCtx.Err() != nil { return sessionCtx.Err() } - return fmt.Errorf("%w: failed to watch RPC %s: %w", errRPCWatchDisconnected, rpcVersion, err) + return fmt.Errorf("%w: failed to watch RPC v2: %w", errRPCWatchDisconnected, err) } - return fmt.Errorf("%w: RPC %s watch closed unexpectedly", errRPCWatchDisconnected, rpcVersion) + return fmt.Errorf("%w: RPC v2 watch closed unexpectedly", errRPCWatchDisconnected) }) group.Go(func() error { @@ -839,13 +827,6 @@ func (worker *Worker) createVM(odn ondiskname.OnDiskName, vmResource v1.VM) { worker.vmm.Put(odn, vm) } -func (worker *Worker) grpcMetadata() metadata.MD { - return metadata.Join( - worker.client.GPRCMetadata(), - metadata.Pairs(rpc.MetadataWorkerNameKey, worker.name), - ) -} - func (worker *Worker) requestVMSyncing() { select { case worker.syncRequested <- true: diff --git a/internal/worker/worker_test.go b/internal/worker/worker_test.go index 4af4ffb..ae9623b 100644 --- a/internal/worker/worker_test.go +++ b/internal/worker/worker_test.go @@ -19,18 +19,11 @@ import ( "github.com/cirruslabs/orchard/internal/worker/vmmanager/tart" "github.com/cirruslabs/orchard/pkg/client" v1 "github.com/cirruslabs/orchard/pkg/resource/v1" - "github.com/cirruslabs/orchard/rpc" "github.com/coder/websocket" "github.com/samber/lo" "github.com/samber/mo" "github.com/stretchr/testify/require" "go.uber.org/zap" - "golang.org/x/net/http2" - "golang.org/x/net/http2/h2c" - "google.golang.org/grpc" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" - "google.golang.org/protobuf/types/known/emptypb" ) const ( @@ -420,35 +413,24 @@ func TestMonitorRPCWatchHealth(t *testing.T) { func TestWorkerBacksOffPersistentRPCWatchFailures(t *testing.T) { testCases := []struct { name string - useRPCV2 bool closeAfter bool }{ - {name: "HTTP rejection", useRPCV2: true}, - {name: "WebSocket closes immediately", useRPCV2: true, closeAfter: true}, - {name: "gRPC first receive fails"}, + {name: "HTTP rejection"}, + {name: "WebSocket closes immediately", closeAfter: true}, } for _, testCase := range testCases { t.Run(testCase.name, func(t *testing.T) { - testWorkerBacksOffPersistentRPCWatchFailure(t, testCase.useRPCV2, testCase.closeAfter) + testWorkerBacksOffPersistentRPCWatchFailure(t, testCase.closeAfter) }) } } -func testWorkerBacksOffPersistentRPCWatchFailure(t *testing.T, useRPCV2 bool, closeAfter bool) { +func testWorkerBacksOffPersistentRPCWatchFailure(t *testing.T, closeAfter bool) { t.Helper() watchAttempts := make(chan time.Time, 4) - grpcServer := grpc.NewServer() - rpc.RegisterControllerServer(grpcServer, &failingRecoveryRPCServer{watchAttempts: watchAttempts}) - t.Cleanup(grpcServer.Stop) - handler := http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) { - if request.Header.Get("Content-Type") == "application/grpc" { - grpcServer.ServeHTTP(writer, request) - return - } - switch { case request.Method == http.MethodPost && request.URL.Path == "/v1/workers": var workerResource v1.Worker @@ -458,11 +440,9 @@ func testWorkerBacksOffPersistentRPCWatchFailure(t *testing.T, useRPCV2 bool, cl } writeRecoveryTestJSON(t, writer, workerResource) case request.Method == http.MethodGet && request.URL.Path == recoveryTestInfoPath: - info := v1.ControllerInfo{} - if useRPCV2 { - info.Capabilities = v1.ControllerCapabilities{v1.ControllerCapabilityRPCV2} - } - writeRecoveryTestJSON(t, writer, info) + writeRecoveryTestJSON(t, writer, v1.ControllerInfo{ + Capabilities: v1.ControllerCapabilities{v1.ControllerCapabilityRPCV2}, + }) case request.Method == http.MethodGet && request.URL.Path == recoveryTestWorkerPath: writeRecoveryTestJSON(t, writer, v1.Worker{Meta: v1.Meta{Name: recoveryTestWorkerName}}) case request.Method == http.MethodPut && request.URL.Path == recoveryTestWorkerPath: @@ -491,7 +471,7 @@ func testWorkerBacksOffPersistentRPCWatchFailure(t *testing.T, useRPCV2 bool, cl http.NotFound(writer, request) } }) - controller := httptest.NewServer(h2c.NewHandler(handler, &http2.Server{})) + controller := httptest.NewServer(handler) t.Cleanup(controller.Close) controllerClient, err := client.New(client.WithAddress(controller.URL)) @@ -528,24 +508,6 @@ func testWorkerBacksOffPersistentRPCWatchFailure(t *testing.T, useRPCV2 bool, cl require.ErrorIs(t, <-runResult, context.Canceled) } -type failingRecoveryRPCServer struct { - rpc.UnimplementedControllerServer - - watchAttempts chan time.Time -} - -func (server *failingRecoveryRPCServer) Watch( - _ *emptypb.Empty, - _ rpc.Controller_WatchServer, -) error { - select { - case server.watchAttempts <- time.Now(): - default: - } - - return status.Error(codes.Unavailable, "RPC watch upstream unavailable") -} - func TestShouldPreserveRecoveredVM(t *testing.T) { onDiskName := ondiskname.New("running-vm", recoveryTestVMUID, 0) now := time.Unix(1_000, 0) diff --git a/pkg/client/client.go b/pkg/client/client.go index a43cd52..4b7b6c9 100644 --- a/pkg/client/client.go +++ b/pkg/client/client.go @@ -19,11 +19,7 @@ import ( "github.com/cirruslabs/orchard/internal/config" "github.com/cirruslabs/orchard/internal/dialer" "github.com/cirruslabs/orchard/internal/version" - "github.com/cirruslabs/orchard/rpc" "github.com/coder/websocket" - "google.golang.org/grpc/credentials" - "google.golang.org/grpc/credentials/insecure" - "google.golang.org/grpc/metadata" ) const ( @@ -124,7 +120,7 @@ func New(opts ...Option) (*Client, error) { client.baseURL = url // Figure out if HTTP (insecure) or HTTPS (secure) was requested, - // so we can further adapt for gRPC and WebSocket usage patterns + // so we can select the corresponding WebSocket scheme switch client.baseURL.Scheme { case "http": client.insecure = true @@ -138,31 +134,6 @@ func New(opts ...Option) (*Client, error) { return client, nil } -func (client *Client) GRPCTarget() string { - return client.baseURL.Host -} - -func (client *Client) GRPCTransportCredentials() credentials.TransportCredentials { - if client.insecure { - return insecure.NewCredentials() - } - - return credentials.NewTLS(client.tlsConfig) -} - -func (client *Client) GPRCMetadata() metadata.MD { - result := map[string]string{} - - if client.serviceAccountName != "" && client.serviceAccountToken != "" { - result = map[string]string{ - rpc.MetadataServiceAccountNameKey: client.serviceAccountName, - rpc.MetadataServiceAccountTokenKey: client.serviceAccountToken, - } - } - - return metadata.New(result) -} - func (client *Client) configureFromDefaultContext() error { configHandle, err := config.NewHandle() if err != nil { diff --git a/pkg/resource/v1/v1.go b/pkg/resource/v1/v1.go index ac2bd1f..6a47ed9 100644 --- a/pkg/resource/v1/v1.go +++ b/pkg/resource/v1/v1.go @@ -377,7 +377,6 @@ const ( type ControllerCapability string const ( - ControllerCapabilityRPCV1 ControllerCapability = "rpc-v1" ControllerCapabilityRPCV2 ControllerCapability = "rpc-v2" ControllerCapabilityVMStateEndpoint ControllerCapability = "vm-state-endpoint" ) diff --git a/rpc/constants.go b/rpc/constants.go deleted file mode 100644 index f9cf3be..0000000 --- a/rpc/constants.go +++ /dev/null @@ -1,10 +0,0 @@ -package rpc - -const MetadataServiceAccountNameKey = "x-orchard-service-account-name" - -//nolint:gosec // G101 check yields a false-positive here, this is not a hard-coded credential -const MetadataServiceAccountTokenKey = "x-orchard-service-account-token" - -const MetadataWorkerNameKey = "x-orchard-worker-name" - -const MetadataWorkerPortForwardingSessionKey = "x-orchard-port-forwarding-session" diff --git a/rpc/orchard.pb.go b/rpc/orchard.pb.go index 6f7cbef..19fd39b 100644 --- a/rpc/orchard.pb.go +++ b/rpc/orchard.pb.go @@ -9,7 +9,6 @@ package rpc import ( protoreflect "google.golang.org/protobuf/reflect/protoreflect" protoimpl "google.golang.org/protobuf/runtime/protoimpl" - emptypb "google.golang.org/protobuf/types/known/emptypb" reflect "reflect" sync "sync" ) @@ -27,7 +26,6 @@ type WatchInstruction struct { unknownFields protoimpl.UnknownFields // Types that are assignable to Action: - // // *WatchInstruction_PortForwardAction // *WatchInstruction_SyncVmsAction // *WatchInstruction_ResolveIpAction @@ -382,52 +380,39 @@ func (x *WatchInstruction_ResolveIP) GetVmUid() string { var File_orchard_proto protoreflect.FileDescriptor var file_orchard_proto_rawDesc = []byte{ - 0x0a, 0x0d, 0x6f, 0x72, 0x63, 0x68, 0x61, 0x72, 0x64, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x1a, - 0x1b, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, - 0x2f, 0x65, 0x6d, 0x70, 0x74, 0x79, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x22, 0x9a, 0x03, 0x0a, - 0x10, 0x57, 0x61, 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, 0x74, 0x72, 0x75, 0x63, 0x74, 0x69, 0x6f, - 0x6e, 0x12, 0x4f, 0x0a, 0x13, 0x70, 0x6f, 0x72, 0x74, 0x5f, 0x66, 0x6f, 0x72, 0x77, 0x61, 0x72, - 0x64, 0x5f, 0x61, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x1d, + 0x0a, 0x0d, 0x6f, 0x72, 0x63, 0x68, 0x61, 0x72, 0x64, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x22, + 0x9a, 0x03, 0x0a, 0x10, 0x57, 0x61, 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, 0x74, 0x72, 0x75, 0x63, + 0x74, 0x69, 0x6f, 0x6e, 0x12, 0x4f, 0x0a, 0x13, 0x70, 0x6f, 0x72, 0x74, 0x5f, 0x66, 0x6f, 0x72, + 0x77, 0x61, 0x72, 0x64, 0x5f, 0x61, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x18, 0x01, 0x20, 0x01, 0x28, + 0x0b, 0x32, 0x1d, 0x2e, 0x57, 0x61, 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, 0x74, 0x72, 0x75, 0x63, + 0x74, 0x69, 0x6f, 0x6e, 0x2e, 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, 0x64, + 0x48, 0x00, 0x52, 0x11, 0x70, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, 0x64, 0x41, + 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x12, 0x43, 0x0a, 0x0f, 0x73, 0x79, 0x6e, 0x63, 0x5f, 0x76, 0x6d, + 0x73, 0x5f, 0x61, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x19, 0x2e, 0x57, 0x61, 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, 0x74, 0x72, 0x75, 0x63, 0x74, 0x69, 0x6f, - 0x6e, 0x2e, 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, 0x64, 0x48, 0x00, 0x52, - 0x11, 0x70, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, 0x64, 0x41, 0x63, 0x74, 0x69, - 0x6f, 0x6e, 0x12, 0x43, 0x0a, 0x0f, 0x73, 0x79, 0x6e, 0x63, 0x5f, 0x76, 0x6d, 0x73, 0x5f, 0x61, - 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x19, 0x2e, 0x57, 0x61, - 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, 0x74, 0x72, 0x75, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x2e, 0x53, - 0x79, 0x6e, 0x63, 0x56, 0x4d, 0x73, 0x48, 0x00, 0x52, 0x0d, 0x73, 0x79, 0x6e, 0x63, 0x56, 0x6d, - 0x73, 0x41, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x12, 0x49, 0x0a, 0x11, 0x72, 0x65, 0x73, 0x6f, 0x6c, - 0x76, 0x65, 0x5f, 0x69, 0x70, 0x5f, 0x61, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x18, 0x03, 0x20, 0x01, - 0x28, 0x0b, 0x32, 0x1b, 0x2e, 0x57, 0x61, 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, 0x74, 0x72, 0x75, - 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x2e, 0x52, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x49, 0x50, 0x48, - 0x00, 0x52, 0x0f, 0x72, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x49, 0x70, 0x41, 0x63, 0x74, 0x69, - 0x6f, 0x6e, 0x1a, 0x52, 0x0a, 0x0b, 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, - 0x64, 0x12, 0x18, 0x0a, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x18, 0x01, 0x20, 0x01, + 0x6e, 0x2e, 0x53, 0x79, 0x6e, 0x63, 0x56, 0x4d, 0x73, 0x48, 0x00, 0x52, 0x0d, 0x73, 0x79, 0x6e, + 0x63, 0x56, 0x6d, 0x73, 0x41, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x12, 0x49, 0x0a, 0x11, 0x72, 0x65, + 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x5f, 0x69, 0x70, 0x5f, 0x61, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x18, + 0x03, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x1b, 0x2e, 0x57, 0x61, 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, + 0x74, 0x72, 0x75, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x2e, 0x52, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, + 0x49, 0x50, 0x48, 0x00, 0x52, 0x0f, 0x72, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x49, 0x70, 0x41, + 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x1a, 0x52, 0x0a, 0x0b, 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, + 0x77, 0x61, 0x72, 0x64, 0x12, 0x18, 0x0a, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x18, + 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x12, 0x15, + 0x0a, 0x06, 0x76, 0x6d, 0x5f, 0x75, 0x69, 0x64, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, + 0x76, 0x6d, 0x55, 0x69, 0x64, 0x12, 0x12, 0x0a, 0x04, 0x70, 0x6f, 0x72, 0x74, 0x18, 0x03, 0x20, + 0x01, 0x28, 0x0d, 0x52, 0x04, 0x70, 0x6f, 0x72, 0x74, 0x1a, 0x09, 0x0a, 0x07, 0x53, 0x79, 0x6e, + 0x63, 0x56, 0x4d, 0x73, 0x1a, 0x3c, 0x0a, 0x09, 0x52, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x49, + 0x50, 0x12, 0x18, 0x0a, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x12, 0x15, 0x0a, 0x06, 0x76, 0x6d, 0x5f, 0x75, 0x69, 0x64, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x76, 0x6d, 0x55, - 0x69, 0x64, 0x12, 0x12, 0x0a, 0x04, 0x70, 0x6f, 0x72, 0x74, 0x18, 0x03, 0x20, 0x01, 0x28, 0x0d, - 0x52, 0x04, 0x70, 0x6f, 0x72, 0x74, 0x1a, 0x09, 0x0a, 0x07, 0x53, 0x79, 0x6e, 0x63, 0x56, 0x4d, - 0x73, 0x1a, 0x3c, 0x0a, 0x09, 0x52, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x49, 0x50, 0x12, 0x18, - 0x0a, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, - 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x12, 0x15, 0x0a, 0x06, 0x76, 0x6d, 0x5f, 0x75, - 0x69, 0x64, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x76, 0x6d, 0x55, 0x69, 0x64, 0x42, - 0x08, 0x0a, 0x06, 0x61, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x22, 0x25, 0x0a, 0x0f, 0x50, 0x6f, 0x72, - 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, 0x64, 0x44, 0x61, 0x74, 0x61, 0x12, 0x12, 0x0a, 0x04, - 0x64, 0x61, 0x74, 0x61, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x04, 0x64, 0x61, 0x74, 0x61, - 0x22, 0x3b, 0x0a, 0x0f, 0x52, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x49, 0x50, 0x52, 0x65, 0x73, - 0x75, 0x6c, 0x74, 0x12, 0x18, 0x0a, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x18, 0x01, - 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x12, 0x0e, 0x0a, - 0x02, 0x69, 0x70, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x02, 0x69, 0x70, 0x32, 0xb0, 0x01, - 0x0a, 0x0a, 0x43, 0x6f, 0x6e, 0x74, 0x72, 0x6f, 0x6c, 0x6c, 0x65, 0x72, 0x12, 0x34, 0x0a, 0x05, - 0x57, 0x61, 0x74, 0x63, 0x68, 0x12, 0x16, 0x2e, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65, 0x2e, 0x70, - 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2e, 0x45, 0x6d, 0x70, 0x74, 0x79, 0x1a, 0x11, 0x2e, - 0x57, 0x61, 0x74, 0x63, 0x68, 0x49, 0x6e, 0x73, 0x74, 0x72, 0x75, 0x63, 0x74, 0x69, 0x6f, 0x6e, - 0x30, 0x01, 0x12, 0x35, 0x0a, 0x0b, 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, - 0x64, 0x12, 0x10, 0x2e, 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, 0x64, 0x44, - 0x61, 0x74, 0x61, 0x1a, 0x10, 0x2e, 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, - 0x64, 0x44, 0x61, 0x74, 0x61, 0x28, 0x01, 0x30, 0x01, 0x12, 0x35, 0x0a, 0x09, 0x52, 0x65, 0x73, - 0x6f, 0x6c, 0x76, 0x65, 0x49, 0x50, 0x12, 0x10, 0x2e, 0x52, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, - 0x49, 0x50, 0x52, 0x65, 0x73, 0x75, 0x6c, 0x74, 0x1a, 0x16, 0x2e, 0x67, 0x6f, 0x6f, 0x67, 0x6c, - 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2e, 0x45, 0x6d, 0x70, 0x74, 0x79, + 0x69, 0x64, 0x42, 0x08, 0x0a, 0x06, 0x61, 0x63, 0x74, 0x69, 0x6f, 0x6e, 0x22, 0x25, 0x0a, 0x0f, + 0x50, 0x6f, 0x72, 0x74, 0x46, 0x6f, 0x72, 0x77, 0x61, 0x72, 0x64, 0x44, 0x61, 0x74, 0x61, 0x12, + 0x12, 0x0a, 0x04, 0x64, 0x61, 0x74, 0x61, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x04, 0x64, + 0x61, 0x74, 0x61, 0x22, 0x3b, 0x0a, 0x0f, 0x52, 0x65, 0x73, 0x6f, 0x6c, 0x76, 0x65, 0x49, 0x50, + 0x52, 0x65, 0x73, 0x75, 0x6c, 0x74, 0x12, 0x18, 0x0a, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, + 0x6e, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, + 0x12, 0x0e, 0x0a, 0x02, 0x69, 0x70, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x02, 0x69, 0x70, 0x42, 0x23, 0x5a, 0x21, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x63, 0x69, 0x72, 0x72, 0x75, 0x73, 0x6c, 0x61, 0x62, 0x73, 0x2f, 0x6f, 0x72, 0x63, 0x68, 0x61, 0x72, 0x64, 0x2f, 0x72, 0x70, 0x63, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, @@ -453,20 +438,13 @@ var file_orchard_proto_goTypes = []any{ (*WatchInstruction_PortForward)(nil), // 3: WatchInstruction.PortForward (*WatchInstruction_SyncVMs)(nil), // 4: WatchInstruction.SyncVMs (*WatchInstruction_ResolveIP)(nil), // 5: WatchInstruction.ResolveIP - (*emptypb.Empty)(nil), // 6: google.protobuf.Empty } var file_orchard_proto_depIdxs = []int32{ 3, // 0: WatchInstruction.port_forward_action:type_name -> WatchInstruction.PortForward 4, // 1: WatchInstruction.sync_vms_action:type_name -> WatchInstruction.SyncVMs 5, // 2: WatchInstruction.resolve_ip_action:type_name -> WatchInstruction.ResolveIP - 6, // 3: Controller.Watch:input_type -> google.protobuf.Empty - 1, // 4: Controller.PortForward:input_type -> PortForwardData - 2, // 5: Controller.ResolveIP:input_type -> ResolveIPResult - 0, // 6: Controller.Watch:output_type -> WatchInstruction - 1, // 7: Controller.PortForward:output_type -> PortForwardData - 6, // 8: Controller.ResolveIP:output_type -> google.protobuf.Empty - 6, // [6:9] is the sub-list for method output_type - 3, // [3:6] is the sub-list for method input_type + 3, // [3:3] is the sub-list for method output_type + 3, // [3:3] is the sub-list for method input_type 3, // [3:3] is the sub-list for extension type_name 3, // [3:3] is the sub-list for extension extendee 0, // [0:3] is the sub-list for field type_name @@ -564,7 +542,7 @@ func file_orchard_proto_init() { NumEnums: 0, NumMessages: 6, NumExtensions: 0, - NumServices: 1, + NumServices: 0, }, GoTypes: file_orchard_proto_goTypes, DependencyIndexes: file_orchard_proto_depIdxs, diff --git a/rpc/orchard.proto b/rpc/orchard.proto index 7328e7a..58855e5 100644 --- a/rpc/orchard.proto +++ b/rpc/orchard.proto @@ -1,21 +1,7 @@ syntax = "proto3"; -import "google/protobuf/empty.proto"; - option go_package = "github.com/cirruslabs/orchard/rpc"; -service Controller { - // message bus between the controller and a worker - rpc Watch(google.protobuf.Empty) returns (stream WatchInstruction); - - // single purpose method when a port forward is requested and running - // session information is passed in the requests metadata - rpc PortForward(stream PortForwardData) returns (stream PortForwardData); - - // worker calls this method when it has successfully resolved the VM's IP - rpc ResolveIP(ResolveIPResult) returns (google.protobuf.Empty); -} - message WatchInstruction { message PortForward { // we can have multiple port forwards for the same vm/port pair diff --git a/rpc/orchard_grpc.pb.go b/rpc/orchard_grpc.pb.go deleted file mode 100644 index d4f5a29..0000000 --- a/rpc/orchard_grpc.pb.go +++ /dev/null @@ -1,255 +0,0 @@ -// Code generated by protoc-gen-go-grpc. DO NOT EDIT. -// versions: -// - protoc-gen-go-grpc v1.4.0 -// - protoc (unknown) -// source: orchard.proto - -package rpc - -import ( - context "context" - grpc "google.golang.org/grpc" - codes "google.golang.org/grpc/codes" - status "google.golang.org/grpc/status" - emptypb "google.golang.org/protobuf/types/known/emptypb" -) - -// This is a compile-time assertion to ensure that this generated file -// is compatible with the grpc package it is being compiled against. -// Requires gRPC-Go v1.62.0 or later. -const _ = grpc.SupportPackageIsVersion8 - -const ( - Controller_Watch_FullMethodName = "/Controller/Watch" - Controller_PortForward_FullMethodName = "/Controller/PortForward" - Controller_ResolveIP_FullMethodName = "/Controller/ResolveIP" -) - -// ControllerClient is the client API for Controller service. -// -// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream. -type ControllerClient interface { - // message bus between the controller and a worker - Watch(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (Controller_WatchClient, error) - // single purpose method when a port forward is requested and running - // session information is passed in the requests metadata - PortForward(ctx context.Context, opts ...grpc.CallOption) (Controller_PortForwardClient, error) - // worker calls this method when it has successfully resolved the VM's IP - ResolveIP(ctx context.Context, in *ResolveIPResult, opts ...grpc.CallOption) (*emptypb.Empty, error) -} - -type controllerClient struct { - cc grpc.ClientConnInterface -} - -func NewControllerClient(cc grpc.ClientConnInterface) ControllerClient { - return &controllerClient{cc} -} - -func (c *controllerClient) Watch(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (Controller_WatchClient, error) { - cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) - stream, err := c.cc.NewStream(ctx, &Controller_ServiceDesc.Streams[0], Controller_Watch_FullMethodName, cOpts...) - if err != nil { - return nil, err - } - x := &controllerWatchClient{ClientStream: stream} - if err := x.ClientStream.SendMsg(in); err != nil { - return nil, err - } - if err := x.ClientStream.CloseSend(); err != nil { - return nil, err - } - return x, nil -} - -type Controller_WatchClient interface { - Recv() (*WatchInstruction, error) - grpc.ClientStream -} - -type controllerWatchClient struct { - grpc.ClientStream -} - -func (x *controllerWatchClient) Recv() (*WatchInstruction, error) { - m := new(WatchInstruction) - if err := x.ClientStream.RecvMsg(m); err != nil { - return nil, err - } - return m, nil -} - -func (c *controllerClient) PortForward(ctx context.Context, opts ...grpc.CallOption) (Controller_PortForwardClient, error) { - cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) - stream, err := c.cc.NewStream(ctx, &Controller_ServiceDesc.Streams[1], Controller_PortForward_FullMethodName, cOpts...) - if err != nil { - return nil, err - } - x := &controllerPortForwardClient{ClientStream: stream} - return x, nil -} - -type Controller_PortForwardClient interface { - Send(*PortForwardData) error - Recv() (*PortForwardData, error) - grpc.ClientStream -} - -type controllerPortForwardClient struct { - grpc.ClientStream -} - -func (x *controllerPortForwardClient) Send(m *PortForwardData) error { - return x.ClientStream.SendMsg(m) -} - -func (x *controllerPortForwardClient) Recv() (*PortForwardData, error) { - m := new(PortForwardData) - if err := x.ClientStream.RecvMsg(m); err != nil { - return nil, err - } - return m, nil -} - -func (c *controllerClient) ResolveIP(ctx context.Context, in *ResolveIPResult, opts ...grpc.CallOption) (*emptypb.Empty, error) { - cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) - out := new(emptypb.Empty) - err := c.cc.Invoke(ctx, Controller_ResolveIP_FullMethodName, in, out, cOpts...) - if err != nil { - return nil, err - } - return out, nil -} - -// ControllerServer is the server API for Controller service. -// All implementations must embed UnimplementedControllerServer -// for forward compatibility -type ControllerServer interface { - // message bus between the controller and a worker - Watch(*emptypb.Empty, Controller_WatchServer) error - // single purpose method when a port forward is requested and running - // session information is passed in the requests metadata - PortForward(Controller_PortForwardServer) error - // worker calls this method when it has successfully resolved the VM's IP - ResolveIP(context.Context, *ResolveIPResult) (*emptypb.Empty, error) - mustEmbedUnimplementedControllerServer() -} - -// UnimplementedControllerServer must be embedded to have forward compatible implementations. -type UnimplementedControllerServer struct { -} - -func (UnimplementedControllerServer) Watch(*emptypb.Empty, Controller_WatchServer) error { - return status.Errorf(codes.Unimplemented, "method Watch not implemented") -} -func (UnimplementedControllerServer) PortForward(Controller_PortForwardServer) error { - return status.Errorf(codes.Unimplemented, "method PortForward not implemented") -} -func (UnimplementedControllerServer) ResolveIP(context.Context, *ResolveIPResult) (*emptypb.Empty, error) { - return nil, status.Errorf(codes.Unimplemented, "method ResolveIP not implemented") -} -func (UnimplementedControllerServer) mustEmbedUnimplementedControllerServer() {} - -// UnsafeControllerServer may be embedded to opt out of forward compatibility for this service. -// Use of this interface is not recommended, as added methods to ControllerServer will -// result in compilation errors. -type UnsafeControllerServer interface { - mustEmbedUnimplementedControllerServer() -} - -func RegisterControllerServer(s grpc.ServiceRegistrar, srv ControllerServer) { - s.RegisterService(&Controller_ServiceDesc, srv) -} - -func _Controller_Watch_Handler(srv interface{}, stream grpc.ServerStream) error { - m := new(emptypb.Empty) - if err := stream.RecvMsg(m); err != nil { - return err - } - return srv.(ControllerServer).Watch(m, &controllerWatchServer{ServerStream: stream}) -} - -type Controller_WatchServer interface { - Send(*WatchInstruction) error - grpc.ServerStream -} - -type controllerWatchServer struct { - grpc.ServerStream -} - -func (x *controllerWatchServer) Send(m *WatchInstruction) error { - return x.ServerStream.SendMsg(m) -} - -func _Controller_PortForward_Handler(srv interface{}, stream grpc.ServerStream) error { - return srv.(ControllerServer).PortForward(&controllerPortForwardServer{ServerStream: stream}) -} - -type Controller_PortForwardServer interface { - Send(*PortForwardData) error - Recv() (*PortForwardData, error) - grpc.ServerStream -} - -type controllerPortForwardServer struct { - grpc.ServerStream -} - -func (x *controllerPortForwardServer) Send(m *PortForwardData) error { - return x.ServerStream.SendMsg(m) -} - -func (x *controllerPortForwardServer) Recv() (*PortForwardData, error) { - m := new(PortForwardData) - if err := x.ServerStream.RecvMsg(m); err != nil { - return nil, err - } - return m, nil -} - -func _Controller_ResolveIP_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { - in := new(ResolveIPResult) - if err := dec(in); err != nil { - return nil, err - } - if interceptor == nil { - return srv.(ControllerServer).ResolveIP(ctx, in) - } - info := &grpc.UnaryServerInfo{ - Server: srv, - FullMethod: Controller_ResolveIP_FullMethodName, - } - handler := func(ctx context.Context, req interface{}) (interface{}, error) { - return srv.(ControllerServer).ResolveIP(ctx, req.(*ResolveIPResult)) - } - return interceptor(ctx, in, info, handler) -} - -// Controller_ServiceDesc is the grpc.ServiceDesc for Controller service. -// It's only intended for direct use with grpc.RegisterService, -// and not to be introspected or modified (even as a copy) -var Controller_ServiceDesc = grpc.ServiceDesc{ - ServiceName: "Controller", - HandlerType: (*ControllerServer)(nil), - Methods: []grpc.MethodDesc{ - { - MethodName: "ResolveIP", - Handler: _Controller_ResolveIP_Handler, - }, - }, - Streams: []grpc.StreamDesc{ - { - StreamName: "Watch", - Handler: _Controller_Watch_Handler, - ServerStreams: true, - }, - { - StreamName: "PortForward", - Handler: _Controller_PortForward_Handler, - ServerStreams: true, - ClientStreams: true, - }, - }, - Metadata: "orchard.proto", -}