diff --git a/Makefile b/Makefile index 3d9ba57ed..86f6b9e53 100644 --- a/Makefile +++ b/Makefile @@ -203,6 +203,9 @@ manager: generate fmt vet ## Generate (kubebuilder) and build manager binary. run: generate fmt vet go run ./main.go +run-dev: + go run ./main.go --disable-cert-signing-support --disable-webhooks 2>&1 | jq -Cc -R 'fromjson? // empty' + # Install CRDs into a cluster by manually creating or replacing the CRD depending on whether is currently existing # Apply is not applicable as the last-applied-configuration annotation would exceed the size limit enforced by the api server install: manifests ## Install generated CRDs into the configured Kubernetes cluster. diff --git a/api/assets/kafka/jmx-exporter.yml b/api/assets/kafka/jmx-exporter.yml index a3857da29..48d953ef9 100644 --- a/api/assets/kafka/jmx-exporter.yml +++ b/api/assets/kafka/jmx-exporter.yml @@ -30,8 +30,8 @@ rules: # END: Export Kraft metrics - pattern: 'kafka.server<>(Version): ([-.~+\w\d]+)' - name: kafka_server_$1_$3 - type: COUNTER + name: kafka_server_app_info_version + type: GAUGE labels: broker_id: $2 version: $4 diff --git a/charts/kafka-operator/templates/alertmanager-peerauthentication.yaml b/charts/kafka-operator/templates/alertmanager-peerauthentication.yaml index 79f8aa3e2..d1ca50bf4 100644 --- a/charts/kafka-operator/templates/alertmanager-peerauthentication.yaml +++ b/charts/kafka-operator/templates/alertmanager-peerauthentication.yaml @@ -17,7 +17,7 @@ spec: selector: matchLabels: control-plane: controller-manager - component: alertmanager + app.kubernetes.io/component: operator portLevelMtls: {{ .Values.alertManager.port | quote }}: mode: PERMISSIVE diff --git a/charts/kafka-operator/templates/operator-deployment-with-webhook.yaml b/charts/kafka-operator/templates/operator-deployment-with-webhook.yaml index 9d82ddd3e..15dde16ba 100644 --- a/charts/kafka-operator/templates/operator-deployment-with-webhook.yaml +++ b/charts/kafka-operator/templates/operator-deployment-with-webhook.yaml @@ -133,8 +133,6 @@ spec: app.kubernetes.io/name: {{ include "kafka-operator.name" . }} app.kubernetes.io/instance: {{ .Release.Name }} app.kubernetes.io/component: operator - app: prometheus - component: alertmanager spec: {{- with .Values.imagePullSecrets }} imagePullSecrets: diff --git a/charts/kafka-operator/templates/podmonitor.yaml b/charts/kafka-operator/templates/podmonitor.yaml new file mode 100644 index 000000000..06b8b87a7 --- /dev/null +++ b/charts/kafka-operator/templates/podmonitor.yaml @@ -0,0 +1,31 @@ +{{- if .Values.prometheusMetrics.podMonitor.enabled }} +kind: PodMonitor +apiVersion: monitoring.coreos.com/v1 +metadata: + name: {{ include "kafka-operator.fullname" . }} + namespace: {{ .Release.Namespace | quote }} + labels: + helm.sh/chart: {{ include "kafka-operator.chart" . }} + app.kubernetes.io/name: {{ include "kafka-operator.name" . }} + app.kubernetes.io/instance: {{ .Release.Name }} + app.kubernetes.io/managed-by: {{ .Release.Service }} + app.kubernetes.io/version: {{ .Chart.AppVersion }} + app.kubernetes.io/component: operator + {{- with .Values.operator.annotations }} + annotations: + {{- toYaml . | nindent 4 }} + {{- end }} +spec: + namespaceSelector: + matchNames: + - {{ .Release.Namespace }} + selector: + matchLabels: + app.kubernetes.io/name: {{ include "kafka-operator.name" . }} + app.kubernetes.io/instance: {{ .Release.Name }} + app.kubernetes.io/component: operator + endpoints: + - interval: {{ .Values.prometheusMetrics.podMonitor.interval }} + port: metrics + path: /metrics +{{- end }} diff --git a/charts/kafka-operator/values.yaml b/charts/kafka-operator/values.yaml index 5d6dfca02..638cfe3a2 100644 --- a/charts/kafka-operator/values.yaml +++ b/charts/kafka-operator/values.yaml @@ -111,6 +111,10 @@ prometheusMetrics: create: true # -- ServiceAccount used by prometheus auth proxy name: kafka-operator-authproxy + podMonitor: + # -- If true, create a PodMonitor for Prometheus metrics + enabled: false + interval: 30s # -- Health probes configuration healthProbes: {} diff --git a/config/samples/kraft/simplekafkacluster_kraft.yaml b/config/samples/kraft/simplekafkacluster_kraft.yaml index 1e993bb49..0cc38b28d 100644 --- a/config/samples/kraft/simplekafkacluster_kraft.yaml +++ b/config/samples/kraft/simplekafkacluster_kraft.yaml @@ -7,11 +7,14 @@ metadata: spec: kRaft: true monitoringConfig: - jmxImage: "ghcr.io/adobe/koperator/jmx-javaagent:1.4.0" + jmxImage: "ghcr.io/adobe/koperator/jmx-javaagent:1.5.0" headlessServiceEnabled: true propagateLabels: false oneBrokerPerNode: false - clusterImage: "ghcr.io/adobe/koperator/kafka:2.13-3.9.1" + clusterImage: "ghcr.io/adobe/koperator/kafka:2.13-3.9.2-jdk21.0.11" + taintedBrokersSelector: + matchLabels: + restart: me readOnlyConfig: | auto.create.topics.enable=false cruise.control.metrics.topic.auto.create=true @@ -30,6 +33,9 @@ spec: requests: storage: 10Gi broker: + kafkaJvmPerfOpts: >- + -Dcom.sun.management.jmxremote.port=1090 -Dcom.sun.management.jmxremote.rmi.port=1090 + -Dcom.sun.management.jmxremote.local.only=false -Djava.rmi.server.hostname=127.0.0.1 -Djute.maxbuffer=0x9fffff processRoles: - broker storageConfigs: @@ -45,27 +51,26 @@ spec: prometheus.io/port: "9020" brokers: - id: 0 - brokerConfigGroup: "broker" - - id: 1 - brokerConfigGroup: "broker" - - id: 2 - brokerConfigGroup: "broker" - - id: 3 brokerConfigGroup: "default" brokerConfig: processRoles: - controller - # - broker - - id: 4 + - id: 1 brokerConfigGroup: "default" brokerConfig: processRoles: - controller - - id: 5 + - id: 2 brokerConfigGroup: "default" brokerConfig: processRoles: - controller + - id: 100 + brokerConfigGroup: "broker" + - id: 101 + brokerConfigGroup: "broker" + - id: 102 + brokerConfigGroup: "broker" rollingUpgradeConfig: failureThreshold: 1 listenersConfig: @@ -300,3 +305,18 @@ spec: { "min.insync.replicas": 3 } + capacityConfig: |- + { + "brokerCapacities":[ + { + "brokerId": "-1", + "capacity": { + "DISK": {"/kafka-logs-broker/kafka": "10240"}, + "CPU": {"num.cores": "1"}, + "NW_IN": "900000", + "NW_OUT": "900000" + }, + "doc": "This is the default capacity. Capacity unit used for disk is in MB, cpu is in cores, network throughput is in KB." + } + ] + } \ No newline at end of file diff --git a/controllers/kafkacluster_controller.go b/controllers/kafkacluster_controller.go index 4b825d773..15daee703 100644 --- a/controllers/kafkacluster_controller.go +++ b/controllers/kafkacluster_controller.go @@ -17,6 +17,7 @@ package controllers import ( "context" + "encoding/json" "fmt" "reflect" "time" @@ -359,6 +360,41 @@ func (r *KafkaClusterReconciler) updateAndFetchLatest(ctx context.Context, clust return cluster, nil } +// ignoreMetadataFields returns a patch.CalculateOption that drops only the given fields +// under metadata (e.g. "managedFields", "resourceVersion") instead of the whole metadata +// object like patch.IgnoreField("metadata") does, so other metadata changes (labels, +// annotations, etc.) still count toward the diff. +func ignoreMetadataFields(fields ...string) patch.CalculateOption { + return func(current, modified []byte) ([]byte, []byte, error) { + current, err := deleteMetadataFields(current, fields...) + if err != nil { + return nil, nil, errors.WrapIf(err, "could not delete metadata fields from current byte sequence") + } + + modified, err = deleteMetadataFields(modified, fields...) + if err != nil { + return nil, nil, errors.WrapIf(err, "could not delete metadata fields from modified byte sequence") + } + + return current, modified, nil + } +} + +func deleteMetadataFields(obj []byte, fields ...string) ([]byte, error) { + var objectMap map[string]interface{} + if err := json.Unmarshal(obj, &objectMap); err != nil { + return nil, errors.WrapIf(err, "could not unmarshal byte sequence") + } + + if metadata, ok := objectMap["metadata"].(map[string]interface{}); ok { + for _, field := range fields { + delete(metadata, field) + } + } + + return json.Marshal(objectMap) +} + // SetupKafkaClusterWithManager registers kafka cluster controller to the manager func SetupKafkaClusterWithManager(mgr ctrl.Manager, contourEnabled bool) *ctrl.Builder { log := mgr.GetLogger() @@ -390,7 +426,7 @@ func SetupKafkaClusterWithManager(mgr ctrl.Manager, contourEnabled bool) *ctrl.B UpdateFunc: func(e event.UpdateEvent) bool { switch newObj := e.ObjectNew.(type) { case *corev1.Pod, *corev1.ConfigMap, *corev1.PersistentVolumeClaim: - patchResult, err := patch.DefaultPatchMaker.Calculate(e.ObjectOld, e.ObjectNew) + patchResult, err := patch.DefaultPatchMaker.Calculate(e.ObjectOld, e.ObjectNew, patch.IgnoreStatusFields(), ignoreMetadataFields("managedFields", "resourceVersion")) if err != nil { log.Error(err, "could not match objects", "kind", e.ObjectOld.GetObjectKind()) } else if patchResult.IsEmpty() { diff --git a/pkg/jmxextractor/extractor.go b/pkg/jmxextractor/extractor.go index 583e06dd7..099fdcdb6 100644 --- a/pkg/jmxextractor/extractor.go +++ b/pkg/jmxextractor/extractor.go @@ -20,6 +20,7 @@ import ( "io" "net/http" "regexp" + "time" "github.com/banzaicloud/koperator/api/v1beta1" "github.com/banzaicloud/koperator/pkg/errorfactory" @@ -30,9 +31,9 @@ import ( ) const ( - headlessServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessServiceTemplate + ".%s.svc.%s:%d" - headlessControllerServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessControllerServiceTemplate + ".%s.svc.%s:%d" - serviceJMXTemplate = "http://%s-%d.%s.svc.%s:%d" + headlessServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessServiceTemplate + ".%s.svc.%s:%d/metrics" + headlessControllerServiceJMXTemplate = "http://%s-%d." + kafka.HeadlessControllerServiceTemplate + ".%s.svc.%s:%d/metrics" + serviceJMXTemplate = "http://%s-%d.%s.svc.%s:%d/metrics" versionRegexGroup = "version" ) @@ -90,7 +91,8 @@ func (exp *jmxExtractor) ExtractDockerImageAndVersion(brokerId int32, brokerConf requestURL = fmt.Sprintf(serviceJMXTemplate, exp.clusterName, brokerId, exp.clusterNamespace, exp.kubernetesClusterDomain, 9020) } - rsp, err := http.Get(requestURL) + client := &http.Client{Timeout: 30 * time.Second} + rsp, err := client.Get(requestURL) if err != nil { exp.log.Error(err, fmt.Sprintf("error during talking to broker-%d", brokerId)) return nil, errorfactory.New(errorfactory.BrokersNotReady{}, err, "unable to talk to ...") diff --git a/pkg/k8sutil/status.go b/pkg/k8sutil/status.go index 4836cb3ab..de8b36194 100644 --- a/pkg/k8sutil/status.go +++ b/pkg/k8sutil/status.go @@ -127,7 +127,7 @@ func UpdateBrokerStatus(c client.Client, brokerIDs []string, cluster *banzaiclou } // update loses the typeMeta of the config that's used later when setting ownerrefs cluster.TypeMeta = typeMeta - logger.Info("Kafka cluster state updated") + logger.V(1).Info("Broker status updated", "brokers", brokerIDs, "status", state) return nil } diff --git a/pkg/kafkaclient/client.go b/pkg/kafkaclient/client.go index 83f388596..8f0a23d68 100644 --- a/pkg/kafkaclient/client.go +++ b/pkg/kafkaclient/client.go @@ -130,7 +130,7 @@ func NewFromCluster(k8sclient client.Client, cluster *v1beta1.KafkaCluster) (Kaf if err := client.Close(); err != nil { log.Error(err, "Error closing Kafka client") } else { - log.Info("Kafka client closed cleanly") + log.V(1).Info("Kafka client closed cleanly") } } return client, close, err diff --git a/pkg/resources/cruisecontrol/configmap.go b/pkg/resources/cruisecontrol/configmap.go index 20735de89..b121f70dc 100644 --- a/pkg/resources/cruisecontrol/configmap.go +++ b/pkg/resources/cruisecontrol/configmap.go @@ -157,7 +157,7 @@ type JBODInvariantCapacityConfig struct { func GenerateCapacityConfig(kafkaCluster *v1beta1.KafkaCluster, log logr.Logger, config *corev1.ConfigMap) (string, error) { var err error - log.Info("generating capacity config") + log.V(2).Info("generating capacity config") var capacityConfig JBODInvariantCapacityConfig var userConfigBrokerIds []string @@ -184,7 +184,7 @@ func GenerateCapacityConfig(kafkaCluster *v1beta1.KafkaCluster, log logr.Logger, } // If the -1 default exists we don't have to do anything else here since all brokers will have values. if brokerId == "-1" { - log.Info("Using user provided capacity config because it has universal default defined", "capacity config", userProvidedCapacityConfig) + log.V(2).Info("Using user provided capacity config because it has universal default defined", "capacity config", userProvidedCapacityConfig) return userProvidedCapacityConfig, nil } userConfigBrokerIds = append(userConfigBrokerIds, brokerId) @@ -211,7 +211,7 @@ func GenerateCapacityConfig(kafkaCluster *v1beta1.KafkaCluster, log logr.Logger, if err != nil { return "", errors.WrapIf(err, "could not marshal cruise control capacity config") } - log.Info("broker capacity config generated successfully", "capacity config", string(result)) + log.V(2).Info("broker capacity config generated successfully", "capacity config", string(result)) return string(result), err } diff --git a/pkg/resources/kafka/kafka.go b/pkg/resources/kafka/kafka.go index d737754f8..a5e364924 100644 --- a/pkg/resources/kafka/kafka.go +++ b/pkg/resources/kafka/kafka.go @@ -211,7 +211,7 @@ func (r *Reconciler) Reconcile(log logr.Logger) error { log.V(1).Info("Reconciling") - log.Info("broker rack map", "kafkaBrokerAvailabilityZoneMap", getBrokerAzMap(r.KafkaCluster)) + log.V(1).Info("broker rack map", "kafkaBrokerAvailabilityZoneMap", getBrokerAzMap(r.KafkaCluster)) ctx := context.Background() if err := k8sutil.UpdateBrokerConfigurationBackup(r.Client, r.KafkaCluster); err != nil { @@ -414,6 +414,7 @@ func (r *Reconciler) Reconcile(log logr.Logger) error { reorderedBrokers := reorderBrokers(runningBrokers, boundPersistentVolumeClaims, r.KafkaCluster.Spec.Brokers, r.KafkaCluster.Status.BrokersState, controllerID, log) allBrokerDynamicConfigSucceeded := true + brokerStatus := make(map[int32]*banzaiv1beta1.BrokerConfig) for _, broker := range reorderedBrokers { brokerConfig, err := broker.GetBrokerConfig(r.KafkaCluster.Spec) if err != nil { @@ -454,9 +455,10 @@ func (r *Reconciler) Reconcile(log logr.Logger) error { if err != nil { return err } - if err = r.updateStatusWithDockerImageAndVersion(broker.Id, brokerConfig, log); err != nil { - return err + if r.brokerNeedsVersionUpdate(broker.Id, brokerConfig) { + brokerStatus[broker.Id] = brokerConfig } + // If dynamic configs can not be set then let the loop continue to the next broker, // after the loop we return error. This solves that case when other brokers could get healthy, // but the loop exits too soon because dynamic configs can not be set. @@ -479,6 +481,10 @@ func (r *Reconciler) Reconcile(log logr.Logger) error { return err } + if err := r.updateStatusWithDockerImageAndVersion(brokerStatus, log); err != nil { + return err + } + // in case HeadlessServiceEnabled is changed, delete the service that was created by the previous // reconcile flow. The services must be deleted at the end of the reconcile flow after the new services // were created and broker configurations reflecting the new services otherwise the Kafka brokers @@ -898,7 +904,7 @@ func (r *Reconciler) reconcileKafkaPod(log logr.Logger, desiredPod *corev1.Pod, brokerId := currentPod.Labels[banzaiv1beta1.BrokerIdLabelKey] if _, ok := r.KafkaCluster.Status.BrokersState[brokerId]; ok { if currentPod.Spec.NodeName == "" { - log.Info(fmt.Sprintf("pod for brokerId %s does not scheduled to node yet", brokerId)) + log.V(1).Info(fmt.Sprintf("pod for brokerId %s not scheduled to node yet", brokerId)) } else if r.KafkaCluster.Spec.RackAwareness != nil { rackAwarenessState, err := k8sutil.UpdateCrWithRackAwarenessConfig(currentPod, r.KafkaCluster, r.Client, r.DirectClient) if err != nil { @@ -922,21 +928,44 @@ func (r *Reconciler) reconcileKafkaPod(log logr.Logger, desiredPod *corev1.Pod, return nil } -func (r *Reconciler) updateStatusWithDockerImageAndVersion(brokerId int32, brokerConfig *banzaiv1beta1.BrokerConfig, - log logr.Logger) error { - jmxExp := jmxextractor.NewJMXExtractor(r.KafkaCluster.GetNamespace(), - r.KafkaCluster.Spec.GetKubernetesClusterDomain(), r.KafkaCluster.GetName(), log) +// brokerNeedsVersionUpdate returns true when the broker's image/version status is absent, +// incomplete, or stale relative to the desired image — i.e. a JMX fetch is warranted. +func (r *Reconciler) brokerNeedsVersionUpdate(brokerID int32, brokerConfig *banzaiv1beta1.BrokerConfig) bool { + desiredImage := util.GetBrokerImage(brokerConfig, r.KafkaCluster.Spec.GetClusterImage()) + state, ok := r.KafkaCluster.Status.BrokersState[strconv.Itoa(int(brokerID))] + return !ok || state.Version == "" || state.Image != desiredImage +} - kafkaVersion, err := jmxExp.ExtractDockerImageAndVersion(brokerId, brokerConfig, - r.KafkaCluster.Spec.GetClusterImage(), r.KafkaCluster.Spec.HeadlessServiceEnabled) - if err != nil { - return err +type brokerVersionResult struct { + brokerID int32 + kafkaVersion *banzaiv1beta1.KafkaVersion + err error +} + +func (r *Reconciler) updateStatusWithDockerImageAndVersion(brokers map[int32]*banzaiv1beta1.BrokerConfig, log logr.Logger) error { + ch := make(chan brokerVersionResult, len(brokers)) + for brokerID, brokerConfig := range brokers { + go func(id int32, cfg *banzaiv1beta1.BrokerConfig) { + jmxExp := jmxextractor.NewJMXExtractor(r.KafkaCluster.GetNamespace(), r.KafkaCluster.Spec.GetKubernetesClusterDomain(), r.KafkaCluster.GetName(), log) + kv, err := jmxExp.ExtractDockerImageAndVersion(id, cfg, r.KafkaCluster.Spec.GetClusterImage(), r.KafkaCluster.Spec.HeadlessServiceEnabled) + if err != nil { + ch <- brokerVersionResult{brokerID: id, err: err} + return + } + ch <- brokerVersionResult{brokerID: id, kafkaVersion: kv} + }(brokerID, brokerConfig) } - err = k8sutil.UpdateBrokerStatus(r.Client, []string{strconv.Itoa(int(brokerId))}, r.KafkaCluster, - *kafkaVersion, log) - if err != nil { - return err + + for range brokers { + result := <-ch + if result.err != nil { + return result.err + } + if err := k8sutil.UpdateBrokerStatus(r.Client, []string{strconv.Itoa(int(result.brokerID))}, r.KafkaCluster, *result.kafkaVersion, log); err != nil { + return err + } } + return nil } @@ -962,7 +991,7 @@ func (r *Reconciler) handleRollingUpgrade(log logr.Logger, desiredPod, currentPo case err != nil: log.Error(err, "could not match objects", "kind", desiredType) case r.isPodTainted(log, currentPod): - log.Info("pod has tainted labels, deleting it", "pod", currentPod) + log.Info("pod has tainted labels, attempting to delete", "pod", currentPod) case patchResult.IsEmpty(): if !k8sutil.IsPodContainsTerminatedContainer(currentPod) && r.KafkaCluster.Status.BrokersState[currentPod.Labels[banzaiv1beta1.BrokerIdLabelKey]].ConfigurationState == banzaiv1beta1.ConfigInSync && @@ -1032,14 +1061,14 @@ func (r *Reconciler) handleRollingUpgrade(log logr.Logger, desiredPod, currentPo return errors.WrapIf(err, "health check failed") } if len(allOfflineReplicas) > 0 { - log.Info("offline replicas", "IDs", allOfflineReplicas) + log.V(1).Info("offline replicas", "IDs", allOfflineReplicas) } outOfSyncReplicas, err := kClient.OutOfSyncReplicas() if err != nil { return errors.WrapIf(err, "health check failed") } if len(outOfSyncReplicas) > 0 { - log.Info("out-of-sync replicas", "IDs", outOfSyncReplicas) + log.V(1).Info("out-of-sync replicas", "IDs", outOfSyncReplicas) } impactedReplicas := make(map[int32]struct{}) for _, brokerID := range allOfflineReplicas { diff --git a/pkg/resources/kafka/kafka_test.go b/pkg/resources/kafka/kafka_test.go index ad9e6db4b..2c2b34dad 100644 --- a/pkg/resources/kafka/kafka_test.go +++ b/pkg/resources/kafka/kafka_test.go @@ -1986,3 +1986,107 @@ func TestGetBrokerAzMap(t *testing.T) { }) } } + +func TestBrokerNeedsVersionUpdate(t *testing.T) { + t.Parallel() + const clusterImage = "apache/kafka:3.4.0" + const updatedImage = "apache/kafka:3.5.0" + + testCases := []struct { + testName string + brokerID int32 + brokerConfig v1beta1.BrokerConfig + clusterImage string + brokersState map[string]v1beta1.BrokerState + expected bool + }{ + { + testName: "no existing status entry triggers update", + brokerID: 0, + brokerConfig: v1beta1.BrokerConfig{}, + clusterImage: clusterImage, + brokersState: map[string]v1beta1.BrokerState{}, + expected: true, + }, + { + testName: "status entry with empty version triggers update", + brokerID: 0, + brokerConfig: v1beta1.BrokerConfig{}, + clusterImage: clusterImage, + brokersState: map[string]v1beta1.BrokerState{ + "0": {Image: clusterImage, Version: ""}, + }, + expected: true, + }, + { + testName: "status entry with different image triggers update", + brokerID: 0, + brokerConfig: v1beta1.BrokerConfig{}, + clusterImage: updatedImage, + brokersState: map[string]v1beta1.BrokerState{ + "0": {Image: clusterImage, Version: "3.4.0"}, + }, + expected: true, + }, + { + testName: "status up to date with cluster image skips update", + brokerID: 0, + brokerConfig: v1beta1.BrokerConfig{}, + clusterImage: clusterImage, + brokersState: map[string]v1beta1.BrokerState{ + "0": {Image: clusterImage, Version: "3.4.0"}, + }, + expected: false, + }, + { + testName: "broker-level image override used instead of cluster image", + brokerID: 0, + brokerConfig: v1beta1.BrokerConfig{Image: "apache/kafka:3.4.1"}, + clusterImage: clusterImage, + brokersState: map[string]v1beta1.BrokerState{ + "0": {Image: "apache/kafka:3.4.1", Version: "3.4.1"}, + }, + expected: false, + }, + { + testName: "broker-level image override differs from recorded image triggers update", + brokerID: 0, + brokerConfig: v1beta1.BrokerConfig{Image: "apache/kafka:3.5.0"}, + clusterImage: clusterImage, + brokersState: map[string]v1beta1.BrokerState{ + "0": {Image: "apache/kafka:3.4.1", Version: "3.4.1"}, + }, + expected: true, + }, + { + testName: "correct state for one broker does not suppress update for another", + brokerID: 1, + brokerConfig: v1beta1.BrokerConfig{}, + clusterImage: clusterImage, + brokersState: map[string]v1beta1.BrokerState{ + "0": {Image: clusterImage, Version: "3.4.0"}, + }, + expected: true, + }, + } + + for _, test := range testCases { + test := test + t.Run(test.testName, func(t *testing.T) { + t.Parallel() + r := Reconciler{ + Reconciler: resources.Reconciler{ + KafkaCluster: &v1beta1.KafkaCluster{ + Spec: v1beta1.KafkaClusterSpec{ + ClusterImage: test.clusterImage, + }, + Status: v1beta1.KafkaClusterStatus{ + BrokersState: test.brokersState, + }, + }, + }, + } + assert.Equal(t, test.expected, r.brokerNeedsVersionUpdate(test.brokerID, &test.brokerConfig)) + }) + } +}