Skip to content

Commit f03dfa7

Browse files
authored
fix: align instanceset2 revision status (#10488)
1 parent 74318f5 commit f03dfa7

13 files changed

Lines changed: 550 additions & 126 deletions

cmd/manager/main.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@ import (
7373
"github.com/apecloud/kubeblocks/pkg/constant"
7474
"github.com/apecloud/kubeblocks/pkg/controller/instanceset"
7575
"github.com/apecloud/kubeblocks/pkg/controller/multicluster"
76+
"github.com/apecloud/kubeblocks/pkg/controller/revisionmap"
7677
intctrlutil "github.com/apecloud/kubeblocks/pkg/controllerutil"
7778
"github.com/apecloud/kubeblocks/pkg/metrics"
7879
viper "github.com/apecloud/kubeblocks/pkg/viperx"
@@ -149,7 +150,7 @@ func init() {
149150
viper.SetDefault(constant.CfgKeyClusterDefaultResources, `{"zero":true}`)
150151
viper.SetDefault(constant.CfgKeyOperationZeroResourceForUnset, true)
151152
viper.SetDefault(constant.KubernetesClusterDomainEnv, constant.DefaultDNSDomain)
152-
viper.SetDefault(instanceset.MaxPlainRevisionCount, 1024)
153+
viper.SetDefault(revisionmap.MaxPlainRevisionCount, 1024)
153154
viper.SetDefault(instanceset.FeatureGateIgnorePodVerticalScaling, false)
154155
viper.SetDefault(intctrlutil.FeatureGateEnableRuntimeMetrics, false)
155156
viper.SetDefault(constant.FeatureGateIgnoreConfigTemplateDefaultMode, false)

pkg/controller/instanceset/instance_util.go

Lines changed: 0 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -28,11 +28,9 @@ import (
2828
"strconv"
2929
"strings"
3030

31-
"github.com/klauspost/compress/zstd"
3231
appsv1 "k8s.io/api/apps/v1"
3332
corev1 "k8s.io/api/core/v1"
3433
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
35-
"k8s.io/apimachinery/pkg/util/runtime"
3634
"k8s.io/apimachinery/pkg/util/sets"
3735
"k8s.io/klog/v2"
3836
"k8s.io/utils/ptr"
@@ -61,19 +59,6 @@ type instanceTemplateExt struct {
6159
VolumeClaimTemplates []corev1.PersistentVolumeClaim
6260
}
6361

64-
var (
65-
reader *zstd.Decoder
66-
writer *zstd.Encoder
67-
)
68-
69-
func init() {
70-
var err error
71-
reader, err = zstd.NewReader(nil)
72-
runtime.Must(err)
73-
writer, err = zstd.NewWriter(nil)
74-
runtime.Must(err)
75-
}
76-
7762
// parseParentNameAndOrdinal parses parent (instance template) Name and ordinal from the give instance name.
7863
// -1 will be returned if no numeric suffix contained.
7964
func parseParentNameAndOrdinal(s string) (string, int) {

pkg/controller/instanceset/revision_util.go

Lines changed: 3 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -20,14 +20,12 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
2020
package instanceset
2121

2222
import (
23-
"encoding/base64"
2423
"encoding/json"
2524
"fmt"
2625
"hash"
2726
"hash/fnv"
2827
"strconv"
2928

30-
jsoniter "github.com/json-iterator/go"
3129
apps "k8s.io/api/apps/v1"
3230
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
3331
"k8s.io/apimachinery/pkg/runtime"
@@ -38,8 +36,7 @@ import (
3836

3937
workloads "github.com/apecloud/kubeblocks/apis/workloads/v1"
4038
"github.com/apecloud/kubeblocks/pkg/controller/model"
41-
"github.com/apecloud/kubeblocks/pkg/lru"
42-
viper "github.com/apecloud/kubeblocks/pkg/viperx"
39+
"github.com/apecloud/kubeblocks/pkg/controller/revisionmap"
4340
)
4441

4542
// controllerRevisionHashLabel is the label used to indicate the hash value of a controllerRevision's Data.
@@ -49,8 +46,6 @@ var codecs = serializer.NewCodecFactory(model.GetScheme())
4946
var patchCodec = codecs.LegacyCodec(workloads.SchemeGroupVersion)
5047
var controllerKind = apps.SchemeGroupVersion.WithKind("StatefulSet")
5148

52-
var jsonIter = jsoniter.ConfigCompatibleWithStandardLibrary
53-
5449
func newRevision(its *workloads.InstanceSet) (*apps.ControllerRevision, error) {
5550
patch, err := getPatch(its)
5651
if err != nil {
@@ -162,46 +157,10 @@ func deepHashObject(hasher hash.Hash, objectToWrite interface{}) {
162157
fmt.Fprintf(hasher, "%v", dump.ForHash(objectToWrite))
163158
}
164159

165-
var revisionsCache = lru.New(1024)
166-
167160
func GetRevisions(revisions map[string]string) (map[string]string, error) {
168-
if revisions == nil {
169-
return nil, nil
170-
}
171-
revisionsStr, ok := revisions[revisionsZSTDKey]
172-
if !ok {
173-
return revisions, nil
174-
}
175-
if revisionsInCache, ok := revisionsCache.Get(revisionsStr); ok {
176-
return revisionsInCache.(map[string]string), nil
177-
}
178-
revisionsData, err := base64.StdEncoding.DecodeString(revisionsStr)
179-
if err != nil {
180-
return nil, err
181-
}
182-
revisionsJSON, err := reader.DecodeAll(revisionsData, nil)
183-
if err != nil {
184-
return nil, err
185-
}
186-
updateRevisions := make(map[string]string)
187-
188-
if err = jsonIter.Unmarshal(revisionsJSON, &updateRevisions); err != nil {
189-
return nil, err
190-
}
191-
revisionsCache.Put(revisionsStr, updateRevisions)
192-
return updateRevisions, nil
161+
return revisionmap.Decode(revisions)
193162
}
194163

195164
func buildRevisions(updateRevisions map[string]string) (map[string]string, error) {
196-
maxPlainRevisionCount := viper.GetInt(MaxPlainRevisionCount)
197-
if len(updateRevisions) <= maxPlainRevisionCount {
198-
return updateRevisions, nil
199-
}
200-
revisionsJSON, err := jsonIter.Marshal(updateRevisions)
201-
if err != nil {
202-
return nil, err
203-
}
204-
revisionsData := writer.EncodeAll(revisionsJSON, nil)
205-
revisionsStr := base64.StdEncoding.EncodeToString(revisionsData)
206-
return map[string]string{revisionsZSTDKey: revisionsStr}, nil
165+
return revisionmap.Encode(updateRevisions)
207166
}

pkg/controller/instanceset/suite_test.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ import (
3030

3131
"github.com/go-logr/logr"
3232
"github.com/golang/mock/gomock"
33+
"github.com/klauspost/compress/zstd"
3334
apps "k8s.io/api/apps/v1"
3435
corev1 "k8s.io/api/core/v1"
3536
"k8s.io/apimachinery/pkg/api/resource"
@@ -170,6 +171,10 @@ func mockCompressedInstanceTemplates(ns, name string) (*corev1.ConfigMap, string
170171
if err != nil {
171172
return nil, "", err
172173
}
174+
writer, err := zstd.NewWriter(nil)
175+
if err != nil {
176+
return nil, "", err
177+
}
173178
templateData := writer.EncodeAll(templateByte, nil)
174179
templateName := fmt.Sprintf("template-ref-%s", name)
175180
templateObj := builder.NewConfigMapBuilder(ns, templateName).

pkg/controller/instanceset/types.go

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -31,13 +31,8 @@ const (
3131
)
3232

3333
const (
34-
// MaxPlainRevisionCount specified max number of plain revision stored in status.updateRevisions.
35-
// All revisions will be compressed if exceeding this value.
36-
MaxPlainRevisionCount = "MAX_PLAIN_REVISION_COUNT"
37-
3834
templateRefAnnotationKey = "kubeblocks.io/template-ref"
3935
templateRefDataKey = "instances"
40-
revisionsZSTDKey = "zstd"
4136

4237
FeatureGateIgnorePodVerticalScaling = "IGNORE_POD_VERTICAL_SCALING"
4338

pkg/controller/instanceset2/instance_util.go

Lines changed: 41 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@ import (
3838
"github.com/apecloud/kubeblocks/pkg/controller/instancetemplate"
3939
"github.com/apecloud/kubeblocks/pkg/controller/kubebuilderx"
4040
"github.com/apecloud/kubeblocks/pkg/controller/model"
41+
"github.com/apecloud/kubeblocks/pkg/controller/revisionmap"
4142
intctrlutil "github.com/apecloud/kubeblocks/pkg/controllerutil"
4243
)
4344

@@ -325,6 +326,37 @@ func getHeadlessSvcName(itsName string) string {
325326
return strings.Join([]string{itsName, "headless"}, "-")
326327
}
327328

329+
func buildDesiredInstancesByName(tree *kubebuilderx.ObjectTree, its *workloads.InstanceSet) (map[string]*workloads.Instance, []string, error) {
330+
itsExt, err := instancetemplate.BuildInstanceSetExt(its, tree)
331+
if err != nil {
332+
return nil, nil, err
333+
}
334+
nameBuilder, err := instancetemplate.NewPodNameBuilder(itsExt, nil)
335+
if err != nil {
336+
return nil, nil, err
337+
}
338+
nameMap, err := nameBuilder.BuildInstanceName2TemplateMap()
339+
if err != nil {
340+
return nil, nil, err
341+
}
342+
343+
names := make([]string, 0, len(nameMap))
344+
for name := range nameMap {
345+
names = append(names, name)
346+
}
347+
sort.Strings(names)
348+
349+
desired := make(map[string]*workloads.Instance, len(names))
350+
for _, name := range names {
351+
inst, err := buildInstanceByTemplate(tree, name, nameMap[name], its)
352+
if err != nil {
353+
return nil, nil, err
354+
}
355+
desired[name] = inst
356+
}
357+
return desired, names, nil
358+
}
359+
328360
// func mergeInPlaceFields(src, dst *corev1.PodTemplateSpec) {
329361
// mergeMap(&src.Annotations, &dst.Annotations)
330362
// mergeMap(&src.Labels, &dst.Labels)
@@ -393,12 +425,17 @@ func getHeadlessSvcName(itsName string) string {
393425
// return requests, limits
394426
// }
395427

396-
func isInstanceUpdated(its *workloads.InstanceSet, inst *workloads.Instance) bool {
397-
generation, ok := inst.Annotations[constant.KubeBlocksGenerationKey]
398-
if !ok {
428+
func isInstanceUpdated(its *workloads.InstanceSet, inst, desired *workloads.Instance) bool {
429+
updateRevisions, err := revisionmap.Decode(its.Status.UpdateRevisions)
430+
if err != nil {
399431
return false
400432
}
401-
if strconv.FormatInt(its.Generation, 10) != generation {
433+
return isInstanceUpdatedWithRevisions(inst, buildCurrentInstanceRevision(inst, desired), updateRevisions)
434+
}
435+
436+
func isInstanceUpdatedWithRevisions(inst *workloads.Instance, currentRevision string, updateRevisions map[string]string) bool {
437+
updateRevision, ok := updateRevisions[inst.Name]
438+
if !ok || currentRevision != updateRevision {
402439
return false
403440
}
404441
return inst.Generation == inst.Status.ObservedGeneration && inst.Status.UpToDate

pkg/controller/instanceset2/reconciler_revision_update.go

Lines changed: 22 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
workloads "github.com/apecloud/kubeblocks/apis/workloads/v1"
2626
"github.com/apecloud/kubeblocks/pkg/controller/kubebuilderx"
2727
"github.com/apecloud/kubeblocks/pkg/controller/model"
28+
"github.com/apecloud/kubeblocks/pkg/controller/revisionmap"
2829
)
2930

3031
func NewRevisionUpdateReconciler() kubebuilderx.Reconciler {
@@ -45,19 +46,37 @@ func (r *revisionUpdateReconciler) PreCondition(tree *kubebuilderx.ObjectTree) *
4546
func (r *revisionUpdateReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilderx.Result, error) {
4647
its, _ := tree.GetRoot().(*workloads.InstanceSet)
4748

48-
updatedReplicas := r.calculateUpdatedReplicas(its, tree.List(&workloads.Instance{}))
49+
desiredInstances, names, err := buildDesiredInstancesByName(tree, its)
50+
if err != nil {
51+
return kubebuilderx.Continue, err
52+
}
53+
54+
updateRevisions := make(map[string]string, len(names))
55+
for _, name := range names {
56+
updateRevisions[name] = buildInstanceRevision(desiredInstances[name])
57+
}
58+
revisions, err := revisionmap.Encode(updateRevisions)
59+
if err != nil {
60+
return kubebuilderx.Continue, err
61+
}
62+
its.Status.UpdateRevisions = revisions
63+
if len(names) > 0 {
64+
its.Status.UpdateRevision = updateRevisions[names[len(names)-1]]
65+
}
66+
67+
updatedReplicas := r.calculateUpdatedReplicas(its, tree.List(&workloads.Instance{}), desiredInstances)
4968
its.Status.UpdatedReplicas = updatedReplicas
5069

5170
its.Status.ObservedGeneration = its.Generation
5271

5372
return kubebuilderx.Continue, nil
5473
}
5574

56-
func (r *revisionUpdateReconciler) calculateUpdatedReplicas(its *workloads.InstanceSet, instances []client.Object) int32 {
75+
func (r *revisionUpdateReconciler) calculateUpdatedReplicas(its *workloads.InstanceSet, instances []client.Object, desiredInstances map[string]*workloads.Instance) int32 {
5776
updatedReplicas := int32(0)
5877
for i := range instances {
5978
inst, _ := instances[i].(*workloads.Instance)
60-
if isInstanceUpdated(its, inst) {
79+
if isInstanceUpdated(its, inst, desiredInstances[inst.Name]) {
6180
updatedReplicas++
6281
}
6382
}

pkg/controller/instanceset2/reconciler_status.go

Lines changed: 27 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -29,9 +29,11 @@ import (
2929
"k8s.io/apimachinery/pkg/util/sets"
3030

3131
workloads "github.com/apecloud/kubeblocks/apis/workloads/v1"
32+
"github.com/apecloud/kubeblocks/pkg/constant"
3233
"github.com/apecloud/kubeblocks/pkg/controller/instancetemplate"
3334
"github.com/apecloud/kubeblocks/pkg/controller/kubebuilderx"
3435
"github.com/apecloud/kubeblocks/pkg/controller/model"
36+
"github.com/apecloud/kubeblocks/pkg/controller/revisionmap"
3537
intctrlutil "github.com/apecloud/kubeblocks/pkg/controllerutil"
3638
)
3739

@@ -53,6 +55,11 @@ func (r *statusReconciler) PreCondition(tree *kubebuilderx.ObjectTree) *kubebuil
5355
func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilderx.Result, error) {
5456
its, _ := tree.GetRoot().(*workloads.InstanceSet)
5557

58+
desiredInstances, _, err := buildDesiredInstancesByName(tree, its)
59+
if err != nil {
60+
return kubebuilderx.Continue, err
61+
}
62+
5663
instances := tree.List(&workloads.Instance{})
5764
var instanceList []*workloads.Instance
5865
for _, object := range instances {
@@ -65,7 +72,11 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
6572
readyReplicas, availableReplicas := int32(0), int32(0)
6673
notReadyNames := sets.New[string]()
6774
notAvailableNames := sets.New[string]()
68-
// currentRevisions := map[string]string{}
75+
currentRevisions := map[string]string{}
76+
updateRevisions, err := revisionmap.Decode(its.Status.UpdateRevisions)
77+
if err != nil {
78+
return kubebuilderx.Continue, err
79+
}
6980

7081
template2TemplatesStatus := map[string]*workloads.InstanceTemplateStatus{}
7182
template2TotalReplicas := map[string]int32{}
@@ -78,7 +89,7 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
7889
}
7990

8091
for _, inst := range instanceList {
81-
templateName := inst.Labels[instancetemplate.TemplateNameLabelKey]
92+
templateName := getInstanceTemplateName(inst)
8293
if template2TemplatesStatus[templateName] == nil {
8394
template2TemplatesStatus[templateName] = &workloads.InstanceTemplateStatus{
8495
Name: templateName,
@@ -100,8 +111,9 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
100111
notAvailableNames.Insert(inst.Name)
101112
}
102113
}
114+
currentRevisions[inst.Name] = buildCurrentInstanceRevision(inst, desiredInstances[inst.Name])
103115
if !intctrlutil.IsInstanceTerminating(inst) {
104-
if isInstanceUpdated(its, inst) {
116+
if isInstanceUpdatedWithRevisions(inst, currentRevisions[inst.Name], updateRevisions) {
105117
updatedReplicas++
106118
template2TemplatesStatus[templateName].UpdatedReplicas++
107119
} else {
@@ -115,15 +127,15 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
115127
its.Status.AvailableReplicas = availableReplicas
116128
its.Status.CurrentReplicas = currentReplicas
117129
its.Status.UpdatedReplicas = updatedReplicas
118-
// its.Status.CurrentRevisions, _ = buildRevisions(currentRevisions)
130+
its.Status.CurrentRevisions, _ = revisionmap.Encode(currentRevisions)
119131
its.Status.TemplatesStatus = buildTemplatesStatus(template2TemplatesStatus)
120132
// all pods have been updated
121133
totalReplicas := int32(1)
122134
if its.Spec.Replicas != nil {
123135
totalReplicas = *its.Spec.Replicas
124136
}
125137
if its.Status.Replicas == totalReplicas && its.Status.UpdatedReplicas == totalReplicas {
126-
// its.Status.CurrentRevision = its.Status.UpdateRevision
138+
its.Status.CurrentRevision = its.Status.UpdateRevision
127139
its.Status.CurrentReplicas = totalReplicas
128140
}
129141
for idx, templateStatus := range its.Status.TemplatesStatus {
@@ -165,6 +177,16 @@ func (r *statusReconciler) Reconcile(tree *kubebuilderx.ObjectTree) (kubebuilder
165177
return kubebuilderx.Continue, nil
166178
}
167179

180+
func getInstanceTemplateName(inst *workloads.Instance) string {
181+
if inst.Labels == nil {
182+
return ""
183+
}
184+
if templateName := inst.Labels[instancetemplate.TemplateNameLabelKey]; templateName != "" {
185+
return templateName
186+
}
187+
return inst.Labels[constant.KBAppInstanceTemplateLabelKey]
188+
}
189+
168190
func buildConditionMessageWithNames(instanceNames []string) ([]byte, error) {
169191
baseSort(instanceNames, func(i int) (string, int) {
170192
return parseParentNameAndOrdinal(instanceNames[i])

0 commit comments

Comments
 (0)