Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
106 changes: 105 additions & 1 deletion modules/common/service/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ package service
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/url"
"strconv"
Expand All @@ -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"
Expand All @@ -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,
Expand Down Expand Up @@ -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)
}
Expand All @@ -82,13 +103,91 @@ 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
}
}

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)
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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
Expand Down Expand Up @@ -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
Expand Down
110 changes: 110 additions & 0 deletions modules/common/service/service_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
})
}
}
34 changes: 32 additions & 2 deletions modules/common/service/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
20 changes: 20 additions & 0 deletions modules/common/service/zz_generated.deepcopy.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading