ceph: configure a stretched cluster

In clusters where only two datacenters (or similar failure domains)
are available, a different mon and osd approach is needed to deal
with the network partitions or some other reason for one of the failure
domains going down. The Ceph stretched cluster makes the mons aware
of the failure domains by configuring one as the arbiter in a third
zone, while keeping two replicas of the data in each of the data
zoens.

Signed-off-by: Travis Nielsen <tnielsen@redhat.com>
This commit is contained in:
Travis Nielsen
2020-11-05 09:17:52 -07:00
parent 3b188a447a
commit b91f4211c9
30 changed files with 978 additions and 115 deletions
+74 -2
View File
@@ -7,10 +7,11 @@ indent: true
# Ceph Cluster CRD
Rook allows creation and customization of storage clusters through the custom resource definitions (CRDs).
There are two different modes to create your cluster, depending on whether storage can be dynamically provisioned on which to base the Ceph cluster.
There are primarily three different modes in which to create your cluster.
1. Specify [host paths and raw devices](#host-based-cluster)
2. Specify the storage class Rook should use to consume storage [via PVCs](#pvc-based-cluster)
2. Dynamically provision storage underneath Rook by specifying the storage class Rook should use to consume storage [via PVCs](#pvc-based-cluster)
3. Create a [Stretch cluster](#stretch-cluster) that distributes Ceph mons across three zones, while storage (OSDs) is only configured in two zones
Following is an example for each of these approaches. More examples are included [later in this doc](#samples).
@@ -91,6 +92,68 @@ spec:
For a more advanced scenario, such as adding a dedicated device you can refer to the [dedicated metadata device for OSD on PVC section](#dedicated-metadata-device-for-osd-on-pvc).
## Stretch Cluster
**Experimental Mode**
For environments that only have two failure domains available where data can be replicated, consider
the case where one failure domain is down and the data is still fully available in the
remaining failure domain. To support this scenario, Ceph has recently integrated support for "stretch" clusters.
Rook requires three zones. Two zones (A and B) will each run all types of Rook pods, which we call the "data" zones.
Two mons run in each of the two data zones, while two replicas of the data are in each zone for a total of four data replicas.
The third zone (arbiter) runs a single mon. No other Rook or Ceph daemons need to be run in the arbiter zone.
For this example, we assume the desired failure domain is a zone. Another failure domain can also be specified with a
known [topology node label](#osd-topology) which is already being used for OSD failure domains.
```yaml
apiVersion: ceph.rook.io/v1
kind: CephCluster
metadata:
name: rook-ceph
namespace: rook-ceph
spec:
dataDirHostPath: /var/lib/rook
mon:
# Five mons must be created for stretch mode
count: 5
allowMultiplePerNode: false
stretchCluster:
failureDomainLabel: topology.kubernetes.io/zone
subFailureDomain: host
zones:
- name: a
arbiter: true
- name: b
- name: c
cephVersion:
# Stretch cluster support upstream is only planned starting in Ceph Pacific.
# Until Pacific is released, the stretch cluster is **experimental**.
image: ceph/daemon-base:latest-master
allowUnsupported: true
# Either storageClassDeviceSets or the storage section can be specified for creating OSDs.
# This example uses all devices for simplicity.
storage:
useAllNodes: true
useAllDevices: true
deviceFilter: ""
# OSD placement is expected to include the non-arbiter zones
placement:
osd:
nodeAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
nodeSelectorTerms:
- matchExpressions:
- key: topology.kubernetes.io/zone
operator: In
values:
- b
- c
```
For more details, see the [Stretch Cluster design doc](https://github.com/rook/rook/blob/master/design/ceph/ceph-stretch-cluster.md).
## Settings
Settings can be specified at the global level to apply to the cluster as a whole, while other settings can be specified at more fine-grained levels. If any setting is unspecified, a suitable default will be used automatically.
@@ -187,6 +250,15 @@ A specific will contain a specific release of Ceph as well as security fixes fro
This setting only applies to new monitors that are created when the requested
number of monitors increases, or when a monitor fails and is recreated. An
[example CRD configuration is provided below](#using-pvc-storage-for-monitors).
* `stretchCluster`: The stretch cluster settings that define the zones (or other failure domain labels) across which to configure the cluster.
* `failureDomainLabel`: The label that is expected on each node where the cluster is expected to be deployed. The labels must be found
in the list of well-known [topology labels](#osd-topology).
* `subFailureDomain`: With a zone, the data replicas must be spread across OSDs in the subFailureDomain. The default is `host`.
* `zones`: The failure domain names where the Mons and OSDs are expected to be deployed. There must be **three zones** specified in the list.
This element is always named `zone` even if a non-default `failureDomainLabel` is specified. The elements have two values:
* `name`: The name of the zone, which is the value of the domain label.
* `arbiter`: Whether the zone is expected to be the arbiter zone which only runs a single mon. Exactly one zone must be labeled `true`.
The two zones that are not the arbiter zone are expected to have OSDs deployed.
If these settings are changed in the CRD the operator will update the number of mons during a periodic check of the mon health, which by default is every 45 seconds.
+3 -2
View File
@@ -51,9 +51,10 @@ the cluster. These examples represent a very small set of the different ways to
* `cluster-with-drive-groups.yaml`: This file contains example configurations for creating advanced
OSD layouts on nodes using Ceph Drive Groups.
[See docs for more](Documentation/ceph-cluster-crd.md#storage-selection-via-ceph-drive-groups)
* `cluster-external`: Connect to an [external Ceph cluster](ceph-cluster-crd.md#external-cluster) with minimal access to monitor the health of the cluster and connect to the storage.
* `cluster-external-management`: Connect to an [external Ceph cluster](ceph-cluster-crd.md#external-cluster) with the admin key of the external cluster to enable
* `cluster-external.yaml`: Connect to an [external Ceph cluster](ceph-cluster-crd.md#external-cluster) with minimal access to monitor the health of the cluster and connect to the storage.
* `cluster-external-management.yaml`: Connect to an [external Ceph cluster](ceph-cluster-crd.md#external-cluster) with the admin key of the external cluster to enable
remote creation of pools and configure services such as an [Object Store](ceph-object.md) or a [Shared Filesystem](ceph-filesystem.md).
* `cluster-stretched.yaml`: Create a cluster in "stretched" mode, with five mons stretched across three zones, and the OSDs across two zones. See the [Stretch documentation](ceph-cluster-crd.md#stretch-cluster).
See the [Cluster CRD](ceph-cluster-crd.md) topic for more details and more examples for the settings.
+1
View File
@@ -18,6 +18,7 @@ v1.5...
### Ceph
* Stretch clusters for mons and OSDs to work reliably across two datacenters (Experimental mode)
* Ceph Block Pool: add mirroring support
* Ceph Block Pool: add `replicasPerFailureDomain` to set the number of replica in a failure domain ([#5591](https://github.com/rook/rook/issues/5591))
* Ceph Cluster: export the storage capacity of the ceph cluster ([#6475](https://github.com/rook/rook/pull/6475))
@@ -72,6 +72,26 @@ spec:
volumeClaimTemplate:
type: object
x-kubernetes-preserve-unknown-fields: true
stretchCluster:
type: object
nullable: true
properties:
failureDomainLabel:
type: string
subFailureDomain:
type: string
zones:
type: array
items:
type: object
properties:
name:
type: string
arbiter:
type: boolean
volumeClaimTemplate:
type: object
x-kubernetes-preserve-unknown-fields: true
mgr:
type: object
properties:
@@ -0,0 +1,63 @@
#################################################################################################################
# Define the settings for the rook-ceph cluster with common settings for a production cluster.
# All nodes with available raw devices will be used for the Ceph cluster. At least three nodes are required
# in this example. See the documentation for more details on storage settings available.
# For example, to create the cluster:
# kubectl create -f common.yaml
# kubectl create -f operator.yaml
# kubectl create -f cluster-stretched.yaml
#################################################################################################################
apiVersion: ceph.rook.io/v1
kind: CephCluster
metadata:
name: rook-ceph
namespace: rook-ceph
spec:
dataDirHostPath: /var/lib/rook
mon:
# Five mons must be created for stretch mode
count: 5
allowMultiplePerNode: false
stretchCluster:
# The ceph failure domain will be extracted from the label, which by default is the zone. The nodes running OSDs must have
# this label in order for the OSDs to be configured in the correct topology. For topology labels, see
# https://rook.github.io/docs/rook/master/ceph-cluster-crd.html#osd-topology.
failureDomainLabel: topology.kubernetes.io/zone
# The sub failure domain is the secondary level at which the data will be placed to maintain data durability and availability.
# The default is "host", which means that each OSD must be on a different node and you would need at least two nodes per zone.
# If the subFailureDomain is set to "osd", the OSDs would be allowed anywhere in the same zone including on the same node.
# If set to "rack" or some other intermediate failure domain, those labels would also need to be set on the nodes where
# the osds are started.
subFailureDomain: host
zones:
- name: a
arbiter: true
- name: b
- name: c
cephVersion:
# Stretch cluster support upstream is only planned starting in Ceph Pacific
image: ceph/daemon-base:latest-master
allowUnsupported: true
skipUpgradeChecks: false
continueUpgradeAfterChecksEvenIfNotHealthy: false
dashboard:
enabled: true
ssl: true
storage:
useAllNodes: true
useAllDevices: true
deviceFilter: ""
# OSD placement is expected to include the non-arbiter zones
placement:
osd:
nodeAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
nodeSelectorTerms:
- matchExpressions:
- key: topology.kubernetes.io/zone
operator: In
values:
- b
- c
@@ -90,6 +90,26 @@ spec:
volumeClaimTemplate:
type: object
x-kubernetes-preserve-unknown-fields: true
stretchCluster:
type: object
nullable: true
properties:
failureDomainLabel:
type: string
subFailureDomain:
type: string
zones:
type: array
items:
type: object
properties:
name:
type: string
arbiter:
type: boolean
volumeClaimTemplate:
type: object
x-kubernetes-preserve-unknown-fields: true
mgr:
type: object
properties:
+4
View File
@@ -30,6 +30,10 @@ import (
// will be registered for the validating webhook.
var _ webhook.Validator = &CephCluster{}
func (c *ClusterSpec) IsStretchCluster() bool {
return c.Mon.StretchCluster != nil && len(c.Mon.StretchCluster.Zones) > 0
}
func (c *CephCluster) ValidateCreate() error {
logger.Infof("validate create cephcluster %q", c.ObjectMeta.Name)
+13
View File
@@ -280,9 +280,22 @@ const (
type MonSpec struct {
Count int `json:"count,omitempty"`
AllowMultiplePerNode bool `json:"allowMultiplePerNode,omitempty"`
StretchCluster *StretchClusterSpec `json:"stretchCluster,omitempty"`
VolumeClaimTemplate *v1.PersistentVolumeClaim `json:"volumeClaimTemplate,omitempty"`
}
type StretchClusterSpec struct {
FailureDomainLabel string `json:"failureDomainLabel,omitempty"`
SubFailureDomain string `json:"subFailureDomain,omitempty"`
Zones []StretchClusterZoneSpec `json:"zones,omitempty"`
}
type StretchClusterZoneSpec struct {
Name string `json:"name,omitempty"`
Arbiter bool `json:"arbiter,omitempty"`
VolumeClaimTemplate *v1.PersistentVolumeClaim `json:"volumeClaimTemplate,omitempty"`
}
// MgrSpec represents options to configure a ceph mgr
type MgrSpec struct {
Modules []Module `json:"modules,omitempty"`
@@ -1585,6 +1585,11 @@ func (in *Module) DeepCopy() *Module {
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in *MonSpec) DeepCopyInto(out *MonSpec) {
*out = *in
if in.StretchCluster != nil {
in, out := &in.StretchCluster, &out.StretchCluster
*out = new(StretchClusterSpec)
(*in).DeepCopyInto(*out)
}
if in.VolumeClaimTemplate != nil {
in, out := &in.VolumeClaimTemplate, &out.VolumeClaimTemplate
*out = new(corev1.PersistentVolumeClaim)
@@ -1962,6 +1967,50 @@ func (in *Status) DeepCopy() *Status {
return out
}
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in *StretchClusterSpec) DeepCopyInto(out *StretchClusterSpec) {
*out = *in
if in.Zones != nil {
in, out := &in.Zones, &out.Zones
*out = make([]StretchClusterZoneSpec, len(*in))
for i := range *in {
(*in)[i].DeepCopyInto(&(*out)[i])
}
}
return
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new StretchClusterSpec.
func (in *StretchClusterSpec) DeepCopy() *StretchClusterSpec {
if in == nil {
return nil
}
out := new(StretchClusterSpec)
in.DeepCopyInto(out)
return out
}
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in *StretchClusterZoneSpec) DeepCopyInto(out *StretchClusterZoneSpec) {
*out = *in
if in.VolumeClaimTemplate != nil {
in, out := &in.VolumeClaimTemplate, &out.VolumeClaimTemplate
*out = new(corev1.PersistentVolumeClaim)
(*in).DeepCopyInto(*out)
}
return
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new StretchClusterZoneSpec.
func (in *StretchClusterZoneSpec) DeepCopy() *StretchClusterZoneSpec {
if in == nil {
return nil
}
out := new(StretchClusterZoneSpec)
in.DeepCopyInto(out)
return out
}
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in SummarySpec) DeepCopyInto(out *SummarySpec) {
{
+9 -9
View File
@@ -23,10 +23,10 @@ import (
)
const (
crushReplicatedType = 1
ruleMinSizeDefault = 1
ruleMaxSizeDefault = 10
stretchedCRUSHRuleTemplate = `
crushReplicatedType = 1
ruleMinSizeDefault = 1
ruleMaxSizeDefault = 10
twoStepCRUSHRuleTemplate = `
rule %s {
id %d
type replicated
@@ -44,9 +44,9 @@ var (
stepEmit = &stepSpec{Operation: "emit"}
)
func buildStretchClusterPlainCrushRule(crushMap CrushMap, ruleName string, pool cephv1.PoolSpec) string {
func buildTwoStepPlainCrushRule(crushMap CrushMap, ruleName string, pool cephv1.PoolSpec) string {
return fmt.Sprintf(
stretchedCRUSHRuleTemplate,
twoStepCRUSHRuleTemplate,
ruleName,
generateRuleID(crushMap.Rules),
ruleMinSizeDefault,
@@ -57,7 +57,7 @@ func buildStretchClusterPlainCrushRule(crushMap CrushMap, ruleName string, pool
)
}
func buildStretchClusterCrushRule(crushMap CrushMap, ruleName string, pool cephv1.PoolSpec) *ruleSpec {
func buildTwoStepCrushRule(crushMap CrushMap, ruleName string, pool cephv1.PoolSpec) *ruleSpec {
/*
The complete CRUSH rule looks like this:
@@ -82,11 +82,11 @@ func buildStretchClusterCrushRule(crushMap CrushMap, ruleName string, pool cephv
Type: crushReplicatedType,
MinSize: ruleMinSizeDefault,
MaxSize: ruleMaxSizeDefault,
Steps: buildStretchClusterCrushSteps(pool),
Steps: buildTwoStepCrushSteps(pool),
}
}
func buildStretchClusterCrushSteps(pool cephv1.PoolSpec) []stepSpec {
func buildTwoStepCrushSteps(pool cephv1.PoolSpec) []stepSpec {
// Create CRUSH rule steps
steps := []stepSpec{}
+2 -2
View File
@@ -40,7 +40,7 @@ func TestBuildStretchClusterCrushRule(t *testing.T) {
},
}
rule := buildStretchClusterCrushRule(crushMap, "stretched", *pool)
rule := buildTwoStepCrushRule(crushMap, "stretched", *pool)
assert.Equal(t, 2, rule.ID)
}
@@ -52,7 +52,7 @@ func TestBuildCrushSteps(t *testing.T) {
ReplicasPerFailureDomain: 2,
},
}
steps := buildStretchClusterCrushSteps(*pool)
steps := buildTwoStepCrushSteps(*pool)
assert.Equal(t, 4, len(steps))
assert.Equal(t, cephv1.DefaultCRUSHRoot, steps[0].ItemName)
assert.Equal(t, "datacenter", steps[1].Type)
+62
View File
@@ -17,9 +17,17 @@ package client
import (
"encoding/json"
"fmt"
"syscall"
"github.com/pkg/errors"
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
"github.com/rook/rook/pkg/clusterd"
"github.com/rook/rook/pkg/util/exec"
)
const (
defaultStretchCrushRuleName = "default_stretch_cluster_rule"
)
// MonStatusResponse represents the response from a quorum_status mon_command (subset of all available fields, only
@@ -66,3 +74,57 @@ func GetMonQuorumStatus(context *clusterd.Context, clusterInfo *ClusterInfo) (Mo
return resp, nil
}
// EnableStretchElectionStrategy enables the mon connectivity algorithm for stretch clusters
func EnableStretchElectionStrategy(context *clusterd.Context, clusterInfo *ClusterInfo) error {
args := []string{"mon", "set", "election_strategy", "connectivity"}
buf, err := NewCephCommand(context, clusterInfo, args).Run()
if err != nil {
return errors.Wrap(err, "failed to enable stretch cluster election strategy")
}
logger.Infof("successfully enabled stretch cluster election strategy. %s", string(buf))
return nil
}
// CreateDefaultStretchCrushRule creates the default CRUSH rule for the stretch cluster
func CreateDefaultStretchCrushRule(context *clusterd.Context, clusterInfo *ClusterInfo, failureDomain, subFailureDomain string) error {
pool := cephv1.PoolSpec{
FailureDomain: failureDomain,
Replicated: cephv1.ReplicatedSpec{SubFailureDomain: subFailureDomain},
}
if err := createTwoStepCrushRule(context, clusterInfo, defaultStretchCrushRuleName, pool); err != nil {
return errors.Wrap(err, "failed to create default stretch crush rule")
}
logger.Info("successfully created the default stretch crush rule")
return nil
}
// SetMonStretchZone sets the location of a mon in the stretch cluster
func SetMonStretchZone(context *clusterd.Context, clusterInfo *ClusterInfo, monName, failureDomain, zone string) error {
args := []string{"mon", "set_location", monName, fmt.Sprintf("%s=%s", failureDomain, zone)}
buf, err := NewCephCommand(context, clusterInfo, args).Run()
if err != nil {
return errors.Wrap(err, "failed to set mon stretch zone")
}
output := string(buf)
logger.Debug(output)
logger.Infof("successfully set mon %q stretch zone to %q", monName, zone)
return nil
}
// SetMonStretchTiebreaker sets the tiebreaker mon in the stretch cluster
func SetMonStretchTiebreaker(context *clusterd.Context, clusterInfo *ClusterInfo, monName, bucketType string) error {
logger.Infof("enabling stretch mode with mon arbiter %q with crush rule %q in failure domain %q", monName, defaultStretchCrushRuleName, bucketType)
args := []string{"mon", "enable_stretch_mode", monName, defaultStretchCrushRuleName, bucketType}
buf, err := NewCephCommand(context, clusterInfo, args).Run()
if err != nil {
if code, ok := exec.ExitStatus(err); ok && code == int(syscall.EINVAL) {
logger.Infof("stretch mode is already enabled. %s", string(buf))
return nil
}
return errors.Wrap(err, "failed to set mon stretch zone")
}
logger.Debug(string(buf))
logger.Infof("successfully set mon tiebreaker %q in failure domain %q", monName, bucketType)
return nil
}
+64
View File
@@ -19,6 +19,9 @@ import (
"fmt"
"testing"
"github.com/pkg/errors"
"github.com/rook/rook/pkg/clusterd"
exectest "github.com/rook/rook/pkg/util/exec/test"
"github.com/stretchr/testify/assert"
)
@@ -61,3 +64,64 @@ func TestCephArgs(t *testing.T) {
assert.Equal(t, "--name=client.admin", args[3])
assert.Equal(t, "--keyring=/var/lib/rook/a/client.admin.keyring", args[4])
}
func TestStretchElectionStrategy(t *testing.T) {
executor := &exectest.MockExecutor{}
executor.MockExecuteCommandWithOutputFile = func(command, outputFile string, args ...string) (string, error) {
logger.Infof("Command: %s %v", command, args)
if args[0] == "mon" && args[1] == "set" && args[2] == "election_strategy" {
assert.Equal(t, "connectivity", args[3])
return "", nil
}
return "", errors.Errorf("unexpected ceph command %q", args)
}
context := &clusterd.Context{Executor: executor}
clusterInfo := AdminClusterInfo("mycluster")
err := EnableStretchElectionStrategy(context, clusterInfo)
assert.NoError(t, err)
}
func TestStretchClusterSettings(t *testing.T) {
monName := "a"
failureDomain := "rack"
zone := "rack-x"
executor := &exectest.MockExecutor{}
executor.MockExecuteCommandWithOutputFile = func(command, outputFile string, args ...string) (string, error) {
logger.Infof("Command: %s %v", command, args)
switch {
case args[0] == "mon" && args[1] == "set_location":
assert.Equal(t, monName, args[2])
assert.Equal(t, fmt.Sprintf("%s=%s", failureDomain, zone), args[3])
return "", nil
}
return "", errors.Errorf("unexpected ceph command %q", args)
}
context := &clusterd.Context{Executor: executor}
clusterInfo := AdminClusterInfo("mycluster")
err := SetMonStretchZone(context, clusterInfo, monName, failureDomain, zone)
assert.NoError(t, err)
}
func TestStretchClusterMonTiebreaker(t *testing.T) {
monName := "a"
failureDomain := "rack"
executor := &exectest.MockExecutor{}
executor.MockExecuteCommandWithOutputFile = func(command, outputFile string, args ...string) (string, error) {
logger.Infof("Command: %s %v", command, args)
switch {
case args[0] == "mon" && args[1] == "enable_stretch_mode":
assert.Equal(t, monName, args[2])
assert.Equal(t, defaultStretchCrushRuleName, args[3])
assert.Equal(t, failureDomain, args[4])
return "", nil
}
return "", errors.Errorf("unexpected ceph command %q", args)
}
context := &clusterd.Context{Executor: executor}
clusterInfo := AdminClusterInfo("mycluster")
err := SetMonStretchTiebreaker(context, clusterInfo, monName, failureDomain)
assert.NoError(t, err)
}
+8
View File
@@ -357,6 +357,14 @@ func createStretchedReplicationCrushRule(context *clusterd.Context, clusterInfo
return errors.Wrap(err, "failed to get crush map")
}
// Check if the crush rule already exists
for _, rule := range crushMap.Rules {
if rule.Name == ruleName {
logger.Debugf("CRUSH rule %q already exists", ruleName)
return nil
}
}
// Fetch the compiled crush map
compiledCRUSHMapFilePath, err := GetCompiledCrushMap(context, clusterInfo)
if err != nil {
+51
View File
@@ -265,3 +265,54 @@ func TestSetPoolReplicatedSizeProperty(t *testing.T) {
err = SetPoolReplicatedSizeProperty(context, AdminClusterInfo("mycluster"), poolName, "1")
assert.NoError(t, err)
}
func TestCreateStretchCrushRule(t *testing.T) {
testCreateStretchCrushRule(t, true)
testCreateStretchCrushRule(t, false)
}
func testCreateStretchCrushRule(t *testing.T, alreadyExists bool) {
executor := &exectest.MockExecutor{}
context := &clusterd.Context{Executor: executor}
executor.MockExecuteCommandWithOutput = func(command string, args ...string) (string, error) {
logger.Infof("Command: %s %v", command, args)
if args[0] == "osd" {
if args[1] == "getcrushmap" {
return "", nil
}
if args[1] == "setcrushmap" {
if alreadyExists {
return "", errors.New("setcrushmap not expected for already existing crush rule")
}
return "", nil
}
}
if command == "crushtool" {
switch {
case args[0] == "--decompile" || args[0] == "--compile":
if alreadyExists {
return "", errors.New("--compile or --decompile not expected for already existing crush rule")
}
return "", nil
}
}
return "", errors.Errorf("unexpected ceph command %q", args)
}
executor.MockExecuteCommandWithOutputFile = func(command, outputFile string, args ...string) (string, error) {
logger.Infof("Command (file): %s %v", command, args)
if args[0] == "osd" && args[1] == "crush" && args[2] == "dump" {
return testCrushMap, nil
}
return "", errors.Errorf("unexpected ceph command %q", args)
}
clusterInfo := AdminClusterInfo("mycluster")
clusterSpec := &cephv1.ClusterSpec{}
poolSpec := cephv1.PoolSpec{}
ruleName := "testrule"
if alreadyExists {
ruleName = "replicated_ruleset"
}
err := createTwoStepCrushRule(context, clusterInfo, clusterSpec, ruleName, poolSpec)
assert.NoError(t, err)
}
+39 -4
View File
@@ -167,6 +167,13 @@ func (c *cluster) doOrchestration(rookImage string, cephVersion cephver.CephVers
return errors.Wrap(err, "failed to start ceph osds")
}
// If a stretch cluster, enable the arbiter after the OSDs are created with the CRUSH map
if c.Spec.IsStretchCluster() {
if err := c.mons.ConfigureArbiter(); err != nil {
return errors.Wrap(err, "failed to configure stretch arbiter")
}
}
logger.Infof("done reconciling ceph cluster in namespace %q", c.Namespace)
// We should be done updating by now
@@ -221,7 +228,7 @@ func (c *ClusterController) initializeCluster(cluster *cluster, clusterObj *ceph
c.configureCephMonitoring(cluster, clusterInfo)
}
err = c.configureLocalCephCluster(cluster, clusterObj)
err = c.configureLocalCephCluster(cluster)
if err != nil {
return errors.Wrap(err, "failed to configure local ceph cluster")
}
@@ -235,9 +242,9 @@ func (c *ClusterController) initializeCluster(cluster *cluster, clusterObj *ceph
return nil
}
func (c *ClusterController) configureLocalCephCluster(cluster *cluster, clusterObj *cephv1.CephCluster) error {
func (c *ClusterController) configureLocalCephCluster(cluster *cluster) error {
// Cluster Spec validation
err := c.preClusterStartValidation(cluster, clusterObj)
err := c.preClusterStartValidation(cluster)
if err != nil {
return errors.Wrap(err, "failed to perform validation before cluster creation")
}
@@ -328,7 +335,7 @@ func (c *cluster) notifyChildControllerOfUpgrade() error {
}
// Validate the cluster Specs
func (c *ClusterController) preClusterStartValidation(cluster *cluster, clusterObj *cephv1.CephCluster) error {
func (c *ClusterController) preClusterStartValidation(cluster *cluster) error {
if cluster.Spec.Mon.Count == 0 {
logger.Warningf("mon count should be at least 1, will use default value of %d", mon.DefaultMonCount)
@@ -347,6 +354,9 @@ func (c *ClusterController) preClusterStartValidation(cluster *cluster, clusterO
if len(cluster.Spec.Storage.Directories) != 0 {
logger.Warning("running osds on directory is not supported anymore, use devices instead.")
}
if err := validateStretchCluster(cluster); err != nil {
return err
}
if cluster.Spec.Network.IsMultus() {
_, isPublic := cluster.Spec.Network.Selectors[config.PublicNetworkSelectorKeyName]
_, isCluster := cluster.Spec.Network.Selectors[config.ClusterNetworkSelectorKeyName]
@@ -385,6 +395,31 @@ func (c *ClusterController) preClusterStartValidation(cluster *cluster, clusterO
return nil
}
func validateStretchCluster(cluster *cluster) error {
if !cluster.Spec.IsStretchCluster() {
return nil
}
if len(cluster.Spec.Mon.StretchCluster.Zones) != 3 {
return errors.Errorf("expecting exactly three zones for the stretch cluster, but found %d", len(cluster.Spec.Mon.StretchCluster.Zones))
}
if cluster.Spec.Mon.Count != 3 && cluster.Spec.Mon.Count != 5 {
return errors.Errorf("invalid number of mons %d for a stretch cluster, expecting 5 (recommended) or 3 (minimal)", cluster.Spec.Mon.Count)
}
arbitersFound := 0
for _, zone := range cluster.Spec.Mon.StretchCluster.Zones {
if zone.Arbiter {
arbitersFound++
}
if zone.Name == "" {
return errors.New("missing zone name for the stretch cluster")
}
}
if arbitersFound != 1 {
return errors.Errorf("expecting to find exactly one arbiter zone, but found %d", arbitersFound)
}
return nil
}
func extractExitCode(err error) (int, bool) {
exitErr, ok := err.(*exec.ExitError)
if ok {
+75
View File
@@ -0,0 +1,75 @@
/*
Copyright 2020 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.
*/
// Package cluster to manage a Ceph cluster.
package cluster
import (
"testing"
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
"github.com/rook/rook/pkg/clusterd"
testop "github.com/rook/rook/pkg/operator/test"
)
func TestPreClusterStartValidation(t *testing.T) {
type args struct {
cluster *cluster
}
tests := []struct {
name string
args args
wantErr bool
}{
{"no settings", args{&cluster{Spec: &cephv1.ClusterSpec{}}}, false},
{"even mons", args{&cluster{Spec: &cephv1.ClusterSpec{Mon: cephv1.MonSpec{Count: 2}}}}, true},
{"missing stretch zones", args{&cluster{Spec: &cephv1.ClusterSpec{Mon: cephv1.MonSpec{StretchCluster: &cephv1.StretchClusterSpec{Zones: []cephv1.StretchClusterZoneSpec{
{Name: "a"},
}}}}}}, true},
{"missing arbiter", args{&cluster{Spec: &cephv1.ClusterSpec{Mon: cephv1.MonSpec{StretchCluster: &cephv1.StretchClusterSpec{Zones: []cephv1.StretchClusterZoneSpec{
{Name: "a"},
{Name: "b"},
{Name: "c"},
}}}}}}, true},
{"missing zone name", args{&cluster{Spec: &cephv1.ClusterSpec{Mon: cephv1.MonSpec{StretchCluster: &cephv1.StretchClusterSpec{Zones: []cephv1.StretchClusterZoneSpec{
{Arbiter: true},
{Name: "b"},
{Name: "c"},
}}}}}}, true},
{"valid stretch cluster", args{&cluster{Spec: &cephv1.ClusterSpec{Mon: cephv1.MonSpec{Count: 3, StretchCluster: &cephv1.StretchClusterSpec{Zones: []cephv1.StretchClusterZoneSpec{
{Name: "a", Arbiter: true},
{Name: "b"},
{Name: "c"},
}}}}}}, false},
{"not enough stretch nodes", args{&cluster{Spec: &cephv1.ClusterSpec{Mon: cephv1.MonSpec{Count: 5, StretchCluster: &cephv1.StretchClusterSpec{Zones: []cephv1.StretchClusterZoneSpec{
{Name: "a", Arbiter: true},
{Name: "b"},
{Name: "c"},
}}}}}}, true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
c := &ClusterController{
context: &clusterd.Context{
Clientset: testop.New(t, 3),
},
}
if err := c.preClusterStartValidation(tt.args.cluster); (err != nil) != tt.wantErr {
t.Errorf("ClusterController.preClusterStartValidation() error = %v, wantErr %v", err, tt.wantErr)
}
})
}
}
+3 -3
View File
@@ -86,7 +86,7 @@ func CreateOrLoadClusterInfo(context *clusterd.Context, namespace string, ownerR
var clusterInfo *cephclient.ClusterInfo
maxMonID := -1
monMapping := &Mapping{
Node: map[string]*NodeInfo{},
Schedule: map[string]*MonScheduleInfo{},
}
secrets, err := context.Clientset.CoreV1().Secrets(namespace).Get(AppName, metav1.GetOptions{})
@@ -211,7 +211,7 @@ func loadMonConfig(clientset kubernetes.Interface, namespace string) (map[string
monEndpointMap := map[string]*cephclient.MonInfo{}
maxMonID := -1
monMapping := &Mapping{
Node: map[string]*NodeInfo{},
Schedule: map[string]*MonScheduleInfo{},
}
cm, err := clientset.CoreV1().ConfigMaps(namespace).Get(EndpointConfigMapName, metav1.GetOptions{})
@@ -249,7 +249,7 @@ func loadMonConfig(clientset kubernetes.Interface, namespace string) (map[string
logger.Errorf("invalid JSON in mon mapping. %v", err)
}
logger.Debugf("loaded: maxMonID=%d, mons=%+v, mapping=%+v", maxMonID, monEndpointMap, monMapping)
logger.Debugf("loaded: maxMonID=%d, mons=%+v, assignment=%+v", maxMonID, monEndpointMap, monMapping)
return monEndpointMap, maxMonID, monMapping, nil
}
+28 -4
View File
@@ -330,8 +330,15 @@ func (c *Cluster) failoverMon(name string) error {
}
}()
// remove the failed mon from a local list of the existing mons for finding a stretch zone
existingMons := c.clusterInfoToMonConfig(name)
zone, err := c.findAvailableZoneIfStretched(existingMons)
if err != nil {
return errors.Wrap(err, "failed to find available stretch zone")
}
// Start a new monitor
m := c.newMonConfig(c.maxMonID + 1)
m := c.newMonConfig(c.maxMonID+1, zone)
logger.Infof("starting new mon: %+v", m)
mConf := []*monConfig{m}
@@ -342,11 +349,11 @@ func (c *Cluster) failoverMon(name string) error {
}
if c.spec.Network.IsHost() {
node, ok := c.mapping.Node[m.DaemonName]
schedule, ok := c.mapping.Schedule[m.DaemonName]
if !ok {
return errors.Errorf("mon %s doesn't exist in assignment map", m.DaemonName)
}
m.PublicIP = node.Address
m.PublicIP = schedule.Address
} else {
// Create the service endpoint
serviceIP, err := c.createService(m)
@@ -362,6 +369,23 @@ func (c *Cluster) failoverMon(name string) error {
return errors.Wrapf(err, "failed to start new mon %s", m.DaemonName)
}
// Assign to a zone if a stretch cluster
if c.spec.IsStretchCluster() {
updateArbiter := false
if name == c.arbiterMon {
updateArbiter = true
}
if err := c.assignStretchMonsToZones([]*monConfig{m}); err != nil {
return errors.Wrap(err, "failed to assign mons to zones")
}
if updateArbiter {
// Update the arbiter mon for the stretch cluster if it changed
if err := c.ConfigureArbiter(); err != nil {
return errors.Wrap(err, "failed to configure stretch arbiter")
}
}
}
// Only increment the max mon id if the new pod started successfully
c.maxMonID++
newMonSucceeded = true
@@ -393,7 +417,7 @@ func (c *Cluster) removeMon(daemonName string) error {
}
delete(c.ClusterInfo.Monitors, daemonName)
delete(c.mapping.Node, daemonName)
delete(c.mapping.Schedule, daemonName)
// Remove the service endpoint
if err := c.context.Clientset.CoreV1().Services(c.Namespace).Delete(resourceName, options); err != nil {
+4 -4
View File
@@ -76,7 +76,7 @@ func TestCheckHealth(t *testing.T) {
c.waitForStart = false
defer os.RemoveAll(c.context.ConfigDir)
c.mapping.Node["f"] = &NodeInfo{
c.mapping.Schedule["f"] = &MonScheduleInfo{
Name: "node0",
Address: "",
}
@@ -201,10 +201,10 @@ func TestCheckHealthNotFound(t *testing.T) {
c.waitForStart = false
defer os.RemoveAll(c.context.ConfigDir)
c.mapping.Node["a"] = &NodeInfo{
c.mapping.Schedule["a"] = &MonScheduleInfo{
Name: "node0",
}
c.mapping.Node["b"] = &NodeInfo{
c.mapping.Schedule["b"] = &MonScheduleInfo{
Name: "node0",
}
c.maxMonID = 4
@@ -220,7 +220,7 @@ func TestCheckHealthNotFound(t *testing.T) {
}
// Because the mon a isn't in the MonInQuorumResponse() this will create a new mon
delete(c.mapping.Node, "b")
delete(c.mapping.Schedule, "b")
err = c.checkHealth()
assert.Nil(t, err)
// No updates in unit tests w/ workaround
+194 -53
View File
@@ -93,11 +93,6 @@ const (
// pods and waiting for kubernetes scheduling to complete.
canaryRetries = 30
canaryRetryDelaySeconds = 5
// Fallback pod anti affinity for mon pods to PreferredDuringSchedulingIgnoredDuringExecution
// if not in a case (e.g. HostNetworking, not AllowMultiplePerHost)
// needing RequiredDuringSchedulingIgnoredDuringExecution pod anti-affinity
PreferredDuringScheduling = true
)
var (
@@ -126,6 +121,7 @@ type Cluster struct {
ownerRef metav1.OwnerReference
csiConfigMutex *sync.Mutex
isUpgrade bool
arbiterMon string
}
// monConfig for a single monitor
@@ -138,6 +134,8 @@ type monConfig struct {
PublicIP string
// Port is the port on which the mon will listen for connections
Port int32
// The zone used for a stretch cluster
Zone string
// DataPathMap is the mapping relationship between mon data stored on the host and mon data
// stored in containers.
DataPathMap *config.DataPathMap
@@ -145,14 +143,17 @@ type monConfig struct {
// Mapping is mon node and port mapping
type Mapping struct {
Node map[string]*NodeInfo `json:"node"`
// This isn't really node info since it could also be for zones, but we leave it as "node" for backward compatibility.
Schedule map[string]*MonScheduleInfo `json:"node"`
}
// NodeInfo contains name and address of a node
type NodeInfo struct {
Name string
Hostname string
Address string
// MonScheduleInfo contains name and address of a node.
type MonScheduleInfo struct {
// Name of the node. **json names are capitalized for backwards compat**
Name string `json:"Name,omitempty"`
Hostname string `json:"Hostname,omitempty"`
Address string `json:"Address,omitempty"`
Zone string `json:"zone,omitempty"`
}
type SchedulingResult struct {
@@ -173,7 +174,7 @@ func New(context *clusterd.Context, namespace string, spec cephv1.ClusterSpec, o
monPodTimeout: 5 * time.Minute,
monTimeoutList: map[string]time.Time{},
mapping: &Mapping{
Node: map[string]*NodeInfo{},
Schedule: map[string]*MonScheduleInfo{},
},
ownerRef: ownerRef,
csiConfigMutex: csiConfigMutex,
@@ -218,7 +219,10 @@ func (c *Cluster) Start(clusterInfo *cephclient.ClusterInfo, rookVersion string,
func (c *Cluster) startMons(targetCount int) error {
// init the mon config
existingCount, mons := c.initMonConfig(targetCount)
existingCount, mons, err := c.initMonConfig(targetCount)
if err != nil {
return errors.Wrap(err, "failed to init mon config")
}
// Assign the mons to nodes
if err := c.assignMons(mons); err != nil {
@@ -281,6 +285,12 @@ func (c *Cluster) startMons(targetCount int) error {
}
}
if c.spec.IsStretchCluster() {
if err := c.configureStretchCluster(mons); err != nil {
return errors.Wrap(err, "failed to configure stretch mons")
}
}
logger.Debugf("mon endpoints used are: %s", FlattenMonEndpoints(c.ClusterInfo.Monitors))
// Check if there are orphaned mon resources that should be cleaned up at the end of a reconcile.
@@ -290,6 +300,61 @@ func (c *Cluster) startMons(targetCount int) error {
return nil
}
func (c *Cluster) configureStretchCluster(mons []*monConfig) error {
if err := c.assignStretchMonsToZones(mons); err != nil {
return errors.Wrap(err, "failed to assign mons to zones")
}
// Enable the mon connectivity strategy
if err := client.EnableStretchElectionStrategy(c.context, c.ClusterInfo); err != nil {
return errors.Wrap(err, "failed to enable stretch cluster")
}
// Create the default crush rule for stretch clusters, that by default will also apply to all pools
if err := client.CreateDefaultStretchCrushRule(c.context, c.ClusterInfo, c.stretchFailureDomainName(), c.spec.Mon.StretchCluster.SubFailureDomain); err != nil {
return errors.Wrap(err, "failed to create default stretch rule")
}
return nil
}
func (c *Cluster) assignStretchMonsToZones(mons []*monConfig) error {
var arbiterZone string
for _, zone := range c.spec.Mon.StretchCluster.Zones {
if zone.Arbiter {
arbiterZone = zone.Name
break
}
}
// Set the location for each mon
domainName := c.stretchFailureDomainName()
for _, mon := range mons {
if mon.Zone == arbiterZone {
// remember the arbiter mon to be set later in the reconcile after the OSDs are configured
c.arbiterMon = mon.DaemonName
}
logger.Infof("setting mon %q to stretch %s=%s", mon.DaemonName, domainName, mon.Zone)
if err := client.SetMonStretchZone(c.context, c.ClusterInfo, mon.DaemonName, domainName, mon.Zone); err != nil {
return errors.Wrapf(err, "failed to set mon %q zone", mon.DaemonName)
}
}
return nil
}
func (c *Cluster) ConfigureArbiter() error {
if c.arbiterMon == "" {
return errors.New("arbiter not specified for the stretch cluster")
}
// Set the mon tiebreaker
if err := client.SetMonStretchTiebreaker(c.context, c.ClusterInfo, c.arbiterMon, c.stretchFailureDomainName()); err != nil {
return errors.Wrap(err, "failed to set mon tiebreaker")
}
return nil
}
// ensureMonsRunning is called in two scenarios:
// 1. To create a new mon and wait for it to join quorum (requireAllInQuorum = true). This method will be called multiple times
// to add a mon until we have reached the desired number of mons.
@@ -365,42 +430,91 @@ func (c *Cluster) initClusterInfo(cephVersion cephver.CephVersion) error {
return nil
}
func (c *Cluster) initMonConfig(size int) (int, []*monConfig) {
mons := []*monConfig{}
func (c *Cluster) initMonConfig(size int) (int, []*monConfig, error) {
// initialize the mon pod info for mons that have been previously created
for _, monitor := range c.ClusterInfo.Monitors {
mons = append(mons, &monConfig{
ResourceName: resourceName(monitor.Name),
DaemonName: monitor.Name,
Port: cephutil.GetPortFromEndpoint(monitor.Endpoint),
DataPathMap: config.NewStatefulDaemonDataPathMap(
c.spec.DataDirHostPath, dataDirRelativeHostPath(monitor.Name), config.MonType, monitor.Name, c.Namespace),
})
}
mons := c.clusterInfoToMonConfig("")
// initialize mon info if we don't have enough mons (at first startup)
existingCount := len(c.ClusterInfo.Monitors)
for i := len(c.ClusterInfo.Monitors); i < size; i++ {
c.maxMonID++
mons = append(mons, c.newMonConfig(c.maxMonID))
zone, err := c.findAvailableZoneIfStretched(mons)
if err != nil {
return existingCount, mons, errors.Wrap(err, "stretch zone not available")
}
mons = append(mons, c.newMonConfig(c.maxMonID, zone))
}
return existingCount, mons
return existingCount, mons, nil
}
func (c *Cluster) newMonConfig(monID int) *monConfig {
func (c *Cluster) clusterInfoToMonConfig(excludedMon string) []*monConfig {
mons := []*monConfig{}
for _, monitor := range c.ClusterInfo.Monitors {
if monitor.Name == excludedMon {
// Skip a mon if it is being failed over
continue
}
var zone string
schedule := c.mapping.Schedule[monitor.Name]
if schedule != nil {
zone = schedule.Zone
}
mons = append(mons, &monConfig{
ResourceName: resourceName(monitor.Name),
DaemonName: monitor.Name,
Port: cephutil.GetPortFromEndpoint(monitor.Endpoint),
Zone: zone,
DataPathMap: config.NewStatefulDaemonDataPathMap(
c.spec.DataDirHostPath, dataDirRelativeHostPath(monitor.Name), config.MonType, monitor.Name, c.Namespace),
})
}
return mons
}
func (c *Cluster) newMonConfig(monID int, zone string) *monConfig {
daemonName := k8sutil.IndexToName(monID)
return &monConfig{
ResourceName: resourceName(daemonName),
DaemonName: daemonName,
Port: DefaultMsgr1Port,
Zone: zone,
DataPathMap: config.NewStatefulDaemonDataPathMap(
c.spec.DataDirHostPath, dataDirRelativeHostPath(daemonName), config.MonType, daemonName, c.Namespace),
}
}
func (c *Cluster) findAvailableZoneIfStretched(mons []*monConfig) (string, error) {
if !c.spec.IsStretchCluster() {
return "", nil
}
// Build the count of current mons per zone
zoneCount := map[string]int{}
for _, m := range mons {
if m.Zone == "" {
return "", errors.Errorf("zone not found on mon %q", m.DaemonName)
}
zoneCount[m.Zone]++
}
// Find a zone in the stretch cluster that still needs an assignment
for _, zone := range c.spec.Mon.StretchCluster.Zones {
count, ok := zoneCount[zone.Name]
if !ok {
// The zone isn't currently assigned to any mon, so return it
return zone.Name, nil
}
if c.spec.Mon.Count == 5 && count == 1 && !zone.Arbiter {
// The zone only has 1 mon assigned, but needs 2 mons since it is not the arbiter
return zone.Name, nil
}
}
return "", errors.New("A zone is not available to assign a new mon")
}
// resourceName ensures the mon name has the rook-ceph-mon prefix
func resourceName(name string) string {
if strings.HasPrefix(name, AppName) {
@@ -433,14 +547,14 @@ func scheduleMonitor(c *Cluster, mon *monConfig) (*apps.Deployment, error) {
// setup affinity settings for pod scheduling
p := cephv1.GetMonPlacement(c.spec.Placement)
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec), PreferredDuringScheduling,
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec),
map[string]string{k8sutil.AppAttr: AppName}, nil)
// setup storage on the canary since scheduling will be affected when
// monitors are configured to use persistent volumes. the pvcName is set to
// the non-empty name of the PVC only when the PVC is created as a result of
// this call to the scheduler.
if c.spec.Mon.VolumeClaimTemplate == nil {
if c.monVolumeClaimTemplate(mon) == nil {
d.Spec.Template.Spec.Volumes = append(d.Spec.Template.Spec.Volumes,
controller.DaemonVolumesDataHostPath(mon.DataPathMap)...)
} else {
@@ -543,7 +657,7 @@ func (c *Cluster) initMonIPs(mons []*monConfig) error {
for _, m := range mons {
if c.spec.Network.IsHost() {
logger.Infof("setting mon endpoints for hostnetwork mode")
node, ok := c.mapping.Node[m.DaemonName]
node, ok := c.mapping.Schedule[m.DaemonName]
if !ok {
return errors.New("mon doesn't exist in assignment map")
}
@@ -601,7 +715,7 @@ func (c *Cluster) assignMons(mons []*monConfig) error {
for _, mon := range mons {
// scheduling for this monitor has already been completed
if _, ok := c.mapping.Node[mon.DaemonName]; ok {
if _, ok := c.mapping.Schedule[mon.DaemonName]; ok {
logger.Debugf("assignmon: mon %s already scheduled", mon.DaemonName)
continue
}
@@ -618,20 +732,20 @@ func (c *Cluster) assignMons(mons []*monConfig) error {
// start waiting for the deployment
monSchedulingWait.Add(1)
go func(deployment *apps.Deployment, monDaemonName string) {
go func(deployment *apps.Deployment, mon *monConfig) {
// signal that the mon is done scheduling
defer monSchedulingWait.Done()
result, err := waitForMonitorScheduling(c, deployment)
if err != nil {
logger.Errorf("failed to schedule mon %q. %v", monDaemonName, err)
logger.Errorf("failed to schedule mon %q. %v", mon.DaemonName, err)
failedMonSchedule = true
return
}
nodeChoice := result.Node
if nodeChoice == nil {
logger.Errorf("assignmon: could not schedule monitor %s", monDaemonName)
logger.Errorf("assignmon: could not schedule monitor %q", mon.DaemonName)
failedMonSchedule = true
return
}
@@ -639,24 +753,32 @@ func (c *Cluster) assignMons(mons []*monConfig) error {
// store nil in the node mapping to indicate that an explicit node
// placement is not being made. otherwise, the node choice will map
// directly to a node selector on the monitor pod.
var nodeInfo *NodeInfo
if c.spec.Network.IsHost() || c.spec.Mon.VolumeClaimTemplate == nil {
logger.Infof("assignmon: mon %s assigned to node %s", monDaemonName, nodeChoice.Name)
nodeInfo, err = getNodeInfoFromNode(*nodeChoice)
var schedule *MonScheduleInfo
if c.spec.Network.IsHost() || c.monVolumeClaimTemplate(mon) == nil {
logger.Infof("assignmon: mon %s assigned to node %s", mon.DaemonName, nodeChoice.Name)
schedule, err = getNodeInfoFromNode(*nodeChoice)
if err != nil {
logger.Errorf("assignmon: couldn't get node info for node %s", nodeChoice.Name)
logger.Errorf("assignmon: couldn't get node info for node %q. %v", nodeChoice.Name, err)
failedMonSchedule = true
return
}
} else {
logger.Infof("assignmon: mon %s placement using native scheduler", monDaemonName)
logger.Infof("assignmon: mon %q placement using native scheduler", mon.DaemonName)
}
if c.spec.IsStretchCluster() {
if schedule == nil {
schedule = &MonScheduleInfo{}
}
logger.Infof("mon %q is assigned to zone %q", mon.DaemonName, mon.Zone)
schedule.Zone = mon.Zone
}
// protect against multiple goroutines updating the status at the same time
resultLock.Lock()
c.mapping.Node[monDaemonName] = nodeInfo
c.mapping.Schedule[mon.DaemonName] = schedule
resultLock.Unlock()
}(deployment, mon.DaemonName)
}(deployment, mon)
}
monSchedulingWait.Wait()
@@ -668,6 +790,25 @@ func (c *Cluster) assignMons(mons []*monConfig) error {
return nil
}
func (c *Cluster) monVolumeClaimTemplate(mon *monConfig) *v1.PersistentVolumeClaim {
if !c.spec.IsStretchCluster() {
return c.spec.Mon.VolumeClaimTemplate
}
// If a stretch cluster, a zone can override the template from the default.
for _, zone := range c.spec.Mon.StretchCluster.Zones {
if zone.Name == mon.Zone {
if zone.VolumeClaimTemplate != nil {
// Found an override for the volume claim template in the zone
return zone.VolumeClaimTemplate
}
break
}
}
// Return the default template since one wasn't found in the zone
return c.spec.Mon.VolumeClaimTemplate
}
func (c *Cluster) startDeployments(mons []*monConfig, requireAllInQuorum bool) error {
if len(mons) == 0 {
return errors.New("cannot start 0 mons")
@@ -705,8 +846,8 @@ func (c *Cluster) startDeployments(mons []*monConfig, requireAllInQuorum bool) e
// Ensure each of the mons have been created. If already created, it will be a no-op.
for i := 0; i < len(mons); i++ {
node := c.mapping.Node[mons[i].DaemonName]
err := c.startMon(mons[i], node)
schedule := c.mapping.Schedule[mons[i].DaemonName]
err := c.startMon(mons[i], schedule)
if err != nil {
if c.isUpgrade {
// if we're upgrading, we don't want to risk the health of the cluster by continuing to upgrade
@@ -874,7 +1015,7 @@ func (c *Cluster) updateMon(m *monConfig, d *apps.Deployment) error {
// - if HostPath -> leave node selector as is
// - if PVC -> remove node selector, if present
//
func (c *Cluster) startMon(m *monConfig, node *NodeInfo) error {
func (c *Cluster) startMon(m *monConfig, schedule *MonScheduleInfo) error {
// check if the monitor deployment already exists. if the deployment does
// exist, also determine if it using pvc storage.
pvcExists := false
@@ -905,7 +1046,7 @@ func (c *Cluster) startMon(m *monConfig, node *NodeInfo) error {
// deployment spec created above does not specify persistent storage. here
// we add in PVC or HostPath storage based on an existing deployment OR on
// the current state of the CRD.
if pvcExists || (!deploymentExists && c.spec.Mon.VolumeClaimTemplate != nil) {
if pvcExists || (!deploymentExists && c.monVolumeClaimTemplate(m) != nil) {
pvcName := m.ResourceName
d.Spec.Template.Spec.Volumes = append(d.Spec.Template.Spec.Volumes, controller.DaemonVolumesDataPVC(pvcName))
controller.AddVolumeMountSubPath(&d.Spec.Template.Spec, "ceph-daemon-data")
@@ -926,16 +1067,16 @@ func (c *Cluster) startMon(m *monConfig, node *NodeInfo) error {
if c.spec.Network.IsHost() || !pvcExists {
p.PodAffinity = nil
p.PodAntiAffinity = nil
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec), PreferredDuringScheduling,
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec),
map[string]string{k8sutil.AppAttr: AppName}, existingDeployment.Spec.Template.Spec.NodeSelector)
} else {
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec), PreferredDuringScheduling,
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec),
map[string]string{k8sutil.AppAttr: AppName}, nil)
}
return c.updateMon(m, d)
}
if c.spec.Mon.VolumeClaimTemplate != nil {
if c.monVolumeClaimTemplate(m) != nil {
pvc, err := c.makeDeploymentPVC(m, false)
if err != nil {
return errors.Wrapf(err, "failed to make mon %s pvc", d.Name)
@@ -950,14 +1091,14 @@ func (c *Cluster) startMon(m *monConfig, node *NodeInfo) error {
}
}
if node == nil {
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec), PreferredDuringScheduling,
if schedule == nil {
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec),
map[string]string{k8sutil.AppAttr: AppName}, nil)
} else {
p.PodAffinity = nil
p.PodAntiAffinity = nil
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec), PreferredDuringScheduling,
map[string]string{k8sutil.AppAttr: AppName}, map[string]string{v1.LabelHostname: node.Hostname})
k8sutil.SetNodeAntiAffinityForPod(&d.Spec.Template.Spec, p, requiredDuringScheduling(&c.spec),
map[string]string{k8sutil.AppAttr: AppName}, map[string]string{v1.LabelHostname: schedule.Hostname})
}
logger.Debugf("Starting mon: %+v", d.Name)
+123 -5
View File
@@ -21,6 +21,7 @@ import (
"io/ioutil"
"os"
"path"
"reflect"
"strconv"
"strings"
"sync"
@@ -28,14 +29,13 @@ import (
"time"
"github.com/pkg/errors"
"github.com/rook/rook/pkg/operator/ceph/config"
"github.com/rook/rook/pkg/operator/k8sutil"
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
"github.com/rook/rook/pkg/clusterd"
"github.com/rook/rook/pkg/daemon/ceph/client"
clienttest "github.com/rook/rook/pkg/daemon/ceph/client/test"
"github.com/rook/rook/pkg/operator/ceph/config"
cephver "github.com/rook/rook/pkg/operator/ceph/version"
"github.com/rook/rook/pkg/operator/k8sutil"
"github.com/rook/rook/pkg/operator/test"
exectest "github.com/rook/rook/pkg/util/exec/test"
"github.com/stretchr/testify/assert"
@@ -120,7 +120,7 @@ func newCluster(context *clusterd.Context, namespace string, allowMultiplePerNod
monPodTimeout: 1 * time.Second,
monTimeoutList: map[string]time.Time{},
mapping: &Mapping{
Node: map[string]*NodeInfo{},
Schedule: map[string]*MonScheduleInfo{},
},
ownerRef: metav1.OwnerReference{},
}
@@ -248,7 +248,7 @@ func TestSaveMonEndpoints(t *testing.T) {
// update the config map
c.ClusterInfo.Monitors["a"].Endpoint = "2.3.4.5:6789"
c.maxMonID = 2
c.mapping.Node["a"] = &NodeInfo{
c.mapping.Schedule["a"] = &MonScheduleInfo{
Name: "node0",
Address: "1.1.1.1",
Hostname: "myhost",
@@ -346,3 +346,121 @@ func TestMonFoundInQuorum(t *testing.T) {
assert.True(t, monFoundInQuorum("c", response))
assert.False(t, monFoundInQuorum("d", response))
}
func TestFindAvailableZoneForStretchedMon(t *testing.T) {
c := &Cluster{spec: cephv1.ClusterSpec{
Mon: cephv1.MonSpec{
StretchCluster: &cephv1.StretchClusterSpec{
Zones: []cephv1.StretchClusterZoneSpec{
{Name: "a", Arbiter: true},
{Name: "b"},
{Name: "c"},
},
},
},
}}
// No mons are assigned to a zone yet
existingMons := []*monConfig{}
availableZone, err := c.findAvailableZoneIfStretched(existingMons)
assert.NoError(t, err)
assert.NotEqual(t, "", availableZone)
// With 3 mons, we have one available zone
existingMons = []*monConfig{
{ResourceName: "x", Zone: "a"},
{ResourceName: "y", Zone: "b"},
}
c.spec.Mon.Count = 3
availableZone, err = c.findAvailableZoneIfStretched(existingMons)
assert.NoError(t, err)
assert.Equal(t, "c", availableZone)
// With 3 mons and no available zones
existingMons = []*monConfig{
{ResourceName: "x", Zone: "a"},
{ResourceName: "y", Zone: "b"},
{ResourceName: "z", Zone: "c"},
}
c.spec.Mon.Count = 3
availableZone, err = c.findAvailableZoneIfStretched(existingMons)
assert.Error(t, err)
assert.Equal(t, "", availableZone)
// With 5 mons and no available zones
existingMons = []*monConfig{
{ResourceName: "w", Zone: "a"},
{ResourceName: "x", Zone: "b"},
{ResourceName: "y", Zone: "b"},
{ResourceName: "z", Zone: "c"},
{ResourceName: "q", Zone: "c"},
}
c.spec.Mon.Count = 5
availableZone, err = c.findAvailableZoneIfStretched(existingMons)
assert.Error(t, err)
assert.Equal(t, "", availableZone)
// With 5 mons and one available zone
existingMons = []*monConfig{
{ResourceName: "w", Zone: "a"},
{ResourceName: "x", Zone: "b"},
{ResourceName: "y", Zone: "b"},
{ResourceName: "z", Zone: "c"},
}
availableZone, err = c.findAvailableZoneIfStretched(existingMons)
assert.NoError(t, err)
assert.Equal(t, "c", availableZone)
// With 5 mons and arbiter zone is available zone
existingMons = []*monConfig{
{ResourceName: "w", Zone: "b"},
{ResourceName: "x", Zone: "b"},
{ResourceName: "y", Zone: "c"},
{ResourceName: "z", Zone: "c"},
}
availableZone, err = c.findAvailableZoneIfStretched(existingMons)
assert.NoError(t, err)
assert.Equal(t, "a", availableZone)
}
func TestStretchMonVolumeClaimTemplate(t *testing.T) {
generalSC := "generalSC"
zoneSC := "zoneSC"
defaultTemplate := &v1.PersistentVolumeClaim{Spec: v1.PersistentVolumeClaimSpec{StorageClassName: &generalSC}}
zoneTemplate := &v1.PersistentVolumeClaim{Spec: v1.PersistentVolumeClaimSpec{StorageClassName: &zoneSC}}
type fields struct {
spec cephv1.ClusterSpec
}
type args struct {
mon *monConfig
}
tests := []struct {
name string
fields fields
args args
want *v1.PersistentVolumeClaim
}{
{"no template", fields{cephv1.ClusterSpec{}}, args{&monConfig{Zone: "z1"}}, nil},
{"default template", fields{cephv1.ClusterSpec{Mon: cephv1.MonSpec{VolumeClaimTemplate: defaultTemplate}}}, args{&monConfig{Zone: "z1"}}, defaultTemplate},
{"default template with 3 zones", fields{cephv1.ClusterSpec{Mon: cephv1.MonSpec{
VolumeClaimTemplate: defaultTemplate,
StretchCluster: &cephv1.StretchClusterSpec{Zones: []cephv1.StretchClusterZoneSpec{{Name: "z1"}, {Name: "z2"}, {Name: "z3"}}}}}},
args{&monConfig{Zone: "z1"}},
defaultTemplate},
{"overridden template", fields{cephv1.ClusterSpec{Mon: cephv1.MonSpec{
VolumeClaimTemplate: defaultTemplate,
StretchCluster: &cephv1.StretchClusterSpec{Zones: []cephv1.StretchClusterZoneSpec{{Name: "z1", VolumeClaimTemplate: zoneTemplate}, {Name: "z2"}, {Name: "z3"}}}}}},
args{&monConfig{Zone: "z1"}},
zoneTemplate},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
c := &Cluster{
spec: tt.fields.spec,
}
if got := c.monVolumeClaimTemplate(tt.args.mon); !reflect.DeepEqual(got, tt.want) {
t.Errorf("Cluster.monVolumeClaimTemplate() = %v, want %v", got, tt.want)
}
})
}
}
+2 -2
View File
@@ -32,8 +32,8 @@ type NodeUsage struct {
MonValid bool
}
func getNodeInfoFromNode(n v1.Node) (*NodeInfo, error) {
nr := &NodeInfo{
func getNodeInfoFromNode(n v1.Node) (*MonScheduleInfo, error) {
nr := &MonScheduleInfo{
Name: n.Name,
Hostname: n.Labels[v1.LabelHostname],
}
+1 -1
View File
@@ -217,7 +217,7 @@ func TestGetNodeInfoFromNode(t *testing.T) {
},
}
var info *NodeInfo
var info *MonScheduleInfo
_, err = getNodeInfoFromNode(*node)
assert.NotNil(t, err)
+35 -4
View File
@@ -20,6 +20,7 @@ import (
"fmt"
"os"
"path"
"strings"
"github.com/pkg/errors"
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
@@ -48,13 +49,35 @@ func (c *Cluster) getLabels(monConfig *monConfig, canary, includeNewLabels bool)
if canary {
labels["mon_canary"] = "true"
}
if c.spec.Mon.VolumeClaimTemplate != nil && includeNewLabels {
labels["pvc_name"] = monConfig.ResourceName
if includeNewLabels {
if c.monVolumeClaimTemplate(monConfig) != nil {
labels["pvc_name"] = monConfig.ResourceName
}
if monConfig.Zone != "" {
labels["stretch-zone"] = monConfig.Zone
}
}
return labels
}
func (c *Cluster) stretchFailureDomainName() string {
label := c.stretchFailureDomainLabel()
index := strings.Index(label, "/")
if index == -1 {
return label
}
return label[index+1:]
}
func (c *Cluster) stretchFailureDomainLabel() string {
if c.spec.Mon.StretchCluster.FailureDomainLabel != "" {
return c.spec.Mon.StretchCluster.FailureDomainLabel
}
// The default topology label is for a zone
return "topology.kubernetes.io/zone"
}
func (c *Cluster) makeDeployment(monConfig *monConfig, canary bool) (*apps.Deployment, error) {
d := &apps.Deployment{
ObjectMeta: metav1.ObjectMeta{
@@ -92,7 +115,7 @@ func (c *Cluster) makeDeployment(monConfig *monConfig, canary bool) (*apps.Deplo
}
func (c *Cluster) makeDeploymentPVC(m *monConfig, canary bool) (*v1.PersistentVolumeClaim, error) {
template := c.spec.Mon.VolumeClaimTemplate
template := c.monVolumeClaimTemplate(m)
volumeMode := v1.PersistentVolumeFilesystem
pvc := &v1.PersistentVolumeClaim{
ObjectMeta: metav1.ObjectMeta{
@@ -156,7 +179,7 @@ func (c *Cluster) makeMonPod(monConfig *monConfig, canary bool) (*v1.Pod, error)
}
// Replace default unreachable node toleration
if c.spec.Mon.VolumeClaimTemplate != nil {
if c.monVolumeClaimTemplate(monConfig) != nil {
k8sutil.AddUnreachableNodeToleration(&podSpec)
}
@@ -179,6 +202,14 @@ func (c *Cluster) makeMonPod(monConfig *monConfig, canary bool) (*v1.Pod, error)
}
}
if c.spec.IsStretchCluster() {
nodeAffinity, err := k8sutil.GenerateNodeAffinity(fmt.Sprintf("%s=%s", c.stretchFailureDomainLabel(), monConfig.Zone))
if err != nil {
return nil, errors.Wrapf(err, "failed to generate mon %q node affinity", monConfig.DaemonName)
}
pod.Spec.Affinity = &v1.Affinity{NodeAffinity: nodeAffinity}
}
return pod, nil
}
+1 -1
View File
@@ -632,7 +632,7 @@ func removeDuplicateEnvVarsFromContainer(container *v1.Container) {
vars := []v1.EnvVar{}
for _, v := range container.Env {
if _, ok := foundVars[v.Name]; ok {
logger.Infof("duplicate env var %q skipped on container %q", v.Name, container.Name)
logger.Debugf("duplicate env var %q skipped on container %q", v.Name, container.Name)
continue
}
+1 -2
View File
@@ -106,9 +106,8 @@ func (c *clusterConfig) makeRGWPodSpec(rgwConfig *rgwConfig) (v1.PodTemplateSpec
}
// If host networking is not enabled, preferred pod anti-affinity is added to the rgw daemons
preferredDuringScheduling := true
labels := getLabels(c.store.Name, c.store.Namespace, false)
k8sutil.SetNodeAntiAffinityForPod(&podSpec, c.store.Spec.Gateway.Placement, c.clusterSpec.Network.IsHost(), preferredDuringScheduling, labels, nil)
k8sutil.SetNodeAntiAffinityForPod(&podSpec, c.store.Spec.Gateway.Placement, c.clusterSpec.Network.IsHost(), labels, nil)
podTemplateSpec := v1.PodTemplateSpec{
ObjectMeta: metav1.ObjectMeta{
+2 -2
View File
@@ -300,7 +300,7 @@ func ClusterDaemonEnvVars(image string) []v1.EnvVar {
}
// SetNodeAntiAffinityForPod assign pod anti-affinity when pod should not be co-located
func SetNodeAntiAffinityForPod(pod *v1.PodSpec, p rookv1.Placement, requiredDuringScheduling, preferredDuringScheduling bool,
func SetNodeAntiAffinityForPod(pod *v1.PodSpec, p rookv1.Placement, requiredDuringScheduling bool,
labels, nodeSelector map[string]string) {
p.ApplyToPodSpec(pod)
pod.NodeSelector = nodeSelector
@@ -329,7 +329,7 @@ func SetNodeAntiAffinityForPod(pod *v1.PodSpec, p rookv1.Placement, requiredDuri
if requiredDuringScheduling {
paa.RequiredDuringSchedulingIgnoredDuringExecution =
append(paa.RequiredDuringSchedulingIgnoredDuringExecution, podAntiAffinity)
} else if preferredDuringScheduling {
} else {
paa.PreferredDuringSchedulingIgnoredDuringExecution =
append(paa.PreferredDuringSchedulingIgnoredDuringExecution, v1.WeightedPodAffinityTerm{
Weight: 50,
+7 -15
View File
@@ -177,22 +177,19 @@ func TestAddUnreachableNodeToleration(t *testing.T) {
}
func testPodSpecPlacement(t *testing.T, requiredDuringScheduling, preferredDuringScheduling bool, req, pref int, placement *rookv1.Placement) {
func testPodSpecPlacement(t *testing.T, requiredDuringScheduling bool, req, pref int, placement *rookv1.Placement) {
spec := v1.PodSpec{
InitContainers: []v1.Container{},
Containers: []v1.Container{},
RestartPolicy: v1.RestartPolicyAlways,
}
SetNodeAntiAffinityForPod(&spec, *placement, requiredDuringScheduling, preferredDuringScheduling, map[string]string{"app": "mon"}, nil)
SetNodeAntiAffinityForPod(&spec, *placement, requiredDuringScheduling, map[string]string{"app": "mon"}, nil)
// should have a required anti-affinity and no preferred anti-affinity
assert.Equal(t,
req,
len(spec.Affinity.PodAntiAffinity.RequiredDuringSchedulingIgnoredDuringExecution))
assert.Equal(t,
pref,
len(spec.Affinity.PodAntiAffinity.PreferredDuringSchedulingIgnoredDuringExecution))
}
func makePlacement() rookv1.Placement {
@@ -217,18 +214,13 @@ func makePlacement() rookv1.Placement {
func TestPodSpecPlacement(t *testing.T) {
// no placement settings in the crd
p := rookv1.Placement{}
testPodSpecPlacement(t, true, true, 1, 0, &p)
testPodSpecPlacement(t, true, false, 1, 0, &p)
testPodSpecPlacement(t, false, true, 0, 1, &p)
testPodSpecPlacement(t, false, false, 0, 0, &p)
testPodSpecPlacement(t, true, 1, 0, &p)
testPodSpecPlacement(t, false, 0, 1, &p)
testPodSpecPlacement(t, false, 0, 0, &p)
// crd has other preferred and required anti-affinity setting
p = makePlacement()
testPodSpecPlacement(t, true, true, 2, 1, &p)
testPodSpecPlacement(t, true, 2, 1, &p)
p = makePlacement()
testPodSpecPlacement(t, true, false, 2, 1, &p)
p = makePlacement()
testPodSpecPlacement(t, false, true, 1, 2, &p)
p = makePlacement()
testPodSpecPlacement(t, false, false, 1, 1, &p)
testPodSpecPlacement(t, false, 1, 2, &p)
}
@@ -170,6 +170,26 @@ spec:
volumeClaimTemplate:
type: object
x-kubernetes-preserve-unknown-fields: true
stretchCluster:
type: object
nullable: true
properties:
failureDomainLabel:
type: string
subFailureDomain:
type: string
zones:
type: array
items:
type: object
properties:
name:
type: string
arbiter:
type: boolean
volumeClaimTemplate:
type: object
x-kubernetes-preserve-unknown-fields: true
mgr:
type: object
properties: