core: add observedGeneration to CR status

adding observedGeneration field in the cephcluster cr
status for having better control on reconciling,
as observedGeneration field will be updated by the controller

Closes: https://github.com/rook/rook/issues/9673

Signed-off-by: parth-gr <paarora@redhat.com>
This commit is contained in:
parth-gr
2022-03-16 19:58:35 +05:30
parent 1bbe39a81b
commit 2dfd64a97c
31 changed files with 405 additions and 133 deletions
@@ -295,6 +295,10 @@ spec:
type: object
type: object
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -474,6 +478,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -621,6 +629,10 @@ spec:
description: The ARN of the topic generated by the RGW
nullable: true
type: string
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -691,6 +703,10 @@ spec:
type: string
nullable: true
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -4222,6 +4238,10 @@ spec:
type: array
message:
type: string
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -4876,6 +4896,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -6234,6 +6258,10 @@ spec:
description: LastChecked is the last time time the status was checked
type: string
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -6363,6 +6391,10 @@ spec:
type: string
nullable: true
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -7025,6 +7057,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -7093,6 +7129,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -8514,6 +8554,10 @@ spec:
type: object
message:
type: string
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -8652,6 +8696,10 @@ spec:
type: string
nullable: true
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -8713,6 +8761,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -9102,6 +9154,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -9752,6 +9808,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
+60
View File
@@ -298,6 +298,10 @@ spec:
type: object
type: object
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -476,6 +480,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -622,6 +630,10 @@ spec:
description: The ARN of the topic generated by the RGW
nullable: true
type: string
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -691,6 +703,10 @@ spec:
type: string
nullable: true
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -4221,6 +4237,10 @@ spec:
type: array
message:
type: string
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -4874,6 +4894,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -6231,6 +6255,10 @@ spec:
description: LastChecked is the last time time the status was checked
type: string
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -6359,6 +6387,10 @@ spec:
type: string
nullable: true
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -7020,6 +7052,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -7087,6 +7123,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -8507,6 +8547,10 @@ spec:
type: object
message:
type: string
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
description: ConditionType represent a resource's status
type: string
@@ -8644,6 +8688,10 @@ spec:
type: string
nullable: true
type: object
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -8704,6 +8752,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -9092,6 +9144,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
@@ -9741,6 +9797,10 @@ spec:
status:
description: Status represents the status of an object
properties:
observedGeneration:
description: ObservedGeneration is the latest generation observed by the controller.
format: int64
type: integer
phase:
type: string
type: object
+27
View File
@@ -315,6 +315,9 @@ type ClusterStatus struct {
CephStatus *CephStatus `json:"ceph,omitempty"`
CephStorage *CephStorage `json:"storage,omitempty"`
CephVersion *ClusterVersion `json:"version,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// CephDaemonsVersions show the current ceph version for different ceph daemons
@@ -685,6 +688,9 @@ type CephBlockPoolStatus struct {
// +optional
// +nullable
Info map[string]string `json:"info,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// MirroringStatusSpec is the status of the pool mirroring
@@ -843,6 +849,9 @@ type SnapshotSchedule struct {
type Status struct {
// +optional
Phase string `json:"phase,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// ReplicatedSpec represents the spec for replication in a pool
@@ -1109,6 +1118,9 @@ type CephFilesystemStatus struct {
// +optional
MirroringStatus *FilesystemMirroringInfoSpec `json:"mirroringStatus,omitempty"`
Conditions []Condition `json:"conditions,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// FilesystemMirroringInfo is the status of the pool mirroring
@@ -1421,6 +1433,9 @@ type ObjectStoreStatus struct {
// +nullable
Info map[string]string `json:"info,omitempty"`
Conditions []Condition `json:"conditions,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// BucketStatus represents the status of a bucket
@@ -1457,6 +1472,9 @@ type ObjectStoreUserStatus struct {
// +optional
// +nullable
Info map[string]string `json:"info,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// +k8s:deepcopy-gen:interfaces=k8s.io/apimachinery/pkg/runtime.Object
@@ -1646,6 +1664,9 @@ type BucketTopicStatus struct {
// +optional
// +nullable
ARN *string `json:"ARN,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// CephBucketTopicList represents a list Ceph Object Store Bucket Notification Topics
@@ -2036,6 +2057,9 @@ type CephClientStatus struct {
// +optional
// +nullable
Info map[string]string `json:"info,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
// CleanupPolicySpec represents a Ceph Cluster cleanup policy
@@ -2400,4 +2424,7 @@ type CephFilesystemSubVolumeGroupStatus struct {
// +optional
// +nullable
Info map[string]string `json:"info,omitempty"`
// ObservedGeneration is the latest generation observed by the controller.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
}
+14 -6
View File
@@ -138,6 +138,10 @@ func (r *ReconcileCephClient) reconcile(request reconcile.Request) (reconcile.Re
// Error reading the object - requeue the request.
return reconcile.Result{}, *cephClient, errors.Wrap(err, "failed to get cephClient")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephClient.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephClient)
@@ -147,7 +151,7 @@ func (r *ReconcileCephClient) reconcile(request reconcile.Request) (reconcile.Re
// The CR was just created, initializing status fields
if cephClient.Status == nil {
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionProgressing)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, cephv1.ConditionProgressing)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -212,12 +216,13 @@ func (r *ReconcileCephClient) reconcile(request reconcile.Request) (reconcile.Re
logger.Info(opcontroller.OperatorNotInitializedMessage)
return opcontroller.WaitForRequeueIfOperatorNotInitialized, *cephClient, nil
}
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, cephv1.ConditionFailure)
return reconcile.Result{}, *cephClient, errors.Wrapf(err, "failed to create or update client %q", cephClient.Name)
}
// update status with latest ObservedGeneration value at the end of reconcile
// Success! Let's update the status
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady)
r.updateStatus(observedGeneration, request.NamespacedName, cephv1.ConditionReady)
// Return and do not requeue
logger.Debug("done reconciling")
@@ -335,9 +340,9 @@ func generateClientName(name string) string {
}
// updateStatus updates an object with a given status
func (r *ReconcileCephClient) updateStatus(client client.Client, name types.NamespacedName, status cephv1.ConditionType) {
func (r *ReconcileCephClient) updateStatus(observedGeneration int64, name types.NamespacedName, status cephv1.ConditionType) {
cephClient := &cephv1.CephClient{}
if err := client.Get(r.opManagerContext, name, cephClient); err != nil {
if err := r.client.Get(r.opManagerContext, name, cephClient); err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephClient resource not found. Ignoring since object must be deleted.")
return
@@ -353,7 +358,10 @@ func (r *ReconcileCephClient) updateStatus(client client.Client, name types.Name
if cephClient.Status.Phase == cephv1.ConditionReady {
cephClient.Status.Info = generateStatusInfo(cephClient)
}
if err := reporting.UpdateStatus(client, cephClient); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
cephClient.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, cephClient); err != nil {
logger.Errorf("failed to set ceph client %q status to %q. %v", name, status, err)
return
}
+1 -1
View File
@@ -198,7 +198,7 @@ func (c *cephStatusChecker) updateCephStatus(status *cephclient.CephStatus, cond
// Update condition
logger.Debugf("updating ceph cluster %q status and condition to %+v, %v, %s, %s", clusterName.Namespace, status, conditionStatus, reason, message)
opcontroller.UpdateClusterCondition(c.context, cephCluster, c.clusterInfo.NamespacedName(), condition, conditionStatus, reason, message, true)
opcontroller.UpdateClusterCondition(c.context, cephCluster, c.clusterInfo.NamespacedName(), k8sutil.ObservedGenerationNotAvailable, condition, conditionStatus, reason, message, true)
}
// toCustomResourceStatus converts the ceph status to the struct expected for the CephCluster CR status
+13 -8
View File
@@ -59,6 +59,7 @@ type cluster struct {
ownerInfo *k8sutil.OwnerInfo
isUpgrade bool
monitoringRoutines map[string]*clusterHealth
observedGeneration int64
}
type clusterHealth struct {
@@ -79,6 +80,10 @@ func newCluster(c *cephv1.CephCluster, context *clusterd.Context, ownerInfo *k8s
monitoringRoutines: make(map[string]*clusterHealth),
ownerInfo: ownerInfo,
mons: mon.New(context, c.Namespace, c.Spec, ownerInfo),
// update observedGeneration with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration: c.ObjectMeta.Generation,
}
}
@@ -100,7 +105,7 @@ func (c *cluster) reconcileCephDaemons(rookImage string, cephVersion cephver.Cep
}
// Start the mon pods
controller.UpdateCondition(c.ClusterInfo.Context, c.context, c.namespacedName, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring Ceph Mons")
controller.UpdateCondition(c.ClusterInfo.Context, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring Ceph Mons")
clusterInfo, err := c.mons.Start(c.ClusterInfo, rookImage, cephVersion, *c.Spec)
if err != nil {
return errors.Wrap(err, "failed to start ceph monitors")
@@ -128,7 +133,7 @@ func (c *cluster) reconcileCephDaemons(rookImage string, cephVersion cephver.Cep
}
// Start Ceph manager
controller.UpdateCondition(c.ClusterInfo.Context, c.context, c.namespacedName, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring Ceph Mgr(s)")
controller.UpdateCondition(c.ClusterInfo.Context, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring Ceph Mgr(s)")
mgrs := mgr.New(c.context, c.ClusterInfo, *c.Spec, rookImage)
err = mgrs.Start()
if err != nil {
@@ -136,7 +141,7 @@ func (c *cluster) reconcileCephDaemons(rookImage string, cephVersion cephver.Cep
}
// Start the OSDs
controller.UpdateCondition(c.ClusterInfo.Context, c.context, c.namespacedName, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring Ceph OSDs")
controller.UpdateCondition(c.ClusterInfo.Context, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring Ceph OSDs")
osds := osd.New(c.context, c.ClusterInfo, *c.Spec, rookImage)
err = osds.Start()
if err != nil {
@@ -177,7 +182,7 @@ func (c *ClusterController) initializeCluster(cluster *cluster) error {
if cluster.Spec.External.Enable {
err := c.configureExternalCephCluster(cluster)
if err != nil {
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, cephv1.ConditionProgressing, v1.ConditionFalse, cephv1.ClusterProgressingReason, err.Error())
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionFalse, cephv1.ClusterProgressingReason, err.Error())
return errors.Wrap(err, "failed to configure external ceph cluster")
}
} else {
@@ -205,7 +210,7 @@ func (c *ClusterController) initializeCluster(cluster *cluster) error {
err = c.configureLocalCephCluster(cluster)
if err != nil {
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, cephv1.ConditionProgressing, v1.ConditionFalse, cephv1.ClusterProgressingReason, err.Error())
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionFalse, cephv1.ClusterProgressingReason, err.Error())
return errors.Wrap(err, "failed to configure local ceph cluster")
}
}
@@ -227,7 +232,7 @@ func (c *ClusterController) configureLocalCephCluster(cluster *cluster) error {
}
// Run image validation job
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Detecting Ceph version")
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Detecting Ceph version")
cephVersion, isUpgrade, err := c.detectAndValidateCephVersion(cluster)
if err != nil {
return errors.Wrap(err, "failed the ceph version check")
@@ -242,7 +247,7 @@ func (c *ClusterController) configureLocalCephCluster(cluster *cluster) error {
}
}
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring the Ceph cluster")
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, "Configuring the Ceph cluster")
cluster.ClusterInfo.Context = c.OpManagerCtx
// Run the orchestration
@@ -252,7 +257,7 @@ func (c *ClusterController) configureLocalCephCluster(cluster *cluster) error {
}
// Set the condition to the cluster object
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, cephv1.ConditionReady, v1.ConditionTrue, cephv1.ClusterCreatedReason, "Cluster created successfully")
controller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, cluster.observedGeneration, cephv1.ConditionReady, v1.ConditionTrue, cephv1.ClusterCreatedReason, "Cluster created successfully")
return nil
}
@@ -44,7 +44,7 @@ func (c *ClusterController) configureExternalCephCluster(cluster *cluster) error
return errors.Wrap(err, "failed to validate external cluster specs")
}
opcontroller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, cephv1.ConditionConnecting, v1.ConditionTrue, cephv1.ClusterConnectingReason, "Attempting to connect to an external Ceph cluster")
opcontroller.UpdateCondition(c.OpManagerCtx, c.context, c.namespacedName, k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionConnecting, v1.ConditionTrue, cephv1.ClusterConnectingReason, "Attempting to connect to an external Ceph cluster")
// loop until we find the secret necessary to connect to the external cluster
// then populate clusterInfo
+3 -1
View File
@@ -269,7 +269,7 @@ func (r *ReconcileCephCluster) reconcileDelete(cephCluster *cephv1.CephCluster)
// Set the deleting status
opcontroller.UpdateClusterCondition(r.context, cephCluster, nsName,
cephv1.ConditionDeleting, corev1.ConditionTrue, cephv1.ClusterDeletingReason, "Deleting the CephCluster",
k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionDeleting, corev1.ConditionTrue, cephv1.ClusterDeletingReason, "Deleting the CephCluster",
true /* keep all other conditions to be safe */)
deps, err := CephClusterDependents(r.context, cephCluster.Namespace)
@@ -343,6 +343,8 @@ func (c *ClusterController) reconcileCephCluster(clusterObj *cephv1.CephCluster,
cluster = newCluster(clusterObj, c.context, ownerInfo)
}
cluster.namespacedName = c.namespacedName
// updating observedGeneration in cluster if it's not the first reconcile
cluster.observedGeneration = clusterObj.ObjectMeta.Generation
// Pass down the client to interact with Kubernetes objects
// This will be used later down by spec code to create objects like deployment, services etc
+2 -2
View File
@@ -384,7 +384,7 @@ func createDaemonOnPVC(c *Cluster, osd OSDInfo, pvcName string, config *provisio
}
message := fmt.Sprintf("Processing OSD %d on PVC %q", osd.ID, pvcName)
updateConditionFunc(c.clusterInfo.Context, c.context, c.clusterInfo.NamespacedName(), cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
updateConditionFunc(c.clusterInfo.Context, c.context, c.clusterInfo.NamespacedName(), k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
_, err = k8sutil.CreateDeployment(c.clusterInfo.Context, c.context.Clientset, d)
return errors.Wrapf(err, "failed to create deployment for OSD %d on PVC %q", osd.ID, pvcName)
@@ -397,7 +397,7 @@ func createDaemonOnNode(c *Cluster, osd OSDInfo, nodeName string, config *provis
}
message := fmt.Sprintf("Processing OSD %d on node %q", osd.ID, nodeName)
updateConditionFunc(c.clusterInfo.Context, c.context, c.clusterInfo.NamespacedName(), cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
updateConditionFunc(c.clusterInfo.Context, c.context, c.clusterInfo.NamespacedName(), k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
_, err = k8sutil.CreateDeployment(c.clusterInfo.Context, c.context.Clientset, d)
return errors.Wrapf(err, "failed to create deployment for OSD %d on node %q", osd.ID, nodeName)
@@ -110,7 +110,7 @@ func testOSDIntegration(t *testing.T) {
}()
// stub out the conditionExportFunc to do nothing. we do not have a fake Rook interface that
// allows us to interact with a CephCluster resource like the fake K8s clientset.
updateConditionFunc = func(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, conditionType cephv1.ConditionType, status corev1.ConditionStatus, reason cephv1.ConditionReason, message string) {
updateConditionFunc = func(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, observedGeneration int64, conditionType cephv1.ConditionType, status corev1.ConditionStatus, reason cephv1.ConditionReason, message string) {
// do nothing
}
+1 -1
View File
@@ -132,7 +132,7 @@ func TestAddRemoveNode(t *testing.T) {
}()
// stub out the conditionExportFunc to do nothing. we do not have a fake Rook interface that
// allows us to interact with a CephCluster resource like the fake K8s clientset.
updateConditionFunc = func(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, conditionType cephv1.ConditionType, status corev1.ConditionStatus, reason cephv1.ConditionReason, message string) {
updateConditionFunc = func(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, observedGeneration int64, conditionType cephv1.ConditionType, status corev1.ConditionStatus, reason cephv1.ConditionReason, message string) {
// do nothing
}
+2 -2
View File
@@ -150,7 +150,7 @@ func (c *updateConfig) updateExistingOSDs(errs *provisionErrors) {
updatedDep, err = deploymentOnPVCFunc(c.cluster, osdInfo, nodeOrPVCName, c.provisionConfig)
message := fmt.Sprintf("Processing OSD %d on PVC %q", osdID, nodeOrPVCName)
updateConditionFunc(c.cluster.clusterInfo.Context, c.cluster.context, c.cluster.clusterInfo.NamespacedName(), cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
updateConditionFunc(c.cluster.clusterInfo.Context, c.cluster.context, c.cluster.clusterInfo.NamespacedName(), k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
} else {
if !c.cluster.ValidStorage.NodeExists(nodeOrPVCName) {
// node will not reconcile, so don't update the deployment
@@ -167,7 +167,7 @@ func (c *updateConfig) updateExistingOSDs(errs *provisionErrors) {
updatedDep, err = deploymentOnNodeFunc(c.cluster, osdInfo, nodeOrPVCName, c.provisionConfig)
message := fmt.Sprintf("Processing OSD %d on node %q", osdID, nodeOrPVCName)
updateConditionFunc(c.cluster.clusterInfo.Context, c.cluster.context, c.cluster.clusterInfo.NamespacedName(), cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
updateConditionFunc(c.cluster.clusterInfo.Context, c.cluster.context, c.cluster.clusterInfo.NamespacedName(), k8sutil.ObservedGenerationNotAvailable, cephv1.ConditionProgressing, v1.ConditionTrue, cephv1.ClusterProgressingReason, message)
}
if err != nil {
errs.addError("%v", errors.Wrapf(err, "failed to update OSD %d", osdID))
+1 -1
View File
@@ -121,7 +121,7 @@ func Test_updateExistingOSDs(t *testing.T) {
// stub out the conditionExportFunc to do nothing. we do not have a fake Rook interface that
// allows us to interact with a CephCluster resource like the fake K8s clientset.
updateConditionFunc = func(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, conditionType cephv1.ConditionType, status corev1.ConditionStatus, reason cephv1.ConditionReason, message string) {
updateConditionFunc = func(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, observedGeneration int64, conditionType cephv1.ConditionType, status corev1.ConditionStatus, reason cephv1.ConditionReason, message string) {
// do nothing
}
shouldCheckOkToStopFunc = func(context *clusterd.Context, clusterInfo *cephclient.ClusterInfo) bool {
+14 -6
View File
@@ -144,7 +144,7 @@ func (r *ReconcileCephRBDMirror) Reconcile(context context.Context, request reco
// workaround because the rook logging mechanism is not compatible with the controller-runtime logging interface
reconcileResponse, cephRBDMirror, err := r.reconcile(request)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.FailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.FailedStatus)
logger.Errorf("failed to reconcile %v", err)
}
@@ -166,8 +166,12 @@ func (r *ReconcileCephRBDMirror) reconcile(request reconcile.Request) (reconcile
// The CR was just created, initializing status fields
if cephRBDMirror.Status == nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus)
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephRBDMirror.ObjectMeta.Generation
// validate the pool settings
if err := validateSpec(&cephRBDMirror.Spec); err != nil {
@@ -232,8 +236,9 @@ func (r *ReconcileCephRBDMirror) reconcile(request reconcile.Request) (reconcile
return opcontroller.ImmediateRetryResult, *cephRBDMirror, errors.Wrap(err, "failed to create ceph rbd mirror deployments")
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReadyStatus)
r.updateStatus(observedGeneration, request.NamespacedName, k8sutil.ReadyStatus)
// Return and do not requeue
logger.Debug("done reconciling ceph rbd mirror")
@@ -260,9 +265,9 @@ func (r *ReconcileCephRBDMirror) reconcileCreateCephRBDMirror(cephRBDMirror *cep
}
// updateStatus updates an object with a given status
func (r *ReconcileCephRBDMirror) updateStatus(client client.Client, name types.NamespacedName, status string) {
func (r *ReconcileCephRBDMirror) updateStatus(observedGeneration int64, name types.NamespacedName, status string) {
rbdMirror := &cephv1.CephRBDMirror{}
err := client.Get(r.opManagerContext, name, rbdMirror)
err := r.client.Get(r.opManagerContext, name, rbdMirror)
if err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephRBDMirror resource not found. Ignoring since object must be deleted.")
@@ -277,7 +282,10 @@ func (r *ReconcileCephRBDMirror) updateStatus(client client.Client, name types.N
}
rbdMirror.Status.Phase = status
if err := reporting.UpdateStatus(client, rbdMirror); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
rbdMirror.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, rbdMirror); err != nil {
logger.Errorf("failed to set rbd mirror %q status to %q. %v", rbdMirror.Name, status, err)
return
}
+8 -3
View File
@@ -24,13 +24,14 @@ import (
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
"github.com/rook/rook/pkg/clusterd"
"github.com/rook/rook/pkg/operator/ceph/reporting"
"github.com/rook/rook/pkg/operator/k8sutil"
v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
)
// UpdateCondition function will export each condition into the cluster custom resource
func UpdateCondition(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, conditionType cephv1.ConditionType, status v1.ConditionStatus, reason cephv1.ConditionReason, message string) {
func UpdateCondition(ctx context.Context, c *clusterd.Context, namespaceName types.NamespacedName, observedGeneration int64, conditionType cephv1.ConditionType, status v1.ConditionStatus, reason cephv1.ConditionReason, message string) {
// use client.Client unit test this more easily with updating statuses which must use the client
cluster := &cephv1.CephCluster{}
if err := c.Client.Get(ctx, namespaceName, cluster); err != nil {
@@ -38,11 +39,11 @@ func UpdateCondition(ctx context.Context, c *clusterd.Context, namespaceName typ
return
}
UpdateClusterCondition(c, cluster, namespaceName, conditionType, status, reason, message, false)
UpdateClusterCondition(c, cluster, namespaceName, observedGeneration, conditionType, status, reason, message, false)
}
// UpdateClusterCondition function will export each condition into the cluster custom resource
func UpdateClusterCondition(c *clusterd.Context, cluster *cephv1.CephCluster, namespaceName types.NamespacedName, conditionType cephv1.ConditionType, status v1.ConditionStatus,
func UpdateClusterCondition(c *clusterd.Context, cluster *cephv1.CephCluster, namespaceName types.NamespacedName, observedGeneration int64, conditionType cephv1.ConditionType, status v1.ConditionStatus,
reason cephv1.ConditionReason, message string, preserveAllConditions bool) {
// Keep the conditions that already existed if they are in the list of long-term conditions,
@@ -88,6 +89,10 @@ func UpdateClusterCondition(c *clusterd.Context, cluster *cephv1.CephCluster, na
}
conditions = append(conditions, *currentCondition)
cluster.Status.Conditions = conditions
// update observed generation
if observedGeneration != k8sutil.ObservedGenerationNotAvailable && conditionType == cephv1.ConditionReady {
cluster.Status.ObservedGeneration = observedGeneration
}
// Once the cluster begins deleting, the phase should not revert back to any other phase
if cluster.Status.Phase != cephv1.ConditionDeleting {
+12 -5
View File
@@ -186,6 +186,11 @@ func (r *ReconcileCephFilesystem) reconcile(request reconcile.Request) (reconcil
return reconcile.Result{}, *cephFilesystem, errors.Wrap(err, "failed to get cephFilesystem")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephFilesystem.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephFilesystem)
if err != nil {
@@ -194,7 +199,7 @@ func (r *ReconcileCephFilesystem) reconcile(request reconcile.Request) (reconcil
// The CR was just created, initializing status fields
if cephFilesystem.Status == nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus, nil)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus, nil)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -324,7 +329,7 @@ func (r *ReconcileCephFilesystem) reconcile(request reconcile.Request) (reconcil
logger.Debug("reconciling ceph filesystem store deployments")
reconcileResponse, err = r.reconcileCreateFilesystem(cephFilesystem)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure, nil)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, cephv1.ConditionFailure, nil)
return reconcileResponse, *cephFilesystem, err
}
@@ -352,7 +357,7 @@ func (r *ReconcileCephFilesystem) reconcile(request reconcile.Request) (reconcil
logger.Info("reconciling create cephfs-mirror peer configuration")
reconcileResponse, err = opcontroller.CreateBootstrapPeerSecret(r.context, r.clusterInfo, cephFilesystem, k8sutil.NewOwnerInfo(cephFilesystem, r.scheme))
if err != nil {
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure, nil)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, cephv1.ConditionFailure, nil)
return reconcileResponse, *cephFilesystem,
errors.Wrapf(err, "failed to create cephfs-mirror bootstrap peer for filesystem %q.", cephFilesystem.Name)
}
@@ -364,8 +369,9 @@ func (r *ReconcileCephFilesystem) reconcile(request reconcile.Request) (reconcil
errors.Wrapf(err, "failed to configure mirroring for filesystem %q.", cephFilesystem.Name)
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady, opcontroller.GenerateStatusInfo(cephFilesystem))
r.updateStatus(observedGeneration, request.NamespacedName, cephv1.ConditionReady, opcontroller.GenerateStatusInfo(cephFilesystem))
statusUpdated = true
// Run go routine check for mirroring status
@@ -383,9 +389,10 @@ func (r *ReconcileCephFilesystem) reconcile(request reconcile.Request) (reconcil
}
}
if !statusUpdated {
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling$
// TODO: set status to Ready **only** if the filesystem is ready
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady, nil)
r.updateStatus(observedGeneration, request.NamespacedName, cephv1.ConditionReady, nil)
}
return reconcile.Result{}, *cephFilesystem, nil
+16 -6
View File
@@ -135,7 +135,8 @@ func (r *ReconcileFilesystemMirror) Reconcile(context context.Context, request r
// workaround because the rook logging mechanism is not compatible with the controller-runtime logging interface
reconcileResponse, cephFilesystemMirror, err := r.reconcile(request)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.FailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.FailedStatus)
logger.Errorf("failed to reconcile %v", err)
}
return reporting.ReportReconcileResult(logger, r.recorder, request, &cephFilesystemMirror, reconcileResponse, err)
@@ -154,9 +155,14 @@ func (r *ReconcileFilesystemMirror) reconcile(request reconcile.Request) (reconc
return reconcile.Result{}, *filesystemMirror, errors.Wrap(err, "failed to get CephFilesystemMirror")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := filesystemMirror.ObjectMeta.Generation
// The CR was just created, initializing status fields
if filesystemMirror.Status == nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -218,8 +224,9 @@ func (r *ReconcileFilesystemMirror) reconcile(request reconcile.Request) (reconc
return opcontroller.ImmediateRetryResult, *filesystemMirror, errors.Wrap(err, "failed to create ceph filesystem mirror deployments")
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReadyStatus)
r.updateStatus(observedGeneration, request.NamespacedName, k8sutil.ReadyStatus)
// Return and do not requeue
logger.Debug("done reconciling ceph filesystem mirror")
@@ -246,9 +253,9 @@ func (r *ReconcileFilesystemMirror) reconcileFilesystemMirror(filesystemMirror *
}
// updateStatus updates an object with a given status
func (r *ReconcileFilesystemMirror) updateStatus(client client.Client, name types.NamespacedName, status string) {
func (r *ReconcileFilesystemMirror) updateStatus(observedGeneration int64, name types.NamespacedName, status string) {
fsMirror := &cephv1.CephFilesystemMirror{}
err := client.Get(r.opManagerContext, name, fsMirror)
err := r.client.Get(r.opManagerContext, name, fsMirror)
if err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephFilesystemMirror resource not found. Ignoring since object must be deleted.")
@@ -263,7 +270,10 @@ func (r *ReconcileFilesystemMirror) updateStatus(client client.Client, name type
}
fsMirror.Status.Phase = status
if err := reporting.UpdateStatus(client, fsMirror); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
fsMirror.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, fsMirror); err != nil {
logger.Errorf("failed to set filesystem mirror %q status to %q. %v", fsMirror.Name, status, err)
return
}
+7 -4
View File
@@ -22,15 +22,15 @@ import (
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
"github.com/rook/rook/pkg/operator/ceph/reporting"
"github.com/rook/rook/pkg/operator/k8sutil"
kerrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
)
// updateStatus updates a fs CR with the given status
func (r *ReconcileCephFilesystem) updateStatus(client client.Client, namespacedName types.NamespacedName, status cephv1.ConditionType, info map[string]string) {
func (r *ReconcileCephFilesystem) updateStatus(observedGeneration int64, namespacedName types.NamespacedName, status cephv1.ConditionType, info map[string]string) {
fs := &cephv1.CephFilesystem{}
err := client.Get(r.opManagerContext, namespacedName, fs)
err := r.client.Get(r.opManagerContext, namespacedName, fs)
if err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephFilesystem resource not found. Ignoring since object must be deleted.")
@@ -46,7 +46,10 @@ func (r *ReconcileCephFilesystem) updateStatus(client client.Client, namespacedN
fs.Status.Phase = status
fs.Status.Info = info
if err := reporting.UpdateStatus(client, fs); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
fs.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, fs); err != nil {
logger.Warningf("failed to set filesystem %q status to %q. %v", fs.Name, status, err)
return
}
@@ -131,6 +131,10 @@ func (r *ReconcileCephFilesystemSubVolumeGroup) reconcile(request reconcile.Requ
// Error reading the object - requeue the request.
return reconcile.Result{}, errors.Wrap(err, "failed to get cephFilesystemSubVolumeGroup")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephFilesystemSubVolumeGroup.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephFilesystemSubVolumeGroup)
@@ -140,7 +144,7 @@ func (r *ReconcileCephFilesystemSubVolumeGroup) reconcile(request reconcile.Requ
// The CR was just created, initializing status fields
if cephFilesystemSubVolumeGroup.Status == nil {
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionProgressing)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, cephv1.ConditionProgressing)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -239,7 +243,7 @@ func (r *ReconcileCephFilesystemSubVolumeGroup) reconcile(request reconcile.Requ
logger.Info(opcontroller.OperatorNotInitializedMessage)
return opcontroller.WaitForRequeueIfOperatorNotInitialized, nil
}
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, cephv1.ConditionFailure)
return reconcile.Result{}, errors.Wrapf(err, "failed to create or update ceph filesystem subvolume group %q", cephFilesystemSubVolumeGroup.Name)
}
}
@@ -258,11 +262,12 @@ func (r *ReconcileCephFilesystemSubVolumeGroup) reconcile(request reconcile.Requ
return reconcile.Result{}, errors.Wrap(err, "failed to save cluster config")
}
// update ObservedGeneration in status at te end of reconcile
// Success! Let's update the status
if cephCluster.Spec.External.Enable {
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionConnected)
r.updateStatus(observedGeneration, request.NamespacedName, cephv1.ConditionConnected)
} else {
r.updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady)
r.updateStatus(observedGeneration, request.NamespacedName, cephv1.ConditionReady)
}
// Return and do not requeue
@@ -306,9 +311,9 @@ func (r *ReconcileCephFilesystemSubVolumeGroup) deleteSubVolumeGroup(cephFilesys
}
// updateStatus updates an object with a given status
func (r *ReconcileCephFilesystemSubVolumeGroup) updateStatus(client client.Client, name types.NamespacedName, status cephv1.ConditionType) {
func (r *ReconcileCephFilesystemSubVolumeGroup) updateStatus(observedGeneration int64, name types.NamespacedName, status cephv1.ConditionType) {
cephFilesystemSubVolumeGroup := &cephv1.CephFilesystemSubVolumeGroup{}
if err := client.Get(r.opManagerContext, name, cephFilesystemSubVolumeGroup); err != nil {
if err := r.client.Get(r.opManagerContext, name, cephFilesystemSubVolumeGroup); err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephFilesystemSubVolumeGroup resource not found. Ignoring since object must be deleted.")
return
@@ -322,11 +327,14 @@ func (r *ReconcileCephFilesystemSubVolumeGroup) updateStatus(client client.Clien
cephFilesystemSubVolumeGroup.Status.Phase = status
cephFilesystemSubVolumeGroup.Status.Info = map[string]string{"clusterID": buildClusterID(cephFilesystemSubVolumeGroup)}
if err := reporting.UpdateStatus(client, cephFilesystemSubVolumeGroup); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
cephFilesystemSubVolumeGroup.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, cephFilesystemSubVolumeGroup); err != nil {
logger.Errorf("failed to set ceph filesystem subvolume group %q status to %q. %v", name, status, err)
return
}
logger.Debugf("ceph ceph filesystem subvolume group %q status updated to %q", name, status)
logger.Debugf("ceph filesystem subvolume group %q status updated to %q", name, status)
}
func buildClusterID(cephFilesystemSubVolumeGroup *cephv1.CephFilesystemSubVolumeGroup) string {
+12 -4
View File
@@ -155,6 +155,10 @@ func (r *ReconcileCephNFS) reconcile(request reconcile.Request) (reconcile.Resul
// Error reading the object - requeue the request.
return reconcile.Result{}, *cephNFS, errors.Wrap(err, "failed to get cephNFS")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephNFS.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephNFS)
@@ -164,7 +168,7 @@ func (r *ReconcileCephNFS) reconcile(request reconcile.Request) (reconcile.Resul
// The CR was just created, initializing status fields
if cephNFS.Status == nil {
updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus)
updateStatus(k8sutil.ObservedGenerationNotAvailable, r.client, request.NamespacedName, k8sutil.EmptyStatus)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -292,12 +296,13 @@ func (r *ReconcileCephNFS) reconcile(request reconcile.Request) (reconcile.Resul
logger.Debug("reconciling ceph nfs deployments")
_, err = r.reconcileCreateCephNFS(cephNFS)
if err != nil {
updateStatus(r.client, request.NamespacedName, k8sutil.FailedStatus)
updateStatus(k8sutil.ObservedGenerationNotAvailable, r.client, request.NamespacedName, k8sutil.FailedStatus)
return reconcile.Result{}, *cephNFS, errors.Wrap(err, "failed to create ceph nfs deployments")
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
updateStatus(r.client, request.NamespacedName, k8sutil.ReadyStatus)
updateStatus(observedGeneration, r.client, request.NamespacedName, k8sutil.ReadyStatus)
// Return and do not requeue
logger.Debug("done reconciling ceph nfs")
@@ -348,7 +353,7 @@ func (r *ReconcileCephNFS) reconcileCreateCephNFS(cephNFS *cephv1.CephNFS) (reco
}
// updateStatus updates an object with a given status
func updateStatus(client client.Client, name types.NamespacedName, status string) {
func updateStatus(observedGeneration int64, client client.Client, name types.NamespacedName, status string) {
nfs := &cephv1.CephNFS{}
err := client.Get(context.TODO(), name, nfs)
if err != nil {
@@ -364,6 +369,9 @@ func updateStatus(client client.Client, name types.NamespacedName, status string
}
nfs.Status.Phase = status
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
nfs.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(client, nfs); err != nil {
logger.Errorf("failed to set nfs %q status to %q. %v", nfs.Name, status, err)
}
+16 -11
View File
@@ -178,6 +178,10 @@ func (r *ReconcileCephObjectStore) reconcile(request reconcile.Request) (reconci
// Error reading the object - requeue the request.
return reconcile.Result{}, *cephObjectStore, errors.Wrap(err, "failed to get cephObjectStore")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephObjectStore.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephObjectStore)
@@ -188,7 +192,7 @@ func (r *ReconcileCephObjectStore) reconcile(request reconcile.Request) (reconci
// The CR was just created, initializing status fields
if cephObjectStore.Status == nil {
// The store is not available so let's not build the status Info yet
updateStatus(r.client, request.NamespacedName, cephv1.ConditionProgressing, map[string]string{})
updateStatus(k8sutil.ObservedGenerationNotAvailable, r.client, request.NamespacedName, cephv1.ConditionProgressing, map[string]string{})
}
// Make sure a CephCluster is present otherwise do nothing
@@ -237,7 +241,7 @@ func (r *ReconcileCephObjectStore) reconcile(request reconcile.Request) (reconci
// DELETE: the CR was deleted
if !cephObjectStore.GetDeletionTimestamp().IsZero() {
updateStatus(r.client, request.NamespacedName, cephv1.ConditionDeleting, buildStatusInfo(cephObjectStore))
updateStatus(k8sutil.ObservedGenerationNotAvailable, r.client, request.NamespacedName, cephv1.ConditionDeleting, buildStatusInfo(cephObjectStore))
// Detect running Ceph version
runningCephVersion, err := cephclient.LeastUptodateDaemonVersion(r.context, r.clusterInfo, config.MonType)
@@ -341,12 +345,13 @@ func (r *ReconcileCephObjectStore) reconcile(request reconcile.Request) (reconci
logger.Info(opcontroller.OperatorNotInitializedMessage)
return opcontroller.WaitForRequeueIfOperatorNotInitialized, *cephObjectStore, nil
} else if err != nil {
result, err := r.setFailedStatus(request.NamespacedName, "failed to create object store deployments", err)
result, err := r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, "failed to create object store deployments", err)
return result, *cephObjectStore, err
}
// update ObservedGeneration in status at the end of reconcile
// Set Progressing status, we are done reconciling, the health check go routine will update the status
updateStatus(r.client, request.NamespacedName, cephv1.ConditionProgressing, buildStatusInfo(cephObjectStore))
updateStatus(observedGeneration, r.client, request.NamespacedName, cephv1.ConditionProgressing, buildStatusInfo(cephObjectStore))
// Return and do not requeue
logger.Debug("done reconciling")
@@ -378,7 +383,7 @@ func (r *ReconcileCephObjectStore) reconcileCreateObjectStore(cephObjectStore *c
logger.Info("reconciling object store service")
_, err = cfg.reconcileService(cephObjectStore)
if err != nil {
return r.setFailedStatus(namespacedName, "failed to reconcile service", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, namespacedName, "failed to reconcile service", err)
}
// RECONCILE ENDPOINTS
@@ -386,11 +391,11 @@ func (r *ReconcileCephObjectStore) reconcileCreateObjectStore(cephObjectStore *c
logger.Info("reconciling external object store endpoint")
err = cfg.reconcileExternalEndpoint(cephObjectStore)
if err != nil {
return r.setFailedStatus(namespacedName, "failed to reconcile external endpoint", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, namespacedName, "failed to reconcile external endpoint", err)
}
if err := UpdateEndpoint(objContext, &cephObjectStore.Spec); err != nil {
return r.setFailedStatus(namespacedName, "failed to set endpoint", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, namespacedName, "failed to set endpoint", err)
}
} else {
logger.Info("reconciling object store deployments")
@@ -418,11 +423,11 @@ func (r *ReconcileCephObjectStore) reconcileCreateObjectStore(cephObjectStore *c
logger.Debug("reconciling object store service")
serviceIP, err := cfg.reconcileService(cephObjectStore)
if err != nil {
return r.setFailedStatus(namespacedName, "failed to reconcile service", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, namespacedName, "failed to reconcile service", err)
}
if err := UpdateEndpoint(objContext, &cephObjectStore.Spec); err != nil {
return r.setFailedStatus(namespacedName, "failed to set endpoint", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, namespacedName, "failed to set endpoint", err)
}
// Reconcile Pool Creation
@@ -430,7 +435,7 @@ func (r *ReconcileCephObjectStore) reconcileCreateObjectStore(cephObjectStore *c
logger.Info("reconciling object store pools")
err = CreatePools(objContext, r.clusterSpec, cephObjectStore.Spec.MetadataPool, cephObjectStore.Spec.DataPool)
if err != nil {
return r.setFailedStatus(namespacedName, "failed to create object pools", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, namespacedName, "failed to create object pools", err)
}
}
@@ -440,7 +445,7 @@ func (r *ReconcileCephObjectStore) reconcileCreateObjectStore(cephObjectStore *c
if err != nil && kerrors.IsNotFound(err) {
return reconcile.Result{}, err
} else if err != nil {
return r.setFailedStatus(namespacedName, "failed to configure multisite for object store", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, namespacedName, "failed to configure multisite for object store", err)
}
// Create or Update Store
+19 -11
View File
@@ -138,10 +138,14 @@ func (r *ReconcileObjectRealm) reconcile(request reconcile.Request) (reconcile.R
// Error reading the object - requeue the request.
return reconcile.Result{}, *cephObjectRealm, errors.Wrap(err, "failed to get CephObjectRealm")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephObjectRealm.ObjectMeta.Generation
// The CR was just created, initializing status fields
if cephObjectRealm.Status == nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -173,12 +177,12 @@ func (r *ReconcileObjectRealm) reconcile(request reconcile.Request) (reconcile.R
// validate the realm settings
err = validateRealmCR(cephObjectRealm)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcileFailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus)
return reconcile.Result{}, *cephObjectRealm, errors.Wrapf(err, "invalid CephObjectRealm CR %q", cephObjectRealm.Name)
}
// Start object reconciliation, updating status for this
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcilingStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcilingStatus)
// Create/Pull Ceph Realm
if cephObjectRealm.Spec.IsPullRealm() {
@@ -190,17 +194,18 @@ func (r *ReconcileObjectRealm) reconcile(request reconcile.Request) (reconcile.R
} else {
_, err = r.createRealmKeys(cephObjectRealm)
if err != nil {
return r.setFailedStatus(cephObjectRealm, request.NamespacedName, "failed to create keys for realm", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, cephObjectRealm, request.NamespacedName, "failed to create keys for realm", err)
}
_, err = r.createCephRealm(cephObjectRealm)
if err != nil {
return r.setFailedStatus(cephObjectRealm, request.NamespacedName, "failed to create ceph realm", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, cephObjectRealm, request.NamespacedName, "failed to create ceph realm", err)
}
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReadyStatus)
r.updateStatus(observedGeneration, request.NamespacedName, k8sutil.ReadyStatus)
// Return and do not requeue
logger.Debug("realm %q done reconciling", request.NamespacedName)
@@ -322,15 +327,15 @@ func validateRealmCR(u *cephv1.CephObjectRealm) error {
return nil
}
func (r *ReconcileObjectRealm) setFailedStatus(cephObjectRealm *cephv1.CephObjectRealm, name types.NamespacedName, errMessage string, err error) (reconcile.Result, cephv1.CephObjectRealm, error) {
r.updateStatus(r.client, name, k8sutil.ReconcileFailedStatus)
func (r *ReconcileObjectRealm) setFailedStatus(observedGeneration int64, cephObjectRealm *cephv1.CephObjectRealm, name types.NamespacedName, errMessage string, err error) (reconcile.Result, cephv1.CephObjectRealm, error) {
r.updateStatus(observedGeneration, name, k8sutil.ReconcileFailedStatus)
return reconcile.Result{}, *cephObjectRealm, errors.Wrapf(err, "%s", errMessage)
}
// updateStatus updates an realm with a given status
func (r *ReconcileObjectRealm) updateStatus(client client.Client, name types.NamespacedName, status string) {
func (r *ReconcileObjectRealm) updateStatus(observedGeneration int64, name types.NamespacedName, status string) {
objectRealm := &cephv1.CephObjectRealm{}
if err := client.Get(r.opManagerContext, name, objectRealm); err != nil {
if err := r.client.Get(r.opManagerContext, name, objectRealm); err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephObjectRealm %q resource not found. Ignoring since object must be deleted", name)
return
@@ -343,7 +348,10 @@ func (r *ReconcileObjectRealm) updateStatus(client client.Client, name types.Nam
}
objectRealm.Status.Phase = status
if err := reporting.UpdateStatus(client, objectRealm); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
objectRealm.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, objectRealm); err != nil {
logger.Errorf("failed to set object realm %q status to %q. %v", name, status, err)
return
}
+7 -3
View File
@@ -22,6 +22,7 @@ import (
"github.com/pkg/errors"
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
"github.com/rook/rook/pkg/operator/ceph/reporting"
"github.com/rook/rook/pkg/operator/k8sutil"
kerrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
@@ -29,13 +30,13 @@ import (
"sigs.k8s.io/controller-runtime/pkg/reconcile"
)
func (r *ReconcileCephObjectStore) setFailedStatus(name types.NamespacedName, errMessage string, err error) (reconcile.Result, error) {
updateStatus(r.client, name, cephv1.ConditionFailure, map[string]string{})
func (r *ReconcileCephObjectStore) setFailedStatus(observedGeneration int64, name types.NamespacedName, errMessage string, err error) (reconcile.Result, error) {
updateStatus(observedGeneration, r.client, name, cephv1.ConditionFailure, map[string]string{})
return reconcile.Result{}, errors.Wrapf(err, "%s", errMessage)
}
// updateStatus updates an object with a given status
func updateStatus(client client.Client, namespacedName types.NamespacedName, status cephv1.ConditionType, info map[string]string) {
func updateStatus(observedGeneration int64, client client.Client, namespacedName types.NamespacedName, status cephv1.ConditionType, info map[string]string) {
// Updating the status is important to users, but we can still keep operating if there is a
// failure. Retry a few times to give it our best effort attempt.
err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
@@ -58,6 +59,9 @@ func updateStatus(client client.Client, namespacedName types.NamespacedName, sta
objectStore.Status.Phase = status
objectStore.Status.Info = info
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
objectStore.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(client, objectStore); err != nil {
return errors.Wrapf(err, "failed to set object store %q status to %q", namespacedName.String(), status)
+16 -8
View File
@@ -108,6 +108,10 @@ func (r *ReconcileBucketTopic) reconcile(request reconcile.Request) (reconcile.R
// Error reading the object - requeue the request.
return reconcile.Result{}, errors.Wrapf(err, "failed to get CephBucketTopic %q", request.NamespacedName)
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephBucketTopic.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephBucketTopic)
@@ -117,7 +121,7 @@ func (r *ReconcileBucketTopic) reconcile(request reconcile.Request) (reconcile.R
// The CR was just created, initializing status fields
if cephBucketTopic.Status == nil {
r.updateStatus(request.NamespacedName, k8sutil.EmptyStatus, nil)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus, nil)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -169,21 +173,22 @@ func (r *ReconcileBucketTopic) reconcile(request reconcile.Request) (reconcile.R
// validate the topic settings
err = cephBucketTopic.ValidateCreate()
if err != nil {
r.updateStatus(request.NamespacedName, k8sutil.ReconcileFailedStatus, nil)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus, nil)
return reconcile.Result{}, errors.Wrapf(err, "invalid CephBucketTopic %q", request.NamespacedName)
}
// Start object reconciliation, updating status for this
r.updateStatus(request.NamespacedName, k8sutil.ReconcilingStatus, nil)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcilingStatus, nil)
// create topic
topicARN, err := r.createCephBucketTopic(cephBucketTopic)
if err != nil {
return r.setFailedStatus(request.NamespacedName, "failed to create topic for bucket notifications", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, "failed to create topic for bucket notifications", err)
}
// update ObservedGeneration in status a the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(request.NamespacedName, k8sutil.ReadyStatus, topicARN)
r.updateStatus(observedGeneration, request.NamespacedName, k8sutil.ReadyStatus, topicARN)
// Return and do not requeue
return reconcile.Result{}, nil
@@ -216,13 +221,13 @@ func (r *ReconcileBucketTopic) deleteCephBucketTopic(topic *cephv1.CephBucketTop
)
}
func (r *ReconcileBucketTopic) setFailedStatus(name types.NamespacedName, errMessage string, err error) (reconcile.Result, error) {
r.updateStatus(name, k8sutil.ReconcileFailedStatus, nil)
func (r *ReconcileBucketTopic) setFailedStatus(observedGeneration int64, name types.NamespacedName, errMessage string, err error) (reconcile.Result, error) {
r.updateStatus(observedGeneration, name, k8sutil.ReconcileFailedStatus, nil)
return reconcile.Result{}, errors.Wrapf(err, "%s", errMessage)
}
// updateStatus updates the topic with a given status
func (r *ReconcileBucketTopic) updateStatus(nsName types.NamespacedName, status string, topicARN *string) {
func (r *ReconcileBucketTopic) updateStatus(observedGeneration int64, nsName types.NamespacedName, status string, topicARN *string) {
topic := &cephv1.CephBucketTopic{}
if err := r.client.Get(r.opManagerContext, nsName, topic); err != nil {
if kerrors.IsNotFound(err) {
@@ -238,6 +243,9 @@ func (r *ReconcileBucketTopic) updateStatus(nsName types.NamespacedName, status
topic.Status.ARN = topicARN
topic.Status.Phase = status
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
topic.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, topic); err != nil {
logger.Errorf("failed to set CephBucketTopic %q status to %q. error %v", nsName, status, err)
return
+17 -9
View File
@@ -148,6 +148,10 @@ func (r *ReconcileObjectStoreUser) reconcile(request reconcile.Request) (reconci
// Error reading the object - requeue the request.
return reconcile.Result{}, *cephObjectStoreUser, errors.Wrap(err, "failed to get CephObjectStoreUser")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephObjectStoreUser.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephObjectStoreUser)
@@ -157,7 +161,7 @@ func (r *ReconcileObjectStoreUser) reconcile(request reconcile.Request) (reconci
// The CR was just created, initializing status fields
if cephObjectStoreUser.Status == nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -205,7 +209,7 @@ func (r *ReconcileObjectStoreUser) reconcile(request reconcile.Request) (reconci
}
logger.Debugf("ObjectStore resource not ready in namespace %q, retrying in %q. %v",
request.NamespacedName.Namespace, opcontroller.WaitForRequeueIfCephClusterNotReady.RequeueAfter.String(), err)
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcileFailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus)
return opcontroller.WaitForRequeueIfCephClusterNotReady, *cephObjectStoreUser, nil
}
@@ -237,26 +241,27 @@ func (r *ReconcileObjectStoreUser) reconcile(request reconcile.Request) (reconci
// validate the user settings
err = r.validateUser(cephObjectStoreUser)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcileFailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus)
return reconcile.Result{}, *cephObjectStoreUser, errors.Wrapf(err, "invalid pool CR %q spec", cephObjectStoreUser.Name)
}
// CREATE/UPDATE CEPH USER
reconcileResponse, err = r.reconcileCephUser(cephObjectStoreUser)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcileFailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus)
return reconcileResponse, *cephObjectStoreUser, err
}
// CREATE/UPDATE KUBERNETES SECRET
reconcileResponse, err = r.reconcileCephUserSecret(cephObjectStoreUser)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcileFailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus)
return reconcileResponse, *cephObjectStoreUser, err
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReadyStatus)
r.updateStatus(observedGeneration, request.NamespacedName, k8sutil.ReadyStatus)
// Return and do not requeue
logger.Debug("done reconciling")
@@ -578,9 +583,9 @@ func labelsForRgw(name string) map[string]string {
}
// updateStatus updates an object with a given status
func (r *ReconcileObjectStoreUser) updateStatus(client client.Client, name types.NamespacedName, status string) {
func (r *ReconcileObjectStoreUser) updateStatus(observedGeneration int64, name types.NamespacedName, status string) {
user := &cephv1.CephObjectStoreUser{}
if err := client.Get(r.opManagerContext, name, user); err != nil {
if err := r.client.Get(r.opManagerContext, name, user); err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephObjectStoreUser resource not found. Ignoring since object must be deleted.")
return
@@ -596,7 +601,10 @@ func (r *ReconcileObjectStoreUser) updateStatus(client client.Client, name types
if user.Status.Phase == k8sutil.ReadyStatus {
user.Status.Info = generateStatusInfo(user)
}
if err := reporting.UpdateStatus(client, user); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
user.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, user); err != nil {
logger.Errorf("failed to set object store user %q status to %q. %v", name, status, err)
return
}
+18 -10
View File
@@ -135,10 +135,14 @@ func (r *ReconcileObjectZone) reconcile(request reconcile.Request) (reconcile.Re
// Error reading the object - requeue the request.
return reconcile.Result{}, *cephObjectZone, errors.Wrap(err, "failed to get CephObjectZone")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephObjectZone.ObjectMeta.Generation
// The CR was just created, initializing status fields
if cephObjectZone.Status == nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -172,12 +176,12 @@ func (r *ReconcileObjectZone) reconcile(request reconcile.Request) (reconcile.Re
// validate the zone settings
err = r.validateZoneCR(cephObjectZone)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcileFailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus)
return reconcile.Result{}, *cephObjectZone, errors.Wrapf(err, "invalid CephObjectZone CR %q", cephObjectZone.Name)
}
// Start object reconciliation, updating status for this
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcilingStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcilingStatus)
// Make sure an ObjectZoneGroup is present
realmName, reconcileResponse, err := r.reconcileObjectZoneGroup(cephObjectZone)
@@ -194,11 +198,12 @@ func (r *ReconcileObjectZone) reconcile(request reconcile.Request) (reconcile.Re
// Create Ceph Zone
_, err = r.createCephZone(cephObjectZone, realmName)
if err != nil {
return r.setFailedStatus(cephObjectZone, request.NamespacedName, "failed to create ceph zone", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, cephObjectZone, request.NamespacedName, "failed to create ceph zone", err)
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReadyStatus)
r.updateStatus(observedGeneration, request.NamespacedName, k8sutil.ReadyStatus)
// Return and do not requeue
logger.Debug("zone done reconciling")
@@ -341,15 +346,15 @@ func (r *ReconcileObjectZone) validateZoneCR(z *cephv1.CephObjectZone) error {
return nil
}
func (r *ReconcileObjectZone) setFailedStatus(cephObjectZone *cephv1.CephObjectZone, name types.NamespacedName, errMessage string, err error) (reconcile.Result, cephv1.CephObjectZone, error) {
r.updateStatus(r.client, name, k8sutil.ReconcileFailedStatus)
func (r *ReconcileObjectZone) setFailedStatus(observedGeneration int64, cephObjectZone *cephv1.CephObjectZone, name types.NamespacedName, errMessage string, err error) (reconcile.Result, cephv1.CephObjectZone, error) {
r.updateStatus(observedGeneration, name, k8sutil.ReconcileFailedStatus)
return reconcile.Result{}, *cephObjectZone, errors.Wrapf(err, "%s", errMessage)
}
// updateStatus updates an zone with a given status
func (r *ReconcileObjectZone) updateStatus(client client.Client, name types.NamespacedName, status string) {
func (r *ReconcileObjectZone) updateStatus(observedGeneration int64, name types.NamespacedName, status string) {
objectZone := &cephv1.CephObjectZone{}
if err := client.Get(r.opManagerContext, name, objectZone); err != nil {
if err := r.client.Get(r.opManagerContext, name, objectZone); err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephObjectZone resource not found. Ignoring since object must be deleted.")
return
@@ -362,7 +367,10 @@ func (r *ReconcileObjectZone) updateStatus(client client.Client, name types.Name
}
objectZone.Status.Phase = status
if err := reporting.UpdateStatus(client, objectZone); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
objectZone.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, objectZone); err != nil {
logger.Errorf("failed to set object zone %q status to %q. %v", name, status, err)
return
}
@@ -141,8 +141,9 @@ func TestCephObjectZoneController(t *testing.T) {
dataPool := cephv1.PoolSpec{}
objectZone := &cephv1.CephObjectZone{
ObjectMeta: metav1.ObjectMeta{
Name: name,
Namespace: namespace,
Name: name,
Namespace: namespace,
Generation: 0,
},
TypeMeta: metav1.TypeMeta{
Kind: "CephObjectZone",
@@ -132,10 +132,14 @@ func (r *ReconcileObjectZoneGroup) reconcile(request reconcile.Request) (reconci
// Error reading the object - requeue the request.
return reconcile.Result{}, errors.Wrap(err, "failed to get CephObjectZoneGroup")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephObjectZoneGroup.ObjectMeta.Generation
// The CR was just created, initializing status fields
if cephObjectZoneGroup.Status == nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.EmptyStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.EmptyStatus)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -166,12 +170,12 @@ func (r *ReconcileObjectZoneGroup) reconcile(request reconcile.Request) (reconci
// validate the zone group settings
err = validateZoneGroup(cephObjectZoneGroup)
if err != nil {
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcileFailedStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcileFailedStatus)
return reconcile.Result{}, errors.Wrapf(err, "invalid CephObjectZoneGroup CR %q", cephObjectZoneGroup.Name)
}
// Start object reconciliation, updating status for this
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReconcilingStatus)
r.updateStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, k8sutil.ReconcilingStatus)
// Make sure an ObjectRealm Resource is present
reconcileResponse, err = r.reconcileObjectRealm(cephObjectZoneGroup)
@@ -188,11 +192,12 @@ func (r *ReconcileObjectZoneGroup) reconcile(request reconcile.Request) (reconci
// Create/Update Ceph Zone Group
_, err = r.createCephZoneGroup(cephObjectZoneGroup)
if err != nil {
return r.setFailedStatus(request.NamespacedName, "failed to create ceph zone group", err)
return r.setFailedStatus(k8sutil.ObservedGenerationNotAvailable, request.NamespacedName, "failed to create ceph zone group", err)
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
r.updateStatus(r.client, request.NamespacedName, k8sutil.ReadyStatus)
r.updateStatus(observedGeneration, request.NamespacedName, k8sutil.ReadyStatus)
// Return and do not requeue
logger.Debug("zone group done reconciling")
@@ -290,15 +295,15 @@ func (r *ReconcileObjectZoneGroup) reconcileCephRealm(zoneGroup *cephv1.CephObje
return reconcile.Result{}, nil
}
func (r *ReconcileObjectZoneGroup) setFailedStatus(name types.NamespacedName, errMessage string, err error) (reconcile.Result, error) {
r.updateStatus(r.client, name, k8sutil.ReconcileFailedStatus)
func (r *ReconcileObjectZoneGroup) setFailedStatus(observedGeneration int64, name types.NamespacedName, errMessage string, err error) (reconcile.Result, error) {
r.updateStatus(observedGeneration, name, k8sutil.ReconcileFailedStatus)
return reconcile.Result{}, errors.Wrapf(err, "%s", errMessage)
}
// updateStatus updates an zone group with a given status
func (r *ReconcileObjectZoneGroup) updateStatus(client client.Client, name types.NamespacedName, status string) {
func (r *ReconcileObjectZoneGroup) updateStatus(observedGeneration int64, name types.NamespacedName, status string) {
objectZoneGroup := &cephv1.CephObjectZoneGroup{}
if err := client.Get(r.opManagerContext, name, objectZoneGroup); err != nil {
if err := r.client.Get(r.opManagerContext, name, objectZoneGroup); err != nil {
if kerrors.IsNotFound(err) {
logger.Debug("CephObjectZoneGroup resource not found. Ignoring since object must be deleted.")
return
@@ -311,7 +316,10 @@ func (r *ReconcileObjectZoneGroup) updateStatus(client client.Client, name types
}
objectZoneGroup.Status.Phase = status
if err := reporting.UpdateStatus(client, objectZoneGroup); err != nil {
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
objectZoneGroup.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(r.client, objectZoneGroup); err != nil {
logger.Errorf("failed to set object zone group %q status to %q. %v", name, status, err)
return
}
+11 -5
View File
@@ -160,6 +160,10 @@ func (r *ReconcileCephBlockPool) reconcile(request reconcile.Request) (reconcile
// Error reading the object - requeue the request.
return opcontroller.ImmediateRetryResult, *cephBlockPool, errors.Wrap(err, "failed to get CephBlockPool")
}
// update observedGeneration local variable with current generation value,
// because generation can be changed before reconile got completed
// CR status will be updated at end of reconcile, so to reflect the reconcile has finished
observedGeneration := cephBlockPool.ObjectMeta.Generation
// Set a finalizer so we can do cleanup before the object goes away
err = opcontroller.AddFinalizerIfNotPresent(r.opManagerContext, r.client, cephBlockPool)
@@ -169,7 +173,7 @@ func (r *ReconcileCephBlockPool) reconcile(request reconcile.Request) (reconcile
// The CR was just created, initializing status fields
if cephBlockPool.Status == nil {
updateStatus(r.client, request.NamespacedName, cephv1.ConditionProgressing, nil)
updateStatus(r.client, request.NamespacedName, cephv1.ConditionProgressing, nil, k8sutil.ObservedGenerationNotAvailable)
}
// Make sure a CephCluster is present otherwise do nothing
@@ -277,7 +281,7 @@ func (r *ReconcileCephBlockPool) reconcile(request reconcile.Request) (reconcile
logger.Info(opcontroller.OperatorNotInitializedMessage)
return opcontroller.WaitForRequeueIfOperatorNotInitialized, *cephBlockPool, nil
}
updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure, nil)
updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure, nil, k8sutil.ObservedGenerationNotAvailable)
return reconcileResponse, *cephBlockPool, errors.Wrapf(err, "failed to create pool %q.", cephBlockPool.GetName())
}
@@ -294,7 +298,7 @@ func (r *ReconcileCephBlockPool) reconcile(request reconcile.Request) (reconcile
// Always create a bootstrap peer token in case another cluster wants to add us as a peer
reconcileResponse, err = opcontroller.CreateBootstrapPeerSecret(r.context, clusterInfo, cephBlockPool, k8sutil.NewOwnerInfo(cephBlockPool, r.scheme))
if err != nil {
updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure, nil)
updateStatus(r.client, request.NamespacedName, cephv1.ConditionFailure, nil, k8sutil.ObservedGenerationNotAvailable)
return reconcileResponse, *cephBlockPool, errors.Wrapf(err, "failed to create rbd-mirror bootstrap peer for pool %q.", cephBlockPool.GetName())
}
@@ -324,13 +328,15 @@ func (r *ReconcileCephBlockPool) reconcile(request reconcile.Request) (reconcile
return reconcileResponse, *cephBlockPool, errors.Wrapf(err, "failed to update pool ID mapping config for the pool %q", cephBlockPool.Name)
}
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady, opcontroller.GenerateStatusInfo(cephBlockPool))
updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady, opcontroller.GenerateStatusInfo(cephBlockPool), observedGeneration)
// If not mirrored there is no Status Info field to fulfil
} else {
// update ObservedGeneration in status at the end of reconcile
// Set Ready status, we are done reconciling
updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady, nil)
updateStatus(r.client, request.NamespacedName, cephv1.ConditionReady, nil, observedGeneration)
// Stop monitoring the mirroring status of this pool
if blockPoolContextsExists && r.blockPoolContexts[blockPoolChannelKey].started {
+5 -1
View File
@@ -23,13 +23,14 @@ import (
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
"github.com/rook/rook/pkg/operator/ceph/reporting"
"github.com/rook/rook/pkg/operator/k8sutil"
kerrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
)
// updateStatus updates a pool CR with the given status
func updateStatus(client client.Client, poolName types.NamespacedName, status cephv1.ConditionType, info map[string]string) {
func updateStatus(client client.Client, poolName types.NamespacedName, status cephv1.ConditionType, info map[string]string, observedGeneration int64) {
pool := &cephv1.CephBlockPool{}
err := client.Get(context.TODO(), poolName, pool)
if err != nil {
@@ -47,6 +48,9 @@ func updateStatus(client client.Client, poolName types.NamespacedName, status ce
pool.Status.Phase = status
pool.Status.Info = info
if observedGeneration != k8sutil.ObservedGenerationNotAvailable {
pool.Status.ObservedGeneration = observedGeneration
}
if err := reporting.UpdateStatus(client, pool); err != nil {
logger.Warningf("failed to set pool %q status to %q. %v", pool.Name, status, err)
return
+4 -3
View File
@@ -49,9 +49,10 @@ const (
// ConfigOverrideName config override name
ConfigOverrideName = "rook-config-override"
// ConfigOverrideVal config override value
ConfigOverrideVal = "config"
configMountDir = "/etc/rook/config"
overrideFilename = "override.conf"
ConfigOverrideVal = "config"
configMountDir = "/etc/rook/config"
overrideFilename = "override.conf"
ObservedGenerationNotAvailable int64 = -1
)
// ConfigOverrideMount is an override mount