From dd398d918b96934e5f800b515fa7efe2c0b4ef89 Mon Sep 17 00:00:00 2001 From: Arvind Thirumurugan Date: Wed, 7 Oct 2026 16:39:00 -0700 Subject: [PATCH 1/2] Add mutating, validating webhooks for Jobs to allow reconciliation, allow job controller to create pods based on label Signed-off-by: Arvind Thirumurugan --- .../manager_integration_test.go | 2 +- .../podsnreplicasets.go | 3 + pkg/webhook/add_handler.go | 3 + pkg/webhook/job/mutating_webhook.go | 85 +++++++++ pkg/webhook/job/mutating_webhook_test.go | 174 ++++++++++++++++++ pkg/webhook/job/validating_webhook.go | 83 +++++++++ pkg/webhook/job/validating_webhook_test.go | 130 +++++++++++++ pkg/webhook/webhook.go | 31 ++++ pkg/webhook/webhook_test.go | 76 +++++++- 9 files changed, 582 insertions(+), 5 deletions(-) create mode 100644 pkg/webhook/job/mutating_webhook.go create mode 100644 pkg/webhook/job/mutating_webhook_test.go create mode 100644 pkg/webhook/job/validating_webhook.go create mode 100644 pkg/webhook/job/validating_webhook_test.go diff --git a/pkg/admissionpolicymanager/manager_integration_test.go b/pkg/admissionpolicymanager/manager_integration_test.go index f00298969..269be078f 100644 --- a/pkg/admissionpolicymanager/manager_integration_test.go +++ b/pkg/admissionpolicymanager/manager_integration_test.go @@ -111,7 +111,7 @@ var _ = Describe("Policies, Policy Bindings and their Effects", Ordered, func() }, Validations: []admissionregistrationv1.Validation{ { - Expression: `((request.namespace.startsWith("fleet-")) || (request.namespace.startsWith("kube-"))) || ((object.metadata.labels["fleet.azure.com/reconcile"] == "managed") && ((request.userInfo.username == "system:serviceaccount:kube-system:deployment-controller") || (request.userInfo.username == "system:serviceaccount:kube-system:replicaset-controller") || (request.userInfo.username == "system:kube-controller-manager")))`, + Expression: `((request.namespace.startsWith("fleet-")) || (request.namespace.startsWith("kube-"))) || ((object.metadata.labels["fleet.azure.com/reconcile"] == "managed") && ((request.userInfo.username == "system:serviceaccount:kube-system:deployment-controller") || (request.userInfo.username == "system:serviceaccount:kube-system:job-controller") || (request.userInfo.username == "system:serviceaccount:kube-system:replicaset-controller") || (request.userInfo.username == "system:kube-controller-manager")))`, Message: "creating pods and replicas is disallowed in the fleet hub cluster", Reason: ptr.To(metav1.StatusReasonForbidden), }, diff --git a/pkg/admissionpolicymanager/podsnreplicasets.go b/pkg/admissionpolicymanager/podsnreplicasets.go index 3ca13de80..b24792b2d 100644 --- a/pkg/admissionpolicymanager/podsnreplicasets.go +++ b/pkg/admissionpolicymanager/podsnreplicasets.go @@ -43,6 +43,7 @@ const ( reconcileIfManagedLabelValue = "managed" deploymentControllerUserName = "system:serviceaccount:kube-system:deployment-controller" + jobControllerUserName = "system:serviceaccount:kube-system:job-controller" replicaSetControllerUserName = "system:serviceaccount:kube-system:replicaset-controller" ) @@ -91,12 +92,14 @@ func (g *PodsAndReplicaSetsValidatingAdmissionPolicyGenerator) PoliciesWithBindi // controllers (or the controller manager, just in case per controller service account is not enabled). hasReconcileIfManagedLabel := RawCELExpr(fmt.Sprintf(`object.metadata.labels["%s"] == "%s"`, reconcileIfManagedLabelKey, reconcileIfManagedLabelValue)) isCreatedByDeploymentController := isFromUsername(deploymentControllerUserName) + isCreatedByJobController := isFromUsername(jobControllerUserName) isCreatedByReplicaSetController := isFromUsername(replicaSetControllerUserName) isCreatedByControllerManager := isFromUsername(kubeControllerManagerUserName) allowIfManagedByAzure := LogicalAnd( hasReconcileIfManagedLabel, LogicalOr( isCreatedByDeploymentController, + isCreatedByJobController, isCreatedByReplicaSetController, isCreatedByControllerManager, ), diff --git a/pkg/webhook/add_handler.go b/pkg/webhook/add_handler.go index dc1e81736..fe8e654ee 100644 --- a/pkg/webhook/add_handler.go +++ b/pkg/webhook/add_handler.go @@ -7,6 +7,7 @@ import ( "go.goms.io/fleet/pkg/webhook/clusterresourceplacementeviction" "go.goms.io/fleet/pkg/webhook/deployment" "go.goms.io/fleet/pkg/webhook/fleetresourcehandler" + "go.goms.io/fleet/pkg/webhook/job" "go.goms.io/fleet/pkg/webhook/membercluster" "go.goms.io/fleet/pkg/webhook/pdb" "go.goms.io/fleet/pkg/webhook/pod" @@ -32,4 +33,6 @@ func init() { AddToManagerFuncs = append(AddToManagerFuncs, clusterresourceplacementdisruptionbudget.Add) AddToManagerFuncs = append(AddToManagerFuncs, deployment.AddMutating) AddToManagerFuncs = append(AddToManagerFuncs, deployment.Add) + AddToManagerFuncs = append(AddToManagerFuncs, job.AddMutating) + AddToManagerFuncs = append(AddToManagerFuncs, job.Add) } diff --git a/pkg/webhook/job/mutating_webhook.go b/pkg/webhook/job/mutating_webhook.go new file mode 100644 index 000000000..08b81777b --- /dev/null +++ b/pkg/webhook/job/mutating_webhook.go @@ -0,0 +1,85 @@ +/* +Copyright 2026 The KubeFleet Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package job + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + + batchv1 "k8s.io/api/batch/v1" + "k8s.io/klog/v2" + "sigs.k8s.io/controller-runtime/pkg/manager" + "sigs.k8s.io/controller-runtime/pkg/webhook" + "sigs.k8s.io/controller-runtime/pkg/webhook/admission" + + "go.goms.io/fleet/pkg/utils" +) + +// MutatingPath is the webhook service path for mutating Job resources. +var MutatingPath = fmt.Sprintf(utils.MutatingPathFmt, batchv1.SchemeGroupVersion.Group, batchv1.SchemeGroupVersion.Version, "job") + +type jobMutator struct { + decoder webhook.AdmissionDecoder +} + +// AddMutating registers the mutating webhook for Jobs with the manager. +func AddMutating(mgr manager.Manager) error { + hookServer := mgr.GetWebhookServer() + hookServer.Register(MutatingPath, &webhook.Admission{Handler: &jobMutator{decoder: admission.NewDecoder(mgr.GetScheme())}}) + return nil +} + +// Handle injects the fleet reconcile label onto the Job and its pod template +// when the request originated from the aksService user. +func (m *jobMutator) Handle(_ context.Context, req admission.Request) admission.Response { + klog.V(2).InfoS("handling job mutating webhook", + "operation", req.Operation, "namespace", req.Namespace, "name", req.Name, "user", req.UserInfo.Username) + + if utils.IsReservedNamespace(req.Namespace) { + return admission.Allowed(fmt.Sprintf("namespace %s is a reserved system namespace, no mutation needed", req.Namespace)) + } + + if !utils.IsAKSService(req.UserInfo) { + return admission.Allowed("user is not aksService, no mutation needed") + } + + var job batchv1.Job + if err := m.decoder.Decode(req, &job); err != nil { + return admission.Errored(http.StatusBadRequest, err) + } + + if job.Labels == nil { + job.Labels = map[string]string{} + } + job.Labels[utils.ReconcileLabelKey] = utils.ReconcileLabelValue + + if job.Spec.Template.Labels == nil { + job.Spec.Template.Labels = map[string]string{} + } + job.Spec.Template.Labels[utils.ReconcileLabelKey] = utils.ReconcileLabelValue + + marshaled, err := json.Marshal(job) + if err != nil { + return admission.Errored(http.StatusInternalServerError, err) + } + + klog.V(2).InfoS("mutated job with reconcile label", + "operation", req.Operation, "namespace", req.Namespace, "name", req.Name) + return admission.PatchResponseFromRaw(req.Object.Raw, marshaled) +} diff --git a/pkg/webhook/job/mutating_webhook_test.go b/pkg/webhook/job/mutating_webhook_test.go new file mode 100644 index 000000000..ef2741718 --- /dev/null +++ b/pkg/webhook/job/mutating_webhook_test.go @@ -0,0 +1,174 @@ +/* +Copyright 2026 The KubeFleet Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package job + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "net/http" + "testing" + + "github.com/google/go-cmp/cmp" + "github.com/google/go-cmp/cmp/cmpopts" + "gomodules.xyz/jsonpatch/v2" + admissionv1 "k8s.io/api/admission/v1" + authenticationv1 "k8s.io/api/authentication/v1" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/utils/ptr" + "sigs.k8s.io/controller-runtime/pkg/webhook/admission" + + "go.goms.io/fleet/pkg/utils" +) + +func TestMutatingPath(t *testing.T) { + want := "/mutate-batch-v1-job" + if MutatingPath != want { + t.Errorf("MutatingPath = %q, want %q", MutatingPath, want) + } +} + +func TestMutatingHandle(t *testing.T) { + scheme := runtime.NewScheme() + if err := batchv1.AddToScheme(scheme); err != nil { + t.Fatalf("batchv1.AddToScheme() = %v, want nil", err) + } + decoder := admission.NewDecoder(scheme) + + aksServiceUser := authenticationv1.UserInfo{ + Username: utils.AKSServiceUserName, + Groups: []string{utils.SystemMastersGroup}, + } + regularUser := authenticationv1.UserInfo{ + Username: "regular-user", + Groups: []string{"system:authenticated"}, + } + + job := newTestJob("test-job", "default", nil, nil) + jobBytes := marshalOrFatal(t, job) + reservedJob := newTestJob("test-job", "kube-system", nil, nil) + reservedJobBytes := marshalOrFatal(t, reservedJob) + + testCases := map[string]struct { + req admission.Request + wantResponse admission.Response + }{ + "mutate job and pod template for aksService user": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobBytes, aksServiceUser), + wantResponse: admission.Response{ + AdmissionResponse: admissionv1.AdmissionResponse{ + Allowed: true, + PatchType: ptr.To(admissionv1.PatchTypeJSONPatch), + }, + Patches: []jsonpatch.JsonPatchOperation{ + { + Operation: "add", + Path: "/metadata/labels", + Value: map[string]any{ + utils.ReconcileLabelKey: utils.ReconcileLabelValue, + }, + }, + { + Operation: "add", + Path: "/spec/template/metadata/labels", + Value: map[string]any{ + utils.ReconcileLabelKey: utils.ReconcileLabelValue, + }, + }, + }, + }, + }, + "skip non-aksService user": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobBytes, regularUser), + wantResponse: admission.Allowed("user is not aksService, no mutation needed"), + }, + "skip reserved namespace": { + req: newAdmissionRequest("test-job", "kube-system", admissionv1.Create, reservedJobBytes, aksServiceUser), + wantResponse: admission.Allowed( + fmt.Sprintf("namespace %s is a reserved system namespace, no mutation needed", "kube-system"), + ), + }, + "error on malformed request object": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, []byte("not valid json"), aksServiceUser), + wantResponse: admission.Errored(http.StatusBadRequest, errors.New("")), + }, + } + + for testName, tc := range testCases { + t.Run(testName, func(t *testing.T) { + mutator := &jobMutator{decoder: decoder} + gotResponse := mutator.Handle(context.Background(), tc.req) + cmpOptions := []cmp.Option{ + cmpopts.IgnoreFields(metav1.Status{}, "Message"), + cmpopts.SortSlices(func(a, b jsonpatch.JsonPatchOperation) bool { + return a.Path < b.Path + }), + } + if diff := cmp.Diff(tc.wantResponse, gotResponse, cmpOptions...); diff != "" { + t.Errorf("Handle() mismatch (-want +got):\n%s", diff) + } + }) + } +} + +func newAdmissionRequest(name, namespace string, operation admissionv1.Operation, raw []byte, userInfo authenticationv1.UserInfo) admission.Request { + return admission.Request{ + AdmissionRequest: admissionv1.AdmissionRequest{ + Name: name, + Namespace: namespace, + Operation: operation, + Object: runtime.RawExtension{Raw: raw}, + UserInfo: userInfo, + }, + } +} + +func marshalOrFatal(t *testing.T, object any) []byte { + t.Helper() + raw, err := json.Marshal(object) + if err != nil { + t.Fatalf("json.Marshal() = %v, want nil", err) + } + return raw +} + +func newTestJob(name, namespace string, jobLabels, podTemplateLabels map[string]string) *batchv1.Job { + return &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: namespace, + Labels: jobLabels, + }, + Spec: batchv1.JobSpec{ + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: podTemplateLabels, + }, + Spec: corev1.PodSpec{ + RestartPolicy: corev1.RestartPolicyNever, + Containers: []corev1.Container{ + {Name: "test", Image: "busybox"}, + }, + }, + }, + }, + } +} diff --git a/pkg/webhook/job/validating_webhook.go b/pkg/webhook/job/validating_webhook.go new file mode 100644 index 000000000..79e4bd3ad --- /dev/null +++ b/pkg/webhook/job/validating_webhook.go @@ -0,0 +1,83 @@ +/* +Copyright 2026 The KubeFleet Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Package job implements admission webhooks for Job resources. +package job + +import ( + "context" + "fmt" + "net/http" + + admissionv1 "k8s.io/api/admission/v1" + batchv1 "k8s.io/api/batch/v1" + "k8s.io/apimachinery/pkg/types" + "k8s.io/klog/v2" + "sigs.k8s.io/controller-runtime/pkg/manager" + "sigs.k8s.io/controller-runtime/pkg/webhook" + "sigs.k8s.io/controller-runtime/pkg/webhook/admission" + + "go.goms.io/fleet/pkg/utils" +) + +const deniedReconcileLabelFmt = "the %s label is reserved for aksService and cannot be set by user %q" + +// ValidationPath is the webhook service path for validating Job resources. +var ValidationPath = fmt.Sprintf(utils.ValidationPathFmt, batchv1.SchemeGroupVersion.Group, batchv1.SchemeGroupVersion.Version, "job") + +type jobValidator struct { + decoder webhook.AdmissionDecoder +} + +// Add registers the validating webhook for Jobs with the manager. +func Add(mgr manager.Manager) error { + hookServer := mgr.GetWebhookServer() + hookServer.Register(ValidationPath, &webhook.Admission{Handler: &jobValidator{decoder: admission.NewDecoder(mgr.GetScheme())}}) + return nil +} + +// Handle rejects Jobs that carry the fleet reconcile label unless the request +// was made by the aksService user with system:masters group membership. +func (v *jobValidator) Handle(_ context.Context, req admission.Request) admission.Response { + namespacedName := types.NamespacedName{Name: req.Name, Namespace: req.Namespace} + klog.V(2).InfoS("handling job validating webhook", + "operation", req.Operation, "namespacedName", namespacedName, "user", req.UserInfo.Username) + + if req.Operation != admissionv1.Create && req.Operation != admissionv1.Update { + return admission.Allowed("operation is not CREATE or UPDATE, no validation needed") + } + + var job batchv1.Job + if err := v.decoder.Decode(req, &job); err != nil { + return admission.Errored(http.StatusBadRequest, err) + } + + hasLabelOnJob := utils.HasReconcileLabel(job.Labels) + hasLabelOnPodTemplate := utils.HasReconcileLabel(job.Spec.Template.Labels) + if !hasLabelOnJob && !hasLabelOnPodTemplate { + return admission.Allowed("job does not have the reconcile label, no validation needed") + } + + if utils.IsAKSService(req.UserInfo) { + klog.V(2).InfoS("aksService user allowed to set reconcile label", + "namespacedName", namespacedName) + return admission.Allowed("aksService user is allowed to set the reconcile label") + } + + klog.V(2).InfoS("denied non-aksService user from setting reconcile label", + "user", req.UserInfo.Username, "groups", req.UserInfo.Groups, "namespacedName", namespacedName) + return admission.Denied(fmt.Sprintf(deniedReconcileLabelFmt, utils.ReconcileLabelKey, req.UserInfo.Username)) +} diff --git a/pkg/webhook/job/validating_webhook_test.go b/pkg/webhook/job/validating_webhook_test.go new file mode 100644 index 000000000..0a82b315b --- /dev/null +++ b/pkg/webhook/job/validating_webhook_test.go @@ -0,0 +1,130 @@ +/* +Copyright 2026 The KubeFleet Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package job + +import ( + "context" + "errors" + "net/http" + "testing" + + "github.com/google/go-cmp/cmp" + "github.com/google/go-cmp/cmp/cmpopts" + admissionv1 "k8s.io/api/admission/v1" + authenticationv1 "k8s.io/api/authentication/v1" + batchv1 "k8s.io/api/batch/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/webhook/admission" + + "go.goms.io/fleet/pkg/utils" +) + +func TestValidationPath(t *testing.T) { + want := "/validate-batch-v1-job" + if ValidationPath != want { + t.Errorf("ValidationPath = %q, want %q", ValidationPath, want) + } +} + +func TestValidatingHandle(t *testing.T) { + scheme := runtime.NewScheme() + if err := batchv1.AddToScheme(scheme); err != nil { + t.Fatalf("batchv1.AddToScheme() = %v, want nil", err) + } + decoder := admission.NewDecoder(scheme) + + aksServiceUser := authenticationv1.UserInfo{ + Username: utils.AKSServiceUserName, + Groups: []string{utils.SystemMastersGroup}, + } + regularUser := authenticationv1.UserInfo{ + Username: "regular-user", + Groups: []string{"system:authenticated"}, + } + + unlabeledJob := marshalOrFatal(t, newTestJob("test-job", "default", nil, nil)) + jobMetadataLabeled := marshalOrFatal(t, newTestJob( + "test-job", + "default", + map[string]string{utils.ReconcileLabelKey: utils.ReconcileLabelValue}, + nil, + )) + podTemplateLabeled := marshalOrFatal(t, newTestJob( + "test-job", + "default", + nil, + map[string]string{utils.ReconcileLabelKey: utils.ReconcileLabelValue}, + )) + + testCases := map[string]struct { + req admission.Request + wantAllowed bool + wantErrCode int32 + }{ + "allow unlabeled job from regular user": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, unlabeledJob, regularUser), + wantAllowed: true, + }, + "allow update that removes labels from both locations": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, unlabeledJob, regularUser), + wantAllowed: true, + }, + "deny regular user with label on job metadata": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobMetadataLabeled, regularUser), + wantAllowed: false, + }, + "deny regular user with label on pod template": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, podTemplateLabeled, regularUser), + wantAllowed: false, + }, + "allow aksService user with labels": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, podTemplateLabeled, aksServiceUser), + wantAllowed: true, + }, + "allow delete operation": { + req: newAdmissionRequest("test-job", "default", admissionv1.Delete, jobMetadataLabeled, regularUser), + wantAllowed: true, + }, + "error on malformed request object": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, []byte("not valid json"), aksServiceUser), + wantAllowed: false, + wantErrCode: http.StatusBadRequest, + }, + } + + for testName, tc := range testCases { + t.Run(testName, func(t *testing.T) { + validator := &jobValidator{decoder: decoder} + gotResponse := validator.Handle(context.Background(), tc.req) + if gotResponse.Allowed != tc.wantAllowed { + t.Errorf("Handle() Allowed = %v, want %v, reason = %v", gotResponse.Allowed, tc.wantAllowed, gotResponse.Result) + } + + if tc.wantErrCode != 0 { + wantResponse := admission.Errored(tc.wantErrCode, errors.New("")) + if diff := cmp.Diff(wantResponse, gotResponse, cmpopts.IgnoreFields(metav1.Status{}, "Message")); diff != "" { + t.Errorf("Handle() error response mismatch (-want +got):\n%s", diff) + } + } + + if !tc.wantAllowed && tc.wantErrCode == 0 && (gotResponse.Result == nil || gotResponse.Result.Message == "") { + t.Error("Handle() denied response should include a non-empty reason message") + } + }) + } +} diff --git a/pkg/webhook/webhook.go b/pkg/webhook/webhook.go index ed84fd4d1..449e3903c 100644 --- a/pkg/webhook/webhook.go +++ b/pkg/webhook/webhook.go @@ -62,6 +62,7 @@ import ( "go.goms.io/fleet/pkg/webhook/clusterresourceplacementeviction" "go.goms.io/fleet/pkg/webhook/deployment" "go.goms.io/fleet/pkg/webhook/fleetresourcehandler" + "go.goms.io/fleet/pkg/webhook/job" "go.goms.io/fleet/pkg/webhook/membercluster" "go.goms.io/fleet/pkg/webhook/pdb" "go.goms.io/fleet/pkg/webhook/pod" @@ -432,6 +433,23 @@ func (w *Config) buildFleetMutatingWebhooks() []admv1.MutatingWebhook { }, TimeoutSeconds: longWebhookTimeout, }, + { + Name: "fleet.job.mutating", + ClientConfig: w.createClientConfig(job.MutatingPath), + FailurePolicy: &ignoreFailurePolicy, + SideEffects: &sideEffortsNone, + AdmissionReviewVersions: admissionReviewVersions, + Rules: []admv1.RuleWithOperations{ + { + Operations: []admv1.OperationType{ + admv1.Create, + admv1.Update, + }, + Rule: createRule([]string{batchv1.SchemeGroupVersion.Group}, []string{batchv1.SchemeGroupVersion.Version}, []string{jobResourceName}, &namespacedScope), + }, + }, + TimeoutSeconds: longWebhookTimeout, + }, } return webHooks } @@ -620,6 +638,19 @@ func (w *Config) buildFleetValidatingWebhooks() []admv1.ValidatingWebhook { TimeoutSeconds: longWebhookTimeout, }) + webHooks = append(webHooks, admv1.ValidatingWebhook{ + Name: "fleet.job.validating", + ClientConfig: w.createClientConfig(job.ValidationPath), + FailurePolicy: &failFailurePolicy, + SideEffects: &sideEffortsNone, + AdmissionReviewVersions: admissionReviewVersions, + Rules: []admv1.RuleWithOperations{{ + Operations: []admv1.OperationType{admv1.Create, admv1.Update}, + Rule: createRule([]string{batchv1.SchemeGroupVersion.Group}, []string{batchv1.SchemeGroupVersion.Version}, []string{jobResourceName}, &namespacedScope), + }}, + TimeoutSeconds: longWebhookTimeout, + }) + return webHooks } diff --git a/pkg/webhook/webhook_test.go b/pkg/webhook/webhook_test.go index 04dc86111..d02d0d650 100644 --- a/pkg/webhook/webhook_test.go +++ b/pkg/webhook/webhook_test.go @@ -9,6 +9,7 @@ import ( "github.com/google/go-cmp/cmp/cmpopts" "github.com/stretchr/testify/assert" admv1 "k8s.io/api/admissionregistration/v1" + batchv1 "k8s.io/api/batch/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "sigs.k8s.io/controller-runtime/pkg/client" @@ -17,6 +18,7 @@ import ( "go.goms.io/fleet/cmd/hubagent/options" "go.goms.io/fleet/pkg/utils" + jobwebhook "go.goms.io/fleet/pkg/webhook/job" testmanager "go.goms.io/fleet/test/utils/manager" ) @@ -33,7 +35,7 @@ func TestBuildFleetMutatingWebhooks(t *testing.T) { serviceURL: "test-url", clientConnectionType: &url, }, - wantLength: 2, + wantLength: 3, }, } @@ -60,7 +62,7 @@ func TestBuildFleetValidatingWebhooks(t *testing.T) { serviceURL: "test-url", clientConnectionType: &url, }, - wantLength: 10, + wantLength: 11, }, "enable workload": { config: Config{ @@ -70,7 +72,7 @@ func TestBuildFleetValidatingWebhooks(t *testing.T) { clientConnectionType: &url, enableWorkload: true, }, - wantLength: 8, + wantLength: 9, }, "enable PDBs": { config: Config{ @@ -80,7 +82,7 @@ func TestBuildFleetValidatingWebhooks(t *testing.T) { clientConnectionType: &url, enablePDBs: true, }, - wantLength: 9, + wantLength: 10, }, } @@ -92,6 +94,72 @@ func TestBuildFleetValidatingWebhooks(t *testing.T) { } } +func TestBuildFleetJobWebhooks(t *testing.T) { + url := options.WebhookClientConnectionType("url") + config := Config{ + serviceNamespace: "test-namespace", + servicePort: 8080, + serviceURL: "test-url", + clientConnectionType: &url, + } + + wantMutating := admv1.MutatingWebhook{ + Name: "fleet.job.mutating", + ClientConfig: config.createClientConfig(jobwebhook.MutatingPath), + FailurePolicy: &ignoreFailurePolicy, + SideEffects: &sideEffortsNone, + AdmissionReviewVersions: admissionReviewVersions, + Rules: []admv1.RuleWithOperations{{ + Operations: []admv1.OperationType{admv1.Create, admv1.Update}, + Rule: createRule([]string{batchv1.SchemeGroupVersion.Group}, []string{batchv1.SchemeGroupVersion.Version}, []string{jobResourceName}, &namespacedScope), + }}, + TimeoutSeconds: longWebhookTimeout, + } + var gotMutating *admv1.MutatingWebhook + mutatingWebhooks := config.buildFleetMutatingWebhooks() + for i := range mutatingWebhooks { + webhook := &mutatingWebhooks[i] + if webhook.Name == wantMutating.Name { + gotMutating = webhook + break + } + } + if gotMutating == nil { + t.Fatalf("buildFleetMutatingWebhooks() did not include %q", wantMutating.Name) + } + if diff := cmp.Diff(wantMutating, *gotMutating); diff != "" { + t.Errorf("Job mutating webhook mismatch (-want +got):\n%s", diff) + } + + wantValidating := admv1.ValidatingWebhook{ + Name: "fleet.job.validating", + ClientConfig: config.createClientConfig(jobwebhook.ValidationPath), + FailurePolicy: &failFailurePolicy, + SideEffects: &sideEffortsNone, + AdmissionReviewVersions: admissionReviewVersions, + Rules: []admv1.RuleWithOperations{{ + Operations: []admv1.OperationType{admv1.Create, admv1.Update}, + Rule: createRule([]string{batchv1.SchemeGroupVersion.Group}, []string{batchv1.SchemeGroupVersion.Version}, []string{jobResourceName}, &namespacedScope), + }}, + TimeoutSeconds: longWebhookTimeout, + } + var gotValidating *admv1.ValidatingWebhook + validatingWebhooks := config.buildFleetValidatingWebhooks() + for i := range validatingWebhooks { + webhook := &validatingWebhooks[i] + if webhook.Name == wantValidating.Name { + gotValidating = webhook + break + } + } + if gotValidating == nil { + t.Fatalf("buildFleetValidatingWebhooks() did not include %q", wantValidating.Name) + } + if diff := cmp.Diff(wantValidating, *gotValidating); diff != "" { + t.Errorf("Job validating webhook mismatch (-want +got):\n%s", diff) + } +} + func TestBuildFleetGuardRailValidatingWebhooks(t *testing.T) { url := options.WebhookClientConnectionType("url") testCases := map[string]struct { From 1124879f3fb3d9e3c103770523a496249beb9f4e Mon Sep 17 00:00:00 2001 From: Arvind Thirumurugan Date: Thu, 8 Oct 2026 16:16:26 -0700 Subject: [PATCH 2/2] add UTs Signed-off-by: Arvind Thirumurugan --- pkg/webhook/job/mutating_webhook_test.go | 102 ++++++++++++++++++--- pkg/webhook/job/validating_webhook_test.go | 80 +++++++++++++++- 2 files changed, 165 insertions(+), 17 deletions(-) diff --git a/pkg/webhook/job/mutating_webhook_test.go b/pkg/webhook/job/mutating_webhook_test.go index ef2741718..06514aa5c 100644 --- a/pkg/webhook/job/mutating_webhook_test.go +++ b/pkg/webhook/job/mutating_webhook_test.go @@ -57,22 +57,77 @@ func TestMutatingHandle(t *testing.T) { Username: utils.AKSServiceUserName, Groups: []string{utils.SystemMastersGroup}, } + aksServiceUserNoMasters := authenticationv1.UserInfo{ + Username: utils.AKSServiceUserName, + Groups: []string{"system:authenticated"}, + } regularUser := authenticationv1.UserInfo{ Username: "regular-user", Groups: []string{"system:authenticated"}, } + regularUserWithMasters := authenticationv1.UserInfo{ + Username: "regular-user", + Groups: []string{utils.SystemMastersGroup}, + } job := newTestJob("test-job", "default", nil, nil) jobBytes := marshalOrFatal(t, job) + jobWithLabels := newTestJob( + "test-job", + "default", + map[string]string{"app": "test"}, + map[string]string{"app": "test"}, + ) + jobWithLabelsBytes := marshalOrFatal(t, jobWithLabels) + alreadyLabeledJob := newTestJob( + "test-job", + "default", + map[string]string{"app": "test", utils.ReconcileLabelKey: utils.ReconcileLabelValue}, + map[string]string{"app": "test", utils.ReconcileLabelKey: utils.ReconcileLabelValue}, + ) + alreadyLabeledJobBytes := marshalOrFatal(t, alreadyLabeledJob) reservedJob := newTestJob("test-job", "kube-system", nil, nil) reservedJobBytes := marshalOrFatal(t, reservedJob) + fleetSystemJob := newTestJob("test-job", utils.FleetSystemNamespace, nil, nil) + fleetSystemJobBytes := marshalOrFatal(t, fleetSystemJob) + + wantMutatedResponse := admission.Response{ + AdmissionResponse: admissionv1.AdmissionResponse{ + Allowed: true, + PatchType: ptr.To(admissionv1.PatchTypeJSONPatch), + }, + Patches: []jsonpatch.JsonPatchOperation{ + { + Operation: "add", + Path: "/metadata/labels", + Value: map[string]any{ + utils.ReconcileLabelKey: utils.ReconcileLabelValue, + }, + }, + { + Operation: "add", + Path: "/spec/template/metadata/labels", + Value: map[string]any{ + utils.ReconcileLabelKey: utils.ReconcileLabelValue, + }, + }, + }, + } testCases := map[string]struct { req admission.Request wantResponse admission.Response }{ - "mutate job and pod template for aksService user": { - req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobBytes, aksServiceUser), + "mutate job and pod template on create for aksService user": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobBytes, aksServiceUser), + wantResponse: wantMutatedResponse, + }, + "mutate job and pod template on update for aksService user": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, jobBytes, aksServiceUser), + wantResponse: wantMutatedResponse, + }, + "mutate job and pod template while preserving existing labels": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobWithLabelsBytes, aksServiceUser), wantResponse: admission.Response{ AdmissionResponse: admissionv1.AdmissionResponse{ Allowed: true, @@ -81,31 +136,54 @@ func TestMutatingHandle(t *testing.T) { Patches: []jsonpatch.JsonPatchOperation{ { Operation: "add", - Path: "/metadata/labels", - Value: map[string]any{ - utils.ReconcileLabelKey: utils.ReconcileLabelValue, - }, + Path: "/metadata/labels/fleet.azure.com~1reconcile", + Value: utils.ReconcileLabelValue, }, { Operation: "add", - Path: "/spec/template/metadata/labels", - Value: map[string]any{ - utils.ReconcileLabelKey: utils.ReconcileLabelValue, - }, + Path: "/spec/template/metadata/labels/fleet.azure.com~1reconcile", + Value: utils.ReconcileLabelValue, }, }, }, }, - "skip non-aksService user": { + "return no-op patch when both labels are already present": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, alreadyLabeledJobBytes, aksServiceUser), + wantResponse: admission.Response{ + AdmissionResponse: admissionv1.AdmissionResponse{ + Allowed: true, + }, + Patches: []jsonpatch.JsonPatchOperation{}, + }, + }, + "skip non-aksService user on create": { req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobBytes, regularUser), wantResponse: admission.Allowed("user is not aksService, no mutation needed"), }, - "skip reserved namespace": { + "skip non-aksService user on update": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, jobBytes, regularUser), + wantResponse: admission.Allowed("user is not aksService, no mutation needed"), + }, + "skip aksService user without system masters": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobBytes, aksServiceUserNoMasters), + wantResponse: admission.Allowed("user is not aksService, no mutation needed"), + }, + "skip non-aksService user with system masters": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, jobBytes, regularUserWithMasters), + wantResponse: admission.Allowed("user is not aksService, no mutation needed"), + }, + "skip kube-system namespace": { req: newAdmissionRequest("test-job", "kube-system", admissionv1.Create, reservedJobBytes, aksServiceUser), wantResponse: admission.Allowed( fmt.Sprintf("namespace %s is a reserved system namespace, no mutation needed", "kube-system"), ), }, + "skip fleet-system namespace": { + req: newAdmissionRequest("test-job", utils.FleetSystemNamespace, admissionv1.Update, fleetSystemJobBytes, aksServiceUser), + wantResponse: admission.Allowed( + fmt.Sprintf("namespace %s is a reserved system namespace, no mutation needed", utils.FleetSystemNamespace), + ), + }, "error on malformed request object": { req: newAdmissionRequest("test-job", "default", admissionv1.Create, []byte("not valid json"), aksServiceUser), wantResponse: admission.Errored(http.StatusBadRequest, errors.New("")), diff --git a/pkg/webhook/job/validating_webhook_test.go b/pkg/webhook/job/validating_webhook_test.go index 0a82b315b..87183833c 100644 --- a/pkg/webhook/job/validating_webhook_test.go +++ b/pkg/webhook/job/validating_webhook_test.go @@ -52,10 +52,18 @@ func TestValidatingHandle(t *testing.T) { Username: utils.AKSServiceUserName, Groups: []string{utils.SystemMastersGroup}, } + aksServiceUserNoMasters := authenticationv1.UserInfo{ + Username: utils.AKSServiceUserName, + Groups: []string{"system:authenticated"}, + } regularUser := authenticationv1.UserInfo{ Username: "regular-user", Groups: []string{"system:authenticated"}, } + regularUserWithMasters := authenticationv1.UserInfo{ + Username: "regular-user", + Groups: []string{utils.SystemMastersGroup}, + } unlabeledJob := marshalOrFatal(t, newTestJob("test-job", "default", nil, nil)) jobMetadataLabeled := marshalOrFatal(t, newTestJob( @@ -70,6 +78,24 @@ func TestValidatingHandle(t *testing.T) { nil, map[string]string{utils.ReconcileLabelKey: utils.ReconcileLabelValue}, )) + bothLabeled := marshalOrFatal(t, newTestJob( + "test-job", + "default", + map[string]string{utils.ReconcileLabelKey: utils.ReconcileLabelValue}, + map[string]string{utils.ReconcileLabelKey: utils.ReconcileLabelValue}, + )) + jobMetadataWithDifferentReconcileValue := marshalOrFatal(t, newTestJob( + "test-job", + "default", + map[string]string{utils.ReconcileLabelKey: "other-value"}, + nil, + )) + reservedNamespaceLabeled := marshalOrFatal(t, newTestJob( + "test-job", + "kube-system", + map[string]string{utils.ReconcileLabelKey: utils.ReconcileLabelValue}, + nil, + )) testCases := map[string]struct { req admission.Request @@ -84,22 +110,66 @@ func TestValidatingHandle(t *testing.T) { req: newAdmissionRequest("test-job", "default", admissionv1.Update, unlabeledJob, regularUser), wantAllowed: true, }, - "deny regular user with label on job metadata": { + "allow aksService user to create with label on job metadata": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobMetadataLabeled, aksServiceUser), + wantAllowed: true, + }, + "allow aksService user to create with label on pod template": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, podTemplateLabeled, aksServiceUser), + wantAllowed: true, + }, + "allow aksService user to update with labels on both locations": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, bothLabeled, aksServiceUser), + wantAllowed: true, + }, + "allow aksService user to create labeled job in reserved namespace": { + req: newAdmissionRequest("test-job", "kube-system", admissionv1.Create, reservedNamespaceLabeled, aksServiceUser), + wantAllowed: true, + }, + "deny regular user create with label on job metadata": { req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobMetadataLabeled, regularUser), wantAllowed: false, }, - "deny regular user with label on pod template": { + "deny regular user create with different reconcile label value": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, jobMetadataWithDifferentReconcileValue, regularUser), + wantAllowed: false, + }, + "deny regular user create with label on pod template": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, podTemplateLabeled, regularUser), + wantAllowed: false, + }, + "deny regular user update with label on job metadata": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, jobMetadataLabeled, regularUser), + wantAllowed: false, + }, + "deny regular user update with label on pod template": { req: newAdmissionRequest("test-job", "default", admissionv1.Update, podTemplateLabeled, regularUser), wantAllowed: false, }, - "allow aksService user with labels": { - req: newAdmissionRequest("test-job", "default", admissionv1.Update, podTemplateLabeled, aksServiceUser), - wantAllowed: true, + "deny regular user update with labels on both locations": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, bothLabeled, regularUser), + wantAllowed: false, + }, + "deny regular user create with label in reserved namespace": { + req: newAdmissionRequest("test-job", "kube-system", admissionv1.Create, reservedNamespaceLabeled, regularUser), + wantAllowed: false, + }, + "deny aksService user without system masters": { + req: newAdmissionRequest("test-job", "default", admissionv1.Create, bothLabeled, aksServiceUserNoMasters), + wantAllowed: false, + }, + "deny non-aksService user with system masters": { + req: newAdmissionRequest("test-job", "default", admissionv1.Update, bothLabeled, regularUserWithMasters), + wantAllowed: false, }, "allow delete operation": { req: newAdmissionRequest("test-job", "default", admissionv1.Delete, jobMetadataLabeled, regularUser), wantAllowed: true, }, + "allow connect operation": { + req: newAdmissionRequest("test-job", "default", admissionv1.Connect, bothLabeled, regularUser), + wantAllowed: true, + }, "error on malformed request object": { req: newAdmissionRequest("test-job", "default", admissionv1.Create, []byte("not valid json"), aksServiceUser), wantAllowed: false,