mirror of
https://github.com/cirruslabs/orchard.git
synced 2026-09-30 03:51:43 +02:00
Remove legacy RPC v1 support
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -3,6 +3,3 @@ plugins:
|
||||
- plugin: go
|
||||
out: rpc
|
||||
opt: paths=source_relative
|
||||
- plugin: go-grpc
|
||||
out: rpc
|
||||
opt: paths=source_relative
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
+1
-30
@@ -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 {
|
||||
|
||||
@@ -377,7 +377,6 @@ const (
|
||||
type ControllerCapability string
|
||||
|
||||
const (
|
||||
ControllerCapabilityRPCV1 ControllerCapability = "rpc-v1"
|
||||
ControllerCapabilityRPCV2 ControllerCapability = "rpc-v2"
|
||||
ControllerCapabilityVMStateEndpoint ControllerCapability = "vm-state-endpoint"
|
||||
)
|
||||
|
||||
@@ -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"
|
||||
+33
-55
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
}
|
||||
Reference in New Issue
Block a user