From 8442e21621a575441fd3f3dee1c09fac0a9f8c5c Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Thu, 14 Apr 2022 11:35:58 +0530 Subject: [PATCH] Use common httpserver across fission (#2409) * Defining httpserver package to capture httpserver shutdown and introduces uniform running of http server across codebase. * Add unit tests for httpserver Signed-off-by: Sanket Sudake --- cmd/builder/app/server.go | 6 ++- cmd/builder/main.go | 10 ++-- cmd/fetcher/app/server.go | 6 +-- cmd/fetcher/main.go | 3 +- cmd/fission-bundle/main.go | 5 +- pkg/controller/api.go | 7 +-- pkg/executor/api.go | 10 ++-- pkg/executor/executor.go | 2 +- pkg/mqtrigger/metrics.go | 1 - pkg/mqtrigger/mqtmanager.go | 11 +---- pkg/router/mutablemux_test.go | 17 +++---- pkg/router/router.go | 8 ++- pkg/storagesvc/client/storagesvc_test.go | 2 +- pkg/storagesvc/storagesvc.go | 10 ++-- pkg/utils/httpserver/server.go | 32 ++++++++++++ pkg/utils/httpserver/server_test.go | 62 ++++++++++++++++++++++++ pkg/utils/metrics/server.go | 26 ++-------- pkg/utils/profile/profile.go | 14 ++---- 18 files changed, 139 insertions(+), 93 deletions(-) create mode 100644 pkg/utils/httpserver/server.go create mode 100644 pkg/utils/httpserver/server_test.go diff --git a/cmd/builder/app/server.go b/cmd/builder/app/server.go index 7f599855..0fc5a316 100644 --- a/cmd/builder/app/server.go +++ b/cmd/builder/app/server.go @@ -17,15 +17,17 @@ limitations under the License. package app import ( + "context" "net/http" "go.uber.org/zap" builder "github.com/fission/fission/pkg/builder" + "github.com/fission/fission/pkg/utils/httpserver" ) // Usage: builder -func Run(logger *zap.Logger, shareVolume string) error { +func Run(ctx context.Context, logger *zap.Logger, shareVolume string) { builder := builder.MakeBuilder(logger, shareVolume) mux := http.NewServeMux() mux.HandleFunc("/", builder.Handler) @@ -33,5 +35,5 @@ func Run(logger *zap.Logger, shareVolume string) error { mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) }) - return http.ListenAndServe(":8001", mux) + httpserver.StartServer(ctx, logger, "builder", "8001", mux) } diff --git a/cmd/builder/main.go b/cmd/builder/main.go index 37684aa4..fe49fd7e 100644 --- a/cmd/builder/main.go +++ b/cmd/builder/main.go @@ -24,15 +24,15 @@ import ( "github.com/fission/fission/cmd/builder/app" "github.com/fission/fission/pkg/utils/loggerfactory" "github.com/fission/fission/pkg/utils/profile" + "github.com/fission/fission/pkg/utils/signals" ) // Usage: builder func main() { logger := loggerfactory.GetLogger() defer logger.Sync() - - profile.ProfileIfEnabled(logger) - + ctx := signals.SetupSignalHandlerWithContext(logger) + profile.ProfileIfEnabled(ctx, logger) shareVolume := os.Args[1] if _, err := os.Stat(shareVolume); err != nil { if os.IsNotExist(err) { @@ -42,7 +42,5 @@ func main() { } } } - - err := app.Run(logger, shareVolume) - logger.Error("error running builder", zap.Error(err)) + app.Run(ctx, logger, shareVolume) } diff --git a/cmd/fetcher/app/server.go b/cmd/fetcher/app/server.go index a1225a7c..48966b7a 100644 --- a/cmd/fetcher/app/server.go +++ b/cmd/fetcher/app/server.go @@ -21,7 +21,6 @@ import ( "encoding/json" "flag" "fmt" - "log" "net/http" "os" "sync/atomic" @@ -31,6 +30,7 @@ import ( "go.uber.org/zap" "github.com/fission/fission/pkg/fetcher" + "github.com/fission/fission/pkg/utils/httpserver" otelUtils "github.com/fission/fission/pkg/utils/otel" "github.com/fission/fission/pkg/utils/tracing" ) @@ -135,9 +135,7 @@ func Run(ctx context.Context, logger *zap.Logger) { } else { handler = otelUtils.GetHandlerWithOTEL(mux, "fission-fetcher", otelUtils.UrlsToIgnore("/healthz", "/readiness-healthz")) } - if err = http.ListenAndServe(":8000", handler); err != nil { - log.Fatal(err) - } + httpserver.StartServer(ctx, logger, "fetcher", "8000", handler) } func fetcherUsage() { diff --git a/cmd/fetcher/main.go b/cmd/fetcher/main.go index 25f8b8cb..6fc2cc0b 100644 --- a/cmd/fetcher/main.go +++ b/cmd/fetcher/main.go @@ -28,8 +28,7 @@ func main() { logger := loggerfactory.GetLogger() defer logger.Sync() - profile.ProfileIfEnabled(logger) - ctx := signals.SetupSignalHandlerWithContext(logger) + profile.ProfileIfEnabled(ctx, logger) app.Run(ctx, logger) } diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index 50e7355c..a7e18a0f 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -200,7 +200,8 @@ Options: logger := loggerfactory.GetLogger() defer exitWithSync(logger) - profile.ProfileIfEnabled(logger) + ctx := signals.SetupSignalHandlerWithContext(logger) + profile.ProfileIfEnabled(ctx, logger) version := fmt.Sprintf("Fission Bundle Version: %v", info.BuildInfo().String()) arguments, err := docopt.ParseArgs(usage, nil, version) @@ -209,8 +210,6 @@ Options: return } - ctx := signals.SetupSignalHandlerWithContext(logger) - openTracingEnabled := tracing.TracingEnabled(logger) if openTracingEnabled { err = tracing.RegisterTraceExporter(logger, os.Getenv("TRACE_JAEGER_COLLECTOR_ENDPOINT"), getServiceName(arguments)) diff --git a/pkg/controller/api.go b/pkg/controller/api.go index c4fea0e6..b8a44bb2 100644 --- a/pkg/controller/api.go +++ b/pkg/controller/api.go @@ -35,6 +35,7 @@ import ( ferror "github.com/fission/fission/pkg/error" "github.com/fission/fission/pkg/fission-cli/logdb" "github.com/fission/fission/pkg/info" + "github.com/fission/fission/pkg/utils/httpserver" "github.com/fission/fission/pkg/utils/metrics" "github.com/fission/fission/pkg/utils/otel" ) @@ -270,9 +271,6 @@ func (api *API) GetHandler() http.Handler { } func (api *API) Serve(ctx context.Context, port int, openTracingEnabled bool) { - address := fmt.Sprintf(":%v", port) - api.logger.Info("server started", zap.Int("port", port)) - var handler http.Handler if openTracingEnabled { handler = &ochttp.Handler{Handler: api.GetHandler()} @@ -281,6 +279,5 @@ func (api *API) Serve(ctx context.Context, port int, openTracingEnabled bool) { } go metrics.ServeMetrics(ctx, api.logger) - err := http.ListenAndServe(address, handler) - api.logger.Fatal("done listening", zap.Error(err)) + httpserver.StartServer(ctx, api.logger, "controller", fmt.Sprintf("%d", port), handler) } diff --git a/pkg/executor/api.go b/pkg/executor/api.go index b5a07222..34cfe792 100644 --- a/pkg/executor/api.go +++ b/pkg/executor/api.go @@ -34,6 +34,7 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" ferror "github.com/fission/fission/pkg/error" "github.com/fission/fission/pkg/executor/client" + "github.com/fission/fission/pkg/utils/httpserver" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -262,10 +263,7 @@ func (executor *Executor) GetHandler() http.Handler { } // Serve starts an HTTP server. -func (executor *Executor) Serve(port int, openTracingEnabled bool) { - executor.logger.Info("starting executor API", zap.Int("port", port)) - address := fmt.Sprintf(":%v", port) - +func (executor *Executor) Serve(ctx context.Context, port int, openTracingEnabled bool) { var handler http.Handler if openTracingEnabled { handler = &ochttp.Handler{Handler: executor.GetHandler()} @@ -273,6 +271,6 @@ func (executor *Executor) Serve(port int, openTracingEnabled bool) { handler = otelUtils.GetHandlerWithOTEL(executor.GetHandler(), "fission-executor", otelUtils.UrlsToIgnore("/healthz")) } - err := http.ListenAndServe(address, handler) - executor.logger.Fatal("done listening", zap.Error(err)) + httpserver.StartServer(ctx, executor.logger, "executor", fmt.Sprintf("%d", port), handler) + } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 8b7a60ae..c58b438a 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -370,7 +370,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st } go reaper.CleanupRoleBindings(ctx, logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30) go metrics.ServeMetrics(ctx, logger) - go api.Serve(port, openTracingEnabled) + go api.Serve(ctx, port, openTracingEnabled) return nil } diff --git a/pkg/mqtrigger/metrics.go b/pkg/mqtrigger/metrics.go index 4bfdb282..b8204b4f 100644 --- a/pkg/mqtrigger/metrics.go +++ b/pkg/mqtrigger/metrics.go @@ -22,7 +22,6 @@ import ( ) var ( - metricsAddr = ":8080" labels = []string{"trigger_name", "trigger_namespace"} subscriptionCount = promauto.NewGaugeVec( prometheus.GaugeOpts{ diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 2807f71c..4e3a1509 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -19,10 +19,8 @@ package mqtrigger import ( "context" "errors" - "net/http" "time" - "github.com/prometheus/client_golang/prometheus/promhttp" "go.uber.org/zap" k8sCache "k8s.io/client-go/tools/cache" @@ -30,6 +28,7 @@ import ( "github.com/fission/fission/pkg/crd" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/mqtrigger/messageQueue" + "github.com/fission/fission/pkg/utils/metrics" ) const ( @@ -88,7 +87,7 @@ func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) { if ok := k8sCache.WaitForCacheSync(ctx.Done(), mqTriggerInformer.HasSynced); !ok { mqt.logger.Fatal("failed to wait for caches to sync") } - go mqt.serveMetrics() + go metrics.ServeMetrics(ctx, mqt.logger) } func (mqt *MessageQueueTriggerManager) service() { @@ -128,12 +127,6 @@ func (mqt *MessageQueueTriggerManager) service() { } } -func (mqt *MessageQueueTriggerManager) serveMetrics() { - http.Handle("/metrics", promhttp.Handler()) - err := http.ListenAndServe(metricsAddr, nil) - mqt.logger.Fatal("done listening on metrics endpoint", zap.Error(err)) -} - func (mqt *MessageQueueTriggerManager) makeRequest(requestType requestType, triggerSub *triggerSubscription) response { respChan := make(chan response) mqt.reqChan <- request{requestType, triggerSub, respChan} diff --git a/pkg/router/mutablemux_test.go b/pkg/router/mutablemux_test.go index 344a46c7..9dd99d22 100644 --- a/pkg/router/mutablemux_test.go +++ b/pkg/router/mutablemux_test.go @@ -17,6 +17,7 @@ limitations under the License. package router import ( + "context" "log" "net/http" "testing" @@ -26,6 +27,7 @@ import ( "go.uber.org/zap" "go.uber.org/zap/zapcore" + "github.com/fission/fission/pkg/utils/httpserver" "github.com/fission/fission/pkg/utils/metrics" ) @@ -47,13 +49,6 @@ func verifyRequest(expectedResponse string) { testRequest(targetURL, expectedResponse) } -func startServer(mr *mutableRouter) { - err := http.ListenAndServe(":3333", mr) - if err != nil { - log.Fatal(err) - } -} - func spamServer(quit chan bool) { i := 0 for { @@ -83,10 +78,10 @@ func TestMutableMux(t *testing.T) { panicIf(err) mr := newMutableRouter(logger, muxRouter) - + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() // start http server - log.Print("Start http server") - go startServer(mr) + go httpserver.StartServer(ctx, logger, "router", "3333", mr) // continuously make requests, panic if any fails time.Sleep(100 * time.Millisecond) @@ -102,7 +97,7 @@ func TestMutableMux(t *testing.T) { // change the muxer log.Print("Change mux router") newMuxRouter := mux.NewRouter() - muxRouter.Use(metrics.HTTPMetricMiddleware()) + newMuxRouter.Use(metrics.HTTPMetricMiddleware()) newMuxRouter.HandleFunc("/", NewHandler) mr.updateRouter(newMuxRouter) diff --git a/pkg/router/router.go b/pkg/router/router.go index 3f7d7b7c..bd77bacc 100644 --- a/pkg/router/router.go +++ b/pkg/router/router.go @@ -58,6 +58,7 @@ import ( "github.com/fission/fission/pkg/crd" executorClient "github.com/fission/fission/pkg/executor/client" "github.com/fission/fission/pkg/throttler" + "github.com/fission/fission/pkg/utils/httpserver" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -86,7 +87,6 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger func serve(ctx context.Context, logger *zap.Logger, port int, tracingSamplingRate float64, httpTriggerSet *HTTPTriggerSet, displayAccessLog bool, openTracingEnabled bool) { mr := router(ctx, logger, httpTriggerSet) - url := fmt.Sprintf(":%v", port) var handler http.Handler if openTracingEnabled { @@ -114,10 +114,8 @@ func serve(ctx context.Context, logger *zap.Logger, port int, tracingSamplingRat } else { handler = otelUtils.GetHandlerWithOTEL(mr, "fission-router", otelUtils.UrlsToIgnore("/router-healthz")) } - err := http.ListenAndServe(url, handler) - if err != nil { - logger.Error("HTTP server error", zap.Error(err)) - } + + httpserver.StartServer(ctx, logger, "router", fmt.Sprintf("%d", port), handler) } // Start starts a router diff --git a/pkg/storagesvc/client/storagesvc_test.go b/pkg/storagesvc/client/storagesvc_test.go index 562a754a..bd36fe67 100644 --- a/pkg/storagesvc/client/storagesvc_test.go +++ b/pkg/storagesvc/client/storagesvc_test.go @@ -209,7 +209,7 @@ func TestLocalStorageService(t *testing.T) { storage := storagesvc.NewLocalStorage(localPath) ctx, cancel := context.WithCancel(context.Background()) defer cancel() - os.Setenv("METRICS_ADDR", ":8083") + os.Setenv("METRICS_ADDR", "8083") _ = storagesvc.Start(ctx, logger, storage, port, true) time.Sleep(time.Second) diff --git a/pkg/storagesvc/storagesvc.go b/pkg/storagesvc/storagesvc.go index 7500bb0c..a5105da5 100644 --- a/pkg/storagesvc/storagesvc.go +++ b/pkg/storagesvc/storagesvc.go @@ -31,6 +31,7 @@ import ( "go.opencensus.io/plugin/ochttp" "go.uber.org/zap" + "github.com/fission/fission/pkg/utils/httpserver" "github.com/fission/fission/pkg/utils/metrics" "github.com/fission/fission/pkg/utils/otel" ) @@ -214,7 +215,7 @@ func MakeStorageService(logger *zap.Logger, storageClient *StowClient, port int) } } -func (ss *StorageService) Start(port int, openTracingEnabled bool) { +func (ss *StorageService) Start(ctx context.Context, port int, openTracingEnabled bool) { r := mux.NewRouter() r.Use(metrics.HTTPMetricMiddleware()) r.HandleFunc("/v1/archive", ss.uploadHandler).Methods("POST") @@ -222,8 +223,6 @@ func (ss *StorageService) Start(port int, openTracingEnabled bool) { r.HandleFunc("/v1/archive", ss.deleteHandler).Methods("DELETE") r.HandleFunc("/healthz", ss.healthHandler).Methods("GET") - address := fmt.Sprintf(":%v", port) - var handler http.Handler if openTracingEnabled { handler = &ochttp.Handler{ @@ -232,8 +231,7 @@ func (ss *StorageService) Start(port int, openTracingEnabled bool) { } else { handler = otel.GetHandlerWithOTEL(r, "fission-storagesvc", otel.UrlsToIgnore("/healthz")) } - err := http.ListenAndServe(address, handler) - ss.logger.Fatal("done listening", zap.Error(err)) + httpserver.StartServer(ctx, ss.logger, "storagesvc", fmt.Sprintf("%d", port), handler) } // Start runs storage service @@ -248,7 +246,7 @@ func Start(ctx context.Context, logger *zap.Logger, storage Storage, port int, o // create http handlers storageService := MakeStorageService(logger, storageClient, port) go metrics.ServeMetrics(ctx, logger) - go storageService.Start(port, openTracingEnabled) + go storageService.Start(ctx, port, openTracingEnabled) // enablePruner prevents storagesvc unit test from needing to talk to kubernetes if enablePruner { diff --git a/pkg/utils/httpserver/server.go b/pkg/utils/httpserver/server.go new file mode 100644 index 00000000..4580ffe2 --- /dev/null +++ b/pkg/utils/httpserver/server.go @@ -0,0 +1,32 @@ +package httpserver + +import ( + "context" + "fmt" + "net/http" + + "go.uber.org/zap" +) + +func StartServer(ctx context.Context, log *zap.Logger, svc string, port string, handler http.Handler) { + server := http.Server{ + Addr: fmt.Sprintf(":%s", port), + Handler: handler, + } + l := log.With(zap.String("service", svc), zap.String("addr", server.Addr)) + l.Info("starting server") + go func() { + if err := server.ListenAndServe(); err != nil { + if err != http.ErrServerClosed { + l.Error("server error", zap.Error(err)) + } + } + }() + <-ctx.Done() + l.Info("shutting down server") + if err := server.Shutdown(ctx); err != nil { + if err != context.Canceled && err != context.DeadlineExceeded { + l.Error("server shutdown error", zap.Error(err)) + } + } +} diff --git a/pkg/utils/httpserver/server_test.go b/pkg/utils/httpserver/server_test.go new file mode 100644 index 00000000..ec219043 --- /dev/null +++ b/pkg/utils/httpserver/server_test.go @@ -0,0 +1,62 @@ +package httpserver + +import ( + "context" + "io/ioutil" + "net/http" + "testing" + + "github.com/gorilla/mux" + "go.uber.org/zap" + + "github.com/fission/fission/pkg/utils/loggerfactory" +) + +func TestStartServer(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + logger := loggerfactory.GetLogger() + m := mux.NewRouter() + m.Handle("/", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + _, err := w.Write([]byte("test handler")) + if err != nil { + logger.Error("failed to write response", zap.Error(err)) + } + })) + go StartServer(ctx, logger, "test", "8999", m) + + tests := []struct { + URL string + StatusCode int + Body string + }{ + { + URL: "http://localhost:8999", + StatusCode: http.StatusOK, + Body: "test handler", + }, + { + URL: "http://localhost:8999/notfound", + StatusCode: http.StatusNotFound, + Body: "404 page not found\n", + }, + } + for _, test := range tests { + resp, err := http.Get(test.URL) + if err != nil { + t.Errorf("failed to make get request %v: %v", test.URL, err) + } + defer resp.Body.Close() + if resp.StatusCode != test.StatusCode { + t.Errorf("expected status code %v, got %v", test.StatusCode, resp.StatusCode) + } + body, err := ioutil.ReadAll(resp.Body) + if err != nil { + t.Errorf("failed to read response body: %v", err) + } + if string(body) != test.Body { + t.Errorf("expected body \"%v\", got \"%v\"", test.Body, string(body)) + } + } +} diff --git a/pkg/utils/metrics/server.go b/pkg/utils/metrics/server.go index e5205160..642b7390 100644 --- a/pkg/utils/metrics/server.go +++ b/pkg/utils/metrics/server.go @@ -23,34 +23,16 @@ import ( "github.com/prometheus/client_golang/prometheus/promhttp" "go.uber.org/zap" + + "github.com/fission/fission/pkg/utils/httpserver" ) func ServeMetrics(ctx context.Context, logger *zap.Logger) { metricsAddr := os.Getenv("METRICS_ADDR") if metricsAddr == "" { - metricsAddr = ":8080" + metricsAddr = "8080" } mux := http.NewServeMux() mux.Handle("/metrics", promhttp.Handler()) - s := &http.Server{ - Addr: metricsAddr, - Handler: mux, - } - logger.Info("Starting metrics server", zap.String("address", metricsAddr)) - go func() { - if err := s.ListenAndServe(); err != nil { - if err != http.ErrServerClosed { - logger.Error("Metrics server error", zap.Error(err)) - } - } - }() - <-ctx.Done() - logger.Info("Shutting down metrics server") - err := s.Shutdown(ctx) - if err == context.DeadlineExceeded || err == context.Canceled { - return - } - if err != nil { - logger.Error("Failed to shutdown metrics server", zap.Error(err)) - } + httpserver.StartServer(ctx, logger, "metrics", metricsAddr, mux) } diff --git a/pkg/utils/profile/profile.go b/pkg/utils/profile/profile.go index 2b3f91f4..fc796c47 100644 --- a/pkg/utils/profile/profile.go +++ b/pkg/utils/profile/profile.go @@ -24,12 +24,15 @@ limitations under the License. package profile import ( + "context" "fmt" "net/http" _ "net/http/pprof" "os" "go.uber.org/zap" + + "github.com/fission/fission/pkg/utils/httpserver" ) func getPprofAddr() string { @@ -44,7 +47,7 @@ func getPprofAddr() string { return fmt.Sprintf("%s:%s", pprofHost, pprofPort) } -func ProfileIfEnabled(logger *zap.Logger) { +func ProfileIfEnabled(ctx context.Context, logger *zap.Logger) { enablePprof := os.Getenv("PPROF_ENABLED") if enablePprof != "true" { return @@ -53,12 +56,5 @@ func ProfileIfEnabled(logger *zap.Logger) { pprofMux := http.DefaultServeMux http.DefaultServeMux = http.NewServeMux() - addr := getPprofAddr() - logger.Info("Running pprof server", zap.String("addr", addr)) - go func() { - err := http.ListenAndServe(addr, pprofMux) - if err != nil { - logger.Fatal("pprof http server failed", zap.Error(err)) - } - }() + go httpserver.StartServer(ctx, logger, "pprof", getPprofAddr(), pprofMux) }