mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-06 22:12:27 +02:00
Periodically sync roles with the running clusters. (#102)
The sync adds or alters database roles based on the roles defined in the cluster's TPR, Team API and operator's infrastructure roles. At the moment, roles are not deleted, as it would be dangerous for the robot roles in case TPR is misconfigured. In addition, ALTER ROLE does not remove role options, i.e. SUPERUSER or CREATEROLE, neither it removes role membership: only new options are added and new role membership is granted. So far, options like NOSUPERUSER and NOCREATEROLE won't be handed correctly, when mixed with the non-negative counterparts, also NOLOGIN should be processed correctly. The code assumes that only MD5 passwords are stored in the DB and will likely break with the new SCRAM auth in PostgreSQL 10. On the implementation side, create the new interface to abstract roles merge and creation, move most of the role-based functionality from cluster/pg into the new 'users' module, strip create user code of special cases related to human-based users (moving them to init instead) and fixed the password md5 generator to avoid processing already encrypted passwords. In addition, moved the system roles off the slice containing all other roles in order to avoid extra efforts to avoid creating them. Also, fix a leak in DB connections when the new connection is not considered healthy and discarded without being closed. Initialize the database during the sync phase before syncing users.
This commit is contained in:
committed by
Murat Kabilov
parent
411487e66d
commit
6983f444ed
+19
-4
@@ -26,6 +26,7 @@ import (
|
||||
"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/teams"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/users"
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -58,6 +59,7 @@ type Cluster struct {
|
||||
Config
|
||||
logger *logrus.Entry
|
||||
pgUsers map[string]spec.PgUser
|
||||
systemUsers map[string]spec.PgUser
|
||||
podEvents chan spec.PodEvent
|
||||
podSubscribers map[spec.NamespacedName]chan spec.PodEvent
|
||||
podSubscribersMu sync.RWMutex
|
||||
@@ -65,6 +67,7 @@ type Cluster struct {
|
||||
mu sync.Mutex
|
||||
masterLess bool
|
||||
podDispatcherRunning bool
|
||||
userSyncStrategy spec.UserSyncer
|
||||
deleteOptions *v1.DeleteOptions
|
||||
}
|
||||
|
||||
@@ -78,11 +81,13 @@ func New(cfg Config, pgSpec spec.Postgresql, logger *logrus.Entry) *Cluster {
|
||||
Postgresql: pgSpec,
|
||||
logger: lg,
|
||||
pgUsers: make(map[string]spec.PgUser),
|
||||
systemUsers: make(map[string]spec.PgUser),
|
||||
podEvents: make(chan spec.PodEvent),
|
||||
podSubscribers: make(map[spec.NamespacedName]chan spec.PodEvent),
|
||||
kubeResources: kubeResources,
|
||||
masterLess: false,
|
||||
podDispatcherRunning: false,
|
||||
userSyncStrategy: users.DefaultUserSyncStrategy{},
|
||||
deleteOptions: &v1.DeleteOptions{OrphanDependents: &orphanDependents},
|
||||
}
|
||||
|
||||
@@ -426,12 +431,15 @@ func (c *Cluster) ReceivePodEvent(event spec.PodEvent) {
|
||||
}
|
||||
|
||||
func (c *Cluster) initSystemUsers() {
|
||||
c.pgUsers[c.OpConfig.SuperUsername] = spec.PgUser{
|
||||
// We don't actually use that to create users, delegating this
|
||||
// task to Patroni. Those definitions are only used to create
|
||||
// secrets, therefore, setting flags like SUPERUSER or REPLICATION
|
||||
// is not necessary here
|
||||
c.systemUsers[constants.SuperuserKeyName] = spec.PgUser{
|
||||
Name: c.OpConfig.SuperUsername,
|
||||
Password: util.RandomPassword(constants.PasswordLength),
|
||||
}
|
||||
|
||||
c.pgUsers[c.OpConfig.ReplicationUsername] = spec.PgUser{
|
||||
c.systemUsers[constants.ReplicationUserKeyName] = spec.PgUser{
|
||||
Name: c.OpConfig.ReplicationUsername,
|
||||
Password: util.RandomPassword(constants.PasswordLength),
|
||||
}
|
||||
@@ -464,7 +472,9 @@ func (c *Cluster) initHumanUsers() error {
|
||||
return fmt.Errorf("Can't get list of team members: %s", err)
|
||||
} else {
|
||||
for _, username := range teamMembers {
|
||||
c.pgUsers[username] = spec.PgUser{Name: username}
|
||||
flags := []string{constants.RoleFlagLogin, constants.RoleFlagSuperuser}
|
||||
memberOf := []string{c.OpConfig.PamRoleName}
|
||||
c.pgUsers[username] = spec.PgUser{Name: username, Flags: flags, MemberOf: memberOf}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -477,6 +487,11 @@ func (c *Cluster) initInfrastructureRoles() error {
|
||||
if !isValidUsername(username) {
|
||||
return fmt.Errorf("Invalid username: '%s'", username)
|
||||
}
|
||||
if flags, err := normalizeUserFlags(data.Flags); err != nil {
|
||||
return fmt.Errorf("Invalid flags for user '%s': %s", username, err)
|
||||
} else {
|
||||
data.Flags = flags
|
||||
}
|
||||
c.pgUsers[username] = data
|
||||
}
|
||||
return nil
|
||||
|
||||
+31
-15
@@ -228,32 +228,48 @@ func persistentVolumeClaimTemplate(volumeSize, volumeStorageClass string) *v1.Pe
|
||||
return volumeClaim
|
||||
}
|
||||
|
||||
func (c *Cluster) genUserSecrets() (secrets map[string]*v1.Secret, err error) {
|
||||
func (c *Cluster) genUserSecrets() (secrets map[string]*v1.Secret) {
|
||||
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 := c.genSingleUserSecret(namespace, pgUser)
|
||||
if secret != nil {
|
||||
secrets[username] = secret
|
||||
}
|
||||
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),
|
||||
},
|
||||
}
|
||||
/* special case for the system user */
|
||||
for _, systemUser := range c.systemUsers {
|
||||
secret := c.genSingleUserSecret(namespace, systemUser)
|
||||
if secret != nil {
|
||||
secrets[systemUser.Name] = secret
|
||||
}
|
||||
secrets[username] = &secret
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) genSingleUserSecret(namespace string, pgUser spec.PgUser) *v1.Secret {
|
||||
//Skip users with no password i.e. human users (they'll be authenticated using pam)
|
||||
if pgUser.Password == "" {
|
||||
return nil
|
||||
}
|
||||
username := pgUser.Name
|
||||
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),
|
||||
},
|
||||
}
|
||||
return &secret
|
||||
}
|
||||
|
||||
func (c *Cluster) genService(allowedSourceRanges []string) *v1.Service {
|
||||
service := &v1.Service{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
|
||||
+55
-38
@@ -8,18 +8,28 @@ import (
|
||||
_ "github.com/lib/pq"
|
||||
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/spec"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util"
|
||||
"github.com/lib/pq"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/constants"
|
||||
)
|
||||
|
||||
var createUserSQL = `SET LOCAL synchronous_commit = 'local'; CREATE ROLE "%s" %s %s;`
|
||||
var getUserSQL = `SELECT a.rolname, COALESCE(a.rolpassword, ''), a.rolsuper, a.rolinherit,
|
||||
a.rolcreaterole, a.rolcreatedb, a.rolcanlogin,
|
||||
ARRAY(SELECT b.rolname
|
||||
FROM pg_catalog.pg_auth_members m
|
||||
JOIN pg_catalog.pg_authid b ON (m.roleid = b.oid)
|
||||
WHERE m.member = a.oid) as memberof
|
||||
FROM pg_catalog.pg_authid a
|
||||
WHERE a.rolname = ANY($1)
|
||||
ORDER BY 1;`
|
||||
|
||||
func (c *Cluster) pgConnectionString() string {
|
||||
hostname := fmt.Sprintf("%s.%s.svc.cluster.local", c.Metadata.Name, c.Metadata.Namespace)
|
||||
password := c.pgUsers[c.OpConfig.SuperUsername].Password
|
||||
username := c.systemUsers[constants.SuperuserKeyName].Name
|
||||
password := c.systemUsers[constants.SuperuserKeyName].Password
|
||||
|
||||
return fmt.Sprintf("host='%s' dbname=postgres sslmode=require user='%s' password='%s'",
|
||||
hostname,
|
||||
c.OpConfig.SuperUsername,
|
||||
username,
|
||||
strings.Replace(password, "$", "\\$", -1))
|
||||
}
|
||||
|
||||
@@ -33,6 +43,7 @@ func (c *Cluster) initDbConn() error {
|
||||
}
|
||||
err = conn.Ping()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -43,42 +54,48 @@ func (c *Cluster) initDbConn() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) createPgUser(user spec.PgUser) (isHuman bool, err error) {
|
||||
var flags []string = user.Flags
|
||||
|
||||
if user.Password == "" {
|
||||
isHuman = true
|
||||
flags = append(flags, "SUPERUSER")
|
||||
flags = append(flags, fmt.Sprintf("IN ROLE \"%s\"", c.OpConfig.PamRoleName))
|
||||
} else {
|
||||
isHuman = false
|
||||
func (c *Cluster) readPgUsersFromDatabase(userNames []string) (users spec.PgUserMap, err error) {
|
||||
var rows *sql.Rows
|
||||
users = make(spec.PgUserMap)
|
||||
if rows, err = c.pgDb.Query(getUserSQL, pq.Array(userNames)); err != nil {
|
||||
return nil, fmt.Errorf("Error when querying users: %s", err)
|
||||
}
|
||||
|
||||
addLoginFlag := true
|
||||
for _, v := range flags {
|
||||
if v == "NOLOGIN" {
|
||||
addLoginFlag = false
|
||||
break
|
||||
defer rows.Close()
|
||||
for rows.Next() {
|
||||
var (
|
||||
rolname, rolpassword string
|
||||
rolsuper, rolinherit, rolcreaterole, rolcreatedb, rolcanlogin bool
|
||||
memberof []string
|
||||
)
|
||||
err := rows.Scan(&rolname, &rolpassword, &rolsuper, &rolinherit,
|
||||
&rolcreaterole, &rolcreatedb, &rolcanlogin, pq.Array(&memberof))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("Error when processing user rows: %s", err)
|
||||
}
|
||||
}
|
||||
if addLoginFlag {
|
||||
flags = append(flags, "LOGIN")
|
||||
}
|
||||
if !isHuman && user.MemberOf != "" {
|
||||
flags = append(flags, fmt.Sprintf("IN ROLE \"%s\"", user.MemberOf))
|
||||
}
|
||||
userFlags := strings.Join(flags, " ")
|
||||
userPassword := fmt.Sprintf("ENCRYPTED PASSWORD '%s'", util.PGUserPassword(user))
|
||||
if user.Password == "" {
|
||||
userPassword = "PASSWORD NULL"
|
||||
}
|
||||
query := fmt.Sprintf(createUserSQL, user.Name, userFlags, userPassword)
|
||||
|
||||
_, err = c.pgDb.Query(query) // TODO: Try several times
|
||||
if err != nil {
|
||||
err = fmt.Errorf("DB error: %s", err)
|
||||
return
|
||||
flags := makeUserFlags(rolsuper, rolinherit, rolcreaterole, rolcreatedb, rolcanlogin)
|
||||
// XXX: the code assumes the password we get from pg_authid is always MD5
|
||||
users[rolname] = spec.PgUser{Name: rolname, Password: rolpassword, Flags: flags, MemberOf: memberof}
|
||||
}
|
||||
|
||||
return
|
||||
return users, nil
|
||||
}
|
||||
|
||||
func makeUserFlags(rolsuper, rolinherit, rolcreaterole, rolcreatedb, rolcanlogin bool) (result []string) {
|
||||
if rolsuper {
|
||||
result = append(result, constants.RoleFlagSuperuser)
|
||||
}
|
||||
if rolinherit {
|
||||
result = append(result, constants.RoleFlagInherit)
|
||||
}
|
||||
if rolcreaterole {
|
||||
result = append(result, constants.RoleFlagCreateRole)
|
||||
}
|
||||
if rolcreatedb {
|
||||
result = append(result, constants.RoleFlagCreateDB)
|
||||
}
|
||||
if rolcanlogin {
|
||||
result = append(result, constants.RoleFlagLogin)
|
||||
}
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
+22
-25
@@ -6,9 +6,11 @@ import (
|
||||
"k8s.io/client-go/pkg/api"
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
"k8s.io/client-go/pkg/apis/apps/v1beta1"
|
||||
|
||||
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/k8sutil"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/spec"
|
||||
"github.bus.zalan.do/acid/postgres-operator/pkg/util/constants"
|
||||
)
|
||||
|
||||
func (c *Cluster) loadResources() error {
|
||||
@@ -183,7 +185,7 @@ func (c *Cluster) createService() (*v1.Service, error) {
|
||||
return service, nil
|
||||
}
|
||||
|
||||
func (c *Cluster) updateService(newService *v1.Service) error {
|
||||
func (c *Cluster) updateService(newService *v1.Service) error {
|
||||
if c.Service == nil {
|
||||
return fmt.Errorf("There is no Service in the cluster")
|
||||
}
|
||||
@@ -262,23 +264,29 @@ func (c *Cluster) deleteEndpoint() error {
|
||||
}
|
||||
|
||||
func (c *Cluster) applySecrets() error {
|
||||
secrets, err := c.genUserSecrets()
|
||||
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get user Secrets")
|
||||
}
|
||||
secrets := c.genUserSecrets()
|
||||
|
||||
for secretUsername, secretSpec := range secrets {
|
||||
secret, err := c.KubeClient.Secrets(secretSpec.Namespace).Create(secretSpec)
|
||||
if k8sutil.ResourceAlreadyExists(err) {
|
||||
var userMap map[string]spec.PgUser
|
||||
curSecret, err := c.KubeClient.Secrets(secretSpec.Namespace).Get(secretSpec.Name)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Can't get current Secret: %s", err)
|
||||
}
|
||||
c.logger.Debugf("Secret '%s' already exists, fetching it's password", util.NameFromMeta(curSecret.ObjectMeta))
|
||||
pwdUser := c.pgUsers[secretUsername]
|
||||
if secretUsername == c.systemUsers[constants.SuperuserKeyName].Name {
|
||||
secretUsername = constants.SuperuserKeyName
|
||||
userMap = c.systemUsers
|
||||
} else if secretUsername == c.systemUsers[constants.ReplicationUserKeyName].Name {
|
||||
secretUsername = constants.ReplicationUserKeyName
|
||||
userMap = c.systemUsers
|
||||
} else {
|
||||
userMap = c.pgUsers
|
||||
}
|
||||
pwdUser := userMap[secretUsername]
|
||||
pwdUser.Password = string(curSecret.Data["password"])
|
||||
c.pgUsers[secretUsername] = pwdUser
|
||||
userMap[secretUsername] = pwdUser
|
||||
|
||||
continue
|
||||
} else {
|
||||
@@ -305,23 +313,12 @@ func (c *Cluster) deleteSecret(secret *v1.Secret) error {
|
||||
return err
|
||||
}
|
||||
|
||||
func (c *Cluster) createUsers() error {
|
||||
func (c *Cluster) createUsers() (err error) {
|
||||
// TODO: figure out what to do with duplicate names (humans and robots) among pgUsers
|
||||
for username, user := range c.pgUsers {
|
||||
if username == c.OpConfig.SuperUsername || username == c.OpConfig.ReplicationUsername {
|
||||
continue
|
||||
}
|
||||
|
||||
isHuman, err := c.createPgUser(user)
|
||||
var userType string
|
||||
if isHuman {
|
||||
userType = "human"
|
||||
} else {
|
||||
userType = "robot"
|
||||
}
|
||||
if err != nil {
|
||||
c.logger.Warnf("Can't create %s user '%s': %s", userType, username, err)
|
||||
}
|
||||
reqs := c.userSyncStrategy.ProduceSyncRequests(nil, c.pgUsers)
|
||||
err = c.userSyncStrategy.ExecuteSyncRequests(reqs, c.pgDb)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -36,6 +36,14 @@ func (c *Cluster) SyncCluster(stopCh <-chan struct{}) {
|
||||
if err := c.syncStatefulSet(); err != nil {
|
||||
c.logger.Errorf("Can't sync StatefulSets: %s", err)
|
||||
}
|
||||
if err := c.initDbConn(); err != nil {
|
||||
c.logger.Errorf("Can't init db connection: %s", err)
|
||||
} else {
|
||||
c.logger.Debugf("Syncing Roles")
|
||||
if err := c.SyncRoles(); err != nil {
|
||||
c.logger.Errorf("Can't sync Roles: %s", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Cluster) syncSecrets() error {
|
||||
@@ -150,3 +158,23 @@ func (c *Cluster) syncStatefulSet() error {
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) SyncRoles() error {
|
||||
var userNames []string
|
||||
|
||||
if err := c.initUsers(); err != nil {
|
||||
return err
|
||||
}
|
||||
for _, u := range c.pgUsers {
|
||||
userNames = append(userNames, u.Name)
|
||||
}
|
||||
dbUsers, err := c.readPgUsersFromDatabase(userNames)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error getting users from the database: %s", err)
|
||||
}
|
||||
pgSyncRequests := c.userSyncStrategy.ProduceSyncRequests(dbUsers, c.pgUsers)
|
||||
if err := c.userSyncStrategy.ExecuteSyncRequests(pgSyncRequests, c.pgDb); err != nil {
|
||||
return fmt.Errorf("Error executing sync statements: %s", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -24,6 +24,7 @@ func isValidUsername(username string) bool {
|
||||
|
||||
func normalizeUserFlags(userFlags []string) (flags []string, err error) {
|
||||
uniqueFlags := make(map[string]bool)
|
||||
addLogin := true
|
||||
|
||||
for _, flag := range userFlags {
|
||||
if !alphaNumericRegexp.MatchString(flag) {
|
||||
@@ -36,11 +37,25 @@ func normalizeUserFlags(userFlags []string) (flags []string, err error) {
|
||||
}
|
||||
}
|
||||
}
|
||||
if uniqueFlags[constants.RoleFlagLogin] && uniqueFlags[constants.RoleFlagNoLogin] {
|
||||
return nil, fmt.Errorf("Conflicting or redundant flags: LOGIN and NOLOGIN")
|
||||
}
|
||||
|
||||
flags = []string{}
|
||||
for k := range uniqueFlags {
|
||||
if k == constants.RoleFlagNoLogin || k == constants.RoleFlagLogin {
|
||||
addLogin = false
|
||||
if k == constants.RoleFlagNoLogin {
|
||||
// we don't add NOLOGIN to the list of flags to be consistent with what we get
|
||||
// from the readPgUsersFromDatabase in SyncUsers
|
||||
continue
|
||||
}
|
||||
}
|
||||
flags = append(flags, k)
|
||||
}
|
||||
if addLogin {
|
||||
flags = append(flags, constants.RoleFlagLogin)
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user