From 8de5a5b0f317b94525ff1d41cddecbeb5cb68899 Mon Sep 17 00:00:00 2001 From: soharab-ic <156293296+soharab-ic@users.noreply.github.com> Date: Mon, 15 Jul 2024 18:28:42 +0530 Subject: [PATCH] Fix: Storage leak in Builder and Fetcher (#2979) * Fix storage leak in builder ``` builder pod keeps old src and deployment packages irrespective of build status. delete src package after every build request is completed. delete deployment package after package is uploaded. ``` * Optimized src/deploy cleanup pkg code * Fix high severity security issue * Add a test for builder's Clean API * Resolve review comments --------- Signed-off-by: Md Soharab Ansari --- cmd/builder/app/server.go | 1 + pkg/builder/builder.go | 33 ++++++++++++++++++++++++++ pkg/builder/builder_test.go | 46 ++++++++++++++++++++++++++++++++++++ pkg/builder/client/client.go | 32 +++++++++++++++++++++++++ pkg/buildermgr/common.go | 18 ++++++++++++++ pkg/fetcher/fetcher.go | 9 +++++++ pkg/utils/utils.go | 29 +++++++++++++++++++++++ 7 files changed, 168 insertions(+) diff --git a/cmd/builder/app/server.go b/cmd/builder/app/server.go index fc55b173..7f7c4cc5 100644 --- a/cmd/builder/app/server.go +++ b/cmd/builder/app/server.go @@ -32,6 +32,7 @@ func Run(ctx context.Context, logger *zap.Logger, mgr manager.Interface, shareVo builder := builder.MakeBuilder(logger, shareVolume) mux := http.NewServeMux() mux.HandleFunc("/", builder.Handler) + mux.HandleFunc("/clean", builder.Clean) mux.HandleFunc("/version", builder.VersionHandler) mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) diff --git a/pkg/builder/builder.go b/pkg/builder/builder.go index b6edddc1..b2281012 100644 --- a/pkg/builder/builder.go +++ b/pkg/builder/builder.go @@ -35,6 +35,7 @@ import ( "go.uber.org/zap" "github.com/fission/fission/pkg/info" + "github.com/fission/fission/pkg/utils" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -154,6 +155,38 @@ func (builder *Builder) Handler(w http.ResponseWriter, r *http.Request) { builder.reply(r.Context(), w, deployPkgFilename, buildLogs, http.StatusOK) } +func (builder *Builder) Clean(w http.ResponseWriter, r *http.Request) { + logger := otelUtils.LoggerWithTraceID(r.Context(), builder.logger) + + if r.Method != "DELETE" { + e := "method not allowed" + logger.Error(e, zap.String("http_method", r.Method)) + builder.reply(r.Context(), w, "", fmt.Sprintf("%s: %s", e, r.Method), http.StatusMethodNotAllowed) + return + } + + startTime := time.Now() + defer func() { + elapsed := time.Since(startTime) + logger.Info("clean request complete", zap.Duration("elapsed_time", elapsed)) + }() + + srcPkgFilename := r.URL.Query().Get("name") + srcPkgPath := filepath.Join(builder.sharedVolumePath, srcPkgFilename) + + logger.Info("builder received clean request", zap.Any("source_package", srcPkgFilename)) + + err := utils.DeleteOldPackages(srcPkgPath, envSrcPkg) + if err != nil { + e := "error deleting src package after build" + logger.Error(e, zap.Error(err)) + builder.reply(r.Context(), w, srcPkgFilename, "", http.StatusInternalServerError) + return + } + + builder.reply(r.Context(), w, srcPkgFilename, "", http.StatusOK) +} + func (builder *Builder) reply(ctx context.Context, w http.ResponseWriter, pkgFilename string, buildLogs string, statusCode int) { logger := otelUtils.LoggerWithTraceID(ctx, builder.logger) resp := PackageBuildResponse{ diff --git a/pkg/builder/builder_test.go b/pkg/builder/builder_test.go index 106cd257..8439c72b 100644 --- a/pkg/builder/builder_test.go +++ b/pkg/builder/builder_test.go @@ -168,4 +168,50 @@ func TestBuilder(t *testing.T) { }) } }) + + // Test CleanHandler + t.Run("CleanHandler", func(t *testing.T) { + for _, test := range []struct { + name string + srcPkgFilename string + handler func(w http.ResponseWriter, r *http.Request) + status int + }{ + { + name: "should fail deleting src pkg: invalid shared volume path", + srcPkgFilename: "test2", + handler: func(w http.ResponseWriter, r *http.Request) { + builder.Clean(w, r) + }, + status: http.StatusInternalServerError, + }, + { + name: "should fail deleting src pkg: method not allowed", + srcPkgFilename: "test3", + handler: func(w http.ResponseWriter, r *http.Request) { + builder.Handler(w, r) + }, + status: http.StatusMethodNotAllowed, + }, + } { + t.Run(test.name, func(t *testing.T) { + _, err := os.MkdirTemp(dir, test.srcPkgFilename) + if err != nil { + t.Fatal(err) + } + srcFile, err := os.Create(dir + "/" + test.srcPkgFilename) + if err != nil { + t.Fatal(err) + } + defer srcFile.Close() + w := httptest.NewRecorder() + r := httptest.NewRequest(http.MethodDelete, "/clean/"+test.srcPkgFilename, http.NoBody) + test.handler(w, r) + resp := w.Result() + if resp.StatusCode != test.status { + t.Errorf("expected status code %d, got %d", test.status, resp.StatusCode) + } + }) + } + }) } diff --git a/pkg/builder/client/client.go b/pkg/builder/client/client.go index 5977f6d5..2b6764e9 100644 --- a/pkg/builder/client/client.go +++ b/pkg/builder/client/client.go @@ -21,6 +21,7 @@ import ( "context" "encoding/json" "io" + "net/http" "strings" "github.com/hashicorp/go-retryablehttp" @@ -37,6 +38,7 @@ import ( type ( ClientInterface interface { Build(context.Context, *builder.PackageBuildRequest) (*builder.PackageBuildResponse, error) + Clean(context.Context, string) error } client struct { @@ -57,6 +59,10 @@ func MakeClient(logger *zap.Logger, builderUrl string) ClientInterface { } } +func (c *client) getCleanUrl(srcPkgFilename string) string { + return c.url + "/clean" + "?name=" + srcPkgFilename +} + func (c *client) Build(ctx context.Context, req *builder.PackageBuildRequest) (*builder.PackageBuildResponse, error) { logger := otelUtils.LoggerWithTraceID(ctx, c.logger) @@ -86,3 +92,29 @@ func (c *client) Build(ctx context.Context, req *builder.PackageBuildRequest) (* return &pkgBuildResp, ferror.MakeErrorFromHTTP(resp) } + +func (c *client) Clean(ctx context.Context, srcPkgFilename string) error { + logger := otelUtils.LoggerWithTraceID(ctx, c.logger) + + req, err := http.NewRequest(http.MethodDelete, c.getCleanUrl(srcPkgFilename), http.NoBody) + if err != nil { + return errors.Wrap(err, "failed to create http request for clean api") + } + + resp, err := ctxhttp.Do(ctx, c.httpClient.StandardClient(), req) + if err != nil { + logger.Error("error sending clean request", zap.Error(err)) + return err + } + defer resp.Body.Close() + + if resp.StatusCode == http.StatusMethodNotAllowed { + return nil + } + + if resp.StatusCode != http.StatusOK { + return ferror.MakeErrorFromHTTP(resp) + } + + return nil +} diff --git a/pkg/buildermgr/common.go b/pkg/buildermgr/common.go index c3bbbf5d..f9c6000b 100644 --- a/pkg/buildermgr/common.go +++ b/pkg/buildermgr/common.go @@ -60,6 +60,15 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient version fetcherC := fetcherClient.MakeClient(logger, fmt.Sprintf("http://%v:8000", svcName)) builderC := builderClient.MakeClient(logger, fmt.Sprintf("http://%v:8001", svcName)) + defer func() { + logger.Info("cleaning src pkg from builder storage", zap.String("source_package", srcPkgFilename)) + errC := cleanPackage(ctx, builderC, srcPkgFilename) + if errC != nil { + m := "error cleaning src pkg from builder storage" + logger.Error(m, zap.Error(errC)) + } + }() + fetchReq := &fetcher.FunctionFetchRequest{ FetchType: fv1.FETCH_SOURCE, Package: pkg.ObjectMeta, @@ -121,6 +130,15 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient version return uploadResp, buildResp.BuildLogs, nil } +func cleanPackage(ctx context.Context, builderClient builderClient.ClientInterface, srcPkgFileName string) error { + err := builderClient.Clean(ctx, srcPkgFileName) + if err != nil { + return err + } + + return nil +} + func updatePackage(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, pkg *fv1.Package, status fv1.BuildStatus, buildLogs string, uploadResp *fetcher.ArchiveUploadResponse) (*fv1.Package, error) { diff --git a/pkg/fetcher/fetcher.go b/pkg/fetcher/fetcher.go index 93de0bb6..6a417f86 100644 --- a/pkg/fetcher/fetcher.go +++ b/pkg/fetcher/fetcher.go @@ -525,6 +525,14 @@ func (fetcher *Fetcher) UploadHandler(w http.ResponseWriter, r *http.Request) { srcFilepath := filepath.Join(fetcher.sharedVolumePath, req.Filename) dstFilepath := filepath.Join(fetcher.sharedVolumePath, zipFilename) + defer func() { + errC := utils.DeleteOldPackages(srcFilepath, "DEPLOY_PKG") + if errC != nil { + m := "error deleting deploy package after upload" + logger.Error(m, zap.Error(errC)) + } + }() + if req.ArchivePackage { err = fetcher.archive(srcFilepath, dstFilepath) if err != nil { @@ -584,6 +592,7 @@ func (fetcher *Fetcher) UploadHandler(w http.ResponseWriter, r *http.Request) { logger.Error(e, zap.Error(err)) http.Error(w, fmt.Sprintf("%s: %v", e, err), http.StatusInternalServerError) } + } func (fetcher *Fetcher) rename(src string, dst string) error { diff --git a/pkg/utils/utils.go b/pkg/utils/utils.go index 00ca97e6..89d0a135 100644 --- a/pkg/utils/utils.go +++ b/pkg/utils/utils.go @@ -267,3 +267,32 @@ func FindFreePort() (int, error) { return port, nil } + +// DeleteOldPackages deletes src and built deployment packages from builder's storage. +// The function also verifies that sharedVolumePath for builder and fetcher containers +// is /packages. A source_package contains a directory and a .tmp file while a deployment +// package contains a directory and a .zip file. +func DeleteOldPackages(pkgPath, pkgType string) error { + sharedVolumePath := "/packages" + if !strings.HasPrefix(pkgPath, sharedVolumePath) { + return fmt.Errorf("invalid shared volume path: %s", pkgPath) + } + + var file string + if pkgType == "DEPLOY_PKG" { + file = pkgPath + ".zip" + } else if pkgType == "SRC_PKG" { + file = pkgPath + ".tmp" + } + + err := os.RemoveAll(pkgPath) + if err != nil { + return err + } + err = os.Remove(file) + if err != nil { + return err + } + + return nil +}