use kubernetes client-go from master branch o_O

This commit is contained in:
Murat Kabilov
2017-05-22 11:34:20 +02:00
parent d9b0ee9198
commit 614c270c55
17 changed files with 155 additions and 297 deletions
+5 -5
View File
@@ -11,11 +11,11 @@ import (
"sync"
"github.com/Sirupsen/logrus"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/pkg/api"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/apis/apps/v1beta1"
"k8s.io/client-go/pkg/types"
"k8s.io/client-go/rest"
"github.com/zalando-incubator/postgres-operator/pkg/spec"
@@ -65,7 +65,7 @@ type Cluster struct {
masterLess bool
podDispatcherRunning bool
userSyncStrategy spec.UserSyncer
deleteOptions *v1.DeleteOptions
deleteOptions *meta_v1.DeleteOptions
}
func New(cfg Config, pgSpec spec.Postgresql, logger *logrus.Entry) *Cluster {
@@ -85,7 +85,7 @@ func New(cfg Config, pgSpec spec.Postgresql, logger *logrus.Entry) *Cluster {
masterLess: false,
podDispatcherRunning: false,
userSyncStrategy: users.DefaultUserSyncStrategy{},
deleteOptions: &v1.DeleteOptions{OrphanDependents: &orphanDependents},
deleteOptions: &meta_v1.DeleteOptions{OrphanDependents: &orphanDependents},
}
return cluster
@@ -108,7 +108,7 @@ func (c *Cluster) setStatus(status spec.PostgresStatus) {
}
request := []byte(fmt.Sprintf(`{"status": %s}`, string(b))) //TODO: Look into/wait for k8s go client methods
_, err = c.RestClient.Patch(api.MergePatchType).
_, err = c.RestClient.Patch(types.MergePatchType).
RequestURI(c.Metadata.GetSelfLink()).
Body(request).
DoRaw()
+9 -8
View File
@@ -5,10 +5,11 @@ import (
"sort"
"encoding/json"
"k8s.io/client-go/pkg/api/resource"
"k8s.io/apimachinery/pkg/api/resource"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/intstr"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/apis/apps/v1beta1"
"k8s.io/client-go/pkg/util/intstr"
"github.com/zalando-incubator/postgres-operator/pkg/spec"
"github.com/zalando-incubator/postgres-operator/pkg/util/constants"
@@ -310,7 +311,7 @@ func (c *Cluster) genPodTemplate(resourceRequirements *v1.ResourceRequirements,
}
template := v1.PodTemplateSpec{
ObjectMeta: v1.ObjectMeta{
ObjectMeta: meta_v1.ObjectMeta{
Labels: c.labelsSet(),
Namespace: c.Metadata.Name,
},
@@ -336,7 +337,7 @@ func (c *Cluster) genStatefulSet(spec spec.PostgresSpec) (*v1beta1.StatefulSet,
}
statefulSet := &v1beta1.StatefulSet{
ObjectMeta: v1.ObjectMeta{
ObjectMeta: meta_v1.ObjectMeta{
Name: c.Metadata.Name,
Namespace: c.Metadata.Namespace,
Labels: c.labelsSet(),
@@ -353,7 +354,7 @@ func (c *Cluster) genStatefulSet(spec spec.PostgresSpec) (*v1beta1.StatefulSet,
}
func persistentVolumeClaimTemplate(volumeSize, volumeStorageClass string) (*v1.PersistentVolumeClaim, error) {
metadata := v1.ObjectMeta{
metadata := meta_v1.ObjectMeta{
Name: constants.DataVolumeName,
}
if volumeStorageClass != "" {
@@ -411,7 +412,7 @@ func (c *Cluster) genSingleUserSecret(namespace string, pgUser spec.PgUser) *v1.
}
username := pgUser.Name
secret := v1.Secret{
ObjectMeta: v1.ObjectMeta{
ObjectMeta: meta_v1.ObjectMeta{
Name: c.credentialSecretName(username),
Namespace: namespace,
Labels: c.labelsSet(),
@@ -427,7 +428,7 @@ func (c *Cluster) genSingleUserSecret(namespace string, pgUser spec.PgUser) *v1.
func (c *Cluster) genService(allowedSourceRanges []string) *v1.Service {
service := &v1.Service{
ObjectMeta: v1.ObjectMeta{
ObjectMeta: meta_v1.ObjectMeta{
Name: c.Metadata.Name,
Namespace: c.Metadata.Namespace,
Labels: c.labelsSet(),
@@ -448,7 +449,7 @@ func (c *Cluster) genService(allowedSourceRanges []string) *v1.Service {
func (c *Cluster) genEndpoints() *v1.Endpoints {
endpoints := &v1.Endpoints{
ObjectMeta: v1.ObjectMeta{
ObjectMeta: meta_v1.ObjectMeta{
Name: c.Metadata.Name,
Namespace: c.Metadata.Namespace,
Labels: c.labelsSet(),
+4 -3
View File
@@ -3,6 +3,7 @@ package cluster
import (
"fmt"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/pkg/api/v1"
"github.com/zalando-incubator/postgres-operator/pkg/spec"
@@ -12,7 +13,7 @@ import (
func (c *Cluster) listPods() ([]v1.Pod, error) {
ns := c.Metadata.Namespace
listOptions := v1.ListOptions{
listOptions := meta_v1.ListOptions{
LabelSelector: c.labelsSet().String(),
}
@@ -26,7 +27,7 @@ func (c *Cluster) listPods() ([]v1.Pod, error) {
func (c *Cluster) listPersistentVolumeClaims() ([]v1.PersistentVolumeClaim, error) {
ns := c.Metadata.Namespace
listOptions := v1.ListOptions{
listOptions := meta_v1.ListOptions{
LabelSelector: c.labelsSet().String(),
}
@@ -167,7 +168,7 @@ func (c *Cluster) recreatePods() error {
ls := c.labelsSet()
namespace := c.Metadata.Namespace
listOptions := v1.ListOptions{
listOptions := meta_v1.ListOptions{
LabelSelector: ls.String(),
}
+8 -7
View File
@@ -3,7 +3,8 @@ package cluster
import (
"fmt"
"k8s.io/client-go/pkg/api"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/apis/apps/v1beta1"
@@ -16,7 +17,7 @@ import (
func (c *Cluster) loadResources() error {
ns := c.Metadata.Namespace
listOptions := v1.ListOptions{
listOptions := meta_v1.ListOptions{
LabelSelector: c.labelsSet().String(),
}
@@ -139,7 +140,7 @@ func (c *Cluster) updateStatefulSet(newStatefulSet *v1beta1.StatefulSet) error {
statefulSet, err := c.KubeClient.StatefulSets(c.Statefulset.Namespace).Patch(
c.Statefulset.Name,
api.MergePatchType,
types.MergePatchType,
patchData, "")
if err != nil {
return fmt.Errorf("Can't patch StatefulSet '%s': %s", statefulSetName, err)
@@ -162,7 +163,7 @@ func (c *Cluster) replaceStatefulSet(newStatefulSet *v1beta1.StatefulSet) error
orphanDepencies := true
oldStatefulset := c.Statefulset
options := v1.DeleteOptions{OrphanDependents: &orphanDepencies}
options := meta_v1.DeleteOptions{OrphanDependents: &orphanDepencies}
if err := c.KubeClient.StatefulSets(oldStatefulset.Namespace).Delete(oldStatefulset.Name, &options); err != nil {
return fmt.Errorf("Can't delete statefulset '%s': %s", statefulSetName, err)
}
@@ -173,7 +174,7 @@ func (c *Cluster) replaceStatefulSet(newStatefulSet *v1beta1.StatefulSet) error
err := retryutil.Retry(constants.StatefulsetDeletionInterval, constants.StatefulsetDeletionTimeout,
func() (bool, error) {
_, err := c.KubeClient.StatefulSets(oldStatefulset.Namespace).Get(oldStatefulset.Name)
_, err := c.KubeClient.StatefulSets(oldStatefulset.Namespace).Get(oldStatefulset.Name, meta_v1.GetOptions{})
return err != nil, nil
})
@@ -251,7 +252,7 @@ func (c *Cluster) updateService(newService *v1.Service) error {
svc, err := c.KubeClient.Services(c.Service.Namespace).Patch(
c.Service.Name,
api.MergePatchType,
types.MergePatchType,
patchData, "")
if err != nil {
return fmt.Errorf("Can't patch Service '%s': %s", serviceName, err)
@@ -317,7 +318,7 @@ func (c *Cluster) applySecrets() error {
secret, err := c.KubeClient.Secrets(secretSpec.Namespace).Create(secretSpec)
if k8sutil.ResourceAlreadyExists(err) {
var userMap map[string]spec.PgUser
curSecret, err := c.KubeClient.Secrets(secretSpec.Namespace).Get(secretSpec.Name)
curSecret, err := c.KubeClient.Secrets(secretSpec.Namespace).Get(secretSpec.Name, meta_v1.GetOptions{})
if err != nil {
return fmt.Errorf("Can't get current Secret: %s", err)
}
+6 -5
View File
@@ -6,9 +6,10 @@ import (
"strings"
"time"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/apis/apps/v1beta1"
"k8s.io/client-go/pkg/labels"
"github.com/zalando-incubator/postgres-operator/pkg/spec"
"github.com/zalando-incubator/postgres-operator/pkg/util"
@@ -151,7 +152,7 @@ func (c *Cluster) waitForPodDeletion(podEvents chan spec.PodEvent) error {
func (c *Cluster) waitStatefulsetReady() error {
return retryutil.Retry(c.OpConfig.ResourceCheckInterval, c.OpConfig.ResourceCheckTimeout,
func() (bool, error) {
listOptions := v1.ListOptions{
listOptions := meta_v1.ListOptions{
LabelSelector: c.labelsSet().String(),
}
ss, err := c.KubeClient.StatefulSets(c.Metadata.Namespace).List(listOptions)
@@ -171,15 +172,15 @@ func (c *Cluster) waitPodLabelsReady() error {
ls := c.labelsSet()
namespace := c.Metadata.Namespace
listOptions := v1.ListOptions{
listOptions := meta_v1.ListOptions{
LabelSelector: ls.String(),
}
masterListOption := v1.ListOptions{
masterListOption := meta_v1.ListOptions{
LabelSelector: labels.Merge(ls, labels.Set{
c.OpConfig.PodRoleLabel: constants.PodRoleMaster,
}).String(),
}
replicaListOption := v1.ListOptions{
replicaListOption := meta_v1.ListOptions{
LabelSelector: labels.Merge(ls, labels.Set{
c.OpConfig.PodRoleLabel: constants.PodRoleReplica,
}).String(),
+2 -2
View File
@@ -26,8 +26,8 @@ type Config struct {
type Controller struct {
Config
opConfig *config.Config
logger *logrus.Entry
opConfig *config.Config
logger *logrus.Entry
clustersMu sync.RWMutex
clusters map[spec.NamespacedName]*cluster.Cluster
+3 -3
View File
@@ -4,9 +4,10 @@ import (
"bytes"
"fmt"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
remotecommandconsts "k8s.io/apimachinery/pkg/util/remotecommand"
"k8s.io/client-go/pkg/api"
"k8s.io/kubernetes/pkg/client/unversioned/remotecommand"
"k8s.io/client-go/tools/remotecommand"
"github.com/zalando-incubator/postgres-operator/pkg/spec"
)
@@ -16,8 +17,7 @@ func (c *Controller) ExecCommand(podName spec.NamespacedName, command []string)
execOut bytes.Buffer
execErr bytes.Buffer
)
pod, err := c.KubeClient.Pods(podName.Namespace).Get(podName.Name)
pod, err := c.KubeClient.Pods(podName.Namespace).Get(podName.Name, meta_v1.GetOptions{})
if err != nil {
return "", fmt.Errorf("Can't get Pod info: %s", err)
}
+7 -44
View File
@@ -1,58 +1,21 @@
package controller
import (
"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) {
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{
LabelSelector: labelSelector,
FieldSelector: fieldSelector,
Watch: options.Watch,
ResourceVersion: options.ResourceVersion,
TimeoutSeconds: options.TimeoutSeconds,
}
return c.KubeClient.CoreV1().Pods(c.opConfig.Namespace).List(opts)
func (c *Controller) podListFunc(options meta_v1.ListOptions) (runtime.Object, error) {
return c.KubeClient.CoreV1().Pods(c.opConfig.Namespace).List(options)
}
func (c *Controller) podWatchFunc(options api.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{
LabelSelector: labelSelector,
FieldSelector: fieldSelector,
Watch: options.Watch,
ResourceVersion: options.ResourceVersion,
TimeoutSeconds: options.TimeoutSeconds,
}
return c.KubeClient.CoreV1Client.Pods(c.opConfig.Namespace).Watch(opts)
func (c *Controller) podWatchFunc(options meta_v1.ListOptions) (watch.Interface, error) {
return c.KubeClient.CoreV1().Pods(c.opConfig.Namespace).Watch(options)
}
func (c *Controller) podAdd(obj interface{}) {
+8 -7
View File
@@ -4,12 +4,13 @@ import (
"fmt"
"reflect"
"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"
@@ -18,7 +19,7 @@ import (
"github.com/zalando-incubator/postgres-operator/pkg/util/constants"
)
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().
@@ -68,7 +69,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 -3
View File
@@ -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"
@@ -33,7 +34,7 @@ func (c *Controller) getOAuthToken() (string, error) {
// Temporary getting postgresql-operator secret from the NamespaceDefault
credentialsSecret, err := c.KubeClient.
Secrets(c.opConfig.OAuthTokenSecretName.Namespace).
Get(c.opConfig.OAuthTokenSecretName.Name)
Get(c.opConfig.OAuthTokenSecretName.Name, meta_v1.GetOptions{})
if err != nil {
c.logger.Debugf("Oauth token secret name: %s", c.opConfig.OAuthTokenSecretName)
@@ -50,7 +51,7 @@ func (c *Controller) getOAuthToken() (string, error) {
func thirdPartyResource(TPRName string) *extv1beta.ThirdPartyResource {
return &extv1beta.ThirdPartyResource{
ObjectMeta: v1.ObjectMeta{
ObjectMeta: meta_v1.ObjectMeta{
//ThirdPartyResources are cluster-wide
Name: TPRName,
},
@@ -90,7 +91,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("Can't get infrastructure roles Secret: %s", err)
+11 -11
View File
@@ -7,9 +7,9 @@ import (
"strings"
"time"
"k8s.io/client-go/pkg/api/meta"
"k8s.io/client-go/pkg/api/unversioned"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"
)
var alphaRegexp = regexp.MustCompile("^[a-zA-Z]*$")
@@ -67,8 +67,8 @@ const (
// PostgreSQL Third Party (resource) Object
type Postgresql struct {
unversioned.TypeMeta `json:",inline"`
Metadata v1.ObjectMeta `json:"metadata"`
v1.TypeMeta `json:",inline"`
Metadata meta_v1.ObjectMeta `json:"metadata"`
Spec PostgresSpec `json:"spec"`
Status PostgresStatus `json:"status"`
@@ -90,8 +90,8 @@ type PostgresSpec struct {
}
type PostgresqlList struct {
unversioned.TypeMeta `json:",inline"`
Metadata unversioned.ListMeta `json:"metadata"`
v1.TypeMeta `json:",inline"`
Metadata v1.ListMeta `json:"metadata"`
Items []Postgresql `json:"items"`
}
@@ -175,19 +175,19 @@ func (m *MaintenanceWindow) UnmarshalJSON(data []byte) error {
return nil
}
func (p *Postgresql) GetObjectKind() unversioned.ObjectKind {
func (p *Postgresql) GetObjectKind() schema.ObjectKind {
return &p.TypeMeta
}
func (p *Postgresql) GetObjectMeta() meta.Object {
func (p *Postgresql) GetObjectMeta() v1.Object {
return &p.Metadata
}
func (pl *PostgresqlList) GetObjectKind() unversioned.ObjectKind {
func (pl *PostgresqlList) GetObjectKind() schema.ObjectKind {
return &pl.TypeMeta
}
func (pl *PostgresqlList) GetListMeta() unversioned.List {
func (pl *PostgresqlList) GetListMeta() v1.List {
return &pl.Metadata
}
+1 -1
View File
@@ -3,8 +3,8 @@ package spec
import (
"database/sql"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/types"
)
type EventType string
+9 -8
View File
@@ -4,12 +4,13 @@ import (
"fmt"
"time"
apierrors "k8s.io/apimachinery/pkg/api/errors"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apimachinery/pkg/runtime/serializer"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/pkg/api"
apierrors "k8s.io/client-go/pkg/api/errors"
"k8s.io/client-go/pkg/api/unversioned"
"k8s.io/client-go/pkg/runtime"
"k8s.io/client-go/pkg/runtime/serializer"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
@@ -38,21 +39,21 @@ func ResourceNotFound(err error) bool {
}
func KubernetesRestClient(c *rest.Config) (*rest.RESTClient, error) {
c.GroupVersion = &unversioned.GroupVersion{Version: constants.K8sVersion}
c.GroupVersion = &schema.GroupVersion{Version: constants.K8sVersion}
c.APIPath = constants.K8sAPIPath
c.NegotiatedSerializer = serializer.DirectCodecFactory{CodecFactory: api.Codecs}
schemeBuilder := runtime.NewSchemeBuilder(
func(scheme *runtime.Scheme) error {
scheme.AddKnownTypes(
unversioned.GroupVersion{
schema.GroupVersion{
Group: constants.TPRVendor,
Version: constants.TPRApiVersion,
},
&spec.Postgresql{},
&spec.PostgresqlList{},
&api.ListOptions{},
&api.DeleteOptions{},
&meta_v1.ListOptions{},
&meta_v1.DeleteOptions{},
)
return nil
})
+2 -2
View File
@@ -9,7 +9,7 @@ import (
"time"
"github.com/motomux/pretty"
"k8s.io/client-go/pkg/api/v1"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"github.com/zalando-incubator/postgres-operator/pkg/spec"
)
@@ -33,7 +33,7 @@ func RandomPassword(n int) string {
return string(b)
}
func NameFromMeta(meta v1.ObjectMeta) spec.NamespacedName {
func NameFromMeta(meta meta_v1.ObjectMeta) spec.NamespacedName {
return spec.NamespacedName{
Namespace: meta.Namespace,
Name: meta.Name,