Author SHA1 Message Date
sneak 4f513bede8 Exit with the shutdown's code and wait for image processing (closes #86)
check / check (push) Waiting to run
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
2026-10-03 14:42:41 +00:00
clawbot b6b48bc94a Test shutdown exit codes, Sentry startup failure and the processing wait
These tests fail until the change that follows: runApp and
WaitForProcessing do not exist yet, the server has no shutdowner, and a
Sentry DSN that cannot be used exits the process from a goroutine
instead of failing the server's start hook.

runApp must return the exit code a shutdown request carries, 0 without
one, and 1 when the app fails to start or stop. A listen error must ask
fx to shut down with exit code 1. WaitForProcessing must wait for an
image being processed and report it when its context ends first.

Model: opus-5-5
2026-10-03 14:06:56 +00:00
clawbot 869b5ba67f Stamp the tag or short commit in a plain docker build (closes #166)
check / check (push) Successful in 4m4s
.dockerignore now lets .git into the build context, without
.git/config, which can hold a remote URL with a credential. ARG VERSION
has no default: given none, the build stage takes the version from
git describe --tags --always, and fails if the context carries .git
and no version comes out. pixad now logs its name, version and
architecture as its first log line, through the existing
Logger.Identify, which nothing called.

Model: opus-5-5
2026-10-02 06:01:40 +02:00
11 changed files with 393 additions and 72 deletions
+4 -1
View File
@@ -1,4 +1,7 @@
# .git is sent so the build can stamp the version, without its config.
# .git is sent without its config. Without a VERSION build argument the
# stage that compiles runs `git describe --tags --always` on .git, which
# does not need .git/config; that file can hold a credential, such as a
# password in a remote URL or the token the CI checkout step stores there.
.git/config
.gitignore
.DS_Store
+11
View File
@@ -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
View File
@@ -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
}
+80
View File
@@ -0,0 +1,80 @@
package main
import (
"errors"
"testing"
"go.uber.org/fx"
)
// errTestHook is the error returned by the test hooks that fail.
var errTestHook = errors.New("test hook failed")
// TestRunAppExitCode checks the exit code runApp returns: the one a
// shutdown request carries, 0 for a request without one (as for SIGINT or
// SIGTERM), and 1 when the app fails to start or to stop.
func TestRunAppExitCode(t *testing.T) {
t.Parallel()
cases := []struct {
name string
hook func(shutdowner fx.Shutdowner) fx.Hook
want int
}{
{
name: "shutdown requested with exit code 1",
hook: func(shutdowner fx.Shutdowner) fx.Hook {
return fx.StartHook(func() error {
return shutdowner.Shutdown(fx.ExitCode(1))
})
},
want: 1,
},
{
name: "shutdown requested without an exit code",
hook: func(shutdowner fx.Shutdowner) fx.Hook {
return fx.StartHook(func() error {
return shutdowner.Shutdown()
})
},
want: 0,
},
{
name: "start fails",
hook: func(fx.Shutdowner) fx.Hook {
return fx.StartHook(func() error { return errTestHook })
},
want: 1,
},
{
name: "stop fails",
hook: func(shutdowner fx.Shutdowner) fx.Hook {
return fx.StartStopHook(
func() error { return shutdowner.Shutdown() },
func() error { return errTestHook },
)
},
want: 1,
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
app := fx.New(
fx.NopLogger,
fx.Invoke(func(lc fx.Lifecycle, shutdowner fx.Shutdowner) {
lc.Append(tc.hook(shutdowner))
}),
)
got := runApp(app)
t.Logf("runApp() = %d", got)
if got != tc.want {
t.Errorf("runApp() = %d, want %d", got, tc.want)
}
})
}
}
+6
View File
@@ -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
View File
@@ -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
@@ -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)
}
}
+6
View File
@@ -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)
+8 -6
View File
@@ -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
View File
@@ -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
}
+96
View File
@@ -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")
}
}