Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 28 additions & 12 deletions apis/druid/v1alpha1/druid_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
24 changes: 24 additions & 0 deletions apis/druid/v1alpha1/zz_generated.deepcopy.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

23 changes: 20 additions & 3 deletions chart/crds/druid.apache.org_druids.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions chart/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
23 changes: 20 additions & 3 deletions config/crd/bases/druid.apache.org_druids.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
66 changes: 55 additions & 11 deletions controllers/druid/druid_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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"),
}
}

Expand All @@ -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
Expand All @@ -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
Expand All @@ -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 {
Expand All @@ -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"

Expand Down
19 changes: 19 additions & 0 deletions controllers/druid/druid_controller_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{}
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
15 changes: 9 additions & 6 deletions controllers/druid/interface.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading