mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-11 16:45:40 +02:00
move types from spec into separate package
This commit is contained in:
@@ -10,8 +10,8 @@ import (
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/cluster"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/spec"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/types"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util/config"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util/constants"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util/teams"
|
||||
@@ -22,7 +22,7 @@ type Config struct {
|
||||
KubeClient *kubernetes.Clientset
|
||||
RestClient *rest.RESTClient
|
||||
TeamsAPIClient *teams.API
|
||||
InfrastructureRoles map[string]spec.PgUser
|
||||
InfrastructureRoles map[string]types.PgUser
|
||||
}
|
||||
|
||||
type Controller struct {
|
||||
@@ -31,12 +31,12 @@ type Controller struct {
|
||||
logger *logrus.Entry
|
||||
|
||||
clustersMu sync.RWMutex
|
||||
clusters map[spec.NamespacedName]cluster.Interface
|
||||
stopChs map[spec.NamespacedName]chan struct{}
|
||||
clusters map[types.NamespacedName]types.Cluster
|
||||
stopChs map[types.NamespacedName]chan struct{}
|
||||
|
||||
postgresqlInformer cache.SharedIndexInformer
|
||||
podInformer cache.SharedIndexInformer
|
||||
podCh chan spec.PodEvent
|
||||
podCh chan types.PodEvent
|
||||
|
||||
clusterEventQueues []*cache.FIFO
|
||||
|
||||
@@ -56,9 +56,9 @@ func New(controllerConfig *Config, operatorConfig *config.Config) *Controller {
|
||||
Config: *controllerConfig,
|
||||
opConfig: operatorConfig,
|
||||
logger: logger.WithField("pkg", "controller"),
|
||||
clusters: make(map[spec.NamespacedName]cluster.Interface),
|
||||
stopChs: make(map[spec.NamespacedName]chan struct{}),
|
||||
podCh: make(chan spec.PodEvent),
|
||||
clusters: make(map[types.NamespacedName]types.Cluster),
|
||||
stopChs: make(map[types.NamespacedName]chan struct{}),
|
||||
podCh: make(chan types.PodEvent),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -130,7 +130,7 @@ func (c *Controller) initController() {
|
||||
c.clusterEventQueues = make([]*cache.FIFO, c.opConfig.Workers)
|
||||
for i := range c.clusterEventQueues {
|
||||
c.clusterEventQueues[i] = cache.NewFIFO(func(obj interface{}) (string, error) {
|
||||
e, ok := obj.(spec.ClusterEvent)
|
||||
e, ok := obj.(types.ClusterEvent)
|
||||
if !ok {
|
||||
return "", fmt.Errorf("could not cast to ClusterEvent")
|
||||
}
|
||||
|
||||
@@ -6,7 +6,7 @@ import (
|
||||
"k8s.io/client-go/pkg/runtime"
|
||||
"k8s.io/client-go/pkg/watch"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/spec"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/types"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util"
|
||||
)
|
||||
|
||||
@@ -61,11 +61,11 @@ func (c *Controller) podAdd(obj interface{}) {
|
||||
return
|
||||
}
|
||||
|
||||
podEvent := spec.PodEvent{
|
||||
podEvent := types.PodEvent{
|
||||
ClusterName: c.podClusterName(pod),
|
||||
PodName: util.NameFromMeta(pod.ObjectMeta),
|
||||
CurPod: pod,
|
||||
EventType: spec.EventAdd,
|
||||
EventType: types.EventAdd,
|
||||
ResourceVersion: pod.ResourceVersion,
|
||||
}
|
||||
|
||||
@@ -83,12 +83,12 @@ func (c *Controller) podUpdate(prev, cur interface{}) {
|
||||
return
|
||||
}
|
||||
|
||||
podEvent := spec.PodEvent{
|
||||
podEvent := types.PodEvent{
|
||||
ClusterName: c.podClusterName(curPod),
|
||||
PodName: util.NameFromMeta(curPod.ObjectMeta),
|
||||
PrevPod: prevPod,
|
||||
CurPod: curPod,
|
||||
EventType: spec.EventUpdate,
|
||||
EventType: types.EventUpdate,
|
||||
ResourceVersion: curPod.ResourceVersion,
|
||||
}
|
||||
|
||||
@@ -101,11 +101,11 @@ func (c *Controller) podDelete(obj interface{}) {
|
||||
return
|
||||
}
|
||||
|
||||
podEvent := spec.PodEvent{
|
||||
podEvent := types.PodEvent{
|
||||
ClusterName: c.podClusterName(pod),
|
||||
PodName: util.NameFromMeta(pod.ObjectMeta),
|
||||
CurPod: pod,
|
||||
EventType: spec.EventDelete,
|
||||
EventType: types.EventDelete,
|
||||
ResourceVersion: pod.ResourceVersion,
|
||||
}
|
||||
|
||||
|
||||
@@ -10,12 +10,13 @@ import (
|
||||
"k8s.io/client-go/pkg/api/meta"
|
||||
"k8s.io/client-go/pkg/fields"
|
||||
"k8s.io/client-go/pkg/runtime"
|
||||
"k8s.io/client-go/pkg/types"
|
||||
ktypes "k8s.io/client-go/pkg/types"
|
||||
"k8s.io/client-go/pkg/watch"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/cluster"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/spec"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/types"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util/constants"
|
||||
)
|
||||
@@ -68,7 +69,7 @@ func (c *Controller) clusterListFunc(options api.ListOptions) (runtime.Object, e
|
||||
failedClustersCnt++
|
||||
continue
|
||||
}
|
||||
c.queueClusterEvent(nil, pg, spec.EventSync)
|
||||
c.queueClusterEvent(nil, pg, types.EventSync)
|
||||
activeClustersCnt++
|
||||
}
|
||||
if len(objList) > 0 {
|
||||
@@ -97,15 +98,15 @@ func (c *Controller) clusterWatchFunc(options api.ListOptions) (watch.Interface,
|
||||
}
|
||||
|
||||
func (c *Controller) processEvent(obj interface{}) error {
|
||||
var clusterName spec.NamespacedName
|
||||
var clusterName types.NamespacedName
|
||||
|
||||
event, ok := obj.(spec.ClusterEvent)
|
||||
event, ok := obj.(types.ClusterEvent)
|
||||
if !ok {
|
||||
return fmt.Errorf("could not cast to ClusterEvent")
|
||||
}
|
||||
logger := c.logger.WithField("worker", event.WorkerID)
|
||||
|
||||
if event.EventType == spec.EventAdd || event.EventType == spec.EventSync {
|
||||
if event.EventType == types.EventAdd || event.EventType == types.EventSync {
|
||||
clusterName = util.NameFromMeta(event.NewSpec.Metadata)
|
||||
} else {
|
||||
clusterName = util.NameFromMeta(event.OldSpec.Metadata)
|
||||
@@ -116,7 +117,7 @@ func (c *Controller) processEvent(obj interface{}) error {
|
||||
c.clustersMu.RUnlock()
|
||||
|
||||
switch event.EventType {
|
||||
case spec.EventAdd:
|
||||
case types.EventAdd:
|
||||
if clusterFound {
|
||||
logger.Debugf("Cluster '%s' already exists", clusterName)
|
||||
return nil
|
||||
@@ -142,7 +143,7 @@ func (c *Controller) processEvent(obj interface{}) error {
|
||||
}
|
||||
|
||||
logger.Infof("Cluster '%s' has been created", clusterName)
|
||||
case spec.EventUpdate:
|
||||
case types.EventUpdate:
|
||||
logger.Infof("Update of the '%s' cluster started", clusterName)
|
||||
|
||||
if !clusterFound {
|
||||
@@ -158,7 +159,7 @@ func (c *Controller) processEvent(obj interface{}) error {
|
||||
}
|
||||
cl.SetFailed(nil)
|
||||
logger.Infof("Cluster '%s' has been updated", clusterName)
|
||||
case spec.EventDelete:
|
||||
case types.EventDelete:
|
||||
logger.Infof("Deletion of the '%s' cluster started", clusterName)
|
||||
if !clusterFound {
|
||||
logger.Errorf("Unknown cluster: %s", clusterName)
|
||||
@@ -177,7 +178,7 @@ func (c *Controller) processEvent(obj interface{}) error {
|
||||
c.clustersMu.Unlock()
|
||||
|
||||
logger.Infof("Cluster '%s' has been deleted", clusterName)
|
||||
case spec.EventSync:
|
||||
case types.EventSync:
|
||||
logger.Infof("Syncing of the '%s' cluster started", clusterName)
|
||||
|
||||
// no race condition because a cluster is always processed by single worker
|
||||
@@ -214,18 +215,18 @@ func (c *Controller) processClusterEventsQueue(idx int) {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Controller) queueClusterEvent(old, new *spec.Postgresql, eventType spec.EventType) {
|
||||
func (c *Controller) queueClusterEvent(old, new *spec.Postgresql, eventType types.EventType) {
|
||||
var (
|
||||
uid types.UID
|
||||
clusterName spec.NamespacedName
|
||||
uid ktypes.UID
|
||||
clusterName types.NamespacedName
|
||||
clusterError error
|
||||
)
|
||||
|
||||
if old != nil { //update, delete
|
||||
uid = old.Metadata.GetUID()
|
||||
clusterName = util.NameFromMeta(old.Metadata)
|
||||
if eventType == spec.EventUpdate && new.Error == nil && old.Error != nil {
|
||||
eventType = spec.EventSync
|
||||
if eventType == types.EventUpdate && new.Error == nil && old.Error != nil {
|
||||
eventType = types.EventSync
|
||||
clusterError = new.Error
|
||||
} else {
|
||||
clusterError = old.Error
|
||||
@@ -236,13 +237,13 @@ func (c *Controller) queueClusterEvent(old, new *spec.Postgresql, eventType spec
|
||||
clusterError = new.Error
|
||||
}
|
||||
|
||||
if clusterError != nil && eventType != spec.EventDelete {
|
||||
if clusterError != nil && eventType != types.EventDelete {
|
||||
c.logger.Debugf("Skipping %s event for invalid cluster %s (reason: %v)", eventType, clusterName, clusterError)
|
||||
return
|
||||
}
|
||||
|
||||
workerID := c.clusterWorkerID(clusterName)
|
||||
clusterEvent := spec.ClusterEvent{
|
||||
clusterEvent := types.ClusterEvent{
|
||||
EventType: eventType,
|
||||
UID: uid,
|
||||
OldSpec: old,
|
||||
@@ -265,7 +266,7 @@ func (c *Controller) postgresqlAdd(obj interface{}) {
|
||||
}
|
||||
|
||||
// We will not get multiple Add events for the same cluster
|
||||
c.queueClusterEvent(nil, pg, spec.EventAdd)
|
||||
c.queueClusterEvent(nil, pg, types.EventAdd)
|
||||
}
|
||||
|
||||
func (c *Controller) postgresqlUpdate(prev, cur interface{}) {
|
||||
@@ -284,7 +285,7 @@ func (c *Controller) postgresqlUpdate(prev, cur interface{}) {
|
||||
return
|
||||
}
|
||||
|
||||
c.queueClusterEvent(pgOld, pgNew, spec.EventUpdate)
|
||||
c.queueClusterEvent(pgOld, pgNew, types.EventUpdate)
|
||||
}
|
||||
|
||||
func (c *Controller) postgresqlDelete(obj interface{}) {
|
||||
@@ -294,5 +295,5 @@ func (c *Controller) postgresqlDelete(obj interface{}) {
|
||||
return
|
||||
}
|
||||
|
||||
c.queueClusterEvent(pg, nil, spec.EventDelete)
|
||||
c.queueClusterEvent(pg, nil, types.EventDelete)
|
||||
}
|
||||
|
||||
+10
-10
@@ -8,14 +8,14 @@ import (
|
||||
extv1beta "k8s.io/client-go/pkg/apis/extensions/v1beta1"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/cluster"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/spec"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/types"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util/config"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util/constants"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util/k8sutil"
|
||||
)
|
||||
|
||||
func (c *Controller) makeClusterConfig() cluster.Config {
|
||||
infrastructureRoles := make(map[string]spec.PgUser)
|
||||
infrastructureRoles := make(map[string]types.PgUser)
|
||||
for k, v := range c.InfrastructureRoles {
|
||||
infrastructureRoles[k] = v
|
||||
}
|
||||
@@ -43,7 +43,7 @@ func thirdPartyResource(TPRName string) *extv1beta.ThirdPartyResource {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Controller) clusterWorkerID(clusterName spec.NamespacedName) uint32 {
|
||||
func (c *Controller) clusterWorkerID(clusterName types.NamespacedName) uint32 {
|
||||
return crc32.ChecksumIEEE([]byte(clusterName.String())) % c.opConfig.Workers
|
||||
}
|
||||
|
||||
@@ -64,8 +64,8 @@ func (c *Controller) createTPR() error {
|
||||
return k8sutil.WaitTPRReady(c.RestClient, c.opConfig.TPR.ReadyWaitInterval, c.opConfig.TPR.ReadyWaitTimeout, c.opConfig.Namespace)
|
||||
}
|
||||
|
||||
func (c *Controller) getInfrastructureRoles() (result map[string]spec.PgUser, err error) {
|
||||
if c.opConfig.InfrastructureRolesSecretName == (spec.NamespacedName{}) {
|
||||
func (c *Controller) getInfrastructureRoles() (result map[string]types.PgUser, err error) {
|
||||
if c.opConfig.InfrastructureRolesSecretName == (types.NamespacedName{}) {
|
||||
// we don't have infrastructure roles defined, bail out
|
||||
return nil, nil
|
||||
}
|
||||
@@ -79,12 +79,12 @@ func (c *Controller) getInfrastructureRoles() (result map[string]spec.PgUser, er
|
||||
}
|
||||
|
||||
data := infraRolesSecret.Data
|
||||
result = make(map[string]spec.PgUser)
|
||||
result = make(map[string]types.PgUser)
|
||||
Users:
|
||||
// in worst case we would have one line per user
|
||||
for i := 1; i <= len(data); i++ {
|
||||
properties := []string{"user", "password", "inrole"}
|
||||
t := spec.PgUser{}
|
||||
t := types.PgUser{}
|
||||
for _, p := range properties {
|
||||
key := fmt.Sprintf("%s%d", p, i)
|
||||
if val, present := data[key]; !present {
|
||||
@@ -115,13 +115,13 @@ Users:
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (c *Controller) podClusterName(pod *v1.Pod) spec.NamespacedName {
|
||||
func (c *Controller) podClusterName(pod *v1.Pod) types.NamespacedName {
|
||||
if name, ok := pod.Labels[c.opConfig.ClusterNameLabel]; ok {
|
||||
return spec.NamespacedName{
|
||||
return types.NamespacedName{
|
||||
Namespace: pod.Namespace,
|
||||
Name: name,
|
||||
}
|
||||
}
|
||||
|
||||
return spec.NamespacedName{}
|
||||
return types.NamespacedName{}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user