Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
62bce544d3 | ||
|
|
56ec847948 |
+5
-11
@@ -18,14 +18,9 @@ COPY . .
|
|||||||
# Tells script/lint it is inside a container, so it runs the linter.
|
# Tells script/lint it is inside a container, so it runs the linter.
|
||||||
ENV container=docker
|
ENV container=docker
|
||||||
|
|
||||||
# Run formatting check and linter. script/cibuild and script/docker pass
|
# Run formatting check and linter
|
||||||
# a new CHECK_EPOCH on every run, and each check step names it in its
|
RUN make fmt-check
|
||||||
# command, so a new value reruns the step instead of reusing a cached
|
RUN make lint
|
||||||
# success that checked nothing. A plain `docker build .` leaves it empty
|
|
||||||
# and reuses the check steps only for an identical build context.
|
|
||||||
ARG CHECK_EPOCH
|
|
||||||
RUN echo "check epoch: ${CHECK_EPOCH}" && make fmt-check
|
|
||||||
RUN echo "check epoch: ${CHECK_EPOCH}" && make lint
|
|
||||||
|
|
||||||
# Build stage
|
# Build stage
|
||||||
# golang:1.25.4-alpine, 2026-02-25
|
# golang:1.25.4-alpine, 2026-02-25
|
||||||
@@ -44,9 +39,8 @@ RUN script/bootstrap
|
|||||||
# Copy source code
|
# Copy source code
|
||||||
COPY . .
|
COPY . .
|
||||||
|
|
||||||
# Run tests; a new CHECK_EPOCH reruns them, as in the lint stage.
|
# Run tests
|
||||||
ARG CHECK_EPOCH
|
RUN make test
|
||||||
RUN echo "check epoch: ${CHECK_EPOCH}" && make test
|
|
||||||
|
|
||||||
# VERSION is declared here, not earlier: a new value reruns only the
|
# VERSION is declared here, not earlier: a new value reruns only the
|
||||||
# build, not script/bootstrap or the tests. Given none, the version is
|
# build, not script/bootstrap or the tests. Given none, the version is
|
||||||
|
|||||||
@@ -107,14 +107,6 @@ or proxy cache keeps the image after pixa would refuse the URL. A URL with no
|
|||||||
expiry gets one year. `immutable` only stops a client revalidating while its
|
expiry gets one year. `immutable` only stops a client revalidating while its
|
||||||
copy is fresh.
|
copy is fresh.
|
||||||
|
|
||||||
When several requests for the same image, size, format, quality and fit miss
|
|
||||||
the cache at once, they share one upstream fetch (or one read of the cached
|
|
||||||
source) and one transcode: the first request does the work, and the others wait
|
|
||||||
for its image or its error, holding no upstream connection or processing slot
|
|
||||||
of their own. A waiting request stops waiting when its own client goes away.
|
|
||||||
The work goes on for the others even if the first request's client goes away,
|
|
||||||
until that request's `downstream_timeout` ends.
|
|
||||||
|
|
||||||
The login form (`POST /`) is limited to 5 attempts per minute per client
|
The login form (`POST /`) is limited to 5 attempts per minute per client
|
||||||
address, counting an IPv6 client by its /64; an attempt over the limit is
|
address, counting an IPv6 client by its /64; an attempt over the limit is
|
||||||
refused with 429 and a `Retry-After` header. Behind a reverse proxy the client
|
refused with 429 and a `Retry-After` header. Behind a reverse proxy the client
|
||||||
@@ -353,9 +345,8 @@ them. We provide:
|
|||||||
- `script/check` — run test, lint, and fmt-check
|
- `script/check` — run test, lint, and fmt-check
|
||||||
- `script/docker` — build the Docker image tagged via `script/projectname`
|
- `script/docker` — build the Docker image tagged via `script/projectname`
|
||||||
- `script/docker-smoke` — build the image, start it, wait for it to be healthy
|
- `script/docker-smoke` — build the image, start it, wait for it to be healthy
|
||||||
- `script/cibuild` — CI entrypoint: `docker build .` with a new
|
- `script/cibuild` — CI entrypoint: `docker build .` (the Dockerfile
|
||||||
`CHECK_EPOCH` on every run, so the Dockerfile's checks run instead of
|
runs the checks, so a green build implies a green repo)
|
||||||
coming from the build cache, and a green run implies a green repo
|
|
||||||
- `script/precommit` — pre-commit checks (`go mod tidy` guard, then
|
- `script/precommit` — pre-commit checks (`go mod tidy` guard, then
|
||||||
`script/check`)
|
`script/check`)
|
||||||
- `script/install-precommit` — install the git pre-commit hook that
|
- `script/install-precommit` — install the git pre-commit hook that
|
||||||
|
|||||||
@@ -29,27 +29,17 @@ P2: security: referer blacklist
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
- 2026-10-03 every `script/cibuild` and `script/docker` run executes the checks
|
- 2026-10-03 shutdown sets the exit code and waits for image processing
|
||||||
(closes #101): the `Dockerfile` declares `CHECK_EPOCH` above `make fmt-check`
|
(closes #86): fx alone handles SIGINT and SIGTERM, and the server's own
|
||||||
and `make lint` in the lint stage and above `make test` in the build stage,
|
signal handler is gone; fx's `Run` in `cmd/pixad` exits with the shutdown's
|
||||||
and each of those steps names it in its command; both scripts pass a new value
|
code: 0 for a signal, 1 when the HTTP server cannot listen or the app fails
|
||||||
on every run, so Docker runs the checks instead of reusing cached results,
|
to start or to stop; the server's stop hook, which fx waits for, stops the
|
||||||
while the `script/bootstrap` steps stay cached; a plain `docker build .` still
|
HTTP server, waits for the images still being processed, both within 5
|
||||||
works, leaves it empty, and reuses the check steps only for an identical build
|
seconds, then flushes Sentry; images still being processed after that are
|
||||||
context; the `script/cibuild` comment and `README.md` no longer say that any
|
logged with their count and make the exit code 1; a Sentry DSN that cannot be
|
||||||
successful build implies a green repo.
|
used fails startup, so the stop hooks of what had already started run,
|
||||||
- 2026-09-29 share concurrent misses (closes #65): requests that miss the same
|
instead of exiting the process from a goroutine; the eviction loop is left to
|
||||||
variant at once (the same cache key, so quality and fit included) share one
|
#102.
|
||||||
upstream fetch or cached source read and one transcode through
|
|
||||||
`golang.org/x/sync/singleflight`; the first request's processing ignores its
|
|
||||||
cancellation but keeps its deadline, and the others wait for its image or
|
|
||||||
error holding no upstream connection or processing slot, and stop waiting
|
|
||||||
when their own context ends; the request doing the processing waits for it
|
|
||||||
even then, up to its deadline; a request whose context has already ended
|
|
||||||
starts nothing; each request counts one miss, and the processing counts its
|
|
||||||
fetch and transcode once; a panic while processing is reported to Sentry when
|
|
||||||
`sentry_dsn` is set and becomes an error for every waiting request instead of
|
|
||||||
stopping pixad; documented in `README.md`.
|
|
||||||
- 2026-09-29 only the image routes send CORS headers (closes #98): the CORS
|
- 2026-09-29 only the image routes send CORS headers (closes #98): the CORS
|
||||||
middleware, with the `access_control_allow_origin` origin, moved from the
|
middleware, with the `access_control_allow_origin` origin, moved from the
|
||||||
router root onto a `/v1` subrouter holding `/v1/image/` and `/v1/e/`, where it
|
router root onto a `/v1` subrouter holding `/v1/image/` and `/v1/e/`, where it
|
||||||
|
|||||||
@@ -4,6 +4,8 @@ package main
|
|||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"os"
|
"os"
|
||||||
|
"os/signal"
|
||||||
|
"syscall"
|
||||||
|
|
||||||
"github.com/spf13/cobra"
|
"github.com/spf13/cobra"
|
||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
@@ -45,6 +47,9 @@ func run(_ *cobra.Command, _ []string) {
|
|||||||
_ = os.Setenv("PIXA_CONFIG_PATH", configPath)
|
_ = 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.New(
|
||||||
fx.Provide(
|
fx.Provide(
|
||||||
config.New,
|
config.New,
|
||||||
|
|||||||
@@ -20,7 +20,6 @@ require (
|
|||||||
github.com/spf13/cobra v1.10.2
|
github.com/spf13/cobra v1.10.2
|
||||||
go.uber.org/fx v1.24.0
|
go.uber.org/fx v1.24.0
|
||||||
golang.org/x/crypto v0.41.0
|
golang.org/x/crypto v0.41.0
|
||||||
golang.org/x/sync v0.19.0
|
|
||||||
modernc.org/sqlite v1.42.2
|
modernc.org/sqlite v1.42.2
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -135,6 +134,7 @@ require (
|
|||||||
golang.org/x/image v0.34.0 // indirect
|
golang.org/x/image v0.34.0 // indirect
|
||||||
golang.org/x/net v0.43.0 // indirect
|
golang.org/x/net v0.43.0 // indirect
|
||||||
golang.org/x/oauth2 v0.30.0 // indirect
|
golang.org/x/oauth2 v0.30.0 // indirect
|
||||||
|
golang.org/x/sync v0.19.0 // indirect
|
||||||
golang.org/x/sys v0.36.0 // indirect
|
golang.org/x/sys v0.36.0 // indirect
|
||||||
golang.org/x/term v0.34.0 // indirect
|
golang.org/x/term v0.34.0 // indirect
|
||||||
golang.org/x/text v0.32.0 // indirect
|
golang.org/x/text v0.32.0 // indirect
|
||||||
|
|||||||
@@ -81,6 +81,12 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) {
|
|||||||
return s, nil
|
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.
|
// initImageService initializes the image cache and service.
|
||||||
func (s *Handlers) initImageService() error {
|
func (s *Handlers) initImageService() error {
|
||||||
// Create the cache. cache_max_bytes: 0 disables the disk cache
|
// 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
|
// acquireSlot takes a slot in processingSemaphore, waiting at most
|
||||||
// processingWaitTimeout for one to free up, and returns the func that gives
|
// 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
|
// 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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -471,8 +471,7 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
|
|||||||
return &stats, nil
|
return &stats, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// IncrementStats counts a cache hit or miss, and an upstream fetch that read
|
// IncrementStats increments cache statistics.
|
||||||
// fetchBytes bytes, as IncrementUpstreamFetch does.
|
|
||||||
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
|
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
@@ -496,17 +495,8 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
|
|||||||
c.log.Warn("failed to count cache hit or miss", "hit", hit, "error", err)
|
c.log.Warn("failed to count cache hit or miss", "hit", hit, "error", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
c.IncrementUpstreamFetch(ctx, fetchBytes)
|
if fetchBytes > 0 {
|
||||||
}
|
_, err = c.db.ExecContext(ctx, `
|
||||||
|
|
||||||
// IncrementUpstreamFetch counts one upstream fetch that read fetchBytes bytes.
|
|
||||||
// A fetch that read no bytes is not counted.
|
|
||||||
func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
|
|
||||||
if fetchBytes <= 0 {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err := c.db.ExecContext(ctx, `
|
|
||||||
UPDATE cache_stats
|
UPDATE cache_stats
|
||||||
SET upstream_fetch_count = upstream_fetch_count + 1,
|
SET upstream_fetch_count = upstream_fetch_count + 1,
|
||||||
upstream_fetch_bytes = upstream_fetch_bytes + ?,
|
upstream_fetch_bytes = upstream_fetch_bytes + ?,
|
||||||
@@ -517,6 +507,7 @@ func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
|
|||||||
c.log.Warn("failed to count upstream fetch",
|
c.log.Warn("failed to count upstream fetch",
|
||||||
"fetch_bytes", fetchBytes, "error", err)
|
"fetch_bytes", fetchBytes, "error", err)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// IncrementTransformCount counts one image transcoded by the image processor.
|
// IncrementTransformCount counts one image transcoded by the image processor.
|
||||||
|
|||||||
@@ -1,509 +0,0 @@
|
|||||||
package imgcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"bytes"
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"image/jpeg"
|
|
||||||
"io"
|
|
||||||
"io/fs"
|
|
||||||
"strings"
|
|
||||||
"sync"
|
|
||||||
"sync/atomic"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/getsentry/sentry-go"
|
|
||||||
"sneak.berlin/go/pixa/internal/httpfetcher"
|
|
||||||
"sneak.berlin/go/pixa/internal/magic"
|
|
||||||
)
|
|
||||||
|
|
||||||
// arrivalWait is how long a test gives requests it has started to reach the
|
|
||||||
// point where they wait for a held fetch.
|
|
||||||
const arrivalWait = 100 * time.Millisecond
|
|
||||||
|
|
||||||
// heldFetcher counts the fetches it is asked for and holds each one until
|
|
||||||
// releaseFetches is called, so that a test can have requests arrive while a
|
|
||||||
// fetch is in progress. started receives once for every fetch.
|
|
||||||
type heldFetcher struct {
|
|
||||||
upstream httpfetcher.Fetcher
|
|
||||||
fetches atomic.Int32
|
|
||||||
started chan struct{}
|
|
||||||
release chan struct{}
|
|
||||||
releaseOnce sync.Once
|
|
||||||
}
|
|
||||||
|
|
||||||
func (f *heldFetcher) Fetch(
|
|
||||||
ctx context.Context, url string,
|
|
||||||
) (*httpfetcher.FetchResult, error) {
|
|
||||||
f.fetches.Add(1)
|
|
||||||
|
|
||||||
f.started <- struct{}{}
|
|
||||||
|
|
||||||
select {
|
|
||||||
case <-f.release:
|
|
||||||
case <-ctx.Done():
|
|
||||||
return nil, ctx.Err()
|
|
||||||
}
|
|
||||||
|
|
||||||
return f.upstream.Fetch(ctx, url)
|
|
||||||
}
|
|
||||||
|
|
||||||
// releaseFetches lets every held fetch, and every later one, go on.
|
|
||||||
func (f *heldFetcher) releaseFetches() {
|
|
||||||
f.releaseOnce.Do(func() { close(f.release) })
|
|
||||||
}
|
|
||||||
|
|
||||||
// setupHeldFetchService returns a test service whose fetches go through a
|
|
||||||
// heldFetcher. Its database is limited to one connection: each connection to
|
|
||||||
// an in-memory SQLite database opens a new, empty one, so requests running at
|
|
||||||
// once must share the connection that holds the schema.
|
|
||||||
func setupHeldFetchService(t *testing.T) (*Service, *TestFixtures, *heldFetcher) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
svc, fixtures := SetupTestService(t)
|
|
||||||
svc.cache.db.SetMaxOpenConns(1)
|
|
||||||
|
|
||||||
fetcher := &heldFetcher{
|
|
||||||
upstream: svc.fetcher,
|
|
||||||
started: make(chan struct{}, 100),
|
|
||||||
release: make(chan struct{}),
|
|
||||||
}
|
|
||||||
svc.fetcher = fetcher
|
|
||||||
|
|
||||||
t.Cleanup(fetcher.releaseFetches)
|
|
||||||
|
|
||||||
return svc, fixtures, fetcher
|
|
||||||
}
|
|
||||||
|
|
||||||
// photoVariant asks for the test photo, 100x100, at 50x25 as a JPEG of the
|
|
||||||
// given quality and fit mode. Each call returns a new request, as Get writes
|
|
||||||
// to the request it is given.
|
|
||||||
func photoVariant(fixtures *TestFixtures, quality int, fit FitMode) *ImageRequest {
|
|
||||||
return &ImageRequest{
|
|
||||||
SourceHost: fixtures.GoodHost,
|
|
||||||
SourcePath: testPathPhoto,
|
|
||||||
Size: Size{Width: 50, Height: 25},
|
|
||||||
Format: FormatJPEG,
|
|
||||||
Quality: quality,
|
|
||||||
FitMode: fit,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// getResult is what one Get call returned, with the image read out.
|
|
||||||
type getResult struct {
|
|
||||||
image []byte
|
|
||||||
err error
|
|
||||||
}
|
|
||||||
|
|
||||||
// startGet calls Get in a goroutine of its own and delivers what it returned
|
|
||||||
// on the channel.
|
|
||||||
func startGet(
|
|
||||||
ctx context.Context, svc *Service, req *ImageRequest,
|
|
||||||
) <-chan getResult {
|
|
||||||
results := make(chan getResult, 1)
|
|
||||||
|
|
||||||
go func() {
|
|
||||||
resp, err := svc.Get(ctx, req)
|
|
||||||
if err != nil {
|
|
||||||
results <- getResult{err: err}
|
|
||||||
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
defer func() { _ = resp.Content.Close() }()
|
|
||||||
|
|
||||||
image, err := io.ReadAll(resp.Content)
|
|
||||||
results <- getResult{image: image, err: err}
|
|
||||||
}()
|
|
||||||
|
|
||||||
return results
|
|
||||||
}
|
|
||||||
|
|
||||||
// jpegSize returns the width and height of the JPEG image in data, or 0 and 0
|
|
||||||
// if data is not one.
|
|
||||||
func jpegSize(data []byte) (int, int) {
|
|
||||||
config, err := jpeg.DecodeConfig(bytes.NewReader(data))
|
|
||||||
if err != nil {
|
|
||||||
return 0, 0
|
|
||||||
}
|
|
||||||
|
|
||||||
return config.Width, config.Height
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_ConcurrentMissesShareOneFetch starts several requests for
|
|
||||||
// one uncached variant while the first one's fetch is held. Between them they
|
|
||||||
// must fetch the source once and transcode it once, every one must be answered
|
|
||||||
// with the same 50x25 JPEG, and each must count one miss.
|
|
||||||
func TestService_Get_ConcurrentMissesShareOneFetch(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures, fetcher := setupHeldFetchService(t)
|
|
||||||
|
|
||||||
const requests = 8
|
|
||||||
|
|
||||||
pending := make([]<-chan getResult, 0, requests)
|
|
||||||
for range requests {
|
|
||||||
pending = append(pending,
|
|
||||||
startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover)))
|
|
||||||
}
|
|
||||||
|
|
||||||
<-fetcher.started
|
|
||||||
time.Sleep(arrivalWait)
|
|
||||||
fetcher.releaseFetches()
|
|
||||||
|
|
||||||
var first []byte
|
|
||||||
|
|
||||||
for i, results := range pending {
|
|
||||||
got := <-results
|
|
||||||
if got.err != nil {
|
|
||||||
t.Fatalf("request %d: Get() error = %v", i, got.err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if width, height := jpegSize(got.image); width != 50 || height != 25 {
|
|
||||||
t.Errorf("request %d: image is %dx%d, want a 50x25 JPEG", i, width, height)
|
|
||||||
}
|
|
||||||
|
|
||||||
if first == nil {
|
|
||||||
first = got.image
|
|
||||||
} else if !bytes.Equal(got.image, first) {
|
|
||||||
t.Errorf("request %d: image differs from request 0's", i)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if fetches := fetcher.fetches.Load(); fetches != 1 {
|
|
||||||
t.Errorf("%d requests made %d upstream fetches, want 1", requests, fetches)
|
|
||||||
}
|
|
||||||
|
|
||||||
// NewTestFS builds the same files the test service's fetcher serves.
|
|
||||||
testFS, _ := NewTestFS(t)
|
|
||||||
|
|
||||||
photo, err := fs.ReadFile(testFS, fixtures.GoodHostJPEG)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
want := cacheStatsCounters{0, requests, 1, int64(len(photo)), 1}
|
|
||||||
|
|
||||||
if got := readCacheStatsCounters(t, svc.cache); got != want {
|
|
||||||
t.Errorf("counters = %+v, want %+v", got, want)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_ConcurrentVariantsStayApart requests three variants of the
|
|
||||||
// test photo at once that differ only in quality or fit. Each must be made by
|
|
||||||
// a fetch and a transcode of its own, and each answer must be its own variant.
|
|
||||||
func TestService_Get_ConcurrentVariantsStayApart(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures, fetcher := setupHeldFetchService(t)
|
|
||||||
|
|
||||||
cover := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover))
|
|
||||||
lowQuality := startGet(t.Context(), svc, photoVariant(fixtures, 40, FitCover))
|
|
||||||
contain := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitContain))
|
|
||||||
|
|
||||||
for range 3 {
|
|
||||||
select {
|
|
||||||
case <-fetcher.started:
|
|
||||||
case <-time.After(5 * time.Second):
|
|
||||||
t.Fatal("fewer fetches started than variants requested: " +
|
|
||||||
"variants differing in quality or fit were merged")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fetcher.releaseFetches()
|
|
||||||
|
|
||||||
images := make(map[string][]byte)
|
|
||||||
|
|
||||||
for _, variant := range []struct {
|
|
||||||
name string
|
|
||||||
results <-chan getResult
|
|
||||||
width, height int
|
|
||||||
}{
|
|
||||||
{"q=85 fit=cover", cover, 50, 25},
|
|
||||||
{"q=40 fit=cover", lowQuality, 50, 25},
|
|
||||||
{"q=85 fit=contain", contain, 25, 25},
|
|
||||||
} {
|
|
||||||
got := <-variant.results
|
|
||||||
if got.err != nil {
|
|
||||||
t.Fatalf("%s: Get() error = %v", variant.name, got.err)
|
|
||||||
}
|
|
||||||
|
|
||||||
width, height := jpegSize(got.image)
|
|
||||||
if width != variant.width || height != variant.height {
|
|
||||||
t.Errorf("%s: image is %dx%d, want a %dx%d JPEG", variant.name,
|
|
||||||
width, height, variant.width, variant.height)
|
|
||||||
}
|
|
||||||
|
|
||||||
images[variant.name] = got.image
|
|
||||||
}
|
|
||||||
|
|
||||||
if bytes.Equal(images["q=85 fit=cover"], images["q=40 fit=cover"]) {
|
|
||||||
t.Error("q=40 was answered with the q=85 image")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_WaiterStopsWhenItsContextEnds has a second request for a
|
|
||||||
// variant join the first one's held fetch, then ends the second request's
|
|
||||||
// context. The second request must return at once with the context's error,
|
|
||||||
// while the fetch is still held, and the first must still be answered.
|
|
||||||
func TestService_Get_WaiterStopsWhenItsContextEnds(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures, fetcher := setupHeldFetchService(t)
|
|
||||||
|
|
||||||
first := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover))
|
|
||||||
|
|
||||||
<-fetcher.started
|
|
||||||
|
|
||||||
waiterCtx, cancelWaiter := context.WithCancel(t.Context())
|
|
||||||
waiter := startGet(waiterCtx, svc, photoVariant(fixtures, 85, FitCover))
|
|
||||||
|
|
||||||
time.Sleep(arrivalWait)
|
|
||||||
cancelWaiter()
|
|
||||||
|
|
||||||
select {
|
|
||||||
case got := <-waiter:
|
|
||||||
if !errors.Is(got.err, context.Canceled) {
|
|
||||||
t.Errorf("waiting request: Get() error = %v, want %v",
|
|
||||||
got.err, context.Canceled)
|
|
||||||
}
|
|
||||||
case <-time.After(time.Second):
|
|
||||||
t.Fatal("waiting request did not return when its context ended")
|
|
||||||
}
|
|
||||||
|
|
||||||
fetcher.releaseFetches()
|
|
||||||
|
|
||||||
got := <-first
|
|
||||||
if got.err != nil {
|
|
||||||
t.Fatalf("first request: Get() error = %v", got.err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if width, height := jpegSize(got.image); width != 50 || height != 25 {
|
|
||||||
t.Errorf("first request: image is %dx%d, want a 50x25 JPEG", width, height)
|
|
||||||
}
|
|
||||||
|
|
||||||
if fetches := fetcher.fetches.Load(); fetches != 1 {
|
|
||||||
t.Errorf("upstream fetches = %d, want 1", fetches)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_FirstRequestLeavingKeepsTheWork ends the context of the
|
|
||||||
// request whose fetch is held, after a second request has joined it. The fetch
|
|
||||||
// and transcode must go on and answer the second request. The first request
|
|
||||||
// waits for its own work, as every request did before misses were shared, and
|
|
||||||
// is answered too.
|
|
||||||
func TestService_Get_FirstRequestLeavingKeepsTheWork(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures, fetcher := setupHeldFetchService(t)
|
|
||||||
|
|
||||||
firstCtx, cancelFirst := context.WithCancel(t.Context())
|
|
||||||
first := startGet(firstCtx, svc, photoVariant(fixtures, 85, FitCover))
|
|
||||||
|
|
||||||
<-fetcher.started
|
|
||||||
|
|
||||||
second := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover))
|
|
||||||
|
|
||||||
time.Sleep(arrivalWait)
|
|
||||||
cancelFirst()
|
|
||||||
fetcher.releaseFetches()
|
|
||||||
|
|
||||||
for name, results := range map[string]<-chan getResult{
|
|
||||||
"first request": first, "second request": second,
|
|
||||||
} {
|
|
||||||
got := <-results
|
|
||||||
if got.err != nil {
|
|
||||||
t.Fatalf("%s: Get() error = %v", name, got.err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if width, height := jpegSize(got.image); width != 50 || height != 25 {
|
|
||||||
t.Errorf("%s: image is %dx%d, want a 50x25 JPEG", name, width, height)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if fetches := fetcher.fetches.Load(); fetches != 1 {
|
|
||||||
t.Errorf("upstream fetches = %d, want 1", fetches)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_ConcurrentMissesShareAFailure has several requests for an
|
|
||||||
// image that cannot be served arrive while its fetch is held: the one fetch
|
|
||||||
// answers all of them with its error. The request after them is answered from
|
|
||||||
// the negative cache when the failure is kept there, and fetches again when it
|
|
||||||
// is not.
|
|
||||||
func TestService_Get_ConcurrentMissesShareAFailure(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
tests := []struct {
|
|
||||||
name string
|
|
||||||
path string
|
|
||||||
wantErr error // what every request at once gets
|
|
||||||
wantNextErr error // what the request after them gets
|
|
||||||
wantFetches int32 // fetches once the request after them is answered
|
|
||||||
}{
|
|
||||||
{"upstream answers 404, kept in the negative cache",
|
|
||||||
"/images/missing.jpg", httpfetcher.ErrUpstreamError,
|
|
||||||
ErrNegativeCached, 1},
|
|
||||||
{"source fails the magic byte check, not kept",
|
|
||||||
"/images/text.png", magic.ErrUnknownFormat, magic.ErrUnknownFormat, 2},
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, tc := range tests {
|
|
||||||
t.Run(tc.name, func(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures, fetcher := setupHeldFetchService(t)
|
|
||||||
|
|
||||||
request := func() *ImageRequest {
|
|
||||||
req := photoVariant(fixtures, 85, FitCover)
|
|
||||||
req.SourcePath = tc.path
|
|
||||||
|
|
||||||
return req
|
|
||||||
}
|
|
||||||
|
|
||||||
pending := make([]<-chan getResult, 0, 4)
|
|
||||||
for range 4 {
|
|
||||||
pending = append(pending, startGet(t.Context(), svc, request()))
|
|
||||||
}
|
|
||||||
|
|
||||||
<-fetcher.started
|
|
||||||
time.Sleep(arrivalWait)
|
|
||||||
fetcher.releaseFetches()
|
|
||||||
|
|
||||||
for i, results := range pending {
|
|
||||||
if got := <-results; !errors.Is(got.err, tc.wantErr) {
|
|
||||||
t.Errorf("request %d: Get() error = %v, want %v", i, got.err, tc.wantErr)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err := svc.Get(t.Context(), request())
|
|
||||||
if !errors.Is(err, tc.wantNextErr) {
|
|
||||||
t.Errorf("next request: Get() error = %v, want %v", err, tc.wantNextErr)
|
|
||||||
}
|
|
||||||
|
|
||||||
if fetches := fetcher.fetches.Load(); fetches != tc.wantFetches {
|
|
||||||
t.Errorf("upstream fetches = %d, want %d", fetches, tc.wantFetches)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_EndedRequestFetchesNothing checks that a request whose
|
|
||||||
// context has already ended when it misses the cache starts no fetch.
|
|
||||||
func TestService_Get_EndedRequestFetchesNothing(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures, fetcher := setupHeldFetchService(t)
|
|
||||||
|
|
||||||
ctx, cancel := context.WithCancel(t.Context())
|
|
||||||
cancel()
|
|
||||||
|
|
||||||
_, err := svc.Get(ctx, photoVariant(fixtures, 85, FitCover))
|
|
||||||
if !errors.Is(err, context.Canceled) {
|
|
||||||
t.Errorf("Get() error = %v, want %v", err, context.Canceled)
|
|
||||||
}
|
|
||||||
|
|
||||||
if fetches := fetcher.fetches.Load(); fetches != 0 {
|
|
||||||
t.Errorf("upstream fetches = %d, want 0", fetches)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_ReturnsByItsDeadline gives a request a deadline and holds
|
|
||||||
// its fetch until the fetch's context ends, as the fetcher does while it waits
|
|
||||||
// for a free connection to the host. The request must return by its deadline
|
|
||||||
// with the deadline's error.
|
|
||||||
func TestService_Get_ReturnsByItsDeadline(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures, _ := setupHeldFetchService(t)
|
|
||||||
|
|
||||||
const timeout = 200 * time.Millisecond
|
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(t.Context(), timeout)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
results := startGet(ctx, svc, photoVariant(fixtures, 85, FitCover))
|
|
||||||
|
|
||||||
select {
|
|
||||||
case got := <-results:
|
|
||||||
if !errors.Is(got.err, context.DeadlineExceeded) {
|
|
||||||
t.Errorf("Get() error = %v, want %v", got.err, context.DeadlineExceeded)
|
|
||||||
}
|
|
||||||
case <-time.After(timeout + time.Second):
|
|
||||||
t.Fatal("request did not return by its deadline")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// panickingFetcher panics on every fetch.
|
|
||||||
type panickingFetcher struct{}
|
|
||||||
|
|
||||||
func (panickingFetcher) Fetch(
|
|
||||||
context.Context, string,
|
|
||||||
) (*httpfetcher.FetchResult, error) {
|
|
||||||
panic("upstream fetcher panicked")
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_PanicBecomesAnError checks that a panic while a variant is
|
|
||||||
// being made reaches its request as an error naming the panic, instead of
|
|
||||||
// being raised again.
|
|
||||||
func TestService_Get_PanicBecomesAnError(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures := SetupTestService(t)
|
|
||||||
svc.fetcher = panickingFetcher{}
|
|
||||||
|
|
||||||
var err error
|
|
||||||
|
|
||||||
func() {
|
|
||||||
defer func() {
|
|
||||||
if recovered := recover(); recovered != nil {
|
|
||||||
t.Fatalf("Get() panicked: %v", recovered)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
_, err = svc.Get(t.Context(), photoVariant(fixtures, 85, FitCover))
|
|
||||||
}()
|
|
||||||
|
|
||||||
if err == nil || !strings.Contains(err.Error(), "upstream fetcher panicked") {
|
|
||||||
t.Errorf("Get() error = %v, want one naming the panic", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestService_Get_PanicIsReportedToSentry checks that a panic while a variant
|
|
||||||
// is being made is reported through the Sentry hub on the request's context,
|
|
||||||
// where the Sentry middleware puts one when sentry_dsn is set.
|
|
||||||
func TestService_Get_PanicIsReportedToSentry(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
svc, fixtures := SetupTestService(t)
|
|
||||||
svc.fetcher = panickingFetcher{}
|
|
||||||
|
|
||||||
transport := &sentry.MockTransport{}
|
|
||||||
|
|
||||||
client, err := sentry.NewClient(sentry.ClientOptions{
|
|
||||||
Dsn: "https://abc123@sentry.example.com/42",
|
|
||||||
Transport: transport,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
ctx := sentry.SetHubOnContext(t.Context(),
|
|
||||||
sentry.NewHub(client, sentry.NewScope()))
|
|
||||||
|
|
||||||
func() {
|
|
||||||
defer func() {
|
|
||||||
if recovered := recover(); recovered != nil {
|
|
||||||
t.Fatalf("Get() panicked: %v", recovered)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
_, _ = svc.Get(ctx, photoVariant(fixtures, 85, FitCover))
|
|
||||||
}()
|
|
||||||
|
|
||||||
events := transport.Events()
|
|
||||||
if len(events) != 1 || events[0].Message != "upstream fetcher panicked" {
|
|
||||||
t.Errorf("Sentry events = %+v, want one for the panic", events)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+33
-125
@@ -8,12 +8,9 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/url"
|
"net/url"
|
||||||
"runtime/debug"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/dustin/go-humanize"
|
"github.com/dustin/go-humanize"
|
||||||
"github.com/getsentry/sentry-go"
|
|
||||||
"golang.org/x/sync/singleflight"
|
|
||||||
"sneak.berlin/go/pixa/internal/allowlist"
|
"sneak.berlin/go/pixa/internal/allowlist"
|
||||||
"sneak.berlin/go/pixa/internal/httpfetcher"
|
"sneak.berlin/go/pixa/internal/httpfetcher"
|
||||||
"sneak.berlin/go/pixa/internal/imageprocessor"
|
"sneak.berlin/go/pixa/internal/imageprocessor"
|
||||||
@@ -32,9 +29,6 @@ type Service struct {
|
|||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
allowHTTP bool
|
allowHTTP bool
|
||||||
maxResponseSize int64
|
maxResponseSize int64
|
||||||
// variantsInProgress lets the requests that miss the same variant at the
|
|
||||||
// same time share one fetch and one transcode.
|
|
||||||
variantsInProgress singleflight.Group
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// ServiceConfig holds configuration for the image service.
|
// ServiceConfig holds configuration for the image service.
|
||||||
@@ -166,12 +160,14 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Cache miss - get the variant, processed once for all the requests that
|
// Cache miss - process the cached source or fetch it, then count the
|
||||||
// miss it at the same time, then count this request's miss, also when it
|
// miss with the bytes it fetched from upstream, also when it failed or
|
||||||
// failed or the request context has ended meanwhile
|
// the request context has ended meanwhile
|
||||||
response, err := s.processOrWait(ctx, req)
|
cacheKey := CacheKey(req)
|
||||||
|
|
||||||
s.cache.IncrementStats(context.WithoutCancel(ctx), false, 0)
|
response, fetchedBytes, err := s.processFromSourceOrFetch(ctx, req, cacheKey)
|
||||||
|
|
||||||
|
s.cache.IncrementStats(context.WithoutCancel(ctx), false, fetchedBytes)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -199,6 +195,12 @@ func (s *Service) Stats(ctx context.Context) (*CacheStats, error) {
|
|||||||
return s.cache.Stats(ctx)
|
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.
|
// ValidateRequest validates the request signature if required.
|
||||||
func (s *Service) ValidateRequest(req *ImageRequest) error {
|
func (s *Service) ValidateRequest(req *ImageRequest) error {
|
||||||
// Check if host is allowed (no signature required)
|
// Check if host is allowed (no signature required)
|
||||||
@@ -245,96 +247,6 @@ func (s *Service) GenerateSignedURL(
|
|||||||
baseURL, path, sig, exp, req.Quality, req.FitMode), nil
|
baseURL, path, sig, exp, req.Quality, req.FitMode), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// errPanicked is returned when processing a variant panicked.
|
|
||||||
var errPanicked = errors.New("panic while processing image")
|
|
||||||
|
|
||||||
// processOrWait returns the variant req asks for. The first of the requests
|
|
||||||
// that miss a variant at the same time processes it, and singleflight hands
|
|
||||||
// its result, or its error, to the others: they fetch nothing, read no source
|
|
||||||
// and take no upstream connection or processing slot. The processing ignores
|
|
||||||
// the first request's cancellation, so the others are still served if that
|
|
||||||
// client goes away, but keeps its deadline.
|
|
||||||
func (s *Service) processOrWait(
|
|
||||||
ctx context.Context, req *ImageRequest,
|
|
||||||
) (*ImageResponse, error) {
|
|
||||||
// A request that has already ended starts no processing
|
|
||||||
if ctx.Err() != nil {
|
|
||||||
return nil, ctx.Err()
|
|
||||||
}
|
|
||||||
|
|
||||||
cacheKey := CacheKey(req)
|
|
||||||
|
|
||||||
// Closed when this request's own function runs, which singleflight does
|
|
||||||
// only when no other request is processing the variant
|
|
||||||
processing := make(chan struct{})
|
|
||||||
|
|
||||||
results := s.variantsInProgress.DoChan(string(cacheKey),
|
|
||||||
func() (_ any, err error) {
|
|
||||||
close(processing)
|
|
||||||
|
|
||||||
// singleflight would raise a panic again in a goroutine of its
|
|
||||||
// own, where no handler recovers it, and stop pixad. It is
|
|
||||||
// reported through the Sentry hub that the Sentry middleware puts
|
|
||||||
// on the request's context when sentry_dsn is set.
|
|
||||||
defer func() {
|
|
||||||
recovered := recover()
|
|
||||||
if recovered != nil {
|
|
||||||
s.log.Error("panic while processing image",
|
|
||||||
"host", req.SourceHost, "path", req.SourcePath,
|
|
||||||
"panic", recovered, "stack", string(debug.Stack()))
|
|
||||||
|
|
||||||
if hub := sentry.GetHubFromContext(ctx); hub != nil {
|
|
||||||
hub.RecoverWithContext(ctx, recovered)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = fmt.Errorf("%w: %v", errPanicked, recovered)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
processingCtx := context.WithoutCancel(ctx)
|
|
||||||
|
|
||||||
if deadline, ok := ctx.Deadline(); ok {
|
|
||||||
var cancel context.CancelFunc
|
|
||||||
|
|
||||||
processingCtx, cancel = context.WithDeadline(processingCtx, deadline)
|
|
||||||
defer cancel()
|
|
||||||
}
|
|
||||||
|
|
||||||
return s.processFromSourceOrFetch(processingCtx, req, cacheKey)
|
|
||||||
})
|
|
||||||
|
|
||||||
var result singleflight.Result
|
|
||||||
|
|
||||||
select {
|
|
||||||
case result = <-results:
|
|
||||||
case <-ctx.Done():
|
|
||||||
select {
|
|
||||||
case <-processing:
|
|
||||||
// This request is processing the variant: it waits for the
|
|
||||||
// result, as every request did before misses were shared
|
|
||||||
result = <-results
|
|
||||||
default:
|
|
||||||
// Another request is processing the variant, or this request's
|
|
||||||
// function has not started yet; the processing goes on without it
|
|
||||||
return nil, ctx.Err()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if result.Err != nil {
|
|
||||||
return nil, result.Err
|
|
||||||
}
|
|
||||||
|
|
||||||
variant, _ := result.Val.(*processedVariant)
|
|
||||||
|
|
||||||
return &ImageResponse{
|
|
||||||
Content: io.NopCloser(bytes.NewReader(variant.data)),
|
|
||||||
ContentLength: int64(len(variant.data)),
|
|
||||||
ContentType: variant.contentType,
|
|
||||||
FetchedBytes: variant.fetchedBytes,
|
|
||||||
ETag: formatETag(cacheKey),
|
|
||||||
}, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// loadCachedSource opens source content from cache, without reading it, and
|
// loadCachedSource opens source content from cache, without reading it, and
|
||||||
// returns it with its size; nil if the cached data is unavailable, empty or
|
// returns it with its size; nil if the cached data is unavailable, empty or
|
||||||
// exceeds maxResponseSize.
|
// exceeds maxResponseSize.
|
||||||
@@ -369,12 +281,13 @@ func (s *Service) loadCachedSource(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// processFromSourceOrFetch processes an image, using cached source content
|
// processFromSourceOrFetch processes an image, using cached source content
|
||||||
// if available.
|
// if available. It also returns the number of bytes fetched from upstream,
|
||||||
|
// as fetchAndProcess does, or 0 when the cached source was used.
|
||||||
func (s *Service) processFromSourceOrFetch(
|
func (s *Service) processFromSourceOrFetch(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
req *ImageRequest,
|
req *ImageRequest,
|
||||||
cacheKey VariantKey,
|
cacheKey VariantKey,
|
||||||
) (*processedVariant, error) {
|
) (*ImageResponse, int64, error) {
|
||||||
// Check if we have cached source content
|
// Check if we have cached source content
|
||||||
contentHash, _, err := s.cache.LookupSource(ctx, req)
|
contentHash, _, err := s.cache.LookupSource(ctx, req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -401,17 +314,19 @@ func (s *Service) processFromSourceOrFetch(
|
|||||||
// Process using cached source; nothing was fetched from upstream. The
|
// Process using cached source; nothing was fetched from upstream. The
|
||||||
// image processor reads the source only once it has a processing slot,
|
// image processor reads the source only once it has a processing slot,
|
||||||
// so a request waiting for one holds none of it in memory.
|
// so a request waiting for one holds none of it in memory.
|
||||||
return s.processAndStore(ctx, req, cacheKey, source, sourceSize)
|
resp, err := s.processAndStore(ctx, req, cacheKey, source, sourceSize)
|
||||||
|
|
||||||
|
return resp, 0, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// fetchAndProcess fetches from upstream, processes, and caches the result.
|
// fetchAndProcess fetches from upstream, processes, and caches the result.
|
||||||
// It counts the fetch with the bytes read from upstream, including when
|
// It also returns the number of bytes read from upstream, including when
|
||||||
// reading the response or a later step fails.
|
// reading the response or a later step fails.
|
||||||
func (s *Service) fetchAndProcess(
|
func (s *Service) fetchAndProcess(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
req *ImageRequest,
|
req *ImageRequest,
|
||||||
cacheKey VariantKey,
|
cacheKey VariantKey,
|
||||||
) (*processedVariant, error) {
|
) (*ImageResponse, int64, error) {
|
||||||
// Fetch from upstream
|
// Fetch from upstream
|
||||||
sourceURL := req.SourceURL()
|
sourceURL := req.SourceURL()
|
||||||
|
|
||||||
@@ -430,7 +345,7 @@ func (s *Service) fetchAndProcess(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil, fmt.Errorf("upstream fetch failed: %w", err)
|
return nil, 0, fmt.Errorf("upstream fetch failed: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Closing the body frees the upstream connection. It is closed only
|
// Closing the body frees the upstream connection. It is closed only
|
||||||
@@ -443,11 +358,8 @@ func (s *Service) fetchAndProcess(
|
|||||||
sourceData, err := io.ReadAll(fetchResult.Content)
|
sourceData, err := io.ReadAll(fetchResult.Content)
|
||||||
fetchBytes := int64(len(sourceData))
|
fetchBytes := int64(len(sourceData))
|
||||||
|
|
||||||
// Counted also when the request context has ended meanwhile
|
|
||||||
s.cache.IncrementUpstreamFetch(context.WithoutCancel(ctx), fetchBytes)
|
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("failed to read upstream response: %w", err)
|
return nil, fetchBytes, fmt.Errorf("failed to read upstream response: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Calculate download bitrate
|
// Calculate download bitrate
|
||||||
@@ -475,7 +387,7 @@ func (s *Service) fetchAndProcess(
|
|||||||
// Validate magic bytes match content type
|
// Validate magic bytes match content type
|
||||||
err = magic.ValidateMagicBytes(sourceData, fetchResult.ContentType)
|
err = magic.ValidateMagicBytes(sourceData, fetchResult.ContentType)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("content validation failed: %w", err)
|
return nil, fetchBytes, fmt.Errorf("content validation failed: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Store source content
|
// Store source content
|
||||||
@@ -485,17 +397,11 @@ func (s *Service) fetchAndProcess(
|
|||||||
// Continue even if caching fails
|
// Continue even if caching fails
|
||||||
}
|
}
|
||||||
|
|
||||||
return s.processAndStore(
|
resp, err := s.processAndStore(
|
||||||
ctx, req, cacheKey, bytes.NewReader(sourceData), fetchBytes,
|
ctx, req, cacheKey, bytes.NewReader(sourceData), fetchBytes,
|
||||||
)
|
)
|
||||||
}
|
|
||||||
|
|
||||||
// processedVariant is a variant as processAndStore made it. Each request that
|
return resp, fetchBytes, err
|
||||||
// shared its processing serves it through a reader of its own.
|
|
||||||
type processedVariant struct {
|
|
||||||
data []byte
|
|
||||||
contentType string
|
|
||||||
fetchedBytes int64
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// processAndStore processes the image read from source and stores the
|
// processAndStore processes the image read from source and stores the
|
||||||
@@ -506,7 +412,7 @@ func (s *Service) processAndStore(
|
|||||||
cacheKey VariantKey,
|
cacheKey VariantKey,
|
||||||
source io.Reader,
|
source io.Reader,
|
||||||
fetchBytes int64,
|
fetchBytes int64,
|
||||||
) (*processedVariant, error) {
|
) (*ImageResponse, error) {
|
||||||
// Process the image
|
// Process the image
|
||||||
processStart := time.Now()
|
processStart := time.Now()
|
||||||
|
|
||||||
@@ -570,10 +476,12 @@ func (s *Service) processAndStore(
|
|||||||
// Continue even if caching fails
|
// Continue even if caching fails
|
||||||
}
|
}
|
||||||
|
|
||||||
return &processedVariant{
|
return &ImageResponse{
|
||||||
data: processedData,
|
Content: io.NopCloser(bytes.NewReader(processedData)),
|
||||||
contentType: processResult.ContentType,
|
ContentLength: outputSize,
|
||||||
fetchedBytes: fetchBytes,
|
ContentType: processResult.ContentType,
|
||||||
|
FetchedBytes: fetchBytes,
|
||||||
|
ETag: formatETag(cacheKey),
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"go.uber.org/fx"
|
||||||
)
|
)
|
||||||
|
|
||||||
// HTTP server configuration constants.
|
// 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() {
|
func (s *Server) serveUntilShutdown() {
|
||||||
s.httpServer = s.newHTTPServer()
|
|
||||||
|
|
||||||
s.SetupRoutes()
|
|
||||||
|
|
||||||
s.log.Info("http begin listen", "listenaddr", s.httpServer.Addr)
|
s.log.Info("http begin listen", "listenaddr", s.httpServer.Addr)
|
||||||
|
|
||||||
err := s.httpServer.ListenAndServe()
|
err := s.httpServer.ListenAndServe()
|
||||||
if err != nil && !errors.Is(err, http.ErrServerClosed) {
|
if err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||||
s.log.Error("listen error", "error", err)
|
s.log.Error("listen error", "error", err)
|
||||||
|
|
||||||
if s.cancelFunc != nil {
|
err = s.shutdowner.Shutdown(fx.ExitCode(1))
|
||||||
s.cancelFunc()
|
if err != nil {
|
||||||
|
s.log.Error("shutdown request failed", "error", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+44
-55
@@ -3,12 +3,10 @@ package server
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
|
||||||
"os/signal"
|
|
||||||
"syscall"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/getsentry/sentry-go"
|
"github.com/getsentry/sentry-go"
|
||||||
@@ -27,6 +25,10 @@ const (
|
|||||||
SentryFlushTimeout = 2 * time.Second
|
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.
|
// Params defines dependencies for Server.
|
||||||
type Params struct {
|
type Params struct {
|
||||||
fx.In
|
fx.In
|
||||||
@@ -36,6 +38,7 @@ type Params struct {
|
|||||||
Config *config.Config
|
Config *config.Config
|
||||||
Middleware *middleware.Middleware
|
Middleware *middleware.Middleware
|
||||||
Handlers *handlers.Handlers
|
Handlers *handlers.Handlers
|
||||||
|
Shutdowner fx.Shutdowner
|
||||||
}
|
}
|
||||||
|
|
||||||
// Server is the main HTTP server.
|
// Server is the main HTTP server.
|
||||||
@@ -45,15 +48,16 @@ type Server struct {
|
|||||||
globals *globals.Globals
|
globals *globals.Globals
|
||||||
mw *middleware.Middleware
|
mw *middleware.Middleware
|
||||||
h *handlers.Handlers
|
h *handlers.Handlers
|
||||||
|
shutdowner fx.Shutdowner
|
||||||
startupTime time.Time
|
startupTime time.Time
|
||||||
exitCode int
|
|
||||||
sentryEnabled bool
|
sentryEnabled bool
|
||||||
cancelFunc context.CancelFunc
|
|
||||||
httpServer *http.Server
|
httpServer *http.Server
|
||||||
router *chi.Mux
|
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) {
|
func New(lc fx.Lifecycle, params Params) (*Server, error) {
|
||||||
s := &Server{
|
s := &Server{
|
||||||
log: params.Logger.Get(),
|
log: params.Logger.Get(),
|
||||||
@@ -61,43 +65,41 @@ func New(lc fx.Lifecycle, params Params) (*Server, error) {
|
|||||||
globals: params.Globals,
|
globals: params.Globals,
|
||||||
mw: params.Middleware,
|
mw: params.Middleware,
|
||||||
h: params.Handlers,
|
h: params.Handlers,
|
||||||
|
shutdowner: params.Shutdowner,
|
||||||
}
|
}
|
||||||
|
|
||||||
lc.Append(fx.Hook{
|
lc.Append(fx.Hook{
|
||||||
OnStart: func(ctx context.Context) error {
|
OnStart: func(_ context.Context) error {
|
||||||
s.startupTime = time.Now()
|
s.startupTime = time.Now()
|
||||||
go s.Run(context.WithoutCancel(ctx))
|
|
||||||
|
|
||||||
return nil
|
err := s.enableSentry()
|
||||||
},
|
if err != nil {
|
||||||
OnStop: func(_ context.Context) error {
|
return err
|
||||||
if s.cancelFunc != nil {
|
|
||||||
s.cancelFunc()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
s.SetupRoutes()
|
||||||
|
s.httpServer = s.newHTTPServer()
|
||||||
|
|
||||||
|
go s.serveUntilShutdown()
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
},
|
},
|
||||||
|
OnStop: s.cleanShutdown,
|
||||||
})
|
})
|
||||||
|
|
||||||
return s, nil
|
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.
|
// MaintenanceMode returns whether maintenance mode is enabled.
|
||||||
func (s *Server) MaintenanceMode() bool {
|
func (s *Server) MaintenanceMode() bool {
|
||||||
return s.config.MaintenanceMode
|
return s.config.MaintenanceMode
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) enableSentry() {
|
func (s *Server) enableSentry() error {
|
||||||
s.sentryEnabled = false
|
s.sentryEnabled = false
|
||||||
|
|
||||||
if s.config.SentryDSN == "" {
|
if s.config.SentryDSN == "" {
|
||||||
return
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
err := sentry.Init(sentry.ClientOptions{
|
err := sentry.Init(sentry.ClientOptions{
|
||||||
@@ -105,55 +107,42 @@ func (s *Server) enableSentry() {
|
|||||||
Release: fmt.Sprintf("%s-%s", s.globals.Appname, s.globals.Version),
|
Release: fmt.Sprintf("%s-%s", s.globals.Appname, s.globals.Version),
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.log.Error("sentry init failure", "error", err)
|
return fmt.Errorf("sentry init failure: %w", err)
|
||||||
os.Exit(1)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
s.log.Info("sentry error reporting activated")
|
s.log.Info("sentry error reporting activated")
|
||||||
s.sentryEnabled = true
|
s.sentryEnabled = true
|
||||||
|
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) serve(ctx context.Context) int {
|
// cleanShutdown stops the HTTP server, waits for the images still being
|
||||||
ctx, cancelFunc := context.WithCancel(ctx)
|
// processed, then flushes Sentry. The first two share ShutdownTimeout. It
|
||||||
s.cancelFunc = cancelFunc
|
// 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() {
|
ctxShutdown, shutdownCancel := context.WithTimeout(ctx, ShutdownTimeout)
|
||||||
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)
|
|
||||||
defer shutdownCancel()
|
defer shutdownCancel()
|
||||||
|
|
||||||
if s.httpServer != nil {
|
|
||||||
err := s.httpServer.Shutdown(ctxShutdown)
|
err := s.httpServer.Shutdown(ctxShutdown)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.log.Error("server clean shutdown failed", "error", err)
|
s.log.Error("server clean shutdown failed", "error", err)
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
stillProcessing := s.h.WaitForProcessing(ctxShutdown)
|
||||||
|
|
||||||
if s.sentryEnabled {
|
if s.sentryEnabled {
|
||||||
sentry.Flush(SentryFlushTimeout)
|
sentry.Flush(SentryFlushTimeout)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if stillProcessing > 0 {
|
||||||
|
s.log.Error("images still being processed at shutdown",
|
||||||
|
"count", stillProcessing)
|
||||||
|
|
||||||
|
return errStillProcessing
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
+4
-7
@@ -1,18 +1,15 @@
|
|||||||
#!/bin/sh
|
#!/bin/sh
|
||||||
# script/cibuild: run the CI build. The Dockerfile runs the checks
|
# script/cibuild: run the CI build. The Dockerfile runs the checks
|
||||||
# (make fmt-check, lint, test) as build steps. This script passes a new
|
# (make fmt-check, lint, test), so a successful build implies a green
|
||||||
# CHECK_EPOCH on every run, so Docker runs those steps instead of
|
# repo. Generic: needs no adaptation. The Gitea workflow runs this on
|
||||||
# reusing cached results: a successful run means the checks ran and
|
# push.
|
||||||
# passed on this tree. Generic: needs no adaptation. The Gitea workflow
|
|
||||||
# runs this on push.
|
|
||||||
set -eu
|
set -eu
|
||||||
|
|
||||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||||
|
|
||||||
main() {
|
main() {
|
||||||
cd "$ROOT"
|
cd "$ROOT"
|
||||||
epoch="$(date +%s)$$"
|
docker build .
|
||||||
docker build --build-arg CHECK_EPOCH="$epoch" .
|
|
||||||
}
|
}
|
||||||
|
|
||||||
main "$@"
|
main "$@"
|
||||||
|
|||||||
+3
-7
@@ -1,9 +1,7 @@
|
|||||||
#!/bin/sh
|
#!/bin/sh
|
||||||
# script/docker: build the Docker image tagged with the project name.
|
# script/docker: build the Docker image tagged with the project name.
|
||||||
# Identical in all repos; the tag comes from script/projectname. Like
|
# Identical in all repos; the tag comes from script/projectname.
|
||||||
# script/cibuild, it passes a new CHECK_EPOCH, so the build runs the
|
# Generic: needs no adaptation.
|
||||||
# checks instead of reusing cached results. Generic: needs no
|
|
||||||
# adaptation.
|
|
||||||
set -eu
|
set -eu
|
||||||
|
|
||||||
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd -P)"
|
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd -P)"
|
||||||
@@ -11,9 +9,7 @@ ROOT="$(cd "$SCRIPT_DIR/.." && pwd -P)"
|
|||||||
|
|
||||||
main() {
|
main() {
|
||||||
cd "$ROOT"
|
cd "$ROOT"
|
||||||
epoch="$(date +%s)$$"
|
docker build -t "$("$SCRIPT_DIR/projectname")" .
|
||||||
docker build --build-arg CHECK_EPOCH="$epoch" \
|
|
||||||
-t "$("$SCRIPT_DIR/projectname")" .
|
|
||||||
}
|
}
|
||||||
|
|
||||||
main "$@"
|
main "$@"
|
||||||
|
|||||||
Reference in New Issue
Block a user