Skip to content

Commit 5d03c6d

Browse files
lmicciniclaude
andcommitted
Propagate RabbitMQ service overrides via strategic merge
Replace manual cherry-picking of individual ServiceSpec override fields (Type, IPFamilyPolicy) with lib-common's OverrideServiceSpec whitelist and strategic merge patch. This fixes LoadBalancerClass not being propagated (OSPRH-30887) and also enables SessionAffinity, ExternalTrafficPolicy, LoadBalancerSourceRanges, InternalTrafficPolicy, ExternalName, and SessionAffinityConfig overrides — all of which were silently ignored before. The override's full corev1.ServiceSpec is converted to the restricted OverrideServiceSpec via JSON marshal/unmarshal, ensuring unsafe fields (Selector, Ports, ClusterIP) cannot be overridden. This matches the pattern used by lib-common's service.NewService() and by the RabbitMQ per-pod services in this same controller. The headless service only applies metadata overrides and IPFamilyPolicy — other spec fields are invalid or meaningless for ClusterIP/None services. Related: OSPRH-30887 Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 797736a commit 5d03c6d

3 files changed

Lines changed: 345 additions & 63 deletions

File tree

internal/controller/rabbitmq/rabbitmq_controller.go

Lines changed: 18 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -632,15 +632,16 @@ func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (result ct
632632
},
633633
}
634634
hsop, err := controllerutil.CreateOrPatch(ctx, r.Client, headlessSvc, func() error {
635-
desired := rabbitmq.HeadlessService(instance)
636-
headlessSvc.Spec.ClusterIP = desired.Spec.ClusterIP
637-
headlessSvc.Spec.PublishNotReadyAddresses = desired.Spec.PublishNotReadyAddresses
638-
mergeServicePorts(&headlessSvc.Spec.Ports, desired.Spec.Ports)
639-
headlessSvc.Spec.Selector = desired.Spec.Selector
640-
headlessSvc.Labels = desired.Labels
641-
if desired.Spec.IPFamilyPolicy != nil {
642-
headlessSvc.Spec.IPFamilyPolicy = desired.Spec.IPFamilyPolicy
635+
desired, err := rabbitmq.HeadlessService(instance)
636+
if err != nil {
637+
return err
643638
}
639+
savedPorts := headlessSvc.Spec.Ports
640+
headlessSvc.Spec = desired.Spec
641+
mergeServicePorts(&savedPorts, desired.Spec.Ports)
642+
headlessSvc.Spec.Ports = savedPorts
643+
headlessSvc.Labels = desired.Labels
644+
headlessSvc.Annotations = util.MergeStringMaps(desired.Annotations, headlessSvc.Annotations)
644645
return r.setOwnership(headlessSvc, instance, isMigration)
645646
})
646647
if err != nil {
@@ -658,18 +659,17 @@ func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (result ct
658659
},
659660
}
660661
csop, err := controllerutil.CreateOrPatch(ctx, r.Client, clientSvc, func() error {
661-
desired := rabbitmq.ClientService(instance)
662-
mergeServicePorts(&clientSvc.Spec.Ports, desired.Spec.Ports)
663-
clientSvc.Spec.Selector = desired.Spec.Selector
664-
clientSvc.Spec.Type = desired.Spec.Type
662+
desired, err := rabbitmq.ClientService(instance)
663+
if err != nil {
664+
return err
665+
}
666+
savedPorts := clientSvc.Spec.Ports
667+
clientSvc.Spec = desired.Spec
668+
mergeServicePorts(&savedPorts, desired.Spec.Ports)
669+
clientSvc.Spec.Ports = savedPorts
665670
// Merge annotations: set our desired ones without removing externally-added
666671
// annotations (e.g. from MetalLB) to avoid reconcile loops
667-
if clientSvc.Annotations == nil {
668-
clientSvc.Annotations = map[string]string{}
669-
}
670-
for k, v := range desired.Annotations {
671-
clientSvc.Annotations[k] = v
672-
}
672+
clientSvc.Annotations = util.MergeStringMaps(desired.Annotations, clientSvc.Annotations)
673673
clientSvc.Labels = desired.Labels
674674
return r.setOwnership(clientSvc, instance, isMigration)
675675
})

internal/rabbitmq/service.go

Lines changed: 71 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -1,27 +1,57 @@
11
package rabbitmq
22

33
import (
4+
"encoding/json"
45
"fmt"
56

67
rabbitmqv1 "github.com/openstack-k8s-operators/infra-operator/apis/rabbitmq/v1beta1"
8+
"github.com/openstack-k8s-operators/lib-common/modules/common/service"
9+
"github.com/openstack-k8s-operators/lib-common/modules/common/util"
710
corev1 "k8s.io/api/core/v1"
811
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
912
"k8s.io/apimachinery/pkg/util/intstr"
13+
"k8s.io/apimachinery/pkg/util/strategicpatch"
1014
"k8s.io/utils/ptr"
1115
)
1216

13-
// serviceOverrideType returns the service type from override, or empty string
14-
func serviceOverrideType(r *rabbitmqv1.RabbitMq) corev1.ServiceType {
15-
if r.Spec.Override.Service != nil && r.Spec.Override.Service.Spec != nil {
16-
return r.Spec.Override.Service.Spec.Type
17-
}
18-
return ""
19-
}
17+
// applyServiceOverride applies override metadata and spec fields to a Service.
18+
// The full corev1.ServiceSpec override is converted to lib-common's restricted
19+
// service.OverrideServiceSpec before merging via strategic merge patch, so only
20+
// the whitelisted fields are applied and unsafe fields like Selector, Ports, or
21+
// ClusterIP cannot be overridden.
22+
func applyServiceOverride(svc *corev1.Service, override *rabbitmqv1.RabbitMQServiceOverride) error {
23+
if override == nil {
24+
return nil
25+
}
26+
if override.EmbeddedLabelsAnnotations != nil {
27+
svc.Labels = util.MergeStringMaps(override.Labels, svc.Labels)
28+
svc.Annotations = util.MergeStringMaps(override.Annotations, svc.Annotations)
29+
}
30+
if override.Spec != nil {
31+
overrideSpecJSON, err := json.Marshal(override.Spec)
32+
if err != nil {
33+
return fmt.Errorf("error marshalling service override spec: %w", err)
34+
}
35+
whitelisted := &service.OverrideServiceSpec{}
36+
if err := json.Unmarshal(overrideSpecJSON, whitelisted); err != nil {
37+
return fmt.Errorf("error converting to OverrideServiceSpec: %w", err)
38+
}
2039

21-
// serviceOverrideIPFamilyPolicy returns the IPFamilyPolicy from override, or nil
22-
func serviceOverrideIPFamilyPolicy(r *rabbitmqv1.RabbitMq) *corev1.IPFamilyPolicy {
23-
if r.Spec.Override.Service != nil && r.Spec.Override.Service.Spec != nil {
24-
return r.Spec.Override.Service.Spec.IPFamilyPolicy
40+
originalJSON, err := json.Marshal(svc.Spec)
41+
if err != nil {
42+
return fmt.Errorf("error marshalling service spec: %w", err)
43+
}
44+
patchJSON, err := json.Marshal(whitelisted)
45+
if err != nil {
46+
return fmt.Errorf("error marshalling override patch: %w", err)
47+
}
48+
patchedJSON, err := strategicpatch.StrategicMergePatch(originalJSON, patchJSON, corev1.ServiceSpec{})
49+
if err != nil {
50+
return fmt.Errorf("error applying service spec override: %w", err)
51+
}
52+
if err := json.Unmarshal(patchedJSON, &svc.Spec); err != nil {
53+
return fmt.Errorf("error unmarshalling patched service spec: %w", err)
54+
}
2555
}
2656
return nil
2757
}
@@ -64,7 +94,7 @@ const (
6494

6595
// HeadlessService creates the headless service for StatefulSet pod DNS
6696
// matching the old rabbitmq-cluster-operator layout
67-
func HeadlessService(r *rabbitmqv1.RabbitMq) *corev1.Service {
97+
func HeadlessService(r *rabbitmqv1.RabbitMq) (*corev1.Service, error) {
6898
ls := CommonLabels(r.Name)
6999
selector := SelectorLabels(r.Name)
70100

@@ -98,61 +128,57 @@ func HeadlessService(r *rabbitmqv1.RabbitMq) *corev1.Service {
98128
},
99129
}
100130

101-
// Apply override settings if specified
102-
if ipfp := serviceOverrideIPFamilyPolicy(r); ipfp != nil {
103-
svc.Spec.IPFamilyPolicy = ipfp
131+
if override := r.Spec.Override.Service; override != nil {
132+
if override.EmbeddedLabelsAnnotations != nil {
133+
svc.Labels = util.MergeStringMaps(override.Labels, svc.Labels)
134+
svc.Annotations = util.MergeStringMaps(override.Annotations, svc.Annotations)
135+
}
136+
// Only IPFamilyPolicy is applicable to headless services — other spec
137+
// fields (Type, ExternalTrafficPolicy, LoadBalancerClass, etc.) are
138+
// either invalid or meaningless for ClusterIP/None services.
139+
if override.Spec != nil && override.Spec.IPFamilyPolicy != nil {
140+
svc.Spec.IPFamilyPolicy = override.Spec.IPFamilyPolicy
141+
}
104142
}
105143

106-
return svc
144+
return svc, nil
107145
}
108146

109147
// ClientService creates the client-facing service for RabbitMQ
110148
// matching the old rabbitmq-cluster-operator layout
111-
func ClientService(r *rabbitmqv1.RabbitMq) *corev1.Service {
149+
func ClientService(r *rabbitmqv1.RabbitMq) (*corev1.Service, error) {
112150
ls := CommonLabels(r.Name)
113151
selector := SelectorLabels(r.Name)
114152

115153
ports := buildClientServicePorts(r)
116154

117-
serviceType := corev1.ServiceTypeClusterIP
118-
annotations := make(map[string]string)
119-
120-
// Apply override settings if specified
121-
if overrideType := serviceOverrideType(r); overrideType != "" {
122-
serviceType = overrideType
123-
}
124-
if r.Spec.Override.Service != nil && r.Spec.Override.Service.EmbeddedLabelsAnnotations != nil {
125-
for k, v := range r.Spec.Override.Service.Annotations {
126-
annotations[k] = v
127-
}
128-
}
129-
130-
if serviceType == corev1.ServiceTypeLoadBalancer {
131-
if _, exists := annotations["dnsmasq.network.openstack.org/hostname"]; !exists {
132-
annotations["dnsmasq.network.openstack.org/hostname"] = fmt.Sprintf("%s.%s.svc", r.Name, r.Namespace)
133-
}
134-
}
135-
136155
svc := &corev1.Service{
137156
ObjectMeta: metav1.ObjectMeta{
138-
Name: r.Name,
139-
Namespace: r.Namespace,
140-
Labels: ls,
141-
Annotations: annotations,
157+
Name: r.Name,
158+
Namespace: r.Namespace,
159+
Labels: ls,
142160
},
143161
Spec: corev1.ServiceSpec{
144-
Type: serviceType,
162+
Type: corev1.ServiceTypeClusterIP,
145163
Selector: selector,
146164
Ports: ports,
147165
},
148166
}
149167

150-
// Apply IPFamilyPolicy from override if specified
151-
if ipfp := serviceOverrideIPFamilyPolicy(r); ipfp != nil {
152-
svc.Spec.IPFamilyPolicy = ipfp
168+
if err := applyServiceOverride(svc, r.Spec.Override.Service); err != nil {
169+
return nil, fmt.Errorf("client service override: %w", err)
170+
}
171+
172+
if svc.Spec.Type == corev1.ServiceTypeLoadBalancer {
173+
if svc.Annotations == nil {
174+
svc.Annotations = make(map[string]string)
175+
}
176+
if _, exists := svc.Annotations["dnsmasq.network.openstack.org/hostname"]; !exists {
177+
svc.Annotations["dnsmasq.network.openstack.org/hostname"] = fmt.Sprintf("%s.%s.svc", r.Name, r.Namespace)
178+
}
153179
}
154180

155-
return svc
181+
return svc, nil
156182
}
157183

158184
// buildClientServicePorts builds the list of ports for the client service

0 commit comments

Comments
 (0)