From 2dfd64a97c3c9bc8285417307c41cee402dcf168 Mon Sep 17 00:00:00 2001 From: parth-gr Date: Wed, 16 Mar 2022 19:58:35 +0530 Subject: [PATCH] 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 --- .../charts/rook-ceph/templates/resources.yaml | 60 +++++++++++++++++++ deploy/examples/crds.yaml | 60 +++++++++++++++++++ pkg/apis/ceph.rook.io/v1/types.go | 27 +++++++++ pkg/operator/ceph/client/controller.go | 20 +++++-- pkg/operator/ceph/cluster/cephstatus.go | 2 +- pkg/operator/ceph/cluster/cluster.go | 21 ++++--- pkg/operator/ceph/cluster/cluster_external.go | 2 +- pkg/operator/ceph/cluster/controller.go | 4 +- pkg/operator/ceph/cluster/osd/create.go | 4 +- .../ceph/cluster/osd/integration_test.go | 2 +- pkg/operator/ceph/cluster/osd/osd_test.go | 2 +- pkg/operator/ceph/cluster/osd/update.go | 4 +- pkg/operator/ceph/cluster/osd/update_test.go | 2 +- pkg/operator/ceph/cluster/rbd/controller.go | 20 +++++-- pkg/operator/ceph/controller/conditions.go | 11 +++- pkg/operator/ceph/file/controller.go | 17 ++++-- pkg/operator/ceph/file/mirror/controller.go | 22 +++++-- pkg/operator/ceph/file/status.go | 11 ++-- .../ceph/file/subvolumegroup/controller.go | 24 +++++--- pkg/operator/ceph/nfs/controller.go | 16 +++-- pkg/operator/ceph/object/controller.go | 27 +++++---- pkg/operator/ceph/object/realm/controller.go | 30 ++++++---- pkg/operator/ceph/object/status.go | 10 +++- pkg/operator/ceph/object/topic/controller.go | 24 +++++--- pkg/operator/ceph/object/user/controller.go | 26 +++++--- pkg/operator/ceph/object/zone/controller.go | 28 +++++---- .../ceph/object/zone/controller_test.go | 5 +- .../ceph/object/zonegroup/controller.go | 28 +++++---- pkg/operator/ceph/pool/controller.go | 16 +++-- pkg/operator/ceph/pool/status.go | 6 +- pkg/operator/k8sutil/pod.go | 7 ++- 31 files changed, 405 insertions(+), 133 deletions(-) diff --git a/deploy/charts/rook-ceph/templates/resources.yaml b/deploy/charts/rook-ceph/templates/resources.yaml index ab4534973..8576080af 100644 --- a/deploy/charts/rook-ceph/templates/resources.yaml +++ b/deploy/charts/rook-ceph/templates/resources.yaml @@ -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 diff --git a/deploy/examples/crds.yaml b/deploy/examples/crds.yaml index b24f255ef..609fd7f3b 100644 --- a/deploy/examples/crds.yaml +++ b/deploy/examples/crds.yaml @@ -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 diff --git a/pkg/apis/ceph.rook.io/v1/types.go b/pkg/apis/ceph.rook.io/v1/types.go index 60f0358ce..3cfbb6176 100755 --- a/pkg/apis/ceph.rook.io/v1/types.go +++ b/pkg/apis/ceph.rook.io/v1/types.go @@ -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"` } diff --git a/pkg/operator/ceph/client/controller.go b/pkg/operator/ceph/client/controller.go index 6828f0e71..f7238fd1b 100644 --- a/pkg/operator/ceph/client/controller.go +++ b/pkg/operator/ceph/client/controller.go @@ -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 } diff --git a/pkg/operator/ceph/cluster/cephstatus.go b/pkg/operator/ceph/cluster/cephstatus.go index 6cb190121..a7739bee3 100644 --- a/pkg/operator/ceph/cluster/cephstatus.go +++ b/pkg/operator/ceph/cluster/cephstatus.go @@ -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 diff --git a/pkg/operator/ceph/cluster/cluster.go b/pkg/operator/ceph/cluster/cluster.go index 88eb585b7..b689ed43d 100755 --- a/pkg/operator/ceph/cluster/cluster.go +++ b/pkg/operator/ceph/cluster/cluster.go @@ -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 } diff --git a/pkg/operator/ceph/cluster/cluster_external.go b/pkg/operator/ceph/cluster/cluster_external.go index a77a6f184..aea264f51 100644 --- a/pkg/operator/ceph/cluster/cluster_external.go +++ b/pkg/operator/ceph/cluster/cluster_external.go @@ -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 diff --git a/pkg/operator/ceph/cluster/controller.go b/pkg/operator/ceph/cluster/controller.go index 7a7bab4cb..f8a1a5c19 100644 --- a/pkg/operator/ceph/cluster/controller.go +++ b/pkg/operator/ceph/cluster/controller.go @@ -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 diff --git a/pkg/operator/ceph/cluster/osd/create.go b/pkg/operator/ceph/cluster/osd/create.go index cd357e50a..fd3075588 100644 --- a/pkg/operator/ceph/cluster/osd/create.go +++ b/pkg/operator/ceph/cluster/osd/create.go @@ -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) diff --git a/pkg/operator/ceph/cluster/osd/integration_test.go b/pkg/operator/ceph/cluster/osd/integration_test.go index 37a2d91b8..8d6a27cfb 100644 --- a/pkg/operator/ceph/cluster/osd/integration_test.go +++ b/pkg/operator/ceph/cluster/osd/integration_test.go @@ -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 } diff --git a/pkg/operator/ceph/cluster/osd/osd_test.go b/pkg/operator/ceph/cluster/osd/osd_test.go index 1cddf4f35..07b14df94 100644 --- a/pkg/operator/ceph/cluster/osd/osd_test.go +++ b/pkg/operator/ceph/cluster/osd/osd_test.go @@ -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 } diff --git a/pkg/operator/ceph/cluster/osd/update.go b/pkg/operator/ceph/cluster/osd/update.go index 763aaa86a..52c985c5f 100644 --- a/pkg/operator/ceph/cluster/osd/update.go +++ b/pkg/operator/ceph/cluster/osd/update.go @@ -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)) diff --git a/pkg/operator/ceph/cluster/osd/update_test.go b/pkg/operator/ceph/cluster/osd/update_test.go index 9ca86324e..1f76f18d2 100644 --- a/pkg/operator/ceph/cluster/osd/update_test.go +++ b/pkg/operator/ceph/cluster/osd/update_test.go @@ -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 { diff --git a/pkg/operator/ceph/cluster/rbd/controller.go b/pkg/operator/ceph/cluster/rbd/controller.go index 698afbec8..a841e552d 100644 --- a/pkg/operator/ceph/cluster/rbd/controller.go +++ b/pkg/operator/ceph/cluster/rbd/controller.go @@ -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 } diff --git a/pkg/operator/ceph/controller/conditions.go b/pkg/operator/ceph/controller/conditions.go index 64f00d071..327d692ed 100644 --- a/pkg/operator/ceph/controller/conditions.go +++ b/pkg/operator/ceph/controller/conditions.go @@ -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 { diff --git a/pkg/operator/ceph/file/controller.go b/pkg/operator/ceph/file/controller.go index c9c988358..9dd90d6ff 100644 --- a/pkg/operator/ceph/file/controller.go +++ b/pkg/operator/ceph/file/controller.go @@ -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 diff --git a/pkg/operator/ceph/file/mirror/controller.go b/pkg/operator/ceph/file/mirror/controller.go index da6fb865b..c4cf593f5 100644 --- a/pkg/operator/ceph/file/mirror/controller.go +++ b/pkg/operator/ceph/file/mirror/controller.go @@ -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 } diff --git a/pkg/operator/ceph/file/status.go b/pkg/operator/ceph/file/status.go index 952c98a00..975d9749c 100644 --- a/pkg/operator/ceph/file/status.go +++ b/pkg/operator/ceph/file/status.go @@ -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 } diff --git a/pkg/operator/ceph/file/subvolumegroup/controller.go b/pkg/operator/ceph/file/subvolumegroup/controller.go index f5f538ef2..df1055f97 100644 --- a/pkg/operator/ceph/file/subvolumegroup/controller.go +++ b/pkg/operator/ceph/file/subvolumegroup/controller.go @@ -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 { diff --git a/pkg/operator/ceph/nfs/controller.go b/pkg/operator/ceph/nfs/controller.go index f7acf3d18..f90b46c21 100644 --- a/pkg/operator/ceph/nfs/controller.go +++ b/pkg/operator/ceph/nfs/controller.go @@ -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) } diff --git a/pkg/operator/ceph/object/controller.go b/pkg/operator/ceph/object/controller.go index 1857ee90f..d3a9f5838 100644 --- a/pkg/operator/ceph/object/controller.go +++ b/pkg/operator/ceph/object/controller.go @@ -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 diff --git a/pkg/operator/ceph/object/realm/controller.go b/pkg/operator/ceph/object/realm/controller.go index b766a37bb..b221cfec0 100644 --- a/pkg/operator/ceph/object/realm/controller.go +++ b/pkg/operator/ceph/object/realm/controller.go @@ -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 } diff --git a/pkg/operator/ceph/object/status.go b/pkg/operator/ceph/object/status.go index ddf4c5ad4..34a839ff9 100644 --- a/pkg/operator/ceph/object/status.go +++ b/pkg/operator/ceph/object/status.go @@ -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) diff --git a/pkg/operator/ceph/object/topic/controller.go b/pkg/operator/ceph/object/topic/controller.go index 9b27d0829..46d0b8a2d 100644 --- a/pkg/operator/ceph/object/topic/controller.go +++ b/pkg/operator/ceph/object/topic/controller.go @@ -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 diff --git a/pkg/operator/ceph/object/user/controller.go b/pkg/operator/ceph/object/user/controller.go index b93e79712..4ff1f7896 100644 --- a/pkg/operator/ceph/object/user/controller.go +++ b/pkg/operator/ceph/object/user/controller.go @@ -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 } diff --git a/pkg/operator/ceph/object/zone/controller.go b/pkg/operator/ceph/object/zone/controller.go index 7a57fd438..e1b3fde19 100644 --- a/pkg/operator/ceph/object/zone/controller.go +++ b/pkg/operator/ceph/object/zone/controller.go @@ -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 } diff --git a/pkg/operator/ceph/object/zone/controller_test.go b/pkg/operator/ceph/object/zone/controller_test.go index e668c1d42..d0a8814f9 100644 --- a/pkg/operator/ceph/object/zone/controller_test.go +++ b/pkg/operator/ceph/object/zone/controller_test.go @@ -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", diff --git a/pkg/operator/ceph/object/zonegroup/controller.go b/pkg/operator/ceph/object/zonegroup/controller.go index 6d2e22e2c..feeb3bcf4 100644 --- a/pkg/operator/ceph/object/zonegroup/controller.go +++ b/pkg/operator/ceph/object/zonegroup/controller.go @@ -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 } diff --git a/pkg/operator/ceph/pool/controller.go b/pkg/operator/ceph/pool/controller.go index 1dc0b965a..2a6a1ff37 100644 --- a/pkg/operator/ceph/pool/controller.go +++ b/pkg/operator/ceph/pool/controller.go @@ -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 { diff --git a/pkg/operator/ceph/pool/status.go b/pkg/operator/ceph/pool/status.go index ffef8b1db..8a70cfdfd 100644 --- a/pkg/operator/ceph/pool/status.go +++ b/pkg/operator/ceph/pool/status.go @@ -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 diff --git a/pkg/operator/k8sutil/pod.go b/pkg/operator/k8sutil/pod.go index fcb3aaefd..e13301c41 100644 --- a/pkg/operator/k8sutil/pod.go +++ b/pkg/operator/k8sutil/pod.go @@ -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