fix: correct tls reconciler and add tenantowners (#1946)

* fix(controller): decode old object for delete requests

Signed-off-by: Oliver Bähler <oliverbaehler@hotmail.com>

* chore: modernize golang

Signed-off-by: Oliver Bähler <oliverbaehler@hotmail.com>

* chore: modernize golang

Signed-off-by: Oliver Bähler <oliverbaehler@hotmail.com>

* chore: modernize golang

Signed-off-by: Oliver Bähler <oliverbaehler@hotmail.com>

* fix: tls controller

Signed-off-by: Oliver Baehler <oliver@sudo-i.net>

* feat: add tenantowner tenant status reference

Signed-off-by: Oliver Baehler <oliver@sudo-i.net>

* fix: tlsreconciler only patches cabundles

Signed-off-by: Oliver Baehler <oliver@sudo-i.net>

* chore: refactor logger usage

Signed-off-by: Oliver Baehler <oliver@sudo-i.net>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* fix: tlsreconciler only patches cabundles

Signed-off-by: Oliver Baehler <oliver@sudo-i.net>

* fix: tlsreconciler only patches cabundles

Signed-off-by: Oliver Baehler <oliver@sudo-i.net>

---------

Signed-off-by: Oliver Bähler <oliverbaehler@hotmail.com>
Signed-off-by: Oliver Baehler <oliver@sudo-i.net>
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
Oliver Bähler
2026-06-03 16:02:17 +02:00
committed by GitHub
co-authored by Copilot Autofix powered by AI
parent f62b7748ae
commit 3d1158ed23
63 changed files with 2377 additions and 601 deletions
-1
View File
@@ -64,7 +64,6 @@ jobs:
fail-fast: false
matrix:
k8s-version:
- 'v1.31.0'
- 'v1.32.0'
- 'v1.33.0'
- 'v1.34.0'
+7 -6
View File
@@ -186,6 +186,7 @@ dev-setup:
--create-namespace \
--version=$(CHART_VERSION) \
--set 'proxy.enabled=true' \
--set 'proxy.certManager.generateCertificates=false' \
--set 'crds.install=true' \
--set 'crds.exclusive=true'\
--set 'crds.createConfig=true'\
@@ -194,8 +195,8 @@ dev-setup:
--set "monitoring.diagnostics.enabled=true"\
--set "manager.rbac.minimal=true"\
--set 'certManager.generateCertificates=false' \
--set 'tls.enableController=true' \
--set 'tls.create=true' \
--set 'tls.enableController=false' \
--set 'tls.create=false' \
--set "webhooks.exclusive=true"\
--set "webhooks.service.url=$${WEBHOOK_URL}" \
--set "webhooks.service.caBundle=$${CA_BUNDLE}" \
@@ -453,7 +454,7 @@ e2e-build: kind
$(MAKE) e2e-install
.PHONY: e2e-install
e2e-install: helm-controller-version ko-build-all
e2e-install: helm-controller-version ko-build-all dev-install-gw-api-crds
$(MAKE) e2e-load-image CLUSTER_NAME=$(CLUSTER_NAME) IMAGE=$(CAPSULE_IMG) VERSION=$(VERSION)
$(KUBECTL) label clusterrole admin projectcapsule.dev/aggregate-to-controller=true
$(HELM) upgrade \
@@ -470,10 +471,10 @@ e2e-install: helm-controller-version ko-build-all
--set 'manager.resources=null'\
--set "manager.image.tag=$(VERSION)" \
--set 'manager.livenessProbe.failureThreshold=10' \
--set 'manager.options.logLevel=debug' \
--set-string 'manager.options.logLevel=5' \
--set 'manager.options.workers=4' \
--set 'manager.options.clientConnectionQPS=2000' \
--set 'manager.options.clientConnectionQPS=1000' \
--set 'manager.options.clientConnectionBurst=1000' \
--set 'manager.rbac.minimal=true' \
--set 'webhooks.hooks.nodes.enabled=true' \
--set "webhooks.exclusive=true"\
@@ -604,7 +605,7 @@ helm-doc:
# -- Tools
####################
CONTROLLER_GEN := $(LOCALBIN)/controller-gen
CONTROLLER_GEN_VERSION ?= v0.20.0
CONTROLLER_GEN_VERSION ?= v0.21.0
CONTROLLER_GEN_LOOKUP := kubernetes-sigs/controller-tools
controller-gen:
@test -s $(CONTROLLER_GEN) && $(CONTROLLER_GEN) --version | grep -q $(CONTROLLER_GEN_VERSION) || \
+13
View File
@@ -6,6 +6,7 @@ package v1beta2
import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/api/rbac"
)
@@ -26,11 +27,23 @@ type TenantOwnerStatus struct {
// ObservedGeneration is the most recent generation the controller has observed.
// +optional
ObservedGeneration int64 `json:"observedGeneration,omitempty"`
// Tenants lists the names of all Tenants that this TenantOwner is currently matched to
// via the Tenant's spec.permissions.matchOwners selectors.
// +optional
// +listType=atomic
Tenants []string `json:"tenants,omitempty"`
// Conditions contains the reconciliation conditions for this TenantOwner.
// +optional
Conditions meta.ConditionList `json:"conditions,omitempty"`
}
// +kubebuilder:object:root=true
// +kubebuilder:subresource:status
// +kubebuilder:resource:scope=Cluster,shortName=to
// +kubebuilder:printcolumn:name="Ready",type="string",JSONPath=".status.conditions[?(@.type==\"Ready\")].status",description="Reconcile status of this TenantOwner"
// +kubebuilder:printcolumn:name="Age",type="date",JSONPath=".metadata.creationTimestamp",description="Age"
// TenantOwner is the Schema for the tenantowners API.
type TenantOwner struct {
+13 -1
View File
@@ -1811,7 +1811,7 @@ func (in *TenantOwner) DeepCopyInto(out *TenantOwner) {
out.TypeMeta = in.TypeMeta
in.ObjectMeta.DeepCopyInto(&out.ObjectMeta)
in.Spec.DeepCopyInto(&out.Spec)
out.Status = in.Status
in.Status.DeepCopyInto(&out.Status)
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new TenantOwner.
@@ -1883,6 +1883,18 @@ func (in *TenantOwnerSpec) DeepCopy() *TenantOwnerSpec {
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in *TenantOwnerStatus) DeepCopyInto(out *TenantOwnerStatus) {
*out = *in
if in.Tenants != nil {
in, out := &in.Tenants, &out.Tenants
*out = make([]string, len(*in))
copy(*out, *in)
}
if in.Conditions != nil {
in, out := &in.Conditions, &out.Conditions
*out = make(meta.ConditionList, len(*in))
for i := range *in {
(*in)[i].DeepCopyInto(&(*out)[i])
}
}
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new TenantOwnerStatus.
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: capsuleconfigurations.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: customquotas.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: globalcustomquotas.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: globaltenantresources.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: quantityledgers.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: resourcepoolclaims.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: resourcepools.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: rulestatuses.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: tenantowners.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -16,7 +16,16 @@ spec:
singular: tenantowner
scope: Cluster
versions:
- name: v1beta2
- additionalPrinterColumns:
- description: Reconcile status of this TenantOwner
jsonPath: .status.conditions[?(@.type=="Ready")].status
name: Ready
type: string
- description: Age
jsonPath: .metadata.creationTimestamp
name: Age
type: date
name: v1beta2
schema:
openAPIV3Schema:
description: TenantOwner is the Schema for the tenantowners API.
@@ -75,11 +84,77 @@ spec:
status:
description: status defines the observed state of TenantOwner.
properties:
conditions:
description: Conditions contains the reconciliation conditions for
this TenantOwner.
items:
description: Condition contains details for one aspect of the current
state of this API Resource.
properties:
lastTransitionTime:
description: |-
lastTransitionTime is the last time the condition transitioned from one status to another.
This should be when the underlying condition changed. If that is not known, then using the time when the API field changed is acceptable.
format: date-time
type: string
message:
description: |-
message is a human readable message indicating details about the transition.
This may be an empty string.
maxLength: 32768
type: string
observedGeneration:
description: |-
observedGeneration represents the .metadata.generation that the condition was set based upon.
For instance, if .metadata.generation is currently 12, but the .status.conditions[x].observedGeneration is 9, the condition is out of date
with respect to the current state of the instance.
format: int64
minimum: 0
type: integer
reason:
description: |-
reason contains a programmatic identifier indicating the reason for the condition's last transition.
Producers of specific condition types may define expected values and meanings for this field,
and whether the values are considered a guaranteed API.
The value should be a CamelCase string.
This field may not be empty.
maxLength: 1024
minLength: 1
pattern: ^[A-Za-z]([A-Za-z0-9_,:]*[A-Za-z0-9_])?$
type: string
status:
description: status of the condition, one of True, False, Unknown.
enum:
- "True"
- "False"
- Unknown
type: string
type:
description: type of condition in CamelCase or in foo.example.com/CamelCase.
maxLength: 316
pattern: ^([a-z0-9]([-a-z0-9]*[a-z0-9])?(\.[a-z0-9]([-a-z0-9]*[a-z0-9])?)*/)?(([A-Za-z0-9][-A-Za-z0-9_.]*)?[A-Za-z0-9])$
type: string
required:
- lastTransitionTime
- message
- reason
- status
- type
type: object
type: array
observedGeneration:
description: ObservedGeneration is the most recent generation the
controller has observed.
format: int64
type: integer
tenants:
description: |-
Tenants lists the names of all Tenants that this TenantOwner is currently matched to
via the Tenant's spec.permissions.matchOwners selectors.
items:
type: string
type: array
x-kubernetes-list-type: atomic
type: object
required:
- spec
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: tenantresources.capsule.clastix.io
spec:
group: capsule.clastix.io
@@ -3,7 +3,7 @@ apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.20.0
controller-gen.kubebuilder.io/version: v0.21.0
name: tenants.capsule.clastix.io
spec:
group: capsule.clastix.io
+10 -1
View File
@@ -53,6 +53,7 @@ import (
rulestatuscontroller "github.com/projectcapsule/capsule/internal/controllers/rulestatus"
servicelabelscontroller "github.com/projectcapsule/capsule/internal/controllers/servicelabels"
tenantcontroller "github.com/projectcapsule/capsule/internal/controllers/tenant"
tenantownercontroller "github.com/projectcapsule/capsule/internal/controllers/tenantowner"
tlscontroller "github.com/projectcapsule/capsule/internal/controllers/tls"
utilscontroller "github.com/projectcapsule/capsule/internal/controllers/utils"
"github.com/projectcapsule/capsule/internal/metrics"
@@ -321,7 +322,7 @@ func main() {
}
// Reconcile TLS certificates before starting controllers and webhooks
if err = tlsReconciler.ReconcileCertificates(ctx, tlsCert); err != nil {
if err = tlsReconciler.ReconcileCertificates(ctx, ctrl.Log.WithName("capsule.setup").WithName("tls"), tlsCert); err != nil {
setupLog.Error(err, "unable to reconcile Capsule TLS secret")
os.Exit(1)
}
@@ -687,6 +688,14 @@ func main() {
os.Exit(1)
}
if err = (&tenantownercontroller.TenantOwnerManager{
Log: ctrl.Log.WithName("capsule.ctrl").WithName("tenantowners"),
Client: manager.GetClient(),
}).SetupWithManager(manager, controllerConfig); err != nil {
setupLog.Error(err, "unable to create controller", "controller", "TenantOwners")
os.Exit(1)
}
if err = (&servicelabelscontroller.ServicesLabelsReconciler{
Log: ctrl.Log.WithName("capsule.ctrl").WithName("services"),
}).SetupWithManager(ctx, manager); err != nil {
@@ -21,7 +21,7 @@ import (
"github.com/projectcapsule/capsule/pkg/api/rbac"
)
var _ = Describe("when Tenant owner interacts with the webhooks", Ordered, Label("tenant"), func() {
var _ = Describe("when Tenant owner interacts with the webhooks", Ordered, Label("tenant", "permissions", "owners"), func() {
tnt := &capsulev1beta2.Tenant{
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-owner-admission",
+385
View File
@@ -0,0 +1,385 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package e2e
import (
"context"
"fmt"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
capmeta "github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/api/rbac"
)
// tenantOwnerReady waits until the TenantOwner's Ready condition has the expected status.
func tenantOwnerReady(to *capsulev1beta2.TenantOwner, expected metav1.ConditionStatus) {
Eventually(func(g Gomega) {
current := &capsulev1beta2.TenantOwner{}
g.Expect(k8sClient.Get(context.TODO(), types.NamespacedName{Name: to.GetName()}, current)).To(Succeed())
cond := current.Status.Conditions.GetConditionByType(capmeta.ReadyCondition)
g.Expect(cond).NotTo(BeNil(), "TenantOwner %q should have a Ready condition", to.GetName())
g.Expect(cond.Status).To(
Equal(expected),
"TenantOwner %q Ready condition should be %s, got %s: %s",
to.GetName(), expected, cond.Status, cond.Message,
)
}, defaultTimeoutInterval, defaultPollInterval).Should(Succeed())
}
// tenantOwnerMatchedTenants waits until the TenantOwner's status.matchedTenantNames
// equals the expected sorted list. The controller calls sort.Strings before persisting,
// so the assertion is order-sensitive to validate stable ordering for API consumers.
func tenantOwnerMatchedTenants(to *capsulev1beta2.TenantOwner, expected []string) {
Eventually(func(g Gomega) {
current := &capsulev1beta2.TenantOwner{}
g.Expect(k8sClient.Get(context.TODO(), types.NamespacedName{Name: to.GetName()}, current)).To(Succeed())
actual := current.Status.Tenants
if actual == nil {
actual = []string{}
}
g.Expect(actual).To(
Equal(expected),
"TenantOwner %q status.matchedTenantNames mismatch", to.GetName(),
)
}, defaultTimeoutInterval, defaultPollInterval).Should(Succeed())
}
var _ = Describe("TenantOwner status tracks matched Tenants", Ordered, Label("tenantowner", "tenant", "permissions", "owners"), func() {
// to is a shared TenantOwner used across the sub-tests in this Ordered block.
// Two Tenants are created/deleted to verify count changes.
to := &capsulev1beta2.TenantOwner{
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-to-status",
Labels: map[string]string{
"e2e-to-status": "true",
"e2e.suite.capsule/name": "owner_status_test",
},
},
Spec: capsulev1beta2.TenantOwnerSpec{
Aggregate: true,
CoreOwnerSpec: rbac.CoreOwnerSpec{
UserSpec: rbac.UserSpec{
Kind: rbac.UserOwner,
Name: "e2e-to-status-user",
},
ClusterRoles: []string{"view"},
},
},
}
tnt1 := &capsulev1beta2.Tenant{
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-to-status-tnt1",
},
Spec: capsulev1beta2.TenantSpec{
Permissions: capsulev1beta2.Permissions{
MatchOwners: []*metav1.LabelSelector{
{MatchLabels: map[string]string{"e2e.suite.capsule/name": "owner_status_test"}},
},
},
},
}
tnt2 := &capsulev1beta2.Tenant{
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-to-status-tnt2",
},
Spec: capsulev1beta2.TenantSpec{
Permissions: capsulev1beta2.Permissions{
MatchOwners: []*metav1.LabelSelector{
{MatchLabels: map[string]string{"e2e.suite.capsule/name": "owner_status_test"}},
},
},
},
}
JustBeforeEach(func() {
EventuallyCreation(func() error {
to.ResourceVersion = ""
return k8sClient.Create(context.TODO(), to)
}).Should(Succeed())
})
JustAfterEach(func() {
EventuallyDeletion(to)
EventuallyDeletion(tnt1)
EventuallyDeletion(tnt2)
})
It("has Ready=True and empty tenant list when no Tenant matches", func() {
tenantOwnerReady(to, metav1.ConditionTrue)
tenantOwnerMatchedTenants(to, []string{})
})
It("reflects one matched Tenant in status", func() {
EventuallyCreation(func() error {
tnt1.ResourceVersion = ""
return k8sClient.Create(context.TODO(), tnt1)
}).Should(Succeed())
TenantReadyTrue(tnt1)
tenantOwnerReady(to, metav1.ConditionTrue)
tenantOwnerMatchedTenants(to, []string{tnt1.Name})
})
It("reflects two matched Tenants in status", func() {
EventuallyCreation(func() error {
tnt1.ResourceVersion = ""
return k8sClient.Create(context.TODO(), tnt1)
}).Should(Succeed())
EventuallyCreation(func() error {
tnt2.ResourceVersion = ""
return k8sClient.Create(context.TODO(), tnt2)
}).Should(Succeed())
TenantReadyTrue(tnt1)
TenantReadyTrue(tnt2)
tenantOwnerReady(to, metav1.ConditionTrue)
tenantOwnerMatchedTenants(to, []string{tnt1.Name, tnt2.Name})
})
It("clears the matched Tenant from status when the Tenant is deleted", func() {
EventuallyCreation(func() error {
tnt1.ResourceVersion = ""
return k8sClient.Create(context.TODO(), tnt1)
}).Should(Succeed())
TenantReadyTrue(tnt1)
tenantOwnerMatchedTenants(to, []string{tnt1.Name})
EventuallyDeletion(tnt1)
tenantOwnerReady(to, metav1.ConditionTrue)
tenantOwnerMatchedTenants(to, []string{})
})
It("updates status when a Tenant's matchOwners selector is removed", func() {
EventuallyCreation(func() error {
tnt1.ResourceVersion = ""
return k8sClient.Create(context.TODO(), tnt1)
}).Should(Succeed())
TenantReadyTrue(tnt1)
tenantOwnerMatchedTenants(to, []string{tnt1.Name})
// Remove the matchOwners selector so the TenantOwner is no longer matched.
Eventually(func(g Gomega) {
current := &capsulev1beta2.Tenant{}
g.Expect(k8sClient.Get(context.TODO(), types.NamespacedName{Name: tnt1.Name}, current)).To(Succeed())
current.Spec.Permissions.MatchOwners = nil
g.Expect(k8sClient.Update(context.TODO(), current)).To(Succeed())
}, defaultTimeoutInterval, defaultPollInterval).Should(Succeed())
tenantOwnerReady(to, metav1.ConditionTrue)
tenantOwnerMatchedTenants(to, []string{})
})
It("tracks observedGeneration in TenantOwner status", func() {
Eventually(func(g Gomega) {
current := &capsulev1beta2.TenantOwner{}
g.Expect(k8sClient.Get(context.TODO(), types.NamespacedName{Name: to.Name}, current)).To(Succeed())
g.Expect(current.Status.ObservedGeneration).To(
Equal(current.GetGeneration()),
"status.observedGeneration should equal metadata.generation",
)
}, defaultTimeoutInterval, defaultPollInterval).Should(Succeed())
})
})
var _ = Describe("TenantOwner status with 10 matched Tenants", Ordered, Label("tenantowner", "status", "scale"), func() {
const tenantCount = 10
to := &capsulev1beta2.TenantOwner{
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-to-status-10match",
Labels: map[string]string{
"e2e-to-status-10match": "true",
"e2e.suite.capsule/name": "owner_status_test",
},
},
Spec: capsulev1beta2.TenantOwnerSpec{
Aggregate: true,
CoreOwnerSpec: rbac.CoreOwnerSpec{
UserSpec: rbac.UserSpec{
Kind: rbac.UserOwner,
Name: "e2e-to-status-10match-user",
},
ClusterRoles: []string{"view"},
},
},
}
tenants := func() []*capsulev1beta2.Tenant {
tnts := make([]*capsulev1beta2.Tenant, tenantCount)
for i := range tnts {
tnts[i] = &capsulev1beta2.Tenant{
ObjectMeta: metav1.ObjectMeta{
Name: fmt.Sprintf("e2e-to-status-10match-tnt-%d", i),
Labels: map[string]string{
"e2e.suite.capsule/name": "owner_status_test",
},
},
Spec: capsulev1beta2.TenantSpec{
Permissions: capsulev1beta2.Permissions{
MatchOwners: []*metav1.LabelSelector{
{MatchLabels: map[string]string{"e2e-to-status-10match": "true"}},
},
},
},
}
}
return tnts
}
JustBeforeEach(func() {
EventuallyCreation(func() error {
to.ResourceVersion = ""
return k8sClient.Create(context.TODO(), to)
}).Should(Succeed())
for _, tnt := range tenants() {
tnt := tnt
EventuallyCreation(func() error {
tnt.ResourceVersion = ""
return k8sClient.Create(context.TODO(), tnt)
}).Should(Succeed())
}
})
JustAfterEach(func() {
EventuallyDeletion(to)
for _, tnt := range tenants() {
EventuallyDeletion(tnt)
}
})
It("reports all 10 matched Tenants and Ready=True", func() {
expectedNames := make([]string, tenantCount)
for i := range expectedNames {
expectedNames[i] = fmt.Sprintf("e2e-to-status-10match-tnt-%d", i)
}
tenantOwnerReady(to, metav1.ConditionTrue)
tenantOwnerMatchedTenants(to, expectedNames)
})
})
var _ = Describe("TenantOwner status fan-out: 10 TenantOwners matched by 1 Tenant", Ordered, Label("tenantowner", "status", "fanout"), func() {
const ownerCount = 10
owners := func() []*capsulev1beta2.TenantOwner {
tos := make([]*capsulev1beta2.TenantOwner, ownerCount)
for i := range tos {
tos[i] = &capsulev1beta2.TenantOwner{
ObjectMeta: metav1.ObjectMeta{
Name: fmt.Sprintf("e2e-to-fanout-%d", i),
Labels: map[string]string{
"e2e-to-fanout": "true",
"e2e.suite.capsule/name": "owner_status_test",
},
},
Spec: capsulev1beta2.TenantOwnerSpec{
Aggregate: true,
CoreOwnerSpec: rbac.CoreOwnerSpec{
UserSpec: rbac.UserSpec{
Kind: rbac.UserOwner,
Name: fmt.Sprintf("e2e-to-fanout-user-%d", i),
},
ClusterRoles: []string{"view"},
},
},
}
}
return tos
}
tnt := &capsulev1beta2.Tenant{
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-tnt-fanout-single",
Labels: map[string]string{
"e2e.suite.capsule/name": "owner_status_test",
},
},
Spec: capsulev1beta2.TenantSpec{
Permissions: capsulev1beta2.Permissions{
MatchOwners: []*metav1.LabelSelector{
{MatchLabels: map[string]string{
"e2e-to-fanout": "true",
"e2e.suite.capsule/name": "owner_status_test",
}},
},
},
},
}
JustBeforeEach(func() {
for _, to := range owners() {
to := to
EventuallyCreation(func() error {
to.ResourceVersion = ""
return k8sClient.Create(context.TODO(), to)
}).Should(Succeed())
}
EventuallyCreation(func() error {
tnt.ResourceVersion = ""
return k8sClient.Create(context.TODO(), tnt)
}).Should(Succeed())
})
JustAfterEach(func() {
for _, to := range owners() {
EventuallyDeletion(to)
}
EventuallyDeletion(tnt)
})
It("updates all 10 TenantOwners when 1 Tenant is created", func() {
// All 10 TenantOwners should each report matchedTenants=1 and tenants=[tnt.Name].
for _, to := range owners() {
to := to
tenantOwnerReady(to, metav1.ConditionTrue)
tenantOwnerMatchedTenants(to, []string{tnt.Name})
}
})
It("clears all 10 TenantOwners when the Tenant is deleted", func() {
// Establish baseline: all matched.
for _, to := range owners() {
tenantOwnerMatchedTenants(to, []string{tnt.Name})
}
EventuallyDeletion(tnt)
// After deletion every TenantOwner should report zero matches.
for _, to := range owners() {
tenantOwnerMatchedTenants(to, []string{})
}
})
})
+19 -11
View File
@@ -17,7 +17,7 @@ import (
"github.com/projectcapsule/capsule/pkg/api/rbac"
)
var _ = Describe("Owners", Ordered, Label("config", "tenant", "permissions", "owners"), func() {
var _ = Describe("Owners", Ordered, Label("config", "tenantowner", "tenant", "permissions", "owners"), func() {
originConfig := &capsulev1beta2.CapsuleConfiguration{}
tnt1 := &capsulev1beta2.Tenant{
@@ -32,12 +32,14 @@ var _ = Describe("Owners", Ordered, Label("config", "tenant", "permissions", "ow
MatchOwners: []*metav1.LabelSelector{
{
MatchLabels: map[string]string{
"customer": "x",
"customer": "x",
"e2e.suite.capsule/name": "owner_test",
},
},
{
MatchLabels: map[string]string{
"team": "devops",
"team": "devops",
"e2e.suite.capsule/name": "owner_test",
},
},
},
@@ -83,12 +85,14 @@ var _ = Describe("Owners", Ordered, Label("config", "tenant", "permissions", "ow
MatchOwners: []*metav1.LabelSelector{
{
MatchLabels: map[string]string{
"customer": "x",
"customer": "x",
"e2e.suite.capsule/name": "owner_test",
},
},
{
MatchLabels: map[string]string{
"team": "infrastructure",
"team": "infrastructure",
"e2e.suite.capsule/name": "owner_test",
},
},
},
@@ -126,7 +130,8 @@ var _ = Describe("Owners", Ordered, Label("config", "tenant", "permissions", "ow
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-owners-infra",
Labels: map[string]string{
"team": "infrastructure",
"team": "infrastructure",
"e2e.suite.capsule/name": "owner_test",
},
},
Spec: capsulev1beta2.TenantOwnerSpec{
@@ -147,7 +152,8 @@ var _ = Describe("Owners", Ordered, Label("config", "tenant", "permissions", "ow
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-owners-devops",
Labels: map[string]string{
"team": "devops",
"team": "devops",
"e2e.suite.capsule/name": "owner_test",
},
},
Spec: capsulev1beta2.TenantOwnerSpec{
@@ -168,8 +174,9 @@ var _ = Describe("Owners", Ordered, Label("config", "tenant", "permissions", "ow
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-owners-common",
Labels: map[string]string{
"team": "infrastructure",
"customer": "x",
"team": "infrastructure",
"customer": "x",
"e2e.suite.capsule/name": "owner_test",
},
},
Spec: capsulev1beta2.TenantOwnerSpec{
@@ -211,8 +218,9 @@ var _ = Describe("Owners", Ordered, Label("config", "tenant", "permissions", "ow
ObjectMeta: metav1.ObjectMeta{
Name: "e2e-owners-common-user",
Labels: map[string]string{
"team": "infrastructure",
"customer": "x",
"team": "infrastructure",
"customer": "x",
"e2e.suite.capsule/name": "owner_test",
},
},
Spec: capsulev1beta2.TenantOwnerSpec{
+13 -24
View File
@@ -25,6 +25,7 @@ import (
"sigs.k8s.io/controller-runtime/pkg/reconcile"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/internal/controllers/tls"
"github.com/projectcapsule/capsule/internal/controllers/utils"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/runtime/admission"
@@ -143,13 +144,15 @@ func (r *mutatingReconciler) reconcileConfiguration(
obj.SetAnnotations(annotations)
// Do not overwrite caBundle. cert-manager or the legacy TLS reconciler owns it.
if cfg.Client.CABundle == nil {
obj.Webhooks = preserveMutatingWebhookCABundles(obj.Webhooks, desiredHooks)
} else {
obj.Webhooks = desiredHooks
obj.Webhooks = desiredHooks
caCert, err := tls.FetchCurrentCaBundleForAdmission(ctx, r.client, r.configuration, cfg.Client.CABundle)
if err != nil {
return err
}
preserveMutatingWebhookCABundles(obj.Webhooks, caCert)
return err
})
if err != nil {
@@ -220,24 +223,10 @@ func (r *mutatingReconciler) webhooks(
}
func preserveMutatingWebhookCABundles(
existing []admissionv1.MutatingWebhook,
desired []admissionv1.MutatingWebhook,
) []admissionv1.MutatingWebhook {
existingByName := make(map[string][]byte, len(existing))
for _, hook := range existing {
if len(hook.ClientConfig.CABundle) == 0 {
continue
}
existingByName[hook.Name] = append([]byte(nil), hook.ClientConfig.CABundle...)
hooks []admissionv1.MutatingWebhook,
caBundle []byte,
) {
for i := range hooks {
hooks[i].ClientConfig.CABundle = append([]byte(nil), caBundle...)
}
for i := range desired {
if caBundle, ok := existingByName[desired[i].Name]; ok {
desired[i].ClientConfig.CABundle = append([]byte(nil), caBundle...)
}
}
return desired
}
+13 -28
View File
@@ -25,6 +25,7 @@ import (
"sigs.k8s.io/controller-runtime/pkg/reconcile"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/internal/controllers/tls"
"github.com/projectcapsule/capsule/internal/controllers/utils"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/runtime/admission"
@@ -145,13 +146,15 @@ func (r *validatingReconciler) reconcileValidatingConfiguration(
obj.SetAnnotations(annotations)
// Do not overwrite caBundle. cert-manager or the legacy TLS reconciler owns it.
if cfg.Client.CABundle == nil {
obj.Webhooks = preserveValidatingWebhookCABundles(obj.Webhooks, desiredHooks)
} else {
obj.Webhooks = desiredHooks
obj.Webhooks = desiredHooks
caCert, err := tls.FetchCurrentCaBundleForAdmission(ctx, r.client, r.configuration, cfg.Client.CABundle)
if err != nil {
return err
}
preserveValidatingWebhookCABundles(obj.Webhooks, caCert)
return err
})
if err != nil {
@@ -222,28 +225,10 @@ func (r *validatingReconciler) validatingWebhooks(
}
func preserveValidatingWebhookCABundles(
existing []admissionv1.ValidatingWebhook,
desired []admissionv1.ValidatingWebhook,
) []admissionv1.ValidatingWebhook {
existingByName := make(map[string][]byte, len(existing))
for _, hook := range existing {
if len(hook.ClientConfig.CABundle) == 0 {
continue
}
existingByName[hook.Name] = append([]byte(nil), hook.ClientConfig.CABundle...)
hooks []admissionv1.ValidatingWebhook,
caBundle []byte,
) {
for i := range hooks {
hooks[i].ClientConfig.CABundle = append([]byte(nil), caBundle...)
}
for i := range desired {
if len(desired[i].ClientConfig.CABundle) > 0 {
continue
}
if caBundle, ok := existingByName[desired[i].Name]; ok {
desired[i].ClientConfig.CABundle = append([]byte(nil), caBundle...)
}
}
return desired
}
+63
View File
@@ -0,0 +1,63 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package tenant
import (
"context"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/workqueue"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
)
func (r *Manager) enqueueForTenantsWithCondition(
ctx context.Context,
obj client.Object,
q workqueue.TypedRateLimitingInterface[reconcile.Request],
fn func(*capsulev1beta2.Tenant, client.Object) bool,
) {
var tenants capsulev1beta2.TenantList
if err := r.List(ctx, &tenants); err != nil {
r.Log.Error(err, "failed to list Tenants for class event")
return
}
for i := range tenants.Items {
tnt := &tenants.Items[i]
if !fn(tnt, obj) {
continue
}
q.Add(reconcile.Request{
NamespacedName: types.NamespacedName{
Name: tnt.Name,
},
})
}
}
func (r *Manager) enqueueAllTenants(ctx context.Context, _ client.Object) []reconcile.Request {
var tenants capsulev1beta2.TenantList
if err := r.List(ctx, &tenants); err != nil {
r.Log.Error(err, "failed to list Tenants for class event")
return nil
}
reqs := make([]reconcile.Request, 0, len(tenants.Items))
for i := range tenants.Items {
reqs = append(reqs, reconcile.Request{
NamespacedName: types.NamespacedName{
Name: tenants.Items[i].Name,
},
})
}
return reqs
}
+13 -9
View File
@@ -1,7 +1,6 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
//nolint:dupl
package tenant
import (
@@ -14,6 +13,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"github.com/go-logr/logr"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/meta"
)
@@ -21,26 +21,30 @@ import (
// Ensuring all the LimitRange are applied to each Namespace handled by the Tenant.
//
func (r *Manager) syncLimitRanges(ctx context.Context, tenant *capsulev1beta2.Tenant) error {
func (r *Manager) syncLimitRanges(ctx context.Context, log logr.Logger, tenant *capsulev1beta2.Tenant) error {
// getting requested LimitRange keys
keys := make([]string, 0, len(tenant.Spec.LimitRanges.Items)) //nolint:staticcheck
keys := make([]string, 0, len(tenant.Spec.LimitRanges.Items))
//nolint:staticcheck
for i := range tenant.Spec.LimitRanges.Items {
keys = append(keys, strconv.Itoa(i))
}
return runForTenantNamespaces(ctx, tenant, func(ctx context.Context, namespace string) error {
return r.syncLimitRange(ctx, tenant, namespace, keys)
return r.syncLimitRange(ctx, log, tenant, namespace, keys)
})
}
func (r *Manager) syncLimitRange(ctx context.Context, tenant *capsulev1beta2.Tenant, namespace string, keys []string) (err error) {
func (r *Manager) syncLimitRange(
ctx context.Context,
log logr.Logger,
tenant *capsulev1beta2.Tenant,
namespace string,
keys []string,
) (err error) {
if err = r.pruningResources(ctx, namespace, keys, &corev1.LimitRange{}); err != nil {
return err
}
//nolint:staticcheck
for i, spec := range tenant.Spec.LimitRanges.Items {
target := &corev1.LimitRange{
ObjectMeta: metav1.ObjectMeta{
@@ -71,7 +75,7 @@ func (r *Manager) syncLimitRange(ctx context.Context, tenant *capsulev1beta2.Ten
})
if err != nil {
if apierrors.HasStatusCause(err, corev1.NamespaceTerminatingCause) {
r.Log.V(4).Info(
log.Info(
"skipping LimitRange sync because namespace is terminating",
"name", target.Name,
"namespace", target.Namespace,
@@ -84,7 +88,7 @@ func (r *Manager) syncLimitRange(ctx context.Context, tenant *capsulev1beta2.Ten
return err
}
r.Log.V(4).Info("LimitRange sync result: "+string(res), "name", target.Name, "namespace", target.Namespace)
log.Info("LimitRange sync result: "+string(res), "name", target.Name, "namespace", target.Namespace)
}
return nil
+17 -17
View File
@@ -236,13 +236,13 @@ func (r *Manager) SetupWithManager(mgr ctrl.Manager, ctrlConfig utils.Controller
}
func (r *Manager) Reconcile(ctx context.Context, request ctrl.Request) (result ctrl.Result, err error) {
r.Log = r.Log.WithValues("Request.Name", request.Name)
log := r.Log.WithValues("Request.Name", request.Name)
// Fetch the Tenant instance
instance := &capsulev1beta2.Tenant{}
if err = r.Get(ctx, request.NamespacedName, instance); err != nil {
if apierrors.IsNotFound(err) {
r.Log.V(5).Info("request object not found, could have been deleted after reconcile request")
log.V(5).Info("request object not found, could have been deleted after reconcile request")
// If tenant was deleted or cannot be found, clean up metrics
r.Metrics.DeleteAllMetricsForTenant(request.Name)
@@ -266,7 +266,7 @@ func (r *Manager) Reconcile(ctx context.Context, request ctrl.Request) (result c
return reconcile.Result{}, updateErr
}
reconcileError := r.reconcile(ctx, instance)
reconcileError := r.reconcile(ctx, log, instance)
defer func() {
r.syncTenantStatusMetrics(instance)
@@ -293,7 +293,7 @@ func (r *Manager) Reconcile(ctx context.Context, request ctrl.Request) (result c
}
// Collect available resources
if err = r.collectAvailableResources(ctx, instance); err != nil {
if err = r.collectAvailableResources(ctx, log, instance); err != nil {
err = fmt.Errorf("cannot collect available resources: %w", err)
return reconcile.Result{}, err
@@ -306,7 +306,7 @@ func (r *Manager) Reconcile(ctx context.Context, request ctrl.Request) (result c
return reconcile.Result{}, reconcileError
}
func (r *Manager) reconcile(ctx context.Context, instance *capsulev1beta2.Tenant) (err error) {
func (r *Manager) reconcile(ctx context.Context, log logr.Logger, instance *capsulev1beta2.Tenant) (err error) {
var errs []error
// Collect Ownership/Promotions for Status
@@ -315,9 +315,9 @@ func (r *Manager) reconcile(ctx context.Context, instance *capsulev1beta2.Tenant
}
// Reconcile Namespaces
r.Log.V(4).Info("starting processing of Namespaces", "items", len(instance.Status.Namespaces))
log.V(4).Info("starting processing of Namespaces", "items", len(instance.Status.Namespaces))
if err = r.reconcileNamespaces(ctx, instance); err != nil {
if err = r.reconcileNamespaces(ctx, log, instance); err != nil {
errs = append(errs, fmt.Errorf("namespace(s) had reconciliation errors: %w", err))
}
@@ -328,37 +328,37 @@ func (r *Manager) reconcile(ctx context.Context, instance *capsulev1beta2.Tenant
}
// Ensuring ResourceQuota
r.Log.V(4).Info("ensuring limit resources count is updated")
log.V(4).Info("ensuring limit resources count is updated")
if err = r.syncCustomResourceQuotaUsages(ctx, instance); err != nil {
errs = append(errs, fmt.Errorf("cannot count limited resources: %w", err))
}
// Ensuring NetworkPolicy resources
r.Log.V(4).Info("starting processing of Network Policies")
log.V(4).Info("starting processing of Network Policies")
if err = r.syncNetworkPolicies(ctx, instance); err != nil {
if err = r.syncNetworkPolicies(ctx, log, instance); err != nil {
errs = append(errs, fmt.Errorf("cannot sync networkPolicy items: %w", err))
}
// Ensuring LimitRange resources
r.Log.V(4).Info("Starting processing of Limit Ranges", "items", len(instance.Spec.LimitRanges.Items)) //nolint:staticcheck
r.Log.V(4).Info("Starting processing of Limit Ranges", "items", len(instance.Spec.LimitRanges.Items))
if err = r.syncLimitRanges(ctx, instance); err != nil {
if err = r.syncLimitRanges(ctx, log, instance); err != nil {
errs = append(errs, fmt.Errorf("cannot sync limitrange items: %w", err))
}
// Ensuring ResourceQuota resources
r.Log.V(4).Info("Starting processing of Resource Quotas", "items", len(instance.Spec.ResourceQuota.Items))
log.V(4).Info("Starting processing of Resource Quotas", "items", len(instance.Spec.ResourceQuota.Items))
if err = r.syncResourceQuotas(ctx, instance); err != nil {
if err = r.syncResourceQuotas(ctx, log, instance); err != nil {
errs = append(errs, fmt.Errorf("cannot sync resourcequota items: %w", err))
}
// Ensuring RoleBinding resources
r.Log.V(4).Info("Ensuring RoleBindings for Owners and Tenant")
log.V(4).Info("Ensuring RoleBindings for Owners and Tenant")
if err = r.syncRoleBindings(ctx, r.Log, instance); err != nil {
if err = r.syncRoleBindings(ctx, log, instance); err != nil {
errs = append(errs, fmt.Errorf("cannot sync rolebindings items: %w", err))
}
@@ -366,7 +366,7 @@ func (r *Manager) reconcile(ctx context.Context, instance *capsulev1beta2.Tenant
return err
}
r.Log.V(4).Info("Tenant reconciling completed")
log.V(4).Info("Tenant reconciling completed")
return err
}
+7 -6
View File
@@ -9,6 +9,7 @@ import (
"fmt"
"maps"
"github.com/go-logr/logr"
"golang.org/x/sync/errgroup"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
@@ -24,7 +25,7 @@ import (
)
// Ensuring all annotations are applied to each Namespace handled by the Tenant.
func (r *Manager) reconcileNamespaces(ctx context.Context, tnt *capsulev1beta2.Tenant) (err error) {
func (r *Manager) reconcileNamespaces(ctx context.Context, log logr.Logger, tnt *capsulev1beta2.Tenant) (err error) {
if tnt.DeletionTimestamp != nil {
for _, ns := range tnt.Status.Spaces {
ns := &corev1.Namespace{
@@ -36,7 +37,7 @@ func (r *Manager) reconcileNamespaces(ctx context.Context, tnt *capsulev1beta2.T
if err := r.Delete(ctx, ns, &client.DeleteOptions{
PropagationPolicy: ptr.To(metav1.DeletePropagationBackground),
}); err != nil && !apierrors.IsNotFound(err) {
r.Log.Error(err, "unable to delete tenant namespace",
log.Error(err, "unable to delete tenant namespace",
"tenant", tnt.GetName(),
"namespace", ns.Name,
)
@@ -66,13 +67,13 @@ func (r *Manager) reconcileNamespaces(ctx context.Context, tnt *capsulev1beta2.T
ns := list.Items[i].DeepCopy()
group.Go(func() error {
stat, err := r.reconcileNamespace(ctx, ns, tnt)
stat, err := r.reconcileNamespace(ctx, log, ns, tnt)
if stat != nil {
results <- stat
}
if err != nil {
r.Log.Error(err, "failed to reconcile namespace",
log.Error(err, "failed to reconcile namespace",
"tenant", tnt.GetName(),
"namespace", ns.GetName(),
)
@@ -125,7 +126,7 @@ func (r *Manager) reconcileNamespaces(ctx context.Context, tnt *capsulev1beta2.T
return err
}
func (r *Manager) reconcileNamespace(ctx context.Context, namespace *corev1.Namespace, tnt *capsulev1beta2.Tenant) (
func (r *Manager) reconcileNamespace(ctx context.Context, log logr.Logger, namespace *corev1.Namespace, tnt *capsulev1beta2.Tenant) (
stat *capsulev1beta2.TenantStatusNamespaceItem,
err error,
) {
@@ -246,7 +247,7 @@ func (r *Manager) reconcileNamespace(ctx context.Context, namespace *corev1.Name
}
// Collect Rules for namespace
err = r.reconcileRuleStatus(ctx, tnt, namespace)
err = r.reconcileRuleStatus(ctx, log, tnt, namespace)
if err != nil {
return stat, err
}
@@ -1,7 +1,6 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
//nolint:dupl
package tenant
import (
@@ -9,6 +8,7 @@ import (
"fmt"
"strconv"
"github.com/go-logr/logr"
corev1 "k8s.io/api/core/v1"
networkingv1 "k8s.io/api/networking/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
@@ -22,25 +22,23 @@ import (
// Ensuring all the NetworkPolicies are applied to each Namespace handled by the Tenant.
//
func (r *Manager) syncNetworkPolicies(ctx context.Context, tenant *capsulev1beta2.Tenant) error {
keys := make([]string, 0, len(tenant.Spec.NetworkPolicies.Items)) //nolint:staticcheck
func (r *Manager) syncNetworkPolicies(ctx context.Context, log logr.Logger, tenant *capsulev1beta2.Tenant) error {
keys := make([]string, 0, len(tenant.Spec.NetworkPolicies.Items))
//nolint:staticcheck
for i := range tenant.Spec.NetworkPolicies.Items {
keys = append(keys, strconv.Itoa(i))
}
return runForTenantNamespaces(ctx, tenant, func(ctx context.Context, namespace string) error {
return r.syncNetworkPolicy(ctx, tenant, namespace, keys)
return r.syncNetworkPolicy(ctx, log, tenant, namespace, keys)
})
}
func (r *Manager) syncNetworkPolicy(ctx context.Context, tenant *capsulev1beta2.Tenant, namespace string, keys []string) (err error) {
func (r *Manager) syncNetworkPolicy(ctx context.Context, log logr.Logger, tenant *capsulev1beta2.Tenant, namespace string, keys []string) (err error) {
if err = r.pruningResources(ctx, namespace, keys, &networkingv1.NetworkPolicy{}); err != nil {
return err
}
//nolint:staticcheck
for i, spec := range tenant.Spec.NetworkPolicies.Items {
target := &networkingv1.NetworkPolicy{
ObjectMeta: metav1.ObjectMeta{
@@ -71,7 +69,7 @@ func (r *Manager) syncNetworkPolicy(ctx context.Context, tenant *capsulev1beta2.
})
if err != nil {
if apierrors.HasStatusCause(err, corev1.NamespaceTerminatingCause) {
r.Log.V(4).Info(
log.V(4).Info(
"skipping NetworkPolicy sync because namespace is terminating",
"name", target.Name,
"namespace", target.Namespace,
@@ -84,7 +82,7 @@ func (r *Manager) syncNetworkPolicy(ctx context.Context, tenant *capsulev1beta2.
return err
}
r.Log.V(4).Info("network Policy sync result: "+string(res), "name", target.Name, "namespace", target.Namespace)
log.V(4).Info("network Policy sync result: "+string(res), "name", target.Name, "namespace", target.Namespace)
}
return nil
+23 -14
View File
@@ -10,6 +10,7 @@ import (
"strconv"
"strings"
"github.com/go-logr/logr"
"golang.org/x/sync/errgroup"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
@@ -43,7 +44,7 @@ import (
// In case of Namespace-scoped Resource Budget, we're just replicating the resources across all registered Namespaces.
//nolint:cyclop
func (r *Manager) syncResourceQuotas(ctx context.Context, tenant *capsulev1beta2.Tenant) (err error) { //nolint:gocognit
func (r *Manager) syncResourceQuotas(ctx context.Context, log logr.Logger, tenant *capsulev1beta2.Tenant) (err error) { //nolint:gocognit
// Remove prior metrics, to avoid cleaning up for metrics of deleted ResourceQuotas
r.Metrics.DeleteTenantResourceMetrics(tenant.Name)
// Expose the namespace quota and usage as metrics for the tenant
@@ -78,20 +79,20 @@ func (r *Manager) syncResourceQuotas(ctx context.Context, tenant *capsulev1beta2
var tntRequirement *labels.Requirement
if tntRequirement, scopeErr = labels.NewRequirement(meta.NewTenantLabel, selection.Equals, []string{tenant.Name}); scopeErr != nil {
r.Log.Error(scopeErr, "cannot build ResourceQuota Tenant requirement")
log.Error(scopeErr, "cannot build ResourceQuota Tenant requirement")
}
// Requirement to list ResourceQuota for the current index
var indexRequirement *labels.Requirement
if indexRequirement, scopeErr = labels.NewRequirement(meta.ResourceQuotaLabel, selection.Equals, []string{strconv.Itoa(index)}); scopeErr != nil {
r.Log.Error(scopeErr, "cannot build ResourceQuota index requirement")
log.Error(scopeErr, "cannot build ResourceQuota index requirement")
}
// Listing all the ResourceQuota according to the said requirements.
// These are required since Capsule is going to sum all the used quota to
// sum them and get the Tenant one.
list := &corev1.ResourceQuotaList{}
if scopeErr = r.reader.List(ctx, list, &client.ListOptions{LabelSelector: labels.NewSelector().Add(*tntRequirement).Add(*indexRequirement)}); scopeErr != nil {
r.Log.Error(scopeErr, "cannot list ResourceQuota", "tenantFilter", tntRequirement.String(), "indexFilter", indexRequirement.String())
log.Error(scopeErr, "cannot list ResourceQuota", "tenantFilter", tntRequirement.String(), "indexFilter", indexRequirement.String())
return scopeErr
}
@@ -102,7 +103,7 @@ func (r *Manager) syncResourceQuotas(ctx context.Context, tenant *capsulev1beta2
// For this case, we're going to block the Quota setting the Hard as the
// used one.
for name, hardQuota := range resourceQuota.Hard {
r.Log.V(4).Info("desired hard " + name.String() + " quota is " + hardQuota.String())
log.V(4).Info("desired hard " + name.String() + " quota is " + hardQuota.String())
// Getting the whole usage across all the Tenant Namespaces
var quantity resource.Quantity
@@ -110,7 +111,7 @@ func (r *Manager) syncResourceQuotas(ctx context.Context, tenant *capsulev1beta2
quantity.Add(item.Status.Used[name])
}
r.Log.V(4).Info("computed " + name.String() + " quota for the whole Tenant is " + quantity.String())
log.V(4).Info("computed " + name.String() + " quota for the whole Tenant is " + quantity.String())
// Expose usage and limit metrics for the resource (name) of the ResourceQuota (index)
r.Metrics.TenantResourceUsageGauge.WithLabelValues(
@@ -171,8 +172,8 @@ func (r *Manager) syncResourceQuotas(ctx context.Context, tenant *capsulev1beta2
}
}
if scopeErr = r.resourceQuotasUpdate(ctx, name, quantity, toKeep, resourceQuota.Hard[name], list.Items...); scopeErr != nil {
r.Log.Error(scopeErr, "cannot proceed with outer ResourceQuota")
if scopeErr = r.resourceQuotasUpdate(ctx, log, name, quantity, toKeep, resourceQuota.Hard[name], list.Items...); scopeErr != nil {
log.Error(scopeErr, "cannot proceed with outer ResourceQuota")
return scopeErr
}
@@ -214,14 +215,14 @@ func (r *Manager) syncResourceQuotas(ctx context.Context, tenant *capsulev1beta2
}
group.Go(func() error {
return r.syncResourceQuota(ctx, tenant, namespace, keys)
return r.syncResourceQuota(ctx, log, tenant, namespace, keys)
})
}
return group.Wait()
}
func (r *Manager) syncResourceQuota(ctx context.Context, tenant *capsulev1beta2.Tenant, namespace string, keys []string) (err error) {
func (r *Manager) syncResourceQuota(ctx context.Context, log logr.Logger, tenant *capsulev1beta2.Tenant, namespace string, keys []string) (err error) {
// getting ResourceQuota labels for the mutateFn
var typeLabel string
@@ -274,7 +275,7 @@ func (r *Manager) syncResourceQuota(ctx context.Context, tenant *capsulev1beta2.
})
if err != nil {
if apierrors.HasStatusCause(err, corev1.NamespaceTerminatingCause) {
r.Log.V(4).Info(
log.V(4).Info(
"skipping ResourceQuota sync because namespace is terminating",
"name", target.Name,
"namespace", target.Namespace,
@@ -287,7 +288,7 @@ func (r *Manager) syncResourceQuota(ctx context.Context, tenant *capsulev1beta2.
return err
}
r.Log.V(4).Info("resource Quota sync result: "+string(res), "name", target.Name, "namespace", target.Namespace)
log.V(4).Info("resource Quota sync result: "+string(res), "name", target.Name, "namespace", target.Namespace)
}
return nil
@@ -296,7 +297,15 @@ func (r *Manager) syncResourceQuota(ctx context.Context, tenant *capsulev1beta2.
// Serial ResourceQuota processing is expensive: using Go routines we can speed it up.
// In case of multiple errors these are logged properly, returning a generic error since we have to repush back the
// reconciliation loop.
func (r *Manager) resourceQuotasUpdate(ctx context.Context, resourceName corev1.ResourceName, actual resource.Quantity, toKeep sets.Set[corev1.ResourceName], limit resource.Quantity, list ...corev1.ResourceQuota) (err error) {
func (r *Manager) resourceQuotasUpdate(
ctx context.Context,
log logr.Logger,
resourceName corev1.ResourceName,
actual resource.Quantity,
toKeep sets.Set[corev1.ResourceName],
limit resource.Quantity,
list ...corev1.ResourceQuota,
) (err error) {
group := new(errgroup.Group)
annotationsToKeep := sets.New[string]()
@@ -372,7 +381,7 @@ func (r *Manager) resourceQuotasUpdate(ctx context.Context, resourceName corev1.
if err = group.Wait(); err != nil {
// We had an error and we mark the whole transaction as failed
// to process it another time according to the Tenant controller back-off factor.
r.Log.Error(err, "cannot update outer ResourceQuotas", "resourceName", resourceName.String())
log.Error(err, "cannot update outer ResourceQuotas", "resourceName", resourceName.String())
err = fmt.Errorf("update of outer ResourceQuota items has failed: %w", err)
}
+4 -3
View File
@@ -57,12 +57,13 @@ func (r *Manager) syncRoleBindings(ctx context.Context, log logr.Logger, tenant
}
return runForTenantNamespaces(ctx, tenant, func(ctx context.Context, namespace string) error {
return r.syncAdditionalRoleBinding(ctx, tenant, namespace, namespaceBindings[namespace])
return r.syncAdditionalRoleBinding(ctx, log, tenant, namespace, namespaceBindings[namespace])
})
}
func (r *Manager) syncAdditionalRoleBinding(
ctx context.Context,
log logr.Logger,
tenant *capsulev1beta2.Tenant,
ns string,
bindings map[string]rbac.AdditionalRoleBindingsSpec,
@@ -114,7 +115,7 @@ func (r *Manager) syncAdditionalRoleBinding(
})
if err != nil {
if apierrors.HasStatusCause(err, corev1.NamespaceTerminatingCause) {
r.Log.V(4).Info(
log.V(4).Info(
"skipping RoleBinding sync because namespace is terminating",
"name", target.Name,
"namespace", target.Namespace,
@@ -127,7 +128,7 @@ func (r *Manager) syncAdditionalRoleBinding(
return fmt.Errorf("%w (role: %s)", err, roleBinding.ClusterRoleName)
}
r.Log.V(4).Info(fmt.Sprintf("roleBinding sync result: %s", string(res)), "name", target.Name, "namespace", target.Namespace)
log.V(4).Info(fmt.Sprintf("roleBinding sync result: %s", string(res)), "name", target.Name, "namespace", target.Namespace)
}
// Prune at finish to prevent gaps
+5 -1
View File
@@ -6,6 +6,7 @@ package tenant
import (
"context"
"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"
@@ -19,6 +20,7 @@ import (
func (r *Manager) reconcileRuleStatus(
ctx context.Context,
log logr.Logger,
tnt *capsulev1beta2.Tenant,
ns *corev1.Namespace,
) error {
@@ -30,6 +32,7 @@ func (r *Manager) reconcileRuleStatus(
return r.ensureRuleStatus(
ctx,
log,
tnt,
ns,
ruleBody,
@@ -38,6 +41,7 @@ func (r *Manager) reconcileRuleStatus(
func (r *Manager) ensureRuleStatus(
ctx context.Context,
log logr.Logger,
tnt *capsulev1beta2.Tenant,
namespace *corev1.Namespace,
body *api.NamespaceRuleBodyNamespace,
@@ -68,7 +72,7 @@ func (r *Manager) ensureRuleStatus(
})
if err != nil {
if apierrors.HasStatusCause(err, corev1.NamespaceTerminatingCause) {
r.Log.V(4).Info(
log.V(4).Info(
"skipping RuleStatus sync because namespace is terminating",
"name", rule.Name,
"namespace", rule.Namespace,
+2 -5
View File
@@ -8,6 +8,7 @@ import (
"regexp"
"sort"
"github.com/go-logr/logr"
nodev1 "k8s.io/api/node/v1"
resources "k8s.io/api/resource/v1"
schedulingv1 "k8s.io/api/scheduling/v1"
@@ -19,7 +20,6 @@ import (
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/log"
gatewayv1 "sigs.k8s.io/gateway-api/apis/v1"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
@@ -148,9 +148,7 @@ func (r *Manager) collectRBAC(ctx context.Context, tnt *capsulev1beta2.Tenant) (
return nil
}
func (r *Manager) collectAvailableResources(ctx context.Context, tnt *capsulev1beta2.Tenant) (err error) {
log := log.FromContext(ctx)
func (r *Manager) collectAvailableResources(ctx context.Context, log logr.Logger, tnt *capsulev1beta2.Tenant) (err error) {
if r.classes.device {
log.V(5).Info("collecting available deviceclasses")
@@ -344,7 +342,6 @@ func listObjectNamesBySelector(
var regex *regexp.Regexp
//nolint:staticcheck
if allowed.Regex != "" {
regex, err = regexp.Compile(allowed.Regex)
if err != nil {
-51
View File
@@ -14,11 +14,8 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/selection"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
"k8s.io/client-go/util/workqueue"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/meta"
@@ -75,54 +72,6 @@ func runForTenantNamespaces(
return errors.Join(joined...)
}
func (r *Manager) enqueueForTenantsWithCondition(
ctx context.Context,
obj client.Object,
q workqueue.TypedRateLimitingInterface[reconcile.Request],
fn func(*capsulev1beta2.Tenant, client.Object) bool,
) {
var tenants capsulev1beta2.TenantList
if err := r.List(ctx, &tenants); err != nil {
r.Log.Error(err, "failed to list Tenants for class event")
return
}
for i := range tenants.Items {
tnt := &tenants.Items[i]
if !fn(tnt, obj) {
continue
}
q.Add(reconcile.Request{
NamespacedName: types.NamespacedName{
Name: tnt.Name,
},
})
}
}
func (r *Manager) enqueueAllTenants(ctx context.Context, _ client.Object) []reconcile.Request {
var tenants capsulev1beta2.TenantList
if err := r.List(ctx, &tenants); err != nil {
r.Log.Error(err, "failed to list Tenants for class event")
return nil
}
reqs := make([]reconcile.Request, 0, len(tenants.Items))
for i := range tenants.Items {
reqs = append(reqs, reconcile.Request{
NamespacedName: types.NamespacedName{
Name: tenants.Items[i].Name,
},
})
}
return reqs
}
// pruningResources is taking care of removing the no more requested sub-resources as LimitRange, ResourceQuota or
// NetworkPolicy using the "exists" and "notin" LabelSelector to perform an outer-join removal.
func (r *Manager) pruningResources(ctx context.Context, ns string, keys []string, obj client.Object) (err error) {
@@ -0,0 +1,68 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package tenantowners
import (
"context"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/workqueue"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/pkg/api/rbac"
indexer "github.com/projectcapsule/capsule/pkg/runtime/indexers/tenantowner"
)
func (r *TenantOwnerManager) enqueueTenantOwnerRequests(
ctx context.Context,
q workqueue.TypedRateLimitingInterface[ctrl.Request],
owners rbac.OwnerStatusListSpec,
) {
seen := make(map[types.NamespacedName]struct{}, len(owners))
for _, owner := range owners {
if owner.Name == "" {
continue
}
tenantOwnerList := &capsulev1beta2.TenantOwnerList{}
if err := r.List(
ctx,
tenantOwnerList,
client.MatchingFields{
indexer.NameIndexerFieldName: owner.Name,
},
); err != nil {
r.Log.Error(
err,
"Failed to list TenantOwners by spec.name",
"ownerName", owner.Name,
"ownerKind", owner.Kind,
)
continue
}
for _, tenantOwner := range tenantOwnerList.Items {
if tenantOwner.Spec.Kind != owner.Kind {
continue
}
key := types.NamespacedName{
Name: tenantOwner.Name,
}
if _, ok := seen[key]; ok {
continue
}
seen[key] = struct{}{}
q.Add(ctrl.Request{
NamespacedName: key,
})
}
}
}
+231
View File
@@ -0,0 +1,231 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package tenantowners
import (
"context"
"fmt"
"sort"
"github.com/go-logr/logr"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
"k8s.io/client-go/util/workqueue"
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/event"
"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"
"github.com/projectcapsule/capsule/internal/controllers/utils"
capmeta "github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/api/rbac"
indexer "github.com/projectcapsule/capsule/pkg/runtime/indexers/tenant"
"github.com/projectcapsule/capsule/pkg/tenant"
)
// TenantOwnerManager reconciles TenantOwner objects and keeps
// status.tenants / status.matchedTenants / status.conditions in sync.
//
// It watches both TenantOwner objects (for spec changes) and Tenant objects
// (for selector changes), re-enqueuing every TenantOwner whenever any Tenant
// changes, because a single Tenant change can affect any subset of TenantOwners.
type TenantOwnerManager struct {
client.Client
reader client.Reader
Log logr.Logger
}
func (r *TenantOwnerManager) SetupWithManager(mgr ctrl.Manager, ctrlConfig utils.ControllerOptions) error {
r.reader = mgr.GetAPIReader()
return ctrl.NewControllerManagedBy(mgr).
Named("capsule/tenant-owner-status").
For(
&capsulev1beta2.TenantOwner{},
builder.WithPredicates(
predicate.Or(
predicate.GenerationChangedPredicate{},
predicate.LabelChangedPredicate{},
),
),
).
Watches(
&capsulev1beta2.Tenant{},
handler.TypedFuncs[client.Object, ctrl.Request]{
CreateFunc: func(
ctx context.Context,
e event.TypedCreateEvent[client.Object],
q workqueue.TypedRateLimitingInterface[reconcile.Request],
) {
tnt, ok := e.Object.(*capsulev1beta2.Tenant)
if !ok {
return
}
r.enqueueTenantOwnerRequests(ctx, q, tnt.Status.Owners)
},
UpdateFunc: func(
ctx context.Context,
e event.TypedUpdateEvent[client.Object],
q workqueue.TypedRateLimitingInterface[reconcile.Request],
) {
oldTnt, ok1 := e.ObjectOld.(*capsulev1beta2.Tenant)
newTnt, ok2 := e.ObjectNew.(*capsulev1beta2.Tenant)
if !ok1 || !ok2 {
return
}
owners := make(
rbac.OwnerStatusListSpec,
0,
len(oldTnt.Status.Owners)+len(newTnt.Status.Owners),
)
owners = append(owners, oldTnt.Status.Owners...)
owners = append(owners, newTnt.Status.Owners...)
r.enqueueTenantOwnerRequests(ctx, q, owners)
},
DeleteFunc: func(
ctx context.Context,
e event.TypedDeleteEvent[client.Object],
q workqueue.TypedRateLimitingInterface[reconcile.Request],
) {
tnt, ok := e.Object.(*capsulev1beta2.Tenant)
if !ok {
return
}
r.enqueueTenantOwnerRequests(ctx, q, tnt.Status.Owners)
},
},
).
Complete(r)
}
func (r *TenantOwnerManager) Reconcile(ctx context.Context, req ctrl.Request) (result ctrl.Result, err error) {
log := r.Log.WithValues("tenantowner", req.Name)
instance := &capsulev1beta2.TenantOwner{}
if err = r.Get(ctx, req.NamespacedName, instance); err != nil {
if apierrors.IsNotFound(err) {
return reconcile.Result{}, nil
}
return reconcile.Result{}, err
}
matchedTenants, reconcileErr := r.reconcileMatchedTenants(ctx, log, instance)
log.V(5).Info("matched tenants", "count", len(matchedTenants))
if statusErr := r.updateTenantOwnerStatus(ctx, instance, matchedTenants, reconcileErr); statusErr != nil {
if apierrors.IsNotFound(statusErr) {
return reconcile.Result{}, nil
}
return reconcile.Result{}, fmt.Errorf("cannot update TenantOwner status: %w", statusErr)
}
return reconcile.Result{}, reconcileErr
}
// reconcileMatchedTenants lists all Tenants and returns the sorted names of
// those whose spec.permissions.matchOwners selectors select this TenantOwner.
// All work is done against the in-memory cache — no direct API server calls.
func (r *TenantOwnerManager) reconcileMatchedTenants(ctx context.Context, log logr.Logger, to *capsulev1beta2.TenantOwner) ([]string, error) {
if to.Spec.Name == "" || to.Spec.Kind.String() == "" {
return nil, fmt.Errorf(
"TenantOwner %s has incomplete owner reference: kind=%q name=%q",
to.Name,
to.Spec.Kind,
to.Spec.Name,
)
}
ownerKey := tenant.OwnerKindIndexKey(to.Spec.Kind.String(), to.Spec.Name)
tnts := &capsulev1beta2.TenantList{}
if err := r.List(
ctx,
tnts,
client.MatchingFields{
indexer.OwnerKindIndexerFieldName: ownerKey,
},
); err != nil {
return nil, fmt.Errorf(
"listing Tenants by owner %s/%s: %w",
to.Spec.Kind,
to.Spec.Name,
err,
)
}
log.V(4).Info("found tenant references", "tenants", len(tnts.Items))
matched := make([]string, 0, len(tnts.Items))
for i := range tnts.Items {
matched = append(matched, tnts.Items[i].Name)
}
sort.Strings(matched)
return matched, nil
}
// updateTenantOwnerStatus writes the reconciled status back to the API server
// using a RetryOnConflict re-GET loop to avoid clobbering concurrent updates.
func (r *TenantOwnerManager) updateTenantOwnerStatus(
ctx context.Context,
instance *capsulev1beta2.TenantOwner,
matchedTenants []string,
reconcileError error,
) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
latest := &capsulev1beta2.TenantOwner{}
if err := r.reader.Get(ctx, types.NamespacedName{Name: instance.Name}, latest); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return err
}
latest.Status.ObservedGeneration = latest.GetGeneration()
latest.Status.Tenants = matchedTenants
readyCondition := capmeta.NewReadyCondition(latest)
if reconcileError != nil {
readyCondition.Status = metav1.ConditionFalse
readyCondition.Reason = capmeta.FailedReason
readyCondition.Message = reconcileError.Error()
}
latest.Status.Conditions.UpdateConditionByType(readyCondition)
if err := r.Client.Status().Update(ctx, latest); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return err
}
instance.Status = latest.Status
return nil
})
}
+232 -261
View File
@@ -1,11 +1,12 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
//nolint:nestif
package tls
import (
"bytes"
"context"
"crypto/rsa"
"crypto/x509"
"fmt"
"time"
@@ -27,7 +28,6 @@ import (
"sigs.k8s.io/controller-runtime/pkg/reconcile"
"github.com/projectcapsule/capsule/pkg/runtime/cert"
capsuleclient "github.com/projectcapsule/capsule/pkg/runtime/client"
"github.com/projectcapsule/capsule/pkg/runtime/configuration"
"github.com/projectcapsule/capsule/pkg/runtime/predicates"
)
@@ -35,11 +35,6 @@ import (
const (
certificateExpirationThreshold = 3 * 24 * time.Hour
certificateValidity = 6 * 30 * 24 * time.Hour
// caPrivateKeyKey is intentionally not a Kubernetes core constant.
// The TLS Secret remains type kubernetes.io/tls, but we persist the CA key
// so serving cert renewal does not require CA rotation.
caPrivateKeyKey = "ca.key"
)
type Reconciler struct {
@@ -56,7 +51,7 @@ func (r *Reconciler) SetupWithManager(mgr ctrl.Manager) error {
return []reconcile.Request{
{
NamespacedName: types.NamespacedName{
Namespace: r.Namespace,
Namespace: configuration.ControllerNamespace(),
Name: r.Configuration.TLSSecretName(),
},
},
@@ -104,7 +99,7 @@ func (r *Reconciler) SetupWithManager(mgr ctrl.Manager) error {
}
func (r *Reconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
r.Log = r.Log.WithValues(
log := r.Log.WithValues(
"Request.Namespace", request.Namespace,
"Request.Name", request.Name,
)
@@ -117,6 +112,8 @@ func (r *Reconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.
request.Name = r.Configuration.TLSSecretName()
}
log.V(4).Info("TLS reconciliation started")
certSecret := &corev1.Secret{}
if err := r.Get(ctx, request.NamespacedName, certSecret); err != nil {
if !apierrors.IsNotFound(err) {
@@ -129,7 +126,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.
certSecret.Data = map[string][]byte{}
}
if err := r.ReconcileCertificates(ctx, certSecret); err != nil {
if err := r.ReconcileCertificates(ctx, log, certSecret); err != nil {
return ctrl.Result{}, err
}
@@ -142,7 +139,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.
requeueAfter := max(time.Until(requeueTime), 0)
r.Log.V(4).Info("TLS reconciliation completed", "requeueAfter", requeueAfter.String())
log.V(4).Info("TLS reconciliation completed", "requeueAfter", requeueAfter.String())
return ctrl.Result{
Requeue: true,
@@ -150,27 +147,43 @@ func (r *Reconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.
}, nil
}
func (r *Reconciler) ReconcileCertificates(ctx context.Context, certSecret *corev1.Secret) error {
dnsName := r.webhookDNSName()
ca, caBundle, rotateServingCert, err := r.ensureCertificateMaterial(certSecret, dnsName)
func (r *Reconciler) ReconcileCertificates(
ctx context.Context,
log logr.Logger,
certSecret *corev1.Secret,
) error {
sans, err := r.desiredWebhookSANs(ctx)
if err != nil {
return err
}
log.V(4).Info(
"Resolved desired webhook certificate SANs",
"dnsNames", sans.DNSNames,
"ipAddresses", cert.IPsToStrings(sans.IPAddrs),
)
ca, caBundle, rotateServingCert, err := r.ensureCertificateMaterial(log, certSecret, sans)
if err != nil {
return err
}
log.V(4).Info(
"certificate requires rotation",
"rotation", rotateServingCert,
)
if rotateServingCert {
if ca == nil {
return fmt.Errorf("cannot rotate serving certificate without CA private key")
}
r.Log.V(3).Info("Generating new serving TLS certificate", "dnsName", dnsName)
crt, key, err := ca.GenerateCertificate(cert.NewCertOpts(
time.Now().Add(certificateValidity),
dnsName,
sans,
))
if err != nil {
r.Log.Error(err, "cannot generate serving TLS certificate")
log.Error(err, "cannot generate serving TLS certificate")
return err
}
@@ -178,7 +191,7 @@ func (r *Reconciler) ReconcileCertificates(ctx context.Context, certSecret *core
certSecret.Data[corev1.TLSCertKey] = crt.Bytes()
certSecret.Data[corev1.TLSPrivateKeyKey] = key.Bytes()
if err := r.validateSecretCertificate(certSecret, dnsName); err != nil {
if err := r.validateSecretCertificate(certSecret, sans); err != nil {
return err
}
@@ -192,18 +205,10 @@ func (r *Reconciler) ReconcileCertificates(ctx context.Context, certSecret *core
return fmt.Errorf("missing %q field in %q secret", corev1.ServiceAccountRootCAKey, r.Configuration.TLSSecretName())
}
r.Log.V(4).Info("Patching caBundle in webhooks and managed CRD conversions")
log.V(5).Info("Patching caBundle in managed CRD conversions")
patchGroup, groupCtx := errgroup.WithContext(ctx)
patchGroup.Go(func() error {
return r.patchMutatingWebhookConfigurationCABundle(groupCtx, caBundle)
})
patchGroup.Go(func() error {
return r.patchValidatingWebhookConfigurationCABundle(groupCtx, caBundle)
})
for key, managed := range r.conversionManagedCRDs() {
patchGroup.Go(func() error {
if err := r.updateManagedCustomResourceDefinition(groupCtx, managed, caBundle); err != nil {
@@ -227,112 +232,157 @@ func (r *Reconciler) ReconcileCertificates(ctx context.Context, certSecret *core
// - Serving certificate renewal never rotates the CA.
// - Legacy Secrets without ca.key rotate once into the stable format.
func (r *Reconciler) ensureCertificateMaterial(
log logr.Logger,
certSecret *corev1.Secret,
dnsName string,
sans cert.CertificateSANs,
) (*cert.CapsuleCA, []byte, bool, error) {
sans = sans.Normalize()
if sans.Empty() {
return nil, nil, false, fmt.Errorf("cannot ensure TLS material without SANs")
}
if certSecret.Data == nil {
certSecret.Data = map[string][]byte{}
}
caBundle := certSecret.Data[corev1.ServiceAccountRootCAKey]
caKey := certSecret.Data[caPrivateKeyKey]
tlsCrt := certSecret.Data[corev1.TLSCertKey]
tlsKey := certSecret.Data[corev1.TLSPrivateKeyKey]
caKey := certSecret.Data["ca.key"]
hasCA := len(caBundle) > 0
hasCABundle := len(caBundle) > 0
hasCAKey := len(caKey) > 0
hasServingCert := len(tlsCrt) > 0 && len(tlsKey) > 0
// Fresh empty Secret or completely broken Secret.
if !hasCA || !hasServingCert {
r.Log.Info(
"Generating new certificate authority and serving certificate",
"reason", "missing ca.crt or serving certificate",
"secret", certSecret.Name,
"namespace", certSecret.Namespace,
)
var ca *cert.CapsuleCA
ca, newCABundle, newCAKey, err := generateCertificateAuthorityMaterial()
rotateServingCert := false
switch {
case hasCABundle && hasCAKey:
loadedCA, err := cert.NewCertificateAuthorityFromBytes(caBundle, caKey)
if err != nil {
return nil, nil, false, err
}
certSecret.Data[corev1.ServiceAccountRootCAKey] = newCABundle
certSecret.Data[caPrivateKeyKey] = newCAKey
return ca, newCABundle, true, nil
}
// Legacy mode:
// The Secret has a CA and serving cert, but not the CA private key.
// If the serving cert is still valid and chains to ca.crt, do NOT rotate now.
// Rotating here causes a temporary caBundle/server-cert mismatch during startup.
if !hasCAKey {
if err := r.validateSecretCertificate(certSecret, dnsName); err == nil {
r.Log.Info(
"TLS Secret is using legacy CA material without ca.key; keeping existing CA and serving certificate",
"secret", certSecret.Name,
"namespace", certSecret.Namespace,
log.V(3).Info(
"Existing CA material is invalid, generating new CA",
"error", err.Error(),
)
return nil, caBundle, false, nil
generatedCA, generatedCABundle, generatedCAKey, err := generateCertificateAuthorityMaterial()
if err != nil {
return nil, nil, false, err
}
ca = generatedCA
caBundle = generatedCABundle
certSecret.Data[corev1.ServiceAccountRootCAKey] = generatedCABundle
certSecret.Data["ca.key"] = generatedCAKey
rotateServingCert = true
} else {
ca = loadedCA
}
r.Log.Info(
"TLS Secret is missing ca.key and serving certificate is invalid or expiring; rotating CA",
"secret", certSecret.Name,
"namespace", certSecret.Namespace,
case hasCABundle && !hasCAKey:
// Legacy mode: we can validate and patch caBundle, but we cannot issue
// a new serving certificate without the CA private key.
log.V(10).Info(
"TLS Secret contains CA bundle but no CA private key; running in legacy CA mode",
"secret", client.ObjectKeyFromObject(certSecret).String(),
)
ca, newCABundle, newCAKey, err := generateCertificateAuthorityMaterial()
if err := r.validateSecretCertificate(certSecret, sans); err != nil {
return nil, nil, false, fmt.Errorf(
"TLS Secret %s contains legacy CA material without ca.key and the serving certificate is invalid: %w",
client.ObjectKeyFromObject(certSecret).String(),
err,
)
}
return nil, caBundle, false, nil
default:
log.V(10).Info(
"TLS Secret is missing CA material, generating new CA",
"secret", client.ObjectKeyFromObject(certSecret).String(),
)
generatedCA, generatedCABundle, generatedCAKey, err := generateCertificateAuthorityMaterial()
if err != nil {
return nil, nil, false, err
}
certSecret.Data[corev1.ServiceAccountRootCAKey] = newCABundle
certSecret.Data[caPrivateKeyKey] = newCAKey
ca = generatedCA
caBundle = generatedCABundle
return ca, newCABundle, true, nil
certSecret.Data[corev1.ServiceAccountRootCAKey] = generatedCABundle
certSecret.Data["ca.key"] = generatedCAKey
rotateServingCert = true
}
ca, err := cert.NewCertificateAuthorityFromBytes(caBundle, caKey)
if err != nil {
r.Log.Error(err, "existing CA material is invalid, regenerating CA")
servingCertPEM := certSecret.Data[corev1.TLSCertKey]
servingKeyPEM := certSecret.Data[corev1.TLSPrivateKeyKey]
newCA, newCABundle, newCAKey, err := generateCertificateAuthorityMaterial()
if len(servingCertPEM) == 0 || len(servingKeyPEM) == 0 {
log.V(10).Info(
"TLS Secret is missing serving certificate material, rotating serving certificate",
"secret", client.ObjectKeyFromObject(certSecret).String(),
)
rotateServingCert = true
} else {
servingCert, err := cert.GetCertificateFromBytes(servingCertPEM)
if err != nil {
return nil, nil, false, err
log.V(10).Info(
"Failed to parse serving certificate, rotating serving certificate",
"secret", client.ObjectKeyFromObject(certSecret).String(),
"error", err.Error(),
)
rotateServingCert = true
} else {
if time.Until(servingCert.NotAfter) <= certificateExpirationThreshold {
log.V(10).Info(
"Serving certificate is close to expiry, rotating serving certificate",
"secret", client.ObjectKeyFromObject(certSecret).String(),
"notAfter", servingCert.NotAfter,
)
rotateServingCert = true
}
if !sans.MatchesCertificate(servingCert) {
log.V(3).Info(
"Serving certificate SANs differ from desired SANs, rotating serving certificate",
"secret", client.ObjectKeyFromObject(certSecret).String(),
"desiredDNSNames", sans.DNSNames,
"desiredIPAddresses", cert.IPsToStrings(sans.IPAddrs),
"currentDNSNames", servingCert.DNSNames,
"currentIPAddresses", cert.IPsToStrings(servingCert.IPAddresses),
)
rotateServingCert = true
}
if err := r.validateSecretCertificate(certSecret, sans); err != nil {
log.V(10).Info(
"Serving certificate failed validation, rotating serving certificate",
"secret", client.ObjectKeyFromObject(certSecret).String(),
"error", err.Error(),
)
rotateServingCert = true
}
}
certSecret.Data[corev1.ServiceAccountRootCAKey] = newCABundle
certSecret.Data[caPrivateKeyKey] = newCAKey
return newCA, newCABundle, true, nil
}
if err := validateCAKeyPair(caBundle, caKey); err != nil {
r.Log.Error(err, "existing CA certificate/key pair is invalid, regenerating CA")
newCA, newCABundle, newCAKey, err := generateCertificateAuthorityMaterial()
if err != nil {
return nil, nil, false, err
}
certSecret.Data[corev1.ServiceAccountRootCAKey] = newCABundle
certSecret.Data[caPrivateKeyKey] = newCAKey
return newCA, newCABundle, true, nil
if rotateServingCert && ca == nil {
return nil, nil, false, fmt.Errorf(
"cannot rotate serving certificate for TLS Secret %s without CA private key",
client.ObjectKeyFromObject(certSecret).String(),
)
}
if err := r.validateSecretCertificate(certSecret, dnsName); err != nil {
r.Log.Info("serving certificate requires renewal", "reason", err.Error())
return ca, caBundle, true, nil
}
r.Log.V(4).Info("Skipping TLS certificate generation as existing certificate is valid")
return ca, caBundle, false, nil
return ca, caBundle, rotateServingCert, nil
}
func generateCertificateAuthorityMaterial() (*cert.CapsuleCA, []byte, []byte, error) {
@@ -385,138 +435,91 @@ func (r *Reconciler) upsertTLSSecret(ctx context.Context, certSecret *corev1.Sec
return nil
}
func (r *Reconciler) validateSecretCertificate(secret *corev1.Secret, dnsName string) error {
if secret == nil {
return fmt.Errorf("secret is nil")
}
if secret.Data == nil {
return fmt.Errorf("secret data is nil")
}
caBundle := secret.Data[corev1.ServiceAccountRootCAKey]
if len(caBundle) == 0 {
return fmt.Errorf("missing %q", corev1.ServiceAccountRootCAKey)
}
leafPEM := secret.Data[corev1.TLSCertKey]
if len(leafPEM) == 0 {
return fmt.Errorf("missing %q", corev1.TLSCertKey)
}
keyPEM := secret.Data[corev1.TLSPrivateKeyKey]
if len(keyPEM) == 0 {
return fmt.Errorf("missing %q", corev1.TLSPrivateKeyKey)
}
leaf, key, err := cert.GetCertificateWithPrivateKeyFromBytes(leafPEM, keyPEM)
if err != nil {
return fmt.Errorf("cannot parse serving certificate/key pair: %w", err)
}
if err := cert.ValidateCertificate(leaf, key, certificateExpirationThreshold); err != nil {
return fmt.Errorf("serving certificate is invalid or expiring: %w", err)
func (r *Reconciler) validateSecretCertificate(
certSecret *corev1.Secret,
sans cert.CertificateSANs,
) error {
caPEM := certSecret.Data[corev1.ServiceAccountRootCAKey]
if len(caPEM) == 0 {
return fmt.Errorf("missing %q in TLS Secret %s/%s",
corev1.ServiceAccountRootCAKey,
certSecret.Namespace,
certSecret.Name,
)
}
roots := x509.NewCertPool()
if !roots.AppendCertsFromPEM(caBundle) {
return fmt.Errorf("cannot parse caBundle")
if ok := roots.AppendCertsFromPEM(caPEM); !ok {
return fmt.Errorf("failed to parse %q in TLS Secret %s/%s",
corev1.ServiceAccountRootCAKey,
certSecret.Namespace,
certSecret.Name,
)
}
if _, err := leaf.Verify(x509.VerifyOptions{
DNSName: dnsName,
Roots: roots,
KeyUsages: []x509.ExtKeyUsage{
x509.ExtKeyUsageServerAuth,
},
}); err != nil {
return fmt.Errorf("serving certificate does not verify against caBundle: %w", err)
}
return nil
}
func validateCAKeyPair(caCertPEM, caKeyPEM []byte) error {
caCert, caKey, err := cert.GetCertificateWithPrivateKeyFromBytes(caCertPEM, caKeyPEM)
leaf, err := cert.GetCertificateFromBytes(certSecret.Data[corev1.TLSCertKey])
if err != nil {
return fmt.Errorf("cannot parse CA certificate/key pair: %w", err)
return fmt.Errorf("parse serving certificate from TLS Secret %s/%s: %w",
certSecret.Namespace,
certSecret.Name,
err,
)
}
if !caCert.IsCA {
return fmt.Errorf("ca.crt is not a CA certificate")
keyPEM := certSecret.Data[corev1.TLSPrivateKeyKey]
if len(keyPEM) == 0 {
return fmt.Errorf("missing %q in TLS Secret %s/%s",
corev1.TLSPrivateKeyKey,
certSecret.Namespace,
certSecret.Name,
)
}
if !publicKeysEqual(caCert.PublicKey, &caKey.PublicKey) {
return fmt.Errorf("ca.crt does not match ca.key")
key, err := cert.GetPrivateKeyFromBytes(keyPEM)
if err != nil {
return fmt.Errorf("parse serving private key from TLS Secret %s/%s: %w",
certSecret.Namespace,
certSecret.Name,
err,
)
}
now := time.Now()
if now.Before(caCert.NotBefore) {
return fmt.Errorf("CA certificate is not valid yet")
if err := cert.ValidateCertificate(leaf, key, certificateExpirationThreshold); err != nil {
return fmt.Errorf("serving certificate/key pair is invalid or expiring: %w", err)
}
if now.After(caCert.NotAfter.Add(-certificateExpirationThreshold)) {
return fmt.Errorf("CA certificate expired or expires soon")
normalized := sans.Normalize()
if normalized.Empty() {
return fmt.Errorf("cannot validate serving certificate without desired SANs")
}
for _, dnsName := range normalized.DNSNames {
if _, err := leaf.Verify(x509.VerifyOptions{
DNSName: dnsName,
Roots: roots,
KeyUsages: []x509.ExtKeyUsage{
x509.ExtKeyUsageServerAuth,
},
}); err != nil {
return fmt.Errorf("serving certificate is not valid for DNS SAN %q: %w", dnsName, err)
}
}
for _, ip := range normalized.IPAddrs {
if _, err := leaf.Verify(x509.VerifyOptions{
DNSName: ip.String(),
Roots: roots,
KeyUsages: []x509.ExtKeyUsage{
x509.ExtKeyUsageServerAuth,
},
}); err != nil {
return fmt.Errorf("serving certificate is not valid for IP SAN %q: %w", ip.String(), err)
}
}
return nil
}
func publicKeysEqual(a any, b *rsa.PublicKey) bool {
pub, ok := a.(*rsa.PublicKey)
if !ok {
return false
}
return pub.Equal(b)
}
//nolint:dupl
func (r *Reconciler) patchValidatingWebhookConfigurationCABundle(ctx context.Context, caBundle []byte) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
vw := &admissionregistrationv1.ValidatingWebhookConfiguration{}
if err := r.Get(ctx, types.NamespacedName{
Name: string(r.Configuration.Admission().Validating.Name),
}, vw); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return err
}
patches := r.validatingWebhookCABundlePatches(vw.Webhooks, caBundle)
if len(patches) == 0 {
return nil
}
return capsuleclient.ApplyPatches(ctx, r.Client, vw, patches, "capsule-tls-controller")
})
}
//nolint:dupl
func (r *Reconciler) patchMutatingWebhookConfigurationCABundle(ctx context.Context, caBundle []byte) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() error {
mw := &admissionregistrationv1.MutatingWebhookConfiguration{}
if err := r.Get(ctx, types.NamespacedName{
Name: string(r.Configuration.Admission().Mutating.Name),
}, mw); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return err
}
patches := r.mutatingWebhookCABundlePatches(mw.Webhooks, caBundle)
if len(patches) == 0 {
return nil
}
return capsuleclient.ApplyPatches(ctx, r.Client, mw, patches, "capsule-tls-controller")
})
}
func (r *Reconciler) updateManagedCustomResourceDefinition(
ctx context.Context,
managed ManagedCRD,
@@ -536,44 +539,26 @@ func (r *Reconciler) updateManagedCustomResourceDefinition(
return err
}
// Only patch CRDs that already use webhook conversion.
if crd.Spec.Conversion == nil ||
crd.Spec.Conversion.Webhook == nil ||
crd.Spec.Conversion.Webhook.ClientConfig == nil {
return nil
}
current := crd.Spec.Conversion.Webhook.ClientConfig.CABundle
if bytes.Equal(current, caBundle) {
return nil
}
before := crd.DeepCopy()
path := managed.ConversionPath
if path == "" {
path = "/convert"
}
versions := managed.ConversionReviewVersions
if len(versions) == 0 {
versions = []string{"v1", "v1beta1"}
}
port := int32(443)
crd.Spec.Conversion = &apiextensionsv1.CustomResourceConversion{
Strategy: apiextensionsv1.WebhookConverter,
Webhook: &apiextensionsv1.WebhookConversion{
ClientConfig: &apiextensionsv1.WebhookClientConfig{
Service: &apiextensionsv1.ServiceReference{
Namespace: r.Namespace,
Name: r.Configuration.Admission().ServiceName,
Path: &path,
Port: &port,
},
CABundle: caBundle,
},
ConversionReviewVersions: versions,
},
}
crd.Spec.Conversion.Webhook.ClientConfig.CABundle = append([]byte(nil), caBundle...)
return r.Patch(ctx, crd, client.MergeFrom(before))
})
}
func (r *Reconciler) webhookDNSName() string {
return fmt.Sprintf("%s.%s.svc", r.Configuration.Admission().ServiceName, r.Namespace)
}
func copySecretData(in map[string][]byte) map[string][]byte {
out := make(map[string][]byte, len(in))
@@ -583,17 +568,3 @@ func copySecretData(in map[string][]byte) map[string][]byte {
return out
}
func equalBytes(a, b []byte) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
+81 -63
View File
@@ -4,13 +4,15 @@
package tls
import (
"crypto/sha256"
"encoding/hex"
"context"
"fmt"
admissionregistrationv1 "k8s.io/api/admissionregistration/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
capsuleclient "github.com/projectcapsule/capsule/pkg/runtime/client"
"github.com/projectcapsule/capsule/pkg/runtime/cert"
"github.com/projectcapsule/capsule/pkg/runtime/configuration"
)
type ManagedCRD struct {
@@ -90,74 +92,90 @@ func (r Reconciler) conversionManagedCRDs() map[string]ManagedCRD {
return out
}
//nolint:dupl
func (r *Reconciler) validatingWebhookCABundlePatches(
webhooks []admissionregistrationv1.ValidatingWebhook,
caBundle []byte,
) []capsuleclient.JSONPatch {
patches := make([]capsuleclient.JSONPatch, 0, len(webhooks))
// Collects required SANs for certificate
// desiredWebhookSANs collects required SANs for the webhook serving certificate.
func (r *Reconciler) desiredWebhookSANs(ctx context.Context) (cert.CertificateSANs, error) {
sans := cert.CertificateSANs{}
for i := range webhooks {
if webhooks[i].ClientConfig.Service == nil {
continue
sans.AddDNSNames(
r.Configuration.Admission().ServiceName,
fmt.Sprintf("%s.%s", r.Configuration.Admission().ServiceName, r.Namespace),
fmt.Sprintf("%s.%s.svc", r.Configuration.Admission().ServiceName, r.Namespace),
fmt.Sprintf("%s.%s.svc.cluster.local", r.Configuration.Admission().ServiceName, r.Namespace),
)
mutating := r.Configuration.Admission().Mutating.Client
if mutating != nil {
sans.AddServiceReference(mutating.Service, configuration.ControllerNamespace())
if err := sans.AddURL(mutating.URL); err != nil {
return cert.CertificateSANs{}, fmt.Errorf("mutating admission client URL: %w", err)
}
if equalBytes(webhooks[i].ClientConfig.CABundle, caBundle) {
continue
}
r.Log.V(3).Info(
"Patching webhook caBundle",
"webhook", webhooks[i].Name,
"old", certFingerprint(webhooks[i].ClientConfig.CABundle),
"new", certFingerprint(caBundle),
)
patches = append(patches, capsuleclient.JSONPatch{
Operation: capsuleclient.JSONPatchAdd,
Path: fmt.Sprintf("/webhooks/%d/clientConfig/caBundle", i),
Value: caBundle,
})
}
return patches
}
validating := r.Configuration.Admission().Validating.Client
if validating != nil {
sans.AddServiceReference(validating.Service, configuration.ControllerNamespace())
//nolint:dupl
func (r *Reconciler) mutatingWebhookCABundlePatches(
webhooks []admissionregistrationv1.MutatingWebhook,
caBundle []byte,
) []capsuleclient.JSONPatch {
patches := make([]capsuleclient.JSONPatch, 0, len(webhooks))
for i := range webhooks {
if webhooks[i].ClientConfig.Service == nil {
continue
if err := sans.AddURL(validating.URL); err != nil {
return cert.CertificateSANs{}, fmt.Errorf("validating admission client URL: %w", err)
}
if equalBytes(webhooks[i].ClientConfig.CABundle, caBundle) {
continue
}
r.Log.V(3).Info(
"Patching webhook caBundle",
"webhook", webhooks[i].Name,
"old", certFingerprint(webhooks[i].ClientConfig.CABundle),
"new", certFingerprint(caBundle),
)
patches = append(patches, capsuleclient.JSONPatch{
Operation: capsuleclient.JSONPatchAdd,
Path: fmt.Sprintf("/webhooks/%d/clientConfig/caBundle", i),
Value: caBundle,
})
}
return patches
sans = sans.Normalize()
if sans.Empty() {
return cert.CertificateSANs{}, fmt.Errorf("no webhook SANs could be resolved")
}
r.Log.V(5).Info(
"Evaluated required SANs for TLS controller",
"dnsNames", sans.DNSNames,
"ipAddresses", cert.IPsToStrings(sans.IPAddrs),
)
return sans, nil
}
func certFingerprint(pemBytes []byte) string {
sum := sha256.Sum256(pemBytes)
func FetchCurrentCaBundleForAdmission(
ctx context.Context,
c client.Reader,
cfg configuration.Configuration,
configuredCABundle []byte,
) ([]byte, error) {
// Explicit configuration wins.
if len(configuredCABundle) > 0 {
return append([]byte(nil), configuredCABundle...), nil
}
return hex.EncodeToString(sum[:8])
// Internal Capsule TLS enabled: source of truth is the TLS Secret.
if cfg.EnableTLSConfiguration() {
secret := &corev1.Secret{}
if err := c.Get(ctx, types.NamespacedName{
Namespace: configuration.ControllerNamespace(),
Name: cfg.TLSSecretName(),
}, secret); err != nil {
return nil, fmt.Errorf("get TLS Secret %s/%s: %w",
configuration.ControllerNamespace(),
cfg.TLSSecretName(),
err,
)
}
caBundle := secret.Data[corev1.ServiceAccountRootCAKey]
if len(caBundle) == 0 {
return nil, fmt.Errorf("TLS Secret %s/%s missing %q",
secret.Namespace,
secret.Name,
corev1.ServiceAccountRootCAKey,
)
}
return append([]byte(nil), caBundle...), nil
}
// cert-manager / external injector mode:
// return nil and preserve current webhook caBundle.
return nil, nil
}
+38 -11
View File
@@ -10,6 +10,7 @@ import (
"crypto/x509"
"crypto/x509/pkix"
"encoding/pem"
"fmt"
"math/big"
"time"
@@ -141,20 +142,40 @@ func GenerateCertificateAuthority() (s *CapsuleCA, err error) {
return s, err
}
func GetCertificateFromBytes(certBytes []byte) (*x509.Certificate, error) {
var b *pem.Block
func GetCertificateFromBytes(raw []byte) (*x509.Certificate, error) {
block, _ := pem.Decode(raw)
if block == nil {
return nil, fmt.Errorf("failed to decode certificate PEM")
}
b, _ = pem.Decode(certBytes)
if block.Type != "CERTIFICATE" {
return nil, fmt.Errorf("expected CERTIFICATE PEM block, got %q", block.Type)
}
return x509.ParseCertificate(b.Bytes)
certificate, err := x509.ParseCertificate(block.Bytes)
if err != nil {
return nil, fmt.Errorf("parse certificate: %w", err)
}
return certificate, nil
}
func GetPrivateKeyFromBytes(keyBytes []byte) (*rsa.PrivateKey, error) {
var b *pem.Block
func GetPrivateKeyFromBytes(raw []byte) (*rsa.PrivateKey, error) {
block, _ := pem.Decode(raw)
if block == nil {
return nil, fmt.Errorf("failed to decode private key PEM")
}
b, _ = pem.Decode(keyBytes)
if block.Type != "RSA PRIVATE KEY" {
return nil, fmt.Errorf("expected RSA PRIVATE KEY PEM block, got %q", block.Type)
}
return x509.ParsePKCS1PrivateKey(b.Bytes)
privateKey, err := x509.ParsePKCS1PrivateKey(block.Bytes)
if err != nil {
return nil, fmt.Errorf("parse RSA private key: %w", err)
}
return privateKey, nil
}
func GetCertificateWithPrivateKeyFromBytes(certBytes, keyBytes []byte) (*x509.Certificate, *rsa.PrivateKey, error) {
@@ -171,7 +192,7 @@ func GetCertificateWithPrivateKeyFromBytes(certBytes, keyBytes []byte) (*x509.Ce
return cert, key, nil
}
func (c *CapsuleCA) GenerateCertificate(opts CertificateOptions) (certificatePem *bytes.Buffer, certificateKey *bytes.Buffer, err error) {
func (c *CapsuleCA) GenerateCertificate(opts CertOpts) (certificatePem *bytes.Buffer, certificateKey *bytes.Buffer, err error) {
var certPrivKey *rsa.PrivateKey
certPrivKey, err = rsa.GenerateKey(rand.Reader, 4096)
@@ -179,6 +200,11 @@ func (c *CapsuleCA) GenerateCertificate(opts CertificateOptions) (certificatePem
return nil, nil, err
}
sans := opts.SAN
if sans.Empty() {
return nil, nil, fmt.Errorf("cannot generate certificate without SANs")
}
cert := &x509.Certificate{
SerialNumber: big.NewInt(1658),
Subject: pkix.Name{
@@ -189,9 +215,10 @@ func (c *CapsuleCA) GenerateCertificate(opts CertificateOptions) (certificatePem
StreetAddress: []string{"27, Old Gloucester Street"},
PostalCode: []string{"WC1N 3AX"},
},
DNSNames: opts.DNSNames(),
DNSNames: sans.DNSNames,
IPAddresses: sans.IPAddrs,
NotBefore: time.Now().AddDate(0, 0, -1),
NotAfter: opts.ExpirationDate(),
NotAfter: opts.GetExpirationDate(),
SubjectKeyId: []byte{1, 2, 3, 4, 6},
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageClientAuth, x509.ExtKeyUsageServerAuth},
KeyUsage: x509.KeyUsageDigitalSignature,
+267 -50
View File
@@ -8,71 +8,288 @@ import (
"crypto/tls"
"crypto/x509"
"encoding/pem"
"net"
"reflect"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/projectcapsule/capsule/pkg/runtime/cert"
)
func TestNewCertificateAuthorityFromBytes(t *testing.T) {
var ca *cert.CapsuleCA
t.Parallel()
var err error
ca, err = cert.GenerateCertificateAuthority()
assert.Nil(t, err)
var crt *bytes.Buffer
crt, err = ca.CACertificatePem()
assert.Nil(t, err)
var key *bytes.Buffer
key, err = ca.CAPrivateKeyPem()
assert.Nil(t, err)
_, err = cert.NewCertificateAuthorityFromBytes(crt.Bytes(), key.Bytes())
assert.Nil(t, err)
}
func TestCapsuleCa_GenerateCertificate(t *testing.T) {
type testCase struct {
dnsNames []string
ca, err := cert.GenerateCertificateAuthority()
if err != nil {
t.Fatalf("expected CA generation to succeed, got %v", err)
}
for name, c := range map[string]testCase{
"foo.tld": {[]string{"foo.tld"}},
"SAN": {[]string{"capsule-webhook-service.capsule-system.svc", "capsule-webhook-service.capsule-system.default.cluster"}},
} {
crt, err := ca.CACertificatePem()
if err != nil {
t.Fatalf("expected CA certificate PEM encoding to succeed, got %v", err)
}
key, err := ca.CAPrivateKeyPem()
if err != nil {
t.Fatalf("expected CA private key PEM encoding to succeed, got %v", err)
}
loadedCA, err := cert.NewCertificateAuthorityFromBytes(crt.Bytes(), key.Bytes())
if err != nil {
t.Fatalf("expected loading CA from PEM bytes to succeed, got %v", err)
}
loadedCrt, err := loadedCA.CACertificatePem()
if err != nil {
t.Fatalf("expected loaded CA certificate PEM encoding to succeed, got %v", err)
}
if !bytes.Equal(crt.Bytes(), loadedCrt.Bytes()) {
t.Fatal("expected loaded CA certificate PEM to match original CA certificate PEM")
}
}
func TestNewCertificateAuthorityFromBytesRejectsInvalidCertificatePEM(t *testing.T) {
t.Parallel()
ca, err := cert.GenerateCertificateAuthority()
if err != nil {
t.Fatalf("expected CA generation to succeed, got %v", err)
}
key, err := ca.CAPrivateKeyPem()
if err != nil {
t.Fatalf("expected CA private key PEM encoding to succeed, got %v", err)
}
if _, err := cert.NewCertificateAuthorityFromBytes([]byte("invalid certificate"), key.Bytes()); err == nil {
t.Fatal("expected invalid certificate PEM to be rejected")
}
}
func TestNewCertificateAuthorityFromBytesRejectsInvalidPrivateKeyPEM(t *testing.T) {
t.Parallel()
ca, err := cert.GenerateCertificateAuthority()
if err != nil {
t.Fatalf("expected CA generation to succeed, got %v", err)
}
crt, err := ca.CACertificatePem()
if err != nil {
t.Fatalf("expected CA certificate PEM encoding to succeed, got %v", err)
}
if _, err := cert.NewCertificateAuthorityFromBytes(crt.Bytes(), []byte("invalid private key")); err == nil {
t.Fatal("expected invalid private key PEM to be rejected")
}
}
func TestCapsuleCAGenerateCertificate(t *testing.T) {
t.Parallel()
tests := map[string]struct {
sans cert.CertificateSANs
}{
"dns name": {
sans: cert.CertificateSANs{
DNSNames: []string{"foo.tld"},
},
},
"multiple dns names": {
sans: cert.CertificateSANs{
DNSNames: []string{
"capsule-webhook-service.capsule-system.svc",
"capsule-webhook-service.capsule-system.svc.cluster.local",
},
},
},
"dns names are normalized": {
sans: cert.CertificateSANs{
DNSNames: []string{
" Capsule-Webhook-Service.Capsule-System.Svc ",
"capsule-webhook-service.capsule-system.svc",
},
},
},
"ip addresses": {
sans: cert.CertificateSANs{
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
net.ParseIP("fd00::1"),
},
},
},
"dns names and ip addresses": {
sans: cert.CertificateSANs{
DNSNames: []string{
"capsule-webhook-service.capsule-system.svc",
},
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
},
},
},
}
for name, tt := range tests {
tt := tt
t.Run(name, func(t *testing.T) {
var ca *cert.CapsuleCA
var err error
t.Parallel()
e := time.Now().AddDate(1, 0, 0)
expirationDate := time.Now().AddDate(1, 0, 0).UTC()
ca, err = cert.GenerateCertificateAuthority()
assert.Nil(t, err)
var crt *bytes.Buffer
var key *bytes.Buffer
crt, key, err = ca.GenerateCertificate(cert.NewCertOpts(e, c.dnsNames...))
assert.Nil(t, err)
var b *pem.Block
var c *x509.Certificate
b, _ = pem.Decode(crt.Bytes())
c, err = x509.ParseCertificate(b.Bytes)
assert.Nil(t, err)
assert.Equal(t, e.Unix(), c.NotAfter.Unix())
for _, i := range c.DNSNames {
assert.Contains(t, c.DNSNames, i)
ca, err := cert.GenerateCertificateAuthority()
if err != nil {
t.Fatalf("expected CA generation to succeed, got %v", err)
}
_, err = tls.X509KeyPair(crt.Bytes(), key.Bytes())
assert.Nil(t, err)
crt, key, err := ca.GenerateCertificate(cert.NewCertOpts(expirationDate, tt.sans))
if err != nil {
t.Fatalf("expected serving certificate generation to succeed, got %v", err)
}
servingCert := parseCertificatePEM(t, crt.Bytes())
if servingCert.NotAfter.Unix() != expirationDate.Unix() {
t.Fatalf("expected certificate NotAfter %d, got %d", expirationDate.Unix(), servingCert.NotAfter.Unix())
}
expectedSANs := tt.sans.Normalize()
actualSANs := cert.CertificateSANs{
DNSNames: servingCert.DNSNames,
IPAddrs: servingCert.IPAddresses,
}.Normalize()
assertCertificateSANsEqual(t, expectedSANs, actualSANs)
if !tt.sans.MatchesCertificate(servingCert) {
t.Fatalf(
"expected generated certificate to match requested SANs, desiredDNSNames=%v desiredIPAddresses=%v actualDNSNames=%v actualIPAddresses=%v",
expectedSANs.DNSNames,
cert.IPsToStrings(expectedSANs.IPAddrs),
servingCert.DNSNames,
cert.IPsToStrings(servingCert.IPAddresses),
)
}
if _, err := tls.X509KeyPair(crt.Bytes(), key.Bytes()); err != nil {
t.Fatalf("expected generated certificate/key pair to be valid, got %v", err)
}
})
}
}
func TestCapsuleCAGenerateCertificateIsSignedByCA(t *testing.T) {
t.Parallel()
ca, err := cert.GenerateCertificateAuthority()
if err != nil {
t.Fatalf("expected CA generation to succeed, got %v", err)
}
expirationDate := time.Now().AddDate(1, 0, 0).UTC()
crt, _, err := ca.GenerateCertificate(cert.NewCertOpts(expirationDate, cert.CertificateSANs{
DNSNames: []string{"capsule-webhook-service.capsule-system.svc"},
}))
if err != nil {
t.Fatalf("expected serving certificate generation to succeed, got %v", err)
}
servingCert := parseCertificatePEM(t, crt.Bytes())
caCrt, err := ca.CACertificatePem()
if err != nil {
t.Fatalf("expected CA certificate PEM encoding to succeed, got %v", err)
}
caCert := parseCertificatePEM(t, caCrt.Bytes())
roots := x509.NewCertPool()
roots.AddCert(caCert)
verifyOptions := x509.VerifyOptions{
DNSName: "capsule-webhook-service.capsule-system.svc",
Roots: roots,
}
if _, err := servingCert.Verify(verifyOptions); err != nil {
t.Fatalf("expected serving certificate to verify against generated CA, got %v", err)
}
}
func TestCapsuleCAGenerateCertificateWithIPAddressSANVerifiesForIP(t *testing.T) {
t.Parallel()
ca, err := cert.GenerateCertificateAuthority()
if err != nil {
t.Fatalf("expected CA generation to succeed, got %v", err)
}
ip := net.ParseIP("10.96.0.10")
expirationDate := time.Now().AddDate(1, 0, 0).UTC()
crt, _, err := ca.GenerateCertificate(cert.NewCertOpts(expirationDate, cert.CertificateSANs{
IPAddrs: []net.IP{ip},
}))
if err != nil {
t.Fatalf("expected serving certificate generation to succeed, got %v", err)
}
servingCert := parseCertificatePEM(t, crt.Bytes())
caCrt, err := ca.CACertificatePem()
if err != nil {
t.Fatalf("expected CA certificate PEM encoding to succeed, got %v", err)
}
caCert := parseCertificatePEM(t, caCrt.Bytes())
roots := x509.NewCertPool()
roots.AddCert(caCert)
verifyOptions := x509.VerifyOptions{
DNSName: ip.String(),
Roots: roots,
}
if _, err := servingCert.Verify(verifyOptions); err != nil {
t.Fatalf("expected serving certificate to verify against generated CA for IP SAN, got %v", err)
}
}
func parseCertificatePEM(t *testing.T, certificatePEM []byte) *x509.Certificate {
t.Helper()
block, _ := pem.Decode(certificatePEM)
if block == nil {
t.Fatal("expected certificate PEM block, got nil")
}
certificate, err := x509.ParseCertificate(block.Bytes)
if err != nil {
t.Fatalf("expected certificate parsing to succeed, got %v", err)
}
return certificate
}
func assertCertificateSANsEqual(t *testing.T, expected, actual cert.CertificateSANs) {
t.Helper()
expected = expected.Normalize()
actual = actual.Normalize()
if !reflect.DeepEqual(expected.DNSNames, actual.DNSNames) {
t.Fatalf("expected DNS names %v, got %v", expected.DNSNames, actual.DNSNames)
}
expectedIPs := cert.IPsToStrings(expected.IPAddrs)
actualIPs := cert.IPsToStrings(actual.IPAddrs)
if !reflect.DeepEqual(expectedIPs, actualIPs) {
t.Fatalf("expected IP addresses %v, got %v", expectedIPs, actualIPs)
}
}
+14 -11
View File
@@ -6,23 +6,26 @@ package cert
import "time"
type CertificateOptions interface {
DNSNames() []string
ExpirationDate() time.Time
GetDNSNames() CertificateSANs
GetExpirationDate() time.Time
}
type certOpts struct {
dnsNames []string
expirationDate time.Time
type CertOpts struct {
SAN CertificateSANs
ExpirationDate time.Time
}
func (c certOpts) DNSNames() []string {
return c.dnsNames
func NewCertOpts(expirationDate time.Time, sans CertificateSANs) CertOpts {
return CertOpts{
ExpirationDate: expirationDate,
SAN: sans.Normalize(),
}
}
func (c certOpts) ExpirationDate() time.Time {
return c.expirationDate
func (c CertOpts) GetDNSNames() CertificateSANs {
return c.SAN
}
func NewCertOpts(expirationDate time.Time, dnsNames ...string) CertificateOptions {
return &certOpts{dnsNames: dnsNames, expirationDate: expirationDate}
func (c CertOpts) GetExpirationDate() time.Time {
return c.ExpirationDate
}
+51
View File
@@ -0,0 +1,51 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package cert
import (
"net"
"testing"
"time"
)
func TestNewCertOpts(t *testing.T) {
t.Parallel()
expirationDate := time.Now().Add(24 * time.Hour).UTC()
opts := NewCertOpts(expirationDate, CertificateSANs{
DNSNames: []string{
" Webhook.Capsule-System.Svc ",
"webhook.capsule-system.svc",
"api.example.com",
},
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
net.ParseIP("10.96.0.10"),
nil,
},
})
if !opts.GetExpirationDate().Equal(expirationDate) {
t.Fatalf("expected expiration date %s, got %s", expirationDate, opts.GetExpirationDate())
}
expectedSANs := CertificateSANs{
DNSNames: []string{
"api.example.com",
"webhook.capsule-system.svc",
},
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
},
}
assertCertificateSANsEqual(t, expectedSANs, opts.GetDNSNames())
}
func TestCertOptsImplementsCertificateOptions(t *testing.T) {
t.Parallel()
var _ CertificateOptions = CertOpts{}
}
+161
View File
@@ -0,0 +1,161 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package cert
import (
"crypto/x509"
"fmt"
"net"
"net/url"
"sort"
"strings"
admissionregistrationv1 "k8s.io/api/admissionregistration/v1"
)
type CertificateSANs struct {
DNSNames []string
IPAddrs []net.IP
}
func (s CertificateSANs) Empty() bool {
return len(s.DNSNames) == 0 && len(s.IPAddrs) == 0
}
func (s CertificateSANs) Normalize() CertificateSANs {
dnsSet := make(map[string]struct{}, len(s.DNSNames))
dnsNames := make([]string, 0, len(s.DNSNames))
for _, name := range s.DNSNames {
name = strings.TrimSpace(strings.ToLower(name))
if name == "" {
continue
}
if _, ok := dnsSet[name]; ok {
continue
}
dnsSet[name] = struct{}{}
dnsNames = append(dnsNames, name)
}
ipSet := make(map[string]struct{}, len(s.IPAddrs))
ipAddrs := make([]net.IP, 0, len(s.IPAddrs))
for _, ip := range s.IPAddrs {
if ip == nil {
continue
}
normalized := ip.String()
if _, ok := ipSet[normalized]; ok {
continue
}
ipSet[normalized] = struct{}{}
ipAddrs = append(ipAddrs, ip)
}
sort.Strings(dnsNames)
sort.Slice(ipAddrs, func(i, j int) bool {
return ipAddrs[i].String() < ipAddrs[j].String()
})
return CertificateSANs{
DNSNames: dnsNames,
IPAddrs: ipAddrs,
}
}
func (s *CertificateSANs) AddDNSNames(names ...string) {
s.DNSNames = append(s.DNSNames, names...)
}
func (s *CertificateSANs) AddIPAddrs(ips ...net.IP) {
s.IPAddrs = append(s.IPAddrs, ips...)
}
func (s *CertificateSANs) AddServiceReference(
service *admissionregistrationv1.ServiceReference,
defaultNamespace string,
) {
if service == nil || service.Name == "" {
return
}
namespace := service.Namespace
if namespace == "" {
namespace = defaultNamespace
}
s.AddDNSNames(
service.Name,
fmt.Sprintf("%s.%s", service.Name, namespace),
fmt.Sprintf("%s.%s.svc", service.Name, namespace),
fmt.Sprintf("%s.%s.svc.cluster.local", service.Name, namespace),
)
}
func (s *CertificateSANs) AddURL(rawURL *string) error {
if rawURL == nil || strings.TrimSpace(*rawURL) == "" {
return nil
}
parsed, err := url.Parse(*rawURL)
if err != nil {
return fmt.Errorf("parse webhook URL %q: %w", *rawURL, err)
}
host := parsed.Hostname()
if host == "" {
return fmt.Errorf("webhook URL %q has empty host", *rawURL)
}
if ip := net.ParseIP(host); ip != nil {
s.AddIPAddrs(ip)
return nil
}
s.AddDNSNames(host)
return nil
}
func (s CertificateSANs) MatchesCertificate(certificate *x509.Certificate) bool {
if certificate == nil {
return false
}
desired := s.Normalize()
actual := CertificateSANs{
DNSNames: append([]string(nil), certificate.DNSNames...),
IPAddrs: append([]net.IP(nil), certificate.IPAddresses...),
}.Normalize()
if len(desired.DNSNames) != len(actual.DNSNames) {
return false
}
for i := range desired.DNSNames {
if desired.DNSNames[i] != actual.DNSNames[i] {
return false
}
}
if len(desired.IPAddrs) != len(actual.IPAddrs) {
return false
}
for i := range desired.IPAddrs {
if desired.IPAddrs[i].String() != actual.IPAddrs[i].String() {
return false
}
}
return true
}
+399
View File
@@ -0,0 +1,399 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package cert
import (
"crypto/x509"
"net"
"reflect"
"testing"
admissionregistrationv1 "k8s.io/api/admissionregistration/v1"
"k8s.io/utils/ptr"
)
func TestCertificateSANsEmpty(t *testing.T) {
t.Parallel()
tests := []struct {
name string
sans CertificateSANs
want bool
}{
{
name: "empty",
sans: CertificateSANs{},
want: true,
},
{
name: "dns names present",
sans: CertificateSANs{
DNSNames: []string{"webhook.capsule-system.svc"},
},
want: false,
},
{
name: "ip addresses present",
sans: CertificateSANs{
IPAddrs: []net.IP{net.ParseIP("10.96.0.10")},
},
want: false,
},
}
for _, tt := range tests {
tt := tt
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
if got := tt.sans.Empty(); got != tt.want {
t.Fatalf("expected Empty() to be %t, got %t", tt.want, got)
}
})
}
}
func TestCertificateSANsNormalize(t *testing.T) {
t.Parallel()
sans := CertificateSANs{
DNSNames: []string{
" Webhook.Capsule-System.Svc ",
"",
"webhook.capsule-system.svc",
"API.Example.Com",
" api.example.com ",
},
IPAddrs: []net.IP{
net.ParseIP("10.96.0.20"),
nil,
net.ParseIP("10.96.0.10"),
net.ParseIP("10.96.0.20"),
},
}
expected := CertificateSANs{
DNSNames: []string{
"api.example.com",
"webhook.capsule-system.svc",
},
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
net.ParseIP("10.96.0.20"),
},
}
assertCertificateSANsEqual(t, expected, sans.Normalize())
}
func TestCertificateSANsAddDNSNames(t *testing.T) {
t.Parallel()
var sans CertificateSANs
sans.AddDNSNames("webhook", "webhook.capsule-system.svc")
expected := []string{
"webhook",
"webhook.capsule-system.svc",
}
if !reflect.DeepEqual(expected, sans.DNSNames) {
t.Fatalf("expected DNS names %v, got %v", expected, sans.DNSNames)
}
}
func TestCertificateSANsAddIPAddrs(t *testing.T) {
t.Parallel()
var sans CertificateSANs
sans.AddIPAddrs(net.ParseIP("10.96.0.10"), net.ParseIP("fd00::1"))
expected := []string{
"10.96.0.10",
"fd00::1",
}
if got := IPsToStrings(sans.IPAddrs); !reflect.DeepEqual(expected, got) {
t.Fatalf("expected IP addresses %v, got %v", expected, got)
}
}
func TestCertificateSANsAddServiceReference(t *testing.T) {
t.Parallel()
tests := []struct {
name string
service *admissionregistrationv1.ServiceReference
defaultNamespace string
expected []string
}{
{
name: "nil service is ignored",
},
{
name: "empty service name is ignored",
service: &admissionregistrationv1.ServiceReference{
Namespace: "capsule-system",
},
},
{
name: "uses service namespace",
service: &admissionregistrationv1.ServiceReference{
Name: "capsule-webhook-service",
Namespace: "capsule-system",
},
defaultNamespace: "default",
expected: []string{
"capsule-webhook-service",
"capsule-webhook-service.capsule-system",
"capsule-webhook-service.capsule-system.svc",
"capsule-webhook-service.capsule-system.svc.cluster.local",
},
},
{
name: "uses default namespace when service namespace is empty",
service: &admissionregistrationv1.ServiceReference{
Name: "capsule-webhook-service",
},
defaultNamespace: "capsule-system",
expected: []string{
"capsule-webhook-service",
"capsule-webhook-service.capsule-system",
"capsule-webhook-service.capsule-system.svc",
"capsule-webhook-service.capsule-system.svc.cluster.local",
},
},
}
for _, tt := range tests {
tt := tt
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
var sans CertificateSANs
sans.AddServiceReference(tt.service, tt.defaultNamespace)
if !reflect.DeepEqual(tt.expected, sans.DNSNames) {
t.Fatalf("expected DNS names %v, got %v", tt.expected, sans.DNSNames)
}
})
}
}
func TestCertificateSANsAddURL(t *testing.T) {
t.Parallel()
tests := []struct {
name string
rawURL *string
expectedDNSNames []string
expectedIPAddrs []net.IP
wantErr bool
}{
{
name: "nil URL is ignored",
},
{
name: "blank URL is ignored",
rawURL: ptr.To(" "),
},
{
name: "adds dns hostname",
rawURL: ptr.To("https://webhook.example.com/mutate"),
expectedDNSNames: []string{"webhook.example.com"},
},
{
name: "adds dns hostname without port",
rawURL: ptr.To("https://webhook.example.com:9443/mutate"),
expectedDNSNames: []string{"webhook.example.com"},
},
{
name: "adds ipv4 address",
rawURL: ptr.To("https://10.96.0.10:9443/mutate"),
expectedIPAddrs: []net.IP{net.ParseIP("10.96.0.10")},
},
{
name: "adds ipv6 address",
rawURL: ptr.To("https://[fd00::1]:9443/mutate"),
expectedIPAddrs: []net.IP{net.ParseIP("fd00::1")},
},
{
name: "rejects invalid URL",
rawURL: ptr.To("://bad-url"),
wantErr: true,
},
{
name: "rejects URL without host",
rawURL: ptr.To("https:///mutate"),
wantErr: true,
},
}
for _, tt := range tests {
tt := tt
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
var sans CertificateSANs
err := sans.AddURL(tt.rawURL)
if tt.wantErr {
if err == nil {
t.Fatal("expected error, got nil")
}
return
}
if err != nil {
t.Fatalf("expected no error, got %v", err)
}
expected := CertificateSANs{
DNSNames: tt.expectedDNSNames,
IPAddrs: tt.expectedIPAddrs,
}
assertCertificateSANsEqual(t, expected, sans)
})
}
}
func TestCertificateSANsMatchesCertificate(t *testing.T) {
t.Parallel()
tests := []struct {
name string
sans CertificateSANs
certificate *x509.Certificate
want bool
}{
{
name: "nil certificate does not match",
sans: CertificateSANs{
DNSNames: []string{"webhook.capsule-system.svc"},
},
want: false,
},
{
name: "matches normalized dns and ip sans",
sans: CertificateSANs{
DNSNames: []string{
" Webhook.Capsule-System.Svc ",
"api.example.com",
},
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
},
},
certificate: &x509.Certificate{
DNSNames: []string{
"api.example.com",
"webhook.capsule-system.svc",
},
IPAddresses: []net.IP{
net.ParseIP("10.96.0.10"),
},
},
want: true,
},
{
name: "does not match when certificate misses dns san",
sans: CertificateSANs{
DNSNames: []string{
"api.example.com",
"webhook.capsule-system.svc",
},
},
certificate: &x509.Certificate{
DNSNames: []string{
"webhook.capsule-system.svc",
},
},
want: false,
},
{
name: "does not match when certificate has stale extra dns san",
sans: CertificateSANs{
DNSNames: []string{
"webhook.capsule-system.svc",
},
},
certificate: &x509.Certificate{
DNSNames: []string{
"old-webhook.capsule-system.svc",
"webhook.capsule-system.svc",
},
},
want: false,
},
{
name: "does not match when certificate misses ip san",
sans: CertificateSANs{
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
},
},
certificate: &x509.Certificate{},
want: false,
},
{
name: "does not match when certificate has stale extra ip san",
sans: CertificateSANs{
IPAddrs: []net.IP{
net.ParseIP("10.96.0.10"),
},
},
certificate: &x509.Certificate{
IPAddresses: []net.IP{
net.ParseIP("10.96.0.10"),
net.ParseIP("10.96.0.20"),
},
},
want: false,
},
{
name: "matches empty sans against certificate without sans",
sans: CertificateSANs{},
certificate: &x509.Certificate{},
want: true,
},
}
for _, tt := range tests {
tt := tt
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
if got := tt.sans.MatchesCertificate(tt.certificate); got != tt.want {
t.Fatalf("expected MatchesCertificate() to be %t, got %t", tt.want, got)
}
})
}
}
func assertCertificateSANsEqual(t *testing.T, expected, actual CertificateSANs) {
t.Helper()
expected = expected.Normalize()
actual = actual.Normalize()
if !reflect.DeepEqual(expected.DNSNames, actual.DNSNames) {
t.Fatalf("expected DNS names %v, got %v", expected.DNSNames, actual.DNSNames)
}
expectedIPs := IPsToStrings(expected.IPAddrs)
actualIPs := IPsToStrings(actual.IPAddrs)
if !reflect.DeepEqual(expectedIPs, actualIPs) {
t.Fatalf("expected IP addresses %v, got %v", expectedIPs, actualIPs)
}
}
+25
View File
@@ -0,0 +1,25 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package cert
import (
"net"
"sort"
)
func IPsToStrings(ips []net.IP) []string {
out := make([]string, 0, len(ips))
for _, ip := range ips {
if ip == nil {
continue
}
out = append(out, ip.String())
}
sort.Strings(out)
return out
}
+31
View File
@@ -0,0 +1,31 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package cert
import (
"net"
"reflect"
"testing"
)
func TestIPsToStrings(t *testing.T) {
t.Parallel()
ips := []net.IP{
net.ParseIP("10.96.0.20"),
nil,
net.ParseIP("10.96.0.10"),
net.ParseIP("fd00::1"),
}
expected := []string{
"10.96.0.10",
"10.96.0.20",
"fd00::1",
}
if got := IPsToStrings(ips); !reflect.DeepEqual(expected, got) {
t.Fatalf("expected IP strings %v, got %v", expected, got)
}
}
+2
View File
@@ -20,6 +20,7 @@ import (
"github.com/projectcapsule/capsule/pkg/runtime/indexers/namespace"
"github.com/projectcapsule/capsule/pkg/runtime/indexers/resourcepool"
"github.com/projectcapsule/capsule/pkg/runtime/indexers/tenant"
"github.com/projectcapsule/capsule/pkg/runtime/indexers/tenantowner"
"github.com/projectcapsule/capsule/pkg/runtime/indexers/tenantresource"
"github.com/projectcapsule/capsule/pkg/utils"
)
@@ -32,6 +33,7 @@ type CustomIndexer interface {
func AddToManager(ctx context.Context, log logr.Logger, mgr manager.Manager) error {
indexers := []CustomIndexer{
tenantowner.OwnerNameReference{},
tenant.NamespacesReference{Obj: &capsulev1beta2.Tenant{}},
tenantresource.GlobalServiceAccount{},
tenantresource.GlobalProcessedItems{},
@@ -0,0 +1,8 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package tenantowner
const (
NameIndexerFieldName string = ".spec.name"
)
+35
View File
@@ -0,0 +1,35 @@
// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package tenantowner
import (
"sigs.k8s.io/controller-runtime/pkg/client"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
)
type OwnerNameReference struct{}
func (o OwnerNameReference) Object() client.Object {
return &capsulev1beta2.TenantOwner{}
}
func (o OwnerNameReference) Field() string {
return NameIndexerFieldName
}
func (o OwnerNameReference) Func() client.IndexerFunc {
return func(object client.Object) []string {
instance, ok := object.(*capsulev1beta2.TenantOwner)
if !ok {
return nil
}
if instance.Spec.Name == "" {
return nil
}
return []string{instance.Spec.Name}
}
}
+27 -4
View File
@@ -12,9 +12,32 @@ import (
type TenantStatusOwnersChangedPredicate struct{}
func (TenantStatusOwnersChangedPredicate) Create(event.CreateEvent) bool { return false }
func (TenantStatusOwnersChangedPredicate) Delete(event.DeleteEvent) bool { return false }
func (TenantStatusOwnersChangedPredicate) Generic(event.GenericEvent) bool { return false }
func (TenantStatusOwnersChangedPredicate) Create(e event.CreateEvent) bool {
tenant, ok := e.Object.(*capsulev1beta2.Tenant)
if !ok {
return false
}
return len(tenant.Status.Owners) > 0
}
func (TenantStatusOwnersChangedPredicate) Delete(e event.DeleteEvent) bool {
tenant, ok := e.Object.(*capsulev1beta2.Tenant)
if !ok {
return false
}
return len(tenant.Status.Owners) > 0
}
func (TenantStatusOwnersChangedPredicate) Generic(e event.GenericEvent) bool {
tenant, ok := e.Object.(*capsulev1beta2.Tenant)
if !ok {
return false
}
return len(tenant.Status.Owners) > 0
}
// TenantCountChangedPredicate fires on Tenant create/delete only.
// Update events are ignored because only create and delete change the set of
@@ -45,7 +68,7 @@ func ownersChanged(a, b rbac.OwnerStatusListSpec) bool {
}
for i := range a {
if a[i].Name == b[i].Name && a[i].Kind == b[i].Kind {
if a[i].Name != b[i].Name || a[i].Kind != b[i].Kind {
return true
}
}
+5 -1
View File
@@ -77,8 +77,12 @@ func CollectOwners(
func GetOwnersWithKinds(tenant *capsulev1beta2.Tenant) (owners []string) {
for _, owner := range tenant.Status.Owners {
owners = append(owners, fmt.Sprintf("%s:%s", owner.Kind.String(), owner.Name))
owners = append(owners, OwnerKindIndexKey(owner.Kind.String(), owner.Name))
}
return owners
}
func OwnerKindIndexKey(kind, name string) string {
return fmt.Sprintf("%s:%s", kind, name)
}