Exit with the shutdown's code and wait for image processing (closes #86)
check / check (push) Successful in 5m14s
check / check (push) Successful in 5m14s
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. runApp starts the fx app, waits for a signal or a shutdown request, stops the app and returns the exit code, which main passes to os.Exit. 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
This commit is contained in:
@@ -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; `cmd/pixad` starts the fx app, waits for a signal or
|
||||
a shutdown request, stops the app and exits with its 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-02 a plain `docker build .` stamps the tag or short commit, not
|
||||
`dev` (closes #166): `.dockerignore` lets `.git` into the build context,
|
||||
without `.git/config`; with no `VERSION` build argument the `Dockerfile`
|
||||
|
||||
+39
-2
@@ -2,8 +2,11 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
"go.uber.org/fx"
|
||||
@@ -45,7 +48,10 @@ func run(_ *cobra.Command, _ []string) {
|
||||
_ = os.Setenv("PIXA_CONFIG_PATH", configPath)
|
||||
}
|
||||
|
||||
fx.New(
|
||||
// A write to a closed stdout or stderr must not end the process.
|
||||
signal.Ignore(syscall.SIGPIPE)
|
||||
|
||||
app := fx.New(
|
||||
fx.Provide(
|
||||
config.New,
|
||||
database.New,
|
||||
@@ -60,5 +66,36 @@ func run(_ *cobra.Command, _ []string) {
|
||||
func(log *logger.Logger) { log.Identify() },
|
||||
func(*server.Server) {},
|
||||
),
|
||||
).Run()
|
||||
)
|
||||
|
||||
os.Exit(runApp(app))
|
||||
}
|
||||
|
||||
// runApp starts app, waits for SIGINT, SIGTERM or a shutdown request, then
|
||||
// stops app. It returns the exit code: the one the shutdown request carries,
|
||||
// 0 for a signal, or 1 when app fails to start or to stop; fx logs that
|
||||
// error itself. It does what fx's App.Run does, but returns the exit code
|
||||
// instead of exiting.
|
||||
func runApp(app *fx.App) int {
|
||||
startCtx, cancelStart := context.WithTimeout(
|
||||
context.Background(), app.StartTimeout())
|
||||
defer cancelStart()
|
||||
|
||||
err := app.Start(startCtx)
|
||||
if err != nil {
|
||||
return 1
|
||||
}
|
||||
|
||||
shutdown := <-app.Wait()
|
||||
|
||||
stopCtx, cancelStop := context.WithTimeout(
|
||||
context.Background(), app.StopTimeout())
|
||||
defer cancelStop()
|
||||
|
||||
err = app.Stop(stopCtx)
|
||||
if err != nil {
|
||||
return 1
|
||||
}
|
||||
|
||||
return shutdown.ExitCode
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+52
-63
@@ -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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user