diff --git a/pkg/executor/executortype/container/common.go b/pkg/executor/executortype/container/common.go index d5dfa0ea..8a79e26a 100644 --- a/pkg/executor/executortype/container/common.go +++ b/pkg/executor/executortype/container/common.go @@ -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), diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index 4c507cb2..58cac49c 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -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 diff --git a/pkg/executor/executortype/container/hpa.go b/pkg/executor/executortype/container/hpa.go deleted file mode 100644 index 9be8763d..00000000 --- a/pkg/executor/executortype/container/hpa.go +++ /dev/null @@ -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{}) -} diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index 2c269831..56cd7f50 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -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), diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 5308b9ca..fe74b892 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -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 diff --git a/pkg/executor/reaper/reaper.go b/pkg/executor/reaper/reaper.go index da828783..2c028e0c 100644 --- a/pkg/executor/reaper/reaper.go +++ b/pkg/executor/reaper/reaper.go @@ -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 diff --git a/pkg/executor/util/hpa/hpa.go b/pkg/executor/util/hpa/hpa.go new file mode 100644 index 00000000..6729f3ac --- /dev/null +++ b/pkg/executor/util/hpa/hpa.go @@ -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{}) +} diff --git a/pkg/executor/util/hpa/hpa_test.go b/pkg/executor/util/hpa/hpa_test.go new file mode 100644 index 00000000..d648887b --- /dev/null +++ b/pkg/executor/util/hpa/hpa_test.go @@ -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") + } +} diff --git a/pkg/fission-cli/cmd/support/resources/kubernetes.go b/pkg/fission-cli/cmd/support/resources/kubernetes.go index ad655fc7..381b31bc 100644 --- a/pkg/fission-cli/cmd/support/resources/kubernetes.go +++ b/pkg/fission-cli/cmd/support/resources/kubernetes.go @@ -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 diff --git a/pkg/fission-cli/util/portforward.go b/pkg/fission-cli/util/portforward.go index b394e631..2a48fa10 100644 --- a/pkg/fission-cli/util/portforward.go +++ b/pkg/fission-cli/util/portforward.go @@ -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) } diff --git a/pkg/utils/otel/attributes.go b/pkg/utils/otel/attributes.go index 162a8122..2c8177b1 100644 --- a/pkg/utils/otel/attributes.go +++ b/pkg/utils/otel/attributes.go @@ -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{} }