Ensuring passing context across fission (#2555)
Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
@@ -153,7 +153,7 @@ func (c *Client) service() {
|
||||
svcReqs = append(svcReqs, req)
|
||||
}
|
||||
c.logger.Debug("tapped services in batch", zap.Int("service_count", len(urls)))
|
||||
err := c._tapService(context.TODO(), svcReqs)
|
||||
err := c._tapService(context.Background(), svcReqs)
|
||||
if err != nil {
|
||||
c.logger.Error("error tapping function service address", zap.Error(err))
|
||||
}
|
||||
|
||||
@@ -295,7 +295,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st
|
||||
}
|
||||
gpmPodInformer := gpmInformerFactory.Core().V1().Pods()
|
||||
gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets()
|
||||
gpm, err := poolmgr.MakeGenericPoolManager(
|
||||
gpm, err := poolmgr.MakeGenericPoolManager(ctx,
|
||||
logger,
|
||||
fissionClient, kubernetesClient, metricsClient,
|
||||
functionNamespace, fetcherConfig, executorInstanceID,
|
||||
@@ -311,7 +311,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st
|
||||
}
|
||||
ndmDeplInformer := ndmInformerFactory.Apps().V1().Deployments()
|
||||
ndmSvcInformer := ndmInformerFactory.Core().V1().Services()
|
||||
ndm, err := newdeploy.MakeNewDeploy(
|
||||
ndm, err := newdeploy.MakeNewDeploy(ctx,
|
||||
logger,
|
||||
fissionClient, kubernetesClient,
|
||||
functionNamespace, fetcherConfig, executorInstanceID,
|
||||
|
||||
@@ -51,8 +51,8 @@ func panicIf(err error) {
|
||||
}
|
||||
|
||||
// return the number of pods in the given namespace matching the given labels
|
||||
func countPods(kubeClient kubernetes.Interface, ns string, labelz map[string]string) int {
|
||||
pods, err := kubeClient.CoreV1().Pods(ns).List(context.TODO(), metav1.ListOptions{
|
||||
func countPods(ctx context.Context, kubeClient kubernetes.Interface, ns string, labelz map[string]string) int {
|
||||
pods, err := kubeClient.CoreV1().Pods(ns).List(ctx, metav1.ListOptions{
|
||||
LabelSelector: labels.Set(labelz).AsSelector().String(),
|
||||
})
|
||||
if err != nil {
|
||||
@@ -61,8 +61,8 @@ func countPods(kubeClient kubernetes.Interface, ns string, labelz map[string]str
|
||||
return len(pods.Items)
|
||||
}
|
||||
|
||||
func createTestNamespace(kubeClient kubernetes.Interface, ns string) {
|
||||
_, err := kubeClient.CoreV1().Namespaces().Create(context.TODO(), &apiv1.Namespace{
|
||||
func createTestNamespace(ctx context.Context, kubeClient kubernetes.Interface, ns string) {
|
||||
_, err := kubeClient.CoreV1().Namespaces().Create(ctx, &apiv1.Namespace{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: ns,
|
||||
},
|
||||
@@ -74,8 +74,8 @@ func createTestNamespace(kubeClient kubernetes.Interface, ns string) {
|
||||
}
|
||||
|
||||
// create a nodeport service
|
||||
func createSvc(kubeClient kubernetes.Interface, ns string, name string, targetPort int, nodePort int32, labels map[string]string) *apiv1.Service {
|
||||
svc, err := kubeClient.CoreV1().Services(ns).Create(context.TODO(), &apiv1.Service{
|
||||
func createSvc(ctx context.Context, kubeClient kubernetes.Interface, ns string, name string, targetPort int, nodePort int32, labels map[string]string) *apiv1.Service {
|
||||
svc, err := kubeClient.CoreV1().Services(ns).Create(ctx, &apiv1.Service{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: name,
|
||||
},
|
||||
@@ -120,18 +120,19 @@ func TestExecutor(t *testing.T) {
|
||||
log.Panicf("failed to connect: %v", err)
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
// create the test's namespaces
|
||||
createTestNamespace(kubeClient, fissionNs)
|
||||
createTestNamespace(ctx, kubeClient, fissionNs)
|
||||
defer func() {
|
||||
err := kubeClient.CoreV1().Namespaces().Delete(context.TODO(), fissionNs, metav1.DeleteOptions{})
|
||||
err := kubeClient.CoreV1().Namespaces().Delete(ctx, fissionNs, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
log.Fatalf("failed to delete namespace: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
createTestNamespace(kubeClient, functionNs)
|
||||
createTestNamespace(ctx, kubeClient, functionNs)
|
||||
defer func() {
|
||||
err := kubeClient.CoreV1().Namespaces().Delete(context.TODO(), functionNs, metav1.DeleteOptions{})
|
||||
err := kubeClient.CoreV1().Namespaces().Delete(ctx, functionNs, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
log.Fatalf("failed to delete namespace: %v", err)
|
||||
}
|
||||
@@ -143,18 +144,18 @@ func TestExecutor(t *testing.T) {
|
||||
panicIf(err)
|
||||
|
||||
// make sure CRD types exist on cluster
|
||||
err = crd.EnsureFissionCRDs(context.TODO(), logger, apiExtClient)
|
||||
err = crd.EnsureFissionCRDs(ctx, logger, apiExtClient)
|
||||
if err != nil {
|
||||
log.Panicf("failed to ensure crds: %v", err)
|
||||
}
|
||||
|
||||
err = crd.WaitForCRDs(context.TODO(), fissionClient)
|
||||
err = crd.WaitForCRDs(ctx, fissionClient)
|
||||
if err != nil {
|
||||
log.Panicf("failed to wait crds: %v", err)
|
||||
}
|
||||
|
||||
// create an env on the cluster
|
||||
env, err := fissionClient.CoreV1().Environments(fissionNs).Create(context.TODO(), &fv1.Environment{
|
||||
env, err := fissionClient.CoreV1().Environments(fissionNs).Create(ctx, &fv1.Environment{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: "nodejs",
|
||||
Namespace: fissionNs,
|
||||
@@ -173,7 +174,6 @@ func TestExecutor(t *testing.T) {
|
||||
|
||||
// create poolmgr
|
||||
port := 9999
|
||||
ctx := context.Background()
|
||||
err = StartExecutor(ctx, logger, functionNs, "fission-builder", port)
|
||||
if err != nil {
|
||||
log.Panicf("failed to start poolmgr: %v", err)
|
||||
@@ -208,7 +208,7 @@ func TestExecutor(t *testing.T) {
|
||||
Deployment: deployment,
|
||||
},
|
||||
}
|
||||
p, err = fissionClient.CoreV1().Packages(fissionNs).Create(context.TODO(), p, metav1.CreateOptions{})
|
||||
p, err = fissionClient.CoreV1().Packages(fissionNs).Create(ctx, p, metav1.CreateOptions{})
|
||||
if err != nil {
|
||||
log.Panicf("failed to create package: %v", err)
|
||||
}
|
||||
@@ -230,7 +230,7 @@ func TestExecutor(t *testing.T) {
|
||||
},
|
||||
},
|
||||
}
|
||||
_, err = fissionClient.CoreV1().Functions(fissionNs).Create(context.TODO(), f, metav1.CreateOptions{})
|
||||
_, err = fissionClient.CoreV1().Functions(fissionNs).Create(ctx, f, metav1.CreateOptions{})
|
||||
if err != nil {
|
||||
log.Panicf("failed to create function: %v", err)
|
||||
}
|
||||
@@ -238,18 +238,18 @@ func TestExecutor(t *testing.T) {
|
||||
// create a service to call fetcher and the env container
|
||||
labels := map[string]string{"functionName": f.ObjectMeta.Name}
|
||||
var fetcherPort int32 = 30001
|
||||
fetcherSvc := createSvc(kubeClient, functionNs, fmt.Sprintf("%v-%v", f.ObjectMeta.Name, "fetcher"), 8000, fetcherPort, labels)
|
||||
fetcherSvc := createSvc(ctx, kubeClient, functionNs, fmt.Sprintf("%v-%v", f.ObjectMeta.Name, "fetcher"), 8000, fetcherPort, labels)
|
||||
defer func() {
|
||||
err := kubeClient.CoreV1().Services(functionNs).Delete(context.TODO(), fetcherSvc.ObjectMeta.Name, metav1.DeleteOptions{})
|
||||
err := kubeClient.CoreV1().Services(functionNs).Delete(ctx, fetcherSvc.ObjectMeta.Name, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
log.Fatalf("failed to delete service: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
var funcSvcPort int32 = 30002
|
||||
functionSvc := createSvc(kubeClient, functionNs, f.ObjectMeta.Name, 8888, funcSvcPort, labels)
|
||||
functionSvc := createSvc(ctx, kubeClient, functionNs, f.ObjectMeta.Name, 8888, funcSvcPort, labels)
|
||||
defer func() {
|
||||
err := kubeClient.CoreV1().Services(functionNs).Delete(context.TODO(), functionSvc.ObjectMeta.Name, metav1.DeleteOptions{})
|
||||
err := kubeClient.CoreV1().Services(functionNs).Delete(ctx, functionSvc.ObjectMeta.Name, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
log.Fatalf("failed to delete service: %v", err)
|
||||
}
|
||||
@@ -257,14 +257,14 @@ func TestExecutor(t *testing.T) {
|
||||
|
||||
// the main test: get a service for a given function
|
||||
t1 := time.Now()
|
||||
svc, err := poolmgrClient.GetServiceForFunction(context.TODO(), f)
|
||||
svc, err := poolmgrClient.GetServiceForFunction(ctx, f)
|
||||
if err != nil {
|
||||
log.Panicf("failed to get func svc: %v", err)
|
||||
}
|
||||
log.Printf("svc for function created at: %v (in %v)", svc, time.Since(t1))
|
||||
|
||||
// ensure that a pod with the label functionName=f.ObjectMeta.Name exists
|
||||
podCount := countPods(kubeClient, functionNs, map[string]string{"functionName": f.ObjectMeta.Name})
|
||||
podCount := countPods(ctx, kubeClient, functionNs, map[string]string{"functionName": f.ObjectMeta.Name})
|
||||
if podCount != 1 {
|
||||
log.Panicf("expected 1 function pod, found %v", podCount)
|
||||
}
|
||||
|
||||
@@ -25,14 +25,13 @@ import (
|
||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||
)
|
||||
|
||||
func (deploy *NewDeploy) EnvEventHandlers() k8sCache.ResourceEventHandlerFuncs {
|
||||
func (deploy *NewDeploy) EnvEventHandlers(ctx context.Context) k8sCache.ResourceEventHandlerFuncs {
|
||||
return k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: func(obj interface{}) {},
|
||||
DeleteFunc: func(obj interface{}) {},
|
||||
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))
|
||||
|
||||
@@ -24,14 +24,13 @@ import (
|
||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||
)
|
||||
|
||||
func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFuncs {
|
||||
func (deploy *NewDeploy) FunctionEventHandlers(ctx context.Context) k8sCache.ResourceEventHandlerFuncs {
|
||||
return k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: func(obj interface{}) {
|
||||
// TODO: A workaround to process items in parallel. We should use workqueue ("k8s.io/client-go/util/workqueue")
|
||||
// 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(ctx, fn)
|
||||
@@ -46,7 +45,6 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu
|
||||
DeleteFunc: func(obj interface{}) {
|
||||
fn := obj.(*fv1.Function)
|
||||
go func() {
|
||||
ctx := context.Background()
|
||||
err := deploy.deleteFunction(ctx, fn)
|
||||
if err != nil {
|
||||
deploy.logger.Error("error deleting function",
|
||||
@@ -59,7 +57,6 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu
|
||||
oldFn := oldObj.(*fv1.Function)
|
||||
newFn := newObj.(*fv1.Function)
|
||||
go func() {
|
||||
ctx := context.Background()
|
||||
err := deploy.updateFunction(ctx, oldFn, newFn)
|
||||
if err != nil {
|
||||
deploy.logger.Error("error updating function",
|
||||
|
||||
@@ -120,7 +120,7 @@ func (deploy *NewDeploy) createOrGetDeployment(ctx context.Context, fn *fv1.Func
|
||||
|
||||
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)
|
||||
err := deploy.fetcherConfig.SetupServiceAccount(ctx, deploy.kubernetesClient, deployNamespace, fn.ObjectMeta)
|
||||
if err != nil {
|
||||
deploy.logger.Error("error creating fission fetcher service account for function",
|
||||
zap.Error(err),
|
||||
|
||||
@@ -96,6 +96,7 @@ type (
|
||||
|
||||
// MakeNewDeploy initializes and returns an instance of NewDeploy.
|
||||
func MakeNewDeploy(
|
||||
ctx context.Context,
|
||||
logger *zap.Logger,
|
||||
fissionClient versioned.Interface,
|
||||
kubernetesClient kubernetes.Interface,
|
||||
@@ -145,8 +146,8 @@ func MakeNewDeploy(
|
||||
nd.svcLister = svcInformer.Lister()
|
||||
nd.svcListerSynced = svcInformer.Informer().HasSynced
|
||||
|
||||
funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers())
|
||||
envInformer.Informer().AddEventHandler(nd.EnvEventHandlers())
|
||||
funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx))
|
||||
envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx))
|
||||
|
||||
return nd, nil
|
||||
}
|
||||
@@ -402,8 +403,7 @@ func (deploy *NewDeploy) deleteFunction(ctx context.Context, fn *fv1.Function) e
|
||||
}
|
||||
|
||||
func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
|
||||
cleanupFunc := func(ns string, name string) {
|
||||
ctx := context.Background()
|
||||
cleanupFunc := func(ctx context.Context, ns string, name string) {
|
||||
err := deploy.cleanupNewdeploy(ctx, ns, name)
|
||||
if err != nil {
|
||||
deploy.logger.Error("received error while cleaning function resources",
|
||||
@@ -436,7 +436,7 @@ func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscac
|
||||
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)
|
||||
go cleanupFunc(context.Background(), ns, objName)
|
||||
return nil, errors.Wrapf(err, "error creating service %v", objName)
|
||||
}
|
||||
svcAddress := fmt.Sprintf("%v.%v", svc.Name, svc.Namespace)
|
||||
@@ -444,14 +444,14 @@ func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscac
|
||||
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)
|
||||
go cleanupFunc(context.Background(), ns, objName)
|
||||
return nil, errors.Wrapf(err, "error creating deployment %v", objName)
|
||||
}
|
||||
|
||||
hpa, err := deploy.hpaops.CreateOrGetHpa(ctx, objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
|
||||
if err != nil {
|
||||
deploy.logger.Error("error creating HPA", zap.Error(err), zap.String("hpa", objName))
|
||||
go cleanupFunc(ns, objName)
|
||||
go cleanupFunc(context.Background(), ns, objName)
|
||||
return nil, errors.Wrapf(err, "error creating the HPA %v", objName)
|
||||
}
|
||||
|
||||
|
||||
@@ -76,7 +76,7 @@ func TestRefreshFuncPods(t *testing.T) {
|
||||
t.Fatalf("Error creating fetcher config: %s", err)
|
||||
}
|
||||
|
||||
executor, err := MakeNewDeploy(logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, "test",
|
||||
executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, "test",
|
||||
funcInformer, envInformer, deployInformer, svcInformer, podSpecPatch)
|
||||
if err != nil {
|
||||
t.Fatalf("new deploy manager creation failed: %s", err)
|
||||
|
||||
@@ -41,10 +41,9 @@ func getIstioServiceLabels(fnName string) map[string]string {
|
||||
// Based on function create/update/delete event, we create role binding
|
||||
// for the secret/configmap access which is used by fetcher component.
|
||||
// If istio is enabled, we create a service for the function.
|
||||
func FunctionEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs {
|
||||
func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, 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,
|
||||
@@ -133,7 +132,6 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Inter
|
||||
},
|
||||
|
||||
DeleteFunc: func(obj interface{}) {
|
||||
ctx := context.Background()
|
||||
fn := obj.(*fv1.Function)
|
||||
|
||||
fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType
|
||||
@@ -183,7 +181,6 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Inter
|
||||
if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault {
|
||||
envNs = newFunc.Spec.Environment.Namespace
|
||||
}
|
||||
ctx := context.Background()
|
||||
err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.SecretConfigMapGetterRB,
|
||||
newFunc.ObjectMeta.Namespace, utils.GetSecretConfigMapGetterCR(), fv1.ClusterRole,
|
||||
fv1.FissionFetcherSA, envNs)
|
||||
|
||||
@@ -142,7 +142,7 @@ func MakeGenericPool(
|
||||
|
||||
func (gp *GenericPool) setup(ctx context.Context) error {
|
||||
// create fetcher SA in this ns, if not already created
|
||||
err := gp.fetcherConfig.SetupServiceAccount(gp.kubernetesClient, gp.namespace, nil)
|
||||
err := gp.fetcherConfig.SetupServiceAccount(ctx, gp.kubernetesClient, gp.namespace, nil)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "error creating fetcher service account in namespace %q", gp.namespace)
|
||||
}
|
||||
@@ -156,7 +156,7 @@ func (gp *GenericPool) setup(ctx context.Context) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
go gp.updateCPUUtilizationSvc()
|
||||
go gp.updateCPUUtilizationSvc(ctx)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -185,7 +185,7 @@ func (gp *GenericPool) checkMetricsApi() bool {
|
||||
return utils.SupportedMetricsAPIVersionAvailable(apiGroups)
|
||||
}
|
||||
|
||||
func (gp *GenericPool) updateCPUUtilizationSvc() {
|
||||
func (gp *GenericPool) updateCPUUtilizationSvc(ctx context.Context) {
|
||||
var metricsApiAvailabe bool
|
||||
checkDuration := 30
|
||||
|
||||
@@ -194,8 +194,8 @@ func (gp *GenericPool) updateCPUUtilizationSvc() {
|
||||
gp.logger.Warn("Metrics API not available")
|
||||
}
|
||||
|
||||
serviceFunc := func() {
|
||||
podMetricsList, err := gp.metricsClient.MetricsV1beta1().PodMetricses(gp.namespace).List(context.TODO(), metav1.ListOptions{
|
||||
serviceFunc := func(ctx context.Context) {
|
||||
podMetricsList, err := gp.metricsClient.MetricsV1beta1().PodMetricses(gp.namespace).List(ctx, metav1.ListOptions{
|
||||
LabelSelector: "managed=false",
|
||||
})
|
||||
if err != nil {
|
||||
@@ -220,7 +220,7 @@ func (gp *GenericPool) updateCPUUtilizationSvc() {
|
||||
|
||||
for {
|
||||
if metricsApiAvailabe {
|
||||
serviceFunc()
|
||||
serviceFunc(ctx)
|
||||
} else {
|
||||
if gp.checkMetricsApi() {
|
||||
metricsApiAvailabe = true
|
||||
@@ -342,22 +342,20 @@ func (gp *GenericPool) labelsForFunction(metadata *metav1.ObjectMeta) map[string
|
||||
return label
|
||||
}
|
||||
|
||||
func (gp *GenericPool) scheduleDeletePod(name string) {
|
||||
go func() {
|
||||
// The sleep allows debugging or collecting logs from the pod before it's
|
||||
// cleaned up. (We need a better solutions for both those things; log
|
||||
// aggregation and storage will help.)
|
||||
gp.logger.Error("error in pod - scheduling cleanup", zap.String("pod", name))
|
||||
err := gp.kubernetesClient.CoreV1().Pods(gp.namespace).Delete(context.TODO(), name, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
gp.logger.Error(
|
||||
"error deleting pod",
|
||||
zap.String("name", name),
|
||||
zap.String("namespace", gp.namespace),
|
||||
zap.Error(err),
|
||||
)
|
||||
}
|
||||
}()
|
||||
func (gp *GenericPool) scheduleDeletePod(ctx context.Context, name string) {
|
||||
// The sleep allows debugging or collecting logs from the pod before it's
|
||||
// cleaned up. (We need a better solutions for both those things; log
|
||||
// aggregation and storage will help.)
|
||||
gp.logger.Error("error in pod - scheduling cleanup", zap.String("pod", name))
|
||||
err := gp.kubernetesClient.CoreV1().Pods(gp.namespace).Delete(ctx, name, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
gp.logger.Error(
|
||||
"error deleting pod",
|
||||
zap.String("name", name),
|
||||
zap.String("namespace", gp.namespace),
|
||||
zap.Error(err),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// IsIPv6 validates if the podIP follows to IPv6 protocol
|
||||
@@ -502,7 +500,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
|
||||
gp.readyPodQueue.Done(key)
|
||||
err = gp.specializePod(ctx, pod, fn)
|
||||
if err != nil {
|
||||
gp.scheduleDeletePod(pod.ObjectMeta.Name)
|
||||
go gp.scheduleDeletePod(context.Background(), pod.ObjectMeta.Name)
|
||||
return nil, err
|
||||
}
|
||||
logger.Info("specialized pod", zap.String("pod", pod.ObjectMeta.Name), zap.String("podNamespace", pod.ObjectMeta.Namespace), zap.String("podIP", pod.Status.PodIP))
|
||||
@@ -516,11 +514,11 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
|
||||
|
||||
svc, err := gp.createSvc(ctx, svcName, funcLabels)
|
||||
if err != nil {
|
||||
gp.scheduleDeletePod(pod.ObjectMeta.Name)
|
||||
go gp.scheduleDeletePod(context.Background(), pod.ObjectMeta.Name)
|
||||
return nil, err
|
||||
}
|
||||
if svc.ObjectMeta.Name != svcName {
|
||||
gp.scheduleDeletePod(pod.ObjectMeta.Name)
|
||||
go gp.scheduleDeletePod(context.Background(), pod.ObjectMeta.Name)
|
||||
return nil, errors.Errorf("sanity check failed for svc %v", svc.ObjectMeta.Name)
|
||||
}
|
||||
|
||||
|
||||
@@ -111,7 +111,7 @@ type (
|
||||
}
|
||||
)
|
||||
|
||||
func MakeGenericPoolManager(
|
||||
func MakeGenericPoolManager(ctx context.Context,
|
||||
logger *zap.Logger,
|
||||
fissionClient versioned.Interface,
|
||||
kubernetesClient kubernetes.Interface,
|
||||
@@ -138,7 +138,7 @@ func MakeGenericPoolManager(
|
||||
enableIstio = istio
|
||||
}
|
||||
|
||||
poolPodC := NewPoolPodController(gpmLogger, kubernetesClient, functionNamespace,
|
||||
poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient, functionNamespace,
|
||||
enableIstio, funcInformer, pkgInformer, envInformer, rsInformer, podInformer)
|
||||
|
||||
gpm := &GenericPoolManager{
|
||||
@@ -171,10 +171,10 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) {
|
||||
}
|
||||
go gpm.service()
|
||||
gpm.poolPodC.InjectGpm(gpm)
|
||||
go gpm.WebsocketStartEventChecker(gpm.kubernetesClient)
|
||||
go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient)
|
||||
go gpm.WebsocketStartEventChecker(ctx, gpm.kubernetesClient)
|
||||
go gpm.NoActiveConnectionEventChecker(ctx, gpm.kubernetesClient)
|
||||
go gpm.idleObjectReaper(ctx)
|
||||
go gpm.poolPodC.Run(ctx.Done())
|
||||
go gpm.poolPodC.Run(ctx, ctx.Done())
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType {
|
||||
@@ -664,17 +664,17 @@ func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) {
|
||||
}
|
||||
|
||||
// WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event
|
||||
func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient kubernetes.Interface) {
|
||||
func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, kubeClient kubernetes.Interface) {
|
||||
|
||||
informer := k8sCache.NewSharedInformer(
|
||||
&k8sCache.ListWatch{
|
||||
ListFunc: func(options metav1.ListOptions) (runtime.Object, error) {
|
||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=WsConnectionStarted"
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(context.TODO(), options)
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(ctx, options)
|
||||
},
|
||||
WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) {
|
||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=WsConnectionStarted"
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(context.TODO(), options)
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(ctx, options)
|
||||
},
|
||||
},
|
||||
&apiv1.Event{},
|
||||
@@ -705,17 +705,17 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient kubernetes.
|
||||
}
|
||||
|
||||
// NoActiveConnectionEventChecker checks if the pod has emitted an inactive event
|
||||
func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient kubernetes.Interface) {
|
||||
func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Context, kubeClient kubernetes.Interface) {
|
||||
|
||||
informer := k8sCache.NewSharedInformer(
|
||||
&k8sCache.ListWatch{
|
||||
ListFunc: func(options metav1.ListOptions) (runtime.Object, error) {
|
||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=NoActiveConnections"
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(context.TODO(), options)
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(ctx, options)
|
||||
},
|
||||
WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) {
|
||||
options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=NoActiveConnections"
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(context.TODO(), options)
|
||||
return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(ctx, options)
|
||||
},
|
||||
},
|
||||
&apiv1.Event{},
|
||||
@@ -737,7 +737,6 @@ func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient kuberne
|
||||
gpm.logger.Error("could not convert value from PodToFsvc")
|
||||
return
|
||||
}
|
||||
ctx := context.Background()
|
||||
gpm.fsCache.DeleteFunctionSvc(ctx, fsvc)
|
||||
for i := range fsvc.KubernetesObjects {
|
||||
gpm.logger.Info("release idle function resources due to inactivity",
|
||||
|
||||
@@ -31,7 +31,7 @@ import (
|
||||
// PackageEventHandlers provides handlers for package events.
|
||||
// Based on package create/update event, we create role binding
|
||||
// for the package which is used by fetcher component.
|
||||
func PackageEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs {
|
||||
func PackageEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs {
|
||||
return k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: func(obj interface{}) {
|
||||
pkg := obj.(*fv1.Package)
|
||||
@@ -44,7 +44,6 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interf
|
||||
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(ctx, logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, utils.GetPackageGetterCR(), fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
|
||||
@@ -81,7 +80,6 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interf
|
||||
envNs = newPkg.Spec.Environment.Namespace
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.PackageGetterRB,
|
||||
newPkg.ObjectMeta.Namespace, utils.GetPackageGetterCR(), fv1.ClusterRole,
|
||||
fv1.FissionFetcherSA, envNs)
|
||||
|
||||
@@ -66,7 +66,7 @@ type (
|
||||
}
|
||||
)
|
||||
|
||||
func NewPoolPodController(logger *zap.Logger,
|
||||
func NewPoolPodController(ctx context.Context, logger *zap.Logger,
|
||||
kubernetesClient kubernetes.Interface,
|
||||
namespace string,
|
||||
enableIstio bool,
|
||||
@@ -86,8 +86,8 @@ func NewPoolPodController(logger *zap.Logger,
|
||||
envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"),
|
||||
spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"),
|
||||
}
|
||||
funcInformer.Informer().AddEventHandler(FunctionEventHandlers(p.logger, p.kubernetesClient, p.namespace, p.enableIstio))
|
||||
pkgInformer.Informer().AddEventHandler(PackageEventHandlers(p.logger, p.kubernetesClient, p.namespace))
|
||||
funcInformer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio))
|
||||
pkgInformer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace))
|
||||
envInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: p.enqueueEnvAdd,
|
||||
UpdateFunc: p.enqueueEnvUpdate,
|
||||
@@ -216,7 +216,7 @@ func (p *PoolPodController) enqueueEnvDelete(obj interface{}) {
|
||||
p.envDeleteQueue.Add(env)
|
||||
}
|
||||
|
||||
func (p *PoolPodController) Run(stopCh <-chan struct{}) {
|
||||
func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) {
|
||||
defer utilruntime.HandleCrash()
|
||||
defer p.envCreateUpdateQueue.ShutDown()
|
||||
defer p.envDeleteQueue.ShutDown()
|
||||
@@ -228,20 +228,20 @@ func (p *PoolPodController) Run(stopCh <-chan struct{}) {
|
||||
p.logger.Fatal("failed to wait for caches to sync")
|
||||
}
|
||||
for i := 0; i < 4; i++ {
|
||||
go wait.Until(p.workerRun("envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh)
|
||||
go wait.Until(p.workerRun(ctx, "envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh)
|
||||
}
|
||||
go wait.Until(p.workerRun("envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh)
|
||||
go wait.Until(p.workerRun("spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh)
|
||||
go wait.Until(p.workerRun(ctx, "envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh)
|
||||
go wait.Until(p.workerRun(ctx, "spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh)
|
||||
p.logger.Info("Started workers for poolPodController")
|
||||
<-stopCh
|
||||
p.logger.Info("Shutting down workers for poolPodController")
|
||||
}
|
||||
|
||||
func (p *PoolPodController) workerRun(name string, processFunc func() bool) func() {
|
||||
func (p *PoolPodController) workerRun(ctx context.Context, name string, processFunc func(ctx context.Context) bool) func() {
|
||||
return func() {
|
||||
p.logger.Debug("Starting worker with func", zap.String("name", name))
|
||||
for {
|
||||
if quit := processFunc(); quit {
|
||||
if quit := processFunc(ctx); quit {
|
||||
p.logger.Info("Shutting down worker", zap.String("name", name))
|
||||
return
|
||||
}
|
||||
@@ -249,7 +249,7 @@ func (p *PoolPodController) workerRun(name string, processFunc func() bool) func
|
||||
}
|
||||
}
|
||||
|
||||
func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool {
|
||||
func (p *PoolPodController) envCreateUpdateQueueProcessFunc(ctx context.Context) bool {
|
||||
maxRetries := 3
|
||||
handleEnv := func(ctx context.Context, env *fv1.Environment) error {
|
||||
log := p.logger.With(zap.String("env", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace))
|
||||
@@ -310,7 +310,6 @@ func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool {
|
||||
return false
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
err = handleEnv(ctx, env)
|
||||
if err != nil {
|
||||
if p.envCreateUpdateQueue.NumRequeues(key) < maxRetries {
|
||||
@@ -326,7 +325,7 @@ func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func (p *PoolPodController) envDeleteQueueProcessFunc() bool {
|
||||
func (p *PoolPodController) envDeleteQueueProcessFunc(ctx context.Context) bool {
|
||||
obj, quit := p.envDeleteQueue.Get()
|
||||
if quit {
|
||||
return true
|
||||
@@ -338,7 +337,6 @@ func (p *PoolPodController) envDeleteQueueProcessFunc() bool {
|
||||
p.envDeleteQueue.Forget(obj)
|
||||
return false
|
||||
}
|
||||
ctx := context.Background()
|
||||
p.logger.Debug("env delete request processing")
|
||||
p.gpm.cleanupPool(ctx, env)
|
||||
specializePodLables := getSpecializedPodLabels(env)
|
||||
@@ -368,7 +366,7 @@ func (p *PoolPodController) envDeleteQueueProcessFunc() bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool {
|
||||
func (p *PoolPodController) spCleanupPodQueueProcessFunc(ctx context.Context) bool {
|
||||
maxRetries := 3
|
||||
obj, quit := p.spCleanupPodQueue.Get()
|
||||
if quit {
|
||||
@@ -403,7 +401,6 @@ func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool {
|
||||
}
|
||||
return false
|
||||
}
|
||||
ctx := context.Background()
|
||||
podName := strings.SplitAfter(pod.GetName(), ".")
|
||||
if fsvc, ok := p.gpm.fsCache.PodToFsvc.Load(strings.TrimSuffix(podName[0], ".")); ok {
|
||||
fsvc, ok := fsvc.(*fscache.FuncSvc)
|
||||
|
||||
@@ -44,6 +44,8 @@ func runInformers(ctx context.Context, informers []k8sCache.SharedIndexInformer)
|
||||
}
|
||||
|
||||
func TestPoolPodControllerPodCleanup(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
logger := loggerfactory.GetLogger()
|
||||
kubernetesClient := fake.NewSimpleClientset()
|
||||
fissionClient := fClient.NewSimpleClientset()
|
||||
@@ -60,7 +62,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
|
||||
gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets()
|
||||
|
||||
fnNamespace := "fission-function"
|
||||
ppc := NewPoolPodController(logger, kubernetesClient, fnNamespace, false,
|
||||
ppc := NewPoolPodController(ctx, logger, kubernetesClient, fnNamespace, false,
|
||||
funcInformer,
|
||||
pkgInformer,
|
||||
envInformer,
|
||||
@@ -73,7 +75,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatalf("Error creating fetcher config: %v", err)
|
||||
}
|
||||
executor, err := MakeGenericPoolManager(
|
||||
executor, err := MakeGenericPoolManager(ctx,
|
||||
logger,
|
||||
fissionClient, kubernetesClient, metricsClient,
|
||||
fnNamespace, fetcherConfig, executorInstanceID,
|
||||
@@ -85,10 +87,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
|
||||
gpm := executor.(*GenericPoolManager)
|
||||
ppc.InjectGpm(gpm)
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
go ppc.Run(ctx.Done())
|
||||
go ppc.Run(ctx, ctx.Done())
|
||||
|
||||
podInformer := gpmPodInformer.Informer()
|
||||
|
||||
|
||||
@@ -184,7 +184,9 @@ func TestFunctionServiceNewCache(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
fsc.AddFunc(ctx, *fsvc)
|
||||
_, active, err := fsc.GetFuncSvc(ctx, fsvc.Function, 5)
|
||||
if err != nil {
|
||||
|
||||
@@ -60,7 +60,9 @@ securityContext:
|
||||
Data: configMapData,
|
||||
}
|
||||
|
||||
configmap, err := kubeClient.CoreV1().ConfigMaps("fission").Create(context.Background(), &testConfigMap, metav1.CreateOptions{})
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
configmap, err := kubeClient.CoreV1().ConfigMaps("fission").Create(ctx, &testConfigMap, metav1.CreateOptions{})
|
||||
if err != nil {
|
||||
t.Errorf("Error creating configmap %v", err)
|
||||
}
|
||||
@@ -106,7 +108,7 @@ securityContext:
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
got, err := GetSpecFromConfigMap(context.Background(), kubeClient, tt.cm, tt.cmns)
|
||||
got, err := GetSpecFromConfigMap(ctx, kubeClient, tt.cm, tt.cmns)
|
||||
if (err != nil) != tt.wantErr {
|
||||
t.Errorf("GetSpecFromConfigMap() error = %v, wantErr %v", err, tt.wantErr)
|
||||
return
|
||||
|
||||
Reference in New Issue
Block a user