mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-08 23:06:45 +02:00
Merge branch 'client-go-v4' into fix/graceful-shutdown
# Conflicts: # pkg/controller/pod.go
This commit is contained in:
@@ -102,13 +102,11 @@ func (c *Controller) initController() {
|
||||
constants.QueueResyncPeriodTPR,
|
||||
cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc})
|
||||
|
||||
if err := c.postgresqlInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
c.postgresqlInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
AddFunc: c.postgresqlAdd,
|
||||
UpdateFunc: c.postgresqlUpdate,
|
||||
DeleteFunc: c.postgresqlDelete,
|
||||
}); err != nil {
|
||||
c.logger.Fatalf("could not add event handlers: %v", err)
|
||||
}
|
||||
})
|
||||
|
||||
// Pods
|
||||
podLw := &cache.ListWatch{
|
||||
@@ -122,13 +120,11 @@ func (c *Controller) initController() {
|
||||
constants.QueueResyncPeriodPod,
|
||||
cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc})
|
||||
|
||||
if err := c.podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
c.podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
AddFunc: c.podAdd,
|
||||
UpdateFunc: c.podUpdate,
|
||||
DeleteFunc: c.podDelete,
|
||||
}); err != nil {
|
||||
c.logger.Fatalf("could not add event handlers: %v", err)
|
||||
}
|
||||
})
|
||||
|
||||
c.clusterEventQueues = make([]*cache.FIFO, c.opConfig.Workers)
|
||||
for i := range c.clusterEventQueues {
|
||||
|
||||
+7
-22
@@ -3,27 +3,20 @@ package controller
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"k8s.io/client-go/pkg/api"
|
||||
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/watch"
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
"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/util"
|
||||
)
|
||||
|
||||
func (c *Controller) podListFunc(options api.ListOptions) (runtime.Object, error) {
|
||||
func (c *Controller) podListFunc(options meta_v1.ListOptions) (runtime.Object, error) {
|
||||
var labelSelector string
|
||||
var fieldSelector string
|
||||
|
||||
if options.LabelSelector != nil {
|
||||
labelSelector = options.LabelSelector.String()
|
||||
}
|
||||
|
||||
if options.FieldSelector != nil {
|
||||
fieldSelector = options.FieldSelector.String()
|
||||
}
|
||||
opts := v1.ListOptions{
|
||||
opts := meta_v1.ListOptions{
|
||||
LabelSelector: labelSelector,
|
||||
FieldSelector: fieldSelector,
|
||||
Watch: options.Watch,
|
||||
@@ -34,19 +27,11 @@ func (c *Controller) podListFunc(options api.ListOptions) (runtime.Object, error
|
||||
return c.KubeClient.CoreV1().Pods(c.opConfig.Namespace).List(opts)
|
||||
}
|
||||
|
||||
func (c *Controller) podWatchFunc(options api.ListOptions) (watch.Interface, error) {
|
||||
func (c *Controller) podWatchFunc(options meta_v1.ListOptions) (watch.Interface, error) {
|
||||
var labelSelector string
|
||||
var fieldSelector string
|
||||
|
||||
if options.LabelSelector != nil {
|
||||
labelSelector = options.LabelSelector.String()
|
||||
}
|
||||
|
||||
if options.FieldSelector != nil {
|
||||
fieldSelector = options.FieldSelector.String()
|
||||
}
|
||||
|
||||
opts := v1.ListOptions{
|
||||
opts := meta_v1.ListOptions{
|
||||
LabelSelector: labelSelector,
|
||||
FieldSelector: fieldSelector,
|
||||
Watch: options.Watch,
|
||||
|
||||
@@ -7,12 +7,13 @@ import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"k8s.io/apimachinery/pkg/api/meta"
|
||||
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/fields"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
"k8s.io/apimachinery/pkg/watch"
|
||||
"k8s.io/client-go/pkg/api"
|
||||
"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"
|
||||
"k8s.io/client-go/pkg/watch"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/cluster"
|
||||
@@ -28,14 +29,14 @@ func (c *Controller) clusterResync(stopCh <-chan struct{}, wg *sync.WaitGroup) {
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
c.clusterListFunc(api.ListOptions{ResourceVersion: "0"})
|
||||
c.clusterListFunc(meta_v1.ListOptions{ResourceVersion: "0"})
|
||||
case <-stopCh:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Controller) clusterListFunc(options api.ListOptions) (runtime.Object, error) {
|
||||
func (c *Controller) clusterListFunc(options meta_v1.ListOptions) (runtime.Object, error) {
|
||||
c.logger.Info("Getting list of currently running clusters")
|
||||
|
||||
req := c.RestClient.Get().
|
||||
@@ -90,7 +91,7 @@ func (c *Controller) clusterListFunc(options api.ListOptions) (runtime.Object, e
|
||||
return object, err
|
||||
}
|
||||
|
||||
func (c *Controller) clusterWatchFunc(options api.ListOptions) (watch.Interface, error) {
|
||||
func (c *Controller) clusterWatchFunc(options meta_v1.ListOptions) (watch.Interface, error) {
|
||||
req := c.RestClient.Get().
|
||||
RequestURI(fmt.Sprintf(constants.WatchClustersURITemplate, c.opConfig.Namespace)).
|
||||
VersionedParams(&options, api.ParameterCodec).
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"fmt"
|
||||
"hash/crc32"
|
||||
|
||||
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
extv1beta "k8s.io/client-go/pkg/apis/extensions/v1beta1"
|
||||
|
||||
@@ -32,7 +33,7 @@ func (c *Controller) makeClusterConfig() cluster.Config {
|
||||
|
||||
func thirdPartyResource(TPRName string) *extv1beta.ThirdPartyResource {
|
||||
return &extv1beta.ThirdPartyResource{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
ObjectMeta: meta_v1.ObjectMeta{
|
||||
//ThirdPartyResources are cluster-wide
|
||||
Name: TPRName,
|
||||
},
|
||||
@@ -72,7 +73,7 @@ func (c *Controller) getInfrastructureRoles() (result map[string]spec.PgUser, er
|
||||
|
||||
infraRolesSecret, err := c.KubeClient.
|
||||
Secrets(c.opConfig.InfrastructureRolesSecretName.Namespace).
|
||||
Get(c.opConfig.InfrastructureRolesSecretName.Name)
|
||||
Get(c.opConfig.InfrastructureRolesSecretName.Name, meta_v1.GetOptions{})
|
||||
if err != nil {
|
||||
c.logger.Debugf("Infrastructure roles secret name: %s", c.opConfig.InfrastructureRolesSecretName)
|
||||
return nil, fmt.Errorf("could not get infrastructure roles secret: %v", err)
|
||||
|
||||
Reference in New Issue
Block a user