From 6b49c194ca52f6f4519bc83fb3ef81fb0a7cd3b7 Mon Sep 17 00:00:00 2001 From: Ta-Ching Chen Date: Wed, 19 Sep 2018 14:35:51 +0800 Subject: [PATCH] Fission supportability: Add dump command to dump information for debugging (#754) * Add v1 support tool proposal * Fission supportability: Add dump command to dump information for debugging --- Documentation/wip/support-tool.md | 62 ++++ common.go | 48 +++ fission/buildwatch.go | 3 +- fission/common.go | 343 ------------------- fission/environment.go | 25 +- fission/function.go | 105 ++++-- fission/httptrigger.go | 21 +- fission/main.go | 87 +---- fission/mqtrigger.go | 21 +- fission/package.go | 256 ++++++++++++-- fission/recorder.go | 23 +- fission/records.go | 19 +- fission/replay.go | 5 +- fission/spec.go | 37 +- fission/support/dump.go | 141 ++++++++ fission/support/resources/crd.go | 157 +++++++++ fission/support/resources/fissionversion.go | 40 +++ fission/support/resources/kubernetes.go | 249 ++++++++++++++ fission/support/resources/resource.go | 57 +++ fission/timetrigger.go | 29 +- fission/upgrade.go | 57 +-- fission/{portforward => util}/portforward.go | 16 +- fission/util/util.go | 142 ++++++++ fission/util/version.go | 48 +++ fission/watch.go | 15 +- 25 files changed, 1397 insertions(+), 609 deletions(-) create mode 100644 Documentation/wip/support-tool.md delete mode 100644 fission/common.go create mode 100644 fission/support/dump.go create mode 100644 fission/support/resources/crd.go create mode 100644 fission/support/resources/fissionversion.go create mode 100644 fission/support/resources/kubernetes.go create mode 100644 fission/support/resources/resource.go rename fission/{portforward => util}/portforward.go (92%) create mode 100644 fission/util/util.go create mode 100644 fission/util/version.go diff --git a/Documentation/wip/support-tool.md b/Documentation/wip/support-tool.md new file mode 100644 index 00000000..a9c260ca --- /dev/null +++ b/Documentation/wip/support-tool.md @@ -0,0 +1,62 @@ +# Fission Support Tool +Fission now has rich functionality supported by multiple services, however, it brings the complexity of troubleshooting. +This proposal tends to give a picture of fission support tool that can help both user and developer to locate the problem in short time. +To achieve this, the support tool will dump related kubernetes objects, fission resources and pod logs from the given cluster. + +# Functionality + +## Environment Information Collection + +Before troubleshooting, some of the basic information is needed to give others an overview of kubernetes/fission user test with so that we can locate the problem in short time. + +* Fission version + * Client/Server version + +* Kubernetes cluster version + * Cluster version (i.e v1.9.7-gke.0) + * Running environment (i.e GKE, AKS and minikube) + * Nodes version and other information + +## Service Logs Collection + +The component logs and the logs of interaction between components are important for people to understand what really happened in cluster. Following are components need to collect logs from. + +* All fission component pods +* Function pods +* Builder pods +* Environment pods + +## Object dumping + +Fission is deeply coupled with Kubernetes, most of the objects are created and maintained by it. There is two major type of objects need to be dumped from kubernetes: + +* K8S objects +* CRD resources + +All objects should be dumped into a readable file format. It will be great if people can reproduce similar environment with these files. + +## Information upload + +Upload dump files to the specific backend server for support channel to analysis + +# CLI Interface + +``` +$ fission support collect + NAME: + fission support collect - Collect pod logs, fission resources and related kubernetes objects for troubleshooting + + USAGE: + fission support collect [command options] [arguments...] + + OPTIONS: + --dumpdir value Directory to save dump kubernetes objects and fission resources (default: "fission-dump") + --fissionns value Namespace of fission installation (default: "fission") + --builderns value Namespace of fission package builder (default: "fission-builder") + --funcns value Namespace of fission function pod (default: "fission-function") +``` + +# Thoughts? + +1. What to do with sensitive objects like secrets and configmap? Ignore the dump for such objects? +2. The functionality is necessary but not listed above? diff --git a/common.go b/common.go index e925490e..e73bb237 100644 --- a/common.go +++ b/common.go @@ -18,18 +18,24 @@ package fission import ( "fmt" + "io/ioutil" "net" "net/http" "os" "os/signal" + "path/filepath" "runtime/debug" "strings" "syscall" "github.com/gorilla/handlers" "github.com/imdario/mergo" + "github.com/mholt/archiver" + uuid "github.com/satori/go.uuid" apiv1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + "github.com/fission/fission/fission/log" ) func UrlForFunction(name, namespace string) string { @@ -130,3 +136,45 @@ func IsReadyPod(pod *apiv1.Pod) bool { return true } + +// GetTempDir creates and return a temporary directory +func GetTempDir() (string, error) { + tmpDir := uuid.NewV4().String() + dir, err := ioutil.TempDir("", tmpDir) + return dir, err +} + +func MakeArchive(targetName string, globs ...string) (string, error) { + files := make([]string, 0) + for _, glob := range globs { + f, err := filepath.Glob(glob) + if err != nil { + log.Warn(fmt.Sprintf("Invalid glob %v: %v", glob, err)) + return "", err + } + files = append(files, f...) + } + + // zip up the file list + err := archiver.Zip.Make(targetName, files) + if err != nil { + return "", err + } + + return filepath.Abs(targetName) +} + +// RemoveZeroBytes remove empty byte(\x00) from input byte slice and return a new byte slice +// This function is trying to fix the problem that empty byte will fail os.Openfile +// For more information, please visit: +// 1. https://github.com/golang/go/issues/24195 +// 2. https://play.golang.org/p/5F9ykC2tlbc +func RemoveZeroBytes(src []byte) []byte { + var bs []byte + for _, v := range src { + if v != 0 { + bs = append(bs, v) + } + } + return bs +} diff --git a/fission/buildwatch.go b/fission/buildwatch.go index 74c40b12..f653129c 100644 --- a/fission/buildwatch.go +++ b/fission/buildwatch.go @@ -26,6 +26,7 @@ import ( "github.com/fission/fission" "github.com/fission/fission/controller/client" "github.com/fission/fission/crd" + "github.com/fission/fission/fission/util" ) type ( @@ -67,7 +68,7 @@ func (w *packageBuildWatcher) watch(ctx context.Context) { // poll list of packages (TODO: convert to watch) pkgs, err := w.fclient.PackageList(metav1.NamespaceAll) - checkErr(err, "Getting list of packages") + util.CheckErr(err, "Getting list of packages") // find packages that (a) are in the app spec and (b) have an interesting // build status (either succeeded or failed; not "none") diff --git a/fission/common.go b/fission/common.go deleted file mode 100644 index 925456cf..00000000 --- a/fission/common.go +++ /dev/null @@ -1,343 +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 main - -import ( - "crypto/sha256" - "encoding/hex" - "errors" - "fmt" - "io" - "io/ioutil" - "net/http" - "os" - "path" - "path/filepath" - "regexp" - "strings" - - "github.com/dchest/uniuri" - uuid "github.com/satori/go.uuid" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - - "github.com/fission/fission" - "github.com/fission/fission/controller/client" - "github.com/fission/fission/crd" - "github.com/fission/fission/fission/log" - storageSvcClient "github.com/fission/fission/storagesvc/client" -) - -func getClient(serverUrl string) *client.Client { - if len(serverUrl) == 0 { - // starts local portforwarder etc. - serverUrl = getServerUrl() - } - - isHTTPS := strings.Index(serverUrl, "https://") == 0 - isHTTP := strings.Index(serverUrl, "http://") == 0 - - if !(isHTTP || isHTTPS) { - serverUrl = "http://" + serverUrl - } - - return client.MakeClient(serverUrl) -} - -func checkErr(err error, msg string) { - if err != nil { - log.Fatal(fmt.Sprintf("Failed to %v: %v", msg, err)) - } -} - -func httpRequest(method, url, body string, headers []string) *http.Response { - if method == "" { - method = "GET" - } - - if method != http.MethodGet && - method != http.MethodDelete && - method != http.MethodPost && - method != http.MethodPut { - log.Fatal(fmt.Sprintf("Invalid HTTP method '%s'.", method)) - } - - req, err := http.NewRequest(method, url, strings.NewReader(body)) - checkErr(err, "create HTTP request") - - for _, header := range headers { - headerKeyValue := strings.SplitN(header, ":", 2) - if len(headerKeyValue) != 2 { - checkErr(errors.New(""), "create request without appropriate headers") - } - req.Header.Set(headerKeyValue[0], headerKeyValue[1]) - } - - client := &http.Client{} - resp, err := client.Do(req) - checkErr(err, "execute HTTP request") - - return resp -} - -func fileSize(filePath string) int64 { - info, err := os.Stat(filePath) - checkErr(err, fmt.Sprintf("stat %v", filePath)) - return info.Size() -} - -func fileChecksum(fileName string) (*fission.Checksum, error) { - f, err := os.Open(fileName) - if err != nil { - return nil, fmt.Errorf("failed to open file %v: %v", fileName, err) - } - defer f.Close() - - h := sha256.New() - _, err = io.Copy(h, f) - if err != nil { - return nil, fmt.Errorf("failed to calculate checksum for %v", fileName) - } - - return &fission.Checksum{ - Type: fission.ChecksumTypeSHA256, - Sum: hex.EncodeToString(h.Sum(nil)), - }, nil -} - -// upload a file and return a fission.Archive -func createArchive(client *client.Client, fileName string, specFile string) *fission.Archive { - var archive fission.Archive - - // fetch archive from arbitrary url if fileName is a url - if strings.HasPrefix(fileName, "http://") || strings.HasPrefix(fileName, "https://") { - fileName = downloadToTempFile(fileName) - } - - if len(specFile) > 0 { - // create an ArchiveUploadSpec and reference it from the archive - aus := &ArchiveUploadSpec{ - Name: kubifyName(path.Base(fileName)), - IncludeGlobs: []string{fileName}, - } - // save the uploadspec - err := specSave(*aus, specFile) - checkErr(err, fmt.Sprintf("write spec file %v", specFile)) - // create the archive - ar := &fission.Archive{ - Type: fission.ArchiveTypeUrl, - URL: fmt.Sprintf("%v%v", ARCHIVE_URL_PREFIX, aus.Name), - } - return ar - } - - if fileSize(fileName) < fission.ArchiveLiteralSizeLimit { - contents := getContents(fileName) - archive.Type = fission.ArchiveTypeLiteral - archive.Literal = contents - } else { - u := strings.TrimSuffix(client.Url, "/") + "/proxy/storage" - ssClient := storageSvcClient.MakeClient(u) - - // TODO add a progress bar - id, err := ssClient.Upload(fileName, nil) - checkErr(err, fmt.Sprintf("upload file %v", fileName)) - - storageSvc, err := client.GetSvcURL("application=fission-storage") - storageSvcURL := "http://" + storageSvc - checkErr(err, "get fission storage service name") - - // We make a new client with actual URL of Storage service so that the URL is not - // pointing to 127.0.0.1 i.e. proxy. DON'T reuse previous ssClient - pkgClient := storageSvcClient.MakeClient(storageSvcURL) - archiveURL := pkgClient.GetUrl(id) - - archive.Type = fission.ArchiveTypeUrl - archive.URL = archiveURL - - csum, err := fileChecksum(fileName) - checkErr(err, fmt.Sprintf("calculate checksum for file %v", fileName)) - - archive.Checksum = *csum - } - return &archive -} - -func createPackage(client *client.Client, pkgNamespace, envName, envNamespace, srcArchiveName, deployArchiveName, buildcmd string, specFile string) *metav1.ObjectMeta { - pkgSpec := fission.PackageSpec{ - Environment: fission.EnvironmentReference{ - Namespace: envNamespace, - Name: envName, - }, - } - var pkgStatus fission.BuildStatus = fission.BuildStatusSucceeded - - var pkgName string - if len(deployArchiveName) > 0 { - if len(specFile) > 0 { // we should do this in all cases, i think - pkgStatus = fission.BuildStatusNone - } - pkgSpec.Deployment = *createArchive(client, deployArchiveName, specFile) - pkgName = kubifyName(fmt.Sprintf("%v-%v", path.Base(deployArchiveName), uniuri.NewLen(4))) - } - if len(srcArchiveName) > 0 { - pkgSpec.Source = *createArchive(client, srcArchiveName, specFile) - // set pending status to package - pkgStatus = fission.BuildStatusPending - pkgName = kubifyName(fmt.Sprintf("%v-%v", path.Base(srcArchiveName), uniuri.NewLen(4))) - } - - if len(buildcmd) > 0 { - pkgSpec.BuildCommand = buildcmd - } - - if len(pkgName) == 0 { - pkgName = strings.ToLower(uuid.NewV4().String()) - } - pkg := &crd.Package{ - Metadata: metav1.ObjectMeta{ - Name: pkgName, - Namespace: pkgNamespace, - }, - Spec: pkgSpec, - Status: fission.PackageStatus{ - BuildStatus: pkgStatus, - }, - } - - if len(specFile) > 0 { - err := specSave(*pkg, specFile) - checkErr(err, "save package spec") - return &pkg.Metadata - } else { - pkgMetadata, err := client.PackageCreate(pkg) - checkErr(err, "create package") - return pkgMetadata - } -} - -func getContents(filePath string) []byte { - var code []byte - var err error - - code, err = ioutil.ReadFile(filePath) - checkErr(err, fmt.Sprintf("read %v", filePath)) - return code -} - -func getTempDir() (string, error) { - tmpDir := uuid.NewV4().String() - tmpPath := filepath.Join(os.TempDir(), tmpDir) - err := os.Mkdir(tmpPath, 0744) - return tmpPath, err -} - -func writeArchiveToFile(fileName string, reader io.Reader) error { - tmpDir, err := getTempDir() - if err != nil { - return err - } - - path := filepath.Join(tmpDir, fileName+".tmp") - w, err := os.Create(path) - if err != nil { - return err - } - _, err = io.Copy(w, reader) - if err != nil { - return err - } - err = os.Chmod(path, 0644) - if err != nil { - return err - } - - err = os.Rename(path, fileName) - if err != nil { - return err - } - - return nil -} - -// downloadToTempFile fetches archive file from arbitrary url -// and write it to temp file for further usage -func downloadToTempFile(fileUrl string) string { - reader, err := downloadURL(fileUrl) - defer reader.Close() - checkErr(err, fmt.Sprintf("download from url: %v", fileUrl)) - - tmpDir, err := getTempDir() - checkErr(err, "create temp directory") - - tmpFilename := uuid.NewV4().String() - destination := filepath.Join(tmpDir, tmpFilename) - err = os.Mkdir(tmpDir, 0744) - checkErr(err, "create temp directory") - - err = writeArchiveToFile(destination, reader) - checkErr(err, "write archive to file") - - return destination -} - -// downloadURL downloads file from given url -func downloadURL(fileUrl string) (io.ReadCloser, error) { - resp, err := http.Get(fileUrl) - if err != nil { - return nil, err - } - if resp.StatusCode != http.StatusOK { - return nil, fmt.Errorf("%v - HTTP response returned non 200 status", resp.StatusCode) - } - return resp.Body, nil -} - -// make a kubernetes compliant name out of an arbitrary string -func kubifyName(old string) string { - // Kubernetes maximum name length (for some names; others can be 253 chars) - maxLen := 63 - - newName := strings.ToLower(old) - - // replace disallowed chars with '-' - inv, err := regexp.Compile("[^-a-z0-9]") - checkErr(err, "compile regexp") - newName = string(inv.ReplaceAll([]byte(newName), []byte("-"))) - - // trim leading non-alphabetic - leadingnonalpha, err := regexp.Compile("^[^a-z]+") - checkErr(err, "compile regexp") - newName = string(leadingnonalpha.ReplaceAll([]byte(newName), []byte{})) - - // trim trailing - trailing, err := regexp.Compile("[^a-z0-9]+$") - checkErr(err, "compile regexp") - newName = string(trailing.ReplaceAll([]byte(newName), []byte{})) - - // truncate to length - if len(newName) > maxLen { - newName = newName[0:maxLen] - } - - // if we removed everything, call this thing "default". maybe - // we should generate a unique name... - if len(newName) == 0 { - newName = "default" - } - - return newName -} diff --git a/fission/environment.go b/fission/environment.go index dc79fe27..72904dfb 100644 --- a/fission/environment.go +++ b/fission/environment.go @@ -22,6 +22,7 @@ import ( "strconv" "text/tabwriter" + "github.com/fission/fission/fission/util" "github.com/urfave/cli" "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/resource" @@ -48,7 +49,7 @@ func getFunctionsByEnvironment(client *client.Client, envName, envNamespace stri } func envCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) envName := c.String("name") if len(envName) == 0 { @@ -129,19 +130,19 @@ func envCreate(c *cli.Context) error { if c.Bool("spec") { specFile := fmt.Sprintf("env-%v.yaml", envName) err := specSave(*env, specFile) - checkErr(err, "create environment spec") + util.CheckErr(err, "create environment spec") return nil } _, err = client.EnvironmentCreate(env) - checkErr(err, "create environment") + util.CheckErr(err, "create environment") fmt.Printf("environment '%v' created\n", envName) return err } func envGet(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) envName := c.String("name") if len(envName) == 0 { @@ -154,7 +155,7 @@ func envGet(c *cli.Context) error { Namespace: envNamespace, } env, err := client.EnvironmentGet(m) - checkErr(err, "get environment") + util.CheckErr(err, "get environment") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) fmt.Fprintf(w, "%v\t%v\t%v\n", "NAME", "UID", "IMAGE") @@ -165,7 +166,7 @@ func envGet(c *cli.Context) error { } func envUpdate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) envName := c.String("name") if len(envName) == 0 { @@ -186,7 +187,7 @@ func envUpdate(c *cli.Context) error { Name: envName, Namespace: envNamespace, }) - checkErr(err, "find environment") + util.CheckErr(err, "find environment") if len(envImg) > 0 { env.Spec.Runtime.Image = envImg @@ -218,14 +219,14 @@ func envUpdate(c *cli.Context) error { env.Spec.AllowAccessToExternalNetwork = envExternalNetwork _, err = client.EnvironmentUpdate(env) - checkErr(err, "update environment") + util.CheckErr(err, "update environment") fmt.Printf("environment '%v' updated\n", envName) return nil } func envDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) envName := c.String("name") if len(envName) == 0 { @@ -238,18 +239,18 @@ func envDelete(c *cli.Context) error { Namespace: envNamespace, } err := client.EnvironmentDelete(m) - checkErr(err, "delete environment") + util.CheckErr(err, "delete environment") fmt.Printf("environment '%v' deleted\n", envName) return nil } func envList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) envNamespace := c.String("envNamespace") envs, err := client.EnvironmentList(envNamespace) - checkErr(err, "list environments") + util.CheckErr(err, "list environments") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\t%v\t%v\t%v\t%v\t%v\n", "NAME", "UID", "IMAGE", "POOLSIZE", "MINCPU", "MAXCPU", "MINMEMORY", "MAXMEMORY", "EXTNET", "GRACETIME") diff --git a/fission/function.go b/fission/function.go index 59600b60..1530070d 100644 --- a/fission/function.go +++ b/fission/function.go @@ -38,7 +38,7 @@ import ( "github.com/fission/fission/crd" "github.com/fission/fission/fission/log" "github.com/fission/fission/fission/logdb" - "github.com/fission/fission/fission/portforward" + "github.com/fission/fission/fission/util" ) func printPodLogs(c *cli.Context) error { @@ -47,16 +47,16 @@ func printPodLogs(c *cli.Context) error { log.Fatal("Need --name argument.") } - queryURL, err := url.Parse(getServerUrl()) - checkErr(err, "parse the base URL") + queryURL, err := url.Parse(util.GetServerUrl()) + util.CheckErr(err, "parse the base URL") queryURL.Path = fmt.Sprintf("/proxy/logs/%s", fnName) req, err := http.NewRequest("POST", queryURL.String(), nil) - checkErr(err, "create logs request") + util.CheckErr(err, "create logs request") httpClient := http.Client{} resp, err := httpClient.Do(req) - checkErr(err, "execute get logs request") + util.CheckErr(err, "execute get logs request") defer resp.Body.Close() if resp.StatusCode != http.StatusOK { @@ -64,7 +64,7 @@ func printPodLogs(c *cli.Context) error { } body, err := ioutil.ReadAll(resp.Body) - checkErr(err, "read the response body") + util.CheckErr(err, "read the response body") fmt.Println(string(body)) return nil } @@ -120,7 +120,7 @@ func getTargetCPU(c *cli.Context) int { // From this change onwards, we mandate that a function should reference a secret, config map and package in its own ns func fnCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) fnNamespace := c.String("fnNamespace") envNamespace := c.String("envNamespace") @@ -140,7 +140,7 @@ func fnCreate(c *cli.Context) error { // check for unique function names within a namespace fnList, err := client.FunctionList(fnNamespace) - checkErr(err, "get function list") + util.CheckErr(err, "get function list") // check function existence before creating package for _, fn := range fnList { if fn.Metadata.Name == fnName { @@ -162,7 +162,7 @@ func fnCreate(c *cli.Context) error { Namespace: fnNamespace, Name: pkgName, }) - checkErr(err, fmt.Sprintf("read package in '%v' in Namespace: %s. Package needs to be present in the same namespace as function", pkgName, fnNamespace)) + util.CheckErr(err, fmt.Sprintf("read package in '%v' in Namespace: %s. Package needs to be present in the same namespace as function", pkgName, fnNamespace)) pkgMetadata = &pkg.Metadata envName = pkg.Spec.Environment.Name envNamespace = pkg.Spec.Environment.Namespace @@ -182,7 +182,7 @@ func fnCreate(c *cli.Context) error { if e, ok := err.(fission.Error); ok && e.Code == fission.ErrorNotFound { fmt.Printf("Environment \"%v\" does not exist. Please create the environment before executing the function. \nFor example: `fission env create --name %v --envns %v --image `\n", envName, envName, envNamespace) } else { - checkErr(err, "retrieve environment information") + util.CheckErr(err, "retrieve environment information") } } @@ -274,13 +274,13 @@ func fnCreate(c *cli.Context) error { // if we're writing a spec, don't create the function if spec { err = specSave(*function, specFile) - checkErr(err, "create function spec") + util.CheckErr(err, "create function spec") return nil } _, err = client.FunctionCreate(function) - checkErr(err, "create function") + util.CheckErr(err, "create function") fmt.Printf("function '%v' created\n", fnName) @@ -295,7 +295,7 @@ func fnCreate(c *cli.Context) error { method := c.String("method") if len(method) == 0 { - method = "GET" + method = http.MethodGet } triggerName := uuid.NewV4().String() ht := &crd.HTTPTrigger{ @@ -313,14 +313,14 @@ func fnCreate(c *cli.Context) error { }, } _, err = client.HTTPTriggerCreate(ht) - checkErr(err, "create HTTP trigger") + util.CheckErr(err, "create HTTP trigger") fmt.Printf("route created: %v %v -> %v\n", method, triggerUrl, fnName) return err } func fnGet(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) fnName := c.String("name") if len(fnName) == 0 { @@ -332,20 +332,20 @@ func fnGet(c *cli.Context) error { Namespace: fnNamespace, } fn, err := client.FunctionGet(m) - checkErr(err, "get function") + util.CheckErr(err, "get function") pkg, err := client.PackageGet(&metav1.ObjectMeta{ Name: fn.Spec.Package.PackageRef.Name, Namespace: fn.Spec.Package.PackageRef.Namespace, }) - checkErr(err, "get package") + util.CheckErr(err, "get package") os.Stdout.Write(pkg.Spec.Deployment.Literal) return err } func fnGetMeta(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) fnName := c.String("name") if len(fnName) == 0 { @@ -359,7 +359,7 @@ func fnGetMeta(c *cli.Context) error { } f, err := client.FunctionGet(m) - checkErr(err, "get function") + util.CheckErr(err, "get function") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) fmt.Fprintf(w, "%v\t%v\t%v\n", "NAME", "UID", "ENV") @@ -370,7 +370,7 @@ func fnGetMeta(c *cli.Context) error { } func fnUpdate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) if len(c.String("package")) > 0 { log.Fatal("--package is deprecated, please use --deploy instead.") @@ -390,7 +390,7 @@ func fnUpdate(c *cli.Context) error { Name: fnName, Namespace: fnNamespace, }) - checkErr(err, fmt.Sprintf("read function '%v'", fnName)) + util.CheckErr(err, fmt.Sprintf("read function '%v'", fnName)) envName := c.String("env") envNamespace := c.String("envNamespace") @@ -483,20 +483,20 @@ func fnUpdate(c *cli.Context) error { Namespace: fnNamespace, Name: pkgName, }) - checkErr(err, fmt.Sprintf("read package '%v.%v'. Pkg should be present in the same ns as the function", pkgName, fnNamespace)) + util.CheckErr(err, fmt.Sprintf("read package '%v.%v'. Pkg should be present in the same ns as the function", pkgName, fnNamespace)) pkgMetadata := &pkg.Metadata if len(deployArchiveName) != 0 || len(srcArchiveName) != 0 || len(buildcmd) != 0 || len(envName) != 0 || len(envNamespace) != 0 { fnList, err := getFunctionsByPackage(client, pkg.Metadata.Name, pkg.Metadata.Namespace) - checkErr(err, "get function list") + util.CheckErr(err, "get function list") if !force && len(fnList) > 1 { log.Fatal("Package is used by multiple functions, use --force to force update") } pkgMetadata, err = updatePackage(client, pkg, envName, envNamespace, srcArchiveName, deployArchiveName, buildcmd, false) - checkErr(err, fmt.Sprintf("update package '%v'", pkgName)) + util.CheckErr(err, fmt.Sprintf("update package '%v'", pkgName)) fmt.Printf("package '%v' updated\n", pkgMetadata.GetName()) @@ -506,7 +506,7 @@ func fnUpdate(c *cli.Context) error { if fn.Metadata.Name != fnName { fn.Spec.Package.PackageRef.ResourceVersion = pkgMetadata.ResourceVersion _, err := client.FunctionUpdate(&fn) - checkErr(err, "update function") + util.CheckErr(err, "update function") } } } @@ -570,14 +570,14 @@ func fnUpdate(c *cli.Context) error { } _, err = client.FunctionUpdate(function) - checkErr(err, "update function") + util.CheckErr(err, "update function") fmt.Printf("function '%v' updated\n", fnName) return err } func fnDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) fnName := c.String("name") if len(fnName) == 0 { @@ -591,18 +591,18 @@ func fnDelete(c *cli.Context) error { } err := client.FunctionDelete(m) - checkErr(err, fmt.Sprintf("delete function '%v'", fnName)) + util.CheckErr(err, fmt.Sprintf("delete function '%v'", fnName)) fmt.Printf("function '%v' deleted\n", fnName) return err } func fnList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) ns := c.String("fnNamespace") fns, err := client.FunctionList(ns) - checkErr(err, "list functions") + util.CheckErr(err, "list functions") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) @@ -628,7 +628,7 @@ func fnList(c *cli.Context) error { func fnLogs(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) fnName := c.String("name") if len(fnName) == 0 { @@ -653,10 +653,10 @@ func fnLogs(c *cli.Context) error { } f, err := client.FunctionGet(m) - checkErr(err, "get function") + util.CheckErr(err, "get function") // request the controller to establish a proxy server to the database. - logDB, err := logdb.GetLogDB(dbType, getServerUrl()) + logDB, err := logdb.GetLogDB(dbType, util.GetServerUrl()) if err != nil { log.Fatal("failed to connect log database") } @@ -718,8 +718,8 @@ func fnTest(c *cli.Context) error { routerURL := os.Getenv("FISSION_ROUTER") if len(routerURL) == 0 { // Portforward to the fission router - localRouterPort := portforward.Setup(getKubeConfigPath(), - getFissionNamespace(), "application=fission-router") + localRouterPort := util.SetupPortForward(util.GetKubeConfigPath(), + util.GetFissionNamespace(), "application=fission-router") routerURL = "127.0.0.1:" + localRouterPort } else { routerURL = strings.TrimPrefix(routerURL, "http://") @@ -757,14 +757,14 @@ func fnTest(c *cli.Context) error { resp := httpRequest(c.String("method"), functionUrl.String(), c.String("body"), c.StringSlice("header")) if resp.StatusCode < 400 { body, err := ioutil.ReadAll(resp.Body) - checkErr(err, "Function test") + util.CheckErr(err, "Function test") fmt.Print(string(body)) defer resp.Body.Close() return nil } body, err := ioutil.ReadAll(resp.Body) - checkErr(err, "read log response from pod") + util.CheckErr(err, "read log response from pod") fmt.Printf("Error calling function %s: %d %s", fnName, resp.StatusCode, string(body)) defer resp.Body.Close() err = printPodLogs(c) @@ -774,3 +774,34 @@ func fnTest(c *cli.Context) error { return nil } + +func httpRequest(method, url, body string, headers []string) *http.Response { + if method == "" { + method = "GET" + } + + if method != http.MethodGet && + method != http.MethodDelete && + method != http.MethodPost && + method != http.MethodPut && + method != http.MethodOptions { + log.Fatal(fmt.Sprintf("Invalid HTTP method '%s'.", method)) + } + + req, err := http.NewRequest(method, url, strings.NewReader(body)) + util.CheckErr(err, "create HTTP request") + + for _, header := range headers { + headerKeyValue := strings.SplitN(header, ":", 2) + if len(headerKeyValue) != 2 { + log.Fatal("Failed to create request without appropriate headers") + } + req.Header.Set(headerKeyValue[0], headerKeyValue[1]) + } + + client := &http.Client{} + resp, err := client.Do(req) + util.CheckErr(err, "execute HTTP request") + + return resp +} diff --git a/fission/httptrigger.go b/fission/httptrigger.go index cd88951a..2d12d75f 100644 --- a/fission/httptrigger.go +++ b/fission/httptrigger.go @@ -23,6 +23,7 @@ import ( "strings" "text/tabwriter" + "github.com/fission/fission/fission/util" "github.com/satori/go.uuid" "github.com/urfave/cli" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -72,7 +73,7 @@ func checkFunctionExistence(fissionClient *client.Client, fnName string, fnNames } func htCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) fnName := c.String("function") if len(fnName) == 0 { @@ -125,12 +126,12 @@ func htCreate(c *cli.Context) error { if c.Bool("spec") { specFile := fmt.Sprintf("route-%v.yaml", triggerName) err := specSave(*ht, specFile) - checkErr(err, "create HTTP trigger spec") + util.CheckErr(err, "create HTTP trigger spec") return nil } _, err := client.HTTPTriggerCreate(ht) - checkErr(err, "create HTTP trigger") + util.CheckErr(err, "create HTTP trigger") fmt.Printf("trigger '%v' created\n", triggerName) return err @@ -141,7 +142,7 @@ func htGet(c *cli.Context) error { } func htUpdate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) htName := c.String("name") if len(htName) == 0 { log.Fatal("Need name of trigger, use --name") @@ -152,7 +153,7 @@ func htUpdate(c *cli.Context) error { Name: htName, Namespace: triggerNamespace, }) - checkErr(err, "get HTTP trigger") + util.CheckErr(err, "get HTTP trigger") if c.IsSet("function") { ht.Spec.FunctionReference.Name = c.String("function") @@ -168,14 +169,14 @@ func htUpdate(c *cli.Context) error { } _, err = client.HTTPTriggerUpdate(ht) - checkErr(err, "update HTTP trigger") + util.CheckErr(err, "update HTTP trigger") fmt.Printf("trigger '%v' updated\n", htName) return nil } func htDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) htName := c.String("name") if len(htName) == 0 { log.Fatal("Need name of trigger to delete, use --name") @@ -186,18 +187,18 @@ func htDelete(c *cli.Context) error { Name: htName, Namespace: triggerNamespace, }) - checkErr(err, "delete trigger") + util.CheckErr(err, "delete trigger") fmt.Printf("trigger '%v' deleted\n", htName) return nil } func htList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) triggerNamespace := c.String("triggerNamespace") hts, err := client.HTTPTriggerList(triggerNamespace) - checkErr(err, "list HTTP triggers") + util.CheckErr(err, "list HTTP triggers") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) diff --git a/fission/main.go b/fission/main.go index fbf7e5c6..8de29d73 100644 --- a/fission/main.go +++ b/fission/main.go @@ -19,57 +19,18 @@ package main import ( "fmt" "os" - "path/filepath" "strings" - "github.com/ghodss/yaml" "github.com/urfave/cli" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "github.com/fission/fission" "github.com/fission/fission/fission/log" "github.com/fission/fission/fission/plugin" - "github.com/fission/fission/fission/portforward" + "github.com/fission/fission/fission/support" + "github.com/fission/fission/fission/util" ) -func getFissionNamespace() string { - fissionNamespace := os.Getenv("FISSION_NAMESPACE") - return fissionNamespace -} - -func getKubeConfigPath() string { - kubeConfig := os.Getenv("KUBECONFIG") - if len(kubeConfig) == 0 { - home := os.Getenv("HOME") - kubeConfig = filepath.Join(home, ".kube", "config") - - if _, err := os.Stat(kubeConfig); os.IsNotExist(err) { - log.Fatal("Couldn't find kubeconfig file. " + - "Set the KUBECONFIG environment variable to your kubeconfig's path.") - } - } - return kubeConfig -} - -func getServerUrl() string { - return getApplicationUrl("application=fission-api") -} - -func getApplicationUrl(selector string) string { - var serverUrl string - // Use FISSION_URL env variable if set; otherwise, port-forward to controller. - fissionUrl := os.Getenv("FISSION_URL") - if len(fissionUrl) == 0 { - fissionNamespace := getFissionNamespace() - kubeConfig := getKubeConfigPath() - localPort := portforward.Setup(kubeConfig, fissionNamespace, "application=fission-api") - serverUrl = "http://127.0.0.1:" + localPort - } else { - serverUrl = fissionUrl - } - return serverUrl -} - func cliHook(c *cli.Context) error { log.Verbosity = c.Int("verbosity") log.Verbose(2, "Verbosity = 2") @@ -303,6 +264,13 @@ func main() { {Name: "helm", Usage: "Create a helm chart from the app specification", Flags: []cli.Flag{specDirFlag}, Action: specHelm, Hidden: true}, } + // support + supportOutputFlag := cli.StringFlag{Name: "output, o", Value: support.DEFAULT_OUTPUT_DIR, Usage: "Output directory to save dump archive/files"} + supportNoZipFlag := cli.BoolFlag{Name: "nozip", Usage: "Save dump information into multiple files instead of single zip file"} + supportSubCommands := []cli.Command{ + {Name: "dump", Usage: "Collect & dump all necessary for troubleshooting", Flags: []cli.Flag{supportOutputFlag, supportNoZipFlag}, Action: support.DumpInfo}, + } + app.Commands = []cli.Command{ {Name: "function", Aliases: []string{"fn"}, Usage: "Create, update and manage functions", Subcommands: fnSubcommands}, {Name: "httptrigger", Aliases: []string{"ht", "route"}, Usage: "Manage HTTP triggers (routes) for functions", Subcommands: htSubcommands}, @@ -316,6 +284,7 @@ func main() { {Name: "package", Aliases: []string{"pkg"}, Usage: "Manage packages", Subcommands: pkgSubCommands}, {Name: "spec", Aliases: []string{"specs"}, Usage: "Manage a declarative app specification", Subcommands: specSubCommands}, {Name: "upgrade", Aliases: []string{}, Usage: "Upgrade tool from fission v0.1", Subcommands: upgradeSubCommands}, + {Name: "support", Usage: "Collect an archive of diagnostic information for support", Subcommands: supportSubCommands}, cmdPlugin, } app.Before = cliHook @@ -361,40 +330,10 @@ To install it for your local Fission CLI: } } -// Versions is a container of versions of the client (and its plugins) and server (and its plugins). -type Versions struct { - Client map[string]fission.BuildMeta `json:"client"` - Server map[string]fission.BuildMeta `json:"server"` -} - func versionPrinter(_ *cli.Context) { - serverInfo, err := getClient(getServerUrl()).ServerInfo() - if err != nil { - log.Warn(fmt.Sprintf("Error getting Fission API version: %v", err)) - } - - // Fetch client versions - versions := Versions{ - Client: map[string]fission.BuildMeta{ - "fission/core": fission.BuildInfo(), - }, - } - for _, pmd := range plugin.FindAll() { - versions.Client[pmd.Name] = fission.BuildMeta{ - Version: pmd.Version, - } - } - - // Fetch server versions - versions.Server = map[string]fission.BuildMeta{ - "fission/core": serverInfo.Build, - } - // FUTURE: fetch versions of plugins server-side - bs, err := yaml.Marshal(versions) - if err != nil { - log.Fatal("Failed to format versions: " + err.Error()) - } - fmt.Print(string(bs)) + client := util.GetApiClient(util.GetServerUrl()) + ver := util.GetVersion(client) + fmt.Print(string(ver)) } var helpTemplate = `NAME: diff --git a/fission/mqtrigger.go b/fission/mqtrigger.go index 780f9ca5..a3192885 100644 --- a/fission/mqtrigger.go +++ b/fission/mqtrigger.go @@ -21,6 +21,7 @@ import ( "os" "text/tabwriter" + "github.com/fission/fission/fission/util" "github.com/satori/go.uuid" "github.com/urfave/cli" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -32,7 +33,7 @@ import ( ) func mqtCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) mqtName := c.String("name") if len(mqtName) == 0 { @@ -107,12 +108,12 @@ func mqtCreate(c *cli.Context) error { if c.Bool("spec") { specFile := fmt.Sprintf("mqtrigger-%v.yaml", mqtName) err := specSave(*mqt, specFile) - checkErr(err, "create message queue trigger spec") + util.CheckErr(err, "create message queue trigger spec") return nil } _, err := client.MessageQueueTriggerCreate(mqt) - checkErr(err, "create message queue trigger") + util.CheckErr(err, "create message queue trigger") fmt.Printf("trigger '%s' created\n", mqtName) return err @@ -123,7 +124,7 @@ func mqtGet(c *cli.Context) error { } func mqtUpdate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) mqtName := c.String("name") if len(mqtName) == 0 { log.Fatal("Need name of trigger, use --name") @@ -141,7 +142,7 @@ func mqtUpdate(c *cli.Context) error { Name: mqtName, Namespace: mqtNs, }) - checkErr(err, "get Time trigger") + util.CheckErr(err, "get Time trigger") // TODO : Find out if we can make a call to checkIfFunctionExists, in the same ns more importantly. @@ -178,14 +179,14 @@ func mqtUpdate(c *cli.Context) error { } _, err = client.MessageQueueTriggerUpdate(mqt) - checkErr(err, "update Time trigger") + util.CheckErr(err, "update Time trigger") fmt.Printf("trigger '%v' updated\n", mqtName) return nil } func mqtDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) mqtName := c.String("name") if len(mqtName) == 0 { log.Fatal("Need name of trigger to delete, use --name") @@ -196,18 +197,18 @@ func mqtDelete(c *cli.Context) error { Name: mqtName, Namespace: mqtNs, }) - checkErr(err, "delete trigger") + util.CheckErr(err, "delete trigger") fmt.Printf("trigger '%v' deleted\n", mqtName) return nil } func mqtList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) mqtNs := c.String("triggerns") mqts, err := client.MessageQueueTriggerList(c.String("mqtype"), mqtNs) - checkErr(err, "list message queue triggers") + util.CheckErr(err, "list message queue triggers") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) diff --git a/fission/package.go b/fission/package.go index 3f65c255..d404b556 100644 --- a/fission/package.go +++ b/fission/package.go @@ -18,13 +18,23 @@ package main import ( "bytes" + "crypto/sha256" + "encoding/hex" "fmt" "io" + "io/ioutil" + "net/http" "net/url" "os" + "path" + "path/filepath" "strings" "text/tabwriter" + "github.com/dchest/uniuri" + "github.com/fission/fission/fission/util" + storageSvcClient "github.com/fission/fission/storagesvc/client" + "github.com/satori/go.uuid" "github.com/urfave/cli" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -59,12 +69,12 @@ func downloadStoragesvcURL(client *client.Client, fileUrl string) io.ReadCloser fileDownloadUrl := strings.TrimSuffix(client.Url, "/") + "/proxy/storage/" + u.RequestURI() reader, err := downloadURL(fileDownloadUrl) - checkErr(err, fmt.Sprintf("download from storage service url: %v", fileUrl)) + util.CheckErr(err, fmt.Sprintf("download from storage service url: %v", fileUrl)) return reader } func pkgCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) pkgNamespace := c.String("pkgNamespace") envName := c.String("env") @@ -87,7 +97,7 @@ func pkgCreate(c *cli.Context) error { } func pkgUpdate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) pkgName := c.String("name") if len(pkgName) == 0 { @@ -115,7 +125,7 @@ func pkgUpdate(c *cli.Context) error { Namespace: pkgNamespace, Name: pkgName, }) - checkErr(err, "get package") + util.CheckErr(err, "get package") // if the new env specified is the same as the old one, no need to update package // same is true for all update parameters, but, for now, we dont check all of them - because, its ok to @@ -129,7 +139,7 @@ func pkgUpdate(c *cli.Context) error { } fnList, err := getFunctionsByPackage(client, pkg.Metadata.Name, pkg.Metadata.Namespace) - checkErr(err, "get function list") + util.CheckErr(err, "get function list") if !force && len(fnList) > 1 { log.Fatal("Package is used by multiple functions, use --force to force update") @@ -138,14 +148,14 @@ func pkgUpdate(c *cli.Context) error { newPkgMeta, err := updatePackage(client, pkg, envName, envNamespace, srcArchiveName, deployArchiveName, buildcmd, false) if err != nil { - checkErr(err, "update package") + util.CheckErr(err, "update package") } // update resource version of package reference of functions that shared the same package for _, fn := range fnList { fn.Spec.Package.PackageRef.ResourceVersion = newPkgMeta.ResourceVersion _, err := client.FunctionUpdate(&fn) - checkErr(err, "update function") + util.CheckErr(err, "update function") } fmt.Printf("Package '%v' updated\n", newPkgMeta.GetName()) @@ -197,13 +207,13 @@ func updatePackage(client *client.Client, pkg *crd.Package, envName, envNamespac } newPkgMeta, err := client.PackageUpdate(pkg) - checkErr(err, "update package") + util.CheckErr(err, "update package") return newPkgMeta, err } func pkgSourceGet(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) pkgName := c.String("name") if len(pkgName) == 0 { @@ -240,7 +250,7 @@ func pkgSourceGet(c *cli.Context) error { } func pkgDeployGet(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) pkgName := c.String("name") if len(pkgName) == 0 { @@ -277,7 +287,7 @@ func pkgDeployGet(c *cli.Context) error { } func pkgInfo(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) pkgName := c.String("name") if len(pkgName) == 0 { @@ -304,7 +314,7 @@ func pkgInfo(c *cli.Context) error { } func pkgList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) // option for the user to list all orphan packages (not referenced by any function) listOrphans := c.Bool("orphan") pkgNamespace := c.String("pkgNamespace") @@ -319,7 +329,7 @@ func pkgList(c *cli.Context) error { if listOrphans { for _, pkg := range pkgList { fnList, err := getFunctionsByPackage(client, pkg.Metadata.Name, pkg.Metadata.Namespace) - checkErr(err, fmt.Sprintf("get functions sharing package %s", pkg.Metadata.Name)) + util.CheckErr(err, fmt.Sprintf("get functions sharing package %s", pkg.Metadata.Name)) if len(fnList) == 0 { fmt.Fprintf(w, "%v\t%v\t%v\n", pkg.Metadata.Name, pkg.Status.BuildStatus, pkg.Spec.Environment.Name) } @@ -345,7 +355,7 @@ func deleteOrphanPkgs(client *client.Client, pkgNamespace string) error { // range through all packages and find out the ones not referenced by any function for _, pkg := range pkgList { fnList, err := getFunctionsByPackage(client, pkg.Metadata.Name, pkgNamespace) - checkErr(err, fmt.Sprintf("get functions sharing package %s", pkg.Metadata.Name)) + util.CheckErr(err, fmt.Sprintf("get functions sharing package %s", pkg.Metadata.Name)) if len(fnList) == 0 { err = deletePackage(client, pkg.Metadata.Name, pkgNamespace) if err != nil { @@ -364,7 +374,7 @@ func deletePackage(client *client.Client, pkgName string, pkgNamespace string) e } func pkgDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) pkgName := c.String("name") pkgNamespace := c.String("pkgNamespace") @@ -386,7 +396,7 @@ func pkgDelete(c *cli.Context) error { Namespace: pkgNamespace, Name: pkgName, }) - checkErr(err, "find package") + util.CheckErr(err, "find package") fnList, err := getFunctionsByPackage(client, pkgName, pkgNamespace) @@ -402,7 +412,7 @@ func pkgDelete(c *cli.Context) error { fmt.Printf("Package '%v' deleted\n", pkgName) } else { err := deleteOrphanPkgs(client, pkgNamespace) - checkErr(err, "error deleting orphan packages") + util.CheckErr(err, "error deleting orphan packages") fmt.Println("Orphan packages deleted") } @@ -410,7 +420,7 @@ func pkgDelete(c *cli.Context) error { } func pkgRebuild(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) pkgName := c.String("name") if len(pkgName) == 0 { @@ -422,7 +432,7 @@ func pkgRebuild(c *cli.Context) error { Name: pkgName, Namespace: pkgNamespace, }) - checkErr(err, "find package") + util.CheckErr(err, "find package") if pkg.Status.BuildStatus != fission.BuildStatusFailed { log.Fatal(fmt.Sprintf("Package %v is not in %v state.", @@ -430,9 +440,215 @@ func pkgRebuild(c *cli.Context) error { } _, err = updatePackage(client, pkg, "", "", "", "", "", true) - checkErr(err, "update package") + util.CheckErr(err, "update package") fmt.Printf("Retrying build for pkg %v. Use \"fission pkg info --name %v\" to view status.\n", pkg.Metadata.Name, pkg.Metadata.Name) return nil } + +func fileSize(filePath string) int64 { + info, err := os.Stat(filePath) + util.CheckErr(err, fmt.Sprintf("stat %v", filePath)) + return info.Size() +} + +func fileChecksum(fileName string) (*fission.Checksum, error) { + f, err := os.Open(fileName) + if err != nil { + return nil, fmt.Errorf("failed to open file %v: %v", fileName, err) + } + defer f.Close() + + h := sha256.New() + _, err = io.Copy(h, f) + if err != nil { + return nil, fmt.Errorf("failed to calculate checksum for %v", fileName) + } + + return &fission.Checksum{ + Type: fission.ChecksumTypeSHA256, + Sum: hex.EncodeToString(h.Sum(nil)), + }, nil +} + +// upload a file and return a fission.Archive +func createArchive(client *client.Client, fileName string, specFile string) *fission.Archive { + var archive fission.Archive + + // fetch archive from arbitrary url if fileName is a url + if strings.HasPrefix(fileName, "http://") || strings.HasPrefix(fileName, "https://") { + fileName = downloadToTempFile(fileName) + } + + if len(specFile) > 0 { + // create an ArchiveUploadSpec and reference it from the archive + aus := &ArchiveUploadSpec{ + Name: util.KubifyName(path.Base(fileName)), + IncludeGlobs: []string{fileName}, + } + // save the uploadspec + err := specSave(*aus, specFile) + util.CheckErr(err, fmt.Sprintf("write spec file %v", specFile)) + // create the archive + ar := &fission.Archive{ + Type: fission.ArchiveTypeUrl, + URL: fmt.Sprintf("%v%v", ARCHIVE_URL_PREFIX, aus.Name), + } + return ar + } + + if fileSize(fileName) < fission.ArchiveLiteralSizeLimit { + contents := getContents(fileName) + archive.Type = fission.ArchiveTypeLiteral + archive.Literal = contents + } else { + u := strings.TrimSuffix(client.Url, "/") + "/proxy/storage" + ssClient := storageSvcClient.MakeClient(u) + + // TODO add a progress bar + id, err := ssClient.Upload(fileName, nil) + util.CheckErr(err, fmt.Sprintf("upload file %v", fileName)) + + storageSvc, err := client.GetSvcURL("application=fission-storage") + storageSvcURL := "http://" + storageSvc + util.CheckErr(err, "get fission storage service name") + + // We make a new client with actual URL of Storage service so that the URL is not + // pointing to 127.0.0.1 i.e. proxy. DON'T reuse previous ssClient + pkgClient := storageSvcClient.MakeClient(storageSvcURL) + archiveURL := pkgClient.GetUrl(id) + + archive.Type = fission.ArchiveTypeUrl + archive.URL = archiveURL + + csum, err := fileChecksum(fileName) + util.CheckErr(err, fmt.Sprintf("calculate checksum for file %v", fileName)) + + archive.Checksum = *csum + } + return &archive +} + +func createPackage(client *client.Client, pkgNamespace, envName, envNamespace, srcArchiveName, deployArchiveName, buildcmd string, specFile string) *metav1.ObjectMeta { + pkgSpec := fission.PackageSpec{ + Environment: fission.EnvironmentReference{ + Namespace: envNamespace, + Name: envName, + }, + } + var pkgStatus fission.BuildStatus = fission.BuildStatusSucceeded + + var pkgName string + if len(deployArchiveName) > 0 { + if len(specFile) > 0 { // we should do this in all cases, i think + pkgStatus = fission.BuildStatusNone + } + pkgSpec.Deployment = *createArchive(client, deployArchiveName, specFile) + pkgName = util.KubifyName(fmt.Sprintf("%v-%v", path.Base(deployArchiveName), uniuri.NewLen(4))) + } + if len(srcArchiveName) > 0 { + pkgSpec.Source = *createArchive(client, srcArchiveName, specFile) + // set pending status to package + pkgStatus = fission.BuildStatusPending + pkgName = util.KubifyName(fmt.Sprintf("%v-%v", path.Base(srcArchiveName), uniuri.NewLen(4))) + } + + if len(buildcmd) > 0 { + pkgSpec.BuildCommand = buildcmd + } + + if len(pkgName) == 0 { + pkgName = strings.ToLower(uuid.NewV4().String()) + } + pkg := &crd.Package{ + Metadata: metav1.ObjectMeta{ + Name: pkgName, + Namespace: pkgNamespace, + }, + Spec: pkgSpec, + Status: fission.PackageStatus{ + BuildStatus: pkgStatus, + }, + } + + if len(specFile) > 0 { + err := specSave(*pkg, specFile) + util.CheckErr(err, "save package spec") + return &pkg.Metadata + } else { + pkgMetadata, err := client.PackageCreate(pkg) + util.CheckErr(err, "create package") + return pkgMetadata + } +} + +func getContents(filePath string) []byte { + var code []byte + var err error + + code, err = ioutil.ReadFile(filePath) + util.CheckErr(err, fmt.Sprintf("read %v", filePath)) + return code +} + +func writeArchiveToFile(fileName string, reader io.Reader) error { + tmpDir, err := fission.GetTempDir() + if err != nil { + return err + } + + path := filepath.Join(tmpDir, fileName+".tmp") + w, err := os.Create(path) + if err != nil { + return err + } + _, err = io.Copy(w, reader) + if err != nil { + return err + } + err = os.Chmod(path, 0644) + if err != nil { + return err + } + + err = os.Rename(path, fileName) + if err != nil { + return err + } + + return nil +} + +// downloadToTempFile fetches archive file from arbitrary url +// and write it to temp file for further usage +func downloadToTempFile(fileUrl string) string { + reader, err := downloadURL(fileUrl) + defer reader.Close() + util.CheckErr(err, fmt.Sprintf("download from url: %v", fileUrl)) + + tmpDir, err := fission.GetTempDir() + util.CheckErr(err, "create temp directory") + + tmpFilename := uuid.NewV4().String() + destination := filepath.Join(tmpDir, tmpFilename) + err = os.Mkdir(tmpDir, 0744) + util.CheckErr(err, "create temp directory") + + err = writeArchiveToFile(destination, reader) + util.CheckErr(err, "write archive to file") + + return destination +} + +// downloadURL downloads file from given url +func downloadURL(fileUrl string) (io.ReadCloser, error) { + resp, err := http.Get(fileUrl) + if err != nil { + return nil, err + } + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("%v - HTTP response returned non 200 status", resp.StatusCode) + } + return resp.Body, nil +} diff --git a/fission/recorder.go b/fission/recorder.go index be5d3f01..f439ea2d 100644 --- a/fission/recorder.go +++ b/fission/recorder.go @@ -29,10 +29,11 @@ import ( "github.com/fission/fission" "github.com/fission/fission/crd" + "github.com/fission/fission/fission/util" ) func recorderCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) recName := c.String("name") if len(recName) == 0 { @@ -81,19 +82,19 @@ func recorderCreate(c *cli.Context) error { if c.Bool("spec") { specFile := fmt.Sprintf("recorder-%v.yaml", recName) err := specSave(*recorder, specFile) - checkErr(err, "create recorder spec") + util.CheckErr(err, "create recorder spec") return nil } _, err := client.RecorderCreate(recorder) - checkErr(err, "create recorder") + util.CheckErr(err, "create recorder") fmt.Printf("recorder '%s' created\n", recName) return err } func recorderGet(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) recName := c.String("name") @@ -102,7 +103,7 @@ func recorderGet(c *cli.Context) error { Namespace: "default", }) - checkErr(err, "get recorder") + util.CheckErr(err, "get recorder") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) @@ -116,7 +117,7 @@ func recorderGet(c *cli.Context) error { } func recorderUpdate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) recName := c.String("name") enable := c.Bool("enable") @@ -190,14 +191,14 @@ func recorderUpdate(c *cli.Context) error { } _, err = client.RecorderUpdate(recorder) - checkErr(err, "update recorder") + util.CheckErr(err, "update recorder") fmt.Printf("recorder '%v' updated\n", recName) return nil } func recorderDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) recName := c.String("name") @@ -212,17 +213,17 @@ func recorderDelete(c *cli.Context) error { Namespace: recNs, }) - checkErr(err, "delete recorder") + util.CheckErr(err, "delete recorder") fmt.Printf("recorder '%v' deleted\n", recName) return nil } func recorderList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) recorders, err := client.RecorderList("default") - checkErr(err, "list recorders") + util.CheckErr(err, "list recorders") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) diff --git a/fission/records.go b/fission/records.go index 2114b908..73205596 100644 --- a/fission/records.go +++ b/fission/records.go @@ -24,6 +24,7 @@ import ( "github.com/urfave/cli" + "github.com/fission/fission/fission/util" "github.com/fission/fission/redis/build/gen" ) @@ -59,15 +60,15 @@ func recordsView(c *cli.Context) error { return recordsByTime(from, to, verbosity, c) } err := recordsAll(verbosity, c) - checkErr(err, "view records") + util.CheckErr(err, "view records") return nil } func recordsAll(verbosity int, c *cli.Context) error { - fc := getClient(c.GlobalString("server")) + fc := util.GetApiClient(c.GlobalString("server")) records, err := fc.RecordsAll() - checkErr(err, "view records") + util.CheckErr(err, "view records") showRecords(records, verbosity) @@ -75,10 +76,10 @@ func recordsAll(verbosity int, c *cli.Context) error { } func recordsByTrigger(trigger string, verbosity int, c *cli.Context) error { - fc := getClient(c.GlobalString("server")) + fc := util.GetApiClient(c.GlobalString("server")) records, err := fc.RecordsByTrigger(trigger) - checkErr(err, "view records") + util.CheckErr(err, "view records") showRecords(records, verbosity) @@ -87,10 +88,10 @@ func recordsByTrigger(trigger string, verbosity int, c *cli.Context) error { // TODO: More accurate function name (function filter) func recordsByFunction(function string, verbosity int, c *cli.Context) error { - fc := getClient(c.GlobalString("server")) + fc := util.GetApiClient(c.GlobalString("server")) records, err := fc.RecordsByFunction(function) - checkErr(err, "view records") + util.CheckErr(err, "view records") showRecords(records, verbosity) @@ -98,10 +99,10 @@ func recordsByFunction(function string, verbosity int, c *cli.Context) error { } func recordsByTime(from string, to string, verbosity int, c *cli.Context) error { - fc := getClient(c.GlobalString("server")) + fc := util.GetApiClient(c.GlobalString("server")) records, err := fc.RecordsByTime(from, to) - checkErr(err, "view records") + util.CheckErr(err, "view records") showRecords(records, verbosity) diff --git a/fission/replay.go b/fission/replay.go index 3513e627..e6868188 100644 --- a/fission/replay.go +++ b/fission/replay.go @@ -22,11 +22,12 @@ import ( "os" "text/tabwriter" + "github.com/fission/fission/fission/util" "github.com/urfave/cli" ) func replay(c *cli.Context) error { - fc := getClient(c.GlobalString("server")) + fc := util.GetApiClient(c.GlobalString("server")) reqUID := c.String("reqUID") if len(reqUID) == 0 { @@ -34,7 +35,7 @@ func replay(c *cli.Context) error { } responses, err := fc.ReplayByReqUID(reqUID) - checkErr(err, "replay records") + util.CheckErr(err, "replay records") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) diff --git a/fission/spec.go b/fission/spec.go index aed26f9f..d677f89f 100644 --- a/fission/spec.go +++ b/fission/spec.go @@ -40,6 +40,7 @@ import ( "github.com/fission/fission/controller/client" "github.com/fission/fission/crd" "github.com/fission/fission/fission/log" + "github.com/fission/fission/fission/util" ) const SPEC_API_VERSION = "fission.io/v1" @@ -158,15 +159,15 @@ func specInit(c *cli.Context) error { if len(name) == 0 { // come up with a name using the current dir dir, err := filepath.Abs(".") - checkErr(err, "get current working directory") + util.CheckErr(err, "get current working directory") basename := filepath.Base(dir) - name = kubifyName(basename) + name = util.KubifyName(basename) } // Create spec dir fmt.Printf("Creating fission spec directory '%v'\n", specDir) err := os.MkdirAll(specDir, 0755) - checkErr(err, fmt.Sprintf("create spec directory '%v'", specDir)) + util.CheckErr(err, fmt.Sprintf("create spec directory '%v'", specDir)) // Add a bit of documentation to the spec dir here err = ioutil.WriteFile(filepath.Join(specDir, "README"), []byte(SPEC_README), 0644) @@ -188,7 +189,7 @@ func specInit(c *cli.Context) error { UID: uuid.NewV4().String(), } err = writeDeploymentConfig(specDir, &dc) - checkErr(err, "write deployment config") + util.CheckErr(err, "write deployment config") // Other possible things to do here: // - add example specs to the dir to make it easy to manually @@ -226,7 +227,7 @@ func specValidate(c *cli.Context) error { // this will error on parse errors and on duplicates specDir := getSpecDir(c) fr, err := readSpecs(specDir) - checkErr(err, "read specs") + util.CheckErr(err, "read specs") // this does the rest of the checks, like dangling refs err = fr.validate() @@ -631,7 +632,7 @@ func waitForFileWatcherToSettleDown(watcher *fsnotify.Watcher) error { // etc, while doing an apply, they will get a partially applied deployment. However, // they can retry their apply command once they're back online. func specApply(c *cli.Context) error { - fclient := getClient(c.GlobalString("server")) + fclient := util.GetApiClient(c.GlobalString("server")) specDir := getSpecDir(c) deleteResources := c.Bool("delete") @@ -649,35 +650,35 @@ func specApply(c *cli.Context) error { if watchResources { var err error watcher, err = fsnotify.NewWatcher() - checkErr(err, "create file watcher") + util.CheckErr(err, "create file watcher") // add watches rootDir := filepath.Clean(specDir + "/..") err = filepath.Walk(rootDir, func(path string, info os.FileInfo, err error) error { - checkErr(err, "scan project files") + util.CheckErr(err, "scan project files") if ignoreFile(path) { return nil } err = watcher.Add(path) - checkErr(err, fmt.Sprintf("watch path %v", path)) + util.CheckErr(err, fmt.Sprintf("watch path %v", path)) return nil }) - checkErr(err, "scan files to watch") + util.CheckErr(err, "scan files to watch") } for { // read all specs fr, err := readSpecs(specDir) - checkErr(err, "read specs") + util.CheckErr(err, "read specs") // validate err = fr.validate() - checkErr(err, "validate specs") + util.CheckErr(err, "validate specs") // make changes to the cluster based on the specs pkgMetas, as, err := applyResources(fclient, specDir, fr, deleteResources) - checkErr(err, "apply specs") + util.CheckErr(err, "apply specs") printApplyStatus(as) if watchResources || waitForBuild { @@ -718,11 +719,11 @@ func specApply(c *cli.Context) error { pkgWatchCancel() err = waitForFileWatcherToSettleDown(watcher) - checkErr(err, "watching files") + util.CheckErr(err, "watching files") break waitloop case err := <-watcher.Errors: - checkErr(err, "watching files") + util.CheckErr(err, "watching files") } } } @@ -775,14 +776,14 @@ func pluralize(num int, word string) string { // specDestroy destroys everything in the spec. func specDestroy(c *cli.Context) error { - fclient := getClient(c.GlobalString("server")) + fclient := util.GetApiClient(c.GlobalString("server")) // get specdir specDir := getSpecDir(c) // read everything fr, err := readSpecs(specDir) - checkErr(err, "read specs") + util.CheckErr(err, "read specs") // set desired state to nothing, but keep the UID so "apply" can find it emptyFr := FissionResources{} @@ -790,7 +791,7 @@ func specDestroy(c *cli.Context) error { // "apply" the empty state _, _, err = applyResources(fclient, specDir, &emptyFr, true) - checkErr(err, "delete resources") + util.CheckErr(err, "delete resources") return nil } diff --git a/fission/support/dump.go b/fission/support/dump.go new file mode 100644 index 00000000..85386b28 --- /dev/null +++ b/fission/support/dump.go @@ -0,0 +1,141 @@ +/* +Copyright 2018 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 support + +import ( + "fmt" + "os" + "path/filepath" + "sync" + "time" + + "github.com/pkg/errors" + "github.com/urfave/cli" + + "github.com/fission/fission" + "github.com/fission/fission/fission/support/resources" + "github.com/fission/fission/fission/util" +) + +const ( + DUMP_ARCHIVE_PREFIX = "fission-dump" + DEFAULT_OUTPUT_DIR = "fission-dump" +) + +func DumpInfo(c *cli.Context) error { + + fmt.Println("Start dumping process...") + + nozip := c.Bool("nozip") + outputDir := c.String("output") + + // check whether the dump directory exists. + _, err := os.Stat(outputDir) + if err != nil && os.IsNotExist(err) { + err = os.Mkdir(outputDir, 0755) + if err != nil { + panic(err) + } + } else if err != nil { + panic(errors.Wrap(err, "Error checking dump directory status")) + } + + outputDir, err = filepath.Abs(outputDir) + if err != nil { + panic(errors.Wrap(err, "Error creating dump directory for dumping files")) + } + + client := util.GetApiClient(util.GetServerUrl()) + _, k8sClient := util.GetKubernetesClient(util.GetKubeConfigPath()) + + ress := map[string]resources.Resource{ + // kubernetes info + "kubernetes-version": resources.NewKubernetesVersion(k8sClient), + "kubernetes-nodes": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesNode, ""), + + // fission info + "fission-version": resources.NewFissionVersion(client), + + // fission component logs & spec + "fission-components-svc-sepc": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesService, + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, redis, router, storagesvc, timer)"), + "fission-components-deployment-sepc": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesDeployment, + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, redis, router, storagesvc, timer)"), + "fission-components-pod-sepc": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesPod, + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, redis, router, storagesvc, timer)"), + "fission-components-pod-log": resources.NewKubernetesPodLogDumper(k8sClient, + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, redis, router, storagesvc, timer)"), + + // fission builder logs & spec + "fission-builder-svc-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesService, "owner=buildermgr"), + "fission-builder-deployment-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesDeployment, "owner=buildermgr"), + "fission-builder-pod-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesPod, "owner=buildermgr"), + "fission-builder-pod-log": resources.NewKubernetesPodLogDumper(k8sClient, "owner=buildermgr"), + + // fission function logs & spec + "fission-function-svc-sepc": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesService, "executorType=newdeploy"), + "fission-function-deployment-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesDeployment, "executorType in (poolmgr, newdeploy)"), + "fission-function-pod-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesPod, "executorType in (poolmgr, newdeploy)"), + "fission-function-pod-log": resources.NewKubernetesPodLogDumper(k8sClient, "executorType in (poolmgr, newdeploy)"), + + // CRD resources + "fission-crd-packages": resources.NewCrdDumper(client, resources.CrdPackage), + "fission-crd-environments": resources.NewCrdDumper(client, resources.CrdEnvironment), + "fission-crd-functions": resources.NewCrdDumper(client, resources.CrdFunction), + "fission-crd-httptriggers": resources.NewCrdDumper(client, resources.CrdHttpTrigger), + "fission-crd-kubewatchers": resources.NewCrdDumper(client, resources.CrdKubeWatcher), + "fission-crd-mqtriggers": resources.NewCrdDumper(client, resources.CrdMessageQueueTrigger), + "fission-crd-timetriggers": resources.NewCrdDumper(client, resources.CrdTimeTrigger), + } + + dumpName := fmt.Sprintf("%v_%v", DUMP_ARCHIVE_PREFIX, time.Now().Unix()) + dumpDir := filepath.Join(outputDir, dumpName) + + wg := &sync.WaitGroup{} + + for key, res := range ress { + dir := fmt.Sprintf("%v/%v/", dumpDir, key) + if _, err := os.Stat(dir); os.IsNotExist(err) { + err = os.MkdirAll(dir, 0755) + if err != nil { + panic(err) + } + } + go func(res resources.Resource, dir string) { + wg.Add(1) + defer wg.Done() + res.Dump(dir) + }(res, dir) + } + + wg.Wait() + + if !nozip { + defer os.Remove(dumpDir) + path := filepath.Join(outputDir, fmt.Sprintf("%v.zip", dumpName)) + _, err := fission.MakeArchive(path, dumpDir) + if err != nil { + fmt.Printf("Error creating archive for dump files: %v", err) + return nil + } + fmt.Printf("The archive dump file is %v\n", path) + } else { + fmt.Printf("The dump files are placed at %v\n", dumpDir) + } + + return nil +} diff --git a/fission/support/resources/crd.go b/fission/support/resources/crd.go new file mode 100644 index 00000000..b2286a68 --- /dev/null +++ b/fission/support/resources/crd.go @@ -0,0 +1,157 @@ +/* +Copyright 2018 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 resources + +import ( + "log" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + "github.com/fission/fission" + "github.com/fission/fission/controller/client" + "github.com/fission/fission/crd" +) + +const ( + CrdEnvironment = "Environment" + CrdFunction = "Function" + CrdPackage = "Packages" + + CrdHttpTrigger = "HTTPTrigger" + CrdKubeWatcher = "KubeWatcher" + CrdMessageQueueTrigger = "MessageQueue" + CrdTimeTrigger = "TimeTrigger" +) + +type CrdDumper struct { + client *client.Client + crdType string +} + +func NewCrdDumper(client *client.Client, crdType string) Resource { + return CrdDumper{client: client, crdType: crdType} +} + +func (res CrdDumper) Dump(dumpDir string) { + + switch res.crdType { + case CrdEnvironment: + items, err := res.client.EnvironmentList(metav1.NamespaceAll) + if err != nil { + log.Printf("Error getting %v list: %v", res.crdType, err) + return + } + + for _, item := range items { + f := getFileName(dumpDir, item.Metadata) + writeToFile(f, item) + } + + case CrdFunction: + items, err := res.client.FunctionList(metav1.NamespaceAll) + if err != nil { + log.Printf("Error getting %v list: %v", res.crdType, err) + return + } + + for _, item := range items { + f := getFileName(dumpDir, item.Metadata) + writeToFile(f, item) + } + + case CrdPackage: + items, err := res.client.PackageList(metav1.NamespaceAll) + if err != nil { + log.Printf("Error getting %v list: %v", res.crdType, err) + return + } + + for _, item := range items { + item = pkgClean(item) + f := getFileName(dumpDir, item.Metadata) + writeToFile(f, item) + } + + case CrdHttpTrigger: + items, err := res.client.HTTPTriggerList(metav1.NamespaceAll) + if err != nil { + log.Printf("Error getting %v list: %v", res.crdType, err) + return + } + + for _, item := range items { + f := getFileName(dumpDir, item.Metadata) + writeToFile(f, item) + } + + case CrdKubeWatcher: + items, err := res.client.WatchList(metav1.NamespaceAll) + if err != nil { + log.Printf("Error getting %v list: %v", res.crdType, err) + return + } + + for _, item := range items { + f := getFileName(dumpDir, item.Metadata) + writeToFile(f, item) + } + + case CrdMessageQueueTrigger: + var triggers []crd.MessageQueueTrigger + + for _, mqType := range []string{fission.MessageQueueTypeNats, fission.MessageQueueTypeASQ} { + l, err := res.client.MessageQueueTriggerList(mqType, metav1.NamespaceAll) + if err != nil { + log.Printf("Error getting %v list: %v", res.crdType, err) + break + } + triggers = append(triggers, l...) + } + + for _, item := range triggers { + f := getFileName(dumpDir, item.Metadata) + writeToFile(f, item) + } + + case CrdTimeTrigger: + items, err := res.client.TimeTriggerList(metav1.NamespaceAll) + if err != nil { + log.Printf("Error getting %v list: %v", res.crdType, err) + return + } + + for _, item := range items { + f := getFileName(dumpDir, item.Metadata) + writeToFile(f, item) + } + + default: + log.Printf("Unknown type: %v", res.crdType) + } +} + +func pkgClean(pkg crd.Package) crd.Package { + // mask the sensitive information + // use "-" as mask value to indicate the field wasn't empty + if pkg.Spec.Source.Literal != nil { + pkg.Spec.Source.Literal = []byte("-") + } + if pkg.Spec.Deployment.Literal != nil { + pkg.Spec.Deployment.Literal = []byte("-") + } + return pkg +} diff --git a/fission/support/resources/fissionversion.go b/fission/support/resources/fissionversion.go new file mode 100644 index 00000000..ae58e6bb --- /dev/null +++ b/fission/support/resources/fissionversion.go @@ -0,0 +1,40 @@ +/* +Copyright 2018 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 resources + +import ( + "fmt" + "path/filepath" + + "github.com/fission/fission/controller/client" + "github.com/fission/fission/fission/util" +) + +type FissionVersion struct { + client *client.Client + namespace string +} + +func NewFissionVersion(client *client.Client) Resource { + return FissionVersion{client: client} +} + +func (res FissionVersion) Dump(dumpDir string) { + ver := util.GetVersion(res.client) + file := filepath.Clean(fmt.Sprintf("%v/%v", dumpDir, "fission-version.txt")) + writeToFile(file, ver) +} diff --git a/fission/support/resources/kubernetes.go b/fission/support/resources/kubernetes.go new file mode 100644 index 00000000..6d76ef18 --- /dev/null +++ b/fission/support/resources/kubernetes.go @@ -0,0 +1,249 @@ +/* +Copyright 2018 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 resources + +import ( + "bufio" + "bytes" + "fmt" + "io" + "log" + "path/filepath" + "sync" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + + "github.com/fission/fission" +) + +const ( + KubernetesService = "Service" + KubernetesDeployment = "Deployment" + KubernetesPod = "Pod" + KubernetesHPA = "HPA" + KubernetesNode = "Node" +) + +// Kubernetes Version +type KubernetesVersion struct { + client *kubernetes.Clientset +} + +func NewKubernetesVersion(clientset *kubernetes.Clientset) Resource { + return KubernetesVersion{client: clientset} +} + +func (res KubernetesVersion) Dump(dumpDir string) { + serverVer, err := res.client.ServerVersion() + if err != nil { + log.Printf("Error setting up kubernetes client: %v", err) + return + } + + file := fmt.Sprintf("%v/%v", dumpDir, "kubernetes-version.txt") + writeToFile(file, serverVer) +} + +// Kubernetes Object Dumper +type KubernetesObjectDumper struct { + client *kubernetes.Clientset + objType string + selector string +} + +func NewKubernetesObjectDumper(clientset *kubernetes.Clientset, objType string, selector string) Resource { + return KubernetesObjectDumper{ + client: clientset, + objType: objType, + selector: selector, + } +} + +func (res KubernetesObjectDumper) Dump(dumpDir string) { + switch res.objType { + case KubernetesService: + objs, err := res.client.CoreV1().Services(metav1.NamespaceAll).List(metav1.ListOptions{LabelSelector: res.selector}) + if err != nil { + log.Printf("Error getting %v list with selector %v: %v", res.objType, res.selector, err) + return + } + + for _, item := range objs.Items { + item = serviceClean(item) + f := getFileName(dumpDir, item.ObjectMeta) + writeToFile(f, item) + } + + case KubernetesDeployment: + objs, err := res.client.AppsV1().Deployments(metav1.NamespaceAll).List(metav1.ListOptions{LabelSelector: res.selector}) + if err != nil { + log.Printf("Error getting %v list with selector %v: %v", res.objType, res.selector, err) + return + } + + for _, item := range objs.Items { + f := getFileName(dumpDir, item.ObjectMeta) + writeToFile(f, item) + } + + case KubernetesPod: + objs, err := res.client.CoreV1().Pods(metav1.NamespaceAll).List(metav1.ListOptions{LabelSelector: res.selector}) + if err != nil { + log.Printf("Error getting %v list with selector %v: %v", res.objType, res.selector, err) + return + } + + for _, item := range objs.Items { + f := getFileName(dumpDir, item.ObjectMeta) + writeToFile(f, item) + } + + case KubernetesHPA: + objs, err := res.client.AutoscalingV2beta1().HorizontalPodAutoscalers(metav1.NamespaceAll).List(metav1.ListOptions{LabelSelector: res.selector}) + if err != nil { + log.Printf("Error getting %v list with selector %v: %v", res.objType, res.selector, err) + return + } + + for _, item := range objs.Items { + f := getFileName(dumpDir, item.ObjectMeta) + writeToFile(f, item) + } + + case KubernetesNode: + objs, err := res.client.CoreV1().Nodes().List(metav1.ListOptions{LabelSelector: res.selector}) + if err != nil { + log.Printf("Error getting %v list with selector %v: %v", res.objType, res.selector, err) + return + } + + for _, item := range objs.Items { + item = nodeClean(item) + // Node doesn't have namespace value, use name here + f := filepath.Clean(fmt.Sprintf("%v/%v", dumpDir, item.Name)) + getFileName(dumpDir, item.ObjectMeta) + writeToFile(f, item) + } + + default: + log.Printf("Unknown type: %v", res.objType) + } +} + +// serviceClean remove sensitive data(e.g. public IP, external name) from service objects +func serviceClean(svc corev1.Service) corev1.Service { + svc.Spec.ExternalIPs = []string{} + svc.Spec.LoadBalancerIP = "" + svc.Spec.LoadBalancerSourceRanges = []string{} + svc.Spec.ExternalName = "" + svc.Status.LoadBalancer = corev1.LoadBalancerStatus{} + return svc +} + +func nodeClean(node corev1.Node) corev1.Node { + + var nodeAddresses []corev1.NodeAddress + for _, address := range node.Status.Addresses { + // use whitelist to filter the necessary information for debugging + if address.Type == "InternalIP" || address.Type == "Hostname" { + nodeAddresses = append(nodeAddresses, address) + } + } + node.Status.Addresses = nodeAddresses + + return node +} + +type KubernetesPodLogDumper struct { + client *kubernetes.Clientset + labelSelector string +} + +func NewKubernetesPodLogDumper(clientset *kubernetes.Clientset, selector string) Resource { + return KubernetesPodLogDumper{ + client: clientset, + labelSelector: selector, + } +} + +func (res KubernetesPodLogDumper) Dump(dumpDir string) { + l, err := res.client.CoreV1(). + Pods(metav1.NamespaceAll). + List(metav1.ListOptions{LabelSelector: res.labelSelector}) + if err != nil { + log.Printf("Error getting controller list: %v", err) + return + } + + wg := &sync.WaitGroup{} + + for _, p := range l.Items { + wg.Add(1) + + go func(pod corev1.Pod) { + defer wg.Done() + + if !fission.IsReadyPod(&pod) { + fmt.Printf("Pod %v is not in ready state, ignore it\n", pod.Name) + return + } + + // dump logs from each containers + for _, container := range append(pod.Spec.Containers, pod.Spec.InitContainers...) { + req := res.client.CoreV1().Pods(pod.Namespace). + GetLogs(pod.Name, &corev1.PodLogOptions{Container: container.Name}) + + stream, err := req.Stream() + if err != nil { + log.Printf("Error streaming logs for pod %v: %v", pod.Name, err) + return + } + + reader := bufio.NewReader(stream) + var buffer bytes.Buffer + + for { + line, _, err := reader.ReadLine() + if err != nil { + if err == io.EOF { + stream.Close() + break + } + log.Printf("Error reading logs from buffer: %v", err) + return + } + + _, err = buffer.WriteString(string(line) + "\n") + if err != nil { + log.Printf("Error writing bytes to buffer: %v", err) + } + } + + f := getFileName(dumpDir, pod.ObjectMeta) + writeToFile(f, buffer.String()) + + stream.Close() + } + }(p) + } + + wg.Wait() + + return +} diff --git a/fission/support/resources/resource.go b/fission/support/resources/resource.go new file mode 100644 index 00000000..f052354f --- /dev/null +++ b/fission/support/resources/resource.go @@ -0,0 +1,57 @@ +/* +Copyright 2018 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 resources + +import ( + "fmt" + "io/ioutil" + "log" + "path/filepath" + + "github.com/ghodss/yaml" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + "github.com/fission/fission" +) + +type Resource interface { + Dump(string) +} + +func getFileName(dumpdir string, meta metav1.ObjectMeta) string { + f := fmt.Sprintf("%v/%v_%v_%v.txt", dumpdir, meta.Namespace, meta.Name, meta.ResourceVersion) + return filepath.Clean(f) +} + +func writeToFile(file string, obj interface{}) { + bs, err := yaml.Marshal(obj) + if err != nil { + log.Printf("Error encoding object: %v", err) + return + } + + // Due to unknown reason, the kubernetes objectMeta fields contain + // empty byte and will fail os.Create/os.Openfile with error message + // "open invalid argument". To fix the problem, we need to + // remove the empty byte from string. + file = string(fission.RemoveZeroBytes([]byte(file))) + + err = ioutil.WriteFile(file, bs, 0644) + if err != nil { + log.Printf("Error writing file %v: %v", file, err) + } +} diff --git a/fission/timetrigger.go b/fission/timetrigger.go index 630bc8e2..6e31015a 100644 --- a/fission/timetrigger.go +++ b/fission/timetrigger.go @@ -31,6 +31,7 @@ import ( "github.com/fission/fission/controller/client" "github.com/fission/fission/crd" "github.com/fission/fission/fission/log" + "github.com/fission/fission/fission/util" ) func getAPITimeInfo(client *client.Client) time.Time { @@ -58,7 +59,7 @@ func getCronNextNActivationTime(cronSpec string, serverTime time.Time, round int } func ttCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) name := c.String("name") if len(name) == 0 { @@ -94,17 +95,17 @@ func ttCreate(c *cli.Context) error { if c.Bool("spec") { specFile := fmt.Sprintf("timetrigger-%v.yaml", name) err := specSave(*tt, specFile) - checkErr(err, "create time trigger spec") + util.CheckErr(err, "create time trigger spec") return nil } _, err := client.TimeTriggerCreate(tt) - checkErr(err, "create Time trigger") + util.CheckErr(err, "create Time trigger") fmt.Printf("trigger '%v' created\n", name) err = getCronNextNActivationTime(cronSpec, getAPITimeInfo(client), 1) - checkErr(err, "pass cron spec examination") + util.CheckErr(err, "pass cron spec examination") return err } @@ -114,7 +115,7 @@ func ttGet(c *cli.Context) error { } func ttUpdate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) ttName := c.String("name") if len(ttName) == 0 { log.Fatal("Need name of trigger, use --name") @@ -125,7 +126,7 @@ func ttUpdate(c *cli.Context) error { Name: ttName, Namespace: ttNs, }) - checkErr(err, "get time trigger") + util.CheckErr(err, "get time trigger") updated := false newCron := c.String("cron") @@ -148,18 +149,18 @@ func ttUpdate(c *cli.Context) error { } _, err = client.TimeTriggerUpdate(tt) - checkErr(err, "update Time trigger") + util.CheckErr(err, "update Time trigger") fmt.Printf("trigger '%v' updated\n", ttName) err = getCronNextNActivationTime(newCron, getAPITimeInfo(client), 1) - checkErr(err, "pass cron spec examination") + util.CheckErr(err, "pass cron spec examination") return nil } func ttDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) ttName := c.String("name") if len(ttName) == 0 { log.Fatal("Need name of trigger to delete, use --name") @@ -170,18 +171,18 @@ func ttDelete(c *cli.Context) error { Name: ttName, Namespace: ttNs, }) - checkErr(err, "delete trigger") + util.CheckErr(err, "delete trigger") fmt.Printf("trigger '%v' deleted\n", ttName) return nil } func ttList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) ttNs := c.String("triggerns") tts, err := client.TimeTriggerList(ttNs) - checkErr(err, "list Time triggers") + util.CheckErr(err, "list Time triggers") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) @@ -196,7 +197,7 @@ func ttList(c *cli.Context) error { } func ttTest(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) round := c.Int("round") cronSpec := c.String("cron") @@ -205,7 +206,7 @@ func ttTest(c *cli.Context) error { } err := getCronNextNActivationTime(cronSpec, getAPITimeInfo(client), round) - checkErr(err, "pass cron spec examination") + util.CheckErr(err, "pass cron spec examination") return nil } diff --git a/fission/upgrade.go b/fission/upgrade.go index b2ad8937..e879c027 100644 --- a/fission/upgrade.go +++ b/fission/upgrade.go @@ -17,6 +17,7 @@ import ( "github.com/fission/fission" "github.com/fission/fission/crd" "github.com/fission/fission/fission/log" + "github.com/fission/fission/fission/util" "github.com/fission/fission/v1" ) @@ -51,11 +52,11 @@ func getV1URL(serverUrl string) string { func get(url string) []byte { resp, err := http.Get(url) - checkErr(err, "get fission v0.1 state") + util.CheckErr(err, "get fission v0.1 state") defer resp.Body.Close() body, err := ioutil.ReadAll(resp.Body) - checkErr(err, "reading server response") + util.CheckErr(err, "reading server response") if resp.StatusCode != 200 { log.Fatal(fmt.Sprintf("Failed to fetch fission v0.1 state: %v", string(body))) @@ -70,7 +71,7 @@ func (nr *nameRemapper) trackName(old string) { maxLen := 63 ok, err := regexp.MatchString(kubeNameRegex, old) - checkErr(err, "match name regexp") + util.CheckErr(err, "match name regexp") if ok && len(old) < maxLen { // no rename nr.oldToNew[old] = old @@ -82,17 +83,17 @@ func (nr *nameRemapper) trackName(old string) { // remove disallowed inv, err := regexp.Compile("[^-a-z0-9]") - checkErr(err, "compile regexp") + util.CheckErr(err, "compile regexp") newName = string(inv.ReplaceAll([]byte(newName), []byte("-"))) // trim leading non-alphabetic leadingnonalpha, err := regexp.Compile("^[^a-z]+") - checkErr(err, "compile regexp") + util.CheckErr(err, "compile regexp") newName = string(leadingnonalpha.ReplaceAll([]byte(newName), []byte{})) // trim trailing trailing, err := regexp.Compile("[^a-z0-9]+$") - checkErr(err, "compile regexp") + util.CheckErr(err, "compile regexp") newName = string(trailing.ReplaceAll([]byte(newName), []byte{})) // truncate to length @@ -125,32 +126,32 @@ func upgradeDumpV1State(v1url string, filename string) { fmt.Println("Getting environments") resp := get(v1url + "/environments") err := json.Unmarshal(resp, &v1state.Environments) - checkErr(err, "parse server response") + util.CheckErr(err, "parse server response") fmt.Println("Getting watches") resp = get(v1url + "/watches") err = json.Unmarshal(resp, &v1state.Watches) - checkErr(err, "parse server response") + util.CheckErr(err, "parse server response") fmt.Println("Getting routes") resp = get(v1url + "/triggers/http") err = json.Unmarshal(resp, &v1state.HTTPTriggers) - checkErr(err, "parse server response") + util.CheckErr(err, "parse server response") fmt.Println("Getting message queue triggers") resp = get(v1url + "/triggers/messagequeue") err = json.Unmarshal(resp, &v1state.Mqtriggers) - checkErr(err, "parse server response") + util.CheckErr(err, "parse server response") fmt.Println("Getting time triggers") resp = get(v1url + "/triggers/time") err = json.Unmarshal(resp, &v1state.TimeTriggers) - checkErr(err, "parse server response") + util.CheckErr(err, "parse server response") fmt.Println("Getting function list") resp = get(v1url + "/functions") err = json.Unmarshal(resp, &v1state.Functions) - checkErr(err, "parse server response") + util.CheckErr(err, "parse server response") // we have to change names that are disallowed in kubernetes nr := nameRemapper{ @@ -199,7 +200,7 @@ func upgradeDumpV1State(v1url string, filename string) { // unmarshal err = json.Unmarshal(resp, &f) - checkErr(err, "parse server response") + util.CheckErr(err, "parse server response") // load into a map to remove duplicates funcs[f.Metadata] = f @@ -216,14 +217,14 @@ func upgradeDumpV1State(v1url string, filename string) { // serialize v1state out, err := json.MarshalIndent(v1state, "", " ") - checkErr(err, "serialize v0.1 state") + util.CheckErr(err, "serialize v0.1 state") // dump to file fission-v01-state.json if len(filename) == 0 { filename = "fission-v01-state.json" } err = ioutil.WriteFile(filename, out, 0644) - checkErr(err, "write file") + util.CheckErr(err, "write file") fmt.Printf("Done: Saved %v functions, %v HTTP triggers, %v watches, %v message queue triggers, %v time triggers.\n", len(v1state.Functions), len(v1state.HTTPTriggers), len(v1state.Watches), len(v1state.Mqtriggers), len(v1state.TimeTriggers)) @@ -249,7 +250,7 @@ func upgradeDumpState(c *cli.Context) error { // check v1 resp, err := http.Get(u + "/environments") - checkErr(err, "reach fission server") + util.CheckErr(err, "reach fission server") if resp.StatusCode == http.StatusNotFound { msg := fmt.Sprintf("Server %v isn't a v1 Fission server. Use --server to point at a pre-0.2.x Fission server.", u) log.Fatal(msg) @@ -266,14 +267,14 @@ func upgradeRestoreState(c *cli.Context) error { } contents, err := ioutil.ReadFile(filename) - checkErr(err, fmt.Sprintf("open file %v", filename)) + util.CheckErr(err, fmt.Sprintf("open file %v", filename)) var v1state V1FissionState err = json.Unmarshal(contents, &v1state) - checkErr(err, "parse dumped v1 state") + util.CheckErr(err, "parse dumped v1 state") // create a regular v2 client - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) // create functions for _, f := range v1state.Functions { @@ -284,9 +285,9 @@ func upgradeRestoreState(c *cli.Context) error { // write function to file tmpfile, err := ioutil.TempFile("", pkgName) - checkErr(err, "create temporary file") + util.CheckErr(err, "create temporary file") code, err := base64.StdEncoding.DecodeString(f.Code) - checkErr(err, "decode base64 function contents") + util.CheckErr(err, "decode base64 function contents") tmpfile.Write(code) tmpfile.Sync() tmpfile.Close() @@ -310,7 +311,7 @@ func upgradeRestoreState(c *cli.Context) error { }, Spec: pkgSpec, }) - checkErr(err, fmt.Sprintf("create package %v", pkgName)) + util.CheckErr(err, fmt.Sprintf("create package %v", pkgName)) _, err = client.FunctionCreate(&crd.Function{ Metadata: *crdMetadataFromV1Metadata(&f.Metadata, v1state.NameChanges), Spec: fission.FunctionSpec{ @@ -324,7 +325,7 @@ func upgradeRestoreState(c *cli.Context) error { }, }, }) - checkErr(err, fmt.Sprintf("create function %v", v1state.NameChanges[f.Metadata.Name])) + util.CheckErr(err, fmt.Sprintf("create function %v", v1state.NameChanges[f.Metadata.Name])) } @@ -339,7 +340,7 @@ func upgradeRestoreState(c *cli.Context) error { }, }, }) - checkErr(err, fmt.Sprintf("create environment %v", e.Metadata.Name)) + util.CheckErr(err, fmt.Sprintf("create environment %v", e.Metadata.Name)) } // create httptriggers @@ -352,7 +353,7 @@ func upgradeRestoreState(c *cli.Context) error { FunctionReference: *functionRefFromV1Metadata(&t.Function, v1state.NameChanges), }, }) - checkErr(err, fmt.Sprintf("create http trigger %v", t.Metadata.Name)) + util.CheckErr(err, fmt.Sprintf("create http trigger %v", t.Metadata.Name)) } // create mqtriggers @@ -366,7 +367,7 @@ func upgradeRestoreState(c *cli.Context) error { ResponseTopic: t.ResponseTopic, }, }) - checkErr(err, fmt.Sprintf("create http trigger %v", t.Metadata.Name)) + util.CheckErr(err, fmt.Sprintf("create http trigger %v", t.Metadata.Name)) } // create time triggers @@ -378,7 +379,7 @@ func upgradeRestoreState(c *cli.Context) error { Cron: t.Cron, }, }) - checkErr(err, fmt.Sprintf("create time trigger %v", t.Metadata.Name)) + util.CheckErr(err, fmt.Sprintf("create time trigger %v", t.Metadata.Name)) } // create watches @@ -391,7 +392,7 @@ func upgradeRestoreState(c *cli.Context) error { FunctionReference: *functionRefFromV1Metadata(&t.Function, v1state.NameChanges), }, }) - checkErr(err, fmt.Sprintf("create kubernetes watch trigger %v", t.Metadata.Name)) + util.CheckErr(err, fmt.Sprintf("create kubernetes watch trigger %v", t.Metadata.Name)) } return nil diff --git a/fission/portforward/portforward.go b/fission/util/portforward.go similarity index 92% rename from fission/portforward/portforward.go rename to fission/util/portforward.go index 0ac52bbf..101f4a1b 100644 --- a/fission/portforward/portforward.go +++ b/fission/util/portforward.go @@ -14,7 +14,7 @@ See the License for the specific language governing permissions and limitations under the License. */ -package portforward +package util import ( "fmt" @@ -27,8 +27,6 @@ import ( "k8s.io/api/core/v1" meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/client-go/kubernetes" - "k8s.io/client-go/tools/clientcmd" "k8s.io/client-go/tools/portforward" "k8s.io/client-go/transport/spdy" @@ -41,7 +39,7 @@ import ( // is found by looking for a service in the same namespace and using // its targetPort. Once the port forward is started, wait for it to // start accepting connections before returning. -func Setup(kubeConfig, namespace, labelSelector string) string { +func SetupPortForward(kubeConfig, namespace, labelSelector string) string { log.Verbose(2, "Setting up port forward to %s in namespace %s using the kubeconfig at %s", labelSelector, namespace, kubeConfig) @@ -104,15 +102,7 @@ func findFreePort() (string, error) { // runPortForward creates a local port forward to the specified pod func runPortForward(kubeConfig string, labelSelector string, localPort string, ns string) error { - config, err := clientcmd.BuildConfigFromFlags("", kubeConfig) - if err != nil { - log.Fatal(fmt.Sprintf("Failed to connect to Kubernetes: %s", err)) - } - - clientset, err := kubernetes.NewForConfig(config) - if err != nil { - log.Fatal(fmt.Sprintf("Failed to connect to Kubernetes: %s", err)) - } + config, clientset := GetKubernetesClient(kubeConfig) log.Verbose(2, "Connected to Kubernetes API") diff --git a/fission/util/util.go b/fission/util/util.go new file mode 100644 index 00000000..c7a28c28 --- /dev/null +++ b/fission/util/util.go @@ -0,0 +1,142 @@ +/* +Copyright 2018 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 util + +import ( + "fmt" + "os" + "path/filepath" + "regexp" + "strings" + + "k8s.io/client-go/kubernetes" + restclient "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" + + "github.com/fission/fission/controller/client" + "github.com/fission/fission/fission/log" +) + +func GetApiClient(serverUrl string) *client.Client { + if len(serverUrl) == 0 { + // starts local portforwarder etc. + serverUrl = GetServerUrl() + } + + isHTTPS := strings.Index(serverUrl, "https://") == 0 + isHTTP := strings.Index(serverUrl, "http://") == 0 + + if !(isHTTP || isHTTPS) { + serverUrl = "http://" + serverUrl + } + + return client.MakeClient(serverUrl) +} + +func GetFissionNamespace() string { + fissionNamespace := os.Getenv("FISSION_NAMESPACE") + return fissionNamespace +} + +func GetKubeConfigPath() string { + kubeConfig := os.Getenv("KUBECONFIG") + if len(kubeConfig) == 0 { + home := os.Getenv("HOME") + kubeConfig = filepath.Join(home, ".kube", "config") + + if _, err := os.Stat(kubeConfig); os.IsNotExist(err) { + log.Fatal("Couldn't find kubeconfig file. " + + "Set the KUBECONFIG environment variable to your kubeconfig's path.") + } + } + return kubeConfig +} + +func GetServerUrl() string { + return GetApplicationUrl("application=fission-api") +} + +func GetApplicationUrl(selector string) string { + var serverUrl string + // Use FISSION_URL env variable if set; otherwise, port-forward to controller. + fissionUrl := os.Getenv("FISSION_URL") + if len(fissionUrl) == 0 { + fissionNamespace := GetFissionNamespace() + kubeConfig := GetKubeConfigPath() + localPort := SetupPortForward(kubeConfig, fissionNamespace, "application=fission-api") + serverUrl = "http://127.0.0.1:" + localPort + } else { + serverUrl = fissionUrl + } + return serverUrl +} + +func CheckErr(err error, msg string) { + if err != nil { + log.Fatal(fmt.Sprintf("Failed to %v: %v", msg, err)) + } +} + +// KubifyName make a kubernetes compliant name out of an arbitrary string +func KubifyName(old string) string { + // Kubernetes maximum name length (for some names; others can be 253 chars) + maxLen := 63 + + newName := strings.ToLower(old) + + // replace disallowed chars with '-' + inv, err := regexp.Compile("[^-a-z0-9]") + CheckErr(err, "compile regexp") + newName = string(inv.ReplaceAll([]byte(newName), []byte("-"))) + + // trim leading non-alphabetic + leadingnonalpha, err := regexp.Compile("^[^a-z]+") + CheckErr(err, "compile regexp") + newName = string(leadingnonalpha.ReplaceAll([]byte(newName), []byte{})) + + // trim trailing + trailing, err := regexp.Compile("[^a-z0-9]+$") + CheckErr(err, "compile regexp") + newName = string(trailing.ReplaceAll([]byte(newName), []byte{})) + + // truncate to length + if len(newName) > maxLen { + newName = newName[0:maxLen] + } + + // if we removed everything, call this thing "default". maybe + // we should generate a unique name... + if len(newName) == 0 { + newName = "default" + } + + return newName +} + +func GetKubernetesClient(kubeConfig string) (*restclient.Config, *kubernetes.Clientset) { + config, err := clientcmd.BuildConfigFromFlags("", kubeConfig) + if err != nil { + log.Fatal(fmt.Sprintf("Failed to connect to Kubernetes: %s", err)) + } + + clientset, err := kubernetes.NewForConfig(config) + if err != nil { + log.Fatal(fmt.Sprintf("Failed to connect to Kubernetes: %s", err)) + } + + return config, clientset +} diff --git a/fission/util/version.go b/fission/util/version.go new file mode 100644 index 00000000..20280516 --- /dev/null +++ b/fission/util/version.go @@ -0,0 +1,48 @@ +package util + +import ( + "fmt" + + "github.com/fission/fission" + "github.com/fission/fission/controller/client" + "github.com/fission/fission/fission/log" + "github.com/fission/fission/fission/plugin" + yaml "gopkg.in/yaml.v2" +) + +// Versions is a container of versions of the client (and its plugins) and server (and its plugins). +type Versions struct { + Client map[string]fission.BuildMeta `json:"client"` + Server map[string]fission.BuildMeta `json:"server"` +} + +func GetVersion(client *client.Client) []byte { + serverInfo, err := client.ServerInfo() + if err != nil { + log.Warn(fmt.Sprintf("Error getting Fission API version: %v", err)) + } + + // Fetch client versions + versions := Versions{ + Client: map[string]fission.BuildMeta{ + "fission/core": fission.BuildInfo(), + }, + } + for _, pmd := range plugin.FindAll() { + versions.Client[pmd.Name] = fission.BuildMeta{ + Version: pmd.Version, + } + } + + // Fetch server versions + versions.Server = map[string]fission.BuildMeta{ + "fission/core": serverInfo.Build, + } + // FUTURE: fetch versions of plugins server-side + bs, err := yaml.Marshal(versions) + if err != nil { + log.Fatal("Failed to format versions: " + err.Error()) + } + + return bs +} diff --git a/fission/watch.go b/fission/watch.go index 14c5cebf..53d9d240 100644 --- a/fission/watch.go +++ b/fission/watch.go @@ -28,10 +28,11 @@ import ( "github.com/fission/fission" "github.com/fission/fission/crd" "github.com/fission/fission/fission/log" + "github.com/fission/fission/fission/util" ) func wCreate(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) fnName := c.String("function") if len(fnName) == 0 { @@ -83,12 +84,12 @@ func wCreate(c *cli.Context) error { if c.Bool("spec") { specFile := fmt.Sprintf("kubewatch-%v.yaml", watchName) err := specSave(*w, specFile) - checkErr(err, "create kubernetes watch spec") + util.CheckErr(err, "create kubernetes watch spec") return nil } _, err := client.WatchCreate(w) - checkErr(err, "create watch") + util.CheckErr(err, "create watch") fmt.Printf("watch '%v' created\n", w.Metadata.Name) return err @@ -107,7 +108,7 @@ func wUpdate(c *cli.Context) error { } func wDelete(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) wName := c.String("name") if len(wName) == 0 { @@ -119,19 +120,19 @@ func wDelete(c *cli.Context) error { Name: wName, Namespace: wNs, }) - checkErr(err, "delete watch") + util.CheckErr(err, "delete watch") fmt.Printf("watch '%v' deleted\n", wName) return nil } func wList(c *cli.Context) error { - client := getClient(c.GlobalString("server")) + client := util.GetApiClient(c.GlobalString("server")) wNs := c.String("triggerns") ws, err := client.WatchList(wNs) - checkErr(err, "list watches") + util.CheckErr(err, "list watches") w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0)