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
127 changes: 101 additions & 26 deletions core/cmd/backup/run.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ SPDX-License-Identifier: Apache-2.0
package backup

import (
"context"
"encoding/json"
"fmt"
"os"
Expand Down Expand Up @@ -95,6 +96,15 @@ func runBackup(cmd *cobra.Command, _ []string) error {
}
defer kopiaClient.Close(cmd.Context())

// Set the per-cluster tier1 compression policy before uploading so that
// the base backup snapshot is compressed. This overrides the repository
// global policy for this cluster's source.
if err := setTier1CompressionPolicy(cmd.Context(), kopiaClient, &configuration); err != nil {
return cli.NewCodedError(
fmt.Errorf("while setting the tier1 compression policy: %w", err),
backupfailure.RepositoryError.ExitCode)
}

conn, err := pgx.Connect(cmd.Context(), configuration.Source.StandardDSN)
if err != nil {
return cli.NewCodedError(
Expand Down Expand Up @@ -139,34 +149,19 @@ func runBackup(cmd *cobra.Command, _ []string) error {
}

for {
var tier2RetentionPolicy string
if configuration.Tier2RetentionPolicy != nil {
policy := kopiaWrapper.RetentionPolicy{
KeepLatest: configuration.Tier2RetentionPolicy.KeepLatest,
KeepHourly: configuration.Tier2RetentionPolicy.KeepHourly,
KeepDaily: configuration.Tier2RetentionPolicy.KeepDaily,
KeepWeekly: configuration.Tier2RetentionPolicy.KeepWeekly,
KeepMonthly: configuration.Tier2RetentionPolicy.KeepMonthly,
KeepAnnual: configuration.Tier2RetentionPolicy.KeepAnnual,
}

content, err := json.Marshal(policy)
if err != nil {
contextLogger.Error(err, "Error while serializing the tier2 retention policy, skipping")
} else {
tier2RetentionPolicy = string(content)
}
}
//nolint:gosec // postgres timeline is uint32 in practice, fits int32
timeline := int32(metadata.Timeline)

result, err := grpcClient.CloseBackup(cmd.Context(), &grpc.CloseBackupRequest{
ClusterName: kopiaClient.GetHostname(),
BackupName: metadata.Name,
Timeline: int32(metadata.Timeline), //nolint:gosec // postgres timeline is uint32 in practice, fits int32
StartWal: metadata.StartWAL,
EndWal: metadata.EndWAL,
SegmentSize: metadata.SegmentSize,
SendToTier2: tier2,
Tier2RetentionPolicy: tier2RetentionPolicy,
ClusterName: kopiaClient.GetHostname(),
BackupName: metadata.Name,
Timeline: timeline,
StartWal: metadata.StartWAL,
EndWal: metadata.EndWAL,
SegmentSize: metadata.SegmentSize,
SendToTier2: tier2,
Tier2RetentionPolicy: marshalTier2RetentionPolicy(cmd.Context(), &configuration),
Tier2CompressionPolicy: marshalTier2CompressionPolicy(cmd.Context(), &configuration),
})
if err != nil {
return cli.NewCodedError(
Expand Down Expand Up @@ -199,6 +194,86 @@ func runBackup(cmd *cobra.Command, _ []string) error {
return nil
}

// setTier1CompressionPolicy applies the per-cluster tier1 compression policy,
// if configured, to the cluster's source in the tier1 repository.
func setTier1CompressionPolicy(
ctx context.Context,
client *kopia.MultiConnection,
configuration *config.Data,
) error {
policy := toKopiaCompressionPolicy(configuration.Tier1CompressionPolicy)
if policy.IsZero() {
return nil
}

target := kopiaWrapper.Target{
Username: client.GetUsername(),
Hostname: client.GetHostname(),
}

return client.SetCompressionPolicy(ctx, target, policy)
}

// toKopiaCompressionPolicy converts a config compression policy into the Kopia
// wrapper representation. A nil input yields the zero policy.
func toKopiaCompressionPolicy(p *config.CompressionPolicy) kopiaWrapper.CompressionPolicy {
if p == nil {
return kopiaWrapper.CompressionPolicy{}
}

return kopiaWrapper.CompressionPolicy{
Algorithm: p.Algorithm,
MinSize: p.MinSize,
MaxSize: p.MaxSize,
}
}

// marshalTier2RetentionPolicy serializes the tier2 retention policy to the
// JSON representation expected by the WAL server. It returns an empty string
// when no policy is configured or serialization fails.
func marshalTier2RetentionPolicy(ctx context.Context, configuration *config.Data) string {
if configuration.Tier2RetentionPolicy == nil {
return ""
}

policy := kopiaWrapper.RetentionPolicy{
KeepLatest: configuration.Tier2RetentionPolicy.KeepLatest,
KeepHourly: configuration.Tier2RetentionPolicy.KeepHourly,
KeepDaily: configuration.Tier2RetentionPolicy.KeepDaily,
KeepWeekly: configuration.Tier2RetentionPolicy.KeepWeekly,
KeepMonthly: configuration.Tier2RetentionPolicy.KeepMonthly,
KeepAnnual: configuration.Tier2RetentionPolicy.KeepAnnual,
}

content, err := json.Marshal(policy)
if err != nil {
log.FromContext(ctx).Error(err, "Error while serializing the tier2 retention policy, skipping")

return ""
}

return string(content)
}

// marshalTier2CompressionPolicy serializes the tier2 compression policy to the
// JSON representation expected by the WAL server. It returns an empty string
// when no policy is configured or serialization fails.
func marshalTier2CompressionPolicy(ctx context.Context, configuration *config.Data) string {
policy := toKopiaCompressionPolicy(configuration.Tier2CompressionPolicy)
if policy.IsZero() {
return ""
}

content, err := json.Marshal(policy)
if err != nil {
log.FromContext(ctx).Error(err, "Error while serializing the tier2 compression policy, skipping")

return ""
}

return string(content)
}

//nolint:gochecknoinits
func init() {
// Here you will define your flags and configuration settings.
Expand Down
59 changes: 53 additions & 6 deletions core/cmd/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,56 @@ import (
"github.com/cloudnative-pg/machinery/pkg/log"
"github.com/thejerf/suture/v4"

"github.com/cloudnative-pg/klio/core/internal/kopia"
"github.com/cloudnative-pg/klio/core/internal/server"
"github.com/cloudnative-pg/klio/core/internal/server/kopiaconfig"
"github.com/cloudnative-pg/klio/core/pkg/config"
)

// applyGlobalCompressionPolicy sets the repository-wide (global) Kopia
// compression policy using the passed persistent config file. It is a no-op
// when the policy carries no settings. This runs before the Kopia servers
// start, so the direct write to the repository predates any server cache.
func applyGlobalCompressionPolicy(
ctx context.Context,
configFile string,
compression config.CompressionServerConfig,
) error {
if compression.IsZero() {
return nil
}

kopiaBinary, err := kopia.LookupBinary()
if err != nil {
return err
}

client := &kopia.Client{
KopiaBinary: kopiaBinary,
ConfigFile: configFile,
}

return client.SetKopiaGlobalCompressionPolicy(ctx, kopia.CompressionPolicy{
Algorithm: compression.Algorithm,
MinSize: compression.MinSize,
MaxSize: compression.MaxSize,
})
}

// setupTier1KopiaConfig connects the tier1 config file to the repository and
// applies the tier1 repository-wide compression policy.
func setupTier1KopiaConfig(ctx context.Context, configFile string, cfg *config.Tier1Config) error {
if err := kopiaconfig.CreateTier1KopiaConfigFile(ctx, configFile, cfg); err != nil {
return fmt.Errorf("error creating tier1 kopia config file: %w", err)
}

if err := applyGlobalCompressionPolicy(ctx, configFile, cfg.Compression); err != nil {
return fmt.Errorf("error setting tier1 global compression policy: %w", err)
}

return nil
}

type serverOpts struct {
tier1 bool
tier2 bool
Expand Down Expand Up @@ -192,12 +237,8 @@ func runServer(ctx context.Context, opts serverOpts) error {
}
}()

if err := kopiaconfig.CreateTier1KopiaConfigFile(
ctx,
tier1ConfigFileName,
&opts.cfg.Tier1,
); err != nil {
return fmt.Errorf("error creating tier1 kopia config file: %w", err)
if err := setupTier1KopiaConfig(ctx, tier1ConfigFileName, &opts.cfg.Tier1); err != nil {
return err
}

tier1 := suture.NewSimple("tier1")
Expand Down Expand Up @@ -229,6 +270,12 @@ func runServer(ctx context.Context, opts serverOpts) error {
tier2RWConfigFileName = tier2Configs.rwConfigFileName
tier2ROConfigFileName = tier2Configs.roConfigFileName

if err := applyGlobalCompressionPolicy(
ctx, tier2RWConfigFileName, opts.cfg.Tier2.Compression,
); err != nil {
return fmt.Errorf("error setting tier2 global compression policy: %w", err)
}

tier2 := suture.NewSimple("tier2")
tier2.Add(&server.Tier2KopiaServer{
Config: opts.cfg,
Expand Down
3 changes: 3 additions & 0 deletions core/internal/client/klioclient/interfaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,9 @@ type Client interface {
// SetRetentionPolicy sets the retention policy for backups of this cluster.
SetRetentionPolicy(ctx context.Context, t kopia.Target, p kopia.RetentionPolicy) error

// SetCompressionPolicy sets the compression policy for backups of this cluster.
SetCompressionPolicy(ctx context.Context, t kopia.Target, policy kopia.CompressionPolicy) error

// GetRetentionPolicy gets the currently applied retention policy for this cluster.
GetRetentionPolicy(ctx context.Context, t kopia.Target) (*kopia.RetentionPolicy, error)

Expand Down
13 changes: 13 additions & 0 deletions core/internal/client/klioclient/kopia/multiconnect.go
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,19 @@ func (s *MultiConnection) SetRetentionPolicy(
return s.Tier1.SetRetentionPolicy(ctx, t, p)
}

// SetCompressionPolicy implements the Client interface.
func (s *MultiConnection) SetCompressionPolicy(
ctx context.Context,
t kopia.Target,
policy kopia.CompressionPolicy,
) error {
if s.Tier1 == nil {
return ErrUnsupportedWriteOperation
}

return s.Tier1.SetCompressionPolicy(ctx, t, policy)
}

// GetRetentionPolicy implements the Client interface.
func (s *MultiConnection) GetRetentionPolicy(
ctx context.Context,
Expand Down
5 changes: 5 additions & 0 deletions core/internal/client/klioclient/kopia/retention.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,11 @@ func (s *Connection) SetRetentionPolicy(ctx context.Context, t kopia.Target, p k
return s.kopia.SetKopiaPolicy(ctx, t, &p)
}

// SetCompressionPolicy sets the compression policy for backups of this cluster.
func (s *Connection) SetCompressionPolicy(ctx context.Context, t kopia.Target, policy kopia.CompressionPolicy) error {
return s.kopia.SetKopiaCompressionPolicy(ctx, t, policy)
}

// GetRetentionPolicy gets the currently applied retention policy for this cluster.
func (s *Connection) GetRetentionPolicy(ctx context.Context, t kopia.Target) (*kopia.RetentionPolicy, error) {
policy, err := s.kopia.GetCurrentKopiaPolicy(ctx, t)
Expand Down
15 changes: 15 additions & 0 deletions core/internal/consumer/backup.go
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,21 @@ func (d *Backup) relayAndMaintain(ctx context.Context, task *queue.BackupTask, e
func (d *Backup) relayTier2(ctx context.Context, task *queue.BackupTask, entries []kopia.Manifest) error {
sources := manifestListToDescriptors(entries)

// Set the per-cluster tier2 compression policy before migrating so that
// the data relayed to tier2 is compressed. This overrides the tier2
// repository global policy for this cluster's source. A direct write is
// unavoidable here (the consumer has no tier2 server connection) and is
// safe: it only writes a policy manifest, never rewrites a live backup.
if p := task.Tier2CompressionPolicy; p != nil && !p.IsZero() && len(entries) > 0 {
target := kopia.Target{
Username: entries[0].Source.UserName,
Hostname: task.ClusterName,
}
if err := d.tier2Kopia.SetKopiaCompressionPolicy(ctx, target, *p); err != nil {
return fmt.Errorf("while setting the tier2 compression policy: %w", err)
}
}

if err := d.tier2Kopia.MigrateSnapshots(ctx, kopia.SnapshotMigrateOpts{
SourceConfig: d.opts.Tier1KopiaConfig,
Sources: sources,
Expand Down
19 changes: 15 additions & 4 deletions core/internal/grpc/klio_wal.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

20 changes: 20 additions & 0 deletions core/internal/kopia/data.go
Original file line number Diff line number Diff line change
Expand Up @@ -162,6 +162,26 @@ type RetentionPolicy struct {
KeepAnnual *int `json:"keepAnnual,omitempty"`
}

// CompressionPolicy describes the compression policy for a source.
type CompressionPolicy struct {
// Algorithm is the name of the Kopia compression algorithm to use.
// The special value "none" disables compression.
Algorithm string `json:"compressionAlgorithm,omitempty"`

// MinSize is the minimum file size, in bytes, to attempt compression for.
// Files smaller than this are stored uncompressed. Zero means no minimum.
MinSize int64 `json:"compressionMinSize,omitempty"`

// MaxSize is the maximum file size, in bytes, to attempt compression for.
// Files larger than this are stored uncompressed. Zero means no maximum.
MaxSize int64 `json:"compressionMaxSize,omitempty"`
}

// IsZero reports whether the policy carries no compression settings.
func (p CompressionPolicy) IsZero() bool {
return p.Algorithm == "" && p.MinSize == 0 && p.MaxSize == 0
}

// Target is used to point a Kopia transaction to the set of snapshots
// having the specified Hostname and Username.
type Target struct {
Expand Down
Loading
Loading