Compare commits

...

10 Commits

Author SHA1 Message Date
Armando Ruocco
0c70cd3dbe
Merge 6a55a361a3 into dd799d8ed1 2026-07-13 18:24:27 +00:00
Niccolò Fei
dd799d8ed1
ci: fix image tag mismatch when running via workflow_dispatch (#995)
Closes #994

Signed-off-by: Niccolò Fei <niccolo.fei@enterprisedb.com>
2026-07-13 14:26:17 +08:00
renovate[bot]
2cc4e9905d
fix(deps): update k8s.io/utils digest to be93311 (#986)
This PR contains the following updates:

| Package | Type | Update | Change |
|---|---|---|---|
| [k8s.io/utils](https://redirect.github.com/kubernetes/utils) | require
| digest | `a95e086` → `be93311` |

---

> [!WARNING]
> Some dependencies could not be looked up. Check the [Dependency
Dashboard](../issues/7) for more information.

---

### Configuration

📅 **Schedule**: (UTC)

- Branch creation
  - At any time (no schedule defined)
- Automerge
  - At any time (no schedule defined)

🚦 **Automerge**: Enabled.

♻ **Rebasing**: Never, or you tick the rebase/retry checkbox.

🔕 **Ignore**: Close this PR and you won't be reminded about this update
again.

---

- [ ] <!-- rebase-check -->If you want to rebase/retry this PR, check
this box

---

This PR was generated by [Mend Renovate](https://mend.io/renovate/).
View the [repository job
log](https://developer.mend.io/github/cloudnative-pg/plugin-barman-cloud).

<!--renovate-debug:eyJjcmVhdGVkSW5WZXIiOiI0My4yNDIuMiIsInVwZGF0ZWRJblZlciI6IjQzLjI0Mi4yIiwidGFyZ2V0QnJhbmNoIjoibWFpbiIsImxhYmVscyI6WyJhdXRvbWF0ZWQiLCJuby1pc3N1ZSJdfQ==-->

Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com>
Co-authored-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com>
Co-authored-by: Leonardo Cecchi <leonardo.cecchi@enterprisedb.com>
2026-07-09 16:45:23 +02:00
renovate[bot]
39881b0d9c
chore(deps): refresh pip-compile outputs (#989)
This PR contains the following updates:

| Update | Change |
|---|---|
| lockFileMaintenance | All locks refreshed |

🔧 This Pull Request updates lock files to use the latest dependency
versions.

---

### Configuration

📅 **Schedule**: (UTC)

- Branch creation
  - "before 4am on monday"
- Automerge
  - At any time (no schedule defined)

🚦 **Automerge**: Enabled.

♻ **Rebasing**: Never, or you tick the rebase/retry checkbox.

👻 **Immortal**: This PR will be recreated if closed unmerged. Get
[config
help](https://redirect.github.com/renovatebot/renovate/discussions) if
that's undesired.

---

- [ ] <!-- rebase-check -->If you want to rebase/retry this PR, check
this box

---

This PR was generated by [Mend Renovate](https://mend.io/renovate/).
View the [repository job
log](https://developer.mend.io/github/cloudnative-pg/plugin-barman-cloud).

<!--renovate-debug:eyJjcmVhdGVkSW5WZXIiOiI0My4yNDIuMiIsInVwZGF0ZWRJblZlciI6IjQzLjI0Mi4yIiwidGFyZ2V0QnJhbmNoIjoibWFpbiIsImxhYmVscyI6WyJhdXRvbWF0ZWQiLCJuby1pc3N1ZSJdfQ==-->

Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com>
Co-authored-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com>
Co-authored-by: Leonardo Cecchi <leonardo.cecchi@enterprisedb.com>
2026-07-09 15:45:20 +02:00
renovate[bot]
1a701f66eb
chore(deps): update dependency serialize-javascript to v7.0.7 (#990)
This PR contains the following updates:

| Package | Change |
[Age](https://docs.renovatebot.com/merge-confidence/) |
[Confidence](https://docs.renovatebot.com/merge-confidence/) |
|---|---|---|---|
|
[serialize-javascript](https://redirect.github.com/yahoo/serialize-javascript)
| [`7.0.6` →
`7.0.7`](https://renovatebot.com/diffs/npm/serialize-javascript/7.0.6/7.0.7)
|
![age](https://developer.mend.io/api/mc/badges/age/npm/serialize-javascript/7.0.7?slim=true)
|
![confidence](https://developer.mend.io/api/mc/badges/confidence/npm/serialize-javascript/7.0.6/7.0.7?slim=true)
|

---

### Release Notes

<details>
<summary>yahoo/serialize-javascript (serialize-javascript)</summary>

###
[`v7.0.7`](https://redirect.github.com/yahoo/serialize-javascript/compare/v7.0.6...01bec60c108143e69f6b105e443d62a13b61cbbe)

[Compare
Source](https://redirect.github.com/yahoo/serialize-javascript/compare/v7.0.6...v7.0.7)

</details>

---

### Configuration

📅 **Schedule**: (UTC)

- Branch creation
  - At any time (no schedule defined)
- Automerge
  - At any time (no schedule defined)

🚦 **Automerge**: Enabled.

♻ **Rebasing**: Never, or you tick the rebase/retry checkbox.

🔕 **Ignore**: Close this PR and you won't be reminded about this update
again.

---

- [ ] <!-- rebase-check -->If you want to rebase/retry this PR, check
this box

---

This PR was generated by [Mend Renovate](https://mend.io/renovate/).
View the [repository job
log](https://developer.mend.io/github/cloudnative-pg/plugin-barman-cloud).

<!--renovate-debug:eyJjcmVhdGVkSW5WZXIiOiI0My4yNDIuMiIsInVwZGF0ZWRJblZlciI6IjQzLjI0Mi4yIiwidGFyZ2V0QnJhbmNoIjoibWFpbiIsImxhYmVscyI6WyJhdXRvbWF0ZWQiLCJuby1pc3N1ZSJdfQ==-->

Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com>
Co-authored-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com>
Co-authored-by: Leonardo Cecchi <leonardo.cecchi@enterprisedb.com>
2026-07-09 14:35:13 +02:00
Armando Ruocco
21aadbfd11
test(e2e): cover parallel WAL restore via the plugin (#978)
Recreate the parallel WAL-restore coverage in the plugin repo: a
2-instance cluster archiving to minio with wal.maxParallel=3, forged WAL
segments on the object store, and assertions on the plugin's
prefetch/spool/end-of-wal-stream state machine driven through
`/controller/manager wal-restore` on the standby.

Part of cloudnative-pg/cloudnative-pg#10954.

Signed-off-by: Armando Ruocco <armando.ruocco@enterprisedb.com>
Signed-off-by: Niccolò Fei <niccolo.fei@enterprisedb.com>
Signed-off-by: Marco Nenciarini <marco.nenciarini@enterprisedb.com>
Co-authored-by: Niccolò Fei <niccolo.fei@enterprisedb.com>
Co-authored-by: Marco Nenciarini <marco.nenciarini@enterprisedb.com>
2026-07-09 11:26:08 +02:00
Marco Nenciarini
b0e4538147
chore(renovate): stop ignoring test/e2e/internal/tests/** (#996)
ignorePaths still had '**/tests/**' left over from the recommended
preset, even though a comment right above documents removing
'**/test/**' specifically to let renovate scan test/e2e for emulator
image dependencies. Every e2e Ginkgo package lives under
test/e2e/internal/tests/**, so the plural pattern was silently excluding
all of them; any image pinned by digest in those packages (e.g. minio,
aws-cli) was never going to get a bump PR.

No other tracked path in the repo matches '**/tests/**' once
node_modules is excluded by its own rule, so this narrows nothing else.

Signed-off-by: Marco Nenciarini <marco.nenciarini@enterprisedb.com>
2026-07-09 11:00:04 +02:00
Niccolò Fei
bf955430cb
fix: reduce startupProbe periodSeconds without losing failure tolerance (#992)
The startup probe for the injected plugin-barman-cloud sidecar
previously left `periodSeconds` unset, so the API server defaulted it to
10s. Because the sidecar is a native init container that gates the main
postgres container on reaching `Started`, this added roughly one full
period to every pod's startup, even though the probe itself (a local
unix-socket health check) normally succeeds in milliseconds.

`periodSeconds` is now 1s, so the probe reports success almost
immediately in the common case. To avoid trading away failure tolerance
for that faster common case, `failureThreshold` is raised to 30 and
`timeoutSeconds` is lowered to 5s: a unix-socket call essentially never
times out under mere load, so hitting the timeout means the sidecar is
genuinely unresponsive rather than just slow, and it's fine to give that
rare case more attempts before restarting the container.

Closes #991

Signed-off-by: Niccolò Fei <niccolo.fei@enterprisedb.com>
Signed-off-by: Marco Nenciarini <marco.nenciarini@enterprisedb.com>
Co-authored-by: Marco Nenciarini <marco.nenciarini@enterprisedb.com>
2026-07-09 10:36:09 +02:00
Marco Nenciarini
6a55a361a3 test: replace sleep-based test with deterministic channel verification
The cleanup routine test used time.Sleep() without actually verifying
the goroutine stopped. Added a done channel to provide deterministic
verification of goroutine termination.

Signed-off-by: Marco Nenciarini <marco.nenciarini@enterprisedb.com>
2025-12-23 17:06:00 +01:00
Armando Ruocco
62b579101f fix: prevent memory leak by periodically cleaning up expired cache entries
Signed-off-by: Armando Ruocco <armando.ruocco@enterprisedb.com>
2025-12-23 17:06:00 +01:00
14 changed files with 757 additions and 23 deletions

View File

@ -360,7 +360,9 @@ tasks:
SIDECAR_IMAGE_NAME: ghcr.io/{{.GITHUB_REPOSITORY}}-sidecar{{if not (hasPrefix "refs/tags/v" .GITHUB_REF)}}-testing{{end}}
# remove /merge suffix from the branch name. This is a workaround for the GitHub workflow on PRs,
# where the branch name is suffixed with /merge. Prepend pr- to the branch name on PRs.
IMAGE_VERSION: '{{regexReplaceAll "(\\d+)/merge" .GITHUB_REF_NAME "pr-${1}"}}'
# Any remaining "/" (e.g. from namespaced branches like "dev/foo") is replaced with "-" since
# "/" is not a valid character in a Docker tag.
IMAGE_VERSION: '{{replace "/" "-" (regexReplaceAll "(\\d+)/merge" .GITHUB_REF_NAME "pr-${1}")}}'
env:
# renovate: datasource=git-refs depName=docker lookupName=https://github.com/purpleclay/daggerverse currentValue=main
DAGGER_DOCKER_SHA: ee12c1a4a2630e194ec20c5a9959183e3a78c192
@ -450,7 +452,9 @@ tasks:
SIDECAR_IMAGE_NAME: ghcr.io/{{.GITHUB_REPOSITORY}}-sidecar{{if not (hasPrefix "refs/tags/v" .GITHUB_REF)}}-testing{{end}}
# remove /merge suffix from the branch name. This is a workaround for the GitHub workflow on PRs,
# where the branch name is suffixed with /merge. Prepend pr- to the branch name on PRs.
IMAGE_VERSION: '{{regexReplaceAll "(\\d+)/merge" .GITHUB_REF_NAME "pr-${1}"}}'
# Any remaining "/" (e.g. from namespaced branches like "dev/foo") is replaced with "-" since
# "/" is not a valid character in a Docker tag.
IMAGE_VERSION: '{{replace "/" "-" (regexReplaceAll "(\\d+)/merge" .GITHUB_REF_NAME "pr-${1}")}}'
env:
# renovate: datasource=git-refs depName=kustomize lookupName=https://github.com/sagikazarmark/daggerverse currentValue=main
DAGGER_KUSTOMIZE_SHA: ff27cd50f6b4eed2e3753c520632cd6099e1ce52

View File

@ -449,9 +449,9 @@ google-api-core==2.31.0 \
# via
# google-cloud-core
# google-cloud-storage
google-auth==2.55.0 \
--hash=sha256:a17cef9dedf98c4ebae2fb0c48c8f75952c877cbc2efe09f329ef16c2783d88a \
--hash=sha256:fcd3a130f575fa36403d38774af1c64a4fbfbca09215f0589d2372b5119697cb
google-auth==2.55.1 \
--hash=sha256:eada68dfd52b3b81191827601e2a0c3fa12540c818534b630ddc5355769c3995 \
--hash=sha256:fb2d9b730f2c9b8d326ec8d7222f21aef2ead15bf0513793d6442485d87af0a1
# via
# google-api-core
# google-cloud-core

2
go.mod
View File

@ -20,7 +20,7 @@ require (
k8s.io/apiextensions-apiserver v0.36.2
k8s.io/apimachinery v0.36.2
k8s.io/client-go v0.36.2
k8s.io/utils v0.0.0-20260617174310-a95e086a2553
k8s.io/utils v0.0.0-20260626114624-be93311217bd
sigs.k8s.io/controller-runtime v0.24.1
sigs.k8s.io/kustomize/api v0.21.1
sigs.k8s.io/kustomize/kyaml v0.21.1

4
go.sum
View File

@ -326,8 +326,8 @@ k8s.io/kube-openapi v0.0.0-20260502001324-b7f5293f4787 h1:kHv8PETbPIVHfqKBYwTNNS
k8s.io/kube-openapi v0.0.0-20260502001324-b7f5293f4787/go.mod h1:Cyq7UE0QtGe+Zo+/6XFrxiS4Mq0tLyQEONkFzSkfp9o=
k8s.io/streaming v0.36.2 h1:NSKthPPg9UFSKsRauVJUVGH2Dvn8fhKmY4qrMkw/p98=
k8s.io/streaming v0.36.2/go.mod h1:z6fV3D+NVkoeqRMtWwlUZK6U17SY/LqNzOxWL6GyR/s=
k8s.io/utils v0.0.0-20260617174310-a95e086a2553 h1:hmGqDecjc8d7HVzWzRFl0QD9bYuYKbBEG7t8xwnVxfI=
k8s.io/utils v0.0.0-20260617174310-a95e086a2553/go.mod h1:xDxuJ0whA3d0I4mf/C4ppKHxXynQ+fxnkmQH0vTHnuk=
k8s.io/utils v0.0.0-20260626114624-be93311217bd h1:Ea7fgQ5we8Y9T0OX5o0dAHzQOBRI07D/dEYRaB9ZZEs=
k8s.io/utils v0.0.0-20260626114624-be93311217bd/go.mod h1:xDxuJ0whA3d0I4mf/C4ppKHxXynQ+fxnkmQH0vTHnuk=
sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.34.0 h1:hSfpvjjTQXQY2Fol2CS0QHMNs/WI1MOSGzCm1KhM5ec=
sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.34.0/go.mod h1:Ve9uj1L+deCXFrPOk1LpFXqTg7LCFzFso6PA48q/XZw=
sigs.k8s.io/controller-runtime v0.24.1 h1:miPEwrmirImAvgME1L9qebGHrOnGJoVmVdtOU9fRfo4=

View File

@ -36,6 +36,9 @@ import (
// DefaultTTLSeconds is the default TTL in seconds of cache entries
const DefaultTTLSeconds = 10
// DefaultCleanupIntervalSeconds is the default interval in seconds for cache cleanup
const DefaultCleanupIntervalSeconds = 30
type cachedEntry struct {
entry client.Object
fetchUnixTime int64
@ -49,18 +52,30 @@ func (e *cachedEntry) isExpired() bool {
// ExtendedClient is an extended client that is capable of caching multiple secrets without relying on informers
type ExtendedClient struct {
client.Client
cachedObjects []cachedEntry
mux *sync.Mutex
cachedObjects []cachedEntry
mux *sync.Mutex
cleanupInterval time.Duration
cleanupDone chan struct{} // Signals when cleanup routine exits
}
// NewExtendedClient returns an extended client capable of caching secrets on the 'Get' operation
// NewExtendedClient returns an extended client capable of caching secrets on the 'Get' operation.
// It starts a background goroutine that periodically cleans up expired cache entries.
// The cleanup routine will stop when the provided context is cancelled.
func NewExtendedClient(
ctx context.Context,
baseClient client.Client,
) client.Client {
return &ExtendedClient{
Client: baseClient,
mux: &sync.Mutex{},
ec := &ExtendedClient{
Client: baseClient,
mux: &sync.Mutex{},
cleanupInterval: DefaultCleanupIntervalSeconds * time.Second,
cleanupDone: make(chan struct{}),
}
// Start the background cleanup routine
go ec.startCleanupRoutine(ctx)
return ec
}
func (e *ExtendedClient) isObjectCached(obj client.Object) bool {
@ -208,3 +223,55 @@ func (e *ExtendedClient) Patch(
return e.Client.Patch(ctx, obj, patch, opts...)
}
// startCleanupRoutine periodically removes expired entries from the cache.
// It runs until the context is cancelled.
func (e *ExtendedClient) startCleanupRoutine(ctx context.Context) {
defer close(e.cleanupDone)
contextLogger := log.FromContext(ctx).WithName("extended_client_cleanup")
ticker := time.NewTicker(e.cleanupInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
contextLogger.Debug("stopping cache cleanup routine")
return
case <-ticker.C:
// Check context before cleanup to avoid unnecessary work during shutdown
if ctx.Err() != nil {
return
}
e.cleanupExpiredEntries(ctx)
}
}
}
// cleanupExpiredEntries removes all expired entries from the cache.
func (e *ExtendedClient) cleanupExpiredEntries(ctx context.Context) {
contextLogger := log.FromContext(ctx).WithName("extended_client_cleanup")
e.mux.Lock()
defer e.mux.Unlock()
initialCount := len(e.cachedObjects)
if initialCount == 0 {
return
}
// Create a new slice with only non-expired entries
validEntries := make([]cachedEntry, 0, initialCount)
for _, entry := range e.cachedObjects {
if !entry.isExpired() {
validEntries = append(validEntries, entry)
}
}
removedCount := initialCount - len(validEntries)
if removedCount > 0 {
e.cachedObjects = validEntries
contextLogger.Debug("cleaned up expired cache entries",
"removedCount", removedCount,
"remainingCount", len(validEntries))
}
}

View File

@ -20,6 +20,7 @@ SPDX-License-Identifier: Apache-2.0
package client
import (
"context"
"time"
corev1 "k8s.io/api/core/v1"
@ -59,6 +60,7 @@ var _ = Describe("ExtendedClient Get", func() {
extendedClient *ExtendedClient
secretInClient *corev1.Secret
objectStore *barmancloudv1.ObjectStore
cancelCtx context.CancelFunc
)
BeforeEach(func() {
@ -79,7 +81,14 @@ var _ = Describe("ExtendedClient Get", func() {
baseClient := fake.NewClientBuilder().
WithScheme(scheme).
WithObjects(secretInClient, objectStore).Build()
extendedClient = NewExtendedClient(baseClient).(*ExtendedClient)
ctx, cancel := context.WithCancel(context.Background())
cancelCtx = cancel
extendedClient = NewExtendedClient(ctx, baseClient).(*ExtendedClient)
})
AfterEach(func() {
// Cancel the context to stop the cleanup routine
cancelCtx()
})
It("returns secret from cache if not expired", func(ctx SpecContext) {
@ -164,3 +173,141 @@ var _ = Describe("ExtendedClient Get", func() {
Expect(objectStore.GetResourceVersion()).To(Equal("from cache"))
})
})
var _ = Describe("ExtendedClient Cache Cleanup", func() {
var (
extendedClient *ExtendedClient
cancelCtx context.CancelFunc
)
BeforeEach(func() {
baseClient := fake.NewClientBuilder().
WithScheme(scheme).
Build()
ctx, cancel := context.WithCancel(context.Background())
cancelCtx = cancel
extendedClient = NewExtendedClient(ctx, baseClient).(*ExtendedClient)
})
AfterEach(func() {
cancelCtx()
})
It("cleans up expired entries", func(ctx SpecContext) {
// Add some expired entries
expiredSecret1 := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "expired-secret-1",
},
}
expiredSecret2 := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "expired-secret-2",
},
}
validSecret := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "valid-secret",
},
}
// Add expired entries (2 minutes ago)
addToCache(extendedClient, expiredSecret1, time.Now().Add(-2*time.Minute).Unix())
addToCache(extendedClient, expiredSecret2, time.Now().Add(-2*time.Minute).Unix())
// Add valid entry (just now)
addToCache(extendedClient, validSecret, time.Now().Unix())
Expect(extendedClient.cachedObjects).To(HaveLen(3))
// Trigger cleanup
extendedClient.cleanupExpiredEntries(ctx)
// Only the valid entry should remain
Expect(extendedClient.cachedObjects).To(HaveLen(1))
Expect(extendedClient.cachedObjects[0].entry.GetName()).To(Equal("valid-secret"))
})
It("does nothing when all entries are valid", func(ctx SpecContext) {
validSecret1 := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "valid-secret-1",
},
}
validSecret2 := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "valid-secret-2",
},
}
addToCache(extendedClient, validSecret1, time.Now().Unix())
addToCache(extendedClient, validSecret2, time.Now().Unix())
Expect(extendedClient.cachedObjects).To(HaveLen(2))
// Trigger cleanup
extendedClient.cleanupExpiredEntries(ctx)
// Both entries should remain
Expect(extendedClient.cachedObjects).To(HaveLen(2))
})
It("does nothing when cache is empty", func(ctx SpecContext) {
Expect(extendedClient.cachedObjects).To(BeEmpty())
// Trigger cleanup
extendedClient.cleanupExpiredEntries(ctx)
Expect(extendedClient.cachedObjects).To(BeEmpty())
})
It("removes all entries when all are expired", func(ctx SpecContext) {
expiredSecret1 := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "expired-secret-1",
},
}
expiredSecret2 := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "expired-secret-2",
},
}
addToCache(extendedClient, expiredSecret1, time.Now().Add(-2*time.Minute).Unix())
addToCache(extendedClient, expiredSecret2, time.Now().Add(-2*time.Minute).Unix())
Expect(extendedClient.cachedObjects).To(HaveLen(2))
// Trigger cleanup
extendedClient.cleanupExpiredEntries(ctx)
Expect(extendedClient.cachedObjects).To(BeEmpty())
})
It("stops cleanup routine when context is cancelled", func() {
// Create a new client with a short cleanup interval for testing
baseClient := fake.NewClientBuilder().
WithScheme(scheme).
Build()
ctx, cancel := context.WithCancel(context.Background())
ec := NewExtendedClient(ctx, baseClient).(*ExtendedClient)
ec.cleanupInterval = 10 * time.Millisecond
// Cancel the context immediately
cancel()
// Verify the cleanup routine actually stops by waiting for the done channel
select {
case <-ec.cleanupDone:
// Success: cleanup routine exited as expected
case <-time.After(1 * time.Second):
Fail("cleanup routine did not stop within timeout")
}
})
})

View File

@ -80,7 +80,7 @@ func Start(ctx context.Context) error {
return err
}
customCacheClient := extendedclient.NewExtendedClient(mgr.GetClient())
customCacheClient := extendedclient.NewExtendedClient(ctx, mgr.GetClient())
if err := mgr.Add(&CNPGI{
Client: customCacheClient,

View File

@ -403,8 +403,9 @@ func reconcilePodSpec(
envs = append(envs, config.env...)
baseProbe := &corev1.Probe{
FailureThreshold: 10,
TimeoutSeconds: 10,
PeriodSeconds: 1,
FailureThreshold: 30,
TimeoutSeconds: 5,
ProbeHandler: corev1.ProbeHandler{
Exec: &corev1.ExecAction{
Command: []string{"/manager", "healthcheck", "unix"},

View File

@ -12,14 +12,14 @@
rebaseWhen: 'never',
prConcurrentLimit: 5,
// Override default ignorePaths to scan test/e2e for emulator image dependencies
// Removed: '**/test/**'
// Removed: '**/test/**', '**/tests/**' (this repo's e2e Ginkgo packages live
// under test/e2e/internal/tests/**, which the plural pattern was excluding)
ignorePaths: [
'**/node_modules/**',
'**/bower_components/**',
'**/vendor/**',
'**/examples/**',
'**/__tests__/**',
'**/tests/**',
'**/__fixtures__/**',
],
lockFileMaintenance: {

View File

@ -39,6 +39,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/replicacluster"
_ "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/tests/walrestore"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"

View File

@ -0,0 +1,24 @@
/*
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 walrestore contains the end-to-end test for the parallel WAL restore
// behaviour of the Barman Cloud Plugin: prefetching upcoming segments into the
// spool directory, serving later requests from the spool, and tracking the
// end-of-wal-stream sentinel.
package walrestore

View File

@ -0,0 +1,178 @@
/*
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 walrestore
import (
cloudnativepgv1 "github.com/cloudnative-pg/api/pkg/api/v1"
barmanapi "github.com/cloudnative-pg/barman-cloud/pkg/api"
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/utils/ptr"
pluginBarmanCloudV1 "github.com/cloudnative-pg/plugin-barman-cloud/api/v1"
"github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/objectstore"
)
const (
minioName = "minio"
objectStoreName = "source"
clusterName = "source"
s3ClientName = "s3-client"
storageSize = "1Gi"
// walMaxParallel is the prefetch parallelism under test: for a regular WAL
// request the plugin restores the requested segment and prefetches the next
// ones, up to this many segments in total.
walMaxParallel = 3
)
// newObjectStoreResources returns the minio server Deployment/Service/Secret/PVC.
func newObjectStoreResources(namespace string) *objectstore.Resources {
return objectstore.NewMinioObjectStoreResources(namespace, minioName)
}
// newObjectStore returns a minio-backed ObjectStore configured with the WAL
// prefetch parallelism (maxParallel) under test. Archiving with gzip makes the
// archived segments carry the ".gz" suffix that forged segments are copied from.
func newObjectStore(namespace string) *pluginBarmanCloudV1.ObjectStore {
store := objectstore.NewMinioObjectStore(namespace, objectStoreName, minioName)
store.Spec.Configuration.Wal = &barmanapi.WalBackupConfiguration{
MaxParallel: walMaxParallel,
Compression: barmanapi.CompressionTypeGzip,
}
return store
}
// newCluster returns a 2-instance cluster that uses the plugin as its WAL
// archiver, so the standby drives WAL restore (and its prefetch/spool/
// end-of-wal-stream state machine) through the plugin.
func newCluster(namespace string) *cloudnativepgv1.Cluster {
return &cloudnativepgv1.Cluster{
TypeMeta: metav1.TypeMeta{
Kind: "Cluster",
APIVersion: "postgresql.cnpg.io/v1",
},
ObjectMeta: metav1.ObjectMeta{
Name: clusterName,
Namespace: namespace,
},
Spec: cloudnativepgv1.ClusterSpec{
Instances: 2,
ImagePullPolicy: corev1.PullAlways,
Plugins: []cloudnativepgv1.PluginConfiguration{
{
Name: "barman-cloud.cloudnative-pg.io",
Parameters: map[string]string{
"barmanObjectName": objectStoreName,
},
IsWALArchiver: ptr.To(true),
},
},
PostgresConfiguration: cloudnativepgv1.PostgresConfiguration{
Parameters: map[string]string{
"log_min_messages": "DEBUG4",
},
},
StorageConfiguration: cloudnativepgv1.StorageConfiguration{
Size: storageSize,
},
},
}
}
// newS3ClientDeployment returns a deployment running the AWS CLI configured to
// talk to the in-namespace minio service. The test execs `aws s3` commands in
// it to forge WAL segments on the object store and to check their presence.
func newS3ClientDeployment(namespace string) *appsv1.Deployment {
labels := map[string]string{"app": s3ClientName}
return &appsv1.Deployment{
TypeMeta: metav1.TypeMeta{
Kind: "Deployment",
APIVersion: "apps/v1",
},
ObjectMeta: metav1.ObjectMeta{
Name: s3ClientName,
Namespace: namespace,
},
Spec: appsv1.DeploymentSpec{
Replicas: ptr.To(int32(1)),
Selector: &metav1.LabelSelector{MatchLabels: labels},
Template: corev1.PodTemplateSpec{
ObjectMeta: metav1.ObjectMeta{Labels: labels},
Spec: corev1.PodSpec{
Containers: []corev1.Container{
{
Name: s3ClientName,
// renovate: datasource=docker depName=amazon/aws-cli versioning=docker
// Version: 2.35.11
Image: "docker.io/amazon/aws-cli@sha256:749bfaf91d690b9a1768083822d620f96c19defdf9ca2dc227eb3695281fda5b",
Command: []string{"sleep", "infinity"},
Env: []corev1.EnvVar{
{
Name: "AWS_ENDPOINT_URL",
Value: "http://" + minioName + ":9000",
},
{
Name: "AWS_ACCESS_KEY_ID",
ValueFrom: &corev1.EnvVarSource{
SecretKeyRef: &corev1.SecretKeySelector{
LocalObjectReference: corev1.LocalObjectReference{Name: minioName},
Key: "ACCESS_KEY_ID",
},
},
},
{
Name: "AWS_SECRET_ACCESS_KEY",
ValueFrom: &corev1.EnvVarSource{
SecretKeyRef: &corev1.SecretKeySelector{
LocalObjectReference: corev1.LocalObjectReference{Name: minioName},
Key: "ACCESS_SECRET_KEY",
},
},
},
{
Name: "AWS_DEFAULT_REGION",
Value: "us-east-1",
},
// The CRC-based default checksums introduced in AWS
// CLI 2.23 are not supported by every S3-compatible
// object store, minio included.
{
Name: "AWS_REQUEST_CHECKSUM_CALCULATION",
Value: "when_required",
},
{
Name: "AWS_RESPONSE_CHECKSUM_VALIDATION",
Value: "when_required",
},
},
SecurityContext: &corev1.SecurityContext{
AllowPrivilegeEscalation: ptr.To(false),
SeccompProfile: &corev1.SeccompProfile{
Type: corev1.SeccompProfileTypeRuntimeDefault,
},
},
},
},
},
},
},
}
}

View File

@ -0,0 +1,312 @@
/*
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 walrestore
import (
"context"
"errors"
"fmt"
"strconv"
"strings"
"time"
corev1 "k8s.io/api/core/v1"
apitypes "k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
executil "k8s.io/client-go/util/exec"
"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"
"github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/command"
"github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/deployment"
nmsp "github.com/cloudnative-pg/plugin-barman-cloud/test/e2e/internal/namespace"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
)
const (
// spoolDirectory is where the plugin sidecar prefetches WAL segments and
// records the end-of-wal-stream sentinel. It must match the SPOOL_DIRECTORY
// the operator injects on the sidecar (internal/cnpgi/operator/lifecycle.go),
// which lives on the CloudNativePG scratch-data volume shared with the
// postgres container, so the path is visible from the postgres container too.
spoolDirectory = "/controller/wal-restore-spool"
// pgWalPath is PgDataPath + "/pg_wal" in the CloudNativePG operand image.
pgWalPath = "/var/lib/postgresql/data/pgdata/pg_wal"
// endOfWALStreamFlag is the sentinel file the barman-cloud restorer writes in
// the spool to record that the archive ran out of segments.
endOfWALStreamFlag = "end-of-wal-stream"
// managerExecutable is the CloudNativePG instance manager. Its wal-restore
// subcommand delegates to the plugin when the cluster uses it as WAL archiver.
managerExecutable = "/controller/manager"
// postgresContainer is the container we exec into on instance pods.
postgresContainer = "postgres"
// walLogDir is the WALs subdirectory (timeline + log id) the forged segments
// live under; a freshly bootstrapped, idle cluster stays within it.
walLogDir = "0000000100000000"
// bucket is the destination bucket of the minio ObjectStore.
bucket = "backups"
)
// walFile returns the name of the n-th forged WAL segment (segment 0xF0+n on
// timeline 1, log 0, for n <= 15). The high segment number keeps it out of the
// range an idle PostgreSQL would archive on its own.
func walFile(n int) string {
return fmt.Sprintf("0000000100000000%08X", 0xF0+n)
}
// walObjectURI returns the s3 URI of a wals object (by file name) in the store.
func walObjectURI(name string) string {
return fmt.Sprintf("s3://%s/%s/wals/%s/%s", bucket, clusterName, walLogDir, name)
}
// execInPod runs a command in the given container and returns stdout, stderr
// and the error (non-nil for a non-zero exit code).
func execInPod(
ctx context.Context,
clientSet *kubernetes.Clientset,
cfg *rest.Config,
namespace, pod, container string,
args ...string,
) (string, string, error) {
return command.ExecuteInContainer(
ctx,
*clientSet,
cfg,
command.ContainerLocator{
NamespaceName: namespace,
PodName: pod,
ContainerName: container,
},
nil,
args,
)
}
// This test drives the plugin's parallel WAL restore directly: it invokes the
// instance-manager wal-restore command on the standby (which delegates to the
// plugin) and asserts the prefetch/spool/end-of-wal-stream state machine with
// maxParallel = 3. To control the archive deterministically, it forges WAL
// segments on the object store by copying a real archived segment under new
// names.
var _ = Describe("Parallel WAL restore", func() {
var (
namespace *corev1.Namespace
cl client.Client
clientSet *kubernetes.Clientset
cfg *rest.Config
)
BeforeEach(func(ctx SpecContext) {
var err error
cl, _, err = internalClient.NewClient()
Expect(err).NotTo(HaveOccurred())
clientSet, cfg, err = internalClient.NewClientSet()
Expect(err).NotTo(HaveOccurred())
namespace, err = nmsp.CreateUniqueNamespace(ctx, cl, "wal-restore-parallel")
Expect(err).NotTo(HaveOccurred())
})
AfterEach(func(ctx SpecContext) {
Expect(cl.Delete(ctx, namespace)).To(Succeed())
})
It("prefetches segments, serves them from the spool and tracks end-of-wal-stream",
func(ctx SpecContext) {
ns := namespace.Name
By("creating the object store backing resources")
Expect(newObjectStoreResources(ns).Create(ctx, cl)).To(Succeed())
By("creating the ObjectStore with WAL maxParallel")
Expect(cl.Create(ctx, newObjectStore(ns))).To(Succeed())
By("deploying the S3 client used to forge and inspect WAL segments")
Expect(cl.Create(ctx, newS3ClientDeployment(ns))).To(Succeed())
By("creating the cluster using the plugin as WAL archiver")
cluster := newCluster(ns)
Expect(cl.Create(ctx, cluster)).To(Succeed())
By("waiting for the cluster to become ready")
Eventually(func(g Gomega) {
g.Expect(cl.Get(ctx,
apitypes.NamespacedName{Name: clusterName, Namespace: ns},
cluster)).To(Succeed())
g.Expect(internalCluster.IsReady(*cluster)).To(BeTrue())
}).WithTimeout(10 * time.Minute).WithPolling(10 * time.Second).Should(Succeed())
By("waiting for the S3 client to become ready")
Eventually(func(g Gomega) {
ready, err := deployment.IsReady(ctx, cl,
apitypes.NamespacedName{Name: s3ClientName, Namespace: ns})
g.Expect(err).NotTo(HaveOccurred())
g.Expect(ready).To(BeTrue())
}).WithTimeout(2 * time.Minute).WithPolling(5 * time.Second).Should(Succeed())
primary := cluster.Status.CurrentPrimary
Expect(primary).NotTo(BeEmpty())
standby := clusterName + "-2"
var s3ClientPods corev1.PodList
Expect(cl.List(ctx, &s3ClientPods,
client.InNamespace(ns),
client.MatchingLabels{"app": s3ClientName})).To(Succeed())
Expect(s3ClientPods.Items).NotTo(BeEmpty())
s3Client := s3ClientPods.Items[0].Name
// Operations scoped to the fixed pods/clients, kept as closures so the
// step assertions below read like the original state-machine table.
restore := func(name string) error {
_, _, err := execInPod(ctx, clientSet, cfg, ns, standby, postgresContainer,
managerExecutable, "wal-restore", name, pgWalPath+"/"+name)
return err
}
existsIn := func(dir, name string) bool {
_, _, err := execInPod(ctx, clientSet, cfg, ns, standby, postgresContainer,
"test", "-f", dir+"/"+name)
return err == nil
}
flagSet := func() bool { return existsIn(spoolDirectory, endOfWALStreamFlag) }
spoolSegments := func() int {
out, _, _ := execInPod(ctx, clientSet, cfg, ns, standby, postgresContainer,
"sh", "-c",
"ls -1 "+spoolDirectory+" 2>/dev/null | grep -Ec '^[0-9A-F]{24}$' || true")
n, _ := strconv.Atoi(strings.TrimSpace(out))
return n
}
purgeSpool := func() {
_, _, _ = execInPod(ctx, clientSet, cfg, ns, standby, postgresContainer,
"sh", "-c", "rm -f "+spoolDirectory+"/* 2>/dev/null; true")
}
forge := func(src, dst string) {
// ExecuteInContainer drops stdout/stderr on a non-zero exit, so on
// failure this only reports the exit code, not the aws CLI's error text.
_, _, err := execInPod(ctx, clientSet, cfg, ns, s3Client, s3ClientName,
"aws", "s3", "cp", walObjectURI(src), walObjectURI(dst))
Expect(err).NotTo(HaveOccurred(), "forging %s -> %s", src, dst)
}
objectExists := func(name string) bool {
out, _, err := execInPod(ctx, clientSet, cfg, ns, s3Client, s3ClientName,
"aws", "s3", "ls", walObjectURI(name))
return err == nil && strings.TrimSpace(out) != ""
}
By("archiving a real WAL on the primary and learning its name")
_, _, err := execInPod(ctx, clientSet, cfg, ns, primary, postgresContainer,
"psql", "-tAc", "CHECKPOINT")
Expect(err).NotTo(HaveOccurred(), "CHECKPOINT on the primary failed")
out, _, err := execInPod(ctx, clientSet, cfg, ns, primary, postgresContainer,
"psql", "-tAc", "SELECT pg_walfile_name(pg_switch_wal())")
Expect(err).NotTo(HaveOccurred(), "switching WAL on the primary failed")
latestWAL := strings.TrimSpace(out)
Expect(latestWAL).To(HavePrefix(walLogDir),
"the freshly bootstrapped cluster should still be on the first WAL log")
By("waiting for the archived WAL to land on the object store")
Eventually(func() bool {
return objectExists(latestWAL + ".gz")
}).WithTimeout(2 * time.Minute).WithPolling(5 * time.Second).Should(BeTrue())
By("forging WAL segments #1 to #5 from the archived WAL")
for n := 1; n <= 5; n++ {
forge(latestWAL+".gz", walFile(n))
}
By("ensuring the spool directory is empty on the standby")
purgeSpool()
// #1: served fresh; #2 and #3 prefetched into the spool; flag unset.
By("requesting WAL #1: #1 restored, #2 and #3 prefetched")
Expect(restore(walFile(1))).To(Succeed())
Eventually(func(g Gomega) {
g.Expect(existsIn(pgWalPath, walFile(1))).To(BeTrue(), "#1 in pg_wal")
g.Expect(existsIn(spoolDirectory, walFile(2))).To(BeTrue(), "#2 in spool")
g.Expect(existsIn(spoolDirectory, walFile(3))).To(BeTrue(), "#3 in spool")
g.Expect(flagSet()).To(BeFalse(), "end-of-wal-stream unset")
}).WithTimeout(time.Minute).WithPolling(2 * time.Second).Should(Succeed())
// #2: served from the spool; #3 stays prefetched; no new prefetch; flag unset.
By("requesting WAL #2: served from the spool, #3 still prefetched")
Expect(restore(walFile(2))).To(Succeed())
Eventually(func(g Gomega) {
g.Expect(existsIn(pgWalPath, walFile(2))).To(BeTrue(), "#2 in pg_wal")
g.Expect(existsIn(spoolDirectory, walFile(3))).To(BeTrue(), "#3 in spool")
g.Expect(flagSet()).To(BeFalse(), "end-of-wal-stream unset")
}).WithTimeout(time.Minute).WithPolling(2 * time.Second).Should(Succeed())
// #3: served from the spool; spool now empty; flag unset.
By("requesting WAL #3: served from the spool, spool now empty")
Expect(restore(walFile(3))).To(Succeed())
Eventually(func(g Gomega) {
g.Expect(existsIn(pgWalPath, walFile(3))).To(BeTrue(), "#3 in pg_wal")
g.Expect(spoolSegments()).To(Equal(0), "no WAL segments in spool")
g.Expect(flagSet()).To(BeFalse(), "end-of-wal-stream unset")
}).WithTimeout(time.Minute).WithPolling(2 * time.Second).Should(Succeed())
// #4: served fresh; #5 prefetched; #6 absent so end-of-wal-stream is set.
By("requesting WAL #4: #4 restored, #5 prefetched, end-of-wal-stream set")
Expect(restore(walFile(4))).To(Succeed())
Eventually(func(g Gomega) {
g.Expect(existsIn(pgWalPath, walFile(4))).To(BeTrue(), "#4 in pg_wal")
g.Expect(existsIn(spoolDirectory, walFile(5))).To(BeTrue(), "#5 in spool")
g.Expect(flagSet()).To(BeTrue(), "end-of-wal-stream set (#6 absent)")
}).WithTimeout(time.Minute).WithPolling(2 * time.Second).Should(Succeed())
By("forging WAL segment #6 on the object store")
forge(latestWAL+".gz", walFile(6))
// #5: served from the spool; flag untouched (served before it is checked).
By("requesting WAL #5: served from the spool, end-of-wal-stream still set")
Expect(restore(walFile(5))).To(Succeed())
Eventually(func(g Gomega) {
g.Expect(existsIn(pgWalPath, walFile(5))).To(BeTrue(), "#5 in pg_wal")
g.Expect(spoolSegments()).To(Equal(0), "no WAL segments in spool")
g.Expect(flagSet()).To(BeTrue(), "end-of-wal-stream still set")
}).WithTimeout(time.Minute).WithPolling(2 * time.Second).Should(Succeed())
// #6 (first): flag is set, so the request fails fast (exit 1) and the
// flag is consumed, leaving an empty spool.
By("requesting WAL #6: fails fast on the end-of-wal-stream flag, spool cleared")
restoreErr := restore(walFile(6))
var exitErr executil.CodeExitError
Expect(errors.As(restoreErr, &exitErr)).To(BeTrue(),
"expected a CodeExitError, got %T: %v", restoreErr, restoreErr)
Expect(exitErr.ExitStatus()).To(Equal(1), "exit code should be 1")
Eventually(func(g Gomega) {
g.Expect(existsIn(pgWalPath, walFile(6))).To(BeFalse(), "#6 not restored")
g.Expect(spoolSegments()).To(Equal(0), "no WAL segments in spool")
g.Expect(flagSet()).To(BeFalse(), "end-of-wal-stream consumed")
}).WithTimeout(time.Minute).WithPolling(2 * time.Second).Should(Succeed())
// #6 (second): now present, restored; #7 and #8 absent so the flag is
// set again.
By("requesting WAL #6 again: #6 restored, end-of-wal-stream set again")
Expect(restore(walFile(6))).To(Succeed())
Eventually(func(g Gomega) {
g.Expect(existsIn(pgWalPath, walFile(6))).To(BeTrue(), "#6 in pg_wal")
g.Expect(spoolSegments()).To(Equal(0), "no WAL segments in spool")
g.Expect(flagSet()).To(BeTrue(), "end-of-wal-stream set (#7/#8 absent)")
}).WithTimeout(time.Minute).WithPolling(2 * time.Second).Should(Succeed())
})
})

View File

@ -8396,9 +8396,9 @@ send@~0.19.0, send@~0.19.1:
statuses "~2.0.2"
serialize-javascript@>=7.0.5, serialize-javascript@^6.0.0, serialize-javascript@^6.0.1:
version "7.0.6"
resolved "https://registry.yarnpkg.com/serialize-javascript/-/serialize-javascript-7.0.6.tgz#f2f20c8af0757e4d8fa329d0210636da0682ddef"
integrity sha512-ATTK5Q4gFVg0YDp1my2vqygyvhcklD/UV5GIlYHooGTn/NogJqIzpetkD6E5kmuVULqz/S9inUL25XcAgDRJQg==
version "7.0.7"
resolved "https://registry.yarnpkg.com/serialize-javascript/-/serialize-javascript-7.0.7.tgz#06ec40576d4cea96d68010a534520bff1f948a72"
integrity sha512-YAy8Od6KV+uuwUuU50np8fGB/Aues6Y0nAhA9y/hId74PlKUcme4pXcBD46NWKr1Q4osN/iseZ17YqO1XfmI8g==
serve-handler@^6.1.7:
version "6.1.7"