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 <soharab.ansari@infracloud.io>
This commit is contained in:
@@ -32,6 +32,7 @@ func Run(ctx context.Context, logger *zap.Logger, mgr manager.Interface, shareVo
|
|||||||
builder := builder.MakeBuilder(logger, shareVolume)
|
builder := builder.MakeBuilder(logger, shareVolume)
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
mux.HandleFunc("/", builder.Handler)
|
mux.HandleFunc("/", builder.Handler)
|
||||||
|
mux.HandleFunc("/clean", builder.Clean)
|
||||||
mux.HandleFunc("/version", builder.VersionHandler)
|
mux.HandleFunc("/version", builder.VersionHandler)
|
||||||
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
|
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
|
|||||||
@@ -35,6 +35,7 @@ import (
|
|||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
|
|
||||||
"github.com/fission/fission/pkg/info"
|
"github.com/fission/fission/pkg/info"
|
||||||
|
"github.com/fission/fission/pkg/utils"
|
||||||
otelUtils "github.com/fission/fission/pkg/utils/otel"
|
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)
|
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) {
|
func (builder *Builder) reply(ctx context.Context, w http.ResponseWriter, pkgFilename string, buildLogs string, statusCode int) {
|
||||||
logger := otelUtils.LoggerWithTraceID(ctx, builder.logger)
|
logger := otelUtils.LoggerWithTraceID(ctx, builder.logger)
|
||||||
resp := PackageBuildResponse{
|
resp := PackageBuildResponse{
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,6 +21,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"io"
|
"io"
|
||||||
|
"net/http"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/hashicorp/go-retryablehttp"
|
"github.com/hashicorp/go-retryablehttp"
|
||||||
@@ -37,6 +38,7 @@ import (
|
|||||||
type (
|
type (
|
||||||
ClientInterface interface {
|
ClientInterface interface {
|
||||||
Build(context.Context, *builder.PackageBuildRequest) (*builder.PackageBuildResponse, error)
|
Build(context.Context, *builder.PackageBuildRequest) (*builder.PackageBuildResponse, error)
|
||||||
|
Clean(context.Context, string) error
|
||||||
}
|
}
|
||||||
|
|
||||||
client struct {
|
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) {
|
func (c *client) Build(ctx context.Context, req *builder.PackageBuildRequest) (*builder.PackageBuildResponse, error) {
|
||||||
logger := otelUtils.LoggerWithTraceID(ctx, c.logger)
|
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)
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -60,6 +60,15 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient version
|
|||||||
fetcherC := fetcherClient.MakeClient(logger, fmt.Sprintf("http://%v:8000", svcName))
|
fetcherC := fetcherClient.MakeClient(logger, fmt.Sprintf("http://%v:8000", svcName))
|
||||||
builderC := builderClient.MakeClient(logger, fmt.Sprintf("http://%v:8001", 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{
|
fetchReq := &fetcher.FunctionFetchRequest{
|
||||||
FetchType: fv1.FETCH_SOURCE,
|
FetchType: fv1.FETCH_SOURCE,
|
||||||
Package: pkg.ObjectMeta,
|
Package: pkg.ObjectMeta,
|
||||||
@@ -121,6 +130,15 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient version
|
|||||||
return uploadResp, buildResp.BuildLogs, nil
|
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,
|
func updatePackage(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface,
|
||||||
pkg *fv1.Package, status fv1.BuildStatus, buildLogs string,
|
pkg *fv1.Package, status fv1.BuildStatus, buildLogs string,
|
||||||
uploadResp *fetcher.ArchiveUploadResponse) (*fv1.Package, error) {
|
uploadResp *fetcher.ArchiveUploadResponse) (*fv1.Package, error) {
|
||||||
|
|||||||
@@ -525,6 +525,14 @@ func (fetcher *Fetcher) UploadHandler(w http.ResponseWriter, r *http.Request) {
|
|||||||
srcFilepath := filepath.Join(fetcher.sharedVolumePath, req.Filename)
|
srcFilepath := filepath.Join(fetcher.sharedVolumePath, req.Filename)
|
||||||
dstFilepath := filepath.Join(fetcher.sharedVolumePath, zipFilename)
|
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 {
|
if req.ArchivePackage {
|
||||||
err = fetcher.archive(srcFilepath, dstFilepath)
|
err = fetcher.archive(srcFilepath, dstFilepath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -584,6 +592,7 @@ func (fetcher *Fetcher) UploadHandler(w http.ResponseWriter, r *http.Request) {
|
|||||||
logger.Error(e, zap.Error(err))
|
logger.Error(e, zap.Error(err))
|
||||||
http.Error(w, fmt.Sprintf("%s: %v", e, err), http.StatusInternalServerError)
|
http.Error(w, fmt.Sprintf("%s: %v", e, err), http.StatusInternalServerError)
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fetcher *Fetcher) rename(src string, dst string) error {
|
func (fetcher *Fetcher) rename(src string, dst string) error {
|
||||||
|
|||||||
@@ -267,3 +267,32 @@ func FindFreePort() (int, error) {
|
|||||||
|
|
||||||
return port, nil
|
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
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user