Exit with the shutdown's code and wait for image processing (closes #86) #169

已合并
clawbot 于 2026-10-04 04:59:31 +02:00 将 2 次代码提交从 issue-86-shutdown-correctness合并至 next
共修改 7 个文件,包含 113 行新增和 69 行删除
仅显示提交 464ba8ccda 的更改 - 显示所有提交
+11
查看文件
@@ -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,
+5
查看文件
@@ -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,
+6
查看文件
@@ -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
+25
查看文件
@@ -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
+6
查看文件
@@ -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)
+8 -6
查看文件
@@ -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
}