mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-05 10:21:45 +02:00
merge with master
This commit is contained in:
@@ -1,9 +1,12 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -13,6 +16,7 @@ import (
|
||||
"github.com/zalando/postgres-operator/pkg/cluster"
|
||||
acidv1informer "github.com/zalando/postgres-operator/pkg/generated/informers/externalversions/acid.zalan.do/v1"
|
||||
"github.com/zalando/postgres-operator/pkg/spec"
|
||||
"github.com/zalando/postgres-operator/pkg/teams"
|
||||
"github.com/zalando/postgres-operator/pkg/util"
|
||||
"github.com/zalando/postgres-operator/pkg/util/config"
|
||||
"github.com/zalando/postgres-operator/pkg/util/constants"
|
||||
@@ -31,8 +35,9 @@ import (
|
||||
|
||||
// Controller represents operator controller
|
||||
type Controller struct {
|
||||
config spec.ControllerConfig
|
||||
opConfig *config.Config
|
||||
config spec.ControllerConfig
|
||||
opConfig *config.Config
|
||||
pgTeamMap teams.PostgresTeamMap
|
||||
|
||||
logger *logrus.Entry
|
||||
KubeClient k8sutil.KubernetesClient
|
||||
@@ -53,10 +58,11 @@ type Controller struct {
|
||||
clusterHistory map[spec.NamespacedName]ringlog.RingLogger // history of the cluster changes
|
||||
teamClusters map[string][]spec.NamespacedName
|
||||
|
||||
postgresqlInformer cache.SharedIndexInformer
|
||||
podInformer cache.SharedIndexInformer
|
||||
nodesInformer cache.SharedIndexInformer
|
||||
podCh chan cluster.PodEvent
|
||||
postgresqlInformer cache.SharedIndexInformer
|
||||
postgresTeamInformer cache.SharedIndexInformer
|
||||
podInformer cache.SharedIndexInformer
|
||||
nodesInformer cache.SharedIndexInformer
|
||||
podCh chan cluster.PodEvent
|
||||
|
||||
clusterEventQueues []*cache.FIFO // [workerID]Queue
|
||||
lastClusterSyncTime int64
|
||||
@@ -71,6 +77,13 @@ type Controller struct {
|
||||
// NewController creates a new controller
|
||||
func NewController(controllerConfig *spec.ControllerConfig, controllerId string) *Controller {
|
||||
logger := logrus.New()
|
||||
if controllerConfig.EnableJsonLogging {
|
||||
logger.SetFormatter(&logrus.JSONFormatter{})
|
||||
} else {
|
||||
if os.Getenv("LOG_NOQUOTE") != "" {
|
||||
logger.SetFormatter(&logrus.TextFormatter{PadLevelText: true, DisableQuote: true})
|
||||
}
|
||||
}
|
||||
|
||||
var myComponentName = "postgres-operator"
|
||||
if controllerId != "" {
|
||||
@@ -78,7 +91,10 @@ func NewController(controllerConfig *spec.ControllerConfig, controllerId string)
|
||||
}
|
||||
|
||||
eventBroadcaster := record.NewBroadcaster()
|
||||
eventBroadcaster.StartLogging(logger.Infof)
|
||||
|
||||
// disabling the sending of events also to the logoutput
|
||||
// the operator currently duplicates a lot of log entries with this setup
|
||||
// eventBroadcaster.StartLogging(logger.Infof)
|
||||
recorder := eventBroadcaster.NewRecorder(scheme.Scheme, v1.EventSource{Component: myComponentName})
|
||||
|
||||
c := &Controller{
|
||||
@@ -187,10 +203,18 @@ func (c *Controller) warnOnDeprecatedOperatorParameters() {
|
||||
}
|
||||
}
|
||||
|
||||
func compactValue(v string) string {
|
||||
var compact bytes.Buffer
|
||||
if err := json.Compact(&compact, []byte(v)); err != nil {
|
||||
panic("Hard coded json strings broken!")
|
||||
}
|
||||
return compact.String()
|
||||
}
|
||||
|
||||
func (c *Controller) initPodServiceAccount() {
|
||||
|
||||
if c.opConfig.PodServiceAccountDefinition == "" {
|
||||
c.opConfig.PodServiceAccountDefinition = `
|
||||
stringValue := `
|
||||
{
|
||||
"apiVersion": "v1",
|
||||
"kind": "ServiceAccount",
|
||||
@@ -198,6 +222,9 @@ func (c *Controller) initPodServiceAccount() {
|
||||
"name": "postgres-pod"
|
||||
}
|
||||
}`
|
||||
|
||||
c.opConfig.PodServiceAccountDefinition = compactValue(stringValue)
|
||||
|
||||
}
|
||||
|
||||
// re-uses k8s internal parsing. See k8s client-go issue #193 for explanation
|
||||
@@ -227,7 +254,7 @@ func (c *Controller) initRoleBinding() {
|
||||
// operator binds it to the cluster role with sufficient privileges
|
||||
// we assume the role is created by the k8s administrator
|
||||
if c.opConfig.PodServiceAccountRoleBindingDefinition == "" {
|
||||
c.opConfig.PodServiceAccountRoleBindingDefinition = fmt.Sprintf(`
|
||||
stringValue := fmt.Sprintf(`
|
||||
{
|
||||
"apiVersion": "rbac.authorization.k8s.io/v1",
|
||||
"kind": "RoleBinding",
|
||||
@@ -246,6 +273,7 @@ func (c *Controller) initRoleBinding() {
|
||||
}
|
||||
]
|
||||
}`, c.PodServiceAccount.Name, c.PodServiceAccount.Name, c.PodServiceAccount.Name)
|
||||
c.opConfig.PodServiceAccountRoleBindingDefinition = compactValue(stringValue)
|
||||
}
|
||||
c.logger.Info("Parse role bindings")
|
||||
// re-uses k8s internal parsing. See k8s client-go issue #193 for explanation
|
||||
@@ -264,7 +292,14 @@ func (c *Controller) initRoleBinding() {
|
||||
|
||||
}
|
||||
|
||||
// actual roles bindings are deployed at the time of Postgres/Spilo cluster creation
|
||||
// actual roles bindings ar*logrus.Entrye deployed at the time of Postgres/Spilo cluster creation
|
||||
}
|
||||
|
||||
func logMultiLineConfig(log *logrus.Entry, config string) {
|
||||
lines := strings.Split(config, "\n")
|
||||
for _, l := range lines {
|
||||
log.Infof("%s", l)
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Controller) initController() {
|
||||
@@ -292,14 +327,19 @@ func (c *Controller) initController() {
|
||||
c.logger.Fatalf("could not register Postgres CustomResourceDefinition: %v", err)
|
||||
}
|
||||
|
||||
c.initPodServiceAccount()
|
||||
c.initSharedInformers()
|
||||
|
||||
if c.opConfig.EnablePostgresTeamCRD {
|
||||
c.loadPostgresTeams()
|
||||
} else {
|
||||
c.pgTeamMap = teams.PostgresTeamMap{}
|
||||
}
|
||||
|
||||
if c.opConfig.DebugLogging {
|
||||
c.logger.Logger.Level = logrus.DebugLevel
|
||||
}
|
||||
|
||||
c.logger.Infof("config: %s", c.opConfig.MustMarshal())
|
||||
logMultiLineConfig(c.logger, c.opConfig.MustMarshal())
|
||||
|
||||
roleDefs := c.getInfrastructureRoleDefinitions()
|
||||
if infraRoles, err := c.getInfrastructureRoles(roleDefs); err != nil {
|
||||
@@ -326,6 +366,7 @@ func (c *Controller) initController() {
|
||||
|
||||
func (c *Controller) initSharedInformers() {
|
||||
|
||||
// Postgresqls
|
||||
c.postgresqlInformer = acidv1informer.NewPostgresqlInformer(
|
||||
c.KubeClient.AcidV1ClientSet,
|
||||
c.opConfig.WatchedNamespace,
|
||||
@@ -338,6 +379,20 @@ func (c *Controller) initSharedInformers() {
|
||||
DeleteFunc: c.postgresqlDelete,
|
||||
})
|
||||
|
||||
// PostgresTeams
|
||||
if c.opConfig.EnablePostgresTeamCRD {
|
||||
c.postgresTeamInformer = acidv1informer.NewPostgresTeamInformer(
|
||||
c.KubeClient.AcidV1ClientSet,
|
||||
c.opConfig.WatchedNamespace,
|
||||
constants.QueueResyncPeriodTPR*6, // 30 min
|
||||
cache.Indexers{})
|
||||
|
||||
c.postgresTeamInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
AddFunc: c.postgresTeamAdd,
|
||||
UpdateFunc: c.postgresTeamUpdate,
|
||||
})
|
||||
}
|
||||
|
||||
// Pods
|
||||
podLw := &cache.ListWatch{
|
||||
ListFunc: c.podListFunc,
|
||||
@@ -398,6 +453,10 @@ func (c *Controller) Run(stopCh <-chan struct{}, wg *sync.WaitGroup) {
|
||||
go c.apiserver.Run(stopCh, wg)
|
||||
go c.kubeNodesInformer(stopCh, wg)
|
||||
|
||||
if c.opConfig.EnablePostgresTeamCRD {
|
||||
go c.runPostgresTeamInformer(stopCh, wg)
|
||||
}
|
||||
|
||||
c.logger.Info("started working in background")
|
||||
}
|
||||
|
||||
@@ -413,6 +472,12 @@ func (c *Controller) runPostgresqlInformer(stopCh <-chan struct{}, wg *sync.Wait
|
||||
c.postgresqlInformer.Run(stopCh)
|
||||
}
|
||||
|
||||
func (c *Controller) runPostgresTeamInformer(stopCh <-chan struct{}, wg *sync.WaitGroup) {
|
||||
defer wg.Done()
|
||||
|
||||
c.postgresTeamInformer.Run(stopCh)
|
||||
}
|
||||
|
||||
func queueClusterKey(eventType EventType, uid types.UID) string {
|
||||
return fmt.Sprintf("%s-%s", eventType, uid)
|
||||
}
|
||||
|
||||
@@ -42,7 +42,7 @@ func (c *Controller) nodeAdd(obj interface{}) {
|
||||
return
|
||||
}
|
||||
|
||||
c.logger.Debugf("new node has been added: %q (%s)", util.NameFromMeta(node.ObjectMeta), node.Spec.ProviderID)
|
||||
c.logger.Debugf("new node has been added: %s (%s)", util.NameFromMeta(node.ObjectMeta), node.Spec.ProviderID)
|
||||
|
||||
// check if the node became not ready while the operator was down (otherwise we would have caught it in nodeUpdate)
|
||||
if !c.nodeIsReady(node) {
|
||||
@@ -76,7 +76,7 @@ func (c *Controller) nodeUpdate(prev, cur interface{}) {
|
||||
}
|
||||
|
||||
func (c *Controller) nodeIsReady(node *v1.Node) bool {
|
||||
return (!node.Spec.Unschedulable || util.MapContains(node.Labels, c.opConfig.NodeReadinessLabel) ||
|
||||
return (!node.Spec.Unschedulable || (len(c.opConfig.NodeReadinessLabel) > 0 && util.MapContains(node.Labels, c.opConfig.NodeReadinessLabel)) ||
|
||||
util.MapContains(node.Labels, map[string]string{"master": "true"}))
|
||||
}
|
||||
|
||||
|
||||
+41
-11
@@ -15,7 +15,6 @@ const (
|
||||
|
||||
func newNodeTestController() *Controller {
|
||||
var controller = NewController(&spec.ControllerConfig{}, "node-test")
|
||||
controller.opConfig.NodeReadinessLabel = map[string]string{readyLabel: readyValue}
|
||||
return controller
|
||||
}
|
||||
|
||||
@@ -36,27 +35,58 @@ var nodeTestController = newNodeTestController()
|
||||
func TestNodeIsReady(t *testing.T) {
|
||||
testName := "TestNodeIsReady"
|
||||
var testTable = []struct {
|
||||
in *v1.Node
|
||||
out bool
|
||||
in *v1.Node
|
||||
out bool
|
||||
readinessLabel map[string]string
|
||||
}{
|
||||
{
|
||||
in: makeNode(map[string]string{"foo": "bar"}, true),
|
||||
out: true,
|
||||
in: makeNode(map[string]string{"foo": "bar"}, true),
|
||||
out: true,
|
||||
readinessLabel: map[string]string{readyLabel: readyValue},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{"foo": "bar"}, false),
|
||||
out: false,
|
||||
in: makeNode(map[string]string{"foo": "bar"}, false),
|
||||
out: false,
|
||||
readinessLabel: map[string]string{readyLabel: readyValue},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{readyLabel: readyValue}, false),
|
||||
out: true,
|
||||
in: makeNode(map[string]string{readyLabel: readyValue}, false),
|
||||
out: true,
|
||||
readinessLabel: map[string]string{readyLabel: readyValue},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{"foo": "bar", "master": "true"}, false),
|
||||
out: true,
|
||||
in: makeNode(map[string]string{"foo": "bar", "master": "true"}, false),
|
||||
out: true,
|
||||
readinessLabel: map[string]string{readyLabel: readyValue},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{"foo": "bar", "master": "true"}, false),
|
||||
out: true,
|
||||
readinessLabel: map[string]string{readyLabel: readyValue},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{"foo": "bar"}, true),
|
||||
out: true,
|
||||
readinessLabel: map[string]string{},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{"foo": "bar"}, false),
|
||||
out: false,
|
||||
readinessLabel: map[string]string{},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{readyLabel: readyValue}, false),
|
||||
out: false,
|
||||
readinessLabel: map[string]string{},
|
||||
},
|
||||
{
|
||||
in: makeNode(map[string]string{"foo": "bar", "master": "true"}, false),
|
||||
out: true,
|
||||
readinessLabel: map[string]string{},
|
||||
},
|
||||
}
|
||||
for _, tt := range testTable {
|
||||
nodeTestController.opConfig.NodeReadinessLabel = tt.readinessLabel
|
||||
if isReady := nodeTestController.nodeIsReady(tt.in); isReady != tt.out {
|
||||
t.Errorf("%s: expected response %t doesn't match the actual %t for the node %#v",
|
||||
testName, tt.out, isReady, tt.in)
|
||||
|
||||
@@ -163,6 +163,8 @@ func (c *Controller) importConfigurationFromCRD(fromCRD *acidv1.OperatorConfigur
|
||||
result.PamConfiguration = util.Coalesce(fromCRD.TeamsAPI.PamConfiguration, "https://info.example.com/oauth2/tokeninfo?access_token= uid realm=/employees")
|
||||
result.ProtectedRoles = util.CoalesceStrArr(fromCRD.TeamsAPI.ProtectedRoles, []string{"admin"})
|
||||
result.PostgresSuperuserTeams = fromCRD.TeamsAPI.PostgresSuperuserTeams
|
||||
result.EnablePostgresTeamCRD = fromCRD.TeamsAPI.EnablePostgresTeamCRD
|
||||
result.EnablePostgresTeamCRDSuperusers = fromCRD.TeamsAPI.EnablePostgresTeamCRDSuperusers
|
||||
|
||||
// logging REST API config
|
||||
result.APIPort = util.CoalesceInt(fromCRD.LoggingRESTAPI.APIPort, 8080)
|
||||
|
||||
@@ -225,7 +225,7 @@ func (c *Controller) processEvent(event ClusterEvent) {
|
||||
switch event.EventType {
|
||||
case EventAdd:
|
||||
if clusterFound {
|
||||
lg.Debugf("cluster already exists")
|
||||
lg.Debugf("Recieved add event for existing cluster")
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -15,6 +15,7 @@ import (
|
||||
acidv1 "github.com/zalando/postgres-operator/pkg/apis/acid.zalan.do/v1"
|
||||
"github.com/zalando/postgres-operator/pkg/cluster"
|
||||
"github.com/zalando/postgres-operator/pkg/spec"
|
||||
"github.com/zalando/postgres-operator/pkg/teams"
|
||||
"github.com/zalando/postgres-operator/pkg/util"
|
||||
"github.com/zalando/postgres-operator/pkg/util/config"
|
||||
"github.com/zalando/postgres-operator/pkg/util/k8sutil"
|
||||
@@ -30,6 +31,7 @@ func (c *Controller) makeClusterConfig() cluster.Config {
|
||||
return cluster.Config{
|
||||
RestConfig: c.config.RestConfig,
|
||||
OpConfig: config.Copy(c.opConfig),
|
||||
PgTeamMap: c.pgTeamMap,
|
||||
InfrastructureRoles: infrastructureRoles,
|
||||
PodServiceAccount: c.PodServiceAccount,
|
||||
}
|
||||
@@ -394,6 +396,37 @@ func (c *Controller) getInfrastructureRole(
|
||||
return roles, nil
|
||||
}
|
||||
|
||||
func (c *Controller) loadPostgresTeams() {
|
||||
// reset team map
|
||||
c.pgTeamMap = teams.PostgresTeamMap{}
|
||||
|
||||
pgTeams, err := c.KubeClient.AcidV1ClientSet.AcidV1().PostgresTeams(c.opConfig.WatchedNamespace).List(context.TODO(), metav1.ListOptions{})
|
||||
if err != nil {
|
||||
c.logger.Errorf("could not list postgres team objects: %v", err)
|
||||
}
|
||||
|
||||
c.pgTeamMap.Load(pgTeams)
|
||||
c.logger.Debugf("Internal Postgres Team Cache: %#v", c.pgTeamMap)
|
||||
}
|
||||
|
||||
func (c *Controller) postgresTeamAdd(obj interface{}) {
|
||||
pgTeam, ok := obj.(*acidv1.PostgresTeam)
|
||||
if !ok {
|
||||
c.logger.Errorf("could not cast to PostgresTeam spec")
|
||||
}
|
||||
c.logger.Debugf("PostgreTeam %q added. Reloading postgres team CRDs and overwriting cached map", pgTeam.Name)
|
||||
c.loadPostgresTeams()
|
||||
}
|
||||
|
||||
func (c *Controller) postgresTeamUpdate(prev, obj interface{}) {
|
||||
pgTeam, ok := obj.(*acidv1.PostgresTeam)
|
||||
if !ok {
|
||||
c.logger.Errorf("could not cast to PostgresTeam spec")
|
||||
}
|
||||
c.logger.Debugf("PostgreTeam %q updated. Reloading postgres team CRDs and overwriting cached map", pgTeam.Name)
|
||||
c.loadPostgresTeams()
|
||||
}
|
||||
|
||||
func (c *Controller) podClusterName(pod *v1.Pod) spec.NamespacedName {
|
||||
if name, ok := pod.Labels[c.opConfig.ClusterNameLabel]; ok {
|
||||
return spec.NamespacedName{
|
||||
|
||||
Reference in New Issue
Block a user