4 Commits
Author SHA1 Message Date
clawbot d98d4dc119 Test a fetch whose context ends waiting for a shared connection
check / check (push) Successful in 11m43s
A fetch from a host with nothing open waits for the connection shared by
all hosts, and its context ends long before the wait timeout; once every
response is closed, no semaphore may be left in hostSems.

Model: opus-5-5
2026-10-04 03:47:37 +00:00
clawbot c9867d801a Remove idle host semaphores and delete .meta with its variant (closes #87)
Each upstream host's semaphore now counts the fetches holding or waiting
for one of its slots, and is removed from hostSems when the last of them
gives its slot back or stops waiting, so a long-running pixad no longer
keeps one semaphore per host it ever fetched from. The semLen test helper
reads hostSems directly, as getHostSemaphore now counts its caller.

VariantStorage.Delete removes the variant's .meta file too, a missing one
not being an error; DeleteWithMeta, which eviction called for that, is
gone.

Model: opus-5-5
2026-10-04 03:43:09 +00:00
clawbot 96be48f127 Test that idle host semaphores and .meta files are removed
New tests, failing before the fix: fetches from many hosts, and fetches
that end without a connection, must leave no semaphore in hostSems once
they finish; VariantStorage.Delete must remove the variant's .meta file,
and must succeed when that file is already missing.

Model: opus-5-5
2026-10-04 03:42:53 +00:00
clawbot 00da62db2c Exit with the shutdown's code and wait for image processing (closes #86)
check / check (push) Successful in 9m56s
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
2026-10-04 04:59:30 +02:00
15 changed files with 510 additions and 92 deletions
+18
View File
@@ -29,6 +29,24 @@ P2: security: referer blacklist
# Completed Steps
- 2026-10-04 upstream host semaphores and variant `.meta` files no longer
outlive their use (closes #87): the fetcher counts the fetches holding or
waiting for a slot of each upstream host's semaphore and removes the host's
semaphore once none is left, so fetches from many hosts no longer leave one
semaphore each until restart; `VariantStorage.Delete` removes the variant's
`.meta` file along with it, a missing `.meta` file not being an error, and
`DeleteWithMeta`, which eviction called for that, is gone.
- 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
View File
@@ -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
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
+17 -1
View File
@@ -191,7 +191,23 @@ func fetchBody(t *testing.T, f *HTTPFetcher, path string) string {
// semLen reports how many per-host semaphore slots are currently held.
func semLen(f *HTTPFetcher, host string) int {
return len(f.getHostSemaphore(host))
f.hostSemMu.Lock()
defer f.hostSemMu.Unlock()
sem, ok := f.hostSems[host]
if !ok {
return 0
}
return len(sem.slots)
}
// hostSemCount reports how many hosts have a semaphore in hostSems.
func hostSemCount(f *HTTPFetcher) int {
f.hostSemMu.Lock()
defer f.hostSemMu.Unlock()
return len(f.hostSems)
}
func TestFetchRedirectToPrivateIPBlocked(t *testing.T) {
+44 -8
View File
@@ -160,10 +160,13 @@ func DefaultConfig() *Config {
// HTTPFetcher implements Fetcher with SSRF protection and connection limits
// per host and for all hosts together.
type HTTPFetcher struct {
client *http.Client
config *Config
hostSems map[string]chan struct{} // per-host semaphores
hostSemMu sync.Mutex // protects hostSems map
client *http.Client
config *Config
// hostSems holds the semaphore of each host with a fetch holding or
// waiting for one of its slots; the entry is removed when the host's
// last such fetch gives its slot back or stops waiting.
hostSems map[string]*hostSemaphore
hostSemMu sync.Mutex // protects hostSems and each entry's count
// allHostsSemaphore has one slot per connection allowed to all hosts
// together (config.MaxConnections).
allHostsSemaphore chan struct{}
@@ -171,6 +174,14 @@ type HTTPFetcher struct {
connectionWaitTimeout time.Duration
}
// hostSemaphore is one host's connection slots
// (config.MaxConnectionsPerHost) and the number of fetches holding or
// waiting for one of them.
type hostSemaphore struct {
slots chan struct{}
count int
}
// New creates a new HTTPFetcher with SSRF protection.
func New(config *Config) *HTTPFetcher {
if config == nil {
@@ -211,7 +222,7 @@ func New(config *Config) *HTTPFetcher {
return &HTTPFetcher{
client: client,
config: config,
hostSems: make(map[string]chan struct{}),
hostSems: make(map[string]*hostSemaphore),
allHostsSemaphore: make(chan struct{}, config.MaxConnections),
connectionWaitTimeout: ConnectionWaitTimeout,
}
@@ -307,6 +318,8 @@ func (f *HTTPFetcher) acquireConnection(
select {
case hostSem <- struct{}{}:
case <-ctx.Done():
f.putHostSemaphore(host)
return nil, ctx.Err()
}
@@ -314,32 +327,55 @@ func (f *HTTPFetcher) acquireConnection(
case f.allHostsSemaphore <- struct{}{}:
case <-time.After(f.connectionWaitTimeout):
<-hostSem
f.putHostSemaphore(host)
return nil, ErrTooManyConnections
case <-ctx.Done():
<-hostSem
f.putHostSemaphore(host)
return nil, ctx.Err()
}
return func() {
<-hostSem
f.putHostSemaphore(host)
<-f.allHostsSemaphore
}, nil
}
// getHostSemaphore returns the semaphore for a host, creating it if necessary.
// getHostSemaphore returns the semaphore for a host, creating it if
// necessary, and counts the caller among the fetches using it. The caller
// calls putHostSemaphore once it holds no slot and waits for none.
func (f *HTTPFetcher) getHostSemaphore(host string) chan struct{} {
f.hostSemMu.Lock()
defer f.hostSemMu.Unlock()
sem, ok := f.hostSems[host]
if !ok {
sem = make(chan struct{}, f.config.MaxConnectionsPerHost)
sem = &hostSemaphore{
slots: make(chan struct{}, f.config.MaxConnectionsPerHost),
}
f.hostSems[host] = sem
}
return sem
sem.count++
return sem.slots
}
// putHostSemaphore stops counting the caller among the fetches using the
// host's semaphore, and removes the semaphore when no fetch uses it.
func (f *HTTPFetcher) putHostSemaphore(host string) {
f.hostSemMu.Lock()
defer f.hostSemMu.Unlock()
sem := f.hostSems[host]
sem.count--
if sem.count == 0 {
delete(f.hostSems, host)
}
}
// buildResult validates the upstream response and assembles a FetchResult
@@ -5,6 +5,7 @@ import (
"errors"
"net"
"strconv"
"sync"
"testing"
"time"
)
@@ -119,6 +120,98 @@ func TestFetchFreesHostSlotWhenContextEndsWaitingForConnection(t *testing.T) {
}
}
// TestFetchRemovesIdleHostSemaphores checks that a host's semaphore is
// removed once no fetch holds or waits for one of its slots: after 100
// concurrent fetches from 50 hosts have all finished, no semaphore is left.
func TestFetchRemovesIdleHostSemaphores(t *testing.T) {
t.Parallel()
srv := startUpstream(t)
f, _ := newServerFetcher(t, srv, nil)
ctx := testContext(t)
var wg sync.WaitGroup
for i := range 100 {
wg.Go(func() {
res, err := f.Fetch(ctx, imageURLOnPort(1+i%50))
if err != nil {
t.Errorf("Fetch() error = %v", err)
return
}
_ = res.Content.Close()
})
}
wg.Wait()
if n := hostSemCount(f); n != 0 {
t.Errorf("%d host semaphores left after every fetch finished, want 0", n)
}
}
// TestFetchRemovesHostSemaphoreWhenNoConnection checks that a fetch that
// ends without a connection leaves no semaphore behind: when it is refused
// after waiting for a connection shared by all hosts, when its context ends
// while it waits for its host's slot, and when its context ends while it
// waits for a connection shared by all hosts, long before the 10 second
// wait timeout.
func TestFetchRemovesHostSemaphoreWhenNoConnection(t *testing.T) {
t.Parallel()
srv := startUpstream(t)
cfg := DefaultConfig()
cfg.MaxConnections = 1
cfg.MaxConnectionsPerHost = 1
f, _ := newServerFetcher(t, srv, cfg)
f.connectionWaitTimeout = 100 * time.Millisecond
open, err := f.Fetch(testContext(t), imageURLOnPort(81))
if err != nil {
t.Fatalf("first Fetch() error = %v", err)
}
_, err = f.Fetch(testContext(t), imageURLOnPort(82))
if !errors.Is(err, ErrTooManyConnections) {
t.Fatalf("Fetch() from another host: error = %v, "+
"want ErrTooManyConnections", err)
}
ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
defer cancel()
_, err = f.Fetch(ctx, imageURLOnPort(81))
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("Fetch() from the busy host: error = %v, "+
"want context.DeadlineExceeded", err)
}
// Back to the 10 second wait, so the next fetch's context ends first.
f.connectionWaitTimeout = ConnectionWaitTimeout
ctx, cancel = context.WithTimeout(t.Context(), 100*time.Millisecond)
defer cancel()
_, err = f.Fetch(ctx, imageURLOnPort(83))
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("Fetch() from a host with nothing open: error = %v, "+
"want context.DeadlineExceeded", err)
}
err = open.Content.Close()
if err != nil {
t.Fatalf("close first body: %v", err)
}
if n := hostSemCount(f); n != 0 {
t.Errorf("%d host semaphores left after every fetch finished, want 0", n)
}
}
// TestFetchReleasesConnectionOnError checks that a fetch that fails after
// taking its connection gives it back: with MaxConnections at 1, the slot
// must be free after the failure and the next fetch must succeed.
+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 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
@@ -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)
}
}
+1 -1
View File
@@ -283,7 +283,7 @@ func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error {
c.metaCache.Remove(cacheKey)
err = c.variants.DeleteWithMeta(cacheKey)
err = c.variants.Delete(cacheKey)
if err != nil {
return err
}
+6
View File
@@ -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)
+3 -13
View File
@@ -564,7 +564,8 @@ func (s *VariantStorage) Exists(key VariantKey) bool {
return err == nil
}
// Delete removes content at the given key.
// Delete removes the content at the given key together with its .meta
// sidecar file. A missing file is not an error.
func (s *VariantStorage) Delete(key VariantKey) error {
path := s.keyToPath(key)
@@ -573,18 +574,7 @@ func (s *VariantStorage) Delete(key VariantKey) error {
return fmt.Errorf("failed to delete content: %w", err)
}
return nil
}
// DeleteWithMeta removes the content at the given key together with
// its .meta sidecar file. A missing file is not an error.
func (s *VariantStorage) DeleteWithMeta(key VariantKey) error {
err := s.Delete(key)
if err != nil {
return err
}
metaPath := s.keyToPath(key) + ".meta"
metaPath := path + ".meta"
err = os.Remove(metaPath)
if err != nil && !os.IsNotExist(err) {
@@ -438,3 +438,73 @@ func TestVariantStorage_StoreLogsFailedMetaWrite(t *testing.T) {
t.Errorf("log missing %s; got %q", want, logBuf.String())
}
}
// storeTestVariant stores one variant, with its .meta file, in a new
// VariantStorage and returns the storage and the variant's key.
func storeTestVariant(t *testing.T) (*VariantStorage, VariantKey) {
t.Helper()
storage, err := NewVariantStorage(t.TempDir(), slog.New(slog.DiscardHandler))
if err != nil {
t.Fatalf("NewVariantStorage() error = %v", err)
}
key := CacheKey(&ImageRequest{SourceHost: testHostCDN, SourcePath: testPathCat})
_, err = storage.Store(key, bytes.NewReader([]byte("variant data")), "image/webp")
if err != nil {
t.Fatalf("Store() error = %v", err)
}
return storage, key
}
// TestVariantStorage_DeleteRemovesMeta verifies that Delete removes the
// variant's .meta file along with the variant file.
func TestVariantStorage_DeleteRemovesMeta(t *testing.T) {
t.Parallel()
storage, key := storeTestVariant(t)
metaPath := storage.keyToPath(key) + ".meta"
_, err := os.Stat(metaPath)
if err != nil {
t.Fatalf("Store() wrote no .meta file: %v", err)
}
err = storage.Delete(key)
if err != nil {
t.Fatalf("Delete() error = %v", err)
}
if storage.Exists(key) {
t.Error("Exists() = true after delete, want false")
}
_, err = os.Stat(metaPath)
if !os.IsNotExist(err) {
t.Errorf(".meta file left after Delete() (stat err=%v)", err)
}
}
// TestVariantStorage_DeleteWithoutMeta verifies that Delete succeeds for a
// variant whose .meta file is missing.
func TestVariantStorage_DeleteWithoutMeta(t *testing.T) {
t.Parallel()
storage, key := storeTestVariant(t)
err := os.Remove(storage.keyToPath(key) + ".meta")
if err != nil {
t.Fatalf("removing .meta file: %v", err)
}
err = storage.Delete(key)
if err != nil {
t.Fatalf("Delete() error = %v, want nil", err)
}
if storage.Exists(key) {
t.Error("Exists() = true after delete, want false")
}
}
+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")
}
}