forked from rook/rook
ceph: proxy ceph commands when multus is configured
When the CephCluster is configured with Multus and multiple networks are used to deploy Ceph some commands are failing to be executed from the Operator. These commands, in particular, `radosgw-admin` ones need access to the "ceph public network" to talk to OSDs. Unfortunately, the Rook-Ceph Operator does not have the network annotations and thus doesn't have the networks available and cannot reach OSDs. So the commands end up hanging and eventually time out. Applying the annotations to the Operator pod is possible but will result in restarting the operator too and this should be avoided at all costs. Also, applying the annotations beforehand is not possible since the Multus declaration is in the CephCluster specification. So we would have no idea what to do. So the current approach runs a new sidecar container in the mgr pod to act as a proxy for "some" ceph commands, only the `radosgw-admin` ones for multi-site setup. This is a small container with admin access running idle waiting for commands to be executed. In a sense, it is similar to the toolbox but we didn't want to clearly expose it, so running as a sidecar is quite nice. Proxying command is obviously not always recommended since we add an extra hop in the network path. Now each request has to go from the operator pod to the API server to the remote pod to Ceph. Previously, the command only goes from the operator to Ceph. It's worth noting that external mode is not impacted since no rgw pod is configured. This scenario is flexible and allows us to scale pretty well since any CephCluster with Multus will see its mgr sidecar deployed and can then talk to Ceph. We are not limited. Signed-off-by: Sébastien Han <seb@redhat.com>
This commit is contained in:
@@ -157,4 +157,23 @@ rules:
|
||||
- apiGroups: ["coordination.k8s.io"]
|
||||
resources: ["leases"]
|
||||
verbs: ["get", "watch", "list", "delete", "update", "create"]
|
||||
---
|
||||
kind: ClusterRole
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
metadata:
|
||||
name: rook-ceph-system
|
||||
labels:
|
||||
operator: rook
|
||||
storage-backend: ceph
|
||||
rules:
|
||||
# Most resources are represented by a string representation of their name, such as “pods”, just as it appears in the URL for the relevant API endpoint.
|
||||
# However, some Kubernetes APIs involve a “subresource”, such as the logs for a pod. [...]
|
||||
# To represent this in an RBAC role, use a slash to delimit the resource and subresource.
|
||||
# https://kubernetes.io/docs/reference/access-authn-authz/rbac/#referring-to-resources
|
||||
- apiGroups: [""]
|
||||
resources: ["pods", "pods/log"]
|
||||
verbs: ["get", "list"]
|
||||
- apiGroups: [""]
|
||||
resources: ["pods/exec"]
|
||||
verbs: ["create"]
|
||||
{{- end }}
|
||||
|
||||
@@ -118,4 +118,20 @@ roleRef:
|
||||
kind: Role
|
||||
name: rbd-external-provisioner-cfg
|
||||
apiGroup: rbac.authorization.k8s.io
|
||||
---
|
||||
kind: ClusterRoleBinding
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
metadata:
|
||||
name: rook-ceph-system
|
||||
labels:
|
||||
operator: rook
|
||||
storage-backend: ceph
|
||||
roleRef:
|
||||
apiGroup: rbac.authorization.k8s.io
|
||||
kind: ClusterRole
|
||||
name: rook-ceph-system
|
||||
subjects:
|
||||
- kind: ServiceAccount
|
||||
name: rook-ceph-system
|
||||
namespace: rook-ceph
|
||||
{{- end }}
|
||||
|
||||
@@ -91,6 +91,25 @@ rules:
|
||||
- update
|
||||
- delete
|
||||
---
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
kind: ClusterRole
|
||||
metadata:
|
||||
name: rook-ceph-system
|
||||
labels:
|
||||
operator: rook
|
||||
storage-backend: ceph
|
||||
rules:
|
||||
# Most resources are represented by a string representation of their name, such as “pods”, just as it appears in the URL for the relevant API endpoint.
|
||||
# However, some Kubernetes APIs involve a “subresource”, such as the logs for a pod. [...]
|
||||
# To represent this in an RBAC role, use a slash to delimit the resource and subresource.
|
||||
# https://kubernetes.io/docs/reference/access-authn-authz/rbac/#referring-to-resources
|
||||
- apiGroups: [""]
|
||||
resources: ["pods", "pods/log"]
|
||||
verbs: ["get", "list"]
|
||||
- apiGroups: [""]
|
||||
resources: ["pods/exec"]
|
||||
verbs: ["create"]
|
||||
---
|
||||
# The role for the operator to manage resources in its own namespace
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
kind: Role
|
||||
@@ -351,6 +370,22 @@ subjects:
|
||||
name: rook-ceph-system
|
||||
namespace: rook-ceph # namespace:operator
|
||||
---
|
||||
kind: ClusterRoleBinding
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
metadata:
|
||||
name: rook-ceph-system
|
||||
labels:
|
||||
operator: rook
|
||||
storage-backend: ceph
|
||||
roleRef:
|
||||
apiGroup: rbac.authorization.k8s.io
|
||||
kind: ClusterRole
|
||||
name: rook-ceph-system
|
||||
subjects:
|
||||
- kind: ServiceAccount
|
||||
name: rook-ceph-system
|
||||
namespace: rook-ceph # namespace:operator
|
||||
---
|
||||
# Grant the rook system daemons cluster-wide access to manage the Rook CRDs, PVCs, and storage classes
|
||||
kind: ClusterRoleBinding
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
|
||||
@@ -154,6 +154,9 @@ func NewContext() *clusterd.Context {
|
||||
context.Clientset, err = kubernetes.NewForConfig(context.KubeConfig)
|
||||
TerminateOnError(err, "failed to create k8s clientset")
|
||||
|
||||
context.RemoteExecutor.ClientSet = context.Clientset
|
||||
context.RemoteExecutor.RestClient = context.KubeConfig
|
||||
|
||||
// Dynamic clientset allows dealing with resources that aren't statically typed but determined
|
||||
// at runtime.
|
||||
context.DynamicClientset, err = dynamic.NewForConfig(context.KubeConfig)
|
||||
|
||||
@@ -30,11 +30,11 @@ require (
|
||||
golang.org/x/sync v0.0.0-20201207232520-09787c993a3a
|
||||
gopkg.in/ini.v1 v1.57.0
|
||||
gopkg.in/yaml.v2 v2.4.0
|
||||
k8s.io/api v0.21.1
|
||||
k8s.io/api v0.21.2
|
||||
k8s.io/apiextensions-apiserver v0.21.1
|
||||
k8s.io/apimachinery v0.21.1
|
||||
k8s.io/apimachinery v0.21.2
|
||||
k8s.io/apiserver v0.21.1
|
||||
k8s.io/client-go v0.21.1
|
||||
k8s.io/client-go v0.21.2
|
||||
k8s.io/cloud-provider v0.21.1
|
||||
k8s.io/component-helpers v0.21.1
|
||||
k8s.io/kube-controller-manager v0.21.1
|
||||
|
||||
@@ -362,6 +362,7 @@ github.com/elastic/go-windows v1.0.1/go.mod h1:FoVvqWSun28vaDQPbj2Elfc0JahhPB7WQ
|
||||
github.com/elazarl/go-bindata-assetfs v1.0.0 h1:G/bYguwHIzWq9ZoyUQqrjTmJbbYn3j3CKKpKinvZLFk=
|
||||
github.com/elazarl/go-bindata-assetfs v1.0.0/go.mod h1:v+YaWX3bdea5J/mo8dSETolEo7R71Vk1u8bnjau5yw4=
|
||||
github.com/elazarl/goproxy v0.0.0-20170405201442-c4fc26588b6e/go.mod h1:/Zj4wYkgs4iZTTu3o/KG3Itv/qCCa8VVMlb3i9OVuzc=
|
||||
github.com/elazarl/goproxy v0.0.0-20180725130230-947c36da3153 h1:yUdfgN0XgIJw7foRItutHYUIhlcKzcSf5vDpdhQAKTc=
|
||||
github.com/elazarl/goproxy v0.0.0-20180725130230-947c36da3153/go.mod h1:/Zj4wYkgs4iZTTu3o/KG3Itv/qCCa8VVMlb3i9OVuzc=
|
||||
github.com/ema/qdisc v0.0.0-20190904071900-b82c76788043/go.mod h1:ix4kG2zvdUd8kEKSW0ZTr1XLks0epFpI4j745DXxlNE=
|
||||
github.com/emicklei/go-restful v0.0.0-20170410110728-ff4f55a20633/go.mod h1:otzb+WCGbkyDHkqmQmT5YD2WR4BBwUdeQoFo8l/7tVs=
|
||||
@@ -1107,6 +1108,7 @@ github.com/mitchellh/pointerstructure v0.0.0-20190430161007-f252a8fd71c8/go.mod
|
||||
github.com/mitchellh/reflectwalk v1.0.0/go.mod h1:mSTlrgnPZtwu0c4WaC2kGObEpuNDbx0jmZXqmk4esnw=
|
||||
github.com/mitchellh/reflectwalk v1.0.1 h1:FVzMWA5RllMAKIdUSC8mdWo3XtwoecrH79BY70sEEpE=
|
||||
github.com/mitchellh/reflectwalk v1.0.1/go.mod h1:mSTlrgnPZtwu0c4WaC2kGObEpuNDbx0jmZXqmk4esnw=
|
||||
github.com/moby/spdystream v0.2.0 h1:cjW1zVyyoiM0T7b6UoySUFqzXMoqRckQtXwGPiBhOM8=
|
||||
github.com/moby/spdystream v0.2.0/go.mod h1:f7i0iNDQJ059oMTcWxx8MA/zKFIuD/lY+0GqbN2Wy8c=
|
||||
github.com/moby/term v0.0.0-20200312100748-672ec06f55cd/go.mod h1:DdlQx2hp0Ss5/fLikoLlEeIYiATotOjgB//nb973jeo=
|
||||
github.com/moby/term v0.0.0-20201216013528-df9cb8a40635/go.mod h1:FBS0z0QWA44HXygs7VXDUOGoN/1TV3RuWkLO04am3wc=
|
||||
@@ -1812,6 +1814,7 @@ golang.org/x/sys v0.0.0-20210119212857-b64e53b001e4/go.mod h1:h1NjWce9XRLGQEsW7w
|
||||
golang.org/x/sys v0.0.0-20210124154548-22da62e12c0c/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20210225134936-a50acf3fe073/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20210426230700-d19ff857e887/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20210603081109-ebe580a85c40 h1:JWgyZ1qgdTaF3N3oxC+MdTV7qvEEgHo3otj+HB5CM7Q=
|
||||
golang.org/x/sys v0.0.0-20210603081109-ebe580a85c40/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/term v0.0.0-20201117132131-f5c789dd3221/go.mod h1:Nr5EML6q2oocZ2LXRh80K7BxOlk5/8JxuGnuhpl+muw=
|
||||
@@ -2124,8 +2127,9 @@ k8s.io/api v0.19.0/go.mod h1:I1K45XlvTrDjmj5LoM5LuP/KYrhWbjUKT/SoPG0qTjw=
|
||||
k8s.io/api v0.19.1/go.mod h1:+u/k4/K/7vp4vsfdT7dyl8Oxk1F26Md4g5F26Tu85PU=
|
||||
k8s.io/api v0.19.2/go.mod h1:IQpK0zFQ1xc5iNIQPqzgoOwuFugaYHK4iCknlAQP9nI=
|
||||
k8s.io/api v0.19.3/go.mod h1:VF+5FT1B74Pw3KxMdKyinLo+zynBaMBiAfGMuldcNDs=
|
||||
k8s.io/api v0.21.1 h1:94bbZ5NTjdINJEdzOkpS4vdPhkb1VFpTYC9zh43f75c=
|
||||
k8s.io/api v0.21.1/go.mod h1:FstGROTmsSHBarKc8bylzXih8BLNYTiS3TZcsoEDg2s=
|
||||
k8s.io/api v0.21.2 h1:vz7DqmRsXTCSa6pNxXwQ1IYeAZgdIsua+DZU+o+SX3Y=
|
||||
k8s.io/api v0.21.2/go.mod h1:Lv6UGJZ1rlMI1qusN8ruAp9PUBFyBwpEHAdG24vIsiU=
|
||||
k8s.io/apiextensions-apiserver v0.0.0-20190409022649-727a075fdec8/go.mod h1:IxkesAMoaCRoLrPJdZNZUQp9NfZnzqaVzLhb2VEQzXE=
|
||||
k8s.io/apiextensions-apiserver v0.0.0-20190918161926-8f644eb6e783/go.mod h1:xvae1SZB3E17UpV59AWc271W/Ph25N+bjPyR63X6tPY=
|
||||
k8s.io/apiextensions-apiserver v0.15.7/go.mod h1:ctb/NYtsiBt6CGN42Z+JrOkxi9nJYaKZYmatJ6SUy0Y=
|
||||
@@ -2149,8 +2153,9 @@ k8s.io/apimachinery v0.19.0/go.mod h1:DnPGDnARWFvYa3pMHgSxtbZb7gpzzAZ1pTfaUNDVlm
|
||||
k8s.io/apimachinery v0.19.1/go.mod h1:DnPGDnARWFvYa3pMHgSxtbZb7gpzzAZ1pTfaUNDVlmA=
|
||||
k8s.io/apimachinery v0.19.2/go.mod h1:DnPGDnARWFvYa3pMHgSxtbZb7gpzzAZ1pTfaUNDVlmA=
|
||||
k8s.io/apimachinery v0.19.3/go.mod h1:DnPGDnARWFvYa3pMHgSxtbZb7gpzzAZ1pTfaUNDVlmA=
|
||||
k8s.io/apimachinery v0.21.1 h1:Q6XuHGlj2xc+hlMCvqyYfbv3H7SRGn2c8NycxJquDVs=
|
||||
k8s.io/apimachinery v0.21.1/go.mod h1:jbreFvJo3ov9rj7eWT7+sYiRx+qZuCYXwWT1bcDswPY=
|
||||
k8s.io/apimachinery v0.21.2 h1:vezUc/BHqWlQDnZ+XkrpXSmnANSLbpnlpwo0Lhk0gpc=
|
||||
k8s.io/apimachinery v0.21.2/go.mod h1:CdTY8fU/BlvAbJ2z/8kBwimGki5Zp8/fbVuLY8gJumM=
|
||||
k8s.io/apiserver v0.0.0-20190918160949-bfa5e2e684ad/go.mod h1:XPCXEwhjaFN29a8NldXA901ElnKeKLrLtREO9ZhFyhg=
|
||||
k8s.io/apiserver v0.0.0-20191122221311-9d521947b1e1/go.mod h1:RbsZY5zzBIWnz4KbctZsTVjwIuOpTp4Z8oCgFHN4kZQ=
|
||||
k8s.io/apiserver v0.15.7/go.mod h1:d5Dbyt588GbBtUnbx9fSK+pYeqgZa32op+I1BmXiNuE=
|
||||
@@ -2169,8 +2174,9 @@ k8s.io/client-go v0.19.0/go.mod h1:H9E/VT95blcFQnlyShFgnFT9ZnJOAceiUHM3MlRC+mU=
|
||||
k8s.io/client-go v0.19.1/go.mod h1:AZOIVSI9UUtQPeJD3zJFp15CEhSjRgAuQP5PWRJrCIQ=
|
||||
k8s.io/client-go v0.19.2/go.mod h1:S5wPhCqyDNAlzM9CnEdgTGV4OqhsW3jGO1UM1epwfJA=
|
||||
k8s.io/client-go v0.19.3/go.mod h1:+eEMktZM+MG0KO+PTkci8xnbCZHvj9TqR6Q1XDUIJOM=
|
||||
k8s.io/client-go v0.21.1 h1:bhblWYLZKUu+pm50plvQF8WpY6TXdRRtcS/K9WauOj4=
|
||||
k8s.io/client-go v0.21.1/go.mod h1:/kEw4RgW+3xnBGzvp9IWxKSNA+lXn3A7AuH3gdOAzLs=
|
||||
k8s.io/client-go v0.21.2 h1:Q1j4L/iMN4pTw6Y4DWppBoUxgKO8LbffEMVEV00MUp0=
|
||||
k8s.io/client-go v0.21.2/go.mod h1:HdJ9iknWpbl3vMGtib6T2PyI/VYxiZfq936WNVHBRrA=
|
||||
k8s.io/cloud-provider v0.21.1 h1:V7ro0ZuxMBNYVH4lJKxCdI+h2bQ7EApC5f7sQYrQLVE=
|
||||
k8s.io/cloud-provider v0.21.1/go.mod h1:GgiRu7hOsZh3+VqMMbfLJJS9ZZM9A8k/YiZG8zkWpX4=
|
||||
k8s.io/code-generator v0.0.0-20190912054826-cd179ad6a269/go.mod h1:V5BD6M4CyaN5m+VthcclXWsVcT1Hu+glwa1bi3MIsyE=
|
||||
|
||||
@@ -54,6 +54,9 @@ type Context struct {
|
||||
// The implementation of executing a console command
|
||||
Executor exec.Executor
|
||||
|
||||
// The implementation of executing remotely a console command to a given pod
|
||||
RemoteExecutor exec.RemotePodCommandExecutor
|
||||
|
||||
// The root configuration directory used by services
|
||||
ConfigDir string
|
||||
|
||||
|
||||
@@ -24,6 +24,7 @@ import (
|
||||
|
||||
"github.com/pkg/errors"
|
||||
"github.com/rook/rook/pkg/clusterd"
|
||||
"github.com/rook/rook/pkg/util/exec"
|
||||
)
|
||||
|
||||
// RunAllCephCommandsInToolboxPod - when running the e2e tests, all ceph commands need to be run in the toolbox.
|
||||
@@ -41,8 +42,7 @@ const (
|
||||
// Kubectl is the name of the CLI tool for 'kubectl'
|
||||
Kubectl = "kubectl"
|
||||
// CrushTool is the name of the CLI tool for 'crushtool'
|
||||
CrushTool = "crushtool"
|
||||
CephCommandTimeout = 15 * time.Second
|
||||
CrushTool = "crushtool"
|
||||
// DefaultPGCount will cause Ceph to use the internal default PG count
|
||||
DefaultPGCount = "0"
|
||||
)
|
||||
@@ -62,7 +62,7 @@ func FinalizeCephCommandArgs(command string, clusterInfo *ClusterInfo, args []st
|
||||
|
||||
// we could use a slice and iterate over it but since we have only 3 elements
|
||||
// I don't think this is worth a loop
|
||||
timeout := strconv.Itoa(int(CephCommandTimeout.Seconds()))
|
||||
timeout := strconv.Itoa(int(exec.CephCommandTimeout.Seconds()))
|
||||
if command != "rbd" && command != "crushtool" && command != "radosgw-admin" {
|
||||
args = append(args, "--connect-timeout="+timeout)
|
||||
}
|
||||
@@ -174,7 +174,7 @@ func (c *CephToolCommand) RunWithTimeout(timeout time.Duration) ([]byte, error)
|
||||
// configured its arguments. It is future work to integrate this case into the
|
||||
// generalization.
|
||||
func ExecuteRBDCommandWithTimeout(context *clusterd.Context, args []string) (string, error) {
|
||||
output, err := context.Executor.ExecuteCommandWithTimeout(CephCommandTimeout, RBDTool, args...)
|
||||
output, err := context.Executor.ExecuteCommandWithTimeout(exec.CephCommandTimeout, RBDTool, args...)
|
||||
return output, err
|
||||
}
|
||||
|
||||
|
||||
@@ -20,6 +20,7 @@ import (
|
||||
"strconv"
|
||||
"testing"
|
||||
|
||||
"github.com/rook/rook/pkg/util/exec"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
@@ -30,7 +31,7 @@ func TestFinalizeCephCommandArgs(t *testing.T) {
|
||||
args := []string{"quorum_status"}
|
||||
expectedArgs := []string{
|
||||
"quorum_status",
|
||||
"--connect-timeout=" + strconv.Itoa(int(CephCommandTimeout.Seconds())),
|
||||
"--connect-timeout=" + strconv.Itoa(int(exec.CephCommandTimeout.Seconds())),
|
||||
"--cluster=rook",
|
||||
"--conf=/var/lib/rook/rook-ceph/rook/rook.config",
|
||||
"--name=client.admin",
|
||||
|
||||
@@ -31,6 +31,7 @@ import (
|
||||
"github.com/rook/rook/pkg/daemon/ceph/client"
|
||||
cephver "github.com/rook/rook/pkg/operator/ceph/version"
|
||||
"github.com/rook/rook/pkg/operator/k8sutil"
|
||||
"github.com/rook/rook/pkg/util/exec"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
kerrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
@@ -178,7 +179,7 @@ func (c *Cluster) createSelfSignedCert() (bool, error) {
|
||||
|
||||
// retry a few times in the case that the mgr module is not ready to accept commands
|
||||
for i := 0; i < 5; i++ {
|
||||
_, err := client.NewCephCommand(c.context, c.clusterInfo, args).RunWithTimeout(client.CephCommandTimeout)
|
||||
_, err := client.NewCephCommand(c.context, c.clusterInfo, args).RunWithTimeout(exec.CephCommandTimeout)
|
||||
if err == context.DeadlineExceeded {
|
||||
logger.Warning("cert creation timed out. trying again")
|
||||
continue
|
||||
@@ -252,7 +253,7 @@ func (c *Cluster) setLoginCredentials(password string) error {
|
||||
}
|
||||
|
||||
_, err := client.ExecuteCephCommandWithRetry(func() (string, []byte, error) {
|
||||
output, err := client.NewCephCommand(c.context, c.clusterInfo, args).RunWithTimeout(client.CephCommandTimeout)
|
||||
output, err := client.NewCephCommand(c.context, c.clusterInfo, args).RunWithTimeout(exec.CephCommandTimeout)
|
||||
return "set dashboard creds", output, err
|
||||
}, c.exitCode, 5, invalidArgErrorCode, dashboardInitWaitTime)
|
||||
if err != nil {
|
||||
|
||||
@@ -22,6 +22,7 @@ import (
|
||||
|
||||
"github.com/pkg/errors"
|
||||
"github.com/rook/rook/pkg/daemon/ceph/client"
|
||||
"github.com/rook/rook/pkg/util/exec"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -53,7 +54,7 @@ func (c *Cluster) setRookOrchestratorBackend() error {
|
||||
// retry a few times in the case that the mgr module is not ready to accept commands
|
||||
_, err := client.ExecuteCephCommandWithRetry(func() (string, []byte, error) {
|
||||
args := []string{orchestratorCLIName, "set", "backend", "rook"}
|
||||
output, err := client.NewCephCommand(c.context, c.clusterInfo, args).RunWithTimeout(client.CephCommandTimeout)
|
||||
output, err := client.NewCephCommand(c.context, c.clusterInfo, args).RunWithTimeout(exec.CephCommandTimeout)
|
||||
return "set rook backend", output, err
|
||||
}, c.exitCode, 5, invalidArgErrorCode, orchestratorInitWaitTime)
|
||||
if err != nil {
|
||||
|
||||
@@ -26,6 +26,7 @@ import (
|
||||
"github.com/rook/rook/pkg/apis/rook.io"
|
||||
"github.com/rook/rook/pkg/operator/ceph/cluster/mon"
|
||||
"github.com/rook/rook/pkg/operator/ceph/config"
|
||||
"github.com/rook/rook/pkg/operator/ceph/config/keyring"
|
||||
"github.com/rook/rook/pkg/operator/ceph/controller"
|
||||
"github.com/rook/rook/pkg/operator/k8sutil"
|
||||
apps "k8s.io/api/apps/v1"
|
||||
@@ -35,12 +36,20 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
podIPEnvVar = "ROOK_POD_IP"
|
||||
serviceMetricName = "http-metrics"
|
||||
podIPEnvVar = "ROOK_POD_IP"
|
||||
serviceMetricName = "http-metrics"
|
||||
CommandProxyInitContainerName = "cmd-proxy"
|
||||
)
|
||||
|
||||
func (c *Cluster) makeDeployment(mgrConfig *mgrConfig) (*apps.Deployment, error) {
|
||||
logger.Debugf("mgrConfig: %+v", mgrConfig)
|
||||
|
||||
volumes := controller.DaemonVolumes(mgrConfig.DataPathMap, mgrConfig.ResourceName)
|
||||
if c.spec.Network.IsMultus() {
|
||||
adminKeyringVol, _ := keyring.Volume().Admin(), keyring.VolumeMount().Admin()
|
||||
volumes = append(volumes, adminKeyringVol)
|
||||
}
|
||||
|
||||
podSpec := v1.PodTemplateSpec{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: mgrConfig.ResourceName,
|
||||
@@ -55,7 +64,7 @@ func (c *Cluster) makeDeployment(mgrConfig *mgrConfig) (*apps.Deployment, error)
|
||||
},
|
||||
ServiceAccountName: serviceAccountName,
|
||||
RestartPolicy: v1.RestartPolicyAlways,
|
||||
Volumes: controller.DaemonVolumes(mgrConfig.DataPathMap, mgrConfig.ResourceName),
|
||||
Volumes: volumes,
|
||||
HostNetwork: c.spec.Network.IsHost(),
|
||||
PriorityClassName: cephv1.GetMgrPriorityClassName(c.spec.PriorityClassNames),
|
||||
},
|
||||
@@ -91,6 +100,7 @@ func (c *Cluster) makeDeployment(mgrConfig *mgrConfig) (*apps.Deployment, error)
|
||||
if err := k8sutil.ApplyMultus(c.spec.Network, &podSpec.ObjectMeta); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
podSpec.Spec.Containers = append(podSpec.Spec.Containers, c.makeCmdProxySidecarContainer(mgrConfig))
|
||||
}
|
||||
|
||||
cephv1.GetMgrAnnotations(c.spec.Annotations).ApplyToObjectMeta(&podSpec.ObjectMeta)
|
||||
@@ -228,6 +238,22 @@ func (c *Cluster) makeMgrSidecarContainer(mgrConfig *mgrConfig) v1.Container {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Cluster) makeCmdProxySidecarContainer(mgrConfig *mgrConfig) v1.Container {
|
||||
_, adminKeyringVolMount := keyring.Volume().Admin(), keyring.VolumeMount().Admin()
|
||||
container := v1.Container{
|
||||
Name: CommandProxyInitContainerName,
|
||||
Command: []string{"sleep"},
|
||||
Args: []string{"infinity"},
|
||||
Image: c.spec.CephVersion.Image,
|
||||
VolumeMounts: append(controller.DaemonVolumeMounts(mgrConfig.DataPathMap, mgrConfig.ResourceName), adminKeyringVolMount),
|
||||
Env: append(controller.DaemonEnvVars(c.spec.CephVersion.Image), v1.EnvVar{Name: "CEPH_ARGS", Value: fmt.Sprintf("-m $(ROOK_CEPH_MON_HOST) -k %s", keyring.VolumeMount().AdminKeyringFilePath())}),
|
||||
Resources: cephv1.GetMgrResources(c.spec.Resources),
|
||||
SecurityContext: controller.PodSecurityContext(),
|
||||
}
|
||||
|
||||
return container
|
||||
}
|
||||
|
||||
func getDefaultMgrLivenessProbe() *v1.Probe {
|
||||
return &v1.Probe{
|
||||
Handler: v1.Handler{
|
||||
|
||||
@@ -62,21 +62,37 @@ func TestPodSpec(t *testing.T) {
|
||||
DataPathMap: config.NewStatelessDaemonDataPathMap(config.MgrType, "a", "rook-ceph", "/var/lib/rook/"),
|
||||
}
|
||||
|
||||
d, err := c.makeDeployment(&mgrTestConfig)
|
||||
assert.NoError(t, err)
|
||||
t.Run("traditional deployment", func(t *testing.T) {
|
||||
d, err := c.makeDeployment(&mgrTestConfig)
|
||||
assert.NoError(t, err)
|
||||
|
||||
// Deployment should have Ceph labels
|
||||
test.AssertLabelsContainCephRequirements(t, d.ObjectMeta.Labels,
|
||||
config.MgrType, "a", AppName, "ns")
|
||||
// Deployment should have Ceph labels
|
||||
test.AssertLabelsContainCephRequirements(t, d.ObjectMeta.Labels,
|
||||
config.MgrType, "a", AppName, "ns")
|
||||
|
||||
podTemplate := test.NewPodTemplateSpecTester(t, &d.Spec.Template)
|
||||
podTemplate.Spec().Containers().RequireAdditionalEnvVars(
|
||||
"ROOK_OPERATOR_NAMESPACE", "ROOK_CEPH_CLUSTER_CRD_VERSION",
|
||||
"ROOK_CEPH_CLUSTER_CRD_NAME")
|
||||
podTemplate.RunFullSuite(config.MgrType, "a", AppName, "ns", "ceph/ceph:myceph",
|
||||
"200", "100", "500", "250", /* resources */
|
||||
"my-priority-class")
|
||||
assert.Equal(t, 2, len(d.Spec.Template.Annotations))
|
||||
podTemplate := test.NewPodTemplateSpecTester(t, &d.Spec.Template)
|
||||
podTemplate.Spec().Containers().RequireAdditionalEnvVars(
|
||||
"ROOK_OPERATOR_NAMESPACE", "ROOK_CEPH_CLUSTER_CRD_VERSION",
|
||||
"ROOK_CEPH_CLUSTER_CRD_NAME")
|
||||
podTemplate.RunFullSuite(config.MgrType, "a", AppName, "ns", "ceph/ceph:myceph",
|
||||
"200", "100", "500", "250", /* resources */
|
||||
"my-priority-class")
|
||||
assert.Equal(t, 2, len(d.Spec.Template.Annotations))
|
||||
assert.Equal(t, 1, len(d.Spec.Template.Spec.Containers))
|
||||
assert.Equal(t, 5, len(d.Spec.Template.Spec.Containers[0].VolumeMounts))
|
||||
})
|
||||
|
||||
t.Run("deployment with multus with new sidecar proxy command container", func(t *testing.T) {
|
||||
c.spec.Network.Provider = "multus"
|
||||
d, err := c.makeDeployment(&mgrTestConfig)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 3, len(d.Spec.Template.Annotations)) // Multus annotations
|
||||
assert.Equal(t, 2, len(d.Spec.Template.Spec.Containers)) // mgr pod + sidecar
|
||||
assert.Equal(t, CommandProxyInitContainerName, d.Spec.Template.Spec.Containers[1].Name) // sidecar pod
|
||||
assert.Equal(t, 6, len(d.Spec.Template.Spec.Containers[1].VolumeMounts)) // + admin keyring
|
||||
assert.Equal(t, "CEPH_ARGS", d.Spec.Template.Spec.Containers[1].Env[len(d.Spec.Template.Spec.Containers[1].Env)-1].Name) // connection info to the cluster
|
||||
assert.Equal(t, "-m $(ROOK_CEPH_MON_HOST) -k /etc/ceph/admin-keyring-store/keyring", d.Spec.Template.Spec.Containers[1].Env[len(d.Spec.Template.Spec.Containers[1].Env)-1].Value) // connection info to the cluster
|
||||
})
|
||||
}
|
||||
|
||||
func TestServiceSpec(t *testing.T) {
|
||||
|
||||
@@ -27,6 +27,7 @@ import (
|
||||
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
|
||||
"github.com/rook/rook/pkg/clusterd"
|
||||
cephclient "github.com/rook/rook/pkg/daemon/ceph/client"
|
||||
"github.com/rook/rook/pkg/operator/ceph/cluster/mgr"
|
||||
"github.com/rook/rook/pkg/util/exec"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
@@ -34,14 +35,15 @@ import (
|
||||
|
||||
// Context holds the context for the object store.
|
||||
type Context struct {
|
||||
Context *clusterd.Context
|
||||
clusterInfo *cephclient.ClusterInfo
|
||||
Name string
|
||||
UID string
|
||||
Endpoint string
|
||||
Realm string
|
||||
ZoneGroup string
|
||||
Zone string
|
||||
Context *clusterd.Context
|
||||
clusterInfo *cephclient.ClusterInfo
|
||||
CephClusterSpec cephv1.ClusterSpec
|
||||
Name string
|
||||
UID string
|
||||
Endpoint string
|
||||
Realm string
|
||||
ZoneGroup string
|
||||
Zone string
|
||||
}
|
||||
|
||||
// AdminOpsContext holds the object store context as well as information for connecting to the admin
|
||||
@@ -161,13 +163,21 @@ func extractJSON(output string) (string, error) {
|
||||
|
||||
// RunAdminCommandNoMultisite is for running radosgw-admin commands in scenarios where an object-store has not been created yet or for commands on the realm or zonegroup (ex: radosgw-admin zonegroup get)
|
||||
// This function times out after a fixed interval if no response is received.
|
||||
// The function will return a Kubernetes error "NotFound" when exec fails when the pod does not exist
|
||||
func RunAdminCommandNoMultisite(c *Context, expectJSON bool, args ...string) (string, error) {
|
||||
command, args := cephclient.FinalizeCephCommandArgs("radosgw-admin", c.clusterInfo, args, c.Context.ConfigDir)
|
||||
var output, stderr string
|
||||
var err error
|
||||
|
||||
// If Multus is enabled we proxy all the command to the mgr sidecar
|
||||
if c.CephClusterSpec.Network.IsMultus() {
|
||||
output, stderr, err = c.Context.RemoteExecutor.ExecCommandInContainerWithFullOutputWithTimeout(mgr.AppName, mgr.CommandProxyInitContainerName, c.clusterInfo.Namespace, append([]string{"radosgw-admin"}, args...)...)
|
||||
} else {
|
||||
command, args := cephclient.FinalizeCephCommandArgs("radosgw-admin", c.clusterInfo, args, c.Context.ConfigDir)
|
||||
output, err = c.Context.Executor.ExecuteCommandWithTimeout(exec.CephCommandTimeout, command, args...)
|
||||
}
|
||||
|
||||
// start the rgw admin command
|
||||
output, err := c.Context.Executor.ExecuteCommandWithTimeout(cephclient.CephCommandTimeout, command, args...)
|
||||
if err != nil {
|
||||
return output, err
|
||||
return fmt.Sprintf("%s. %s", output, stderr), err
|
||||
}
|
||||
if expectJSON {
|
||||
match, err := extractJSON(output)
|
||||
@@ -203,7 +213,7 @@ func runAdminCommand(c *Context, expectJSON bool, args ...string) (string, error
|
||||
// installed in Rook operator and RGW version in Ceph cluster (#7573)
|
||||
result, err := RunAdminCommandNoMultisite(c, expectJSON, args...)
|
||||
if err != nil && isFifoFileIOError(err) {
|
||||
logger.Debug("retrying 'radosgw-admin' command with OMAP backend to work around FIFO file I/O issue")
|
||||
logger.Debugf("retrying 'radosgw-admin' command with OMAP backend to work around FIFO file I/O issue. %v", result)
|
||||
|
||||
// We can either run 'ceph --version' to determine the Ceph version running in the operator
|
||||
// and then pick a flag to use, or we can just try to use both flags and return the one that
|
||||
@@ -224,7 +234,7 @@ func runAdminCommand(c *Context, expectJSON bool, args ...string) (string, error
|
||||
func isFifoFileIOError(err error) bool {
|
||||
exitCode, extractErr := exec.ExtractExitCode(err)
|
||||
if extractErr != nil {
|
||||
logger.Errorf("failed to determine return code of 'radosgw-admin' command. assuming this could be a FIFO file I/O issue. %#v", err)
|
||||
logger.Errorf("failed to determine return code of 'radosgw-admin' command. assuming this could be a FIFO file I/O issue. %#v", extractErr)
|
||||
return true
|
||||
}
|
||||
// exit code 5 (EIO) is returned when there is a FIFO file I/O issue
|
||||
@@ -234,7 +244,7 @@ func isFifoFileIOError(err error) bool {
|
||||
func isInvalidFlagError(err error) bool {
|
||||
exitCode, extractErr := exec.ExtractExitCode(err)
|
||||
if extractErr != nil {
|
||||
logger.Errorf("failed to determine return code of 'radosgw-admin' command. assuming this could be an invalid flag error. %#v", err)
|
||||
logger.Errorf("failed to determine return code of 'radosgw-admin' command. assuming this could be an invalid flag error. %#v", extractErr)
|
||||
}
|
||||
// exit code 22 (EINVAL) is returned when there is an invalid flag
|
||||
// it's also returned from some other failures, but this should be rare for Rook
|
||||
|
||||
@@ -321,7 +321,10 @@ func (r *ReconcileCephObjectStore) reconcile(request reconcile.Request) (reconci
|
||||
|
||||
// CREATE/UPDATE
|
||||
_, err = r.reconcileCreateObjectStore(cephObjectStore, request.NamespacedName, cephCluster.Spec)
|
||||
if err != nil {
|
||||
if err != nil && kerrors.IsNotFound(err) {
|
||||
logger.Info(opcontroller.OperatorNotInitializedMessage)
|
||||
return opcontroller.WaitForRequeueIfOperatorNotInitialized, cephObjectStore, nil
|
||||
} else if err != nil {
|
||||
result, err := r.setFailedStatus(request.NamespacedName, "failed to create object store deployments", err)
|
||||
return result, cephObjectStore, err
|
||||
}
|
||||
@@ -348,6 +351,7 @@ func (r *ReconcileCephObjectStore) reconcileCreateObjectStore(cephObjectStore *c
|
||||
}
|
||||
objContext := NewContext(r.context, r.clusterInfo, cephObjectStore.Name)
|
||||
objContext.UID = string(cephObjectStore.UID)
|
||||
objContext.CephClusterSpec = cluster
|
||||
|
||||
var err error
|
||||
|
||||
@@ -413,7 +417,9 @@ func (r *ReconcileCephObjectStore) reconcileCreateObjectStore(cephObjectStore *c
|
||||
// Reconcile Multisite Creation
|
||||
logger.Infof("setting multisite settings for object store %q", cephObjectStore.Name)
|
||||
err = setMultisite(objContext, cephObjectStore, serviceIP)
|
||||
if err != nil {
|
||||
if err != nil && kerrors.IsNotFound(err) {
|
||||
return reconcile.Result{}, err
|
||||
} else if err != nil {
|
||||
return r.setFailedStatus(namespacedName, "failed to configure multisite for object store", err)
|
||||
}
|
||||
|
||||
@@ -440,7 +446,8 @@ func (r *ReconcileCephObjectStore) reconcileCephZone(store *cephv1.CephObjectSto
|
||||
|
||||
_, err := RunAdminCommandNoMultisite(objContext, true, "zone", "get", realmArg, zoneGroupArg, zoneArg)
|
||||
if err != nil {
|
||||
if code, ok := exec.ExitStatus(err); ok && code == int(syscall.ENOENT) {
|
||||
// ENOENT mean “No such file or directory”
|
||||
if code, err := exec.ExtractExitCode(err); err == nil && code == int(syscall.ENOENT) {
|
||||
return waitForRequeueIfObjectStoreNotReady, errors.Wrapf(err, "ceph zone %q not found", store.Spec.Zone.Name)
|
||||
} else {
|
||||
return waitForRequeueIfObjectStoreNotReady, errors.Wrapf(err, "radosgw-admin zone get failed with code %d", code)
|
||||
|
||||
@@ -149,7 +149,7 @@ func (c *bucketChecker) checkObjectStoreHealth() error {
|
||||
userConfig := genUserCheckerConfig(c.objContext.UID)
|
||||
|
||||
// Create checker user
|
||||
logger.Debugf("creating s3 user object %q for object store %q", userConfig.ID, c.namespacedName.Name)
|
||||
logger.Debugf("creating s3 user object %q for object store %q health check", userConfig.ID, c.namespacedName.Name)
|
||||
var user admin.User
|
||||
user, err = opsCtx.AdminOpsClient.CreateUser(context.TODO(), userConfig)
|
||||
if err != nil {
|
||||
|
||||
@@ -37,6 +37,7 @@ import (
|
||||
"github.com/rook/rook/pkg/util/exec"
|
||||
"golang.org/x/sync/errgroup"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
kerrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
)
|
||||
|
||||
@@ -120,6 +121,10 @@ func removeObjectStoreFromMultisite(objContext *Context, spec cephv1.ObjectStore
|
||||
if zoneIsMaster {
|
||||
_, err = RunAdminCommandNoMultisite(objContext, false, "zonegroup", "modify", realmArg, zoneGroupArg, endpointArg)
|
||||
if err != nil {
|
||||
|
||||
if kerrors.IsNotFound(err) {
|
||||
return err
|
||||
}
|
||||
return errors.Wrapf(err, "failed to remove object store %q endpoint from rgw zone group %q", objContext.Name, objContext.ZoneGroup)
|
||||
}
|
||||
logger.Debugf("endpoint %q was removed from zone group %q. the remaining endpoints in the zone group are %q", objContext.Endpoint, objContext.ZoneGroup, zoneEndpoints)
|
||||
@@ -217,6 +222,11 @@ func checkZoneIsMaster(objContext *Context) (bool, error) {
|
||||
|
||||
zoneGroupJson, err := RunAdminCommandNoMultisite(objContext, true, "zonegroup", "get", realmArg, zoneGroupArg)
|
||||
if err != nil {
|
||||
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
||||
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
||||
if kerrors.IsNotFound(err) {
|
||||
return false, err
|
||||
}
|
||||
return false, errors.Wrap(err, "failed to get rgw zone group")
|
||||
}
|
||||
zoneGroupOutput, err := DecodeZoneGroupConfig(zoneGroupJson)
|
||||
@@ -227,6 +237,11 @@ func checkZoneIsMaster(objContext *Context) (bool, error) {
|
||||
|
||||
zoneOutput, err := RunAdminCommandNoMultisite(objContext, true, "zone", "get", realmArg, zoneGroupArg, zoneArg)
|
||||
if err != nil {
|
||||
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
||||
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
||||
if kerrors.IsNotFound(err) {
|
||||
return false, err
|
||||
}
|
||||
return false, errors.Wrap(err, "failed to get rgw zone")
|
||||
}
|
||||
zoneID, err := decodeID(zoneOutput)
|
||||
@@ -251,6 +266,11 @@ func checkZoneGroupIsMaster(objContext *Context) (bool, error) {
|
||||
|
||||
zoneGroupOutput, err := RunAdminCommandNoMultisite(objContext, true, "zonegroup", "get", realmArg, zoneGroupArg)
|
||||
if err != nil {
|
||||
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
||||
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
||||
if kerrors.IsNotFound(err) {
|
||||
return false, err
|
||||
}
|
||||
return false, errors.Wrap(err, "failed to get rgw zone group")
|
||||
}
|
||||
|
||||
@@ -311,7 +331,9 @@ func getZoneEndpoints(objContext *Context, serviceEndpoint string) ([]string, er
|
||||
|
||||
zoneGroupOutput, err := RunAdminCommandNoMultisite(objContext, true, "zonegroup", "get", realmArg, zoneGroupArg)
|
||||
if err != nil {
|
||||
return []string{}, errors.Wrap(err, "failed to get rgw zone group")
|
||||
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
||||
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
||||
return []string{}, errorOrIsNotFound(err, "failed to get rgw zone group %q", objContext.Name)
|
||||
}
|
||||
zoneGroupJson, err := DecodeZoneGroupConfig(zoneGroupOutput)
|
||||
if err != nil {
|
||||
@@ -344,45 +366,48 @@ func createMultisite(objContext *Context, endpointArg string) error {
|
||||
// create the realm if it doesn't exist yet
|
||||
output, err := RunAdminCommandNoMultisite(objContext, true, "realm", "get", realmArg)
|
||||
if err != nil {
|
||||
if code, ok := exec.ExitStatus(err); ok && code == int(syscall.ENOENT) {
|
||||
// ENOENT means “No such file or directory”
|
||||
if code, err := exec.ExtractExitCode(err); err == nil && code == int(syscall.ENOENT) {
|
||||
updatePeriod = true
|
||||
output, err = RunAdminCommandNoMultisite(objContext, false, "realm", "create", realmArg)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed to create ceph realm %q, for reason %q", objContext.ZoneGroup, output)
|
||||
return errorOrIsNotFound(err, "failed to create ceph realm %q, for reason %q", objContext.ZoneGroup, output)
|
||||
}
|
||||
logger.Debugf("created realm %v", objContext.Realm)
|
||||
} else {
|
||||
return errors.Wrapf(err, "radosgw-admin realm get failed with code %d, for reason %q", code, output)
|
||||
return errorOrIsNotFound(err, "radosgw-admin realm get failed with code %d, for reason %q. %v", strconv.Itoa(code), output, string(kerrors.ReasonForError(err)))
|
||||
}
|
||||
}
|
||||
|
||||
// create the zonegroup if it doesn't exist yet
|
||||
output, err = RunAdminCommandNoMultisite(objContext, true, "zonegroup", "get", realmArg, zoneGroupArg)
|
||||
if err != nil {
|
||||
if code, ok := exec.ExitStatus(err); ok && code == int(syscall.ENOENT) {
|
||||
// ENOENT means “No such file or directory”
|
||||
if code, err := exec.ExtractExitCode(err); err == nil && code == int(syscall.ENOENT) {
|
||||
updatePeriod = true
|
||||
output, err = RunAdminCommandNoMultisite(objContext, false, "zonegroup", "create", "--master", realmArg, zoneGroupArg, endpointArg)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed to create ceph zone group %q, for reason %q", objContext.ZoneGroup, output)
|
||||
return errorOrIsNotFound(err, "failed to create ceph zone group %q, for reason %q", objContext.ZoneGroup, output)
|
||||
}
|
||||
logger.Debugf("created zone group %v", objContext.ZoneGroup)
|
||||
} else {
|
||||
return errors.Wrapf(err, "radosgw-admin zonegroup get failed with code %d, for reason %q", code, output)
|
||||
return errorOrIsNotFound(err, "radosgw-admin zonegroup get failed with code %d, for reason %q", strconv.Itoa(code), output)
|
||||
}
|
||||
}
|
||||
|
||||
// create the zone if it doesn't exist yet
|
||||
output, err = runAdminCommand(objContext, true, "zone", "get")
|
||||
if err != nil {
|
||||
if code, ok := exec.ExitStatus(err); ok && code == int(syscall.ENOENT) {
|
||||
// ENOENT means “No such file or directory”
|
||||
if code, err := exec.ExtractExitCode(err); err == nil && code == int(syscall.ENOENT) {
|
||||
updatePeriod = true
|
||||
output, err = runAdminCommand(objContext, false, "zone", "create", "--master", endpointArg)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed to create ceph zone %q, for reason %q", objContext.Zone, output)
|
||||
return errorOrIsNotFound(err, "failed to create ceph zone %q, for reason %q", objContext.Zone, output)
|
||||
}
|
||||
logger.Debugf("created zone %v", objContext.Zone)
|
||||
} else {
|
||||
return errors.Wrapf(err, "radosgw-admin zone get failed with code %d, for reason %q", code, output)
|
||||
return errorOrIsNotFound(err, "radosgw-admin zone get failed with code %d, for reason %q", strconv.Itoa(code), output)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -390,7 +415,7 @@ func createMultisite(objContext *Context, endpointArg string) error {
|
||||
// the period will help notify other zones of changes if there are multi-zones
|
||||
_, err := runAdminCommand(objContext, false, "period", "update", "--commit")
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "failed to update period")
|
||||
return errorOrIsNotFound(err, "failed to update period")
|
||||
}
|
||||
logger.Debugf("updated period for realm %v", objContext.Realm)
|
||||
}
|
||||
@@ -416,7 +441,7 @@ func joinMultisite(objContext *Context, endpointArg, zoneEndpoints, namespace st
|
||||
// endpoints that are part of a master zone are supposed to be the endpoints for a zone group
|
||||
_, err := RunAdminCommandNoMultisite(objContext, false, "zonegroup", "modify", realmArg, zoneGroupArg, endpointArg)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed to add object store %q in rgw zone group %q", objContext.Name, objContext.ZoneGroup)
|
||||
return errorOrIsNotFound(err, "failed to add object store %q in rgw zone group %q", objContext.Name, objContext.ZoneGroup)
|
||||
}
|
||||
logger.Debugf("endpoints for zonegroup %q are now %q", objContext.ZoneGroup, zoneEndpoints)
|
||||
|
||||
@@ -428,14 +453,14 @@ func joinMultisite(objContext *Context, endpointArg, zoneEndpoints, namespace st
|
||||
}
|
||||
_, err = RunAdminCommandNoMultisite(objContext, false, "zone", "modify", realmArg, zoneGroupArg, zoneArg, endpointArg)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed to add object store %q in rgw zone %q", objContext.Name, objContext.Zone)
|
||||
return errorOrIsNotFound(err, "failed to add object store %q in rgw zone %q", objContext.Name, objContext.Zone)
|
||||
}
|
||||
logger.Debugf("endpoints for zone %q are now %q", objContext.Zone, zoneEndpoints)
|
||||
|
||||
// the period will help notify other zones of changes if there are multi-zones
|
||||
_, err = RunAdminCommandNoMultisite(objContext, false, "period", "update", "--commit", realmArg, zoneGroupArg, zoneArg)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "failed to update period")
|
||||
return errorOrIsNotFound(err, "failed to update period")
|
||||
}
|
||||
logger.Infof("added object store %q to realm %q, zonegroup %q, zone %q", objContext.Name, objContext.Realm, objContext.ZoneGroup, objContext.Zone)
|
||||
|
||||
@@ -474,11 +499,11 @@ func createSystemUser(objContext *Context, namespace string) error {
|
||||
displayNameArg := fmt.Sprintf("--display-name=%s.user", objContext.Realm)
|
||||
output, err = RunAdminCommandNoMultisite(objContext, false, "user", "create", realmArg, zoneGroupArg, zoneArg, uidArg, displayNameArg, accessKeyArg, secretKeyArg, systemArg)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed to create realm system user %q for reason: %q", uid, output)
|
||||
return errorOrIsNotFound(err, "failed to create realm system user %q for reason: %q", uid, output)
|
||||
}
|
||||
logger.Debugf("created realm system user %v", uid)
|
||||
} else {
|
||||
return errors.Wrapf(err, "radosgw-admin user info for system user failed with code %d and output %q", code, output)
|
||||
return errorOrIsNotFound(err, "radosgw-admin user info for system user failed with code %d and output %q", strconv.Itoa(code), output)
|
||||
}
|
||||
|
||||
return nil
|
||||
@@ -510,7 +535,7 @@ func setMultisite(objContext *Context, store *cephv1.CephObjectStore, serviceIP
|
||||
endpointArg := fmt.Sprintf("--endpoints=%s", serviceEndpoint)
|
||||
err := createMultisite(objContext, endpointArg)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed create ceph multisite for object-store %q", objContext.Name)
|
||||
return errorOrIsNotFound(err, "failed create ceph multisite for object-store %q", objContext.Name)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -563,6 +588,11 @@ func DecodeZoneGroupConfig(data string) (zoneGroupType, error) {
|
||||
func getObjectStores(context *Context) ([]string, error) {
|
||||
output, err := RunAdminCommandNoMultisite(context, true, "realm", "list")
|
||||
if err != nil {
|
||||
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
||||
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
||||
if kerrors.IsNotFound(err) {
|
||||
return []string{}, err
|
||||
}
|
||||
// exit status 2 indicates the object store does not exist, so return nothing
|
||||
if strings.Index(err.Error(), "exit status 2") == 0 {
|
||||
return []string{}, nil
|
||||
@@ -881,7 +911,7 @@ func enableRGWDashboard(context *Context) error {
|
||||
// starting in ceph v15.2.8. We run it in a goroutine until the fix
|
||||
// is found. We expect the ceph command to timeout so at least the goroutine exits.
|
||||
logger.Info("setting the dashboard api secret key")
|
||||
_, err = cephCmd.RunWithTimeout(cephclient.CephCommandTimeout)
|
||||
_, err = cephCmd.RunWithTimeout(exec.CephCommandTimeout)
|
||||
if err != nil {
|
||||
logger.Errorf("failed to set user %q secretkey. %v", DashboardUser, err)
|
||||
}
|
||||
@@ -913,16 +943,25 @@ func disableRGWDashboard(context *Context) {
|
||||
|
||||
args := []string{"dashboard", "reset-rgw-api-access-key"}
|
||||
cephCmd := cephclient.NewCephCommand(context.Context, context.clusterInfo, args)
|
||||
_, err = cephCmd.RunWithTimeout(cephclient.CephCommandTimeout)
|
||||
_, err = cephCmd.RunWithTimeout(exec.CephCommandTimeout)
|
||||
if err != nil {
|
||||
logger.Warningf("failed to reset user accesskey for user %q. %v", DashboardUser, err)
|
||||
}
|
||||
|
||||
args = []string{"dashboard", "reset-rgw-api-secret-key"}
|
||||
cephCmd = cephclient.NewCephCommand(context.Context, context.clusterInfo, args)
|
||||
_, err = cephCmd.RunWithTimeout(cephclient.CephCommandTimeout)
|
||||
_, err = cephCmd.RunWithTimeout(exec.CephCommandTimeout)
|
||||
if err != nil {
|
||||
logger.Warningf("failed to reset user secretkey for user %q. %v", DashboardUser, err)
|
||||
}
|
||||
logger.Info("done disabling the dashboard api secret key")
|
||||
}
|
||||
|
||||
func errorOrIsNotFound(err error, msg string, args ...string) error {
|
||||
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
||||
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
||||
if kerrors.IsNotFound(err) {
|
||||
return err
|
||||
}
|
||||
return errors.Wrapf(err, msg, args)
|
||||
}
|
||||
|
||||
+30
-6
@@ -25,11 +25,19 @@ import (
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"os/exec"
|
||||
"reflect"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/coreos/pkg/capnslog"
|
||||
"github.com/pkg/errors"
|
||||
kerrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
kexec "k8s.io/utils/exec"
|
||||
)
|
||||
|
||||
var (
|
||||
CephCommandTimeout = 15 * time.Second
|
||||
)
|
||||
|
||||
// Executor is the main interface for all the exec commands
|
||||
@@ -44,8 +52,7 @@ type Executor interface {
|
||||
}
|
||||
|
||||
// CommandExecutor is the type of the Executor
|
||||
type CommandExecutor struct {
|
||||
}
|
||||
type CommandExecutor struct{}
|
||||
|
||||
// ExecuteCommand starts a process and wait for its completion
|
||||
func (c *CommandExecutor) ExecuteCommand(command string, arg ...string) error {
|
||||
@@ -322,10 +329,27 @@ func assertErrorType(err error) string {
|
||||
}
|
||||
|
||||
// ExtractExitCode attempts to get the exit code from the error returned by an Executor function.
|
||||
// This should also work for any errors returned by the golang os/exec package.
|
||||
// This should also work for any errors returned by the golang os/exec package and "k8s.io/utils/exec"
|
||||
func ExtractExitCode(err error) (int, error) {
|
||||
if ee, ok := err.(*exec.ExitError); ok {
|
||||
return ee.ExitCode(), nil
|
||||
switch errType := err.(type) {
|
||||
case *exec.ExitError:
|
||||
return errType.ExitCode(), nil
|
||||
|
||||
case *kexec.CodeExitError:
|
||||
return errType.ExitStatus(), nil
|
||||
|
||||
case *kerrors.StatusError:
|
||||
return int(errType.ErrStatus.Code), nil
|
||||
|
||||
default:
|
||||
logger.Debugf(err.Error())
|
||||
// This is ugly but I don't know why the type assertion does not work...
|
||||
// Whatever I've tried I can see the type "exec.CodeExitError" but none of the "case" nor other attempts with "errors.As()" worked :(
|
||||
// So I'm parsing the Error string until we have a solution
|
||||
if strings.Contains(err.Error(), "command terminated with exit code") {
|
||||
a := strings.SplitAfter(err.Error(), "command terminated with exit code")
|
||||
return strconv.Atoi(strings.TrimSpace(a[1]))
|
||||
}
|
||||
return 0, errors.Errorf("error %#v is not an ExitError nor CodeExitError but is %v", err, reflect.TypeOf(err))
|
||||
}
|
||||
return 0, errors.Errorf("error %#v is not an ExitError", err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,135 @@
|
||||
/*
|
||||
Copyright 2021 The Rook Authors. All rights reserved.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
// Heavily inspired by https://github.com/kubernetes/kubernetes/blob/master/test/e2e/framework/exec_util.go
|
||||
|
||||
package exec
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/pkg/errors"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/kubernetes/scheme"
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/remotecommand"
|
||||
)
|
||||
|
||||
// ExecOptions passed to ExecWithOptions
|
||||
type ExecOptions struct {
|
||||
Command []string
|
||||
Namespace string
|
||||
PodName string
|
||||
ContainerName string
|
||||
Stdin io.Reader
|
||||
CaptureStdout bool
|
||||
CaptureStderr bool
|
||||
// If false, whitespace in std{err,out} will be removed.
|
||||
PreserveWhitespace bool
|
||||
}
|
||||
|
||||
// RemotePodCommandExecutor is an exec.Executor that execs every command in a remote container
|
||||
// This is especially useful when the CephCluster networking type is Multus and when the Operator pod
|
||||
// does not have the right network annotations.
|
||||
type RemotePodCommandExecutor struct {
|
||||
ClientSet kubernetes.Interface
|
||||
RestClient *rest.Config
|
||||
}
|
||||
|
||||
// ExecWithOptions executes a command in the specified container,
|
||||
// returning stdout, stderr and error. `options` allowed for
|
||||
// additional parameters to be passed.
|
||||
func (e *RemotePodCommandExecutor) ExecWithOptions(options ExecOptions) (string, string, error) {
|
||||
const tty = false
|
||||
|
||||
logger.Debugf("ExecWithOptions %+v", options)
|
||||
|
||||
req := e.ClientSet.CoreV1().RESTClient().Post().
|
||||
Resource("pods").
|
||||
Name(options.PodName).
|
||||
Namespace(options.Namespace).
|
||||
SubResource("exec").
|
||||
Param("container", options.ContainerName)
|
||||
req.VersionedParams(&v1.PodExecOptions{
|
||||
Container: options.ContainerName,
|
||||
Command: options.Command,
|
||||
Stdin: options.Stdin != nil,
|
||||
Stdout: options.CaptureStdout,
|
||||
Stderr: options.CaptureStderr,
|
||||
TTY: tty,
|
||||
}, scheme.ParameterCodec)
|
||||
|
||||
var stdout, stderr bytes.Buffer
|
||||
err := execute(http.MethodPost, req.URL(), e.RestClient, options.Stdin, &stdout, &stderr, tty)
|
||||
|
||||
if options.PreserveWhitespace {
|
||||
return stdout.String(), stderr.String(), err
|
||||
}
|
||||
return strings.TrimSpace(stdout.String()), strings.TrimSpace(stderr.String()), err
|
||||
}
|
||||
|
||||
// ExecCommandInContainerWithFullOutput executes a command in the
|
||||
// specified container and return stdout, stderr and error
|
||||
func (e *RemotePodCommandExecutor) ExecCommandInContainerWithFullOutput(appLabel, containerName, namespace string, cmd ...string) (string, string, error) {
|
||||
options := metav1.ListOptions{LabelSelector: fmt.Sprintf("app=%s", appLabel)}
|
||||
pods, err := e.ClientSet.CoreV1().Pods(namespace).List(context.TODO(), options)
|
||||
if err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
|
||||
if len(pods.Items) == 0 {
|
||||
return "", "", errors.Errorf("no pods found with selector %q", appLabel)
|
||||
}
|
||||
|
||||
return e.ExecWithOptions(ExecOptions{
|
||||
Command: cmd,
|
||||
Namespace: namespace,
|
||||
// Always pick the first pod, it's always 1 unless stretched cluster is enabled
|
||||
// TODO: if we have 2 pods we could try each result if the command fails to run due to a network partition-related error.
|
||||
PodName: pods.Items[0].Name,
|
||||
ContainerName: containerName,
|
||||
Stdin: nil,
|
||||
CaptureStdout: true,
|
||||
CaptureStderr: true,
|
||||
PreserveWhitespace: false,
|
||||
})
|
||||
}
|
||||
|
||||
func execute(method string, url *url.URL, config *rest.Config, stdin io.Reader, stdout, stderr io.Writer, tty bool) error {
|
||||
exec, err := remotecommand.NewSPDYExecutor(config, method, url)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return exec.Stream(remotecommand.StreamOptions{
|
||||
Stdin: stdin,
|
||||
Stdout: stdout,
|
||||
Stderr: stderr,
|
||||
Tty: tty,
|
||||
})
|
||||
}
|
||||
|
||||
func (e *RemotePodCommandExecutor) ExecCommandInContainerWithFullOutputWithTimeout(appLabel, containerName, namespace string, cmd ...string) (string, string, error) {
|
||||
return e.ExecCommandInContainerWithFullOutput(appLabel, containerName, namespace, append([]string{"timeout", strconv.Itoa(int(CephCommandTimeout.Seconds()))}, cmd...)...)
|
||||
}
|
||||
@@ -50,6 +50,7 @@ import (
|
||||
// K8sHelper is a helper for common kubectl commands
|
||||
type K8sHelper struct {
|
||||
executor *exec.CommandExecutor
|
||||
remoteExecutor *exec.RemotePodCommandExecutor
|
||||
Clientset *kubernetes.Clientset
|
||||
RookClientset *rookclient.Clientset
|
||||
RunningInCluster bool
|
||||
@@ -91,7 +92,12 @@ func CreateK8sHelper(t func() *testing.T) (*K8sHelper, error) {
|
||||
return nil, fmt.Errorf("failed to get rook clientset. %+v", err)
|
||||
}
|
||||
|
||||
h := &K8sHelper{executor: executor, Clientset: clientset, RookClientset: rookClientset, T: t}
|
||||
remoteExecutor := &exec.RemotePodCommandExecutor{
|
||||
ClientSet: clientset,
|
||||
RestClient: config,
|
||||
}
|
||||
|
||||
h := &K8sHelper{executor: executor, Clientset: clientset, RookClientset: rookClientset, T: t, remoteExecutor: remoteExecutor}
|
||||
if strings.Contains(config.Host, "//10.") {
|
||||
h.RunningInCluster = true
|
||||
}
|
||||
@@ -198,6 +204,28 @@ func (k8sh *K8sHelper) Exec(namespace, podName, command string, commandArgs []st
|
||||
return k8sh.ExecWithRetry(1, namespace, podName, command, commandArgs)
|
||||
}
|
||||
|
||||
func (k8sh *K8sHelper) ExecRemote(namespace, command string, commandArgs []string) (string, error) {
|
||||
return k8sh.ExecRemoteWithRetry(1, namespace, command, commandArgs)
|
||||
}
|
||||
|
||||
// ExecRemoteWithRetry will attempt to remotely (in toolbox) run a command "retries" times, waiting 3s between each call. Upon success, returns the output.
|
||||
func (k8sh *K8sHelper) ExecRemoteWithRetry(retries int, namespace, command string, commandArgs []string) (string, error) {
|
||||
var err error
|
||||
var output, stderr string
|
||||
cliFinal := append([]string{command}, commandArgs...)
|
||||
for i := 0; i < retries; i++ {
|
||||
output, stderr, err = k8sh.remoteExecutor.ExecCommandInContainerWithFullOutput("rook-ceph-tools", "rook-ceph-tools", namespace, cliFinal...)
|
||||
if err == nil {
|
||||
return output, nil
|
||||
}
|
||||
if i < retries-1 {
|
||||
logger.Warningf("remote command %v execution failed trying again... %v", cliFinal, kerrors.ReasonForError(err))
|
||||
time.Sleep(3 * time.Second)
|
||||
}
|
||||
}
|
||||
return "", fmt.Errorf("remote exec command %v failed on pod in namespace %s. %s. %s. %+v", cliFinal, namespace, output, stderr, err)
|
||||
}
|
||||
|
||||
// ExecWithRetry will attempt to run a command "retries" times, waiting 3s between each call. Upon success, returns the output.
|
||||
func (k8sh *K8sHelper) ExecWithRetry(retries int, namespace, podName, command string, commandArgs []string) (string, error) {
|
||||
var err error
|
||||
|
||||
@@ -502,13 +502,13 @@ func createFilesystemMountCephCredentials(helper *clients.TestClient, k8sh *util
|
||||
require.Nil(s.T(), err)
|
||||
// Mount CephFS in toolbox and create /foo directory on it
|
||||
logger.Info("Creating /foo directory on CephFS")
|
||||
_, err = k8sh.Exec(settings.Namespace, client.RunAllCephCommandsInToolboxPod, "mkdir", []string{"-p", utils.TestMountPath})
|
||||
_, err = k8sh.ExecRemote(settings.Namespace, "mkdir", []string{"-p", utils.TestMountPath})
|
||||
require.Nil(s.T(), err)
|
||||
_, err = k8sh.ExecWithRetry(10, settings.Namespace, client.RunAllCephCommandsInToolboxPod, "bash", []string{"-c", fmt.Sprintf("mount -t ceph -o mds_namespace=%s,name=admin,secret=$(grep key /etc/ceph/keyring | awk '{print $3}') $(grep mon_host /etc/ceph/ceph.conf | awk '{print $3}'):/ %s", filesystemName, utils.TestMountPath)})
|
||||
_, err = k8sh.ExecRemoteWithRetry(10, settings.Namespace, "bash", []string{"-c", fmt.Sprintf("mount -t ceph -o mds_namespace=%s,name=admin,secret=$(grep key /etc/ceph/keyring | awk '{print $3}') $(grep mon_host /etc/ceph/ceph.conf | awk '{print $3}'):/ %s", filesystemName, utils.TestMountPath)})
|
||||
require.Nil(s.T(), err)
|
||||
_, err = k8sh.Exec(settings.Namespace, client.RunAllCephCommandsInToolboxPod, "mkdir", []string{"-p", fmt.Sprintf("%s/foo", utils.TestMountPath)})
|
||||
_, err = k8sh.ExecRemote(settings.Namespace, "mkdir", []string{"-p", fmt.Sprintf("%s/foo", utils.TestMountPath)})
|
||||
require.Nil(s.T(), err)
|
||||
_, err = k8sh.Exec(settings.Namespace, client.RunAllCephCommandsInToolboxPod, "umount", []string{utils.TestMountPath})
|
||||
_, err = k8sh.ExecRemote(settings.Namespace, "umount", []string{utils.TestMountPath})
|
||||
require.Nil(s.T(), err)
|
||||
logger.Info("Created /foo directory on CephFS")
|
||||
|
||||
@@ -522,7 +522,7 @@ func createFilesystemMountCephCredentials(helper *clients.TestClient, k8sh *util
|
||||
),
|
||||
}
|
||||
logger.Infof("ceph credentials command args: %s", commandArgs[1])
|
||||
result, err := k8sh.Exec(settings.Namespace, client.RunAllCephCommandsInToolboxPod, "bash", commandArgs)
|
||||
result, err := k8sh.ExecRemote(settings.Namespace, "bash", commandArgs)
|
||||
logger.Infof("Ceph filesystem credentials output: %s", result)
|
||||
logger.Info("Created Ceph credentials")
|
||||
require.Nil(s.T(), err)
|
||||
|
||||
Reference in New Issue
Block a user