fix(walrestore): serve pg_rewind without prefetching and flag machinery

The operator now tells WAL restore plugins when a restore request is
made on behalf of pg_rewind (cloudnative-pg/cnpg-i#351). pg_rewind walks
the timeline backwards, fetches every WAL file it needs exactly once,
and treats any restore failure as fatal, so both optimizations meant for
an instance in recovery must stay off: prefetching the following
segments is wasted work that ends in an archive miss, and the
end-of-wal-stream flag recorded by such a miss makes a later invocation
fail on a segment that is available in the archive, aborting the whole
rewind.

The cnpg-i dependency points to a pseudo-version of that pull request
and will be moved to the next tagged release once it is available.

Ref: cloudnative-pg/cloudnative-pg#11200
Signed-off-by: Armando Ruocco <armando.ruocco@enterprisedb.com>
This commit is contained in:
Armando Ruocco 2026-07-15 15:20:26 +02:00 committed by Leonardo Cecchi
parent fe82a311ec
commit 176f19cc00
2 changed files with 90 additions and 10 deletions

View File

@ -27,6 +27,7 @@ import (
"path" "path"
"time" "time"
barmanapi "github.com/cloudnative-pg/barman-cloud/pkg/api"
"github.com/cloudnative-pg/barman-cloud/pkg/archiver" "github.com/cloudnative-pg/barman-cloud/pkg/archiver"
barmanCommand "github.com/cloudnative-pg/barman-cloud/pkg/command" barmanCommand "github.com/cloudnative-pg/barman-cloud/pkg/command"
barmanCredentials "github.com/cloudnative-pg/barman-cloud/pkg/credentials" barmanCredentials "github.com/cloudnative-pg/barman-cloud/pkg/credentials"
@ -249,9 +250,11 @@ func (w WALServiceImplementation) Restore(
"Restoring WAL file", "Restoring WAL file",
"objectStore", objectStore.Name, "objectStore", objectStore.Name,
"serverName", serverName, "serverName", serverName,
"walName", walName) "walName", walName,
"mode", request.GetMode())
return &wal.WALRestoreResult{}, w.restoreFromBarmanObjectStore( return &wal.WALRestoreResult{}, w.restoreFromBarmanObjectStore(
ctx, configuration.Cluster, &objectStore, serverName, walName, destinationPath) ctx, configuration.Cluster, &objectStore, serverName, walName, destinationPath,
request.GetMode() == wal.WALRestoreRequest_MODE_REWIND)
} }
// resolveRestoreObjectStore selects the object store and server name to use when // resolveRestoreObjectStore selects the object store and server name to use when
@ -288,6 +291,7 @@ func (w WALServiceImplementation) restoreFromBarmanObjectStore(
serverName string, serverName string,
walName string, walName string,
destinationPath string, destinationPath string,
rewindMode bool,
) error { ) error {
contextLogger := log.FromContext(ctx) contextLogger := log.FromContext(ctx)
startTime := time.Now() startTime := time.Now()
@ -331,8 +335,10 @@ func (w WALServiceImplementation) restoreFromBarmanObjectStore(
return nil return nil
} }
// We skip this step if streaming connection is not available // Step 2: return error if the end-of-wal-stream flag is set.
if isStreamingAvailable(cluster, w.InstanceName) { // We skip this step if the flag machinery does not apply to this invocation
useEndOfWALStreamFlag := shouldUseEndOfWALStreamFlag(cluster, w.InstanceName, rewindMode)
if useEndOfWALStreamFlag {
if err := checkEndOfWALStreamFlag(walRestorer); err != nil { if err := checkEndOfWALStreamFlag(walRestorer); err != nil {
return err return err
} }
@ -340,10 +346,7 @@ func (w WALServiceImplementation) restoreFromBarmanObjectStore(
// Step 3: gather the WAL files names to restore. If the required file isn't a regular WAL, we download it directly. // Step 3: gather the WAL files names to restore. If the required file isn't a regular WAL, we download it directly.
var walFilesList []string var walFilesList []string
maxParallel := 1 maxParallel := maxWALFilesPerInvocation(barmanConfiguration, rewindMode)
if barmanConfiguration.Wal != nil && barmanConfiguration.Wal.MaxParallel > 1 {
maxParallel = barmanConfiguration.Wal.MaxParallel
}
if IsWALFile(walName) { if IsWALFile(walName) {
// If this is a regular WAL file, we try to prefetch // If this is a regular WAL file, we try to prefetch
if walFilesList, err = gatherWALFilesToRestore(walName, maxParallel); err != nil { if walFilesList, err = gatherWALFilesToRestore(walName, maxParallel); err != nil {
@ -365,9 +368,9 @@ func (w WALServiceImplementation) restoreFromBarmanObjectStore(
return classifyWALRestoreError(walStatus[0].WalName, walStatus[0].Err) return classifyWALRestoreError(walStatus[0].WalName, walStatus[0].Err)
} }
// We skip this step if streaming connection is not available // We skip this step if the flag machinery does not apply to this invocation
endOfWALStream := isEndOfWALStream(walStatus) endOfWALStream := isEndOfWALStream(walStatus)
if isStreamingAvailable(cluster, w.InstanceName) && endOfWALStream { if useEndOfWALStreamFlag && endOfWALStream {
contextLogger.Info( contextLogger.Info(
"Set end-of-wal-stream flag as one of the WAL files to be prefetched was not found") "Set end-of-wal-stream flag as one of the WAL files to be prefetched was not found")
@ -415,6 +418,38 @@ func (w WALServiceImplementation) SetFirstRequired(
panic("implement me") panic("implement me")
} }
// maxWALFilesPerInvocation returns how many WAL files a single restore
// invocation is allowed to fetch, the requested one included. Prefetching is
// disabled when restoring on behalf of pg_rewind, which walks the WAL
// backward: the following segments would never be requested, and past the end
// of the timeline they do not even exist
func maxWALFilesPerInvocation(barmanConfiguration *barmanapi.BarmanObjectStoreConfiguration, rewindMode bool) int {
if rewindMode {
return 1
}
if barmanConfiguration.Wal != nil && barmanConfiguration.Wal.MaxParallel > 1 {
return barmanConfiguration.Wal.MaxParallel
}
return 1
}
// shouldUseEndOfWALStreamFlag returns true when the end-of-wal-stream flag
// machinery applies to the current invocation. The flag makes the following
// invocation fail, so that PostgreSQL stops polling the WAL archive and
// switches to streaming replication. It does not apply when no streaming
// connection is available, nor when restoring on behalf of pg_rewind:
// pg_rewind cannot fall back to streaming replication, and a stale flag would
// make it abort on a segment that is available in the archive
func shouldUseEndOfWALStreamFlag(cluster *cnpgv1.Cluster, podName string, rewindMode bool) bool {
if rewindMode {
return false
}
return isStreamingAvailable(cluster, podName)
}
// isStreamingAvailable checks if this pod can replicate via streaming connection. // isStreamingAvailable checks if this pod can replicate via streaming connection.
func isStreamingAvailable(cluster *cnpgv1.Cluster, podName string) bool { func isStreamingAvailable(cluster *cnpgv1.Cluster, podName string) bool {
if cluster == nil { if cluster == nil {

View File

@ -20,6 +20,7 @@ SPDX-License-Identifier: Apache-2.0
package common package common
import ( import (
barmanapi "github.com/cloudnative-pg/barman-cloud/pkg/api"
cnpgv1 "github.com/cloudnative-pg/cloudnative-pg/api/v1" cnpgv1 "github.com/cloudnative-pg/cloudnative-pg/api/v1"
. "github.com/onsi/ginkgo/v2" . "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega" . "github.com/onsi/gomega"
@ -96,3 +97,47 @@ var _ = Describe("resolveRestoreObjectStore", func() {
"cluster-server", "cluster-store"), "cluster-server", "cluster-store"),
) )
}) })
var _ = Describe("maxWALFilesPerInvocation", func() {
configWithMaxParallel := func(maxParallel int) *barmanapi.BarmanObjectStoreConfiguration {
return &barmanapi.BarmanObjectStoreConfiguration{
Wal: &barmanapi.WalBackupConfiguration{MaxParallel: maxParallel},
}
}
DescribeTable(
"computes how many WAL files a single invocation may fetch",
func(cfg *barmanapi.BarmanObjectStoreConfiguration, rewindMode bool, want int) {
Expect(maxWALFilesPerInvocation(cfg, rewindMode)).To(Equal(want))
},
Entry("no WAL configuration", &barmanapi.BarmanObjectStoreConfiguration{}, false, 1),
Entry("parallel restore configured", configWithMaxParallel(8), false, 8),
// pg_rewind walks the timeline backwards: prefetching must stay off no
// matter what the object store configuration asks for
Entry("rewind mode overrides the configured parallelism", configWithMaxParallel(8), true, 1),
)
})
var _ = Describe("shouldUseEndOfWALStreamFlag", func() {
clusterWithPrimary := func(currentPrimary string) *cnpgv1.Cluster {
return &cnpgv1.Cluster{
Status: cnpgv1.ClusterStatus{CurrentPrimary: currentPrimary},
}
}
DescribeTable(
"decides whether the end-of-wal-stream flag machinery applies",
func(cluster *cnpgv1.Cluster, podName string, rewindMode bool, want bool) {
Expect(shouldUseEndOfWALStreamFlag(cluster, podName, rewindMode)).To(Equal(want))
},
Entry("replica with streaming available", clusterWithPrimary("cluster-1"), "cluster-2", false, true),
Entry("primary cannot stream from anyone", clusterWithPrimary("cluster-1"), "cluster-1", false, false),
// pg_rewind cannot fall back to streaming replication: the flag machinery
// must stay off even where a standby would use it
Entry("rewind mode wins over streaming availability", clusterWithPrimary("cluster-1"), "cluster-2", true, false),
)
})