Add correct context required in executor (#2175)

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2021-08-25 13:16:26 +05:30
committed by GitHub
parent d6b47c1a4e
commit 1df59316e7
20 changed files with 243 additions and 227 deletions
+1 -1
View File
@@ -47,7 +47,7 @@ import (
func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient *crd.FissionClient, envBuilderNamespace string,
storageSvcUrl string, pkg *fv1.Package) (uploadResp *fetcher.ArchiveUploadResponse, buildLogs string, err error) {
env, err := fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(context.TODO(), pkg.Spec.Environment.Name, metav1.GetOptions{})
env, err := fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(ctx, pkg.Spec.Environment.Name, metav1.GetOptions{})
if err != nil {
e := "error getting environment CRD info"
logger.Error(e, zap.Error(err))
+7 -7
View File
@@ -73,7 +73,7 @@ func makePackageWatcher(logger *zap.Logger, fissionClient *crd.FissionClient, k8
// 5. Update package resource in package ref of functions that share the same package
// 6. Update package status to succeed state
// *. Update package status to failed state,if any one of steps above failed/time out
func (pkgw *packageWatcher) build(srcpkg *fv1.Package) {
func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
// Ignore duplicate build requests
key := fmt.Sprintf("%v-%v", srcpkg.ObjectMeta.Name, srcpkg.ObjectMeta.ResourceVersion)
_, err := pkgw.buildCache.Set(key, srcpkg)
@@ -95,7 +95,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) {
return
}
env, err := pkgw.fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(context.TODO(), pkg.Spec.Environment.Name, metav1.GetOptions{})
env, err := pkgw.fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(ctx, pkg.Spec.Environment.Name, metav1.GetOptions{})
if k8serrors.IsNotFound(err) {
e := "environment does not exist"
pkgw.logger.Error(e, zap.String("environment", pkg.Spec.Environment.Name))
@@ -167,7 +167,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) {
// Add the package getter rolebinding to builder sa
// we continue here if role binding was not setup successfully. this is because without this, the fetcher wont be able to fetch the source pkg into the container and
// the build will fail eventually
err := utils.SetupRoleBinding(pkgw.logger, pkgw.k8sClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionBuilderSA, builderNs)
err := utils.SetupRoleBinding(ctx, pkgw.logger, pkgw.k8sClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionBuilderSA, builderNs)
if err != nil {
pkgw.logger.Error("error setting up role binding for package",
zap.Error(err),
@@ -181,7 +181,6 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) {
zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace)))
}
ctx := context.Background()
uploadResp, buildLogs, err := buildPackage(ctx, pkgw.logger, pkgw.fissionClient, builderNs, pkgw.storageSvcUrl, pkg)
if err != nil {
pkgw.logger.Error("error building package", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name))
@@ -200,7 +199,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) {
pkgw.logger.Info("starting package info update", zap.String("package_name", pkg.ObjectMeta.Name))
fnList, err := pkgw.fissionClient.CoreV1().
Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil {
e := "error getting function list"
pkgw.logger.Error(e, zap.Error(err))
@@ -224,7 +223,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) {
fn.Spec.Package.PackageRef.ResourceVersion != pkg.ObjectMeta.ResourceVersion {
fn.Spec.Package.PackageRef.ResourceVersion = pkg.ObjectMeta.ResourceVersion
// update CRD
_, err = pkgw.fissionClient.CoreV1().Functions(fn.ObjectMeta.Namespace).Update(context.TODO(), &fn, metav1.UpdateOptions{})
_, err = pkgw.fissionClient.CoreV1().Functions(fn.ObjectMeta.Namespace).Update(ctx, &fn, metav1.UpdateOptions{})
if err != nil {
e := "error updating function package resource version"
pkgw.logger.Error(e, zap.Error(err))
@@ -296,7 +295,8 @@ func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandl
}
// Only build pending state packages.
if pkg.Status.BuildStatus == fv1.BuildStatusPending {
go pkgw.build(pkg)
ctx := context.Background()
go pkgw.build(ctx, pkg)
}
}
return k8sCache.ResourceEventHandlerFuncs{
@@ -65,10 +65,10 @@ func (cn *Container) getResources(fn *fv1.Function) apiv1.ResourceRequirements {
}
// cleanupContainer cleans all kubernetes objects related to function
func (cn *Container) cleanupContainer(ns string, name string) error {
func (cn *Container) cleanupContainer(ctx context.Context, ns string, name string) error {
result := &multierror.Error{}
err := cn.deleteSvc(ns, name)
err := cn.deleteSvc(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
cn.logger.Error("error deleting service for Container function",
zap.Error(err),
@@ -77,7 +77,7 @@ func (cn *Container) cleanupContainer(ns string, name string) error {
result = multierror.Append(result, err)
}
err = cn.deleteHpa(ns, name)
err = cn.deleteHpa(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
cn.logger.Error("error deleting HPA for Container function",
zap.Error(err),
@@ -86,7 +86,7 @@ func (cn *Container) cleanupContainer(ns string, name string) error {
result = multierror.Append(result, err)
}
err = cn.deleteDeployment(ns, name)
err = cn.deleteDeployment(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
cn.logger.Error("error deleting deployment for Container function",
zap.Error(err),
@@ -174,7 +174,7 @@ func (caaf *Container) TapService(ctx context.Context, svcHost string) error {
return nil
}
func (caaf *Container) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service, error) {
func (caaf *Container) getServiceInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Service, error) {
item, exists, err := utils.GetCachedItem(obj, caaf.serviceInformer)
if err != nil || !exists {
@@ -183,7 +183,7 @@ func (caaf *Container) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service
zap.Bool("exists", exists),
zap.Error(err),
)
service, err := caaf.kubernetesClient.CoreV1().Services(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{})
service, err := caaf.kubernetesClient.CoreV1().Services(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
return service, err
}
@@ -191,7 +191,7 @@ func (caaf *Container) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service
return service, nil
}
func (caaf *Container) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Deployment, error) {
func (caaf *Container) getDeploymentInfo(ctx context.Context, obj apiv1.ObjectReference) (*appsv1.Deployment, error) {
item, exists, err := utils.GetCachedItem(obj, caaf.deploymentInformer)
if err != nil || !exists {
@@ -200,7 +200,7 @@ func (caaf *Container) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Dep
zap.Bool("exists", exists),
zap.Error(err),
)
deployment, err := caaf.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{})
deployment, err := caaf.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
return deployment, err
}
@@ -222,7 +222,7 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool
}
for _, obj := range fsvc.KubernetesObjects {
if strings.ToLower(obj.Kind) == "service" {
_, err := caaf.getServiceInfo(obj)
_, err := caaf.getServiceInfo(ctx, obj)
if err != nil {
if !k8sErrs.IsNotFound(err) {
caaf.logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err))
@@ -231,7 +231,7 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool
}
} else if strings.ToLower(obj.Kind) == "deployment" {
currentDeploy, err := caaf.getDeploymentInfo(obj)
currentDeploy, err := caaf.getDeploymentInfo(ctx, obj)
if err != nil {
if !k8sErrs.IsNotFound(err) {
caaf.logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err))
@@ -317,17 +317,17 @@ func (caaf *Container) CleanupOldExecutorObjects(ctx context.Context) {
LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypeContainer)}).AsSelector().String(),
}
err := reaper.CleanupHpa(caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts)
err := reaper.CleanupHpa(ctx, caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupDeployments(caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts)
err = reaper.CleanupDeployments(ctx, caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupServices(caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts)
err = reaper.CleanupServices(ctx, caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
@@ -366,11 +366,11 @@ func (caaf *Container) createFunction(ctx context.Context, fn *fv1.Function) (*f
return fsvc, err
}
func (caaf *Container) deleteFunction(fn *fv1.Function) error {
func (caaf *Container) deleteFunction(ctx context.Context, fn *fv1.Function) error {
if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypeContainer {
return nil
}
err := caaf.fnDelete(fn)
err := caaf.fnDelete(ctx, fn)
if err != nil {
err = errors.Wrapf(err, "error deleting kubernetes objects of function %v", fn.ObjectMeta)
}
@@ -379,7 +379,7 @@ func (caaf *Container) deleteFunction(fn *fv1.Function) error {
func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
cleanupFunc := func(ns string, name string) {
err := caaf.cleanupContainer(ns, name)
err := caaf.cleanupContainer(ctx, ns, name)
if err != nil {
caaf.logger.Error("received error while cleaning function resources",
zap.String("namespace", ns), zap.String("name", name))
@@ -401,7 +401,7 @@ func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache
// Since Container waits for pods of deployment to be ready,
// change the order of kubeObject creation (create service first,
// then deployment) to take advantage of waiting time.
svc, err := caaf.createOrGetSvc(fn, deployLabels, deployAnnotations, objName, ns)
svc, err := caaf.createOrGetSvc(ctx, fn, deployLabels, deployAnnotations, objName, ns)
if err != nil {
caaf.logger.Error("error creating service", zap.Error(err), zap.String("service", objName))
go cleanupFunc(ns, objName)
@@ -416,7 +416,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(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
hpa, err := caaf.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)
@@ -488,7 +488,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function,
caaf.logger.Info("function does not use new deployment executor anymore, deleting resources",
zap.Any("function", newFn))
// IMP - pass the oldFn, as the new/modified function is not in cache
return caaf.deleteFunction(oldFn)
return caaf.deleteFunction(ctx, oldFn)
}
// Executor type changed to Container from something else
@@ -519,7 +519,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function,
return err
}
hpa, err := caaf.getHpa(ns, fsvc.Name)
hpa, err := caaf.getHpa(ctx, ns, fsvc.Name)
if err != nil {
caaf.updateStatus(oldFn, err, "error getting HPA while updating function")
return err
@@ -545,7 +545,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function,
}
if hpaChanged {
err := caaf.updateHpa(hpa)
err := caaf.updateHpa(ctx, hpa)
if err != nil {
caaf.updateStatus(oldFn, err, "error updating HPA while updating function")
return err
@@ -623,7 +623,7 @@ func (caaf *Container) updateFuncDeployment(ctx context.Context, fn *fv1.Functio
return err
}
err = caaf.updateDeployment(newDeployment, ns)
err = caaf.updateDeployment(ctx, newDeployment, ns)
if err != nil {
caaf.updateStatus(fn, err, "failed to update deployment while updating function")
return err
@@ -632,7 +632,7 @@ func (caaf *Container) updateFuncDeployment(ctx context.Context, fn *fv1.Functio
return nil
}
func (caaf *Container) fnDelete(fn *fv1.Function) error {
func (caaf *Container) fnDelete(ctx context.Context, fn *fv1.Function) error {
multierr := &multierror.Error{}
// GetByFunction uses resource version as part of cache key, however,
@@ -661,7 +661,7 @@ func (caaf *Container) fnDelete(fn *fv1.Function) error {
ns = fn.ObjectMeta.Namespace
}
err = caaf.cleanupContainer(ns, objName)
err = caaf.cleanupContainer(ctx, ns, objName)
multierr = multierror.Append(multierr, err)
return multierr.ErrorOrNil()
@@ -716,6 +716,7 @@ func (caaf *Container) updateStatus(fn *fv1.Function, err error, message string)
// idleObjectReaper reaps objects after certain idle time
func (caaf *Container) idleObjectReaper() {
ctx := context.Background()
pollSleep := 5 * time.Second
for {
@@ -734,7 +735,7 @@ func (caaf *Container) idleObjectReaper() {
continue
}
fn, err := caaf.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(context.TODO(), fsvc.Function.Name, metav1.GetOptions{})
fn, err := caaf.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(ctx, fsvc.Function.Name, metav1.GetOptions{})
if err != nil {
// CaaF manager handles the function delete event and clean cache/kubeobjs itself,
// so we ignore the not found error for functions with CaaF executor type here.
@@ -762,7 +763,7 @@ func (caaf *Container) idleObjectReaper() {
}
currentDeploy, err := caaf.kubernetesClient.AppsV1().
Deployments(deployObj.Namespace).Get(context.TODO(), deployObj.Name, metav1.GetOptions{})
Deployments(deployObj.Namespace).Get(ctx, deployObj.Name, metav1.GetOptions{})
if err != nil {
caaf.logger.Error("error getting function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
return
@@ -775,7 +776,7 @@ func (caaf *Container) idleObjectReaper() {
return
}
err = caaf.scaleDeployment(deployObj.Namespace, deployObj.Name, minScale)
err = caaf.scaleDeployment(ctx, deployObj.Namespace, deployObj.Name, minScale)
if err != nil {
caaf.logger.Error("error scaling down function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
}
@@ -75,14 +75,14 @@ func (cn *Container) createOrGetDeployment(ctx context.Context, fn *fv1.Function
}
if *existingDepl.Spec.Replicas < minScale {
err = cn.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale)
err = cn.scaleDeployment(ctx, existingDepl.Namespace, existingDepl.Name, minScale)
if err != nil {
cn.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.ObjectMeta.Name))
return nil, err
}
}
if existingDepl.Status.AvailableReplicas < minScale {
existingDepl, err = cn.waitForDeploy(existingDepl, minScale, specializationTimeout)
existingDepl, err = cn.waitForDeploy(ctx, existingDepl, minScale, specializationTimeout)
}
return existingDepl, err
@@ -102,28 +102,28 @@ func (cn *Container) createOrGetDeployment(ctx context.Context, fn *fv1.Function
}
}
if minScale > 0 {
depl, err = cn.waitForDeploy(depl, minScale, specializationTimeout)
depl, err = cn.waitForDeploy(ctx, depl, minScale, specializationTimeout)
}
return depl, err
}
return nil, err
}
func (cn *Container) updateDeployment(deployment *appsv1.Deployment, ns string) error {
_, err := cn.kubernetesClient.AppsV1().Deployments(ns).Update(context.TODO(), deployment, metav1.UpdateOptions{})
func (cn *Container) updateDeployment(ctx context.Context, deployment *appsv1.Deployment, ns string) error {
_, err := cn.kubernetesClient.AppsV1().Deployments(ns).Update(ctx, deployment, metav1.UpdateOptions{})
return err
}
func (cn *Container) deleteDeployment(ns string, name string) error {
func (cn *Container) deleteDeployment(ctx context.Context, ns string, name string) error {
// DeletePropagationBackground deletes the object immediately and dependent are deleted later
// DeletePropagationForeground not advisable; it marks for deleteion and API can still serve those objects
deletePropagation := metav1.DeletePropagationBackground
return cn.kubernetesClient.AppsV1().Deployments(ns).Delete(context.TODO(), name, metav1.DeleteOptions{
return cn.kubernetesClient.AppsV1().Deployments(ns).Delete(ctx, name, metav1.DeleteOptions{
PropagationPolicy: &deletePropagation,
})
}
func (cn *Container) waitForDeploy(depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) {
func (cn *Container) waitForDeploy(ctx context.Context, depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) {
oldStatus := depl.Status
// if no specializationTimeout is set, use default value
@@ -132,7 +132,7 @@ func (cn *Container) waitForDeploy(depl *appsv1.Deployment, replicas int32, spec
}
for i := 0; i < specializationTimeout; i++ {
latestDepl, err := cn.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(context.TODO(), depl.Name, metav1.GetOptions{})
latestDepl, err := cn.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(ctx, depl.Name, metav1.GetOptions{})
if err != nil {
return nil, err
}
@@ -196,7 +196,7 @@ func (cn *Container) getDeploymentSpec(ctx context.Context, fn *fv1.Function, ta
resources := cn.getResources(fn)
// Other executor types rely on Environments to add configmaps and secrets
envFromSources, err := util.ConvertConfigSecrets(fn, cn.kubernetesClient)
envFromSources, err := util.ConvertConfigSecrets(ctx, fn, cn.kubernetesClient)
if err != nil {
return nil, err
}
@@ -277,12 +277,12 @@ func (cn *Container) getDeploymentSpec(ctx context.Context, fn *fv1.Function, ta
return deployment, nil
}
func (caaf *Container) scaleDeployment(deplNS string, deplName string, replicas int32) error {
func (caaf *Container) scaleDeployment(ctx context.Context, deplNS string, deplName string, replicas int32) error {
caaf.logger.Info("scaling deployment",
zap.String("deployment", deplName),
zap.String("namespace", deplNS),
zap.Int32("replicas", replicas))
_, err := caaf.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(context.TODO(), deplName, &autoscalingv1.Scale{
_, err := caaf.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(ctx, deplName, &autoscalingv1.Scale{
ObjectMeta: metav1.ObjectMeta{
Name: deplName,
Namespace: deplNS,
@@ -55,7 +55,7 @@ func (caaf *Container) FuncInformerHandler(ctx context.Context) k8sCache.Resourc
go func() {
log := caaf.logger.With(zap.String("function_name", fn.ObjectMeta.Name), zap.String("function_namespace", fn.ObjectMeta.Namespace))
log.Debug("start function delete handler")
err := caaf.deleteFunction(fn)
err := caaf.deleteFunction(ctx, fn)
if err != nil {
log.Error("error deleting function", zap.Error(err))
}
+11 -11
View File
@@ -34,7 +34,7 @@ const (
DeploymentVersion = "apps/v1"
)
func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy,
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 {
@@ -69,14 +69,14 @@ func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionS
},
}
existingHpa, err := cn.getHpa(depl.ObjectMeta.Namespace, hpaName)
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(context.TODO(), existingHpa, metav1.UpdateOptions{})
existingHpa, err = cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(ctx, existingHpa, metav1.UpdateOptions{})
if err != nil {
cn.logger.Warn("error adopting HPA", zap.Error(err),
zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace))
@@ -85,10 +85,10 @@ func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionS
}
return existingHpa, err
} else if k8s_err.IsNotFound(err) {
cHpa, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(context.TODO(), hpa, metav1.CreateOptions{})
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(depl.ObjectMeta.Namespace, hpaName)
cHpa, err = cn.getHpa(ctx, depl.ObjectMeta.Namespace, hpaName)
}
if err != nil {
return nil, err
@@ -99,15 +99,15 @@ func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionS
return nil, err
}
func (cn *Container) getHpa(ns, name string) (*asv1.HorizontalPodAutoscaler, error) {
return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(context.TODO(), name, metav1.GetOptions{})
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(hpa *asv1.HorizontalPodAutoscaler) error {
_, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(context.TODO(), hpa, metav1.UpdateOptions{})
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(ns string, name string) error {
return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(context.TODO(), name, metav1.DeleteOptions{})
func (cn *Container) deleteHpa(ctx context.Context, ns string, name string) error {
return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(ctx, name, metav1.DeleteOptions{})
}
+7 -7
View File
@@ -42,7 +42,7 @@ func (cn *Container) getSvPort(fn *fv1.Function) (port int32, err error) {
return fn.Spec.PodSpec.Containers[0].Ports[0].ContainerPort, nil
}
func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) {
func (cn *Container) createOrGetSvc(ctx context.Context, fn *fv1.Function, deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) {
targetPort, err := cn.getSvPort(fn)
if err != nil {
return nil, err
@@ -66,7 +66,7 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st
},
}
existingSvc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{})
existingSvc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{})
if err == nil {
// to adopt orphan service
if existingSvc.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != cn.instanceID {
@@ -75,7 +75,7 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st
existingSvc.Spec.Ports = service.Spec.Ports
existingSvc.Spec.Selector = service.Spec.Selector
existingSvc.Spec.Type = service.Spec.Type
existingSvc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Update(context.TODO(), existingSvc, metav1.UpdateOptions{})
existingSvc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Update(ctx, existingSvc, metav1.UpdateOptions{})
if err != nil {
cn.logger.Warn("error adopting service", zap.Error(err),
zap.String("service", svcName), zap.String("ns", svcNamespace))
@@ -84,10 +84,10 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st
}
return existingSvc, err
} else if k8s_err.IsNotFound(err) {
svc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Create(context.TODO(), service, metav1.CreateOptions{})
svc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Create(ctx, service, metav1.CreateOptions{})
if err != nil {
if k8s_err.IsAlreadyExists(err) {
svc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{})
svc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{})
}
if err != nil {
return nil, err
@@ -98,6 +98,6 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st
return nil, err
}
func (cn *Container) deleteSvc(ns string, name string) error {
return cn.kubernetesClient.CoreV1().Services(ns).Delete(context.TODO(), name, metav1.DeleteOptions{})
func (cn *Container) deleteSvc(ctx context.Context, ns string, name string) error {
return cn.kubernetesClient.CoreV1().Services(ns).Delete(ctx, name, metav1.DeleteOptions{})
}
@@ -32,17 +32,18 @@ func (deploy *NewDeploy) EnvEventHandlers() k8sCache.ResourceEventHandlerFuncs {
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
newEnv := newObj.(*fv1.Environment)
oldEnv := oldObj.(*fv1.Environment)
ctx := context.Background()
// Currently only an image update in environment calls for function's deployment recreation. In future there might be more attributes which would want to do it
if oldEnv.Spec.Runtime.Image != newEnv.Spec.Runtime.Image {
deploy.logger.Debug("Updating all function of the environment that changed, old env:", zap.Any("environment", oldEnv))
funcs := deploy.getEnvFunctions(&newEnv.ObjectMeta)
funcs := deploy.getEnvFunctions(ctx, &newEnv.ObjectMeta)
for _, f := range funcs {
function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(context.TODO(), f.ObjectMeta.Name, metav1.GetOptions{})
function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(ctx, f.ObjectMeta.Name, metav1.GetOptions{})
if err != nil {
deploy.logger.Error("Error getting function", zap.Error(err), zap.Any("function", function))
continue
}
err = deploy.updateFuncDeployment(function, newEnv)
err = deploy.updateFuncDeployment(ctx, function, newEnv)
if err != nil {
deploy.logger.Error("Error updating function", zap.Error(err), zap.Any("function", function))
continue
@@ -16,6 +16,8 @@ limitations under the License.
package newdeploy
import (
"context"
"go.uber.org/zap"
k8sCache "k8s.io/client-go/tools/cache"
@@ -29,9 +31,10 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu
// and worker pattern to process items instead of moving process to another goroutine.
// example: https://github.com/kubernetes/kubernetes/blob/master/pkg/controller/job/job_controller.go
go func() {
ctx := context.Background()
fn := obj.(*fv1.Function)
deploy.logger.Debug("create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec))
_, err := deploy.createFunction(fn)
_, err := deploy.createFunction(ctx, fn)
if err != nil {
deploy.logger.Error("error eager creating function",
zap.Error(err),
@@ -43,7 +46,8 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu
DeleteFunc: func(obj interface{}) {
fn := obj.(*fv1.Function)
go func() {
err := deploy.deleteFunction(fn)
ctx := context.Background()
err := deploy.deleteFunction(ctx, fn)
if err != nil {
deploy.logger.Error("error deleting function",
zap.Error(err),
@@ -55,7 +59,8 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu
oldFn := oldObj.(*fv1.Function)
newFn := newObj.(*fv1.Function)
go func() {
err := deploy.updateFunction(oldFn, newFn)
ctx := context.Background()
err := deploy.updateFunction(ctx, oldFn, newFn)
if err != nil {
deploy.logger.Error("error updating function",
zap.Error(err),
@@ -45,7 +45,7 @@ const (
DeploymentVersion = "apps/v1"
)
func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Environment,
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) {
specializationTimeout := fn.Spec.InvokeStrategy.ExecutionStrategy.SpecializationTimeout
@@ -58,12 +58,12 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
minScale = 1
}
deployment, err := deploy.getDeploymentSpec(fn, env, &minScale, deployName, deployNamespace, deployLabels, deployAnnotations)
deployment, err := deploy.getDeploymentSpec(ctx, fn, env, &minScale, deployName, deployNamespace, deployLabels, deployAnnotations)
if err != nil {
return nil, err
}
existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(context.TODO(), deployName, metav1.GetOptions{})
existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(ctx, deployName, metav1.GetOptions{})
if err == nil {
// Try to adopt orphan deployment created by the old executor.
if existingDepl.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID {
@@ -75,7 +75,7 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
// Update with the latest deployment spec. Kubernetes will trigger
// rolling update if spec is different from the one in the cluster.
existingDepl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Update(context.TODO(), existingDepl, metav1.UpdateOptions{})
existingDepl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Update(ctx, existingDepl, metav1.UpdateOptions{})
if err != nil {
deploy.logger.Warn("error adopting deploy", zap.Error(err),
zap.String("deploy", deployName), zap.String("ns", deployNamespace))
@@ -86,27 +86,27 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
}
if *existingDepl.Spec.Replicas < minScale {
err = deploy.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale)
err = deploy.scaleDeployment(ctx, existingDepl.Namespace, existingDepl.Name, minScale)
if err != nil {
deploy.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.ObjectMeta.Name))
return nil, err
}
}
if existingDepl.Status.AvailableReplicas < minScale {
existingDepl, err = deploy.waitForDeploy(existingDepl, minScale, specializationTimeout)
existingDepl, err = deploy.waitForDeploy(ctx, existingDepl, minScale, specializationTimeout)
}
return existingDepl, err
} else if k8s_err.IsNotFound(err) {
err := deploy.setupRBACObjs(deployNamespace, fn)
err := deploy.setupRBACObjs(ctx, deployNamespace, fn)
if err != nil {
return nil, err
}
depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(context.TODO(), deployment, metav1.CreateOptions{})
depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(ctx, deployment, metav1.CreateOptions{})
if err != nil {
if k8s_err.IsAlreadyExists(err) {
depl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(context.TODO(), deployName, metav1.GetOptions{})
depl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(ctx, deployName, metav1.GetOptions{})
}
if err != nil {
deploy.logger.Error("error while creating function deployment",
@@ -118,14 +118,14 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
}
}
if minScale > 0 {
depl, err = deploy.waitForDeploy(depl, minScale, specializationTimeout)
depl, err = deploy.waitForDeploy(ctx, depl, minScale, specializationTimeout)
}
return depl, err
}
return nil, err
}
func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) error {
func (deploy *NewDeploy) setupRBACObjs(ctx context.Context, deployNamespace string, fn *fv1.Function) error {
// create fetcher SA in this ns, if not already created
err := deploy.fetcherConfig.SetupServiceAccount(deploy.kubernetesClient, deployNamespace, fn.ObjectMeta)
if err != nil {
@@ -139,7 +139,7 @@ func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function)
}
// create a cluster role binding for the fetcher SA, if not already created, granting access to do a get on packages in any ns
err = utils.SetupRoleBinding(deploy.logger, deploy.kubernetesClient, fv1.PackageGetterRB, fn.Spec.Package.PackageRef.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace)
err = utils.SetupRoleBinding(ctx, deploy.logger, deploy.kubernetesClient, fv1.PackageGetterRB, fn.Spec.Package.PackageRef.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace)
if err != nil {
deploy.logger.Error("error creating role binding for function",
zap.Error(err),
@@ -150,7 +150,7 @@ func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function)
}
// create rolebinding in function namespace for fetcherSA.envNamespace to be able to get secrets and configmaps
err = utils.SetupRoleBinding(deploy.logger, deploy.kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace)
err = utils.SetupRoleBinding(ctx, deploy.logger, deploy.kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace)
if err != nil {
deploy.logger.Error("error creating role binding for function",
zap.Error(err),
@@ -166,21 +166,21 @@ func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function)
return nil
}
func (deploy *NewDeploy) updateDeployment(deployment *appsv1.Deployment, ns string) error {
_, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Update(context.TODO(), deployment, metav1.UpdateOptions{})
func (deploy *NewDeploy) updateDeployment(ctx context.Context, deployment *appsv1.Deployment, ns string) error {
_, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Update(ctx, deployment, metav1.UpdateOptions{})
return err
}
func (deploy *NewDeploy) deleteDeployment(ns string, name string) error {
func (deploy *NewDeploy) deleteDeployment(ctx context.Context, ns string, name string) error {
// DeletePropagationBackground deletes the object immediately and dependent are deleted later
// DeletePropagationForeground not advisable; it marks for deleteion and API can still serve those objects
deletePropagation := metav1.DeletePropagationBackground
return deploy.kubernetesClient.AppsV1().Deployments(ns).Delete(context.TODO(), name, metav1.DeleteOptions{
return deploy.kubernetesClient.AppsV1().Deployments(ns).Delete(ctx, name, metav1.DeleteOptions{
PropagationPolicy: &deletePropagation,
})
}
func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environment, targetReplicas *int32,
func (deploy *NewDeploy) getDeploymentSpec(ctx context.Context, fn *fv1.Function, env *fv1.Environment, targetReplicas *int32,
deployName string, deployNamespace string, deployLabels map[string]string, deployAnnotations map[string]string) (*appsv1.Deployment, error) {
replicas := int32(fn.Spec.InvokeStrategy.ExecutionStrategy.MinScale)
@@ -233,7 +233,7 @@ func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environmen
// rollback, set RevisionHistoryLimit to 0 to disable this feature.
revisionHistoryLimit := int32(0)
rvCount, err := referencedResourcesRVSum(deploy.kubernetesClient, fn.ObjectMeta.Namespace, fn.Spec.Secrets, fn.Spec.ConfigMaps)
rvCount, err := referencedResourcesRVSum(ctx, deploy.kubernetesClient, fn.ObjectMeta.Namespace, fn.Spec.Secrets, fn.Spec.ConfigMaps)
if err != nil {
return nil, err
}
@@ -366,7 +366,7 @@ func (deploy *NewDeploy) getResources(env *fv1.Environment, fn *fv1.Function) ap
return resources
}
func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy,
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 {
@@ -401,14 +401,14 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut
},
}
existingHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(context.TODO(), hpaName, metav1.GetOptions{})
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(context.TODO(), existingHpa, metav1.UpdateOptions{})
existingHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(ctx, existingHpa, metav1.UpdateOptions{})
if err != nil {
deploy.logger.Warn("error adopting HPA", zap.Error(err),
zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace))
@@ -417,10 +417,10 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut
}
return existingHpa, err
} else if k8s_err.IsNotFound(err) {
cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(context.TODO(), hpa, metav1.CreateOptions{})
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(context.TODO(), hpaName, metav1.GetOptions{})
cHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(ctx, hpaName, metav1.GetOptions{})
}
if err != nil {
return nil, err
@@ -431,20 +431,20 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut
return nil, err
}
func (deploy *NewDeploy) getHpa(ns, name string) (*asv1.HorizontalPodAutoscaler, error) {
return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(context.TODO(), name, metav1.GetOptions{})
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(hpa *asv1.HorizontalPodAutoscaler) error {
_, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(context.TODO(), hpa, metav1.UpdateOptions{})
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(ns string, name string) error {
return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(context.TODO(), name, metav1.DeleteOptions{})
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(deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) {
func (deploy *NewDeploy) createOrGetSvc(ctx context.Context, deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) {
service := &apiv1.Service{
ObjectMeta: metav1.ObjectMeta{
Name: svcName,
@@ -465,7 +465,7 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn
},
}
existingSvc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{})
existingSvc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{})
if err == nil {
// to adopt orphan service
if existingSvc.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID {
@@ -474,7 +474,7 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn
existingSvc.Spec.Ports = service.Spec.Ports
existingSvc.Spec.Selector = service.Spec.Selector
existingSvc.Spec.Type = service.Spec.Type
existingSvc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Update(context.TODO(), existingSvc, metav1.UpdateOptions{})
existingSvc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Update(ctx, existingSvc, metav1.UpdateOptions{})
if err != nil {
deploy.logger.Warn("error adopting service", zap.Error(err),
zap.String("service", svcName), zap.String("ns", svcNamespace))
@@ -483,10 +483,10 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn
}
return existingSvc, err
} else if k8s_err.IsNotFound(err) {
svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(context.TODO(), service, metav1.CreateOptions{})
svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(ctx, service, metav1.CreateOptions{})
if err != nil {
if k8s_err.IsAlreadyExists(err) {
svc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{})
svc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{})
}
if err != nil {
return nil, err
@@ -497,11 +497,11 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn
return nil, err
}
func (deploy *NewDeploy) deleteSvc(ns string, name string) error {
return deploy.kubernetesClient.CoreV1().Services(ns).Delete(context.TODO(), name, metav1.DeleteOptions{})
func (deploy *NewDeploy) deleteSvc(ctx context.Context, ns string, name string) error {
return deploy.kubernetesClient.CoreV1().Services(ns).Delete(ctx, name, metav1.DeleteOptions{})
}
func (deploy *NewDeploy) waitForDeploy(depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) {
func (deploy *NewDeploy) waitForDeploy(ctx context.Context, depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) {
oldStatus := depl.Status
// if no specializationTimeout is set, use default value
@@ -510,7 +510,7 @@ func (deploy *NewDeploy) waitForDeploy(depl *appsv1.Deployment, replicas int32,
}
for i := 0; i < specializationTimeout; i++ {
latestDepl, err = deploy.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(context.TODO(), depl.Name, metav1.GetOptions{})
latestDepl, err = deploy.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(ctx, depl.Name, metav1.GetOptions{})
if err != nil {
return nil, err
}
@@ -533,10 +533,10 @@ func (deploy *NewDeploy) waitForDeploy(depl *appsv1.Deployment, replicas int32,
}
// cleanupNewdeploy cleans all kubernetes objects related to function
func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error {
func (deploy *NewDeploy) cleanupNewdeploy(ctx context.Context, ns string, name string) error {
result := &multierror.Error{}
err := deploy.deleteSvc(ns, name)
err := deploy.deleteSvc(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
deploy.logger.Error("error deleting service for newdeploy function",
zap.Error(err),
@@ -545,7 +545,7 @@ func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error {
result = multierror.Append(result, err)
}
err = deploy.deleteHpa(ns, name)
err = deploy.deleteHpa(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
deploy.logger.Error("error deleting HPA for newdeploy function",
zap.Error(err),
@@ -554,7 +554,7 @@ func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error {
result = multierror.Append(result, err)
}
err = deploy.deleteDeployment(ns, name)
err = deploy.deleteDeployment(ctx, ns, name)
if err != nil && !k8s_err.IsNotFound(err) {
deploy.logger.Error("error deleting deployment for newdeploy function",
zap.Error(err),
@@ -574,11 +574,11 @@ func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error {
// identical way to get a value that can reflect resources changed without affecting by the time.
// To achieve this goal, the sum of the resource version of all referenced resources is a good fit for our
// scenario since the sum of the resource version is always the same as long as no resources changed.
func referencedResourcesRVSum(client *kubernetes.Clientset, namespace string, secrets []fv1.SecretReference, cfgmaps []fv1.ConfigMapReference) (int, error) {
func referencedResourcesRVSum(ctx context.Context, client *kubernetes.Clientset, namespace string, secrets []fv1.SecretReference, cfgmaps []fv1.ConfigMapReference) (int, error) {
rvCount := 0
if len(secrets) > 0 {
list, err := client.CoreV1().Secrets(namespace).List(context.TODO(), metav1.ListOptions{})
list, err := client.CoreV1().Secrets(namespace).List(ctx, metav1.ListOptions{})
if err != nil {
return 0, err
}
@@ -598,7 +598,7 @@ func referencedResourcesRVSum(client *kubernetes.Clientset, namespace string, se
}
if len(cfgmaps) > 0 {
list, err := client.CoreV1().ConfigMaps(namespace).List(context.TODO(), metav1.ListOptions{})
list, err := client.CoreV1().ConfigMaps(namespace).List(ctx, metav1.ListOptions{})
if err != nil {
return 0, err
}
@@ -144,7 +144,7 @@ func (deploy *NewDeploy) GetFuncSvc(ctx context.Context, fn *fv1.Function) (*fsc
// TODO: client-go doesn't support to pass in context.
// Once it supports context, we should change the signature of method.
// https://github.com/kubernetes/kubernetes/issues/46503
return deploy.createFunction(fn)
return deploy.createFunction(ctx, fn)
}
// GetFuncSvcFromCache returns a function service from cache; error otherwise.
@@ -194,7 +194,7 @@ func (deploy *NewDeploy) getServiceInfo(ctx context.Context, obj apiv1.ObjectRef
return service, nil
}
func (deploy *NewDeploy) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Deployment, error) {
func (deploy *NewDeploy) getDeploymentInfo(ctx context.Context, obj apiv1.ObjectReference) (*appsv1.Deployment, error) {
item, exists, err := utils.GetCachedItem(obj, deploy.deploymentInformer)
if err != nil || !exists {
@@ -203,7 +203,7 @@ func (deploy *NewDeploy) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.D
zap.Bool("exists", exists),
zap.Error(err),
)
deployment, err := deploy.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{})
deployment, err := deploy.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
return deployment, err
}
@@ -234,7 +234,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo
}
} else if strings.ToLower(obj.Kind) == "deployment" {
currentDeploy, err := deploy.getDeploymentInfo(obj)
currentDeploy, err := deploy.getDeploymentInfo(ctx, obj)
if err != nil {
if !k8sErrs.IsNotFound(err) {
deploy.logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err))
@@ -273,7 +273,7 @@ func (deploy *NewDeploy) RefreshFuncPods(ctx context.Context, logger *zap.Logger
// Ideally there should be only one deployment but for now we rely on label/selector to ensure that condition
for _, deployment := range dep.Items {
rvCount, err := referencedResourcesRVSum(deploy.kubernetesClient, deployment.Namespace, f.Spec.Secrets, f.Spec.ConfigMaps)
rvCount, err := referencedResourcesRVSum(ctx, deploy.kubernetesClient, deployment.Namespace, f.Spec.Secrets, f.Spec.ConfigMaps)
if err != nil {
return err
}
@@ -308,7 +308,7 @@ func (deploy *NewDeploy) AdoptExistingResources(ctx context.Context) {
go func() {
defer wg.Done()
_, err = deploy.fnCreate(fn)
_, err = deploy.fnCreate(ctx, fn)
if err != nil {
deploy.logger.Warn("failed to adopt resources for function", zap.Error(err))
return
@@ -330,17 +330,17 @@ func (deploy *NewDeploy) CleanupOldExecutorObjects(ctx context.Context) {
LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy)}).AsSelector().String(),
}
err := reaper.CleanupHpa(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
err := reaper.CleanupHpa(ctx, deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupDeployments(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
err = reaper.CleanupDeployments(ctx, deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupServices(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
err = reaper.CleanupServices(ctx, deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
@@ -351,8 +351,8 @@ func (deploy *NewDeploy) CleanupOldExecutorObjects(ctx context.Context) {
}
}
func (deploy *NewDeploy) getEnvFunctions(m *metav1.ObjectMeta) []fv1.Function {
funcList, err := deploy.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
func (deploy *NewDeploy) getEnvFunctions(ctx context.Context, m *metav1.ObjectMeta) []fv1.Function {
funcList, err := deploy.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil {
deploy.logger.Error("Error getting functions for env", zap.Error(err), zap.Any("environment", m))
}
@@ -365,14 +365,14 @@ func (deploy *NewDeploy) getEnvFunctions(m *metav1.ObjectMeta) []fv1.Function {
return relatedFunctions
}
func (deploy *NewDeploy) createFunction(fn *fv1.Function) (*fscache.FuncSvc, error) {
func (deploy *NewDeploy) createFunction(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypeNewdeploy {
return nil, nil
}
fsvcObj, err := deploy.throttler.RunOnce(string(fn.ObjectMeta.UID), func(ableToCreate bool) (interface{}, error) {
if ableToCreate {
return deploy.fnCreate(fn)
return deploy.fnCreate(ctx, fn)
}
return deploy.fsCache.GetByFunctionUID(fn.ObjectMeta.UID)
})
@@ -394,20 +394,21 @@ func (deploy *NewDeploy) createFunction(fn *fv1.Function) (*fscache.FuncSvc, err
return fsvc, err
}
func (deploy *NewDeploy) deleteFunction(fn *fv1.Function) error {
func (deploy *NewDeploy) deleteFunction(ctx context.Context, fn *fv1.Function) error {
if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypeNewdeploy {
return nil
}
err := deploy.fnDelete(fn)
err := deploy.fnDelete(ctx, fn)
if err != nil {
err = errors.Wrapf(err, "error deleting kubernetes objects of function %v", fn.ObjectMeta)
}
return err
}
func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) {
func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
cleanupFunc := func(ns string, name string) {
err := deploy.cleanupNewdeploy(ns, name)
ctx := context.Background()
err := deploy.cleanupNewdeploy(ctx, ns, name)
if err != nil {
deploy.logger.Error("received error while cleaning function resources",
zap.String("namespace", ns), zap.String("name", name))
@@ -415,7 +416,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) {
}
env, err := deploy.fissionClient.CoreV1().
Environments(fn.Spec.Environment.Namespace).
Get(context.TODO(), fn.Spec.Environment.Name, metav1.GetOptions{})
Get(ctx, fn.Spec.Environment.Name, metav1.GetOptions{})
if err != nil {
return nil, err
}
@@ -436,7 +437,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) {
// Since newdeploy waits for pods of deployment to be ready,
// change the order of kubeObject creation (create service first,
// then deployment) to take advantage of waiting time.
svc, err := deploy.createOrGetSvc(deployLabels, deployAnnotations, objName, ns)
svc, err := deploy.createOrGetSvc(ctx, deployLabels, deployAnnotations, objName, ns)
if err != nil {
deploy.logger.Error("error creating service", zap.Error(err), zap.String("service", objName))
go cleanupFunc(ns, objName)
@@ -444,14 +445,14 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) {
}
svcAddress := fmt.Sprintf("%v.%v", svc.Name, svc.Namespace)
depl, err := deploy.createOrGetDeployment(fn, env, objName, deployLabels, deployAnnotations, ns)
depl, err := deploy.createOrGetDeployment(ctx, fn, env, objName, deployLabels, deployAnnotations, ns)
if err != nil {
deploy.logger.Error("error creating deployment", zap.Error(err), zap.String("deployment", objName))
go cleanupFunc(ns, objName)
return nil, errors.Wrapf(err, "error creating deployment %v", objName)
}
hpa, err := deploy.createOrGetHpa(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
hpa, err := deploy.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)
@@ -506,7 +507,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) {
return fsvc, nil
}
func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function) error {
func (deploy *NewDeploy) updateFunction(ctx context.Context, oldFn *fv1.Function, newFn *fv1.Function) error {
if oldFn.ObjectMeta.ResourceVersion == newFn.ObjectMeta.ResourceVersion {
return nil
@@ -524,7 +525,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function
deploy.logger.Info("function does not use new deployment executor anymore, deleting resources",
zap.Any("function", newFn))
// IMP - pass the oldFn, as the new/modified function is not in cache
return deploy.deleteFunction(oldFn)
return deploy.deleteFunction(ctx, oldFn)
}
// Executor type changed to New Deployment from something else
@@ -533,7 +534,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function
deploy.logger.Info("function type changed to new deployment, creating resources",
zap.Any("old_function", oldFn.ObjectMeta),
zap.Any("new_function", newFn.ObjectMeta))
_, err := deploy.createFunction(newFn)
_, err := deploy.createFunction(ctx, newFn)
if err != nil {
deploy.updateStatus(oldFn, err, "error changing the function's type to newdeploy")
}
@@ -557,7 +558,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function
return err
}
hpa, err := deploy.getHpa(ns, fsvc.Name)
hpa, err := deploy.getHpa(ctx, ns, fsvc.Name)
if err != nil {
deploy.updateStatus(oldFn, err, "error getting HPA while updating function")
return err
@@ -583,7 +584,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function
}
if hpaChanged {
err := deploy.updateHpa(hpa)
err := deploy.updateHpa(ctx, hpa)
if err != nil {
deploy.updateStatus(oldFn, err, "error updating HPA while updating function")
return err
@@ -621,18 +622,18 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function
if deployChanged {
env, err := deploy.fissionClient.CoreV1().Environments(newFn.Spec.Environment.Namespace).
Get(context.TODO(), newFn.Spec.Environment.Name, metav1.GetOptions{})
Get(ctx, newFn.Spec.Environment.Name, metav1.GetOptions{})
if err != nil {
deploy.updateStatus(oldFn, err, "failed to get environment while updating function")
return err
}
return deploy.updateFuncDeployment(newFn, env)
return deploy.updateFuncDeployment(ctx, newFn, env)
}
return nil
}
func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environment) error {
func (deploy *NewDeploy) updateFuncDeployment(ctx context.Context, fn *fv1.Function, env *fv1.Environment) error {
fsvc, err := deploy.fsCache.GetByFunctionUID(fn.ObjectMeta.UID)
if err != nil {
err = errors.Wrapf(err, "error updating function due to unable to find function service cache: %v", fn)
@@ -651,7 +652,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ
ns = fn.ObjectMeta.Namespace
}
existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Get(context.TODO(), fnObjName, metav1.GetOptions{})
existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Get(ctx, fnObjName, metav1.GetOptions{})
if err != nil {
return err
}
@@ -659,7 +660,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ
// the resource version inside function packageRef is changed,
// so the content of fetchRequest in deployment cmd is different.
// Therefore, the deployment update will trigger a rolling update.
newDeployment, err := deploy.getDeploymentSpec(fn, env,
newDeployment, err := deploy.getDeploymentSpec(ctx, fn, env,
existingDepl.Spec.Replicas, // use current replicas instead of minscale in the ExecutionStrategy.
fnObjName, ns, deployLabels, deploy.getDeployAnnotations(fn.ObjectMeta, env.ObjectMeta))
if err != nil {
@@ -667,7 +668,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ
return err
}
err = deploy.updateDeployment(newDeployment, ns)
err = deploy.updateDeployment(ctx, newDeployment, ns)
if err != nil {
deploy.updateStatus(fn, err, "failed to update deployment while updating function")
return err
@@ -676,7 +677,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ
return nil
}
func (deploy *NewDeploy) fnDelete(fn *fv1.Function) error {
func (deploy *NewDeploy) fnDelete(ctx context.Context, fn *fv1.Function) error {
multierr := &multierror.Error{}
// GetByFunction uses resource version as part of cache key, however,
@@ -705,7 +706,7 @@ func (deploy *NewDeploy) fnDelete(fn *fv1.Function) error {
ns = fn.ObjectMeta.Namespace
}
err = deploy.cleanupNewdeploy(ns, objName)
err = deploy.cleanupNewdeploy(ctx, ns, objName)
multierr = multierror.Append(multierr, err)
return multierr.ErrorOrNil()
@@ -767,12 +768,12 @@ func (deploy *NewDeploy) updateStatus(fn *fv1.Function, err error, message strin
// idleObjectReaper reaps objects after certain idle time
func (deploy *NewDeploy) idleObjectReaper() {
ctx := context.Background()
pollSleep := 5 * time.Second
for {
time.Sleep(pollSleep)
envs, err := deploy.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
envs, err := deploy.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil {
deploy.logger.Fatal("failed to get environment list", zap.Error(err))
}
@@ -803,7 +804,7 @@ func (deploy *NewDeploy) idleObjectReaper() {
zap.String("function", fsvc.Name))
}
fn, err := deploy.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(context.TODO(), fsvc.Function.Name, metav1.GetOptions{})
fn, err := deploy.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(ctx, fsvc.Function.Name, metav1.GetOptions{})
if err != nil {
// Newdeploy manager handles the function delete event and clean cache/kubeobjs itself,
// so we ignore the not found error for functions with newdeploy executor type here.
@@ -834,7 +835,7 @@ func (deploy *NewDeploy) idleObjectReaper() {
}
currentDeploy, err := deploy.kubernetesClient.AppsV1().
Deployments(deployObj.Namespace).Get(context.TODO(), deployObj.Name, metav1.GetOptions{})
Deployments(deployObj.Namespace).Get(ctx, deployObj.Name, metav1.GetOptions{})
if err != nil {
deploy.logger.Error("error getting function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
return
@@ -847,7 +848,7 @@ func (deploy *NewDeploy) idleObjectReaper() {
return
}
err = deploy.scaleDeployment(deployObj.Namespace, deployObj.Name, minScale)
err = deploy.scaleDeployment(ctx, deployObj.Namespace, deployObj.Name, minScale)
if err != nil {
deploy.logger.Error("error scaling down function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
}
@@ -867,12 +868,12 @@ func getDeploymentObj(kubeobjs []apiv1.ObjectReference) *apiv1.ObjectReference {
return nil
}
func (deploy *NewDeploy) scaleDeployment(deplNS string, deplName string, replicas int32) error {
func (deploy *NewDeploy) scaleDeployment(ctx context.Context, deplNS string, deplName string, replicas int32) error {
deploy.logger.Info("scaling deployment",
zap.String("deployment", deplName),
zap.String("namespace", deplNS),
zap.Int32("replicas", replicas))
_, err := deploy.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(context.TODO(), deplName, &autoscalingv1.Scale{
_, err := deploy.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(ctx, deplName, &autoscalingv1.Scale{
ObjectMeta: metav1.ObjectMeta{
Name: deplName,
Namespace: deplNS,
@@ -44,6 +44,7 @@ func getIstioServiceLabels(fnName string) map[string]string {
func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ctx := context.Background()
fn := obj.(*fv1.Function)
// Since istio only allows accessing pod through k8s service,
@@ -69,7 +70,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie
// setup rolebinding is tried, if it fails, we don't return. we just log an error and move on, because :
// 1. not all functions have secrets and/or configmaps, so things will work without this rolebinding in that case.
// 2. on the contrary, when the route is tried, the env fetcher logs will show a 403 forbidden message and same will be relayed to executor.
err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
if err != nil {
logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB))
} else {
@@ -120,7 +121,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie
}
// create function istio service if it does not exist
_, err = kubernetesClient.CoreV1().Services(envNs).Create(context.TODO(), &svc, metav1.CreateOptions{})
_, err = kubernetesClient.CoreV1().Services(envNs).Create(ctx, &svc, metav1.CreateOptions{})
if err != nil && !kerrors.IsAlreadyExists(err) {
logger.Error("error creating istio service for function",
zap.Error(err),
@@ -132,6 +133,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie
},
DeleteFunc: func(obj interface{}) {
ctx := context.Background()
fn := obj.(*fv1.Function)
fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType
@@ -147,7 +149,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie
if istioEnabled {
svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace)
// delete function istio service
err := kubernetesClient.CoreV1().Services(envNs).Delete(context.TODO(), svcName, metav1.DeleteOptions{})
err := kubernetesClient.CoreV1().Services(envNs).Delete(ctx, svcName, metav1.DeleteOptions{})
if err != nil && !kerrors.IsNotFound(err) {
logger.Error("error deleting istio service for function",
zap.Error(err),
@@ -181,8 +183,8 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie
if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = newFunc.Spec.Environment.Namespace
}
err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.SecretConfigMapGetterRB,
ctx := context.Background()
err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.SecretConfigMapGetterRB,
newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole,
fv1.FissionFetcherSA, envNs)
+6 -7
View File
@@ -143,7 +143,6 @@ func MakeGenericPool(
if err != nil {
return nil, err
}
gpLogger.Info("deployment created", zap.Any("environment", env.ObjectMeta))
go gp.startReadyPodController()
go gp.updateCPUUtilizationSvc()
@@ -223,7 +222,7 @@ func (gp *GenericPool) updateCPUUtilizationSvc() {
// choosePod picks a ready pod from the pool and relabels it, waiting if necessary.
// returns the key and pod API object.
func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Pod, error) {
func (gp *GenericPool) choosePod(ctx context.Context, newLabels map[string]string) (string, *apiv1.Pod, error) {
startTime := time.Now()
expoDelay := 100 * time.Millisecond
for {
@@ -277,7 +276,7 @@ func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Po
patch := fmt.Sprintf(`{"metadata":{"annotations":%v, "labels":%v}}`, string(annotationPatch), string(labelPatch))
gp.logger.Info("relabel pod", zap.String("pod", patch))
newPod, err := gp.kubernetesClient.CoreV1().Pods(chosenPod.Namespace).Patch(context.TODO(), chosenPod.Name, k8sTypes.StrategicMergePatchType, []byte(patch), metav1.PatchOptions{})
newPod, err := gp.kubernetesClient.CoreV1().Pods(chosenPod.Namespace).Patch(ctx, chosenPod.Name, k8sTypes.StrategicMergePatchType, []byte(patch), metav1.PatchOptions{})
if err != nil {
gp.logger.Error("failed to relabel pod", zap.Error(err), zap.String("pod", chosenPod.Name), zap.Duration("delay", expoDelay))
gp.readyPodQueue.Done(key)
@@ -397,7 +396,7 @@ func (gp *GenericPool) specializePod(ctx context.Context, pod *apiv1.Pod, fn *fv
return nil
}
func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1.Service, error) {
func (gp *GenericPool) createSvc(ctx context.Context, name string, labels map[string]string) (*apiv1.Service, error) {
service := apiv1.Service{
ObjectMeta: metav1.ObjectMeta{
Name: name,
@@ -415,7 +414,7 @@ func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1.
Selector: labels,
},
}
svc, err := gp.kubernetesClient.CoreV1().Services(gp.namespace).Create(context.TODO(), &service, metav1.CreateOptions{})
svc, err := gp.kubernetesClient.CoreV1().Services(gp.namespace).Create(ctx, &service, metav1.CreateOptions{})
return svc, err
}
@@ -466,7 +465,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
}
}
key, pod, err := gp.choosePod(funcLabels)
key, pod, err := gp.choosePod(ctx, funcLabels)
if err != nil {
return nil, err
}
@@ -485,7 +484,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
svcName = fmt.Sprintf("%s-%v", svcName, fn.ObjectMeta.UID)
}
svc, err := gp.createSvc(svcName, funcLabels)
svc, err := gp.createSvc(ctx, svcName, funcLabels)
if err != nil {
gp.scheduleDeletePod(pod.ObjectMeta.Name)
return nil, err
@@ -199,5 +199,7 @@ func (gp *GenericPool) createPoolDeployment(ctx context.Context, env *fv1.Enviro
}
gp.deployment = depl
gp.logger.Info("deployment created", zap.String("deployment", depl.Name), zap.String("ns", depl.Namespace), zap.Any("environment", env.ObjectMeta))
return nil
}
+14 -13
View File
@@ -163,7 +163,7 @@ func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType
func (gpm *GenericPoolManager) GetFuncSvc(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
// from Func -> get Env
gpm.logger.Debug("getting environment for function", zap.String("function", fn.ObjectMeta.Name))
env, err := gpm.getFunctionEnv(fn)
env, err := gpm.getFunctionEnv(ctx, fn)
if err != nil {
return nil, err
}
@@ -207,7 +207,7 @@ func (gpm *GenericPoolManager) TapService(ctx context.Context, svcHost string) e
return nil
}
func (gpm *GenericPoolManager) getPodInfo(obj apiv1.ObjectReference) (*apiv1.Pod, error) {
func (gpm *GenericPoolManager) getPodInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Pod, error) {
store := gpm.podInformer.GetStore()
item, exists, err := store.Get(obj)
@@ -217,7 +217,7 @@ func (gpm *GenericPoolManager) getPodInfo(obj apiv1.ObjectReference) (*apiv1.Pod
if err != nil || !exists {
gpm.logger.Debug("Falling back to getting pod info from k8s API -- this may cause performance issues for your function.")
pod, err := gpm.kubernetesClient.CoreV1().Pods(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{})
pod, err := gpm.kubernetesClient.CoreV1().Pods(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
return pod, err
}
@@ -230,7 +230,7 @@ func (gpm *GenericPoolManager) getPodInfo(obj apiv1.ObjectReference) (*apiv1.Pod
func (gpm *GenericPoolManager) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool {
for _, obj := range fsvc.KubernetesObjects {
if strings.ToLower(obj.Kind) == "pod" {
pod, err := gpm.getPodInfo(obj)
pod, err := gpm.getPodInfo(ctx, obj)
if err == nil && utils.IsReadyPod(pod) {
// Normally, the address format is http://[pod-ip]:[port], however, if the
// Istio is enabled the address format changes to http://[svc-name]:[port].
@@ -437,12 +437,12 @@ func (gpm *GenericPoolManager) CleanupOldExecutorObjects(ctx context.Context) {
LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr)}).AsSelector().String(),
}
err := reaper.CleanupDeployments(gpm.logger, gpm.kubernetesClient, gpm.instanceID, listOpts)
err := reaper.CleanupDeployments(ctx, gpm.logger, gpm.kubernetesClient, gpm.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupPods(gpm.logger, gpm.kubernetesClient, gpm.instanceID, listOpts)
err = reaper.CleanupPods(ctx, gpm.logger, gpm.kubernetesClient, gpm.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
@@ -521,7 +521,7 @@ func (gpm *GenericPoolManager) cleanupPools(envs []fv1.Environment) {
}
}
func (gpm *GenericPoolManager) getFunctionEnv(fn *fv1.Function) (*fv1.Environment, error) {
func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Function) (*fv1.Environment, error) {
var env *fv1.Environment
// Cached ?
@@ -533,7 +533,7 @@ func (gpm *GenericPoolManager) getFunctionEnv(fn *fv1.Function) (*fv1.Environmen
}
// Get env from controller
env, err = gpm.fissionClient.CoreV1().Environments(fn.Spec.Environment.Namespace).Get(context.TODO(), fn.Spec.Environment.Name, metav1.GetOptions{})
env, err = gpm.fissionClient.CoreV1().Environments(fn.Spec.Environment.Namespace).Get(ctx, fn.Spec.Environment.Name, metav1.GetOptions{})
if err != nil {
return nil, err
}
@@ -600,13 +600,13 @@ func (gpm *GenericPoolManager) eagerPoolCreator() {
// idleObjectReaper reaps objects after certain idle time
func (gpm *GenericPoolManager) idleObjectReaper() {
ctx := context.Background()
pollSleep := 5 * time.Second
for {
time.Sleep(pollSleep)
envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil {
gpm.logger.Error("failed to get environment list", zap.Error(err))
continue
@@ -617,7 +617,7 @@ func (gpm *GenericPoolManager) idleObjectReaper() {
envList[env.ObjectMeta.UID] = struct{}{}
}
fns, err := gpm.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
fns, err := gpm.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil {
gpm.logger.Error("failed to get environment list", zap.Error(err))
continue
@@ -685,7 +685,7 @@ func (gpm *GenericPoolManager) idleObjectReaper() {
zap.String("executor", string(fsvc.Executor)),
zap.String("pod", fsvc.Name),
)
reaper.CleanupKubeObject(gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i])
reaper.CleanupKubeObject(ctx, gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i])
time.Sleep(50 * time.Millisecond)
gpm.fsCache.ReapTime(fsvc.Function.Name, fsvc.Address, time.Since(startTime).Seconds())
}
@@ -777,7 +777,8 @@ func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient *kubern
zap.String("executor", string(fsvc.Executor)),
zap.String("pod", fsvc.Name),
)
reaper.CleanupKubeObject(gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i])
ctx := context.Background()
reaper.CleanupKubeObject(ctx, gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i])
time.Sleep(50 * time.Millisecond)
}
}
@@ -17,6 +17,8 @@ limitations under the License.
package poolmgr
import (
"context"
"go.uber.org/zap"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
@@ -42,10 +44,10 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clien
if pkg.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = pkg.Spec.Environment.Namespace
}
ctx := context.Background()
// here, we return if we hit an error during rolebinding setup. this is because this rolebinding is mandatory for
// every function's package to be loaded into its env. without that, there's no point to move forward.
err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
if err != nil {
logger.Error("error creating rolebinding for package",
zap.Error(err),
@@ -79,7 +81,8 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clien
envNs = newPkg.Spec.Environment.Namespace
}
err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.PackageGetterRB,
ctx := context.Background()
err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.PackageGetterRB,
newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole,
fv1.FissionFetcherSA, envNs)
if err != nil {
+22 -21
View File
@@ -37,28 +37,28 @@ var (
)
// CleanupKubeObject deletes given kubernetes object
func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) {
func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) {
switch strings.ToLower(kubeobj.Kind) {
case "pod":
err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(context.TODO(), kubeobj.Name, meta_v1.DeleteOptions{})
err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.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(context.TODO(), kubeobj.Name, meta_v1.DeleteOptions{})
err := kubeClient.CoreV1().Services(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up service", zap.Error(err), zap.String("service", kubeobj.Name))
}
case "deployment":
err := kubeClient.AppsV1().Deployments(kubeobj.Namespace).Delete(context.TODO(), kubeobj.Name, delOpt)
err := kubeClient.AppsV1().Deployments(kubeobj.Namespace).Delete(ctx, kubeobj.Name, delOpt)
if err != nil {
logger.Error("error cleaning up deployment", zap.Error(err), zap.String("deployment", kubeobj.Name))
}
case "horizontalpodautoscaler":
err := kubeClient.AutoscalingV1().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(context.TODO(), kubeobj.Name, meta_v1.DeleteOptions{})
err := kubeClient.AutoscalingV1().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up horizontalpodautoscaler", zap.Error(err), zap.String("horizontalpodautoscaler", kubeobj.Name))
}
@@ -70,8 +70,8 @@ func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kub
}
// CleanupDeployments deletes deployment(s) for a given instanceID
func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(context.TODO(), listOps)
func CleanupDeployments(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -83,7 +83,7 @@ func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan
}
if ok && id != instanceID {
logger.Info("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name))
err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(context.TODO(), dep.ObjectMeta.Name, delOpt)
err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(ctx, dep.ObjectMeta.Name, delOpt)
if err != nil {
logger.Error("error cleaning up deployment",
zap.Error(err),
@@ -97,8 +97,8 @@ func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan
}
// CleanupPods deletes pod(s) for a given instanceID
func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(context.TODO(), listOps)
func CleanupPods(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -110,7 +110,7 @@ func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceID st
}
if ok && id != instanceID {
logger.Info("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name))
err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(context.TODO(), pod.ObjectMeta.Name, meta_v1.DeleteOptions{})
err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(ctx, pod.ObjectMeta.Name, meta_v1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up pod",
zap.Error(err),
@@ -124,8 +124,8 @@ func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceID st
}
// CleanupServices deletes service(s) for a given instanceID
func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(context.TODO(), listOps)
func CleanupServices(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -137,7 +137,7 @@ func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI
}
if ok && id != instanceID {
logger.Info("cleaning up service", zap.String("service", svc.ObjectMeta.Name))
err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(context.TODO(), svc.ObjectMeta.Name, meta_v1.DeleteOptions{})
err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(ctx, svc.ObjectMeta.Name, meta_v1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up service",
zap.Error(err),
@@ -151,8 +151,8 @@ func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI
}
// CleanupHpa deletes horizontal pod autoscaler(s) for a given instanceID
func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(context.TODO(), listOps)
func CleanupHpa(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
}
@@ -165,7 +165,7 @@ func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceID str
}
if ok && id != instanceID {
logger.Info("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name))
err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(context.TODO(), hpa.ObjectMeta.Name, meta_v1.DeleteOptions{})
err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(ctx, hpa.ObjectMeta.Name, meta_v1.DeleteOptions{})
if err != nil {
logger.Error("error cleaning up HPA",
zap.Error(err),
@@ -181,13 +181,14 @@ func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceID str
// CleanupRoleBindings periodically lists rolebindings across all namespaces and removes Service Accounts from them or
// deletes the rolebindings completely if there are no Service Accounts in a rolebinding object.
func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissionClient *crd.FissionClient, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) {
ctx := context.Background()
for {
// some sleep before the next reaper iteration
time.Sleep(cleanupRoleBindingInterval)
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(context.TODO(), meta_v1.ListOptions{})
rbList, err := client.RbacV1().RoleBindings(meta_v1.NamespaceAll).List(ctx, meta_v1.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 +209,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi
// 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(context.TODO(), meta_v1.ListOptions{})
funcList, err := fissionClient.CoreV1().Functions(roleBinding.Namespace).List(ctx, meta_v1.ListOptions{})
if err != nil {
logger.Error("error fetching function list in namespace", zap.Error(err), zap.String("namespace", roleBinding.Namespace))
continue
@@ -259,7 +260,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi
// 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(context.TODO(), meta_v1.ListOptions{})
envList, err := fissionClient.CoreV1().Environments(saNs).List(ctx, meta_v1.ListOptions{})
if err != nil {
logger.Error("error fetching environment list in service account namespace", zap.Error(err), zap.String("namespace", saNs))
continue
@@ -305,7 +306,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi
zap.String("role_binding_namespace", roleBinding.Namespace))
// call this once in the end for each role-binding
err = utils.RemoveSAFromRoleBindingWithRetries(logger, client, roleBinding.Name, roleBinding.Namespace, saToRemove)
err = utils.RemoveSAFromRoleBindingWithRetries(ctx, logger, client, roleBinding.Name, roleBinding.Namespace, saToRemove)
if err != nil {
// if there's an error, we just log it and proceed with the next role-binding, hoping that this role-binding
// will be processed in next iteration.
+5 -5
View File
@@ -57,7 +57,7 @@ func WaitTimeout(wg *sync.WaitGroup, timeout time.Duration) {
}
// ConvertConfigSecrets returns envFromSource which can be passed directly into the pod spec
func ConvertConfigSecrets(fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.EnvFromSource, error) {
func ConvertConfigSecrets(ctx context.Context, fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.EnvFromSource, error) {
cmList := fn.Spec.ConfigMaps
secList := fn.Spec.Secrets
@@ -65,9 +65,9 @@ func ConvertConfigSecrets(fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.E
secEnvSources := make([]*apiv1.SecretEnvSource, 0)
for _, cm := range cmList {
if cm.Namespace != fn.Namespace {
return nil, errors.New("Function should not reference config map of different namespace")
return nil, errors.New("function should not reference config map of different namespace")
}
_, err := kc.CoreV1().ConfigMaps(cm.Namespace).Get(context.TODO(), cm.Name, metav1.GetOptions{})
_, err := kc.CoreV1().ConfigMaps(cm.Namespace).Get(ctx, cm.Name, metav1.GetOptions{})
if err != nil {
return nil, err
}
@@ -81,9 +81,9 @@ func ConvertConfigSecrets(fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.E
for _, sec := range secList {
if sec.Namespace != fn.Namespace {
return nil, errors.New("Function should not reference secret of different namespace")
return nil, errors.New("function should not reference secret of different namespace")
}
_, err := kc.CoreV1().Secrets(sec.Namespace).Get(context.TODO(), sec.Name, metav1.GetOptions{})
_, err := kc.CoreV1().Secrets(sec.Namespace).Get(ctx, sec.Name, metav1.GetOptions{})
if err != nil {
return nil, err
}
+14 -14
View File
@@ -104,7 +104,7 @@ type PatchSpec struct {
}
// AddSaToRoleBindingWithRetries adds a service account to a rolebinding object. IT retries on already exists and conflict errors.
func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind string) (err error) {
func AddSaToRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind string) (err error) {
patch := PatchSpec{}
patch.Op = "add"
patch.Path = "/subjects/-"
@@ -121,7 +121,7 @@ func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Cli
}
for i := 0; i < maxRetries; i++ {
_, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Patch(context.TODO(), roleBinding, types.JSONPatchType, patchJson, metav1.PatchOptions{})
_, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Patch(ctx, roleBinding, types.JSONPatchType, patchJson, metav1.PatchOptions{})
if err == nil {
logger.Debug("patched rolebinding",
zap.String("role_binding", roleBinding),
@@ -137,7 +137,7 @@ func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Cli
// someone may have deleted the object between us checking if the object is present and deciding to patch
// so just create the object again
rbObj := makeRoleBindingObj(roleBinding, roleBindingNs, role, roleKind, sa, saNamespace)
_, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(context.TODO(), rbObj, metav1.CreateOptions{})
_, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(ctx, rbObj, metav1.CreateOptions{})
if err == nil {
logger.Debug("created rolebinding",
zap.String("role_binding", roleBinding),
@@ -176,10 +176,10 @@ func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Cli
// RemoveSAFromRoleBindingWithRetries removes an SA from the rolebinding passed as parameter. If this is the only SA in
// the rolebinding, then it deletes the rolebinding object.
func RemoveSAFromRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string, saToRemove map[string]bool) (err error) {
func RemoveSAFromRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string, saToRemove map[string]bool) (err error) {
for i := 0; i < maxRetries; i++ {
rbObj, err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Get(
context.TODO(),
ctx,
roleBinding, metav1.GetOptions{})
if err != nil {
// silently ignoring the error. there's no need for us to remove sa anymore.
@@ -206,13 +206,13 @@ func RemoveSAFromRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernete
})
}
if len(newSubjects) == 0 {
return DeleteRoleBinding(k8sClient, roleBinding, roleBindingNs)
return DeleteRoleBinding(ctx, k8sClient, roleBinding, roleBindingNs)
}
rbObj.Subjects = newSubjects
// cant use patch for deletes, the results become in-deterministic, so using update.
_, err = k8sClient.RbacV1().RoleBindings(rbObj.Namespace).Update(context.TODO(), rbObj, metav1.UpdateOptions{})
_, err = k8sClient.RbacV1().RoleBindings(rbObj.Namespace).Update(ctx, rbObj, metav1.UpdateOptions{})
switch {
case err == nil:
logger.Debug("removed service accounts from rolebinding",
@@ -236,10 +236,10 @@ func RemoveSAFromRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernete
// SetupRoleBinding adds a role to a service account if the rolebinding object is already present in the namespace.
// if not, it creates a rolebinding object granting the role to the SA in the namespace.
func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, role, roleKind, sa, saNamespace string) error {
func SetupRoleBinding(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, role, roleKind, sa, saNamespace string) error {
// get the role binding object
rbObj, err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Get(
context.TODO(),
ctx,
roleBinding, metav1.GetOptions{})
if err == nil {
@@ -249,7 +249,7 @@ func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleB
zap.String("service_account_namespace", saNamespace),
zap.String("role_binding", roleBinding),
zap.String("role_binding_namespace", roleBindingNs))
return AddSaToRoleBindingWithRetries(logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind)
return AddSaToRoleBindingWithRetries(ctx, logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind)
}
logger.Debug("service account already present in rolebinding so nothing to add",
zap.String("service_account_name", sa),
@@ -266,14 +266,14 @@ func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleB
zap.String("role_binding", roleBinding),
zap.String("role_binding_namespace", roleBindingNs))
rbObj = makeRoleBindingObj(roleBinding, roleBindingNs, role, roleKind, sa, saNamespace)
_, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(context.TODO(), rbObj, metav1.CreateOptions{})
_, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(ctx, rbObj, metav1.CreateOptions{})
if k8serrors.IsAlreadyExists(err) {
logger.Debug("rolebinding already exists in namespace - adding service account to rolebinding",
zap.String("service_account_name", sa),
zap.String("service_account_namespace", saNamespace),
zap.String("role_binding", roleBinding),
zap.String("role_binding_namespace", roleBindingNs))
err = AddSaToRoleBindingWithRetries(logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind)
err = AddSaToRoleBindingWithRetries(ctx, logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind)
}
}
@@ -282,10 +282,10 @@ func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleB
// DeleteRoleBinding deletes a rolebinding object. if k8s throws an error that the rolebinding is not there, it just
// returns silently.
func DeleteRoleBinding(k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string) error {
func DeleteRoleBinding(ctx context.Context, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string) error {
// if deleteRoleBinding is invoked by 2 fission services at the same time for the same rolebinding,
// the first call will succeed while the 2nd will fail with isNotFound. but we dont want to error out then.
err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Delete(context.TODO(), roleBinding, metav1.DeleteOptions{})
err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Delete(ctx, roleBinding, metav1.DeleteOptions{})
if err == nil || k8serrors.IsNotFound(err) {
return nil
}