RockteMQ-AI commented on code in PR #63:
URL: https://github.com/apache/rocketmq-operator/pull/63#discussion_r3902468051


##########
pkg/controller/rocketmq/rocketmq_controller.go:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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 rocketmq contains the implementation of the Rocketmq CRD reconcile 
function
+package rocketmq
+
+import (
+       "context"
+       rocketmqv1alpha1 
"github.com/apache/rocketmq-operator/pkg/apis/rocketmq/v1alpha1"
+       cons "github.com/apache/rocketmq-operator/pkg/constants"
+       corev1 "k8s.io/api/core/v1"
+       "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/runtime"
+       "k8s.io/apimachinery/pkg/types"
+       "reflect"
+       "sigs.k8s.io/controller-runtime/pkg/client"
+       "sigs.k8s.io/controller-runtime/pkg/controller"
+       "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
+       "sigs.k8s.io/controller-runtime/pkg/handler"
+       "sigs.k8s.io/controller-runtime/pkg/manager"
+       "sigs.k8s.io/controller-runtime/pkg/reconcile"
+       logf "sigs.k8s.io/controller-runtime/pkg/runtime/log"
+       "sigs.k8s.io/controller-runtime/pkg/source"
+       "time"
+)
+
+var log = logf.Log.WithName("controller_rocketmq")
+
+// Add creates a new Rocketmq Controller and adds it to the Manager. The 
Manager will set fields on the Controller
+// and Start it when the Manager is Started.
+func Add(mgr manager.Manager) error {
+       return add(mgr, newReconciler(mgr))
+}
+
+// newReconciler returns a new reconcile.Reconciler
+func newReconciler(mgr manager.Manager) reconcile.Reconciler {
+       return &ReconcileRocketmq{client: mgr.GetClient(), scheme: 
mgr.GetScheme()}
+}
+
+// add adds a new Controller to mgr with r as the reconcile.Reconciler
+func add(mgr manager.Manager, r reconcile.Reconciler) error {
+       c, err := controller.New("rocketmq-controller", mgr, 
controller.Options{Reconciler: r})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to primary resource Rocketmq
+       err = c.Watch(&source.Kind{Type: &rocketmqv1alpha1.Rocketmq{}}, 
&handler.EnqueueRequestForObject{})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to secondary resource Pods and requeue the owner 
Rocketmq
+       err = c.Watch(&source.Kind{Type: &corev1.Pod{}}, 
&handler.EnqueueRequestForOwner{
+               IsController: true,
+               OwnerType: &rocketmqv1alpha1.Rocketmq{},
+       })
+       if err != nil {
+               return err
+       }
+       return nil
+}
+
+// blank assignment to verify that ReconcileRocketmq implements 
reconcile.Reconciler
+var _ reconcile.Reconciler = &ReconcileRocketmq{}
+
+// ReconcileRocketmq reconciles a Rocketmq object
+type ReconcileRocketmq struct {
+       // This client, initialized using mgr.Client() above, is a split client
+       // that reads objects from the cache and writes to the apiserver
+       client client.Client
+       scheme *runtime.Scheme
+}
+
+// Reconcile reads that state of the cluster for a NameService object and 
makes changes based on the state read
+// and what is in the NameService.Spec
+// TODO(user): Modify this Reconcile function to implement your Controller 
logic.  This example creates
+// a Pod as an example
+// Note:
+// The Controller will requeue the Request to be processed again if the 
returned error is non-nil or
+// Result.Requeue is true, otherwise upon completion it will remove the work 
from the queue.
+func (r *ReconcileRocketmq) Reconcile(request reconcile.Request) 
(reconcile.Result, error) {
+       reqLogger := log.WithValues("Request.Namespace", request.Namespace, 
"Request.Name", request.Name)
+       reqLogger.Info("Reconciling Rocketmq")
+       // Fetch the Rocketmq instance
+       instance := &rocketmqv1alpha1.Rocketmq{}
+       err := r.client.Get(context.TODO(), request.NamespacedName, instance)
+       if err != nil {
+               if errors.IsNotFound(err) {
+                       // 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
+                       return reconcile.Result{}, nil
+               }
+               // Error reading the object, requeue the request
+               return reconcile.Result{}, err
+       }
+
+       // Check if nameserver already exist, if not create a new one
+       nameServiceFound := &rocketmqv1alpha1.NameService{}
+       nameServiceSts := r.nameServiceForRocketmq(instance)
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
nameServiceSts.Name, Namespace: nameServiceSts.Namespace}, nameServiceFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), nameServiceSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new nameService 
of rocketmq", "nameservice.namespace",
+                               nameServiceSts.Namespace, "nameservice.Name", 
nameServiceSts.Name)
+                       return reconcile.Result{}, err
+               }
+       } else if err != nil {
+               reqLogger.Error(err, "Failed to get rocketmq nameservice.")
+               return reconcile.Result{}, err
+       } else {
+               // Resource NameService will change; Only size nameServiceImage 
imagePullPolicy can update
+               if !reflect.DeepEqual(nameServiceSts.Spec.Size, 
nameServiceFound.Spec.Size) ||
+                       
!reflect.DeepEqual(nameServiceSts.Spec.NameServiceImage, 
nameServiceFound.Spec.NameServiceImage) ||
+                       !reflect.DeepEqual(nameServiceSts.Spec.ImagePullPolicy, 
nameServiceFound.Spec.ImagePullPolicy) {
+                       nameServiceFound.Spec.ImagePullPolicy = 
nameServiceSts.Spec.ImagePullPolicy
+                       nameServiceFound.Spec.NameServiceImage = 
nameServiceSts.Spec.NameServiceImage
+                       nameServiceFound.Spec.Size = nameServiceSts.Spec.Size
+                       err = r.client.Update(context.TODO(), nameServiceFound)
+                       if err != nil {
+                               reqLogger.Error(err, "should update nameservice 
resource", "spec.nameService",
+                                       nameServiceSts.Spec, 
"nameServiceFound.spec", nameServiceFound.Spec)
+                               return reconcile.Result{}, err
+                       }
+               }
+       }
+
+       // Check if broker already exists, if not create a new one
+       brokerFound := &rocketmqv1alpha1.Broker{}
+       brokerSts := r.brokerForRocketmq(instance)
+
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
brokerSts.Name, Namespace: brokerSts.Namespace}, brokerFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), brokerSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new broker of 
rocketmq", "broker.namespace",
+                               brokerSts.Namespace, "broker.Name", 
brokerSts.Name)

