diff --git a/cmd/fission-cli/app/app.go b/cmd/fission-cli/app/app.go index 6a2d530f..96f3a746 100644 --- a/cmd/fission-cli/app/app.go +++ b/cmd/fission-cli/app/app.go @@ -22,6 +22,7 @@ import ( wrapper "github.com/fission/fission/pkg/fission-cli/cliwrapper/driver/cobra" "github.com/fission/fission/pkg/fission-cli/cliwrapper/driver/cobra/helptemplate" "github.com/fission/fission/pkg/fission-cli/cmd" + "github.com/fission/fission/pkg/fission-cli/cmd/archive" "github.com/fission/fission/pkg/fission-cli/cmd/canaryconfig" "github.com/fission/fission/pkg/fission-cli/cmd/check" "github.com/fission/fission/pkg/fission-cli/cmd/environment" @@ -91,7 +92,7 @@ func App() *cobra.Command { groups := helptemplate.CommandGroups{} groups = append(groups, helptemplate.CreateCmdGroup("Auth Commands(Note: Authentication should be enabled to use a command in this group.)", token.Commands())) - groups = append(groups, helptemplate.CreateCmdGroup("Basic Commands", environment.Commands(), _package.Commands(), function.Commands())) + groups = append(groups, helptemplate.CreateCmdGroup("Basic Commands", environment.Commands(), _package.Commands(), function.Commands(), archive.Commands())) groups = append(groups, helptemplate.CreateCmdGroup("Trigger Commands", httptrigger.Commands(), mqtrigger.Commands(), timetrigger.Commands(), kubewatch.Commands())) groups = append(groups, helptemplate.CreateCmdGroup("Deploy Strategies Commands", canaryconfig.Commands())) groups = append(groups, helptemplate.CreateCmdGroup("Declarative Application Commands", spec.Commands())) diff --git a/pkg/fission-cli/cmd/archive/command.go b/pkg/fission-cli/cmd/archive/command.go new file mode 100644 index 00000000..c1fdeb56 --- /dev/null +++ b/pkg/fission-cli/cmd/archive/command.go @@ -0,0 +1,85 @@ +/* +Copyright 2022 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 archive + +import ( + "github.com/spf13/cobra" + + wrapper "github.com/fission/fission/pkg/fission-cli/cliwrapper/driver/cobra" + "github.com/fission/fission/pkg/fission-cli/flag" +) + +func Commands() *cobra.Command { + + uploadCmd := &cobra.Command{ + Use: "upload", + Short: "Upload an archive", + RunE: wrapper.Wrapper(Upload), + } + wrapper.SetFlags(uploadCmd, flag.FlagSet{ + Required: []flag.Flag{flag.ArchiveName}, + Optional: []flag.Flag{flag.KubeContext}, + }) + + listCmd := &cobra.Command{ + Use: "list", + Short: "List all uploaded archives", + RunE: wrapper.Wrapper(List), + } + wrapper.SetFlags(listCmd, flag.FlagSet{ + Optional: []flag.Flag{flag.KubeContext}, + }) + + deleteCmd := &cobra.Command{ + Use: "delete", + Short: "Delete an archive", + RunE: wrapper.Wrapper(Delete), + } + wrapper.SetFlags(deleteCmd, flag.FlagSet{ + Required: []flag.Flag{flag.ArchiveID}, + Optional: []flag.Flag{flag.ArchiveOutput}, + }) + + geturlCmd := &cobra.Command{ + Use: "get-url", + Short: "Get URL of an uploaded archive", + RunE: wrapper.Wrapper(GetURL), + } + wrapper.SetFlags(geturlCmd, flag.FlagSet{ + Required: []flag.Flag{flag.ArchiveID}, + Optional: []flag.Flag{flag.KubeContext}, + }) + + downloadCmd := &cobra.Command{ + Use: "download", + Short: "Download an archive", + RunE: wrapper.Wrapper(Download), + } + wrapper.SetFlags(downloadCmd, flag.FlagSet{ + Required: []flag.Flag{flag.ArchiveID}, + Optional: []flag.Flag{flag.KubeContext, flag.ArchiveOutput}, + }) + + command := &cobra.Command{ + Use: "archive", + Short: "Manage archives stored with Fission Storage Service.", + Aliases: []string{"ar"}, + } + + command.AddCommand(uploadCmd, listCmd, deleteCmd, geturlCmd, downloadCmd) + return command +} diff --git a/pkg/fission-cli/cmd/archive/delete.go b/pkg/fission-cli/cmd/archive/delete.go new file mode 100644 index 00000000..d6d28bd5 --- /dev/null +++ b/pkg/fission-cli/cmd/archive/delete.go @@ -0,0 +1,58 @@ +/* +Copyright 2022 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 archive + +import ( + "context" + "fmt" + + "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" + "github.com/fission/fission/pkg/fission-cli/cmd" + flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" + "github.com/fission/fission/pkg/fission-cli/util" + storagesvcClient "github.com/fission/fission/pkg/storagesvc/client" +) + +type DeleteSubCommand struct { + cmd.CommandActioner +} + +func Delete(input cli.Input) error { + return (&DeleteSubCommand{}).do(input) +} + +func (opts *DeleteSubCommand) do(input cli.Input) error { + + kubeContext := input.String(flagkey.KubeContext) + archiveID := input.String(flagkey.ArchiveID) + + storagesvcURL, err := util.GetStorageURL(kubeContext) + if err != nil { + return err + } + + client := storagesvcClient.MakeClient(storagesvcURL.String()) + + err = client.Delete(context.Background(), archiveID) + if err != nil { + return err + } + + fmt.Printf("Deleted archive with id: %s", archiveID) + + return nil +} diff --git a/pkg/fission-cli/cmd/archive/download.go b/pkg/fission-cli/cmd/archive/download.go new file mode 100644 index 00000000..974f52cd --- /dev/null +++ b/pkg/fission-cli/cmd/archive/download.go @@ -0,0 +1,62 @@ +/* +Copyright 2022 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 archive + +import ( + "context" + "fmt" + "strings" + + "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" + "github.com/fission/fission/pkg/fission-cli/cmd" + flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" + "github.com/fission/fission/pkg/fission-cli/util" + storagesvcClient "github.com/fission/fission/pkg/storagesvc/client" +) + +type DownloadSubCommand struct { + cmd.CommandActioner +} + +func Download(input cli.Input) error { + return (&DownloadSubCommand{}).do(input) +} + +func (opts *DownloadSubCommand) do(input cli.Input) error { + + kubeContext := input.String(flagkey.KubeContext) + archiveID := input.String(flagkey.ArchiveID) + archiveOutput := input.String(flagkey.ArchiveOutput) + + if len(archiveOutput) == 0 { + archiveOutput = strings.TrimPrefix(archiveID, "/fission/fission-functions/") + } + + storageAccessURL, err := util.GetStorageURL(kubeContext) + if err != nil { + return err + } + + client := storagesvcClient.MakeClient(storageAccessURL.String()) + err = client.Download(context.Background(), archiveID, archiveOutput) + if err != nil { + return err + } + + fmt.Printf("File download complete. File name: %s", archiveOutput) + return nil +} diff --git a/pkg/fission-cli/cmd/archive/geturl.go b/pkg/fission-cli/cmd/archive/geturl.go new file mode 100644 index 00000000..906a1566 --- /dev/null +++ b/pkg/fission-cli/cmd/archive/geturl.go @@ -0,0 +1,85 @@ +/* +Copyright 2019 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 archive + +import ( + "fmt" + "net/http" + "net/url" + + "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" + "github.com/fission/fission/pkg/fission-cli/cmd" + flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" + "github.com/fission/fission/pkg/fission-cli/util" + storagesvcClient "github.com/fission/fission/pkg/storagesvc/client" +) + +type GetURLSubCommand struct { + cmd.CommandActioner +} + +func GetURL(input cli.Input) error { + return (&GetURLSubCommand{}).do(input) +} + +func (opts *GetURLSubCommand) do(input cli.Input) error { + + kubeContext := input.String(flagkey.KubeContext) + archiveID := input.String(flagkey.ArchiveID) + + serverURL, err := util.GetStorageURL(kubeContext) + if err != nil { + return err + } + + relativeURL, _ := url.Parse(util.FISSION_STORAGE_URI) + + queryString := relativeURL.Query() + queryString.Set("id", archiveID) + relativeURL.RawQuery = queryString.Encode() + + storageAccessURL := serverURL.ResolveReference(relativeURL) + + resp, err := http.Head(storageAccessURL.String()) + if err != nil { + return err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("Error getting URL. Exited with Status: %v", resp.Status) + } + + archiveURL, err := url.Parse(resp.Header.Get("X-FISSION-ARCHIVEURL")) + if err != nil { + return err + } + + if archiveURL.Scheme == "file" { + storageSvc, err := opts.Client().V1().Misc().GetSvcURL("application=fission-storage") + if err != nil { + return err + } + storagesvcURL := "http://" + storageSvc + client := storagesvcClient.MakeClient(storagesvcURL) + fmt.Printf("URL: %s", client.GetUrl(archiveID)) + } else { + fmt.Printf("URL: %s", archiveURL.String()) + } + + return nil +} diff --git a/pkg/fission-cli/cmd/archive/list.go b/pkg/fission-cli/cmd/archive/list.go new file mode 100644 index 00000000..109cce7e --- /dev/null +++ b/pkg/fission-cli/cmd/archive/list.go @@ -0,0 +1,60 @@ +/* +Copyright 2022 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 archive + +import ( + "context" + "fmt" + + "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" + "github.com/fission/fission/pkg/fission-cli/cmd" + flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" + "github.com/fission/fission/pkg/fission-cli/util" + storagesvcClient "github.com/fission/fission/pkg/storagesvc/client" +) + +type ListSubCommand struct { + cmd.CommandActioner +} + +func List(input cli.Input) error { + return (&ListSubCommand{}).do(input) +} + +func (opts *ListSubCommand) do(input cli.Input) error { + + kubeContext := input.String(flagkey.KubeContext) + + storageAccessURL, err := util.GetStorageURL(kubeContext) + if err != nil { + return err + } + + client := storagesvcClient.MakeClient(storageAccessURL.String()) + files, err := client.List(context.Background()) + if err != nil { + return err + } + + fmt.Println("ARCHIVES") + for _, file := range files { + fmt.Println(file) + } + + return nil + +} diff --git a/pkg/fission-cli/cmd/archive/upload.go b/pkg/fission-cli/cmd/archive/upload.go new file mode 100644 index 00000000..782085b7 --- /dev/null +++ b/pkg/fission-cli/cmd/archive/upload.go @@ -0,0 +1,57 @@ +/* +Copyright 2022 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 archive + +import ( + "context" + "fmt" + + "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" + "github.com/fission/fission/pkg/fission-cli/cmd" + flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" + "github.com/fission/fission/pkg/fission-cli/util" + storagesvcClient "github.com/fission/fission/pkg/storagesvc/client" +) + +type UploadSubCommand struct { + cmd.CommandActioner +} + +func Upload(input cli.Input) error { + return (&UploadSubCommand{}).do(input) +} + +func (opts *UploadSubCommand) do(input cli.Input) error { + + kubeContext := input.String(flagkey.KubeContext) + archiveName := input.String(flagkey.ArchiveName) + + storagesvcURL, err := util.GetStorageURL(kubeContext) + if err != nil { + return err + } + + client := storagesvcClient.MakeClient(storagesvcURL.String()) + archiveID, err := client.Upload(context.Background(), archiveName, nil) + if err != nil { + return err + } + + fmt.Printf("File successfully uploaded with ID: %s ", archiveID) + + return nil +} diff --git a/pkg/fission-cli/flag/flag.go b/pkg/fission-cli/flag/flag.go index f17245ce..b51126a1 100644 --- a/pkg/fission-cli/flag/flag.go +++ b/pkg/fission-cli/flag/flag.go @@ -227,4 +227,8 @@ var ( CanaryWeightIncrement = Flag{Type: Int, Name: flagkey.CanaryWeightIncrement, Aliases: []string{"step"}, Usage: "Weight increment step for function", DefaultValue: 20} CanaryIncrementInterval = Flag{Type: String, Name: flagkey.CanaryIncrementInterval, Aliases: []string{"internal"}, Usage: "Weight increment interval, string representation of time.Duration, ex : 1m, 2h, 2d", DefaultValue: "2m"} CanaryFailureThreshold = Flag{Type: Int, Name: flagkey.CanaryFailureThreshold, Aliases: []string{"threshold"}, Usage: "Threshold in percentage beyond which the new version of the function is considered unstable", DefaultValue: 10} + + ArchiveName = Flag{Type: String, Name: flagkey.ArchiveName, Usage: "Name of the archive file"} + ArchiveID = Flag{Type: String, Name: flagkey.ArchiveID, Usage: "Id for the archive file"} + ArchiveOutput = Flag{Type: String, Name: flagkey.ArchiveOutput, Usage: "Download file with this name", Aliases: []string{"o"}, DefaultValue: ""} ) diff --git a/pkg/fission-cli/flag/key/key.go b/pkg/fission-cli/flag/key/key.go index 97607eec..955259d4 100644 --- a/pkg/fission-cli/flag/key/key.go +++ b/pkg/fission-cli/flag/key/key.go @@ -178,5 +178,9 @@ const ( CanaryIncrementInterval = "increment-interval" CanaryFailureThreshold = "failure-threshold" + ArchiveName = resourceName + ArchiveID = "id" + ArchiveOutput = Output + DefaultSpecOutputDir = "fission-dump" ) diff --git a/pkg/fission-cli/util/constants.go b/pkg/fission-cli/util/constants.go index bca1af5d..ec084686 100644 --- a/pkg/fission-cli/util/constants.go +++ b/pkg/fission-cli/util/constants.go @@ -18,8 +18,9 @@ package util // fission-cli options const ( - SPEC_IGNORE_FILE = ".specignore" - COMMIT_LABEL = "commit" - FISSION_AUTH_URI = "/auth/login" - FISSION_AUTH_TOKEN = "FISSION_AUTH_TOKEN" + SPEC_IGNORE_FILE = ".specignore" + COMMIT_LABEL = "commit" + FISSION_AUTH_URI = "/auth/login" + FISSION_AUTH_TOKEN = "FISSION_AUTH_TOKEN" + FISSION_STORAGE_URI = "/v1/archive" ) diff --git a/pkg/fission-cli/util/util.go b/pkg/fission-cli/util/util.go index 4bd3d28a..ed512d98 100644 --- a/pkg/fission-cli/util/util.go +++ b/pkg/fission-cli/util/util.go @@ -18,6 +18,7 @@ package util import ( "fmt" + "net/url" "os" "os/user" "path/filepath" @@ -443,3 +444,17 @@ func ApplyLabelsAndAnnotations(input cli.Input, objectMeta *metav1.ObjectMeta) e } return nil } + +func GetStorageURL(kubeContext string) (*url.URL, error) { + storageLocalPort, err := SetupPortForward(GetFissionNamespace(), "application=fission-storage", kubeContext) + if err != nil { + return nil, err + } + + serverURL, err := url.Parse("http://127.0.0.1:" + storageLocalPort) + if err != nil { + return nil, err + } + + return serverURL, nil +} diff --git a/pkg/storagesvc/client/client.go b/pkg/storagesvc/client/client.go index e598386d..83bdd4a4 100644 --- a/pkg/storagesvc/client/client.go +++ b/pkg/storagesvc/client/client.go @@ -123,6 +123,33 @@ func (c *Client) GetUrl(id string) string { return fmt.Sprintf("%v/archive?id=%v", c.url, url.PathEscape(id)) } +func (c *Client) List(ctx context.Context) ([]string, error) { + req, err := http.NewRequest(http.MethodGet, c.url+"/archive", nil) + if err != nil { + return []string{}, err + } + resp, err := ctxhttp.Do(ctx, c.httpClient, req) + if err != nil { + return []string{}, err + } + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + if err != nil { + return []string{}, err + } + if resp.StatusCode != http.StatusOK { + msg := fmt.Sprintf("List error %v", resp.Status) + return []string{}, errors.New(msg) + } + + var ids []string + err = json.Unmarshal(body, &ids) + if err != nil { + return []string{}, err + } + return ids, nil +} + // Download fetches the file identified by ID to the local file path. // filePath must not exist. func (c *Client) Download(ctx context.Context, id string, filePath string) error { diff --git a/pkg/storagesvc/client/storagesvc_test.go b/pkg/storagesvc/client/storagesvc_test.go index bd36fe67..ddf31f04 100644 --- a/pkg/storagesvc/client/storagesvc_test.go +++ b/pkg/storagesvc/client/storagesvc_test.go @@ -42,20 +42,25 @@ const ( minioRegion = "ap-south-1" ) -func panicIf(err error) { +func failTest(t *testing.T, err error) { if err != nil { - log.Panicf("Error: %v", err) + t.Fatalf("%v", err) } } -func MakeTestFile(size int) *os.File { +func MakeTestFile(size int) (*os.File, error) { f, err := os.CreateTemp("", "storagesvc_test_") - panicIf(err) + + if err != nil { + return nil, err + } _, err = f.Write(bytes.Repeat([]byte("."), size)) - panicIf(err) - return f + if err != nil { + return nil, err + } + return f, nil } func runMinioDockerContainer(pool *dockertest.Pool) *dockertest.Resource { @@ -120,7 +125,7 @@ func TestS3StorageService(t *testing.T) { config := zap.NewDevelopmentConfig() config.EncoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder logger, err := config.Build() - panicIf(err) + failTest(t, err) log.Println("starting storage svc") os.Setenv("STORAGE_S3_ENDPOINT", endpoint) @@ -139,27 +144,28 @@ func TestS3StorageService(t *testing.T) { client := MakeClient(fmt.Sprintf("http://localhost:%v/", 8081)) // generate a test file - tmpfile := MakeTestFile(10 * 1024) + tmpfile, err := MakeTestFile(10 * 1024) + failTest(t, err) defer os.Remove(tmpfile.Name()) // store it metadata := make(map[string]string) fileID, err := client.Upload(ctx, tmpfile.Name(), &metadata) - panicIf(err) + failTest(t, err) time.Sleep(10 * time.Second) // Retrieve file through minioClient reader, err := minioClient.GetObject(bucketName, fileID, minio.GetObjectOptions{}) - panicIf(err) + failTest(t, err) defer reader.Close() retThroughMinio, err := os.CreateTemp("", "storagesvc_verify_minio_") - panicIf(err) + failTest(t, err) defer os.Remove(retThroughMinio.Name()) stat, err := reader.Stat() - panicIf(err) + failTest(t, err) if _, err := io.CopyN(retThroughMinio, reader, stat.Size); err != nil { log.Fatalln(err) @@ -167,25 +173,25 @@ func TestS3StorageService(t *testing.T) { // Retrieve file through API retThroughAPI, err := os.CreateTemp("", "storagesvc_verify_") - panicIf(err) + failTest(t, err) os.Remove(retThroughAPI.Name()) err = client.Download(ctx, fileID, retThroughAPI.Name()) - panicIf(err) + failTest(t, err) defer os.Remove(retThroughAPI.Name()) // compare contents contentsMinio, err := os.ReadFile(retThroughMinio.Name()) - panicIf(err) + failTest(t, err) contentsAPI, err := os.ReadFile(retThroughAPI.Name()) - panicIf(err) + failTest(t, err) if !bytes.Equal(contentsMinio, contentsAPI) { log.Panic("Contents don't match") } // delete uploaded file err = client.Delete(ctx, fileID) - panicIf(err) + failTest(t, err) // make sure download fails err = client.Download(ctx, fileID, "xxx") @@ -201,7 +207,7 @@ func TestLocalStorageService(t *testing.T) { config := zap.NewDevelopmentConfig() config.EncoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder logger, err := config.Build() - panicIf(err) + failTest(t, err) log.Println("starting storage svc") localPath := fmt.Sprintf("/tmp/%v", testID) @@ -216,36 +222,45 @@ func TestLocalStorageService(t *testing.T) { client := MakeClient(fmt.Sprintf("http://localhost:%v/", port)) // generate a test file - tmpfile := MakeTestFile(10 * 1024) + tmpfile, err := MakeTestFile(10 * 1024) + failTest(t, err) defer os.Remove(tmpfile.Name()) // store it metadata := make(map[string]string) fileID, err := client.Upload(ctx, tmpfile.Name(), &metadata) - panicIf(err) + failTest(t, err) + + ids, err := client.List(ctx) + if err != nil { + t.Fatalf("Could not list files: %s", err) + } + if len(ids) != 1 { + t.Fatalf("Expected 1 file, got %v", len(ids)) + } // make a temp file for verification retrievedfile, err := os.CreateTemp("", "storagesvc_verify_") - panicIf(err) + failTest(t, err) os.Remove(retrievedfile.Name()) // retrieve uploaded file err = client.Download(ctx, fileID, retrievedfile.Name()) - panicIf(err) + failTest(t, err) defer os.Remove(retrievedfile.Name()) // compare contents contents1, err := os.ReadFile(tmpfile.Name()) - panicIf(err) + failTest(t, err) contents2, err := os.ReadFile(retrievedfile.Name()) - panicIf(err) + failTest(t, err) if !bytes.Equal(contents1, contents2) { log.Panic("Contents don't match") } // delete uploaded file err = client.Delete(ctx, fileID) - panicIf(err) + failTest(t, err) // make sure download fails err = client.Download(ctx, fileID, "xxx") diff --git a/pkg/storagesvc/storagesvc.go b/pkg/storagesvc/storagesvc.go index c74538a1..18fc9160 100644 --- a/pkg/storagesvc/storagesvc.go +++ b/pkg/storagesvc/storagesvc.go @@ -67,6 +67,32 @@ func getStorageLocation(config *storageConfig) (stow.Location, error) { return config.storage.dial() } +func (ss *StorageService) listItems(w http.ResponseWriter, r *http.Request) { + // get all archives on storage + // out of them, there may be some just created but not referenced by packages yet. + // need to filter them out. + archivesInStorage, err := ss.storageClient.getItemIDsWithFilter(ss.storageClient.filterAllItems, false) + if err != nil { + ss.logger.Error("error getting items from storage", zap.Error(err)) + return + } + ss.logger.Debug("archives in storage", zap.Strings("archives", archivesInStorage)) + + // respond with the list of items + resp, err := json.Marshal(archivesInStorage) + if err != nil { + http.Error(w, "error marshaling item list", http.StatusInternalServerError) + return + } + _, err = w.Write(resp) + if err != nil { + ss.logger.Error( + "error writing HTTP response", + zap.Error(err), + ) + } +} + // Handle multipart file uploads. func (ss *StorageService) uploadHandler(w http.ResponseWriter, r *http.Request) { // handle upload @@ -203,6 +229,23 @@ func (ss *StorageService) downloadHandler(w http.ResponseWriter, r *http.Request } } +func (ss *StorageService) infoHandler(w http.ResponseWriter, r *http.Request) { + + fileId, err := ss.getIdFromRequest(r) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + itemURL, err := ss.storageClient.getURL(fileId) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + } + + w.Header().Add("X-FISSION-ARCHIVEURL", itemURL.String()) + +} + func (ss *StorageService) healthHandler(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) } @@ -219,8 +262,10 @@ func (ss *StorageService) Start(ctx context.Context, port int, openTracingEnable r := mux.NewRouter() r.Use(metrics.HTTPMetricMiddleware) r.HandleFunc("/v1/archive", ss.uploadHandler).Methods("POST") - r.HandleFunc("/v1/archive", ss.downloadHandler).Methods("GET") + r.HandleFunc("/v1/archive", ss.downloadHandler).Queries("id", "{id}").Methods("GET") + r.HandleFunc("/v1/archive", ss.listItems).Methods("GET") r.HandleFunc("/v1/archive", ss.deleteHandler).Methods("DELETE") + r.HandleFunc("/v1/archive", ss.infoHandler).Methods("HEAD") r.HandleFunc("/healthz", ss.healthHandler).Methods("GET") var handler http.Handler diff --git a/pkg/storagesvc/stowClient.go b/pkg/storagesvc/stowClient.go index be02393d..ad5c2ce9 100644 --- a/pkg/storagesvc/stowClient.go +++ b/pkg/storagesvc/stowClient.go @@ -20,6 +20,7 @@ import ( "fmt" "io" "mime/multipart" + "net/url" "os" "strings" "time" @@ -182,6 +183,18 @@ func (client *StowClient) removeFileByID(itemID string) error { return client.container.RemoveItem(itemID) } +func (client *StowClient) getURL(itemID string) (*url.URL, error) { + item, err := client.container.Item(itemID) + if err != nil { + if err == stow.ErrNotFound { + return nil, ErrNotFound + } else { + return nil, ErrRetrievingItem + } + } + return item.URL(), nil +} + func (client *StowClient) getFileSize(itemID string) (int64, error) { item, err := client.container.Item(itemID) if err != nil { @@ -240,3 +253,12 @@ func (client StowClient) filterItemCreatedAMinuteAgo(item stow.Item, currentTime } return false } + +func (client StowClient) filterAllItems(item stow.Item, _ interface{}) bool { + itemLastModTime, _ := item.LastMod() + client.logger.Debug("item info", + zap.String("item", item.ID()), + zap.Time("last_modified_time", itemLastModTime)) + return false + +} diff --git a/test/kind_CI.sh b/test/kind_CI.sh index 1853820d..d8bb536a 100755 --- a/test/kind_CI.sh +++ b/test/kind_CI.sh @@ -57,6 +57,7 @@ main() { # run tests without newdeploy in parallel. export JOBS=6 source $ROOT/test/run_test.sh \ + $ROOT/test/tests/test_archive_cli.sh \ $ROOT/test/tests/test_canary.sh \ $ROOT/test/tests/test_fn_update/test_idle_objects_reaper.sh \ $ROOT/test/tests/mqtrigger/kafka/test_kafka.sh \ diff --git a/test/tests/test_archive_cli.sh b/test/tests/test_archive_cli.sh new file mode 100755 index 00000000..dfb889a9 --- /dev/null +++ b/test/tests/test_archive_cli.sh @@ -0,0 +1,63 @@ +#!/bin/bash + +set -euo pipefail +source "$(dirname $0)"/../utils.sh + +TEST_ID=$(generate_test_id) +echo "TEST_ID = $TEST_ID" + +tmp_dir="/tmp/test-$TEST_ID" +mkdir -p "$tmp_dir" + +podname=$(kubectl get pods -n fission | grep "storagesvc" |awk '{print $1}') + +cleanup() { + log "Cleaning up..." + clean_resource_by_id $TEST_ID + rm -rf $tmp_dir +} + +if [ -z "${TEST_NOCLEANUP:-}" ]; then + trap cleanup EXIT +else + log "TEST_NOCLEANUP is set; not cleaning up test artifacts afterwards." +fi + +create_archive() { + log "Creating an archive" + mkdir -p "$tmp_dir"/archive + dd if=/dev/urandom of="$tmp_dir"/archive/dynamically_generated_file bs=256k count=1 + printf 'def main():\n return "Hello, world!"' > "$tmp_dir"/archive/hello.py + zip -jr "$tmp_dir"/test-deploy-pkg.zip "$tmp_dir"/archive/ +} + +# Test for upload +create_archive +uploadResp=$(fission ar upload --name "$tmp_dir"/test-deploy-pkg.zip) +filename=$(echo "$uploadResp" | cut -d':' -f2 | tr -d ' ') + +kubectl exec -i "$podname" -n fission -- /bin/sh -c "ls $filename" + +# Test for list +listResp=$(fission ar list) + +echo "$listResp" | grep "$filename" + +# Test for download +fission ar download --id "$filename" + +fileID=$(echo "$filename" | cut -d'/' -f4) + +ls | grep "$fileID" + +# Test for get-url +getURLResp=$(fission ar get-url --id "$filename") + +echo "$getURLResp" | grep "http://storagesvc.fission/v1/archive?id=%2Ffission%2Ffission-functions%2F$fileID" + +# Test for delete +fission ar delete --id "$filename" + +fission ar list | grep -v "$filename" + +log "Archive CLI tests done." \ No newline at end of file