Migrate HPA v1 to v2beta2 (#2421)

* Migrate HPA v1 to v2beta2
HPA v2beta2 is defined and supported from 1.19+ onwards.
Also HPA v2 is stable from 1.23 onwards. As we support 1.19+
onwards using HPA v2beta2.
This change is base for custom metrics support we want to add
later by modifying Function spec.
* Add unit tests for hpa operations
* Use constants instead of strings

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2022-05-02 13:30:21 +05:30
committed by GitHub
parent e72641c0f3
commit ed4bd2573b
11 changed files with 353 additions and 242 deletions
@@ -77,7 +77,7 @@ func (cn *Container) cleanupContainer(ctx context.Context, ns string, name strin
result = multierror.Append(result, err)
}
err = cn.deleteHpa(ctx, ns, name)
err = cn.hpaops.DeleteHpa(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
cn.logger.Error("error deleting HPA for Container function",
zap.Error(err),
@@ -29,6 +29,7 @@ import (
multierror "github.com/hashicorp/go-multierror"
"github.com/pkg/errors"
"go.uber.org/zap"
asv2beta2 "k8s.io/api/autoscaling/v2beta2"
apiv1 "k8s.io/api/core/v1"
k8sErrs "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -46,6 +47,7 @@ import (
"github.com/fission/fission/pkg/executor/fscache"
"github.com/fission/fission/pkg/executor/metrics"
"github.com/fission/fission/pkg/executor/reaper"
hpautils "github.com/fission/fission/pkg/executor/util/hpa"
"github.com/fission/fission/pkg/generated/clientset/versioned"
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
"github.com/fission/fission/pkg/throttler"
@@ -81,6 +83,8 @@ type (
deplListerSynced k8sCache.InformerSynced
svcListerSynced k8sCache.InformerSynced
hpaops *hpautils.HpaOperations
}
)
@@ -120,6 +124,8 @@ func MakeContainer(
useIstio: enableIstio,
// Time is set slightly higher than NewDeploy as cold starts are longer for CaaF
defaultIdlePodReapTime: 1 * time.Minute,
hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID),
}
caaf.deplLister = deplInformer.Lister()
caaf.deplListerSynced = deplInformer.Informer().HasSynced
@@ -395,7 +401,7 @@ func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache
return nil, errors.Wrapf(err, "error creating deployment %v", objName)
}
hpa, err := caaf.createOrGetHpa(ctx, objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
hpa, err := caaf.hpaops.CreateOrGetHpa(ctx, objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
if err != nil {
caaf.logger.Error("error creating HPA", zap.Error(err), zap.String("hpa", objName))
go cleanupFunc(ns, objName)
@@ -499,7 +505,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function,
return err
}
hpa, err := caaf.getHpa(ctx, ns, fsvc.Name)
hpa, err := caaf.hpaops.GetHpa(ctx, ns, fsvc.Name)
if err != nil {
caaf.updateStatus(oldFn, err, "error getting HPA while updating function")
return err
@@ -520,12 +526,13 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function,
if newFn.Spec.InvokeStrategy.ExecutionStrategy.TargetCPUPercent != oldFn.Spec.InvokeStrategy.ExecutionStrategy.TargetCPUPercent {
targetCpupercent := int32(newFn.Spec.InvokeStrategy.ExecutionStrategy.TargetCPUPercent)
hpa.Spec.TargetCPUUtilizationPercentage = &targetCpupercent
hpaMetric := hpautils.ConvertTargetCPUToCustomMetric(targetCpupercent)
hpa.Spec.Metrics = []asv2beta2.MetricSpec{hpaMetric}
hpaChanged = true
}
if hpaChanged {
err := caaf.updateHpa(ctx, hpa)
err := caaf.hpaops.UpdateHpa(ctx, hpa)
if err != nil {
caaf.updateStatus(oldFn, err, "error updating HPA while updating function")
return err
-116
View File
@@ -1,116 +0,0 @@
/*
Copyright 2020 The Fission Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package container
import (
"context"
"github.com/pkg/errors"
"go.uber.org/zap"
appsv1 "k8s.io/api/apps/v1"
asv1 "k8s.io/api/autoscaling/v1"
k8s_err "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
otelUtils "github.com/fission/fission/pkg/utils/otel"
)
const (
DeploymentKind = "Deployment"
DeploymentVersion = "apps/v1"
)
func (cn *Container) createOrGetHpa(ctx context.Context, hpaName string, execStrategy *fv1.ExecutionStrategy,
depl *appsv1.Deployment, deployLabels map[string]string, deployAnnotations map[string]string) (*asv1.HorizontalPodAutoscaler, error) {
if depl == nil {
return nil, errors.New("failed to create HPA, found empty deployment")
}
logger := otelUtils.LoggerWithTraceID(ctx, cn.logger)
minRepl := int32(execStrategy.MinScale)
if minRepl == 0 {
minRepl = 1
}
maxRepl := int32(execStrategy.MaxScale)
if maxRepl == 0 {
maxRepl = minRepl
}
targetCPU := int32(execStrategy.TargetCPUPercent)
hpa := &asv1.HorizontalPodAutoscaler{
ObjectMeta: metav1.ObjectMeta{
Name: hpaName,
Labels: deployLabels,
Annotations: deployAnnotations,
},
Spec: asv1.HorizontalPodAutoscalerSpec{
ScaleTargetRef: asv1.CrossVersionObjectReference{
Kind: DeploymentKind,
Name: depl.ObjectMeta.Name,
APIVersion: DeploymentVersion,
},
MinReplicas: &minRepl,
MaxReplicas: maxRepl,
TargetCPUUtilizationPercentage: &targetCPU,
},
}
existingHpa, err := cn.getHpa(ctx, depl.ObjectMeta.Namespace, hpaName)
if err == nil {
// to adopt orphan service
if existingHpa.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != cn.instanceID {
existingHpa.Annotations = hpa.Annotations
existingHpa.Labels = hpa.Labels
existingHpa.Spec = hpa.Spec
existingHpa, err = cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(ctx, existingHpa, metav1.UpdateOptions{})
if err != nil {
logger.Warn("error adopting HPA", zap.Error(err),
zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace))
return nil, err
}
}
return existingHpa, err
} else if k8s_err.IsNotFound(err) {
cHpa, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(ctx, hpa, metav1.CreateOptions{})
if err != nil {
if k8s_err.IsAlreadyExists(err) {
cHpa, err = cn.getHpa(ctx, depl.ObjectMeta.Namespace, hpaName)
}
if err != nil {
return nil, err
}
}
otelUtils.SpanTrackEvent(ctx, "hpaCreated", otelUtils.GetAttributesForHPA(cHpa)...)
return cHpa, nil
}
return nil, err
}
func (cn *Container) getHpa(ctx context.Context, ns, name string) (*asv1.HorizontalPodAutoscaler, error) {
return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(ctx, name, metav1.GetOptions{})
}
func (cn *Container) updateHpa(ctx context.Context, hpa *asv1.HorizontalPodAutoscaler) error {
_, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(ctx, hpa, metav1.UpdateOptions{})
return err
}
func (cn *Container) deleteHpa(ctx context.Context, ns string, name string) error {
return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(ctx, name, metav1.DeleteOptions{})
}
@@ -23,10 +23,8 @@ import (
"time"
multierror "github.com/hashicorp/go-multierror"
"github.com/pkg/errors"
"go.uber.org/zap"
appsv1 "k8s.io/api/apps/v1"
asv1 "k8s.io/api/autoscaling/v1"
apiv1 "k8s.io/api/core/v1"
k8s_err "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/resource"
@@ -40,12 +38,6 @@ import (
otelUtils "github.com/fission/fission/pkg/utils/otel"
)
// Deployment Constants
const (
DeploymentKind = "Deployment"
DeploymentVersion = "apps/v1"
)
func (deploy *NewDeploy) createOrGetDeployment(ctx context.Context, fn *fv1.Function, env *fv1.Environment,
deployName string, deployLabels map[string]string, deployAnnotations map[string]string, deployNamespace string) (*appsv1.Deployment, error) {
@@ -367,86 +359,6 @@ func (deploy *NewDeploy) getResources(env *fv1.Environment, fn *fv1.Function) ap
return resources
}
func (deploy *NewDeploy) createOrGetHpa(ctx context.Context, hpaName string, execStrategy *fv1.ExecutionStrategy,
depl *appsv1.Deployment, deployLabels map[string]string, deployAnnotations map[string]string) (*asv1.HorizontalPodAutoscaler, error) {
if depl == nil {
return nil, errors.New("failed to create HPA, found empty deployment")
}
logger := otelUtils.LoggerWithTraceID(ctx, deploy.logger)
minRepl := int32(execStrategy.MinScale)
if minRepl == 0 {
minRepl = 1
}
maxRepl := int32(execStrategy.MaxScale)
if maxRepl == 0 {
maxRepl = minRepl
}
targetCPU := int32(execStrategy.TargetCPUPercent)
hpa := &asv1.HorizontalPodAutoscaler{
ObjectMeta: metav1.ObjectMeta{
Name: hpaName,
Labels: deployLabels,
Annotations: deployAnnotations,
},
Spec: asv1.HorizontalPodAutoscalerSpec{
ScaleTargetRef: asv1.CrossVersionObjectReference{
Kind: DeploymentKind,
Name: depl.ObjectMeta.Name,
APIVersion: DeploymentVersion,
},
MinReplicas: &minRepl,
MaxReplicas: maxRepl,
TargetCPUUtilizationPercentage: &targetCPU,
},
}
existingHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(ctx, hpaName, metav1.GetOptions{})
if err == nil {
// to adopt orphan service
if existingHpa.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID {
existingHpa.Annotations = hpa.Annotations
existingHpa.Labels = hpa.Labels
existingHpa.Spec = hpa.Spec
existingHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(ctx, existingHpa, metav1.UpdateOptions{})
if err != nil {
logger.Warn("error adopting HPA", zap.Error(err),
zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace))
return nil, err
}
}
return existingHpa, err
} else if k8s_err.IsNotFound(err) {
cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(ctx, hpa, metav1.CreateOptions{})
if err != nil {
if k8s_err.IsAlreadyExists(err) {
cHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(ctx, hpaName, metav1.GetOptions{})
}
if err != nil {
return nil, err
}
}
otelUtils.SpanTrackEvent(ctx, "createdService", otelUtils.GetAttributesForHPA(cHpa)...)
return cHpa, nil
}
return nil, err
}
func (deploy *NewDeploy) getHpa(ctx context.Context, ns, name string) (*asv1.HorizontalPodAutoscaler, error) {
return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(ctx, name, metav1.GetOptions{})
}
func (deploy *NewDeploy) updateHpa(ctx context.Context, hpa *asv1.HorizontalPodAutoscaler) error {
_, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(ctx, hpa, metav1.UpdateOptions{})
return err
}
func (deploy *NewDeploy) deleteHpa(ctx context.Context, ns string, name string) error {
return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(ctx, name, metav1.DeleteOptions{})
}
func (deploy *NewDeploy) createOrGetSvc(ctx context.Context, deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) {
logger := otelUtils.LoggerWithTraceID(ctx, deploy.logger)
service := &apiv1.Service{
@@ -554,7 +466,7 @@ func (deploy *NewDeploy) cleanupNewdeploy(ctx context.Context, ns string, name s
result = multierror.Append(result, err)
}
err = deploy.deleteHpa(ctx, ns, name)
err = deploy.hpaops.DeleteHpa(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
deploy.logger.Error("error deleting HPA for newdeploy function",
zap.Error(err),
@@ -29,6 +29,7 @@ import (
"github.com/pkg/errors"
"go.uber.org/zap"
autoscalingv1 "k8s.io/api/autoscaling/v1"
asv2beta2 "k8s.io/api/autoscaling/v2beta2"
apiv1 "k8s.io/api/core/v1"
k8sErrs "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -46,6 +47,7 @@ import (
"github.com/fission/fission/pkg/executor/fscache"
"github.com/fission/fission/pkg/executor/metrics"
"github.com/fission/fission/pkg/executor/reaper"
hpautils "github.com/fission/fission/pkg/executor/util/hpa"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/generated/clientset/versioned"
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
@@ -82,6 +84,8 @@ type (
deplListerSynced k8sCache.InformerSynced
svcListerSynced k8sCache.InformerSynced
hpaops *hpautils.HpaOperations
}
)
@@ -123,6 +127,8 @@ func MakeNewDeploy(
useIstio: enableIstio,
defaultIdlePodReapTime: 2 * time.Minute,
hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID),
}
nd.deplLister = deplInformer.Lister()
@@ -434,7 +440,7 @@ func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscac
return nil, errors.Wrapf(err, "error creating deployment %v", objName)
}
hpa, err := deploy.createOrGetHpa(ctx, objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
hpa, err := deploy.hpaops.CreateOrGetHpa(ctx, objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
if err != nil {
deploy.logger.Error("error creating HPA", zap.Error(err), zap.String("hpa", objName))
go cleanupFunc(ns, objName)
@@ -541,7 +547,7 @@ func (deploy *NewDeploy) updateFunction(ctx context.Context, oldFn *fv1.Function
return err
}
hpa, err := deploy.getHpa(ctx, ns, fsvc.Name)
hpa, err := deploy.hpaops.GetHpa(ctx, ns, fsvc.Name)
if err != nil {
deploy.updateStatus(oldFn, err, "error getting HPA while updating function")
return err
@@ -562,12 +568,13 @@ func (deploy *NewDeploy) updateFunction(ctx context.Context, oldFn *fv1.Function
if newFn.Spec.InvokeStrategy.ExecutionStrategy.TargetCPUPercent != oldFn.Spec.InvokeStrategy.ExecutionStrategy.TargetCPUPercent {
targetCpupercent := int32(newFn.Spec.InvokeStrategy.ExecutionStrategy.TargetCPUPercent)
hpa.Spec.TargetCPUUtilizationPercentage = &targetCpupercent
hpaMetric := hpautils.ConvertTargetCPUToCustomMetric(targetCpupercent)
hpa.Spec.Metrics = []asv2beta2.MetricSpec{hpaMetric}
hpaChanged = true
}
if hpaChanged {
err := deploy.updateHpa(ctx, hpa)
err := deploy.hpaops.UpdateHpa(ctx, hpa)
if err != nil {
deploy.updateStatus(oldFn, err, "error updating HPA while updating function")
return err
+21 -21
View File
@@ -23,7 +23,7 @@ import (
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
@@ -32,21 +32,21 @@ import (
)
var (
deletePropagation = meta_v1.DeletePropagationBackground
delOpt = meta_v1.DeleteOptions{PropagationPolicy: &deletePropagation}
deletePropagation = metav1.DeletePropagationBackground
delOpt = metav1.DeleteOptions{PropagationPolicy: &deletePropagation}
)
// CleanupKubeObject deletes given kubernetes object
func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient kubernetes.Interface, kubeobj *apiv1.ObjectReference) {
switch strings.ToLower(kubeobj.Kind) {
case "pod":
err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{})
err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up pod", zap.Error(err), zap.String("pod", kubeobj.Name))
}
case "service":
err := kubeClient.CoreV1().Services(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{})
err := kubeClient.CoreV1().Services(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up service", zap.Error(err), zap.String("service", kubeobj.Name))
}
@@ -58,7 +58,7 @@ func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient kuber
}
case "horizontalpodautoscaler":
err := kubeClient.AutoscalingV1().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{})
err := kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(ctx, kubeobj.Name, metav1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up horizontalpodautoscaler", zap.Error(err), zap.String("horizontalpodautoscaler", kubeobj.Name))
}
@@ -70,8 +70,8 @@ func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient kuber
}
// CleanupDeployments deletes deployment(s) for a given instanceID
func CleanupDeployments(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(ctx, listOps)
func CleanupDeployments(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps metav1.ListOptions) error {
deploymentList, err := client.AppsV1().Deployments(metav1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -97,8 +97,8 @@ func CleanupDeployments(ctx context.Context, logger *zap.Logger, client kubernet
}
// CleanupPods deletes pod(s) for a given instanceID
func CleanupPods(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(ctx, listOps)
func CleanupPods(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps metav1.ListOptions) error {
podList, err := client.CoreV1().Pods(metav1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -110,7 +110,7 @@ func CleanupPods(ctx context.Context, logger *zap.Logger, client kubernetes.Inte
}
if ok && id != instanceID {
logger.Info("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name))
err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(ctx, pod.ObjectMeta.Name, meta_v1.DeleteOptions{})
err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(ctx, pod.ObjectMeta.Name, metav1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up pod",
zap.Error(err),
@@ -124,8 +124,8 @@ func CleanupPods(ctx context.Context, logger *zap.Logger, client kubernetes.Inte
}
// CleanupServices deletes service(s) for a given instanceID
func CleanupServices(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(ctx, listOps)
func CleanupServices(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps metav1.ListOptions) error {
svcList, err := client.CoreV1().Services(metav1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -137,7 +137,7 @@ func CleanupServices(ctx context.Context, logger *zap.Logger, client kubernetes.
}
if ok && id != instanceID {
logger.Info("cleaning up service", zap.String("service", svc.ObjectMeta.Name))
err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(ctx, svc.ObjectMeta.Name, meta_v1.DeleteOptions{})
err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(ctx, svc.ObjectMeta.Name, metav1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up service",
zap.Error(err),
@@ -151,8 +151,8 @@ func CleanupServices(ctx context.Context, logger *zap.Logger, client kubernetes.
}
// CleanupHpa deletes horizontal pod autoscaler(s) for a given instanceID
func CleanupHpa(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(ctx, listOps)
func CleanupHpa(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps metav1.ListOptions) error {
hpaList, err := client.AutoscalingV2beta2().HorizontalPodAutoscalers(metav1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -165,7 +165,7 @@ func CleanupHpa(ctx context.Context, logger *zap.Logger, client kubernetes.Inter
}
if ok && id != instanceID {
logger.Info("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name))
err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(ctx, hpa.ObjectMeta.Name, meta_v1.DeleteOptions{})
err := client.AutoscalingV2beta2().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(ctx, hpa.ObjectMeta.Name, metav1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up HPA",
zap.Error(err),
@@ -187,7 +187,7 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne
logger.Debug("starting cleanupRoleBindings cycle")
// get all rolebindings ( just to be efficient, one call to kubernetes )
rbList, err := client.RbacV1().RoleBindings(meta_v1.NamespaceAll).List(ctx, meta_v1.ListOptions{})
rbList, err := client.RbacV1().RoleBindings(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil {
// something wrong, but next iteration hopefully succeeds
logger.Error("error listing role bindings in all namespaces", zap.Error(err))
@@ -208,7 +208,7 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne
// in order to find out if there are any functions that need this role-binding in role-binding namespace,
// we can list the functions once per role-binding.
funcList, err := fissionClient.CoreV1().Functions(roleBinding.Namespace).List(ctx, meta_v1.ListOptions{})
funcList, err := fissionClient.CoreV1().Functions(roleBinding.Namespace).List(ctx, metav1.ListOptions{})
if err != nil {
logger.Error("error fetching function list in namespace", zap.Error(err), zap.String("namespace", roleBinding.Namespace))
continue
@@ -235,7 +235,7 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne
isInReservedNS := false
if subj.Namespace == functionNs ||
subj.Namespace == envBuilderNs {
saNs = meta_v1.NamespaceDefault
saNs = metav1.NamespaceDefault
isInReservedNS = true
}
@@ -259,7 +259,7 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne
// else if its a secret-configmap-rb, we have only one SA which is fission-fetcher
if roleBinding.Name == fv1.PackageGetterRB {
// check if there is an env obj in saNs
envList, err := fissionClient.CoreV1().Environments(saNs).List(ctx, meta_v1.ListOptions{})
envList, err := fissionClient.CoreV1().Environments(saNs).List(ctx, metav1.ListOptions{})
if err != nil {
logger.Error("error fetching environment list in service account namespace", zap.Error(err), zap.String("namespace", saNs))
continue
+153
View File
@@ -0,0 +1,153 @@
/*
Copyright 2022 The Fission Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package hpa
import (
"context"
"errors"
"go.uber.org/zap"
appsv1 "k8s.io/api/apps/v1"
asv2beta2 "k8s.io/api/autoscaling/v2beta2"
corev1 "k8s.io/api/core/v1"
k8s_err "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
otelUtils "github.com/fission/fission/pkg/utils/otel"
)
// Deployment Constants
const (
DeploymentKind = "Deployment"
DeploymentVersion = "apps/v1"
)
type HpaOperations struct {
logger *zap.Logger
kubernetesClient kubernetes.Interface
instanceID string
}
func NewHpaOperations(logger *zap.Logger, kubernetesClient kubernetes.Interface, instanceID string) *HpaOperations {
return &HpaOperations{
logger: logger,
kubernetesClient: kubernetesClient,
instanceID: instanceID,
}
}
func ConvertTargetCPUToCustomMetric(targetCPUVal int32) asv2beta2.MetricSpec {
return asv2beta2.MetricSpec{
Type: asv2beta2.ResourceMetricSourceType,
Resource: &asv2beta2.ResourceMetricSource{
Name: corev1.ResourceCPU,
Target: asv2beta2.MetricTarget{
Type: asv2beta2.UtilizationMetricType,
AverageUtilization: &targetCPUVal,
},
},
}
}
func getScaleTargetRef(deployment *appsv1.Deployment) asv2beta2.CrossVersionObjectReference {
return asv2beta2.CrossVersionObjectReference{
APIVersion: DeploymentVersion,
Kind: DeploymentKind,
Name: deployment.ObjectMeta.Name,
}
}
func (hpaops *HpaOperations) CreateOrGetHpa(ctx context.Context, hpaName string, execStrategy *fv1.ExecutionStrategy,
depl *appsv1.Deployment, deployLabels map[string]string, deployAnnotations map[string]string) (*asv2beta2.HorizontalPodAutoscaler, error) {
if depl == nil {
return nil, errors.New("failed to create HPA, found empty deployment")
}
logger := otelUtils.LoggerWithTraceID(ctx, hpaops.logger)
minRepl := int32(execStrategy.MinScale)
if minRepl == 0 {
minRepl = 1
}
maxRepl := int32(execStrategy.MaxScale)
if maxRepl == 0 {
maxRepl = minRepl
}
targetCPU := int32(execStrategy.TargetCPUPercent)
var hpaMetrics []asv2beta2.MetricSpec
if targetCPU > 0 {
hpaMetrics = append(hpaMetrics, ConvertTargetCPUToCustomMetric(targetCPU))
}
hpa := &asv2beta2.HorizontalPodAutoscaler{
ObjectMeta: metav1.ObjectMeta{
Name: hpaName,
Labels: deployLabels,
Annotations: deployAnnotations,
},
Spec: asv2beta2.HorizontalPodAutoscalerSpec{
ScaleTargetRef: getScaleTargetRef(depl),
MinReplicas: &minRepl,
MaxReplicas: maxRepl,
Metrics: hpaMetrics,
},
}
existingHpa, err := hpaops.GetHpa(ctx, depl.ObjectMeta.Namespace, hpaName)
if err == nil {
// to adopt orphan service
if existingHpa.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != hpaops.instanceID {
existingHpa.Annotations = hpa.Annotations
existingHpa.Labels = hpa.Labels
existingHpa.Spec = hpa.Spec
existingHpa, err = hpaops.kubernetesClient.AutoscalingV2beta2().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(ctx, existingHpa, metav1.UpdateOptions{})
if err != nil {
logger.Warn("error adopting HPA", zap.Error(err),
zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace))
return nil, err
}
}
return existingHpa, err
} else if k8s_err.IsNotFound(err) {
cHpa, err := hpaops.kubernetesClient.AutoscalingV2beta2().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(ctx, hpa, metav1.CreateOptions{})
if err != nil {
if k8s_err.IsAlreadyExists(err) {
cHpa, err = hpaops.kubernetesClient.AutoscalingV2beta2().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(ctx, hpaName, metav1.GetOptions{})
}
if err != nil {
return nil, err
}
}
otelUtils.SpanTrackEvent(ctx, "hpaCreated", otelUtils.GetAttributesForHPA(cHpa)...)
return cHpa, nil
}
return nil, err
}
func (hpaops *HpaOperations) GetHpa(ctx context.Context, ns, name string) (*asv2beta2.HorizontalPodAutoscaler, error) {
return hpaops.kubernetesClient.AutoscalingV2beta2().HorizontalPodAutoscalers(ns).Get(ctx, name, metav1.GetOptions{})
}
func (hpaops *HpaOperations) UpdateHpa(ctx context.Context, hpa *asv2beta2.HorizontalPodAutoscaler) error {
_, err := hpaops.kubernetesClient.AutoscalingV2beta2().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(ctx, hpa, metav1.UpdateOptions{})
return err
}
func (hpaops *HpaOperations) DeleteHpa(ctx context.Context, ns string, name string) error {
return hpaops.kubernetesClient.AutoscalingV2beta2().HorizontalPodAutoscalers(ns).Delete(ctx, name, metav1.DeleteOptions{})
}
+147
View File
@@ -0,0 +1,147 @@
/*
Copyright 2022 The Fission Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package hpa
import (
"context"
"strings"
"testing"
"github.com/dchest/uniuri"
appsv1 "k8s.io/api/apps/v1"
asv2beta2 "k8s.io/api/autoscaling/v2beta2"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes/fake"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/utils/loggerfactory"
)
func TestConvertTargetCPUToCustomMetric(t *testing.T) {
metricSpec := ConvertTargetCPUToCustomMetric(50)
if metricSpec.Type != asv2beta2.ResourceMetricSourceType {
t.Errorf("Expected metric type to be Resource, got %v", metricSpec.Type)
}
if metricSpec.Resource.Name != corev1.ResourceCPU {
t.Errorf("Expected metric name to be cpu, got %v", metricSpec.Resource.Name)
}
if metricSpec.Resource.Target.Type != asv2beta2.UtilizationMetricType {
t.Errorf("Expected metric target type to be Utilization, got %v", metricSpec.Resource.Target.Type)
}
if metricSpec.Resource.Target.AverageUtilization == nil {
t.Errorf("Expected metric target average utilization to be set, got nil")
}
}
func TestHpaOps(t *testing.T) {
logger := loggerfactory.GetLogger()
kubernetesClient := fake.NewSimpleClientset()
instanceID := strings.ToLower(uniuri.NewLen(8))
ns := "test-namespace"
hpaops := NewHpaOperations(logger, kubernetesClient, instanceID)
if hpaops.instanceID != instanceID {
t.Errorf("Expected instanceID to be %v, got %v", instanceID, hpaops.instanceID)
}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
deployLabels := map[string]string{
"test-label": "test-label-value",
}
deployAnnotations := map[string]string{
"test-annotation": "test-annotation-value",
}
// Test CreateHPA
hpa, err := hpaops.CreateOrGetHpa(ctx, "test-hpa",
&fv1.ExecutionStrategy{
ExecutorType: fv1.ExecutorTypeNewdeploy,
MinScale: 1,
MaxScale: 5,
TargetCPUPercent: 50,
SpecializationTimeout: 300,
},
&appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: "test-deployment",
Namespace: ns,
},
},
deployLabels,
deployAnnotations)
if err != nil {
t.Errorf("Unexpected error: %v", err)
}
if *hpa.Spec.MinReplicas != 1 {
t.Errorf("Expected min replicas to be 1, got %v", hpa.Spec.MinReplicas)
}
if hpa.Spec.MaxReplicas != 5 {
t.Errorf("Expected max replicas to be 5, got %v", hpa.Spec.MaxReplicas)
}
if hpa.Spec.Metrics[0].Type != asv2beta2.ResourceMetricSourceType {
t.Errorf("Expected metric type to be Resource, got %v", hpa.Spec.Metrics[0].Type)
}
if hpa.Spec.Metrics[0].Resource.Name != corev1.ResourceCPU {
t.Errorf("Expected metric name to be cpu, got %v", hpa.Spec.Metrics[0].Resource.Name)
}
if hpa.Spec.Metrics[0].Resource.Target.Type != asv2beta2.UtilizationMetricType {
t.Errorf("Expected metric target type to be Utilization, got %v", hpa.Spec.Metrics[0].Resource.Target.Type)
}
if hpa.Spec.Metrics[0].Resource.Target.AverageUtilization == nil {
t.Errorf("Expected metric target average utilization to be set, got nil")
}
if *hpa.Spec.Metrics[0].Resource.Target.AverageUtilization != 50 {
t.Errorf("Expected metric target average utilization to be 50, got %v", *hpa.Spec.Metrics[0].Resource.Target.AverageUtilization)
}
if hpa.ObjectMeta.Labels["test-label"] != "test-label-value" {
t.Errorf("Expected label to be set, got %v", hpa.ObjectMeta.Labels["test-label"])
}
if hpa.ObjectMeta.Annotations["test-annotation"] != "test-annotation-value" {
t.Errorf("Expected annotation to be set, got %v", hpa.ObjectMeta.Annotations["test-annotation"])
}
hpa, err = hpaops.GetHpa(ctx, ns, "test-hpa")
if err != nil {
t.Errorf("Unexpected error: %v", err)
}
hpa.Spec.MaxReplicas = 10
// Test UpdateHPA
err = hpaops.UpdateHpa(ctx, hpa)
if err != nil {
t.Errorf("Unexpected error: %v", err)
}
hpa, err = hpaops.GetHpa(ctx, ns, "test-hpa")
if err != nil {
t.Errorf("Unexpected error: %v", err)
}
if hpa.Spec.MaxReplicas != 10 {
t.Errorf("Expected max replicas to be 10, got %v", hpa.Spec.MaxReplicas)
}
// Test DeleteHPA
err = hpaops.DeleteHpa(ctx, ns, "test-hpa")
if err != nil {
t.Errorf("Unexpected error: %v", err)
}
_, err = hpaops.GetHpa(ctx, ns, "test-hpa")
if err == nil {
t.Errorf("Expected error, got nil")
}
}
@@ -116,7 +116,7 @@ func (res KubernetesObjectDumper) Dump(dumpDir string) {
}
case KubernetesHPA:
objs, err := res.client.AutoscalingV2beta1().HorizontalPodAutoscalers(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector})
objs, err := res.client.AutoscalingV2beta2().HorizontalPodAutoscalers(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector})
if err != nil {
console.Error(fmt.Sprintf("Error getting %v list with selector %v: %v", res.objType, res.selector, err))
return
+4 -4
View File
@@ -28,7 +28,7 @@ import (
"github.com/pkg/errors"
v1 "k8s.io/api/core/v1"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/tools/portforward"
"k8s.io/client-go/transport/spdy"
@@ -120,12 +120,12 @@ func runPortForward(ctx context.Context, labelSelector string, localPort string,
// if namespace is unset, try to find a pod in any namespace
if len(ns) == 0 {
ns = meta_v1.NamespaceAll
ns = metav1.NamespaceAll
}
// get the pod; if there is more than one, ask the user to disambiguate
podList, err := clientset.CoreV1().Pods(ns).
List(ctx, meta_v1.ListOptions{LabelSelector: labelSelector})
List(ctx, metav1.ListOptions{LabelSelector: labelSelector})
if err != nil {
return nil, nil, errors.Wrapf(err, "error getting pod for port-forwarding with label selector %v", labelSelector)
} else if len(podList.Items) == 0 {
@@ -171,7 +171,7 @@ func runPortForward(ctx context.Context, labelSelector string, localPort string,
// get the service and the target port
svcs, err := clientset.CoreV1().Services(podNameSpace).
List(ctx, meta_v1.ListOptions{LabelSelector: labelSelector})
List(ctx, metav1.ListOptions{LabelSelector: labelSelector})
if err != nil {
return nil, nil, errors.Wrapf(err, "Error getting %v service", labelSelector)
}
+3 -2
View File
@@ -6,7 +6,8 @@ import (
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
appsv1 "k8s.io/api/apps/v1"
asv1 "k8s.io/api/autoscaling/v1"
asv2beta2 "k8s.io/api/autoscaling/v2beta2"
apiv1 "k8s.io/api/core/v1"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
@@ -74,7 +75,7 @@ func GetAttributesForDeployment(deployment *appsv1.Deployment) []attribute.KeyVa
}
}
func GetAttributesForHPA(hpa *asv1.HorizontalPodAutoscaler) []attribute.KeyValue {
func GetAttributesForHPA(hpa *asv2beta2.HorizontalPodAutoscaler) []attribute.KeyValue {
if hpa == nil {
return []attribute.KeyValue{}
}