diff --git a/cmd/fetcher/app/server.go b/cmd/fetcher/app/server.go index 4f005163..04c80a6e 100644 --- a/cmd/fetcher/app/server.go +++ b/cmd/fetcher/app/server.go @@ -30,7 +30,6 @@ import ( "go.uber.org/zap" "github.com/fission/fission/pkg/fetcher" - "github.com/fission/fission/pkg/types" ) func registerTraceExporter(collectorEndpoint string) error { @@ -94,7 +93,7 @@ func Run(logger *zap.Logger) { // do specialization in other goroutine to prevent blocking in newdeploy go func() { if *specializeOnStart { - var specializeReq types.FunctionSpecializeRequest + var specializeReq fetcher.FunctionSpecializeRequest err := json.Unmarshal([]byte(*specializePayload), &specializeReq) if err != nil { diff --git a/cmd/preupgradechecks/preupgradechecks.go b/cmd/preupgradechecks/preupgradechecks.go index 39001612..4fac3912 100644 --- a/cmd/preupgradechecks/preupgradechecks.go +++ b/cmd/preupgradechecks/preupgradechecks.go @@ -29,7 +29,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/crd" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -190,33 +189,33 @@ func (client *PreUpgradeTaskClient) SetupRoleBindings() { // the fact that we're here implies that there had been a prior installation of fission and objects are present still // so, we go ahead and create the role-bindings necessary for the fission-fetcher and fission-builder Service Accounts. - err := utils.SetupRoleBinding(client.logger, client.k8sClient, types.PackageGetterRB, metav1.NamespaceDefault, types.PackageGetterCR, types.ClusterRole, types.FissionFetcherSA, client.fnPodNs) + err := utils.SetupRoleBinding(client.logger, client.k8sClient, fv1.PackageGetterRB, metav1.NamespaceDefault, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, client.fnPodNs) if err != nil { client.logger.Fatal("error setting up rolebinding for service account", zap.Error(err), - zap.String("role_binding", types.PackageGetterRB), - zap.String("service_account", types.FissionFetcherSA), + zap.String("role_binding", fv1.PackageGetterRB), + zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namespace", client.fnPodNs)) } - err = utils.SetupRoleBinding(client.logger, client.k8sClient, types.PackageGetterRB, metav1.NamespaceDefault, types.PackageGetterCR, types.ClusterRole, types.FissionBuilderSA, client.envBuilderNs) + err = utils.SetupRoleBinding(client.logger, client.k8sClient, fv1.PackageGetterRB, metav1.NamespaceDefault, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionBuilderSA, client.envBuilderNs) if err != nil { client.logger.Fatal("error setting up rolebinding for service account", zap.Error(err), - zap.String("role_binding", types.PackageGetterRB), - zap.String("service_account", types.FissionBuilderSA), + zap.String("role_binding", fv1.PackageGetterRB), + zap.String("service_account", fv1.FissionBuilderSA), zap.String("service_account_namespace", client.envBuilderNs)) } - err = utils.SetupRoleBinding(client.logger, client.k8sClient, types.SecretConfigMapGetterRB, metav1.NamespaceDefault, types.SecretConfigMapGetterCR, types.ClusterRole, types.FissionFetcherSA, client.fnPodNs) + err = utils.SetupRoleBinding(client.logger, client.k8sClient, fv1.SecretConfigMapGetterRB, metav1.NamespaceDefault, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, client.fnPodNs) if err != nil { client.logger.Fatal("error setting up rolebinding for service account", zap.Error(err), - zap.String("role_binding", types.SecretConfigMapGetterRB), - zap.String("service_account", types.FissionFetcherSA), + zap.String("role_binding", fv1.SecretConfigMapGetterRB), + zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namespace", client.fnPodNs)) } client.logger.Info("created rolebindings in default namespace", - zap.Strings("role_bindings", []string{types.PackageGetterRB, types.SecretConfigMapGetterRB})) + zap.Strings("role_bindings", []string{fv1.PackageGetterRB, fv1.SecretConfigMapGetterRB})) } diff --git a/pkg/apis/fission.io/v1/const.go b/pkg/apis/fission.io/v1/const.go index b85d77f1..b5706eab 100644 --- a/pkg/apis/fission.io/v1/const.go +++ b/pkg/apis/fission.io/v1/const.go @@ -105,3 +105,42 @@ const ( const ( DefaultSpecializationTimeOut = 120 ) + +const ( + FETCH_SOURCE = iota + FETCH_DEPLOYMENT + FETCH_URL +) + +// executor kubernetes object label key +const ( + ENVIRONMENT_NAMESPACE = "environmentNamespace" + ENVIRONMENT_NAME = "environmentName" + ENVIRONMENT_UID = "environmentUid" + FUNCTION_NAMESPACE = "functionNamespace" + FUNCTION_NAME = "functionName" + FUNCTION_UID = "functionUid" + FUNCTION_RESOURCE_VERSION = "functionResourceVersion" + EXECUTOR_TYPE = "executorType" +) + +const ( + ANNOTATION_SVC_HOST = "svcHost" +) + +const ( + ArchiveLiteralSizeLimit int64 = 256 * 1024 +) + +const ( + FissionBuilderSA = "fission-builder" + FissionFetcherSA = "fission-fetcher" + + SecretConfigMapGetterCR = "secret-configmap-getter" + SecretConfigMapGetterRB = "secret-configmap-getter-binding" + + PackageGetterCR = "package-getter" + PackageGetterRB = "package-getter-binding" + + ClusterRole = "ClusterRole" +) diff --git a/pkg/buildermgr/common.go b/pkg/buildermgr/common.go index 1d405528..2c256d61 100644 --- a/pkg/buildermgr/common.go +++ b/pkg/buildermgr/common.go @@ -19,7 +19,6 @@ package buildermgr import ( "context" "fmt" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "net/http" "strings" "time" @@ -27,14 +26,15 @@ import ( "github.com/dchest/uniuri" "github.com/pkg/errors" "go.uber.org/zap" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/builder" builderClient "github.com/fission/fission/pkg/builder/client" "github.com/fission/fission/pkg/crd" ferror "github.com/fission/fission/pkg/error" + "github.com/fission/fission/pkg/fetcher" fetcherClient "github.com/fission/fission/pkg/fetcher/client" - "github.com/fission/fission/pkg/types" ) // buildPackage helps to build source package into deployment package. @@ -45,7 +45,7 @@ import ( // 4. Return upload response and build logs. // *. Return build logs and error if any one of steps above failed. func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient *crd.FissionClient, envBuilderNamespace string, - storageSvcUrl string, pkg *fv1.Package) (uploadResp *types.ArchiveUploadResponse, buildLogs string, err error) { + storageSvcUrl string, pkg *fv1.Package) (uploadResp *fetcher.ArchiveUploadResponse, buildLogs string, err error) { env, err := fissionClient.V1().Environments(pkg.Spec.Environment.Namespace).Get(pkg.Spec.Environment.Name, metav1.GetOptions{}) if err != nil { @@ -60,8 +60,8 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient *crd.Fi fetcherC := fetcherClient.MakeClient(logger, fmt.Sprintf("http://%v:8000", svcName)) builderC := builderClient.MakeClient(logger, fmt.Sprintf("http://%v:8001", svcName)) - fetchReq := &types.FunctionFetchRequest{ - FetchType: types.FETCH_SOURCE, + fetchReq := &fetcher.FunctionFetchRequest{ + FetchType: fv1.FETCH_SOURCE, Package: pkg.ObjectMeta, Filename: srcPkgFilename, KeepArchive: false, @@ -103,7 +103,7 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient *crd.Fi archivePackage := !env.Spec.KeepArchive - uploadReq := &types.ArchiveUploadRequest{ + uploadReq := &fetcher.ArchiveUploadRequest{ Filename: buildResp.ArtifactFilename, StorageSvcUrl: storageSvcUrl, ArchivePackage: archivePackage, @@ -123,7 +123,7 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient *crd.Fi func updatePackage(logger *zap.Logger, fissionClient *crd.FissionClient, pkg *fv1.Package, status fv1.BuildStatus, buildLogs string, - uploadResp *types.ArchiveUploadResponse) (*fv1.Package, error) { + uploadResp *fetcher.ArchiveUploadResponse) (*fv1.Package, error) { pkg.Status = fv1.PackageStatus{ BuildStatus: status, @@ -133,7 +133,7 @@ func updatePackage(logger *zap.Logger, fissionClient *crd.FissionClient, if uploadResp != nil { pkg.Spec.Deployment = fv1.Archive{ - Type: types.ArchiveTypeUrl, + Type: fv1.ArchiveTypeUrl, URL: uploadResp.ArchiveDownloadUrl, Checksum: uploadResp.Checksum, } diff --git a/pkg/buildermgr/envwatcher.go b/pkg/buildermgr/envwatcher.go index 1953296c..c50c6412 100644 --- a/pkg/buildermgr/envwatcher.go +++ b/pkg/buildermgr/envwatcher.go @@ -36,7 +36,6 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/util" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -346,9 +345,9 @@ func (envw *environmentWatcher) createBuilder(env *fv1.Environment, ns string) ( // there should be only one deploy in deployList if len(deployList) == 0 { // create builder SA in this ns, if not already created - _, err := utils.SetupSA(envw.kubernetesClient, types.FissionBuilderSA, ns) + _, err := utils.SetupSA(envw.kubernetesClient, fv1.FissionBuilderSA, ns) if err != nil { - return nil, errors.Wrapf(err, "error creating %q in ns: %s", types.FissionBuilderSA, ns) + return nil, errors.Wrapf(err, "error creating %q in ns: %s", fv1.FissionBuilderSA, ns) } deploy, err = envw.createBuilderDeployment(env, ns) diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index 12e87b8f..38367a6a 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -32,7 +32,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/cache" "github.com/fission/fission/pkg/crd" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -157,17 +156,17 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package) // Add the package getter rolebinding to builder sa // we continue here if role binding was not setup succeesffully. 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, types.PackageGetterRB, pkg.ObjectMeta.Namespace, types.PackageGetterCR, types.ClusterRole, types.FissionBuilderSA, builderNs) + err := utils.SetupRoleBinding(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), - zap.String("role_binding", types.PackageGetterRB), + zap.String("role_binding", fv1.PackageGetterRB), zap.String("package_name", pkg.ObjectMeta.Name), zap.String("package_namespace", pkg.ObjectMeta.Namespace)) continue } else { pkgw.logger.Info("setup rolebinding for sa package", - zap.String("sa", fmt.Sprintf("%s.%s", types.FissionBuilderSA, builderNs)), + zap.String("sa", fmt.Sprintf("%s.%s", fv1.FissionBuilderSA, builderNs)), zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace))) } @@ -175,7 +174,7 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package) 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)) - updatePackage(pkgw.logger, pkgw.fissionClient, pkg, types.BuildStatusFailed, buildLogs, nil) + updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) return } @@ -210,10 +209,10 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package) } _, err = updatePackage(pkgw.logger, pkgw.fissionClient, pkg, - types.BuildStatusSucceeded, buildLogs, uploadResp) + fv1.BuildStatusSucceeded, buildLogs, uploadResp) if err != nil { pkgw.logger.Error("error updating package info", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name)) - updatePackage(pkgw.logger, pkgw.fissionClient, pkg, types.BuildStatusFailed, buildLogs, nil) + updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) return } @@ -223,7 +222,7 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package) } // build timeout updatePackage(pkgw.logger, pkgw.fissionClient, pkg, - types.BuildStatusFailed, "Build timeout due to environment builder not ready", nil) + fv1.BuildStatusFailed, "Build timeout due to environment builder not ready", nil) pkgw.logger.Error("max retries exceeded in building source package, timeout due to environment builder not ready", zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace))) diff --git a/pkg/canaryconfigmgr/canaryConfigMgr.go b/pkg/canaryconfigmgr/canaryConfigMgr.go index ab43c44f..474ac5ce 100644 --- a/pkg/canaryconfigmgr/canaryConfigMgr.go +++ b/pkg/canaryconfigmgr/canaryConfigMgr.go @@ -35,7 +35,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/crd" - "github.com/fission/fission/pkg/types" ) const ( @@ -112,7 +111,7 @@ func (canaryCfgMgr *canaryConfigMgr) initCanaryConfigController() (k8sCache.Stor k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { canaryConfig := obj.(*fv1.CanaryConfig) - if canaryConfig.Status.Status == types.CanaryConfigStatusPending { + if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { go canaryCfgMgr.addCanaryConfig(canaryConfig) } }, @@ -124,7 +123,7 @@ func (canaryCfgMgr *canaryConfigMgr) initCanaryConfigController() (k8sCache.Stor oldConfig := oldObj.(*fv1.CanaryConfig) newConfig := newObj.(*fv1.CanaryConfig) if oldConfig.ObjectMeta.ResourceVersion != newConfig.ObjectMeta.ResourceVersion && - newConfig.Status.Status == types.CanaryConfigStatusPending { + newConfig.Status.Status == fv1.CanaryConfigStatusPending { canaryCfgMgr.logger.Info("update canary config invoked", zap.String("name", newConfig.ObjectMeta.Name), zap.String("namespace", newConfig.ObjectMeta.Namespace), @@ -267,7 +266,7 @@ func (canaryCfgMgr *canaryConfigMgr) RollForwardOrBack(canaryConfig *fv1.CanaryC } // handle a race between ticker.Stop and receiving a notification on ticker.C - if canaryConfig.Status.Status != types.CanaryConfigStatusPending { + if canaryConfig.Status.Status != fv1.CanaryConfigStatusPending { canaryCfgMgr.logger.Info("no need of processing the config, not pending anymore", zap.String("name", canaryConfig.ObjectMeta.Name), zap.String("namespace", canaryConfig.ObjectMeta.Namespace), @@ -275,7 +274,7 @@ func (canaryCfgMgr *canaryConfigMgr) RollForwardOrBack(canaryConfig *fv1.CanaryC return } - if triggerObj.Spec.FunctionReference.Type == types.FunctionReferenceTypeFunctionWeights && + if triggerObj.Spec.FunctionReference.Type == fv1.FunctionReferenceTypeFunctionWeights && triggerObj.Spec.FunctionReference.FunctionWeights[canaryConfig.Spec.NewFunction] != 0 { failurePercent, err := canaryCfgMgr.promClient.GetFunctionFailurePercentage(triggerObj.Spec.RelativeURL, triggerObj.Spec.Method, canaryConfig.Spec.NewFunction, canaryConfig.ObjectMeta.Namespace, canaryConfig.Spec.WeightIncrementDuration) @@ -341,7 +340,7 @@ func (canaryCfgMgr *canaryConfigMgr) RollForwardOrBack(canaryConfig *fv1.CanaryC // update the status of canary config as done processing, we dont care if we arent able to update because // resync takes care of the update err = canaryCfgMgr.updateCanaryConfigStatusWithRetries(canaryConfig.ObjectMeta.Name, canaryConfig.ObjectMeta.Namespace, - types.CanaryConfigStatusSucceeded) + fv1.CanaryConfigStatusSucceeded) if err != nil { // cant do much after max retries other than logging it. canaryCfgMgr.logger.Error("error updating canary config after max retries", @@ -452,7 +451,7 @@ func (canaryCfgMgr *canaryConfigMgr) rollback(canaryConfig *fv1.CanaryConfig, tr } err = canaryCfgMgr.updateCanaryConfigStatusWithRetries(canaryConfig.ObjectMeta.Name, canaryConfig.ObjectMeta.Namespace, - types.CanaryConfigStatusFailed) + fv1.CanaryConfigStatusFailed) return err } @@ -487,7 +486,7 @@ func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs() { for _, obj := range canaryCfgMgr.canaryConfigStore.List() { canaryConfig := obj.(*fv1.CanaryConfig) _, err := canaryCfgMgr.canaryCfgCancelFuncMap.lookup(&canaryConfig.ObjectMeta) - if err != nil && canaryConfig.Status.Status == types.CanaryConfigStatusPending { + if err != nil && canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { canaryCfgMgr.logger.Debug("adding canary config from resync loop", zap.String("name", canaryConfig.ObjectMeta.Name), zap.String("namespace", canaryConfig.ObjectMeta.Namespace), diff --git a/pkg/controller/functionApi.go b/pkg/controller/functionApi.go index 09d96289..e384b9d1 100644 --- a/pkg/controller/functionApi.go +++ b/pkg/controller/functionApi.go @@ -41,7 +41,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" ferror "github.com/fission/fission/pkg/error" - "github.com/fission/fission/pkg/types" ) func RegisterFunctionRoute(ws *restful.WebService) { @@ -307,9 +306,9 @@ func (a *API) FunctionPodLogs(w http.ResponseWriter, r *http.Request) { // Get function Pods first selector := map[string]string{ - types.FUNCTION_UID: string(f.ObjectMeta.UID), - types.ENVIRONMENT_NAME: f.Spec.Environment.Name, - types.ENVIRONMENT_NAMESPACE: f.Spec.Environment.Namespace, + fv1.FUNCTION_UID: string(f.ObjectMeta.UID), + fv1.ENVIRONMENT_NAME: f.Spec.Environment.Name, + fv1.ENVIRONMENT_NAMESPACE: f.Spec.Environment.Namespace, } podList, err := a.kubernetesClient.CoreV1().Pods(podNs).List(metav1.ListOptions{ LabelSelector: labels.Set(selector).AsSelector().String(), diff --git a/pkg/controller/packageApi.go b/pkg/controller/packageApi.go index 979bb8ce..1f175552 100644 --- a/pkg/controller/packageApi.go +++ b/pkg/controller/packageApi.go @@ -25,7 +25,6 @@ import ( "github.com/dustin/go-humanize" "github.com/emicklei/go-restful" restfulspec "github.com/emicklei/go-restful-openapi" - "github.com/fission/fission/pkg/types" "github.com/go-openapi/spec" "github.com/gorilla/mux" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -136,15 +135,15 @@ func (a *API) PackageApiCreate(w http.ResponseWriter, r *http.Request) { } // Ensure size limits - if len(f.Spec.Source.Literal) > int(types.ArchiveLiteralSizeLimit) { + if len(f.Spec.Source.Literal) > int(fv1.ArchiveLiteralSizeLimit) { err := ferror.MakeError(ferror.ErrorInvalidArgument, - fmt.Sprintf("Package literal larger than %s", humanize.Bytes(uint64(types.ArchiveLiteralSizeLimit)))) + fmt.Sprintf("Package literal larger than %s", humanize.Bytes(uint64(fv1.ArchiveLiteralSizeLimit)))) a.respondWithError(w, err) return } - if len(f.Spec.Deployment.Literal) > int(types.ArchiveLiteralSizeLimit) { + if len(f.Spec.Deployment.Literal) > int(fv1.ArchiveLiteralSizeLimit) { err := ferror.MakeError(ferror.ErrorInvalidArgument, - fmt.Sprintf("Package literal larger than %s", humanize.Bytes(uint64(types.ArchiveLiteralSizeLimit)))) + fmt.Sprintf("Package literal larger than %s", humanize.Bytes(uint64(fv1.ArchiveLiteralSizeLimit)))) a.respondWithError(w, err) return } diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index 641869a9..6df5629a 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -35,7 +35,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/executor/util" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -127,7 +126,7 @@ func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) if err != nil { deploy.logger.Error("error creating fission fetcher service account for function", zap.Error(err), - zap.String("service_account_name", types.FissionFetcherSA), + zap.String("service_account_name", fv1.FissionFetcherSA), zap.String("service_account_namespace", deployNamespace), zap.String("function_name", fn.ObjectMeta.Name), zap.String("function_namespace", fn.ObjectMeta.Namespace)) @@ -135,22 +134,22 @@ 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, types.PackageGetterRB, fn.Spec.Package.PackageRef.Namespace, types.PackageGetterCR, types.ClusterRole, types.FissionFetcherSA, deployNamespace) + err = utils.SetupRoleBinding(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), - zap.String("role_binding", types.PackageGetterRB), + zap.String("role_binding", fv1.PackageGetterRB), zap.String("function_name", fn.ObjectMeta.Name), zap.String("function_namespace", fn.ObjectMeta.Namespace)) return err } // create rolebinding in function namespace for fetcherSA.envNamespace to be able to get secrets and configmaps - err = utils.SetupRoleBinding(deploy.logger, deploy.kubernetesClient, types.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, types.SecretConfigMapGetterCR, types.ClusterRole, types.FissionFetcherSA, deployNamespace) + err = utils.SetupRoleBinding(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), - zap.String("role_binding", types.SecretConfigMapGetterRB), + zap.String("role_binding", fv1.SecretConfigMapGetterRB), zap.String("function_name", fn.ObjectMeta.Name), zap.String("function_namespace", fn.ObjectMeta.Namespace)) return err diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index fb8ca389..4797a5d4 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -46,7 +46,6 @@ import ( "github.com/fission/fission/pkg/executor/reaper" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/throttler" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -278,7 +277,7 @@ func (deploy *NewDeploy) CleanupOldExecutorObjects() { errs := &multierror.Error{} listOpts := metav1.ListOptions{ - LabelSelector: labels.Set(map[string]string{types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy)}).AsSelector().String(), + LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy)}).AsSelector().String(), } err := reaper.CleanupHpa(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts) @@ -744,20 +743,20 @@ func (deploy *NewDeploy) getObjName(fn *fv1.Function) string { func (deploy *NewDeploy) getDeployLabels(fnMeta metav1.ObjectMeta, envMeta metav1.ObjectMeta) map[string]string { return map[string]string{ - types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy), - types.ENVIRONMENT_NAME: envMeta.Name, - types.ENVIRONMENT_NAMESPACE: envMeta.Namespace, - types.ENVIRONMENT_UID: string(envMeta.UID), - types.FUNCTION_NAME: fnMeta.Name, - types.FUNCTION_NAMESPACE: fnMeta.Namespace, - types.FUNCTION_UID: string(fnMeta.UID), + fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy), + fv1.ENVIRONMENT_NAME: envMeta.Name, + fv1.ENVIRONMENT_NAMESPACE: envMeta.Namespace, + fv1.ENVIRONMENT_UID: string(envMeta.UID), + fv1.FUNCTION_NAME: fnMeta.Name, + fv1.FUNCTION_NAMESPACE: fnMeta.Namespace, + fv1.FUNCTION_UID: string(fnMeta.UID), } } func (deploy *NewDeploy) getDeployAnnotations(fnMeta metav1.ObjectMeta) map[string]string { return map[string]string{ - types.EXECUTOR_INSTANCEID_LABEL: deploy.instanceID, - types.FUNCTION_RESOURCE_VERSION: fnMeta.ResourceVersion, + fv1.EXECUTOR_INSTANCEID_LABEL: deploy.instanceID, + fv1.FUNCTION_RESOURCE_VERSION: fnMeta.ResourceVersion, } } diff --git a/pkg/executor/executortype/poolmgr/funcwatcher.go b/pkg/executor/executortype/poolmgr/funcwatcher.go index 007c0f97..548a63f2 100644 --- a/pkg/executor/executortype/poolmgr/funcwatcher.go +++ b/pkg/executor/executortype/poolmgr/funcwatcher.go @@ -19,7 +19,6 @@ package poolmgr import ( "time" - "github.com/fission/fission/pkg/types" "go.uber.org/zap" apiv1 "k8s.io/api/core/v1" kerrors "k8s.io/apimachinery/pkg/api/errors" @@ -74,12 +73,12 @@ func (gpm *GenericPoolManager) makeFuncController(fissionClient *crd.FissionClie // setup rolebinding is tried, if it fails, we dont 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(gpm.logger, kubernetesClient, types.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, types.SecretConfigMapGetterCR, types.ClusterRole, types.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { - gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", types.SecretConfigMapGetterRB)) + gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) } else { gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", - zap.String("service_account", types.FissionFetcherSA), + zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namepsace", envNs), zap.String("function_name", fn.ObjectMeta.Name), zap.String("function_namespace", fn.ObjectMeta.Namespace)) @@ -187,15 +186,15 @@ func (gpm *GenericPoolManager) makeFuncController(fissionClient *crd.FissionClie envNs = newFunc.Spec.Environment.Namespace } - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, types.SecretConfigMapGetterRB, - newFunc.ObjectMeta.Namespace, types.SecretConfigMapGetterCR, types.ClusterRole, - types.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, + newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, + fv1.FissionFetcherSA, envNs) if err != nil { - gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", types.SecretConfigMapGetterRB)) + gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) } else { gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", - zap.String("service_account", types.FissionFetcherSA), + zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namepsace", envNs), zap.String("function_name", newFunc.ObjectMeta.Name), zap.String("function_namespace", newFunc.ObjectMeta.Namespace)) diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index a5c512a1..968e6f06 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -26,7 +26,6 @@ import ( "time" "github.com/dchest/uniuri" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" multierror "github.com/hashicorp/go-multierror" "github.com/pkg/errors" @@ -148,11 +147,11 @@ func MakeGenericPool( func (gp *GenericPool) getEnvironmentPoolLabels() map[string]string { return map[string]string{ - types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), - types.ENVIRONMENT_NAME: gp.env.ObjectMeta.Name, - types.ENVIRONMENT_NAMESPACE: gp.env.ObjectMeta.Namespace, - types.ENVIRONMENT_UID: string(gp.env.ObjectMeta.UID), - "managed": "true", // this allows us to easily find pods managed by the deployment + fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), + fv1.ENVIRONMENT_NAME: gp.env.ObjectMeta.Name, + fv1.ENVIRONMENT_NAMESPACE: gp.env.ObjectMeta.Namespace, + fv1.ENVIRONMENT_UID: string(gp.env.ObjectMeta.UID), + "managed": "true", // this allows us to easily find pods managed by the deployment } } @@ -241,7 +240,7 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*apiv1.Pod, erro // and make a good scheduling decision. chosenPod := readyPods[rand.Intn(len(readyPods))] - if gp.env.Spec.AllowedFunctionsPerContainer != types.AllowedFunctionsPerContainerInfinite { + if gp.env.Spec.AllowedFunctionsPerContainer != fv1.AllowedFunctionsPerContainerInfinite { // Relabel. If the pod already got picked and // modified, this should fail; in that case just // retry. @@ -265,10 +264,10 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*apiv1.Pod, erro func (gp *GenericPool) labelsForFunction(metadata *metav1.ObjectMeta) map[string]string { label := gp.getEnvironmentPoolLabels() - label[types.FUNCTION_NAME] = metadata.Name - label[types.FUNCTION_UID] = string(metadata.UID) - label[types.FUNCTION_NAMESPACE] = metadata.Namespace // function CRD must stay within same namespace of environment CRD - label["managed"] = "false" // this allows us to easily find pods not managed by the deployment + label[fv1.FUNCTION_NAME] = metadata.Name + label[fv1.FUNCTION_UID] = string(metadata.UID) + label[fv1.FUNCTION_NAMESPACE] = metadata.Namespace // function CRD must stay within same namespace of environment CRD + label["managed"] = "false" // this allows us to easily find pods not managed by the deployment return label } @@ -634,7 +633,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac // patch svc-host and resource version to the pod annotations for new executor to adopt the pod patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v","%v":"%v"}}}`, - types.ANNOTATION_SVC_HOST, svcHost, types.FUNCTION_RESOURCE_VERSION, fn.ObjectMeta.ResourceVersion) + fv1.ANNOTATION_SVC_HOST, svcHost, fv1.FUNCTION_RESOURCE_VERSION, fn.ObjectMeta.ResourceVersion) p, err := gp.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch)) if err != nil { // just log the error since it won't affect the function serving diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index d21cfff1..81023bb0 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -43,7 +43,6 @@ import ( "github.com/fission/fission/pkg/executor/fscache" "github.com/fission/fission/pkg/executor/reaper" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -278,7 +277,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources() { } l := map[string]string{ - types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), + fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), } podList, err := gpm.kubernetesClient.CoreV1().Pods(metav1.NamespaceAll).List(metav1.ListOptions{ @@ -303,7 +302,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources() { // avoid too many requests arrive Kubernetes API server at the same time. time.Sleep(time.Duration(rand.Intn(30)) * time.Millisecond) - patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, gpm.instanceId) + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, fv1.EXECUTOR_INSTANCEID_LABEL, gpm.instanceId) pod, err = gpm.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch)) if err != nil { // just log the error since it won't affect the function serving @@ -317,13 +316,13 @@ func (gpm *GenericPoolManager) AdoptExistingResources() { return } - fnName, ok1 := pod.Labels[types.FUNCTION_NAME] - fnNS, ok2 := pod.Labels[types.FUNCTION_NAMESPACE] - fnUID, ok3 := pod.Labels[types.FUNCTION_UID] - fnRV, ok4 := pod.Annotations[types.FUNCTION_RESOURCE_VERSION] - envName, ok5 := pod.Labels[types.ENVIRONMENT_NAME] - envNS, ok6 := pod.Labels[types.ENVIRONMENT_NAMESPACE] - svcHost, ok7 := pod.Annotations[types.ANNOTATION_SVC_HOST] + fnName, ok1 := pod.Labels[fv1.FUNCTION_NAME] + fnNS, ok2 := pod.Labels[fv1.FUNCTION_NAMESPACE] + fnUID, ok3 := pod.Labels[fv1.FUNCTION_UID] + fnRV, ok4 := pod.Annotations[fv1.FUNCTION_RESOURCE_VERSION] + envName, ok5 := pod.Labels[fv1.ENVIRONMENT_NAME] + envNS, ok6 := pod.Labels[fv1.ENVIRONMENT_NAMESPACE] + svcHost, ok7 := pod.Annotations[fv1.ANNOTATION_SVC_HOST] env, ok8 := envMap[fmt.Sprintf("%v/%v", envNS, envName)] if !(ok1 && ok2 && ok3 && ok4 && ok5 && ok6 && ok7 && ok8) { @@ -382,7 +381,7 @@ func (gpm *GenericPoolManager) CleanupOldExecutorObjects() { errs := &multierror.Error{} listOpts := metav1.ListOptions{ - LabelSelector: labels.Set(map[string]string{types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr)}).AsSelector().String(), + LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr)}).AsSelector().String(), } err := reaper.CleanupDeployments(gpm.logger, gpm.kubernetesClient, gpm.instanceId, listOpts) @@ -412,7 +411,7 @@ func (gpm *GenericPoolManager) service() { if !ok { poolsize := gpm.getEnvPoolsize(req.env) switch req.env.Spec.AllowedFunctionsPerContainer { - case types.AllowedFunctionsPerContainerInfinite: + case fv1.AllowedFunctionsPerContainerInfinite: poolsize = 1 } @@ -589,7 +588,7 @@ func (gpm *GenericPoolManager) idleObjectReaper() { zap.String("function", fsvc.Name)) } - if fsvc.Environment.Spec.AllowedFunctionsPerContainer == types.AllowedFunctionsPerContainerInfinite { + if fsvc.Environment.Spec.AllowedFunctionsPerContainer == fv1.AllowedFunctionsPerContainerInfinite { continue } diff --git a/pkg/executor/executortype/poolmgr/packagewatcher.go b/pkg/executor/executortype/poolmgr/packagewatcher.go index db0ab201..15c9bb12 100644 --- a/pkg/executor/executortype/poolmgr/packagewatcher.go +++ b/pkg/executor/executortype/poolmgr/packagewatcher.go @@ -19,7 +19,6 @@ package poolmgr import ( "time" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -53,18 +52,18 @@ func (gpm *GenericPoolManager) makePkgController(fissionClient *crd.FissionClien // 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(gpm.logger, kubernetesClient, types.PackageGetterRB, pkg.ObjectMeta.Namespace, types.PackageGetterCR, types.ClusterRole, types.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { gpm.logger.Error("error creating rolebinding for package", zap.Error(err), - zap.String("role_binding", types.PackageGetterRB), + zap.String("role_binding", fv1.PackageGetterRB), zap.String("package_name", pkg.ObjectMeta.Name), zap.String("package_namespace", pkg.ObjectMeta.Namespace)) return } gpm.logger.Debug("successfully set up rolebinding for fetcher service account", - zap.String("service_account", types.FissionFetcherSA), + zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namespace", envNs), zap.String("package_name", pkg.ObjectMeta.Name), zap.String("package_namespace", pkg.ObjectMeta.Namespace)) @@ -87,20 +86,20 @@ func (gpm *GenericPoolManager) makePkgController(fissionClient *crd.FissionClien envNs = newPkg.Spec.Environment.Namespace } - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, types.PackageGetterRB, - newPkg.ObjectMeta.Namespace, types.PackageGetterCR, types.ClusterRole, - types.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, + newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, + fv1.FissionFetcherSA, envNs) if err != nil { gpm.logger.Error("error updating rolebinding for package", zap.Error(err), - zap.String("role_binding", types.PackageGetterRB), + zap.String("role_binding", fv1.PackageGetterRB), zap.String("package_name", newPkg.ObjectMeta.Name), zap.String("package_namespace", newPkg.ObjectMeta.Namespace)) return } gpm.logger.Debug("successfully updated rolebinding for fetcher service account", - zap.String("service_account", types.FissionFetcherSA), + zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namespace", envNs), zap.String("package_name", newPkg.ObjectMeta.Name), zap.String("package_namespace", newPkg.ObjectMeta.Namespace)) diff --git a/pkg/executor/reaper/reaper.go b/pkg/executor/reaper/reaper.go index 51e8d7ff..eedac454 100644 --- a/pkg/executor/reaper/reaper.go +++ b/pkg/executor/reaper/reaper.go @@ -27,7 +27,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/crd" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -75,10 +74,10 @@ func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan return err } for _, dep := range deploymentList.Items { - id, ok := dep.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := dep.ObjectMeta.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] if !ok { // Backward compatibility with older label name - id, ok = dep.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok = dep.ObjectMeta.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { logger.Info("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name)) @@ -101,10 +100,10 @@ func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st return err } for _, pod := range podList.Items { - id, ok := pod.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := pod.ObjectMeta.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] if !ok { // Backward compatibility with older label name - id, ok = pod.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok = pod.ObjectMeta.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { logger.Info("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name)) @@ -127,10 +126,10 @@ func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI return err } for _, svc := range svcList.Items { - id, ok := svc.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := svc.ObjectMeta.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] if !ok { // Backward compatibility with older label name - id, ok = svc.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok = svc.ObjectMeta.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { logger.Info("cleaning up service", zap.String("service", svc.ObjectMeta.Name)) @@ -154,10 +153,10 @@ func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str } for _, hpa := range hpaList.Items { - id, ok := hpa.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := hpa.ObjectMeta.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] if !ok { // Backward compatibility with older label name - id, ok = hpa.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok = hpa.ObjectMeta.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { logger.Info("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name)) @@ -199,7 +198,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi } // ignore role-bindings not created by fission - if roleBinding.Name != types.PackageGetterRB && roleBinding.Name != types.SecretConfigMapGetterRB { + if roleBinding.Name != fv1.PackageGetterRB && roleBinding.Name != fv1.SecretConfigMapGetterRB { continue } @@ -250,7 +249,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi // if its a package-getter-rb, we have 2 kinds of SAs and each of them is handled differently // else if its a secret-configmap-rb, we have only one SA which is fission-fetcher - if roleBinding.Name == types.PackageGetterRB { + if roleBinding.Name == fv1.PackageGetterRB { // check if there is an env obj in saNs envList, err := fissionClient.V1().Environments(saNs).List(meta_v1.ListOptions{}) if err != nil { @@ -263,7 +262,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi // if there's at least one function in the role-binding namespace with env reference // to the SA's namespace. // if neither, then we can remove this SA from this role-binding - if subj.Name == types.FissionBuilderSA { + if subj.Name == fv1.FissionBuilderSA { if len(envList.Items) == 0 && !funcEnvReference { saToRemove[utils.MakeSAMapKey(subj.Name, subj.Namespace)] = true } @@ -273,13 +272,13 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi // we also need to check if there's at least one function with executor type New deploy // in the rolebinding's namespace. // if none of them are true, then remove this SA from this role-binding - if subj.Name == types.FissionFetcherSA { + if subj.Name == fv1.FissionFetcherSA { if len(envList.Items) == 0 && !ndmFunc && !funcEnvReference { // remove SA from rolebinding saToRemove[utils.MakeSAMapKey(subj.Name, subj.Namespace)] = true } } - } else if roleBinding.Name == types.SecretConfigMapGetterRB { + } else if roleBinding.Name == fv1.SecretConfigMapGetterRB { // if there's not even one function in the role-binding's namespace and there's not even // one function with env reference to the SA's namespace, then remove that SA // from this role-binding diff --git a/pkg/fetcher/client/client.go b/pkg/fetcher/client/client.go index fa8457c0..ca6551f9 100644 --- a/pkg/fetcher/client/client.go +++ b/pkg/fetcher/client/client.go @@ -15,7 +15,7 @@ import ( "golang.org/x/net/context/ctxhttp" ferror "github.com/fission/fission/pkg/error" - "github.com/fission/fission/pkg/types" + "github.com/fission/fission/pkg/fetcher" ) type ( @@ -48,23 +48,23 @@ func (c *Client) getUploadUrl() string { return c.url + "/upload" } -func (c *Client) Specialize(ctx context.Context, req *types.FunctionSpecializeRequest) error { +func (c *Client) Specialize(ctx context.Context, req *fetcher.FunctionSpecializeRequest) error { _, err := sendRequest(c.logger, ctx, c.httpClient, req, c.getSpecializeUrl()) return err } -func (c *Client) Fetch(ctx context.Context, fr *types.FunctionFetchRequest) error { +func (c *Client) Fetch(ctx context.Context, fr *fetcher.FunctionFetchRequest) error { _, err := sendRequest(c.logger, ctx, c.httpClient, fr, c.getFetchUrl()) return err } -func (c *Client) Upload(ctx context.Context, fr *types.ArchiveUploadRequest) (*types.ArchiveUploadResponse, error) { +func (c *Client) Upload(ctx context.Context, fr *fetcher.ArchiveUploadRequest) (*fetcher.ArchiveUploadResponse, error) { body, err := sendRequest(c.logger, ctx, c.httpClient, fr, c.getUploadUrl()) if err != nil { return nil, err } - uploadResp := types.ArchiveUploadResponse{} + uploadResp := fetcher.ArchiveUploadResponse{} err = json.Unmarshal(body, &uploadResp) if err != nil { return nil, err diff --git a/pkg/fetcher/config/config.go b/pkg/fetcher/config/config.go index 01360b91..d6e02fcf 100644 --- a/pkg/fetcher/config/config.go +++ b/pkg/fetcher/config/config.go @@ -16,7 +16,7 @@ import ( "k8s.io/client-go/kubernetes" fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" - "github.com/fission/fission/pkg/types" + "github.com/fission/fission/pkg/fetcher" "github.com/fission/fission/pkg/utils" ) @@ -88,14 +88,14 @@ func MakeFetcherConfig(sharedMountPath string) (*Config, error) { sharedSecretPath: "/secrets", sharedCfgMapPath: "/configs", jaegerCollectorEndpoint: os.Getenv("TRACE_JAEGER_COLLECTOR_ENDPOINT"), - serviceAccount: types.FissionFetcherSA, + serviceAccount: fv1.FissionFetcherSA, }, nil } func (cfg *Config) SetupServiceAccount(kubernetesClient *kubernetes.Clientset, namespace string, context interface{}) error { - _, err := utils.SetupSA(kubernetesClient, types.FissionFetcherSA, namespace) + _, err := utils.SetupSA(kubernetesClient, fv1.FissionFetcherSA, namespace) if err != nil { - log.Printf("Error : %v creating %s in ns : %s for: %#v", err, types.FissionFetcherSA, namespace, context) + log.Printf("Error : %v creating %s in ns : %s for: %#v", err, fv1.FissionFetcherSA, namespace, context) return err } @@ -106,7 +106,7 @@ func (cfg *Config) SharedMountPath() string { return cfg.sharedMountPath } -func (cfg *Config) NewSpecializeRequest(fn *fv1.Function, env *fv1.Environment) types.FunctionSpecializeRequest { +func (cfg *Config) NewSpecializeRequest(fn *fv1.Function, env *fv1.Environment) fetcher.FunctionSpecializeRequest { // for backward compatibility, since most v1 env // still try to load user function from hard coded // path /userfunc/user @@ -115,9 +115,9 @@ func (cfg *Config) NewSpecializeRequest(fn *fv1.Function, env *fv1.Environment) targetFilename = string(fn.ObjectMeta.UID) } - return types.FunctionSpecializeRequest{ - FetchReq: types.FunctionFetchRequest{ - FetchType: types.FETCH_DEPLOYMENT, + return fetcher.FunctionSpecializeRequest{ + FetchReq: fetcher.FunctionFetchRequest{ + FetchType: fv1.FETCH_DEPLOYMENT, Package: metav1.ObjectMeta{ Namespace: fn.Spec.Package.PackageRef.Namespace, Name: fn.Spec.Package.PackageRef.Name, @@ -127,7 +127,7 @@ func (cfg *Config) NewSpecializeRequest(fn *fv1.Function, env *fv1.Environment) ConfigMaps: fn.Spec.ConfigMaps, KeepArchive: env.Spec.KeepArchive, }, - LoadReq: types.FunctionLoadRequest{ + LoadReq: fetcher.FunctionLoadRequest{ FilePath: filepath.Join(cfg.sharedMountPath, targetFilename), FunctionName: fn.Spec.Package.FunctionName, FunctionMetadata: &fn.ObjectMeta, @@ -172,19 +172,19 @@ func (cfg *Config) fetcherCommand(extraArgs ...string) []string { func (cfg *Config) volumesWithMounts() ([]apiv1.Volume, []apiv1.VolumeMount) { volumes := []apiv1.Volume{ { - Name: types.SharedVolumeUserfunc, + Name: fv1.SharedVolumeUserfunc, VolumeSource: apiv1.VolumeSource{ EmptyDir: &apiv1.EmptyDirVolumeSource{}, }, }, { - Name: types.SharedVolumeSecrets, + Name: fv1.SharedVolumeSecrets, VolumeSource: apiv1.VolumeSource{ EmptyDir: &apiv1.EmptyDirVolumeSource{}, }, }, { - Name: types.SharedVolumeConfigmaps, + Name: fv1.SharedVolumeConfigmaps, VolumeSource: apiv1.VolumeSource{ EmptyDir: &apiv1.EmptyDirVolumeSource{}, }, @@ -192,15 +192,15 @@ func (cfg *Config) volumesWithMounts() ([]apiv1.Volume, []apiv1.VolumeMount) { } mounts := []apiv1.VolumeMount{ { - Name: types.SharedVolumeUserfunc, + Name: fv1.SharedVolumeUserfunc, MountPath: cfg.sharedMountPath, }, { - Name: types.SharedVolumeSecrets, + Name: fv1.SharedVolumeSecrets, MountPath: cfg.sharedSecretPath, }, { - Name: types.SharedVolumeConfigmaps, + Name: fv1.SharedVolumeConfigmaps, MountPath: cfg.sharedCfgMapPath, }, } @@ -289,7 +289,7 @@ func (cfg *Config) addFetcherToPodSpecWithCommand(podSpec *apiv1.PodSpec, mainCo podSpec.Volumes = append(podSpec.Volumes, volumes...) podSpec.Containers = append(podSpec.Containers, c) if podSpec.ServiceAccountName == "" { - podSpec.ServiceAccountName = types.FissionFetcherSA + podSpec.ServiceAccountName = fv1.FissionFetcherSA } return nil diff --git a/pkg/fetcher/fetcher.go b/pkg/fetcher/fetcher.go index d1049356..b97ed4a2 100644 --- a/pkg/fetcher/fetcher.go +++ b/pkg/fetcher/fetcher.go @@ -42,7 +42,6 @@ import ( "github.com/fission/fission/pkg/error/network" "github.com/fission/fission/pkg/info" storageSvcClient "github.com/fission/fission/pkg/storagesvc/client" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -139,7 +138,7 @@ func (fetcher *Fetcher) FetchHandler(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusInternalServerError) return } - var req types.FunctionFetchRequest + var req FunctionFetchRequest err = json.Unmarshal(body, &req) if err != nil { fetcher.logger.Error("error parsing request body", zap.Error(err)) @@ -187,7 +186,7 @@ func (fetcher *Fetcher) SpecializeHandler(w http.ResponseWriter, r *http.Request http.Error(w, err.Error(), http.StatusInternalServerError) return } - var req types.FunctionSpecializeRequest + var req FunctionSpecializeRequest err = json.Unmarshal(body, &req) if err != nil { fetcher.logger.Error("error parsing request body", zap.Error(err)) @@ -208,7 +207,7 @@ func (fetcher *Fetcher) SpecializeHandler(w http.ResponseWriter, r *http.Request // Fetch takes FetchRequest and makes the fetch call // It returns the HTTP code and error if any -func (fetcher *Fetcher) Fetch(ctx context.Context, pkg *fv1.Package, req types.FunctionFetchRequest) (int, error) { +func (fetcher *Fetcher) Fetch(ctx context.Context, pkg *fv1.Package, req FunctionFetchRequest) (int, error) { // check that the requested filename is not an empty string and error out if so if len(req.Filename) == 0 { e := "fetch request received for an empty file name" @@ -227,7 +226,7 @@ func (fetcher *Fetcher) Fetch(ctx context.Context, pkg *fv1.Package, req types.F tmpFile := req.Filename + ".tmp" tmpPath := filepath.Join(fetcher.sharedVolumePath, tmpFile) - if req.FetchType == types.FETCH_URL { + if req.FetchType == fv1.FETCH_URL { // fetch the file and save it to the tmp path err := utils.DownloadUrl(ctx, fetcher.httpClient, req.Url, tmpPath) if err != nil { @@ -237,15 +236,15 @@ func (fetcher *Fetcher) Fetch(ctx context.Context, pkg *fv1.Package, req types.F } } else { var archive *fv1.Archive - if req.FetchType == types.FETCH_SOURCE { + if req.FetchType == fv1.FETCH_SOURCE { archive = &pkg.Spec.Source - } else if req.FetchType == types.FETCH_DEPLOYMENT { + } else if req.FetchType == fv1.FETCH_DEPLOYMENT { // sometimes, the user may invoke the function even before the source code is built into a deploy pkg. // this results in executor sending a fetch request of type FETCH_DEPLOYMENT and since pkg.Spec.Deployment.Url will be empty, // we hit this "Get : unsupported protocol scheme "" error. // it may be useful to the user if we can send a more meaningful error in such a scenario. - if pkg.Status.BuildStatus != types.BuildStatusSucceeded && pkg.Status.BuildStatus != types.BuildStatusNone { - e := fmt.Sprintf("cannot fetch deployment: package build status was not %q", types.BuildStatusSucceeded) + if pkg.Status.BuildStatus != fv1.BuildStatusSucceeded && pkg.Status.BuildStatus != fv1.BuildStatusNone { + e := fmt.Sprintf("cannot fetch deployment: package build status was not %q", fv1.BuildStatusSucceeded) fetcher.logger.Error(e, zap.String("package_name", pkg.ObjectMeta.Name), zap.String("package_namespace", pkg.ObjectMeta.Namespace), @@ -441,7 +440,7 @@ func (fetcher *Fetcher) UploadHandler(w http.ResponseWriter, r *http.Request) { return } - var req types.ArchiveUploadRequest + var req ArchiveUploadRequest err = json.Unmarshal(body, &req) if err != nil { fetcher.logger.Error("error parsing request body", zap.Error(err)) @@ -491,7 +490,7 @@ func (fetcher *Fetcher) UploadHandler(w http.ResponseWriter, r *http.Request) { return } - resp := types.ArchiveUploadResponse{ + resp := ArchiveUploadResponse{ ArchiveDownloadUrl: ssClient.GetUrl(fileID), Checksum: *sum, } @@ -548,7 +547,7 @@ func (fetcher *Fetcher) unarchive(src string, dst string) error { } // getPkgInformation gets package information from k8s api server. -func (fetcher *Fetcher) getPkgInformation(req types.FunctionFetchRequest) (pkg *fv1.Package, err error) { +func (fetcher *Fetcher) getPkgInformation(req FunctionFetchRequest) (pkg *fv1.Package, err error) { maxRetries := 5 for i := 0; i < maxRetries; i++ { pkg, err = fetcher.fissionClient.V1().Packages(req.Package.Namespace).Get(req.Package.Name, metav1.GetOptions{}) @@ -569,7 +568,7 @@ func (fetcher *Fetcher) getPkgInformation(req types.FunctionFetchRequest) (pkg * return nil, err } -func (fetcher *Fetcher) SpecializePod(ctx context.Context, fetchReq types.FunctionFetchRequest, loadReq types.FunctionLoadRequest) error { +func (fetcher *Fetcher) SpecializePod(ctx context.Context, fetchReq FunctionFetchRequest, loadReq FunctionLoadRequest) error { startTime := time.Now() defer func() { elapsed := time.Since(startTime) diff --git a/pkg/fetcher/types.go b/pkg/fetcher/types.go new file mode 100644 index 00000000..3203d0af --- /dev/null +++ b/pkg/fetcher/types.go @@ -0,0 +1,85 @@ +/* +Copyright 2016 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package fetcher + +import ( + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" +) + +// +// Fission-Environment interface. The following types are not +// exposed in the Fission API, but rather used by Fission to +// talk to environments. +// +type ( + FetchRequestType int + + FunctionSpecializeRequest struct { + FetchReq FunctionFetchRequest + LoadReq FunctionLoadRequest + } + + FunctionFetchRequest struct { + FetchType FetchRequestType `json:"fetchType"` + Package metav1.ObjectMeta `json:"package"` + Url string `json:"url"` + StorageSvcUrl string `json:"storagesvcurl"` + Filename string `json:"filename"` + Secrets []fv1.SecretReference `json:"secretList"` + ConfigMaps []fv1.ConfigMapReference `json:"configMapList"` + KeepArchive bool `json:"keeparchive"` + } + + FunctionLoadRequest struct { + // FilePath is an absolute filesystem path to the + // function. What exactly is stored here is + // env-specific. Optional. + FilePath string `json:"filepath"` + + // FunctionName has an environment-specific meaning; + // usually, it defines a function within a module + // containing multiple functions. Optional; default is + // environment-specific. + FunctionName string `json:"functionName"` + + // URL to expose this function at. Optional; defaults + // to "/". + URL string `json:"url"` + + // Metatdata + FunctionMetadata *metav1.ObjectMeta + + EnvVersion int `json:"envVersion"` + } + + // ArchiveUploadRequest send from builder manager describes which + // deployment package should be upload to storage service. + ArchiveUploadRequest struct { + Filename string `json:"filename"` + StorageSvcUrl string `json:"storagesvcurl"` + ArchivePackage bool `json:"archivepackage"` + } + + // ArchiveUploadResponse defines the download url of an archive and + // its checksum. + ArchiveUploadResponse struct { + ArchiveDownloadUrl string `json:"archiveDownloadUrl"` + Checksum fv1.Checksum `json:"checksum"` + } +) diff --git a/pkg/fission-cli/cmd/canaryconfig/create.go b/pkg/fission-cli/cmd/canaryconfig/create.go index 39905868..8d5a85c5 100644 --- a/pkg/fission-cli/cmd/canaryconfig/create.go +++ b/pkg/fission-cli/cmd/canaryconfig/create.go @@ -28,7 +28,6 @@ import ( "github.com/fission/fission/pkg/fission-cli/cmd" flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" "github.com/fission/fission/pkg/fission-cli/util" - "github.com/fission/fission/pkg/types" ) type CreateSubCommand struct { @@ -76,7 +75,7 @@ func (opts *CreateSubCommand) complete(input cli.Input) error { } // check that the trigger has function reference type function weights - if htTrigger.Spec.FunctionReference.Type != types.FunctionReferenceTypeFunctionWeights { + if htTrigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionWeights { return errors.New("canary config cannot be created for http triggers that do not reference functions by weights") } diff --git a/pkg/fission-cli/cmd/mqtrigger/create.go b/pkg/fission-cli/cmd/mqtrigger/create.go index 59296ee1..d1004312 100644 --- a/pkg/fission-cli/cmd/mqtrigger/create.go +++ b/pkg/fission-cli/cmd/mqtrigger/create.go @@ -30,7 +30,6 @@ import ( "github.com/fission/fission/pkg/fission-cli/console" flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" "github.com/fission/fission/pkg/fission-cli/util" - "github.com/fission/fission/pkg/types" ) type CreateSubCommand struct { @@ -62,13 +61,13 @@ func (opts *CreateSubCommand) complete(input cli.Input) error { var mqType fv1.MessageQueueType switch input.String(flagkey.MqtMQType) { case "": - mqType = types.MessageQueueTypeNats - case types.MessageQueueTypeNats: - mqType = types.MessageQueueTypeNats - case types.MessageQueueTypeASQ: - mqType = types.MessageQueueTypeASQ - case types.MessageQueueTypeKafka: - mqType = types.MessageQueueTypeKafka + mqType = fv1.MessageQueueTypeNats + case fv1.MessageQueueTypeNats: + mqType = fv1.MessageQueueTypeNats + case fv1.MessageQueueTypeASQ: + mqType = fv1.MessageQueueTypeASQ + case fv1.MessageQueueTypeKafka: + mqType = fv1.MessageQueueTypeKafka default: return errors.New("Unknown message queue type, currently only \"nats-streaming, azure-storage-queue, kafka \" is supported") } @@ -131,7 +130,7 @@ func (opts *CreateSubCommand) complete(input cli.Input) error { }, Spec: fv1.MessageQueueTriggerSpec{ FunctionReference: fv1.FunctionReference{ - Type: types.FunctionReferenceTypeFunctionName, + Type: fv1.FunctionReferenceTypeFunctionName, Name: fnName, }, MessageQueueType: mqType, diff --git a/pkg/fission-cli/cmd/package/util/util.go b/pkg/fission-cli/cmd/package/util/util.go index d2aefeb0..7508ac58 100644 --- a/pkg/fission-cli/cmd/package/util/util.go +++ b/pkg/fission-cli/cmd/package/util/util.go @@ -33,7 +33,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/controller/client" storageSvcClient "github.com/fission/fission/pkg/storagesvc/client" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -45,7 +44,7 @@ func UploadArchiveFile(ctx context.Context, client client.Interface, fileName st return nil, err } - if size < types.ArchiveLiteralSizeLimit { + if size < fv1.ArchiveLiteralSizeLimit { archive.Type = fv1.ArchiveTypeLiteral archive.Literal, err = GetContents(fileName) if err != nil { diff --git a/pkg/fission-cli/cmd/spec/apply.go b/pkg/fission-cli/cmd/spec/apply.go index e6cd3e32..afaece81 100644 --- a/pkg/fission-cli/cmd/spec/apply.go +++ b/pkg/fission-cli/cmd/spec/apply.go @@ -40,7 +40,6 @@ import ( "github.com/fission/fission/pkg/fission-cli/console" flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" "github.com/fission/fission/pkg/fission-cli/util" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -474,7 +473,7 @@ func localArchiveFromSpec(specDir string, aus *spectypes.ArchiveUploadSpec) (*fv } // figure out if we're making a literal or a URL-based archive - if size < types.ArchiveLiteralSizeLimit { + if size < fv1.ArchiveLiteralSizeLimit { contents, err := pkgutil.GetContents(archiveFileName) if err != nil { return nil, err diff --git a/pkg/fission-cli/cmd/spec/buildwatch.go b/pkg/fission-cli/cmd/spec/buildwatch.go index 222b3610..a4c611dc 100644 --- a/pkg/fission-cli/cmd/spec/buildwatch.go +++ b/pkg/fission-cli/cmd/spec/buildwatch.go @@ -27,7 +27,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/controller/client" "github.com/fission/fission/pkg/fission-cli/cmd/package/util" - "github.com/fission/fission/pkg/types" ) type ( @@ -83,11 +82,11 @@ func (w *packageBuildWatcher) watch(ctx context.Context) { if !ok { continue } - if pkg.Status.BuildStatus == types.BuildStatusNone { + if pkg.Status.BuildStatus == fv1.BuildStatusNone { continue } - if pkg.Status.BuildStatus == types.BuildStatusPending || - pkg.Status.BuildStatus == types.BuildStatusRunning { + if pkg.Status.BuildStatus == fv1.BuildStatusPending || + pkg.Status.BuildStatus == fv1.BuildStatusRunning { keepWaiting = true } buildpkgs = append(buildpkgs, pkg) @@ -99,8 +98,8 @@ func (w *packageBuildWatcher) watch(ctx context.Context) { if _, printed := w.finished[k]; printed { continue } - if pkg.Status.BuildStatus == types.BuildStatusFailed || - pkg.Status.BuildStatus == types.BuildStatusSucceeded { + if pkg.Status.BuildStatus == fv1.BuildStatusFailed || + pkg.Status.BuildStatus == fv1.BuildStatusSucceeded { w.finished[k] = true fmt.Printf("------\n") util.PrintPackageSummary(os.Stdout, &pkg) diff --git a/pkg/fission-cli/cmd/support/resources/crd.go b/pkg/fission-cli/cmd/support/resources/crd.go index 7056f2dc..9522f2bd 100644 --- a/pkg/fission-cli/cmd/support/resources/crd.go +++ b/pkg/fission-cli/cmd/support/resources/crd.go @@ -24,7 +24,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/controller/client" "github.com/fission/fission/pkg/fission-cli/console" - "github.com/fission/fission/pkg/types" ) const ( @@ -114,7 +113,7 @@ func (res CrdDumper) Dump(dumpDir string) { case CrdMessageQueueTrigger: var triggers []fv1.MessageQueueTrigger - for _, mqType := range []string{types.MessageQueueTypeNats, types.MessageQueueTypeASQ} { + for _, mqType := range []string{fv1.MessageQueueTypeNats, fv1.MessageQueueTypeASQ, fv1.MessageQueueTypeKafka} { l, err := res.client.V1().MessageQueueTrigger().List(mqType, metav1.NamespaceAll) if err != nil { console.Warn(fmt.Sprintf("Error getting %v list: %v", res.crdType, err)) diff --git a/pkg/logger/logger.go b/pkg/logger/logger.go index b8f1ed5e..1447db18 100644 --- a/pkg/logger/logger.go +++ b/pkg/logger/logger.go @@ -24,8 +24,6 @@ import ( "strings" "time" - "github.com/fission/fission/pkg/types" - "github.com/fission/fission/pkg/utils" "go.uber.org/zap" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -33,7 +31,9 @@ import ( "k8s.io/client-go/kubernetes" k8sCache "k8s.io/client-go/tools/cache" + fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" "github.com/fission/fission/pkg/crd" + "github.com/fission/fission/pkg/utils" ) var nodeName = os.Getenv("NODE_NAME") @@ -55,7 +55,7 @@ func makePodLoggerController(zapLogger *zap.Logger, k8sClientSet *kubernetes.Cli } err := createLogSymlinks(zapLogger, pod) if err != nil { - funcName := pod.Labels[types.FUNCTION_NAME] + funcName := pod.Labels[fv1.FUNCTION_NAME] zapLogger.Error("error creating symlink", zap.String("function", funcName), zap.Error(err)) } @@ -67,7 +67,7 @@ func makePodLoggerController(zapLogger *zap.Logger, k8sClientSet *kubernetes.Cli } err := createLogSymlinks(zapLogger, pod) if err != nil { - funcName := pod.Labels[types.FUNCTION_NAME] + funcName := pod.Labels[fv1.FUNCTION_NAME] zapLogger.Error("error creating symlink", zap.String("function", funcName), zap.Error(err)) } @@ -115,8 +115,8 @@ func isValidFunctionPodOnNode(pod *corev1.Pod) bool { if pod.Spec.NodeName != nodeName { return false } - labels := []string{types.ENVIRONMENT_NAMESPACE, types.ENVIRONMENT_NAME, types.ENVIRONMENT_UID, - types.FUNCTION_NAMESPACE, types.FUNCTION_NAME, types.FUNCTION_UID, types.EXECUTOR_TYPE} + labels := []string{fv1.ENVIRONMENT_NAMESPACE, fv1.ENVIRONMENT_NAME, fv1.ENVIRONMENT_UID, + fv1.FUNCTION_NAMESPACE, fv1.FUNCTION_NAME, fv1.FUNCTION_UID, fv1.EXECUTOR_TYPE} for _, l := range labels { if len(pod.Labels[l]) == 0 { return false diff --git a/pkg/mqtrigger/messageQueue/asq.go b/pkg/mqtrigger/messageQueue/asq.go index 3d3c7e48..df4ace01 100644 --- a/pkg/mqtrigger/messageQueue/asq.go +++ b/pkg/mqtrigger/messageQueue/asq.go @@ -29,7 +29,6 @@ import ( "time" "github.com/Azure/azure-sdk-for-go/storage" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" "github.com/pkg/errors" "go.uber.org/zap" @@ -204,7 +203,7 @@ func newAzureStorageConnection(logger *zap.Logger, routerURL string, config Mess func (asc AzureStorageConnection) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubscription, error) { asc.logger.Info("subscribing to Azure storage queue", zap.String("queue", trigger.Spec.Topic)) - if trigger.Spec.FunctionReference.Type != types.FunctionReferenceTypeFunctionName { + if trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { return nil, fmt.Errorf("unsupported function reference type (%v) for trigger %q", trigger.Spec.FunctionReference.Type, trigger.ObjectMeta.Name) } diff --git a/pkg/mqtrigger/messageQueue/asq_test.go b/pkg/mqtrigger/messageQueue/asq_test.go index cd4e9901..d28d6134 100644 --- a/pkg/mqtrigger/messageQueue/asq_test.go +++ b/pkg/mqtrigger/messageQueue/asq_test.go @@ -26,7 +26,6 @@ import ( "time" "github.com/Azure/azure-sdk-for-go/storage" - "github.com/fission/fission/pkg/types" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" "go.uber.org/zap" @@ -114,7 +113,7 @@ func TestNewStorageConnectionMissingAccountName(t *testing.T) { panicIf(err) connection, err := newAzureStorageConnection(logger, DummyRouterURL, MessageQueueConfig{ - MQType: types.MessageQueueTypeASQ, + MQType: fv1.MessageQueueTypeASQ, Url: "", }) require.Nil(t, connection) @@ -127,7 +126,7 @@ func TestNewStorageConnectionMissingAccessKey(t *testing.T) { _ = os.Setenv("AZURE_STORAGE_ACCOUNT_NAME", "accountname") connection, err := newAzureStorageConnection(logger, DummyRouterURL, MessageQueueConfig{ - MQType: types.MessageQueueTypeASQ, + MQType: fv1.MessageQueueTypeASQ, Url: "", }) _ = os.Unsetenv("AZURE_STORAGE_ACCOUNT_NAME") @@ -309,10 +308,10 @@ func TestAzureStorageQueuePoisonMessage(t *testing.T) { }, Spec: fv1.MessageQueueTriggerSpec{ FunctionReference: fv1.FunctionReference{ - Type: types.FunctionReferenceTypeFunctionName, + Type: fv1.FunctionReferenceTypeFunctionName, Name: FunctionName, }, - MessageQueueType: types.MessageQueueTypeASQ, + MessageQueueType: fv1.MessageQueueTypeASQ, Topic: QueueName, ContentType: ContentType, }, @@ -457,10 +456,10 @@ func runAzureStorageQueueTest(t *testing.T, count int, output bool) { }, Spec: fv1.MessageQueueTriggerSpec{ FunctionReference: fv1.FunctionReference{ - Type: types.FunctionReferenceTypeFunctionName, + Type: fv1.FunctionReferenceTypeFunctionName, Name: FunctionName, }, - MessageQueueType: types.MessageQueueTypeASQ, + MessageQueueType: fv1.MessageQueueTypeASQ, Topic: QueueName, ResponseTopic: responseTopic, ContentType: ContentType, diff --git a/pkg/mqtrigger/messageQueue/kafka.go b/pkg/mqtrigger/messageQueue/kafka.go index e52adaee..8d91e1bb 100644 --- a/pkg/mqtrigger/messageQueue/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka.go @@ -28,12 +28,11 @@ import ( sarama "github.com/Shopify/sarama" cluster "github.com/bsm/sarama-cluster" - "github.com/fission/fission/pkg/types" - "github.com/fission/fission/pkg/utils" "github.com/pkg/errors" "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" + "github.com/fission/fission/pkg/utils" ) type ( @@ -202,7 +201,7 @@ func (kafka Kafka) unsubscribe(subscription messageQueueSubscription) error { func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.MessageQueueTrigger, msg *sarama.ConsumerMessage, consumer *cluster.Consumer) { var value string = string(msg.Value[:]) // Support other function ref types - if trigger.Spec.FunctionReference.Type != types.FunctionReferenceTypeFunctionName { + if trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { kafka.logger.Fatal("unsupported function reference type for trigger", zap.Any("function_reference_type", trigger.Spec.FunctionReference.Type), zap.String("trigger", trigger.ObjectMeta.Name)) diff --git a/pkg/mqtrigger/messageQueue/messageQueue.go b/pkg/mqtrigger/messageQueue/messageQueue.go index 5747b251..9098e1bd 100644 --- a/pkg/mqtrigger/messageQueue/messageQueue.go +++ b/pkg/mqtrigger/messageQueue/messageQueue.go @@ -21,7 +21,6 @@ import ( "fmt" "time" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -87,11 +86,11 @@ func MakeMessageQueueTriggerManager(logger *zap.Logger, fissionClient *crd.Fissi fissionClient: fissionClient, } switch mqConfig.MQType { - case types.MessageQueueTypeNats: + case fv1.MessageQueueTypeNats: messageQueue, err = makeNatsMessageQueue(logger, routerUrl, mqConfig) - case types.MessageQueueTypeASQ: + case fv1.MessageQueueTypeASQ: messageQueue, err = newAzureStorageConnection(logger, routerUrl, mqConfig) - case types.MessageQueueTypeKafka: + case fv1.MessageQueueTypeKafka: messageQueue, err = makeKafkaMessageQueue(logger, routerUrl, mqConfig) default: err = fmt.Errorf("no supported message queue type found for %q", mqConfig.MQType) diff --git a/pkg/mqtrigger/messageQueue/nats.go b/pkg/mqtrigger/messageQueue/nats.go index d8ee28a1..99a047aa 100644 --- a/pkg/mqtrigger/messageQueue/nats.go +++ b/pkg/mqtrigger/messageQueue/nats.go @@ -28,7 +28,6 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" - "github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/utils" ) @@ -106,7 +105,7 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) { return func(msg *ns.Msg) { // Support other function ref types - if trigger.Spec.FunctionReference.Type != types.FunctionReferenceTypeFunctionName { + if trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { nats.logger.Fatal("unsupported function reference type for trigger", zap.Any("function_reference_type", trigger.Spec.FunctionReference.Type), zap.String("trigger", trigger.ObjectMeta.Name)) diff --git a/pkg/router/functionHandler.go b/pkg/router/functionHandler.go index ee32dc1c..2a44450f 100644 --- a/pkg/router/functionHandler.go +++ b/pkg/router/functionHandler.go @@ -40,7 +40,6 @@ import ( "github.com/fission/fission/pkg/error/network" executorClient "github.com/fission/fission/pkg/executor/client" "github.com/fission/fission/pkg/throttler" - "github.com/fission/fission/pkg/types" ) const ( @@ -367,7 +366,7 @@ func (fh *functionHandler) tapService(fn *fv1.Function, serviceUrl *url.URL) { } func (fh functionHandler) handler(responseWriter http.ResponseWriter, request *http.Request) { - if fh.httpTrigger != nil && fh.httpTrigger.Spec.FunctionReference.Type == types.FunctionReferenceTypeFunctionWeights { + if fh.httpTrigger != nil && fh.httpTrigger.Spec.FunctionReference.Type == fv1.FunctionReferenceTypeFunctionWeights { // canary deployment. need to determine the function to send request to now fn := getCanaryBackend(fh.functionMap, fh.fnWeightDistributionList) if fn == nil { diff --git a/pkg/router/functionHandler_test.go b/pkg/router/functionHandler_test.go index 57af70a8..ec83bb4d 100644 --- a/pkg/router/functionHandler_test.go +++ b/pkg/router/functionHandler_test.go @@ -31,7 +31,6 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" - "github.com/fission/fission/pkg/types" ) func createBackendService(testResponseString string) *url.URL { @@ -72,7 +71,7 @@ func TestFunctionProxying(t *testing.T) { }, Spec: fv1.HTTPTriggerSpec{ FunctionReference: fv1.FunctionReference{ - Type: types.FunctionReferenceTypeFunctionName, + Type: fv1.FunctionReferenceTypeFunctionName, }, }, } diff --git a/pkg/router/router_test.go b/pkg/router/router_test.go index 06f56356..3c2345cc 100644 --- a/pkg/router/router_test.go +++ b/pkg/router/router_test.go @@ -22,7 +22,6 @@ import ( "testing" "time" - "github.com/fission/fission/pkg/types" "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -36,7 +35,7 @@ func TestRouter(t *testing.T) { // and a reference to it fr := fv1.FunctionReference{ - Type: types.FunctionReferenceTypeFunctionName, + Type: fv1.FunctionReferenceTypeFunctionName, Name: fnMeta.Name, } diff --git a/pkg/types/types.go b/pkg/types/types.go deleted file mode 100644 index fbb7dace..00000000 --- a/pkg/types/types.go +++ /dev/null @@ -1,188 +0,0 @@ -/* -Copyright 2016 The Fission Authors. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ - -package types - -import ( - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - - fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" -) - -// -// Fission-Environment interface. The following types are not -// exposed in the Fission API, but rather used by Fission to -// talk to environments. -// -type ( - FetchRequestType int - - FunctionSpecializeRequest struct { - FetchReq FunctionFetchRequest - LoadReq FunctionLoadRequest - } - - FunctionFetchRequest struct { - FetchType FetchRequestType `json:"fetchType"` - Package metav1.ObjectMeta `json:"package"` - Url string `json:"url"` - StorageSvcUrl string `json:"storagesvcurl"` - Filename string `json:"filename"` - Secrets []fv1.SecretReference `json:"secretList"` - ConfigMaps []fv1.ConfigMapReference `json:"configMapList"` - KeepArchive bool `json:"keeparchive"` - } - - FunctionLoadRequest struct { - // FilePath is an absolute filesystem path to the - // function. What exactly is stored here is - // env-specific. Optional. - FilePath string `json:"filepath"` - - // FunctionName has an environment-specific meaning; - // usually, it defines a function within a module - // containing multiple functions. Optional; default is - // environment-specific. - FunctionName string `json:"functionName"` - - // URL to expose this function at. Optional; defaults - // to "/". - URL string `json:"url"` - - // Metatdata - FunctionMetadata *metav1.ObjectMeta - - EnvVersion int `json:"envVersion"` - } - - // ArchiveUploadRequest send from builder manager describes which - // deployment package should be upload to storage service. - ArchiveUploadRequest struct { - Filename string `json:"filename"` - StorageSvcUrl string `json:"storagesvcurl"` - ArchivePackage bool `json:"archivepackage"` - } - - // ArchiveUploadResponse defines the download url of an archive and - // its checksum. - ArchiveUploadResponse struct { - ArchiveDownloadUrl string `json:"archiveDownloadUrl"` - Checksum fv1.Checksum `json:"checksum"` - } -) - -const ( - FETCH_SOURCE = iota - FETCH_DEPLOYMENT - FETCH_URL // remove this? -) - -const EXECUTOR_INSTANCEID_LABEL = fv1.EXECUTOR_INSTANCEID_LABEL - -const ( - ChecksumTypeSHA256 = fv1.ChecksumTypeSHA256 -) - -const ( - // ArchiveTypeLiteral means the package contents are specified in the Literal field of - // resource itself. - ArchiveTypeLiteral = fv1.ArchiveTypeLiteral - - // ArchiveTypeUrl means the package contents are at the specified URL. - ArchiveTypeUrl = fv1.ArchiveTypeUrl -) - -const ( - BuildStatusPending = fv1.BuildStatusPending - BuildStatusRunning = fv1.BuildStatusRunning - BuildStatusSucceeded = fv1.BuildStatusSucceeded - BuildStatusFailed = fv1.BuildStatusFailed - BuildStatusNone = fv1.BuildStatusNone -) - -const ( - AllowedFunctionsPerContainerSingle = fv1.AllowedFunctionsPerContainerSingle - AllowedFunctionsPerContainerInfinite = fv1.AllowedFunctionsPerContainerInfinite -) - -// executor kubernetes object label key -const ( - ENVIRONMENT_NAMESPACE = "environmentNamespace" - ENVIRONMENT_NAME = "environmentName" - ENVIRONMENT_UID = "environmentUid" - FUNCTION_NAMESPACE = "functionNamespace" - FUNCTION_NAME = "functionName" - FUNCTION_UID = "functionUid" - FUNCTION_RESOURCE_VERSION = "functionResourceVersion" - EXECUTOR_TYPE = "executorType" -) - -const ( - ANNOTATION_SVC_HOST = "svcHost" -) - -const ( - SharedVolumeUserfunc = fv1.SharedVolumeUserfunc - SharedVolumePackages = fv1.SharedVolumePackages - SharedVolumeSecrets = fv1.SharedVolumeSecrets - SharedVolumeConfigmaps = fv1.SharedVolumeConfigmaps -) - -const ( - MessageQueueTypeNats = fv1.MessageQueueTypeNats - MessageQueueTypeASQ = fv1.MessageQueueTypeASQ - MessageQueueTypeKafka = fv1.MessageQueueTypeKafka -) - -const ( - // FunctionReferenceFunctionName means that the function - // reference is simply by function name. - FunctionReferenceTypeFunctionName = fv1.FunctionReferenceTypeFunctionName - - // Set of function references (recursively), by percentage of traffic - FunctionReferenceTypeFunctionWeights = fv1.FunctionReferenceTypeFunctionWeights - - // Other function reference types we'd like to support: - // Versioned function, latest version - // Versioned function. by semver "latest compatible" - -) - -const ( - ArchiveLiteralSizeLimit int64 = 256 * 1024 -) - -const ( - FissionBuilderSA = "fission-builder" - FissionFetcherSA = "fission-fetcher" - - SecretConfigMapGetterCR = "secret-configmap-getter" - SecretConfigMapGetterRB = "secret-configmap-getter-binding" - - PackageGetterCR = "package-getter" - PackageGetterRB = "package-getter-binding" - - ClusterRole = "ClusterRole" -) - -const ( - FailureTypeStatusCode = fv1.FailureTypeStatusCode - CanaryConfigStatusPending = fv1.CanaryConfigStatusPending - CanaryConfigStatusSucceeded = fv1.CanaryConfigStatusSucceeded - CanaryConfigStatusFailed = fv1.CanaryConfigStatusFailed - CanaryConfigStatusAborted = fv1.CanaryConfigStatusAborted - MaxIterationsForCanaryConfig = fv1.MaxIterationsForCanaryConfig -) diff --git a/pkg/v1/types.go b/pkg/v1/types.go deleted file mode 100644 index fbc48599..00000000 --- a/pkg/v1/types.go +++ /dev/null @@ -1,100 +0,0 @@ -/* -Copyright 2016 The Fission Authors. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ - -package v1 - -// -// These are types from the v1 API, and are only preserved for -// compatibility. They should never be changed. -// - -type ( - // ObjectMeta is used as the general identifier for all kinds of - // resources managed by the controller. - Metadata struct { - Name string `json:"name"` - Uid string `json:"uid,omitempty"` - } - - // Function is a unit of executable code. Though it's called - // a function, the code may have more than one function; it's - // usually some sort of module or package. - Function struct { - Metadata `json:"metadata"` - Environment Metadata `json:"environment"` - Code string `json:"code"` - } - - // Environment identifies the language and OS specific - // resources that a function depends on. For now this - // includes only the function run container image. Later, - // this will also include build containers, as well as support - // tools like debuggers, profilers, etc. - Environment struct { - Metadata `json:"metadata"` - RunContainerImageUrl string `json:"runContainerImageUrl"` - } - - // HTTPTrigger maps URL patterns to functions. Function.UID - // is optional; if absent, the latest version of the function - // will automatically be selected. - HTTPTrigger struct { - Metadata `json:"metadata"` - UrlPattern string `json:"urlpattern"` - Method string `json:"method"` - Function Metadata `json:"function"` - } - - MessageQueueTrigger struct { - Metadata `json:"metadata"` - Function Metadata `json:"function"` - MessageQueueType string `json:"messageQueueType"` - Topic string `json:"topic"` - ResponseTopic string `json:"respTopic,omitempty"` - } - - // Watch is a specification of Kubernetes watch along with a URL to post events to. - Watch struct { - Metadata `json:"metadata"` - - Namespace string `json:"namespace"` - ObjType string `json:"objtype"` - LabelSelector string `json:"labelselector"` - FieldSelector string `json:"fieldselector"` - - Function Metadata `json:"function"` - - Target string `json:"target"` // Watch publish target (URL, NATS stream, etc) - } - - // TimeTrigger invokes the specific function at a time or - // times specified by a cron string. - TimeTrigger struct { - Metadata `json:"metadata"` - - Cron string `json:"cron"` - - Function Metadata `json:"function"` - } - - // Errors returned by the Fission API. - Error struct { - Code errorCode `json:"code"` - Message string `json:"message"` - } - - errorCode int -)