Make use of kubernetes client-go v4

* client-go v4.0.0-beta0
* remove unnecessary methods for tpr object
* rest client: use interface instead of structure pointer
* proper names for constants; some clean up for log messages
* remove teams api client from controller and make it per cluster
This commit is contained in:
Murat Kabilov
2017-07-25 15:25:17 +02:00
committed by GitHub
parent 4455f1b639
commit 1f8b37f33d
31 changed files with 579 additions and 666 deletions
+86 -63
View File
@@ -1,17 +1,16 @@
package controller
import (
"encoding/json"
"fmt"
"reflect"
"sync/atomic"
"time"
"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"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/watch"
"k8s.io/client-go/tools/cache"
"github.com/zalando-incubator/postgres-operator/pkg/cluster"
@@ -26,52 +25,43 @@ func (c *Controller) clusterResync(stopCh <-chan struct{}) {
for {
select {
case <-ticker.C:
c.clusterListFunc(api.ListOptions{ResourceVersion: "0"})
c.clusterListFunc(metav1.ListOptions{ResourceVersion: "0"})
case <-stopCh:
return
}
}
}
func (c *Controller) clusterListFunc(options api.ListOptions) (runtime.Object, error) {
c.logger.Info("Getting list of currently running clusters")
func (c *Controller) clusterListFunc(options metav1.ListOptions) (runtime.Object, error) {
var list spec.PostgresqlList
var activeClustersCnt, failedClustersCnt int
req := c.RestClient.Get().
RequestURI(fmt.Sprintf(constants.ListClustersURITemplate, c.opConfig.Namespace)).
VersionedParams(&options, api.ParameterCodec).
FieldsSelectorParam(fields.Everything())
object, err := req.Do().Get()
req := c.RestClient.
Get().
Namespace(c.opConfig.Namespace).
Resource(constants.ResourceName).
VersionedParams(&options, metav1.ParameterCodec)
b, err := req.DoRaw()
if err != nil {
return nil, fmt.Errorf("could not get list of postgresql objects: %v", err)
}
objList, err := meta.ExtractList(object)
if err != nil {
return nil, fmt.Errorf("could not extract list of postgresql objects: %v", err)
return nil, err
}
err = json.Unmarshal(b, &list)
if time.Now().Unix()-atomic.LoadInt64(&c.lastClusterSyncTime) <= int64(c.opConfig.ResyncPeriod.Seconds()) {
c.logger.Debugln("skipping resync of clusters")
return object, err
return &list, err
}
var activeClustersCnt, failedClustersCnt int
for _, obj := range objList {
pg, ok := obj.(*spec.Postgresql)
if !ok {
return nil, fmt.Errorf("could not cast object to postgresql")
}
for _, pg := range list.Items {
if pg.Error != nil {
failedClustersCnt++
continue
}
c.queueClusterEvent(nil, pg, spec.EventSync)
c.queueClusterEvent(nil, &pg, spec.EventSync)
activeClustersCnt++
}
if len(objList) > 0 {
if len(list.Items) > 0 {
if failedClustersCnt > 0 && activeClustersCnt == 0 {
c.logger.Infof("There are no clusters running. %d are in the failed state", failedClustersCnt)
} else if failedClustersCnt == 0 && activeClustersCnt > 0 {
@@ -85,15 +75,48 @@ func (c *Controller) clusterListFunc(options api.ListOptions) (runtime.Object, e
atomic.StoreInt64(&c.lastClusterSyncTime, time.Now().Unix())
return object, err
return &list, err
}
func (c *Controller) clusterWatchFunc(options api.ListOptions) (watch.Interface, error) {
req := c.RestClient.Get().
RequestURI(fmt.Sprintf(constants.WatchClustersURITemplate, c.opConfig.Namespace)).
VersionedParams(&options, api.ParameterCodec).
FieldsSelectorParam(fields.Everything())
return req.Watch()
type tprDecoder struct {
dec *json.Decoder
close func() error
}
func (d *tprDecoder) Close() {
d.close()
}
func (d *tprDecoder) Decode() (action watch.EventType, object runtime.Object, err error) {
var e struct {
Type watch.EventType
Object spec.Postgresql
}
if err := d.dec.Decode(&e); err != nil {
return watch.Error, nil, err
}
return e.Type, &e.Object, nil
}
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, metav1.ParameterCodec).
FieldsSelectorParam(nil).
Stream()
if err != nil {
return nil, err
}
return watch.NewStreamWatcher(&tprDecoder{
dec: json.NewDecoder(r),
close: r.Close,
}), nil
}
func (c *Controller) processEvent(obj interface{}) error {
@@ -106,9 +129,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()
@@ -118,14 +141,14 @@ func (c *Controller) processEvent(obj interface{}) error {
switch event.EventType {
case spec.EventAdd:
if clusterFound {
logger.Debugf("Cluster '%s' already exists", clusterName)
logger.Debugf("Cluster %q already exists", clusterName)
return nil
}
logger.Infof("Creation of the '%s' cluster started", clusterName)
logger.Infof("Creation of the %q cluster started", clusterName)
stopCh := make(chan struct{})
cl = cluster.New(c.makeClusterConfig(), *event.NewSpec, logger)
cl = cluster.New(c.makeClusterConfig(), c.KubeClient, *event.NewSpec, logger)
cl.Run(stopCh)
c.clustersMu.Lock()
@@ -140,31 +163,31 @@ func (c *Controller) processEvent(obj interface{}) error {
return nil
}
logger.Infof("Cluster '%s' has been created", clusterName)
logger.Infof("Cluster %q has been created", clusterName)
case spec.EventUpdate:
logger.Infof("Update of the '%s' cluster started", clusterName)
logger.Infof("Update of the %q cluster started", clusterName)
if !clusterFound {
logger.Warnf("Cluster '%s' does not exist", clusterName)
logger.Warnf("Cluster %q does not exist", clusterName)
return nil
}
if err := cl.Update(event.NewSpec); err != nil {
cl.Error = fmt.Errorf("could not update cluster: %s", err)
cl.Error = fmt.Errorf("could not update cluster: %v", err)
logger.Errorf("%v", cl.Error)
return nil
}
cl.Error = nil
logger.Infof("Cluster '%s' has been updated", clusterName)
logger.Infof("Cluster %q has been updated", clusterName)
case spec.EventDelete:
logger.Infof("Deletion of the '%s' cluster started", clusterName)
logger.Infof("Deletion of the %q cluster started", clusterName)
if !clusterFound {
logger.Errorf("Unknown cluster: %s", clusterName)
logger.Errorf("Unknown cluster: %q", clusterName)
return nil
}
if err := cl.Delete(); err != nil {
logger.Errorf("could not delete cluster '%s': %s", clusterName, err)
logger.Errorf("could not delete cluster %q: %v", clusterName, err)
return nil
}
close(c.stopChs[clusterName])
@@ -174,14 +197,14 @@ func (c *Controller) processEvent(obj interface{}) error {
delete(c.stopChs, clusterName)
c.clustersMu.Unlock()
logger.Infof("Cluster '%s' has been deleted", clusterName)
logger.Infof("Cluster %q has been deleted", clusterName)
case spec.EventSync:
logger.Infof("Syncing of the '%s' cluster started", clusterName)
logger.Infof("Syncing of the %q cluster started", clusterName)
// no race condition because a cluster is always processed by single worker
if !clusterFound {
stopCh := make(chan struct{})
cl = cluster.New(c.makeClusterConfig(), *event.NewSpec, logger)
cl = cluster.New(c.makeClusterConfig(), c.KubeClient, *event.NewSpec, logger)
cl.Run(stopCh)
c.clustersMu.Lock()
@@ -191,13 +214,13 @@ func (c *Controller) processEvent(obj interface{}) error {
}
if err := cl.Sync(); err != nil {
cl.Error = fmt.Errorf("could not sync cluster '%s': %v", clusterName, err)
cl.Error = fmt.Errorf("could not sync cluster %q: %v", clusterName, err)
logger.Errorf("%v", cl.Error)
return nil
}
cl.Error = nil
logger.Infof("Cluster '%s' has been synced", clusterName)
logger.Infof("Cluster %q has been synced", clusterName)
}
return nil
@@ -219,8 +242,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
@@ -228,13 +251,13 @@ 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
}
if clusterError != nil && eventType != spec.EventDelete {
c.logger.Debugf("Skipping %s event for invalid cluster %s (reason: %v)", eventType, clusterName, clusterError)
c.logger.Debugf("Skipping %q event for invalid cluster %q (reason: %v)", eventType, clusterName, clusterError)
return
}
@@ -251,7 +274,7 @@ func (c *Controller) queueClusterEvent(old, new *spec.Postgresql, eventType spec
if err := c.clusterEventQueues[workerID].Add(clusterEvent); err != nil {
c.logger.WithField("worker", workerID).Errorf("error when queueing cluster event: %v", clusterEvent)
}
c.logger.WithField("worker", workerID).Infof("%s of the '%s' cluster has been queued", eventType, clusterName)
c.logger.WithField("worker", workerID).Infof("%q of the %q cluster has been queued", eventType, clusterName)
}
func (c *Controller) postgresqlAdd(obj interface{}) {
@@ -274,7 +297,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) {