diff --git a/TODO.md b/TODO.md index 539e37e..7d2391e 100644 --- a/TODO.md +++ b/TODO.md @@ -29,6 +29,17 @@ P2: security: referer blacklist # Completed Steps +- 2026-10-03 shutdown sets the exit code and waits for image processing + (closes #86): fx alone handles SIGINT and SIGTERM, and the server's own + signal handler is gone; fx's `Run` in `cmd/pixad` exits with the shutdown's + code: 0 for a signal, 1 when the HTTP server cannot listen or the app fails + to start or to stop; the server's stop hook, which fx waits for, stops the + HTTP server, waits for the images still being processed, both within 5 + seconds, then flushes Sentry; images still being processed after that are + logged with their count and make the exit code 1; a Sentry DSN that cannot be + used fails startup, so the stop hooks of what had already started run, + instead of exiting the process from a goroutine; the eviction loop is left to + #102. - 2026-09-29 only the image routes send CORS headers (closes #98): the CORS middleware, with the `access_control_allow_origin` origin, moved from the router root onto a `/v1` subrouter holding `/v1/image/` and `/v1/e/`, where it diff --git a/cmd/pixad/main.go b/cmd/pixad/main.go index 96442dc..be56d5f 100644 --- a/cmd/pixad/main.go +++ b/cmd/pixad/main.go @@ -4,6 +4,8 @@ package main import ( "fmt" "os" + "os/signal" + "syscall" "github.com/spf13/cobra" "go.uber.org/fx" @@ -45,6 +47,9 @@ func run(_ *cobra.Command, _ []string) { _ = os.Setenv("PIXA_CONFIG_PATH", configPath) } + // A write to a closed stdout or stderr must not end the process. + signal.Ignore(syscall.SIGPIPE) + fx.New( fx.Provide( config.New, diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index b350595..fafb27b 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -81,6 +81,12 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) { return s, nil } +// WaitForProcessing waits until no image is being processed, or until ctx +// ends, and returns how many images were still being processed then. +func (s *Handlers) WaitForProcessing(ctx context.Context) int { + return s.imgSvc.WaitForProcessing(ctx) +} + // initImageService initializes the image cache and service. func (s *Handlers) initImageService() error { // Create the cache. cache_max_bytes: 0 disables the disk cache diff --git a/internal/imageprocessor/imageprocessor.go b/internal/imageprocessor/imageprocessor.go index 6a9a92d..2129ece 100644 --- a/internal/imageprocessor/imageprocessor.go +++ b/internal/imageprocessor/imageprocessor.go @@ -331,6 +331,31 @@ func FormatToMIME(format Format) string { } } +// WaitForProcessing waits until no image is being processed, or until ctx +// ends, and returns how many images were still being processed then. It +// waits by taking every slot in processingSemaphore as it frees up, so no +// new image starts meanwhile, and gives them all back before it returns. +func (p *ImageProcessor) WaitForProcessing(ctx context.Context) int { + taken := 0 + + defer func() { + for range taken { + <-p.processingSemaphore + } + }() + + for taken < cap(p.processingSemaphore) { + select { + case p.processingSemaphore <- struct{}{}: + taken++ + case <-ctx.Done(): + return len(p.processingSemaphore) - taken + } + } + + return 0 +} + // acquireSlot takes a slot in processingSemaphore, waiting at most // processingWaitTimeout for one to free up, and returns the func that gives // it back. A free slot is taken even when ctx has ended; only the wait for diff --git a/internal/imgcache/service.go b/internal/imgcache/service.go index e1152dc..209c3ed 100644 --- a/internal/imgcache/service.go +++ b/internal/imgcache/service.go @@ -195,6 +195,12 @@ func (s *Service) Stats(ctx context.Context) (*CacheStats, error) { return s.cache.Stats(ctx) } +// WaitForProcessing waits until no image is being processed, or until ctx +// ends, and returns how many images were still being processed then. +func (s *Service) WaitForProcessing(ctx context.Context) int { + return s.processor.WaitForProcessing(ctx) +} + // ValidateRequest validates the request signature if required. func (s *Service) ValidateRequest(req *ImageRequest) error { // Check if host is allowed (no signature required) diff --git a/internal/server/http.go b/internal/server/http.go index 13162e0..33ce153 100644 --- a/internal/server/http.go +++ b/internal/server/http.go @@ -5,6 +5,8 @@ import ( "fmt" "net/http" "time" + + "go.uber.org/fx" ) // HTTP server configuration constants. @@ -36,19 +38,19 @@ func (s *Server) newHTTPServer() *http.Server { } } +// serveUntilShutdown serves on s.httpServer until it is shut down. When it +// stops for any other reason, such as its port being in use, it asks fx to +// shut down with exit code 1. func (s *Server) serveUntilShutdown() { - s.httpServer = s.newHTTPServer() - - s.SetupRoutes() - s.log.Info("http begin listen", "listenaddr", s.httpServer.Addr) err := s.httpServer.ListenAndServe() if err != nil && !errors.Is(err, http.ErrServerClosed) { s.log.Error("listen error", "error", err) - if s.cancelFunc != nil { - s.cancelFunc() + err = s.shutdowner.Shutdown(fx.ExitCode(1)) + if err != nil { + s.log.Error("shutdown request failed", "error", err) } } } diff --git a/internal/server/server.go b/internal/server/server.go index bbaf44f..7be7d67 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -3,12 +3,10 @@ package server import ( "context" + "errors" "fmt" "log/slog" "net/http" - "os" - "os/signal" - "syscall" "time" "github.com/getsentry/sentry-go" @@ -27,6 +25,10 @@ const ( SentryFlushTimeout = 2 * time.Second ) +// errStillProcessing is returned by the server's stop hook when images are +// still being processed once ShutdownTimeout has passed. +var errStillProcessing = errors.New("images still being processed at shutdown") + // Params defines dependencies for Server. type Params struct { fx.In @@ -36,6 +38,7 @@ type Params struct { Config *config.Config Middleware *middleware.Middleware Handlers *handlers.Handlers + Shutdowner fx.Shutdowner } // Server is the main HTTP server. @@ -45,59 +48,58 @@ type Server struct { globals *globals.Globals mw *middleware.Middleware h *handlers.Handlers + shutdowner fx.Shutdowner startupTime time.Time - exitCode int sentryEnabled bool - cancelFunc context.CancelFunc httpServer *http.Server router *chi.Mux } -// New creates a new Server instance. +// New creates a new Server instance. Its start hook starts Sentry and the +// HTTP server; its stop hook, which fx runs on SIGINT, SIGTERM or a +// shutdown request, shuts them down. func New(lc fx.Lifecycle, params Params) (*Server, error) { s := &Server{ - log: params.Logger.Get(), - config: params.Config, - globals: params.Globals, - mw: params.Middleware, - h: params.Handlers, + log: params.Logger.Get(), + config: params.Config, + globals: params.Globals, + mw: params.Middleware, + h: params.Handlers, + shutdowner: params.Shutdowner, } lc.Append(fx.Hook{ - OnStart: func(ctx context.Context) error { + OnStart: func(_ context.Context) error { s.startupTime = time.Now() - go s.Run(context.WithoutCancel(ctx)) - return nil - }, - OnStop: func(_ context.Context) error { - if s.cancelFunc != nil { - s.cancelFunc() + err := s.enableSentry() + if err != nil { + return err } + s.SetupRoutes() + s.httpServer = s.newHTTPServer() + + go s.serveUntilShutdown() + return nil }, + OnStop: s.cleanShutdown, }) return s, nil } -// Run starts the server. -func (s *Server) Run(ctx context.Context) { - s.enableSentry() - s.serve(ctx) -} - // MaintenanceMode returns whether maintenance mode is enabled. func (s *Server) MaintenanceMode() bool { return s.config.MaintenanceMode } -func (s *Server) enableSentry() { +func (s *Server) enableSentry() error { s.sentryEnabled = false if s.config.SentryDSN == "" { - return + return nil } err := sentry.Init(sentry.ClientOptions{ @@ -105,55 +107,42 @@ func (s *Server) enableSentry() { Release: fmt.Sprintf("%s-%s", s.globals.Appname, s.globals.Version), }) if err != nil { - s.log.Error("sentry init failure", "error", err) - os.Exit(1) + return fmt.Errorf("sentry init failure: %w", err) } s.log.Info("sentry error reporting activated") s.sentryEnabled = true + + return nil } -func (s *Server) serve(ctx context.Context) int { - ctx, cancelFunc := context.WithCancel(ctx) - s.cancelFunc = cancelFunc +// cleanShutdown stops the HTTP server, waits for the images still being +// processed, then flushes Sentry. The first two share ShutdownTimeout. It +// returns errStillProcessing when images are still being processed after +// that, as their work is abandoned. +func (s *Server) cleanShutdown(ctx context.Context) error { + s.log.Info("shutting down") - go func() { - c := make(chan os.Signal, 1) - - signal.Ignore(syscall.SIGPIPE) - signal.Notify(c, os.Interrupt, syscall.SIGTERM) - - sig := <-c - s.log.Info("signal received", "signal", sig) - - if s.cancelFunc != nil { - s.cancelFunc() - } - }() - - go s.serveUntilShutdown() - - <-ctx.Done() - s.cleanShutdown(ctx) - - return s.exitCode -} - -func (s *Server) cleanShutdown(ctx context.Context) { - s.exitCode = 0 - - ctxShutdown, shutdownCancel := context.WithTimeout( - context.WithoutCancel(ctx), ShutdownTimeout) + ctxShutdown, shutdownCancel := context.WithTimeout(ctx, ShutdownTimeout) defer shutdownCancel() - if s.httpServer != nil { - err := s.httpServer.Shutdown(ctxShutdown) - if err != nil { - s.log.Error("server clean shutdown failed", "error", err) - } + err := s.httpServer.Shutdown(ctxShutdown) + if err != nil { + s.log.Error("server clean shutdown failed", "error", err) } + stillProcessing := s.h.WaitForProcessing(ctxShutdown) + if s.sentryEnabled { sentry.Flush(SentryFlushTimeout) } + + if stillProcessing > 0 { + s.log.Error("images still being processed at shutdown", + "count", stillProcessing) + + return errStillProcessing + } + + return nil }