mirror of
https://github.com/cloudnative-pg/plugin-barman-cloud.git
synced 2026-09-05 22:52:19 +02:00
Merge fdc056a867 into 561289e434
This commit is contained in:
commit
dffb33e7c5
@ -52,6 +52,11 @@ const (
|
|||||||
// BarmanEndpointCACertificateFileName is the name of the file in which the barman endpoint
|
// BarmanEndpointCACertificateFileName is the name of the file in which the barman endpoint
|
||||||
// CA certificate is stored.
|
// CA certificate is stored.
|
||||||
BarmanEndpointCACertificateFileName = "barman-ca.crt"
|
BarmanEndpointCACertificateFileName = "barman-ca.crt"
|
||||||
|
|
||||||
|
// PgWalVolumePgWalPath is the path of the pg_wal directory inside the WAL volume,
|
||||||
|
// used when a separate WAL storage is configured. During a restore the pg_wal
|
||||||
|
// directory is moved here and symlinked back into PGDATA.
|
||||||
|
PgWalVolumePgWalPath = "/var/lib/postgresql/wal/pg_wal"
|
||||||
)
|
)
|
||||||
|
|
||||||
// GetRestoreCABundleEnv gets the enveronment variables to be used when custom
|
// GetRestoreCABundleEnv gets the enveronment variables to be used when custom
|
||||||
|
|||||||
@ -70,6 +70,13 @@ func (i IdentityImplementation) GetPluginCapabilities(
|
|||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
{
|
||||||
|
Type: &identity.PluginCapability_Service_{
|
||||||
|
Service: &identity.PluginCapability_Service{
|
||||||
|
Type: identity.PluginCapability_Service_TYPE_RESTORE_JOB,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
},
|
},
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|||||||
52
internal/cnpgi/instance/identity_test.go
Normal file
52
internal/cnpgi/instance/identity_test.go
Normal file
@ -0,0 +1,52 @@
|
|||||||
|
/*
|
||||||
|
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 instance
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/cloudnative-pg/cnpg-i/pkg/identity"
|
||||||
|
|
||||||
|
. "github.com/onsi/ginkgo/v2"
|
||||||
|
. "github.com/onsi/gomega"
|
||||||
|
)
|
||||||
|
|
||||||
|
var _ = Describe("IdentityImplementation", func() {
|
||||||
|
Describe("GetPluginCapabilities", func() {
|
||||||
|
It("declares the WAL, backup, metrics and restore-job services", func(ctx SpecContext) {
|
||||||
|
impl := IdentityImplementation{}
|
||||||
|
response, err := impl.GetPluginCapabilities(ctx, &identity.GetPluginCapabilitiesRequest{})
|
||||||
|
Expect(err).NotTo(HaveOccurred())
|
||||||
|
Expect(response).NotTo(BeNil())
|
||||||
|
|
||||||
|
var serviceTypes []identity.PluginCapability_Service_Type
|
||||||
|
for _, capability := range response.GetCapabilities() {
|
||||||
|
serviceTypes = append(serviceTypes, capability.GetService().GetType())
|
||||||
|
}
|
||||||
|
|
||||||
|
// The instance sidecar now runs the phase-0 restore in-process, so it must
|
||||||
|
// advertise TYPE_RESTORE_JOB alongside the services it already served.
|
||||||
|
Expect(serviceTypes).To(ConsistOf(
|
||||||
|
identity.PluginCapability_Service_TYPE_WAL_SERVICE,
|
||||||
|
identity.PluginCapability_Service_TYPE_BACKUP_SERVICE,
|
||||||
|
identity.PluginCapability_Service_TYPE_METRICS,
|
||||||
|
identity.PluginCapability_Service_TYPE_RESTORE_JOB,
|
||||||
|
))
|
||||||
|
})
|
||||||
|
})
|
||||||
|
})
|
||||||
@ -25,11 +25,13 @@ import (
|
|||||||
"github.com/cloudnative-pg/cnpg-i-machinery/pkg/pluginhelper/http"
|
"github.com/cloudnative-pg/cnpg-i-machinery/pkg/pluginhelper/http"
|
||||||
"github.com/cloudnative-pg/cnpg-i/pkg/backup"
|
"github.com/cloudnative-pg/cnpg-i/pkg/backup"
|
||||||
"github.com/cloudnative-pg/cnpg-i/pkg/metrics"
|
"github.com/cloudnative-pg/cnpg-i/pkg/metrics"
|
||||||
|
restore "github.com/cloudnative-pg/cnpg-i/pkg/restore/job"
|
||||||
"github.com/cloudnative-pg/cnpg-i/pkg/wal"
|
"github.com/cloudnative-pg/cnpg-i/pkg/wal"
|
||||||
"google.golang.org/grpc"
|
"google.golang.org/grpc"
|
||||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
|
||||||
"github.com/cloudnative-pg/plugin-barman-cloud/internal/cnpgi/common"
|
"github.com/cloudnative-pg/plugin-barman-cloud/internal/cnpgi/common"
|
||||||
|
barmanrestore "github.com/cloudnative-pg/plugin-barman-cloud/internal/cnpgi/restore"
|
||||||
)
|
)
|
||||||
|
|
||||||
// CNPGI is the implementation of the PostgreSQL sidecar
|
// CNPGI is the implementation of the PostgreSQL sidecar
|
||||||
@ -60,6 +62,15 @@ func (c *CNPGI) Start(ctx context.Context) error {
|
|||||||
metrics.RegisterMetricsServer(server, &metricsImpl{
|
metrics.RegisterMetricsServer(server, &metricsImpl{
|
||||||
Client: c.Client,
|
Client: c.Client,
|
||||||
})
|
})
|
||||||
|
// The instance pod runs the phase-0 bootstrap in-process (no separate
|
||||||
|
// recovery Job), so the same sidecar must answer the Restore RPC that
|
||||||
|
// initializes PGDATA from the object store before PostgreSQL starts.
|
||||||
|
restore.RegisterRestoreJobHooksServer(server, &barmanrestore.JobHookImpl{
|
||||||
|
Client: c.Client,
|
||||||
|
SpoolDirectory: c.SpoolDirectory,
|
||||||
|
PgDataPath: c.PGDataPath,
|
||||||
|
PgWalFolderToSymlink: common.PgWalVolumePgWalPath,
|
||||||
|
})
|
||||||
common.AddHealthCheck(server)
|
common.AddHealthCheck(server)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@ -102,6 +102,15 @@ func (config *PluginConfiguration) GetReplicaSourceBarmanObjectKey() types.Names
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// HasAnyBarmanObjectStore returns true if the configuration references at least
|
||||||
|
// one barman object store, be it for backup/archiving, recovery, or as a
|
||||||
|
// replica source.
|
||||||
|
func (config *PluginConfiguration) HasAnyBarmanObjectStore() bool {
|
||||||
|
return len(config.BarmanObjectName) > 0 ||
|
||||||
|
len(config.RecoveryBarmanObjectName) > 0 ||
|
||||||
|
len(config.ReplicaSourceBarmanObjectName) > 0
|
||||||
|
}
|
||||||
|
|
||||||
// GetReferredBarmanObjectsKey gets the list of barman objects referred by this
|
// GetReferredBarmanObjectsKey gets the list of barman objects referred by this
|
||||||
// plugin configuration
|
// plugin configuration
|
||||||
func (config *PluginConfiguration) GetReferredBarmanObjectsKey() []types.NamespacedName {
|
func (config *PluginConfiguration) GetReferredBarmanObjectsKey() []types.NamespacedName {
|
||||||
@ -263,9 +272,7 @@ func getReplicaSourcePlugin(cluster *cnpgv1.Cluster) *cnpgv1.PluginConfiguration
|
|||||||
func (config *PluginConfiguration) Validate() error {
|
func (config *PluginConfiguration) Validate() error {
|
||||||
err := NewConfigurationError()
|
err := NewConfigurationError()
|
||||||
|
|
||||||
if len(config.BarmanObjectName) == 0 &&
|
if !config.HasAnyBarmanObjectStore() {
|
||||||
len(config.RecoveryBarmanObjectName) == 0 &&
|
|
||||||
len(config.ReplicaSourceBarmanObjectName) == 0 {
|
|
||||||
return err.WithMessage("no reference to barmanObjectName have been included")
|
return err.WithMessage("no reference to barmanObjectName have been included")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -322,6 +322,45 @@ func (impl LifecycleImplementation) collectAdditionalInstanceArgs(
|
|||||||
return nil, nil
|
return nil, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// shouldInjectBarmanSidecar decides whether an instance pod needs the
|
||||||
|
// plugin-barman-cloud sidecar.
|
||||||
|
//
|
||||||
|
// A cluster doing backup/archiving or serving as a replica source needs the
|
||||||
|
// sidecar in every instance pod for as long as the cluster exists, so those
|
||||||
|
// two cases always inject it. A recovery-only cluster (only
|
||||||
|
// RecoveryBarmanObjectName set, mirroring what pluginConfiguration.Validate()
|
||||||
|
// accepts) only ever needs the sidecar for its one-time bootstrap restore, so
|
||||||
|
// it's gated on cluster.Status.CurrentPrimary instead.
|
||||||
|
//
|
||||||
|
// CurrentPrimary is set by the instance manager itself, from inside the
|
||||||
|
// primary pod, only once it's up and has completed its own bootstrap (see
|
||||||
|
// instance_startup.go in cloudnative-pg) - unlike cluster.Status.Instances /
|
||||||
|
// IsInitialized(), which flips as soon as the instance's PVC exists, well
|
||||||
|
// before the pod is even created. Using IsInitialized() here would mean the
|
||||||
|
// sidecar is never added to the one pod that needs it to perform its restore.
|
||||||
|
//
|
||||||
|
// Once CurrentPrimary is set, this makes the operator's own drift-check
|
||||||
|
// (checkPodSpecIsOutdated) see the running pod's spec as outdated and roll it
|
||||||
|
// out to drop the sidecar. That's deliberately accepted rather than
|
||||||
|
// engineered around: it's one deterministic rollout using the same machinery
|
||||||
|
// the operator already uses for every other pod-spec change (a switchover if
|
||||||
|
// a replica is available, an in-place restart otherwise), not a new or
|
||||||
|
// fragile risk.
|
||||||
|
func shouldInjectBarmanSidecar(
|
||||||
|
cluster *cnpgv1.Cluster,
|
||||||
|
pluginConfiguration *config.PluginConfiguration,
|
||||||
|
) bool {
|
||||||
|
if len(pluginConfiguration.BarmanObjectName) != 0 || len(pluginConfiguration.ReplicaSourceBarmanObjectName) != 0 {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(pluginConfiguration.RecoveryBarmanObjectName) == 0 {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
return cluster.Status.CurrentPrimary == ""
|
||||||
|
}
|
||||||
|
|
||||||
func reconcileInstancePod(
|
func reconcileInstancePod(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
cluster *cnpgv1.Cluster,
|
cluster *cnpgv1.Cluster,
|
||||||
@ -339,8 +378,7 @@ func reconcileInstancePod(
|
|||||||
|
|
||||||
mutatedPod := pod.DeepCopy()
|
mutatedPod := pod.DeepCopy()
|
||||||
|
|
||||||
if len(pluginConfiguration.BarmanObjectName) != 0 ||
|
if shouldInjectBarmanSidecar(cluster, pluginConfiguration) {
|
||||||
len(pluginConfiguration.ReplicaSourceBarmanObjectName) != 0 {
|
|
||||||
if err := reconcilePodSpec(
|
if err := reconcilePodSpec(
|
||||||
cluster,
|
cluster,
|
||||||
&mutatedPod.Spec,
|
&mutatedPod.Spec,
|
||||||
@ -353,7 +391,7 @@ func reconcileInstancePod(
|
|||||||
return nil, fmt.Errorf("while reconciling pod spec for pod: %w", err)
|
return nil, fmt.Errorf("while reconciling pod spec for pod: %w", err)
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
contextLogger.Debug("No need to mutate instance with no backup & archiving configuration")
|
contextLogger.Debug("No need to mutate instance, sidecar not required for this configuration and pod")
|
||||||
}
|
}
|
||||||
|
|
||||||
patch, err := object.CreatePatch(mutatedPod, pod)
|
patch, err := object.CreatePatch(mutatedPod, pod)
|
||||||
|
|||||||
@ -242,6 +242,70 @@ var _ = Describe("LifecycleImplementation", func() {
|
|||||||
HaveKey("value")))
|
HaveKey("value")))
|
||||||
})
|
})
|
||||||
|
|
||||||
|
It("injects the sidecar for a recovery-only cluster", func(ctx SpecContext) {
|
||||||
|
recoveryOnlyConfig := &config.PluginConfiguration{
|
||||||
|
RecoveryBarmanObjectName: "minio-store-recovery",
|
||||||
|
}
|
||||||
|
pod := &corev1.Pod{
|
||||||
|
TypeMeta: podTypeMeta,
|
||||||
|
ObjectMeta: metav1.ObjectMeta{Name: "test-pod"},
|
||||||
|
Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "postgres"}}},
|
||||||
|
}
|
||||||
|
podJSON, _ := json.Marshal(pod)
|
||||||
|
request := &lifecycle.OperatorLifecycleRequest{
|
||||||
|
ObjectDefinition: podJSON,
|
||||||
|
}
|
||||||
|
|
||||||
|
response, err := reconcileInstancePod(ctx, cluster, request, recoveryOnlyConfig, sidecarConfiguration{})
|
||||||
|
Expect(err).NotTo(HaveOccurred())
|
||||||
|
Expect(response).NotTo(BeNil())
|
||||||
|
Expect(response.JsonPatch).NotTo(BeEmpty())
|
||||||
|
var patch []map[string]interface{}
|
||||||
|
Expect(json.Unmarshal(response.JsonPatch, &patch)).To(Succeed())
|
||||||
|
Expect(patch).To(ContainElement(HaveKeyWithValue("path", "/spec/initContainers")))
|
||||||
|
})
|
||||||
|
|
||||||
|
It("does not inject the sidecar for a recovery-only cluster that has "+
|
||||||
|
"already completed its initial bootstrap", func(ctx SpecContext) {
|
||||||
|
recoveryOnlyConfig := &config.PluginConfiguration{
|
||||||
|
RecoveryBarmanObjectName: "minio-store-recovery",
|
||||||
|
}
|
||||||
|
cluster.Status.CurrentPrimary = "test-pod"
|
||||||
|
pod := &corev1.Pod{
|
||||||
|
TypeMeta: podTypeMeta,
|
||||||
|
ObjectMeta: metav1.ObjectMeta{Name: "test-pod"},
|
||||||
|
Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "postgres"}}},
|
||||||
|
}
|
||||||
|
podJSON, _ := json.Marshal(pod)
|
||||||
|
request := &lifecycle.OperatorLifecycleRequest{
|
||||||
|
ObjectDefinition: podJSON,
|
||||||
|
}
|
||||||
|
|
||||||
|
response, err := reconcileInstancePod(ctx, cluster, request, recoveryOnlyConfig, sidecarConfiguration{})
|
||||||
|
Expect(err).NotTo(HaveOccurred())
|
||||||
|
Expect(response).NotTo(BeNil())
|
||||||
|
Expect(response.JsonPatch).To(BeEmpty())
|
||||||
|
})
|
||||||
|
|
||||||
|
It("does not mutate the pod when no object store is configured", func(ctx SpecContext) {
|
||||||
|
emptyConfig := &config.PluginConfiguration{}
|
||||||
|
pod := &corev1.Pod{
|
||||||
|
TypeMeta: podTypeMeta,
|
||||||
|
ObjectMeta: metav1.ObjectMeta{Name: "test-pod"},
|
||||||
|
Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "postgres"}}},
|
||||||
|
}
|
||||||
|
podJSON, _ := json.Marshal(pod)
|
||||||
|
request := &lifecycle.OperatorLifecycleRequest{
|
||||||
|
ObjectDefinition: podJSON,
|
||||||
|
}
|
||||||
|
|
||||||
|
response, err := reconcileInstancePod(ctx, cluster, request, emptyConfig, sidecarConfiguration{})
|
||||||
|
Expect(err).NotTo(HaveOccurred())
|
||||||
|
Expect(response).NotTo(BeNil())
|
||||||
|
// An empty patch means the pod was left untouched: no sidecar injected.
|
||||||
|
Expect(response.JsonPatch).To(BeEmpty())
|
||||||
|
})
|
||||||
|
|
||||||
It("returns an error for invalid pod definition", func(ctx SpecContext) {
|
It("returns an error for invalid pod definition", func(ctx SpecContext) {
|
||||||
request := &lifecycle.OperatorLifecycleRequest{
|
request := &lifecycle.OperatorLifecycleRequest{
|
||||||
ObjectDefinition: []byte("invalid-json"),
|
ObjectDefinition: []byte("invalid-json"),
|
||||||
|
|||||||
@ -44,9 +44,6 @@ type CNPGI struct {
|
|||||||
|
|
||||||
// Start starts the GRPC service
|
// Start starts the GRPC service
|
||||||
func (c *CNPGI) Start(ctx context.Context) error {
|
func (c *CNPGI) Start(ctx context.Context) error {
|
||||||
// PgWalVolumePgWalPath is the path of pg_wal directory inside the WAL volume when present
|
|
||||||
const PgWalVolumePgWalPath = "/var/lib/postgresql/wal/pg_wal"
|
|
||||||
|
|
||||||
enrich := func(server *grpc.Server) error {
|
enrich := func(server *grpc.Server) error {
|
||||||
wal.RegisterWALServer(server, common.WALServiceImplementation{
|
wal.RegisterWALServer(server, common.WALServiceImplementation{
|
||||||
InstanceName: c.InstanceName,
|
InstanceName: c.InstanceName,
|
||||||
@ -60,7 +57,7 @@ func (c *CNPGI) Start(ctx context.Context) error {
|
|||||||
Client: c.Client,
|
Client: c.Client,
|
||||||
SpoolDirectory: c.SpoolDirectory,
|
SpoolDirectory: c.SpoolDirectory,
|
||||||
PgDataPath: c.PGDataPath,
|
PgDataPath: c.PGDataPath,
|
||||||
PgWalFolderToSymlink: PgWalVolumePgWalPath,
|
PgWalFolderToSymlink: common.PgWalVolumePgWalPath,
|
||||||
})
|
})
|
||||||
|
|
||||||
common.AddHealthCheck(server)
|
common.AddHealthCheck(server)
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user