Skip to content
Merged
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
1 change: 1 addition & 0 deletions config/rbac/role.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,7 @@ rules:
- roles
verbs:
- create
- delete
- get
- list
- update
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ module github.com/rabbitmq/cluster-operator/v2
go 1.26.2

require (
github.com/Masterminds/semver/v3 v3.4.0
github.com/cloudflare/cfssl v1.6.5
github.com/eclipse/paho.mqtt.golang v1.5.1
github.com/go-logr/logr v1.4.3
Expand All @@ -25,7 +26,6 @@ require (

require (
cel.dev/expr v0.25.1 // indirect
github.com/Masterminds/semver/v3 v3.4.0 // indirect
github.com/antlr4-go/antlr/v4 v4.13.1 // indirect
github.com/beorn7/perks v1.0.1 // indirect
github.com/blang/semver/v4 v4.0.0 // indirect
Expand Down
1 change: 0 additions & 1 deletion go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,6 @@ github.com/fxamacker/cbor/v2 v2.9.1/go.mod h1:vM4b+DJCtHn+zz7h3FFp/hDAI9WNWCsZj2
github.com/gkampitakis/ciinfo v0.3.2 h1:JcuOPk8ZU7nZQjdUhctuhQofk7BGHuIy0c9Ez8BNhXs=
github.com/gkampitakis/ciinfo v0.3.2/go.mod h1:1NIwaOcFChN4fa/B0hEBdAb6npDlFL8Bwx4dfRLRqAo=
github.com/gkampitakis/go-diff v1.3.2 h1:Qyn0J9XJSDTgnsgHRdz9Zp24RaJeKMUHg2+PDZZdC4M=
github.com/gkampitakis/go-diff v1.3.2/go.mod h1:LLgOrpqleQe26cte8s36HTWcTmMEur6OPYerdAAS9tk=
github.com/gkampitakis/go-snaps v0.5.15 h1:amyJrvM1D33cPHwVrjo9jQxX8g/7E2wYdZ+01KS3zGE=
github.com/gkampitakis/go-snaps v0.5.15/go.mod h1:HNpx/9GoKisdhw9AFOBT1N7DBs9DiHo/hGheFGBZ+mc=
github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A=
Expand Down
19 changes: 17 additions & 2 deletions internal/controller/rabbitmqcluster_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -90,10 +90,10 @@ type RabbitmqClusterReconciler struct {
// +kubebuilder:rbac:groups="",resources=events,verbs=get;create;patch
// +kubebuilder:rbac:groups="",resources=serviceaccounts,verbs=get;list;watch;create;update
// +kubebuilder:rbac:groups="",resources=persistentvolumeclaims,verbs=get;list;watch;create;update
// +kubebuilder:rbac:groups="rbac.authorization.k8s.io",resources=roles,verbs=get;list;watch;create;update
// +kubebuilder:rbac:groups="rbac.authorization.k8s.io",resources=roles,verbs=get;list;watch;create;update;delete
// +kubebuilder:rbac:groups="discovery.k8s.io",resources=endpointslices,verbs=get;list;watch
// +kubebuilder:rbac:groups="",resources=endpoints,verbs=get;watch;list
// +kubebuilder:rbac:groups="rbac.authorization.k8s.io",resources=rolebindings,verbs=get;list;watch;create;update
// +kubebuilder:rbac:groups="rbac.authorization.k8s.io",resources=rolebindings,verbs=get;list;watch;create;update;delete

func (r *RabbitmqClusterReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
logger := ctrl.LoggerFrom(ctx)
Expand Down Expand Up @@ -183,6 +183,21 @@ func (r *RabbitmqClusterReconciler) Reconcile(ctx context.Context, req ctrl.Requ
Scheme: r.Scheme,
}

if !resource.ShouldCreatePeerDiscoveryRBAC(rabbitmqCluster) {
// Ensure peer-discovery Role and RoleBinding are deleted.
// The ServiceAccount is intentionally kept because other integrations
// (e.g. Vault Kubernetes auth) may rely on it.
for _, obj := range []client.Object{
&rbacv1.Role{ObjectMeta: metav1.ObjectMeta{Name: rabbitmqCluster.ChildResourceName("peer-discovery"), Namespace: rabbitmqCluster.Namespace}},
&rbacv1.RoleBinding{ObjectMeta: metav1.ObjectMeta{Name: rabbitmqCluster.ChildResourceName("server"), Namespace: rabbitmqCluster.Namespace}},
} {
if err := r.Client.Delete(ctx, obj); client.IgnoreNotFound(err) != nil {
logger.Error(err, "Failed to delete peer-discovery RBAC resource")
return ctrl.Result{}, err
}
}
}

builders := resourceBuilder.ResourceBuilders()

for _, builder := range builders {
Expand Down
39 changes: 39 additions & 0 deletions internal/controller/rabbitmqcluster_controller_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -617,6 +617,45 @@ var _ = Describe("RabbitmqClusterController", func() {
})
})

When("the RabbitMQ version is upgraded to 4.1.0 or greater", func() {
It("deletes the peer-discovery Role and RoleBinding but keeps the ServiceAccount", func() {
// First, ensure the resources exist
_, err := clientSet.CoreV1().ServiceAccounts(cluster.Namespace).Get(ctx, cluster.ChildResourceName("server"), metav1.GetOptions{})
Expect(err).NotTo(HaveOccurred())

_, err = clientSet.RbacV1().Roles(cluster.Namespace).Get(ctx, cluster.ChildResourceName("peer-discovery"), metav1.GetOptions{})
Expect(err).NotTo(HaveOccurred())

_, err = clientSet.RbacV1().RoleBindings(cluster.Namespace).Get(ctx, cluster.ChildResourceName("server"), metav1.GetOptions{})
Expect(err).NotTo(HaveOccurred())

// Upgrade to 4.1.5
Expect(updateWithRetry(cluster, func(r *rabbitmqv1beta1.RabbitmqCluster) {
if r.Annotations == nil {
r.Annotations = make(map[string]string)
}
r.Annotations[rabbitmqv1beta1.RabbitmqVersionAnnotation] = "4.1.5"
})).To(Succeed())

// Verify Role and RoleBinding are deleted
Eventually(func() bool {
_, err := clientSet.RbacV1().Roles(cluster.Namespace).Get(ctx, cluster.ChildResourceName("peer-discovery"), metav1.GetOptions{})
return apierrors.IsNotFound(err)
}, 5).Should(BeTrueBecause("Role should be deleted when version >= 4.1.0"))

Eventually(func() bool {
_, err := clientSet.RbacV1().RoleBindings(cluster.Namespace).Get(ctx, cluster.ChildResourceName("server"), metav1.GetOptions{})
return apierrors.IsNotFound(err)
}, 5).Should(BeTrueBecause("RoleBinding should be deleted when version >= 4.1.0"))

// Verify the ServiceAccount is kept (other integrations such as Vault Kubernetes auth may rely on it)
Consistently(func() bool {
_, err := clientSet.CoreV1().ServiceAccounts(cluster.Namespace).Get(ctx, cluster.ChildResourceName("server"), metav1.GetOptions{})
return err == nil
}, 3).Should(BeTrueBecause("ServiceAccount should be kept when version >= 4.1.0"))
})
})

It("service type is updated", func() {
Expect(updateWithRetry(cluster, func(r *rabbitmqv1beta1.RabbitmqCluster) {
r.Spec.Service.Type = "NodePort"
Expand Down
53 changes: 49 additions & 4 deletions internal/resource/rabbitmq_resource_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,12 @@
package resource

import (
"slices"

"github.com/Masterminds/semver/v3"
rabbitmqv1beta1 "github.com/rabbitmq/cluster-operator/v2/api/v1beta1"
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"slices"
)

type RabbitmqResourceBuilder struct {
Expand All @@ -27,6 +29,39 @@ type ResourceBuilder interface {
UpdateMayRequireStsRecreate() bool
}

// peerDiscoveryRBACConstraint is the semver constraint used to determine whether
// peer-discovery RBAC (Role and RoleBinding) should be created. It is created
// once at package scope to avoid repeated allocations.
var peerDiscoveryRBACConstraint = mustNewConstraint(">= 4.1.0")

// mustNewConstraint creates a semver constraint and panics if parsing fails.
// This is only used for constant constraint strings that are always valid.
func mustNewConstraint(c string) *semver.Constraints {
constraint, err := semver.NewConstraint(c)
if err != nil {
panic(err)
}
return constraint
}

// ShouldCreatePeerDiscoveryRBAC returns true if the peer-discovery Role and
// RoleBinding should be created for this RabbitmqCluster. The ServiceAccount is
// always created regardless of RabbitMQ version because other integrations (e.g.
// Vault Kubernetes auth) may rely on it.
func ShouldCreatePeerDiscoveryRBAC(rmq *rabbitmqv1beta1.RabbitmqCluster) bool {
version := rmq.GetRabbitMQVersion()
if version == rabbitmqv1beta1.VersionNotAnnotated {
return true
}

v, err := semver.NewVersion(version)
if err != nil {
return true
}

return !peerDiscoveryRBACConstraint.Check(v)
}

func (builder *RabbitmqResourceBuilder) ResourceBuilders() []ResourceBuilder {

builders := []ResourceBuilder{
Expand All @@ -37,10 +72,20 @@ func (builder *RabbitmqResourceBuilder) ResourceBuilders() []ResourceBuilder {
builder.RabbitmqPluginsConfigMap(),
builder.ServerConfigMap(),
builder.ServiceAccount(),
builder.Role(),
builder.RoleBinding(),
builder.StatefulSet(),
}

if ShouldCreatePeerDiscoveryRBAC(builder.Instance) {
builders = append(builders,
builder.Role(),
builder.RoleBinding(),
)
}

// Appending StatefulSet builder separately because the order of the builders is important
// The SA, ConfigMap, and Secret need to be created before the StatefulSet. Otherwise, Pods
// created by the StatefulSet will block on the creation of dependent resources.
builders = append(builders, builder.StatefulSet())

if builder.Instance.VaultDefaultUserSecretEnabled() || builder.Instance.ExternalSecretEnabled() {
// do not create default-user K8s Secret
builders = slices.Delete(builders, 3, 3+1)
Expand Down
60 changes: 60 additions & 0 deletions internal/resource/rabbitmq_resource_builder_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,51 @@ import (
)

var _ = Describe("RabbitmqResourceBuilder", func() {
Context("ShouldCreatePeerDiscoveryRBAC", func() {
It("returns true if version is not annotated", func() {
rmq := &rabbitmqv1beta1.RabbitmqCluster{
ObjectMeta: v1.ObjectMeta{Annotations: map[string]string{}},
}
Expect(resource.ShouldCreatePeerDiscoveryRBAC(rmq)).To(BeTrueBecause("fallback to old behavior when version is not annotated"))
})

It("returns true if version cannot be parsed", func() {
rmq := &rabbitmqv1beta1.RabbitmqCluster{
ObjectMeta: v1.ObjectMeta{Annotations: map[string]string{
rabbitmqv1beta1.RabbitmqVersionAnnotation: "invalid",
}},
}
Expect(resource.ShouldCreatePeerDiscoveryRBAC(rmq)).To(BeTrueBecause("fallback to old behavior when version cannot be parsed"))
})

It("returns true if version is less than 4.1.0", func() {
rmq := &rabbitmqv1beta1.RabbitmqCluster{
ObjectMeta: v1.ObjectMeta{Annotations: map[string]string{
rabbitmqv1beta1.RabbitmqVersionAnnotation: "3.13.0",
}},
}
Expect(resource.ShouldCreatePeerDiscoveryRBAC(rmq)).To(BeTrueBecause("version is less than 4.1.0"))

rmq.Annotations[rabbitmqv1beta1.RabbitmqVersionAnnotation] = "4.0.0"
Expect(resource.ShouldCreatePeerDiscoveryRBAC(rmq)).To(BeTrueBecause("version is less than 4.1.0"))
})

It("returns false if version is 4.1.0 or greater", func() {
rmq := &rabbitmqv1beta1.RabbitmqCluster{
ObjectMeta: v1.ObjectMeta{Annotations: map[string]string{
rabbitmqv1beta1.RabbitmqVersionAnnotation: "4.1.0",
}},
}
Expect(resource.ShouldCreatePeerDiscoveryRBAC(rmq)).To(BeFalseBecause("peer-discovery RBAC is no longer required for 4.1.0 or greater"))

rmq.Annotations[rabbitmqv1beta1.RabbitmqVersionAnnotation] = "4.1.5"
Expect(resource.ShouldCreatePeerDiscoveryRBAC(rmq)).To(BeFalseBecause("peer-discovery RBAC is no longer required for 4.1.0 or greater"))

rmq.Annotations[rabbitmqv1beta1.RabbitmqVersionAnnotation] = "4.2.0"
Expect(resource.ShouldCreatePeerDiscoveryRBAC(rmq)).To(BeFalseBecause("peer-discovery RBAC is no longer required for 4.1.0 or greater"))
})
})

Context("ResourceBuilders", func() {
var (
instance *rabbitmqv1beta1.RabbitmqCluster
Expand Down Expand Up @@ -79,5 +124,20 @@ var _ = Describe("RabbitmqResourceBuilder", func() {
Expect(resourceBuilders).NotTo(ContainElement(BeAssignableToTypeOf(&resource.DefaultUserSecretBuilder{})))
})
})

When("RabbitMQ version is 4.1.0 or greater", func() {
BeforeEach(func() {
instance.Annotations = map[string]string{
rabbitmqv1beta1.RabbitmqVersionAnnotation: "4.1.0",
}
})
It("returns all resource builders except for peer-discovery Role and RoleBinding", func() {
resourceBuilders := builder.ResourceBuilders()
Expect(resourceBuilders).To(HaveLen(8))
Expect(resourceBuilders).To(ContainElement(BeAssignableToTypeOf(&resource.ServiceAccountBuilder{})))
Expect(resourceBuilders).NotTo(ContainElement(BeAssignableToTypeOf(&resource.RoleBuilder{})))
Expect(resourceBuilders).NotTo(ContainElement(BeAssignableToTypeOf(&resource.RoleBindingBuilder{})))
})
})
})
})
6 changes: 4 additions & 2 deletions internal/resource/statefulset.go
Original file line number Diff line number Diff line change
Expand Up @@ -593,8 +593,6 @@ func (builder *StatefulSetBuilder) podTemplateSpec(previousPodAnnotations map[st
},
ImagePullSecrets: builder.Instance.Spec.ImagePullSecrets,
TerminationGracePeriodSeconds: builder.Instance.Spec.TerminationGracePeriodSeconds,
ServiceAccountName: builder.Instance.ChildResourceName(serviceAccountName),
AutomountServiceAccountToken: ptr.To(true),
Affinity: builder.Instance.Spec.Affinity,
Tolerations: builder.Instance.Spec.Tolerations,
InitContainers: []corev1.Container{setupContainer(builder.Instance)},
Expand Down Expand Up @@ -693,6 +691,10 @@ func (builder *StatefulSetBuilder) podTemplateSpec(previousPodAnnotations map[st
podTemplateSpec.Spec.Containers = append(podTemplateSpec.Spec.Containers,
defaultUserCredentialUpdater(builder.Instance))
}

podTemplateSpec.Spec.ServiceAccountName = builder.Instance.ChildResourceName(serviceAccountName)
podTemplateSpec.Spec.AutomountServiceAccountToken = ptr.To(true)

return podTemplateSpec
}

Expand Down
22 changes: 22 additions & 0 deletions internal/resource/statefulset_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1415,6 +1415,28 @@ default_pass = {{ .Data.data.password }}
Expect(*statefulSet.Spec.Template.Spec.AutomountServiceAccountToken).To(BeTrue())
})

When("RabbitMQ version is 4.1.0 or greater", func() {
BeforeEach(func() {
instance.Annotations = map[string]string{
rabbitmqv1beta1.RabbitmqVersionAnnotation: "4.1.0",
}
})

It("still uses the correct service account", func() {
stsBuilder := builder.StatefulSet()
Expect(stsBuilder.Update(statefulSet)).To(Succeed())

Expect(statefulSet.Spec.Template.Spec.ServiceAccountName).To(Equal(instance.ChildResourceName("server")))
})

It("still mounts the service account token in its pods", func() {
stsBuilder := builder.StatefulSet()
Expect(stsBuilder.Update(statefulSet)).To(Succeed())

Expect(*statefulSet.Spec.Template.Spec.AutomountServiceAccountToken).To(BeTrue())
})
})

It("creates the required SecurityContext", func() {
stsBuilder := builder.StatefulSet()
Expect(stsBuilder.Update(statefulSet)).To(Succeed())
Expand Down
4 changes: 2 additions & 2 deletions test/system/utils_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -654,7 +654,7 @@ func waitForRabbitmqNotRunningWithOffset(cluster *rabbitmqv1beta1.RabbitmqCluste
}

return string(output)
}, podCreationTimeout, 1).Should(Equal("'False'"))
}, podCreationTimeout, 2).Should(Equal("'False'"))

ExpectWithOffset(callStackOffset, err).NotTo(HaveOccurred())
}
Expand All @@ -679,7 +679,7 @@ func waitForRabbitmqRunningWithOffset(cluster *rabbitmqv1beta1.RabbitmqCluster,
}

return string(output)
}, podCreationTimeout, 1).Should(Equal("'True'"))
}, podCreationTimeout, 2).Should(Equal("'True'"))

ExpectWithOffset(callStackOffset, err).NotTo(HaveOccurred())
}
Expand Down
Loading