mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-06 18:30:55 +02:00
merge master
This commit is contained in:
+14
-12
@@ -4,6 +4,7 @@ package cluster
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"reflect"
|
||||
"regexp"
|
||||
@@ -19,8 +20,6 @@ import (
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
|
||||
"encoding/json"
|
||||
|
||||
acidv1 "github.com/zalando/postgres-operator/pkg/apis/acid.zalan.do/v1"
|
||||
"github.com/zalando/postgres-operator/pkg/spec"
|
||||
"github.com/zalando/postgres-operator/pkg/util"
|
||||
@@ -150,21 +149,24 @@ func (c *Cluster) setProcessName(procName string, args ...interface{}) {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Cluster) setStatus(status acidv1.PostgresStatus) {
|
||||
// TODO: eventually switch to updateStatus() for kubernetes 1.11 and above
|
||||
var (
|
||||
err error
|
||||
b []byte
|
||||
)
|
||||
if b, err = json.Marshal(status); err != nil {
|
||||
// SetStatus of Postgres cluster
|
||||
// TODO: eventually switch to updateStatus() for kubernetes 1.11 and above
|
||||
func (c *Cluster) setStatus(status string) {
|
||||
var pgStatus acidv1.PostgresStatus
|
||||
pgStatus.PostgresClusterStatus = status
|
||||
|
||||
patch, err := json.Marshal(struct {
|
||||
PgStatus interface{} `json:"status"`
|
||||
}{&pgStatus})
|
||||
|
||||
if err != nil {
|
||||
c.logger.Errorf("could not marshal status: %v", err)
|
||||
}
|
||||
|
||||
patch := []byte(fmt.Sprintf(`{"status": %s}`, string(b)))
|
||||
// we cannot do a full scale update here without fetching the previous manifest (as the resourceVersion may differ),
|
||||
// however, we could do patch without it. In the future, once /status subresource is there (starting Kubernets 1.11)
|
||||
// we should take advantage of it.
|
||||
newspec, err := c.KubeClient.AcidV1ClientSet.AcidV1().Postgresqls(c.clusterNamespace()).Patch(c.Name, types.MergePatchType, patch)
|
||||
newspec, err := c.KubeClient.AcidV1ClientSet.AcidV1().Postgresqls(c.clusterNamespace()).Patch(c.Name, types.MergePatchType, patch, "status")
|
||||
if err != nil {
|
||||
c.logger.Errorf("could not update status: %v", err)
|
||||
}
|
||||
@@ -173,7 +175,7 @@ func (c *Cluster) setStatus(status acidv1.PostgresStatus) {
|
||||
}
|
||||
|
||||
func (c *Cluster) isNewCluster() bool {
|
||||
return c.Status == acidv1.ClusterStatusCreating
|
||||
return c.Status.Creating()
|
||||
}
|
||||
|
||||
// initUsers populates c.systemUsers and c.pgUsers maps.
|
||||
|
||||
@@ -20,10 +20,20 @@ const (
|
||||
)
|
||||
|
||||
var logger = logrus.New().WithField("test", "cluster")
|
||||
var cl = New(Config{OpConfig: config.Config{ProtectedRoles: []string{"admin"},
|
||||
Auth: config.Auth{SuperUsername: superUserName,
|
||||
ReplicationUsername: replicationUserName}}},
|
||||
k8sutil.KubernetesClient{}, acidv1.Postgresql{}, logger)
|
||||
var cl = New(
|
||||
Config{
|
||||
OpConfig: config.Config{
|
||||
ProtectedRoles: []string{"admin"},
|
||||
Auth: config.Auth{
|
||||
SuperUsername: superUserName,
|
||||
ReplicationUsername: replicationUserName,
|
||||
},
|
||||
},
|
||||
},
|
||||
k8sutil.NewMockKubernetesClient(),
|
||||
acidv1.Postgresql{},
|
||||
logger,
|
||||
)
|
||||
|
||||
func TestInitRobotUsers(t *testing.T) {
|
||||
testName := "TestInitRobotUsers"
|
||||
|
||||
+31
-2
@@ -1016,6 +1016,7 @@ func generatePersistentVolumeClaimTemplate(volumeSize, volumeStorageClass string
|
||||
return nil, fmt.Errorf("could not parse volume size: %v", err)
|
||||
}
|
||||
|
||||
volumeMode := v1.PersistentVolumeFilesystem
|
||||
volumeClaim := &v1.PersistentVolumeClaim{
|
||||
ObjectMeta: metadata,
|
||||
Spec: v1.PersistentVolumeClaimSpec{
|
||||
@@ -1026,6 +1027,7 @@ func generatePersistentVolumeClaimTemplate(volumeSize, volumeStorageClass string
|
||||
},
|
||||
},
|
||||
StorageClassName: storageClassName,
|
||||
VolumeMode: &volumeMode,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -1218,10 +1220,37 @@ func (c *Cluster) generateCloneEnvironment(description *acidv1.CloneDescription)
|
||||
})
|
||||
} else {
|
||||
// cloning with S3, find out the bucket to clone
|
||||
msg := "Clone from S3 bucket"
|
||||
c.logger.Info(msg, description.S3WalPath)
|
||||
|
||||
if description.S3WalPath == "" {
|
||||
msg := "Figure out which S3 bucket to use from env"
|
||||
c.logger.Info(msg, description.S3WalPath)
|
||||
|
||||
envs := []v1.EnvVar{
|
||||
v1.EnvVar{
|
||||
Name: "CLONE_WAL_S3_BUCKET",
|
||||
Value: c.OpConfig.WALES3Bucket,
|
||||
},
|
||||
v1.EnvVar{
|
||||
Name: "CLONE_WAL_BUCKET_SCOPE_SUFFIX",
|
||||
Value: getBucketScopeSuffix(description.UID),
|
||||
},
|
||||
}
|
||||
|
||||
result = append(result, envs...)
|
||||
} else {
|
||||
msg := "Use custom parsed S3WalPath %s from the manifest"
|
||||
c.logger.Warningf(msg, description.S3WalPath)
|
||||
|
||||
result = append(result, v1.EnvVar{
|
||||
Name: "CLONE_WALE_S3_PREFIX",
|
||||
Value: description.S3WalPath,
|
||||
})
|
||||
}
|
||||
|
||||
result = append(result, v1.EnvVar{Name: "CLONE_METHOD", Value: "CLONE_WITH_WALE"})
|
||||
result = append(result, v1.EnvVar{Name: "CLONE_WAL_S3_BUCKET", Value: c.OpConfig.WALES3Bucket})
|
||||
result = append(result, v1.EnvVar{Name: "CLONE_TARGET_TIME", Value: description.EndTimestamp})
|
||||
result = append(result, v1.EnvVar{Name: "CLONE_WAL_BUCKET_SCOPE_SUFFIX", Value: getBucketScopeSuffix(description.UID)})
|
||||
result = append(result, v1.EnvVar{Name: "CLONE_WAL_BUCKET_SCOPE_PREFIX", Value: ""})
|
||||
}
|
||||
|
||||
|
||||
@@ -129,3 +129,82 @@ func TestShmVolume(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCloneEnv(t *testing.T) {
|
||||
testName := "TestCloneEnv"
|
||||
tests := []struct {
|
||||
subTest string
|
||||
cloneOpts *acidv1.CloneDescription
|
||||
env v1.EnvVar
|
||||
envPos int
|
||||
}{
|
||||
{
|
||||
subTest: "custom s3 path",
|
||||
cloneOpts: &acidv1.CloneDescription{
|
||||
ClusterName: "test-cluster",
|
||||
S3WalPath: "s3://some/path/",
|
||||
EndTimestamp: "somewhen",
|
||||
},
|
||||
env: v1.EnvVar{
|
||||
Name: "CLONE_WALE_S3_PREFIX",
|
||||
Value: "s3://some/path/",
|
||||
},
|
||||
envPos: 1,
|
||||
},
|
||||
{
|
||||
subTest: "generated s3 path, bucket",
|
||||
cloneOpts: &acidv1.CloneDescription{
|
||||
ClusterName: "test-cluster",
|
||||
EndTimestamp: "somewhen",
|
||||
UID: "0000",
|
||||
},
|
||||
env: v1.EnvVar{
|
||||
Name: "CLONE_WAL_S3_BUCKET",
|
||||
Value: "wale-bucket",
|
||||
},
|
||||
envPos: 1,
|
||||
},
|
||||
{
|
||||
subTest: "generated s3 path, target time",
|
||||
cloneOpts: &acidv1.CloneDescription{
|
||||
ClusterName: "test-cluster",
|
||||
EndTimestamp: "somewhen",
|
||||
UID: "0000",
|
||||
},
|
||||
env: v1.EnvVar{
|
||||
Name: "CLONE_TARGET_TIME",
|
||||
Value: "somewhen",
|
||||
},
|
||||
envPos: 4,
|
||||
},
|
||||
}
|
||||
|
||||
var cluster = New(
|
||||
Config{
|
||||
OpConfig: config.Config{
|
||||
WALES3Bucket: "wale-bucket",
|
||||
ProtectedRoles: []string{"admin"},
|
||||
Auth: config.Auth{
|
||||
SuperUsername: superUserName,
|
||||
ReplicationUsername: replicationUserName,
|
||||
},
|
||||
},
|
||||
}, k8sutil.KubernetesClient{}, acidv1.Postgresql{}, logger)
|
||||
|
||||
for _, tt := range tests {
|
||||
envs := cluster.generateCloneEnvironment(tt.cloneOpts)
|
||||
|
||||
env := envs[tt.envPos]
|
||||
|
||||
if env.Name != tt.env.Name {
|
||||
t.Errorf("%s %s: Expected env name %s, have %s instead",
|
||||
testName, tt.subTest, tt.env.Name, env.Name)
|
||||
}
|
||||
|
||||
if env.Value != tt.env.Value {
|
||||
t.Errorf("%s %s: Expected env value %s, have %s instead",
|
||||
testName, tt.subTest, tt.env.Value, env.Value)
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -29,7 +29,7 @@ func (c *Cluster) Sync(newSpec *acidv1.Postgresql) error {
|
||||
if err != nil {
|
||||
c.logger.Warningf("error while syncing cluster state: %v", err)
|
||||
c.setStatus(acidv1.ClusterStatusSyncFailed)
|
||||
} else if c.Status != acidv1.ClusterStatusRunning {
|
||||
} else if !c.Status.Running() {
|
||||
c.setStatus(acidv1.ClusterStatusRunning)
|
||||
}
|
||||
}()
|
||||
|
||||
Reference in New Issue
Block a user