Merge branch 'master' into fix/logical-backup-job-cleanup

This commit is contained in:
Felix Kunde
2026-06-17 10:28:20 +02:00
committed by GitHub
52 changed files with 925 additions and 2269 deletions
+12
View File
@@ -748,6 +748,18 @@ var OperatorConfigCRDResourceValidation = apiextv1.CustomResourceValidation{
"enable_replica_pooler_load_balancer": {
Type: "boolean",
},
"enable_master_node_port": {
Type: "boolean",
},
"enable_master_pooler_node_port": {
Type: "boolean",
},
"enable_replica_node_port": {
Type: "boolean",
},
"enable_replica_pooler_node_port": {
Type: "boolean",
},
"external_traffic_policy": {
Type: "string",
Enum: []apiextv1.JSON{
@@ -137,17 +137,24 @@ type OperatorTimeouts struct {
// LoadBalancerConfiguration defines the LB configuration
type LoadBalancerConfiguration struct {
DbHostedZone string `json:"db_hosted_zone,omitempty"`
EnableMasterLoadBalancer bool `json:"enable_master_load_balancer,omitempty"`
EnableMasterPoolerLoadBalancer bool `json:"enable_master_pooler_load_balancer,omitempty"`
EnableReplicaLoadBalancer bool `json:"enable_replica_load_balancer,omitempty"`
EnableReplicaPoolerLoadBalancer bool `json:"enable_replica_pooler_load_balancer,omitempty"`
CustomServiceAnnotations map[string]string `json:"custom_service_annotations,omitempty"`
MasterDNSNameFormat config.StringTemplate `json:"master_dns_name_format,omitempty"`
MasterLegacyDNSNameFormat config.StringTemplate `json:"master_legacy_dns_name_format,omitempty"`
ReplicaDNSNameFormat config.StringTemplate `json:"replica_dns_name_format,omitempty"`
ReplicaLegacyDNSNameFormat config.StringTemplate `json:"replica_legacy_dns_name_format,omitempty"`
ExternalTrafficPolicy string `json:"external_traffic_policy" default:"Cluster"`
DbHostedZone string `json:"db_hosted_zone,omitempty"`
EnableMasterLoadBalancer bool `json:"enable_master_load_balancer,omitempty"`
EnableMasterPoolerLoadBalancer bool `json:"enable_master_pooler_load_balancer,omitempty"`
EnableReplicaLoadBalancer bool `json:"enable_replica_load_balancer,omitempty"`
EnableReplicaPoolerLoadBalancer bool `json:"enable_replica_pooler_load_balancer,omitempty"`
// kept in LoadBalancerConfiguration because all the other parameters apply here too
EnableMasterNodePort bool `json:"enable_master_node_port,omitempty"`
EnableMasterPoolerNodePort bool `json:"enable_master_pooler_node_port,omitempty"`
EnableReplicaNodePort bool `json:"enable_replica_node_port,omitempty"`
EnableReplicaPoolerNodePort bool `json:"enable_replica_pooler_node_port,omitempty"`
CustomServiceAnnotations map[string]string `json:"custom_service_annotations,omitempty"`
MasterDNSNameFormat config.StringTemplate `json:"master_dns_name_format,omitempty"`
MasterLegacyDNSNameFormat config.StringTemplate `json:"master_legacy_dns_name_format,omitempty"`
ReplicaDNSNameFormat config.StringTemplate `json:"replica_dns_name_format,omitempty"`
ReplicaLegacyDNSNameFormat config.StringTemplate `json:"replica_legacy_dns_name_format,omitempty"`
ExternalTrafficPolicy string `json:"external_traffic_policy" default:"Cluster"`
}
// AWSGCPConfiguration defines the configuration for AWS
@@ -168,8 +175,8 @@ type AWSGCPConfiguration struct {
// OperatorDebugConfiguration defines options for the debug mode
type OperatorDebugConfiguration struct {
DebugLogging bool `json:"debug_logging,omitempty"`
EnableDBAccess bool `json:"enable_database_access,omitempty"`
DebugLogging *bool `json:"debug_logging,omitempty"`
EnableDBAccess *bool `json:"enable_database_access,omitempty"`
}
// TeamsAPIConfiguration defines the configuration of TeamsAPI
+25 -1
View File
@@ -108,7 +108,7 @@ spec:
description: load balancers' source ranges are the same for master
and replica services
items:
pattern: ^(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\/(\d|[1-2]\d|3[0-2])$
pattern: ^((\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\/(\d|[1-2]\d|3[0-2])|(([0-9a-fA-F]{1,4}:){7}[0-9a-fA-F]{1,4}|([0-9a-fA-F]{1,4}:){1,7}:|([0-9a-fA-F]{1,4}:){1,6}:[0-9a-fA-F]{1,4}|([0-9a-fA-F]{1,4}:){1,5}(:[0-9a-fA-F]{1,4}){1,2}|([0-9a-fA-F]{1,4}:){1,4}(:[0-9a-fA-F]{1,4}){1,3}|([0-9a-fA-F]{1,4}:){1,3}(:[0-9a-fA-F]{1,4}){1,4}|([0-9a-fA-F]{1,4}:){1,2}(:[0-9a-fA-F]{1,4}){1,5}|[0-9a-fA-F]{1,4}:((:[0-9a-fA-F]{1,4}){1,6})|:((:[0-9a-fA-F]{1,4}){1,7}|:))\/(12[0-8]|1[01][0-9]|[1-9]?[0-9]))$
type: string
nullable: true
type: array
@@ -274,14 +274,26 @@ spec:
vars that enable load balancers are pointers because it is important to know if any of them is omitted from the Postgres manifest
in that case the var evaluates to nil and the value is taken from the operator config
type: boolean
enableMasterNodePort:
description: |-
vars to enable and configure nodeport services
set ports to 0 or nil to let kubernetes decide which port to use
overrides loadbalancer configuration
type: boolean
enableMasterPoolerLoadBalancer:
type: boolean
enableMasterPoolerNodePort:
type: boolean
enableReplicaConnectionPooler:
type: boolean
enableReplicaLoadBalancer:
type: boolean
enableReplicaNodePort:
type: boolean
enableReplicaPoolerLoadBalancer:
type: boolean
enableReplicaPoolerNodePort:
type: boolean
enableShmVolume:
type: boolean
env:
@@ -3409,6 +3421,12 @@ spec:
pattern: '^\ *((Mon|Tue|Wed|Thu|Fri|Sat|Sun):(2[0-3]|[01]?\d):([0-5]?\d)|(2[0-3]|[01]?\d):([0-5]?\d))-((2[0-3]|[01]?\d):([0-5]?\d)|(2[0-3]|[01]?\d):([0-5]?\d))\ *$'
type: string
type: array
masterNodePort:
format: int32
type: integer
masterPoolerNodePort:
format: int32
type: integer
masterServiceAnnotations:
additionalProperties:
type: string
@@ -3713,6 +3731,12 @@ spec:
replicaLoadBalancer:
description: deprecated
type: boolean
replicaNodePort:
format: int32
type: integer
replicaPoolerNodePort:
format: int32
type: integer
replicaServiceAnnotations:
additionalProperties:
type: string
+13 -1
View File
@@ -63,6 +63,18 @@ type PostgresSpec struct {
EnableReplicaLoadBalancer *bool `json:"enableReplicaLoadBalancer,omitempty"`
EnableReplicaPoolerLoadBalancer *bool `json:"enableReplicaPoolerLoadBalancer,omitempty"`
// vars to enable and configure nodeport services
// set ports to 0 or nil to let kubernetes decide which port to use
// overrides loadbalancer configuration
EnableMasterNodePort *bool `json:"enableMasterNodePort,omitempty"`
MasterNodePort *int32 `json:"masterNodePort,omitempty"`
EnableMasterPoolerNodePort *bool `json:"enableMasterPoolerNodePort,omitempty"`
MasterPoolerNodePort *int32 `json:"masterPoolerNodePort,omitempty"`
EnableReplicaNodePort *bool `json:"enableReplicaNodePort,omitempty"`
ReplicaNodePort *int32 `json:"replicaNodePort,omitempty"`
EnableReplicaPoolerNodePort *bool `json:"enableReplicaPoolerNodePort,omitempty"`
ReplicaPoolerNodePort *int32 `json:"replicaPoolerNodePort,omitempty"`
// deprecated load balancer settings maintained for backward compatibility
// see "Load balancers" operator docs
UseLoadBalancer *bool `json:"useLoadBalancer,omitempty"`
@@ -71,7 +83,7 @@ type PostgresSpec struct {
// load balancers' source ranges are the same for master and replica services
// +nullable
// +kubebuilder:validation:items:Pattern=`^(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\/(\d|[1-2]\d|3[0-2])$`
// +kubebuilder:validation:items:Pattern=`^((\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\/(\d|[1-2]\d|3[0-2])|(([0-9a-fA-F]{1,4}:){7}[0-9a-fA-F]{1,4}|([0-9a-fA-F]{1,4}:){1,7}:|([0-9a-fA-F]{1,4}:){1,6}:[0-9a-fA-F]{1,4}|([0-9a-fA-F]{1,4}:){1,5}(:[0-9a-fA-F]{1,4}){1,2}|([0-9a-fA-F]{1,4}:){1,4}(:[0-9a-fA-F]{1,4}){1,3}|([0-9a-fA-F]{1,4}:){1,3}(:[0-9a-fA-F]{1,4}){1,4}|([0-9a-fA-F]{1,4}:){1,2}(:[0-9a-fA-F]{1,4}){1,5}|[0-9a-fA-F]{1,4}:((:[0-9a-fA-F]{1,4}){1,6})|:((:[0-9a-fA-F]{1,4}){1,7}|:))\/(12[0-8]|1[01][0-9]|[1-9]?[0-9]))$`
// +optional
AllowedSourceRanges []string `json:"allowedSourceRanges"`
+45
View File
@@ -5,6 +5,7 @@ import (
"encoding/json"
"errors"
"reflect"
"regexp"
"testing"
"time"
@@ -810,3 +811,47 @@ func TestPostgresqlClone(t *testing.T) {
})
}
}
func TestAllowedSourceRangesPattern(t *testing.T) {
// pattern used in CRD validation for allowedSourceRanges
pattern := `^((\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\.(\d|[1-9]\d|1\d\d|2[0-4]\d|25[0-5])\/(\d|[1-2]\d|3[0-2])|(([0-9a-fA-F]{1,4}:){7}[0-9a-fA-F]{1,4}|([0-9a-fA-F]{1,4}:){1,7}:|([0-9a-fA-F]{1,4}:){1,6}:[0-9a-fA-F]{1,4}|([0-9a-fA-F]{1,4}:){1,5}(:[0-9a-fA-F]{1,4}){1,2}|([0-9a-fA-F]{1,4}:){1,4}(:[0-9a-fA-F]{1,4}){1,3}|([0-9a-fA-F]{1,4}:){1,3}(:[0-9a-fA-F]{1,4}){1,4}|([0-9a-fA-F]{1,4}:){1,2}(:[0-9a-fA-F]{1,4}){1,5}|[0-9a-fA-F]{1,4}:((:[0-9a-fA-F]{1,4}){1,6})|:((:[0-9a-fA-F]{1,4}){1,7}|:))\/(12[0-8]|1[01][0-9]|[1-9]?[0-9]))$`
re := regexp.MustCompile(pattern)
valid := []string{
// IPv4
"192.168.1.0/24",
"0.0.0.0/0",
"127.0.0.1/32",
"10.0.0.0/8",
"185.85.220.0/22",
// IPv6
"fd01::/48",
"::1/128",
"::/0",
"2001:db8::/32",
"fe80::1/64",
"2001:0db8:85a3:0000:0000:8a2e:0370:7334/128",
}
invalid := []string{
"999.999.999.999/24",
"192.168.1.0/33",
"192.168.1.0",
"not-an-ip",
"fd01::/129",
"::gggg/64",
"",
}
for _, cidr := range valid {
if !re.MatchString(cidr) {
t.Errorf("expected %q to match allowedSourceRanges pattern", cidr)
}
}
for _, cidr := range invalid {
if re.MatchString(cidr) {
t.Errorf("expected %q NOT to match allowedSourceRanges pattern", cidr)
}
}
}
@@ -476,7 +476,7 @@ func (in *OperatorConfigurationData) DeepCopyInto(out *OperatorConfigurationData
out.Timeouts = in.Timeouts
in.LoadBalancer.DeepCopyInto(&out.LoadBalancer)
out.AWSGCP = in.AWSGCP
out.OperatorDebug = in.OperatorDebug
in.OperatorDebug.DeepCopyInto(&out.OperatorDebug)
in.TeamsAPI.DeepCopyInto(&out.TeamsAPI)
out.LoggingRESTAPI = in.LoggingRESTAPI
out.Scalyr = in.Scalyr
@@ -532,6 +532,16 @@ func (in *OperatorConfigurationList) DeepCopyObject() runtime.Object {
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in *OperatorDebugConfiguration) DeepCopyInto(out *OperatorDebugConfiguration) {
*out = *in
if in.DebugLogging != nil {
in, out := &in.DebugLogging, &out.DebugLogging
*out = new(bool)
**out = **in
}
if in.EnableDBAccess != nil {
in, out := &in.EnableDBAccess, &out.EnableDBAccess
*out = new(bool)
**out = **in
}
return
}
@@ -725,6 +735,46 @@ func (in *PostgresSpec) DeepCopyInto(out *PostgresSpec) {
*out = new(bool)
**out = **in
}
if in.EnableMasterNodePort != nil {
in, out := &in.EnableMasterNodePort, &out.EnableMasterNodePort
*out = new(bool)
**out = **in
}
if in.MasterNodePort != nil {
in, out := &in.MasterNodePort, &out.MasterNodePort
*out = new(int32)
**out = **in
}
if in.EnableMasterPoolerNodePort != nil {
in, out := &in.EnableMasterPoolerNodePort, &out.EnableMasterPoolerNodePort
*out = new(bool)
**out = **in
}
if in.MasterPoolerNodePort != nil {
in, out := &in.MasterPoolerNodePort, &out.MasterPoolerNodePort
*out = new(int32)
**out = **in
}
if in.EnableReplicaNodePort != nil {
in, out := &in.EnableReplicaNodePort, &out.EnableReplicaNodePort
*out = new(bool)
**out = **in
}
if in.ReplicaNodePort != nil {
in, out := &in.ReplicaNodePort, &out.ReplicaNodePort
*out = new(int32)
**out = **in
}
if in.EnableReplicaPoolerNodePort != nil {
in, out := &in.EnableReplicaPoolerNodePort, &out.EnableReplicaPoolerNodePort
*out = new(bool)
**out = **in
}
if in.ReplicaPoolerNodePort != nil {
in, out := &in.ReplicaPoolerNodePort, &out.ReplicaPoolerNodePort
*out = new(int32)
**out = **in
}
if in.UseLoadBalancer != nil {
in, out := &in.UseLoadBalancer, &out.UseLoadBalancer
*out = new(bool)
+19 -1
View File
@@ -135,7 +135,7 @@ func New(cfg Config, kubeClient k8sutil.KubernetesClient, pgSpec acidv1.Postgres
})
passwordEncryption, ok := pgSpec.Spec.PostgresqlParam.Parameters["password_encryption"]
if !ok {
passwordEncryption = "md5"
passwordEncryption = "scram-sha-256"
}
cluster := &Cluster{
@@ -859,6 +859,14 @@ func (c *Cluster) compareServices(old, new *v1.Service) (bool, string) {
return false, "new service's ExternalTrafficPolicy does not match the current one"
}
if len(old.Spec.Ports) > 0 && len(new.Spec.Ports) > 0 {
// we need to check whether the new port is not zero (=user-defined)
// and only overwrite if it is
if new.Spec.Ports[0].NodePort != 0 && old.Spec.Ports[0].NodePort != new.Spec.Ports[0].NodePort {
return false, "new service's NodePort does not match the current one"
}
}
return true, ""
}
@@ -893,6 +901,16 @@ func (c *Cluster) compareLogicalBackupJob(cur, new *batchv1.CronJob) *compareLog
reasons = append(reasons, fmt.Sprintf("new job's env PG_VERSION %q does not match the current one %q", newPgVersion, curPgVersion))
}
if !reflect.DeepEqual(cur.Labels, new.Labels) {
match = false
reasons = append(reasons, "new job's labels do not match the current ones")
}
if !reflect.DeepEqual(cur.Spec.JobTemplate.Labels, new.Spec.JobTemplate.Labels) {
match = false
reasons = append(reasons, "new job's template labels do not match the current ones")
}
needsReplace := false
contReasons := make([]string, 0)
needsReplace, contReasons = c.compareContainers("cronjob container", cur.Spec.JobTemplate.Spec.Template.Spec.Containers, new.Spec.JobTemplate.Spec.Template.Spec.Containers, needsReplace, contReasons)
+61 -18
View File
@@ -33,8 +33,8 @@ const (
replicationUserName = "standby"
poolerUserName = "pooler"
adminUserName = "admin"
exampleSpiloConfig = `{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 md5","host all all 0.0.0.0/0 md5"]},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_connections":"100","max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`
spiloConfigDiff = `{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 md5","host all all 0.0.0.0/0 md5"]},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`
exampleSpiloConfig = `{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 scram-sha-256","host all all 0.0.0.0/0 scram-sha-256"]},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_connections":"100","max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`
spiloConfigDiff = `{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 scram-sha-256","host all all 0.0.0.0/0 scram-sha-256"]},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`
)
var logger = logrus.New().WithField("test", "cluster")
@@ -1186,11 +1186,11 @@ func TestCompareSpiloConfiguration(t *testing.T) {
ExpectedResult bool
}{
{
`{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 md5","host all all 0.0.0.0/0 md5"]},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_connections":"100","max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`,
`{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 scram-sha-256","host all all 0.0.0.0/0 scram-sha-256"]},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_connections":"100","max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`,
true,
},
{
`{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 md5","host all all 0.0.0.0/0 md5"]},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_connections":"200","max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`,
`{"postgresql":{"bin_dir":"/usr/lib/postgresql/12/bin","parameters":{"autovacuum_analyze_scale_factor":"0.1"},"pg_hba":["hostssl all all 0.0.0.0/0 scram-sha-256","host all all 0.0.0.0/0 scram-sha-256"]},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"postgresql":{"parameters":{"max_connections":"200","max_locks_per_transaction":"64","max_worker_processes":"4"}}}}}`,
true,
},
{
@@ -1334,7 +1334,8 @@ func newService(
svcType v1.ServiceType,
sourceRanges []string,
selector map[string]string,
policy v1.ServiceExternalTrafficPolicyType) *v1.Service {
policy v1.ServiceExternalTrafficPolicyType,
nodePort *int32) *v1.Service {
svc := &v1.Service{
Spec: v1.ServiceSpec{
Selector: selector,
@@ -1344,6 +1345,16 @@ func newService(
},
}
svc.Annotations = annotations
if nodePort != nil {
svc.Spec.Ports = []v1.ServicePort{
{
Name: "port",
NodePort: *nodePort,
},
}
}
return svc
}
@@ -1370,6 +1381,7 @@ func TestCompareServices(t *testing.T) {
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil,
defaultPolicy,
nil,
)
ownerRef := metav1.OwnerReference{
@@ -1381,6 +1393,9 @@ func TestCompareServices(t *testing.T) {
serviceWithOwnerReference.ObjectMeta.OwnerReferences = append(serviceWithOwnerReference.ObjectMeta.OwnerReferences, ownerRef)
portZero := int32(0)
portNotZero := int32(1337)
tests := []struct {
about string
current *v1.Service
@@ -1396,14 +1411,14 @@ func TestCompareServices(t *testing.T) {
},
v1.ServiceTypeClusterIP,
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
new: newService(
map[string]string{
constants.ZalandoDNSNameAnnotation: "clstr.acid.zalan.do",
},
v1.ServiceTypeClusterIP,
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
match: true,
},
{
@@ -1414,14 +1429,14 @@ func TestCompareServices(t *testing.T) {
},
v1.ServiceTypeClusterIP,
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
new: newService(
map[string]string{
constants.ZalandoDNSNameAnnotation: "clstr.acid.zalan.do",
},
v1.ServiceTypeLoadBalancer,
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
match: false,
reason: `new service's type "LoadBalancer" does not match the current one "ClusterIP"`,
},
@@ -1433,14 +1448,14 @@ func TestCompareServices(t *testing.T) {
},
v1.ServiceTypeLoadBalancer,
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
new: newService(
map[string]string{
constants.ZalandoDNSNameAnnotation: "clstr.acid.zalan.do",
},
v1.ServiceTypeLoadBalancer,
[]string{"185.249.56.0/22"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
match: false,
reason: `new service's LoadBalancerSourceRange does not match the current one`,
},
@@ -1452,14 +1467,14 @@ func TestCompareServices(t *testing.T) {
},
v1.ServiceTypeLoadBalancer,
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
new: newService(
map[string]string{
constants.ZalandoDNSNameAnnotation: "clstr.acid.zalan.do",
},
v1.ServiceTypeLoadBalancer,
[]string{},
nil, defaultPolicy),
nil, defaultPolicy, nil),
match: false,
reason: `new service's LoadBalancerSourceRange does not match the current one`,
},
@@ -1471,7 +1486,7 @@ func TestCompareServices(t *testing.T) {
},
v1.ServiceTypeClusterIP,
[]string{"128.141.0.0/16", "137.138.0.0/16"},
nil, defaultPolicy),
nil, defaultPolicy, nil),
new: serviceWithOwnerReference,
match: false,
},
@@ -1481,12 +1496,12 @@ func TestCompareServices(t *testing.T) {
map[string]string{},
v1.ServiceTypeClusterIP,
[]string{},
nil, defaultPolicy),
nil, defaultPolicy, nil),
new: newService(
map[string]string{},
v1.ServiceTypeClusterIP,
[]string{},
map[string]string{"cluster-name": "clstr", "spilo-role": "master"}, defaultPolicy),
map[string]string{"cluster-name": "clstr", "spilo-role": "master"}, defaultPolicy, nil),
match: false,
},
{
@@ -1495,14 +1510,42 @@ func TestCompareServices(t *testing.T) {
map[string]string{},
v1.ServiceTypeClusterIP,
[]string{},
nil, defaultPolicy),
nil, defaultPolicy, nil),
new: newService(
map[string]string{},
v1.ServiceTypeClusterIP,
[]string{},
nil, v1.ServiceExternalTrafficPolicyTypeLocal),
nil, v1.ServiceExternalTrafficPolicyTypeLocal, nil),
match: false,
},
{
about: "services differ on node port",
current: newService(
map[string]string{},
v1.ServiceTypeNodePort,
[]string{},
nil, defaultPolicy, &portZero),
new: newService(
map[string]string{},
v1.ServiceTypeNodePort,
[]string{},
nil, defaultPolicy, &portNotZero),
match: false,
},
{
about: "services do not differ on node port when requesting 0",
current: newService(
map[string]string{},
v1.ServiceTypeNodePort,
[]string{},
nil, defaultPolicy, &portNotZero),
new: newService(
map[string]string{},
v1.ServiceTypeNodePort,
[]string{},
nil, defaultPolicy, &portZero),
match: true,
},
}
for _, tt := range tests {
+37 -2
View File
@@ -566,7 +566,9 @@ func (c *Cluster) generateConnectionPoolerService(connectionPooler *ConnectionPo
},
}
if c.shouldCreateLoadBalancerForPoolerService(poolerRole, spec) {
if ok, port := c.shouldCreateNodePortForPoolerService(poolerRole, spec); ok {
c.configureNodePortService(&serviceSpec, port)
} else if c.shouldCreateLoadBalancerForPoolerService(poolerRole, spec) {
c.configureLoadBalanceService(&serviceSpec, spec.AllowedSourceRanges)
}
@@ -594,7 +596,9 @@ func (c *Cluster) generatePoolerServiceAnnotations(role PostgresRole, spec *acid
var dnsString string
annotations := c.getCustomServiceAnnotations(role, spec)
if c.shouldCreateLoadBalancerForPoolerService(role, spec) {
nodePort, _ := c.shouldCreateNodePortForPoolerService(role, spec)
if !nodePort && c.shouldCreateLoadBalancerForPoolerService(role, spec) {
// -repl suffix will be added by replicaDNSName
clusterNameWithPoolerSuffix := c.connectionPoolerName(Master)
if role == Master {
@@ -635,6 +639,37 @@ func (c *Cluster) shouldCreateLoadBalancerForPoolerService(role PostgresRole, sp
}
}
func (c *Cluster) shouldCreateNodePortForPoolerService(role PostgresRole, spec *acidv1.PostgresSpec) (bool, int32) {
switch role {
case Replica:
// if the value is explicitly set in a Postgresql manifest, follow this setting
if spec.EnableReplicaPoolerNodePort != nil {
port := int32(0)
if spec.ReplicaPoolerNodePort != nil {
port = *spec.ReplicaPoolerNodePort
}
return *spec.EnableReplicaPoolerNodePort, port
}
// otherwise, follow the operator configuration
return c.OpConfig.EnableReplicaPoolerNodePort, 0
case Master:
if spec.EnableMasterPoolerNodePort != nil {
port := int32(0)
if spec.MasterPoolerNodePort != nil {
port = *spec.MasterPoolerNodePort
}
return *spec.EnableMasterPoolerNodePort, port
}
return c.OpConfig.EnableMasterPoolerNodePort, 0
default:
panic(fmt.Sprintf("Unknown role %v", role))
}
}
func (c *Cluster) listPoolerPods(listOptions metav1.ListOptions) ([]v1.Pod, error) {
pods, err := c.KubeClient.Pods(c.Namespace).List(context.TODO(), listOptions)
if err != nil {
+178
View File
@@ -1154,3 +1154,181 @@ func TestConnectionPoolerServiceSpec(t *testing.T) {
}
}
}
func TestConnectionPoolerServiceType(t *testing.T) {
testName := "Test connection pooler service type selection"
cluster := New(
Config{
OpConfig: config.Config{
ProtectedRoles: []string{"admin"},
Auth: config.Auth{
SuperUsername: superUserName,
ReplicationUsername: replicationUserName,
},
ConnectionPooler: config.ConnectionPooler{
ConnectionPoolerDefaultCPURequest: "100m",
ConnectionPoolerDefaultCPULimit: "100m",
ConnectionPoolerDefaultMemoryRequest: "100Mi",
ConnectionPoolerDefaultMemoryLimit: "100Mi",
},
Resources: config.Resources{
EnableOwnerReferences: util.True(),
},
},
},
k8sutil.KubernetesClient{},
acidv1.Postgresql{},
logger,
eventRecorder,
)
cluster.Statefulset = &appsv1.StatefulSet{
ObjectMeta: metav1.ObjectMeta{
Name: "test-sts",
},
}
cluster.ConnectionPooler = map[PostgresRole]*ConnectionPoolerObjects{
Master: {
Deployment: nil,
Service: nil,
LookupFunction: false,
Role: Master,
},
Replica: {
Deployment: nil,
Service: nil,
LookupFunction: false,
Role: Replica,
},
}
tests := []struct {
subTest string
spec *acidv1.PostgresSpec
cluster *Cluster
expectedType map[PostgresRole]v1.ServiceType
}{
{
subTest: "default configuration -> ClusterIP for both",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeClusterIP,
Replica: v1.ServiceTypeClusterIP,
},
},
{
subTest: "LoadBalancer for both roles",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
EnableMasterPoolerLoadBalancer: boolToPointer(true),
EnableReplicaPoolerLoadBalancer: boolToPointer(true),
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeLoadBalancer,
Replica: v1.ServiceTypeLoadBalancer,
},
},
{
subTest: "LoadBalancer for master",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
EnableMasterPoolerLoadBalancer: boolToPointer(true),
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeLoadBalancer,
Replica: v1.ServiceTypeClusterIP,
},
},
{
subTest: "LoadBalancer for replica",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
EnableReplicaPoolerLoadBalancer: boolToPointer(true),
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeClusterIP,
Replica: v1.ServiceTypeLoadBalancer,
},
},
{
subTest: "NodePort for both roles",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
EnableMasterPoolerNodePort: boolToPointer(true),
EnableReplicaPoolerNodePort: boolToPointer(true),
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeNodePort,
Replica: v1.ServiceTypeNodePort,
},
},
{
subTest: "NodePort for master",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
EnableMasterPoolerNodePort: boolToPointer(true),
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeNodePort,
Replica: v1.ServiceTypeClusterIP,
},
},
{
subTest: "NodePort for replica",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
EnableReplicaPoolerNodePort: boolToPointer(true),
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeClusterIP,
Replica: v1.ServiceTypeNodePort,
},
},
{
subTest: "NodePort overrides LoadBalancer for both roles",
spec: &acidv1.PostgresSpec{
ConnectionPooler: &acidv1.ConnectionPooler{},
EnableMasterPoolerLoadBalancer: boolToPointer(true),
EnableReplicaPoolerLoadBalancer: boolToPointer(true),
EnableMasterPoolerNodePort: boolToPointer(true),
EnableReplicaPoolerNodePort: boolToPointer(true),
},
cluster: cluster,
expectedType: map[PostgresRole]v1.ServiceType{
Master: v1.ServiceTypeNodePort,
Replica: v1.ServiceTypeNodePort,
},
},
}
roles := []PostgresRole{Master, Replica}
for _, tt := range tests {
tt.cluster.Spec = *tt.spec
for _, role := range roles {
svc := tt.cluster.generateConnectionPoolerService(tt.cluster.ConnectionPooler[role])
expected, ok := tt.expectedType[role]
if !ok {
t.Fatalf("%s [%s]: missing expectedType for role %v", testName, tt.subTest, role)
}
if svc.Spec.Type != expected {
t.Errorf("%s [%s] role=%s: service Type is incorrect, got %s, expected %s",
testName, tt.subTest, role, svc.Spec.Type, expected)
}
}
}
}
+53 -5
View File
@@ -366,7 +366,7 @@ func generateSpiloJSONConfiguration(pg *acidv1.PostgresqlParam, patroni *acidv1.
config.Bootstrap = pgBootstrap{}
config.Bootstrap.Initdb = []interface{}{map[string]string{"auth-host": "md5"},
config.Bootstrap.Initdb = []interface{}{map[string]string{"auth-host": "scram-sha-256"},
map[string]string{"auth-local": "trust"}}
initdbOptionNames := []string{}
@@ -2004,6 +2004,37 @@ func (c *Cluster) shouldCreateLoadBalancerForService(role PostgresRole, spec *ac
}
func (c *Cluster) shouldCreateNodePortForService(role PostgresRole, spec *acidv1.PostgresSpec) (bool, int32) {
switch role {
case Replica:
// if the value is explicitly set in a Postgresql manifest, follow this setting
if spec.EnableReplicaNodePort != nil {
port := int32(0)
if spec.ReplicaNodePort != nil {
port = *spec.ReplicaNodePort
}
return *spec.EnableReplicaNodePort, port
}
// otherwise, follow the operator configuration
return c.OpConfig.EnableReplicaNodePort, 0
case Master:
if spec.EnableMasterNodePort != nil {
port := int32(0)
if spec.MasterNodePort != nil {
port = *spec.MasterNodePort
}
return *spec.EnableMasterNodePort, port
}
return c.OpConfig.EnableMasterNodePort, 0
default:
panic(fmt.Sprintf("Unknown role %v", role))
}
}
func (c *Cluster) generateService(role PostgresRole, spec *acidv1.PostgresSpec) *v1.Service {
serviceSpec := v1.ServiceSpec{
Ports: []v1.ServicePort{{Name: "postgresql", Port: pgPort, TargetPort: intstr.IntOrString{IntVal: pgPort}}},
@@ -2016,7 +2047,9 @@ func (c *Cluster) generateService(role PostgresRole, spec *acidv1.PostgresSpec)
serviceSpec.Selector = c.roleLabelsSet(false, role)
}
if c.shouldCreateLoadBalancerForService(role, spec) {
if ok, port := c.shouldCreateNodePortForService(role, spec); ok {
c.configureNodePortService(&serviceSpec, port)
} else if c.shouldCreateLoadBalancerForService(role, spec) {
c.configureLoadBalanceService(&serviceSpec, spec.AllowedSourceRanges)
}
@@ -2049,10 +2082,21 @@ func (c *Cluster) configureLoadBalanceService(serviceSpec *v1.ServiceSpec, sourc
serviceSpec.Type = v1.ServiceTypeLoadBalancer
}
func (c *Cluster) configureNodePortService(serviceSpec *v1.ServiceSpec, port int32) {
serviceSpec.ExternalTrafficPolicy = v1.ServiceExternalTrafficPolicyType(c.OpConfig.ExternalTrafficPolicy)
serviceSpec.Type = v1.ServiceTypeNodePort
if port != 0 && len(serviceSpec.Ports) > 0 {
serviceSpec.Ports[0].NodePort = port
}
}
func (c *Cluster) generateServiceAnnotations(role PostgresRole, spec *acidv1.PostgresSpec) map[string]string {
annotations := c.getCustomServiceAnnotations(role, spec)
if c.shouldCreateLoadBalancerForService(role, spec) {
nodePort, _ := c.shouldCreateNodePortForService(role, spec)
if !nodePort && c.shouldCreateLoadBalancerForService(role, spec) {
dnsName := c.dnsName(role)
// External DNS name annotation is not customizable
@@ -2420,6 +2464,10 @@ func (c *Cluster) generateLogicalBackupJob() (*batchv1.CronJob, error) {
// configure a cron job
jobTemplateSpec := batchv1.JobTemplateSpec{
ObjectMeta: metav1.ObjectMeta{
Labels: labels,
Annotations: c.annotationsSet(annotations),
},
Spec: jobSpec,
}
@@ -2444,8 +2492,8 @@ func (c *Cluster) generateLogicalBackupJob() (*batchv1.CronJob, error) {
ObjectMeta: metav1.ObjectMeta{
Name: c.getLogicalBackupJobName(),
Namespace: c.Namespace,
Labels: c.labelsSet(true),
Annotations: c.annotationsSet(nil),
Labels: labels,
Annotations: c.annotationsSet(annotations),
OwnerReferences: c.ownerReferences(),
},
Spec: batchv1.CronJobSpec{
+208 -19
View File
@@ -79,7 +79,7 @@ func TestGenerateSpiloJSONConfiguration(t *testing.T) {
PamRoleName: "zalandos",
},
},
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"}],"dcs":{}}}`,
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"}],"dcs":{}}}`,
},
{
subtest: "Patroni configured",
@@ -90,7 +90,7 @@ func TestGenerateSpiloJSONConfiguration(t *testing.T) {
"locale": "en_US.UTF-8",
"data-checksums": "true",
},
PgHba: []string{"hostssl all all 0.0.0.0/0 md5", "host all all 0.0.0.0/0 md5"},
PgHba: []string{"hostssl all all 0.0.0.0/0 scram-sha-256", "host all all 0.0.0.0/0 scram-sha-256"},
TTL: 30,
LoopWait: 10,
RetryTimeout: 10,
@@ -102,7 +102,7 @@ func TestGenerateSpiloJSONConfiguration(t *testing.T) {
FailsafeMode: util.True(),
},
opConfig: &config.Config{},
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin","pg_hba":["hostssl all all 0.0.0.0/0 md5","host all all 0.0.0.0/0 md5"]},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"synchronous_mode":true,"synchronous_mode_strict":true,"synchronous_node_count":1,"slots":{"permanent_logical_1":{"database":"foo","plugin":"pgoutput","type":"logical"}},"failsafe_mode":true}}}`,
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin","pg_hba":["hostssl all all 0.0.0.0/0 scram-sha-256","host all all 0.0.0.0/0 scram-sha-256"]},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"},"data-checksums",{"encoding":"UTF8"},{"locale":"en_US.UTF-8"}],"dcs":{"ttl":30,"loop_wait":10,"retry_timeout":10,"maximum_lag_on_failover":33554432,"synchronous_mode":true,"synchronous_mode_strict":true,"synchronous_node_count":1,"slots":{"permanent_logical_1":{"database":"foo","plugin":"pgoutput","type":"logical"}},"failsafe_mode":true}}}`,
},
{
subtest: "Patroni failsafe_mode configured globally",
@@ -111,7 +111,7 @@ func TestGenerateSpiloJSONConfiguration(t *testing.T) {
opConfig: &config.Config{
EnablePatroniFailsafeMode: util.True(),
},
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"}],"dcs":{"failsafe_mode":true}}}`,
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"}],"dcs":{"failsafe_mode":true}}}`,
},
{
subtest: "Patroni failsafe_mode configured globally, disabled for cluster",
@@ -122,7 +122,7 @@ func TestGenerateSpiloJSONConfiguration(t *testing.T) {
opConfig: &config.Config{
EnablePatroniFailsafeMode: util.True(),
},
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"}],"dcs":{"failsafe_mode":false}}}`,
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"}],"dcs":{"failsafe_mode":false}}}`,
},
{
subtest: "Patroni failsafe_mode disabled globally, configured for cluster",
@@ -133,7 +133,7 @@ func TestGenerateSpiloJSONConfiguration(t *testing.T) {
opConfig: &config.Config{
EnablePatroniFailsafeMode: util.False(),
},
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"md5"},{"auth-local":"trust"}],"dcs":{"failsafe_mode":true}}}`,
result: `{"postgresql":{"bin_dir":"/usr/lib/postgresql/18/bin"},"bootstrap":{"initdb":[{"auth-host":"scram-sha-256"},{"auth-local":"trust"}],"dcs":{"failsafe_mode":true}}}`,
},
}
for _, tt := range tests {
@@ -2971,32 +2971,32 @@ func newLBFakeClient() (k8sutil.KubernetesClient, *fake.Clientset) {
}, clientSet
}
func getServices(serviceType v1.ServiceType, sourceRanges []string, extTrafficPolicy, clusterName string) []v1.ServiceSpec {
func getServices(serviceType v1.ServiceType, sourceRanges []string, extTrafficPolicy, clusterName string, nodePort int32) []v1.ServiceSpec {
return []v1.ServiceSpec{
{
ExternalTrafficPolicy: v1.ServiceExternalTrafficPolicyType(extTrafficPolicy),
LoadBalancerSourceRanges: sourceRanges,
Ports: []v1.ServicePort{{Name: "postgresql", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}}},
Ports: []v1.ServicePort{{Name: "postgresql", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}, NodePort: nodePort}},
Type: serviceType,
},
{
ExternalTrafficPolicy: v1.ServiceExternalTrafficPolicyType(extTrafficPolicy),
LoadBalancerSourceRanges: sourceRanges,
Ports: []v1.ServicePort{{Name: clusterName + "-pooler", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}}},
Ports: []v1.ServicePort{{Name: clusterName + "-pooler", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}, NodePort: nodePort}},
Selector: map[string]string{"connection-pooler": clusterName + "-pooler"},
Type: serviceType,
},
{
ExternalTrafficPolicy: v1.ServiceExternalTrafficPolicyType(extTrafficPolicy),
LoadBalancerSourceRanges: sourceRanges,
Ports: []v1.ServicePort{{Name: "postgresql", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}}},
Ports: []v1.ServicePort{{Name: "postgresql", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}, NodePort: nodePort}},
Selector: map[string]string{"spilo-role": "replica", "application": "spilo", "cluster-name": clusterName},
Type: serviceType,
},
{
ExternalTrafficPolicy: v1.ServiceExternalTrafficPolicyType(extTrafficPolicy),
LoadBalancerSourceRanges: sourceRanges,
Ports: []v1.ServicePort{{Name: clusterName + "-pooler-repl", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}}},
Ports: []v1.ServicePort{{Name: clusterName + "-pooler-repl", Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}, NodePort: nodePort}},
Selector: map[string]string{"connection-pooler": clusterName + "-pooler-repl"},
Type: serviceType,
},
@@ -3064,7 +3064,7 @@ func TestEnableLoadBalancers(t *testing.T) {
},
},
},
expectedServices: getServices(v1.ServiceTypeClusterIP, nil, "", clusterName),
expectedServices: getServices(v1.ServiceTypeClusterIP, nil, "", clusterName, 0),
},
{
subTest: "LBs enabled in manifest, disabled in config",
@@ -3111,7 +3111,7 @@ func TestEnableLoadBalancers(t *testing.T) {
},
},
},
expectedServices: getServices(v1.ServiceTypeLoadBalancer, sourceRanges, extTrafficPolicy, clusterName),
expectedServices: getServices(v1.ServiceTypeLoadBalancer, sourceRanges, extTrafficPolicy, clusterName, 0),
},
}
@@ -3143,6 +3143,195 @@ func TestEnableLoadBalancers(t *testing.T) {
}
}
func TestEnableNodePorts(t *testing.T) {
clusterName := "acid-test-cluster"
namespace := "default"
clusterNameLabel := "cluster-name"
roleLabel := "spilo-role"
roles := []PostgresRole{Master, Replica}
extTrafficPolicy := "Cluster"
port := int32(1337)
tests := []struct {
subTest string
config config.Config
pgSpec acidv1.Postgresql
expectedServices []v1.ServiceSpec
}{
{
subTest: "NodePorts enabled in config, disabled in manifest",
config: config.Config{
ConnectionPooler: config.ConnectionPooler{
ConnectionPoolerDefaultCPURequest: "100m",
ConnectionPoolerDefaultCPULimit: "100m",
ConnectionPoolerDefaultMemoryRequest: "100Mi",
ConnectionPoolerDefaultMemoryLimit: "100Mi",
NumberOfInstances: k8sutil.Int32ToPointer(1),
},
EnableMasterNodePort: true,
EnableMasterPoolerNodePort: true,
EnableReplicaNodePort: true,
EnableReplicaPoolerNodePort: true,
ExternalTrafficPolicy: extTrafficPolicy,
Resources: config.Resources{
ClusterLabels: map[string]string{"application": "spilo"},
ClusterNameLabel: clusterNameLabel,
PodRoleLabel: roleLabel,
},
},
pgSpec: acidv1.Postgresql{
ObjectMeta: metav1.ObjectMeta{
Name: clusterName,
Namespace: namespace,
},
Spec: acidv1.PostgresSpec{
EnableConnectionPooler: util.True(),
EnableReplicaConnectionPooler: util.True(),
EnableMasterNodePort: util.False(),
EnableMasterPoolerNodePort: util.False(),
EnableReplicaNodePort: util.False(),
EnableReplicaPoolerNodePort: util.False(),
NumberOfInstances: 1,
Resources: &acidv1.Resources{
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("10")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("10")},
},
TeamID: "acid",
Volume: acidv1.Volume{
Size: "1G",
},
},
},
expectedServices: getServices(v1.ServiceTypeClusterIP, nil, "", clusterName, 0),
},
{
subTest: "NodePorts configured in manifest, disabled in config",
config: config.Config{
ConnectionPooler: config.ConnectionPooler{
ConnectionPoolerDefaultCPURequest: "100m",
ConnectionPoolerDefaultCPULimit: "100m",
ConnectionPoolerDefaultMemoryRequest: "100Mi",
ConnectionPoolerDefaultMemoryLimit: "100Mi",
NumberOfInstances: k8sutil.Int32ToPointer(1),
},
EnableMasterNodePort: false,
EnableMasterPoolerNodePort: false,
EnableReplicaNodePort: false,
EnableReplicaPoolerNodePort: false,
ExternalTrafficPolicy: extTrafficPolicy,
Resources: config.Resources{
ClusterLabels: map[string]string{"application": "spilo"},
ClusterNameLabel: clusterNameLabel,
PodRoleLabel: roleLabel,
},
},
pgSpec: acidv1.Postgresql{
ObjectMeta: metav1.ObjectMeta{
Name: clusterName,
Namespace: namespace,
},
Spec: acidv1.PostgresSpec{
EnableConnectionPooler: util.True(),
EnableReplicaConnectionPooler: util.True(),
EnableMasterNodePort: util.True(),
EnableMasterPoolerNodePort: util.True(),
EnableReplicaNodePort: util.True(),
EnableReplicaPoolerNodePort: util.True(),
NumberOfInstances: 1,
Resources: &acidv1.Resources{
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("10")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("10")},
},
TeamID: "acid",
Volume: acidv1.Volume{
Size: "1G",
},
},
},
expectedServices: getServices(v1.ServiceTypeNodePort, nil, extTrafficPolicy, clusterName, 0),
},
{
subTest: "NodePorts configured in manifest, disabled in config, custom port specified",
config: config.Config{
ConnectionPooler: config.ConnectionPooler{
ConnectionPoolerDefaultCPURequest: "100m",
ConnectionPoolerDefaultCPULimit: "100m",
ConnectionPoolerDefaultMemoryRequest: "100Mi",
ConnectionPoolerDefaultMemoryLimit: "100Mi",
NumberOfInstances: k8sutil.Int32ToPointer(1),
},
EnableMasterNodePort: false,
EnableMasterPoolerNodePort: false,
EnableReplicaNodePort: false,
EnableReplicaPoolerNodePort: false,
ExternalTrafficPolicy: extTrafficPolicy,
Resources: config.Resources{
ClusterLabels: map[string]string{"application": "spilo"},
ClusterNameLabel: clusterNameLabel,
PodRoleLabel: roleLabel,
},
},
pgSpec: acidv1.Postgresql{
ObjectMeta: metav1.ObjectMeta{
Name: clusterName,
Namespace: namespace,
},
Spec: acidv1.PostgresSpec{
EnableConnectionPooler: util.True(),
EnableReplicaConnectionPooler: util.True(),
EnableMasterNodePort: util.True(),
MasterNodePort: &port,
EnableMasterPoolerNodePort: util.True(),
MasterPoolerNodePort: &port,
EnableReplicaNodePort: util.True(),
ReplicaNodePort: &port,
EnableReplicaPoolerNodePort: util.True(),
ReplicaPoolerNodePort: &port,
NumberOfInstances: 1,
Resources: &acidv1.Resources{
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("10")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("10")},
},
TeamID: "acid",
Volume: acidv1.Volume{
Size: "1G",
},
},
},
expectedServices: getServices(v1.ServiceTypeNodePort, nil, extTrafficPolicy, clusterName, port),
},
}
for _, tt := range tests {
client, _ := newLBFakeClient()
var cluster = New(
Config{
OpConfig: tt.config,
}, client, tt.pgSpec, logger, eventRecorder)
cluster.Name = clusterName
cluster.Namespace = namespace
cluster.ConnectionPooler = map[PostgresRole]*ConnectionPoolerObjects{}
generatedServices := make([]v1.ServiceSpec, 0)
for _, role := range roles {
cluster.syncService(role)
cluster.ConnectionPooler[role] = &ConnectionPoolerObjects{
Name: cluster.connectionPoolerName(role),
ClusterName: cluster.Name,
Namespace: cluster.Namespace,
Role: role,
}
cluster.syncConnectionPoolerWorker(&tt.pgSpec, &tt.pgSpec, role)
generatedServices = append(generatedServices, cluster.Services[role].Spec)
generatedServices = append(generatedServices, cluster.ConnectionPooler[role].Service.Spec)
}
if !reflect.DeepEqual(tt.expectedServices, generatedServices) {
t.Errorf("%s %s: expected %#v but got %#v", t.Name(), tt.subTest, tt.expectedServices, generatedServices)
}
}
}
func TestGenerateResourceRequirements(t *testing.T) {
client, _ := newFakeK8sTestClient()
clusterName := "acid-test-cluster"
@@ -3875,7 +4064,7 @@ func TestGenerateLogicalBackupJob(t *testing.T) {
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("100m"), Memory: k8sutil.StringToPointer("100Mi")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("500Mi")},
},
expectedLabel: map[string]string{configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedLabel: map[string]string{"application": "spilo-logical-backup", configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedAnnotation: nil,
},
{
@@ -3900,7 +4089,7 @@ func TestGenerateLogicalBackupJob(t *testing.T) {
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("10m"), Memory: k8sutil.StringToPointer("50Mi")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("300m"), Memory: k8sutil.StringToPointer("300Mi")},
},
expectedLabel: map[string]string{configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedLabel: map[string]string{"application": "spilo-logical-backup", configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedAnnotation: nil,
},
{
@@ -3923,7 +4112,7 @@ func TestGenerateLogicalBackupJob(t *testing.T) {
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("50m"), Memory: k8sutil.StringToPointer("100Mi")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("250m"), Memory: k8sutil.StringToPointer("500Mi")},
},
expectedLabel: map[string]string{configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedLabel: map[string]string{"application": "spilo-logical-backup", configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedAnnotation: nil,
},
{
@@ -3946,7 +4135,7 @@ func TestGenerateLogicalBackupJob(t *testing.T) {
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("100m"), Memory: k8sutil.StringToPointer("200Mi")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("200Mi")},
},
expectedLabel: map[string]string{configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedLabel: map[string]string{"application": "spilo-logical-backup", configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedAnnotation: nil,
},
{
@@ -3968,7 +4157,7 @@ func TestGenerateLogicalBackupJob(t *testing.T) {
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("100m"), Memory: k8sutil.StringToPointer("100Mi")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("500Mi")},
},
expectedLabel: map[string]string{"labelKey": "labelValue", "cluster-name": clusterName, "team": teamId},
expectedLabel: map[string]string{"application": "spilo-logical-backup", "labelKey": "labelValue", "cluster-name": clusterName, "team": teamId},
expectedAnnotation: nil,
},
{
@@ -3990,7 +4179,7 @@ func TestGenerateLogicalBackupJob(t *testing.T) {
ResourceRequests: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("100m"), Memory: k8sutil.StringToPointer("100Mi")},
ResourceLimits: acidv1.ResourceDescription{CPU: k8sutil.StringToPointer("1"), Memory: k8sutil.StringToPointer("500Mi")},
},
expectedLabel: map[string]string{configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedLabel: map[string]string{"application": "spilo-logical-backup", configResources.ClusterNameLabel: clusterName, "team": teamId},
expectedAnnotation: map[string]string{"annotationKey": "annotationValue"},
},
}
+7 -2
View File
@@ -324,9 +324,14 @@ func (c *Cluster) updateService(role PostgresRole, oldService *v1.Service, newSe
// patch does not work because of LoadBalancerSourceRanges field (even if set to nil)
oldServiceType := oldService.Spec.Type
newServiceType := newService.Spec.Type
if newServiceType == "ClusterIP" && newServiceType != oldServiceType {
if newServiceType != oldServiceType && oldServiceType == v1.ServiceTypeLoadBalancer {
// Kubernetes rejects updates that change type away from LoadBalancer while
// loadBalancerSourceRanges is still set; clear it before updating
newService.Spec.LoadBalancerSourceRanges = nil
newService.ResourceVersion = oldService.ResourceVersion
newService.Spec.ClusterIP = oldService.Spec.ClusterIP
if newServiceType == v1.ServiceTypeClusterIP {
newService.Spec.ClusterIP = oldService.Spec.ClusterIP
}
}
svc, err = c.KubeClient.Services(serviceName.Namespace).Update(context.TODO(), newService, metav1.UpdateOptions{})
if err != nil {
+10
View File
@@ -1769,6 +1769,16 @@ func (c *Cluster) syncLogicalBackupJob() error {
}
c.logger.Info("the logical backup job is synced")
}
if !reflect.DeepEqual(job.Labels, desiredJob.Labels) {
patchData, err := metaLabelsPatch(desiredJob.Labels)
if err != nil {
return fmt.Errorf("could not form patch for the logical backup job %q labels: %v", jobName, err)
}
_, err = c.KubeClient.CronJobs(c.Namespace).Patch(context.TODO(), jobName, types.MergePatchType, []byte(patchData), metav1.PatchOptions{})
if err != nil {
return fmt.Errorf("could not patch labels of the logical backup job %q: %v", jobName, err)
}
}
if changed, _ := c.compareAnnotations(job.Annotations, desiredJob.Annotations, nil); changed {
patchData, err := metaAnnotationsPatch(desiredJob.Annotations)
if err != nil {
+8
View File
@@ -167,6 +167,14 @@ func metaAnnotationsPatch(annotations map[string]string) ([]byte, error) {
}{&meta})
}
func metaLabelsPatch(labels map[string]string) ([]byte, error) {
var meta metav1.ObjectMeta
meta.Labels = labels
return json.Marshal(struct {
ObjMeta interface{} `json:"metadata"`
}{&meta})
}
func (c *Cluster) logPDBChanges(old, new *policyv1.PodDisruptionBudget, isUpdate bool, reason string) {
if isUpdate {
c.logger.Infof("pod disruption budget %q has been changed", util.NameFromMeta(old.ObjectMeta))
+4 -19
View File
@@ -26,6 +26,7 @@ import (
rbacv1 "k8s.io/api/rbac/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
informers_core_v1 "k8s.io/client-go/informers/core/v1"
"k8s.io/client-go/kubernetes/scheme"
typedcorev1 "k8s.io/client-go/kubernetes/typed/core/v1"
"k8s.io/client-go/tools/cache"
@@ -401,16 +402,8 @@ func (c *Controller) initSharedInformers() {
}
// Pods
podLw := &cache.ListWatch{
ListFunc: c.podListFunc,
WatchFunc: c.podWatchFunc,
}
c.podInformer = cache.NewSharedIndexInformer(
podLw,
&v1.Pod{},
constants.QueueResyncPeriodPod,
cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc})
c.podInformer = informers_core_v1.NewPodInformer(c.KubeClient.Clientset,
c.opConfig.WatchedNamespace, constants.QueueResyncPeriodPod, cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc})
c.podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: c.podAdd,
@@ -419,15 +412,7 @@ func (c *Controller) initSharedInformers() {
})
// Kubernetes Nodes
nodeLw := &cache.ListWatch{
ListFunc: c.nodeListFunc,
WatchFunc: c.nodeWatchFunc,
}
c.nodesInformer = cache.NewSharedIndexInformer(
nodeLw,
&v1.Node{},
constants.QueueResyncPeriodNode,
c.nodesInformer = informers_core_v1.NewNodeInformer(c.KubeClient.Clientset, constants.QueueResyncPeriodNode,
cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc})
c.nodesInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
-22
View File
@@ -9,33 +9,11 @@ import (
v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/watch"
"github.com/zalando/postgres-operator/pkg/cluster"
"github.com/zalando/postgres-operator/pkg/util"
)
func (c *Controller) nodeListFunc(options metav1.ListOptions) (runtime.Object, error) {
opts := metav1.ListOptions{
Watch: options.Watch,
ResourceVersion: options.ResourceVersion,
TimeoutSeconds: options.TimeoutSeconds,
}
return c.KubeClient.Nodes().List(context.TODO(), opts)
}
func (c *Controller) nodeWatchFunc(options metav1.ListOptions) (watch.Interface, error) {
opts := metav1.ListOptions{
Watch: options.Watch,
ResourceVersion: options.ResourceVersion,
TimeoutSeconds: options.TimeoutSeconds,
}
return c.KubeClient.Nodes().Watch(context.TODO(), opts)
}
func (c *Controller) nodeAdd(obj interface{}) {
node, ok := obj.(*v1.Node)
if !ok {
+6 -2
View File
@@ -171,6 +171,10 @@ func (c *Controller) importConfigurationFromCRD(fromCRD *acidv1.OperatorConfigur
result.EnableMasterPoolerLoadBalancer = fromCRD.LoadBalancer.EnableMasterPoolerLoadBalancer
result.EnableReplicaLoadBalancer = fromCRD.LoadBalancer.EnableReplicaLoadBalancer
result.EnableReplicaPoolerLoadBalancer = fromCRD.LoadBalancer.EnableReplicaPoolerLoadBalancer
result.EnableMasterNodePort = fromCRD.LoadBalancer.EnableMasterNodePort
result.EnableMasterPoolerNodePort = fromCRD.LoadBalancer.EnableMasterPoolerNodePort
result.EnableReplicaNodePort = fromCRD.LoadBalancer.EnableReplicaNodePort
result.EnableReplicaPoolerNodePort = fromCRD.LoadBalancer.EnableReplicaPoolerNodePort
result.CustomServiceAnnotations = fromCRD.LoadBalancer.CustomServiceAnnotations
result.MasterDNSNameFormat = fromCRD.LoadBalancer.MasterDNSNameFormat
result.MasterLegacyDNSNameFormat = fromCRD.LoadBalancer.MasterLegacyDNSNameFormat
@@ -218,8 +222,8 @@ func (c *Controller) importConfigurationFromCRD(fromCRD *acidv1.OperatorConfigur
result.LogicalBackupTTLSecondsAfterFinished = fromCRD.LogicalBackup.TTLSecondsAfterFinished
// debug config
result.DebugLogging = fromCRD.OperatorDebug.DebugLogging
result.EnableDBAccess = fromCRD.OperatorDebug.EnableDBAccess
result.DebugLogging = *util.CoalesceBool(fromCRD.OperatorDebug.DebugLogging, util.True())
result.EnableDBAccess = *util.CoalesceBool(fromCRD.OperatorDebug.EnableDBAccess, util.True())
// Teams API config
result.EnableTeamsAPI = fromCRD.TeamsAPI.EnableTeamsAPI
-25
View File
@@ -1,12 +1,7 @@
package controller
import (
"context"
v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/watch"
"github.com/zalando/postgres-operator/pkg/cluster"
"github.com/zalando/postgres-operator/pkg/spec"
@@ -14,26 +9,6 @@ import (
"k8s.io/apimachinery/pkg/types"
)
func (c *Controller) podListFunc(options metav1.ListOptions) (runtime.Object, error) {
opts := metav1.ListOptions{
Watch: options.Watch,
ResourceVersion: options.ResourceVersion,
TimeoutSeconds: options.TimeoutSeconds,
}
return c.KubeClient.Pods(c.opConfig.WatchedNamespace).List(context.TODO(), opts)
}
func (c *Controller) podWatchFunc(options metav1.ListOptions) (watch.Interface, error) {
opts := metav1.ListOptions{
Watch: options.Watch,
ResourceVersion: options.ResourceVersion,
TimeoutSeconds: options.TimeoutSeconds,
}
return c.KubeClient.Pods(c.opConfig.WatchedNamespace).Watch(context.TODO(), opts)
}
func (c *Controller) dispatchPodEvent(clusterName spec.NamespacedName, event cluster.PodEvent) {
c.clustersMu.RLock()
cluster, ok := c.clusters[clusterName]
+4
View File
@@ -216,6 +216,10 @@ type Config struct {
EnableMasterPoolerLoadBalancer bool `name:"enable_master_pooler_load_balancer" default:"false"`
EnableReplicaLoadBalancer bool `name:"enable_replica_load_balancer" default:"false"`
EnableReplicaPoolerLoadBalancer bool `name:"enable_replica_pooler_load_balancer" default:"false"`
EnableMasterNodePort bool `name:"enable_master_node_port" default:"false"`
EnableMasterPoolerNodePort bool `name:"enable_master_pooler_node_port" default:"false"`
EnableReplicaNodePort bool `name:"enable_replica_node_port" default:"false"`
EnableReplicaPoolerNodePort bool `name:"enable_replica_pooler_node_port" default:"false"`
CustomServiceAnnotations map[string]string `name:"custom_service_annotations"`
CustomPodAnnotations map[string]string `name:"custom_pod_annotations"`
EnablePodAntiAffinity bool `name:"enable_pod_antiaffinity" default:"false"`
+2
View File
@@ -67,6 +67,7 @@ type KubernetesClient struct {
zalandov1.FabricEventStreamsGetter
RESTClient rest.Interface
Clientset *kubernetes.Clientset
AcidV1ClientSet *zalandoclient.Clientset
Zalandov1ClientSet *zalandoclient.Clientset
}
@@ -148,6 +149,7 @@ func NewFromConfig(cfg *rest.Config) (KubernetesClient, error) {
return kubeClient, fmt.Errorf("could not get clientset: %v", err)
}
kubeClient.Clientset = client
kubeClient.PodsGetter = client.CoreV1()
kubeClient.ServicesGetter = client.CoreV1()
kubeClient.EndpointsGetter = client.CoreV1()
+1 -1
View File
@@ -87,7 +87,7 @@ func NewEncryptor(encryption string) *Encryptor {
}
hasher, ok := m[encryption]
if !ok {
hasher = e.PGUserPasswordMD5
hasher = e.PGUserPasswordScramSHA256
}
e.encrypt = hasher
return &e