From 12f8017d9bb0b54d48d347ac6a0d31732b8570e1 Mon Sep 17 00:00:00 2001 From: Vardhaman Surana <40083149+vardhaman-surana@users.noreply.github.com> Date: Thu, 21 Dec 2023 17:02:39 +0530 Subject: [PATCH] added fetcher test cases (#2893) * added fetcher test cases * error message formatting fix * error message formatting fix for e2e framework --- cmd/fetcher/app/server.go | 6 +- cmd/fetcher/main.go | 5 +- pkg/fetcher/fetcher.go | 7 +- test/e2e/fetcher/data/test-deploy-pkg.zip | Bin 0 -> 204 bytes .../data/test-specialize-deploy-pkg.zip | Bin 0 -> 188 bytes test/e2e/fetcher/data/test-url-arch.zip | Bin 0 -> 190 bytes test/e2e/fetcher/fetcher_test.go | 618 ++++++++++++++++++ test/e2e/framework/framework.go | 8 +- test/e2e/framework/services/services.go | 77 ++- 9 files changed, 679 insertions(+), 42 deletions(-) create mode 100644 test/e2e/fetcher/data/test-deploy-pkg.zip create mode 100644 test/e2e/fetcher/data/test-specialize-deploy-pkg.zip create mode 100644 test/e2e/fetcher/data/test-url-arch.zip create mode 100644 test/e2e/fetcher/fetcher_test.go diff --git a/cmd/fetcher/app/server.go b/cmd/fetcher/app/server.go index 67fa5a87..5f6a4ab8 100644 --- a/cmd/fetcher/app/server.go +++ b/cmd/fetcher/app/server.go @@ -39,7 +39,7 @@ var ( readyToServe uint32 ) -func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface) { +func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port string, podInfoMountDir string) { flag.Usage = fetcherUsage specializeOnStart := flag.Bool("specialize-on-startup", false, "Flag to activate specialize process at pod startup") specializePayload := flag.String("specialize-request", "", "JSON payload for specialize request") @@ -74,7 +74,7 @@ func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *za ctx, span := tracer.Start(ctx, "fetcher/Run") defer span.End() - f, err := fetcher.MakeFetcher(logger, clientGen, dir, *secretDir, *configDir) + f, err := fetcher.MakeFetcher(logger, clientGen, dir, *secretDir, *configDir, podInfoMountDir) if err != nil { logger.Fatal("error making fetcher", zap.Error(err)) } @@ -121,7 +121,7 @@ func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *za logger.Info("fetcher ready to receive requests") handler := otelUtils.GetHandlerWithOTEL(mux, "fission-fetcher", otelUtils.UrlsToIgnore("/healthz", "/readiness-healthz")) - httpserver.StartServer(ctx, logger, mgr, "fetcher", "8000", handler) + httpserver.StartServer(ctx, logger, mgr, "fetcher", port, handler) } func fetcherUsage() { diff --git a/cmd/fetcher/main.go b/cmd/fetcher/main.go index d14b964b..b724342c 100644 --- a/cmd/fetcher/main.go +++ b/cmd/fetcher/main.go @@ -20,12 +20,15 @@ import ( "sigs.k8s.io/controller-runtime/pkg/manager/signals" "github.com/fission/fission/cmd/fetcher/app" + fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/utils/loggerfactory" "github.com/fission/fission/pkg/utils/manager" "github.com/fission/fission/pkg/utils/profile" ) +const fetcherPort = "8000" + // Usage: fetcher func main() { @@ -38,5 +41,5 @@ func main() { ctx := signals.SetupSignalHandler() profile.ProfileIfEnabled(ctx, logger, mgr) - app.Run(ctx, crd.NewClientGenerator(), logger, mgr) + app.Run(ctx, crd.NewClientGenerator(), logger, mgr, fetcherPort, fv1.PodInfoMount) } diff --git a/pkg/fetcher/fetcher.go b/pkg/fetcher/fetcher.go index 7c1eca95..d7040443 100644 --- a/pkg/fetcher/fetcher.go +++ b/pkg/fetcher/fetcher.go @@ -75,7 +75,8 @@ func makeVolumeDir(dirPath string) error { return os.MkdirAll(dirPath, os.ModeDir|0750) } -func MakeFetcher(logger *zap.Logger, clientGen crd.ClientGeneratorInterface, sharedVolumePath string, sharedSecretPath string, sharedConfigPath string) (*Fetcher, error) { +func MakeFetcher(logger *zap.Logger, clientGen crd.ClientGeneratorInterface, sharedVolumePath string, sharedSecretPath string, + sharedConfigPath string, podInfoMountDir string) (*Fetcher, error) { fLogger := logger.Named("fetcher") err := makeVolumeDir(sharedVolumePath) if err != nil { @@ -99,12 +100,12 @@ func MakeFetcher(logger *zap.Logger, clientGen crd.ClientGeneratorInterface, sha return nil, errors.Wrap(err, "error making the kube client") } - name, err := os.ReadFile(fv1.PodInfoMount + "/name") + name, err := os.ReadFile(podInfoMountDir + "/name") if err != nil { return nil, errors.Wrap(err, "error reading pod name from downward volume") } - namespace, err := os.ReadFile(fv1.PodInfoMount + "/namespace") + namespace, err := os.ReadFile(podInfoMountDir + "/namespace") if err != nil { return nil, errors.Wrap(err, "error reading pod namespace from downward volume") } diff --git a/test/e2e/fetcher/data/test-deploy-pkg.zip b/test/e2e/fetcher/data/test-deploy-pkg.zip new file mode 100644 index 0000000000000000000000000000000000000000..f7a2828044a55e64b696566d3397340dab9810cb GIT binary patch literal 204 zcmWIWW@h1H00IBp>hRmsw}z?#*&xipAj6Q6nv;{SS5O%m!pXp#=bw_A55%Pv+zgB? zFPIq^z(h)FnnG@3W}b$o6_)}K6s4Aw7Ud}@d4TllD3s?H<)kPo1$Z+u$uZ-yNdn{m h21X#>(g}7@6i)LqlH!B-R9U~C>fz*RI3;@JeDvbaD literal 0 HcmV?d00001 diff --git a/test/e2e/fetcher/data/test-specialize-deploy-pkg.zip b/test/e2e/fetcher/data/test-specialize-deploy-pkg.zip new file mode 100644 index 0000000000000000000000000000000000000000..a783d634706f899c235140c48d6b310bf0291c24 GIT binary patch literal 188 zcmWIWW@h1H00EQA-f-6}=_VN<8-!VbWJac5L1kzNCj;}H>dI6QF0J5ZU}Sm0%)kI9 zQc}|tauYN2G&HTa6o8;8wWPEtPeI8eQ&B0vn~_P58JFP_AUhctfp|+Jhy}HZ6=D^d UH38nNY#>F9K@~ literal 0 HcmV?d00001 diff --git a/test/e2e/fetcher/fetcher_test.go b/test/e2e/fetcher/fetcher_test.go new file mode 100644 index 00000000..4dacfd63 --- /dev/null +++ b/test/e2e/fetcher/fetcher_test.go @@ -0,0 +1,618 @@ +package fetcher + +import ( + "context" + "errors" + "fmt" + "net/http" + "net/url" + "os" + "strconv" + "strings" + "testing" + "time" + + "github.com/fission/fission/cmd/fetcher/app" + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/fetcher" + "github.com/fission/fission/pkg/fetcher/client" + storageClient "github.com/fission/fission/pkg/storagesvc/client" + + "github.com/fission/fission/pkg/generated/clientset/versioned" + "github.com/fission/fission/pkg/utils" + "github.com/fission/fission/pkg/utils/httpserver" + "github.com/fission/fission/pkg/utils/loggerfactory" + "github.com/fission/fission/pkg/utils/manager" + "github.com/fission/fission/pkg/utils/profile" + "github.com/fission/fission/test/e2e/framework" + "github.com/fission/fission/test/e2e/framework/cli" + "github.com/fission/fission/test/e2e/framework/services" + "github.com/stretchr/testify/require" + "github.com/stretchr/testify/suite" + "go.uber.org/zap" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/client-go/kubernetes" +) + +const testFileData = ` +module.exports = async function (context) { + return { + status: 200, + body: "Hello, Fission!\n" + }; +} +` + +type FetcherTestSuite struct { + suite.Suite + mgr manager.Interface + logger *zap.Logger + cancel context.CancelFunc + framework *framework.Framework + ctx context.Context + cfgMapsDir string + secretsDir string + sharedVolDir string + podInfoDir string + envName string + fetcherClient client.ClientInterface + fissionClient versioned.Interface + k8sClient kubernetes.Interface + storagesvcURL string + specTestData *specializeTestData +} + +type specializeTestData struct { + throwError bool +} + +func (f *FetcherTestSuite) SetupSuite() { + + var err error + f.envName = "test-pythondeploy" + const testEnvImage = "fission/python-env:latest" + const testBuilderImage = "fission/python-builder:latest" + + f.logger = loggerfactory.GetLogger() + f.framework = framework.NewFramework() + f.mgr = manager.New() + ctx, cancel := context.WithCancel(context.Background()) + f.cancel = cancel + f.ctx = ctx + f.specTestData = &specializeTestData{} + + // create a dummy server to test FETCH_URL source type and specialize handler + mux := http.NewServeMux() + mux.HandleFunc("/specialize", func(w http.ResponseWriter, r *http.Request) { + + if f.specTestData.throwError { + w.WriteHeader(http.StatusBadRequest) + return + } + w.WriteHeader(http.StatusOK) + }) + mux.HandleFunc("/v2/specialize", func(w http.ResponseWriter, r *http.Request) { + + if f.specTestData.throwError { + w.WriteHeader(http.StatusBadRequest) + return + } + w.WriteHeader(http.StatusOK) + }) + mux.Handle("/files/", http.StripPrefix("/files/", http.FileServer(http.Dir(".")))) + + f.mgr.Add(ctx, func(c context.Context) { + httpserver.StartServer(ctx, f.logger, f.mgr, "specialize", "8888", mux) + }) + + err = f.framework.Start(ctx) + require.NoError(f.T(), err) + + err = services.StartServices(ctx, f.framework, f.mgr) + require.NoError(f.T(), err) + + err = wait.PollUntilContextTimeout(ctx, time.Second*5, time.Second*50, true, func(_ context.Context) (bool, error) { + if err := f.framework.CheckService("webhook"); err != nil { + return false, nil + } + return true, nil + }) + require.NoError(f.T(), err) + + defer func() { + if err != nil { + f.TearDownSuite() + } + }() + + err = f.createFetcherTestDirs() + require.NoError(f.T(), err, "error creating fetcher test dirs") + + os.Args = append(os.Args, "-secret-dir", f.secretsDir, "-cfgmap-dir", f.cfgMapsDir) + os.Args = append(os.Args, f.sharedVolDir) + + err = f.setupPodInfoMountDir() + require.NoError(f.T(), err) + + freePort, err := utils.FindFreePort() + require.NoError(f.T(), err) + + fetcherPort := strconv.Itoa(freePort) + + fetcherUrl := "http://localhost:" + fetcherPort + "/" + + profile.ProfileIfEnabled(ctx, f.logger, f.mgr) + f.mgr.Add(ctx, func(c context.Context) { + app.Run(ctx, f.framework.ClientGen(), f.logger, f.mgr, fetcherPort, f.podInfoDir) + }) + + f.fetcherClient = client.MakeClient(f.logger, fetcherUrl) + + _, err = cli.ExecCommand(f.framework, ctx, "env", "create", "--name", f.envName, "--image", testEnvImage, + "--builder", testBuilderImage) + require.NoError(f.T(), err, "error in creating test environment:%s", f.envName) + + f.fissionClient, err = f.framework.ClientGen().GetFissionClient() + require.NoError(f.T(), err) + + f.k8sClient, err = f.framework.ClientGen().GetKubernetesClient() + require.NoError(f.T(), err) + + err = wait.PollUntilContextTimeout(f.ctx, 2*time.Second, 5*time.Minute, true, func(c context.Context) (done bool, err error) { + + deploys, err := f.k8sClient.AppsV1().Deployments(metav1.NamespaceDefault).List(c, metav1.ListOptions{}) + if err != nil || len(deploys.Items) == 0 { + return false, err + } + + for _, deploy := range deploys.Items { + if strings.Contains(deploy.Name, f.envName) && (deploy.Status.AvailableReplicas == deploy.Status.Replicas) { + return true, nil + } + } + + return false, nil + }) + + f.storagesvcURL, err = f.framework.GetServiceURL("storagesvc") + require.NoError(f.T(), err, "error getting storage service URL: %w", err) + +} + +func (f *FetcherTestSuite) TestFetcherDeploymentType() { + + const testDeployPkg = "test-deploy-pkg" + const testDeployArchiveFile = "data/test-deploy-pkg.zip" + const testDeployFuncName = "deploypy" + + _, err := cli.ExecCommand(f.framework, f.ctx, "package", "create", "--name", testDeployPkg, "--deployarchive", testDeployArchiveFile, "--env", f.envName) + require.NoError(f.T(), err, "error in creating package:%s in environment:%s", testDeployPkg, f.envName) + err = wait.PollUntilContextTimeout(f.ctx, 10*time.Second, 5*time.Minute, true, func(c context.Context) (done bool, err error) { + + pkgs, err := f.fissionClient.CoreV1().Packages(metav1.NamespaceDefault).List(c, metav1.ListOptions{}) + + if err != nil { + return false, err + } + + for _, pkg := range pkgs.Items { + if pkg.Name == testDeployPkg && pkg.Status.BuildStatus == fv1.BuildStatusSucceeded { + return true, nil + } else if pkg.Name == testDeployPkg && pkg.Status.BuildStatus == fv1.BuildStatusFailed { + return true, fmt.Errorf("build failed for package:%s", pkg.Name) + } + } + + return false, nil + }) + require.NoError(f.T(), err, "error while waiting for package build to succeed:%w", err) + + _, err = cli.ExecCommand(f.framework, f.ctx, "function", "create", "--name", testDeployFuncName, "--pkg", testDeployPkg, "--entrypoint", "hello.main") + require.NoError(f.T(), err) + + err = f.fetcherClient.Fetch(f.ctx, &fetcher.FunctionFetchRequest{ + Filename: "hello.py", + StorageSvcUrl: f.storagesvcURL, + KeepArchive: true, + FetchType: fv1.FETCH_DEPLOYMENT, + Package: metav1.ObjectMeta{ + Name: testDeployPkg, + Namespace: "default", + }, + }) + require.NoError(f.T(), err) + + file, err := os.OpenFile(f.sharedVolDir+"/hello.py", os.O_RDONLY, 0444) + require.NoError(f.T(), err, "fetched file does not exist in fetcher directory") + defer file.Close() + + _, err = cli.ExecCommand(f.framework, f.ctx, "function", "delete", "--name", testDeployFuncName) + require.NoError(f.T(), err) +} + +func (f *FetcherTestSuite) TestFetcherURLType() { + + const testDeployPkg = "test-url-arch-pkg" + const testDeployArchiveFile = "data/test-url-arch.zip" + const testDeployFuncName = "url-arch-test" + + _, err := cli.ExecCommand(f.framework, f.ctx, "package", "create", "--name", testDeployPkg, "--deployarchive", testDeployArchiveFile, "--env", f.envName) + require.NoError(f.T(), err, "error in creating package:%s in environment:%s", testDeployPkg, f.envName) + err = wait.PollUntilContextTimeout(f.ctx, 10*time.Second, 5*time.Minute, true, func(c context.Context) (done bool, err error) { + + pkgs, err := f.fissionClient.CoreV1().Packages(metav1.NamespaceDefault).List(c, metav1.ListOptions{}) + + if err != nil { + return false, err + } + + for _, pkg := range pkgs.Items { + if pkg.Name == testDeployPkg && pkg.Status.BuildStatus == fv1.BuildStatusSucceeded { + return true, nil + } else if pkg.Name == testDeployPkg && pkg.Status.BuildStatus == fv1.BuildStatusFailed { + return true, fmt.Errorf("build failed for package:%s", pkg.Name) + } + } + + return false, nil + }) + require.NoError(f.T(), err, "error while waiting for package build to succeed:%w", err) + + _, err = cli.ExecCommand(f.framework, f.ctx, "function", "create", "--name", testDeployFuncName, "--pkg", testDeployPkg, "--entrypoint", "hello.main") + require.NoError(f.T(), err) + + err = f.fetcherClient.Fetch(f.ctx, &fetcher.FunctionFetchRequest{ + Filename: "new.py", + StorageSvcUrl: f.storagesvcURL, + KeepArchive: true, + FetchType: fv1.FETCH_URL, + Url: "http://localhost:8888/files/data/test-url-arch.zip", + Package: metav1.ObjectMeta{ + Name: testDeployPkg, + Namespace: "default", + }, + }) + require.NoError(f.T(), err) + + file, err := os.OpenFile(f.sharedVolDir+"/new.py", os.O_RDONLY, 0444) + require.NoError(f.T(), err, "fetched file does not exist in fetcher directory") + defer file.Close() + + _, err = cli.ExecCommand(f.framework, f.ctx, "function", "delete", "--name", testDeployFuncName) + require.NoError(f.T(), err) +} + +func (f *FetcherTestSuite) TestFetcherUpload() { + + storageClient := storageClient.MakeClient(f.storagesvcURL) + + resp, err := f.fetcherClient.Upload(context.Background(), &fetcher.ArchiveUploadRequest{ + Filename: "hello.js", + StorageSvcUrl: f.storagesvcURL, + ArchivePackage: true, + }) + require.NoError(f.T(), err) + + archiveID, err := getArchiveIDFromURL(resp.ArchiveDownloadUrl) + require.NoError(f.T(), err) + + filesIds, err := storageClient.List(f.ctx) + require.NoError(f.T(), err) + + idFound := contains(filesIds, archiveID) + require.True(f.T(), idFound, "archive id not found in storagesvc list") + + resp, err = f.fetcherClient.Upload(context.Background(), &fetcher.ArchiveUploadRequest{ + Filename: "hello.js", + StorageSvcUrl: f.storagesvcURL, + ArchivePackage: false, + }) + require.NoError(f.T(), err) + + archiveID, err = getArchiveIDFromURL(resp.ArchiveDownloadUrl) + require.NoError(f.T(), err) + + filesIds, err = storageClient.List(f.ctx) + require.NoError(f.T(), err) + + idFound = contains(filesIds, archiveID) + require.True(f.T(), idFound, "archive id not found in storagesvc list") +} + +func (f *FetcherTestSuite) TestFetcherSpecialize() { + + const testDeployPkg = "test-deploy-pkg-specialize" + const testDeployArchiveFile = "data/test-specialize-deploy-pkg.zip" + const testDeployFuncName = "deploypy-specialize" + + configMapData := map[string]string{ + "env": "test", + } + + configMapName := "test-configmap" + + configMap := &v1.ConfigMap{ + ObjectMeta: metav1.ObjectMeta{ + Name: configMapName, + Namespace: metav1.NamespaceDefault, + }, + Data: configMapData, + } + + cfgMap, err := f.k8sClient.CoreV1().ConfigMaps(metav1.NamespaceDefault).Create(f.ctx, configMap, metav1.CreateOptions{}) + require.NoError(f.T(), err) + + secretName := "test-secret" + secretData := map[string][]byte{ + "username": []byte("admin"), + } + + secret := &v1.Secret{ + ObjectMeta: metav1.ObjectMeta{ + Name: secretName, + Namespace: metav1.NamespaceDefault, + }, + Data: secretData, + } + + secret, err = f.k8sClient.CoreV1().Secrets(metav1.NamespaceDefault).Create(f.ctx, secret, metav1.CreateOptions{}) + require.NoError(f.T(), err) + + _, err = cli.ExecCommand(f.framework, f.ctx, "package", "create", "--name", testDeployPkg, "--deployarchive", testDeployArchiveFile, "--env", f.envName) + require.NoError(f.T(), err, "error in creating package:%s in environment:%s", testDeployPkg, f.envName) + err = wait.PollUntilContextTimeout(f.ctx, 10*time.Second, 5*time.Minute, true, func(c context.Context) (done bool, err error) { + + pkgs, err := f.fissionClient.CoreV1().Packages(metav1.NamespaceDefault).List(c, metav1.ListOptions{}) + if err != nil { + return false, err + } + + for _, pkg := range pkgs.Items { + if pkg.Name == testDeployPkg && pkg.Status.BuildStatus == fv1.BuildStatusSucceeded { + return true, nil + } else if pkg.Name == testDeployPkg && pkg.Status.BuildStatus == fv1.BuildStatusFailed { + return true, fmt.Errorf("build failed for package:%s", pkg.Name) + } + } + + return false, nil + }) + require.NoError(f.T(), err, "error while waiting for package build to succeed:%w", err) + + _, err = cli.ExecCommand(f.framework, f.ctx, "function", "create", "--name", testDeployFuncName, "--pkg", testDeployPkg, "--entrypoint", "hello.main") + require.NoError(f.T(), err) + + // test with EnvVersion v2 + err = f.fetcherClient.Specialize(f.ctx, &fetcher.FunctionSpecializeRequest{ + FetchReq: fetcher.FunctionFetchRequest{ + Filename: "hi.py", + StorageSvcUrl: f.storagesvcURL, + KeepArchive: true, + FetchType: fv1.FETCH_DEPLOYMENT, + Package: metav1.ObjectMeta{ + Name: testDeployPkg, + Namespace: "default", + }, + ConfigMaps: []fv1.ConfigMapReference{{ + Name: cfgMap.Name, + Namespace: cfgMap.Namespace, + }}, + Secrets: []fv1.SecretReference{{ + Name: secret.Name, + Namespace: secret.Namespace, + }}, + }, + LoadReq: fetcher.FunctionLoadRequest{ + EnvVersion: 2, + FilePath: f.sharedVolDir + "./hi.py", + FunctionName: testDeployFuncName, + FunctionMetadata: &metav1.ObjectMeta{ + Name: testDeployFuncName, + Namespace: metav1.NamespaceDefault, + }, + }, + }) + require.NoError(f.T(), err) + + file, err := os.OpenFile(f.sharedVolDir+"/hi.py", os.O_RDONLY, 0444) + require.NoError(f.T(), err, "fetched file does not exist in fetcher directory") + defer file.Close() + + // test with no EnvVersion + err = f.fetcherClient.Specialize(f.ctx, &fetcher.FunctionSpecializeRequest{ + FetchReq: fetcher.FunctionFetchRequest{ + Filename: "hi.py", + StorageSvcUrl: f.storagesvcURL, + KeepArchive: true, + FetchType: fv1.FETCH_DEPLOYMENT, + Package: metav1.ObjectMeta{ + Name: testDeployPkg, + Namespace: "default", + }, + }, + LoadReq: fetcher.FunctionLoadRequest{ + FilePath: f.sharedVolDir + "./hi.py", + FunctionName: testDeployFuncName, + FunctionMetadata: &metav1.ObjectMeta{ + Name: testDeployFuncName, + Namespace: metav1.NamespaceDefault, + }, + }, + }) + require.NoError(f.T(), err) + + // set throwError to true to test for error case + f.specTestData.throwError = true + + err = f.fetcherClient.Specialize(f.ctx, &fetcher.FunctionSpecializeRequest{ + FetchReq: fetcher.FunctionFetchRequest{ + Filename: "hi.py", + StorageSvcUrl: f.storagesvcURL, + KeepArchive: true, + FetchType: fv1.FETCH_DEPLOYMENT, + Package: metav1.ObjectMeta{ + Name: testDeployPkg, + Namespace: "default", + }, + }, + LoadReq: fetcher.FunctionLoadRequest{ + FilePath: f.sharedVolDir + "./hi.py", + FunctionName: testDeployFuncName, + FunctionMetadata: &metav1.ObjectMeta{ + Name: testDeployFuncName, + Namespace: metav1.NamespaceDefault, + }, + }, + }) + require.Error(f.T(), err) + + _, err = cli.ExecCommand(f.framework, f.ctx, "function", "delete", "--name", testDeployFuncName) + require.NoError(f.T(), err) +} + +func (f *FetcherTestSuite) TearDownSuite() { + _, err := cli.ExecCommand(f.framework, f.ctx, "env", "delete", "--name", f.envName) + require.NoError(f.T(), err) + + err = wait.PollUntilContextTimeout(f.ctx, 2*time.Second, 5*time.Minute, true, func(c context.Context) (done bool, err error) { + + deploys, err := f.k8sClient.AppsV1().Deployments(metav1.NamespaceDefault).List(c, metav1.ListOptions{}) + if err != nil { + return false, err + } + + if len(deploys.Items) == 0 { + return true, nil + } + for _, deploy := range deploys.Items { + if strings.Contains(deploy.Name, f.envName) { + return false, nil + } + } + + return true, nil + }) + require.NoError(f.T(), err) + + f.cancel() + f.logger.Sync() + f.mgr.Wait() + err = f.framework.Stop() + require.NoError(f.T(), err, "error stopping framework") + + err = os.RemoveAll(f.sharedVolDir) + require.NoError(f.T(), err, "error removing fetcher shared vol directory") + + err = os.RemoveAll(f.cfgMapsDir) + require.NoError(f.T(), err, "error removing fetcher cfgMaps directory") + + err = os.RemoveAll(f.secretsDir) + require.NoError(f.T(), err, "error removing fetcher secrets vol directory") +} + +func (f *FetcherTestSuite) createFetcherTestDirs() error { + + cfgMapsDir, err := os.MkdirTemp("/tmp", "fetcher-cfgmaps") + if err != nil { + return fmt.Errorf("error creating temp directory for config maps: %w", err) + } + + f.cfgMapsDir = cfgMapsDir + + secretsDir, err := os.MkdirTemp("/tmp", "fetcher-secrets") + if err != nil { + return fmt.Errorf("error creating temp directory for config maps: %w", err) + } + + f.secretsDir = secretsDir + + shareVolDir, err := os.MkdirTemp("/tmp", "fetcher-shared") + if err != nil { + return fmt.Errorf("error creating temp directory for config maps: %w", err) + } + f.sharedVolDir = shareVolDir + + file, err := os.Create(shareVolDir + "/hello.js") + if err != nil { + return fmt.Errorf("error creating test file in shared volume directory: %w", err) + } + defer file.Close() + + _, err = file.Write([]byte(testFileData)) + if err != nil { + return fmt.Errorf("error writing data to test file in shared volume directory: %w", err) + } + + return nil +} + +func (f *FetcherTestSuite) setupPodInfoMountDir() error { + + podInfoDir, err := os.MkdirTemp("/tmp", "fetcher-pod-info") + if err != nil { + return fmt.Errorf("error creating temp directory for fetcher pod info: %w", err) + } + + f.podInfoDir = podInfoDir + + podNameFile, err := os.Create(podInfoDir + "/name") + if err != nil { + return fmt.Errorf("error creating name file in fetcher pod info directory: %w", err) + } + defer podNameFile.Close() + + _, err = podNameFile.Write([]byte("test")) + if err != nil { + return fmt.Errorf("error writing data in name file in fetcher pod info directory: %w", err) + } + + podNamespaceFile, err := os.Create(podInfoDir + "/namespace") + if err != nil { + return fmt.Errorf("error creating namespace file in fetcher pod info directory: %w", err) + } + defer podNamespaceFile.Close() + + _, err = podNamespaceFile.Write([]byte("test-ns")) + if err != nil { + return fmt.Errorf("error writing data in namespace file in fetcher pod info directory: %w", err) + } + + return nil +} + +func getArchiveIDFromURL(archiveDownloadURL string) (string, error) { + + downloadURL, err := url.Parse(archiveDownloadURL) + if err != nil { + return "", err + } + + querParams, err := url.ParseQuery(downloadURL.RawQuery) + if err != nil { + return "", err + } + + idParamValues := querParams["id"] + if len(idParamValues) == 0 { + return "", errors.New("id param not present in archiveDownloadURL") + } + + return idParamValues[0], nil +} + +func contains(list []string, item string) bool { + + for _, listItem := range list { + if item == listItem { + return true + } + } + + return false +} + +func TestFetcherTestSuite(t *testing.T) { + suite.Run(t, new(FetcherTestSuite)) +} diff --git a/test/e2e/framework/framework.go b/test/e2e/framework/framework.go index 8610040d..13b22a2e 100644 --- a/test/e2e/framework/framework.go +++ b/test/e2e/framework/framework.go @@ -34,7 +34,7 @@ type Framework struct { func NewWebhookOptions() (*envtest.WebhookInstallOptions, error) { webhookPort, err := utils.FindFreePort() if err != nil { - return nil, fmt.Errorf("error finding unused port: %v", err) + return nil, fmt.Errorf("error finding unused port: %w", err) } _, filename, _, _ := runtime.Caller(0) //nolint root := filepath.Dir(filename) @@ -46,7 +46,7 @@ func NewWebhookOptions() (*envtest.WebhookInstallOptions, error) { } err = options.PrepWithoutInstalling() if err != nil { - return nil, fmt.Errorf("error preparing webhook install options: %v", err) + return nil, fmt.Errorf("error preparing webhook install options: %w", err) } return options, nil } @@ -83,7 +83,7 @@ func (f *Framework) Start(ctx context.Context) error { var err error f.config, err = f.env.Start() if err != nil { - return fmt.Errorf("error starting test env: %v", err) + return fmt.Errorf("error starting test env: %w", err) } return nil } @@ -91,7 +91,7 @@ func (f *Framework) Start(ctx context.Context) error { func (f *Framework) ToggleMetricAddr() error { port, err := utils.FindFreePort() if err != nil { - return fmt.Errorf("error finding unused port: %v", err) + return fmt.Errorf("error finding unused port: %w", err) } os.Setenv("METRICS_ADDR", fmt.Sprint(port)) return nil diff --git a/test/e2e/framework/services/services.go b/test/e2e/framework/services/services.go index f201169b..db440937 100644 --- a/test/e2e/framework/services/services.go +++ b/test/e2e/framework/services/services.go @@ -28,7 +28,7 @@ func StartServices(ctx context.Context, f *framework.Framework, mgr manager.Inte webhookPort := env.WebhookInstallOptions.LocalServingPort err := f.ToggleMetricAddr() if err != nil { - return fmt.Errorf("error toggling metric address: %v", err) + return fmt.Errorf("error toggling metric address: %w", err) } mgr.Add(ctx, func(ctx context.Context) { err = webhook.Start(ctx, f.ClientGen(), f.Logger(), cnwebhook.Options{ @@ -43,11 +43,11 @@ func StartServices(ctx context.Context, f *framework.Framework, mgr manager.Inte executorPort, err := utils.FindFreePort() if err != nil { - return fmt.Errorf("error finding unused port: %v", err) + return fmt.Errorf("error finding unused port: %w", err) } err = f.ToggleMetricAddr() if err != nil { - return fmt.Errorf("error toggling metric address: %v", err) + return fmt.Errorf("error toggling metric address: %w", err) } // namespace settings for components @@ -63,42 +63,31 @@ func StartServices(ctx context.Context, f *framework.Framework, mgr manager.Inte os.Setenv("POD_READY_TIMEOUT", "300s") err = executor.StartExecutor(ctx, f.ClientGen(), f.Logger(), mgr, executorPort) if err != nil { - return fmt.Errorf("error starting executor: %v", err) + return fmt.Errorf("error starting executor: %w", err) } f.AddServiceInfo("executor", framework.ServiceInfo{Port: executorPort}) os.Setenv("PRUNE_ENABLED", "true") os.Setenv("PRUNE_INTERVAL", "60") - storageDir, err := os.MkdirTemp("/tmp", "storagesvc") + + err = f.ToggleMetricAddr() if err != nil { - return fmt.Errorf("error creating temp directory: %v", err) + return fmt.Errorf("error toggling metric address: %w", err) } - storageSvcPort, err := utils.FindFreePort() - if err != nil { - return fmt.Errorf("error finding unused port: %v", err) - } err = f.ToggleMetricAddr() if err != nil { - return fmt.Errorf("error toggling metric address: %v", err) + return fmt.Errorf("error toggling metric address: %w", err) } - err = storagesvc.Start(ctx, f.ClientGen(), f.Logger(), storagesvc.NewLocalStorage(storageDir), mgr, storageSvcPort) + + storageSvcPort, err := StartStorageSvc(ctx, f, mgr) if err != nil { - return fmt.Errorf("error starting storage service: %v", err) - } - f.AddServiceInfo("storagesvc", framework.ServiceInfo{Port: storageSvcPort}) - storagesvcURL, err := f.GetServiceURL("storagesvc") - if err != nil { - return fmt.Errorf("error getting storage service URL: %v", err) - } - os.Setenv("FISSION_STORAGESVC_URL", storagesvcURL) - err = f.ToggleMetricAddr() - if err != nil { - return fmt.Errorf("error toggling metric address: %v", err) + return err } + err = buildermgr.Start(ctx, f.ClientGen(), f.Logger(), mgr, fmt.Sprintf("http://localhost:%d", storageSvcPort)) if err != nil { - return fmt.Errorf("error starting builder manager: %v", err) + return fmt.Errorf("error starting builder manager: %w", err) } f.AddServiceInfo("buildermgr", framework.ServiceInfo{}) @@ -115,41 +104,67 @@ func StartServices(ctx context.Context, f *framework.Framework, mgr manager.Inte // os.Setenv("DEBUG_ENV", "false") routerPort, err := utils.FindFreePort() if err != nil { - return fmt.Errorf("error finding unused port: %v", err) + return fmt.Errorf("error finding unused port: %w", err) } err = f.ToggleMetricAddr() if err != nil { - return fmt.Errorf("error toggling metric address: %v", err) + return fmt.Errorf("error toggling metric address: %w", err) } + executor := eclient.MakeClient(f.Logger(), fmt.Sprintf("http://localhost:%d", executorPort)) err = router.Start(ctx, f.ClientGen(), f.Logger(), mgr, routerPort, executor) if err != nil { - return fmt.Errorf("error starting router: %v", err) + return fmt.Errorf("error starting router: %w", err) } f.AddServiceInfo("router", framework.ServiceInfo{Port: routerPort}) routerURL, err := f.GetServiceURL("router") if err != nil { - return fmt.Errorf("error getting router URL: %v", err) + return fmt.Errorf("error getting router URL: %w", err) } os.Setenv("FISSION_ROUTER_URL", routerURL) err = timer.Start(ctx, f.ClientGen(), f.Logger(), mgr, routerURL) if err != nil { - return fmt.Errorf("error starting timer: %v", err) + return fmt.Errorf("error starting timer: %w", err) } f.AddServiceInfo("timer", framework.ServiceInfo{}) err = mqtrigger.StartScalerManager(ctx, f.ClientGen(), f.Logger(), mgr, routerURL) if err != nil { - return fmt.Errorf("error starting mqt scaler manager: %v", err) + return fmt.Errorf("error starting mqt scaler manager: %w", err) } f.AddServiceInfo("mqtrigger-keda", framework.ServiceInfo{}) err = kubewatcher.Start(ctx, f.ClientGen(), f.Logger(), mgr, routerURL) if err != nil { - return fmt.Errorf("error starting kubewatcher: %v", err) + return fmt.Errorf("error starting kubewatcher: %w", err) } f.AddServiceInfo("kubewatcher", framework.ServiceInfo{}) return nil } + +func StartStorageSvc(ctx context.Context, f *framework.Framework, mgr manager.Interface) (storageSvcPort int, err error) { + storageDir, err := os.MkdirTemp("/tmp", "storagesvc") + if err != nil { + return 0, fmt.Errorf("error creating temp directory: %w", err) + } + + storageSvcPort, err = utils.FindFreePort() + if err != nil { + return storageSvcPort, fmt.Errorf("error finding unused port: %w", err) + } + + err = storagesvc.Start(ctx, f.ClientGen(), f.Logger(), storagesvc.NewLocalStorage(storageDir), mgr, storageSvcPort) + if err != nil { + return storageSvcPort, fmt.Errorf("error starting storage service: %w", err) + } + f.AddServiceInfo("storagesvc", framework.ServiceInfo{Port: storageSvcPort}) + storagesvcURL, err := f.GetServiceURL("storagesvc") + if err != nil { + return storageSvcPort, fmt.Errorf("error getting storage service URL: %w", err) + } + os.Setenv("FISSION_STORAGESVC_URL", storagesvcURL) + + return storageSvcPort, nil +}