mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-02 16:27:09 +02:00
move from tpr to crd
This commit is contained in:
@@ -8,7 +8,6 @@ import (
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/apiserver"
|
||||
@@ -27,7 +26,6 @@ type Controller struct {
|
||||
|
||||
logger *logrus.Entry
|
||||
KubeClient k8sutil.KubernetesClient
|
||||
RestClient rest.Interface // kubernetes API group REST client
|
||||
apiserver *apiserver.Server
|
||||
|
||||
stopCh chan struct{}
|
||||
@@ -69,15 +67,11 @@ func NewController(controllerConfig *spec.ControllerConfig) *Controller {
|
||||
}
|
||||
|
||||
func (c *Controller) initClients() {
|
||||
client, err := k8sutil.ClientSet(c.config.RestConfig)
|
||||
if err != nil {
|
||||
c.logger.Fatalf("couldn't create client: %v", err)
|
||||
}
|
||||
c.KubeClient = k8sutil.NewFromKubernetesInterface(client)
|
||||
var err error
|
||||
|
||||
c.RestClient, err = k8sutil.KubernetesRestClient(*c.config.RestConfig)
|
||||
c.KubeClient, err = k8sutil.NewFromConfig(c.config.RestConfig)
|
||||
if err != nil {
|
||||
c.logger.Fatalf("couldn't create rest client: %v", err)
|
||||
c.logger.Fatalf("could not create kubernetes clients: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -119,8 +113,8 @@ func (c *Controller) initController() {
|
||||
c.logger.Logger.Level = logrus.DebugLevel
|
||||
}
|
||||
|
||||
if err := c.createTPR(); err != nil {
|
||||
c.logger.Fatalf("could not register ThirdPartyResource: %v", err)
|
||||
if err := c.createCRD(); err != nil {
|
||||
c.logger.Fatalf("could not register CustomResourceDefinition: %v", err)
|
||||
}
|
||||
|
||||
if infraRoles, err := c.getInfrastructureRoles(&c.opConfig.InfrastructureRolesSecretName); err != nil {
|
||||
|
||||
@@ -44,10 +44,10 @@ func (c *Controller) clusterListFunc(options metav1.ListOptions) (runtime.Object
|
||||
var list spec.PostgresqlList
|
||||
var activeClustersCnt, failedClustersCnt int
|
||||
|
||||
req := c.RestClient.
|
||||
req := c.KubeClient.CRDREST.
|
||||
Get().
|
||||
Namespace(c.opConfig.Namespace).
|
||||
Resource(constants.ResourceName).
|
||||
Resource(constants.CRDResource).
|
||||
VersionedParams(&options, metav1.ParameterCodec)
|
||||
|
||||
b, err := req.DoRaw()
|
||||
@@ -109,10 +109,10 @@ func (d *tprDecoder) Decode() (action watch.EventType, object runtime.Object, er
|
||||
|
||||
func (c *Controller) clusterWatchFunc(options metav1.ListOptions) (watch.Interface, error) {
|
||||
options.Watch = true
|
||||
r, err := c.RestClient.
|
||||
r, err := c.KubeClient.CRDREST.
|
||||
Get().
|
||||
Namespace(c.opConfig.Namespace).
|
||||
Resource(constants.ResourceName).
|
||||
Resource(constants.CRDResource).
|
||||
VersionedParams(&options, metav1.ParameterCodec).
|
||||
FieldsSelectorParam(nil).
|
||||
Stream()
|
||||
|
||||
+49
-25
@@ -4,9 +4,10 @@ import (
|
||||
"fmt"
|
||||
"hash/crc32"
|
||||
|
||||
apiextv1beta1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1beta1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
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"
|
||||
@@ -28,36 +29,59 @@ func (c *Controller) makeClusterConfig() cluster.Config {
|
||||
}
|
||||
}
|
||||
|
||||
func thirdPartyResource(TPRName string) *extv1beta.ThirdPartyResource {
|
||||
return &extv1beta.ThirdPartyResource{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
//ThirdPartyResources are cluster-wide
|
||||
Name: TPRName,
|
||||
},
|
||||
Versions: []extv1beta.APIVersion{
|
||||
{Name: constants.TPRApiVersion},
|
||||
},
|
||||
Description: constants.TPRDescription,
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Controller) clusterWorkerID(clusterName spec.NamespacedName) uint32 {
|
||||
return crc32.ChecksumIEEE([]byte(clusterName.String())) % c.opConfig.Workers
|
||||
}
|
||||
|
||||
func (c *Controller) createTPR() error {
|
||||
tpr := thirdPartyResource(constants.TPRName)
|
||||
|
||||
if _, err := c.KubeClient.ThirdPartyResources().Create(tpr); err != nil {
|
||||
if !k8sutil.ResourceAlreadyExists(err) {
|
||||
return err
|
||||
}
|
||||
c.logger.Infof("thirdPartyResource %q is already registered", constants.TPRName)
|
||||
} else {
|
||||
c.logger.Infof("thirdPartyResource %q' has been registered", constants.TPRName)
|
||||
func (c *Controller) createCRD() error {
|
||||
crd := &apiextv1beta1.CustomResourceDefinition{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: constants.CRDResource + "." + constants.CRDGroup,
|
||||
},
|
||||
Spec: apiextv1beta1.CustomResourceDefinitionSpec{
|
||||
Group: constants.CRDGroup,
|
||||
Version: constants.CRDApiVersion,
|
||||
Names: apiextv1beta1.CustomResourceDefinitionNames{
|
||||
Plural: constants.CRDResource,
|
||||
Singular: constants.CRDKind,
|
||||
ShortNames: []string{constants.CRDShort},
|
||||
Kind: constants.CRDKind,
|
||||
ListKind: constants.CRDKind + "List",
|
||||
},
|
||||
Scope: apiextv1beta1.NamespaceScoped,
|
||||
},
|
||||
}
|
||||
|
||||
return k8sutil.WaitTPRReady(c.RestClient, c.opConfig.TPR.ReadyWaitInterval, c.opConfig.TPR.ReadyWaitTimeout, c.opConfig.Namespace)
|
||||
if _, err := c.KubeClient.CustomResourceDefinitions().Create(crd); err != nil {
|
||||
if !k8sutil.ResourceAlreadyExists(err) {
|
||||
return fmt.Errorf("could not create customResourceDefinition: %v", err)
|
||||
}
|
||||
c.logger.Infof("customResourceDefinition %q is already registered", crd.Name)
|
||||
} else {
|
||||
c.logger.Infof("customResourceDefinition %q has been registered", crd.Name)
|
||||
}
|
||||
|
||||
return wait.Poll(c.opConfig.CRD.ReadyWaitInterval, c.opConfig.CRD.ReadyWaitTimeout, func() (bool, error) {
|
||||
c, err := c.KubeClient.CustomResourceDefinitions().Get(crd.Name, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
for _, cond := range c.Status.Conditions {
|
||||
switch cond.Type {
|
||||
case apiextv1beta1.Established:
|
||||
if cond.Status == apiextv1beta1.ConditionTrue {
|
||||
return true, err
|
||||
}
|
||||
case apiextv1beta1.NamesAccepted:
|
||||
if cond.Status == apiextv1beta1.ConditionFalse {
|
||||
return false, fmt.Errorf("name conflict: %v", cond.Reason)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return false, err
|
||||
})
|
||||
}
|
||||
|
||||
func (c *Controller) getInfrastructureRoles(rolesSecret *spec.NamespacedName) (result map[string]spec.PgUser, err error) {
|
||||
|
||||
Reference in New Issue
Block a user