perf(user): reduce controller cache memory usage (#7256)

* perf(user): reduce controller cache memory

* fix(user): resolve GVK for direct-read finalizers

* perf(user): cache projected user fields

* perf(user): project cached metadata fields

* style(user): satisfy cache lint checks
This commit is contained in:
zijiren
2026-08-27 16:22:19 +08:00
committed by GitHub
parent 8526c3a030
commit a2a338dbcf
14 changed files with 822 additions and 82 deletions
@@ -72,11 +72,8 @@ func (r ReqValidator) ValidateCreate(
) (admission.Warnings, error) {
req, ok := obj.(*Operationrequest)
if !ok {
return admission.Warnings{
"obj convert Operationrequest is error",
}, errors.New(
"obj convert Operationrequest is error",
)
message := "obj convert Operationrequest is error"
return admission.Warnings{message}, errors.New(message)
}
// todo check request, _ := admission.RequestFromContext(ctx), request.UserInfo.Username if legal
@@ -103,11 +100,8 @@ func (r ReqValidator) ValidateCreate(
"phase",
item.Status.Phase,
)
return admission.Warnings{
"there is a request not completed, can not create new request",
}, errors.New(
"there is a request not completed, can not create new request",
)
message := "there is a request not completed, can not create new request"
return admission.Warnings{message}, errors.New(message)
}
}
return admission.Warnings{}, nil
@@ -120,26 +114,17 @@ func (r ReqValidator) ValidateUpdate(
// todo check request, _ := admission.RequestFromContext(ctx), request.UserInfo.Username if legal
oldReq, ok := oldObj.(*Operationrequest)
if !ok {
return admission.Warnings{
"obj convert Operationrequest error",
}, errors.New(
"obj convert Operationrequest error",
)
message := "obj convert Operationrequest error"
return admission.Warnings{message}, errors.New(message)
}
newReq, ok := newObj.(*Operationrequest)
if !ok {
return admission.Warnings{
"obj convert Operationrequest error",
}, errors.New(
"obj convert Operationrequest error",
)
message := "obj convert Operationrequest error"
return admission.Warnings{message}, errors.New(message)
}
if oldReq.Spec != newReq.Spec {
return admission.Warnings{
"operation request spec do not support update",
}, errors.New(
"operation request spec do not support update",
)
message := "operation request spec do not support update"
return admission.Warnings{message}, errors.New(message)
}
return admission.Warnings{}, nil
}
+189
View File
@@ -0,0 +1,189 @@
// Copyright 2026 labring.
//
// Licensed 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 cache
import (
"time"
licensev1 "github.com/labring/sealos/controllers/license/api/v1"
userv1 "github.com/labring/sealos/controllers/user/api/v1"
"github.com/labring/sealos/controllers/user/controllers/helper/config"
corev1 "k8s.io/api/core/v1"
rbacv1 "k8s.io/api/rbac/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
ctrlcache "sigs.k8s.io/controller-runtime/pkg/cache"
"sigs.k8s.io/controller-runtime/pkg/client"
)
// Options keeps only explicitly registered informer data in memory. Large
// fields are removed when controllers only need a smaller object projection.
func Options(syncPeriod *time.Duration) ctrlcache.Options {
return ctrlcache.Options{
SyncPeriod: syncPeriod,
ReaderFailOnMissingInformer: true,
DefaultTransform: ctrlcache.TransformStripManagedFields(),
ByObject: map[client.Object]ctrlcache.ByObject{
&licensev1.License{}: {
Transform: transformKeyMetadata,
},
&userv1.User{}: {
Transform: transformUser,
},
&userv1.DeleteRequest{}: {
Transform: transformKeyMetadata,
},
&corev1.Secret{}: {
Namespaces: map[string]ctrlcache.Config{
config.GetUserSystemNamespace(): {},
},
Transform: transformSecret,
},
&corev1.ServiceAccount{}: {
Namespaces: map[string]ctrlcache.Config{
config.GetUserSystemNamespace(): {},
},
Transform: transformOwnerMetadata,
},
&userv1.Operationrequest{}: {
Namespaces: map[string]ctrlcache.Config{
config.GetUserSystemNamespace(): {},
},
Transform: transformKeyMetadata,
},
&rbacv1.Role{}: {
Transform: transformOwnerMetadata,
},
&rbacv1.RoleBinding{}: {
Transform: transformOwnerMetadata,
},
},
}
}
// UncachedObjects returns objects whose reads require complete, current API data.
func UncachedObjects() []client.Object {
return []client.Object{
&licensev1.License{},
&userv1.DeleteRequest{},
&userv1.Operationrequest{},
&corev1.Namespace{},
&corev1.Secret{},
&corev1.ServiceAccount{},
&rbacv1.ClusterRoleBinding{},
&rbacv1.Role{},
&rbacv1.RoleBinding{},
}
}
func transformUser(obj any) (any, error) {
user, ok := obj.(*userv1.User)
if !ok {
return obj, nil
}
metadata := projectObjectMeta(user.ObjectMeta)
metadata.Finalizers = append([]string(nil), user.Finalizers...)
metadata.Annotations = copyMapValues(
user.Annotations,
userv1.UserAnnotationOwnerKey,
)
metadata.Labels = copyMapValues(
user.Labels,
"user.sealos.io/status",
"user.sealos.io/type",
)
status := *user.Status.DeepCopy()
status.KubeConfig = ""
return &userv1.User{
TypeMeta: user.TypeMeta,
ObjectMeta: metadata,
Spec: *user.Spec.DeepCopy(),
Status: status,
}, nil
}
func transformKeyMetadata(obj any) (any, error) {
return transformMetadata(obj, false)
}
func transformOwnerMetadata(obj any) (any, error) {
return transformMetadata(obj, true)
}
func transformMetadata(obj any, keepOwnerReferences bool) (any, error) {
metadata, ok := obj.(*metav1.PartialObjectMetadata)
if !ok {
return obj, nil
}
projected := projectObjectMeta(metadata.ObjectMeta)
if keepOwnerReferences {
projected.OwnerReferences = append(
[]metav1.OwnerReference(nil),
metadata.OwnerReferences...,
)
}
return &metav1.PartialObjectMetadata{
TypeMeta: metadata.TypeMeta,
ObjectMeta: projected,
}, nil
}
func projectObjectMeta(in metav1.ObjectMeta) metav1.ObjectMeta {
out := metav1.ObjectMeta{
Name: in.Name,
Namespace: in.Namespace,
UID: in.UID,
ResourceVersion: in.ResourceVersion,
Generation: in.Generation,
CreationTimestamp: in.CreationTimestamp,
}
if in.DeletionTimestamp != nil {
out.DeletionTimestamp = in.DeletionTimestamp.DeepCopy()
}
if in.DeletionGracePeriodSeconds != nil {
gracePeriod := *in.DeletionGracePeriodSeconds
out.DeletionGracePeriodSeconds = &gracePeriod
}
return out
}
func copyMapValues(source map[string]string, keys ...string) map[string]string {
var result map[string]string
for _, key := range keys {
if value, ok := source[key]; ok {
if result == nil {
result = make(map[string]string)
}
result[key] = value
}
}
return result
}
func transformSecret(obj any) (any, error) {
secret, ok := obj.(*corev1.Secret)
if !ok {
return obj, nil
}
metadata := projectObjectMeta(secret.ObjectMeta)
if serviceAccountName, ok := secret.Annotations[corev1.ServiceAccountNameKey]; ok {
metadata.Annotations = map[string]string{
corev1.ServiceAccountNameKey: serviceAccountName,
}
}
return &corev1.Secret{TypeMeta: secret.TypeMeta, ObjectMeta: metadata}, nil
}
+259
View File
@@ -0,0 +1,259 @@
// Copyright 2026 labring.
//
// Licensed 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 cache
import (
"reflect"
"testing"
"time"
userv1 "github.com/labring/sealos/controllers/user/api/v1"
"github.com/labring/sealos/controllers/user/controllers/helper/config"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/client"
)
func TestOptionsLimitsSecretCache(t *testing.T) {
syncPeriod := time.Hour
options := Options(&syncPeriod)
if options.SyncPeriod != &syncPeriod {
t.Fatal("sync period was not retained")
}
if !options.ReaderFailOnMissingInformer {
t.Fatal("missing informer reads are allowed")
}
if options.DefaultTransform == nil {
t.Fatal("default managed fields transform is nil")
}
found := false
for obj, byObject := range options.ByObject {
if _, ok := obj.(*corev1.Secret); !ok {
continue
}
found = true
if byObject.Transform == nil {
t.Fatal("secret transform is nil")
}
if len(byObject.Namespaces) != 1 {
t.Fatalf("secret cache namespaces = %d, want 1", len(byObject.Namespaces))
}
if _, ok := byObject.Namespaces[config.GetUserSystemNamespace()]; !ok {
t.Fatal("user system namespace is not cached for secrets")
}
}
if !found {
t.Fatal("secret cache options not found")
}
}
func TestOptionsLimitNamespacedMetadataCaches(t *testing.T) {
options := Options(nil)
for _, required := range []client.Object{
&corev1.ServiceAccount{},
&userv1.Operationrequest{},
} {
found := false
for obj, byObject := range options.ByObject {
if reflect.TypeOf(obj) != reflect.TypeOf(required) {
continue
}
found = true
if len(byObject.Namespaces) != 1 {
t.Fatalf("%T cache namespaces = %d, want 1", required, len(byObject.Namespaces))
}
if _, ok := byObject.Namespaces[config.GetUserSystemNamespace()]; !ok {
t.Fatalf("%T cache is not limited to the user system namespace", required)
}
}
if !found {
t.Fatalf("%T cache options not found", required)
}
}
}
func TestTransformSecretKeepsOnlyIndexMetadata(t *testing.T) {
secret := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: "token-a",
Namespace: config.GetUserSystemNamespace(),
ResourceVersion: "42",
Annotations: map[string]string{
corev1.ServiceAccountNameKey: "user-a",
"unused.example/key": "large-value",
},
ManagedFields: []metav1.ManagedFieldsEntry{{Manager: "test"}},
},
Type: corev1.SecretTypeServiceAccountToken,
Data: map[string][]byte{"token": []byte("sensitive-data")},
}
transformed, err := transformSecret(secret)
if err != nil {
t.Fatalf("transform secret: %v", err)
}
got, ok := transformed.(*corev1.Secret)
if !ok {
t.Fatalf("transformed type = %T, want *corev1.Secret", transformed)
}
if got.Name != secret.Name || got.Namespace != secret.Namespace || got.ResourceVersion != "42" {
t.Fatalf("required metadata was not retained: %#v", got.ObjectMeta)
}
if got.Annotations[corev1.ServiceAccountNameKey] != "user-a" || len(got.Annotations) != 1 {
t.Fatalf("secret index annotations = %#v", got.Annotations)
}
if got.Type != "" || len(got.Data) != 0 || len(got.ManagedFields) != 0 {
t.Fatalf("secret payload was retained: %#v", got)
}
}
func TestTransformUserDropsOnlyLargeUnusedFields(t *testing.T) {
rotateAt := metav1.Now()
user := &userv1.User{
ObjectMeta: metav1.ObjectMeta{
Name: "user-a",
ResourceVersion: "42",
Annotations: map[string]string{
userv1.UserAnnotationOwnerKey: "owner-a",
"unused.example/annotation": "unused",
},
Labels: map[string]string{
"user.sealos.io/status": "active",
"user.sealos.io/type": "Group",
"unused.example/label": "unused",
},
Finalizers: []string{"sealos.io/user.finalizers"},
ManagedFields: []metav1.ManagedFieldsEntry{{Manager: "test"}},
},
Spec: userv1.UserSpec{
CSRExpirationSeconds: 600,
KubeConfigRotateAt: &rotateAt,
},
Status: userv1.UserStatus{
Phase: userv1.UserActive,
KubeConfig: "large-kubeconfig",
ObservedCSRExpirationSeconds: 600,
ObservedKubeConfigRotateAt: &rotateAt,
ObservedGeneration: 7,
Conditions: []userv1.Condition{{
Type: userv1.Ready,
Status: corev1.ConditionTrue,
}},
},
}
transformed, err := transformUser(user)
if err != nil {
t.Fatalf("transform user: %v", err)
}
got, ok := transformed.(*userv1.User)
if !ok {
t.Fatalf("transformed type = %T, want *v1.User", transformed)
}
if got.Status.KubeConfig != "" {
t.Fatal("kubeconfig was retained")
}
if len(got.ManagedFields) != 0 {
t.Fatal("managed fields were retained")
}
wantAnnotations := map[string]string{userv1.UserAnnotationOwnerKey: "owner-a"}
wantLabels := map[string]string{
"user.sealos.io/status": "active",
"user.sealos.io/type": "Group",
}
if got.Name != user.Name || got.ResourceVersion != user.ResourceVersion ||
!reflect.DeepEqual(got.Annotations, wantAnnotations) ||
!reflect.DeepEqual(got.Labels, wantLabels) ||
!reflect.DeepEqual(got.Finalizers, user.Finalizers) ||
!reflect.DeepEqual(got.Spec, user.Spec) {
t.Fatalf("required user fields were not retained: %#v", got)
}
wantStatus := user.Status.DeepCopy()
wantStatus.KubeConfig = ""
if !reflect.DeepEqual(&got.Status, wantStatus) {
t.Fatalf("status = %#v, want %#v", got.Status, *wantStatus)
}
if user.Status.KubeConfig == "" || len(user.ManagedFields) == 0 {
t.Fatal("transform mutated the source user")
}
}
func TestTransformMetadataKeepsOnlyEventFields(t *testing.T) {
metadata := &metav1.PartialObjectMetadata{
ObjectMeta: metav1.ObjectMeta{
Name: "role-a",
Namespace: "ns-a",
ResourceVersion: "42",
Annotations: map[string]string{"unused.example/key": "unused"},
Labels: map[string]string{"unused.example/key": "unused"},
Finalizers: []string{"unused.example/finalizer"},
ManagedFields: []metav1.ManagedFieldsEntry{{Manager: "test"}},
OwnerReferences: []metav1.OwnerReference{{Name: "user-a"}},
},
}
metadata.SetGroupVersionKind(corev1.SchemeGroupVersion.WithKind("ServiceAccount"))
transformed, err := transformOwnerMetadata(metadata)
if err != nil {
t.Fatalf("transform metadata: %v", err)
}
got, ok := transformed.(*metav1.PartialObjectMetadata)
if !ok {
t.Fatalf("transformed type = %T, want *metav1.PartialObjectMetadata", transformed)
}
if got.GroupVersionKind() != metadata.GroupVersionKind() || got.Name != metadata.Name ||
got.Namespace != metadata.Namespace || got.ResourceVersion != metadata.ResourceVersion {
t.Fatalf("required event metadata was not retained: %#v", got)
}
if !reflect.DeepEqual(got.OwnerReferences, metadata.OwnerReferences) {
t.Fatalf("owner references = %#v, want %#v", got.OwnerReferences, metadata.OwnerReferences)
}
if len(got.Annotations) != 0 || len(got.Labels) != 0 || len(got.Finalizers) != 0 ||
len(got.ManagedFields) != 0 {
t.Fatalf("unused metadata was retained: %#v", got.ObjectMeta)
}
keyOnly, err := transformKeyMetadata(metadata)
if err != nil {
t.Fatalf("transform key metadata: %v", err)
}
gotKey, ok := keyOnly.(*metav1.PartialObjectMetadata)
if !ok {
t.Fatalf("key metadata type = %T, want *metav1.PartialObjectMetadata", keyOnly)
}
if len(gotKey.OwnerReferences) != 0 {
t.Fatalf("key-only owner references were retained: %#v", gotKey.OwnerReferences)
}
}
func TestUncachedObjects(t *testing.T) {
types := make(map[reflect.Type]struct{})
for _, obj := range UncachedObjects() {
types[reflect.TypeOf(obj)] = struct{}{}
}
for _, required := range []client.Object{
&corev1.Namespace{},
&corev1.Secret{},
&corev1.ServiceAccount{},
} {
if _, ok := types[reflect.TypeOf(required)]; !ok {
t.Fatalf("%T reads are still cache-backed", required)
}
}
if _, ok := types[reflect.TypeFor[*userv1.User]()]; ok {
t.Fatal("user reads bypass the projected cache")
}
}
@@ -28,6 +28,7 @@ import (
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/client-go/tools/record"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/builder"
"sigs.k8s.io/controller-runtime/pkg/client"
)
@@ -226,7 +227,7 @@ func (r *DeleteRequestReconciler) SetupWithManager(mgr ctrl.Manager) error {
r.expirationTime = time.Minute * 10
r.retentionTime = time.Minute * 30
return ctrl.NewControllerManagedBy(mgr).
For(&userv1.DeleteRequest{}).
For(&userv1.DeleteRequest{}, builder.OnlyMetadata).
Complete(r)
}
@@ -27,6 +27,7 @@ import (
type Finalizer struct {
client client.Client
reader client.Reader
finalizerName string
}
@@ -37,6 +38,9 @@ func (f *Finalizer) AddFinalizer(ctx context.Context, obj client.Object) (bool,
// then lets add the finalizer and update the object. This is equivalent
// registering our finalizer.
notDelete = true
if controllerutil.ContainsFinalizer(obj, f.finalizerName) {
return notDelete, nil
}
controllerutil.AddFinalizer(obj, f.finalizerName)
if err := f.updateFinalizers(
ctx,
@@ -57,10 +61,16 @@ func DefaultFunc(ctx context.Context, obj client.Object) error {
func NewFinalizer(client client.Client, finalizerName string) *Finalizer {
return &Finalizer{
client: client,
reader: client,
finalizerName: finalizerName,
}
}
func (f *Finalizer) WithReader(reader client.Reader) *Finalizer {
f.reader = reader
return f
}
func (f *Finalizer) RemoveFinalizer(
ctx context.Context,
obj client.Object,
@@ -95,12 +105,15 @@ func (f *Finalizer) updateFinalizers(
obj client.Object,
finalizers []string,
) error {
gvk, err := f.client.GroupVersionKindFor(obj)
if err != nil {
return err
}
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
gvk := obj.GetObjectKind().GroupVersionKind()
fetchObject := &unstructured.Unstructured{}
fetchObject.SetAPIVersion(gvk.GroupVersion().String())
fetchObject.SetKind(gvk.Kind)
err := f.client.Get(ctx, objectKey, fetchObject)
err := f.reader.Get(ctx, objectKey, fetchObject)
if err != nil {
// We log this error, but we continue and try to set the ownerRefs on the other resources.
return err
@@ -0,0 +1,114 @@
// Copyright 2026 labring.
//
// Licensed 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 finalizer
import (
"context"
"testing"
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"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
)
func TestAddFinalizerResolvesMissingGVK(t *testing.T) {
t.Parallel()
scheme := runtime.NewScheme()
if err := corev1.AddToScheme(scheme); err != nil {
t.Fatalf("add core scheme: %v", err)
}
stored := &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{Name: "test", Namespace: "default"},
}
cli := fake.NewClientBuilder().WithScheme(scheme).WithObjects(stored).Build()
obj := &corev1.ConfigMap{}
key := client.ObjectKeyFromObject(stored)
if err := cli.Get(context.Background(), key, obj); err != nil {
t.Fatalf("get configmap: %v", err)
}
if !obj.GroupVersionKind().Empty() {
t.Fatalf("fetched object GVK = %s, want empty", obj.GroupVersionKind())
}
const finalizerName = "test.sealos.io/finalizer"
updated, err := NewFinalizer(cli, finalizerName).AddFinalizer(context.Background(), obj)
if err != nil {
t.Fatalf("add finalizer: %v", err)
}
if !updated {
t.Fatal("object was not handled")
}
got := &corev1.ConfigMap{}
if err := cli.Get(context.Background(), key, got); err != nil {
t.Fatalf("get updated configmap: %v", err)
}
if !controllerutil.ContainsFinalizer(got, finalizerName) {
t.Fatalf("finalizers = %v, want %q", got.Finalizers, finalizerName)
}
}
type updateCountingClient struct {
client.Client
updates int
}
func (c *updateCountingClient) Update(
ctx context.Context,
obj client.Object,
opts ...client.UpdateOption,
) error {
c.updates++
return c.Client.Update(ctx, obj, opts...)
}
func TestAddFinalizerSkipsExistingFinalizerUpdate(t *testing.T) {
t.Parallel()
scheme := runtime.NewScheme()
if err := corev1.AddToScheme(scheme); err != nil {
t.Fatalf("add core scheme: %v", err)
}
const finalizerName = "test.sealos.io/finalizer"
stored := &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{
Name: "test",
Namespace: "default",
Finalizers: []string{finalizerName},
},
}
baseClient := fake.NewClientBuilder().WithScheme(scheme).WithObjects(stored).Build()
cli := &updateCountingClient{Client: baseClient}
obj := &corev1.ConfigMap{}
if err := cli.Get(context.Background(), client.ObjectKeyFromObject(stored), obj); err != nil {
t.Fatalf("get configmap: %v", err)
}
handled, err := NewFinalizer(cli, finalizerName).AddFinalizer(context.Background(), obj)
if err != nil {
t.Fatalf("add finalizer: %v", err)
}
if !handled {
t.Fatal("object was not handled")
}
if cli.updates != 0 {
t.Fatalf("updates = %d, want 0", cli.updates)
}
}
@@ -144,11 +144,12 @@ func (sac *ServiceAccountConfig) applyBoundTokenSecret(
// CleanupLegacyBoundTokenSecrets removes stale bound token secrets for a user.
func CleanupLegacyBoundTokenSecrets(
ctx context.Context,
cli client.Client,
reader client.Reader,
writer client.Writer,
userName, keepSecretName string,
) error {
secrets := &v1.SecretList{}
if err := cli.List(
if err := reader.List(
ctx,
secrets,
client.InNamespace(config2.GetUserSystemNamespace()),
@@ -167,7 +168,7 @@ func CleanupLegacyBoundTokenSecrets(
if secret.Annotations == nil || secret.Annotations[v1.ServiceAccountNameKey] != userName {
continue
}
if err := cli.Delete(ctx, secret); err != nil && !apierrors.IsNotFound(err) {
if err := writer.Delete(ctx, secret); err != nil && !apierrors.IsNotFound(err) {
return fmt.Errorf("failed to delete legacy bound token secret %s: %w", secret.Name, err)
}
}
@@ -100,6 +100,7 @@ func TestCleanupLegacyBoundTokenSecrets(t *testing.T) {
if err := CleanupLegacyBoundTokenSecrets(
context.Background(),
cli,
cli,
userName,
currentSecret,
); err != nil {
@@ -19,32 +19,36 @@ import (
licensev1 "github.com/labring/sealos/controllers/license/api/v1"
"github.com/labring/sealos/controllers/user/pkg/licensegate"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/tools/cache"
ctrl "sigs.k8s.io/controller-runtime"
)
func SetupLicenseGate(mgr ctrl.Manager) error {
logger := ctrl.Log.WithName("license-gate")
if err := licensegate.Refresh(context.Background(), mgr.GetAPIReader()); err != nil {
reader := mgr.GetAPIReader()
if err := licensegate.Refresh(context.Background(), reader); err != nil {
logger.Error(err, "initial license gate refresh failed")
}
informer, err := mgr.GetCache().GetInformer(context.Background(), &licensev1.License{})
licenseMetadata := &metav1.PartialObjectMetadata{}
licenseMetadata.SetGroupVersionKind(licensev1.GroupVersion.WithKind("License"))
informer, err := mgr.GetCache().GetInformer(context.Background(), licenseMetadata)
if err != nil {
return err
}
if _, err := informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj any) {
if err := licensegate.Refresh(context.Background(), mgr.GetClient()); err != nil {
if err := licensegate.Refresh(context.Background(), reader); err != nil {
logger.Error(err, "license gate refresh failed on add")
}
},
UpdateFunc: func(oldObj, newObj any) {
if err := licensegate.Refresh(context.Background(), mgr.GetClient()); err != nil {
if err := licensegate.Refresh(context.Background(), reader); err != nil {
logger.Error(err, "license gate refresh failed on update")
}
},
DeleteFunc: func(obj any) {
if err := licensegate.Refresh(context.Background(), mgr.GetClient()); err != nil {
if err := licensegate.Refresh(context.Background(), reader); err != nil {
logger.Error(err, "license gate refresh failed on delete")
}
},
@@ -77,7 +77,11 @@ func (r *OperationReqReconciler) SetupWithManager(
r.userLock = make(map[string]*sync.Mutex)
r.Logger.V(1).Info("init reconcile operationrequest controller")
return ctrl.NewControllerManagedBy(mgr).
For(&userv1.Operationrequest{}, builder.WithPredicates(namespaceOnlyPredicate(config.GetUserSystemNamespace()))).
For(
&userv1.Operationrequest{},
builder.WithPredicates(namespaceOnlyPredicate(config.GetUserSystemNamespace())),
builder.OnlyMetadata,
).
WithOptions(controller.Options{
MaxConcurrentReconciles: ratelimiter.GetConcurrent(opts),
RateLimiter: ratelimiter.GetRateLimiter(opts),
@@ -232,8 +236,7 @@ func (r *OperationReqReconciler) reconcile(
}
if request.Spec.Role == userv1.OwnerRoleType {
// update user annotation
user.Annotations[userv1.UserAnnotationOwnerKey] = request.Spec.User
if err := r.Update(ctx, user); err != nil {
if err := r.patchUserOwner(ctx, user, request.Spec.User); err != nil {
r.Recorder.Eventf(
request,
v1.EventTypeWarning,
@@ -308,8 +311,7 @@ func (r *OperationReqReconciler) reconcile(
}
if request.Spec.Role == userv1.OwnerRoleType {
// update user annotation
user.Annotations[userv1.UserAnnotationOwnerKey] = request.Spec.User
if err := r.Update(ctx, user); err != nil {
if err := r.patchUserOwner(ctx, user, request.Spec.User); err != nil {
r.Recorder.Eventf(
request,
v1.EventTypeWarning,
@@ -341,6 +343,19 @@ func (r *OperationReqReconciler) reconcile(
return ctrl.Result{RequeueAfter: OperationReqRequeueDuration}, nil
}
func (r *OperationReqReconciler) patchUserOwner(
ctx context.Context,
user *userv1.User,
owner string,
) error {
original := user.DeepCopy()
if user.Annotations == nil {
user.Annotations = make(map[string]string)
}
user.Annotations[userv1.UserAnnotationOwnerKey] = owner
return r.Patch(ctx, user, client.MergeFrom(original))
}
// isRetained returns true if the request is isCompleted and exist for retention time
func (r *OperationReqReconciler) isRetained(request *userv1.Operationrequest) bool {
if request.Status.Phase == userv1.RequestCompleted &&
+43 -36
View File
@@ -40,7 +40,6 @@ import (
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
"k8s.io/client-go/tools/clientcmd/api"
@@ -69,7 +68,6 @@ type UserReconciler struct {
Logger logr.Logger
Recorder record.EventRecorder
cache cache.Cache
apiReader client.Reader
userCounter *usercount.Counter
config *rest.Config
*runtime.Scheme
@@ -133,6 +131,15 @@ type ControllerRestartPredicate struct {
checkTime time.Time
}
type OwnerAnnotationChangedPredicate struct {
predicate.Funcs
}
func (OwnerAnnotationChangedPredicate) Update(e event.UpdateEvent) bool {
return e.ObjectOld.GetAnnotations()[userAnnotationOwnerKey] !=
e.ObjectNew.GetAnnotations()[userAnnotationOwnerKey]
}
func NewControllerRestartPredicate(duration time.Duration) *ControllerRestartPredicate {
return &ControllerRestartPredicate{
checkTime: time.Now().Add(-duration),
@@ -159,11 +166,11 @@ func (r *UserReconciler) SetupWithManager(mgr ctrl.Manager, opts ratelimiter.Rat
r.Recorder = mgr.GetEventRecorderFor(controllerName)
}
if r.finalizer == nil {
r.finalizer = finalizer.NewFinalizer(r.Client, "sealos.io/user.finalizers")
r.finalizer = finalizer.NewFinalizer(r.Client, "sealos.io/user.finalizers").
WithReader(mgr.GetAPIReader())
}
r.Scheme = mgr.GetScheme()
r.cache = mgr.GetCache()
r.apiReader = mgr.GetAPIReader()
r.userCounter = userCounter
r.config = mgr.GetConfig()
r.Logger.V(1).Info("init reconcile controller user")
@@ -197,11 +204,21 @@ func (r *UserReconciler) SetupWithManager(mgr ctrl.Manager, opts ratelimiter.Rat
)
return ctrl.NewControllerManagedBy(mgr).
For(&userv1.User{}, builder.WithPredicates(predicate.Or(predicate.GenerationChangedPredicate{}, predicate.AnnotationChangedPredicate{}))).
Watches(&licensev1.License{}, handler.EnqueueRequestsFromMapFunc(r.licenseToUserRequests)).
Watches(&rbacv1.Role{}, ownerEventHandler).
Watches(&rbacv1.RoleBinding{}, ownerEventHandler).
Watches(&v1.ServiceAccount{}, ownerEventHandler).
For(
&userv1.User{},
builder.WithPredicates(predicate.Or(
predicate.GenerationChangedPredicate{},
OwnerAnnotationChangedPredicate{},
)),
).
Watches(
&licensev1.License{},
handler.EnqueueRequestsFromMapFunc(r.licenseToUserRequests),
builder.OnlyMetadata,
).
Watches(&rbacv1.Role{}, ownerEventHandler, builder.OnlyMetadata).
Watches(&rbacv1.RoleBinding{}, ownerEventHandler, builder.OnlyMetadata).
Watches(&v1.ServiceAccount{}, ownerEventHandler, builder.OnlyMetadata).
WithOptions(kubecontroller.Options{
MaxConcurrentReconciles: ratelimiter.GetConcurrent(opts),
RateLimiter: ratelimiter.GetRateLimiter(opts),
@@ -220,7 +237,8 @@ func (r *UserReconciler) reconcile(ctx context.Context, obj client.Object) (ctrl
return ctrl.Result{}, errors.New("obj convert user is error")
}
blocked, err := r.handleLicenseLimit(ctx, user)
originalStatus := user.Status.DeepCopy()
blocked, err := r.handleLicenseLimit(ctx, user, originalStatus)
if err != nil {
return ctrl.Result{}, err
}
@@ -255,6 +273,7 @@ func (r *UserReconciler) reconcile(ctx context.Context, obj client.Object) (ctrl
// Best-effort migration cleanup for legacy service-account-token secrets.
if err := kubeconfig.CleanupLegacyBoundTokenSecrets(
ctx,
r.cache,
r.Client,
user.Name,
state.currentSecretName,
@@ -270,7 +289,7 @@ func (r *UserReconciler) reconcile(ctx context.Context, obj client.Object) (ctrl
r.Logger.Error(err, "cleanup stale bound token secrets", "user", user.Name)
}
}
err = r.updateStatus(ctx, client.ObjectKeyFromObject(obj), user.Status.DeepCopy())
err = r.updateStatus(ctx, user, originalStatus)
if err != nil {
r.Recorder.Eventf(
user,
@@ -882,31 +901,19 @@ func (r *UserReconciler) shouldRotateKubeConfig(user *userv1.User) bool {
func (r *UserReconciler) updateStatus(
ctx context.Context,
nn types.NamespacedName,
status *userv1.UserStatus,
user *userv1.User,
originalStatus *userv1.UserStatus,
) error {
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
original := &userv1.User{}
if err := r.Get(ctx, nn, original); err != nil {
return err
}
original.Status = *status
return r.Client.Status().Update(ctx, original)
})
original := user.DeepCopy()
original.Status = *originalStatus.DeepCopy()
return r.Client.Status().Patch(ctx, user, client.MergeFrom(original))
}
func (r *UserReconciler) handleLicenseLimit(ctx context.Context, user *userv1.User) (bool, error) {
reader := r.apiReader
if reader == nil {
reader = r.Client
}
latest := &userv1.User{}
if err := reader.Get(ctx, client.ObjectKeyFromObject(user), latest); err != nil {
return false, err
}
*user = *latest.DeepCopy()
func (r *UserReconciler) handleLicenseLimit(
ctx context.Context,
user *userv1.User,
originalStatus *userv1.UserStatus,
) (bool, error) {
if !r.isNewUser(user) {
user.Status.Conditions = helper.DeleteCondition(
user.Status.Conditions,
@@ -938,8 +945,8 @@ func (r *UserReconciler) handleLicenseLimit(ctx context.Context, user *userv1.Us
user.Status.Conditions = helper.UpdateCondition(user.Status.Conditions, *limitCondition)
if err := r.updateStatus(
ctx,
client.ObjectKeyFromObject(user),
user.Status.DeepCopy(),
user,
originalStatus,
); err != nil {
return false, err
}
@@ -980,7 +987,7 @@ func (r *UserReconciler) licenseToUserRequests(
obj client.Object,
) []ctrl.Request {
userList := &userv1.UserList{}
if err := r.List(ctx, userList); err != nil {
if err := r.cache.List(ctx, userList); err != nil {
r.Logger.Error(err, "list users for license change failed")
return nil
}
@@ -83,7 +83,8 @@ func (r *UserExpirationReconciler) SetupWithManager(mgr ctrl.Manager) error {
r.Recorder = mgr.GetEventRecorderFor(controllerName)
}
if r.finalizer == nil {
r.finalizer = finalizer.NewFinalizer(r.Client, "sealos.io/user.expiration.finalizers")
r.finalizer = finalizer.NewFinalizer(r.Client, "sealos.io/user.expiration.finalizers").
WithReader(mgr.GetAPIReader())
}
r.Scheme = mgr.GetScheme()
r.config = mgr.GetConfig()
@@ -0,0 +1,148 @@
// Copyright 2026 labring.
//
// Licensed 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 controllers
import (
"context"
"testing"
userv1 "github.com/labring/sealos/controllers/user/api/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
"sigs.k8s.io/controller-runtime/pkg/event"
)
func TestUpdateStatusPreservesUncachedKubeConfig(t *testing.T) {
t.Parallel()
scheme := runtime.NewScheme()
if err := userv1.AddToScheme(scheme); err != nil {
t.Fatalf("add user scheme: %v", err)
}
stored := &userv1.User{
ObjectMeta: metav1.ObjectMeta{Name: "user-a"},
Status: userv1.UserStatus{
Phase: userv1.UserActive,
KubeConfig: "existing-kubeconfig",
},
}
cli := fake.NewClientBuilder().
WithScheme(scheme).
WithStatusSubresource(&userv1.User{}).
WithObjects(stored).
Build()
projected := &userv1.User{}
if err := cli.Get(
context.Background(),
client.ObjectKeyFromObject(stored),
projected,
); err != nil {
t.Fatalf("get user: %v", err)
}
projected.Status.KubeConfig = ""
originalStatus := projected.Status.DeepCopy()
projected.Status.ObservedGeneration = 1
reconciler := &UserReconciler{Client: cli}
if err := reconciler.updateStatus(context.Background(), projected, originalStatus); err != nil {
t.Fatalf("patch status: %v", err)
}
got := &userv1.User{}
if err := cli.Get(context.Background(), client.ObjectKeyFromObject(stored), got); err != nil {
t.Fatalf("get updated user: %v", err)
}
if got.Status.KubeConfig != stored.Status.KubeConfig {
t.Fatalf("kubeconfig = %q, want %q", got.Status.KubeConfig, stored.Status.KubeConfig)
}
if got.Status.ObservedGeneration != 1 {
t.Fatalf("observed generation = %d, want 1", got.Status.ObservedGeneration)
}
}
func TestOwnerAnnotationChangedPredicate(t *testing.T) {
t.Parallel()
oldUser := &userv1.User{ObjectMeta: metav1.ObjectMeta{
Annotations: map[string]string{
userv1.UserAnnotationOwnerKey: "owner-a",
userv1.UserAnnotationDisplayKey: "old-display",
},
}}
newUser := oldUser.DeepCopy()
newUser.Annotations[userv1.UserAnnotationDisplayKey] = "new-display"
predicate := OwnerAnnotationChangedPredicate{}
if predicate.Update(event.UpdateEvent{ObjectOld: oldUser, ObjectNew: newUser}) {
t.Fatal("unrelated annotation change triggered reconciliation")
}
newUser.Annotations[userv1.UserAnnotationOwnerKey] = "owner-b"
if !predicate.Update(event.UpdateEvent{ObjectOld: oldUser, ObjectNew: newUser}) {
t.Fatal("owner annotation change did not trigger reconciliation")
}
}
func TestPatchUserOwnerPreservesUncachedAnnotations(t *testing.T) {
t.Parallel()
scheme := runtime.NewScheme()
if err := userv1.AddToScheme(scheme); err != nil {
t.Fatalf("add user scheme: %v", err)
}
stored := &userv1.User{
ObjectMeta: metav1.ObjectMeta{
Name: "user-a",
Annotations: map[string]string{
userv1.UserAnnotationOwnerKey: "old-owner",
userv1.UserAnnotationDisplayKey: "display-name",
},
},
}
cli := fake.NewClientBuilder().WithScheme(scheme).WithObjects(stored).Build()
projected := &userv1.User{}
if err := cli.Get(
context.Background(),
client.ObjectKeyFromObject(stored),
projected,
); err != nil {
t.Fatalf("get user: %v", err)
}
projected.Annotations = map[string]string{
userv1.UserAnnotationOwnerKey: projected.Annotations[userv1.UserAnnotationOwnerKey],
}
reconciler := &OperationReqReconciler{Client: cli}
if err := reconciler.patchUserOwner(context.Background(), projected, "new-owner"); err != nil {
t.Fatalf("patch user owner: %v", err)
}
got := &userv1.User{}
if err := cli.Get(context.Background(), client.ObjectKeyFromObject(stored), got); err != nil {
t.Fatalf("get updated user: %v", err)
}
if got.Annotations[userv1.UserAnnotationOwnerKey] != "new-owner" {
t.Fatalf("owner = %q, want new-owner", got.Annotations[userv1.UserAnnotationOwnerKey])
}
if got.Annotations[userv1.UserAnnotationDisplayKey] != "display-name" {
t.Fatalf(
"display annotation = %q, want display-name",
got.Annotations[userv1.UserAnnotationDisplayKey],
)
}
}
+6 -4
View File
@@ -26,6 +26,7 @@ import (
licensev1 "github.com/labring/sealos/controllers/license/api/v1"
userv1 "github.com/labring/sealos/controllers/user/api/v1"
"github.com/labring/sealos/controllers/user/controllers"
usercache "github.com/labring/sealos/controllers/user/controllers/cache"
ratelimiter "github.com/labring/sealos/controllers/user/controllers/helper/ratelimiter"
"k8s.io/apimachinery/pkg/runtime"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
@@ -36,7 +37,7 @@ import (
"k8s.io/client-go/rest"
"k8s.io/utils/ptr"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/cache"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/config"
"sigs.k8s.io/controller-runtime/pkg/healthz"
"sigs.k8s.io/controller-runtime/pkg/log/zap"
@@ -200,15 +201,16 @@ func main() {
mgr, err := ctrl.NewManager(cfg, ctrl.Options{
Scheme: scheme,
Metrics: metricsServerOptions,
Cache: usercache.Options(&syncPeriod),
Client: client.Options{Cache: &client.CacheOptions{
DisableFor: usercache.UncachedObjects(),
}},
// WebhookServer: webhook.NewServer(webhook.Options{
// Port: 9443,
// }),
HealthProbeBindAddress: probeAddr,
LeaderElection: enableLeaderElection,
LeaderElectionID: "785548a1.sealos.io",
Cache: cache.Options{
SyncPeriod: &syncPeriod,
},
Controller: config.Controller{
UsePriorityQueue: ptr.To(true),
},