Review Comment:
   Broker creation error is logged but not returned. When `r.client.Create` 
fails, the code falls through to the `return reconcile.Result{Requeue: true}, 
nil` below, silently swallowing the error. This should `return 
reconcile.Result{}, err` on create failure, consistent with the NameService 
create path above.



##########
pkg/controller/rocketmq/rocketmq_controller.go:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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 rocketmq contains the implementation of the Rocketmq CRD reconcile 
function
+package rocketmq
+
+import (
+       "context"
+       rocketmqv1alpha1 
"github.com/apache/rocketmq-operator/pkg/apis/rocketmq/v1alpha1"
+       cons "github.com/apache/rocketmq-operator/pkg/constants"
+       corev1 "k8s.io/api/core/v1"
+       "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/runtime"
+       "k8s.io/apimachinery/pkg/types"
+       "reflect"
+       "sigs.k8s.io/controller-runtime/pkg/client"
+       "sigs.k8s.io/controller-runtime/pkg/controller"
+       "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
+       "sigs.k8s.io/controller-runtime/pkg/handler"
+       "sigs.k8s.io/controller-runtime/pkg/manager"
+       "sigs.k8s.io/controller-runtime/pkg/reconcile"
+       logf "sigs.k8s.io/controller-runtime/pkg/runtime/log"
+       "sigs.k8s.io/controller-runtime/pkg/source"
+       "time"
+)
+
+var log = logf.Log.WithName("controller_rocketmq")
+
+// Add creates a new Rocketmq Controller and adds it to the Manager. The 
Manager will set fields on the Controller
+// and Start it when the Manager is Started.
+func Add(mgr manager.Manager) error {
+       return add(mgr, newReconciler(mgr))
+}
+
+// newReconciler returns a new reconcile.Reconciler
+func newReconciler(mgr manager.Manager) reconcile.Reconciler {
+       return &ReconcileRocketmq{client: mgr.GetClient(), scheme: 
mgr.GetScheme()}
+}
+
+// add adds a new Controller to mgr with r as the reconcile.Reconciler
+func add(mgr manager.Manager, r reconcile.Reconciler) error {
+       c, err := controller.New("rocketmq-controller", mgr, 
controller.Options{Reconciler: r})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to primary resource Rocketmq
+       err = c.Watch(&source.Kind{Type: &rocketmqv1alpha1.Rocketmq{}}, 
&handler.EnqueueRequestForObject{})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to secondary resource Pods and requeue the owner 
Rocketmq
+       err = c.Watch(&source.Kind{Type: &corev1.Pod{}}, 
&handler.EnqueueRequestForOwner{
+               IsController: true,
+               OwnerType: &rocketmqv1alpha1.Rocketmq{},
+       })
+       if err != nil {
+               return err
+       }
+       return nil
+}
+
+// blank assignment to verify that ReconcileRocketmq implements 
reconcile.Reconciler
+var _ reconcile.Reconciler = &ReconcileRocketmq{}
+
+// ReconcileRocketmq reconciles a Rocketmq object
+type ReconcileRocketmq struct {
+       // This client, initialized using mgr.Client() above, is a split client
+       // that reads objects from the cache and writes to the apiserver
+       client client.Client
+       scheme *runtime.Scheme
+}
+
+// Reconcile reads that state of the cluster for a NameService object and 
makes changes based on the state read
+// and what is in the NameService.Spec
+// TODO(user): Modify this Reconcile function to implement your Controller 
logic.  This example creates
+// a Pod as an example
+// Note:
+// The Controller will requeue the Request to be processed again if the 
returned error is non-nil or
+// Result.Requeue is true, otherwise upon completion it will remove the work 
from the queue.
+func (r *ReconcileRocketmq) Reconcile(request reconcile.Request) 
(reconcile.Result, error) {
+       reqLogger := log.WithValues("Request.Namespace", request.Namespace, 
"Request.Name", request.Name)
+       reqLogger.Info("Reconciling Rocketmq")
+       // Fetch the Rocketmq instance
+       instance := &rocketmqv1alpha1.Rocketmq{}
+       err := r.client.Get(context.TODO(), request.NamespacedName, instance)
+       if err != nil {
+               if errors.IsNotFound(err) {
+                       // 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
+                       return reconcile.Result{}, nil
+               }
+               // Error reading the object, requeue the request
+               return reconcile.Result{}, err
+       }
+
+       // Check if nameserver already exist, if not create a new one
+       nameServiceFound := &rocketmqv1alpha1.NameService{}
+       nameServiceSts := r.nameServiceForRocketmq(instance)
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
nameServiceSts.Name, Namespace: nameServiceSts.Namespace}, nameServiceFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), nameServiceSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new nameService 
of rocketmq", "nameservice.namespace",
+                               nameServiceSts.Namespace, "nameservice.Name", 
nameServiceSts.Name)

Review Comment:
   No NameService readiness check before broker creation. The NameService CR is 
created (line 118) and the broker is created immediately after (line 148) in 
the same reconcile pass. The README states 'the name server cluster will be 
created first, after all name server cluster is in running state, the operator 
will create the broker cluster', but this code does not wait for NameService 
pods to become Ready. The broker's `nameServers` endpoint will be unreachable 
on first deploy. Add a check that all NameService pods are Running/Ready before 
proceeding to broker creation — return `reconcile.Result{RequeueAfter: ...}` if 
they aren't.



