From ea9f6a4d669d66527406099c5eeb44f875427fd6 Mon Sep 17 00:00:00 2001 From: Luca Miccini Date: Fri, 2 Oct 2026 09:03:55 +0200 Subject: [PATCH] service: allow overriding Service ports in OverrideServiceSpec Add a Ports field to OverrideServiceSpec using a dedicated OverrideServicePort type that exposes only Name and Port. Each entry changes the Port of the base Service port with the matching Name; TargetPort and every other field are preserved so the exposed port can change without moving backend routing. When the base relies on Kubernetes defaulting TargetPort to Port and the Port changes, TargetPort is pinned to the original base Port. An override whose Name does not match any base port is rejected (ErrUnknownServicePort), and the merged result is validated for duplicate (Port, Protocol) pairs. This supports placing several Services on a single shared LoadBalancer IP by changing their Port values. Co-Authored-By: Claude Opus 4.8 --- modules/common/service/service.go | 106 ++++++++++++++++- modules/common/service/service_test.go | 110 ++++++++++++++++++ modules/common/service/types.go | 34 +++++- .../common/service/zz_generated.deepcopy.go | 20 ++++ 4 files changed, 267 insertions(+), 3 deletions(-) diff --git a/modules/common/service/service.go b/modules/common/service/service.go index d6d4541b..8de7265c 100644 --- a/modules/common/service/service.go +++ b/modules/common/service/service.go @@ -20,6 +20,7 @@ package service import ( "context" "encoding/json" + "errors" "fmt" "net/url" "strconv" @@ -30,6 +31,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/intstr" "k8s.io/apimachinery/pkg/util/strategicpatch" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -39,6 +41,15 @@ import ( ctrl "sigs.k8s.io/controller-runtime" ) +var ( + // ErrUnknownServicePort is returned when a ports override references a Name + // that does not match any base Service port. + ErrUnknownServicePort = errors.New("ports override references unknown service port name") + // ErrDuplicateServicePort is returned when the merged port list contains a + // duplicate (Port, Protocol) pair. + ErrDuplicateServicePort = errors.New("duplicate service port after applying ports override") +) + // NewService returns an initialized Service. func NewService( service *corev1.Service, @@ -67,7 +78,17 @@ func NewService( return svc, fmt.Errorf("error marshalling Service Spec: %w", err) } - patch, err := json.Marshal(override.Spec) + // Ports are handled separately (merged by name below) instead of via + // the strategic merge patch: ServicePort's strategic merge key is + // "port", which makes it impossible to change a port's Port value - + // the patch would append a new entry rather than update the existing + // one. Strip Ports from the patch so the strategic merge only touches + // the other fields. + specWithoutPorts := *override.Spec + overridePorts := specWithoutPorts.Ports + specWithoutPorts.Ports = nil + + patch, err := json.Marshal(&specWithoutPorts) if err != nil { return svc, fmt.Errorf("error marshalling Service Spec override: %w", err) } @@ -82,6 +103,11 @@ func NewService( if err != nil { return svc, fmt.Errorf("error unmarshalling patched Service Spec: %w", err) } + + patchedSpec.Ports, err = mergeServicePortsByName(patchedSpec.Ports, overridePorts) + if err != nil { + return svc, fmt.Errorf("error applying Service ports override: %w", err) + } svc.service.Spec = patchedSpec } } @@ -89,6 +115,79 @@ func NewService( return svc, nil } +// mergeServicePortsByName applies Port overrides to base ports, matched by Name. +// Each override changes only the Port of the base port with the same Name; every +// other field is preserved. An override whose Name does not match any base port +// is an error. The merged result is validated so that invalid combinations +// surface here rather than as a hard-to-trace API server rejection at +// create/update time. +func mergeServicePortsByName(base []corev1.ServicePort, overrides []OverrideServicePort) ([]corev1.ServicePort, error) { + if len(overrides) == 0 { + return base, nil + } + + merged := make([]corev1.ServicePort, len(base)) + copy(merged, base) + + for _, ov := range overrides { + idx := -1 + for i := range merged { + if merged[i].Name == ov.Name { + idx = i + break + } + } + if idx < 0 { + return nil, fmt.Errorf("%w: %q", ErrUnknownServicePort, ov.Name) + } + merged[idx] = overrideServicePort(merged[idx], ov) + } + + if err := validateServicePorts(merged); err != nil { + return nil, err + } + + return merged, nil +} + +// validateServicePorts checks that a merged port list is valid for a Service: no +// two ports may share the same (Port, Protocol) pair, which the API server would +// otherwise reject at create/update time. An empty Protocol is treated as TCP, +// matching Kubernetes defaulting. +func validateServicePorts(ports []corev1.ServicePort) error { + type portKey struct { + port int32 + protocol corev1.Protocol + } + seen := make(map[portKey]struct{}, len(ports)) + for _, p := range ports { + proto := p.Protocol + if proto == "" { + proto = corev1.ProtocolTCP + } + key := portKey{port: p.Port, protocol: proto} + if _, ok := seen[key]; ok { + return fmt.Errorf("%w: %d/%s", ErrDuplicateServicePort, p.Port, proto) + } + seen[key] = struct{}{} + } + return nil +} + +// overrideServicePort returns base with its Port replaced by the override Port. +// When base relies on Kubernetes defaulting TargetPort to Port (TargetPort +// unset) and the Port changes, TargetPort is pinned to the original base Port so +// that changing the exposed port does not silently move backend routing to the +// new Port. +func overrideServicePort(base corev1.ServicePort, ov OverrideServicePort) corev1.ServicePort { + out := base + if ov.Port != out.Port && out.TargetPort == (intstr.IntOrString{}) { + out.TargetPort = intstr.FromInt32(out.Port) + } + out.Port = ov.Port + return out +} + // GetClusterIPs - returns the cluster IPs of the created service func (s *Service) GetClusterIPs() []string { return s.clusterIPs @@ -227,6 +326,11 @@ func (s *Service) ToOverrideServiceSpec() (*OverrideServiceSpec, error) { if err != nil { return nil, fmt.Errorf("error unmarshalling service OverrideSpec: %w", err) } + + // Ports are an override-only input (merged by name in NewService) and are + // not reflected back when converting a live Service spec, to preserve the + // existing behaviour of callers that round-trip through this method. + overrideServiceSpec.Ports = nil } return overrideServiceSpec, nil diff --git a/modules/common/service/service_test.go b/modules/common/service/service_test.go index e0abe92f..4b5a3567 100644 --- a/modules/common/service/service_test.go +++ b/modules/common/service/service_test.go @@ -633,3 +633,113 @@ func TestRoutedOverrideSpecAddAnnotation(t *testing.T) { }) } } + +func TestNewServicePortsOverride(t *testing.T) { + basePorts := []corev1.ServicePort{ + {Name: "amqp", Protocol: corev1.ProtocolTCP, Port: 5672, TargetPort: intstr.FromInt(5672)}, + {Name: "amqps", Protocol: corev1.ProtocolTCP, Port: 5671, TargetPort: intstr.FromInt(5671)}, + } + + tests := []struct { + name string + override *OverrideSpec + want []corev1.ServicePort + }{ + { + name: "no port override keeps base ports", + override: &OverrideSpec{Spec: &OverrideServiceSpec{Type: corev1.ServiceTypeLoadBalancer}}, + want: basePorts, + }, + { + name: "change matching port keeps other fields and other ports", + override: &OverrideSpec{Spec: &OverrideServiceSpec{ + Ports: []OverrideServicePort{{Name: "amqp", Port: 5673}}, + }}, + // amqp Port changes to 5673, TargetPort/Protocol are preserved from + // base; amqps is untouched. + want: []corev1.ServicePort{ + {Name: "amqp", Protocol: corev1.ProtocolTCP, Port: 5673, TargetPort: intstr.FromInt(5672)}, + {Name: "amqps", Protocol: corev1.ProtocolTCP, Port: 5671, TargetPort: intstr.FromInt(5671)}, + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + g := NewWithT(t) + + base := svcClusterIP.DeepCopy() + base.Spec.Ports = append([]corev1.ServicePort{}, basePorts...) + + svc, err := NewService(base, timeout, tt.override) + g.Expect(err).ToNot(HaveOccurred()) + g.Expect(svc.GetSpec().Ports).To(Equal(tt.want)) + }) + } +} + +// TestNewServicePortsOverrideTargetPort verifies that overriding only the Port +// of a base port that omits TargetPort preserves backend routing: TargetPort is +// pinned to the original base Port instead of following the new Port via +// Kubernetes' TargetPort=Port defaulting. +func TestNewServicePortsOverrideTargetPort(t *testing.T) { + g := NewWithT(t) + + base := svcClusterIP.DeepCopy() + // single port, no TargetPort set -> would default to Port (80) + base.Spec.Ports = []corev1.ServicePort{ + {Name: "foo", Protocol: corev1.ProtocolTCP, Port: 80}, + } + + override := &OverrideSpec{Spec: &OverrideServiceSpec{ + Ports: []OverrideServicePort{{Name: "foo", Port: 8080}}, + }} + + svc, err := NewService(base, timeout, override) + g.Expect(err).ToNot(HaveOccurred()) + // Port changes to 8080 but backend routing stays at the original 80. + g.Expect(svc.GetSpec().Ports).To(Equal([]corev1.ServicePort{ + {Name: "foo", Protocol: corev1.ProtocolTCP, Port: 8080, TargetPort: intstr.FromInt32(80)}, + })) +} + +func TestNewServicePortsOverrideInvalid(t *testing.T) { + basePorts := []corev1.ServicePort{ + {Name: "amqp", Protocol: corev1.ProtocolTCP, Port: 5672, TargetPort: intstr.FromInt(5672)}, + {Name: "amqps", Protocol: corev1.ProtocolTCP, Port: 5671, TargetPort: intstr.FromInt(5671)}, + } + + tests := []struct { + name string + override *OverrideSpec + errMatch string + }{ + { + name: "changing a port to collide with another is rejected", + override: &OverrideSpec{Spec: &OverrideServiceSpec{ + Ports: []OverrideServicePort{{Name: "amqp", Port: 5671}}, + }}, + errMatch: "duplicate service port after applying ports override: 5671/TCP", + }, + { + name: "override referencing an unknown port name is rejected", + override: &OverrideSpec{Spec: &OverrideServiceSpec{ + Ports: []OverrideServicePort{{Name: "does-not-exist", Port: 9999}}, + }}, + errMatch: `ports override references unknown service port name: "does-not-exist"`, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + g := NewWithT(t) + + base := svcClusterIP.DeepCopy() + base.Spec.Ports = append([]corev1.ServicePort{}, basePorts...) + + _, err := NewService(base, timeout, tt.override) + g.Expect(err).To(HaveOccurred()) + g.Expect(err.Error()).To(ContainSubstring(tt.errMatch)) + }) + } +} diff --git a/modules/common/service/types.go b/modules/common/service/types.go index a82cb72d..30ece714 100644 --- a/modules/common/service/types.go +++ b/modules/common/service/types.go @@ -180,10 +180,40 @@ type EmbeddedLabelsAnnotations struct { Annotations map[string]string `json:"annotations,omitempty" protobuf:"bytes,12,rep,name=annotations"` } +// OverrideServicePort overrides the Port number of a generated Service port, +// identified by Name. Only the Port can be changed: TargetPort and every other +// field are preserved from the matched base port, so overriding the exposed port +// never changes backend routing. This deliberately exposes a smaller surface +// than the upstream corev1.ServicePort, which would allow changing TargetPort +// (and other fields) and silently break the deployment connection. +type OverrideServicePort struct { + // Name of the base Service port to override. It must match the Name of an + // existing port on the generated Service. + Name string `json:"name"` + + // Port is the new Service port number. + // +kubebuilder:validation:Minimum=1 + // +kubebuilder:validation:Maximum=65535 + Port int32 `json:"port"` +} + // OverrideServiceSpec is a subset of the fields included in https://pkg.go.dev/k8s.io/api@v0.26.6/core/v1#ServiceSpec -// Limited to Type, SessionAffinity, LoadBalancerSourceRanges, ExternalName, ExternalTrafficPolicy, SessionAffinityConfig, -// IPFamilyPolicy, LoadBalancerClass and InternalTrafficPolicy +// Limited to Type, Ports, SessionAffinity, LoadBalancerSourceRanges, ExternalName, ExternalTrafficPolicy, +// SessionAffinityConfig, IPFamilyPolicy, LoadBalancerClass and InternalTrafficPolicy type OverrideServiceSpec struct { + // Ports overrides the Port number of existing generated Service ports, + // matched by Name. Each entry changes only the Port of the base port with the + // same Name; TargetPort and every other field are preserved from the base + // port, so overriding the exposed port never changes backend routing. An + // entry whose Name does not match any base port is rejected. + // + // The intended use is placing several Services on a single shared + // LoadBalancer IP by changing their Port values. To leave a Service's ports + // untouched, omit this field entirely. + // +optional + // +listType=atomic + Ports []OverrideServicePort `json:"ports,omitempty"` + // type determines how the Service is exposed. Defaults to ClusterIP. Valid // options are ExternalName, ClusterIP, NodePort, and LoadBalancer. // "ClusterIP" allocates a cluster-internal IP address for load-balancing diff --git a/modules/common/service/zz_generated.deepcopy.go b/modules/common/service/zz_generated.deepcopy.go index d258e1d9..a61700a6 100644 --- a/modules/common/service/zz_generated.deepcopy.go +++ b/modules/common/service/zz_generated.deepcopy.go @@ -53,9 +53,29 @@ func (in *EmbeddedLabelsAnnotations) DeepCopy() *EmbeddedLabelsAnnotations { return out } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *OverrideServicePort) DeepCopyInto(out *OverrideServicePort) { + *out = *in +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new OverrideServicePort. +func (in *OverrideServicePort) DeepCopy() *OverrideServicePort { + if in == nil { + return nil + } + out := new(OverrideServicePort) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *OverrideServiceSpec) DeepCopyInto(out *OverrideServiceSpec) { *out = *in + if in.Ports != nil { + in, out := &in.Ports, &out.Ports + *out = make([]OverrideServicePort, len(*in)) + copy(*out, *in) + } if in.LoadBalancerSourceRanges != nil { in, out := &in.LoadBalancerSourceRanges, &out.LoadBalancerSourceRanges *out = make([]string, len(*in))