diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index 869c098f..d2c42ce4 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -249,7 +249,7 @@ func (deploy *NewDeploy) getDeploymentSpec(ctx context.Context, fn *fv1.Function Env: []apiv1.EnvVar{ { Name: fv1.ResourceVersionCount, - Value: fmt.Sprintf("%v", rvCount), + Value: fmt.Sprintf("%d", rvCount), }, }, // https://istio.io/docs/setup/kubernetes/additional-setup/requirements/ diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 2c9c8b77..fc15d7ed 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -264,13 +264,13 @@ 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(ctx, deploy.kubernetesClient, deployment.Namespace, f.Spec.Secrets, f.Spec.ConfigMaps) + rvCount, err := referencedResourcesRVSum(ctx, deploy.kubernetesClient, f.ObjectMeta.Namespace, f.Spec.Secrets, f.Spec.ConfigMaps) if err != nil { return err } - patch := fmt.Sprintf(`{"spec" : {"template": {"spec":{"containers":[{"name": "%s", "env":[{"name": "%s", "value": "%v"}]}]}}}}`, - f.ObjectMeta.Name, fv1.ResourceVersionCount, rvCount) + patch := fmt.Sprintf(`{"spec" : {"template": {"spec":{"containers":[{"name": "%s", "image": "%s", "env":[{"name": "%s", "value": "%d"}]}]}}}}`, + env.ObjectMeta.Name, env.Spec.Runtime.Image, fv1.ResourceVersionCount, rvCount) _, err = deploy.kubernetesClient.AppsV1().Deployments(deployment.ObjectMeta.Namespace).Patch(ctx, deployment.ObjectMeta.Name, k8sTypes.StrategicMergePatchType, diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go new file mode 100644 index 00000000..06673067 --- /dev/null +++ b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go @@ -0,0 +1,254 @@ +package newdeploy + +import ( + "context" + "fmt" + "os" + "testing" + "time" + + uuid "github.com/satori/go.uuid" + "github.com/stretchr/testify/assert" + apiv1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/client-go/kubernetes/fake" + k8sCache "k8s.io/client-go/tools/cache" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/executor/util" + fetcherConfig "github.com/fission/fission/pkg/fetcher/config" + fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" + "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/loggerfactory" +) + +const ( + defaultNamespace string = "default" + functionNamespace string = "fission-function" + envName string = "newdeploy-test-env" + functionName string = "newdeploy-test-func" + configmapName string = "newdeploy-test-configmap" +) + +func runInformers(ctx context.Context, informers []k8sCache.SharedIndexInformer) { + // Run all informers + for _, informer := range informers { + go informer.Run(ctx.Done()) + } +} + +func TestRefreshFuncPods(t *testing.T) { + os.Setenv("DEBUG_ENV", "true") + logger := loggerfactory.GetLogger() + kubernetesClient := fake.NewSimpleClientset() + fissionClient := fClient.NewSimpleClientset() + informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) + funcInformer := informerFactory.Core().V1().Functions() + envInformer := informerFactory.Core().V1().Environments() + + newDeployInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30) + if err != nil { + t.Fatalf("Error creating informer factory: %s", err) + } + + deployInformer := newDeployInformerFactory.Apps().V1().Deployments() + svcInformer := newDeployInformerFactory.Core().V1().Services() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + err = BuildConfigMap(ctx, kubernetesClient, functionNamespace, fv1.RuntimePodSpecConfigmap, map[string]string{}) + if err != nil { + t.Fatalf("Error building configmap: %s", err) + } + + podSpecPatch, err := util.GetSpecFromConfigMap(ctx, kubernetesClient, fv1.RuntimePodSpecConfigmap, functionNamespace) + if err != nil { + t.Fatalf("Error creating pod spec: %s", err) + } + + fetcherConfig, err := fetcherConfig.MakeFetcherConfig("/userfunc") + if err != nil { + t.Fatalf("Error creating fetcher config: %s", err) + } + + executor, err := MakeNewDeploy(logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, "test", + funcInformer, envInformer, deployInformer, svcInformer, podSpecPatch) + if err != nil { + t.Fatalf("new deploy manager creation failed: %s", err) + } + + ndm := executor.(*NewDeploy) + + go ndm.Run(ctx) + t.Log("New deploy manager started") + + runInformers(ctx, []k8sCache.SharedIndexInformer{ + envInformer.Informer(), + funcInformer.Informer(), + deployInformer.Informer(), + svcInformer.Informer(), + }) + t.Log("Informers required for new deploy manager started") + + if ok := k8sCache.WaitForCacheSync(ctx.Done(), ndm.deplListerSynced, ndm.svcListerSynced); !ok { + t.Fatal("Timed out waiting for caches to sync") + } + + envSpec := &fv1.Environment{ + ObjectMeta: metav1.ObjectMeta{ + Name: envName, + Namespace: defaultNamespace, + UID: "83c82da2-81e9-4ebd-867e-f383e65e603f", + }, + Spec: fv1.EnvironmentSpec{ + Version: 1, + Runtime: fv1.Runtime{ + Image: "gcr.io/xyz", + }, + }, + } + + _, err = fissionClient.CoreV1().Environments(defaultNamespace).Create(ctx, envSpec, metav1.CreateOptions{}) + if err != nil { + t.Fatalf("creating environment failed : %s", err) + } + + envRes, err := fissionClient.CoreV1().Environments(defaultNamespace).Get(ctx, envName, metav1.GetOptions{}) + if err != nil { + t.Fatalf("Error getting environment: %s", err) + } + assert.Equal(t, envRes.ObjectMeta.Name, envName) + + funcUID, err := uuid.NewV4() + if err != nil { + t.Fatal(err) + } + funcSpec := fv1.Function{ + ObjectMeta: metav1.ObjectMeta{ + Name: functionName, + Namespace: defaultNamespace, + UID: types.UID(funcUID.String()), + }, + Spec: fv1.FunctionSpec{ + Environment: fv1.EnvironmentReference{ + Name: envName, + Namespace: defaultNamespace, + }, + InvokeStrategy: fv1.InvokeStrategy{ + ExecutionStrategy: fv1.ExecutionStrategy{ + ExecutorType: fv1.ExecutorTypeNewdeploy, + }, + }, + }, + } + _, err = fissionClient.CoreV1().Functions(defaultNamespace).Create(ctx, &funcSpec, metav1.CreateOptions{}) + if err != nil { + t.Fatalf("creating function failed : %s", err) + } + + funcRes, err := fissionClient.CoreV1().Functions(defaultNamespace).Get(ctx, functionName, metav1.GetOptions{}) + if err != nil { + t.Fatalf("Error getting function: %s", err) + } + assert.Equal(t, funcRes.ObjectMeta.Name, functionName) + + ctx2, cancel2 := context.WithCancel(context.Background()) + wait.Until(func() { + t.Log("Checking for deployment") + ret, err := kubernetesClient.AppsV1().Deployments(functionNamespace).List(ctx2, metav1.ListOptions{}) + if err != nil { + t.Fatalf("Error getting deployment: %s", err) + } + if len(ret.Items) > 0 { + t.Log("Deployment created", ret.Items[0].Name) + cancel2() + } + }, time.Second*2, ctx2.Done()) + + err = BuildConfigMap(ctx, kubernetesClient, defaultNamespace, configmapName, map[string]string{ + "test-key": "test-value", + }) + if err != nil { + t.Fatalf("Error building configmap: %s", err) + } + + t.Log("Adding configmap to function") + funcRes.Spec.ConfigMaps = []fv1.ConfigMapReference{ + { + Name: configmapName, + Namespace: defaultNamespace, + }, + } + _, err = fissionClient.CoreV1().Functions(defaultNamespace).Update(ctx, funcRes, metav1.UpdateOptions{}) + if err != nil { + t.Fatalf("Error updating function: %s", err) + } + funcRes, err = fissionClient.CoreV1().Functions(defaultNamespace).Get(ctx, functionName, metav1.GetOptions{}) + if err != nil { + t.Fatalf("Error getting function: %s", err) + } + assert.Greater(t, len(funcRes.Spec.ConfigMaps), 0) + + err = ndm.RefreshFuncPods(ctx, logger, *funcRes) + if err != nil { + t.Fatalf("Error refreshing function pods: %s", err) + } + + funcLabels := ndm.getDeployLabels(funcRes.ObjectMeta, envRes.ObjectMeta) + + dep, err := kubernetesClient.AppsV1().Deployments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{ + LabelSelector: labels.Set(funcLabels).AsSelector().String(), + }) + + if err != nil { + t.Fatalf("Error getting deployment: %s", err) + } + assert.Equal(t, len(dep.Items), 1) + + cm, err := kubernetesClient.CoreV1().ConfigMaps(defaultNamespace).Get(ctx, configmapName, metav1.GetOptions{}) + if err != nil { + t.Fatalf("Error getting configmap: %s", err) + } + assert.Equal(t, cm.ObjectMeta.Name, configmapName) + updatedDepl := dep.Items[0] + resourceVersionMatch := false + assert.Equal(t, len(updatedDepl.Spec.Template.Spec.Containers), 2) + for _, v := range updatedDepl.Spec.Template.Spec.Containers { + if v.Name == envName { + assert.Greater(t, len(v.Env), 0) + for _, env := range v.Env { + if env.Name == fv1.ResourceVersionCount { + assert.Equal(t, env.Value, cm.ObjectMeta.ResourceVersion) + resourceVersionMatch = true + } + } + } + } + assert.True(t, resourceVersionMatch) +} + +func FakeResourceVersion() string { + return fmt.Sprint(time.Now().Nanosecond())[:6] +} + +func BuildConfigMap(ctx context.Context, kubernetesClient *fake.Clientset, namespace, name string, data map[string]string) error { + testConfigMap := apiv1.ConfigMap{ + TypeMeta: metav1.TypeMeta{ + Kind: "ConfigMap", + APIVersion: "v1", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: namespace, + ResourceVersion: FakeResourceVersion(), + }, + Data: data, + } + _, err := kubernetesClient.CoreV1().ConfigMaps(namespace).Create(ctx, &testConfigMap, metav1.CreateOptions{}) + return err +} diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 9350a5df..aa983b1a 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -278,11 +278,11 @@ func (gpm *GenericPoolManager) RefreshFuncPods(ctx context.Context, logger *zap. } funcSvc, err := gp.fsCache.GetByFunction(&f.ObjectMeta) - if err != nil { - return err - } - gp.fsCache.DeleteEntry(funcSvc) + // delete function service address from cache only when function service address found in cache + if err == nil { + gp.fsCache.DeleteEntry(funcSvc) + } funcLabels := gp.labelsForFunction(&f.ObjectMeta)