feat: add globalresourcequota api (#2068)

* feat: add globalresourcequota api

Signed-off-by: Oliver Baehler <oliver@sudo-i.net>
This commit is contained in:
Oliver Bähler
2026-08-10 21:25:10 +02:00
committed by GitHub
parent d92453427c
commit bdcdcefe63
60 changed files with 6804 additions and 107 deletions
@@ -0,0 +1,352 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package globalresourcequotas
import (
"context"
"fmt"
"reflect"
"slices"
"github.com/go-logr/logr"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/tools/events"
"k8s.io/client-go/util/retry"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/builder"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/controller-runtime/pkg/handler"
"sigs.k8s.io/controller-runtime/pkg/predicate"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
ctrlutils "github.com/projectcapsule/capsule/internal/controllers/utils"
"github.com/projectcapsule/capsule/internal/metrics"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/runtime/predicates"
"github.com/projectcapsule/capsule/pkg/runtime/selectors"
)
type Controller struct {
client.Client
reader client.Reader
log logr.Logger
recorder events.EventRecorder
metrics *metrics.GlobalResourceQuotaRecorder
}
func (r *Controller) SetupWithManager(mgr ctrl.Manager, options ctrlutils.ControllerOptions) error {
r.reader = mgr.GetAPIReader()
return ctrl.NewControllerManagedBy(mgr).
Named("capsule/global-resource-quotas").
For(
&capsulev1beta2.GlobalResourceQuota{},
builder.WithPredicates(predicate.Or(
predicate.GenerationChangedPredicate{},
predicates.UpdatedMetadataPredicate{},
predicates.DeletionChangedPredicate{},
)),
).
Owns(
&corev1.ResourceQuota{},
builder.WithPredicates(predicate.Or(
predicate.GenerationChangedPredicate{},
predicates.UpdatedMetadataPredicate{},
predicates.DeletionChangedPredicate{},
predicates.ResourceQuotaUsageChangedPredicate{},
)),
).
Owns(&capsulev1beta2.QuantityLedger{}).
Watches(
&corev1.Namespace{},
handler.EnqueueRequestsFromMapFunc(r.globalQuotasForNamespace),
builder.WithPredicates(predicates.UpdatedMetadataPredicate{}),
).
WithOptions(options.Runtime.ToControllerOptions()).
Complete(r)
}
func (r *Controller) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
instance := &capsulev1beta2.GlobalResourceQuota{}
if err := r.Get(ctx, request.NamespacedName, instance); err != nil {
if apierrors.IsNotFound(err) {
r.metrics.Delete(request.Name)
return ctrl.Result{}, nil
}
return ctrl.Result{}, err
}
status, initialized, err := r.reconcile(ctx, instance)
if status == nil {
status = instance.Status.DeepCopy()
}
ready := meta.NewReadyCondition(instance)
if err != nil {
ready.Status = metav1.ConditionFalse
ready.Reason = meta.FailedReason
ready.Message = err.Error()
} else if !initialized {
ready.Status = metav1.ConditionFalse
ready.Reason = meta.ReconcilingReason
ready.Message = "waiting for ResourceQuota usage initialization"
}
status.Conditions.UpdateConditionByType(ready)
status.ObservedGeneration = instance.Generation
if updateErr := r.updateStatus(ctx, instance, *status); updateErr != nil {
return ctrl.Result{}, updateErr
}
instance.Status = *status
r.metrics.Record(instance)
return ctrl.Result{}, err
}
func (r *Controller) reconcile(
ctx context.Context,
instance *capsulev1beta2.GlobalResourceQuota,
) (*capsulev1beta2.GlobalResourceQuotaStatus, bool, error) {
namespaces, err := selectors.GetNamespacesMatchingSelectors(
ctx,
r.reader,
instance.Spec.NamespaceSelectors,
)
if err != nil {
return nil, false, err
}
if err := r.syncResourceQuotas(ctx, instance, namespaces); err != nil {
return nil, false, err
}
status, initialized, err := r.observeUsage(ctx, instance, namespaces)
if err != nil {
return nil, false, err
}
ledger, err := r.ensureLedger(ctx, instance)
if err != nil {
return status, false, err
}
if err := r.reconcileLedger(
ctx,
ledger,
instance.Generation,
status.Namespaces,
status.Total.Used,
initialized,
); err != nil {
return status, false, err
}
return status, initialized, nil
}
func (r *Controller) syncResourceQuotas(
ctx context.Context,
instance *capsulev1beta2.GlobalResourceQuota,
namespaces []corev1.Namespace,
) error {
selected := make(map[string]struct{}, len(namespaces))
for i := range namespaces {
namespace := &namespaces[i]
selected[namespace.Name] = struct{}{}
target := &corev1.ResourceQuota{
ObjectMeta: metav1.ObjectMeta{
Name: instance.GetResourceQuotaName(),
Namespace: namespace.Name,
},
}
if err := retry.RetryOnConflict(retry.DefaultBackoff, func() error {
_, err := controllerutil.CreateOrUpdate(ctx, r.Client, target, func() error {
targetLabels := target.GetLabels()
if targetLabels == nil {
targetLabels = map[string]string{}
}
targetLabels[meta.NewManagedByCapsuleLabel] = meta.ValueController
targetLabels[meta.GlobalResourceQuotaLabel] = instance.Name
target.SetLabels(targetLabels)
target.Spec = *instance.Spec.Quota.DeepCopy()
return controllerutil.SetControllerReference(instance, target, r.Scheme())
})
return err
}); err != nil {
if apierrors.HasStatusCause(err, corev1.NamespaceTerminatingCause) {
continue
}
return fmt.Errorf("sync ResourceQuota in namespace %s: %w", namespace.Name, err)
}
}
list := &corev1.ResourceQuotaList{}
if err := r.List(ctx, list, client.MatchingLabels{
meta.NewManagedByCapsuleLabel: meta.ValueController,
meta.GlobalResourceQuotaLabel: instance.Name,
}); err != nil {
return err
}
for i := range list.Items {
item := &list.Items[i]
if _, keep := selected[item.Namespace]; keep {
continue
}
if err := r.Delete(ctx, item); err != nil && !apierrors.IsNotFound(err) {
return fmt.Errorf("delete stale ResourceQuota %s/%s: %w", item.Namespace, item.Name, err)
}
}
return nil
}
func (r *Controller) observeUsage(
ctx context.Context,
instance *capsulev1beta2.GlobalResourceQuota,
namespaces []corev1.Namespace,
) (*capsulev1beta2.GlobalResourceQuotaStatus, bool, error) {
status := instance.Status.DeepCopy()
status.Total.Hard = instance.Spec.Quota.Hard.DeepCopy()
status.Total.Used = capsulev1beta2.ZeroResourceList(instance.Spec.Quota.Hard)
status.NamespaceUsage = make(capsulev1beta2.GlobalResourceQuotaNamespaceUsage, len(namespaces))
instanceCopy := instance.DeepCopy()
instanceCopy.Status = *status
instanceCopy.AssignNamespaces(namespaces)
status.Namespaces = instanceCopy.Status.Namespaces
status.NamespaceSize = instanceCopy.Status.NamespaceSize
initialized := true
for i := range namespaces {
namespace := namespaces[i].Name
used := capsulev1beta2.ZeroResourceList(instance.Spec.Quota.Hard)
quota := &corev1.ResourceQuota{}
err := r.reader.Get(ctx, types.NamespacedName{
Namespace: namespace,
Name: instance.GetResourceQuotaName(),
}, quota)
if err != nil {
if apierrors.IsNotFound(err) {
initialized = false
status.NamespaceUsage[namespace] = capsulev1beta2.GlobalResourceQuotaNamespaceStatus{Used: used}
continue
}
return status, false, err
}
if !resourceQuotaStatusReady(quota, instance.Spec.Quota.Hard) {
initialized = false
}
for name := range instance.Spec.Quota.Hard {
value := quota.Status.Used[name]
used[name] = value.DeepCopy()
total := status.Total.Used[name]
total.Add(value)
status.Total.Used[name] = total
}
status.NamespaceUsage[namespace] = capsulev1beta2.GlobalResourceQuotaNamespaceStatus{Used: used}
}
statusCopy := instance.DeepCopy()
statusCopy.Status = *status
statusCopy.CalculateAvailable()
return &statusCopy.Status, initialized, nil
}
func (r *Controller) updateStatus(
ctx context.Context,
instance *capsulev1beta2.GlobalResourceQuota,
status capsulev1beta2.GlobalResourceQuotaStatus,
) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
current := &capsulev1beta2.GlobalResourceQuota{}
if err := r.reader.Get(ctx, client.ObjectKeyFromObject(instance), current); err != nil {
return err
}
if reflect.DeepEqual(current.Status, status) {
return nil
}
current.Status = *status.DeepCopy()
return r.Status().Update(ctx, current)
})
}
func (r *Controller) globalQuotasForNamespace(ctx context.Context, object client.Object) []reconcile.Request {
namespace, ok := object.(*corev1.Namespace)
if !ok {
return nil
}
list := &capsulev1beta2.GlobalResourceQuotaList{}
if err := r.List(ctx, list); err != nil {
r.log.Error(err, "failed to list GlobalResourceQuotas", "namespace", namespace.Name)
return nil
}
requests := make([]reconcile.Request, 0)
for i := range list.Items {
item := &list.Items[i]
matched := slices.Contains(item.Status.Namespaces, namespace.Name)
for _, namespaceSelector := range item.Spec.NamespaceSelectors {
if namespaceSelector.LabelSelector == nil {
continue
}
selector, err := metav1.LabelSelectorAsSelector(namespaceSelector.LabelSelector)
if err == nil && selector.Matches(labels.Set(namespace.Labels)) {
matched = true
break
}
}
if matched {
requests = append(requests, reconcile.Request{NamespacedName: client.ObjectKeyFromObject(item)})
}
}
return requests
}
func resourceQuotaStatusReady(quota *corev1.ResourceQuota, hard corev1.ResourceList) bool {
for name := range hard {
if _, ok := quota.Status.Hard[name]; !ok {
return false
}
}
return true
}
@@ -0,0 +1,181 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package globalresourcequotas
import (
"context"
"testing"
"time"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/runtime/selectors"
)
func TestObserveUsageTracksTotalAndNamespaces(t *testing.T) {
t.Parallel()
scheme := runtime.NewScheme()
if err := corev1.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
if err := capsulev1beta2.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
hard := corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("8"),
corev1.ResourceRequestsMemory: resource.MustParse("16Gi"),
}
quota := &capsulev1beta2.GlobalResourceQuota{
ObjectMeta: metav1.ObjectMeta{Name: "shared", UID: types.UID("quota-uid")},
Spec: capsulev1beta2.GlobalResourceQuotaSpec{
Quota: corev1.ResourceQuotaSpec{Hard: hard},
},
}
namespaceA := corev1.Namespace{
ObjectMeta: metav1.ObjectMeta{Name: "a"},
Status: corev1.NamespaceStatus{Phase: corev1.NamespaceActive},
}
namespaceB := corev1.Namespace{
ObjectMeta: metav1.ObjectMeta{Name: "b"},
Status: corev1.NamespaceStatus{Phase: corev1.NamespaceActive},
}
resourceQuotaA := observedResourceQuota(quota, namespaceA.Name, corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("2"),
corev1.ResourceRequestsMemory: resource.MustParse("3Gi"),
})
resourceQuotaB := observedResourceQuota(quota, namespaceB.Name, corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("1"),
corev1.ResourceRequestsMemory: resource.MustParse("4Gi"),
})
cl := fake.NewClientBuilder().
WithScheme(scheme).
WithObjects(quota, resourceQuotaA, resourceQuotaB).
Build()
controller := &Controller{Client: cl, reader: cl}
status, initialized, err := controller.observeUsage(
context.Background(),
quota,
[]corev1.Namespace{namespaceB, namespaceA},
)
if err != nil {
t.Fatalf("observeUsage() error = %v", err)
}
if !initialized {
t.Fatal("observeUsage() initialized = false, want true")
}
if status.NamespaceSize != 2 || len(status.NamespaceUsage) != 2 {
t.Fatalf("namespace status = size %d usage %#v", status.NamespaceSize, status.NamespaceUsage)
}
if status.Namespaces[0] != "a" || status.Namespaces[1] != "b" {
t.Fatalf("ordered namespaces = %#v, want [a b]", status.Namespaces)
}
assertResource(t, status.Total.Used, corev1.ResourceRequestsCPU, "3")
assertResource(t, status.Total.Used, corev1.ResourceRequestsMemory, "7Gi")
assertResource(t, status.Total.Available, corev1.ResourceRequestsCPU, "5")
assertResource(t, status.NamespaceUsage["b"].Used, corev1.ResourceRequestsMemory, "4Gi")
}
func TestReconcileLedgerStatusConsumesObservedUsage(t *testing.T) {
t.Parallel()
now := metav1.Now()
expires := metav1.NewTime(now.Add(time.Minute))
current := &capsulev1beta2.QuantityLedgerResourceQuotaStatus{
Used: corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("0")},
Reservations: []capsulev1beta2.QuantityLedgerResourceQuotaReservation{{
ID: "request",
Delta: corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("2")},
ExpiresAt: &expires,
}},
}
next := reconcileLedgerStatus(
current,
3,
[]string{"a", "b"},
corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("1")},
true,
now,
)
if len(next.Reservations) != 1 {
t.Fatalf("reservations = %d, want 1", len(next.Reservations))
}
assertResource(t, next.Reservations[0].Delta, corev1.ResourceRequestsCPU, "1")
assertResource(t, next.Used, corev1.ResourceRequestsCPU, "1")
assertResource(t, next.Allocated, corev1.ResourceRequestsCPU, "2")
if next.ObservedGeneration != 3 || len(next.Namespaces) != 2 {
t.Fatalf("ledger snapshot = generation %d, namespaces %#v", next.ObservedGeneration, next.Namespaces)
}
}
func TestMatchingNamespaceSelectorsUseOR(t *testing.T) {
t.Parallel()
quota := &capsulev1beta2.GlobalResourceQuota{
Spec: capsulev1beta2.GlobalResourceQuotaSpec{
NamespaceSelectors: []selectors.NamespaceSelector{
{LabelSelector: &metav1.LabelSelector{MatchLabels: map[string]string{"team": "a"}}},
{LabelSelector: &metav1.LabelSelector{MatchLabels: map[string]string{"team": "b"}}},
},
},
}
namespace := &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{
Name: "b", Labels: map[string]string{"team": "b"},
}}
cl := fake.NewClientBuilder().WithObjects(namespace).Build()
matched, err := selectors.GetNamespacesMatchingSelectors(
context.Background(),
cl,
quota.Spec.NamespaceSelectors,
)
if err != nil {
t.Fatal(err)
}
if len(matched) != 1 || matched[0].Name != namespace.Name {
t.Fatalf("matched namespaces = %#v, want b", matched)
}
}
func observedResourceQuota(
quota *capsulev1beta2.GlobalResourceQuota,
namespace string,
used corev1.ResourceList,
) *corev1.ResourceQuota {
return &corev1.ResourceQuota{
ObjectMeta: metav1.ObjectMeta{
Name: quota.GetResourceQuotaName(),
Namespace: namespace,
},
Status: corev1.ResourceQuotaStatus{
Hard: quota.Spec.Quota.Hard.DeepCopy(),
Used: used.DeepCopy(),
},
}
}
func assertResource(
t *testing.T,
resources corev1.ResourceList,
name corev1.ResourceName,
want string,
) {
t.Helper()
got := resources[name]
if got.Cmp(resource.MustParse(want)) != 0 {
t.Fatalf("%s = %s, want %s", name, got.String(), want)
}
}
@@ -0,0 +1,195 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package globalresourcequotas
import (
"context"
"reflect"
"slices"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/runtime/configuration"
runtimequota "github.com/projectcapsule/capsule/pkg/runtime/quota"
)
func (r *Controller) ensureLedger(
ctx context.Context,
quota *capsulev1beta2.GlobalResourceQuota,
) (*capsulev1beta2.QuantityLedger, error) {
ledger := &capsulev1beta2.QuantityLedger{
ObjectMeta: metav1.ObjectMeta{
Name: quota.GetLedgerName(),
Namespace: configuration.ControllerNamespace(),
},
}
_, err := controllerutil.CreateOrUpdate(ctx, r.Client, ledger, func() error {
ledgerLabels := ledger.GetLabels()
if ledgerLabels == nil {
ledgerLabels = map[string]string{}
}
ledgerLabels[meta.NewManagedByCapsuleLabel] = meta.ValueController
ledgerLabels[meta.GlobalResourceQuotaLabel] = quota.Name
ledger.SetLabels(ledgerLabels)
ledger.Spec.TargetRef = capsulev1beta2.QuantityLedgerTargetRef{
APIGroup: capsulev1beta2.GroupVersion.Group,
Kind: "GlobalResourceQuota",
Name: quota.Name,
UID: quota.UID,
}
return controllerutil.SetControllerReference(quota, ledger, r.Scheme())
})
if err != nil {
return nil, err
}
return ledger, nil
}
func (r *Controller) reconcileLedger(
ctx context.Context,
ledger *capsulev1beta2.QuantityLedger,
generation int64,
namespaces []string,
used corev1.ResourceList,
initialized bool,
) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
current := &capsulev1beta2.QuantityLedger{}
if err := r.reader.Get(ctx, types.NamespacedName{
Namespace: ledger.Namespace,
Name: ledger.Name,
}, current); err != nil {
return err
}
before := current.Status.ResourceQuota.DeepCopy()
next := reconcileLedgerStatus(
current.Status.ResourceQuota,
generation,
namespaces,
used,
initialized,
metav1.Now(),
)
if reflect.DeepEqual(before, next) {
return nil
}
current.Status.ResourceQuota = next
return r.Status().Update(ctx, current)
})
}
func reconcileLedgerStatus(
current *capsulev1beta2.QuantityLedgerResourceQuotaStatus,
generation int64,
namespaces []string,
used corev1.ResourceList,
initialized bool,
now metav1.Time,
) *capsulev1beta2.QuantityLedgerResourceQuotaStatus {
if current == nil {
current = &capsulev1beta2.QuantityLedgerResourceQuotaStatus{}
} else {
current = current.DeepCopy()
}
increase := positiveDifference(used, current.Used)
active := make([]capsulev1beta2.QuantityLedgerResourceQuotaReservation, 0, len(current.Reservations))
for _, reservation := range current.Reservations {
if reservation.ExpiresAt != nil && reservation.ExpiresAt.Before(&now) {
continue
}
reservation.Delta = consumeResourceList(reservation.Delta, increase)
if resourceListPositive(reservation.Delta) {
active = append(active, reservation)
}
}
reserved := capsulev1beta2.ZeroResourceList(used)
for _, reservation := range active {
addResourceList(reserved, reservation.Delta)
}
allocated := used.DeepCopy()
addResourceList(allocated, reserved)
current.ObservedGeneration = generation
current.Initialized = initialized
current.Namespaces = slices.Clone(namespaces)
current.Used = used.DeepCopy()
current.Reserved = reserved
current.Allocated = allocated
current.Reservations = active
return current
}
func positiveDifference(next, previous corev1.ResourceList) corev1.ResourceList {
out := make(corev1.ResourceList, len(next))
for name, quantity := range next {
delta := quantity.DeepCopy()
delta.Sub(previous[name])
runtimequota.ClampQuantityToZero(&delta)
out[name] = delta
}
return out
}
func consumeResourceList(delta, available corev1.ResourceList) corev1.ResourceList {
out := delta.DeepCopy()
for name, quantity := range out {
increase := available[name]
if quantity.Sign() <= 0 || increase.Sign() <= 0 {
continue
}
consumed := quantity.DeepCopy()
if consumed.Cmp(increase) > 0 {
consumed = increase.DeepCopy()
}
quantity.Sub(consumed)
out[name] = quantity
increase.Sub(consumed)
available[name] = increase
}
return out
}
func addResourceList(target, addition corev1.ResourceList) {
for name, quantity := range addition {
current := target[name]
current.Add(quantity)
target[name] = current
}
}
func resourceListPositive(resources corev1.ResourceList) bool {
for _, quantity := range resources {
if quantity.Sign() > 0 {
return true
}
}
return false
}
@@ -0,0 +1,34 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package globalresourcequotas
import (
"fmt"
"github.com/go-logr/logr"
"k8s.io/client-go/tools/events"
"sigs.k8s.io/controller-runtime/pkg/manager"
"github.com/projectcapsule/capsule/internal/controllers/utils"
"github.com/projectcapsule/capsule/internal/metrics"
)
func Add(
log logr.Logger,
mgr manager.Manager,
recorder events.EventRecorder,
cfg utils.ControllerOptions,
) error {
controller := &Controller{
Client: mgr.GetClient(),
log: log,
recorder: recorder,
metrics: metrics.MustMakeGlobalResourceQuotaRecorder(),
}
if err := controller.SetupWithManager(mgr, cfg); err != nil {
return fmt.Errorf("unable to create GlobalResourceQuota controller: %w", err)
}
return nil
}
+32 -3
View File
@@ -157,13 +157,17 @@ func (r Manager) reconcile(ctx context.Context, instance *capsulev1beta2.RuleSta
continue
}
enforce := rule.Enforce.DeepCopy()
statusRule := rule.DeepCopy()
// RuleStatus is an enforcement cache. Quota definitions are reconciled
// independently as GlobalResourceQuotas and may include legacy entries
// which predate stable quota names.
statusRule.Quota = nil
enforce := rule.Enforce.DeepCopy()
for i := range enforce.Metadata {
enforce.Metadata[i].APIGroups = enforce.Metadata[i].StatusAPIGroups()
}
statusRule := rule.DeepCopy()
statusRule.Enforce = enforce
ruleStatus = append(ruleStatus, statusRule)
}
@@ -248,8 +252,13 @@ func (r *Manager) updateReconcilingStatus(ctx context.Context, instance *capsule
return err
}
cleanedQuota := removeQuotaDefinitions(&latest.Status)
if latest.Status.ObservedGeneration == instance.GetGeneration() {
return nil
if !cleanedQuota {
return nil
}
return r.Status().Update(ctx, latest)
}
latest.Status.Conditions.UpdateConditionByType(meta.NewReadyConditionReconcilingReason(instance))
@@ -257,3 +266,23 @@ func (r *Manager) updateReconcilingStatus(ctx context.Context, instance *capsule
return r.Status().Update(ctx, latest)
})
}
//nolint:staticcheck // The deprecated flattened Rule must be cleaned for objects written by older Capsule versions.
func removeQuotaDefinitions(status *capsulev1beta2.RuleStatusStatus) bool {
if status == nil {
return false
}
changed := len(status.Rule.Quota) > 0
for _, rule := range status.Rules {
if rule != nil && len(rule.Quota) > 0 {
changed = true
rule.Quota = nil
}
}
status.Rule.Quota = nil
return changed
}
@@ -0,0 +1,70 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package rulestatus
import (
"context"
"testing"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/rules"
)
func TestReconcileExcludesQuotaFromRuleStatus(t *testing.T) {
t.Parallel()
unnamedQuota := rules.ResourceQuotaRule{
ResourceQuotaSpec: corev1.ResourceQuotaSpec{Hard: corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("1"),
}},
}
instance := &capsulev1beta2.RuleStatus{
Spec: []*rules.NamespaceRuleBodyNamespace{
{Quota: []rules.ResourceQuotaRule{unnamedQuota}},
{
Quota: []rules.ResourceQuotaRule{unnamedQuota},
Enforce: &rules.NamespaceRuleEnforceBody{Action: rules.ActionTypeDeny},
},
},
}
if err := (Manager{}).reconcile(context.Background(), instance); err != nil {
t.Fatalf("reconcile() error = %v", err)
}
if len(instance.Status.Rules) != 1 {
t.Fatalf("status rules = %d, want one enforcement rule", len(instance.Status.Rules))
}
if len(instance.Status.Rules[0].Quota) != 0 {
t.Fatalf("status quota = %#v, want none", instance.Status.Rules[0].Quota)
}
if instance.Status.Rules[0].Enforce == nil {
t.Fatal("status enforcement rule was removed")
}
}
func TestRemoveQuotaDefinitionsCleansLegacyStatus(t *testing.T) {
t.Parallel()
status := &capsulev1beta2.RuleStatusStatus{
Rule: rules.NamespaceRuleBodyNamespace{Quota: []rules.ResourceQuotaRule{{}}},
Rules: []*rules.NamespaceRuleBodyNamespace{
nil,
{Quota: []rules.ResourceQuotaRule{{}}},
},
}
if changed := removeQuotaDefinitions(status); !changed {
t.Fatal("legacy quota definitions were not reported as changed")
}
if len(status.Rule.Quota) != 0 || len(status.Rules[1].Quota) != 0 {
t.Fatalf("legacy quota definitions were not removed: %#v", status)
}
if changed := removeQuotaDefinitions(status); changed {
t.Fatal("clean status was reported as changed")
}
}
@@ -0,0 +1,119 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package tenant
import (
"context"
"fmt"
"strings"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/labels"
k8svalidation "k8s.io/apimachinery/pkg/util/validation"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/meta"
tenantutils "github.com/projectcapsule/capsule/pkg/tenant"
)
func (r *Manager) syncGlobalResourceQuotas(
ctx context.Context,
tnt *capsulev1beta2.Tenant,
) error {
desired := make(map[string]struct{})
quotaNames := make(map[string]struct{})
for ruleIndex, rule := range tnt.Spec.Rules {
if rule == nil || rule.NamespaceRuleBodyNamespace == nil {
continue
}
for itemIndex := range rule.Quota {
quotaName := rule.Quota[itemIndex].Name
if errs := k8svalidation.IsDNS1123Label(quotaName); len(errs) > 0 {
return fmt.Errorf(
"rules[%d].quota[%d].name %q is invalid: %s",
ruleIndex,
itemIndex,
quotaName,
strings.Join(errs, "; "),
)
}
if _, duplicate := quotaNames[quotaName]; duplicate {
return fmt.Errorf("rules[%d].quota[%d].name %q is duplicated", ruleIndex, itemIndex, quotaName)
}
quotaNames[quotaName] = struct{}{}
if err := tenantutils.ValidateRuleGlobalResourceQuotaName(tnt, quotaName); err != nil {
return fmt.Errorf("rules[%d].quota[%d]: %w", ruleIndex, itemIndex, err)
}
target := tenantutils.RuleGlobalResourceQuota(tnt, ruleIndex, itemIndex)
desired[target.Name] = struct{}{}
if err := retry.RetryOnConflict(retry.DefaultBackoff, func() error {
_, err := controllerutil.CreateOrUpdate(ctx, r.Client, target, func() error {
currentLabels := target.GetLabels()
if currentLabels == nil {
currentLabels = map[string]string{}
}
currentLabels[meta.NewManagedByCapsuleLabel] = meta.ValueController
currentLabels[meta.NewTenantLabel] = tnt.Name
currentLabels[meta.RuleQuotaLabel] = quotaName
target.SetLabels(currentLabels)
desiredQuota := tenantutils.RuleGlobalResourceQuota(tnt, ruleIndex, itemIndex)
target.Spec = *desiredQuota.Spec.DeepCopy()
return controllerutil.SetControllerReference(tnt, target, r.Scheme())
})
return err
}); err != nil {
return fmt.Errorf("sync GlobalResourceQuota %s: %w", target.Name, err)
}
}
}
return r.pruneGlobalResourceQuotas(ctx, tnt, desired)
}
func (r *Manager) pruneGlobalResourceQuotas(
ctx context.Context,
tnt *capsulev1beta2.Tenant,
desired map[string]struct{},
) error {
list := &capsulev1beta2.GlobalResourceQuotaList{}
selector := labels.SelectorFromSet(labels.Set{
meta.NewManagedByCapsuleLabel: meta.ValueController,
meta.NewTenantLabel: tnt.Name,
})
if err := r.List(ctx, list, &client.ListOptions{LabelSelector: selector}); err != nil {
return err
}
for i := range list.Items {
item := &list.Items[i]
if _, generated := item.Labels[meta.RuleQuotaLabel]; !generated {
continue
}
if _, keep := desired[item.Name]; keep {
continue
}
if err := r.Delete(ctx, item); err != nil && !apierrors.IsNotFound(err) {
return fmt.Errorf("delete stale GlobalResourceQuota %s: %w", item.Name, err)
}
}
return nil
}
@@ -0,0 +1,130 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package tenant
import (
"context"
"testing"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/api/rules"
tenantutils "github.com/projectcapsule/capsule/pkg/tenant"
)
func TestSyncGlobalResourceQuotasGeneratesAndPrunesRuleQuotas(t *testing.T) {
t.Parallel()
scheme := runtime.NewScheme()
if err := capsulev1beta2.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
tnt := &capsulev1beta2.Tenant{
ObjectMeta: metav1.ObjectMeta{Name: "tenant-a", UID: types.UID("tenant-uid")},
Spec: capsulev1beta2.TenantSpec{Rules: []*rules.NamespaceRuleBodyTenant{{
NamespaceRuleBodyNamespace: &rules.NamespaceRuleBodyNamespace{
Quota: []rules.ResourceQuotaRule{{
Name: "shared-compute",
ResourceQuotaSpec: corev1.ResourceQuotaSpec{Hard: corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("8"),
}},
}},
},
NamespaceSelector: &metav1.LabelSelector{MatchLabels: map[string]string{"tier": "paid"}},
}}},
}
cl := fake.NewClientBuilder().WithScheme(scheme).WithObjects(tnt).Build()
manager := &Manager{Client: cl}
if err := manager.syncGlobalResourceQuotas(context.Background(), tnt); err != nil {
t.Fatalf("syncGlobalResourceQuotas() error = %v", err)
}
key := client.ObjectKey{Name: tenantutils.RuleGlobalResourceQuotaName(tnt, "shared-compute")}
generated := &capsulev1beta2.GlobalResourceQuota{}
if err := cl.Get(context.Background(), key, generated); err != nil {
t.Fatalf("get generated GlobalResourceQuota: %v", err)
}
if generated.Labels[meta.RuleQuotaLabel] != "shared-compute" {
t.Fatalf("rule quota label = %q, want shared-compute", generated.Labels[meta.RuleQuotaLabel])
}
selector := generated.Spec.NamespaceSelectors[0].LabelSelector
if selector.MatchLabels[meta.TenantLabel] != tnt.Name || selector.MatchLabels["tier"] != "paid" {
t.Fatalf("generated selector = %#v", selector)
}
generated.Annotations = map[string]string{"test.projectcapsule.dev/identity": "preserved"}
if err := cl.Update(context.Background(), generated); err != nil {
t.Fatalf("annotate generated GlobalResourceQuota: %v", err)
}
sharedRule := tnt.Spec.Rules[0]
sharedRule.NamespaceSelector = &metav1.LabelSelector{MatchLabels: map[string]string{"tier": "enterprise"}}
sharedRule.Quota[0].Hard[corev1.ResourceRequestsCPU] = resource.MustParse("12")
tnt.Spec.Rules = []*rules.NamespaceRuleBodyTenant{
{
NamespaceRuleBodyNamespace: &rules.NamespaceRuleBodyNamespace{
Quota: []rules.ResourceQuotaRule{{
Name: "service-count",
ResourceQuotaSpec: corev1.ResourceQuotaSpec{Hard: corev1.ResourceList{
corev1.ResourceServices: resource.MustParse("5"),
}},
}},
},
},
sharedRule,
}
if err := manager.syncGlobalResourceQuotas(context.Background(), tnt); err != nil {
t.Fatalf("sync reordered GlobalResourceQuotas: %v", err)
}
updated := &capsulev1beta2.GlobalResourceQuota{}
if err := cl.Get(context.Background(), key, updated); err != nil {
t.Fatalf("get stable GlobalResourceQuota after reorder: %v", err)
}
if updated.Annotations["test.projectcapsule.dev/identity"] != "preserved" {
t.Fatal("generated GlobalResourceQuota was replaced after rule reorder")
}
updatedSelector := updated.Spec.NamespaceSelectors[0].LabelSelector
if updatedSelector.MatchLabels["tier"] != "enterprise" {
t.Fatalf("updated selector = %#v, want tier=enterprise", updatedSelector)
}
if got := updated.Spec.Quota.Hard[corev1.ResourceRequestsCPU]; got.Cmp(resource.MustParse("12")) != 0 {
t.Fatalf("updated hard requests.cpu = %s, want 12", got.String())
}
list := &capsulev1beta2.GlobalResourceQuotaList{}
if err := cl.List(context.Background(), list); err != nil {
t.Fatalf("list generated GlobalResourceQuotas: %v", err)
}
if len(list.Items) != 2 {
t.Fatalf("generated GlobalResourceQuotas = %d, want 2", len(list.Items))
}
tnt.Spec.Rules = nil
if err := manager.syncGlobalResourceQuotas(context.Background(), tnt); err != nil {
t.Fatalf("prune GlobalResourceQuota: %v", err)
}
err := cl.Get(context.Background(), key, &capsulev1beta2.GlobalResourceQuota{})
if !apierrors.IsNotFound(err) {
t.Fatalf("get pruned GlobalResourceQuota error = %v, want NotFound", err)
}
list = &capsulev1beta2.GlobalResourceQuotaList{}
if err := cl.List(context.Background(), list); err != nil {
t.Fatalf("list pruned GlobalResourceQuotas: %v", err)
}
if len(list.Items) != 0 {
t.Fatalf("GlobalResourceQuotas after prune = %d, want 0", len(list.Items))
}
}
+14
View File
@@ -83,6 +83,14 @@ func (r *Manager) SetupWithManager(mgr ctrl.Manager, ctrlConfig utils.Controller
).
Owns(&networkingv1.NetworkPolicy{}, builder.WithPredicates(predicates.TenantManagedResourceChangedPredicate{})).
Owns(&corev1.LimitRange{}, builder.WithPredicates(predicates.TenantManagedResourceChangedPredicate{})).
Owns(
&capsulev1beta2.GlobalResourceQuota{},
builder.WithPredicates(predicate.Or(
predicate.GenerationChangedPredicate{},
predicates.UpdatedMetadataPredicate{},
predicates.DeletionChangedPredicate{},
)),
).
Watches(
&corev1.ResourceQuota{},
handler.Funcs{
@@ -403,6 +411,12 @@ func (r *Manager) reconcile(ctx context.Context, log logr.Logger, instance *caps
errs = append(errs, fmt.Errorf("cannot sync resourcequota items: %w", err))
}
log.V(4).Info("starting processing of rule GlobalResourceQuotas")
if err = r.syncGlobalResourceQuotas(ctx, instance); err != nil {
errs = append(errs, fmt.Errorf("cannot sync rule global resource quotas: %w", err))
}
log.V(4).Info("ensuring RoleBindings for Owners and Tenant")
if err = r.syncRoleBindings(ctx, log, instance); err != nil {
+3
View File
@@ -42,6 +42,9 @@ func (r Reconciler) managedCRDs() map[string]ManagedCRD {
"globalcustomquotas": {
Name: "globalcustomquotas.capsule.clastix.io",
},
"globalresourcequotas": {
Name: "globalresourcequotas.capsule.clastix.io",
},
"globaltenantresources": {
Name: "globaltenantresources.capsule.clastix.io",
},
@@ -0,0 +1,148 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package metrics
import (
"github.com/prometheus/client_golang/prometheus"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
crtlmetrics "sigs.k8s.io/controller-runtime/pkg/metrics"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/meta"
)
type GlobalResourceQuotaRecorder struct {
ConditionGauge *prometheus.GaugeVec
ResourceLimitGauge *prometheus.GaugeVec
ResourceUsageGauge *prometheus.GaugeVec
ResourceAvailableGauge *prometheus.GaugeVec
ResourceUsagePercentageGauge *prometheus.GaugeVec
NamespaceUsageGauge *prometheus.GaugeVec
NamespaceUsagePercentageGauge *prometheus.GaugeVec
}
func MustMakeGlobalResourceQuotaRecorder() *GlobalResourceQuotaRecorder {
recorder := NewGlobalResourceQuotaRecorder()
crtlmetrics.Registry.MustRegister(recorder.Collectors()...)
return recorder
}
func NewGlobalResourceQuotaRecorder() *GlobalResourceQuotaRecorder {
const label = "global_resource_quota"
return &GlobalResourceQuotaRecorder{
ConditionGauge: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: metricsPrefix,
Name: "global_resource_quota_condition",
Help: "Current condition for a GlobalResourceQuota.",
}, []string{label, "condition"}),
ResourceLimitGauge: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: metricsPrefix,
Name: "global_resource_quota_limit",
Help: "Shared hard limit for a GlobalResourceQuota resource.",
}, []string{label, "resource"}),
ResourceUsageGauge: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: metricsPrefix,
Name: "global_resource_quota_usage",
Help: "Observed aggregate usage for a GlobalResourceQuota resource.",
}, []string{label, "resource"}),
ResourceAvailableGauge: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: metricsPrefix,
Name: "global_resource_quota_available",
Help: "Available aggregate capacity for a GlobalResourceQuota resource.",
}, []string{label, "resource"}),
ResourceUsagePercentageGauge: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: metricsPrefix,
Name: "global_resource_quota_usage_percentage",
Help: "Observed aggregate usage percentage for a GlobalResourceQuota resource.",
}, []string{label, "resource"}),
NamespaceUsageGauge: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: metricsPrefix,
Name: "global_resource_quota_namespace_usage",
Help: "Observed usage per namespace for a GlobalResourceQuota resource.",
}, []string{label, "target_namespace", "resource"}),
NamespaceUsagePercentageGauge: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: metricsPrefix,
Name: "global_resource_quota_namespace_usage_percentage",
Help: "Observed per-namespace usage as a percentage of the shared limit.",
}, []string{label, "target_namespace", "resource"}),
}
}
func (r *GlobalResourceQuotaRecorder) Collectors() []prometheus.Collector {
return []prometheus.Collector{
r.ConditionGauge,
r.ResourceLimitGauge,
r.ResourceUsageGauge,
r.ResourceAvailableGauge,
r.ResourceUsagePercentageGauge,
r.NamespaceUsageGauge,
r.NamespaceUsagePercentageGauge,
}
}
func (r *GlobalResourceQuotaRecorder) Record(quota *capsulev1beta2.GlobalResourceQuota) {
r.Delete(quota.Name)
for _, conditionType := range []string{meta.ReadyCondition} {
condition := quota.Status.Conditions.GetConditionByType(conditionType)
if condition == nil {
continue
}
value := float64(0)
if condition.Status == metav1.ConditionTrue {
value = 1
}
r.ConditionGauge.WithLabelValues(quota.Name, conditionType).Set(value)
}
for name, hard := range quota.Status.Total.Hard {
used := quota.Status.Total.Used[name]
available := quota.Status.Total.Available[name]
r.ResourceLimitGauge.WithLabelValues(quota.Name, name.String()).Set(quantityMetric(hard))
r.ResourceUsageGauge.WithLabelValues(quota.Name, name.String()).Set(quantityMetric(used))
r.ResourceAvailableGauge.WithLabelValues(quota.Name, name.String()).Set(quantityMetric(available))
r.ResourceUsagePercentageGauge.WithLabelValues(quota.Name, name.String()).Set(quantityPercentage(used, hard))
}
for namespace, usage := range quota.Status.NamespaceUsage {
for name, used := range usage.Used {
hard := quota.Status.Total.Hard[name]
r.NamespaceUsageGauge.WithLabelValues(quota.Name, namespace, name.String()).Set(quantityMetric(used))
r.NamespaceUsagePercentageGauge.WithLabelValues(
quota.Name,
namespace,
name.String(),
).Set(quantityPercentage(used, hard))
}
}
}
func (r *GlobalResourceQuotaRecorder) Delete(name string) {
labels := prometheus.Labels{"global_resource_quota": name}
r.ConditionGauge.DeletePartialMatch(labels)
r.ResourceLimitGauge.DeletePartialMatch(labels)
r.ResourceUsageGauge.DeletePartialMatch(labels)
r.ResourceAvailableGauge.DeletePartialMatch(labels)
r.ResourceUsagePercentageGauge.DeletePartialMatch(labels)
r.NamespaceUsageGauge.DeletePartialMatch(labels)
r.NamespaceUsagePercentageGauge.DeletePartialMatch(labels)
}
func quantityMetric(quantity resource.Quantity) float64 {
return float64(quantity.MilliValue()) / 1000
}
func quantityPercentage(used, hard resource.Quantity) float64 {
if hard.MilliValue() <= 0 {
return 0
}
return float64(used.MilliValue()) / float64(hard.MilliValue()) * 100
}
@@ -0,0 +1,94 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package metrics
import (
"testing"
"github.com/prometheus/client_golang/prometheus"
dto "github.com/prometheus/client_model/go"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
)
func TestGlobalResourceQuotaRecorderTracksAggregateAndNamespaceUsage(t *testing.T) {
t.Parallel()
recorder := NewGlobalResourceQuotaRecorder()
quota := &capsulev1beta2.GlobalResourceQuota{}
quota.Name = "shared"
quota.Status.Total = capsulev1beta2.GlobalResourceQuotaUsage{
Hard: corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("8")},
Used: corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("2")},
Available: corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("6")},
}
quota.Status.NamespaceUsage = capsulev1beta2.GlobalResourceQuotaNamespaceUsage{
"tenant-a": {
Used: corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("1.5")},
},
}
recorder.Record(quota)
assertGauge(t, recorder.ResourceLimitGauge, 8, "shared", string(corev1.ResourceRequestsCPU))
assertGauge(t, recorder.ResourceUsageGauge, 2, "shared", string(corev1.ResourceRequestsCPU))
assertGauge(t, recorder.ResourceAvailableGauge, 6, "shared", string(corev1.ResourceRequestsCPU))
assertGauge(t, recorder.ResourceUsagePercentageGauge, 25, "shared", string(corev1.ResourceRequestsCPU))
assertGauge(
t,
recorder.NamespaceUsageGauge,
1.5,
"shared",
"tenant-a",
string(corev1.ResourceRequestsCPU),
)
assertGauge(
t,
recorder.NamespaceUsagePercentageGauge,
18.75,
"shared",
"tenant-a",
string(corev1.ResourceRequestsCPU),
)
recorder.Delete(quota.Name)
if got := metricCount(recorder.ResourceUsageGauge); got != 0 {
t.Fatalf("usage metric count after delete = %d, want 0", got)
}
}
type gaugeMetric interface {
GetMetricWithLabelValues(lvs ...string) (prometheus.Gauge, error)
}
func assertGauge(t *testing.T, gauge gaugeMetric, want float64, labels ...string) {
t.Helper()
metric, err := gauge.GetMetricWithLabelValues(labels...)
if err != nil {
t.Fatal(err)
}
value := &dto.Metric{}
if err := metric.Write(value); err != nil {
t.Fatal(err)
}
if got := value.GetGauge().GetValue(); got != want {
t.Fatalf("metric %v = %v, want %v", labels, got, want)
}
}
func metricCount(collector prometheus.Collector) int {
metrics := make(chan prometheus.Metric, 32)
collector.Collect(metrics)
close(metrics)
count := 0
for range metrics {
count++
}
return count
}
+637
View File
@@ -0,0 +1,637 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
/*
Copyright 2016 The Kubernetes Authors.
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 evaluator contains the admission-side resource calculations used by
// GlobalResourceQuota. The calculations are adapted from the Kubernetes
// ResourceQuota core evaluators at the version matching this module's
// Kubernetes dependencies.
package evaluator
import (
"encoding/json"
"fmt"
"slices"
"strings"
"time"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apimachinery/pkg/selection"
"k8s.io/apimachinery/pkg/util/sets"
"k8s.io/apiserver/pkg/quota/v1/generic"
resourcehelper "k8s.io/component-helpers/resource"
storagehelpers "k8s.io/component-helpers/storage/volume"
"k8s.io/utils/ptr"
"sigs.k8s.io/controller-runtime/pkg/webhook/admission"
)
const storageClassSuffix = ".storageclass.storage.k8s.io/"
var validationResources = sets.New(
corev1.ResourceCPU,
corev1.ResourceMemory,
corev1.ResourceRequestsCPU,
corev1.ResourceRequestsMemory,
corev1.ResourceLimitsCPU,
corev1.ResourceLimitsMemory,
)
var legacyObjectCountAliases = map[string]corev1.ResourceName{
"configmaps": corev1.ResourceConfigMaps,
"resourcequotas": corev1.ResourceQuotas,
"replicationcontrollers": corev1.ResourceReplicationControllers,
"secrets": corev1.ResourceSecrets,
}
// Result holds one decoded admission request and its native quota usage.
type Result struct {
NewUsage corev1.ResourceList
OldUsage corev1.ResourceList
New runtime.Object
Old runtime.Object
}
// Evaluate decodes an admission request once and calculates the same native
// resource names used by the Kubernetes core quota evaluators.
func Evaluate(req admission.Request) (Result, bool, error) {
if req.SubResource != "" && req.SubResource != "resize" && req.SubResource != "status" {
return Result{}, false, nil
}
resourceName := req.Resource.Resource
switch {
case req.Resource.Group == "" && resourceName == "pods":
return evaluatePod(req)
case req.Resource.Group == "" && resourceName == "services":
return evaluateService(req)
case req.Resource.Group == "" && resourceName == "persistentvolumeclaims":
return evaluatePVC(req)
default:
if req.SubResource != "" || req.Operation != "CREATE" {
return Result{}, false, nil
}
object := &metav1.PartialObjectMetadata{}
if err := json.Unmarshal(req.Object.Raw, object); err != nil {
return Result{}, false, err
}
usage := objectCountUsage(req.Resource.Group, resourceName)
return Result{NewUsage: usage, New: object}, true, nil
}
}
func evaluatePod(req admission.Request) (Result, bool, error) {
if req.Operation != "CREATE" && req.Operation != "UPDATE" {
return Result{}, false, nil
}
newPod := &corev1.Pod{}
if err := json.Unmarshal(req.Object.Raw, newPod); err != nil {
return Result{}, false, fmt.Errorf("decode Pod: %w", err)
}
newUsage := podUsage(newPod, time.Now())
result := Result{NewUsage: newUsage, New: newPod}
if req.Operation == "UPDATE" {
oldPod := &corev1.Pod{}
if err := json.Unmarshal(req.OldObject.Raw, oldPod); err != nil {
return Result{}, false, fmt.Errorf("decode old Pod: %w", err)
}
result.Old = oldPod
result.OldUsage = podUsage(oldPod, time.Now())
}
return result, true, nil
}
func evaluateService(req admission.Request) (Result, bool, error) {
if req.SubResource != "" || (req.Operation != "CREATE" && req.Operation != "UPDATE") {
return Result{}, false, nil
}
service := &corev1.Service{}
if err := json.Unmarshal(req.Object.Raw, service); err != nil {
return Result{}, false, fmt.Errorf("decode Service: %w", err)
}
result := Result{NewUsage: serviceUsage(service), New: service}
if req.Operation == "UPDATE" {
oldService := &corev1.Service{}
if err := json.Unmarshal(req.OldObject.Raw, oldService); err != nil {
return Result{}, false, fmt.Errorf("decode old Service: %w", err)
}
result.Old = oldService
result.OldUsage = serviceUsage(oldService)
}
return result, true, nil
}
func evaluatePVC(req admission.Request) (Result, bool, error) {
if (req.SubResource != "" && req.SubResource != "status") ||
(req.Operation != "CREATE" && req.Operation != "UPDATE") {
return Result{}, false, nil
}
pvc := &corev1.PersistentVolumeClaim{}
if err := json.Unmarshal(req.Object.Raw, pvc); err != nil {
return Result{}, false, fmt.Errorf("decode PersistentVolumeClaim: %w", err)
}
result := Result{NewUsage: pvcUsage(pvc), New: pvc}
if req.Operation == "UPDATE" {
oldPVC := &corev1.PersistentVolumeClaim{}
if err := json.Unmarshal(req.OldObject.Raw, oldPVC); err != nil {
return Result{}, false, fmt.Errorf("decode old PersistentVolumeClaim: %w", err)
}
result.Old = oldPVC
result.OldUsage = pvcUsage(oldPVC)
}
return result, true, nil
}
func objectCountUsage(group, resourceName string) corev1.ResourceList {
countName := generic.ObjectCountQuotaResourceNameFor(
schema.GroupResource{Group: group, Resource: resourceName},
)
one := *resource.NewQuantity(1, resource.DecimalSI)
result := corev1.ResourceList{countName: one}
if group == "" {
if alias, ok := legacyObjectCountAliases[resourceName]; ok {
result[alias] = one
}
}
return result
}
func podUsage(pod *corev1.Pod, now time.Time) corev1.ResourceList {
result := objectCountUsage("", "pods")
if !quotaPod(pod, now) {
return result
}
opts := resourcehelper.PodResourcesOptions{
UseStatusResources: true,
SkipPodLevelResources: false,
}
requests := resourcehelper.PodRequests(pod, opts)
limits := resourcehelper.PodLimits(pod, opts)
addResourceList(result, podComputeUsage(requests, limits))
return result
}
func podComputeUsage(requests, limits corev1.ResourceList) corev1.ResourceList {
result := corev1.ResourceList{
corev1.ResourcePods: *resource.NewQuantity(1, resource.DecimalSI),
}
addRequest := func(name, plain, prefixed corev1.ResourceName) {
if quantity, found := requests[name]; found {
result[plain] = quantity
result[prefixed] = quantity
}
}
addRequest(corev1.ResourceCPU, corev1.ResourceCPU, corev1.ResourceRequestsCPU)
addRequest(corev1.ResourceMemory, corev1.ResourceMemory, corev1.ResourceRequestsMemory)
addRequest(corev1.ResourceEphemeralStorage, corev1.ResourceEphemeralStorage, corev1.ResourceRequestsEphemeralStorage)
if quantity, found := limits[corev1.ResourceCPU]; found {
result[corev1.ResourceLimitsCPU] = quantity
}
if quantity, found := limits[corev1.ResourceMemory]; found {
result[corev1.ResourceLimitsMemory] = quantity
}
if quantity, found := limits[corev1.ResourceEphemeralStorage]; found {
result[corev1.ResourceLimitsEphemeralStorage] = quantity
}
for name, quantity := range requests {
switch {
case strings.HasPrefix(string(name), corev1.ResourceHugePagesPrefix):
result[name] = quantity
result[corev1.ResourceName(corev1.DefaultResourceRequestsPrefix+string(name))] = quantity
case isExtendedResourceName(name):
result[corev1.ResourceName(corev1.DefaultResourceRequestsPrefix+string(name))] = quantity
}
}
return result
}
func serviceUsage(service *corev1.Service) corev1.ResourceList {
result := objectCountUsage("", "services")
result[corev1.ResourceServices] = *resource.NewQuantity(1, resource.DecimalSI)
result[corev1.ResourceServicesLoadBalancers] = *resource.NewQuantity(0, resource.DecimalSI)
result[corev1.ResourceServicesNodePorts] = *resource.NewQuantity(0, resource.DecimalSI)
ports := int64(len(service.Spec.Ports))
switch service.Spec.Type {
case corev1.ServiceTypeClusterIP, corev1.ServiceTypeExternalName:
case corev1.ServiceTypeNodePort:
result[corev1.ResourceServicesNodePorts] = *resource.NewQuantity(ports, resource.DecimalSI)
case corev1.ServiceTypeLoadBalancer:
if ptr.Deref(service.Spec.AllocateLoadBalancerNodePorts, true) {
result[corev1.ResourceServicesNodePorts] = *resource.NewQuantity(ports, resource.DecimalSI)
} else {
var count int64
for _, port := range service.Spec.Ports {
if port.NodePort != 0 {
count++
}
}
result[corev1.ResourceServicesNodePorts] = *resource.NewQuantity(count, resource.DecimalSI)
}
result[corev1.ResourceServicesLoadBalancers] = *resource.NewQuantity(1, resource.DecimalSI)
}
return result
}
func pvcUsage(pvc *corev1.PersistentVolumeClaim) corev1.ResourceList {
result := objectCountUsage("", "persistentvolumeclaims")
one := *resource.NewQuantity(1, resource.DecimalSI)
result[corev1.ResourcePersistentVolumeClaims] = one
storageClass := storagehelpers.GetPersistentVolumeClaimClass(pvc)
if storageClass != "" {
result[resourceByStorageClass(storageClass, corev1.ResourcePersistentVolumeClaims)] = one
}
requested, ok := pvc.Spec.Resources.Requests[corev1.ResourceStorage]
if !ok {
return result
}
rounded := requested.DeepCopy()
_ = rounded.RoundUp(0)
if allocated, ok := pvc.Status.AllocatedResources[corev1.ResourceStorage]; ok && allocated.Cmp(rounded) > 0 {
rounded = allocated.DeepCopy()
_ = rounded.RoundUp(0)
}
result[corev1.ResourceRequestsStorage] = rounded
if storageClass != "" {
result[resourceByStorageClass(storageClass, corev1.ResourceRequestsStorage)] = rounded
}
return result
}
func resourceByStorageClass(storageClass string, name corev1.ResourceName) corev1.ResourceName {
return corev1.ResourceName(storageClass + storageClassSuffix + string(name))
}
// MatchesScopes mirrors generic quota scope matching.
func MatchesScopes(spec corev1.ResourceQuotaSpec, object runtime.Object) (bool, error) {
requirements := make([]corev1.ScopedResourceSelectorRequirement, 0, len(spec.Scopes))
for _, scope := range spec.Scopes {
requirements = append(requirements, corev1.ScopedResourceSelectorRequirement{
ScopeName: scope,
Operator: corev1.ScopeSelectorOpExists,
})
}
if spec.ScopeSelector != nil {
requirements = append(requirements, spec.ScopeSelector.MatchExpressions...)
}
for _, requirement := range requirements {
matches, err := matchesScope(requirement, object)
if err != nil || !matches {
return matches, err
}
}
return true, nil
}
func matchesScope(requirement corev1.ScopedResourceSelectorRequirement, object runtime.Object) (bool, error) {
switch typed := object.(type) {
case *corev1.Pod:
return podMatchesScope(requirement, typed)
case *corev1.PersistentVolumeClaim:
return pvcMatchesScope(requirement, typed)
default:
return false, nil
}
}
func podMatchesScope(requirement corev1.ScopedResourceSelectorRequirement, pod *corev1.Pod) (bool, error) {
switch requirement.ScopeName {
case corev1.ResourceQuotaScopeTerminating:
return isTerminating(pod), nil
case corev1.ResourceQuotaScopeNotTerminating:
return !isTerminating(pod), nil
case corev1.ResourceQuotaScopeBestEffort:
return podQOS(pod) == corev1.PodQOSBestEffort, nil
case corev1.ResourceQuotaScopeNotBestEffort:
return podQOS(pod) != corev1.PodQOSBestEffort, nil
case corev1.ResourceQuotaScopePriorityClass:
if requirement.Operator == corev1.ScopeSelectorOpExists {
return pod.Spec.PriorityClassName != "", nil
}
return requirementMatches(requirement, []string{pod.Spec.PriorityClassName})
case corev1.ResourceQuotaScopeCrossNamespacePodAffinity:
return usesCrossNamespacePodAffinity(pod), nil
case corev1.ResourceQuotaScopeVolumeAttributesClass:
return false, nil
default:
return false, nil
}
}
func pvcMatchesScope(
requirement corev1.ScopedResourceSelectorRequirement,
pvc *corev1.PersistentVolumeClaim,
) (bool, error) {
if requirement.ScopeName != corev1.ResourceQuotaScopeVolumeAttributesClass {
return false, nil
}
values := sets.New[string]()
if value := ptr.Deref(pvc.Spec.VolumeAttributesClassName, ""); value != "" {
values.Insert(value)
}
if value := ptr.Deref(pvc.Status.CurrentVolumeAttributesClassName, ""); value != "" {
values.Insert(value)
}
if pvc.Status.ModifyVolumeStatus != nil && pvc.Status.ModifyVolumeStatus.TargetVolumeAttributesClassName != "" {
values.Insert(pvc.Status.ModifyVolumeStatus.TargetVolumeAttributesClassName)
}
if requirement.Operator == corev1.ScopeSelectorOpExists {
return values.Len() > 0, nil
}
return requirementMatches(requirement, values.UnsortedList())
}
func requirementMatches(
requirement corev1.ScopedResourceSelectorRequirement,
values []string,
) (bool, error) {
operator, err := scopeSelectorOperator(requirement.Operator)
if err != nil {
return false, err
}
labelRequirement, err := labels.NewRequirement(
string(requirement.ScopeName),
operator,
requirement.Values,
)
if err != nil {
return false, err
}
if len(values) == 0 {
return labelRequirement.Matches(labels.Set{}), nil
}
for _, value := range values {
if labelRequirement.Matches(labels.Set{string(requirement.ScopeName): value}) {
return true, nil
}
}
return false, nil
}
func scopeSelectorOperator(operator corev1.ScopeSelectorOperator) (selection.Operator, error) {
switch operator {
case corev1.ScopeSelectorOpIn:
return selection.In, nil
case corev1.ScopeSelectorOpNotIn:
return selection.NotIn, nil
case corev1.ScopeSelectorOpExists:
return selection.Exists, nil
case corev1.ScopeSelectorOpDoesNotExist:
return selection.DoesNotExist, nil
default:
return "", fmt.Errorf("unsupported scope selector operator %q", operator)
}
}
// ValidateConstraints preserves Kubernetes' historical requirement that every
// container explicitly sets CPU/memory resources when those resources are
// quota-controlled.
func ValidateConstraints(hard corev1.ResourceList, object runtime.Object) error {
pod, ok := object.(*corev1.Pod)
if !ok {
return nil
}
// Kubernetes skips the legacy per-container presence check when a supported
// Pod-level request or limit is set. Disabled Pod-level fields are removed by
// the API server before validating admission webhooks receive the Pod.
if resourcehelper.IsPodLevelResourcesSet(pod) {
return nil
}
required := sets.New[corev1.ResourceName]()
for name := range hard {
if validationResources.Has(name) {
required.Insert(name)
}
}
missing := map[corev1.ResourceName][]string{}
containers := append(append([]corev1.Container{}, pod.Spec.Containers...), pod.Spec.InitContainers...)
for _, container := range containers {
usage := podComputeUsage(container.Resources.Requests, container.Resources.Limits)
for name := range required {
if _, ok := usage[name]; !ok {
missing[name] = append(missing[name], container.Name)
}
}
}
if len(missing) == 0 {
return nil
}
parts := make([]string, 0, len(missing))
for _, name := range sets.List(sets.KeySet(missing)) {
parts = append(parts, fmt.Sprintf("%s for: %s", name, strings.Join(missing[name], ",")))
}
return fmt.Errorf("must specify %s", strings.Join(parts, "; "))
}
func quotaPod(pod *corev1.Pod, now time.Time) bool {
if pod.Status.Phase == corev1.PodFailed || pod.Status.Phase == corev1.PodSucceeded {
return false
}
if pod.DeletionTimestamp != nil && pod.DeletionGracePeriodSeconds != nil {
deadline := pod.DeletionTimestamp.Add(time.Duration(*pod.DeletionGracePeriodSeconds) * time.Second)
if now.After(deadline) {
return false
}
}
return true
}
func isTerminating(pod *corev1.Pod) bool {
return pod.Spec.ActiveDeadlineSeconds != nil && *pod.Spec.ActiveDeadlineSeconds >= 0
}
func podQOS(pod *corev1.Pod) corev1.PodQOSClass {
if pod.Status.QOSClass != "" {
return pod.Status.QOSClass
}
requests := corev1.ResourceList{}
limits := corev1.ResourceList{}
guaranteed := true
process := func(target corev1.ResourceList, resources corev1.ResourceList) {
for _, name := range []corev1.ResourceName{corev1.ResourceCPU, corev1.ResourceMemory} {
if quantity, ok := resources[name]; ok && quantity.Sign() > 0 {
addQuantity(target, name, quantity)
}
}
}
hasQoSLimits := func(resources corev1.ResourceList) bool {
for _, name := range []corev1.ResourceName{corev1.ResourceCPU, corev1.ResourceMemory} {
quantity, ok := resources[name]
if !ok || quantity.Sign() <= 0 {
return false
}
}
return true
}
if pod.Spec.Resources != nil {
process(requests, pod.Spec.Resources.Requests)
process(limits, pod.Spec.Resources.Limits)
guaranteed = hasQoSLimits(pod.Spec.Resources.Limits)
} else {
containers := append(append([]corev1.Container{}, pod.Spec.Containers...), pod.Spec.InitContainers...)
for _, container := range containers {
process(requests, container.Resources.Requests)
process(limits, container.Resources.Limits)
if !hasQoSLimits(container.Resources.Limits) {
guaranteed = false
}
}
}
if len(requests) == 0 && len(limits) == 0 {
return corev1.PodQOSBestEffort
}
if guaranteed && len(requests) == len(limits) {
for name, request := range requests {
if limit, ok := limits[name]; !ok || request.Cmp(limit) != 0 {
guaranteed = false
break
}
}
}
if guaranteed && len(requests) == len(limits) {
return corev1.PodQOSGuaranteed
}
return corev1.PodQOSBurstable
}
func addQuantity(target corev1.ResourceList, name corev1.ResourceName, quantity resource.Quantity) {
current := target[name]
current.Add(quantity)
target[name] = current
}
func addResourceList(target, addition corev1.ResourceList) {
for name, quantity := range addition {
addQuantity(target, name, quantity)
}
}
func isExtendedResourceName(name corev1.ResourceName) bool {
value := string(name)
return strings.Contains(value, "/") && !strings.Contains(value, corev1.ResourceDefaultNamespacePrefix)
}
func crossNamespacePodAffinityTerm(term corev1.PodAffinityTerm) bool {
return len(term.Namespaces) != 0 || term.NamespaceSelector != nil
}
func usesCrossNamespacePodAffinity(pod *corev1.Pod) bool {
if pod.Spec.Affinity == nil {
return false
}
check := func(terms []corev1.PodAffinityTerm, weighted []corev1.WeightedPodAffinityTerm) bool {
return slices.ContainsFunc(terms, crossNamespacePodAffinityTerm) ||
slices.ContainsFunc(weighted, func(term corev1.WeightedPodAffinityTerm) bool {
return crossNamespacePodAffinityTerm(term.PodAffinityTerm)
})
}
if affinity := pod.Spec.Affinity.PodAffinity; affinity != nil &&
check(affinity.RequiredDuringSchedulingIgnoredDuringExecution, affinity.PreferredDuringSchedulingIgnoredDuringExecution) {
return true
}
if affinity := pod.Spec.Affinity.PodAntiAffinity; affinity != nil &&
check(affinity.RequiredDuringSchedulingIgnoredDuringExecution, affinity.PreferredDuringSchedulingIgnoredDuringExecution) {
return true
}
return false
}
+471
View File
@@ -0,0 +1,471 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package evaluator
import (
"encoding/json"
"testing"
admissionv1 "k8s.io/api/admission/v1"
autoscalingv2 "k8s.io/api/autoscaling/v2"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/webhook/admission"
)
func TestEvaluatePodUsesUpstreamResourceCalculation(t *testing.T) {
t.Parallel()
pod := &corev1.Pod{
Spec: corev1.PodSpec{
Containers: []corev1.Container{{
Name: "app",
Resources: corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("1"),
corev1.ResourceMemory: resource.MustParse("1Gi"),
},
Limits: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("2"),
corev1.ResourceMemory: resource.MustParse("2Gi"),
},
},
}},
InitContainers: []corev1.Container{{
Name: "init",
Resources: corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("3"),
},
Limits: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("4"),
},
},
}},
Overhead: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("100m"),
},
},
}
result, handled, err := Evaluate(requestFor(t, admissionv1.Create, "pods", "Pod", pod, nil))
if err != nil {
t.Fatalf("Evaluate() error = %v", err)
}
if !handled {
t.Fatal("Evaluate() did not handle Pod")
}
assertQuantity(t, result.NewUsage, corev1.ResourceRequestsCPU, "3100m")
assertQuantity(t, result.NewUsage, corev1.ResourceLimitsCPU, "4100m")
assertQuantity(t, result.NewUsage, corev1.ResourceRequestsMemory, "1Gi")
assertQuantity(t, result.NewUsage, corev1.ResourceLimitsMemory, "2Gi")
assertQuantity(t, result.NewUsage, corev1.ResourcePods, "1")
assertQuantity(t, result.NewUsage, corev1.ResourceName("count/pods"), "1")
}
func TestEvaluatePodUsesEphemeralStorageResources(t *testing.T) {
t.Parallel()
pod := &corev1.Pod{Spec: corev1.PodSpec{
Containers: []corev1.Container{{
Name: "app",
Resources: corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceEphemeralStorage: resource.MustParse("1Gi"),
},
Limits: corev1.ResourceList{
corev1.ResourceEphemeralStorage: resource.MustParse("2Gi"),
},
},
}},
InitContainers: []corev1.Container{{
Name: "init",
Resources: corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceEphemeralStorage: resource.MustParse("3Gi"),
},
Limits: corev1.ResourceList{
corev1.ResourceEphemeralStorage: resource.MustParse("4Gi"),
},
},
}},
Overhead: corev1.ResourceList{
corev1.ResourceEphemeralStorage: resource.MustParse("500Mi"),
},
}}
result, handled, err := Evaluate(requestFor(t, admissionv1.Create, "pods", "Pod", pod, nil))
if err != nil {
t.Fatalf("Evaluate() error = %v", err)
}
if !handled {
t.Fatal("Evaluate() did not handle Pod")
}
assertQuantity(t, result.NewUsage, corev1.ResourceEphemeralStorage, "3572Mi")
assertQuantity(t, result.NewUsage, corev1.ResourceRequestsEphemeralStorage, "3572Mi")
assertQuantity(t, result.NewUsage, corev1.ResourceLimitsEphemeralStorage, "4596Mi")
}
func TestEvaluateTerminalPodOnlyConsumesObjectCount(t *testing.T) {
t.Parallel()
pod := &corev1.Pod{
Spec: corev1.PodSpec{
Containers: []corev1.Container{{
Name: "app",
Resources: corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceEphemeralStorage: resource.MustParse("1Gi"),
},
},
}},
},
Status: corev1.PodStatus{Phase: corev1.PodSucceeded},
}
result, handled, err := Evaluate(requestFor(t, admissionv1.Create, "pods", "Pod", pod, nil))
if err != nil {
t.Fatalf("Evaluate() error = %v", err)
}
if !handled {
t.Fatal("Evaluate() did not handle Pod")
}
assertQuantity(t, result.NewUsage, corev1.ResourceName("count/pods"), "1")
for _, name := range []corev1.ResourceName{
corev1.ResourcePods,
corev1.ResourceEphemeralStorage,
corev1.ResourceRequestsEphemeralStorage,
} {
if _, found := result.NewUsage[name]; found {
t.Fatalf("terminal Pod unexpectedly consumed %q", name)
}
}
}
func TestEvaluatePodUsesPodLevelResourceRequests(t *testing.T) {
t.Parallel()
pod := &corev1.Pod{Spec: corev1.PodSpec{
Containers: []corev1.Container{{Name: "app"}},
Resources: &corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("600m"),
corev1.ResourceMemory: resource.MustParse("512Mi"),
},
Limits: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("800m"),
corev1.ResourceMemory: resource.MustParse("1Gi"),
},
},
}}
result, handled, err := Evaluate(requestFor(t, admissionv1.Create, "pods", "Pod", pod, nil))
if err != nil {
t.Fatalf("Evaluate() error = %v", err)
}
if !handled {
t.Fatal("Evaluate() did not handle Pod")
}
assertQuantity(t, result.NewUsage, corev1.ResourceRequestsCPU, "600m")
assertQuantity(t, result.NewUsage, corev1.ResourceRequestsMemory, "512Mi")
assertQuantity(t, result.NewUsage, corev1.ResourceLimitsCPU, "800m")
assertQuantity(t, result.NewUsage, corev1.ResourceLimitsMemory, "1Gi")
}
func TestValidateConstraintsAllowsPodLevelResources(t *testing.T) {
t.Parallel()
hard := corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("8"),
corev1.ResourceRequestsMemory: resource.MustParse("16Gi"),
corev1.ResourceLimitsCPU: resource.MustParse("8"),
corev1.ResourceLimitsMemory: resource.MustParse("16Gi"),
}
pod := &corev1.Pod{Spec: corev1.PodSpec{
Containers: []corev1.Container{{
Name: "nginx",
Resources: corev1.ResourceRequirements{},
}},
Resources: &corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("100m"),
corev1.ResourceMemory: resource.MustParse("256Mi"),
},
Limits: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("1"),
corev1.ResourceMemory: resource.MustParse("1Gi"),
},
},
}}
if err := ValidateConstraints(hard, pod); err != nil {
t.Fatalf("ValidateConstraints() rejected Pod-level resources: %v", err)
}
}
func TestValidateConstraintsDoesNotTreatUnsupportedPodLevelResourcesAsCompute(t *testing.T) {
t.Parallel()
hard := corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("8"),
}
pod := &corev1.Pod{Spec: corev1.PodSpec{
Containers: []corev1.Container{{
Name: "nginx",
Resources: corev1.ResourceRequirements{},
}},
Resources: &corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceEphemeralStorage: resource.MustParse("1Gi"),
},
},
}}
err := ValidateConstraints(hard, pod)
if err == nil {
t.Fatal("ValidateConstraints() accepted unsupported Pod-level resources")
}
if got, want := err.Error(), "must specify requests.cpu for: nginx"; got != want {
t.Fatalf("ValidateConstraints() error = %q, want %q", got, want)
}
}
func TestEvaluateServiceUpdateReturnsOldAndNewUsage(t *testing.T) {
t.Parallel()
oldService := &corev1.Service{Spec: corev1.ServiceSpec{
Type: corev1.ServiceTypeClusterIP,
Ports: []corev1.ServicePort{{Port: 80}},
}}
newService := oldService.DeepCopy()
newService.Spec.Type = corev1.ServiceTypeLoadBalancer
newService.Spec.Ports = append(newService.Spec.Ports, corev1.ServicePort{Port: 443})
result, handled, err := Evaluate(requestFor(t, admissionv1.Update, "services", "Service", newService, oldService))
if err != nil {
t.Fatalf("Evaluate() error = %v", err)
}
if !handled {
t.Fatal("Evaluate() did not handle Service")
}
assertQuantity(t, result.OldUsage, corev1.ResourceServicesLoadBalancers, "0")
assertQuantity(t, result.OldUsage, corev1.ResourceServices, "1")
assertQuantity(t, result.OldUsage, corev1.ResourceName("count/services"), "1")
assertQuantity(t, result.NewUsage, corev1.ResourceServicesLoadBalancers, "1")
assertQuantity(t, result.NewUsage, corev1.ResourceServicesNodePorts, "2")
}
func TestEvaluateObjectCountNames(t *testing.T) {
t.Parallel()
tests := []struct {
name string
group string
version string
resource string
kind string
expected corev1.ResourceName
legacyName corev1.ResourceName
}{
{
name: "core resource has generic and legacy count",
version: "v1",
resource: "configmaps",
kind: "ConfigMap",
expected: corev1.ResourceName("count/configmaps"),
legacyName: corev1.ResourceConfigMaps,
},
{
name: "grouped resource has qualified generic count",
group: "apps",
version: "v1",
resource: "deployments",
kind: "Deployment",
expected: corev1.ResourceName("count/deployments.apps"),
},
}
for _, test := range tests {
test := test
t.Run(test.name, func(t *testing.T) {
t.Parallel()
object := &metav1.PartialObjectMetadata{
ObjectMeta: metav1.ObjectMeta{Name: "example", Namespace: "tenant-a"},
}
req := requestFor(t, admissionv1.Create, test.resource, test.kind, object, nil)
req.Resource = metav1.GroupVersionResource{
Group: test.group, Version: test.version, Resource: test.resource,
}
req.Kind = metav1.GroupVersionKind{
Group: test.group, Version: test.version, Kind: test.kind,
}
result, handled, err := Evaluate(req)
if err != nil {
t.Fatalf("Evaluate() error = %v", err)
}
if !handled {
t.Fatalf("Evaluate() did not handle %s", test.kind)
}
assertQuantity(t, result.NewUsage, test.expected, "1")
if test.legacyName != "" {
assertQuantity(t, result.NewUsage, test.legacyName, "1")
}
})
}
}
func TestEvaluateCountsHorizontalPodAutoscalers(t *testing.T) {
t.Parallel()
hpa := &autoscalingv2.HorizontalPodAutoscaler{
ObjectMeta: metav1.ObjectMeta{Name: "example", Namespace: "tenant-a"},
Spec: autoscalingv2.HorizontalPodAutoscalerSpec{
ScaleTargetRef: autoscalingv2.CrossVersionObjectReference{
APIVersion: "apps/v1",
Kind: "Deployment",
Name: "example",
},
MaxReplicas: 3,
},
}
req := requestFor(
t,
admissionv1.Create,
"horizontalpodautoscalers",
"HorizontalPodAutoscaler",
hpa,
nil,
)
req.Resource = metav1.GroupVersionResource{
Group: "autoscaling", Version: "v2", Resource: "horizontalpodautoscalers",
}
req.Kind = metav1.GroupVersionKind{
Group: "autoscaling", Version: "v2", Kind: "HorizontalPodAutoscaler",
}
result, handled, err := Evaluate(req)
if err != nil {
t.Fatalf("Evaluate() error = %v", err)
}
if !handled {
t.Fatal("Evaluate() did not handle HorizontalPodAutoscaler")
}
assertQuantity(
t,
result.NewUsage,
corev1.ResourceName("count/horizontalpodautoscalers.autoscaling"),
"1",
)
}
func TestMatchesPodQuotaScopes(t *testing.T) {
t.Parallel()
priority := "high"
pod := &corev1.Pod{Spec: corev1.PodSpec{PriorityClassName: priority}}
spec := corev1.ResourceQuotaSpec{ScopeSelector: &corev1.ScopeSelector{
MatchExpressions: []corev1.ScopedResourceSelectorRequirement{{
ScopeName: corev1.ResourceQuotaScopePriorityClass,
Operator: corev1.ScopeSelectorOpIn,
Values: []string{"high"},
}},
}}
matches, err := MatchesScopes(spec, pod)
if err != nil {
t.Fatalf("MatchesScopes() error = %v", err)
}
if !matches {
t.Fatal("MatchesScopes() = false, want true")
}
}
func TestMatchesBestEffortScopeUsesPodLevelResources(t *testing.T) {
t.Parallel()
pod := &corev1.Pod{Spec: corev1.PodSpec{
Resources: &corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("1"),
corev1.ResourceMemory: resource.MustParse("1Gi"),
},
Limits: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("1"),
corev1.ResourceMemory: resource.MustParse("1Gi"),
},
},
}}
spec := corev1.ResourceQuotaSpec{Scopes: []corev1.ResourceQuotaScope{
corev1.ResourceQuotaScopeBestEffort,
}}
matches, err := MatchesScopes(spec, pod)
if err != nil {
t.Fatalf("MatchesScopes() error = %v", err)
}
if matches {
t.Fatal("pod-level resources were classified as BestEffort")
}
}
func requestFor(
t *testing.T,
operation admissionv1.Operation,
resourceName string,
kind string,
object runtime.Object,
old runtime.Object,
) admission.Request {
t.Helper()
raw, err := json.Marshal(object)
if err != nil {
t.Fatal(err)
}
var oldRaw []byte
if old != nil {
oldRaw, err = json.Marshal(old)
if err != nil {
t.Fatal(err)
}
}
return admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{
Operation: operation,
Resource: metav1.GroupVersionResource{Group: "", Version: "v1", Resource: resourceName},
Kind: metav1.GroupVersionKind{Group: "", Version: "v1", Kind: kind},
Object: runtime.RawExtension{Raw: raw},
OldObject: runtime.RawExtension{Raw: oldRaw},
RequestKind: &metav1.GroupVersionKind{
Group: "", Version: "v1", Kind: kind,
},
RequestResource: &metav1.GroupVersionResource{
Group: "", Version: "v1", Resource: resourceName,
},
}}
}
func assertQuantity(t *testing.T, list corev1.ResourceList, name corev1.ResourceName, want string) {
t.Helper()
got, ok := list[name]
if !ok {
t.Fatalf("resource %q is missing from %#v", name, list)
}
if got.Cmp(resource.MustParse(want)) != 0 {
t.Fatalf("resource %q = %s, want %s", name, got.String(), want)
}
}
@@ -0,0 +1,745 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package globalresourcequota
import (
"context"
"encoding/json"
"fmt"
"slices"
"strings"
"time"
admissionv1 "k8s.io/api/admission/v1"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/webhook/admission"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
quotaevaluator "github.com/projectcapsule/capsule/internal/quota/evaluator"
"github.com/projectcapsule/capsule/pkg/api/meta"
ad "github.com/projectcapsule/capsule/pkg/runtime/admission"
"github.com/projectcapsule/capsule/pkg/runtime/configuration"
"github.com/projectcapsule/capsule/pkg/runtime/events"
"github.com/projectcapsule/capsule/pkg/runtime/handlers"
runtimequota "github.com/projectcapsule/capsule/pkg/runtime/quota"
)
const maxReservations = 1024
var ledgerBackoff = wait.Backoff{
Steps: 8,
Duration: 10 * time.Millisecond,
Factor: 1.6,
Jitter: 0.2,
}
type handler struct{}
type appliedReservation struct {
Key types.NamespacedName
ID string
}
func Handler() handlers.Handler {
return &handler{}
}
func (h *handler) OnCreate(
c client.Client,
reader client.Reader,
_ admission.Decoder,
_ events.EventRecorder,
) handlers.Func {
return h.handle(c, reader)
}
func (h *handler) OnUpdate(
c client.Client,
reader client.Reader,
_ admission.Decoder,
_ events.EventRecorder,
) handlers.Func {
return h.handle(c, reader)
}
func (h *handler) OnDelete(
client.Client,
client.Reader,
admission.Decoder,
events.EventRecorder,
) handlers.Func {
return func(context.Context, admission.Request) *admission.Response {
return nil
}
}
func (h *handler) handle(c client.Client, reader client.Reader) handlers.Func {
return func(ctx context.Context, req admission.Request) *admission.Response {
if isGlobalResourceQuotaRequest(req) {
return validateGlobalResourceQuotaRequest(ctx, reader, req)
}
return enforceGlobalResourceQuotaRequest(ctx, c, reader, req)
}
}
func validateGlobalResourceQuotaRequest(
ctx context.Context,
reader client.Reader,
req admission.Request,
) *admission.Response {
quota := &capsulev1beta2.GlobalResourceQuota{}
if err := json.Unmarshal(req.Object.Raw, quota); err != nil {
return ad.Denyf("GlobalResourceQuota could not be decoded: %v", err)
}
if err := validateGlobalResourceQuota(quota); err != nil {
return ad.Denyf("invalid GlobalResourceQuota: %v", err)
}
if req.Operation != admissionv1.Update {
return nil
}
oldQuota := &capsulev1beta2.GlobalResourceQuota{}
if err := json.Unmarshal(req.OldObject.Raw, oldQuota); err != nil {
return ad.Denyf("previous GlobalResourceQuota could not be decoded: %v", err)
}
if err := validateHardLimit(quota.Spec.Quota.Hard, oldQuota.Status.Total.Used); err != nil {
return ad.Denyf("invalid GlobalResourceQuota: %v", err)
}
ledger := &capsulev1beta2.QuantityLedger{}
err := reader.Get(ctx, types.NamespacedName{
Namespace: configuration.ControllerNamespace(),
Name: oldQuota.GetLedgerName(),
}, ledger)
switch {
case apierrors.IsNotFound(err):
return nil
case err != nil:
return ad.ErroredResponse(err)
case ledger.Spec.TargetRef.UID != oldQuota.UID || ledger.Status.ResourceQuota == nil:
return nil
}
if err := validateHardLimit(
quota.Spec.Quota.Hard,
ledger.Status.ResourceQuota.Allocated,
); err != nil {
return ad.Denyf("invalid GlobalResourceQuota: %v", err)
}
return nil
}
func enforceGlobalResourceQuotaRequest(
ctx context.Context,
c client.Client,
reader client.Reader,
req admission.Request,
) *admission.Response {
if req.Namespace == "" {
return nil
}
namespace := &corev1.Namespace{}
if err := c.Get(ctx, client.ObjectKey{Name: req.Namespace}, namespace); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return ad.ErroredResponse(err)
}
allQuotas := &capsulev1beta2.GlobalResourceQuotaList{}
if err := c.List(ctx, allQuotas); err != nil {
return ad.ErroredResponse(err)
}
quotaList, err := matchingGlobalResourceQuotas(namespace, allQuotas.Items)
if err != nil {
return ad.ErroredResponse(err)
}
if len(quotaList) == 0 {
return nil
}
evaluation, handled, err := quotaevaluator.Evaluate(req)
if err != nil {
return ad.Denyf("GlobalResourceQuota usage could not be calculated: %v", err)
}
if !handled || isManagedResourceQuota(req, evaluation.New) {
// Native ResourceQuotas are implementation details. Counting their
// creation could prevent the quota authorizing them from initializing.
return nil
}
applied := make([]appliedReservation, 0, len(quotaList))
for _, quota := range quotaList {
reservation, response := reserveForGlobalResourceQuota(ctx, c, reader, req, quota, evaluation)
if response != nil {
rollbackReservations(ctx, c, reader, applied)
return response
}
if reservation != nil {
applied = append(applied, *reservation)
}
}
return nil
}
func reserveForGlobalResourceQuota(
ctx context.Context,
c client.Client,
reader client.Reader,
req admission.Request,
quota *capsulev1beta2.GlobalResourceQuota,
evaluation quotaevaluator.Result,
) (*appliedReservation, *admission.Response) {
if quota.DeletionTimestamp != nil {
return nil, nil
}
oldUsage, newUsage, err := usageForQuota(quota.Spec.Quota, evaluation)
if err != nil {
return nil, ad.Denyf(
"resource cannot be evaluated against GlobalResourceQuota %q: %v",
quota.Name,
err,
)
}
if !resourceListPositive(newUsage) && !resourceListPositive(oldUsage) {
return nil, nil
}
delta := positiveDifference(newUsage, oldUsage)
if !resourceListPositive(delta) {
return nil, nil
}
ledgerKey := types.NamespacedName{
Namespace: configuration.ControllerNamespace(),
Name: quota.GetLedgerName(),
}
reservation := newReservation(req, quota.Name, newUsage, delta)
allowed, projected, applied, err := reserve(
ctx,
c,
reader,
ledgerKey,
quota,
reservation,
req.DryRun != nil && *req.DryRun,
)
if err != nil {
if apierrors.IsNotFound(err) {
return nil, ad.Denyf(
"GlobalResourceQuota %q is not ready: QuantityLedger %s does not exist",
quota.Name,
ledgerKey.String(),
)
}
return nil, ad.ErroredResponse(err)
}
if !allowed {
return nil, ad.Denyf(
"resource exceeds GlobalResourceQuota %q: %s",
quota.Name,
formatExceededResources(delta, projected, quota.Spec.Quota.Hard),
)
}
if !applied || (req.DryRun != nil && *req.DryRun) {
return nil, nil
}
return &appliedReservation{Key: ledgerKey, ID: reservation.ID}, nil
}
func isGlobalResourceQuotaRequest(req admission.Request) bool {
return req.Resource.Group == capsulev1beta2.GroupVersion.Group &&
req.Resource.Resource == "globalresourcequotas" &&
req.SubResource == ""
}
func validateGlobalResourceQuota(quota *capsulev1beta2.GlobalResourceQuota) error {
if len(quota.Spec.Quota.Hard) == 0 {
return fmt.Errorf("spec.quota.hard must contain at least one resource")
}
for name, quantity := range quota.Spec.Quota.Hard {
if quantity.Sign() < 0 {
return fmt.Errorf("spec.quota.hard[%q] must not be negative", name)
}
}
for index, namespaceSelector := range quota.Spec.NamespaceSelectors {
if namespaceSelector.LabelSelector == nil {
continue
}
if _, err := metav1.LabelSelectorAsSelector(namespaceSelector.LabelSelector); err != nil {
return fmt.Errorf("spec.namespaceSelectors[%d] is invalid: %w", index, err)
}
}
return nil
}
func validateHardLimit(hard, allocated corev1.ResourceList) error {
for name, usage := range allocated {
if usage.Sign() <= 0 {
continue
}
limit, exists := hard[name]
if !exists {
return fmt.Errorf(
"spec.quota.hard[%q] cannot be removed while %s is allocated",
name,
usage.String(),
)
}
if limit.Cmp(usage) < 0 {
return fmt.Errorf(
"spec.quota.hard[%q] cannot be reduced to %s while %s is allocated",
name,
limit.String(),
usage.String(),
)
}
}
return nil
}
func isManagedResourceQuota(req admission.Request, object any) bool {
if req.Resource.Group != "" || req.Resource.Resource != "resourcequotas" {
return false
}
metadata, ok := object.(metav1.Object)
if !ok {
return false
}
objectLabels := metadata.GetLabels()
return objectLabels[meta.NewManagedByCapsuleLabel] == meta.ValueController &&
objectLabels[meta.GlobalResourceQuotaLabel] != ""
}
func matchingGlobalResourceQuotas(
namespace *corev1.Namespace,
quotas []capsulev1beta2.GlobalResourceQuota,
) ([]*capsulev1beta2.GlobalResourceQuota, error) {
out := make([]*capsulev1beta2.GlobalResourceQuota, 0)
namespaceLabels := labels.Set(namespace.Labels)
for i := range quotas {
quota := &quotas[i]
for _, namespaceSelector := range quota.Spec.NamespaceSelectors {
if namespaceSelector.LabelSelector == nil {
continue
}
selector, err := metav1.LabelSelectorAsSelector(namespaceSelector.LabelSelector)
if err != nil {
return nil, fmt.Errorf("GlobalResourceQuota %q has an invalid namespace selector: %w", quota.Name, err)
}
if selector.Matches(namespaceLabels) {
out = append(out, quota)
break
}
}
}
return out, nil
}
func usageForQuota(
spec corev1.ResourceQuotaSpec,
evaluation quotaevaluator.Result,
) (corev1.ResourceList, corev1.ResourceList, error) {
newUsage := corev1.ResourceList{}
oldUsage := corev1.ResourceList{}
if evaluation.New != nil {
matches, err := quotaevaluator.MatchesScopes(spec, evaluation.New)
if err != nil {
return nil, nil, err
}
if matches {
if err := quotaevaluator.ValidateConstraints(spec.Hard, evaluation.New); err != nil {
return nil, nil, err
}
newUsage = maskResourceList(evaluation.NewUsage, spec.Hard)
}
}
if evaluation.Old != nil {
matches, err := quotaevaluator.MatchesScopes(spec, evaluation.Old)
if err != nil {
return nil, nil, err
}
if matches {
oldUsage = maskResourceList(evaluation.OldUsage, spec.Hard)
}
}
return oldUsage, newUsage, nil
}
func reserve(
ctx context.Context,
c client.Client,
reader client.Reader,
key types.NamespacedName,
quota *capsulev1beta2.GlobalResourceQuota,
reservation capsulev1beta2.QuantityLedgerResourceQuotaReservation,
dryRun bool,
) (allowed bool, allocated corev1.ResourceList, applied bool, err error) {
hard := quota.Spec.Quota.Hard
err = retry.RetryOnConflict(ledgerBackoff, func() error {
applied = false
ledger := &capsulev1beta2.QuantityLedger{}
if getErr := reader.Get(ctx, key, ledger); getErr != nil {
return getErr
}
target := ledger.Spec.TargetRef
if target.Kind != "GlobalResourceQuota" ||
target.Name != quota.Name ||
target.UID != quota.UID {
return fmt.Errorf("QuantityLedger %s has a stale GlobalResourceQuota target", key.String())
}
if ledger.Status.ResourceQuota == nil || !ledger.Status.ResourceQuota.Initialized {
return fmt.Errorf("GlobalResourceQuota QuantityLedger %s is not initialized", key.String())
}
if ledger.Status.ResourceQuota.ObservedGeneration != quota.Generation {
return fmt.Errorf(
"GlobalResourceQuota QuantityLedger %s has not observed generation %d",
key.String(),
quota.Generation,
)
}
if !slices.Contains(ledger.Status.ResourceQuota.Namespaces, reservation.ObjectRef.Namespace) {
return fmt.Errorf(
"GlobalResourceQuota QuantityLedger %s has not observed namespace %s",
key.String(),
reservation.ObjectRef.Namespace,
)
}
now := metav1.Now()
active := make([]capsulev1beta2.QuantityLedgerResourceQuotaReservation, 0, len(ledger.Status.ResourceQuota.Reservations)+1)
found := false
for _, existing := range ledger.Status.ResourceQuota.Reservations {
if existing.ExpiresAt != nil && existing.ExpiresAt.Before(&now) {
continue
}
if existing.ID == reservation.ID {
found = true
applied = true
existing.Usage = reservation.Usage.DeepCopy()
existing.Delta = reservation.Delta.DeepCopy()
existing.ObjectRef = reservation.ObjectRef
existing.UpdatedAt = now
existing.ExpiresAt = reservation.ExpiresAt
}
active = append(active, existing)
}
if !found {
if len(active) >= maxReservations {
return fmt.Errorf("GlobalResourceQuota QuantityLedger %s has too many inflight reservations", key.String())
}
active = append(active, reservation)
applied = true
}
reserved := sumReservations(active, hard)
next := ledger.Status.ResourceQuota.Used.DeepCopy()
addResourceList(next, reserved)
allocated = next.DeepCopy()
if exceeds(next, hard) {
allowed = false
applied = false
return nil
}
if dryRun {
allowed = true
applied = false
return nil
}
ledger.Status.ResourceQuota.Reservations = active
ledger.Status.ResourceQuota.Reserved = reserved
ledger.Status.ResourceQuota.Allocated = next
if updateErr := c.Status().Update(ctx, ledger); updateErr != nil {
applied = false
return updateErr
}
allowed = true
return nil
})
return allowed, allocated, applied, err
}
func rollbackReservations(
ctx context.Context,
c client.Client,
reader client.Reader,
applied []appliedReservation,
) {
for _, item := range slices.Backward(applied) {
_ = rollbackReservation(ctx, c, reader, item.Key, item.ID)
}
}
func rollbackReservation(
ctx context.Context,
c client.Client,
reader client.Reader,
key types.NamespacedName,
id string,
) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
ledger := &capsulev1beta2.QuantityLedger{}
if err := reader.Get(ctx, key, ledger); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return err
}
if ledger.Status.ResourceQuota == nil {
return nil
}
active := make([]capsulev1beta2.QuantityLedgerResourceQuotaReservation, 0, len(ledger.Status.ResourceQuota.Reservations))
removed := false
for _, reservation := range ledger.Status.ResourceQuota.Reservations {
if reservation.ID == id {
removed = true
continue
}
active = append(active, reservation)
}
if !removed {
return nil
}
reserved := sumReservations(active, ledger.Status.ResourceQuota.Used)
allocated := ledger.Status.ResourceQuota.Used.DeepCopy()
addResourceList(allocated, reserved)
ledger.Status.ResourceQuota.Reservations = active
ledger.Status.ResourceQuota.Reserved = reserved
ledger.Status.ResourceQuota.Allocated = allocated
return c.Status().Update(ctx, ledger)
})
}
func newReservation(
req admission.Request,
quotaKey string,
usage corev1.ResourceList,
delta corev1.ResourceList,
) capsulev1beta2.QuantityLedgerResourceQuotaReservation {
now := metav1.Now()
expires := metav1.NewTime(now.Add(2 * time.Minute))
return capsulev1beta2.QuantityLedgerResourceQuotaReservation{
ID: fmt.Sprintf("%s/%s", req.UID, quotaKey),
ObjectRef: capsulev1beta2.QuantityLedgerObjectRef{
APIGroup: req.Kind.Group,
APIVersion: req.Kind.Version,
Kind: req.Kind.Kind,
Namespace: req.Namespace,
Name: req.Name,
},
Usage: usage.DeepCopy(),
Delta: delta.DeepCopy(),
CreatedAt: now,
UpdatedAt: now,
ExpiresAt: &expires,
}
}
func maskResourceList(usage, hard corev1.ResourceList) corev1.ResourceList {
out := make(corev1.ResourceList)
for name := range hard {
if quantity, ok := usage[name]; ok {
out[name] = quantity.DeepCopy()
}
}
return out
}
func positiveDifference(next, previous corev1.ResourceList) corev1.ResourceList {
out := make(corev1.ResourceList)
for name, quantity := range next {
delta := quantity.DeepCopy()
delta.Sub(previous[name])
runtimequota.ClampQuantityToZero(&delta)
out[name] = delta
}
return out
}
func sumReservations(
reservations []capsulev1beta2.QuantityLedgerResourceQuotaReservation,
resources corev1.ResourceList,
) corev1.ResourceList {
out := zeroResourceList(resources)
for _, reservation := range reservations {
addResourceList(out, reservation.Delta)
}
return out
}
func zeroResourceList(resources corev1.ResourceList) corev1.ResourceList {
out := make(corev1.ResourceList, len(resources))
for name := range resources {
out[name] = *resource.NewQuantity(0, resource.DecimalSI)
}
return out
}
func addResourceList(target, addition corev1.ResourceList) {
for name, quantity := range addition {
current := target[name]
current.Add(quantity)
target[name] = current
}
}
func exceeds(usage, hard corev1.ResourceList) bool {
for name, limit := range hard {
quantity := usage[name]
if quantity.Cmp(limit) > 0 {
return true
}
}
return false
}
func resourceListPositive(resources corev1.ResourceList) bool {
for _, quantity := range resources {
if quantity.Sign() > 0 {
return true
}
}
return false
}
func formatExceededResources(requested, projected, hard corev1.ResourceList) string {
names := make([]corev1.ResourceName, 0, len(hard))
for name, limit := range hard {
projectedQuantity := quantityForResource(projected, name)
if projectedQuantity.Cmp(limit) > 0 {
names = append(names, name)
}
}
slices.Sort(names)
details := make([]string, 0, len(names))
for _, name := range names {
requestedQuantity := quantityForResource(requested, name)
projectedQuantity := quantityForResource(projected, name)
hardQuantity := quantityForResource(hard, name)
currentQuantity := projectedQuantity.DeepCopy()
currentQuantity.Sub(requestedQuantity)
runtimequota.ClampQuantityToZero(&currentQuantity)
exceededBy := projectedQuantity.DeepCopy()
exceededBy.Sub(hardQuantity)
details = append(details, fmt.Sprintf(
"%s (requested=%s, current=%s, projected=%s, hard=%s, exceededBy=%s)",
name,
requestedQuantity.String(),
currentQuantity.String(),
projectedQuantity.String(),
hardQuantity.String(),
exceededBy.String(),
))
}
return strings.Join(details, "; ")
}
func quantityForResource(resources corev1.ResourceList, name corev1.ResourceName) resource.Quantity {
if quantity, found := resources[name]; found {
return quantity.DeepCopy()
}
return *resource.NewQuantity(0, resource.DecimalSI)
}
@@ -0,0 +1,417 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package globalresourcequota
import (
"context"
"fmt"
"sync"
"sync/atomic"
"testing"
admissionv1 "k8s.io/api/admission/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
"sigs.k8s.io/controller-runtime/pkg/webhook/admission"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/runtime/selectors"
)
func TestReserveIsAtomicAcrossResources(t *testing.T) {
t.Parallel()
key := types.NamespacedName{Namespace: "capsule-system", Name: "rule-quota"}
hard := corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("8"),
corev1.ResourceRequestsMemory: resource.MustParse("16Gi"),
}
quota := globalQuotaForTest("atomic", hard)
ledger := initializedLedger(key, quota, corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("7"),
corev1.ResourceRequestsMemory: resource.MustParse("10Gi"),
})
cl := ledgerClient(t, ledger)
denied := reservationForTest("denied", corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("1"),
corev1.ResourceRequestsMemory: resource.MustParse("7Gi"),
})
allowed, _, applied, err := reserve(context.Background(), cl, cl, key, quota, denied, false)
if err != nil {
t.Fatalf("reserve(denied) error = %v", err)
}
if allowed || applied {
t.Fatalf("reserve(denied) = allowed %v, applied %v; want false, false", allowed, applied)
}
current := &capsulev1beta2.QuantityLedger{}
if err := cl.Get(context.Background(), key, current); err != nil {
t.Fatal(err)
}
if len(current.Status.ResourceQuota.Reservations) != 0 {
t.Fatalf("denied reservation was persisted: %#v", current.Status.ResourceQuota.Reservations)
}
accepted := reservationForTest("accepted", corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("1"),
corev1.ResourceRequestsMemory: resource.MustParse("6Gi"),
})
allowed, _, applied, err = reserve(context.Background(), cl, cl, key, quota, accepted, false)
if err != nil {
t.Fatalf("reserve(accepted) error = %v", err)
}
if !allowed || !applied {
t.Fatalf("reserve(accepted) = allowed %v, applied %v; want true, true", allowed, applied)
}
}
func TestReserveReportsUpdatedReservationForRollback(t *testing.T) {
t.Parallel()
key := types.NamespacedName{Namespace: "capsule-system", Name: "updated"}
hard := corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("10")}
quota := globalQuotaForTest("updated", hard)
ledger := initializedLedger(key, quota, zeroResourceList(hard))
existing := reservationForTest("same-admission", corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("1"),
})
ledger.Status.ResourceQuota.Reservations = []capsulev1beta2.QuantityLedgerResourceQuotaReservation{existing}
ledger.Status.ResourceQuota.Reserved = existing.Delta.DeepCopy()
ledger.Status.ResourceQuota.Allocated = existing.Delta.DeepCopy()
cl := ledgerClient(t, ledger)
updated := reservationForTest("same-admission", corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("2"),
})
allowed, _, applied, err := reserve(context.Background(), cl, cl, key, quota, updated, false)
if err != nil {
t.Fatalf("reserve(updated) error = %v", err)
}
if !allowed || !applied {
t.Fatalf("reserve(updated) = allowed %v, applied %v; want true, true", allowed, applied)
}
persisted := &capsulev1beta2.QuantityLedger{}
if err := cl.Get(context.Background(), key, persisted); err != nil {
t.Fatal(err)
}
if len(persisted.Status.ResourceQuota.Reservations) != 1 {
t.Fatalf("stored reservations = %d, want 1", len(persisted.Status.ResourceQuota.Reservations))
}
assertLedgerQuantity(
t,
persisted.Status.ResourceQuota.Reservations[0].Delta,
corev1.ResourceRequestsCPU,
"2",
)
if err := rollbackReservation(context.Background(), cl, cl, key, updated.ID); err != nil {
t.Fatalf("rollbackReservation(updated) error = %v", err)
}
current := &capsulev1beta2.QuantityLedger{}
if err := cl.Get(context.Background(), key, current); err != nil {
t.Fatal(err)
}
if len(current.Status.ResourceQuota.Reservations) != 0 {
t.Fatalf("updated reservation was not rolled back: %#v", current.Status.ResourceQuota.Reservations)
}
assertLedgerQuantity(t, current.Status.ResourceQuota.Reserved, corev1.ResourceRequestsCPU, "0")
assertLedgerQuantity(t, current.Status.ResourceQuota.Allocated, corev1.ResourceRequestsCPU, "0")
}
func TestConcurrentReservationsCannotOversubscribe(t *testing.T) {
t.Parallel()
key := types.NamespacedName{Namespace: "capsule-system", Name: "concurrent"}
hard := corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("10")}
quota := globalQuotaForTest("concurrent", hard)
cl := ledgerClient(t, initializedLedger(key, quota, zeroResourceList(hard)))
var allowed atomic.Int32
errs := make(chan error, 20)
var wg sync.WaitGroup
for i := 0; i < 20; i++ {
wg.Add(1)
go func(index int) {
defer wg.Done()
reservation := reservationForTest(
fmt.Sprintf("request-%d", index),
corev1.ResourceList{corev1.ResourceRequestsCPU: resource.MustParse("1")},
)
ok, _, _, err := reserve(context.Background(), cl, cl, key, quota, reservation, false)
if err != nil {
errs <- err
return
}
if ok {
allowed.Add(1)
}
}(i)
}
wg.Wait()
close(errs)
for err := range errs {
t.Errorf("concurrent reserve error = %v", err)
}
if got := allowed.Load(); got != 10 {
t.Fatalf("allowed reservations = %d, want 10", got)
}
current := &capsulev1beta2.QuantityLedger{}
if err := cl.Get(context.Background(), key, current); err != nil {
t.Fatal(err)
}
if len(current.Status.ResourceQuota.Reservations) != 10 {
t.Fatalf("stored reservations = %d, want 10", len(current.Status.ResourceQuota.Reservations))
}
assertLedgerQuantity(t, current.Status.ResourceQuota.Allocated, corev1.ResourceRequestsCPU, "10")
}
func TestManagedResourceQuotaIsExcludedFromAccounting(t *testing.T) {
t.Parallel()
req := admissionRequest("resourcequotas")
object := &metav1.PartialObjectMetadata{ObjectMeta: metav1.ObjectMeta{
Labels: map[string]string{
meta.NewManagedByCapsuleLabel: meta.ValueController,
meta.GlobalResourceQuotaLabel: "shared",
},
}}
if !isManagedResourceQuota(req, object) {
t.Fatal("managed GlobalResourceQuota child was not excluded")
}
object.Labels = nil
if isManagedResourceQuota(req, object) {
t.Fatal("unmanaged ResourceQuota was excluded")
}
}
func TestValidateGlobalResourceQuota(t *testing.T) {
t.Parallel()
tests := []struct {
name string
quota *capsulev1beta2.GlobalResourceQuota
wantErr bool
}{
{
name: "requires hard resources",
quota: globalQuotaForTest("empty", nil),
wantErr: true,
},
{
name: "rejects negative quantities",
quota: globalQuotaForTest("negative", corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("-1"),
}),
wantErr: true,
},
{
name: "accepts an empty selector as all namespaces",
quota: func() *capsulev1beta2.GlobalResourceQuota {
quota := globalQuotaForTest("valid", corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("8"),
})
quota.Spec.NamespaceSelectors = []selectors.NamespaceSelector{{
LabelSelector: &metav1.LabelSelector{},
}}
return quota
}(),
},
}
for _, test := range tests {
test := test
t.Run(test.name, func(t *testing.T) {
t.Parallel()
err := validateGlobalResourceQuota(test.quota)
if (err != nil) != test.wantErr {
t.Fatalf("validateGlobalResourceQuota() error = %v, wantErr %v", err, test.wantErr)
}
})
}
}
func TestValidateHardLimitAgainstAllocatedUsage(t *testing.T) {
t.Parallel()
allocated := corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("3"),
}
if err := validateHardLimit(corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("3"),
}, allocated); err != nil {
t.Fatalf("equal hard limit rejected: %v", err)
}
if err := validateHardLimit(corev1.ResourceList{
corev1.ResourceRequestsCPU: resource.MustParse("2"),
}, allocated); err == nil {
t.Fatal("hard limit below allocated usage was accepted")
}
if err := validateHardLimit(corev1.ResourceList{}, allocated); err == nil {
t.Fatal("allocated resource was removed from hard limit")
}
}
func TestFormatExceededResources(t *testing.T) {
t.Parallel()
requested := corev1.ResourceList{
corev1.ResourceLimitsCPU: resource.MustParse("1"),
corev1.ResourceLimitsMemory: resource.MustParse("1Gi"),
corev1.ResourceRequestsCPU: resource.MustParse("100m"),
corev1.ResourceRequestsMemory: resource.MustParse("256Mi"),
}
projected := corev1.ResourceList{
corev1.ResourceLimitsCPU: resource.MustParse("9"),
corev1.ResourceLimitsMemory: resource.MustParse("9Gi"),
corev1.ResourceRequestsCPU: resource.MustParse("900m"),
corev1.ResourceRequestsMemory: resource.MustParse("2304Mi"),
}
hard := corev1.ResourceList{
corev1.ResourceLimitsCPU: resource.MustParse("8"),
corev1.ResourceLimitsMemory: resource.MustParse("16Gi"),
corev1.ResourceRequestsCPU: resource.MustParse("8"),
corev1.ResourceRequestsMemory: resource.MustParse("16Gi"),
}
want := "limits.cpu (requested=1, current=8, projected=9, hard=8, exceededBy=1)"
if got := formatExceededResources(requested, projected, hard); got != want {
t.Fatalf("formatExceededResources() = %q, want %q", got, want)
}
}
func TestFormatExceededResourcesSortsAndFormatsEveryExceededLimit(t *testing.T) {
t.Parallel()
requested := corev1.ResourceList{
corev1.ResourceLimitsMemory: resource.MustParse("512Mi"),
corev1.ResourceRequestsCPU: resource.MustParse("750m"),
}
projected := corev1.ResourceList{
corev1.ResourceLimitsMemory: resource.MustParse("1536Mi"),
corev1.ResourceRequestsCPU: resource.MustParse("1500m"),
}
hard := corev1.ResourceList{
corev1.ResourceLimitsMemory: resource.MustParse("1Gi"),
corev1.ResourceRequestsCPU: resource.MustParse("1"),
}
want := "limits.memory (requested=512Mi, current=1Gi, projected=1536Mi, hard=1Gi, exceededBy=512Mi); " +
"requests.cpu (requested=750m, current=750m, projected=1500m, hard=1, exceededBy=500m)"
if got := formatExceededResources(requested, projected, hard); got != want {
t.Fatalf("formatExceededResources() = %q, want %q", got, want)
}
}
func admissionRequest(resourceName string) admission.Request {
return admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{
Resource: metav1.GroupVersionResource{Group: "", Version: "v1", Resource: resourceName},
}}
}
func initializedLedger(
key types.NamespacedName,
quota *capsulev1beta2.GlobalResourceQuota,
used corev1.ResourceList,
) *capsulev1beta2.QuantityLedger {
hard := quota.Spec.Quota.Hard
return &capsulev1beta2.QuantityLedger{
ObjectMeta: metav1.ObjectMeta{Name: key.Name, Namespace: key.Namespace},
Spec: capsulev1beta2.QuantityLedgerSpec{
TargetRef: capsulev1beta2.QuantityLedgerTargetRef{
Kind: "GlobalResourceQuota",
Name: quota.Name,
UID: quota.UID,
},
},
Status: capsulev1beta2.QuantityLedgerStatus{
ResourceQuota: &capsulev1beta2.QuantityLedgerResourceQuotaStatus{
ObservedGeneration: quota.Generation,
Initialized: true,
Namespaces: []string{"tenant-a"},
Used: used.DeepCopy(),
Reserved: zeroResourceList(hard),
Allocated: used.DeepCopy(),
},
},
}
}
func globalQuotaForTest(name string, hard corev1.ResourceList) *capsulev1beta2.GlobalResourceQuota {
return &capsulev1beta2.GlobalResourceQuota{
ObjectMeta: metav1.ObjectMeta{Name: name, UID: types.UID(name + "-uid"), Generation: 1},
Spec: capsulev1beta2.GlobalResourceQuotaSpec{
Quota: corev1.ResourceQuotaSpec{Hard: hard.DeepCopy()},
},
}
}
func reservationForTest(
id string,
delta corev1.ResourceList,
) capsulev1beta2.QuantityLedgerResourceQuotaReservation {
now := metav1.Now()
return capsulev1beta2.QuantityLedgerResourceQuotaReservation{
ID: id,
Usage: delta.DeepCopy(),
Delta: delta.DeepCopy(),
ObjectRef: capsulev1beta2.QuantityLedgerObjectRef{
APIVersion: "v1",
Kind: "Pod",
Namespace: "tenant-a",
},
CreatedAt: now,
UpdatedAt: now,
}
}
func ledgerClient(t *testing.T, objects ...client.Object) client.Client {
t.Helper()
scheme := runtime.NewScheme()
if err := capsulev1beta2.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
return fake.NewClientBuilder().
WithScheme(scheme).
WithStatusSubresource(&capsulev1beta2.QuantityLedger{}).
WithObjects(objects...).
Build()
}
func assertLedgerQuantity(
t *testing.T,
list corev1.ResourceList,
name corev1.ResourceName,
want string,
) {
t.Helper()
got := list[name]
if got.Cmp(resource.MustParse(want)) != 0 {
t.Fatalf("%s = %s, want %s", name, got.String(), want)
}
}
@@ -0,0 +1,22 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package route
import "github.com/projectcapsule/capsule/pkg/runtime/handlers"
type globalResourceQuotaCalculation struct {
handlers []handlers.Handler
}
func GlobalResourceQuotaCalculation(handler ...handlers.Handler) handlers.Webhook {
return &globalResourceQuotaCalculation{handlers: handler}
}
func (w *globalResourceQuotaCalculation) GetHandlers() []handlers.Handler {
return w.handlers
}
func (w *globalResourceQuotaCalculation) GetPath() string {
return "/global-resource-quotas/calculations"
}
@@ -5,6 +5,7 @@ package validation
import (
"context"
"fmt"
k8smeta "k8s.io/apimachinery/pkg/api/meta"
"sigs.k8s.io/controller-runtime/pkg/client"
@@ -16,6 +17,7 @@ import (
ad "github.com/projectcapsule/capsule/pkg/runtime/admission"
"github.com/projectcapsule/capsule/pkg/runtime/events"
"github.com/projectcapsule/capsule/pkg/runtime/handlers"
tenantutils "github.com/projectcapsule/capsule/pkg/tenant"
)
type RuleValidationHandler struct {
@@ -89,7 +91,7 @@ func (h *RuleValidationHandler) handle(
}
body := rule.NamespaceRuleBodyNamespace
if body.Enforce == nil {
if body.Enforce == nil && len(body.Quota) == 0 {
continue
}
@@ -104,5 +106,17 @@ func (h *RuleValidationHandler) handle(
return ad.Deny(err.Error())
}
for ruleIndex, rule := range tnt.Spec.Rules {
if rule == nil || rule.NamespaceRuleBodyNamespace == nil {
continue
}
for quotaIndex, quota := range rule.Quota {
if err := tenantutils.ValidateRuleGlobalResourceQuotaName(tnt, quota.Name); err != nil {
return ad.Deny(fmt.Sprintf("rules[%d].quota[%d]: %v", ruleIndex, quotaIndex, err))
}
}
}
return nil
}