diff --git a/adr/2026-10-07-discoverable-endpoint-index.md b/adr/2026-10-07-discoverable-endpoint-index.md new file mode 100644 index 000000000..92cc151ae --- /dev/null +++ b/adr/2026-10-07-discoverable-endpoint-index.md @@ -0,0 +1,56 @@ +# Index discoverable endpoints for admission conflict checks + +**Status**: Accepted +**Date**: 2026-10-07 +**Deciders**: Not recorded + +## Context + +Admission checks for a workspace with discoverable endpoints currently list every DevWorkspace in its namespace, then scan the returned specs for colliding Service names. This copies and examines unrelated workspaces. The existing `discoverableEndpointNames` helper already extracts normalized, unique endpoint names when the discoverable attribute is boolean or string `true` and exposure is not `none`. + +## Decision + +Register a controller-runtime field index on DevWorkspace objects under `controller.devfile.io/discoverable-endpoint`. Its extractor reuses `discoverableEndpointNames`. Register the index before webhook handlers are registered and before the manager starts. + +For each incoming normalized discoverable endpoint name, admission performs one namespace-scoped List using that exact indexed name. Keep the existing update gating, cached-client copy behavior, self-UID exclusion, error handling, and defensive endpoint scan of returned candidates. The scan remains the final admission check for a matching candidate. + +The index covers endpoints visible in `.spec.template`; contributed endpoints remain protected by the controller's Service synchronization ownership check. Admission remains eventually consistent, so concurrent requests can both pass; the controller's fresh Service ownership check remains authoritative. No API, RBAC, dependency, or code-generation changes are required. + +## Considered Alternatives + +### Alternative 1: Keep listing and scanning the whole namespace + +This is simpler, but continues copying every workspace and scanning unrelated specs for each admission check. + +**Rejected because**: The field index targets the matching workspaces through the existing cache. + +### Alternative 2: Disable cache object copying + +This would avoid copies for all listed workspaces but exposes shared cached objects to mutation and races. + +**Rejected because**: Limiting query results with the index avoids copying unrelated workspaces while retaining the client's normal safety behavior. + +### Alternative 3: Add a custom cache or external endpoint index owner + +This would add another mechanism and lifecycle to maintain. + +**Rejected because**: A native controller-runtime field index provides the needed lookup within the existing manager and cache. + +## Consequences + +### Positive + +Only workspaces matching each incoming discoverable endpoint name are copied and scanned. + +### Negative + +The cache uses memory for the index and CPU to extract and update indexed names. Requests with multiple incoming names issue one List per name. + +### Neutral + +The lookup remains a best-effort admission check over cached state. Existing update gating, conflict responses, and the controller's Service conflict guard retain their roles. + +## References + +- [Endpoint validation](../webhook/workspace/handler/validate.go) +- [Webhook configuration](../webhook/workspace/config.go) diff --git a/adr/2026-10-08-discoverable-service-certificates.md b/adr/2026-10-08-discoverable-service-certificates.md new file mode 100644 index 000000000..5ea136f37 --- /dev/null +++ b/adr/2026-10-08-discoverable-service-certificates.md @@ -0,0 +1,59 @@ +# Mount per-service certificates for discoverable endpoints + +**Status**: Accepted +**Date**: 2026-10-08 +**Deciders**: Not recorded + +## Context + +With TLS enabled, cluster routing creates an aggregate workspace Service and one Service per discoverable endpoint. OpenShift issues a serving certificate for each annotated Service. Mounting all certificates at `/var/serving-cert/` overlaps mounts; clients need a distinct certificate location for each discoverable Service. + +## Decision + +Keep the aggregate Service certificate mounted at `/var/serving-cert/`. Mount each discoverable Service's certificate at `/var/serving-cert//`. The solver annotates each Service with its own name, and the Secret reference continues to use that actual Service name. + +Keep `devworkspace-serving-cert-` when the full generated volume name fits within 63 characters; otherwise use `devworkspace-cert-` plus the hex encoding of the first 10 SHA-256 bytes of the Service name. Distinct prefixes prevent a digest name from deterministically colliding with one generated for a short Service name. Hashing avoids truncation collisions; the 80-bit digest can still collide and is not guaranteed unique. Service and Secret names remain unchanged. + +## Considered Alternatives + +### Alternative 1: Share the aggregate certificate mount + +Expose the aggregate certificate to workloads and use it for discoverable Service connections. + +**Rejected because**: OpenShift binds each serving certificate to its Service's internal DNS name. Reusing the aggregate certificate would leave the discoverable Service DNS names uncovered. [OpenShift certificate documentation](https://docs.redhat.com/en/documentation/openshift_container_platform/4.21/html/security_and_compliance/configuring-certificates) + +### Alternative 2: Truncate long volume names or reject long Service names + +Shorten generated names to fit Kubernetes' 63-character volume-name limit, or disallow endpoint names that produce longer values. + +**Rejected because**: Truncation can merge names from distinct Services, while rejecting them would disallow otherwise valid Service names. + +### Alternative 3: Hash every volume name + +Use a fixed-length digest-based name for all Services. + +**Rejected because**: It removes readable names for ordinary cases without avoiding any additional limit. + +## Consequences + +### Positive + +- Each discoverable Service has a distinct certificate mount path alongside the unchanged aggregate mount. +- Long valid Service names can be represented by Kubernetes-compatible volume names without truncation. + +### Negative + +- Each discoverable Service adds a pod volume and mount. +- The 80-bit digest has a theoretical collision risk for distinct long Service names. + +### Neutral + +- Workloads use the actual Service name to identify the corresponding Secret; only the generated volume name may differ for long names. +- The reserved-name guard prevents a discoverable Service from taking the aggregate Service name, and duplicate normalized Service names are rejected. + +## References + +- [Cluster solver](../controllers/controller/devworkspacerouting/solvers/cluster_solver.go) +- [Cluster solver certificate mount tests](../controllers/controller/devworkspacerouting/solvers/cluster_solver_test.go) +- [Serving certificate volume naming](../pkg/common/naming.go) +- [Discoverable Service name guard](../controllers/controller/devworkspacerouting/solvers/common.go) diff --git a/adr/2026-10-08-routing-service-ownership.md b/adr/2026-10-08-routing-service-ownership.md new file mode 100644 index 000000000..e9da51f56 --- /dev/null +++ b/adr/2026-10-08-routing-service-ownership.md @@ -0,0 +1,62 @@ +# Guard routing Service ownership during synchronization + +**Status**: Accepted +**Date**: 2026-10-08 +**Deciders**: Not recorded + +## Context + +DevWorkspaceRouting Services share namespace names, including discoverable endpoints. Admission reads cached workspace state, so concurrent requests can both pass; the cache can also hide a Service missing its workspace ID label. Treating that miss as proof a name is free risks overwriting another routing's Service. + +## Decision + +Identify routing Services by the discoverable annotation or DevWorkspaceRouting controller owner. When configured, read them through the non-caching client (controller setup supplies it) and use this fresh result for ownership. If the desired Service has a controller owner, its UID must match the existing controller owner's UID; that UID remains authoritative when the workspace ID label is stale or missing, allowing synchronization to repair the label. Without a desired controller owner, require a nonempty workspace ID label matching the existing Service. Reject a foreign same-name Service as a permanent conflict. + +If Service creation returns AlreadyExists, retry instead of using the generic update fallback. Updates carry the existing UID, resourceVersion, and ClusterIP; cleanup deletes only a routing-controlled Service with UID and resourceVersion preconditions. These guards reject writes against a changed or replaced Service; update conflicts and missing objects trigger another reconcile. The read and ownership decision are not atomic, and this does not reserve names across reconcilers. + +## Considered Alternatives + +### Alternative 1: Trust admission conflict checks + +Use cached admission as the sole protection for discoverable endpoint names. + +**Rejected because**: Stale cache state lets concurrent requests both pass before a Service exists. + +### Alternative 2: Treat workspace ID labels as ownership proof + +Require the existing Service label to match the desired workspace ID. + +**Rejected because**: A label match does not prove matching controller ownership, and labels can drift. A valid routing owner can repair a stale label without losing its Service. + +### Alternative 3: Keep the generic AlreadyExists update fallback + +Use generic update after Service creation reports AlreadyExists. + +**Rejected because**: The fallback lacks a fresh object for ownership validation and identity preservation. + +## Consequences + +### Positive + +- The fresh check protects foreign Services hidden by cache filtering. +- Owned Services with stale labels can be repaired while preserving identity and ClusterIP. + +### Negative + +- Fresh reads add API calls and latency; changed Services can make guarded writes conflict and require another reconcile. + +### Neutral + +- Admission remains advisory; the controller check is authoritative. Concurrent creates rely on Kubernetes name uniqueness and retry. + +## References + +- [Service synchronization](../pkg/provision/sync/sync.go) +- [Service update fields](../pkg/provision/sync/update.go) +- [Service conflict error](../pkg/provision/sync/service.go) +- [Routing Service cleanup](../controllers/controller/devworkspacerouting/sync_services.go) +- [Routing controller client setup](../controllers/controller/devworkspacerouting/devworkspacerouting_controller.go) +- [Ownership and cache-miss tests](../pkg/provision/sync/service_test.go) +- [Routing cleanup and conflict tests](../controllers/controller/devworkspacerouting/sync_services_test.go) +- [Label repair tests](../controllers/controller/devworkspacerouting/workspace_name_test.go) +- [Admission index decision](2026-10-07-discoverable-endpoint-index.md) diff --git a/controllers/controller/devworkspacerouting/devworkspacerouting_controller.go b/controllers/controller/devworkspacerouting/devworkspacerouting_controller.go index 1a4bbfd42..6d3bb9118 100644 --- a/controllers/controller/devworkspacerouting/devworkspacerouting_controller.go +++ b/controllers/controller/devworkspacerouting/devworkspacerouting_controller.go @@ -21,6 +21,8 @@ import ( "fmt" "time" + dwv2 "github.com/devfile/api/v2/pkg/apis/workspaces/v1alpha2" + "github.com/devfile/devworkspace-operator/controllers/controller/devworkspacerouting/solvers" maputils "github.com/devfile/devworkspace-operator/internal/map" "github.com/devfile/devworkspace-operator/pkg/config" @@ -34,7 +36,9 @@ import ( corev1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" k8sErrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/utils/ptr" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" @@ -54,8 +58,9 @@ const devWorkspaceRoutingFinalizer = "devworkspacerouting.controller.devfile.io" // DevWorkspaceRoutingReconciler reconciles a DevWorkspaceRouting object type DevWorkspaceRoutingReconciler struct { client.Client - Log logr.Logger - Scheme *runtime.Scheme + NonCachingClient client.Client + Log logr.Logger + Scheme *runtime.Scheme // SolverGetter will be used to get solvers for a particular devWorkspaceRouting SolverGetter solvers.RoutingSolverGetter // Enable additional debug logging @@ -97,6 +102,7 @@ func (r *DevWorkspaceRoutingReconciler) Reconcile(ctx context.Context, req ctrl. solver, err := r.SolverGetter.GetSolver(r.Client, instance.Spec.RoutingClass) if err != nil { if errors.Is(err, solvers.RoutingNotSupported) { + reqLogger.Info("Routing class not supported by this controller, skipping reconciliation", "routingClass", instance.Spec.RoutingClass) return reconcile.Result{}, nil } return reconcile.Result{}, r.markRoutingFailed(instance, fmt.Sprintf("Invalid routingClass for DevWorkspace: %s", err)) @@ -126,9 +132,11 @@ func (r *DevWorkspaceRoutingReconciler) Reconcile(ctx context.Context, req ctrl. } workspaceMeta := solvers.DevWorkspaceMetadata{ - DevWorkspaceId: instance.Spec.DevWorkspaceId, - Namespace: instance.Namespace, - PodSelector: instance.Spec.PodSelector, + DevWorkspaceId: instance.Spec.DevWorkspaceId, + DevWorkspaceName: workspaceName(instance), + DevWorkspaceRoutingUID: instance.UID, + Namespace: instance.Namespace, + PodSelector: instance.Spec.PodSelector, } restrictedAccess, setRestrictedAccess := instance.Annotations[constants.DevWorkspaceRestrictedAccessAnnotation] @@ -150,6 +158,18 @@ func (r *DevWorkspaceRoutingReconciler) Reconcile(ctx context.Context, req ctrl. return reconcile.Result{}, r.markRoutingFailed(instance, fmt.Sprintf("Unable to provision networking for DevWorkspace: %s", invalid)) } + var conflict *solvers.ServiceConflictError + if errors.As(err, &conflict) { + reqLogger.Error(conflict, "Routing controller detected a service conflict", "endpointName", conflict.EndpointName, "workspaceName", conflict.WorkspaceName) + return reconcile.Result{}, r.markRoutingFailed(instance, fmt.Sprintf("Unable to provision networking for DevWorkspace: %s", conflict)) + } + + var duplicate *solvers.DuplicateEndpointError + if errors.As(err, &duplicate) { + reqLogger.Error(duplicate, "Routing controller detected a duplicate endpoint name", "endpointName", duplicate.EndpointName) + return reconcile.Result{}, r.markRoutingFailed(instance, fmt.Sprintf("Unable to provision networking for DevWorkspace: %s", duplicate)) + } + // generic error, just fail the reconciliation return reconcile.Result{}, err } @@ -241,6 +261,15 @@ func (r *DevWorkspaceRoutingReconciler) Reconcile(ctx context.Context, req ctrl. return reconcile.Result{}, r.reconcileStatus(instance, &routingObjects, exposedEndpoints, endpointsAreReady, "") } +// Generated by Codex. +func workspaceName(routing *controllerv1alpha1.DevWorkspaceRouting) string { + owner := metav1.GetControllerOf(routing) + if owner == nil || owner.Kind != "DevWorkspace" || schema.FromAPIVersionAndKind(owner.APIVersion, owner.Kind).Group != dwv2.SchemeGroupVersion.Group { + return "" + } + return owner.Name +} + // setFinalizer ensures a finalizer is set on a devWorkspaceRouting instance; no-op if finalizer is already present. func (r *DevWorkspaceRoutingReconciler) setFinalizer(reqLogger logr.Logger, solver solvers.RoutingSolver, m *controllerv1alpha1.DevWorkspaceRouting) error { if !solver.FinalizerRequired(m) || contains(m.GetFinalizers(), devWorkspaceRoutingFinalizer) { @@ -335,6 +364,14 @@ func remove(list []string, s string) []string { } func (r *DevWorkspaceRoutingReconciler) SetupWithManager(mgr ctrl.Manager) error { + if r.NonCachingClient == nil { + freshClient, err := client.New(mgr.GetConfig(), client.Options{Scheme: mgr.GetScheme()}) + if err != nil { + return err + } + r.NonCachingClient = freshClient + } + maxConcurrentReconciles, err := config.GetMaxConcurrentReconciles() if err != nil { return err diff --git a/controllers/controller/devworkspacerouting/devworkspacerouting_controller_test.go b/controllers/controller/devworkspacerouting/devworkspacerouting_controller_test.go index 6c8cf1a30..773dd697c 100644 --- a/controllers/controller/devworkspacerouting/devworkspacerouting_controller_test.go +++ b/controllers/controller/devworkspacerouting/devworkspacerouting_controller_test.go @@ -14,22 +14,86 @@ package devworkspacerouting_test import ( + "context" "fmt" + "testing" - controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" - "github.com/devfile/devworkspace-operator/pkg/common" - "github.com/devfile/devworkspace-operator/pkg/config" - "github.com/devfile/devworkspace-operator/pkg/constants" - "github.com/devfile/devworkspace-operator/pkg/infrastructure" . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" routeV1 "github.com/openshift/api/route/v1" corev1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" k8sErrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/util/intstr" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/controllers/controller/devworkspacerouting" + "github.com/devfile/devworkspace-operator/controllers/controller/devworkspacerouting/solvers" + "github.com/devfile/devworkspace-operator/pkg/common" + "github.com/devfile/devworkspace-operator/pkg/config" + "github.com/devfile/devworkspace-operator/pkg/constants" + "github.com/devfile/devworkspace-operator/pkg/infrastructure" ) +func TestReconcileMarksDuplicateEndpointRoutingFailed(t *testing.T) { + infrastructure.InitializeForTesting(infrastructure.Kubernetes) + g := NewWithT(t) + scheme := runtime.NewScheme() + g.Expect(corev1.AddToScheme(scheme)).To(Succeed()) + g.Expect(controllerv1alpha1.AddToScheme(scheme)).To(Succeed()) + routing := &controllerv1alpha1.DevWorkspaceRouting{ + ObjectMeta: metav1.ObjectMeta{Name: "duplicate-endpoints", Namespace: "test-namespace"}, + Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{ + DevWorkspaceId: "workspace-id", + RoutingClass: controllerv1alpha1.DevWorkspaceRoutingCluster, + PodSelector: map[string]string{constants.DevWorkspaceIDLabel: "workspace-id"}, + Endpoints: map[string]controllerv1alpha1.EndpointList{ + "machine1": {{ + Name: "my--endpoint", + TargetPort: 5000, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }}, + "machine2": {{ + Name: "my-endpoint", + TargetPort: 6000, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }}, + }, + }, + } + fakeClient := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&controllerv1alpha1.DevWorkspaceRouting{}).WithObjects(routing).Build() + reconciler := &devworkspacerouting.DevWorkspaceRoutingReconciler{ + Client: fakeClient, + Log: ctrl.Log, + Scheme: scheme, + SolverGetter: &solvers.SolverGetter{}, + } + testCtx := context.Background() + key := client.ObjectKeyFromObject(routing) + result, err := reconciler.Reconcile(testCtx, ctrl.Request{NamespacedName: key}) + g.Expect(err).To(Succeed()) + g.Expect(result).To(Equal(ctrl.Result{})) + stored := &controllerv1alpha1.DevWorkspaceRouting{} + g.Expect(fakeClient.Get(testCtx, key, stored)).To(Succeed()) + g.Expect(stored.Status.Phase).To(Equal(controllerv1alpha1.RoutingFailed)) + g.Expect(stored.Status.Message).To(ContainSubstring("is declared by more than one component")) + // Container-map iteration determines which raw alias is reported. + g.Expect(stored.Status.Message).To(Or(ContainSubstring("'my--endpoint'"), ContainSubstring("'my-endpoint'"))) + services := &corev1.ServiceList{} + g.Expect(fakeClient.List(testCtx, services)).To(Succeed()) + g.Expect(services.Items).To(BeEmpty()) +} + var _ = Describe("DevWorkspaceRouting Controller", func() { Context("Basic DevWorkspaceRouting Tests", func() { It("Gets Ready Status on OpenShift", func() { diff --git a/controllers/controller/devworkspacerouting/solvers/basic_solver.go b/controllers/controller/devworkspacerouting/solvers/basic_solver.go index 50df27fa1..e7e6f4012 100644 --- a/controllers/controller/devworkspacerouting/solvers/basic_solver.go +++ b/controllers/controller/devworkspacerouting/solvers/basic_solver.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -16,6 +16,8 @@ package solvers import ( + "sigs.k8s.io/controller-runtime/pkg/client" + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" "github.com/devfile/devworkspace-operator/pkg/config" "github.com/devfile/devworkspace-operator/pkg/constants" @@ -47,10 +49,19 @@ var nginxIngressAnnotations = func(endpointName string, endpointAnnotations map[ // According to the current cluster there is different behavior: // Kubernetes: use Ingresses without TLS // OpenShift: use Routes with TLS enabled -type BasicSolver struct{} +type BasicSolver struct { + client client.Client +} var _ RoutingSolver = (*BasicSolver)(nil) +// NewBasicSolver creates a new BasicSolver with the provided dependencies +func NewBasicSolver(client client.Client) *BasicSolver { + return &BasicSolver{ + client: client, + } +} + func (s *BasicSolver) FinalizerRequired(*controllerv1alpha1.DevWorkspaceRouting) bool { return false } @@ -70,7 +81,11 @@ func (s *BasicSolver) GetSpecObjects(routing *controllerv1alpha1.DevWorkspaceRou spec := routing.Spec services := getServicesForEndpoints(spec.Endpoints, workspaceMeta) - services = append(services, GetDiscoverableServicesForEndpoints(spec.Endpoints, workspaceMeta)...) + discoverableServices, err := GetDiscoverableServicesForEndpoints(spec.Endpoints, workspaceMeta, s.client) + if err != nil { + return RoutingObjects{}, err + } + services = append(services, discoverableServices...) routingObjects.Services = services if infrastructure.IsOpenShift() { routingObjects.Routes = getRoutesForSpec(routingSuffix, spec.Endpoints, workspaceMeta) diff --git a/controllers/controller/devworkspacerouting/solvers/cluster_solver.go b/controllers/controller/devworkspacerouting/solvers/cluster_solver.go index 7c5462252..901ab4c4b 100644 --- a/controllers/controller/devworkspacerouting/solvers/cluster_solver.go +++ b/controllers/controller/devworkspacerouting/solvers/cluster_solver.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -23,6 +23,8 @@ import ( corev1 "k8s.io/api/core/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" ) @@ -31,11 +33,20 @@ const ( ) type ClusterSolver struct { - TLS bool + TLS bool + client client.Client } var _ RoutingSolver = (*ClusterSolver)(nil) +// NewClusterSolver creates a new ClusterSolver with the provided dependencies +func NewClusterSolver(client client.Client, tls bool) *ClusterSolver { + return &ClusterSolver{ + TLS: tls, + client: client, + } +} + func (s *ClusterSolver) FinalizerRequired(*controllerv1alpha1.DevWorkspaceRouting) bool { return false } @@ -47,6 +58,11 @@ func (s *ClusterSolver) Finalize(*controllerv1alpha1.DevWorkspaceRouting) error func (s *ClusterSolver) GetSpecObjects(routing *controllerv1alpha1.DevWorkspaceRouting, workspaceMeta DevWorkspaceMetadata) (RoutingObjects, error) { spec := routing.Spec services := getServicesForEndpoints(spec.Endpoints, workspaceMeta) + discoverableServices, err := GetDiscoverableServicesForEndpoints(spec.Endpoints, workspaceMeta, s.client) + if err != nil { + return RoutingObjects{}, err + } + services = append(services, discoverableServices...) podAdditions := &controllerv1alpha1.PodAdditions{} if s.TLS { readOnlyMode := int32(420) @@ -64,10 +80,14 @@ func (s *ClusterSolver) GetSpecObjects(routing *controllerv1alpha1.DevWorkspaceR }, }, }) + mountPath := "/var/serving-cert/" + if service.Annotations[constants.DevWorkspaceDiscoverableServiceAnnotation] == "true" { + mountPath += service.Name + "/" + } podAdditions.VolumeMounts = append(podAdditions.VolumeMounts, corev1.VolumeMount{ Name: common.ServingCertVolumeName(service.Name), ReadOnly: true, - MountPath: "/var/serving-cert/", + MountPath: mountPath, }) } } diff --git a/controllers/controller/devworkspacerouting/solvers/cluster_solver_test.go b/controllers/controller/devworkspacerouting/solvers/cluster_solver_test.go new file mode 100644 index 000000000..2e13c2985 --- /dev/null +++ b/controllers/controller/devworkspacerouting/solvers/cluster_solver_test.go @@ -0,0 +1,206 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 solvers + +import ( + "strings" + "testing" + + "github.com/stretchr/testify/assert" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/pkg/common" +) + +// Generated by Codex. +func TestClusterSolverGetSpecObjectsServingCertMounts(t *testing.T) { + workspaceMeta := DevWorkspaceMetadata{ + DevWorkspaceId: "workspace-id", + DevWorkspaceName: "workspace", + Namespace: "workspace-namespace", + PodSelector: map[string]string{"workspace": "workspace-id"}, + } + endpoints := map[string]controllerv1alpha1.EndpointList{ + "main": { + { + Name: "api", + TargetPort: 8080, + Exposure: controllerv1alpha1.PublicEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }, + { + Name: "database", + TargetPort: 5432, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }, + }, + } + aggregateServiceName := common.ServiceName(workspaceMeta.DevWorkspaceId) + + tests := []struct { + name string + tls bool + expectedMounts map[string]string + }{ + { + name: "TLS discoverable endpoint certificates use distinct mounts", + tls: true, + expectedMounts: map[string]string{ + "/var/serving-cert/": aggregateServiceName, + "/var/serving-cert/api/": "api", + "/var/serving-cert/database/": "database", + }, + }, + { + name: "plaintext routing requests no certificates", + tls: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + scheme := runtime.NewScheme() + assert.NoError(t, corev1.AddToScheme(scheme)) + solver := NewClusterSolver(fake.NewClientBuilder().WithScheme(scheme).Build(), tt.tls) + routing := &controllerv1alpha1.DevWorkspaceRouting{ + Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{ + Endpoints: endpoints, + PodSelector: workspaceMeta.PodSelector, + }, + } + + objects, err := solver.GetSpecObjects(routing, workspaceMeta) + if !assert.NoError(t, err) { + return + } + + if !tt.tls { + assert.Empty(t, objects.PodAdditions.Volumes) + assert.Empty(t, objects.PodAdditions.VolumeMounts) + return + } + + assert.Len(t, objects.Services, 3) + for _, service := range objects.Services { + assert.Equal(t, service.Name, service.Annotations[serviceServingCertAnnot]) + } + + volumeSecrets := make(map[string]string, len(objects.PodAdditions.Volumes)) + for _, volume := range objects.PodAdditions.Volumes { + volumeSecrets[volume.Name] = volume.Secret.SecretName + } + + actualMounts := make(map[string]string, len(objects.PodAdditions.VolumeMounts)) + for _, mount := range objects.PodAdditions.VolumeMounts { + if _, exists := actualMounts[mount.MountPath]; exists { + assert.Failf(t, "duplicate serving certificate mount path", "path %q is used more than once", mount.MountPath) + } + actualMounts[mount.MountPath] = volumeSecrets[mount.Name] + } + + assert.Len(t, actualMounts, 3) + assert.Equal(t, tt.expectedMounts, actualMounts) + }) + } +} + +// Generated by Codex. +func TestClusterSolverGetSpecObjectsLongServingCertVolumeNames(t *testing.T) { + longEndpointName := strings.Repeat("a", 39) + derivedShortEndpointName := "b4d5e56e929ba4cda349" + endpointNames := []string{ + strings.Repeat("a", 37), + strings.Repeat("b", 38), + longEndpointName, + strings.Repeat("d", 63), + strings.Repeat("x", 62) + "a", + strings.Repeat("x", 62) + "b", + derivedShortEndpointName, + } + endpoints := make(controllerv1alpha1.EndpointList, len(endpointNames)) + expectedMounts := map[string]string{ + "/var/serving-cert/": common.ServiceName("workspace-id"), + } + for i, name := range endpointNames { + endpoints[i] = controllerv1alpha1.Endpoint{ + Name: name, + TargetPort: 7000 + i, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + } + expectedMounts["/var/serving-cert/"+name+"/"] = name + } + + workspaceMeta := DevWorkspaceMetadata{ + DevWorkspaceId: "workspace-id", + DevWorkspaceName: "workspace", + Namespace: "workspace-namespace", + PodSelector: map[string]string{"workspace": "workspace-id"}, + } + scheme := runtime.NewScheme() + assert.NoError(t, corev1.AddToScheme(scheme)) + solver := NewClusterSolver(fake.NewClientBuilder().WithScheme(scheme).Build(), true) + routing := &controllerv1alpha1.DevWorkspaceRouting{ + Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{ + Endpoints: map[string]controllerv1alpha1.EndpointList{"main": endpoints}, + PodSelector: workspaceMeta.PodSelector, + }, + } + + objects, err := solver.GetSpecObjects(routing, workspaceMeta) + if !assert.NoError(t, err) { + return + } + + assert.Len(t, objects.Services, len(endpointNames)+1) + assert.Len(t, objects.PodAdditions.Volumes, len(endpointNames)+1) + assert.Len(t, objects.PodAdditions.VolumeMounts, len(endpointNames)+1) + + volumeSecrets := make(map[string]string, len(objects.PodAdditions.Volumes)) + volumeNamesBySecret := make(map[string]string, len(objects.PodAdditions.Volumes)) + for _, volume := range objects.PodAdditions.Volumes { + if !assert.NotContains(t, volumeSecrets, volume.Name) { + continue + } + assert.LessOrEqual(t, len(volume.Name), 63) + if !assert.NotNil(t, volume.Secret) { + continue + } + volumeSecrets[volume.Name] = volume.Secret.SecretName + volumeNamesBySecret[volume.Secret.SecretName] = volume.Name + } + assert.Equal(t, "devworkspace-serving-cert-"+endpointNames[0], volumeNamesBySecret[endpointNames[0]]) + assert.Equal(t, "devworkspace-cert-b4d5e56e929ba4cda349", volumeNamesBySecret[longEndpointName]) + assert.Equal(t, "devworkspace-serving-cert-"+derivedShortEndpointName, volumeNamesBySecret[derivedShortEndpointName]) + + actualMounts := make(map[string]string, len(objects.PodAdditions.VolumeMounts)) + for _, mount := range objects.PodAdditions.VolumeMounts { + if !assert.NotContains(t, actualMounts, mount.MountPath) { + continue + } + secretName, exists := volumeSecrets[mount.Name] + if !assert.True(t, exists, "mount %q references volume %q", mount.MountPath, mount.Name) { + continue + } + actualMounts[mount.MountPath] = secretName + } + assert.Equal(t, expectedMounts, actualMounts) +} diff --git a/controllers/controller/devworkspacerouting/solvers/common.go b/controllers/controller/devworkspacerouting/solvers/common.go index b55367ad8..ac75d1485 100644 --- a/controllers/controller/devworkspacerouting/solvers/common.go +++ b/controllers/controller/devworkspacerouting/solvers/common.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -16,28 +16,37 @@ package solvers import ( - controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" - "github.com/devfile/devworkspace-operator/pkg/common" - "github.com/devfile/devworkspace-operator/pkg/constants" + "context" + "fmt" routeV1 "github.com/openshift/api/route/v1" corev1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" + "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/intstr" "k8s.io/utils/pointer" + "sigs.k8s.io/controller-runtime/pkg/client" + + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/pkg/common" + "github.com/devfile/devworkspace-operator/pkg/constants" ) type DevWorkspaceMetadata struct { - DevWorkspaceId string - Namespace string - PodSelector map[string]string + DevWorkspaceId string + DevWorkspaceName string + DevWorkspaceRoutingUID types.UID + Namespace string + PodSelector map[string]string } // GetDiscoverableServicesForEndpoints converts the endpoint list into a set of services, each corresponding to a single discoverable // endpoint from the list. Endpoints with the NoneEndpointExposure are ignored. -func GetDiscoverableServicesForEndpoints(endpoints map[string]controllerv1alpha1.EndpointList, meta DevWorkspaceMetadata) []corev1.Service { +func GetDiscoverableServicesForEndpoints(endpoints map[string]controllerv1alpha1.EndpointList, meta DevWorkspaceMetadata, cl client.Client) ([]corev1.Service, error) { var services []corev1.Service + seenServiceNames := map[string]bool{} for _, machineEndpoints := range endpoints { for _, endpoint := range machineEndpoints { if endpoint.Exposure == controllerv1alpha1.NoneEndpointExposure { @@ -45,21 +54,49 @@ func GetDiscoverableServicesForEndpoints(endpoints map[string]controllerv1alpha1 } if endpoint.Attributes.GetBoolean(string(controllerv1alpha1.DiscoverableAttribute), nil) { - // Create service with name matching endpoint - // TODO: This could cause a reconcile conflict if multiple workspaces define the same discoverable endpoint - // Also endpoint names may not be valid as service names + serviceName := common.EndpointName(endpoint.Name) + if serviceName == common.ServiceName(meta.DevWorkspaceId) { + return nil, &RoutingInvalid{Reason: fmt.Sprintf("discoverable endpoint '%s' uses reserved Service name '%s'", endpoint.Name, serviceName)} + } + // Two endpoints on different containers can sanitize to the same Service name (e.g. "my--endpoint" and + // "my-endpoint"). Neither Service exists on the cluster yet at this point, so the existing-service check + // below can't catch this; track names claimed within this call instead. + if seenServiceNames[serviceName] { + return nil, &DuplicateEndpointError{EndpointName: endpoint.Name} + } + seenServiceNames[serviceName] = true + + existingService := &corev1.Service{} + err := cl.Get(context.TODO(), client.ObjectKey{Name: serviceName, Namespace: meta.Namespace}, existingService) + if err != nil { + if !errors.IsNotFound(err) { + return nil, err + } + } else { + controllerOwner := metav1.GetControllerOf(existingService) + ownedByRouting := controllerOwner != nil && controllerOwner.Kind == "DevWorkspaceRouting" && + meta.DevWorkspaceRoutingUID != "" && controllerOwner.UID == meta.DevWorkspaceRoutingUID + if !ownedByRouting && existingService.Labels[constants.DevWorkspaceIDLabel] != meta.DevWorkspaceId { + return nil, &ServiceConflictError{ + EndpointName: endpoint.Name, + WorkspaceName: existingService.Labels[constants.DevWorkspaceNameLabel], + } + } + } + servicePort := corev1.ServicePort{ - Name: common.EndpointName(endpoint.Name), + Name: serviceName, Protocol: corev1.ProtocolTCP, Port: int32(endpoint.TargetPort), TargetPort: intstr.FromInt(endpoint.TargetPort), } services = append(services, corev1.Service{ ObjectMeta: metav1.ObjectMeta{ - Name: common.EndpointName(endpoint.Name), + Name: serviceName, Namespace: meta.Namespace, Labels: map[string]string{ - constants.DevWorkspaceIDLabel: meta.DevWorkspaceId, + constants.DevWorkspaceIDLabel: meta.DevWorkspaceId, + constants.DevWorkspaceNameLabel: meta.DevWorkspaceName, }, Annotations: map[string]string{ constants.DevWorkspaceDiscoverableServiceAnnotation: "true", @@ -74,7 +111,7 @@ func GetDiscoverableServicesForEndpoints(endpoints map[string]controllerv1alpha1 } } } - return services + return services, nil } // GetServiceForEndpoints returns a single service that exposes all endpoints of given exposure types, possibly also including the discoverable types. diff --git a/controllers/controller/devworkspacerouting/solvers/common_test.go b/controllers/controller/devworkspacerouting/solvers/common_test.go new file mode 100644 index 000000000..80573bd31 --- /dev/null +++ b/controllers/controller/devworkspacerouting/solvers/common_test.go @@ -0,0 +1,265 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 solvers + +import ( + "testing" + + "github.com/stretchr/testify/assert" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/pkg/constants" +) + +func TestGetDiscoverableServicesForEndpoints(t *testing.T) { + scheme := runtime.NewScheme() + _ = corev1.AddToScheme(scheme) + _ = controllerv1alpha1.AddToScheme(scheme) + + discoverableEndpoint := controllerv1alpha1.Endpoint{ + Name: "test-endpoint", + TargetPort: 8080, + Exposure: controllerv1alpha1.PublicEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + } + endpoints := map[string]controllerv1alpha1.EndpointList{ + "machine1": {discoverableEndpoint}, + } + + meta := DevWorkspaceMetadata{ + DevWorkspaceId: "current-workspace-id", + DevWorkspaceName: "current-workspace", + Namespace: "test-namespace", + } + + tests := []struct { + name string + existing []runtime.Object + expectErr bool + expectErrType error + expectMsg string + }{ + { + name: "No existing service", + existing: []runtime.Object{}, + expectErr: false, + }, + { + name: "Existing service with different owner in same namespace", + existing: []runtime.Object{ + &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-endpoint", + Namespace: "test-namespace", + Labels: map[string]string{ + constants.DevWorkspaceIDLabel: "other-workspace-id", + constants.DevWorkspaceNameLabel: "other-workspace", + }, + }, + }, + }, + expectErr: true, + expectErrType: &ServiceConflictError{}, + expectMsg: "discoverable endpoint 'test-endpoint' is already in use by workspace 'other-workspace'", + }, + { + name: "Existing service with same owner (reconciliation)", + existing: []runtime.Object{ + &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-endpoint", + Namespace: "test-namespace", + Labels: map[string]string{ + constants.DevWorkspaceIDLabel: "current-workspace-id", + constants.DevWorkspaceNameLabel: "current-workspace", + }, + }, + }, + }, + expectErr: false, + }, + { + name: "Service with same name in different namespace (should not conflict)", + existing: []runtime.Object{ + &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-endpoint", + Namespace: "other-namespace", + Labels: map[string]string{ + constants.DevWorkspaceIDLabel: "other-workspace-id", + constants.DevWorkspaceNameLabel: "other-workspace", + }, + }, + }, + }, + expectErr: false, + }, + { + name: "Service without workspace ID label", + existing: []runtime.Object{ + &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-endpoint", + Namespace: "test-namespace", + Labels: map[string]string{}, + }, + }, + }, + expectErr: true, + expectErrType: &ServiceConflictError{}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + fakeClient := fake.NewClientBuilder().WithScheme(scheme).WithRuntimeObjects(tt.existing...).Build() + _, err := GetDiscoverableServicesForEndpoints(endpoints, meta, fakeClient) + + if tt.expectErr { + assert.Error(t, err, "Expected an error but got none") + if tt.expectErrType != nil { + assert.IsType(t, tt.expectErrType, err, "Error is of unexpected type") + } + if tt.expectMsg != "" { + assert.Contains(t, err.Error(), tt.expectMsg) + } + } else { + assert.NoError(t, err, "Got unexpected error") + } + }) + } +} + +func TestGetDiscoverableServicesForEndpoints_MultipleEndpoints(t *testing.T) { + scheme := runtime.NewScheme() + _ = corev1.AddToScheme(scheme) + _ = controllerv1alpha1.AddToScheme(scheme) + + endpoints := map[string]controllerv1alpha1.EndpointList{ + "machine1": { + { + Name: "postgresql", + TargetPort: 5432, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }, + { + Name: "redis", + TargetPort: 6379, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }, + { + Name: "http", + TargetPort: 8080, + Exposure: controllerv1alpha1.PublicEndpointExposure, + // Not discoverable + }, + }, + } + + meta := DevWorkspaceMetadata{ + DevWorkspaceId: "current-workspace-id", + DevWorkspaceName: "current-workspace", + Namespace: "test-namespace", + PodSelector: map[string]string{ + constants.DevWorkspaceIDLabel: "current-workspace", + }, + } + + t.Run("Multiple discoverable endpoints without conflicts", func(t *testing.T) { + fakeClient := fake.NewClientBuilder().WithScheme(scheme).Build() + services, err := GetDiscoverableServicesForEndpoints(endpoints, meta, fakeClient) + assert.NoError(t, err) + assert.Len(t, services, 2, "Should create 2 discoverable services (postgresql and redis, not http)") + + serviceNames := make(map[string]bool) + for _, svc := range services { + serviceNames[svc.Name] = true + } + assert.True(t, serviceNames["postgresql"], "Should have postgresql service") + assert.True(t, serviceNames["redis"], "Should have redis service") + assert.False(t, serviceNames["http"], "Should not have http service (not discoverable)") + }) + + t.Run("Conflict on one of multiple endpoints", func(t *testing.T) { + existingService := &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: "postgresql", + Namespace: "test-namespace", + Labels: map[string]string{ + constants.DevWorkspaceIDLabel: "other-workspace-id", + constants.DevWorkspaceNameLabel: "other-workspace", + }, + }, + } + fakeClient := fake.NewClientBuilder().WithScheme(scheme).WithRuntimeObjects(existingService).Build() + _, err := GetDiscoverableServicesForEndpoints(endpoints, meta, fakeClient) + assert.Error(t, err, "Should error when one endpoint conflicts") + assert.IsType(t, &ServiceConflictError{}, err) + assert.Contains(t, err.Error(), "postgresql") + }) +} + +// Two containers in the same workspace can declare discoverable endpoints named "my--endpoint" and "my-endpoint". +// These distinct, CRD-valid names pass devfilevalidation.ValidateComponents' raw-name uniqueness check but +// collide after common.EndpointName collapses repeated hyphens to the same Service name. +func TestGetDiscoverableServicesForEndpoints_SameWorkspaceDuplicateName(t *testing.T) { + scheme := runtime.NewScheme() + _ = corev1.AddToScheme(scheme) + _ = controllerv1alpha1.AddToScheme(scheme) + + // Two different containers ("machine1", "machine2") each expose their own discoverable endpoint with a + // different raw name, but both sanitize to the same Service name - and they target different ports. + endpoints := map[string]controllerv1alpha1.EndpointList{ + "machine1": { + { + Name: "my--endpoint", + TargetPort: 5000, + Exposure: controllerv1alpha1.PublicEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }, + }, + "machine2": { + { + Name: "my-endpoint", + TargetPort: 6000, + Exposure: controllerv1alpha1.PublicEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }, + }, + } + + meta := DevWorkspaceMetadata{ + DevWorkspaceId: "current-workspace-id", + DevWorkspaceName: "current-workspace", + Namespace: "test-namespace", + } + + fakeClient := fake.NewClientBuilder().WithScheme(scheme).Build() + services, err := GetDiscoverableServicesForEndpoints(endpoints, meta, fakeClient) + + assert.Error(t, err, "Expected a conflict error for two same-workspace endpoints sharing a Service name") + assert.IsType(t, &DuplicateEndpointError{}, err, "Error should be a DuplicateEndpointError") + assert.Nil(t, services, "Expected no services to be produced on conflict") +} diff --git a/controllers/controller/devworkspacerouting/solvers/errors.go b/controllers/controller/devworkspacerouting/solvers/errors.go index ffcdd7b7b..386573e00 100644 --- a/controllers/controller/devworkspacerouting/solvers/errors.go +++ b/controllers/controller/devworkspacerouting/solvers/errors.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -17,11 +17,30 @@ package solvers import ( "errors" + "fmt" "time" + + "github.com/devfile/devworkspace-operator/pkg/provision/sync" ) var _ error = (*RoutingNotReady)(nil) var _ error = (*RoutingInvalid)(nil) +var _ error = (*ServiceConflictError)(nil) +var _ error = (*DuplicateEndpointError)(nil) + +// ServiceConflictError is returned when a discoverable endpoint has a name that is already in use by +// another DevWorkspace's service. +type ServiceConflictError = sync.ServiceConflictError + +// DuplicateEndpointError is returned when a single DevWorkspace declares the same discoverable endpoint name - or +// two endpoint names that sanitize to the same Service name - on more than one container. +type DuplicateEndpointError struct { + EndpointName string +} + +func (e *DuplicateEndpointError) Error() string { + return fmt.Sprintf("discoverable endpoint '%s' is declared by more than one component in this workspace", e.EndpointName) +} // RoutingNotSupported is used by the solvers when they supported the routingclass of the workspace they've been asked to route var RoutingNotSupported = errors.New("routingclass not supported by this controller") diff --git a/controllers/controller/devworkspacerouting/solvers/service_names_test.go b/controllers/controller/devworkspacerouting/solvers/service_names_test.go new file mode 100644 index 000000000..f7dc286bc --- /dev/null +++ b/controllers/controller/devworkspacerouting/solvers/service_names_test.go @@ -0,0 +1,68 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 solvers + +import ( + "testing" + + "github.com/stretchr/testify/assert" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" +) + +// Generated by Codex. +func TestClusterSolverRejectsDiscoverableEndpointUsingWorkspaceServiceName(t *testing.T) { + for _, tt := range []struct { + name string + tls bool + }{ + {name: "plaintext", tls: false}, + {name: "TLS", tls: true}, + } { + t.Run(tt.name, func(t *testing.T) { + scheme := runtime.NewScheme() + assert.NoError(t, corev1.AddToScheme(scheme)) + solver := NewClusterSolver(fake.NewClientBuilder().WithScheme(scheme).Build(), tt.tls) + const reservedServiceName = "mine-service" + routing := &controllerv1alpha1.DevWorkspaceRouting{ + Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{ + Endpoints: map[string]controllerv1alpha1.EndpointList{ + "machine": {{ + Name: reservedServiceName, + TargetPort: 8080, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }}, + }, + }, + } + objects, err := solver.GetSpecObjects(routing, DevWorkspaceMetadata{ + DevWorkspaceId: "mine", + Namespace: "test-namespace", + }) + + var invalid *RoutingInvalid + assert.ErrorAs(t, err, &invalid) + if assert.NotNil(t, invalid) { + assert.Contains(t, invalid.Reason, "discoverable endpoint") + assert.Contains(t, invalid.Reason, reservedServiceName) + } + assert.Empty(t, objects.Services) + }) + } +} diff --git a/controllers/controller/devworkspacerouting/solvers/solver.go b/controllers/controller/devworkspacerouting/solvers/solver.go index c5c5a5c87..fce842c35 100644 --- a/controllers/controller/devworkspacerouting/solvers/solver.go +++ b/controllers/controller/devworkspacerouting/solvers/solver.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -100,18 +100,19 @@ func (_ *SolverGetter) HasSolver(routingClass controllerv1alpha1.DevWorkspaceRou } } -func (_ *SolverGetter) GetSolver(_ client.Client, routingClass controllerv1alpha1.DevWorkspaceRoutingClass) (RoutingSolver, error) { +func (_ *SolverGetter) GetSolver(client client.Client, routingClass controllerv1alpha1.DevWorkspaceRoutingClass) (RoutingSolver, error) { isOpenShift := infrastructure.IsOpenShift() + switch routingClass { case controllerv1alpha1.DevWorkspaceRoutingBasic: - return &BasicSolver{}, nil + return NewBasicSolver(client), nil case controllerv1alpha1.DevWorkspaceRoutingCluster: - return &ClusterSolver{}, nil + return NewClusterSolver(client, false), nil case controllerv1alpha1.DevWorkspaceRoutingClusterTLS, controllerv1alpha1.DevWorkspaceRoutingWebTerminal: if !isOpenShift { return nil, fmt.Errorf("routing class %s only supported on OpenShift", routingClass) } - return &ClusterSolver{TLS: true}, nil + return NewClusterSolver(client, true), nil default: return nil, RoutingNotSupported } diff --git a/controllers/controller/devworkspacerouting/sync_services.go b/controllers/controller/devworkspacerouting/sync_services.go index 41b8925c7..9a299e655 100644 --- a/controllers/controller/devworkspacerouting/sync_services.go +++ b/controllers/controller/devworkspacerouting/sync_services.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -19,12 +19,14 @@ import ( "context" "fmt" - "github.com/devfile/devworkspace-operator/pkg/constants" - "github.com/devfile/devworkspace-operator/pkg/provision/sync" corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "sigs.k8s.io/controller-runtime/pkg/client" + "github.com/devfile/devworkspace-operator/pkg/constants" + "github.com/devfile/devworkspace-operator/pkg/provision/sync" + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" ) @@ -38,7 +40,10 @@ func (r *DevWorkspaceRoutingReconciler) syncServices(routing *controllerv1alpha1 toDelete := getServicesToDelete(clusterServices, specServices) for _, service := range toDelete { - err := r.Delete(context.TODO(), &service) + if !metav1.IsControlledBy(&service, routing) { + continue + } + err := r.Delete(context.TODO(), &service, client.Preconditions{UID: &service.UID, ResourceVersion: &service.ResourceVersion}) if err != nil { return false, nil, err } @@ -46,10 +51,11 @@ func (r *DevWorkspaceRoutingReconciler) syncServices(routing *controllerv1alpha1 } clusterAPI := sync.ClusterAPI{ - Client: r.Client, - Scheme: r.Scheme, - Logger: r.Log.WithValues("Request.Namespace", routing.Namespace, "Request.Name", routing.Name), - Ctx: context.TODO(), + Client: r.Client, + NonCachingClient: r.NonCachingClient, + Scheme: r.Scheme, + Logger: r.Log.WithValues("Request.Namespace", routing.Namespace, "Request.Name", routing.Name), + Ctx: context.TODO(), } var updatedClusterServices []corev1.Service diff --git a/controllers/controller/devworkspacerouting/sync_services_test.go b/controllers/controller/devworkspacerouting/sync_services_test.go new file mode 100644 index 000000000..2615370c1 --- /dev/null +++ b/controllers/controller/devworkspacerouting/sync_services_test.go @@ -0,0 +1,121 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 devworkspacerouting + +import ( + "context" + "testing" + + "github.com/go-logr/logr" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/utils/ptr" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/controllers/controller/devworkspacerouting/solvers" + "github.com/devfile/devworkspace-operator/pkg/constants" + "github.com/devfile/devworkspace-operator/pkg/infrastructure" +) + +// Generated by Codex +func TestReconcileMarksSynchronizationConflictFailed(t *testing.T) { + infrastructure.InitializeForTesting(infrastructure.Kubernetes) + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + if err := controllerv1alpha1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + routing := &controllerv1alpha1.DevWorkspaceRouting{ObjectMeta: metav1.ObjectMeta{Name: "routing-mine", Namespace: "ns", UID: "mine-routing"}, Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{DevWorkspaceId: "mine", RoutingClass: controllerv1alpha1.DevWorkspaceRoutingCluster, Endpoints: map[string]controllerv1alpha1.EndpointList{"main": {{Name: "api", TargetPort: 8080, Exposure: controllerv1alpha1.InternalEndpointExposure, Attributes: controllerv1alpha1.Attributes{}.PutBoolean("discoverable", true)}}}}} + foreign := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "api", Namespace: "ns", UID: "foreign-service"}, Spec: corev1.ServiceSpec{Selector: map[string]string{"app": "foreign"}}} + backing := fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(routing).WithObjects(routing, foreign).Build() + r := DevWorkspaceRoutingReconciler{Client: serviceLabelFilteredClient{backing}, NonCachingClient: backing, Scheme: scheme, Log: logr.Discard(), SolverGetter: &solvers.SolverGetter{}} + _, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(routing)}) + if err != nil { + t.Fatal(err) + } + actualRouting := &controllerv1alpha1.DevWorkspaceRouting{} + if err := backing.Get(context.Background(), client.ObjectKeyFromObject(routing), actualRouting); err != nil { + t.Fatal(err) + } + if actualRouting.Status.Phase != controllerv1alpha1.RoutingFailed { + t.Fatalf("want permanent conflict, got %s: %s", actualRouting.Status.Phase, actualRouting.Status.Message) + } + actual := &corev1.Service{} + if err := backing.Get(context.Background(), client.ObjectKeyFromObject(foreign), actual); err != nil { + t.Fatalf("foreign Service deleted: %v", err) + } + if actual.Spec.Selector["app"] != "foreign" { + t.Fatal("foreign Service modified") + } +} + +// Generated by Codex +func TestSyncServicesDoesNotDeleteForeignServiceWithWorkspaceLabel(t *testing.T) { + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + foreign := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "unused", Namespace: "ns", UID: "foreign-service", Labels: map[string]string{constants.DevWorkspaceIDLabel: "mine"}, OwnerReferences: []metav1.OwnerReference{{Kind: "DevWorkspaceRouting", UID: "other-routing", Controller: ptr.To(true)}}}} + cl := fake.NewClientBuilder().WithScheme(scheme).WithObjects(foreign).Build() + r := DevWorkspaceRoutingReconciler{Client: cl, Scheme: scheme, Log: logr.Discard()} + routing := &controllerv1alpha1.DevWorkspaceRouting{ObjectMeta: metav1.ObjectMeta{Name: "routing-mine", Namespace: "ns", UID: "mine-routing"}, Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{DevWorkspaceId: "mine"}} + _, _, err := r.syncServices(routing, nil) + if err != nil { + t.Fatal(err) + } + if err := cl.Get(context.Background(), client.ObjectKeyFromObject(foreign), &corev1.Service{}); err != nil { + t.Fatalf("foreign Service deleted: %v", err) + } +} + +// Generated by Codex +func TestSyncServicesDeletesOwnedService(t *testing.T) { + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + owned := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "unused", Namespace: "ns", UID: "owned-service", Labels: map[string]string{constants.DevWorkspaceIDLabel: "mine"}, OwnerReferences: []metav1.OwnerReference{{Kind: "DevWorkspaceRouting", UID: "mine-routing", Controller: ptr.To(true)}}}} + cl := fake.NewClientBuilder().WithScheme(scheme).WithObjects(owned).Build() + r := DevWorkspaceRoutingReconciler{Client: cl, Scheme: scheme, Log: logr.Discard()} + routing := &controllerv1alpha1.DevWorkspaceRouting{ObjectMeta: metav1.ObjectMeta{Name: "routing-mine", Namespace: "ns", UID: "mine-routing"}, Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{DevWorkspaceId: "mine"}} + ready, _, err := r.syncServices(routing, nil) + if err != nil || ready { + t.Fatalf("want pending deletion, got ready=%v err=%v", ready, err) + } + if err := cl.Get(context.Background(), client.ObjectKeyFromObject(owned), &corev1.Service{}); !apierrors.IsNotFound(err) { + t.Fatalf("owned Service still exists: %v", err) + } +} + +// Generated by Codex +type serviceLabelFilteredClient struct{ client.Client } + +func (c serviceLabelFilteredClient) Get(ctx context.Context, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + if err := c.Client.Get(ctx, key, obj, opts...); err != nil { + return err + } + _, labeled := obj.GetLabels()[constants.DevWorkspaceIDLabel] + if _, ok := obj.(*corev1.Service); ok && !labeled { + return apierrors.NewNotFound(schema.GroupResource{Resource: "services"}, key.Name) + } + return nil +} diff --git a/controllers/controller/devworkspacerouting/workspace_name_test.go b/controllers/controller/devworkspacerouting/workspace_name_test.go new file mode 100644 index 000000000..7160925dc --- /dev/null +++ b/controllers/controller/devworkspacerouting/workspace_name_test.go @@ -0,0 +1,249 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 devworkspacerouting + +import ( + "context" + "testing" + + "github.com/devfile/api/v2/pkg/apis/workspaces/v1alpha2" + . "github.com/onsi/gomega" + corev1 "k8s.io/api/core/v1" + networkingv1 "k8s.io/api/networking/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/intstr" + "k8s.io/utils/ptr" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + controllerv1alpha1 "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/controllers/controller/devworkspacerouting/solvers" + "github.com/devfile/devworkspace-operator/pkg/common" + "github.com/devfile/devworkspace-operator/pkg/constants" + "github.com/devfile/devworkspace-operator/pkg/infrastructure" +) + +// Generated by Codex. +func TestReconcileUsesDevWorkspaceOwnerName(t *testing.T) { + infrastructure.InitializeForTesting(infrastructure.Kubernetes) + g := NewWithT(t) + scheme := runtime.NewScheme() + g.Expect(corev1.AddToScheme(scheme)).To(Succeed()) + g.Expect(networkingv1.AddToScheme(scheme)).To(Succeed()) + g.Expect(controllerv1alpha1.AddToScheme(scheme)).To(Succeed()) + + firstRouting := workspaceNameTestRouting("routing-first-route", "workspace-id-first", "real-first-workspace") + secondRouting := workspaceNameTestRouting("routing-second-route", "workspace-id-second", "real-second-workspace") + fakeClient := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&controllerv1alpha1.DevWorkspaceRouting{}). + WithObjects(firstRouting, secondRouting).Build() + reconciler := &DevWorkspaceRoutingReconciler{ + Client: fakeClient, + NonCachingClient: fakeClient, + Log: ctrl.Log, + Scheme: scheme, + SolverGetter: &solvers.SolverGetter{}, + } + ctx := context.Background() + + for range 3 { + result, err := reconciler.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(firstRouting)}) + g.Expect(err).To(Succeed()) + if !result.Requeue { + break + } + } + + service := &corev1.Service{} + serviceKey := client.ObjectKey{Namespace: firstRouting.Namespace, Name: common.EndpointName("api")} + g.Expect(fakeClient.Get(ctx, serviceKey, service)).To(Succeed()) + g.Expect(service.Labels).To(HaveKeyWithValue(constants.DevWorkspaceNameLabel, "real-first-workspace")) + + _, err := reconciler.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(secondRouting)}) + g.Expect(err).To(Succeed()) + storedRouting := &controllerv1alpha1.DevWorkspaceRouting{} + g.Expect(fakeClient.Get(ctx, client.ObjectKeyFromObject(secondRouting), storedRouting)).To(Succeed()) + g.Expect(storedRouting.Status.Phase).To(Equal(controllerv1alpha1.RoutingFailed)) + g.Expect(storedRouting.Status.Message).To(ContainSubstring("real-first-workspace")) + g.Expect(storedRouting.Status.Message).NotTo(ContainSubstring(firstRouting.Name)) +} + +// Generated by Codex. +func TestReconcileRepairsOwnedDiscoverableServiceWithStaleWorkspaceID(t *testing.T) { + tests := []struct { + name string + workspaceIDLabel string + includeIDLabel bool + }{ + {name: "wrong workspace ID label", workspaceIDLabel: "stale-workspace-id", includeIDLabel: true}, + {name: "missing workspace ID label"}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + infrastructure.InitializeForTesting(infrastructure.Kubernetes) + g := NewWithT(t) + scheme := runtime.NewScheme() + g.Expect(corev1.AddToScheme(scheme)).To(Succeed()) + g.Expect(networkingv1.AddToScheme(scheme)).To(Succeed()) + g.Expect(controllerv1alpha1.AddToScheme(scheme)).To(Succeed()) + + routing := workspaceNameTestRouting("routing-stale-id", "workspace-id", "real-workspace") + g.Expect(routing.UID).NotTo(BeEmpty()) + serviceLabels := map[string]string{constants.DevWorkspaceNameLabel: "real-workspace"} + if tt.includeIDLabel { + serviceLabels[constants.DevWorkspaceIDLabel] = tt.workspaceIDLabel + } + controller := true + discoverableService := &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: common.EndpointName("api"), + Namespace: routing.Namespace, + UID: "existing-service-uid", + Labels: serviceLabels, + OwnerReferences: []metav1.OwnerReference{{ + APIVersion: controllerv1alpha1.GroupVersion.String(), + Kind: "DevWorkspaceRouting", + Name: routing.Name, + UID: routing.UID, + Controller: &controller, + }}, + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeClusterIP, + ClusterIP: "10.0.0.9", + Selector: routing.Spec.PodSelector, + Ports: []corev1.ServicePort{{ + Name: common.EndpointName("api"), + Protocol: corev1.ProtocolTCP, + Port: 8080, + TargetPort: intstr.FromInt(8080), + }}, + }, + } + fakeClient := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&controllerv1alpha1.DevWorkspaceRouting{}). + WithObjects(routing, discoverableService).Build() + reconciler := &DevWorkspaceRoutingReconciler{ + Client: fakeClient, + NonCachingClient: fakeClient, + Log: ctrl.Log, + Scheme: scheme, + SolverGetter: &solvers.SolverGetter{}, + } + ctx := context.Background() + + for range 3 { + result, err := reconciler.Reconcile(ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(routing)}) + g.Expect(err).To(Succeed()) + if !result.Requeue { + break + } + } + + storedService := &corev1.Service{} + serviceKey := client.ObjectKeyFromObject(discoverableService) + g.Expect(fakeClient.Get(ctx, serviceKey, storedService)).To(Succeed()) + g.Expect(storedService.Labels).To(HaveKeyWithValue(constants.DevWorkspaceIDLabel, routing.Spec.DevWorkspaceId)) + g.Expect(storedService.UID).To(Equal(discoverableService.UID)) + g.Expect(storedService.Spec.ClusterIP).To(Equal(discoverableService.Spec.ClusterIP)) + storedRouting := &controllerv1alpha1.DevWorkspaceRouting{} + g.Expect(fakeClient.Get(ctx, client.ObjectKeyFromObject(routing), storedRouting)).To(Succeed()) + g.Expect(storedRouting.Status.Phase).NotTo(Equal(controllerv1alpha1.RoutingFailed)) + }) + } +} + +// Generated by Codex. +func TestWorkspaceNameRequiresDevWorkspaceControllerOwner(t *testing.T) { + controller := true + nonController := false + tests := []struct { + name string + ownerReferences []metav1.OwnerReference + want string + }{ + {name: "no owner"}, + { + name: "non-controller owner", + ownerReferences: []metav1.OwnerReference{{ + APIVersion: v1alpha2.SchemeGroupVersion.String(), + Kind: "DevWorkspace", + Name: "ignored-workspace", + Controller: &nonController, + }}, + }, + { + name: "wrong kind", + ownerReferences: []metav1.OwnerReference{{ + APIVersion: v1alpha2.SchemeGroupVersion.String(), + Kind: "Pod", + Name: "ignored-pod", + Controller: &controller, + }}, + }, + { + name: "wrong API group", + ownerReferences: []metav1.OwnerReference{{ + APIVersion: controllerv1alpha1.GroupVersion.String(), + Kind: "DevWorkspace", + Name: "ignored-workspace", + Controller: &controller, + }}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + g := NewWithT(t) + routing := &controllerv1alpha1.DevWorkspaceRouting{ObjectMeta: metav1.ObjectMeta{OwnerReferences: tt.ownerReferences}} + g.Expect(workspaceName(routing)).To(Equal(tt.want)) + }) + } +} + +// Generated by Codex. +func workspaceNameTestRouting(name, workspaceID, workspaceName string) *controllerv1alpha1.DevWorkspaceRouting { + return &controllerv1alpha1.DevWorkspaceRouting{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: "test-namespace", + UID: types.UID(name + "-uid"), + OwnerReferences: []metav1.OwnerReference{{ + APIVersion: v1alpha2.SchemeGroupVersion.String(), + Kind: "DevWorkspace", + Name: workspaceName, + UID: types.UID(workspaceID), + Controller: ptr.To(true), + }}, + }, + Spec: controllerv1alpha1.DevWorkspaceRoutingSpec{ + DevWorkspaceId: workspaceID, + RoutingClass: controllerv1alpha1.DevWorkspaceRoutingCluster, + PodSelector: map[string]string{constants.DevWorkspaceIDLabel: workspaceID}, + Endpoints: map[string]controllerv1alpha1.EndpointList{ + "machine": {{ + Name: "api", + TargetPort: 8080, + Exposure: controllerv1alpha1.InternalEndpointExposure, + Attributes: controllerv1alpha1.Attributes{}. + PutBoolean(string(controllerv1alpha1.DiscoverableAttribute), true), + }}, + }, + }, + } +} diff --git a/pkg/common/naming.go b/pkg/common/naming.go index c7f2504f6..0a0440282 100644 --- a/pkg/common/naming.go +++ b/pkg/common/naming.go @@ -22,6 +22,7 @@ import ( "strings" dw "github.com/devfile/api/v2/pkg/apis/workspaces/v1alpha2" + "github.com/devfile/devworkspace-operator/pkg/constants" ) @@ -32,10 +33,31 @@ func DevWorkspaceRoutingName(workspaceId string) string { } func EndpointName(endpointName string) string { - name := strings.ToLower(endpointName) - name = NonAlphaNumRegexp.ReplaceAllString(name, "-") - name = strings.Trim(name, "-") - return name + // Lowercase ASCII letters, digits, and single internal hyphens need no normalization. + // Return these names unchanged to avoid regexp work and allocations during admission scans. + for i := 0; i < len(endpointName); i++ { + c := endpointName[i] + switch { + case c >= 'a' && c <= 'z' || c >= '0' && c <= '9': + case c == '-' && i > 0 && i < len(endpointName)-1 && endpointName[i-1] != '-': + default: + // Generated by Codex. Collapse non-alphanumeric runs without regexp allocations. + separator := false + name := strings.Map(func(c rune) rune { + if c >= 'a' && c <= 'z' || c >= '0' && c <= '9' { + separator = false + return c + } + if separator { + return -1 + } + separator = true + return '-' + }, strings.ToLower(endpointName)) + return strings.Trim(name, "-") + } + } + return endpointName } func PortName(endpoint dw.Endpoint) string { @@ -101,7 +123,13 @@ func DeploymentName(workspaceId string) string { } func ServingCertVolumeName(serviceName string) string { - return fmt.Sprintf("devworkspace-serving-cert-%s", serviceName) + volumeName := fmt.Sprintf("devworkspace-serving-cert-%s", serviceName) + if len(volumeName) <= 63 { + return volumeName + } + // Generated by Codex. Keep hashes apart from unchanged names derived from short services. + hash := sha256.Sum256([]byte(serviceName)) + return fmt.Sprintf("devworkspace-cert-%x", hash[:10]) } func PVCCleanupJobName(workspaceId string) string { diff --git a/pkg/common/naming_test.go b/pkg/common/naming_test.go index d0c493497..90584eb94 100644 --- a/pkg/common/naming_test.go +++ b/pkg/common/naming_test.go @@ -16,6 +16,8 @@ package common import ( + "regexp" + "strings" "testing" "github.com/stretchr/testify/assert" @@ -78,3 +80,78 @@ func TestSanitizeVolumeName(t *testing.T) { }) } } + +func TestEndpointName(t *testing.T) { + originalRegexp := regexp.MustCompile(`[^a-z0-9]+`) + for _, test := range []struct { + name string + want string + }{ + {"", ""}, + {"a", "a"}, + {"0", "0"}, + {"endpoint123", "endpoint123"}, + {"my-endpoint", "my-endpoint"}, + {"a-0-z-9", "a-0-z-9"}, + {strings.Repeat("a", 1000), strings.Repeat("a", 1000)}, + {"my--endpoint", "my-endpoint"}, + {"my---endpoint", "my-endpoint"}, + {"endpoint--0--0", "endpoint-0-0"}, + {"a--é--b", "a-b"}, + {"a\xff--b", "a-b"}, + {"-endpoint", "endpoint"}, + {"endpoint-", "endpoint"}, + {"--endpoint--", "endpoint"}, + {"-", ""}, + {"---", ""}, + {"My-Endpoint", "my-endpoint"}, + {"ENDPOINT", "endpoint"}, + {"my._ /endpoint!?", "my-endpoint"}, + {"!?", ""}, + {"é-endpoint-界", "endpoint"}, + {"端点", ""}, + {"Kelvin", "kelvin"}, + {"a\x00b", "a-b"}, + } { + t.Run(test.name, func(t *testing.T) { + original := strings.Trim(originalRegexp.ReplaceAllString(strings.ToLower(test.name), "-"), "-") + if got := EndpointName(test.name); got != test.want || got != original { + t.Fatalf("EndpointName(%q) = %q, want %q (original %q)", test.name, got, test.want, original) + } + if test.name == test.want { + var got string + allocs := testing.AllocsPerRun(100, func() { got = EndpointName(test.name) }) + if allocs != 0 || got != test.want { + t.Fatalf("Canonical name %q: got %q with %g allocations, want zero", test.name, got, allocs) + } + } + }) + } +} + +// Generated by Codex. +func TestEndpointNameRepeatedHyphensAllocations(t *testing.T) { + for _, name := range []string{"endpoint--0--0", "--endpoint--", "my---endpoint"} { + t.Run(name, func(t *testing.T) { + var got string + allocs := testing.AllocsPerRun(100, func() { got = EndpointName(name) }) + if allocs > 1 { + t.Fatalf("EndpointName(%q) = %q with %g allocations, want at most one", name, got, allocs) + } + }) + } +} + +// Generated by Codex. +func FuzzEndpointName(f *testing.F) { + for _, name := range []string{"", "endpoint--0--0", "--endpoint--", "A._ /b!?", "a--é--b", "Kelvin", "a\xff--b"} { + f.Add(name) + } + originalRegexp := regexp.MustCompile(`[^a-z0-9]+`) + f.Fuzz(func(t *testing.T, name string) { + want := strings.Trim(originalRegexp.ReplaceAllString(strings.ToLower(name), "-"), "-") + if got := EndpointName(name); got != want { + t.Fatalf("EndpointName(%q) = %q, want %q", name, got, want) + } + }) +} diff --git a/pkg/provision/sync/service.go b/pkg/provision/sync/service.go new file mode 100644 index 000000000..c15c6cd32 --- /dev/null +++ b/pkg/provision/sync/service.go @@ -0,0 +1,31 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 sync + +import "fmt" + +// ServiceConflictError reports a discoverable Service owned by another workspace. +// Generated by Codex +type ServiceConflictError struct { + EndpointName string + WorkspaceName string +} + +// Generated by Codex +func (e *ServiceConflictError) Error() string { + if e.WorkspaceName == "" { + return fmt.Sprintf("discoverable endpoint '%s' is already in use by another workspace", e.EndpointName) + } + return fmt.Sprintf("discoverable endpoint '%s' is already in use by workspace '%s'", e.EndpointName, e.WorkspaceName) +} diff --git a/pkg/provision/sync/service_test.go b/pkg/provision/sync/service_test.go new file mode 100644 index 000000000..29887709c --- /dev/null +++ b/pkg/provision/sync/service_test.go @@ -0,0 +1,204 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 sync + +import ( + "context" + "errors" + "reflect" + "testing" + + "github.com/go-logr/logr" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/utils/ptr" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + "github.com/devfile/devworkspace-operator/pkg/constants" +) + +// Generated by Codex +func TestSyncDiscoverableServicePreservesForeignService(t *testing.T) { + for _, scenario := range []string{"hidden by cache", "created after precheck", "matching label with foreign owner", "created during synchronization"} { + t.Run(scenario, func(t *testing.T) { + ctx := context.Background() + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + foreign := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "api", Namespace: "ns", UID: "foreign-service"}, Spec: corev1.ServiceSpec{Selector: map[string]string{"app": "foreign"}, Ports: []corev1.ServicePort{{Port: 8080}}}} + desired := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "api", Namespace: "ns", Labels: map[string]string{constants.DevWorkspaceIDLabel: "mine"}, Annotations: map[string]string{constants.DevWorkspaceDiscoverableServiceAnnotation: "true"}, OwnerReferences: []metav1.OwnerReference{{APIVersion: "controller.devfile.io/v1alpha1", Kind: "DevWorkspaceRouting", Name: "routing-mine", UID: "mine-routing", Controller: ptr.To(true)}}}, Spec: corev1.ServiceSpec{Selector: map[string]string{"app": "mine"}, Ports: []corev1.ServicePort{{Port: 8080}}}} + if scenario != "hidden by cache" { + foreign.Labels = map[string]string{constants.DevWorkspaceIDLabel: "other"} + } + if scenario == "matching label with foreign owner" { + foreign.Labels[constants.DevWorkspaceIDLabel] = "mine" + foreign.OwnerReferences = []metav1.OwnerReference{{Kind: "DevWorkspaceRouting", UID: "other-routing", Controller: ptr.To(true)}} + } + backing := fake.NewClientBuilder().WithScheme(scheme).Build() + var cached client.Client = backing + if scenario == "hidden by cache" { + cached = filteredServiceClient{backing} + } + api := ClusterAPI{Client: cached, NonCachingClient: backing, Ctx: ctx, Logger: logr.Discard()} + if scenario == "created during synchronization" { + api.Client = serviceCreateRaceClient{Client: backing, foreign: foreign} + } else if err := backing.Create(ctx, foreign); err != nil { + t.Fatal(err) + } + _, syncErr := SyncObjectWithCluster(desired, api) + if syncErr == nil { + t.Fatal("expected conflict or retry") + } + actual := &corev1.Service{} + if err := backing.Get(ctx, client.ObjectKeyFromObject(foreign), actual); err != nil { + t.Fatalf("foreign Service deleted: %v", err) + } + if !reflect.DeepEqual(actual.Spec, foreign.Spec) || !reflect.DeepEqual(actual.Labels, foreign.Labels) || !reflect.DeepEqual(actual.OwnerReferences, foreign.OwnerReferences) { + t.Fatalf("foreign Service modified: %#v", actual) + } + if scenario == "created during synchronization" { + var retry *NotInSyncError + if !errors.As(syncErr, &retry) || retry.Reason != NeedRetryReason { + t.Fatalf("want retry on create race, got %v", syncErr) + } + _, syncErr = SyncObjectWithCluster(desired, api) + } + var fatal *UnrecoverableSyncError + if !errors.As(syncErr, &fatal) { + t.Fatalf("want permanent Service conflict, got %v", syncErr) + } + }) + } +} + +// Generated by Codex +func TestSyncDiscoverableServiceUpdatesOwnedService(t *testing.T) { + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + existing := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "api", Namespace: "ns", UID: "owned-service", Labels: map[string]string{constants.DevWorkspaceIDLabel: "mine"}, OwnerReferences: []metav1.OwnerReference{{Kind: "DevWorkspaceRouting", UID: "mine-routing", Controller: ptr.To(true)}}}, Spec: corev1.ServiceSpec{ClusterIP: "10.0.0.1", Selector: map[string]string{"app": "old"}}} + desired := existing.DeepCopy() + desired.Annotations = map[string]string{constants.DevWorkspaceDiscoverableServiceAnnotation: "true"} + desired.Spec.ClusterIP = "" + desired.Spec.Selector = map[string]string{"app": "mine"} + cl := fake.NewClientBuilder().WithScheme(scheme).WithObjects(existing).Build() + _, err := SyncObjectWithCluster(desired, ClusterAPI{Client: cl, NonCachingClient: cl, Ctx: context.Background(), Logger: logr.Discard()}) + var pending *NotInSyncError + if !errors.As(err, &pending) { + t.Fatalf("want pending update, got %v", err) + } + actual := &corev1.Service{} + if err := cl.Get(context.Background(), client.ObjectKeyFromObject(existing), actual); err != nil { + t.Fatal(err) + } + if actual.Spec.Selector["app"] != "mine" || actual.Spec.ClusterIP != "10.0.0.1" || actual.UID != "owned-service" { + t.Fatalf("owned Service was not updated in place: %#v", actual) + } +} + +// Generated by Codex +func TestSyncOwnedServiceRestoresMissingWorkspaceLabel(t *testing.T) { + for _, discoverable := range []bool{false, true} { + t.Run(map[bool]string{false: "aggregate", true: "discoverable"}[discoverable], func(t *testing.T) { + ctx := context.Background() + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + existing := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "api", Namespace: "ns", UID: "owned-service", OwnerReferences: []metav1.OwnerReference{{Kind: "DevWorkspaceRouting", UID: "mine-routing", Controller: ptr.To(true)}}}, Spec: corev1.ServiceSpec{ClusterIP: "10.0.0.1", Selector: map[string]string{"app": "old"}}} + desired := existing.DeepCopy() + desired.UID = "" + desired.Labels = map[string]string{constants.DevWorkspaceIDLabel: "mine"} + desired.Spec.ClusterIP = "" + desired.Spec.Selector = map[string]string{"app": "mine"} + if discoverable { + desired.Annotations = map[string]string{constants.DevWorkspaceDiscoverableServiceAnnotation: "true"} + } + backing := fake.NewClientBuilder().WithScheme(scheme).WithObjects(existing).Build() + api := ClusterAPI{Client: filteredServiceClient{backing}, NonCachingClient: backing, Ctx: ctx, Logger: logr.Discard()} + _, err := SyncObjectWithCluster(desired, api) + var pending *NotInSyncError + if !errors.As(err, &pending) || pending.Reason != UpdatedObjectReason { + t.Fatalf("want in-place update, got %v", err) + } + actual := &corev1.Service{} + if err := backing.Get(ctx, client.ObjectKeyFromObject(existing), actual); err != nil { + t.Fatal(err) + } + if actual.Labels[constants.DevWorkspaceIDLabel] != "mine" || actual.Spec.Selector["app"] != "mine" || actual.Spec.ClusterIP != "10.0.0.1" || actual.UID != "owned-service" { + t.Fatalf("owned Service was not repaired in place: label=%q selector=%q clusterIP=%q uid=%q", actual.Labels[constants.DevWorkspaceIDLabel], actual.Spec.Selector["app"], actual.Spec.ClusterIP, actual.UID) + } + }) + } +} + +// Generated by Codex +func TestSyncAggregateServiceRejectsHiddenForeignService(t *testing.T) { + ctx := context.Background() + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + foreign := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "mine-service", Namespace: "ns", UID: "foreign-service"}, Spec: corev1.ServiceSpec{Selector: map[string]string{"app": "foreign"}}} + desired := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "mine-service", Namespace: "ns", Labels: map[string]string{constants.DevWorkspaceIDLabel: "mine"}, OwnerReferences: []metav1.OwnerReference{{Kind: "DevWorkspaceRouting", UID: "mine-routing", Controller: ptr.To(true)}}}, Spec: corev1.ServiceSpec{Selector: map[string]string{"app": "mine"}}} + backing := fake.NewClientBuilder().WithScheme(scheme).WithObjects(foreign).Build() + api := ClusterAPI{Client: filteredServiceClient{backing}, NonCachingClient: backing, Ctx: ctx, Logger: logr.Discard()} + _, err := SyncObjectWithCluster(desired, api) + var fatal *UnrecoverableSyncError + var conflict *ServiceConflictError + if !errors.As(err, &fatal) || !errors.As(err, &conflict) { + t.Fatalf("want typed permanent conflict, got %v", err) + } + actual := &corev1.Service{} + if err := backing.Get(ctx, client.ObjectKeyFromObject(foreign), actual); err != nil { + t.Fatalf("foreign Service deleted: %v", err) + } + if actual.UID != "foreign-service" || actual.Spec.Selector["app"] != "foreign" || len(actual.Labels) != 0 { + t.Fatalf("foreign Service modified: %#v", actual) + } +} + +// Generated by Codex +// filteredServiceClient models the production cache's workspace-label filter. +type filteredServiceClient struct{ client.Client } + +func (c filteredServiceClient) Get(ctx context.Context, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + if err := c.Client.Get(ctx, key, obj, opts...); err != nil { + return err + } + _, labeled := obj.GetLabels()[constants.DevWorkspaceIDLabel] + if _, ok := obj.(*corev1.Service); ok && !labeled { + return apierrors.NewNotFound(schema.GroupResource{Resource: "services"}, key.Name) + } + return nil +} + +// Generated by Codex +type serviceCreateRaceClient struct { + client.Client + foreign *corev1.Service +} + +func (c serviceCreateRaceClient) Create(ctx context.Context, obj client.Object, opts ...client.CreateOption) error { + if err := c.Client.Create(ctx, c.foreign); err != nil { + return err + } + return c.Client.Create(ctx, obj, opts...) +} diff --git a/pkg/provision/sync/sync.go b/pkg/provision/sync/sync.go index 25cfd9a36..784c7c69f 100644 --- a/pkg/provision/sync/sync.go +++ b/pkg/provision/sync/sync.go @@ -1,4 +1,4 @@ -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -17,8 +17,6 @@ import ( "fmt" "reflect" - "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" - "github.com/devfile/devworkspace-operator/pkg/config" "github.com/go-logr/logr" "github.com/google/go-cmp/cmp" routev1 "github.com/openshift/api/route/v1" @@ -27,8 +25,13 @@ import ( networkingv1 "k8s.io/api/networking/v1" rbacv1 "k8s.io/api/rbac/v1" k8sErrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" crclient "sigs.k8s.io/controller-runtime/pkg/client" + + "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/pkg/config" + "github.com/devfile/devworkspace-operator/pkg/constants" ) // IsRecognizedObject returns whether the provided object kind is recognized by the sync package to support updating @@ -44,11 +47,23 @@ func IsRecognizedObject(specObj crclient.Object) bool { // as required. If specObj is in sync with the cluster, returns the object as it exists on the cluster. Returns a // NotInSyncError if an update is required, UnrecoverableSyncError if object provided is invalid, or generic error // if an unexpected error is encountered +// Generated by Codex func SyncObjectWithCluster(specObj crclient.Object, api ClusterAPI) (crclient.Object, error) { objType := reflect.TypeOf(specObj).Elem() clusterObj := reflect.New(objType).Interface().(crclient.Object) - err := api.Client.Get(api.Ctx, types.NamespacedName{Name: specObj.GetName(), Namespace: specObj.GetNamespace()}, clusterObj) + reader := api.Client + service, isService := specObj.(*corev1.Service) + var owner *metav1.OwnerReference + if isService { + owner = metav1.GetControllerOf(service) + } + // Routing Services can be hidden by the cache when their workspace label is missing. + routingService := isService && (service.Annotations[constants.DevWorkspaceDiscoverableServiceAnnotation] == "true" || owner != nil && owner.Kind == "DevWorkspaceRouting") + if routingService && api.NonCachingClient != nil { + reader = api.NonCachingClient + } + err := reader.Get(api.Ctx, types.NamespacedName{Name: specObj.GetName(), Namespace: specObj.GetNamespace()}, clusterObj) if err != nil { if k8sErrors.IsNotFound(err) { return nil, createObjectGeneric(specObj, api) @@ -56,6 +71,15 @@ func SyncObjectWithCluster(specObj crclient.Object, api ClusterAPI) (crclient.Ob return nil, err } + if routingService { + existing := clusterObj.(*corev1.Service) + existingOwner := metav1.GetControllerOf(existing) + if owner != nil && (existingOwner == nil || owner.UID != existingOwner.UID) || + owner == nil && (service.Labels[constants.DevWorkspaceIDLabel] == "" || existing.Labels[constants.DevWorkspaceIDLabel] != service.Labels[constants.DevWorkspaceIDLabel]) { + return nil, &UnrecoverableSyncError{&ServiceConflictError{EndpointName: service.Name, WorkspaceName: existing.Labels[constants.DevWorkspaceNameLabel]}} + } + } + if !isMutableObject(specObj) { // TODO: we could still update labels here, or treat a need to update as a fatal error return clusterObj, nil } @@ -141,6 +165,10 @@ func createObjectGeneric(specObj crclient.Object, api ClusterAPI) error { api.Logger.Info("Created object", "kind", reflect.TypeOf(specObj).Elem().String(), "name", specObj.GetName()) return NewNotInSync(specObj, CreatedObjectReason) case k8sErrors.IsAlreadyExists(err): + // Services need a fresh instance to preserve ClusterIP and validate ownership. + if _, ok := specObj.(*corev1.Service); ok { + return NewNotInSync(specObj, NeedRetryReason) + } // Need to try to update the object to address an edge case where removing a labelselector // results in the object not being tracked by the controller's cache. return updateObjectGeneric(specObj, nil, api) diff --git a/pkg/provision/sync/update.go b/pkg/provision/sync/update.go index 3a5b75d7c..064a42e6d 100644 --- a/pkg/provision/sync/update.go +++ b/pkg/provision/sync/update.go @@ -1,4 +1,4 @@ -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -48,6 +48,7 @@ func serviceUpdateFunc(spec, cluster crclient.Object) (crclient.Object, error) { specService := spec.DeepCopyObject().(*corev1.Service) clusterService := cluster.(*corev1.Service) specService.ResourceVersion = clusterService.ResourceVersion + specService.UID = clusterService.UID specService.Spec.ClusterIP = clusterService.Spec.ClusterIP return specService, nil } diff --git a/pkg/webhook/cluster_roles.go b/pkg/webhook/cluster_roles.go index ef93554c2..347cd12c6 100755 --- a/pkg/webhook/cluster_roles.go +++ b/pkg/webhook/cluster_roles.go @@ -18,12 +18,13 @@ package webhook import ( "context" - "github.com/devfile/devworkspace-operator/webhook/server" v1 "k8s.io/api/rbac/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" crclient "sigs.k8s.io/controller-runtime/pkg/client" + + "github.com/devfile/devworkspace-operator/webhook/server" ) func CreateWebhookClusterRole(client crclient.Client, @@ -103,6 +104,19 @@ func getSpecClusterRole() (*v1.ClusterRole, error) { "watch", }, }, + { + APIGroups: []string{ + "workspace.devfile.io", + }, + Resources: []string{ + "devworkspaces", + }, + Verbs: []string{ + "get", + "list", + "watch", + }, + }, { APIGroups: []string{ "authentication.k8s.io", diff --git a/webhook/workspace/config.go b/webhook/workspace/config.go index b84eecc78..4eb6efdcb 100644 --- a/webhook/workspace/config.go +++ b/webhook/workspace/config.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -20,24 +20,27 @@ import ( "errors" "fmt" - "sigs.k8s.io/controller-runtime/pkg/manager" - - "github.com/devfile/devworkspace-operator/pkg/config" - "github.com/devfile/devworkspace-operator/pkg/infrastructure" - "github.com/devfile/devworkspace-operator/webhook/server" - admregv1 "k8s.io/api/admissionregistration/v1" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/types" "sigs.k8s.io/controller-runtime/pkg/client" clientConfig "sigs.k8s.io/controller-runtime/pkg/client/config" + "sigs.k8s.io/controller-runtime/pkg/manager" "sigs.k8s.io/controller-runtime/pkg/webhook" + + "github.com/devfile/devworkspace-operator/pkg/config" + "github.com/devfile/devworkspace-operator/pkg/infrastructure" + "github.com/devfile/devworkspace-operator/webhook/server" + "github.com/devfile/devworkspace-operator/webhook/workspace/handler" ) // Configure configures mutate/validating webhooks that provides exec access into workspace for creator only func Configure(ctx context.Context, mgr manager.Manager) error { log.Info("Configuring devworkspace webhooks") + if err := handler.RegisterDiscoverableEndpointIndex(ctx, mgr.GetFieldIndexer()); err != nil { + return fmt.Errorf("register discoverable endpoint index: %w", err) + } c, err := createClient() if err != nil { return err diff --git a/webhook/workspace/config_test.go b/webhook/workspace/config_test.go new file mode 100644 index 000000000..6f49d07ef --- /dev/null +++ b/webhook/workspace/config_test.go @@ -0,0 +1,45 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 workspace + +import ( + "context" + "errors" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/manager" +) + +// Generated by Codex. +func TestConfigurePropagatesEndpointIndexFailure(t *testing.T) { + t.Setenv("KUBECONFIG", filepath.Join(t.TempDir(), "missing-kubeconfig")) + err := errors.New("cannot register endpoint index") + require.ErrorIs(t, Configure(context.Background(), endpointIndexManager{indexer: failingEndpointIndexer{err: err}}), err) +} + +type endpointIndexManager struct { + manager.Manager + indexer client.FieldIndexer +} + +func (m endpointIndexManager) GetFieldIndexer() client.FieldIndexer { return m.indexer } + +type failingEndpointIndexer struct{ err error } + +func (i failingEndpointIndexer) IndexField(context.Context, client.Object, string, client.IndexerFunc) error { + return i.err +} diff --git a/webhook/workspace/handler/endpoint_cache_test.go b/webhook/workspace/handler/endpoint_cache_test.go new file mode 100644 index 000000000..ad634a1c8 --- /dev/null +++ b/webhook/workspace/handler/endpoint_cache_test.go @@ -0,0 +1,182 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 handler + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "strings" + "testing" + "time" + + dwv2 "github.com/devfile/api/v2/pkg/apis/workspaces/v1alpha2" + "github.com/devfile/api/v2/pkg/attributes" + . "github.com/onsi/gomega" + "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/client-go/rest" + restfake "k8s.io/client-go/rest/fake" + "sigs.k8s.io/controller-runtime/pkg/cache" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// newEndpointCache uses the real informer, index, and cache reader. Only the API +// transport is synthetic; queries and copies do not pass through a fake client. +// Generated by Codex. +func newEndpointCache(tb testing.TB, peers []dwv2.DevWorkspace) (cache.Cache, *watch.RaceFreeFakeWatcher) { + tb.Helper() + scheme := runtime.NewScheme() + require.NoError(tb, dwv2.AddToScheme(scheme)) + mapper := meta.NewDefaultRESTMapper([]schema.GroupVersion{dwv2.SchemeGroupVersion}) + mapper.Add(dwv2.SchemeGroupVersion.WithKind("DevWorkspace"), meta.RESTScopeNamespace) + events := watch.NewRaceFreeFake() + ready := make(chan struct{}) + list := &dwv2.DevWorkspaceList{ + TypeMeta: metav1.TypeMeta{APIVersion: dwv2.SchemeGroupVersion.String(), Kind: "DevWorkspaceList"}, + ListMeta: metav1.ListMeta{ResourceVersion: "1"}, + Items: peers, + } + initial, err := json.Marshal(list) + require.NoError(tb, err) + httpClient := restfake.CreateHTTPClient(func(req *http.Request) (*http.Response, error) { + if req.Method != http.MethodGet || !strings.HasSuffix(req.URL.Path, "/devworkspaces") { + return nil, fmt.Errorf("unexpected cache request: %s %s", req.Method, req.URL) + } + body := io.NopCloser(bytes.NewReader(initial)) + if req.URL.Query().Get("watch") == "true" { + reader, writer := io.Pipe() + body = reader + go func() { + defer writer.Close() + encoder := json.NewEncoder(writer) + // Support client-go's streaming initial list as well as List+Watch. + if req.URL.Query().Get("sendInitialEvents") == "true" { + for i := range peers { + if err := encoder.Encode(metav1.WatchEvent{Type: string(watch.Added), Object: runtime.RawExtension{Object: &peers[i]}}); err != nil { + return + } + } + bookmark := &dwv2.DevWorkspace{ + TypeMeta: metav1.TypeMeta{APIVersion: dwv2.SchemeGroupVersion.String(), Kind: "DevWorkspace"}, + ObjectMeta: metav1.ObjectMeta{ResourceVersion: "1", Annotations: map[string]string{"k8s.io/initial-events-end": "true"}}, + } + if err := encoder.Encode(metav1.WatchEvent{Type: string(watch.Bookmark), Object: runtime.RawExtension{Object: bookmark}}); err != nil { + return + } + } + close(ready) + for { + select { + case <-req.Context().Done(): + return + case event, ok := <-events.ResultChan(): + if !ok { + return + } + if err := encoder.Encode(metav1.WatchEvent{Type: string(event.Type), Object: runtime.RawExtension{Object: event.Object}}); err != nil { + return + } + } + } + }() + } + return &http.Response{StatusCode: http.StatusOK, Header: http.Header{"Content-Type": {"application/json"}}, Body: body}, nil + }) + ctx, cancel := context.WithCancel(context.Background()) + cached, err := cache.New(&rest.Config{Host: "https://endpoint-cache.test"}, cache.Options{ + Scheme: scheme, Mapper: mapper, HTTPClient: httpClient, + }) + require.NoError(tb, err) + require.NoError(tb, RegisterDiscoverableEndpointIndex(ctx, cached)) + stopped := make(chan error, 1) + go func() { stopped <- cached.Start(ctx) }() + tb.Cleanup(func() { + cancel() + events.Stop() + select { + case err := <-stopped: + require.NoError(tb, err) + case <-time.After(10 * time.Second): + tb.Error("endpoint cache did not stop") + } + }) + syncCtx, cancelSync := context.WithTimeout(ctx, 10*time.Second) + defer cancelSync() + require.True(tb, cached.WaitForCacheSync(syncCtx), "endpoint cache did not sync") + select { + case <-ready: + case <-syncCtx.Done(): + tb.Fatal("endpoint watch did not start") + } + return cached, events +} + +// Generated by Codex. +func TestDiscoverableEndpointCacheLifecycle(t *testing.T) { + cached, events := newEndpointCache(t, nil) + g := NewWithT(t) + query := func(namespace, name string) []dwv2.DevWorkspace { + list := &dwv2.DevWorkspaceList{} + require.NoError(t, cached.List(context.Background(), list, client.InNamespace(namespace), + client.MatchingFields{discoverableEndpointIndex: name})) + return list.Items + } + peer := setupWorkspace(t, "peer", "peer-uid", "test-namespace") + peer.ResourceVersion = "2" + peer.Spec.Template.Components[0].Container.Env = []dwv2.EnvVar{{Name: "ORIGINAL", Value: "original"}} + events.Add(peer.DeepCopy()) + g.Eventually(func() []dwv2.DevWorkspace { return query("test-namespace", "test-endpoint") }, 5*time.Second).Should(HaveLen(1)) + otherNamespace := peer.DeepCopy() + otherNamespace.Namespace = "other-namespace" + otherNamespace.UID = "other-uid" + events.Add(otherNamespace) + g.Eventually(func() []dwv2.DevWorkspace { return query("other-namespace", "test-endpoint") }, 5*time.Second).Should(HaveLen(1)) + items := query("test-namespace", "test-endpoint") + require.Len(t, items, 1) + require.Equal(t, peer.UID, items[0].UID) + container := items[0].Spec.Template.Components[0].Container + container.Endpoints[0].Name = "mutated" + container.Endpoints[0].Attributes = container.Endpoints[0].Attributes.PutBoolean("discoverable", false) + container.Env[0].Value = "mutated" + original := query("test-namespace", "test-endpoint")[0].Spec.Template.Components[0].Container + require.Equal(t, "test-endpoint", original.Endpoints[0].Name) + require.True(t, original.Endpoints[0].Attributes.GetBoolean("discoverable", nil)) + require.Equal(t, "original", original.Env[0].Value) + + peer.ResourceVersion = "3" + peer.Spec.Template.Components[0].Container.Endpoints[0].Name = "updated--endpoint" + peer.Spec.Template.Components[0].Container.Endpoints[0].Attributes = attributes.Attributes{}.PutString("discoverable", "true") + events.Modify(peer.DeepCopy()) + g.Eventually(func() []dwv2.DevWorkspace { return query("test-namespace", "updated-endpoint") }, 5*time.Second).Should(HaveLen(1)) + require.Empty(t, query("test-namespace", "test-endpoint")) + peer.ResourceVersion = "4" + peer.Spec.Template.Components[0].Container.Endpoints[0].Exposure = dwv2.NoneEndpointExposure + events.Modify(peer.DeepCopy()) + g.Eventually(func() []dwv2.DevWorkspace { return query("test-namespace", "updated-endpoint") }, 5*time.Second).Should(BeEmpty()) + peer.ResourceVersion = "5" + peer.Spec.Template.Components[0].Container.Endpoints[0].Exposure = dwv2.InternalEndpointExposure + events.Modify(peer.DeepCopy()) + g.Eventually(func() []dwv2.DevWorkspace { return query("test-namespace", "updated-endpoint") }, 5*time.Second).Should(HaveLen(1)) + events.Delete(peer.DeepCopy()) + g.Eventually(func() []dwv2.DevWorkspace { return query("test-namespace", "updated-endpoint") }, 5*time.Second).Should(BeEmpty()) + require.Len(t, query("other-namespace", "test-endpoint"), 1) +} diff --git a/webhook/workspace/handler/endpoint_index_test.go b/webhook/workspace/handler/endpoint_index_test.go new file mode 100644 index 000000000..8bf72e2ff --- /dev/null +++ b/webhook/workspace/handler/endpoint_index_test.go @@ -0,0 +1,122 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 handler + +import ( + "context" + "errors" + "fmt" + "testing" + + dwv2 "github.com/devfile/api/v2/pkg/apis/workspaces/v1alpha2" + "github.com/stretchr/testify/require" + admissionv1 "k8s.io/api/admission/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/webhook/admission" +) + +type endpointIndexRegistration struct { + object client.Object + field string + extract client.IndexerFunc + err error +} + +func (r *endpointIndexRegistration) IndexField(_ context.Context, object client.Object, field string, extract client.IndexerFunc) error { + r.object, r.field, r.extract = object, field, extract + return r.err +} + +// Generated by Codex. +func newIndexedWorkspaceClientBuilder(t *testing.T, scheme *runtime.Scheme) *fake.ClientBuilder { + t.Helper() + registration := &endpointIndexRegistration{} + require.NoError(t, RegisterDiscoverableEndpointIndex(context.Background(), registration)) + return fake.NewClientBuilder().WithScheme(scheme).WithIndex(registration.object, registration.field, registration.extract) +} + +func TestRegisterDiscoverableEndpointIndexFailure(t *testing.T) { + err := errors.New("index unavailable") + require.ErrorIs(t, RegisterDiscoverableEndpointIndex(context.Background(), &endpointIndexRegistration{err: err}), err) +} + +// Generated by Codex. +func TestValidateEndpointsListFailure(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, dwv2.AddToScheme(scheme)) + workspace := setupWorkspace(t, "incoming", "incoming-uid", "test-namespace") + err := errors.New("list unavailable") + handler := &WebhookHandler{Client: endpointListFailureClient{err: err}, Decoder: admission.NewDecoder(scheme)} + conflict, got := handler.validateEndpoints(context.Background(), workspace, discoverableEndpointNames(workspace)) + require.ErrorIs(t, got, err) + require.Nil(t, conflict) + response := handler.ValidateDevfile(context.Background(), admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Create, + Object: newRawExtension(t, workspace), + }}) + require.False(t, response.Allowed) + require.EqualValues(t, 500, response.Result.Code) + require.Contains(t, response.Result.Message, err.Error()) +} + +type endpointListFailureClient struct { + client.Client + err error +} + +func (c endpointListFailureClient) List(context.Context, client.ObjectList, ...client.ListOption) error { + return c.err +} + +// Generated by Codex. +func TestValidateEndpointsUsesUniqueIndexedQueries(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, dwv2.AddToScheme(scheme)) + workspace := setupWorkspace(t, "incoming", "incoming-uid", "test-namespace") + endpoint := workspace.Spec.Template.Components[0].Container.Endpoints[0] + endpoint.Name = "unique--endpoint" + workspace.Spec.Template.Components[0].Container.Endpoints = append( + workspace.Spec.Template.Components[0].Container.Endpoints, endpoint, endpoint) + registration := &endpointIndexRegistration{} + require.NoError(t, RegisterDiscoverableEndpointIndex(context.Background(), registration)) + require.ElementsMatch(t, []string{"test-endpoint", "unique-endpoint"}, registration.extract(workspace)) + indexed := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(workspace).Build() + queries := &endpointQueryClient{Client: indexed} + handler := &WebhookHandler{Client: queries} + conflict, err := handler.validateEndpoints(context.Background(), workspace, discoverableEndpointNames(workspace)) + require.NoError(t, err) + require.Nil(t, conflict) + require.ElementsMatch(t, []string{"test-endpoint", "unique-endpoint"}, queries.names) +} + +type endpointQueryClient struct { + client.Client + names []string +} + +// Generated by Codex. +func (c *endpointQueryClient) List(ctx context.Context, out client.ObjectList, opts ...client.ListOption) error { + options := (&client.ListOptions{}).ApplyOptions(opts) + if options.Namespace != "test-namespace" || options.FieldSelector == nil { + return fmt.Errorf("endpoint queries must be namespace-scoped and indexed") + } + name, exact := options.FieldSelector.RequiresExactMatch("controller.devfile.io/discoverable-endpoint") + if !exact { + return fmt.Errorf("endpoint queries must select an exact name") + } + c.names = append(c.names, name) + return c.Client.List(ctx, out, opts...) +} diff --git a/webhook/workspace/handler/testdata/test-devworkspace.yaml b/webhook/workspace/handler/testdata/test-devworkspace.yaml new file mode 100644 index 000000000..82fc943e2 --- /dev/null +++ b/webhook/workspace/handler/testdata/test-devworkspace.yaml @@ -0,0 +1,16 @@ +apiVersion: workspace.devfile.io/v1alpha2 +kind: DevWorkspace +metadata: + name: test-devworkspace +spec: + started: true + template: + components: + - name: test-component + container: + image: test-image + endpoints: + - name: test-endpoint + targetPort: 8080 + attributes: + discoverable: "true" \ No newline at end of file diff --git a/webhook/workspace/handler/validate.go b/webhook/workspace/handler/validate.go index 2faa33c88..59551077d 100644 --- a/webhook/workspace/handler/validate.go +++ b/webhook/workspace/handler/validate.go @@ -1,5 +1,5 @@ // -// Copyright (c) 2019-2025 Red Hat, Inc. +// Copyright (c) 2019-2026 Red Hat, Inc. // 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 @@ -23,9 +23,30 @@ import ( dwv2 "github.com/devfile/api/v2/pkg/apis/workspaces/v1alpha2" devfilevalidation "github.com/devfile/api/v2/pkg/validation" + admissionv1 "k8s.io/api/admission/v1" + "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/webhook/admission" + + "github.com/devfile/devworkspace-operator/apis/controller/v1alpha1" + "github.com/devfile/devworkspace-operator/controllers/controller/devworkspacerouting/solvers" + "github.com/devfile/devworkspace-operator/pkg/common" ) +const discoverableEndpointIndex = "controller.devfile.io/discoverable-endpoint" + +// RegisterDiscoverableEndpointIndex indexes normalized discoverable Service names in the workspace cache. +// Generated by Codex. +func RegisterDiscoverableEndpointIndex(ctx context.Context, indexer client.FieldIndexer) error { + return indexer.IndexField(ctx, &dwv2.DevWorkspace{}, discoverableEndpointIndex, func(obj client.Object) []string { + names := discoverableEndpointNames(obj.(*dwv2.DevWorkspace)) + values := make([]string, 0, len(names)) + for name := range names { + values = append(values, name) + } + return values + }) +} + func (h *WebhookHandler) ValidateDevfile(ctx context.Context, req admission.Request) admission.Response { wksp := &dwv2.DevWorkspace{} @@ -75,9 +96,115 @@ func (h *WebhookHandler) ValidateDevfile(ctx context.Context, req admission.Requ } } + names := discoverableEndpointNames(wksp) + if h.shouldCheckEndpointConflicts(req, names) { + conflict, err := h.validateEndpoints(ctx, wksp, names) + if err != nil { + return admission.Errored(http.StatusInternalServerError, err) + } + if conflict != nil { + devfileErrors = append(devfileErrors, conflict.Error()) + } + } + if len(devfileErrors) > 0 { return admission.Denied(fmt.Sprintf("\n%s\n", strings.Join(devfileErrors, "\n"))) } return admission.Allowed("No Devfile errors were found") } + +// shouldCheckEndpointConflicts reports whether validateEndpoints needs to run for this request. Newly created +// workspaces with discoverable endpoints need the check, since every discoverable endpoint is new. An empty set +// never needs the check. For updates, a new conflict can only be introduced by a discoverable endpoint name that +// wasn't already present, so the check is skipped unless the new set contains a name that wasn't in the old set. +// This also skips pure removals, which can never introduce a conflict and would otherwise re-reject unrelated edits +// blocked only by a pre-existing conflict on an endpoint the update doesn't touch. +func (h *WebhookHandler) shouldCheckEndpointConflicts(req admission.Request, names map[string]bool) bool { + if len(names) == 0 { + return false + } + + if req.Operation != admissionv1.Update { + return true + } + + oldWorkspace := &dwv2.DevWorkspace{} + if err := h.Decoder.DecodeRaw(req.OldObject, oldWorkspace); err != nil { + return true + } + + oldNames := discoverableEndpointNames(oldWorkspace) + for name := range names { + if !oldNames[name] { + return true + } + } + return false +} + +// discoverableEndpointNames returns the sanitized Service names (see common.EndpointName) of the workspace's +// discoverable endpoints whose exposure is not none. Omitted exposure defaults to public and is included. +// Comparing by the sanitized name, rather than the raw devfile endpoint name, matches the +// actual collision surface: two differently-named endpoints that sanitize to the same value produce the same +// Service name and conflict at reconcile time. +// +// Endpoints contributed via workspace.Spec.Contributions (plugins/parent templates) are invisible here, +// permanently: flatten.ResolveDevWorkspace inlines them only in-memory during reconcile, never back to +// .spec.template, so this field never reflects them, no matter how many times a workspace reconciles. Flattening +// here instead isn't safe either - it requires cluster/network calls unsuitable for a webhook's time budget. +// Such conflicts still surface at reconcile time via Service synchronization, which checks the +// actual Service objects rather than any spec - the same behavior that existed before this check was added. +func discoverableEndpointNames(workspace *dwv2.DevWorkspace) map[string]bool { + discoverableEndpoints := map[string]bool{} + for _, component := range workspace.Spec.Template.Components { + if component.Container != nil { + for _, endpoint := range component.Container.Endpoints { + if endpoint.Exposure != dwv2.NoneEndpointExposure && endpoint.Attributes.GetBoolean(string(v1alpha1.DiscoverableAttribute), nil) { + discoverableEndpoints[common.EndpointName(endpoint.Name)] = true + } + } + } + } + return discoverableEndpoints +} + +// validateEndpoints is best-effort: it reads the current cluster state at admission time, so two concurrent updates +// racing to claim the same endpoint name could both be admitted. The controller-level check in +// Service synchronization remains the authoritative ownership guard against that race. +func (h *WebhookHandler) validateEndpoints(ctx context.Context, workspace *dwv2.DevWorkspace, names map[string]bool) (*solvers.ServiceConflictError, error) { + if len(names) == 0 { + return nil, nil + } + + // Field selectors are ANDed, so query each name separately to find their union. + for name := range names { + workspaceList := &dwv2.DevWorkspaceList{} + if err := h.Client.List(ctx, workspaceList, client.InNamespace(workspace.Namespace), + client.MatchingFields{discoverableEndpointIndex: name}); err != nil { + return nil, err + } + for _, otherWorkspace := range workspaceList.Items { + if otherWorkspace.UID == workspace.UID { + continue + } + for _, component := range otherWorkspace.Spec.Template.Components { + if component.Container != nil { + for _, endpoint := range component.Container.Endpoints { + if endpoint.Exposure != dwv2.NoneEndpointExposure && + endpoint.Attributes.Exists(string(v1alpha1.DiscoverableAttribute)) && + names[common.EndpointName(endpoint.Name)] && + endpoint.Attributes.GetBoolean(string(v1alpha1.DiscoverableAttribute), nil) { + return &solvers.ServiceConflictError{ + EndpointName: endpoint.Name, + WorkspaceName: otherWorkspace.Name, + }, nil + } + } + } + } + } + } + + return nil, nil +} diff --git a/webhook/workspace/handler/validate_test.go b/webhook/workspace/handler/validate_test.go new file mode 100644 index 000000000..8c63711b2 --- /dev/null +++ b/webhook/workspace/handler/validate_test.go @@ -0,0 +1,608 @@ +// Copyright (c) 2019-2026 Red Hat, Inc. +// 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 handler + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "testing" + + dwv2 "github.com/devfile/api/v2/pkg/apis/workspaces/v1alpha2" + "github.com/devfile/api/v2/pkg/attributes" + "github.com/stretchr/testify/assert" + admissionv1 "k8s.io/api/admission/v1" + apiextv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/webhook/admission" + "sigs.k8s.io/yaml" +) + +func newRawExtension(t *testing.T, workspace *dwv2.DevWorkspace) runtime.RawExtension { + bytes, err := json.Marshal(workspace) + assert.NoError(t, err, "Failed to marshal workspace") + return runtime.RawExtension{Raw: bytes} +} + +func loadObjectFromFile(objName string, obj client.Object, filename string) error { + path := filepath.Join("testdata", filename) + bytes, err := os.ReadFile(path) + if err != nil { + return err + } + err = yaml.Unmarshal(bytes, obj) + if err != nil { + return err + } + obj.SetName(objName) + return nil +} + +func setupWorkspace(t *testing.T, name, uid, namespace string) *dwv2.DevWorkspace { + workspace := &dwv2.DevWorkspace{} + err := loadObjectFromFile(name, workspace, "test-devworkspace.yaml") + assert.NoError(t, err, "Failed to load workspace") + workspace.SetUID(types.UID(uid)) + workspace.SetNamespace(namespace) + return workspace +} + +func TestValidateEndpoints(t *testing.T) { + scheme := runtime.NewScheme() + _ = dwv2.AddToScheme(scheme) + + t.Run("Conflict in same namespace", func(t *testing.T) { + // Workspace with a discoverable endpoint in namespace "test-namespace" + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + + // Another workspace with a conflicting discoverable endpoint in the SAME namespace + otherWorkspaceSameNS := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + + // Test for conflict in same namespace + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspaceSameNS).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err, "Did not expect an infrastructure error") + assert.NotNil(t, conflict, "Expected a conflict for workspaces in the same namespace") + assert.Equal(t, "test-endpoint", conflict.EndpointName, "Conflict should be on 'test-endpoint'") + assert.Equal(t, "workspace-2", conflict.WorkspaceName, "Conflict should reference 'workspace-2'") + }) + + t.Run("No conflict in different namespace", func(t *testing.T) { + // Workspace in "test-namespace" + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + + // Another workspace with the same endpoint name but in a DIFFERENT namespace + otherWorkspaceDiffNS := setupWorkspace(t, "workspace-3", "uid-3", "other-namespace") + + // Test no conflict in different namespace (workspace only queries its own namespace) + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspaceDiffNS).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err) + assert.Nil(t, conflict, "Did not expect a conflict for workspaces in different namespaces") + }) + + t.Run("No conflict when endpoint name is different", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + workspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "new-endpoint" + + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err) + assert.Nil(t, conflict, "Did not expect a conflict for different endpoint names") + }) + + t.Run("Conflict detected even when workspace is being deleted", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + + // Workspace being deleted with same endpoint name + deletingWorkspace := setupWorkspace(t, "workspace-deleting", "uid-deleting", "test-namespace") + now := metav1.Now() + deletingWorkspace.DeletionTimestamp = &now + // Add finalizer - required by fake client when setting deletionTimestamp + deletingWorkspace.Finalizers = []string{"test-finalizer"} + + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(deletingWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err, "Did not expect an infrastructure error") + assert.NotNil(t, conflict, "Should detect conflict even with workspace being deleted") + assert.Equal(t, "test-endpoint", conflict.EndpointName, "Conflict should be on 'test-endpoint'") + assert.Equal(t, "workspace-deleting", conflict.WorkspaceName, "Conflict should reference 'workspace-deleting'") + }) + + t.Run("No conflict when workspace has no discoverable endpoints", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + // Remove discoverable attribute + workspace.Spec.Template.Components[0].Container.Endpoints[0].Attributes = nil + + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err) + assert.Nil(t, conflict, "Did not expect a conflict when workspace has no discoverable endpoints") + }) + + t.Run("No conflict when other workspace endpoint is not discoverable", func(t *testing.T) { + // Current workspace has a discoverable endpoint + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + + // Other workspace has endpoint with same name but NOT discoverable + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + otherWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Attributes = nil + + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err) + assert.Nil(t, conflict, "Should not conflict when other workspace's endpoint is not discoverable") + }) + + t.Run("Ignores non-discoverable aliases before a discoverable conflict", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + workspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "my-endpoint" + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + endpoint := otherWorkspace.Spec.Template.Components[0].Container.Endpoints[0] + otherWorkspace.Spec.Template.Components[0].Container.Endpoints = []dwv2.Endpoint{ + {Name: "my---endpoint", TargetPort: 8081, Attributes: nil}, + {Name: "my----endpoint", TargetPort: 8082, Attributes: attributes.Attributes{}}, + {Name: "my-----endpoint", TargetPort: 8083, Attributes: attributes.Attributes{}.PutBoolean("discoverable", false)}, + } + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err) + assert.Nil(t, conflict, "Sanitized names alone must not make non-discoverable endpoints conflict") + + endpoint.Name = "my--endpoint" + otherWorkspace.Spec.Template.Components[0].Container.Endpoints = append(otherWorkspace.Spec.Template.Components[0].Container.Endpoints, endpoint) + assert.NoError(t, fakeClient.Update(context.TODO(), otherWorkspace)) + conflict, err = handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err) + if assert.NotNil(t, conflict, "Must continue scanning past non-discoverable aliases") { + assert.Equal(t, "my--endpoint", conflict.EndpointName, "Report the conflicting raw endpoint name") + assert.Equal(t, "workspace-2", conflict.WorkspaceName) + } + }) + + t.Run("Conflict detected when endpoint names differ but sanitize to the same service name", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + workspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "my--endpoint" + + // Different raw name, but common.EndpointName collapses repeated hyphens in "my--endpoint" to "-". + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + otherWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "my-endpoint" + + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err, "Did not expect an infrastructure error") + assert.NotNil(t, conflict, "Expected a conflict for endpoint names that sanitize to the same service name") + assert.Equal(t, "my-endpoint", conflict.EndpointName, "Conflict should report the other workspace's raw endpoint name") + assert.Equal(t, "workspace-2", conflict.WorkspaceName, "Conflict should reference 'workspace-2'") + }) + + t.Run("Ignores exposure none before an eligible sanitized alias", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + workspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "my-endpoint" + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + ignored := otherWorkspace.Spec.Template.Components[0].Container.Endpoints[0] + ignored.Name = "my---endpoint" + ignored.Exposure = dwv2.NoneEndpointExposure + eligible := ignored + eligible.Name = "my--endpoint" + eligible.Exposure = dwv2.InternalEndpointExposure + eligible.TargetPort = 8081 + otherWorkspace.Spec.Template.Components[0].Container.Endpoints = []dwv2.Endpoint{ignored, eligible} + + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient} + conflict, err := handler.validateEndpoints(context.TODO(), workspace, discoverableEndpointNames(workspace)) + assert.NoError(t, err) + if assert.NotNil(t, conflict, "Must continue scanning past exposure none") { + assert.Equal(t, "my--endpoint", conflict.EndpointName, "Report the eligible raw endpoint name") + assert.Equal(t, "workspace-2", conflict.WorkspaceName) + } + }) + + t.Run("Multiple workspaces in different namespaces can have same endpoint", func(t *testing.T) { + // Workspace 1 in namespace-a + workspace1 := setupWorkspace(t, "workspace-ns-a", "uid-ns-a", "namespace-a") + + // Workspace 2 in namespace-b (will be in the fake client as existing) + workspace2 := setupWorkspace(t, "workspace-ns-b", "uid-ns-b", "namespace-b") + + // Workspace 3 in namespace-c (will be in the fake client as existing) + workspace3 := setupWorkspace(t, "workspace-ns-c", "uid-ns-c", "namespace-c") + + // All three workspaces exist, but in different namespaces + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme). + WithObjects(workspace2, workspace3).Build() + handler := &WebhookHandler{Client: fakeClient} + + // Validating workspace1 should succeed (different namespaces) + conflict, err := handler.validateEndpoints(context.TODO(), workspace1, discoverableEndpointNames(workspace1)) + assert.NoError(t, err) + assert.Nil(t, conflict, "Should allow same endpoint name in different namespaces") + }) +} + +func TestValidateDevfilePeerDiscoverability(t *testing.T) { + scheme := runtime.NewScheme() + assert.NoError(t, dwv2.AddToScheme(scheme)) + for _, name := range []string{"test--endpoint", "different-endpoint"} { + for _, test := range []struct { + name string + raw string + discovered bool + }{ + {"missing", "", false}, + {"boolean true", "true", true}, + {"boolean false", "false", false}, + {"string true", `"true"`, true}, + {"string uppercase true", `"TRUE"`, true}, + {"string one", `"1"`, true}, + {"string false", `"false"`, false}, + {"string zero", `"0"`, false}, + {"malformed string", `"invalid"`, false}, + {"number", "1", false}, + {"null", "null", false}, + {"object", "{}", false}, + } { + t.Run(name+"/"+test.name, func(t *testing.T) { + workspace := setupWorkspace(t, "incoming", "incoming-uid", "test-namespace") + peer := setupWorkspace(t, "peer", "peer-uid", "test-namespace") + endpoint := &peer.Spec.Template.Components[0].Container.Endpoints[0] + endpoint.Name = name + endpoint.Attributes = nil + if test.raw != "" { + endpoint.Attributes = attributes.Attributes{"discoverable": apiextv1.JSON{Raw: []byte(test.raw)}} + } + handler := &WebhookHandler{ + Client: newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(peer).Build(), + Decoder: admission.NewDecoder(scheme), + } + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Create, + Object: newRawExtension(t, workspace), + }} + response := handler.ValidateDevfile(context.Background(), req) + wantConflict := name == "test--endpoint" && test.discovered + assert.Equal(t, !wantConflict, response.Allowed) + if wantConflict { + assert.Contains(t, response.Result.Message, name, "Report the raw peer name") + assert.Contains(t, response.Result.Message, peer.Name) + } + }) + } + } +} + +func TestShouldCheckEndpointConflicts(t *testing.T) { + scheme := runtime.NewScheme() + _ = dwv2.AddToScheme(scheme) + handler := &WebhookHandler{Decoder: admission.NewDecoder(scheme)} + + t.Run("Skips decoding on update with no discoverable endpoints", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + workspace.Spec.Template.Components[0].Container.Endpoints[0].Attributes = nil + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + OldObject: runtime.RawExtension{Raw: []byte("not-json")}, + }} + handler := &WebhookHandler{} + + assert.NotPanics(t, func() { + assert.False(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(workspace)), "An empty incoming endpoint set needs no conflict check or old-object decoding") + }) + }) + + t.Run("Always checks on create", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Create, + Object: newRawExtension(t, workspace), + }} + assert.True(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(workspace)), "Create requests must always be checked") + }) + + t.Run("Skips check on update when discoverable endpoints are unchanged", func(t *testing.T) { + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newWorkspace.Spec.Started = !oldWorkspace.Spec.Started + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + assert.False(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(newWorkspace)), "Unrelated update should skip the check") + }) + + t.Run("Checks on update when a discoverable endpoint is added", func(t *testing.T) { + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newEndpoint := newWorkspace.Spec.Template.Components[0].Container.Endpoints[0] + newEndpoint.Name = "another-endpoint" + newWorkspace.Spec.Template.Components[0].Container.Endpoints = append( + newWorkspace.Spec.Template.Components[0].Container.Endpoints, newEndpoint) + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + assert.True(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(newWorkspace)), "Adding a discoverable endpoint must trigger the check") + }) + + t.Run("Skips check on update when a discoverable endpoint is only removed", func(t *testing.T) { + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + secondEndpoint := oldWorkspace.Spec.Template.Components[0].Container.Endpoints[0] + secondEndpoint.Name = "second-endpoint" + oldWorkspace.Spec.Template.Components[0].Container.Endpoints = append( + oldWorkspace.Spec.Template.Components[0].Container.Endpoints, secondEndpoint) + + // newWorkspace keeps only the first endpoint - the second one was removed, nothing was added. + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + assert.False(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(newWorkspace)), "Removing a discoverable endpoint must not trigger the check") + }) + + t.Run("Skips check on update when a discoverable endpoint is renamed to the same sanitized service name", func(t *testing.T) { + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + oldWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "my--endpoint" + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + // Collapsing repeated hyphens gives the same Service name, so the rename introduces no new conflict. + newWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "my-endpoint" + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + assert.False(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(newWorkspace)), "Renaming to a name that sanitizes to the same service name should skip the check") + }) + + t.Run("Checks on update when the old object cannot be decoded", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, workspace), + OldObject: runtime.RawExtension{Raw: []byte("not-json")}, + }} + assert.True(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(workspace)), "Decode failures must fail safe by running the check") + }) + + t.Run("Checks on update when a discoverable endpoint is renamed", func(t *testing.T) { + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Name = "renamed-endpoint" + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + assert.True(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(newWorkspace)), "Renaming a discoverable endpoint must trigger the check") + }) + + t.Run("Checks on update when an endpoint becomes discoverable", func(t *testing.T) { + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + oldWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Attributes = nil + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + assert.True(t, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(newWorkspace)), "An endpoint becoming discoverable must trigger the check") + }) +} + +func TestValidateDevfileEndpointConflictGating(t *testing.T) { + scheme := runtime.NewScheme() + _ = dwv2.AddToScheme(scheme) + + t.Run("Denies invalid events with no discoverable endpoints and a malformed old object", func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + workspace.Spec.Template.Components[0].Container.Endpoints[0].Attributes = nil + workspace.Spec.Template.Events = &dwv2.Events{DevWorkspaceEvents: dwv2.DevWorkspaceEvents{PreStart: []string{"missing-command"}}} + handler := &WebhookHandler{Decoder: admission.NewDecoder(scheme)} + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, workspace), + OldObject: runtime.RawExtension{Raw: []byte("not-json")}, + }} + + resp := handler.ValidateDevfile(context.TODO(), req) + assert.False(t, resp.Allowed, "Skipping endpoint checks must still validate events") + assert.Contains(t, resp.Result.Message, "preStart type events are invalid") + assert.Contains(t, resp.Result.Message, "missing-command does not map to a valid devfile command") + }) + + t.Run("Allows an unrelated update even though another workspace already has a conflicting endpoint", func(t *testing.T) { + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient, Decoder: admission.NewDecoder(scheme)} + + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + // Unrelated change: the discoverable endpoint set is identical to oldWorkspace's. + newWorkspace.Spec.Started = !oldWorkspace.Spec.Started + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + + resp := handler.ValidateDevfile(context.TODO(), req) + assert.True(t, resp.Allowed, "Update unrelated to endpoints must be allowed even though a real conflict exists, proving the check was skipped") + }) + + t.Run("Denies an update that introduces a new conflicting discoverable endpoint", func(t *testing.T) { + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient, Decoder: admission.NewDecoder(scheme)} + + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + // Old workspace has no discoverable endpoints yet. + oldWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Attributes = nil + + // New workspace makes "test-endpoint" discoverable, which now conflicts with otherWorkspace. + newWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + + resp := handler.ValidateDevfile(context.TODO(), req) + assert.False(t, resp.Allowed, "Update introducing a new conflicting discoverable endpoint must be denied") + assert.Contains(t, resp.Result.Message, "test-endpoint") + assert.Contains(t, resp.Result.Message, "workspace-2") + }) + + t.Run("Denies adding a unique endpoint while retaining an existing conflicting endpoint", func(t *testing.T) { + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + fakeClient := newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + handler := &WebhookHandler{Client: fakeClient, Decoder: admission.NewDecoder(scheme)} + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + newWorkspace := oldWorkspace.DeepCopy() + endpoint := newWorkspace.Spec.Template.Components[0].Container.Endpoints[0] + endpoint.Name = "unique-endpoint" + endpoint.TargetPort = 8081 + newWorkspace.Spec.Template.Components[0].Container.Endpoints = append( + newWorkspace.Spec.Template.Components[0].Container.Endpoints, endpoint) + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + + resp := handler.ValidateDevfile(context.TODO(), req) + assert.False(t, resp.Allowed, "Adding a name must validate the full incoming endpoint set") + assert.Contains(t, resp.Result.Message, "test-endpoint") + assert.Contains(t, resp.Result.Message, "workspace-2") + }) +} + +func TestValidateDevfileEndpointExposure(t *testing.T) { + scheme := runtime.NewScheme() + assert.NoError(t, dwv2.AddToScheme(scheme)) + + for _, tc := range []struct { + name string + incoming dwv2.EndpointExposure + existing dwv2.EndpointExposure + allowed bool + withoutClient bool + }{ + {name: "Incoming none allows a matching eligible endpoint", incoming: dwv2.NoneEndpointExposure, existing: dwv2.PublicEndpointExposure, allowed: true}, + {name: "Incoming none skips List", incoming: dwv2.NoneEndpointExposure, existing: dwv2.PublicEndpointExposure, allowed: true, withoutClient: true}, + {name: "Existing none does not block admission", incoming: dwv2.PublicEndpointExposure, existing: dwv2.NoneEndpointExposure, allowed: true}, + {name: "Internal endpoints conflict", incoming: dwv2.InternalEndpointExposure, existing: dwv2.InternalEndpointExposure}, + {name: "Public endpoints conflict", incoming: dwv2.PublicEndpointExposure, existing: dwv2.PublicEndpointExposure}, + {name: "Omitted exposure defaults to public and conflicts"}, + } { + t.Run(tc.name, func(t *testing.T) { + workspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + workspace.Spec.Template.Components[0].Container.Endpoints[0].Exposure = tc.incoming + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + otherWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Exposure = tc.existing + handler := &WebhookHandler{Decoder: admission.NewDecoder(scheme)} + if !tc.withoutClient { + handler.Client = newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + } + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Create, + Object: newRawExtension(t, workspace), + }} + + assert.NotPanics(t, func() { + resp := handler.ValidateDevfile(context.TODO(), req) + assert.Equal(t, tc.allowed, resp.Allowed, "Unexpected admission result: %v", resp.Result) + if !tc.allowed { + assert.Contains(t, resp.Result.Message, "test-endpoint") + assert.Contains(t, resp.Result.Message, "workspace-2") + } + }) + }) + } +} + +func TestValidateDevfileEndpointExposureChanges(t *testing.T) { + scheme := runtime.NewScheme() + assert.NoError(t, dwv2.AddToScheme(scheme)) + + for _, tc := range []struct { + name string + before dwv2.EndpointExposure + after dwv2.EndpointExposure + check bool + }{ + {"none to internal", dwv2.NoneEndpointExposure, dwv2.InternalEndpointExposure, true}, + {"none to public", dwv2.NoneEndpointExposure, dwv2.PublicEndpointExposure, true}, + {"none to omitted", dwv2.NoneEndpointExposure, "", true}, + {"internal to none", dwv2.InternalEndpointExposure, dwv2.NoneEndpointExposure, false}, + {"public to none", dwv2.PublicEndpointExposure, dwv2.NoneEndpointExposure, false}, + {"omitted to none", "", dwv2.NoneEndpointExposure, false}, + {"internal to public", dwv2.InternalEndpointExposure, dwv2.PublicEndpointExposure, false}, + {"public to internal", dwv2.PublicEndpointExposure, dwv2.InternalEndpointExposure, false}, + } { + t.Run(tc.name, func(t *testing.T) { + oldWorkspace := setupWorkspace(t, "workspace-1", "uid-1", "test-namespace") + oldWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Exposure = tc.before + newWorkspace := oldWorkspace.DeepCopy() + newWorkspace.Spec.Template.Components[0].Container.Endpoints[0].Exposure = tc.after + handler := &WebhookHandler{Decoder: admission.NewDecoder(scheme)} + if tc.check { + otherWorkspace := setupWorkspace(t, "workspace-2", "uid-2", "test-namespace") + handler.Client = newIndexedWorkspaceClientBuilder(t, scheme).WithObjects(otherWorkspace).Build() + } + req := admission.Request{AdmissionRequest: admissionv1.AdmissionRequest{ + Operation: admissionv1.Update, + Object: newRawExtension(t, newWorkspace), + OldObject: newRawExtension(t, oldWorkspace), + }} + + assert.Equal(t, tc.check, handler.shouldCheckEndpointConflicts(req, discoverableEndpointNames(newWorkspace))) + assert.NotPanics(t, func() { + resp := handler.ValidateDevfile(context.TODO(), req) + assert.Equal(t, !tc.check, resp.Allowed, "Unexpected admission result: %v", resp.Result) + if tc.check { + assert.Contains(t, resp.Result.Message, "test-endpoint") + assert.Contains(t, resp.Result.Message, "workspace-2") + } + }) + }) + } +}