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))