Merge branch 'master' into fix/graceful-shutdown

# Conflicts:
#	pkg/cluster/cluster.go
#	pkg/cluster/exec.go
#	pkg/cluster/k8sres.go
#	pkg/cluster/pod.go
#	pkg/cluster/resources.go
#	pkg/cluster/util.go
#	pkg/cluster/volumes.go
#	pkg/controller/controller.go
#	pkg/controller/pod.go
#	pkg/controller/postgresql.go
#	pkg/controller/util.go
#	pkg/controller/util_test.go
#	pkg/spec/postgresql.go
#	pkg/spec/postgresql_test.go
#	pkg/util/util.go
#	pkg/util/util_test.go
This commit is contained in:
Murat Kabilov
2017-07-25 15:41:06 +02:00
17 changed files with 106 additions and 132 deletions
+3 -3
View File
@@ -5,7 +5,7 @@ import (
"sync"
"github.com/Sirupsen/logrus"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/cache"
@@ -33,7 +33,7 @@ type Controller struct {
logger *logrus.Entry
KubeClient k8sutil.KubernetesClient
RestClient rest.Interface
RestClient rest.Interface // kubernetes API group REST client
clustersMu sync.RWMutex
clusters map[spec.NamespacedName]*cluster.Cluster
@@ -78,7 +78,7 @@ func (c *Controller) initOperatorConfig() {
if c.config.ConfigMapName != (spec.NamespacedName{}) {
configMap, err := c.KubeClient.ConfigMaps(c.config.ConfigMapName.Namespace).
Get(c.config.ConfigMapName.Name, meta_v1.GetOptions{})
Get(c.config.ConfigMapName.Name, metav1.GetOptions{})
if err != nil {
panic(err)
}
+5 -7
View File
@@ -1,9 +1,7 @@
package controller
import (
"sync"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/watch"
"k8s.io/client-go/pkg/api/v1"
@@ -12,11 +10,11 @@ import (
"github.com/zalando-incubator/postgres-operator/pkg/util"
)
func (c *Controller) podListFunc(options meta_v1.ListOptions) (runtime.Object, error) {
func (c *Controller) podListFunc(options metav1.ListOptions) (runtime.Object, error) {
var labelSelector string
var fieldSelector string
opts := meta_v1.ListOptions{
opts := metav1.ListOptions{
LabelSelector: labelSelector,
FieldSelector: fieldSelector,
Watch: options.Watch,
@@ -27,11 +25,11 @@ func (c *Controller) podListFunc(options meta_v1.ListOptions) (runtime.Object, e
return c.KubeClient.Pods(c.opConfig.Namespace).List(opts)
}
func (c *Controller) podWatchFunc(options meta_v1.ListOptions) (watch.Interface, error) {
func (c *Controller) podWatchFunc(options metav1.ListOptions) (watch.Interface, error) {
var labelSelector string
var fieldSelector string
opts := meta_v1.ListOptions{
opts := metav1.ListOptions{
LabelSelector: labelSelector,
FieldSelector: fieldSelector,
Watch: options.Watch,
+13 -13
View File
@@ -8,7 +8,7 @@ import (
"sync/atomic"
"time"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/watch"
@@ -27,14 +27,14 @@ func (c *Controller) clusterResync(stopCh <-chan struct{}, wg *sync.WaitGroup) {
for {
select {
case <-ticker.C:
c.clusterListFunc(meta_v1.ListOptions{ResourceVersion: "0"})
c.clusterListFunc(metav1.ListOptions{ResourceVersion: "0"})
case <-stopCh:
return
}
}
}
func (c *Controller) clusterListFunc(options meta_v1.ListOptions) (runtime.Object, error) {
func (c *Controller) clusterListFunc(options metav1.ListOptions) (runtime.Object, error) {
var list spec.PostgresqlList
var activeClustersCnt, failedClustersCnt int
@@ -42,7 +42,7 @@ func (c *Controller) clusterListFunc(options meta_v1.ListOptions) (runtime.Objec
Get().
Namespace(c.opConfig.Namespace).
Resource(constants.ResourceName).
VersionedParams(&options, meta_v1.ParameterCodec)
VersionedParams(&options, metav1.ParameterCodec)
b, err := req.DoRaw()
if err != nil {
@@ -101,13 +101,13 @@ func (d *tprDecoder) Decode() (action watch.EventType, object runtime.Object, er
return e.Type, &e.Object, nil
}
func (c *Controller) clusterWatchFunc(options meta_v1.ListOptions) (watch.Interface, error) {
func (c *Controller) clusterWatchFunc(options metav1.ListOptions) (watch.Interface, error) {
options.Watch = true
r, err := c.RestClient.
Get().
Namespace(c.opConfig.Namespace).
Resource(constants.ResourceName).
VersionedParams(&options, meta_v1.ParameterCodec).
VersionedParams(&options, metav1.ParameterCodec).
FieldsSelectorParam(nil).
Stream()
@@ -131,9 +131,9 @@ func (c *Controller) processEvent(obj interface{}) error {
logger := c.logger.WithField("worker", event.WorkerID)
if event.EventType == spec.EventAdd || event.EventType == spec.EventSync {
clusterName = util.NameFromMeta(event.NewSpec.Metadata)
clusterName = util.NameFromMeta(event.NewSpec.ObjectMeta)
} else {
clusterName = util.NameFromMeta(event.OldSpec.Metadata)
clusterName = util.NameFromMeta(event.OldSpec.ObjectMeta)
}
c.clustersMu.RLock()
@@ -248,8 +248,8 @@ func (c *Controller) queueClusterEvent(old, new *spec.Postgresql, eventType spec
)
if old != nil { //update, delete
uid = old.Metadata.GetUID()
clusterName = util.NameFromMeta(old.Metadata)
uid = old.GetUID()
clusterName = util.NameFromMeta(old.ObjectMeta)
if eventType == spec.EventUpdate && new.Error == nil && old.Error != nil {
eventType = spec.EventSync
clusterError = new.Error
@@ -257,8 +257,8 @@ func (c *Controller) queueClusterEvent(old, new *spec.Postgresql, eventType spec
clusterError = old.Error
}
} else { //add, sync
uid = new.Metadata.GetUID()
clusterName = util.NameFromMeta(new.Metadata)
uid = new.GetUID()
clusterName = util.NameFromMeta(new.ObjectMeta)
clusterError = new.Error
}
@@ -303,7 +303,7 @@ func (c *Controller) postgresqlUpdate(prev, cur interface{}) {
if !ok {
c.logger.Errorf("could not cast to postgresql spec")
}
if pgOld.Metadata.ResourceVersion == pgNew.Metadata.ResourceVersion {
if pgOld.ResourceVersion == pgNew.ResourceVersion {
return
}
if reflect.DeepEqual(pgOld.Spec, pgNew.Spec) {
+3 -3
View File
@@ -4,7 +4,7 @@ import (
"fmt"
"hash/crc32"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/pkg/api/v1"
extv1beta "k8s.io/client-go/pkg/apis/extensions/v1beta1"
@@ -30,7 +30,7 @@ func (c *Controller) makeClusterConfig() cluster.Config {
func thirdPartyResource(TPRName string) *extv1beta.ThirdPartyResource {
return &extv1beta.ThirdPartyResource{
ObjectMeta: meta_v1.ObjectMeta{
ObjectMeta: metav1.ObjectMeta{
//ThirdPartyResources are cluster-wide
Name: TPRName,
},
@@ -69,7 +69,7 @@ func (c *Controller) getInfrastructureRoles(rolesSecret *spec.NamespacedName) (r
infraRolesSecret, err := c.KubeClient.
Secrets(rolesSecret.Namespace).
Get(rolesSecret.Name, meta_v1.GetOptions{})
Get(rolesSecret.Name, metav1.GetOptions{})
if err != nil {
c.logger.Debugf("Infrastructure roles secret name: %q", *rolesSecret)
return nil, fmt.Errorf("could not get infrastructure roles secret: %v", err)
+3 -3
View File
@@ -5,7 +5,7 @@ import (
"reflect"
"testing"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
v1core "k8s.io/client-go/kubernetes/typed/core/v1"
"k8s.io/client-go/pkg/api/v1"
@@ -21,7 +21,7 @@ type mockSecret struct {
v1core.SecretInterface
}
func (c *mockSecret) Get(name string, options meta_v1.GetOptions) (*v1.Secret, error) {
func (c *mockSecret) Get(name string, options metav1.GetOptions) (*v1.Secret, error) {
if name != testInfrastructureRolesSecretName {
return nil, fmt.Errorf("NotFound")
}
@@ -70,7 +70,7 @@ func TestPodClusterName(t *testing.T) {
},
{
&v1.Pod{
ObjectMeta: meta_v1.ObjectMeta{
ObjectMeta: metav1.ObjectMeta{
Namespace: v1.NamespaceDefault,
Labels: map[string]string{
mockController.opConfig.ClusterNameLabel: "testcluster",