mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-05 01:34:45 +02:00
conflict resolution for master merge
This commit is contained in:
@@ -781,7 +781,13 @@ func (c *Cluster) Update(oldSpec, newSpec *acidv1.Postgresql) error {
|
||||
}
|
||||
}
|
||||
|
||||
// sync connection pooler
|
||||
// Sync connection pooler. Before actually doing sync reset lookup
|
||||
// installation flag, since manifest updates could add another db which we
|
||||
// need to process. In the future we may want to do this more careful and
|
||||
// check which databases we need to process, but even repeating the whole
|
||||
// installation process should be good enough.
|
||||
c.ConnectionPooler.LookupFunction = false
|
||||
|
||||
if _, err := c.syncConnectionPooler(oldSpec, newSpec,
|
||||
c.installLookupFunction); err != nil {
|
||||
c.logger.Errorf("could not sync connection pooler: %v", err)
|
||||
|
||||
+58
-26
@@ -101,15 +101,20 @@ func (c *Cluster) databaseAccessDisabled() bool {
|
||||
}
|
||||
|
||||
func (c *Cluster) initDbConn() error {
|
||||
return c.initDbConnWithName("")
|
||||
}
|
||||
|
||||
func (c *Cluster) initDbConnWithName(dbname string) error {
|
||||
c.setProcessName("initializing db connection")
|
||||
if c.pgDb != nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return c.initDbConnWithName("")
|
||||
}
|
||||
|
||||
// Worker function for connection initialization. This function does not check
|
||||
// if the connection is already open, if it is then it will be overwritten.
|
||||
// Callers need to make sure no connection is open, otherwise we could leak
|
||||
// connections
|
||||
func (c *Cluster) initDbConnWithName(dbname string) error {
|
||||
c.setProcessName("initializing db connection")
|
||||
|
||||
var conn *sql.DB
|
||||
connstring := c.pgConnectionString(dbname)
|
||||
|
||||
@@ -145,6 +150,12 @@ func (c *Cluster) initDbConnWithName(dbname string) error {
|
||||
conn.SetMaxOpenConns(1)
|
||||
conn.SetMaxIdleConns(-1)
|
||||
|
||||
if c.pgDb != nil {
|
||||
msg := "Closing an existing connection before opening a new one to %s"
|
||||
c.logger.Warningf(msg, dbname)
|
||||
c.closeDbConn()
|
||||
}
|
||||
|
||||
c.pgDb = conn
|
||||
|
||||
return nil
|
||||
@@ -465,8 +476,11 @@ func (c *Cluster) execCreateOrAlterExtension(extName, schemaName, statement, doi
|
||||
// perform remote authentication.
|
||||
func (c *Cluster) installLookupFunction(poolerSchema, poolerUser string) error {
|
||||
var stmtBytes bytes.Buffer
|
||||
|
||||
c.logger.Info("Installing lookup function")
|
||||
|
||||
// Open a new connection if not yet done. This connection will be used only
|
||||
// to get the list of databases, not for the actuall installation.
|
||||
if err := c.initDbConn(); err != nil {
|
||||
return fmt.Errorf("could not init database connection")
|
||||
}
|
||||
@@ -480,37 +494,41 @@ func (c *Cluster) installLookupFunction(poolerSchema, poolerUser string) error {
|
||||
}
|
||||
}()
|
||||
|
||||
// List of databases we failed to process. At the moment it function just
|
||||
// like a flag to retry on the next sync, but in the future we may want to
|
||||
// retry only necessary parts, so let's keep the list.
|
||||
failedDatabases := []string{}
|
||||
currentDatabases, err := c.getDatabases()
|
||||
if err != nil {
|
||||
msg := "could not get databases to install pooler lookup function: %v"
|
||||
return fmt.Errorf(msg, err)
|
||||
}
|
||||
|
||||
// We've got the list of target databases, now close this connection to
|
||||
// open a new one to every each of them.
|
||||
if err := c.closeDbConn(); err != nil {
|
||||
c.logger.Errorf("could not close database connection: %v", err)
|
||||
}
|
||||
|
||||
templater := template.Must(template.New("sql").Parse(connectionPoolerLookup))
|
||||
params := TemplateParams{
|
||||
"pooler_schema": poolerSchema,
|
||||
"pooler_user": poolerUser,
|
||||
}
|
||||
|
||||
if err := templater.Execute(&stmtBytes, params); err != nil {
|
||||
msg := "could not prepare sql statement %+v: %v"
|
||||
return fmt.Errorf(msg, params, err)
|
||||
}
|
||||
|
||||
for dbname := range currentDatabases {
|
||||
|
||||
if dbname == "template0" || dbname == "template1" {
|
||||
continue
|
||||
}
|
||||
|
||||
if err := c.initDbConnWithName(dbname); err != nil {
|
||||
return fmt.Errorf("could not init database connection to %s", dbname)
|
||||
}
|
||||
|
||||
c.logger.Infof("Install pooler lookup function into %s", dbname)
|
||||
|
||||
params := TemplateParams{
|
||||
"pooler_schema": poolerSchema,
|
||||
"pooler_user": poolerUser,
|
||||
}
|
||||
|
||||
if err := templater.Execute(&stmtBytes, params); err != nil {
|
||||
c.logger.Errorf("could not prepare sql statement %+v: %v",
|
||||
params, err)
|
||||
// process other databases
|
||||
continue
|
||||
}
|
||||
|
||||
// golang sql will do retries couple of times if pq driver reports
|
||||
// connections issues (driver.ErrBadConn), but since our query is
|
||||
// idempotent, we can retry in a view of other errors (e.g. due to
|
||||
@@ -520,7 +538,20 @@ func (c *Cluster) installLookupFunction(poolerSchema, poolerUser string) error {
|
||||
constants.PostgresConnectTimeout,
|
||||
constants.PostgresConnectRetryTimeout,
|
||||
func() (bool, error) {
|
||||
if _, err := c.pgDb.Exec(stmtBytes.String()); err != nil {
|
||||
|
||||
// At this moment we are not connected to any database
|
||||
if err := c.initDbConnWithName(dbname); err != nil {
|
||||
msg := "could not init database connection to %s"
|
||||
return false, fmt.Errorf(msg, dbname)
|
||||
}
|
||||
defer func() {
|
||||
if err := c.closeDbConn(); err != nil {
|
||||
msg := "could not close database connection: %v"
|
||||
c.logger.Errorf(msg, err)
|
||||
}
|
||||
}()
|
||||
|
||||
if _, err = c.pgDb.Exec(stmtBytes.String()); err != nil {
|
||||
msg := fmt.Errorf("could not execute sql statement %s: %v",
|
||||
stmtBytes.String(), err)
|
||||
return false, msg
|
||||
@@ -533,15 +564,16 @@ func (c *Cluster) installLookupFunction(poolerSchema, poolerUser string) error {
|
||||
c.logger.Errorf("could not execute after retries %s: %v",
|
||||
stmtBytes.String(), err)
|
||||
// process other databases
|
||||
failedDatabases = append(failedDatabases, dbname)
|
||||
continue
|
||||
}
|
||||
|
||||
c.logger.Infof("pooler lookup function installed into %s", dbname)
|
||||
if err := c.closeDbConn(); err != nil {
|
||||
c.logger.Errorf("could not close database connection: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
c.ConnectionPooler.LookupFunction = true
|
||||
if len(failedDatabases) == 0 {
|
||||
c.ConnectionPooler.LookupFunction = true
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -845,7 +845,7 @@ func (c *Cluster) getPodEnvironmentSecretVariables() ([]v1.EnvVar, error) {
|
||||
return secretPodEnvVarsList, nil
|
||||
}
|
||||
|
||||
secret, err := c.KubeClient.Secrets(c.OpConfig.PodEnvironmentSecret).Get(
|
||||
secret, err := c.KubeClient.Secrets(c.Namespace).Get(
|
||||
context.TODO(),
|
||||
c.OpConfig.PodEnvironmentSecret,
|
||||
metav1.GetOptions{})
|
||||
@@ -2220,6 +2220,13 @@ func (c *Cluster) generateConnectionPoolerPodTemplate(spec *acidv1.PostgresSpec,
|
||||
},
|
||||
},
|
||||
Env: envVars,
|
||||
ReadinessProbe: &v1.Probe{
|
||||
Handler: v1.Handler{
|
||||
TCPSocket: &v1.TCPSocketAction{
|
||||
Port: intstr.IntOrString{IntVal: pgPort},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
podTemplate := &v1.PodTemplateSpec{
|
||||
|
||||
@@ -900,6 +900,11 @@ func (c *Cluster) syncConnectionPooler(oldSpec,
|
||||
c.logger.Errorf("could not sync connection pooler: %v", err)
|
||||
return reason, err
|
||||
}
|
||||
} else {
|
||||
// Lookup function installation seems to be a fragile point, so
|
||||
// let's log for debugging if we skip it
|
||||
msg := "Skip lookup function installation, old: %d, already installed %d"
|
||||
c.logger.Debug(msg, oldNeedConnectionPooler, c.ConnectionPooler.LookupFunction)
|
||||
}
|
||||
|
||||
if oldNeedConnectionPooler && !newNeedConnectionPooler {
|
||||
|
||||
Reference in New Issue
Block a user