From 06530d04676687d60930c99607d1263b3dca02a2 Mon Sep 17 00:00:00 2001 From: Armando Ruocco Date: Thu, 17 Sep 2026 12:47:25 +0200 Subject: [PATCH] feat: support running the plugin with multiple replicas The gRPC server runnable was added to the controller-runtime manager without implementing LeaderElectionRunnable, so it only ran on the leader. With more than one replica the other pods never listened on the plugin port, failed the readiness probe and stayed NotReady, and the Service could only ever route to the leader. The gRPC handlers only read from the informer cache, so every replica can serve them. Make the runnable opt out of leader election; the ObjectStore controller keeps using leader election as before. Add an e2e test that scales the deployment to two replicas, checks both become ready, deletes the leader pod and verifies that a new leader is elected and clusters still reconcile. Document how to run multiple replicas. Closes #620 Signed-off-by: Armando Ruocco --- internal/cnpgi/operator/start.go | 17 ++ internal/cnpgi/operator/start_test.go | 31 +++ test/e2e/e2e_suite_test.go | 1 + .../internal/tests/highavailability/doc.go | 23 ++ .../highavailability/high_availability.go | 228 ++++++++++++++++++ web/docs/installation.mdx | 20 ++ 6 files changed, 320 insertions(+) create mode 100644 internal/cnpgi/operator/start_test.go create mode 100644 test/e2e/internal/tests/highavailability/doc.go create mode 100644 test/e2e/internal/tests/highavailability/high_availability.go diff --git a/internal/cnpgi/operator/start.go b/internal/cnpgi/operator/start.go index fc574e8e..8ebadcf5 100644 --- a/internal/cnpgi/operator/start.go +++ b/internal/cnpgi/operator/start.go @@ -27,6 +27,7 @@ import ( "github.com/cloudnative-pg/cnpg-i/pkg/reconciler" "google.golang.org/grpc" "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/manager" ) // CNPGI is the implementation of the CNPG-i server @@ -39,6 +40,22 @@ type CNPGI struct { ServerAddress string } +// controller-runtime decides how to start a runnable passed to mgr.Add by +// type-asserting it against LeaderElectionRunnable: a runnable that does not +// implement it is started only on the leader, one that returns false here is +// started on every replica right after the caches have synced. +// +// The gRPC handlers only read from the informer cache, so every replica can +// serve them; this lets the Service load-balance across pods and non-leaders +// pass the readiness probe. The ObjectStore controller keeps running on the +// leader only. +var _ manager.LeaderElectionRunnable = &CNPGI{} + +// NeedLeaderElection implements manager.LeaderElectionRunnable. +func (c *CNPGI) NeedLeaderElection() bool { + return false +} + // Start starts the GRPC server // of the operator plugin func (c *CNPGI) Start(ctx context.Context) error { diff --git a/internal/cnpgi/operator/start_test.go b/internal/cnpgi/operator/start_test.go new file mode 100644 index 00000000..a1dc2c40 --- /dev/null +++ b/internal/cnpgi/operator/start_test.go @@ -0,0 +1,31 @@ +/* +Copyright © contributors to CloudNativePG, established as +CloudNativePG a Series of LF Projects, LLC. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. + +SPDX-License-Identifier: Apache-2.0 +*/ + +package operator + +import ( + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var _ = Describe("CNPGI runnable", func() { + It("runs on every replica instead of only on the leader", func() { + Expect((&CNPGI{}).NeedLeaderElection()).To(BeFalse()) + }) +}) diff --git a/test/e2e/e2e_suite_test.go b/test/e2e/e2e_suite_test.go index 0f2bfdf7..635ed307 100644 --- a/test/e2e/e2e_suite_test.go +++ b/test/e2e/e2e_suite_test.go @@ -38,6 +38,7 @@ import ( _ "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/tests/backup" _ "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/tests/credentialrotation" + _ "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/tests/highavailability" _ "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/tests/replicacluster" _ "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/tests/walrestore" diff --git a/test/e2e/internal/tests/highavailability/doc.go b/test/e2e/internal/tests/highavailability/doc.go new file mode 100644 index 00000000..d9a6bac8 --- /dev/null +++ b/test/e2e/internal/tests/highavailability/doc.go @@ -0,0 +1,23 @@ +/* +Copyright © contributors to CloudNativePG, established as +CloudNativePG a Series of LF Projects, LLC. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. + +SPDX-License-Identifier: Apache-2.0 +*/ + +// Package highavailability contains tests validating that the Barman Cloud +// Plugin Deployment keeps serving requests and reconciling when it runs +// with multiple replicas and the leader is lost. +package highavailability diff --git a/test/e2e/internal/tests/highavailability/high_availability.go b/test/e2e/internal/tests/highavailability/high_availability.go new file mode 100644 index 00000000..b75233d9 --- /dev/null +++ b/test/e2e/internal/tests/highavailability/high_availability.go @@ -0,0 +1,228 @@ +/* +Copyright © contributors to CloudNativePG, established as +CloudNativePG a Series of LF Projects, LLC. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. + +SPDX-License-Identifier: Apache-2.0 +*/ + +package highavailability + +import ( + "fmt" + "strings" + "time" + + cloudnativepgv1 "github.com/cloudnative-pg/api/pkg/api/v1" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/kubernetes" + "k8s.io/utils/ptr" + "sigs.k8s.io/controller-runtime/pkg/client" + + internalClient "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/client" + internalCluster "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/cluster" + nmsp "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/namespace" + "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/objectstore" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +const ( + pluginNamespace = "cnpg-system" + pluginDeployment = "barman-cloud" + leaderElectionID = "822e3f5c.cnpg.io" + + objectStoreName = "source" + minioSecretName = "minio" + clusterName = "source" + secondClusterName = "second" +) + +var _ = Describe("Plugin high availability", Serial, func() { + var namespace *corev1.Namespace + var cl client.Client + var clientSet *kubernetes.Clientset + + BeforeEach(func(ctx SpecContext) { + var err error + cl, _, err = internalClient.NewClient() + Expect(err).NotTo(HaveOccurred()) + clientSet, _, err = internalClient.NewClientSet() + Expect(err).NotTo(HaveOccurred()) + namespace, err = nmsp.CreateUniqueNamespace(ctx, cl, "plugin-ha") + Expect(err).NotTo(HaveOccurred()) + }) + + AfterEach(func(ctx SpecContext) { + Expect(cl.Delete(ctx, namespace)).To(Succeed()) + Expect(scaleDeployment(ctx, cl, 1)).To(Succeed()) + Eventually(func(g Gomega) { + var deploy appsv1.Deployment + g.Expect(cl.Get(ctx, types.NamespacedName{ + Name: pluginDeployment, + Namespace: pluginNamespace, + }, &deploy)).To(Succeed()) + g.Expect(deploy.Status.Replicas).To(BeEquivalentTo(1)) + g.Expect(deploy.Status.ReadyReplicas).To(BeEquivalentTo(1)) + }).WithTimeout(2 * time.Minute).WithPolling(5 * time.Second).Should(Succeed()) + }) + + It("should serve requests from every replica and survive losing the leader", func(ctx SpecContext) { + By("scaling the plugin Deployment to 2 replicas") + Expect(scaleDeployment(ctx, cl, 2)).To(Succeed()) + + By("waiting for both replicas to become ready") + Eventually(func(g Gomega) { + var deploy appsv1.Deployment + g.Expect(cl.Get(ctx, types.NamespacedName{ + Name: pluginDeployment, + Namespace: pluginNamespace, + }, &deploy)).To(Succeed()) + g.Expect(deploy.Status.ReadyReplicas).To(BeEquivalentTo(2)) + g.Expect(deploy.Status.UpdatedReplicas).To(BeEquivalentTo(2)) + }).WithTimeout(2 * time.Minute).WithPolling(5 * time.Second).Should(Succeed()) + + By("finding the current leader") + leaderPodName, err := getLeaderPodName(ctx, clientSet) + Expect(err).NotTo(HaveOccurred()) + Expect(podExists(ctx, cl, leaderPodName)).To(BeTrue()) + + By("starting the ObjectStore deployment") + resources := objectstore.NewMinioObjectStoreResources(namespace.Name, minioSecretName) + Expect(resources.Create(ctx, cl)).To(Succeed()) + + By("creating the ObjectStore") + store := objectstore.NewMinioObjectStore(namespace.Name, objectStoreName, minioSecretName) + Expect(cl.Create(ctx, store)).To(Succeed()) + + By("creating the Cluster") + cluster := newCluster(namespace.Name, clusterName, objectStoreName) + Expect(cl.Create(ctx, cluster)).To(Succeed()) + + By("waiting for the Cluster to be ready, exercising gRPC across both replicas") + waitForClusterReady(ctx, cl, cluster) + + By("deleting the leader Pod") + Expect(cl.Delete(ctx, &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: leaderPodName, + Namespace: pluginNamespace, + }, + })).To(Succeed()) + + By("waiting for a new leader to be elected and both replicas to be ready again") + Eventually(func(g Gomega) { + newLeaderPodName, err := getLeaderPodName(ctx, clientSet) + g.Expect(err).NotTo(HaveOccurred()) + g.Expect(newLeaderPodName).NotTo(Equal(leaderPodName)) + g.Expect(podExists(ctx, cl, newLeaderPodName)).To(BeTrue()) + + var deploy appsv1.Deployment + g.Expect(cl.Get(ctx, types.NamespacedName{ + Name: pluginDeployment, + Namespace: pluginNamespace, + }, &deploy)).To(Succeed()) + g.Expect(deploy.Status.ReadyReplicas).To(BeEquivalentTo(2)) + }).WithTimeout(3 * time.Minute).WithPolling(5 * time.Second).Should(Succeed()) + + By("verifying reconciliation still works after the failover") + secondCluster := newCluster(namespace.Name, secondClusterName, objectStoreName) + Expect(cl.Create(ctx, secondCluster)).To(Succeed()) + waitForClusterReady(ctx, cl, secondCluster) + }) +}) + +func scaleDeployment(ctx SpecContext, cl client.Client, replicas int32) error { + var deploy appsv1.Deployment + if err := cl.Get(ctx, types.NamespacedName{ + Name: pluginDeployment, + Namespace: pluginNamespace, + }, &deploy); err != nil { + return err //nolint:wrapcheck + } + deploy.Spec.Replicas = ptr.To(replicas) + + return cl.Update(ctx, &deploy) //nolint:wrapcheck +} + +// getLeaderPodName reads the leader election Lease and returns the name of +// the Pod currently holding it. HolderIdentity has the form _. +func getLeaderPodName(ctx SpecContext, clientSet *kubernetes.Clientset) (string, error) { + lease, err := clientSet.CoordinationV1().Leases(pluginNamespace).Get(ctx, leaderElectionID, metav1.GetOptions{}) + if err != nil { + return "", fmt.Errorf("failed to get %s lease: %w", leaderElectionID, err) + } + if lease.Spec.HolderIdentity == nil || *lease.Spec.HolderIdentity == "" { + return "", nil + } + + return strings.Split(*lease.Spec.HolderIdentity, "_")[0], nil +} + +func podExists(ctx SpecContext, cl client.Client, name string) bool { + var pod corev1.Pod + err := cl.Get(ctx, types.NamespacedName{Name: name, Namespace: pluginNamespace}, &pod) + if err != nil { + if apierrors.IsNotFound(err) { + return false + } + Fail(err.Error()) + } + + return true +} + +func waitForClusterReady(ctx SpecContext, cl client.Client, cluster *cloudnativepgv1.Cluster) { + Eventually(func(g Gomega) { + g.Expect(cl.Get(ctx, types.NamespacedName{ + Name: cluster.Name, + Namespace: cluster.Namespace, + }, cluster)).To(Succeed()) + g.Expect(internalCluster.IsReady(*cluster)).To(BeTrue()) + }).WithTimeout(10 * time.Minute).WithPolling(10 * time.Second).Should(Succeed()) +} + +func newCluster(namespace, name, objectStore string) *cloudnativepgv1.Cluster { + return &cloudnativepgv1.Cluster{ + TypeMeta: metav1.TypeMeta{ + Kind: "Cluster", + APIVersion: "postgresql.cnpg.io/v1", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: namespace, + }, + Spec: cloudnativepgv1.ClusterSpec{ + Instances: 1, + ImagePullPolicy: corev1.PullAlways, + Plugins: []cloudnativepgv1.PluginConfiguration{ + { + Name: "barman-cloud.cloudnative-pg.io", + Parameters: map[string]string{ + "barmanObjectName": objectStore, + }, + IsWALArchiver: ptr.To(true), + }, + }, + StorageConfiguration: cloudnativepgv1.StorageConfiguration{ + Size: "1Gi", + }, + }, + } +} diff --git a/web/docs/installation.mdx b/web/docs/installation.mdx index bff100bc..cab150c1 100644 --- a/web/docs/installation.mdx +++ b/web/docs/installation.mdx @@ -141,6 +141,26 @@ deployment "barman-cloud" successfully rolled out This confirms that the plugin is deployed and ready to use. +## Running Multiple Replicas + +The plugin Deployment can run with more than one replica for high +availability. Every replica serves the gRPC interface used by the CloudNativePG +operator, so the `barman-cloud` Service load-balances requests across all +ready pods. The `ObjectStore` controller uses leader election, so only one +replica reconciles `ObjectStore` resources at a time; the others take over +automatically if the leader is lost. + +To scale the plugin: + +```sh +kubectl -n cnpg-system scale deploy/barman-cloud --replicas=2 +``` + +The Deployment keeps the `Recreate` strategy: during an upgrade all replicas +are replaced together, so the plugin is briefly unavailable regardless of the +replica count. Multiple replicas protect against pod or node loss, not against +upgrade downtime. + ## Testing the latest development snapshot import { DevSnapshotSection } from '@site/src/components/Installation';