From d4a3dca50a99d189f5c3a8595036c093f746a302 Mon Sep 17 00:00:00 2001 From: Razin Bouzar Date: Sun, 5 Jul 2026 23:37:08 -0400 Subject: [PATCH] Add observational deployment lifecycle status --- apis/druid/v1alpha1/druid_types.go | 40 +- apis/druid/v1alpha1/zz_generated.deepcopy.go | 24 + chart/crds/druid.apache.org_druids.yaml | 23 +- chart/values.yaml | 1 + config/crd/bases/druid.apache.org_druids.yaml | 23 +- controllers/druid/druid_controller.go | 66 ++- controllers/druid/druid_controller_test.go | 19 + controllers/druid/interface.go | 15 +- controllers/druid/lifecycle.go | 332 ++++++++++++++ controllers/druid/lifecycle_metrics.go | 86 ++++ controllers/druid/lifecycle_metrics_test.go | 167 +++++++ controllers/druid/lifecycle_test.go | 423 ++++++++++++++++++ docs/api_specifications/druid.md | 128 +++++- docs/features.md | 6 + go.mod | 4 +- 15 files changed, 1300 insertions(+), 57 deletions(-) create mode 100644 controllers/druid/lifecycle.go create mode 100644 controllers/druid/lifecycle_metrics.go create mode 100644 controllers/druid/lifecycle_metrics_test.go create mode 100644 controllers/druid/lifecycle_test.go diff --git a/apis/druid/v1alpha1/druid_types.go b/apis/druid/v1alpha1/druid_types.go index 34c68a1d..30194382 100644 --- a/apis/druid/v1alpha1/druid_types.go +++ b/apis/druid/v1alpha1/druid_types.go @@ -570,20 +570,36 @@ type DruidNodeTypeStatus struct { Reason string `json:"reason,omitempty"` } +type DeploymentLifecyclePhase string + +const ( + DeploymentLifecyclePending DeploymentLifecyclePhase = "Pending" + DeploymentLifecycleInProgress DeploymentLifecyclePhase = "InProgress" + DeploymentLifecycleSucceeded DeploymentLifecyclePhase = "Succeeded" +) + +type DeploymentLifecycleStatus struct { + // +kubebuilder:validation:Enum=Pending;InProgress;Succeeded + Phase DeploymentLifecyclePhase `json:"phase,omitempty"` + Reason string `json:"reason,omitempty"` + ObservedGeneration int64 `json:"observedGeneration,omitempty"` + StartedAt *metav1.Time `json:"startedAt,omitempty"` + CompletedAt *metav1.Time `json:"completedAt,omitempty"` +} + // DruidClusterStatus Defines the observed state of Druid. type DruidClusterStatus struct { - // INSERT ADDITIONAL STATUS FIELD - define observed state of cluster - // Important: Run "make" to regenerate code after modifying this file - DruidNodeStatus DruidNodeTypeStatus `json:"druidNodeStatus,omitempty"` - StatefulSets []string `json:"statefulSets,omitempty"` - Deployments []string `json:"deployments,omitempty"` - Services []string `json:"services,omitempty"` - ConfigMaps []string `json:"configMaps,omitempty"` - PodDisruptionBudgets []string `json:"podDisruptionBudgets,omitempty"` - Ingress []string `json:"ingress,omitempty"` - HPAutoScalers []string `json:"hpAutoscalers,omitempty"` - Pods []string `json:"pods,omitempty"` - PersistentVolumeClaims []string `json:"persistentVolumeClaims,omitempty"` + DeploymentLifecycle DeploymentLifecycleStatus `json:"deploymentLifecycle,omitempty"` + DruidNodeStatus DruidNodeTypeStatus `json:"druidNodeStatus,omitempty"` + StatefulSets []string `json:"statefulSets,omitempty"` + Deployments []string `json:"deployments,omitempty"` + Services []string `json:"services,omitempty"` + ConfigMaps []string `json:"configMaps,omitempty"` + PodDisruptionBudgets []string `json:"podDisruptionBudgets,omitempty"` + Ingress []string `json:"ingress,omitempty"` + HPAutoScalers []string `json:"hpAutoscalers,omitempty"` + Pods []string `json:"pods,omitempty"` + PersistentVolumeClaims []string `json:"persistentVolumeClaims,omitempty"` } // Druid is the Schema for the druids API. diff --git a/apis/druid/v1alpha1/zz_generated.deepcopy.go b/apis/druid/v1alpha1/zz_generated.deepcopy.go index 93eff590..693675ef 100644 --- a/apis/druid/v1alpha1/zz_generated.deepcopy.go +++ b/apis/druid/v1alpha1/zz_generated.deepcopy.go @@ -105,6 +105,29 @@ func (in *DeepStorageSpec) DeepCopy() *DeepStorageSpec { return out } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *DeploymentLifecycleStatus) DeepCopyInto(out *DeploymentLifecycleStatus) { + *out = *in + if in.StartedAt != nil { + in, out := &in.StartedAt, &out.StartedAt + *out = (*in).DeepCopy() + } + if in.CompletedAt != nil { + in, out := &in.CompletedAt, &out.CompletedAt + *out = (*in).DeepCopy() + } +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new DeploymentLifecycleStatus. +func (in *DeploymentLifecycleStatus) DeepCopy() *DeploymentLifecycleStatus { + if in == nil { + return nil + } + out := new(DeploymentLifecycleStatus) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *Druid) DeepCopyInto(out *Druid) { *out = *in @@ -135,6 +158,7 @@ func (in *Druid) DeepCopyObject() runtime.Object { // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *DruidClusterStatus) DeepCopyInto(out *DruidClusterStatus) { *out = *in + in.DeploymentLifecycle.DeepCopyInto(&out.DeploymentLifecycle) out.DruidNodeStatus = in.DruidNodeStatus if in.StatefulSets != nil { in, out := &in.StatefulSets, &out.StatefulSets diff --git a/chart/crds/druid.apache.org_druids.yaml b/chart/crds/druid.apache.org_druids.yaml index e2971b36..43f7acc5 100644 --- a/chart/crds/druid.apache.org_druids.yaml +++ b/chart/crds/druid.apache.org_druids.yaml @@ -11851,14 +11851,31 @@ spec: items: type: string type: array + deploymentLifecycle: + properties: + completedAt: + format: date-time + type: string + observedGeneration: + format: int64 + type: integer + phase: + enum: + - Pending + - InProgress + - Succeeded + type: string + reason: + type: string + startedAt: + format: date-time + type: string + type: object deployments: items: type: string type: array druidNodeStatus: - description: |- - INSERT ADDITIONAL STATUS FIELD - define observed state of cluster - Important: Run "make" to regenerate code after modifying this file properties: druidNode: type: string diff --git a/chart/values.yaml b/chart/values.yaml index b5c489a9..6a127b3e 100644 --- a/chart/values.yaml +++ b/chart/values.yaml @@ -26,6 +26,7 @@ global: env: DENY_LIST: "default,kube-system" # Comma-separated list of namespaces to ignore RECONCILE_WAIT: "10s" # Reconciliation delay + DEPLOYMENT_LIFECYCLE_TIMEOUT: "96h" # Observational timeout for long-running deployment lifecycle status WATCH_NAMESPACE: "" # Namespace to watch or empty string to watch all namespaces, To watch multiple namespaces add , into string. Ex: WATCH_NAMESPACE: "ns1,ns2,ns3" #MAX_CONCURRENT_RECONCILES:: "" # MaxConcurrentReconciles is the maximum number of concurrent Reconciles which can be run. diff --git a/config/crd/bases/druid.apache.org_druids.yaml b/config/crd/bases/druid.apache.org_druids.yaml index e2971b36..43f7acc5 100644 --- a/config/crd/bases/druid.apache.org_druids.yaml +++ b/config/crd/bases/druid.apache.org_druids.yaml @@ -11851,14 +11851,31 @@ spec: items: type: string type: array + deploymentLifecycle: + properties: + completedAt: + format: date-time + type: string + observedGeneration: + format: int64 + type: integer + phase: + enum: + - Pending + - InProgress + - Succeeded + type: string + reason: + type: string + startedAt: + format: date-time + type: string + type: object deployments: items: type: string type: array druidNodeStatus: - description: |- - INSERT ADDITIONAL STATUS FIELD - define observed state of cluster - Important: Run "make" to regenerate code after modifying this file properties: druidNode: type: string diff --git a/controllers/druid/druid_controller.go b/controllers/druid/druid_controller.go index b91dc14a..c668baeb 100644 --- a/controllers/druid/druid_controller.go +++ b/controllers/druid/druid_controller.go @@ -23,7 +23,7 @@ import ( "os" "time" - "k8s.io/apimachinery/pkg/api/errors" + k8serrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/client-go/tools/record" "github.com/go-logr/logr" @@ -42,17 +42,19 @@ type DruidReconciler struct { Log logr.Logger Scheme *runtime.Scheme // reconcile time duration, defaults to 10s - ReconcileWait time.Duration - Recorder record.EventRecorder + ReconcileWait time.Duration + DeploymentLifecycleTimeout time.Duration + Recorder record.EventRecorder } func NewDruidReconciler(mgr ctrl.Manager) *DruidReconciler { return &DruidReconciler{ - Client: mgr.GetClient(), - Log: ctrl.Log.WithName("controllers").WithName("Druid"), - Scheme: mgr.GetScheme(), - ReconcileWait: LookupReconcileTime(), - Recorder: mgr.GetEventRecorderFor("druid-operator"), + Client: mgr.GetClient(), + Log: ctrl.Log.WithName("controllers").WithName("Druid"), + Scheme: mgr.GetScheme(), + ReconcileWait: LookupReconcileTime(), + DeploymentLifecycleTimeout: LookupDeploymentLifecycleTimeout(), + Recorder: mgr.GetEventRecorderFor("druid-operator"), } } @@ -70,14 +72,15 @@ func NewDruidReconciler(mgr ctrl.Manager) *DruidReconciler { // +kubebuilder:rbac:groups=networking.k8s.io,resources=ingresses,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups=storage.k8s.io,resources=storageclasses,verbs=get;list;watch -func (r *DruidReconciler) Reconcile(ctx context.Context, request reconcile.Request) (ctrl.Result, error) { +func (r *DruidReconciler) Reconcile(ctx context.Context, request reconcile.Request) (result ctrl.Result, err error) { _ = r.Log.WithValues("druid", request.NamespacedName) // Fetch the Druid instance instance := &druidv1alpha1.Druid{} - err := r.Get(ctx, request.NamespacedName, instance) + err = r.Get(ctx, request.NamespacedName, instance) if err != nil { - if errors.IsNotFound(err) { + if k8serrors.IsNotFound(err) { + defaultDeploymentLifecycleMetrics.delete(request.Namespace, request.Name) // Request object not found, could have been deleted after reconcile request. // Owned objects are automatically garbage collected. For additional cleanup logic use finalizers. // Return and don't requeue @@ -95,6 +98,8 @@ func (r *DruidReconciler) Reconcile(ctx context.Context, request reconcile.Reque return ctrl.Result{}, err } + r.recordDeploymentLifecycle(ctx, instance, emitEvent) + // Update Druid Dynamic Configs if err := updateDruidDynamicConfigs(ctx, r.Client, instance, emitEvent); err != nil { return ctrl.Result{}, err @@ -114,6 +119,31 @@ func (r *DruidReconciler) SetupWithManager(mgr ctrl.Manager) error { Complete(r) } +func (r *DruidReconciler) recordDeploymentLifecycle(ctx context.Context, instance *druidv1alpha1.Druid, emitEvent EventEmitter) { + deps := r.lifecycleDependencies() + snapshot, err := collectManagedResourceSnapshot(ctx, r.Client, instance) + if err != nil { + r.Log.Error(err, "failed to observe deployment lifecycle", "name", instance.Name, "namespace", instance.Namespace) + return + } + if !snapshot.hasWorkloads() { + return + } + if err := reconcileDeploymentLifecycleWithSnapshot(instance, snapshot, deps, func(updated druidv1alpha1.DeploymentLifecycleStatus) error { + return patchDeploymentLifecycleStatus(ctx, r.Client, instance, updated, emitEvent) + }); err != nil { + r.Log.Error(err, "failed to record deployment lifecycle", "name", instance.Name, "namespace", instance.Namespace) + } + defaultDeploymentLifecycleMetrics.record(instance, deps) +} + +func (r *DruidReconciler) lifecycleDependencies() lifecycleDependencies { + return lifecycleDependencies{ + now: time.Now, + timeout: r.DeploymentLifecycleTimeout, + } +} + func LookupReconcileTime() time.Duration { val, exists := os.LookupEnv("RECONCILE_WAIT") if !exists { @@ -129,6 +159,20 @@ func LookupReconcileTime() time.Duration { } } +func LookupDeploymentLifecycleTimeout() time.Duration { + val, exists := os.LookupEnv("DEPLOYMENT_LIFECYCLE_TIMEOUT") + if !exists { + return time.Hour * 96 + } + + v, err := time.ParseDuration(val) + if err != nil { + logger.Error(err, err.Error()) + os.Exit(1) + } + return v +} + func getMaxConcurrentReconciles() int { var MaxConcurrentReconciles = "MAX_CONCURRENT_RECONCILES" diff --git a/controllers/druid/druid_controller_test.go b/controllers/druid/druid_controller_test.go index 457defa7..916e5734 100644 --- a/controllers/druid/druid_controller_test.go +++ b/controllers/druid/druid_controller_test.go @@ -128,6 +128,23 @@ var _ = Describe("Druid Operator", func() { }) + It("records lifecycle status during real reconcile", func() { + lifecycleDruidCR, err := readDruidClusterSpecFromFile(filePath) + Expect(err).Should(BeNil()) + + lifecycleDruidCR.Name = fmt.Sprintf("lifecycle-status-%d", GinkgoRandomSeed()) + Expect(k8sClient.Create(ctx, lifecycleDruidCR)).To(Succeed()) + + lifecycleStatus := &druidv1alpha1.Druid{} + Eventually(func() string { + err := k8sClient.Get(ctx, types.NamespacedName{Name: lifecycleDruidCR.Name, Namespace: lifecycleDruidCR.Namespace}, lifecycleStatus) + if err != nil { + return "" + } + return string(lifecycleStatus.Status.DeploymentLifecycle.Phase) + }, timeout, interval).Should(Equal(string(druidv1alpha1.DeploymentLifecycleInProgress))) + }) + It("Test broker deployment", func() { componentName := "brokers" createdDeploy := &appsv1.Deployment{} @@ -151,6 +168,7 @@ var _ = Describe("Druid Operator", func() { By("By updating broker deployment replicas") replicaCount := 2 + Expect(k8sClient.Get(ctx, types.NamespacedName{Name: druidCR.Name, Namespace: druidCR.Namespace}, druid)).To(Succeed()) if druidRep, ok := druid.Spec.Nodes[componentName]; ok { druidRep.Replicas = int32(replicaCount) druid.Spec.Nodes[componentName] = druidRep @@ -190,6 +208,7 @@ var _ = Describe("Druid Operator", func() { By(fmt.Sprintf("By updating statefulset replicas %s ", stsName)) replicaCount := 2 + Expect(k8sClient.Get(ctx, types.NamespacedName{Name: druidCR.Name, Namespace: druidCR.Namespace}, druid)).To(Succeed()) if druidRep, ok := druid.Spec.Nodes[componentName]; ok { druidRep.Replicas = int32(replicaCount) druid.Spec.Nodes[componentName] = druidRep diff --git a/controllers/druid/interface.go b/controllers/druid/interface.go index 10966320..379611d6 100644 --- a/controllers/druid/interface.go +++ b/controllers/druid/interface.go @@ -61,12 +61,15 @@ const ( druidFinalizerFailed druidEventReason = "DruidFinalizerFailed" druidFinalizerSuccess druidEventReason = "DruidFinalizerSuccess" - druidGetRouterSvcUrlFailed druidEventReason = "DruidAPIGetRouterSvcUrlFailed" - druidGetAuthCredsFailed druidEventReason = "DruidAPIGetAuthCredsFailed" - druidFetchCurrentConfigsFailed druidEventReason = "DruidAPIFetchCurrentConfigsFailed" - druidConfigComparisonFailed druidEventReason = "DruidAPIConfigComparisonFailed" - druidUpdateConfigsFailed druidEventReason = "DruidAPIUpdateConfigsFailed" - druidUpdateConfigsSuccess druidEventReason = "DruidAPIUpdateConfigsSuccess" + druidGetRouterSvcUrlFailed druidEventReason = "DruidAPIGetRouterSvcUrlFailed" + druidGetAuthCredsFailed druidEventReason = "DruidAPIGetAuthCredsFailed" + druidFetchCurrentConfigsFailed druidEventReason = "DruidAPIFetchCurrentConfigsFailed" + druidConfigComparisonFailed druidEventReason = "DruidAPIConfigComparisonFailed" + druidUpdateConfigsFailed druidEventReason = "DruidAPIUpdateConfigsFailed" + druidUpdateConfigsSuccess druidEventReason = "DruidAPIUpdateConfigsSuccess" + druidDeploymentLifecycleStarted druidEventReason = "DruidDeploymentLifecycleStarted" + druidDeploymentLifecycleTimedOut druidEventReason = "DruidDeploymentLifecycleTimedOut" + druidDeploymentLifecycleSucceeded druidEventReason = "DruidDeploymentLifecycleSucceeded" ) // Reader Interface diff --git a/controllers/druid/lifecycle.go b/controllers/druid/lifecycle.go new file mode 100644 index 00000000..d6699e62 --- /dev/null +++ b/controllers/druid/lifecycle.go @@ -0,0 +1,332 @@ +/* +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +*/ +package druid + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "reflect" + "strings" + "time" + + "github.com/apache/druid-operator/apis/druid/v1alpha1" + appsv1 "k8s.io/api/apps/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +const deploymentLifecycleTimeoutReasonPrefix = "Deployment lifecycle exceeded timeout" + +func isTerminalLifecyclePhase(phase v1alpha1.DeploymentLifecyclePhase) bool { + return phase == v1alpha1.DeploymentLifecycleSucceeded +} + +type lifecycleDependencies struct { + now func() time.Time + timeout time.Duration +} + +func defaultLifecycleDependencies() lifecycleDependencies { + return lifecycleDependencies{ + now: time.Now, + timeout: LookupDeploymentLifecycleTimeout(), + } +} + +func reconcileDeploymentLifecycle( + ctx context.Context, + sdk client.Client, + drd *v1alpha1.Druid, + emitEvent EventEmitter, +) error { + return reconcileDeploymentLifecycleWithDeps( + ctx, + sdk, + drd, + emitEvent, + defaultLifecycleDependencies(), + ) +} + +func reconcileDeploymentLifecycleWithDeps( + ctx context.Context, + sdk client.Client, + drd *v1alpha1.Druid, + emitEvent EventEmitter, + deps lifecycleDependencies, +) error { + snapshot, err := collectManagedResourceSnapshot(ctx, sdk, drd) + if err != nil { + return err + } + if !snapshot.hasWorkloads() { + return nil + } + + return reconcileDeploymentLifecycleWithSnapshot(drd, snapshot, deps, func(updated v1alpha1.DeploymentLifecycleStatus) error { + return patchDeploymentLifecycleStatus(ctx, sdk, drd, updated, emitEvent) + }) +} + +func reconcileDeploymentLifecycleWithSnapshot( + drd *v1alpha1.Druid, + snapshot *managedResourceSnapshot, + deps lifecycleDependencies, + patchStatus func(v1alpha1.DeploymentLifecycleStatus) error, +) error { + current := drd.Status.DeploymentLifecycle + if current.ObservedGeneration == drd.Generation && isTerminalLifecyclePhase(current.Phase) { + return nil + } + + updated := buildLifecycleStatusForGeneration(drd, current) + if updated.StartedAt == nil { + now := metav1.NewTime(deps.now()) + updated.StartedAt = &now + } + + workloadsReady, reason, err := areManagedWorkloadsReadyFromSnapshot(snapshot) + if err != nil { + return err + } + if !workloadsReady { + return patchStatus(withDeploymentLifecycleInProgress(updated, timeoutAwareReason(updated, reason, deps))) + } + + if updated.Phase == v1alpha1.DeploymentLifecycleSucceeded { + return nil + } + + updated = withDeploymentLifecycleSucceeded(updated, deps.now()) + return patchStatus(updated) +} + +func withDeploymentLifecycleInProgress( + status v1alpha1.DeploymentLifecycleStatus, + reason string, +) v1alpha1.DeploymentLifecycleStatus { + status.Phase = v1alpha1.DeploymentLifecycleInProgress + status.Reason = reason + status.CompletedAt = nil + return status +} + +func withDeploymentLifecycleSucceeded( + status v1alpha1.DeploymentLifecycleStatus, + now time.Time, +) v1alpha1.DeploymentLifecycleStatus { + completedAt := metav1.NewTime(now) + status.Phase = v1alpha1.DeploymentLifecycleSucceeded + status.Reason = "Deployment lifecycle completed after managed workloads became ready" + status.CompletedAt = &completedAt + return status +} + +func buildLifecycleStatusForGeneration( + drd *v1alpha1.Druid, + current v1alpha1.DeploymentLifecycleStatus, +) v1alpha1.DeploymentLifecycleStatus { + updated := current + updated.ObservedGeneration = drd.Generation + updated.CompletedAt = nil + if current.ObservedGeneration != drd.Generation { + updated.StartedAt = nil + } + return updated +} + +type managedResourceSnapshot struct { + deployments []appsv1.Deployment + statefulSets []appsv1.StatefulSet +} + +func (s *managedResourceSnapshot) hasWorkloads() bool { + return s != nil && (len(s.deployments) > 0 || len(s.statefulSets) > 0) +} + +func collectManagedResourceSnapshot( + ctx context.Context, + sdk client.Client, + drd *v1alpha1.Druid, +) (*managedResourceSnapshot, error) { + listOpts := []client.ListOption{ + client.InNamespace(drd.Namespace), + client.MatchingLabels(makeLabelsForDruid(drd)), + } + + deployments := &appsv1.DeploymentList{} + if err := sdk.List(ctx, deployments, listOpts...); err != nil { + return nil, err + } + + statefulSets := &appsv1.StatefulSetList{} + if err := sdk.List(ctx, statefulSets, listOpts...); err != nil { + return nil, err + } + + return &managedResourceSnapshot{ + deployments: deployments.Items, + statefulSets: statefulSets.Items, + }, nil +} + +func timeoutAwareReason( + status v1alpha1.DeploymentLifecycleStatus, + reason string, + deps lifecycleDependencies, +) string { + if deps.timeout <= 0 || status.StartedAt == nil || deps.now().Sub(status.StartedAt.Time) < deps.timeout { + return reason + } + + return fmt.Sprintf("%s of %s; still %s", deploymentLifecycleTimeoutReasonPrefix, deps.timeout, waitingReasonFragment(reason)) +} + +func waitingReasonFragment(reason string) string { + if reason == "" { + return "waiting for Druid workloads to roll out" + } + if strings.HasPrefix(reason, "Waiting ") { + return "waiting " + strings.TrimPrefix(reason, "Waiting ") + } + return reason +} + +func deploymentLifecycleTimedOut(status v1alpha1.DeploymentLifecycleStatus) bool { + return status.Phase == v1alpha1.DeploymentLifecycleInProgress && + strings.HasPrefix(status.Reason, deploymentLifecycleTimeoutReasonPrefix) +} + +func areManagedWorkloadsReady(ctx context.Context, sdk client.Client, drd *v1alpha1.Druid) (bool, string, error) { + snapshot, err := collectManagedResourceSnapshot(ctx, sdk, drd) + if err != nil { + return false, "", err + } + return areManagedWorkloadsReadyFromSnapshot(snapshot) +} + +func areManagedWorkloadsReadyFromSnapshot(snapshot *managedResourceSnapshot) (bool, string, error) { + for _, deployment := range snapshot.deployments { + for _, condition := range deployment.Status.Conditions { + if condition.Type == appsv1.DeploymentReplicaFailure { + return false, fmt.Sprintf("Waiting for Deployment [%s] replica failure to clear: %s", deployment.Name, condition.Reason), nil + } + } + specReplicas := int32(1) + if deployment.Spec.Replicas != nil { + specReplicas = *deployment.Spec.Replicas + } + if deployment.Status.ObservedGeneration < deployment.Generation { + return false, fmt.Sprintf("Waiting for Deployment [%s] controller to observe generation", deployment.Name), nil + } + if deployment.Status.UpdatedReplicas != specReplicas || + deployment.Status.ReadyReplicas != specReplicas || + deployment.Status.Replicas != specReplicas || + deployment.Status.UnavailableReplicas != 0 { + return false, fmt.Sprintf("Waiting for Deployment [%s] rollout", deployment.Name), nil + } + } + + for _, statefulSet := range snapshot.statefulSets { + specReplicas := int32(1) + if statefulSet.Spec.Replicas != nil { + specReplicas = *statefulSet.Spec.Replicas + } + if statefulSet.Status.ObservedGeneration < statefulSet.Generation { + return false, fmt.Sprintf("Waiting for StatefulSet [%s] controller to observe generation", statefulSet.Name), nil + } + if statefulSet.Status.CurrentRevision != statefulSet.Status.UpdateRevision || + statefulSet.Status.UpdatedReplicas != specReplicas { + return false, fmt.Sprintf("Waiting for StatefulSet [%s] revision rollout", statefulSet.Name), nil + } + if statefulSet.Status.ReadyReplicas != specReplicas { + return false, fmt.Sprintf("Waiting for StatefulSet [%s] ready replicas", statefulSet.Name), nil + } + } + + return true, "", nil +} + +func patchDeploymentLifecycleStatus( + ctx context.Context, + sdk client.Client, + drd *v1alpha1.Druid, + updated v1alpha1.DeploymentLifecycleStatus, + emitEvent EventEmitter, +) error { + current := drd.Status.DeploymentLifecycle + if reflect.DeepEqual(current, updated) { + return nil + } + + patchBytes, err := json.Marshal(map[string]interface{}{ + "status": map[string]interface{}{ + "deploymentLifecycle": deploymentLifecycleStatusPatch(updated), + }, + }) + if err != nil { + return fmt.Errorf("failed to serialize deployment lifecycle status patch: %v", err) + } + + if err := writers.Patch(ctx, sdk, drd, drd, true, client.RawPatch(types.MergePatchType, patchBytes), emitEvent); err != nil { + return err + } + + drd.Status.DeploymentLifecycle = updated + emitDeploymentLifecycleEvent(drd, emitEvent, current, updated) + return nil +} + +func deploymentLifecycleStatusPatch(status v1alpha1.DeploymentLifecycleStatus) map[string]interface{} { + return map[string]interface{}{ + "phase": status.Phase, + "reason": status.Reason, + "observedGeneration": status.ObservedGeneration, + "startedAt": status.StartedAt, + "completedAt": status.CompletedAt, + } +} + +func emitDeploymentLifecycleEvent( + drd *v1alpha1.Druid, + emitEvent EventEmitter, + previous, current v1alpha1.DeploymentLifecycleStatus, +) { + if previous.Phase == current.Phase && + previous.ObservedGeneration == current.ObservedGeneration && + previous.Reason == current.Reason { + return + } + + msg := fmt.Sprintf("observedGeneration=%d phase=%s reason=%s", current.ObservedGeneration, current.Phase, current.Reason) + if deploymentLifecycleTimedOut(current) && !deploymentLifecycleTimedOut(previous) { + emitEvent.EmitEventGeneric(drd, string(druidDeploymentLifecycleTimedOut), msg, errors.New(current.Reason)) + return + } + + switch current.Phase { + case v1alpha1.DeploymentLifecyclePending, v1alpha1.DeploymentLifecycleInProgress: + emitEvent.EmitEventGeneric(drd, string(druidDeploymentLifecycleStarted), msg, nil) + case v1alpha1.DeploymentLifecycleSucceeded: + emitEvent.EmitEventGeneric(drd, string(druidDeploymentLifecycleSucceeded), msg, nil) + } +} diff --git a/controllers/druid/lifecycle_metrics.go b/controllers/druid/lifecycle_metrics.go new file mode 100644 index 00000000..82dddeab --- /dev/null +++ b/controllers/druid/lifecycle_metrics.go @@ -0,0 +1,86 @@ +/* +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +*/ +package druid + +import ( + "time" + + "github.com/apache/druid-operator/apis/druid/v1alpha1" + "github.com/prometheus/client_golang/prometheus" + ctrlmetrics "sigs.k8s.io/controller-runtime/pkg/metrics" +) + +type deploymentLifecycleMetrics struct { + failed *prometheus.GaugeVec +} + +var defaultDeploymentLifecycleMetrics = newDeploymentLifecycleMetrics(ctrlmetrics.Registry) + +func newDeploymentLifecycleMetrics(registerer prometheus.Registerer) *deploymentLifecycleMetrics { + metrics := &deploymentLifecycleMetrics{ + failed: prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "druid_operator_deployment_lifecycle_failed", + Help: "Whether the current Druid deployment lifecycle is considered failed for metric consumers.", + }, + []string{"namespace", "druid_instance"}, + ), + } + + if registerer != nil { + registerer.MustRegister(metrics.failed) + } + + return metrics +} + +func (m *deploymentLifecycleMetrics) record(drd *v1alpha1.Druid, deps lifecycleDependencies) { + if m == nil || drd == nil { + return + } + + m.failed.WithLabelValues(drd.Namespace, drd.Name).Set(boolToFloat(deploymentLifecycleMetricFailed(drd.Status.DeploymentLifecycle, deps))) +} + +func (m *deploymentLifecycleMetrics) delete(namespace, druidName string) { + if m == nil { + return + } + + m.failed.DeleteLabelValues(namespace, druidName) +} + +func deploymentLifecycleMetricFailed(status v1alpha1.DeploymentLifecycleStatus, deps lifecycleDependencies) bool { + if status.Phase != v1alpha1.DeploymentLifecycleInProgress || status.StartedAt == nil || deps.timeout <= 0 { + return false + } + + now := time.Now + if deps.now != nil { + now = deps.now + } + return !now().Before(status.StartedAt.Time.Add(deps.timeout)) +} + +func boolToFloat(value bool) float64 { + if value { + return 1 + } + return 0 +} diff --git a/controllers/druid/lifecycle_metrics_test.go b/controllers/druid/lifecycle_metrics_test.go new file mode 100644 index 00000000..29b3fc58 --- /dev/null +++ b/controllers/druid/lifecycle_metrics_test.go @@ -0,0 +1,167 @@ +/* +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +*/ +package druid + +import ( + "testing" + "time" + + "github.com/apache/druid-operator/apis/druid/v1alpha1" + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func TestDeploymentLifecycleMetricsRecordFailedWhenTimedOut(t *testing.T) { + registry := prometheus.NewRegistry() + metrics := newDeploymentLifecycleMetrics(registry) + startedAt := metav1.NewTime(time.Unix(100, 0)) + drd := lifecycleMetricDruid(v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleInProgress, + Reason: "Waiting for Deployment [example] rollout", + StartedAt: &startedAt, + }) + + metrics.record(drd, lifecycleDependencies{ + now: func() time.Time { return time.Unix(110, 0) }, + timeout: 5 * time.Second, + }) + + families, err := registry.Gather() + require.NoError(t, err) + assert.Equal(t, 1.0, metricValue(t, families, "druid_operator_deployment_lifecycle_failed", lifecycleMetricLabels())) +} + +func TestDeploymentLifecycleMetricsRecordNotFailedWhenTimeoutReasonButConfiguredTimeoutNotExceeded(t *testing.T) { + registry := prometheus.NewRegistry() + metrics := newDeploymentLifecycleMetrics(registry) + startedAt := metav1.NewTime(time.Unix(100, 0)) + drd := lifecycleMetricDruid(v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleInProgress, + Reason: deploymentLifecycleTimeoutReasonPrefix + " of 5s; still waiting for Deployment [example] rollout", + StartedAt: &startedAt, + }) + + metrics.record(drd, lifecycleDependencies{ + now: func() time.Time { return time.Unix(110, 0) }, + timeout: 30 * time.Second, + }) + + families, err := registry.Gather() + require.NoError(t, err) + assert.Equal(t, 0.0, metricValue(t, families, "druid_operator_deployment_lifecycle_failed", lifecycleMetricLabels())) +} + +func TestDeploymentLifecycleMetricsRecordNotFailedWhenTimeoutDisabled(t *testing.T) { + registry := prometheus.NewRegistry() + metrics := newDeploymentLifecycleMetrics(registry) + startedAt := metav1.NewTime(time.Unix(100, 0)) + drd := lifecycleMetricDruid(v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleInProgress, + StartedAt: &startedAt, + }) + + metrics.record(drd, lifecycleDependencies{ + now: func() time.Time { return time.Unix(110, 0) }, + timeout: 0, + }) + + families, err := registry.Gather() + require.NoError(t, err) + assert.Equal(t, 0.0, metricValue(t, families, "druid_operator_deployment_lifecycle_failed", lifecycleMetricLabels())) +} + +func TestDeploymentLifecycleMetricsDeleteRemovesSeries(t *testing.T) { + registry := prometheus.NewRegistry() + metrics := newDeploymentLifecycleMetrics(registry) + startedAt := metav1.NewTime(time.Unix(100, 0)) + metrics.record(lifecycleMetricDruid(v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleInProgress, + StartedAt: &startedAt, + }), lifecycleDependencies{ + now: func() time.Time { return time.Unix(110, 0) }, + timeout: 5 * time.Second, + }) + + metrics.delete("default", "example") + + families, err := registry.Gather() + require.NoError(t, err) + assert.Equal(t, 0.0, metricValue(t, families, "druid_operator_deployment_lifecycle_failed", lifecycleMetricLabels())) +} + +func lifecycleMetricDruid(status v1alpha1.DeploymentLifecycleStatus) *v1alpha1.Druid { + return &v1alpha1.Druid{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + }, + Status: v1alpha1.DruidClusterStatus{ + DeploymentLifecycle: status, + }, + } +} + +func lifecycleMetricLabels() map[string]string { + return map[string]string{ + "namespace": "default", + "druid_instance": "example", + } +} + +func metricValue(t *testing.T, families []*dto.MetricFamily, name string, labels map[string]string) float64 { + t.Helper() + for _, family := range families { + if family.GetName() != name { + continue + } + for _, metric := range family.Metric { + if metricHasLabels(metric, labels) { + if metric.Gauge != nil { + return metric.Gauge.GetValue() + } + if metric.Counter != nil { + return metric.Counter.GetValue() + } + } + } + } + return 0 +} + +func metricHasLabels(metric *dto.Metric, labels map[string]string) bool { + if len(labels) == 0 { + return true + } + + matched := 0 + for _, label := range metric.Label { + expected, ok := labels[label.GetName()] + if !ok { + continue + } + if label.GetValue() != expected { + return false + } + matched++ + } + return matched == len(labels) +} diff --git a/controllers/druid/lifecycle_test.go b/controllers/druid/lifecycle_test.go new file mode 100644 index 00000000..d52177b3 --- /dev/null +++ b/controllers/druid/lifecycle_test.go @@ -0,0 +1,423 @@ +/* +Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +"License"); you may not use this file except in compliance +with the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License. +*/ +package druid + +import ( + "context" + "encoding/json" + "errors" + "testing" + "time" + + "github.com/apache/druid-operator/apis/druid/v1alpha1" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + appsv1 "k8s.io/api/apps/v1" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +type noopEventEmitter struct{} + +func (noopEventEmitter) EmitEventGeneric(obj object, eventReason, msg string, err error) {} +func (noopEventEmitter) EmitEventRollingDeployWait(obj, k8sObj object, nodeSpecUniqueStr string) {} +func (noopEventEmitter) EmitEventOnGetError(obj, getObj object, err error) {} +func (noopEventEmitter) EmitEventOnUpdate(obj, updateObj object, err error) {} +func (noopEventEmitter) EmitEventOnDelete(obj, deleteObj object, err error) {} +func (noopEventEmitter) EmitEventOnCreate(obj, createObj object, err error) {} +func (noopEventEmitter) EmitEventOnPatch(obj, patchObj object, err error) {} +func (noopEventEmitter) EmitEventOnList(obj object, listObj objectList, err error) {} + +func TestDeploymentLifecycleStatusPatchIncludesNullsForClearedFields(t *testing.T) { + startedAt := metav1.NewTime(time.Unix(100, 0)) + payload, err := json.Marshal(map[string]interface{}{ + "status": map[string]interface{}{ + "deploymentLifecycle": deploymentLifecycleStatusPatch(v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleInProgress, + Reason: "Waiting for rollout", + ObservedGeneration: 3, + StartedAt: &startedAt, + CompletedAt: nil, + }), + }, + }) + assert.NoError(t, err) + + var decoded map[string]interface{} + assert.NoError(t, json.Unmarshal(payload, &decoded)) + + status := decoded["status"].(map[string]interface{}) + lifecycle := status["deploymentLifecycle"].(map[string]interface{}) + assert.Equal(t, nil, lifecycle["completedAt"]) + assert.NotNil(t, lifecycle["startedAt"]) + assert.NotContains(t, lifecycle, "revision") + assert.NotContains(t, lifecycle, "trigger") + assert.NotContains(t, lifecycle, "lastSuccessfulImage") +} + +func TestPatchDeploymentLifecycleStatusDoesNotUpdateLocalStatusWhenPatchFails(t *testing.T) { + drd := &v1alpha1.Druid{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + }, + } + k8sClient := newLifecycleTestClient(t, drd) + + startedAt := metav1.NewTime(time.Unix(100, 0)) + updated := v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecyclePending, + StartedAt: &startedAt, + Reason: "Waiting for Druid workloads to roll out", + } + + previousWriter := writers + writers = failingPatchWriter{err: errors.New("patch failed")} + t.Cleanup(func() { + writers = previousWriter + }) + + err := patchDeploymentLifecycleStatus(context.Background(), k8sClient, drd, updated, noopEventEmitter{}) + require.EqualError(t, err, "patch failed") + assert.Empty(t, drd.Status.DeploymentLifecycle.Phase) +} + +func newLifecycleTestClient(t *testing.T, drd *v1alpha1.Druid, objects ...client.Object) client.Client { + scheme := runtime.NewScheme() + assert.NoError(t, v1alpha1.AddToScheme(scheme)) + assert.NoError(t, appsv1.AddToScheme(scheme)) + assert.NoError(t, v1.AddToScheme(scheme)) + + allObjects := append([]client.Object{drd}, objects...) + return fake.NewClientBuilder(). + WithScheme(scheme). + WithStatusSubresource(drd). + WithObjects(allObjects...). + Build() +} + +func readyManagedDeployment(drd *v1alpha1.Druid) *appsv1.Deployment { + replicas := int32(1) + return &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example-broker", + Namespace: drd.Namespace, + Labels: makeLabelsForDruid(drd), + Generation: 2, + }, + Spec: appsv1.DeploymentSpec{ + Replicas: &replicas, + }, + Status: appsv1.DeploymentStatus{ + ObservedGeneration: 2, + Replicas: 1, + ReadyReplicas: 1, + UpdatedReplicas: 1, + }, + } +} + +func readyManagedStatefulSet(drd *v1alpha1.Druid) *appsv1.StatefulSet { + replicas := int32(1) + return &appsv1.StatefulSet{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example-historical", + Namespace: drd.Namespace, + Labels: makeLabelsForDruid(drd), + Generation: 2, + }, + Spec: appsv1.StatefulSetSpec{ + Replicas: &replicas, + }, + Status: appsv1.StatefulSetStatus{ + ObservedGeneration: 2, + CurrentRevision: "rev-2", + UpdateRevision: "rev-2", + UpdatedReplicas: 1, + ReadyReplicas: 1, + }, + } +} + +func TestReconcileDeploymentLifecycleSucceedsWhenKubernetesResourcesAreReady(t *testing.T) { + + drd := &v1alpha1.Druid{ + TypeMeta: metav1.TypeMeta{ + APIVersion: "druid.apache.org/v1alpha1", + Kind: "Druid", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + Generation: 2, + }, + Spec: v1alpha1.DruidSpec{ + CommonRuntimeProperties: "druid.service=druid/router", + }, + } + + deployment := readyManagedDeployment(drd) + k8sClient := newLifecycleTestClient(t, drd, deployment) + + assert.NoError(t, reconcileDeploymentLifecycle(context.Background(), k8sClient, drd, noopEventEmitter{})) + + stored := &v1alpha1.Druid{} + assert.NoError(t, k8sClient.Get(context.Background(), client.ObjectKeyFromObject(drd), stored)) + assert.Equal(t, v1alpha1.DeploymentLifecycleSucceeded, stored.Status.DeploymentLifecycle.Phase) + assert.Equal(t, "Deployment lifecycle completed after managed workloads became ready", stored.Status.DeploymentLifecycle.Reason) + assert.Equal(t, drd.Generation, stored.Status.DeploymentLifecycle.ObservedGeneration) + assert.NotNil(t, stored.Status.DeploymentLifecycle.CompletedAt) +} + +func TestAreManagedWorkloadsReadyWaitsForDeploymentObservedGeneration(t *testing.T) { + drd := &v1alpha1.Druid{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + }, + } + + replicas := int32(1) + deployment := &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example-broker", + Namespace: drd.Namespace, + Labels: makeLabelsForDruid(drd), + Generation: 2, + }, + Spec: appsv1.DeploymentSpec{ + Replicas: &replicas, + }, + Status: appsv1.DeploymentStatus{ + ObservedGeneration: 1, + Replicas: 1, + ReadyReplicas: 1, + UpdatedReplicas: 1, + }, + } + + k8sClient := newLifecycleTestClient(t, drd, deployment) + ready, reason, err := areManagedWorkloadsReady(context.Background(), k8sClient, drd) + assert.NoError(t, err) + assert.False(t, ready) + assert.Equal(t, "Waiting for Deployment [example-broker] controller to observe generation", reason) +} + +func TestAreManagedWorkloadsReadyTreatsDeploymentReplicaFailureAsWaiting(t *testing.T) { + drd := &v1alpha1.Druid{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + }, + } + + deployment := readyManagedDeployment(drd) + deployment.Status.Conditions = []appsv1.DeploymentCondition{ + { + Type: appsv1.DeploymentReplicaFailure, + Status: v1.ConditionTrue, + Reason: "FailedCreate", + }, + } + + k8sClient := newLifecycleTestClient(t, drd, deployment) + ready, reason, err := areManagedWorkloadsReady(context.Background(), k8sClient, drd) + assert.NoError(t, err) + assert.False(t, ready) + assert.Equal(t, "Waiting for Deployment [example-broker] replica failure to clear: FailedCreate", reason) +} + +func TestAreManagedWorkloadsReadyWaitsForStatefulSetObservedGeneration(t *testing.T) { + drd := &v1alpha1.Druid{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + }, + } + + statefulSet := readyManagedStatefulSet(drd) + statefulSet.Status.ObservedGeneration = 1 + + k8sClient := newLifecycleTestClient(t, drd, statefulSet) + ready, reason, err := areManagedWorkloadsReady(context.Background(), k8sClient, drd) + assert.NoError(t, err) + assert.False(t, ready) + assert.Equal(t, "Waiting for StatefulSet [example-historical] controller to observe generation", reason) +} + +func TestAreManagedWorkloadsReadyAcceptsObservedStatefulSetGeneration(t *testing.T) { + drd := &v1alpha1.Druid{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + }, + } + + statefulSet := readyManagedStatefulSet(drd) + + k8sClient := newLifecycleTestClient(t, drd, statefulSet) + ready, reason, err := areManagedWorkloadsReady(context.Background(), k8sClient, drd) + assert.NoError(t, err) + assert.True(t, ready) + assert.Equal(t, "", reason) +} + +func TestReconcileDeploymentLifecycleKeepsTimedOutRolloutInProgress(t *testing.T) { + startedAt := metav1.NewTime(time.Unix(100, 0)) + drd := &v1alpha1.Druid{ + TypeMeta: metav1.TypeMeta{ + APIVersion: "druid.apache.org/v1alpha1", + Kind: "Druid", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + Generation: 2, + }, + Spec: v1alpha1.DruidSpec{ + CommonRuntimeProperties: "druid.service=druid/router", + }, + Status: v1alpha1.DruidClusterStatus{ + DeploymentLifecycle: v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleInProgress, + ObservedGeneration: 2, + StartedAt: &startedAt, + }, + }, + } + + deployment := readyManagedDeployment(drd) + deployment.Status.ReadyReplicas = 0 + k8sClient := newLifecycleTestClient(t, drd, deployment) + + assert.NoError(t, reconcileDeploymentLifecycleWithDeps(context.Background(), k8sClient, drd, noopEventEmitter{}, lifecycleDependencies{ + now: func() time.Time { return time.Unix(110, 0) }, + timeout: 5 * time.Second, + })) + + stored := &v1alpha1.Druid{} + assert.NoError(t, k8sClient.Get(context.Background(), client.ObjectKeyFromObject(drd), stored)) + assert.Equal(t, v1alpha1.DeploymentLifecycleInProgress, stored.Status.DeploymentLifecycle.Phase) + assert.Contains(t, stored.Status.DeploymentLifecycle.Reason, deploymentLifecycleTimeoutReasonPrefix) + assert.Contains(t, stored.Status.DeploymentLifecycle.Reason, "still waiting for Deployment [example-broker] rollout") + assert.Nil(t, stored.Status.DeploymentLifecycle.CompletedAt) +} + +func TestReconcileDeploymentLifecycleSucceedsAfterTimeoutWithoutNewGeneration(t *testing.T) { + startedAt := metav1.NewTime(time.Unix(100, 0)) + drd := &v1alpha1.Druid{ + TypeMeta: metav1.TypeMeta{ + APIVersion: "druid.apache.org/v1alpha1", + Kind: "Druid", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + Generation: 2, + }, + Spec: v1alpha1.DruidSpec{ + CommonRuntimeProperties: "druid.service=druid/router", + }, + Status: v1alpha1.DruidClusterStatus{ + DeploymentLifecycle: v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleInProgress, + Reason: deploymentLifecycleTimeoutReasonPrefix + " of 5s; still waiting for Deployment [example-broker] rollout", + ObservedGeneration: 2, + StartedAt: &startedAt, + }, + }, + } + + deployment := readyManagedDeployment(drd) + k8sClient := newLifecycleTestClient(t, drd, deployment) + + assert.NoError(t, reconcileDeploymentLifecycleWithDeps(context.Background(), k8sClient, drd, noopEventEmitter{}, lifecycleDependencies{ + now: func() time.Time { return time.Unix(120, 0) }, + timeout: 5 * time.Second, + })) + + stored := &v1alpha1.Druid{} + assert.NoError(t, k8sClient.Get(context.Background(), client.ObjectKeyFromObject(drd), stored)) + assert.Equal(t, v1alpha1.DeploymentLifecycleSucceeded, stored.Status.DeploymentLifecycle.Phase) + assert.Equal(t, "Deployment lifecycle completed after managed workloads became ready", stored.Status.DeploymentLifecycle.Reason) + assert.NotNil(t, stored.Status.DeploymentLifecycle.CompletedAt) +} + +func TestReconcileDeploymentLifecycleIsIdempotentAfterSuccess(t *testing.T) { + startedAt := metav1.NewTime(time.Unix(100, 0)) + completedAt := metav1.NewTime(time.Unix(160, 0)) + drd := &v1alpha1.Druid{ + TypeMeta: metav1.TypeMeta{ + APIVersion: "druid.apache.org/v1alpha1", + Kind: "Druid", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + Generation: 2, + }, + Spec: v1alpha1.DruidSpec{ + CommonRuntimeProperties: "druid.service=druid/router", + }, + Status: v1alpha1.DruidClusterStatus{ + DeploymentLifecycle: v1alpha1.DeploymentLifecycleStatus{ + Phase: v1alpha1.DeploymentLifecycleSucceeded, + Reason: "Deployment lifecycle completed after managed workloads became ready", + ObservedGeneration: 2, + StartedAt: &startedAt, + CompletedAt: &completedAt, + }, + }, + } + + deployment := readyManagedDeployment(drd) + k8sClient := newLifecycleTestClient(t, drd, deployment) + + before := &v1alpha1.Druid{} + assert.NoError(t, k8sClient.Get(context.Background(), client.ObjectKeyFromObject(drd), before)) + + assert.NoError(t, reconcileDeploymentLifecycle(context.Background(), k8sClient, drd, noopEventEmitter{})) + + after := &v1alpha1.Druid{} + assert.NoError(t, k8sClient.Get(context.Background(), client.ObjectKeyFromObject(drd), after)) + assert.Equal(t, before.Status.DeploymentLifecycle, after.Status.DeploymentLifecycle) +} + +type failingPatchWriter struct { + err error +} + +func (w failingPatchWriter) Delete(ctx context.Context, sdk client.Client, drd *v1alpha1.Druid, obj object, emitEvent EventEmitter, deleteOptions ...client.DeleteOption) error { + return nil +} + +func (w failingPatchWriter) Create(ctx context.Context, sdk client.Client, drd *v1alpha1.Druid, obj object, emitEvent EventEmitter) (DruidNodeStatus, error) { + return "", nil +} + +func (w failingPatchWriter) Update(ctx context.Context, sdk client.Client, drd *v1alpha1.Druid, obj object, emitEvent EventEmitter) (DruidNodeStatus, error) { + return "", nil +} + +func (w failingPatchWriter) Patch(ctx context.Context, sdk client.Client, drd *v1alpha1.Druid, obj object, status bool, patch client.Patch, emitEvent EventEmitter) error { + return w.err +} diff --git a/docs/api_specifications/druid.md b/docs/api_specifications/druid.md index 687e50f0..e237eb40 100644 --- a/docs/api_specifications/druid.md +++ b/docs/api_specifications/druid.md @@ -1,21 +1,3 @@ -

Druid API reference

Packages:

druid.apache.org/v1alpha1

+

Licensed to the Apache Software Foundation (ASF) under one +or more contributor license agreements. See the NOTICE file +distributed with this work for additional information +regarding copyright ownership. The ASF licenses this file +to you under the Apache License, Version 2.0 (the +“License”); you may not use this file except in compliance +with the License. You may obtain a copy of the License at

+

http://www.apache.org/licenses/LICENSE-2.0

+

Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on an +“AS IS” BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +KIND, either express or implied. See the License for the +specific language governing permissions and limitations +under the License.

Resource Types:

AdditionalContainer @@ -236,6 +232,88 @@ encoding/json.RawMessage +

DeploymentLifecyclePhase +(string alias)

+

+(Appears on: +DeploymentLifecycleStatus) +

+

DeploymentLifecycleStatus +

+

+(Appears on: +DruidClusterStatus) +

+
+
+ + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
FieldDescription
+phase
+ + +DeploymentLifecyclePhase + + +
+
+reason
+ +string + +
+
+observedGeneration
+ +int64 + +
+
+startedAt
+ + +Kubernetes meta/v1.Time + + +
+
+completedAt
+ + +Kubernetes meta/v1.Time + + +
+
+
+

Druid

Druid is the Schema for the druids API.

@@ -964,6 +1042,18 @@ DruidClusterStatus +deploymentLifecycle
+ + +DeploymentLifecycleStatus + + + + + + + + druidNodeStatus
@@ -972,8 +1062,6 @@ DruidNodeTypeStatus -

INSERT ADDITIONAL STATUS FIELD - define observed state of cluster -Important: Run “make” to regenerate code after modifying this file

diff --git a/docs/features.md b/docs/features.md index 6f8b74b3..324ad02d 100644 --- a/docs/features.md +++ b/docs/features.md @@ -43,6 +43,12 @@ The reconciliation time can be adjusted - in the chart, add `env.RECONCILE_WAIT` in seconds. Examples: "10s", "30s", "120s" +The deployment lifecycle timeout can be adjusted with `env.DEPLOYMENT_LIFECYCLE_TIMEOUT`. +It defaults to "96h". When the timeout is exceeded, the deployment lifecycle remains +`InProgress` with a timeout reason and can still move to `Succeeded` when the existing +deployment becomes ready. The `druid_operator_deployment_lifecycle_failed` metric is set +to `1` while the deployment lifecycle is timed out. + ## Finalizer in Druid CR The Druid operator supports provisioning of StatefulSets and Deployments. When a StatefulSet is created, a PVC is created along. When the Druid CR is deleted, the StatefulSet controller does not delete the PVC's diff --git a/go.mod b/go.mod index 73864b38..7ba2cec5 100644 --- a/go.mod +++ b/go.mod @@ -60,8 +60,8 @@ require ( github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/pmezard/go-difflib v1.0.0 // indirect - github.com/prometheus/client_golang v1.15.1 // indirect - github.com/prometheus/client_model v0.4.0 // indirect + github.com/prometheus/client_golang v1.15.1 + github.com/prometheus/client_model v0.4.0 github.com/prometheus/common v0.42.0 // indirect github.com/prometheus/procfs v0.9.0 // indirect github.com/spf13/pflag v1.0.5 // indirect