From 00da62db2cdd4b20512e0638e331ab14d31b8c49 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Sun, 4 Oct 2026 04:59:30 +0200 Subject: [PATCH] Exit with the shutdown's code and wait for image processing (closes #86) fx alone handles SIGINT and SIGTERM; the server's own handler, which only cancelled a context that fx's stop did not wait for, is gone. fx's Run exits with the shutdown's code: the one a shutdown request carries, 0 for a signal, 1 when the app fails to start or stop. A listen error asks fx to shut down with exit code 1. A Sentry DSN that cannot be used fails the server's start hook, so fx stops what had already started. The server's stop hook stops the HTTP server, then waits for the images still being processed, both within ShutdownTimeout; images still being processed after that are logged and fail the stop, so the exit code is 1. Model: opus-5-5 --- TODO.md | 11 ++ cmd/pixad/main.go | 5 + internal/handlers/handlers.go | 6 + internal/imageprocessor/imageprocessor.go | 25 ++++ ...max_concurrent_processing_internal_test.go | 66 ++++++++++ internal/imgcache/service.go | 6 + internal/server/http.go | 14 ++- internal/server/server.go | 115 ++++++++---------- internal/server/shutdown_internal_test.go | 96 +++++++++++++++ 9 files changed, 275 insertions(+), 69 deletions(-) create mode 100644 internal/server/shutdown_internal_test.go diff --git a/TODO.md b/TODO.md index cebef13..ed8431c 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-10-03 every `script/cibuild` and `script/docker` run executes the checks (closes #101): the `Dockerfile` declares `CHECK_EPOCH` above `make fmt-check` and `make lint` in the lint stage and above `make test` in the build stage, 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..f26ed57 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 each slot in processingSemaphore as it frees up until it +// holds them all, or until ctx ends, then gives back the slots it took. +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/imageprocessor/max_concurrent_processing_internal_test.go b/internal/imageprocessor/max_concurrent_processing_internal_test.go index dc826c4..ba5388c 100644 --- a/internal/imageprocessor/max_concurrent_processing_internal_test.go +++ b/internal/imageprocessor/max_concurrent_processing_internal_test.go @@ -300,3 +300,69 @@ func TestProcessReleasesSlotOnError(t *testing.T) { }) } } + +// TestWaitForProcessing holds a processing slot with a Process call that +// cannot finish reading its input. WaitForProcessing must report that image +// when its context ends first, wait for it otherwise, return 0 once it has +// finished, and give back the slots it took while waiting. +func TestWaitForProcessing(t *testing.T) { + t.Parallel() + + proc := New(Params{MaxConcurrentProcessing: 2}) + + gate := make(chan struct{}) + entered := make(chan struct{}, 1) + results := make(chan error, 1) + + openGate := sync.OnceFunc(func() { close(gate) }) + t.Cleanup(openGate) + + processInBackground(proc, &gatedReader{ + data: bytes.NewReader(createTestJPEG(t, 10, 10)), gate: gate, + entered: entered, counter: &readingCounter{}, + }, results) + waitForEntries(t, entered, 1) + + ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond) + defer cancel() + + stillProcessing := proc.WaitForProcessing(ctx) + t.Logf("WaitForProcessing() after its context ended: %d", stillProcessing) + + if stillProcessing != 1 { + t.Errorf("WaitForProcessing() after its context ended = %d, want 1", + stillProcessing) + } + + waited := make(chan int, 1) + + go func() { waited <- proc.WaitForProcessing(t.Context()) }() + + select { + case got := <-waited: + t.Fatalf("WaitForProcessing() = %d while an image was being processed", + got) + case <-time.After(100 * time.Millisecond): + } + + openGate() + + err := <-results + if err != nil { + t.Errorf("Process() error = %v, want nil", err) + } + + select { + case got := <-waited: + if got != 0 { + t.Errorf("WaitForProcessing() once processing finished = %d, want 0", + got) + } + case <-time.After(5 * time.Second): + t.Fatal("WaitForProcessing() did not return once processing finished") + } + + if held := len(proc.processingSemaphore); held != 0 { + t.Errorf("%d slots still held after WaitForProcessing() returned", held) + } +} diff --git a/internal/imgcache/service.go b/internal/imgcache/service.go index 90acf0b..77c7ce9 100644 --- a/internal/imgcache/service.go +++ b/internal/imgcache/service.go @@ -199,6 +199,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 } diff --git a/internal/server/shutdown_internal_test.go b/internal/server/shutdown_internal_test.go new file mode 100644 index 0000000..4087bb1 --- /dev/null +++ b/internal/server/shutdown_internal_test.go @@ -0,0 +1,96 @@ +package server + +import ( + "log/slog" + "net" + "testing" + "time" + + "go.uber.org/fx" + "go.uber.org/fx/fxtest" + + "sneak.berlin/go/pixa/internal/config" + "sneak.berlin/go/pixa/internal/globals" + "sneak.berlin/go/pixa/internal/logger" +) + +// shutdownRecorder is an fx.Shutdowner that sends the options of each +// shutdown request on requests. +type shutdownRecorder struct { + requests chan []fx.ShutdownOption +} + +func (r shutdownRecorder) Shutdown(opts ...fx.ShutdownOption) error { + r.requests <- opts + + return nil +} + +// TestSentryInitFailureFailsStartup checks that a Sentry DSN that cannot be +// used makes the server's start hook fail, so fx stops what has already +// started, instead of the process exiting from a goroutine. +func TestSentryInitFailureFailsStartup(t *testing.T) { + t.Parallel() + + lc := fxtest.NewLifecycle(t) + + log, err := logger.New(lc, logger.Params{Globals: &globals.Globals{}}) + if err != nil { + t.Fatalf("logger.New() error = %v", err) + } + + _, err = New(lc, Params{ + Logger: log, + Globals: &globals.Globals{Appname: "pixad"}, + Config: &config.Config{SentryDSN: "not-a-dsn"}, + }) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + err = lc.Start(t.Context()) + t.Logf("Start() error = %v", err) + + if err == nil { + t.Fatal("Start() error = nil, want the Sentry initialization error") + } +} + +// TestListenErrorRequestsShutdownWithExitCode1 occupies the server's port +// and checks that the listen error asks fx to shut down with exit code 1. +func TestListenErrorRequestsShutdownWithExitCode1(t *testing.T) { + t.Parallel() + + busy, err := (&net.ListenConfig{}).Listen(t.Context(), "tcp", ":0") + if err != nil { + t.Fatalf("Listen() error = %v", err) + } + + t.Cleanup(func() { _ = busy.Close() }) + + addr, ok := busy.Addr().(*net.TCPAddr) + if !ok { + t.Fatalf("listener address %v is not a TCP address", busy.Addr()) + } + + requests := make(chan []fx.ShutdownOption, 1) + s := &Server{ + log: slog.New(slog.DiscardHandler), + config: &config.Config{Port: addr.Port}, + shutdowner: shutdownRecorder{requests: requests}, + } + s.httpServer = s.newHTTPServer() + + go s.serveUntilShutdown() + + select { + case opts := <-requests: + t.Logf("shutdown options = %v", opts) + + if len(opts) != 1 || opts[0] != fx.ExitCode(1) { + t.Errorf("shutdown options = %v, want [fx.ExitCode(1)]", opts) + } + case <-time.After(5 * time.Second): + t.Fatal("no shutdown was requested after the listen error") + } +}