added manger to keep track of go routines in the services (#2869)
- added manager to wait for all go routines to end before exit - code refactor - renamed Manafer to Interface and GoRoutineManager to GroupManager - replaced some go routine calls with manager Add func - added unit tests for manager
This commit is contained in:
@@ -35,6 +35,7 @@ import (
|
||||
|
||||
"github.com/fission/fission/pkg/crd"
|
||||
"github.com/fission/fission/pkg/storagesvc"
|
||||
"github.com/fission/fission/pkg/utils/manager"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -86,6 +87,9 @@ func TestS3StorageService(t *testing.T) {
|
||||
fmt.Println("Test S3 Storage service")
|
||||
var minioClient *minio.Client
|
||||
|
||||
mgr := manager.New()
|
||||
defer mgr.Wait()
|
||||
|
||||
// Start minio docker container
|
||||
pool, err := dockertest.NewPool("")
|
||||
resource := runMinioDockerContainer(pool)
|
||||
@@ -139,7 +143,7 @@ func TestS3StorageService(t *testing.T) {
|
||||
storage := storagesvc.NewS3Storage()
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
_ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, port)
|
||||
_ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, mgr, port)
|
||||
|
||||
time.Sleep(time.Second)
|
||||
client := MakeClient(fmt.Sprintf("http://localhost:%v/", 8081))
|
||||
@@ -205,6 +209,9 @@ func TestLocalStorageService(t *testing.T) {
|
||||
testID := uniuri.NewLen(8)
|
||||
port := 8082
|
||||
|
||||
mgr := manager.New()
|
||||
defer mgr.Wait()
|
||||
|
||||
config := zap.NewDevelopmentConfig()
|
||||
config.EncoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder
|
||||
logger, err := config.Build()
|
||||
@@ -217,7 +224,7 @@ func TestLocalStorageService(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
os.Setenv("METRICS_ADDR", "8083")
|
||||
_ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, port)
|
||||
_ = storagesvc.Start(ctx, crd.NewClientGenerator(), logger, storage, mgr, port)
|
||||
|
||||
time.Sleep(time.Second)
|
||||
client := MakeClient(fmt.Sprintf("http://localhost:%v/", port))
|
||||
|
||||
@@ -32,6 +32,7 @@ import (
|
||||
|
||||
"github.com/fission/fission/pkg/crd"
|
||||
"github.com/fission/fission/pkg/utils/httpserver"
|
||||
"github.com/fission/fission/pkg/utils/manager"
|
||||
"github.com/fission/fission/pkg/utils/metrics"
|
||||
otelUtils "github.com/fission/fission/pkg/utils/otel"
|
||||
)
|
||||
@@ -267,7 +268,7 @@ func MakeStorageService(logger *zap.Logger, storageClient *StowClient, port int)
|
||||
}
|
||||
}
|
||||
|
||||
func (ss *StorageService) Start(ctx context.Context, port int) {
|
||||
func (ss *StorageService) Start(ctx context.Context, mgr manager.Interface, port int) {
|
||||
r := mux.NewRouter()
|
||||
r.Use(metrics.HTTPMetricMiddleware)
|
||||
r.HandleFunc("/v1/archive", ss.uploadHandler).Methods("POST")
|
||||
@@ -278,11 +279,11 @@ func (ss *StorageService) Start(ctx context.Context, port int) {
|
||||
r.HandleFunc("/healthz", ss.healthHandler).Methods("GET")
|
||||
|
||||
handler := otelUtils.GetHandlerWithOTEL(r, "fission-storagesvc", otelUtils.UrlsToIgnore("/healthz"))
|
||||
httpserver.StartServer(ctx, ss.logger, "storagesvc", fmt.Sprintf("%d", port), handler)
|
||||
httpserver.StartServer(ctx, ss.logger, mgr, "storagesvc", fmt.Sprintf("%d", port), handler)
|
||||
}
|
||||
|
||||
// Start runs storage service
|
||||
func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storage Storage, port int) error {
|
||||
func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, storage Storage, mgr manager.Interface, port int) error {
|
||||
enablePruner, err := strconv.ParseBool(os.Getenv("PRUNE_ENABLED"))
|
||||
if err != nil {
|
||||
logger.Warn("PRUNE_ENABLED value not set. Enabling archive pruner by default.", zap.Error(err))
|
||||
@@ -296,8 +297,13 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *
|
||||
|
||||
// create http handlers
|
||||
storageService := MakeStorageService(logger, storageClient, port)
|
||||
go metrics.ServeMetrics(ctx, "storagesvc", logger)
|
||||
go storageService.Start(ctx, port)
|
||||
mgr.Add(ctx, func(ctx context.Context) {
|
||||
metrics.ServeMetrics(ctx, "storagesvc", logger, mgr)
|
||||
})
|
||||
|
||||
mgr.Add(ctx, func(ctx context.Context) {
|
||||
storageService.Start(ctx, mgr, port)
|
||||
})
|
||||
|
||||
// enablePruner prevents storagesvc unit test from needing to talk to kubernetes
|
||||
if enablePruner {
|
||||
|
||||
Reference in New Issue
Block a user