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
17 changes: 17 additions & 0 deletions internal/cnpgi/operator/start.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 {
Expand Down
31 changes: 31 additions & 0 deletions internal/cnpgi/operator/start_test.go
Original file line number Diff line number Diff line change
@@ -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())
})
})
1 change: 1 addition & 0 deletions test/e2e/e2e_suite_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand Down
23 changes: 23 additions & 0 deletions test/e2e/internal/tests/highavailability/doc.go
Original file line number Diff line number Diff line change
@@ -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
228 changes: 228 additions & 0 deletions test/e2e/internal/tests/highavailability/high_availability.go
Original file line number Diff line number Diff line change
@@ -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 <podName>_<uuid>.
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",
},
},
}
}
20 changes: 20 additions & 0 deletions web/docs/installation.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down
Loading