mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-03 20:14:11 +02:00
add code to sync config maps
This commit is contained in:
+18
-9
@@ -59,6 +59,7 @@ type Config struct {
|
||||
type kubeResources struct {
|
||||
Services map[PostgresRole]*v1.Service
|
||||
Endpoints map[PostgresRole]*v1.Endpoints
|
||||
ConfigMaps map[PostgresRole]*v1.ConfigMap
|
||||
Secrets map[types.UID]*v1.Secret
|
||||
Statefulset *appsv1.StatefulSet
|
||||
PodDisruptionBudget *policybeta1.PodDisruptionBudget
|
||||
@@ -1484,22 +1485,29 @@ func (c *Cluster) GetCurrentProcess() Process {
|
||||
|
||||
// GetStatus provides status of the cluster
|
||||
func (c *Cluster) GetStatus() *ClusterStatus {
|
||||
return &ClusterStatus{
|
||||
Cluster: c.Spec.ClusterName,
|
||||
Team: c.Spec.TeamID,
|
||||
Status: c.Status,
|
||||
Spec: c.Spec,
|
||||
|
||||
status := &ClusterStatus{
|
||||
Cluster: c.Spec.ClusterName,
|
||||
Team: c.Spec.TeamID,
|
||||
Status: c.Status,
|
||||
Spec: c.Spec,
|
||||
MasterService: c.GetServiceMaster(),
|
||||
ReplicaService: c.GetServiceReplica(),
|
||||
MasterEndpoint: c.GetEndpointMaster(),
|
||||
ReplicaEndpoint: c.GetEndpointReplica(),
|
||||
StatefulSet: c.GetStatefulSet(),
|
||||
PodDisruptionBudget: c.GetPodDisruptionBudget(),
|
||||
CurrentProcess: c.GetCurrentProcess(),
|
||||
|
||||
Error: fmt.Errorf("error: %s", c.Error),
|
||||
}
|
||||
|
||||
if c.patroniKubernetesUseConfigMaps() {
|
||||
status.MasterEndpoint = c.GetEndpointMaster()
|
||||
status.ReplicaEndpoint = c.GetEndpointReplica()
|
||||
} else {
|
||||
status.MasterConfigMap = c.GetConfigMapMaster()
|
||||
status.ReplicaConfigMap = c.GetConfigMapReplica()
|
||||
}
|
||||
|
||||
return status
|
||||
}
|
||||
|
||||
// Switchover does a switchover (via Patroni) to a candidate pod
|
||||
@@ -1579,10 +1587,11 @@ func (c *Cluster) deletePatroniClusterObjects() error {
|
||||
}
|
||||
|
||||
if c.patroniKubernetesUseConfigMaps() {
|
||||
actionsList = append(actionsList, c.deletePatroniClusterServices, c.deletePatroniClusterConfigMaps)
|
||||
actionsList = append(actionsList, c.deletePatroniClusterConfigMaps)
|
||||
} else {
|
||||
actionsList = append(actionsList, c.deletePatroniClusterEndpoints)
|
||||
}
|
||||
actionsList = append(actionsList, c.deletePatroniClusterServices)
|
||||
|
||||
c.logger.Debugf("removing leftover Patroni objects (endpoints / services and configmaps)")
|
||||
for _, deleter := range actionsList {
|
||||
|
||||
+15
-6
@@ -76,13 +76,12 @@ func (c *Cluster) statefulSetName() string {
|
||||
return c.Name
|
||||
}
|
||||
|
||||
func (c *Cluster) endpointName(role PostgresRole) string {
|
||||
name := c.Name
|
||||
if role == Replica {
|
||||
name = name + "-repl"
|
||||
}
|
||||
func (c *Cluster) configMapName(role PostgresRole) string {
|
||||
return c.serviceName(role)
|
||||
}
|
||||
|
||||
return name
|
||||
func (c *Cluster) endpointName(role PostgresRole) string {
|
||||
return c.serviceName(role)
|
||||
}
|
||||
|
||||
func (c *Cluster) serviceName(role PostgresRole) string {
|
||||
@@ -1821,6 +1820,16 @@ func (c *Cluster) generateEndpoint(role PostgresRole, subsets []v1.EndpointSubse
|
||||
return endpoints
|
||||
}
|
||||
|
||||
func (c *Cluster) generateConfigMap(role PostgresRole) *v1.ConfigMap {
|
||||
return &v1.ConfigMap{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: c.configMapName(role),
|
||||
Namespace: c.Namespace,
|
||||
Labels: c.roleLabelsSet(true, role),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Cluster) generateCloneEnvironment(description *acidv1.CloneDescription) []v1.EnvVar {
|
||||
result := make([]v1.EnvVar, 0)
|
||||
|
||||
|
||||
@@ -35,8 +35,14 @@ func (c *Cluster) listResources() error {
|
||||
c.logger.Infof("found secret: %q (uid: %q) namesapce: %s", util.NameFromMeta(obj.ObjectMeta), obj.UID, obj.ObjectMeta.Namespace)
|
||||
}
|
||||
|
||||
for role, endpoint := range c.Endpoints {
|
||||
c.logger.Infof("found %s endpoint: %q (uid: %q)", role, util.NameFromMeta(endpoint.ObjectMeta), endpoint.UID)
|
||||
if c.patroniKubernetesUseConfigMaps() {
|
||||
for role, configMap := range c.ConfigMaps {
|
||||
c.logger.Infof("found %s config map: %q (uid: %q)", role, util.NameFromMeta(configMap.ObjectMeta), configMap.UID)
|
||||
}
|
||||
} else {
|
||||
for role, endpoint := range c.Endpoints {
|
||||
c.logger.Infof("found %s endpoint: %q (uid: %q)", role, util.NameFromMeta(endpoint.ObjectMeta), endpoint.UID)
|
||||
}
|
||||
}
|
||||
|
||||
for role, service := range c.Services {
|
||||
@@ -402,6 +408,20 @@ func (c *Cluster) generateEndpointSubsets(role PostgresRole) []v1.EndpointSubset
|
||||
return result
|
||||
}
|
||||
|
||||
func (c *Cluster) createConfigMap(role PostgresRole) (*v1.ConfigMap, error) {
|
||||
c.setProcessName("creating config map")
|
||||
configMapSpec := c.generateConfigMap(role)
|
||||
|
||||
configMap, err := c.KubeClient.ConfigMaps(configMapSpec.Namespace).Create(context.TODO(), configMapSpec, metav1.CreateOptions{})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("could not create %s config map: %v", role, err)
|
||||
}
|
||||
|
||||
c.ConfigMaps[role] = configMap
|
||||
|
||||
return configMap, nil
|
||||
}
|
||||
|
||||
func (c *Cluster) createPodDisruptionBudget() (*policybeta1.PodDisruptionBudget, error) {
|
||||
podDisruptionBudgetSpec := c.generatePodDisruptionBudget()
|
||||
podDisruptionBudget, err := c.KubeClient.
|
||||
@@ -589,11 +609,21 @@ func (c *Cluster) GetEndpointMaster() *v1.Endpoints {
|
||||
return c.Endpoints[Master]
|
||||
}
|
||||
|
||||
// GetEndpointReplica returns cluster's kubernetes master Endpoint
|
||||
// GetEndpointReplica returns cluster's kubernetes replica Endpoint
|
||||
func (c *Cluster) GetEndpointReplica() *v1.Endpoints {
|
||||
return c.Endpoints[Replica]
|
||||
}
|
||||
|
||||
// GetConfigMapMaster returns cluster's kubernetes master ConfigMap
|
||||
func (c *Cluster) GetConfigMapMaster() *v1.ConfigMap {
|
||||
return c.ConfigMaps[Master]
|
||||
}
|
||||
|
||||
// GetConfigMapReplica returns cluster's kubernetes replica ConfigMap
|
||||
func (c *Cluster) GetConfigMapReplica() *v1.ConfigMap {
|
||||
return c.ConfigMaps[Replica]
|
||||
}
|
||||
|
||||
// GetStatefulSet returns cluster's kubernetes StatefulSet
|
||||
func (c *Cluster) GetStatefulSet() *appsv1.StatefulSet {
|
||||
return c.Statefulset
|
||||
|
||||
+39
-1
@@ -144,7 +144,11 @@ func (c *Cluster) syncServices() error {
|
||||
for _, role := range []PostgresRole{Master, Replica} {
|
||||
c.logger.Debugf("syncing %s service", role)
|
||||
|
||||
if !c.patroniKubernetesUseConfigMaps() {
|
||||
if c.patroniKubernetesUseConfigMaps() {
|
||||
if err := c.syncConfigMap(role); err != nil {
|
||||
return fmt.Errorf("could not sync %s config map: %v", role, err)
|
||||
}
|
||||
} else {
|
||||
if err := c.syncEndpoint(role); err != nil {
|
||||
return fmt.Errorf("could not sync %s endpoint: %v", role, err)
|
||||
}
|
||||
@@ -234,6 +238,40 @@ func (c *Cluster) syncEndpoint(role PostgresRole) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) syncConfigMap(role PostgresRole) error {
|
||||
var (
|
||||
cm *v1.ConfigMap
|
||||
err error
|
||||
)
|
||||
c.setProcessName("syncing %s config map", role)
|
||||
|
||||
if cm, err = c.KubeClient.ConfigMaps(c.Namespace).Get(context.TODO(), c.configMapName(role), metav1.GetOptions{}); err == nil {
|
||||
// TODO: No syncing of config map here, is this covered completely by updateService?
|
||||
c.ConfigMaps[role] = cm
|
||||
return nil
|
||||
}
|
||||
if !k8sutil.ResourceNotFound(err) {
|
||||
return fmt.Errorf("could not get %s config map: %v", role, err)
|
||||
}
|
||||
// no existing config map, create new one
|
||||
c.ConfigMaps[role] = nil
|
||||
c.logger.Infof("could not find the cluster's %s config map", role)
|
||||
|
||||
if cm, err = c.createConfigMap(role); err == nil {
|
||||
c.logger.Infof("created missing %s config map %q", role, util.NameFromMeta(cm.ObjectMeta))
|
||||
} else {
|
||||
if !k8sutil.ResourceAlreadyExists(err) {
|
||||
return fmt.Errorf("could not create missing %s config map: %v", role, err)
|
||||
}
|
||||
c.logger.Infof("%s config map %q already exists", role, util.NameFromMeta(cm.ObjectMeta))
|
||||
if cm, err = c.KubeClient.ConfigMaps(c.Namespace).Get(context.TODO(), c.configMapName(role), metav1.GetOptions{}); err != nil {
|
||||
return fmt.Errorf("could not fetch existing %s config map: %v", role, err)
|
||||
}
|
||||
}
|
||||
c.ConfigMaps[role] = cm
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) syncPodDisruptionBudget(isUpdate bool) error {
|
||||
var (
|
||||
pdb *policybeta1.PodDisruptionBudget
|
||||
|
||||
@@ -63,6 +63,8 @@ type ClusterStatus struct {
|
||||
ReplicaService *v1.Service
|
||||
MasterEndpoint *v1.Endpoints
|
||||
ReplicaEndpoint *v1.Endpoints
|
||||
MasterConfigMap *v1.ConfigMap
|
||||
ReplicaConfigMap *v1.ConfigMap
|
||||
StatefulSet *appsv1.StatefulSet
|
||||
PodDisruptionBudget *policybeta1.PodDisruptionBudget
|
||||
|
||||
|
||||
@@ -544,7 +544,8 @@ func (c *Controller) postgresqlCheck(obj interface{}) *acidv1.Postgresql {
|
||||
Ensures the pod service account and role bindings exists in a namespace
|
||||
before a PG cluster is created there so that a user does not have to deploy
|
||||
these credentials manually. StatefulSets require the service account to
|
||||
create pods; Patroni requires relevant RBAC bindings to access endpoints.
|
||||
create pods; Patroni requires relevant RBAC bindings to access endpoints
|
||||
or config maps.
|
||||
|
||||
The operator does not sync accounts/role bindings after creation.
|
||||
*/
|
||||
|
||||
Reference in New Issue
Block a user