Skip to content
Merged
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
100 changes: 60 additions & 40 deletions pkg/controller/disaggregated_cluster_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,12 +61,9 @@ type DisaggregatedClusterReconciler struct {
Recorder record.EventRecorder
Scheme *runtime.Scheme
Scs map[string]sc.DisaggregatedSubController
//record configmap response instance. key: configMap namespacedName, value: DorisDisaggregatedCluster namespacedName
//wcms map[string]string
}

func (dc *DisaggregatedClusterReconciler) Init(mgr ctrl.Manager, options *Options) {
//wcms := make(map[string]string)
scs := make(map[string]sc.DisaggregatedSubController)
msc := metaservice.New(mgr)
scs[msc.GetControllerName()] = msc
Expand All @@ -80,7 +77,6 @@ func (dc *DisaggregatedClusterReconciler) Init(mgr ctrl.Manager, options *Option
Client: mgr.GetClient(),
Recorder: mgr.GetEventRecorderFor(disaggregatedClusterController),
Scs: scs,
//wcms: wcms,
}).SetupWithManager(mgr); err != nil {
klog.Error(err, "unable to create controller ", "disaggregatedClusterReconciler")
os.Exit(1)
Expand All @@ -97,7 +93,7 @@ func (dc *DisaggregatedClusterReconciler) Init(mgr ctrl.Manager, options *Option
func (dc *DisaggregatedClusterReconciler) SetupWithManager(mgr ctrl.Manager) error {
builder := dc.resourceBuilder(ctrl.NewControllerManagedBy(mgr))
builder = dc.watchPodBuilder(builder)
//builder = dc.watchConfigMapBuilder(builder)
builder = dc.watchFDBConfigMapBuilder(builder)
return builder.Complete(dc)
}

Expand Down Expand Up @@ -146,41 +142,65 @@ func (dc *DisaggregatedClusterReconciler) watchPodBuilder(builder *ctrl.Builder)
mapFn, controller_builder.WithPredicates(p))
}

//func (dc *DisaggregatedClusterReconciler) watchConfigMapBuilder(builder *ctrl.Builder) *ctrl.Builder {
// mapFn := handler.EnqueueRequestsFromMapFunc(
// func(a client.Object) []reconcile.Request {
// namespace := a.GetNamespace()
// name := a.GetName()
// cmnn := types.NamespacedName{Namespace: namespace, Name: name}
// cmnnStr := cmnn.String()
// if ddc, ok := dc.wcms[cmnnStr]; ok {
// nna := strings.Split(ddc, "/")
// // not run only for code standard
// if len(nna) != 2 {
// return nil
// }
//
// return []reconcile.Request{{NamespacedName: types.NamespacedName{
// Namespace: nna[0],
// Name: nna[1],
// }}}
// }
// return nil
// })
//
// p := predicate.Funcs{
// UpdateFunc: func(u event.UpdateEvent) bool {
// ns := u.ObjectNew.GetNamespace()
// name := u.ObjectNew.GetName()
// nsn := ns + "/" + name
// _, ok := dc.wcms[nsn]
// return ok
// },
// }
//
// return builder.Watches(&source.Kind{Type: &corev1.ConfigMap{}},
// mapFn, controller_builder.WithPredicates(p))
//}
func (dc *DisaggregatedClusterReconciler) watchFDBConfigMapBuilder(builder *ctrl.Builder) *ctrl.Builder {
mapFn := handler.EnqueueRequestsFromMapFunc(dc.mapFDBConfigMapToDDCs)
return builder.Watches(&corev1.ConfigMap{}, mapFn, controller_builder.WithPredicates(fdbConfigMapPredicate()))
}

func fdbConfigMapPredicate() predicate.Predicate {
return predicate.Funcs{
CreateFunc: func(e event.CreateEvent) bool {
_, ok := fdbClusterFile(e.Object)
return ok
},
UpdateFunc: func(e event.UpdateEvent) bool {
oldValue, oldOK := fdbClusterFile(e.ObjectOld)
newValue, newOK := fdbClusterFile(e.ObjectNew)
return oldOK != newOK || oldValue != newValue
},
DeleteFunc: func(e event.DeleteEvent) bool {
_, ok := fdbClusterFile(e.Object)
return ok
},
GenericFunc: func(event.GenericEvent) bool {
return false
},
}
}

func (dc *DisaggregatedClusterReconciler) mapFDBConfigMapToDDCs(ctx context.Context, obj client.Object) []reconcile.Request {
var ddcList dv1.DorisDisaggregatedClusterList
if err := dc.List(ctx, &ddcList); err != nil {
klog.Errorf("list DorisDisaggregatedClusters for FDB ConfigMap %s/%s failed: %s", obj.GetNamespace(), obj.GetName(), err.Error())
return nil
}

requests := make([]reconcile.Request, 0)
for i := range ddcList.Items {
ddc := &ddcList.Items[i]
if ddc.Spec.MetaService.FDB.Address != "" {
continue
}
ref := ddc.Spec.MetaService.FDB.ConfigMapNamespaceName
if ref.Namespace == obj.GetNamespace() && ref.Name == obj.GetName() {
requests = append(requests, reconcile.Request{NamespacedName: types.NamespacedName{
Namespace: ddc.Namespace,
Name: ddc.Name,
}})
}
}

return requests
}

func fdbClusterFile(obj client.Object) (string, bool) {
cm, ok := obj.(*corev1.ConfigMap)
if !ok {
return "", false
}
value, ok := cm.Data[metaservice.FDBClusterFileKey]
return value, ok
}

func (dc *DisaggregatedClusterReconciler) resourceBuilder(builder *ctrl.Builder) *ctrl.Builder {
return builder.For(&dv1.DorisDisaggregatedCluster{}).
Expand Down
104 changes: 104 additions & 0 deletions pkg/controller/disaggregated_cluster_controller_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
// 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 controller

import (
"context"
"testing"

dv1 "github.com/apache/doris-operator/api/disaggregated/v1"
"github.com/apache/doris-operator/pkg/controller/sub_controller/disaggregated_cluster/metaservice"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
"sigs.k8s.io/controller-runtime/pkg/event"
)

func TestMapFDBConfigMapToDDCs(t *testing.T) {
scheme := runtime.NewScheme()
if err := dv1.AddToScheme(scheme); err != nil {
t.Fatalf("add DDC scheme: %v", err)
}

referencing := newTestDDC("doris-a", "cluster-a", "fdb-system", "fdb-a-config")
unrelated := newTestDDC("doris-b", "cluster-b", "fdb-system", "fdb-b-config")
directAddress := newTestDDC("doris-c", "cluster-c", "fdb-system", "fdb-a-config")
directAddress.Spec.MetaService.FDB.Address = "fdb:direct@127.0.0.1:4500"
reconciler := &DisaggregatedClusterReconciler{
Client: fake.NewClientBuilder().WithScheme(scheme).WithObjects(referencing, unrelated, directAddress).Build(),
}

requests := reconciler.mapFDBConfigMapToDDCs(context.Background(), &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{Namespace: "fdb-system", Name: "fdb-a-config"},
})

if len(requests) != 1 {
t.Fatalf("request count = %d, want 1", len(requests))
}
if requests[0].Namespace != referencing.Namespace || requests[0].Name != referencing.Name {
t.Fatalf("request = %s/%s, want %s/%s", requests[0].Namespace, requests[0].Name, referencing.Namespace, referencing.Name)
}
}

func TestFDBConfigMapPredicate(t *testing.T) {
p := fdbConfigMapPredicate()
oldConfigMap := newTestFDBConfigMap("old-cluster-file")
newConfigMap := newTestFDBConfigMap("new-cluster-file")
sameConfigMap := oldConfigMap.DeepCopy()
sameConfigMap.Annotations = map[string]string{"updated": "true"}
withoutClusterFile := oldConfigMap.DeepCopy()
delete(withoutClusterFile.Data, metaservice.FDBClusterFileKey)

if !p.Create(event.CreateEvent{Object: oldConfigMap}) {
t.Fatal("create event with cluster-file should pass")
}
if p.Update(event.UpdateEvent{ObjectOld: oldConfigMap, ObjectNew: sameConfigMap}) {
t.Fatal("metadata-only update should not pass")
}
if !p.Update(event.UpdateEvent{ObjectOld: oldConfigMap, ObjectNew: newConfigMap}) {
t.Fatal("cluster-file update should pass")
}
if !p.Update(event.UpdateEvent{ObjectOld: oldConfigMap, ObjectNew: withoutClusterFile}) {
t.Fatal("cluster-file removal should pass")
}
if !p.Delete(event.DeleteEvent{Object: oldConfigMap}) {
t.Fatal("delete event with cluster-file should pass")
}
}

func newTestDDC(namespace, name, fdbNamespace, fdbConfigMap string) *dv1.DorisDisaggregatedCluster {
return &dv1.DorisDisaggregatedCluster{
ObjectMeta: metav1.ObjectMeta{Namespace: namespace, Name: name},
Spec: dv1.DorisDisaggregatedClusterSpec{
MetaService: dv1.MetaService{
FDB: dv1.FDB{ConfigMapNamespaceName: dv1.NamespaceName{
Namespace: fdbNamespace,
Name: fdbConfigMap,
}},
},
},
}
}

func newTestFDBConfigMap(clusterFile string) *corev1.ConfigMap {
return &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{Namespace: "fdb-system", Name: "fdb-a-config"},
Data: map[string]string{metaservice.FDBClusterFileKey: clusterFile},
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -157,11 +157,16 @@ func New(mgr ctrl.Manager) *DisaggregatedMSController {

func (dms *DisaggregatedMSController) Sync(ctx context.Context, obj client.Object) error {
ddc := obj.(*v1.DorisDisaggregatedCluster)
fdbEndpoint, err := dms.resolveFDBEndpoint(ctx, ddc)
if err != nil {
dms.K8srecorder.Event(ddc, string(sc.EventWarning), string(sc.FDBAddressNotConfiged), err.Error())
return err
}
msSpec := ddc.Spec.MetaService
confMap := dms.GetConfigValuesFromConfigMaps(ddc.Namespace, resource.MS_RESOLVEKEY, msSpec.ConfigMaps)
svc := dms.newService(ddc, confMap)

st := dms.newStatefulset(ddc, confMap)
st := dms.newStatefulset(ddc, confMap, fdbEndpoint)
dms.initMSStatus(ddc)

dms.CheckSecretMountPath(ddc, ddc.Spec.MetaService.Secrets)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package metaservice

import (
"context"
"fmt"

v1 "github.com/apache/doris-operator/api/disaggregated/v1"
"github.com/apache/doris-operator/pkg/common/utils/k8s"
Expand All @@ -32,10 +33,12 @@ import (

const (
defaultLogPrefixName = "log"
fdbClusterFileKey = "cluster-file"
//DefaultStorageSize int64 = 107374182400
)

// FDBClusterFileKey is the key used by fdb-kubernetes-operator for the cluster file.
const FDBClusterFileKey = "cluster-file"

func (dms *DisaggregatedMSController) newMSPodsSelector(ddcName string) map[string]string {
return map[string]string{
v1.DorisDisaggregatedClusterName: ddcName,
Expand All @@ -51,7 +54,7 @@ func (dms *DisaggregatedMSController) newMSSchedulerLabels(ddcName string) map[s
}
}

func (dms *DisaggregatedMSController) newStatefulset(ddc *v1.DorisDisaggregatedCluster, confMap map[string]interface{}) *appv1.StatefulSet {
func (dms *DisaggregatedMSController) newStatefulset(ddc *v1.DorisDisaggregatedCluster, confMap map[string]interface{}, fdbEndpoint string) *appv1.StatefulSet {
st := dms.NewDefaultStatefulset(ddc)
func() {
st.Name = ddc.GetMSStatefulsetName()
Expand All @@ -71,15 +74,15 @@ func (dms *DisaggregatedMSController) newStatefulset(ddc *v1.DorisDisaggregatedC
st.Spec.Selector = &metav1.LabelSelector{
MatchLabels: matchLabels,
}
st.Spec.Template = dms.NewPodTemplateSpec(ddc, matchLabels, confMap)
st.Spec.Template = dms.NewPodTemplateSpec(ddc, matchLabels, confMap, fdbEndpoint)
st.Spec.ServiceName = ddc.GetMSServiceName()
st.Spec.VolumeClaimTemplates = vcts
}()

return st
}

func (dms *DisaggregatedMSController) NewPodTemplateSpec(ddc *v1.DorisDisaggregatedCluster, selector map[string]string, confMap map[string]interface{}) corev1.PodTemplateSpec {
func (dms *DisaggregatedMSController) NewPodTemplateSpec(ddc *v1.DorisDisaggregatedCluster, selector map[string]string, confMap map[string]interface{}, fdbEndpoint string) corev1.PodTemplateSpec {
pts := resource.NewPodTemplateSpecWithCommonSpec(false, &ddc.Spec.MetaService.CommonSpec, v1.DisaggregatedMS)
//pod template metadata.
func() {
Expand All @@ -88,7 +91,7 @@ func (dms *DisaggregatedMSController) NewPodTemplateSpec(ddc *v1.DorisDisaggrega
pts.Labels = l
}()

c := dms.NewMSContainer(ddc, confMap)
c := dms.NewMSContainer(ddc, confMap, fdbEndpoint)
pts.Spec.Containers = append(pts.Spec.Containers, c)
vs, _, _ := dms.BuildVolumesVolumeMountsAndPVCs(confMap, v1.DisaggregatedMS, &ddc.Spec.MetaService.CommonSpec)
configVolumes, _ := dms.BuildDefaultConfigMapVolumesVolumeMounts(ddc.Spec.MetaService.ConfigMaps)
Expand All @@ -104,7 +107,7 @@ func (dms *DisaggregatedMSController) NewPodTemplateSpec(ddc *v1.DorisDisaggrega
return pts
}

func (dms *DisaggregatedMSController) NewMSContainer(ddc *v1.DorisDisaggregatedCluster, cvs map[string]interface{}) corev1.Container {
func (dms *DisaggregatedMSController) NewMSContainer(ddc *v1.DorisDisaggregatedCluster, cvs map[string]interface{}, fdbEndpoint string) corev1.Container {
c := resource.NewContainerWithCommonSpec(&ddc.Spec.MetaService.CommonSpec)

c.Lifecycle = resource.LifeCycleWithPreStopScript(c.Lifecycle, sc.GetDisaggregatedPreStopScript(v1.DisaggregatedMS))
Expand All @@ -117,7 +120,7 @@ func (dms *DisaggregatedMSController) NewMSContainer(ddc *v1.DorisDisaggregatedC
c.Ports = resource.GetDisaggregatedContainerPorts(cvs, v1.DisaggregatedMS)
c.Env = ddc.Spec.MetaService.CommonSpec.EnvVars
c.Env = append(c.Env, resource.GetPodDefaultEnv()...)
c.Env = append(c.Env, dms.newSpecificEnvs(ddc)...)
c.Env = append(c.Env, dms.newSpecificEnvs(fdbEndpoint)...)
resource.BuildDisaggregatedProbe(&c, &ddc.Spec.MetaService.CommonSpec, v1.DisaggregatedMS)
_, vms, _ := dms.BuildVolumesVolumeMountsAndPVCs(cvs, v1.DisaggregatedMS, &ddc.Spec.MetaService.CommonSpec)
_, cmvms := dms.BuildDefaultConfigMapVolumesVolumeMounts(ddc.Spec.MetaService.ConfigMaps)
Expand All @@ -136,36 +139,29 @@ func (dms *DisaggregatedMSController) NewMSContainer(ddc *v1.DorisDisaggregatedC
return c
}

func (dms *DisaggregatedMSController) newSpecificEnvs(ddc *v1.DorisDisaggregatedCluster) []corev1.EnvVar {
func (dms *DisaggregatedMSController) resolveFDBEndpoint(ctx context.Context, ddc *v1.DorisDisaggregatedCluster) (string, error) {
msSpec := ddc.Spec.MetaService
if msSpec.FDB.Address == "" && (msSpec.FDB.ConfigMapNamespaceName.Namespace == "" || msSpec.FDB.ConfigMapNamespaceName.Name == "") {
dms.K8srecorder.Event(ddc, string(sc.EventWarning), string(sc.FDBAddressNotConfiged), "fdb not configed in spec")
return nil
if msSpec.FDB.Address != "" {
return msSpec.FDB.Address, nil
}

var fdbEndpoint string
if msSpec.FDB.ConfigMapNamespaceName.Namespace != "" && msSpec.FDB.ConfigMapNamespaceName.Name != "" {
cm, err := k8s.GetConfigMap(context.Background(), dms.K8sclient, msSpec.FDB.ConfigMapNamespaceName.Namespace, msSpec.FDB.ConfigMapNamespaceName.Name)
if err != nil {
dms.K8srecorder.Event(ddc, string(sc.EventWarning), string(sc.FDBAddressNotConfiged), "configmap "+"namespace"+msSpec.FDB.ConfigMapNamespaceName.Namespace+" name "+msSpec.FDB.ConfigMapNamespaceName.Name+" find failed "+err.Error())
return nil
}

if cm.Data == nil {
dms.K8srecorder.Event(ddc, string(sc.EventWarning), string(sc.FDBAddressNotConfiged), "configmap "+"namespace"+msSpec.FDB.ConfigMapNamespaceName.Namespace+" name "+msSpec.FDB.ConfigMapNamespaceName.Name+" not have data.")
return nil
}
ref := msSpec.FDB.ConfigMapNamespaceName
if ref.Namespace == "" || ref.Name == "" {
return "", fmt.Errorf("FDB address or ConfigMap reference is not configured")
}

if _, ok := cm.Data[fdbClusterFileKey]; !ok {
dms.K8srecorder.Event(ddc, string(sc.EventWarning), string(sc.FDBAddressNotConfiged), "configmap "+"namespace"+msSpec.FDB.ConfigMapNamespaceName.Namespace+" name "+msSpec.FDB.ConfigMapNamespaceName.Name+" not have cluster-file")
return nil
}
fdbEndpoint = cm.Data[fdbClusterFileKey]
cm, err := k8s.GetConfigMap(ctx, dms.K8sclient, ref.Namespace, ref.Name)
if err != nil {
return "", fmt.Errorf("get FDB ConfigMap %s/%s: %w", ref.Namespace, ref.Name, err)
}
if msSpec.FDB.Address != "" {
fdbEndpoint = msSpec.FDB.Address
fdbEndpoint, ok := cm.Data[FDBClusterFileKey]
if !ok || fdbEndpoint == "" {
return "", fmt.Errorf("FDB ConfigMap %s/%s does not contain a non-empty %q", ref.Namespace, ref.Name, FDBClusterFileKey)
}
return fdbEndpoint, nil
}

func (dms *DisaggregatedMSController) newSpecificEnvs(fdbEndpoint string) []corev1.EnvVar {
return []corev1.EnvVar{{
Name: resource.FDB_ENDPOINT,
Value: fdbEndpoint,
Expand Down
Loading