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
This commit is contained in:
Ta-Ching Chen
2018-09-19 14:35:51 +08:00
committed by GitHub
parent a86c9431d3
commit 6b49c194ca
25 changed files with 1397 additions and 609 deletions
+62
View File
@@ -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?
+48
View File
@@ -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
}
+2 -1
View File
@@ -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")
-343
View File
@@ -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
}
+13 -12
View File
@@ -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")
+68 -37
View File
@@ -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 <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
}
+11 -10
View File
@@ -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)
+13 -74
View File
@@ -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:
+11 -10
View File
@@ -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)
+236 -20
View File
@@ -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
}
+12 -11
View File
@@ -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)
+10 -9
View File
@@ -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)
+3 -2
View File
@@ -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)
+19 -18
View File
@@ -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
}
+141
View File
@@ -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
}
+157
View File
@@ -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
}
@@ -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)
}
+249
View File
@@ -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
}
+57
View File
@@ -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 <file> 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)
}
}
+15 -14
View File
@@ -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
}
+29 -28
View File
@@ -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
@@ -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")
+142
View File
@@ -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
}
+48
View File
@@ -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
}
+8 -7
View File
@@ -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)