Archive pruner (#471)
All functions have a pkg reference. This can be a package with either source and a deploy archives, or, a deploy archive. Everytime a function is updated, a new package is created. With archive pruner, the archives that are pointed to by old pkg reference can be deleted from the storage. * High level spec for package pruning. * Skeleton for archive pruning * Adding meat 1 to skeleton. * Adding meat #2. Separated storage service into a httpHandler component and Storage Layer component. * Adding meat #3. getOrphanedArchives in pruner and getItems on stowClient. * Restructured archivePruner methods. * Commiting the day's work. Ready for testing #1. * Fixing compile errors. * Test ready. added a few logs for debugging. * Adding a filter for getItems in stowClient. * After testing. * Added a test for archivePruner. * Adding helm value pruneInterval for testing. * Modified test. * Final test. * Fixing interval from seconds to minutes. * Small change. * Changing debugs to info. * Removing the WIP design * Ran gofmt on all these files. * Fixing prune_interval as string in ENV var. * Addressing all comments, but one. * changing getFile method in stowClient to stream it into a response. * All comments incorporated. * Introducing a new flag for running archivePruner. 1. This flag is disabled for archivePruner to run in unit test. 2. This flag is enabled for archivePruner to run in production. 3. Also disabling test_archive_pruner.sh in this PR. Follow up with next PR to enable it. * Addressing review comments. * Changing the command to generate a file dynamically. * Enabling arching_pruner_test * giving execute permissions to test_archive_pruner.sh * Making changes of positional parameters after recent commit. Change test case permission and removing kubectlPortForward. * Adding debug to see why test_utils.sh passed junk pruneInterval. * shell needs special handling for positional parameters from 10.
This commit is contained in:
committed by
Ta-Ching Chen
parent
0c92f59c38
commit
31ba992726
+50
-104
@@ -20,33 +20,21 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/handlers"
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/graymeta/stow"
|
||||
_ "github.com/graymeta/stow/local"
|
||||
"github.com/satori/go.uuid"
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
type (
|
||||
StorageType string
|
||||
storageConfig struct {
|
||||
storageType StorageType
|
||||
localPath string
|
||||
containerName string
|
||||
// other stuff, such as google or s3 credentials, bucket names etc
|
||||
}
|
||||
|
||||
StorageService struct {
|
||||
config storageConfig
|
||||
location stow.Location
|
||||
container stow.Container
|
||||
port int
|
||||
storageClient *StowClient
|
||||
port int
|
||||
}
|
||||
|
||||
UploadResponse struct {
|
||||
@@ -54,10 +42,6 @@ type (
|
||||
}
|
||||
)
|
||||
|
||||
const (
|
||||
StorageTypeLocal StorageType = "local"
|
||||
)
|
||||
|
||||
// Handle multipart file uploads.
|
||||
func (ss *StorageService) uploadHandler(w http.ResponseWriter, r *http.Request) {
|
||||
// handle upload
|
||||
@@ -76,39 +60,34 @@ func (ss *StorageService) uploadHandler(w http.ResponseWriter, r *http.Request)
|
||||
|
||||
fileSizeS, ok := r.Header["X-File-Size"]
|
||||
if !ok {
|
||||
log.Printf("Missing X-File-Size")
|
||||
log.Error("Missing X-File-Size")
|
||||
http.Error(w, "missing X-File-Size header", 400)
|
||||
return
|
||||
}
|
||||
|
||||
fileSize, err := strconv.Atoi(fileSizeS[0])
|
||||
if err != nil {
|
||||
log.Printf("Error parsing x-file-size: '%v'", fileSizeS)
|
||||
log.WithError(err).Errorf("Error parsing x-file-size: '%v'", fileSizeS)
|
||||
http.Error(w, "missing or bad X-File-Size header", 400)
|
||||
return
|
||||
}
|
||||
|
||||
// TODO: allow headers to add more metadata (e.g. environment
|
||||
// and function metadata)
|
||||
log.Printf("Handling upload for %v", handler.Filename)
|
||||
log.Infof("Handling upload for %v", handler.Filename)
|
||||
//fileMetadata := make(map[string]interface{})
|
||||
//fileMetadata["filename"] = handler.Filename
|
||||
|
||||
// This is not the item ID (that's returned by Put)
|
||||
// should we just use handler.Filename? what are the constraints here?
|
||||
uploadName := uuid.NewV4().String()
|
||||
|
||||
// save the file to the storage backend
|
||||
item, err := ss.container.Put(uploadName, file, int64(fileSize), nil)
|
||||
id, err := ss.storageClient.putFile(file, int64(fileSize))
|
||||
if err != nil {
|
||||
log.Printf("Error saving uploaded file: '%v'", err)
|
||||
http.Error(w, "Error saving uploaded file", 400)
|
||||
log.WithError(err).Error("Error saving uploaded file")
|
||||
http.Error(w, "Error saving uploaded file", 500)
|
||||
return
|
||||
}
|
||||
|
||||
// respond with an ID that can be used to retrieve the file
|
||||
ur := &UploadResponse{
|
||||
ID: item.ID(),
|
||||
ID: id,
|
||||
}
|
||||
resp, err := json.Marshal(ur)
|
||||
if err != nil {
|
||||
@@ -135,7 +114,7 @@ func (ss *StorageService) deleteHandler(w http.ResponseWriter, r *http.Request)
|
||||
return
|
||||
}
|
||||
|
||||
err = ss.container.RemoveItem(fileId)
|
||||
err = ss.storageClient.removeFileByID(fileId)
|
||||
if err != nil {
|
||||
msg := fmt.Sprintf("Error deleting item: %v", err)
|
||||
http.Error(w, msg, 500)
|
||||
@@ -154,75 +133,27 @@ func (ss *StorageService) downloadHandler(w http.ResponseWriter, r *http.Request
|
||||
|
||||
// Get the file (called "item" in stow's jargon), open it,
|
||||
// stream it to response
|
||||
|
||||
item, err := ss.container.Item(fileId)
|
||||
err = ss.storageClient.copyFileToStream(fileId, w)
|
||||
if err != nil {
|
||||
log.Printf("Error getting item id '%v': %v", fileId, err)
|
||||
if err == stow.ErrNotFound {
|
||||
log.WithError(err).Errorf("Error getting item id '%v'", fileId)
|
||||
if err == ErrNotFound {
|
||||
http.Error(w, "Error retrieving item: not found", 404)
|
||||
} else {
|
||||
} else if err == ErrRetrievingItem {
|
||||
http.Error(w, "Error retrieving item", 400)
|
||||
} else if err == ErrOpeningItem {
|
||||
http.Error(w, "Error opening item", 400)
|
||||
} else if err == ErrWritingFileIntoResponse {
|
||||
http.Error(w, "Error writing response", 500)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
f, err := item.Open()
|
||||
if err != nil {
|
||||
log.Printf("Error opening item %v: %v", fileId, err)
|
||||
// TODO better http errors based on err
|
||||
http.Error(w, "Error opening item", 400)
|
||||
return
|
||||
}
|
||||
defer f.Close()
|
||||
|
||||
_, err = io.Copy(w, f)
|
||||
if err != nil {
|
||||
log.Printf("Error writing response: %v", err)
|
||||
http.Error(w, "Error writing response", 500)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
func MakeStorageService(sc *storageConfig) (*StorageService, error) {
|
||||
ss := &StorageService{
|
||||
config: *sc,
|
||||
func MakeStorageService(storageClient *StowClient, port int) *StorageService {
|
||||
return &StorageService{
|
||||
storageClient: storageClient,
|
||||
port: port,
|
||||
}
|
||||
|
||||
if sc.storageType != StorageTypeLocal {
|
||||
return nil, errors.New("Storage types other than 'local' are not implemented")
|
||||
}
|
||||
|
||||
cfg := stow.ConfigMap{"path": sc.localPath}
|
||||
loc, err := stow.Dial("local", cfg)
|
||||
if err != nil {
|
||||
log.Printf("Error initializing storage: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
ss.location = loc
|
||||
|
||||
con, err := loc.CreateContainer(sc.containerName)
|
||||
if os.IsExist(err) {
|
||||
var cons []stow.Container
|
||||
var cursor string
|
||||
|
||||
// use location.Containers to find containers that match the prefix (container name)
|
||||
cons, cursor, err = loc.Containers(sc.containerName, stow.CursorStart, 1)
|
||||
if err == nil {
|
||||
if !stow.IsCursorEnd(cursor) {
|
||||
// Should only have one storage container
|
||||
err = errors.New("Found more than one matched storage containers")
|
||||
} else {
|
||||
con = cons[0]
|
||||
}
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
log.Printf("Error initializing storage: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
ss.container = con
|
||||
|
||||
return ss, nil
|
||||
}
|
||||
|
||||
func (ss *StorageService) Start(port int) {
|
||||
@@ -235,19 +166,34 @@ func (ss *StorageService) Start(port int) {
|
||||
log.Fatal(http.ListenAndServe(address, handlers.LoggingHandler(os.Stdout, r)))
|
||||
}
|
||||
|
||||
func RunStorageService(storageType StorageType, storagePath string, containerName string, port int) *StorageService {
|
||||
// storage
|
||||
ss, err := MakeStorageService(&storageConfig{
|
||||
storageType: storageType,
|
||||
localPath: storagePath,
|
||||
containerName: containerName,
|
||||
})
|
||||
func RunStorageService(storageType StorageType, storagePath string, containerName string, port int, enablePruner bool) *StorageService {
|
||||
// initialize logger
|
||||
log.SetLevel(log.InfoLevel)
|
||||
|
||||
// create a storage client
|
||||
storageClient, err := MakeStowClient(storageType, storagePath, containerName)
|
||||
if err != nil {
|
||||
log.Panicf("Error initializing storage: %v", err)
|
||||
log.Fatalf("Error creating stowClient: %v", err)
|
||||
}
|
||||
|
||||
// http handlers
|
||||
go ss.Start(port)
|
||||
// create http handlers
|
||||
storageService := MakeStorageService(storageClient, port)
|
||||
go storageService.Start(port)
|
||||
|
||||
return ss
|
||||
// enablePruner prevents storagesvc unit test from needing to talk to kubernetes
|
||||
if enablePruner {
|
||||
// get the prune interval and start the archive pruner
|
||||
pruneInterval, err := strconv.Atoi(os.Getenv("PRUNE_INTERVAL"))
|
||||
if err != nil {
|
||||
pruneInterval = defaultPruneInterval
|
||||
}
|
||||
pruner, err := MakeArchivePruner(storageClient, time.Duration(pruneInterval))
|
||||
if err != nil {
|
||||
log.Fatalf("Error creating archivePruner: %v", err)
|
||||
}
|
||||
go pruner.Start()
|
||||
}
|
||||
|
||||
log.Info("Storage service started")
|
||||
return storageService
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user