mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-03 06:51:53 +02:00
Get config from environment variables;
ignore pg major version change; get rid of resources package;
This commit is contained in:
+11
-14
@@ -21,9 +21,9 @@ import (
|
||||
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/spec"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/config"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/constants"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/k8sutil"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/resources"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/teams"
|
||||
)
|
||||
|
||||
@@ -37,6 +37,7 @@ type Config struct {
|
||||
RestClient *rest.RESTClient
|
||||
EtcdClient etcdclient.KeysAPI
|
||||
TeamsAPIClient *teams.TeamsAPI
|
||||
OpConfig *config.Config
|
||||
}
|
||||
|
||||
type KubeResources struct {
|
||||
@@ -51,10 +52,8 @@ type KubeResources struct {
|
||||
type Cluster struct {
|
||||
KubeResources
|
||||
spec.Postgresql
|
||||
config Config
|
||||
Config
|
||||
logger *logrus.Entry
|
||||
etcdHost string
|
||||
dockerImage string
|
||||
pgUsers map[string]spec.PgUser
|
||||
podEvents chan spec.PodEvent
|
||||
podSubscribers map[spec.PodName]chan spec.PodEvent
|
||||
@@ -67,11 +66,9 @@ func New(cfg Config, pgSpec spec.Postgresql) *Cluster {
|
||||
kubeResources := KubeResources{Secrets: make(map[types.UID]*v1.Secret)}
|
||||
|
||||
cluster := &Cluster{
|
||||
config: cfg,
|
||||
Config: cfg,
|
||||
Postgresql: pgSpec,
|
||||
logger: lg,
|
||||
etcdHost: constants.EtcdHost,
|
||||
dockerImage: constants.SpiloImage,
|
||||
pgUsers: make(map[string]spec.PgUser),
|
||||
podEvents: make(chan spec.PodEvent),
|
||||
podSubscribers: make(map[spec.PodName]chan spec.PodEvent),
|
||||
@@ -106,7 +103,7 @@ func (c *Cluster) SetStatus(status spec.PostgresStatus) {
|
||||
}
|
||||
request := []byte(fmt.Sprintf(`{"status": %s}`, string(b))) //TODO: Look into/wait for k8s go client methods
|
||||
|
||||
_, err = c.config.RestClient.Patch(api.MergePatchType).
|
||||
_, err = c.RestClient.Patch(api.MergePatchType).
|
||||
RequestURI(c.Metadata.GetSelfLink()).
|
||||
Body(request).
|
||||
DoRaw()
|
||||
@@ -237,7 +234,7 @@ func (c *Cluster) Update(newSpec *spec.Postgresql) error {
|
||||
c.logger.Infof("Cluster update from version %s to %s",
|
||||
c.Metadata.ResourceVersion, newSpec.Metadata.ResourceVersion)
|
||||
|
||||
newService := resources.Service(c.ClusterName(), c.TeamName(), newSpec.Spec.AllowedSourceRanges)
|
||||
newService := c.genService(newSpec.Spec.AllowedSourceRanges)
|
||||
if !c.sameServiceWith(newService) {
|
||||
c.logger.Infof("LoadBalancer configuration has changed for Service '%s': %+v -> %+v",
|
||||
util.NameFromMeta(c.Service.ObjectMeta),
|
||||
@@ -255,7 +252,7 @@ func (c *Cluster) Update(newSpec *spec.Postgresql) error {
|
||||
//TODO: update PVC
|
||||
}
|
||||
|
||||
newStatefulSet := genStatefulSet(c.ClusterName(), newSpec.Spec, c.etcdHost, c.dockerImage)
|
||||
newStatefulSet := c.genStatefulSet(newSpec.Spec)
|
||||
sameSS, rollingUpdate := c.compareStatefulSetWith(newStatefulSet)
|
||||
|
||||
if !sameSS {
|
||||
@@ -340,13 +337,13 @@ func (c *Cluster) ReceivePodEvent(event spec.PodEvent) {
|
||||
}
|
||||
|
||||
func (c *Cluster) initSystemUsers() {
|
||||
c.pgUsers[constants.SuperuserName] = spec.PgUser{
|
||||
Name: constants.SuperuserName,
|
||||
c.pgUsers[c.OpConfig.SuperUsername] = spec.PgUser{
|
||||
Name: c.OpConfig.SuperUsername,
|
||||
Password: util.RandomPassword(constants.PasswordLength),
|
||||
}
|
||||
|
||||
c.pgUsers[constants.ReplicationUsername] = spec.PgUser{
|
||||
Name: constants.ReplicationUsername,
|
||||
c.pgUsers[c.OpConfig.ReplicationUsername] = spec.PgUser{
|
||||
Name: c.OpConfig.ReplicationUsername,
|
||||
Password: util.RandomPassword(constants.PasswordLength),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,260 @@
|
||||
package cluster
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"k8s.io/client-go/pkg/api/resource"
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
"k8s.io/client-go/pkg/apis/apps/v1beta1"
|
||||
"k8s.io/client-go/pkg/util/intstr"
|
||||
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/spec"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/constants"
|
||||
)
|
||||
|
||||
func resourceList(resources spec.Resources) *v1.ResourceList {
|
||||
resourceList := v1.ResourceList{}
|
||||
if resources.Cpu != "" {
|
||||
resourceList[v1.ResourceCPU] = resource.MustParse(resources.Cpu)
|
||||
}
|
||||
|
||||
if resources.Memory != "" {
|
||||
resourceList[v1.ResourceMemory] = resource.MustParse(resources.Memory)
|
||||
}
|
||||
|
||||
return &resourceList
|
||||
}
|
||||
|
||||
func (c *Cluster) genPodTemplate(resourceList *v1.ResourceList, pgVersion string) *v1.PodTemplateSpec {
|
||||
envVars := []v1.EnvVar{
|
||||
{
|
||||
Name: "SCOPE",
|
||||
Value: c.Metadata.Name,
|
||||
},
|
||||
{
|
||||
Name: "PGROOT",
|
||||
Value: "/home/postgres/pgdata/pgroot",
|
||||
},
|
||||
{
|
||||
Name: "ETCD_HOST",
|
||||
Value: c.OpConfig.EtcdHost,
|
||||
},
|
||||
{
|
||||
Name: "POD_IP",
|
||||
ValueFrom: &v1.EnvVarSource{
|
||||
FieldRef: &v1.ObjectFieldSelector{
|
||||
APIVersion: "v1",
|
||||
FieldPath: "status.podIP",
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "POD_NAMESPACE",
|
||||
ValueFrom: &v1.EnvVarSource{
|
||||
FieldRef: &v1.ObjectFieldSelector{
|
||||
APIVersion: "v1",
|
||||
FieldPath: "metadata.namespace",
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "PGPASSWORD_SUPERUSER",
|
||||
ValueFrom: &v1.EnvVarSource{
|
||||
SecretKeyRef: &v1.SecretKeySelector{
|
||||
LocalObjectReference: v1.LocalObjectReference{
|
||||
Name: c.credentialSecretName(c.OpConfig.SuperUsername),
|
||||
},
|
||||
Key: "password",
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "PGPASSWORD_STANDBY",
|
||||
ValueFrom: &v1.EnvVarSource{
|
||||
SecretKeyRef: &v1.SecretKeySelector{
|
||||
LocalObjectReference: v1.LocalObjectReference{
|
||||
Name: c.credentialSecretName(c.OpConfig.ReplicationUsername),
|
||||
},
|
||||
Key: "password",
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "PAM_OAUTH2",
|
||||
Value: c.OpConfig.PamConfiguration,
|
||||
},
|
||||
{
|
||||
Name: "SPILO_CONFIGURATION",
|
||||
Value: fmt.Sprintf(`
|
||||
postgresql:
|
||||
bin_dir: /usr/lib/postgresql/%s/bin
|
||||
bootstrap:
|
||||
initdb:
|
||||
- auth-host: md5
|
||||
- auth-local: trust
|
||||
users:
|
||||
%s:
|
||||
password: NULL
|
||||
options:
|
||||
- createdb
|
||||
- nologin
|
||||
pg_hba:
|
||||
- hostnossl all all all reject
|
||||
- hostssl all +%s all pam
|
||||
- hostssl all all all md5`, pgVersion, c.OpConfig.PamRoleName, c.OpConfig.PamRoleName),
|
||||
},
|
||||
}
|
||||
|
||||
container := v1.Container{
|
||||
Name: c.Metadata.Name,
|
||||
Image: c.OpConfig.DockerImage,
|
||||
ImagePullPolicy: v1.PullAlways,
|
||||
Resources: v1.ResourceRequirements{
|
||||
Requests: *resourceList,
|
||||
},
|
||||
Ports: []v1.ContainerPort{
|
||||
{
|
||||
ContainerPort: 8008,
|
||||
Protocol: v1.ProtocolTCP,
|
||||
},
|
||||
{
|
||||
ContainerPort: 5432,
|
||||
Protocol: v1.ProtocolTCP,
|
||||
},
|
||||
{
|
||||
ContainerPort: 8080,
|
||||
Protocol: v1.ProtocolTCP,
|
||||
},
|
||||
},
|
||||
VolumeMounts: []v1.VolumeMount{
|
||||
{
|
||||
Name: constants.DataVolumeName,
|
||||
MountPath: "/home/postgres/pgdata", //TODO: fetch from manifesto
|
||||
},
|
||||
},
|
||||
Env: envVars,
|
||||
}
|
||||
terminateGracePeriodSeconds := int64(30)
|
||||
|
||||
podSpec := v1.PodSpec{
|
||||
ServiceAccountName: c.OpConfig.ServiceAccountName,
|
||||
TerminationGracePeriodSeconds: &terminateGracePeriodSeconds,
|
||||
Containers: []v1.Container{container},
|
||||
}
|
||||
|
||||
template := v1.PodTemplateSpec{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Labels: c.labelsSet(),
|
||||
Namespace: c.Metadata.Name,
|
||||
},
|
||||
Spec: podSpec,
|
||||
}
|
||||
|
||||
return &template
|
||||
}
|
||||
|
||||
func (c *Cluster) genStatefulSet(spec spec.PostgresSpec) *v1beta1.StatefulSet {
|
||||
resourceList := resourceList(spec.Resources)
|
||||
podTemplate := c.genPodTemplate(resourceList, spec.PgVersion)
|
||||
volumeClaimTemplate := persistentVolumeClaimTemplate(spec.Volume.Size, spec.Volume.StorageClass)
|
||||
|
||||
statefulSet := &v1beta1.StatefulSet{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Name: c.Metadata.Name,
|
||||
Namespace: c.Metadata.Namespace,
|
||||
Labels: c.labelsSet(),
|
||||
},
|
||||
Spec: v1beta1.StatefulSetSpec{
|
||||
Replicas: &spec.NumberOfInstances,
|
||||
ServiceName: c.Metadata.Name,
|
||||
Template: *podTemplate,
|
||||
VolumeClaimTemplates: []v1.PersistentVolumeClaim{*volumeClaimTemplate},
|
||||
},
|
||||
}
|
||||
|
||||
return statefulSet
|
||||
}
|
||||
|
||||
func persistentVolumeClaimTemplate(volumeSize, volumeStorageClass string) *v1.PersistentVolumeClaim {
|
||||
metadata := v1.ObjectMeta{
|
||||
Name: constants.DataVolumeName,
|
||||
}
|
||||
if volumeStorageClass != "" {
|
||||
// TODO: check if storage class exists
|
||||
metadata.Annotations = map[string]string{"volume.beta.kubernetes.io/storage-class": volumeStorageClass}
|
||||
} else {
|
||||
metadata.Annotations = map[string]string{"volume.alpha.kubernetes.io/storage-class": "default"}
|
||||
}
|
||||
|
||||
volumeClaim := &v1.PersistentVolumeClaim{
|
||||
ObjectMeta: metadata,
|
||||
Spec: v1.PersistentVolumeClaimSpec{
|
||||
AccessModes: []v1.PersistentVolumeAccessMode{v1.ReadWriteOnce},
|
||||
Resources: v1.ResourceRequirements{
|
||||
Requests: v1.ResourceList{
|
||||
v1.ResourceStorage: resource.MustParse(volumeSize),
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
return volumeClaim
|
||||
}
|
||||
|
||||
func (c *Cluster) genUserSecrets() (secrets map[string]*v1.Secret, err error) {
|
||||
secrets = make(map[string]*v1.Secret, len(c.pgUsers))
|
||||
namespace := c.Metadata.Namespace
|
||||
for username, pgUser := range c.pgUsers {
|
||||
//Skip users with no password i.e. human users (they'll be authenticated using pam)
|
||||
if pgUser.Password == "" {
|
||||
continue
|
||||
}
|
||||
secret := v1.Secret{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Name: c.credentialSecretName(username),
|
||||
Namespace: namespace,
|
||||
Labels: c.labelsSet(),
|
||||
},
|
||||
Type: v1.SecretTypeOpaque,
|
||||
Data: map[string][]byte{
|
||||
"username": []byte(pgUser.Name),
|
||||
"password": []byte(pgUser.Password),
|
||||
},
|
||||
}
|
||||
secrets[username] = &secret
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) genService(allowedSourceRanges []string) *v1.Service {
|
||||
service := &v1.Service{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Name: c.Metadata.Name,
|
||||
Namespace: c.Metadata.Namespace,
|
||||
Labels: c.labelsSet(),
|
||||
Annotations: map[string]string{
|
||||
constants.ZalandoDnsNameAnnotation: util.ClusterDNSName(c.Metadata.Name, c.TeamName(), c.OpConfig.DbHostedZone),
|
||||
},
|
||||
},
|
||||
Spec: v1.ServiceSpec{
|
||||
Type: v1.ServiceTypeLoadBalancer,
|
||||
Ports: []v1.ServicePort{{Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}}},
|
||||
LoadBalancerSourceRanges: allowedSourceRanges,
|
||||
},
|
||||
}
|
||||
|
||||
return service
|
||||
}
|
||||
|
||||
func (c *Cluster) genEndpoints() *v1.Endpoints {
|
||||
endpoints := &v1.Endpoints{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Name: c.Metadata.Name,
|
||||
Namespace: c.Metadata.Namespace,
|
||||
Labels: c.labelsSet(),
|
||||
},
|
||||
}
|
||||
|
||||
return endpoints
|
||||
}
|
||||
+3
-4
@@ -9,18 +9,17 @@ import (
|
||||
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/spec"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/constants"
|
||||
)
|
||||
|
||||
var createUserSQL = `SET LOCAL synchronous_commit = 'local'; CREATE ROLE "%s" %s %s;`
|
||||
|
||||
func (c *Cluster) pgConnectionString() string {
|
||||
hostname := fmt.Sprintf("%s.%s.svc.cluster.local", c.Metadata.Name, c.Metadata.Namespace)
|
||||
password := c.pgUsers[constants.SuperuserName].Password
|
||||
password := c.pgUsers[c.OpConfig.SuperUsername].Password
|
||||
|
||||
return fmt.Sprintf("host='%s' dbname=postgres sslmode=require user='%s' password='%s'",
|
||||
hostname,
|
||||
constants.SuperuserName,
|
||||
c.OpConfig.SuperUsername,
|
||||
strings.Replace(password, "$", "\\$", -1))
|
||||
}
|
||||
|
||||
@@ -52,7 +51,7 @@ func (c *Cluster) createPgUser(user spec.PgUser) (isHuman bool, err error) {
|
||||
if user.Password == "" {
|
||||
isHuman = true
|
||||
flags = append(flags, "SUPERUSER")
|
||||
flags = append(flags, fmt.Sprintf("IN ROLE \"%s\"", constants.PamRoleName))
|
||||
flags = append(flags, fmt.Sprintf("IN ROLE \"%s\"", c.OpConfig.PamRoleName))
|
||||
} else {
|
||||
isHuman = false
|
||||
}
|
||||
|
||||
+6
-6
@@ -15,7 +15,7 @@ func (c *Cluster) listPods() ([]v1.Pod, error) {
|
||||
LabelSelector: c.labelsSet().String(),
|
||||
}
|
||||
|
||||
pods, err := c.config.KubeClient.Pods(ns).List(listOptions)
|
||||
pods, err := c.KubeClient.Pods(ns).List(listOptions)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("Can't get list of Pods: %s", err)
|
||||
}
|
||||
@@ -29,7 +29,7 @@ func (c *Cluster) listPersistentVolumeClaims() ([]v1.PersistentVolumeClaim, erro
|
||||
LabelSelector: c.labelsSet().String(),
|
||||
}
|
||||
|
||||
pvcs, err := c.config.KubeClient.PersistentVolumeClaims(ns).List(listOptions)
|
||||
pvcs, err := c.KubeClient.PersistentVolumeClaims(ns).List(listOptions)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("Can't get list of PersistentVolumeClaims: %s", err)
|
||||
}
|
||||
@@ -60,7 +60,7 @@ func (c *Cluster) deletePersistenVolumeClaims() error {
|
||||
return err
|
||||
}
|
||||
for _, pvc := range pvcs {
|
||||
if err := c.config.KubeClient.PersistentVolumeClaims(ns).Delete(pvc.Name, deleteOptions); err != nil {
|
||||
if err := c.KubeClient.PersistentVolumeClaims(ns).Delete(pvc.Name, deleteOptions); err != nil {
|
||||
c.logger.Warningf("Can't delete PersistentVolumeClaim: %s", err)
|
||||
}
|
||||
}
|
||||
@@ -83,7 +83,7 @@ func (c *Cluster) deletePod(pod *v1.Pod) error {
|
||||
delete(c.podSubscribers, podName)
|
||||
}()
|
||||
|
||||
if err := c.config.KubeClient.Pods(pod.Namespace).Delete(pod.Name, deleteOptions); err != nil {
|
||||
if err := c.KubeClient.Pods(pod.Namespace).Delete(pod.Name, deleteOptions); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -126,7 +126,7 @@ func (c *Cluster) recreatePod(pod v1.Pod, spiloRole string) error {
|
||||
ch := c.registerPodSubscriber(podName)
|
||||
defer c.unregisterPodSubscriber(podName)
|
||||
|
||||
if err := c.config.KubeClient.Pods(pod.Namespace).Delete(pod.Name, deleteOptions); err != nil {
|
||||
if err := c.KubeClient.Pods(pod.Namespace).Delete(pod.Name, deleteOptions); err != nil {
|
||||
return fmt.Errorf("Can't delete Pod: %s", err)
|
||||
}
|
||||
|
||||
@@ -165,7 +165,7 @@ func (c *Cluster) recreatePods() error {
|
||||
listOptions := v1.ListOptions{
|
||||
LabelSelector: ls.String(),
|
||||
}
|
||||
pods, err := c.config.KubeClient.Pods(namespace).List(listOptions)
|
||||
pods, err := c.KubeClient.Pods(namespace).List(listOptions)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get the list of Pods: %s", err)
|
||||
} else {
|
||||
|
||||
+23
-35
@@ -6,11 +6,8 @@ import (
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
"k8s.io/client-go/pkg/apis/apps/v1beta1"
|
||||
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/spec"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/constants"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/k8sutil"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/resources"
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -18,15 +15,6 @@ var (
|
||||
orphanDependents = false
|
||||
)
|
||||
|
||||
func genStatefulSet(clusterName spec.ClusterName, cSpec spec.PostgresSpec, etcdHost, dockerImage string) *v1beta1.StatefulSet {
|
||||
volumeSize := cSpec.Volume.Size
|
||||
volumeStorageClass := cSpec.Volume.StorageClass
|
||||
resourceList := resources.ResourceList(cSpec.Resources)
|
||||
template := resources.PodTemplate(clusterName, resourceList, cSpec.PgVersion, dockerImage, etcdHost)
|
||||
volumeClaimTemplate := resources.VolumeClaimTemplate(volumeSize, volumeStorageClass)
|
||||
|
||||
return resources.StatefulSet(clusterName, template, volumeClaimTemplate, cSpec.NumberOfInstances)
|
||||
}
|
||||
|
||||
func (c *Cluster) LoadResources() error {
|
||||
ns := c.Metadata.Namespace
|
||||
@@ -34,7 +22,7 @@ func (c *Cluster) LoadResources() error {
|
||||
LabelSelector: c.labelsSet().String(),
|
||||
}
|
||||
|
||||
services, err := c.config.KubeClient.Services(ns).List(listOptions)
|
||||
services, err := c.KubeClient.Services(ns).List(listOptions)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get list of Services: %s", err)
|
||||
}
|
||||
@@ -44,7 +32,7 @@ func (c *Cluster) LoadResources() error {
|
||||
c.Service = &services.Items[0]
|
||||
}
|
||||
|
||||
endpoints, err := c.config.KubeClient.Endpoints(ns).List(listOptions)
|
||||
endpoints, err := c.KubeClient.Endpoints(ns).List(listOptions)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get list of Endpoints: %s", err)
|
||||
}
|
||||
@@ -54,7 +42,7 @@ func (c *Cluster) LoadResources() error {
|
||||
c.Endpoint = &endpoints.Items[0]
|
||||
}
|
||||
|
||||
secrets, err := c.config.KubeClient.Secrets(ns).List(listOptions)
|
||||
secrets, err := c.KubeClient.Secrets(ns).List(listOptions)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get list of Secrets: %s", err)
|
||||
}
|
||||
@@ -66,7 +54,7 @@ func (c *Cluster) LoadResources() error {
|
||||
c.logger.Debugf("Secret loaded, uid: %s", secret.UID)
|
||||
}
|
||||
|
||||
statefulSets, err := c.config.KubeClient.StatefulSets(ns).List(listOptions)
|
||||
statefulSets, err := c.KubeClient.StatefulSets(ns).List(listOptions)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get list of StatefulSets: %s", err)
|
||||
}
|
||||
@@ -121,8 +109,8 @@ func (c *Cluster) createStatefulSet() (*v1beta1.StatefulSet, error) {
|
||||
if c.Statefulset != nil {
|
||||
return nil, fmt.Errorf("StatefulSet already exists in the cluster")
|
||||
}
|
||||
statefulSetSpec := genStatefulSet(c.ClusterName(), c.Spec, c.etcdHost, c.dockerImage)
|
||||
statefulSet, err := c.config.KubeClient.StatefulSets(statefulSetSpec.Namespace).Create(statefulSetSpec)
|
||||
statefulSetSpec := c.genStatefulSet(c.Spec)
|
||||
statefulSet, err := c.KubeClient.StatefulSets(statefulSetSpec.Namespace).Create(statefulSetSpec)
|
||||
if k8sutil.ResourceAlreadyExists(err) {
|
||||
return nil, fmt.Errorf("StatefulSet '%s' already exists", util.NameFromMeta(statefulSetSpec.ObjectMeta))
|
||||
}
|
||||
@@ -139,7 +127,7 @@ func (c *Cluster) updateStatefulSet(newStatefulSet *v1beta1.StatefulSet) error {
|
||||
if c.Statefulset == nil {
|
||||
return fmt.Errorf("There is no StatefulSet in the cluster")
|
||||
}
|
||||
statefulSet, err := c.config.KubeClient.StatefulSets(newStatefulSet.Namespace).Update(newStatefulSet)
|
||||
statefulSet, err := c.KubeClient.StatefulSets(newStatefulSet.Namespace).Update(newStatefulSet)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -154,7 +142,7 @@ func (c *Cluster) deleteStatefulSet() error {
|
||||
return fmt.Errorf("There is no StatefulSet in the cluster")
|
||||
}
|
||||
|
||||
err := c.config.KubeClient.StatefulSets(c.Statefulset.Namespace).Delete(c.Statefulset.Name, deleteOptions)
|
||||
err := c.KubeClient.StatefulSets(c.Statefulset.Namespace).Delete(c.Statefulset.Name, deleteOptions)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -167,9 +155,9 @@ func (c *Cluster) createService() (*v1.Service, error) {
|
||||
if c.Service != nil {
|
||||
return nil, fmt.Errorf("Service already exists in the cluster")
|
||||
}
|
||||
serviceSpec := resources.Service(c.ClusterName(), c.TeamName(), c.Spec.AllowedSourceRanges)
|
||||
serviceSpec := c.genService(c.Spec.AllowedSourceRanges)
|
||||
|
||||
service, err := c.config.KubeClient.Services(serviceSpec.Namespace).Create(serviceSpec)
|
||||
service, err := c.KubeClient.Services(serviceSpec.Namespace).Create(serviceSpec)
|
||||
if k8sutil.ResourceAlreadyExists(err) {
|
||||
return nil, fmt.Errorf("Service '%s' already exists", util.NameFromMeta(serviceSpec.ObjectMeta))
|
||||
}
|
||||
@@ -188,7 +176,7 @@ func (c *Cluster) updateService(newService *v1.Service) error {
|
||||
newService.ObjectMeta = c.Service.ObjectMeta
|
||||
newService.Spec.ClusterIP = c.Service.Spec.ClusterIP
|
||||
|
||||
svc, err := c.config.KubeClient.Services(newService.Namespace).Update(newService)
|
||||
svc, err := c.KubeClient.Services(newService.Namespace).Update(newService)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -201,7 +189,7 @@ func (c *Cluster) deleteService() error {
|
||||
if c.Service == nil {
|
||||
return fmt.Errorf("There is no Service in the cluster")
|
||||
}
|
||||
err := c.config.KubeClient.Services(c.Service.Namespace).Delete(c.Service.Name, deleteOptions)
|
||||
err := c.KubeClient.Services(c.Service.Namespace).Delete(c.Service.Name, deleteOptions)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -214,18 +202,18 @@ func (c *Cluster) createEndpoint() (*v1.Endpoints, error) {
|
||||
if c.Endpoint != nil {
|
||||
return nil, fmt.Errorf("Endpoint already exists in the cluster")
|
||||
}
|
||||
endpointSpec := resources.Endpoint(c.ClusterName())
|
||||
endpointsSpec := c.genEndpoints()
|
||||
|
||||
endpoint, err := c.config.KubeClient.Endpoints(endpointSpec.Namespace).Create(endpointSpec)
|
||||
endpoints, err := c.KubeClient.Endpoints(endpointsSpec.Namespace).Create(endpointsSpec)
|
||||
if k8sutil.ResourceAlreadyExists(err) {
|
||||
return nil, fmt.Errorf("Endpoint '%s' already exists", util.NameFromMeta(endpointSpec.ObjectMeta))
|
||||
return nil, fmt.Errorf("Endpoint '%s' already exists", util.NameFromMeta(endpointsSpec.ObjectMeta))
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
c.Endpoint = endpoint
|
||||
c.Endpoint = endpoints
|
||||
|
||||
return endpoint, nil
|
||||
return endpoints, nil
|
||||
}
|
||||
|
||||
func (c *Cluster) updateEndpoint(newEndpoint *v1.Endpoints) error {
|
||||
@@ -238,7 +226,7 @@ func (c *Cluster) deleteEndpoint() error {
|
||||
if c.Endpoint == nil {
|
||||
return fmt.Errorf("There is no Endpoint in the cluster")
|
||||
}
|
||||
err := c.config.KubeClient.Endpoints(c.Endpoint.Namespace).Delete(c.Endpoint.Name, deleteOptions)
|
||||
err := c.KubeClient.Endpoints(c.Endpoint.Namespace).Delete(c.Endpoint.Name, deleteOptions)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -248,16 +236,16 @@ func (c *Cluster) deleteEndpoint() error {
|
||||
}
|
||||
|
||||
func (c *Cluster) applySecrets() error {
|
||||
secrets, err := resources.UserSecrets(c.ClusterName(), c.pgUsers)
|
||||
secrets, err := c.genUserSecrets()
|
||||
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get user Secrets")
|
||||
}
|
||||
|
||||
for secretUsername, secretSpec := range secrets {
|
||||
secret, err := c.config.KubeClient.Secrets(secretSpec.Namespace).Create(secretSpec)
|
||||
secret, err := c.KubeClient.Secrets(secretSpec.Namespace).Create(secretSpec)
|
||||
if k8sutil.ResourceAlreadyExists(err) {
|
||||
curSecrets, err := c.config.KubeClient.Secrets(secretSpec.Namespace).Get(secretSpec.Name)
|
||||
curSecrets, err := c.KubeClient.Secrets(secretSpec.Namespace).Get(secretSpec.Name)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get current Secret: %s", err)
|
||||
}
|
||||
@@ -279,7 +267,7 @@ func (c *Cluster) applySecrets() error {
|
||||
}
|
||||
|
||||
func (c *Cluster) deleteSecret(secret *v1.Secret) error {
|
||||
err := c.config.KubeClient.Secrets(secret.Namespace).Delete(secret.Name, deleteOptions)
|
||||
err := c.KubeClient.Secrets(secret.Namespace).Delete(secret.Name, deleteOptions)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -291,7 +279,7 @@ func (c *Cluster) deleteSecret(secret *v1.Secret) error {
|
||||
func (c *Cluster) createUsers() error {
|
||||
// TODO: figure out what to do with duplicate names (humans and robots) among pgUsers
|
||||
for username, user := range c.pgUsers {
|
||||
if username == constants.SuperuserName || username == constants.ReplicationUsername {
|
||||
if username == c.OpConfig.SuperUsername || username == c.OpConfig.ReplicationUsername {
|
||||
continue
|
||||
}
|
||||
|
||||
|
||||
+3
-4
@@ -7,7 +7,6 @@ import (
|
||||
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/spec"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/resources"
|
||||
)
|
||||
|
||||
func (c *Cluster) SyncCluster() {
|
||||
@@ -55,7 +54,7 @@ func (c *Cluster) syncService() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
desiredSvc := resources.Service(c.ClusterName(), c.Spec.TeamId, cSpec.AllowedSourceRanges)
|
||||
desiredSvc := c.genService(cSpec.AllowedSourceRanges)
|
||||
if c.sameServiceWith(desiredSvc) {
|
||||
return nil
|
||||
}
|
||||
@@ -99,7 +98,7 @@ func (c *Cluster) syncStatefulSet() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
desiredSS := genStatefulSet(c.ClusterName(), cSpec, c.etcdHost, c.dockerImage)
|
||||
desiredSS := c.genStatefulSet(cSpec)
|
||||
equalSS, rollUpdate := c.compareStatefulSetWith(desiredSS)
|
||||
if equalSS {
|
||||
return nil
|
||||
@@ -132,7 +131,7 @@ func (c *Cluster) syncPods() error {
|
||||
listOptions := v1.ListOptions{
|
||||
LabelSelector: ls.String(),
|
||||
}
|
||||
pods, err := c.config.KubeClient.Pods(namespace).List(listOptions)
|
||||
pods, err := c.KubeClient.Pods(namespace).List(listOptions)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get list of Pods: %s", err)
|
||||
}
|
||||
|
||||
+10
-10
@@ -70,7 +70,7 @@ func podMatchesTemplate(pod *v1.Pod, ss *v1beta1.StatefulSet) bool {
|
||||
}
|
||||
|
||||
func (c *Cluster) getTeamMembers() ([]string, error) {
|
||||
teamInfo, err := c.config.TeamsAPIClient.TeamInfo(c.Spec.TeamId)
|
||||
teamInfo, err := c.TeamsAPIClient.TeamInfo(c.Spec.TeamId)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("Can't get team info: %s", err)
|
||||
}
|
||||
@@ -87,7 +87,7 @@ func (c *Cluster) waitForPodLabel(podEvents chan spec.PodEvent, spiloRole string
|
||||
if role == spiloRole { // TODO: newly-created Pods are always replicas => check against empty string only
|
||||
return nil
|
||||
}
|
||||
case <-time.After(constants.PodLabelWaitTimeout):
|
||||
case <-time.After(c.OpConfig.PodLabelWaitTimeout):
|
||||
return fmt.Errorf("Pod label wait timeout")
|
||||
}
|
||||
}
|
||||
@@ -100,19 +100,19 @@ func (c *Cluster) waitForPodDeletion(podEvents chan spec.PodEvent) error {
|
||||
if podEvent.EventType == spec.PodEventDelete {
|
||||
return nil
|
||||
}
|
||||
case <-time.After(constants.PodDeletionWaitTimeout):
|
||||
case <-time.After(c.OpConfig.PodDeletionWaitTimeout):
|
||||
return fmt.Errorf("Pod deletion wait timeout")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Cluster) waitStatefulsetReady() error {
|
||||
return retryutil.Retry(constants.ResourceCheckInterval, constants.ResourceCheckTimeout,
|
||||
return retryutil.Retry(c.OpConfig.ResourceCheckInterval, c.OpConfig.ResourceCheckTimeout,
|
||||
func() (bool, error) {
|
||||
listOptions := v1.ListOptions{
|
||||
LabelSelector: c.labelsSet().String(),
|
||||
}
|
||||
ss, err := c.config.KubeClient.StatefulSets(c.Metadata.Namespace).List(listOptions)
|
||||
ss, err := c.KubeClient.StatefulSets(c.Metadata.Namespace).List(listOptions)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
@@ -138,19 +138,19 @@ func (c *Cluster) waitPodLabelsReady() error {
|
||||
replicaListOption := v1.ListOptions{
|
||||
LabelSelector: labels.Merge(ls, labels.Set{"spilo-role": "replica"}).String(),
|
||||
}
|
||||
pods, err := c.config.KubeClient.Pods(namespace).List(listOptions)
|
||||
pods, err := c.KubeClient.Pods(namespace).List(listOptions)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
podsNumber := len(pods.Items)
|
||||
|
||||
return retryutil.Retry(constants.ResourceCheckInterval, constants.ResourceCheckTimeout,
|
||||
return retryutil.Retry(c.OpConfig.ResourceCheckInterval, c.OpConfig.ResourceCheckTimeout,
|
||||
func() (bool, error) {
|
||||
masterPods, err := c.config.KubeClient.Pods(namespace).List(masterListOption)
|
||||
masterPods, err := c.KubeClient.Pods(namespace).List(masterListOption)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
replicaPods, err := c.config.KubeClient.Pods(namespace).List(replicaListOption)
|
||||
replicaPods, err := c.KubeClient.Pods(namespace).List(replicaListOption)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
@@ -198,7 +198,7 @@ func (c *Cluster) deleteEtcdKey() error {
|
||||
etcdKey := fmt.Sprintf("/service/%s", c.Metadata.Name)
|
||||
|
||||
//TODO: retry multiple times
|
||||
resp, err := c.config.EtcdClient.Delete(context.Background(),
|
||||
resp, err := c.EtcdClient.Delete(context.Background(),
|
||||
etcdKey,
|
||||
&etcdclient.DeleteOptions{Recursive: true})
|
||||
|
||||
|
||||
Reference in New Issue
Block a user