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:
Sébastien Han
2021-07-07 19:08:32 +02:00
parent baaea4a1ea
commit bd58790c31
21 changed files with 451 additions and 81 deletions
@@ -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
+3
View File
@@ -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)
+3 -3
View File
@@ -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
+9 -3
View File
@@ -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=
+3
View File
@@ -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
+4 -4
View File
@@ -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
}
+2 -1
View File
@@ -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",
+3 -2
View File
@@ -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 {
+29 -3
View File
@@ -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{
+29 -13
View File
@@ -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) {
+25 -15
View File
@@ -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
+10 -3
View File
@@ -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)
+1 -1
View File
@@ -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 {
+59 -20
View File
@@ -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
View File
@@ -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)
}
+135
View File
@@ -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...)...)
}
+29 -1
View File
@@ -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
+5 -5
View File
@@ -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)