forked from rook/rook
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:
@@ -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.
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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) {
|
||||
{
|
||||
|
||||
@@ -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{}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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],
|
||||
}
|
||||
|
||||
@@ -217,7 +217,7 @@ func TestGetNodeInfoFromNode(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
var info *NodeInfo
|
||||
var info *MonScheduleInfo
|
||||
_, err = getNodeInfoFromNode(*node)
|
||||
assert.NotNil(t, err)
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user