##########
pkg/controller/rocketmq/rocketmq_controller.go:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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 rocketmq contains the implementation of the Rocketmq CRD reconcile 
function
+package rocketmq
+
+import (
+       "context"
+       rocketmqv1alpha1 
"github.com/apache/rocketmq-operator/pkg/apis/rocketmq/v1alpha1"
+       cons "github.com/apache/rocketmq-operator/pkg/constants"
+       corev1 "k8s.io/api/core/v1"
+       "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/runtime"
+       "k8s.io/apimachinery/pkg/types"
+       "reflect"
+       "sigs.k8s.io/controller-runtime/pkg/client"
+       "sigs.k8s.io/controller-runtime/pkg/controller"
+       "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
+       "sigs.k8s.io/controller-runtime/pkg/handler"
+       "sigs.k8s.io/controller-runtime/pkg/manager"
+       "sigs.k8s.io/controller-runtime/pkg/reconcile"
+       logf "sigs.k8s.io/controller-runtime/pkg/runtime/log"
+       "sigs.k8s.io/controller-runtime/pkg/source"
+       "time"
+)
+
+var log = logf.Log.WithName("controller_rocketmq")
+
+// Add creates a new Rocketmq Controller and adds it to the Manager. The 
Manager will set fields on the Controller
+// and Start it when the Manager is Started.
+func Add(mgr manager.Manager) error {
+       return add(mgr, newReconciler(mgr))
+}
+
+// newReconciler returns a new reconcile.Reconciler
+func newReconciler(mgr manager.Manager) reconcile.Reconciler {
+       return &ReconcileRocketmq{client: mgr.GetClient(), scheme: 
mgr.GetScheme()}
+}
+
+// add adds a new Controller to mgr with r as the reconcile.Reconciler
+func add(mgr manager.Manager, r reconcile.Reconciler) error {
+       c, err := controller.New("rocketmq-controller", mgr, 
controller.Options{Reconciler: r})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to primary resource Rocketmq
+       err = c.Watch(&source.Kind{Type: &rocketmqv1alpha1.Rocketmq{}}, 
&handler.EnqueueRequestForObject{})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to secondary resource Pods and requeue the owner 
Rocketmq
+       err = c.Watch(&source.Kind{Type: &corev1.Pod{}}, 
&handler.EnqueueRequestForOwner{
+               IsController: true,
+               OwnerType: &rocketmqv1alpha1.Rocketmq{},
+       })
+       if err != nil {
+               return err
+       }
+       return nil
+}
+
+// blank assignment to verify that ReconcileRocketmq implements 
reconcile.Reconciler
+var _ reconcile.Reconciler = &ReconcileRocketmq{}
+
+// ReconcileRocketmq reconciles a Rocketmq object
+type ReconcileRocketmq struct {
+       // This client, initialized using mgr.Client() above, is a split client
+       // that reads objects from the cache and writes to the apiserver
+       client client.Client
+       scheme *runtime.Scheme
+}
+
+// Reconcile reads that state of the cluster for a NameService object and 
makes changes based on the state read
+// and what is in the NameService.Spec
+// TODO(user): Modify this Reconcile function to implement your Controller 
logic.  This example creates
+// a Pod as an example
+// Note:
+// The Controller will requeue the Request to be processed again if the 
returned error is non-nil or
+// Result.Requeue is true, otherwise upon completion it will remove the work 
from the queue.
+func (r *ReconcileRocketmq) Reconcile(request reconcile.Request) 
(reconcile.Result, error) {
+       reqLogger := log.WithValues("Request.Namespace", request.Namespace, 
"Request.Name", request.Name)
+       reqLogger.Info("Reconciling Rocketmq")
+       // Fetch the Rocketmq instance
+       instance := &rocketmqv1alpha1.Rocketmq{}
+       err := r.client.Get(context.TODO(), request.NamespacedName, instance)
+       if err != nil {
+               if errors.IsNotFound(err) {
+                       // 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
+                       return reconcile.Result{}, nil
+               }
+               // Error reading the object, requeue the request
+               return reconcile.Result{}, err
+       }
+
+       // Check if nameserver already exist, if not create a new one
+       nameServiceFound := &rocketmqv1alpha1.NameService{}
+       nameServiceSts := r.nameServiceForRocketmq(instance)
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
nameServiceSts.Name, Namespace: nameServiceSts.Namespace}, nameServiceFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), nameServiceSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new nameService 
of rocketmq", "nameservice.namespace",
+                               nameServiceSts.Namespace, "nameservice.Name", 
nameServiceSts.Name)
+                       return reconcile.Result{}, err
+               }
+       } else if err != nil {
+               reqLogger.Error(err, "Failed to get rocketmq nameservice.")
+               return reconcile.Result{}, err
+       } else {
+               // Resource NameService will change; Only size nameServiceImage 
imagePullPolicy can update
+               if !reflect.DeepEqual(nameServiceSts.Spec.Size, 
nameServiceFound.Spec.Size) ||
+                       
!reflect.DeepEqual(nameServiceSts.Spec.NameServiceImage, 
nameServiceFound.Spec.NameServiceImage) ||
+                       !reflect.DeepEqual(nameServiceSts.Spec.ImagePullPolicy, 
nameServiceFound.Spec.ImagePullPolicy) {
+                       nameServiceFound.Spec.ImagePullPolicy = 
nameServiceSts.Spec.ImagePullPolicy
+                       nameServiceFound.Spec.NameServiceImage = 
nameServiceSts.Spec.NameServiceImage
+                       nameServiceFound.Spec.Size = nameServiceSts.Spec.Size
+                       err = r.client.Update(context.TODO(), nameServiceFound)
+                       if err != nil {
+                               reqLogger.Error(err, "should update nameservice 
resource", "spec.nameService",
+                                       nameServiceSts.Spec, 
"nameServiceFound.spec", nameServiceFound.Spec)
+                               return reconcile.Result{}, err
+                       }
+               }
+       }
+
+       // Check if broker already exists, if not create a new one
+       brokerFound := &rocketmqv1alpha1.Broker{}
+       brokerSts := r.brokerForRocketmq(instance)
+
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
brokerSts.Name, Namespace: brokerSts.Namespace}, brokerFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), brokerSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new broker of 
rocketmq", "broker.namespace",
+                               brokerSts.Namespace, "broker.Name", 
brokerSts.Name)
+               }
+               return reconcile.Result{Requeue: true}, nil
+       } else if err != nil {
+               reqLogger.Error(err, "Failed to get rocketmq broker.")
+       } else {
+               // Resource broker will change; Only ReplicaPerGroup Size 
ImagePullPolicy BrokerImage can update
+               if !reflect.DeepEqual(brokerSts.Spec.ReplicaPerGroup, 
brokerFound.Spec.ReplicaPerGroup) ||
+                       !reflect.DeepEqual(brokerSts.Spec.Size, 
brokerFound.Spec.Size) ||
+                       !reflect.DeepEqual(brokerSts.Spec.ImagePullPolicy, 
brokerFound.Spec.ImagePullPolicy) ||
+                       !reflect.DeepEqual(brokerSts.Spec.BrokerImage, 
brokerFound.Spec.BrokerImage) {
+                       brokerFound.Spec.ReplicaPerGroup = 
brokerSts.Spec.ReplicaPerGroup
+                       brokerFound.Spec.Size = brokerSts.Spec.Size
+                       brokerFound.Spec.ImagePullPolicy = 
brokerSts.Spec.ImagePullPolicy
+                       brokerFound.Spec.BrokerImage = 
brokerSts.Spec.BrokerImage
+                       err = r.client.Update(context.TODO(), brokerFound)
+                       if err != nil {
+                               reqLogger.Error(err, "should update broker 
resource", "spec.Broker",
+                                       brokerSts.Spec, "brokerFound.spec", 
brokerFound.Spec)
+                               return reconcile.Result{}, err
+                       }
+               }
+       }
+       if instance.Spec.Console.ConsoleDeployment.Spec.Replicas != nil {
+               consoleFound := &rocketmqv1alpha1.Console{}
+               consoleDep := r.consoleForRocketmq(instance)
+               err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
consoleDep.Name, Namespace: consoleDep.Namespace}, consoleFound)
+               if err != nil && errors.IsNotFound(err) {
+                       err = r.client.Create(context.TODO(), consoleDep)
+                       if err != nil {
+                               reqLogger.Error(err, "Failed to create new 
console of rocketmq")
+                       }
+               } else if err != nil {
+                       reqLogger.Error(err, "Failed to get rocketmq console.")
+               } else {
+                       if 
!reflect.DeepEqual(consoleDep.Spec.ConsoleDeployment.Spec.Replicas, 
consoleFound.Spec.ConsoleDeployment.Spec.Replicas) {
+                               err = r.client.Update(context.TODO(), 
consoleDep)
+                               if err != nil {
+                                       reqLogger.Error(err, "should update 
console resource")
+                               }
+                       }
+               }
+       }
+       // update instance status
+       if !reflect.DeepEqual(instance.Status.Broker, 
brokerFound.ObjectMeta.Name) ||
+               !reflect.DeepEqual(instance.Status.NameService, 
nameServiceFound.ObjectMeta.Name) {
+               instance.Status.Broker = brokerFound.ObjectMeta.Name

Review Comment:
   Status update references `brokerFound.ObjectMeta.Name` and 
`nameServiceFound.ObjectMeta.Name`, but when the broker or NameService was just 
created in this reconcile pass, `brokerFound`/`nameServiceFound` were never 
populated from the API server — they remain zero-value structs with empty 
`ObjectMeta.Name`. This writes empty strings to status. Use the names from the 
constructed objects (`brokerSts.Name`, `nameServiceSts.Name`) instead.



##########
pkg/controller/rocketmq/rocketmq_controller.go:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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 rocketmq contains the implementation of the Rocketmq CRD reconcile 
function
+package rocketmq
+
+import (
+       "context"
+       rocketmqv1alpha1 
"github.com/apache/rocketmq-operator/pkg/apis/rocketmq/v1alpha1"
+       cons "github.com/apache/rocketmq-operator/pkg/constants"
+       corev1 "k8s.io/api/core/v1"
+       "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/runtime"
+       "k8s.io/apimachinery/pkg/types"
+       "reflect"
+       "sigs.k8s.io/controller-runtime/pkg/client"
+       "sigs.k8s.io/controller-runtime/pkg/controller"
+       "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
+       "sigs.k8s.io/controller-runtime/pkg/handler"
+       "sigs.k8s.io/controller-runtime/pkg/manager"
+       "sigs.k8s.io/controller-runtime/pkg/reconcile"
+       logf "sigs.k8s.io/controller-runtime/pkg/runtime/log"
+       "sigs.k8s.io/controller-runtime/pkg/source"
+       "time"
+)
+
+var log = logf.Log.WithName("controller_rocketmq")
+
+// Add creates a new Rocketmq Controller and adds it to the Manager. The 
Manager will set fields on the Controller
+// and Start it when the Manager is Started.
+func Add(mgr manager.Manager) error {
+       return add(mgr, newReconciler(mgr))
+}
+
+// newReconciler returns a new reconcile.Reconciler
+func newReconciler(mgr manager.Manager) reconcile.Reconciler {
+       return &ReconcileRocketmq{client: mgr.GetClient(), scheme: 
mgr.GetScheme()}
+}
+
+// add adds a new Controller to mgr with r as the reconcile.Reconciler
+func add(mgr manager.Manager, r reconcile.Reconciler) error {
+       c, err := controller.New("rocketmq-controller", mgr, 
controller.Options{Reconciler: r})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to primary resource Rocketmq
+       err = c.Watch(&source.Kind{Type: &rocketmqv1alpha1.Rocketmq{}}, 
&handler.EnqueueRequestForObject{})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to secondary resource Pods and requeue the owner 
Rocketmq
+       err = c.Watch(&source.Kind{Type: &corev1.Pod{}}, 
&handler.EnqueueRequestForOwner{
+               IsController: true,
+               OwnerType: &rocketmqv1alpha1.Rocketmq{},
+       })
+       if err != nil {
+               return err
+       }
+       return nil
+}
+
+// blank assignment to verify that ReconcileRocketmq implements 
reconcile.Reconciler
+var _ reconcile.Reconciler = &ReconcileRocketmq{}
+
+// ReconcileRocketmq reconciles a Rocketmq object
+type ReconcileRocketmq struct {
+       // This client, initialized using mgr.Client() above, is a split client
+       // that reads objects from the cache and writes to the apiserver
+       client client.Client
+       scheme *runtime.Scheme
+}
+
+// Reconcile reads that state of the cluster for a NameService object and 
makes changes based on the state read
+// and what is in the NameService.Spec
+// TODO(user): Modify this Reconcile function to implement your Controller 
logic.  This example creates
+// a Pod as an example
+// Note:
+// The Controller will requeue the Request to be processed again if the 
returned error is non-nil or
+// Result.Requeue is true, otherwise upon completion it will remove the work 
from the queue.
+func (r *ReconcileRocketmq) Reconcile(request reconcile.Request) 
(reconcile.Result, error) {
+       reqLogger := log.WithValues("Request.Namespace", request.Namespace, 
"Request.Name", request.Name)
+       reqLogger.Info("Reconciling Rocketmq")
+       // Fetch the Rocketmq instance
+       instance := &rocketmqv1alpha1.Rocketmq{}
+       err := r.client.Get(context.TODO(), request.NamespacedName, instance)
+       if err != nil {
+               if errors.IsNotFound(err) {
+                       // 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
+                       return reconcile.Result{}, nil
+               }
+               // Error reading the object, requeue the request
+               return reconcile.Result{}, err
+       }
+
+       // Check if nameserver already exist, if not create a new one
+       nameServiceFound := &rocketmqv1alpha1.NameService{}
+       nameServiceSts := r.nameServiceForRocketmq(instance)
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
nameServiceSts.Name, Namespace: nameServiceSts.Namespace}, nameServiceFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), nameServiceSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new nameService 
of rocketmq", "nameservice.namespace",
+                               nameServiceSts.Namespace, "nameservice.Name", 
nameServiceSts.Name)
+                       return reconcile.Result{}, err
+               }
+       } else if err != nil {
+               reqLogger.Error(err, "Failed to get rocketmq nameservice.")
+               return reconcile.Result{}, err
+       } else {
+               // Resource NameService will change; Only size nameServiceImage 
imagePullPolicy can update
+               if !reflect.DeepEqual(nameServiceSts.Spec.Size, 
nameServiceFound.Spec.Size) ||
+                       
!reflect.DeepEqual(nameServiceSts.Spec.NameServiceImage, 
nameServiceFound.Spec.NameServiceImage) ||
+                       !reflect.DeepEqual(nameServiceSts.Spec.ImagePullPolicy, 
nameServiceFound.Spec.ImagePullPolicy) {
+                       nameServiceFound.Spec.ImagePullPolicy = 
nameServiceSts.Spec.ImagePullPolicy
+                       nameServiceFound.Spec.NameServiceImage = 
nameServiceSts.Spec.NameServiceImage
+                       nameServiceFound.Spec.Size = nameServiceSts.Spec.Size
+                       err = r.client.Update(context.TODO(), nameServiceFound)
+                       if err != nil {
+                               reqLogger.Error(err, "should update nameservice 
resource", "spec.nameService",
+                                       nameServiceSts.Spec, 
"nameServiceFound.spec", nameServiceFound.Spec)
+                               return reconcile.Result{}, err
+                       }
+               }
+       }
+
+       // Check if broker already exists, if not create a new one
+       brokerFound := &rocketmqv1alpha1.Broker{}
+       brokerSts := r.brokerForRocketmq(instance)
+
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
brokerSts.Name, Namespace: brokerSts.Namespace}, brokerFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), brokerSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new broker of 
rocketmq", "broker.namespace",
+                               brokerSts.Namespace, "broker.Name", 
brokerSts.Name)
+               }
+               return reconcile.Result{Requeue: true}, nil
+       } else if err != nil {
+               reqLogger.Error(err, "Failed to get rocketmq broker.")
+       } else {
+               // Resource broker will change; Only ReplicaPerGroup Size 
ImagePullPolicy BrokerImage can update
+               if !reflect.DeepEqual(brokerSts.Spec.ReplicaPerGroup, 
brokerFound.Spec.ReplicaPerGroup) ||
+                       !reflect.DeepEqual(brokerSts.Spec.Size, 
brokerFound.Spec.Size) ||
+                       !reflect.DeepEqual(brokerSts.Spec.ImagePullPolicy, 
brokerFound.Spec.ImagePullPolicy) ||
+                       !reflect.DeepEqual(brokerSts.Spec.BrokerImage, 
brokerFound.Spec.BrokerImage) {
+                       brokerFound.Spec.ReplicaPerGroup = 
brokerSts.Spec.ReplicaPerGroup
+                       brokerFound.Spec.Size = brokerSts.Spec.Size
+                       brokerFound.Spec.ImagePullPolicy = 
brokerSts.Spec.ImagePullPolicy
+                       brokerFound.Spec.BrokerImage = 
brokerSts.Spec.BrokerImage
+                       err = r.client.Update(context.TODO(), brokerFound)
+                       if err != nil {
+                               reqLogger.Error(err, "should update broker 
resource", "spec.Broker",
+                                       brokerSts.Spec, "brokerFound.spec", 
brokerFound.Spec)
+                               return reconcile.Result{}, err
+                       }
+               }
+       }
+       if instance.Spec.Console.ConsoleDeployment.Spec.Replicas != nil {
+               consoleFound := &rocketmqv1alpha1.Console{}
+               consoleDep := r.consoleForRocketmq(instance)
+               err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
consoleDep.Name, Namespace: consoleDep.Namespace}, consoleFound)
+               if err != nil && errors.IsNotFound(err) {
+                       err = r.client.Create(context.TODO(), consoleDep)
+                       if err != nil {
+                               reqLogger.Error(err, "Failed to create new 
console of rocketmq")
+                       }
+               } else if err != nil {
+                       reqLogger.Error(err, "Failed to get rocketmq console.")

Review Comment:
   Console update calls `r.client.Update(context.TODO(), consoleDep)` with the 
freshly constructed desired object instead of the modified `consoleFound`. This 
means the update sends an object without a ResourceVersion, which the API 
server will reject. Should modify `consoleFound.Spec` and update 
`consoleFound`, matching the pattern used for NameService and Broker updates 
above.



##########
pkg/apis/rocketmq/v1alpha1/broker_types.go:
##########
@@ -33,7 +33,7 @@ type BrokerSpec struct {
        // Add custom validation using kubebuilder tags: 
https://book-v1.book.kubebuilder.io/beyond_basics/generating_crd.html
        Size int `json:"size"`
        // NameServers defines the name service list e.g. 
192.168.1.1:9876;192.168.1.2:9876
-       NameServers string `json:"nameServers,omitempty"`
+       NameServers string `json:"nameServers"`

Review Comment:
   Removing `omitempty` from the `NameServers` json tag makes this field 
required in serialization. This is a breaking change for existing `Broker` CRs 
that relied on the field being optional. Any existing Broker resource without 
`nameServers` set will now fail validation on update. Keep `omitempty` or add a 
migration note.



##########
pkg/controller/console/console_controller.go:
##########
@@ -211,6 +220,6 @@ func newDeploymentForCR(cr *rocketmqv1alpha1.Console) 
*appsv1.Deployment {
                        },
                },
        }
-
+       controllerutil.SetControllerReference(cr, dep, r.scheme)

Review Comment:
   `SetControllerReference` is called here inside `newDeploymentForCR`, and 
also in the caller at line 139 of the Reconcile function. The second call will 
overwrite the first (or error if the owner is already set to a different UID). 
Remove one of the duplicate calls.



##########
pkg/controller/console/console_controller.go:
##########
@@ -178,11 +178,20 @@ func (r *ReconcileConsole) Reconcile(request 
reconcile.Request) (reconcile.Resul
 }
 
 // newDeploymentForCR returns a deployment pod with modifying the ENV
-func newDeploymentForCR(cr *rocketmqv1alpha1.Console) *appsv1.Deployment {
+func (r *ReconcileConsole) newDeploymentForCR(cr *rocketmqv1alpha1.Console) 
*appsv1.Deployment {
        env := corev1.EnvVar{
                Name:  "JAVA_OPTS",
                Value: fmt.Sprintf("-Drocketmq.namesrv.addr=%s 
-Dcom.rocketmq.sendMessageWithVIPChannel=false", share.NameServersStr),
        }

Review Comment:
   The IIFE directly mutates 
`cr.Spec.ConsoleDeployment.Spec.Selector.MatchLabels` and 
`cr.Spec.ConsoleDeployment.Spec.Template.Labels` — these are references to the 
CR's underlying maps, not copies. Adding the `console-cr` key modifies the 
in-memory CR object, which can cause unexpected behavior if the object is read 
from cache later. Copy the maps before mutating.



##########
deploy/crds/rocketmq_v1alpha1_rocketmq_crd.yaml:
##########
@@ -0,0 +1,65 @@
+apiVersion: apiextensions.k8s.io/v1beta1
+kind: CustomResourceDefinition
+metadata:
+  name: rocketmqs.rocketmq.apache.org
+spec:
+  group: rocketmq.apache.org
+  names:
+    kind: Rocketmq
+    listKind: RocketmqList
+    plural: rocketmqs
+    singular: rocketmq
+  scope: Namespaced
+  subresources:
+    status: {}
+  validation:
+    openAPIV3Schema:
+      properties:
+        apiVersion:
+          description: "APIVersion defines the versioned schema of this 
representation
+                                    of an object. Servers should convert 
recognized schemas to the latest
+                                    internal value, and may reject 
unrecognized values. More info: 
https://git.k8s.io/community/contributors/devel/api-conventions.md#resources";
+          type: string
+        kind:
+          description: 'Kind is a string value representing the REST resource 
this
+                                   object represents. Servers may infer this 
from the endpoint the client
+                                   submits requests to. Cannot be updated. In 
CamelCase. More info: 
https://git.k8s.io/community/contributors/devel/api-conventions.md#types-kinds'
+          type: string
+        metadata:
+          type: object
+        spec:
+          properties:
+            broker:
+              description: Broker defines rocketmq broker spec info
+              type: object
+            nameService:
+              description: NameService defines rocketmq name service spec info
+              type: object

Review Comment:
   The `broker`, `nameService`, and `console` fields under `spec` are declared 
as `type: object` with no nested properties or validation. This means any 
arbitrary YAML is accepted under these fields with no schema enforcement. Users 
will get no validation feedback for misspelled or missing required fields. 
Either add property-level validation or use 
`x-kubernetes-preserve-unknown-fields` explicitly to document the intent.



##########
pkg/apis/rocketmq/v1alpha1/rocketmq_types.go:
##########
@@ -0,0 +1,70 @@
+/*
+ * 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 v1alpha1
+
+import metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+
+// EDIT THIS FILE!  THIS IS SCAFFOLDING FOR YOU TO OWN!
+// NOTE: json tags are required.  Any new fields you add must have json tags 
for the fields to be serialized.
+
+// RocketmqSpec defines the desired state of Rocketmq
+// +k8s:openapi-gen=true
+// +kubebuilder:subresource:status
+type RocketmqSpec struct {
+       // Broker defines broker info
+       Broker BrokerSpec `json:"broker"`
+       // NameService defines name Service info
+       NameService NameServiceSpec `json:"nameService"`
+       // Console defines console info
+       Console     ConsoleSpec `json:"console,omitempty"`
+}
+
+// RocketmqStatus defines the observed state of Rocketmq
+// +k8s:openapi-gen=true
+type RocketmqStatus struct {
+       // INSERT ADDITIONAL STATUS FIELD - define observed state of cluster
+       // Important: Run "operator-sdk generate k8s" to regenerate code after 
modifying this file
+       // Add custom validation using kubebuilder tags: 
https://book-v1.book.kubebuilder.io/beyond_basics/generating_crd.html
+       Broker string `json:"broker"`

Review Comment:
   `RocketmqStatus` fields `Broker` and `NameService` are just bare strings 
with no `omitempty`. The CRD YAML marks them as `required` in status. On 
initial creation before the first status update, these will be empty strings, 
which may confuse consumers expecting meaningful values. Consider adding 
`omitempty` and making them optional in the CRD schema.



##########
deploy/crds/rocketmq_v1alpha1_rocketmq_crd.yaml:
##########
@@ -0,0 +1,65 @@
+apiVersion: apiextensions.k8s.io/v1beta1

Review Comment:
   CRD uses `apiextensions.k8s.io/v1beta1` which has been deprecated since 
Kubernetes 1.16 and removed in 1.22. Consider migrating to 
`apiextensions.k8s.io/v1` for forward compatibility.



##########
pkg/controller/rocketmq/rocketmq_controller.go:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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 rocketmq contains the implementation of the Rocketmq CRD reconcile 
function
+package rocketmq
+
+import (
+       "context"
+       rocketmqv1alpha1 
"github.com/apache/rocketmq-operator/pkg/apis/rocketmq/v1alpha1"
+       cons "github.com/apache/rocketmq-operator/pkg/constants"
+       corev1 "k8s.io/api/core/v1"
+       "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/runtime"
+       "k8s.io/apimachinery/pkg/types"
+       "reflect"
+       "sigs.k8s.io/controller-runtime/pkg/client"
+       "sigs.k8s.io/controller-runtime/pkg/controller"
+       "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
+       "sigs.k8s.io/controller-runtime/pkg/handler"
+       "sigs.k8s.io/controller-runtime/pkg/manager"
+       "sigs.k8s.io/controller-runtime/pkg/reconcile"
+       logf "sigs.k8s.io/controller-runtime/pkg/runtime/log"
+       "sigs.k8s.io/controller-runtime/pkg/source"
+       "time"
+)
+
+var log = logf.Log.WithName("controller_rocketmq")
+
+// Add creates a new Rocketmq Controller and adds it to the Manager. The 
Manager will set fields on the Controller
+// and Start it when the Manager is Started.
+func Add(mgr manager.Manager) error {
+       return add(mgr, newReconciler(mgr))
+}
+
+// newReconciler returns a new reconcile.Reconciler
+func newReconciler(mgr manager.Manager) reconcile.Reconciler {
+       return &ReconcileRocketmq{client: mgr.GetClient(), scheme: 
mgr.GetScheme()}
+}
+
+// add adds a new Controller to mgr with r as the reconcile.Reconciler
+func add(mgr manager.Manager, r reconcile.Reconciler) error {
+       c, err := controller.New("rocketmq-controller", mgr, 
controller.Options{Reconciler: r})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to primary resource Rocketmq
+       err = c.Watch(&source.Kind{Type: &rocketmqv1alpha1.Rocketmq{}}, 
&handler.EnqueueRequestForObject{})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to secondary resource Pods and requeue the owner 
Rocketmq
+       err = c.Watch(&source.Kind{Type: &corev1.Pod{}}, 
&handler.EnqueueRequestForOwner{
+               IsController: true,
+               OwnerType: &rocketmqv1alpha1.Rocketmq{},
+       })
+       if err != nil {
+               return err
+       }
+       return nil
+}
+
+// blank assignment to verify that ReconcileRocketmq implements 
reconcile.Reconciler
+var _ reconcile.Reconciler = &ReconcileRocketmq{}
+
+// ReconcileRocketmq reconciles a Rocketmq object
+type ReconcileRocketmq struct {
+       // This client, initialized using mgr.Client() above, is a split client
+       // that reads objects from the cache and writes to the apiserver
+       client client.Client
+       scheme *runtime.Scheme
+}
+
+// Reconcile reads that state of the cluster for a NameService object and 
makes changes based on the state read
+// and what is in the NameService.Spec
+// TODO(user): Modify this Reconcile function to implement your Controller 
logic.  This example creates
+// a Pod as an example
+// Note:
+// The Controller will requeue the Request to be processed again if the 
returned error is non-nil or
+// Result.Requeue is true, otherwise upon completion it will remove the work 
from the queue.
+func (r *ReconcileRocketmq) Reconcile(request reconcile.Request) 
(reconcile.Result, error) {
+       reqLogger := log.WithValues("Request.Namespace", request.Namespace, 
"Request.Name", request.Name)
+       reqLogger.Info("Reconciling Rocketmq")

Review Comment:
   No tests are provided for this new controller. Given the reconciliation 
complexity (orchestrating three sub-resources with ordering constraints), at 
minimum there should be unit tests for `nameServiceForRocketmq`, 
`brokerForRocketmq`, and `consoleForRocketmq` builders, and integration tests 
for the reconcile loop covering create, update, and error paths.



##########
pkg/controller/rocketmq/rocketmq_controller.go:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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 rocketmq contains the implementation of the Rocketmq CRD reconcile 
function
+package rocketmq
+
+import (
+       "context"
+       rocketmqv1alpha1 
"github.com/apache/rocketmq-operator/pkg/apis/rocketmq/v1alpha1"
+       cons "github.com/apache/rocketmq-operator/pkg/constants"
+       corev1 "k8s.io/api/core/v1"
+       "k8s.io/apimachinery/pkg/api/errors"
+       metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/runtime"
+       "k8s.io/apimachinery/pkg/types"
+       "reflect"
+       "sigs.k8s.io/controller-runtime/pkg/client"
+       "sigs.k8s.io/controller-runtime/pkg/controller"
+       "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
+       "sigs.k8s.io/controller-runtime/pkg/handler"
+       "sigs.k8s.io/controller-runtime/pkg/manager"
+       "sigs.k8s.io/controller-runtime/pkg/reconcile"
+       logf "sigs.k8s.io/controller-runtime/pkg/runtime/log"
+       "sigs.k8s.io/controller-runtime/pkg/source"
+       "time"
+)
+
+var log = logf.Log.WithName("controller_rocketmq")
+
+// Add creates a new Rocketmq Controller and adds it to the Manager. The 
Manager will set fields on the Controller
+// and Start it when the Manager is Started.
+func Add(mgr manager.Manager) error {
+       return add(mgr, newReconciler(mgr))
+}
+
+// newReconciler returns a new reconcile.Reconciler
+func newReconciler(mgr manager.Manager) reconcile.Reconciler {
+       return &ReconcileRocketmq{client: mgr.GetClient(), scheme: 
mgr.GetScheme()}
+}
+
+// add adds a new Controller to mgr with r as the reconcile.Reconciler
+func add(mgr manager.Manager, r reconcile.Reconciler) error {
+       c, err := controller.New("rocketmq-controller", mgr, 
controller.Options{Reconciler: r})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to primary resource Rocketmq
+       err = c.Watch(&source.Kind{Type: &rocketmqv1alpha1.Rocketmq{}}, 
&handler.EnqueueRequestForObject{})
+       if err != nil {
+               return err
+       }
+       // Watch for changes to secondary resource Pods and requeue the owner 
Rocketmq
+       err = c.Watch(&source.Kind{Type: &corev1.Pod{}}, 
&handler.EnqueueRequestForOwner{
+               IsController: true,
+               OwnerType: &rocketmqv1alpha1.Rocketmq{},
+       })
+       if err != nil {
+               return err
+       }
+       return nil
+}
+
+// blank assignment to verify that ReconcileRocketmq implements 
reconcile.Reconciler
+var _ reconcile.Reconciler = &ReconcileRocketmq{}
+
+// ReconcileRocketmq reconciles a Rocketmq object
+type ReconcileRocketmq struct {
+       // This client, initialized using mgr.Client() above, is a split client
+       // that reads objects from the cache and writes to the apiserver
+       client client.Client
+       scheme *runtime.Scheme
+}
+
+// Reconcile reads that state of the cluster for a NameService object and 
makes changes based on the state read
+// and what is in the NameService.Spec
+// TODO(user): Modify this Reconcile function to implement your Controller 
logic.  This example creates
+// a Pod as an example
+// Note:
+// The Controller will requeue the Request to be processed again if the 
returned error is non-nil or
+// Result.Requeue is true, otherwise upon completion it will remove the work 
from the queue.
+func (r *ReconcileRocketmq) Reconcile(request reconcile.Request) 
(reconcile.Result, error) {
+       reqLogger := log.WithValues("Request.Namespace", request.Namespace, 
"Request.Name", request.Name)
+       reqLogger.Info("Reconciling Rocketmq")
+       // Fetch the Rocketmq instance
+       instance := &rocketmqv1alpha1.Rocketmq{}
+       err := r.client.Get(context.TODO(), request.NamespacedName, instance)
+       if err != nil {
+               if errors.IsNotFound(err) {
+                       // 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
+                       return reconcile.Result{}, nil
+               }
+               // Error reading the object, requeue the request
+               return reconcile.Result{}, err
+       }
+
+       // Check if nameserver already exist, if not create a new one
+       nameServiceFound := &rocketmqv1alpha1.NameService{}
+       nameServiceSts := r.nameServiceForRocketmq(instance)
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
nameServiceSts.Name, Namespace: nameServiceSts.Namespace}, nameServiceFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), nameServiceSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new nameService 
of rocketmq", "nameservice.namespace",
+                               nameServiceSts.Namespace, "nameservice.Name", 
nameServiceSts.Name)
+                       return reconcile.Result{}, err
+               }
+       } else if err != nil {
+               reqLogger.Error(err, "Failed to get rocketmq nameservice.")
+               return reconcile.Result{}, err
+       } else {
+               // Resource NameService will change; Only size nameServiceImage 
imagePullPolicy can update
+               if !reflect.DeepEqual(nameServiceSts.Spec.Size, 
nameServiceFound.Spec.Size) ||
+                       
!reflect.DeepEqual(nameServiceSts.Spec.NameServiceImage, 
nameServiceFound.Spec.NameServiceImage) ||
+                       !reflect.DeepEqual(nameServiceSts.Spec.ImagePullPolicy, 
nameServiceFound.Spec.ImagePullPolicy) {
+                       nameServiceFound.Spec.ImagePullPolicy = 
nameServiceSts.Spec.ImagePullPolicy
+                       nameServiceFound.Spec.NameServiceImage = 
nameServiceSts.Spec.NameServiceImage
+                       nameServiceFound.Spec.Size = nameServiceSts.Spec.Size
+                       err = r.client.Update(context.TODO(), nameServiceFound)
+                       if err != nil {
+                               reqLogger.Error(err, "should update nameservice 
resource", "spec.nameService",
+                                       nameServiceSts.Spec, 
"nameServiceFound.spec", nameServiceFound.Spec)
+                               return reconcile.Result{}, err
+                       }
+               }
+       }
+
+       // Check if broker already exists, if not create a new one
+       brokerFound := &rocketmqv1alpha1.Broker{}
+       brokerSts := r.brokerForRocketmq(instance)
+
+       err = r.client.Get(context.TODO(), types.NamespacedName{Name: 
brokerSts.Name, Namespace: brokerSts.Namespace}, brokerFound)
+       if err != nil && errors.IsNotFound(err) {
+               err = r.client.Create(context.TODO(), brokerSts)
+               if err != nil {
+                       reqLogger.Error(err, "Failed to create new broker of 
rocketmq", "broker.namespace",
+                               brokerSts.Namespace, "broker.Name", 
brokerSts.Name)
+               }
+               return reconcile.Result{Requeue: true}, nil
+       } else if err != nil {

Review Comment:
   After a successful broker creation, the code unconditionally returns 
`reconcile.Result{Requeue: true}, nil` without waiting for the broker's 
sub-controller to reconcile. This means the broker update block (line 160) and 
console block (line 176) are unreachable on the creation pass. Consider falling 
through or requeueing with a delay to allow the broker controller to act.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to