mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-01 17:47:37 +02:00
Linter-induced code refactoring, run round 2. (#360)
Run more linters in the gometalinter, i.e. deadcode, megacheck, nakedret, dup. More consistent code formatting, remove two dead functions, eliminate naked a bunch of naked returns, refactor a few functions to avoid code duplication.
This commit is contained in:
+11
-9
@@ -450,8 +450,8 @@ func (c *Cluster) compareContainers(setA, setB *v1beta1.StatefulSet) (bool, []st
|
||||
return needsRollUpdate, reasons
|
||||
}
|
||||
|
||||
func compareResources(a *v1.ResourceRequirements, b *v1.ResourceRequirements) (equal bool) {
|
||||
equal = true
|
||||
func compareResources(a *v1.ResourceRequirements, b *v1.ResourceRequirements) bool {
|
||||
equal := true
|
||||
if a != nil {
|
||||
equal = compareResoucesAssumeFirstNotNil(a, b)
|
||||
}
|
||||
@@ -459,7 +459,7 @@ func compareResources(a *v1.ResourceRequirements, b *v1.ResourceRequirements) (e
|
||||
equal = compareResoucesAssumeFirstNotNil(b, a)
|
||||
}
|
||||
|
||||
return
|
||||
return equal
|
||||
}
|
||||
|
||||
func compareResoucesAssumeFirstNotNil(a *v1.ResourceRequirements, b *v1.ResourceRequirements) bool {
|
||||
@@ -786,7 +786,8 @@ func (c *Cluster) initInfrastructureRoles() error {
|
||||
}
|
||||
|
||||
// resolves naming conflicts between existing and new roles by chosing either of them.
|
||||
func (c *Cluster) resolveNameConflict(currentRole, newRole *spec.PgUser) (result spec.PgUser) {
|
||||
func (c *Cluster) resolveNameConflict(currentRole, newRole *spec.PgUser) spec.PgUser {
|
||||
var result spec.PgUser
|
||||
if newRole.Origin >= currentRole.Origin {
|
||||
result = *newRole
|
||||
} else {
|
||||
@@ -794,7 +795,7 @@ func (c *Cluster) resolveNameConflict(currentRole, newRole *spec.PgUser) (result
|
||||
}
|
||||
c.logger.Debugf("resolved a conflict of role %q between %s and %s to %s",
|
||||
newRole.Name, newRole.Origin, currentRole.Origin, result.Origin)
|
||||
return
|
||||
return result
|
||||
}
|
||||
|
||||
func (c *Cluster) shouldAvoidProtectedOrSystemRole(username, purpose string) bool {
|
||||
@@ -838,8 +839,9 @@ func (c *Cluster) GetStatus() *spec.ClusterStatus {
|
||||
}
|
||||
|
||||
// Switchover does a switchover (via Patroni) to a candidate pod
|
||||
func (c *Cluster) Switchover(curMaster *v1.Pod, candidate spec.NamespacedName) (err error) {
|
||||
func (c *Cluster) Switchover(curMaster *v1.Pod, candidate spec.NamespacedName) error {
|
||||
|
||||
var err error
|
||||
c.logger.Debugf("failing over from %q to %q", curMaster.Name, candidate)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
@@ -858,8 +860,8 @@ func (c *Cluster) Switchover(curMaster *v1.Pod, candidate spec.NamespacedName) (
|
||||
|
||||
select {
|
||||
case <-stopCh:
|
||||
case podLabelErr <- func() (err error) {
|
||||
_, err = c.waitForPodLabel(ch, stopCh, &role)
|
||||
case podLabelErr <- func() (err2 error) {
|
||||
_, err2 = c.waitForPodLabel(ch, stopCh, &role)
|
||||
return
|
||||
}():
|
||||
}
|
||||
@@ -882,7 +884,7 @@ func (c *Cluster) Switchover(curMaster *v1.Pod, candidate spec.NamespacedName) (
|
||||
// close the label waiting channel no sooner than the waiting goroutine terminates.
|
||||
close(podLabelErr)
|
||||
|
||||
return
|
||||
return err
|
||||
|
||||
}
|
||||
|
||||
|
||||
+14
-14
@@ -7,13 +7,13 @@ import (
|
||||
|
||||
"github.com/Sirupsen/logrus"
|
||||
|
||||
"k8s.io/api/apps/v1beta1"
|
||||
"k8s.io/api/core/v1"
|
||||
policybeta1 "k8s.io/api/policy/v1beta1"
|
||||
"k8s.io/apimachinery/pkg/api/resource"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
"k8s.io/apimachinery/pkg/util/intstr"
|
||||
"k8s.io/api/core/v1"
|
||||
"k8s.io/api/apps/v1beta1"
|
||||
policybeta1 "k8s.io/api/policy/v1beta1"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/spec"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util"
|
||||
@@ -90,7 +90,7 @@ func (c *Cluster) makeDefaultResources() spec.Resources {
|
||||
defaultRequests := spec.ResourceDescription{CPU: config.DefaultCPURequest, Memory: config.DefaultMemoryRequest}
|
||||
defaultLimits := spec.ResourceDescription{CPU: config.DefaultCPULimit, Memory: config.DefaultMemoryLimit}
|
||||
|
||||
return spec.Resources{ResourceRequest:defaultRequests, ResourceLimits:defaultLimits}
|
||||
return spec.Resources{ResourceRequest: defaultRequests, ResourceLimits: defaultLimits}
|
||||
}
|
||||
|
||||
func generateResourceRequirements(resources spec.Resources, defaultResources spec.Resources) (*v1.ResourceRequirements, error) {
|
||||
@@ -366,7 +366,7 @@ func generateSidecarContainers(sidecars []spec.Sidecar,
|
||||
volumeMounts []v1.VolumeMount, defaultResources spec.Resources,
|
||||
superUserName string, credentialsSecretName string, logger *logrus.Entry) ([]v1.Container, error) {
|
||||
|
||||
if sidecars != nil && len(sidecars) > 0 {
|
||||
if len(sidecars) > 0 {
|
||||
result := make([]v1.Container, 0)
|
||||
for index, sidecar := range sidecars {
|
||||
|
||||
@@ -699,7 +699,7 @@ func (c *Cluster) generateStatefulSet(spec *spec.PostgresSpec) (*v1beta1.Statefu
|
||||
// generate sidecar containers
|
||||
if sidecarContainers, err = generateSidecarContainers(sideCars, volumeMounts, defaultResources,
|
||||
c.OpConfig.SuperUsername, c.credentialSecretName(c.OpConfig.SuperUsername), c.logger); err != nil {
|
||||
return nil, fmt.Errorf("could not generate sidecar containers: %v", err)
|
||||
return nil, fmt.Errorf("could not generate sidecar containers: %v", err)
|
||||
}
|
||||
|
||||
tolerationSpec := tolerations(&spec.Tolerations, c.OpConfig.PodToleration)
|
||||
@@ -716,7 +716,7 @@ func (c *Cluster) generateStatefulSet(spec *spec.PostgresSpec) (*v1beta1.Statefu
|
||||
int64(c.OpConfig.PodTerminateGracePeriod.Seconds()),
|
||||
c.OpConfig.PodServiceAccountName,
|
||||
c.OpConfig.KubeIAMRole,
|
||||
effectivePodPriorityClassName); err != nil{
|
||||
effectivePodPriorityClassName); err != nil {
|
||||
return nil, fmt.Errorf("could not generate pod template: %v", err)
|
||||
}
|
||||
|
||||
@@ -726,7 +726,7 @@ func (c *Cluster) generateStatefulSet(spec *spec.PostgresSpec) (*v1beta1.Statefu
|
||||
|
||||
if volumeClaimTemplate, err = generatePersistentVolumeClaimTemplate(spec.Volume.Size,
|
||||
spec.Volume.StorageClass); err != nil {
|
||||
return nil, fmt.Errorf("could not generate volume claim template: %v", err)
|
||||
return nil, fmt.Errorf("could not generate volume claim template: %v", err)
|
||||
}
|
||||
|
||||
numberOfInstances := c.getNumberOfInstances(spec)
|
||||
@@ -804,11 +804,11 @@ func (c *Cluster) mergeSidecars(sidecars []spec.Sidecar) []spec.Sidecar {
|
||||
return result
|
||||
}
|
||||
|
||||
func (c *Cluster) getNumberOfInstances(spec *spec.PostgresSpec) (newcur int32) {
|
||||
func (c *Cluster) getNumberOfInstances(spec *spec.PostgresSpec) int32 {
|
||||
min := c.OpConfig.MinInstances
|
||||
max := c.OpConfig.MaxInstances
|
||||
cur := spec.NumberOfInstances
|
||||
newcur = cur
|
||||
newcur := cur
|
||||
|
||||
if max >= 0 && newcur > max {
|
||||
newcur = max
|
||||
@@ -820,7 +820,7 @@ func (c *Cluster) getNumberOfInstances(spec *spec.PostgresSpec) (newcur int32) {
|
||||
c.logger.Infof("adjusted number of instances from %d to %d (min: %d, max: %d)", cur, newcur, min, max)
|
||||
}
|
||||
|
||||
return
|
||||
return newcur
|
||||
}
|
||||
|
||||
func generatePersistentVolumeClaimTemplate(volumeSize, volumeStorageClass string) (*v1.PersistentVolumeClaim, error) {
|
||||
@@ -860,8 +860,8 @@ func generatePersistentVolumeClaimTemplate(volumeSize, volumeStorageClass string
|
||||
return volumeClaim, nil
|
||||
}
|
||||
|
||||
func (c *Cluster) generateUserSecrets() (secrets map[string]*v1.Secret) {
|
||||
secrets = make(map[string]*v1.Secret, len(c.pgUsers))
|
||||
func (c *Cluster) generateUserSecrets() map[string]*v1.Secret {
|
||||
secrets := make(map[string]*v1.Secret, len(c.pgUsers))
|
||||
namespace := c.Namespace
|
||||
for username, pgUser := range c.pgUsers {
|
||||
//Skip users with no password i.e. human users (they'll be authenticated using pam)
|
||||
@@ -878,7 +878,7 @@ func (c *Cluster) generateUserSecrets() (secrets map[string]*v1.Secret) {
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
return secrets
|
||||
}
|
||||
|
||||
func (c *Cluster) generateSingleUserSecret(namespace string, pgUser spec.PgUser) *v1.Secret {
|
||||
|
||||
+11
-13
@@ -183,32 +183,30 @@ func (c *Cluster) getDatabases() (dbs map[string]string, err error) {
|
||||
dbs[datname] = owner
|
||||
}
|
||||
|
||||
return
|
||||
return dbs, err
|
||||
}
|
||||
|
||||
// executeCreateDatabase creates new database with the given owner.
|
||||
// The caller is responsible for openinging and closing the database connection.
|
||||
func (c *Cluster) executeCreateDatabase(datname, owner string) error {
|
||||
if !c.databaseNameOwnerValid(datname, owner) {
|
||||
return nil
|
||||
}
|
||||
c.logger.Infof("creating database %q with owner %q", datname, owner)
|
||||
|
||||
if _, err := c.pgDb.Exec(fmt.Sprintf(createDatabaseSQL, datname, owner)); err != nil {
|
||||
return fmt.Errorf("could not execute create database: %v", err)
|
||||
}
|
||||
return nil
|
||||
return c.execCreateOrAlterDatabase(datname, owner, createDatabaseSQL,
|
||||
"creating database", "create database")
|
||||
}
|
||||
|
||||
// executeCreateDatabase changes the owner of the given database.
|
||||
// The caller is responsible for openinging and closing the database connection.
|
||||
func (c *Cluster) executeAlterDatabaseOwner(datname string, owner string) error {
|
||||
return c.execCreateOrAlterDatabase(datname, owner, alterDatabaseOwnerSQL,
|
||||
"changing owner for database", "alter database owner")
|
||||
}
|
||||
|
||||
func (c *Cluster) execCreateOrAlterDatabase(datname, owner, statement, doing, operation string) error {
|
||||
if !c.databaseNameOwnerValid(datname, owner) {
|
||||
return nil
|
||||
}
|
||||
c.logger.Infof("changing database %q owner to %q", datname, owner)
|
||||
if _, err := c.pgDb.Exec(fmt.Sprintf(alterDatabaseOwnerSQL, datname, owner)); err != nil {
|
||||
return fmt.Errorf("could not execute alter database owner: %v", err)
|
||||
c.logger.Infof("%s %q owner %q", doing, datname, owner)
|
||||
if _, err := c.pgDb.Exec(fmt.Sprintf(statement, datname, owner)); err != nil {
|
||||
return fmt.Errorf("could not execute %s: %v", operation, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
+3
-4
@@ -183,8 +183,8 @@ func (c *Cluster) masterCandidate(oldNodeName string) (*v1.Pod, error) {
|
||||
func (c *Cluster) MigrateMasterPod(podName spec.NamespacedName) error {
|
||||
var (
|
||||
masterCandidatePod *v1.Pod
|
||||
err error
|
||||
eol bool
|
||||
err error
|
||||
eol bool
|
||||
)
|
||||
|
||||
oldMaster, err := c.KubeClient.Pods(podName.Namespace).Get(podName.Name, metav1.GetOptions{})
|
||||
@@ -212,7 +212,7 @@ func (c *Cluster) MigrateMasterPod(podName spec.NamespacedName) error {
|
||||
var sset *v1beta1.StatefulSet
|
||||
if sset, err = c.KubeClient.StatefulSets(c.Namespace).Get(c.statefulSetName(),
|
||||
metav1.GetOptions{}); err != nil {
|
||||
return fmt.Errorf("could not retrieve cluster statefulset: %v", err)
|
||||
return fmt.Errorf("could not retrieve cluster statefulset: %v", err)
|
||||
}
|
||||
c.Statefulset = sset
|
||||
}
|
||||
@@ -225,7 +225,6 @@ func (c *Cluster) MigrateMasterPod(podName spec.NamespacedName) error {
|
||||
c.logger.Warningf("single master pod for cluster %q, migration will cause longer downtime of the master instance", c.clusterName())
|
||||
}
|
||||
|
||||
|
||||
// there are two cases for each postgres cluster that has its master pod on the node to migrate from:
|
||||
// - the cluster has some replicas - migrate one of those if necessary and failover to it
|
||||
// - there are no replicas - just terminate the master and wait until it respawns
|
||||
|
||||
@@ -633,27 +633,3 @@ func (c *Cluster) GetStatefulSet() *v1beta1.StatefulSet {
|
||||
func (c *Cluster) GetPodDisruptionBudget() *policybeta1.PodDisruptionBudget {
|
||||
return c.PodDisruptionBudget
|
||||
}
|
||||
|
||||
func (c *Cluster) createDatabases() error {
|
||||
c.setProcessName("creating databases")
|
||||
|
||||
if len(c.Spec.Databases) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
if err := c.initDbConn(); err != nil {
|
||||
return fmt.Errorf("could not init database connection")
|
||||
}
|
||||
defer func() {
|
||||
if err := c.closeDbConn(); err != nil {
|
||||
c.logger.Errorf("could not close database connection: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
for datname, owner := range c.Spec.Databases {
|
||||
if err := c.executeCreateDatabase(datname, owner); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
+21
-29
@@ -2,11 +2,9 @@ package cluster
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"reflect"
|
||||
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
policybeta1 "k8s.io/api/policy/v1beta1"
|
||||
"k8s.io/api/core/v1"
|
||||
policybeta1 "k8s.io/api/policy/v1beta1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/spec"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util"
|
||||
@@ -17,7 +15,8 @@ import (
|
||||
|
||||
// Sync syncs the cluster, making sure the actual Kubernetes objects correspond to what is defined in the manifest.
|
||||
// Unlike the update, sync does not error out if some objects do not exist and takes care of creating them.
|
||||
func (c *Cluster) Sync(newSpec *spec.Postgresql) (err error) {
|
||||
func (c *Cluster) Sync(newSpec *spec.Postgresql) error {
|
||||
var err error
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
@@ -34,7 +33,7 @@ func (c *Cluster) Sync(newSpec *spec.Postgresql) (err error) {
|
||||
|
||||
if err = c.initUsers(); err != nil {
|
||||
err = fmt.Errorf("could not init users: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
c.logger.Debugf("syncing secrets")
|
||||
@@ -42,13 +41,13 @@ func (c *Cluster) Sync(newSpec *spec.Postgresql) (err error) {
|
||||
//TODO: mind the secrets of the deleted/new users
|
||||
if err = c.syncSecrets(); err != nil {
|
||||
err = fmt.Errorf("could not sync secrets: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
c.logger.Debugf("syncing services")
|
||||
if err = c.syncServices(); err != nil {
|
||||
err = fmt.Errorf("could not sync services: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
// potentially enlarge volumes before changing the statefulset. By doing that
|
||||
@@ -60,14 +59,14 @@ func (c *Cluster) Sync(newSpec *spec.Postgresql) (err error) {
|
||||
c.logger.Debugf("syncing persistent volumes")
|
||||
if err = c.syncVolumes(); err != nil {
|
||||
err = fmt.Errorf("could not sync persistent volumes: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
c.logger.Debugf("syncing statefulsets")
|
||||
if err = c.syncStatefulSet(); err != nil {
|
||||
if !k8sutil.ResourceAlreadyExists(err) {
|
||||
err = fmt.Errorf("could not sync statefulsets: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
@@ -76,22 +75,22 @@ func (c *Cluster) Sync(newSpec *spec.Postgresql) (err error) {
|
||||
c.logger.Debugf("syncing roles")
|
||||
if err = c.syncRoles(); err != nil {
|
||||
err = fmt.Errorf("could not sync roles: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
c.logger.Debugf("syncing databases")
|
||||
if err = c.syncDatabases(); err != nil {
|
||||
err = fmt.Errorf("could not sync databases: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
c.logger.Debug("syncing pod disruption budgets")
|
||||
if err = c.syncPodDisruptionBudget(false); err != nil {
|
||||
err = fmt.Errorf("could not sync pod disruption budget: %v", err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
return
|
||||
return err
|
||||
}
|
||||
|
||||
func (c *Cluster) syncServices() error {
|
||||
@@ -153,7 +152,7 @@ func (c *Cluster) syncService(role PostgresRole) error {
|
||||
|
||||
func (c *Cluster) syncEndpoint(role PostgresRole) error {
|
||||
var (
|
||||
ep *v1.Endpoints
|
||||
ep *v1.Endpoints
|
||||
err error
|
||||
)
|
||||
c.setProcessName("syncing %s endpoint", role)
|
||||
@@ -187,8 +186,8 @@ func (c *Cluster) syncEndpoint(role PostgresRole) error {
|
||||
|
||||
func (c *Cluster) syncPodDisruptionBudget(isUpdate bool) error {
|
||||
var (
|
||||
pdb *policybeta1.PodDisruptionBudget
|
||||
err error
|
||||
pdb *policybeta1.PodDisruptionBudget
|
||||
err error
|
||||
)
|
||||
if pdb, err = c.KubeClient.PodDisruptionBudgets(c.Namespace).Get(c.podDisruptionBudgetName(), metav1.GetOptions{}); err == nil {
|
||||
c.PodDisruptionBudget = pdb
|
||||
@@ -257,7 +256,9 @@ func (c *Cluster) syncStatefulSet() error {
|
||||
podsRollingUpdateRequired = (len(pods) > 0)
|
||||
if podsRollingUpdateRequired {
|
||||
c.logger.Warningf("found pods from the previous statefulset: trigger rolling update")
|
||||
c.applyRollingUpdateFlagforStatefulSet(podsRollingUpdateRequired)
|
||||
if err := c.applyRollingUpdateFlagforStatefulSet(podsRollingUpdateRequired); err != nil {
|
||||
return fmt.Errorf("could not set rolling update flag for the statefulset: %v", err)
|
||||
}
|
||||
}
|
||||
c.logger.Infof("created missing statefulset %q", util.NameFromMeta(sset.ObjectMeta))
|
||||
|
||||
@@ -318,7 +319,7 @@ func (c *Cluster) syncStatefulSet() error {
|
||||
// (like max_connections) has changed and if necessary sets it via the Patroni API
|
||||
func (c *Cluster) checkAndSetGlobalPostgreSQLConfiguration() error {
|
||||
var (
|
||||
err error
|
||||
err error
|
||||
pods []v1.Pod
|
||||
)
|
||||
|
||||
@@ -394,7 +395,7 @@ func (c *Cluster) syncSecrets() error {
|
||||
// if this secret belongs to the infrastructure role and the password has changed - replace it in the secret
|
||||
if pwdUser.Password != string(secret.Data["password"]) && pwdUser.Origin == spec.RoleOriginInfrastructure {
|
||||
c.logger.Debugf("updating the secret %q from the infrastructure roles", secretSpec.Name)
|
||||
if secret, err = c.KubeClient.Secrets(secretSpec.Namespace).Update(secretSpec); err != nil {
|
||||
if _, err = c.KubeClient.Secrets(secretSpec.Namespace).Update(secretSpec); err != nil {
|
||||
return fmt.Errorf("could not update infrastructure role secret for role %q: %v", secretUsername, err)
|
||||
}
|
||||
} else {
|
||||
@@ -469,15 +470,6 @@ func (c *Cluster) syncVolumes() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) samePDBWith(pdb *policybeta1.PodDisruptionBudget) (match bool, reason string) {
|
||||
match = reflect.DeepEqual(pdb.Spec, c.PodDisruptionBudget.Spec)
|
||||
if !match {
|
||||
reason = "new service spec doesn't match the current one"
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) syncDatabases() error {
|
||||
c.setProcessName("syncing databases")
|
||||
|
||||
|
||||
Reference in New Issue
Block a user