Author SHA1 Message Date
sneak 410c167c70 Keep counts not written by the request's deadline in memory (closes #224)
check / check (push) Waiting to run
TestService_Get_ReturnsByItsDeadline failed on a busy host because its
request, past its deadline, still waited to write its miss count, a write
with no deadline; in pixad that write waits for the one database
connection every request shares. Each count write now keeps the request's
deadline but not its cancellation. A count not written by then is kept in
memory, where Stats includes it, and one goroutine of the cache writes
those counts one UPDATE at a time, and once more at shutdown before the
database closes. The test phase also runs go test with -parallel 4: on a
busy host, as many tests at once as there are CPUs wait so long to be
scheduled that a timed request can fail.

Model: opus-5-5
2026-10-08 15:44:39 +00:00
clawbot bdde021b45 Save PNG compressed, WebP at effort 4 and AVIF at effort 1 (closes #232)
check / check (push) Waiting to run
Each output format now has its own govips export with its settings
named, in place of govips' generic Export, which sent libvips a zero for
some settings it was not given: PNG had no compression and WebP effort
0. PNG now gets libvips' default compression, 6, and WebP its default
effort, 4. GIF and JPEG output is unchanged.

AVIF was at libvips' default effort, 4, which takes minutes for an
8192x8192 image with one libvips thread, far past the default
downstream_timeout. Effort 1, the lowest govips can set, takes about 51
seconds for an image of random pixels, the worst case, and WebP at 4
about 43. A 16-bit source gets 8 bits per sample, not libvips' 12, which
take over four times as long.

Model: opus-5-5
2026-10-08 14:13:25 +02:00
clawbot db91ab29e6 Serve JPEG XL when a request names no format (closes #222)
check / check (push) Waiting to run
A /v1/image/ URL whose last segment is a size with no format, such as
800x600 or orig, is served as JPEG XL and signed as jxl, so it shares
the signature of the same URL ending in .jxl. An encrypted URL whose
token holds no format is served as JPEG XL, as encurl.DefaultFormat is
now jxl. The generator page selects JPEG XL by default, and a form
with an empty format, or none, makes a URL whose name ends in .jxl.

The image processor no longer takes an empty format as orig: both
routes give every request a format, so it refuses a request with none
instead of keeping a second default. auto still ends with JPEG.

Model: opus-5-5
2026-10-08 12:49:56 +02:00
clawbot 8597253ffd Serve and accept JPEG XL as an image format (part of #222)
check / check (push) Waiting to run
A JPEG XL source is accepted, and orig of one is JPEG XL. The format
jxl works in plain and encrypted URLs and on the generator page,
served as image/jxl; auto chooses it first when Accept names
image/jxl.

govips sends libvips a JPEG XL distance, which overrides the quality,
so q becomes a distance, keeping 100 lossy. Metadata is removed from
the image before the JPEG XL save, as govips cannot have libvips strip
it, and the image is given 72 dpi so that the EXIF block libvips 8.16
adds holds nothing from the source. The sRGB conversion moved ahead of
the format switch, and a CMYK image with no ICC profile is converted
to sRGB, as libvips cannot save CMYK as JPEG XL.

Model: opus-5-5
2026-10-08 11:01:12 +02:00
18 changed files with 1089 additions and 100 deletions
+6 -3
View File
@@ -32,10 +32,13 @@ RUN script/bootstrap --cgo
COPY . .
# Without -v first; on a failure, again with -v for the details, and
# the step fails even if the second run passes.
RUN go test -count=1 -timeout 90s -race -cover ./... || \
# the step fails even if the second run passes. -parallel 4: by default
# a package runs as many of its tests at once as the host has CPUs, and
# on a busy host they then wait so long to be scheduled that a test
# that times a request can see it return a second late.
RUN go test -count=1 -timeout 90s -race -parallel 4 -cover ./... || \
{ echo "--- Rerunning with -v for details ---"; \
go test -count=1 -timeout 90s -race -v ./...; exit 1; }
go test -count=1 -timeout 90s -race -parallel 4 -v ./...; exit 1; }
# Build stage. Nothing is wanted from the two phases above: these copies
# make BuildKit build them first, so this stage runs only when lint and
+35 -26
View File
@@ -83,10 +83,12 @@ which answers 200 whenever pixa is running, in maintenance mode too (see
`maintenance_mode`).
On SIGTERM or SIGINT pixa stops accepting connections, gives the requests in
progress and the images being processed 5 seconds to finish, and exits: with 0,
or with 1 when images were still being processed after those 5 seconds or
another part of pixa failed to stop. A request not finished by then is cut off.
`docker stop` waits 10 seconds before it kills the container.
progress and the images being processed 5 seconds to finish, then writes to the
database the counts of cache hits, misses, fetches and conversions that requests
could not write by their deadline, and exits: with 0, or with 1 when images were
still being processed after those 5 seconds, some of those counts could not be
written, or another part of pixa failed to stop. A request not finished by then
is cut off. `docker stop` waits 10 seconds before it kills the container.
Outside Docker, pixa needs libvips (the image has 8.16) and libheif to run, as
it uses libvips through CGO. pixad does not start unless libvips has its JPEG XL
@@ -175,20 +177,21 @@ path under `/v1/` answers 200, in maintenance mode too.
with the page naming a field that is not valid; 500 when the URL cannot be
made.
- `GET /logout` — end the login session. Needs: nothing. Answers: 303 to `/`.
- `GET` or `HEAD` `/v1/image/<host>/<path>/<size>.<format>` — an image, fetched,
resized and converted (below). Needs: a signature, unless the host is
allowlisted (see Source Hosts). Answers: 200; 304 when `If-None-Match` matches
the image's `ETag`; 400 for a URL or parameter that is not valid, or for the
format `auto` an `Accept` header that is not valid; 406 for the format `auto`
when `Accept` allows none of the formats it chooses from; 401 for a missing or
wrong signature, a missing `exp` or an `exp` in the past; 403 when the
request's `Referer` names a host in `referer_blocklist`, checked before the
signature, the cache and the upstream fetch; 403 when the upstream host, or a
host it redirects to, is `localhost`, ends in `.localhost` or `.local`, or has
an address in a blocked network (see `blocked_networks`); 502 when the
upstream answered with an error status, and for 5 minutes after that for the
same source URL; 503 when pixa is busy or in maintenance mode; 500 for any
other failure.
- `GET` or `HEAD` `/v1/image/<host>/<path>/<size>.<format>`, or
`/v1/image/<host>/<path>/<size>` with no format — an image, fetched, resized
and converted (below). Needs: a signature, unless the host is allowlisted (see
Source Hosts). Answers: 200; 304 when `If-None-Match` matches the image's
`ETag`; 400 for a URL or parameter that is not valid, or for the format `auto`
an `Accept` header that is not valid; 406 for the format `auto` when `Accept`
allows none of the formats it chooses from; 401 for a missing or wrong
signature, a missing `exp` or an `exp` in the past; 403 when the request's
`Referer` names a host in `referer_blocklist`, checked before the signature,
the cache and the upstream fetch; 403 when the upstream host, or a host it
redirects to, is `localhost`, ends in `.localhost` or `.local`, or has an
address in a blocked network (see `blocked_networks`); 502 when the upstream
answered with an error status, and for 5 minutes after that for the same
source URL; 503 when pixa is busy or in maintenance mode; 500 for any other
failure.
- `GET` or `HEAD` `/v1/e/<token>/<name>` — an image through an encrypted URL
(see Encrypted URLs). Needs: nothing but the URL. Answers: 200; 304 when
`If-None-Match` matches the image's `ETag`; 400 for a token that does not
@@ -231,10 +234,11 @@ proxy in front of pixa must pass that header on unchanged. A form body over 1
MiB is refused with 413. The image routes answer the errors listed for them with
JSON holding `error`, `status` and `timestamp`.
An image URL has this form:
An image URL has one of these forms, the second with no format:
```
/v1/image/<host>/<path>/<size>.<format>?sig=<signature>&exp=<expiration>&q=<quality>&fit=<fit>
/v1/image/<host>/<path>/<size>?sig=<signature>&exp=<expiration>&q=<quality>&fit=<fit>
```
Images are only fetched from origins using TLS with valid certificates, unless
@@ -245,7 +249,9 @@ A request whose query string cannot be decoded, or gives any parameter more than
once, is refused with 400.
- `<format>`: one of `orig` (or `original`), `jpeg` (or `jpg`), `png`, `webp`,
`avif`, `jxl` (JPEG XL), `gif`, or `auto` (below)
`avif`, `jxl` (JPEG XL), `gif`, or `auto` (below). A URL with no format (the
second form, with no dot after the size) is served as JPEG XL, the default, as
with `jxl`
- `<size>`: `orig` or `<width>x<height>` (e.g. `800x600`)
- `sig` and `exp`: the signature and its expiry, needed unless the host is
allowlisted (see Signature Specification)
@@ -321,13 +327,15 @@ nor change what it asks for.
lasts 30 days, or until `/logout`.
2. On the generator page, give the source image's URL, the width and height, the
format, quality and fit, and how long the URL lasts, then submit the form
(`POST /generate`). Width and height both empty or `0` keep the original
size; if only one of them is empty or `0`, that side is scaled to keep the
image's proportions.
(`POST /generate`). The format is JPEG XL unless another is chosen; a form
sent with an empty format, or none, also makes a JPEG XL URL. Width and
height both empty or `0` keep the original size; if only one of them is empty
or `0`, that side is scaled to keep the image's proportions.
3. The page shows the URL, `https://<host>/v1/e/<token>/img.<format>`, and when
it expires. `<host>` is the host the page was opened on, and the URL starts
with `http` instead while `debug` is on. The name after the token is ignored
and only gives the URL a file extension, `jpg` for `orig` and `auto`.
and only gives the URL a file extension, `jpg` for `orig` and `auto`, and
`jxl` for a form with no format.
The token holds the source's host, path and query and the size, format, quality,
fit and expiry, encrypted with a key derived from `signing_key`. The source
@@ -387,8 +395,9 @@ Where:
- `width` — requested width in pixels, `0` for original
- `height` — requested height in pixels, `0` for original
- `format` — output format, one of those listed under Routes, with `original`
signed as `orig` and `jpg` as `jpeg`; `auto` is signed as `auto`, not as the
format chosen for the request
signed as `orig`, `jpg` as `jpeg`, and no format as `jxl`, so a URL with no
format has the signature of the same URL ending in `.jxl`; `auto` is signed as
`auto`, not as the format chosen for the request
- `expiration` — the URL's `exp` query parameter, the Unix timestamp when the
signature expires; a request whose `exp` is not a whole number, an empty
`exp=` included, is refused with 400
+37
View File
@@ -30,6 +30,43 @@ P2: security: per-IP rate limiting on the image routes
# Completed Steps
- 2026-10-08 a request past its deadline no longer waits for the database to
count it (closes #224). On a busy host, `TestService_Get_ReturnsByItsDeadline`
failed because its request, past its deadline, still waited for its miss count
to be written, a write with no deadline, which in pixad waits for the one
database connection every request shares. Each count write now keeps the
request's deadline but not its cancellation. A count not written by then is
kept in memory, where `Stats` sees it, and one goroutine of the cache writes
those counts, one write at a time, and once more at shutdown before the
database closes. The test phase of the `Dockerfile` also runs at most 4 tests
of a package at once (`go test -parallel 4`), where it ran as many as the host
has CPUs: on a busy host they then wait so long to be scheduled that a test
that times a request can fail. `-p`, how many packages are tested at once, is
unchanged, as capping it also slows compiling.
- 2026-10-08 every output format is saved with settings pixa sets on purpose
(closes #232): each format has its own govips export, as JPEG XL does, in
place of govips' generic `Export`, which sent libvips a zero for some settings
it was not given. PNG gets libvips' default compression, 6, where it had none,
and WebP libvips' default effort, 4, where it had 0. GIF was already at
libvips' default effort, 7, and JPEG output is unchanged. AVIF is saved at
effort 1, the lowest govips can set, where it had libvips' default, 4. With
one libvips thread, an 8192x8192 image of random pixels, the worst case, takes
about 51 seconds as AVIF at effort 1 and 43 as WebP at effort 4, against the
default `downstream_timeout` of 60 seconds. On an image of milder noise, which
AVIF at effort 1 saves in about 12 seconds, effort 4 takes minutes. AVIF is
also saved with 8 bits per sample from a 16-bit source, which libvips would
save with 12: a 16-bit 8192x8192 image of milder noise takes about 54 seconds
at effort 1 with 12 bits, nearly all of the default `downstream_timeout`, and
about 12 with 8.
- 2026-10-08 JPEG XL is the default output (closes #222): a `/v1/image/` URL
whose last segment is a size with no format, such as `800x600` or `orig`, is
served as JPEG XL and signed as `jxl`, so it has the signature of the same URL
ending in `.jxl`. An encrypted URL whose token holds no format is served as
JPEG XL (`encurl.DefaultFormat`). The generator page selects JPEG XL until
another format is chosen, and a form with an empty format, or none, makes a
JPEG XL URL whose name ends in `.jxl`. The image processor no longer takes an
empty format as `orig`: both routes give every request a format, and it
refuses a request with none. `auto` still ends with JPEG.
- 2026-10-08 JPEG XL as an input and output format (part of #222): a source
whose bytes start with either JPEG XL signature, the bare codestream's `FF 0A`
or the container's, is detected as `image/jxl`, which the upstream fetch
+2 -2
View File
@@ -14,7 +14,7 @@ import (
// Default values for optional fields.
const (
DefaultQuality = 85
DefaultFormat = imgcache.FormatOriginal
DefaultFormat = imgcache.FormatJXL
DefaultFitMode = imgcache.FitCover
// HKDF salt for URL encryption key derivation
@@ -37,7 +37,7 @@ type Payload struct {
SourceQuery string `cbor:"q,omitempty"` // optional
Width int `cbor:"w,omitempty"` // 0 = original
Height int `cbor:"ht,omitempty"` // 0 = original
Format imgcache.ImageFormat `cbor:"f,omitempty"` // default: orig
Format imgcache.ImageFormat `cbor:"f,omitempty"` // default: jxl
Quality int `cbor:"ql,omitempty"` // default: 85
FitMode imgcache.FitMode `cbor:"fm,omitempty"` // default: cover
ExpiresAt int64 `cbor:"e,omitempty"` // 0 = never expires
+8 -3
View File
@@ -369,10 +369,15 @@ func (s *Handlers) buildGeneratedURL(r *http.Request, token, format string) stri
scheme = "http"
}
// Determine file extension for the trailing filename
// Determine file extension for the trailing filename. A form with no
// format makes a token with none, which is served as encurl.DefaultFormat.
ext := format
if ext == "" || ext == "orig" || ext == "auto" {
ext = "jpg" // Default extension
switch format {
case "":
ext = string(encurl.DefaultFormat)
case "orig", "auto":
ext = "jpg"
}
return scheme + "://" + r.Host + "/v1/e/" + url.PathEscape(token) + "/img." + ext
+11 -2
View File
@@ -4,6 +4,7 @@ package handlers
import (
"context"
"encoding/json"
"errors"
"log/slog"
"net/http"
"time"
@@ -71,16 +72,21 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) {
}
lc.Append(fx.Hook{
//nolint:contextcheck // the eviction loop outlives OnStart; OnStop cancels it
//nolint:contextcheck // the cache's goroutines outlive OnStart; OnStop stops them
OnStart: func(_ context.Context) error {
return s.initImageService()
},
// The pending counts are written here, before the database's stop
// hook closes it.
OnStop: func(ctx context.Context) error {
if s.imgCache == nil {
return nil
}
return s.imgCache.StopEviction(ctx)
return errors.Join(
s.imgCache.StopEviction(ctx),
s.imgCache.StopPendingCountWrites(ctx),
)
},
})
@@ -122,6 +128,9 @@ func (s *Handlers) initImageService() error {
// write-pressure passes. No-op when the disk cache is disabled.
cache.StartEviction(imgcache.DefaultEvictionInterval)
// Writes the counts requests could not write by their deadline
cache.StartPendingCountWrites()
// Create the fetcher config
fetcherCfg := httpfetcher.DefaultConfig()
fetcherCfg.AllowHTTP = s.config.AllowHTTP
+2 -1
View File
@@ -18,7 +18,8 @@ import (
)
// HandleImage handles the main image proxy route:
// /v1/image/<host>/<path>/<width>x<height>.<format>
// /v1/image/<host>/<path>/<width>x<height>.<format>, or with no format
// /v1/image/<host>/<path>/<width>x<height>
func (s *Handlers) HandleImage() http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if s.refuseBlockedReferer(w, r) {
@@ -0,0 +1,162 @@
package handlers
import (
"fmt"
"log/slog"
"maps"
"net/http"
"net/http/httptest"
"net/url"
"strings"
"testing"
"time"
"github.com/davidbyttow/govips/v2/vips"
"sneak.berlin/go/pixa/internal/encurl"
"sneak.berlin/go/pixa/internal/imgcache"
"sneak.berlin/go/pixa/internal/signature"
)
// requireJPEGXL requires that rec answers 200 with a JPEG XL image.
func requireJPEGXL(t *testing.T, rec *httptest.ResponseRecorder) {
t.Helper()
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d; body %q",
rec.Code, http.StatusOK, rec.Body.String())
}
if got := rec.Header().Get("Content-Type"); got != jxlType {
t.Errorf("Content-Type = %q, want %s", got, jxlType)
}
if got := vips.DetermineImageType(rec.Body.Bytes()); got != vips.ImageTypeJXL {
t.Errorf("body is %s, want jxl", vips.ImageTypes[got])
}
}
// TestImageWithoutFormat_ServesJPEGXL verifies that a /v1/image/ URL whose
// last segment is a size with no format, 50x50 or orig, answers JPEG XL.
func TestImageWithoutFormat_ServesJPEGXL(t *testing.T) {
t.Parallel()
route := newImageRoute(t, newPhotoFetcher(t, allowlistedHost))
for _, size := range []string{"50x50", "orig"} {
target := "/v1/image/" + allowlistedHost + photoPath + "/" + size
requireJPEGXL(t, sendGet(t, route, target))
}
}
// TestImageWithoutFormat_SignedAsJXL verifies that a /v1/image/ URL with no
// format is signed as jxl: the signature made for the URL ending in .jxl is
// accepted for the same URL without .jxl.
func TestImageWithoutFormat_SignedAsJXL(t *testing.T) {
t.Parallel()
route := newImageRoute(t, newPhotoFetcher(t, signedHost))
expires := time.Now().Add(time.Hour)
sig := signature.New(testSigningKey).Sign(&signature.Request{
SourceHost: signedHost,
SourcePath: photoPath,
Width: 50,
Height: 50,
Format: string(imgcache.FormatJXL),
Quality: encurl.DefaultQuality,
FitMode: string(imgcache.FitCover),
Expires: expires,
})
query := fmt.Sprintf("?sig=%s&exp=%d", sig, expires.Unix())
for _, size := range []string{"50x50.jxl", "50x50"} {
target := "/v1/image/" + signedHost + photoPath + "/" + size + query
requireJPEGXL(t, sendGet(t, route, target))
}
}
// TestImageEncWithoutFormat_ServesJPEGXL verifies that an encrypted URL whose
// token holds no format answers JPEG XL.
func TestImageEncWithoutFormat_ServesJPEGXL(t *testing.T) {
t.Parallel()
h, srv := newSignedHostServer(t, slog.New(slog.DiscardHandler))
token, err := h.encGen.Generate(&encurl.Payload{
SourceHost: signedHost,
SourcePath: photoPath,
Width: 50,
Height: 50,
})
if err != nil {
t.Fatalf("Generate() error = %v", err)
}
requireJPEGXL(t, getEncToken(srv, token))
}
// TestGeneratorPage_SelectsJPEGXL verifies that the generator page's format
// choice is JPEG XL until another is chosen.
func TestGeneratorPage_SelectsJPEGXL(t *testing.T) {
t.Parallel()
h, srv := newCSRFTestRouter(t)
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/", nil)
req.AddCookie(newSessionCookie(t, h))
rec := httptest.NewRecorder()
srv.ServeHTTP(rec, req)
if !strings.Contains(rec.Body.String(), `<option value="jxl" selected>`) {
t.Errorf("generator page does not select JPEG XL: %s", rec.Body.String())
}
}
// TestGeneratePost_NoFormat_MakesJPEGXLURL verifies that the generator form
// sent with an empty format field, or with none, makes a URL whose name ends
// in .jxl and which answers JPEG XL.
func TestGeneratePost_NoFormat_MakesJPEGXLURL(t *testing.T) {
t.Parallel()
photo := url.Values{
sourceURLField: {"https://" + signedHost + photoPath},
widthField: {"50"},
heightField: {"50"},
}
emptyFormat := maps.Clone(photo)
emptyFormat.Set(formatField, "")
for name, form := range map[string]url.Values{
"empty format field": emptyFormat,
"no format field": photo,
} {
t.Run(name, func(t *testing.T) {
t.Parallel()
_, imageSrv := newSignedHostServer(t, slog.New(slog.DiscardHandler))
rec := generatePost(t, form)
match := generatedURLPattern.FindStringSubmatch(rec.Body.String())
if match == nil {
t.Fatalf("generator page shows no URL: %d %s",
rec.Code, rec.Body.String())
}
t.Logf("generated URL path: %s", match[1])
if !strings.HasSuffix(match[1], "/img.jxl") {
t.Errorf("generated URL %s does not end in /img.jxl", match[1])
}
imageRec := httptest.NewRecorder()
imageSrv.ServeHTTP(imageRec, httptest.NewRequestWithContext(
t.Context(), http.MethodGet, match[1], nil))
requireJPEGXL(t, imageRec)
})
}
}
@@ -0,0 +1,21 @@
package imageprocessor
import (
"bytes"
"errors"
"testing"
)
// TestImageProcessor_EmptyFormatRefused verifies that a request with no format
// is refused, as both image routes give every request a format before it is
// processed.
func TestImageProcessor_EmptyFormatRefused(t *testing.T) {
t.Parallel()
_, err := New(Params{}).Process(
t.Context(), bytes.NewReader(createTestJPEG(t, 20, 20)), &Request{},
)
if !errors.Is(err, ErrUnsupportedOutputFormat) {
t.Errorf("Process() error = %v, want %v", err, ErrUnsupportedOutputFormat)
}
}
@@ -0,0 +1,145 @@
package imageprocessor
import (
"bytes"
"image"
"image/color"
"image/png"
"testing"
"github.com/davidbyttow/govips/v2/vips"
)
// processedSize runs input through Process and returns the output's size in
// bytes.
func processedSize(t *testing.T, input []byte, req *Request) int64 {
t.Helper()
result, err := New(Params{}).Process(t.Context(), bytes.NewReader(input), req)
if err != nil {
t.Fatalf("Process() error = %v", err)
}
_ = result.Content.Close()
return result.ContentLength
}
// encodePNG encodes img as PNG with Go's encoder.
func encodePNG(t *testing.T, img image.Image) []byte {
t.Helper()
var buf bytes.Buffer
err := png.Encode(&buf, img)
if err != nil {
t.Fatalf("failed to encode test PNG: %v", err)
}
return buf.Bytes()
}
// decode decodes input with vips, as Process does.
func decode(t *testing.T, input []byte) *vips.ImageRef {
t.Helper()
img, err := vips.NewImageFromBuffer(input)
if err != nil {
t.Fatalf("failed to decode input: %v", err)
}
t.Cleanup(img.Close)
return img
}
// TestImageProcessor_PNGIsCompressed verifies that a PNG output is
// compressed: a flat 200x150 image, 90,000 bytes of raw pixels, comes out at
// a small fraction of that.
func TestImageProcessor_PNGIsCompressed(t *testing.T) {
t.Parallel()
const width, height = 200, 150
flat := image.NewRGBA(image.Rect(0, 0, width, height))
for y := range height {
for x := range width {
flat.Set(x, y, color.RGBA{R: 40, G: 120, B: 200, A: 255})
}
}
size := processedSize(t, encodePNG(t, flat), &Request{Format: FormatPNG})
const rawBytes = width * height * 3
if size > rawBytes/10 {
t.Errorf("PNG output is %d bytes, want under a tenth of its %d raw",
size, rawBytes)
}
}
// TestImageProcessor_WebPAtDefaultEffort verifies that WebP is saved at
// libvips' default effort, 4: the output is the size of the same image saved
// at that effort.
func TestImageProcessor_WebPAtDefaultEffort(t *testing.T) {
t.Parallel()
input := createTestJPEG(t, 200, 150)
size := processedSize(t, input, &Request{Format: FormatWebP, Quality: 85})
want, _, err := decode(t, input).ExportWebp(&vips.WebpExportParams{
StripMetadata: true,
Quality: 85,
ReductionEffort: 4,
})
if err != nil {
t.Fatalf("ExportWebp() error = %v", err)
}
if size != int64(len(want)) {
t.Errorf("WebP output is %d bytes, want %d, the size at effort 4",
size, len(want))
}
}
// TestImageProcessor_AVIFAtEffort1 verifies that AVIF is saved at effort 1
// with 8 bits per sample: the output is the size of the same image saved
// with those settings. The source is 16-bit, which libvips saves with 12
// bits when it is not given a bit depth, and 640x480: on a small image,
// such as 64x48, efforts 1 and 2 give the same output.
func TestImageProcessor_AVIFAtEffort1(t *testing.T) {
t.Parallel()
const width, height = 640, 480
source := image.NewRGBA64(image.Rect(0, 0, width, height))
for y := range height {
for x := range width {
source.Set(x, y, color.RGBA64{
R: uint16((x * 65535 / width) & 0xffff),
G: uint16((y * 65535 / height) & 0xffff),
B: 32768,
A: 65535,
})
}
}
input := encodePNG(t, source)
size := processedSize(t, input, &Request{Format: FormatAVIF, Quality: 85})
want, _, err := decode(t, input).ExportAvif(&vips.AvifExportParams{
StripMetadata: true,
Quality: 85,
Effort: 1,
Bitdepth: 8,
})
if err != nil {
t.Fatalf("ExportAvif() error = %v", err)
}
if size != int64(len(want)) {
t.Errorf("AVIF output is %d bytes, want %d, the size at effort 1 "+
"and 8 bits", size, len(want))
}
}
+88 -30
View File
@@ -256,9 +256,9 @@ func (p *ImageProcessor) Process(
}
}
// Determine output format
// orig is the source's own format; encode refuses an empty format
outputFormat := req.Format
if outputFormat == FormatOriginal || outputFormat == "" {
if outputFormat == FormatOriginal {
outputFormat = p.formatFromString(inputFormat)
}
@@ -504,36 +504,21 @@ func (p *ImageProcessor) encode(
}
}
var params vips.ExportParams
switch format {
case FormatJPEG:
params = vips.ExportParams{
Format: vips.ImageTypeJPEG,
Quality: quality,
}
return exportJPEG(img, quality)
case FormatPNG:
params = vips.ExportParams{
Format: vips.ImageTypePNG,
}
return exportPNG(img)
case FormatGIF:
params = vips.ExportParams{
Format: vips.ImageTypeGIF,
}
return exportGIF(img)
case FormatWebP:
params = vips.ExportParams{
Format: vips.ImageTypeWEBP,
Quality: quality,
}
return exportWebP(img, quality)
case FormatAVIF:
params = vips.ExportParams{
Format: vips.ImageTypeAVIF,
Quality: quality,
}
return exportAVIF(img, quality)
case FormatJXL:
return exportJXL(img, quality)
@@ -544,17 +529,90 @@ func (p *ImageProcessor) encode(
default:
return nil, fmt.Errorf("%w: %s", ErrUnsupportedOutputFormat, format)
}
}
// Drop EXIF, XMP, IPTC and the ICC profile. govips ignores this for
// GIF, which carries none of them.
params.StripMetadata = true
// govips sends libvips Go's zero value for some settings an export leaves
// out, such as no compression at all for PNG, so each export below sets
// every setting whose zero value is not what pixa wants. Stripping metadata
// drops EXIF, XMP, IPTC and the ICC profile.
output, _, err := img.Export(&params)
if err != nil {
return nil, err
}
// exportJPEG encodes img as JPEG at quality, without metadata. The settings
// it leaves out are at libvips' defaults.
func exportJPEG(img *vips.ImageRef, quality int) ([]byte, error) {
output, _, err := img.ExportJpeg(&vips.JpegExportParams{
StripMetadata: true,
Quality: quality,
})
return output, nil
return output, err
}
// pngCompression is libvips' default PNG compression, from 0 (none) to 9.
const pngCompression = 6
// exportPNG encodes img as PNG at libvips' default compression and row
// filter, without metadata.
func exportPNG(img *vips.ImageRef) ([]byte, error) {
output, _, err := img.ExportPng(&vips.PngExportParams{
StripMetadata: true,
Compression: pngCompression,
Filter: vips.PngFilterNone,
})
return output, err
}
// gifEffort is libvips' default GIF effort, from 1 to 10.
const gifEffort = 7
// exportGIF encodes img as GIF at libvips' default effort. govips cannot
// have libvips strip metadata from GIF, which carries none.
func exportGIF(img *vips.ImageRef) ([]byte, error) {
output, _, err := img.ExportGIF(&vips.GifExportParams{Effort: gifEffort})
return output, err
}
// webpEffort is libvips' default WebP effort, from 0 (fastest) to 6.
const webpEffort = 4
// exportWebP encodes img as lossy WebP at quality and libvips' default
// effort, without metadata.
func exportWebP(img *vips.ImageRef, quality int) ([]byte, error) {
output, _, err := img.ExportWebp(&vips.WebpExportParams{
StripMetadata: true,
Quality: quality,
ReductionEffort: webpEffort,
})
return output, err
}
// avifEffort is the AVIF effort, from 0 (fastest) to 9; 1 is the lowest
// govips can set. With one thread, as pixad runs libvips, 1 takes about 51
// seconds to save an 8192x8192 image of random pixels, the worst case,
// against the default downstream_timeout of 60 seconds. On an image of
// milder noise, which 1 saves in about 12 seconds, 2 takes nearly a minute
// and libvips' default, 4, takes minutes.
const avifEffort = 1
// avifBitdepth is the AVIF bit depth, 8 bits per sample for every image.
// libvips would save a 16-bit image with 12, but at avifEffort that takes
// about 54 seconds for a 16-bit 8192x8192 image of milder noise, nearly all
// of the default downstream_timeout, and about 12 seconds with 8.
const avifBitdepth = 8
// exportAVIF encodes img as lossy AVIF at quality, avifEffort and
// avifBitdepth, without metadata.
func exportAVIF(img *vips.ImageRef, quality int) ([]byte, error) {
output, _, err := img.ExportAvif(&vips.AvifExportParams{
StripMetadata: true,
Quality: quality,
Effort: avifEffort,
Bitdepth: avifBitdepth,
})
return output, err
}
// jxlResolution is the resolution every JPEG XL image is saved with, in
+71 -8
View File
@@ -11,6 +11,8 @@ import (
"io"
"log/slog"
"path/filepath"
"sync"
"sync/atomic"
"time"
lru "github.com/hashicorp/golang-lru/v2"
@@ -79,6 +81,26 @@ type Cache struct {
evictionDone chan struct{}
evictionCancel context.CancelFunc
// The pending counts: hits, misses, upstream fetches and transforms
// whose write to the database missed its request's deadline, kept
// here until they are written by the goroutine that
// StartPendingCountWrites starts. As for eviction, the channels are
// created in NewCache: pendingCountsAdded wakes that goroutine, and
// pendingCountsDone is closed when it returns. pendingCountsCancel,
// set by StartPendingCountWrites, cancels its context, which tells it
// to finish.
pendingHits atomic.Int64
pendingMisses atomic.Int64
pendingUpstreamFetches atomic.Int64
pendingUpstreamFetchBytes atomic.Int64
pendingTransforms atomic.Int64
pendingCountsAdded chan struct{}
pendingCountsDone chan struct{}
pendingCountsCancel context.CancelFunc
// Held through each write of the pending counts, and by Stats to read them
pendingCountsWriteMutex sync.Mutex
// metaCache holds the content types of the variants most recently
// stored or served, so a hit does not read the variant's .meta file.
// It never stands in for the variant file, which is always opened.
@@ -139,6 +161,9 @@ func newCache(
metaCache: metaCache,
contentLocks: newContentLock(),
pendingCountsAdded: make(chan struct{}, 1),
pendingCountsDone: make(chan struct{}),
reconciliationPageSize: defaultReconciliationPageSize,
}
@@ -474,12 +499,21 @@ func (c *Cache) CleanExpired(ctx context.Context) error {
func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
var stats CacheStats
// So that no write of the pending counts falls between the two reads
c.pendingCountsWriteMutex.Lock()
// Fetch hit/miss counts from the stats table
err := c.db.QueryRowContext(ctx, `
SELECT hit_count, miss_count
FROM cache_stats WHERE id = 1
`).Scan(&stats.HitCount, &stats.MissCount)
// Hits and misses not yet written count too (see StartPendingCountWrites)
stats.HitCount += c.pendingHits.Load()
stats.MissCount += c.pendingMisses.Load()
c.pendingCountsWriteMutex.Unlock()
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return nil, fmt.Errorf("failed to get cache stats: %w", err)
}
@@ -511,19 +545,30 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
}
// IncrementStats counts a cache hit or miss, and an upstream fetch that read
// fetchBytes bytes, as IncrementUpstreamFetch does.
// fetchBytes bytes, as IncrementUpstreamFetch does. Like the other Increment
// methods, it writes to the database with ctx's deadline but not its
// cancellation, so an ended request is still counted but does not wait for
// the database past its deadline; a count not written by then becomes a
// pending count (see StartPendingCountWrites).
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
var err error
pendingCount := &c.pendingMisses
if hit {
_, err = c.db.ExecContext(ctx, `
pendingCount = &c.pendingHits
_, err = c.db.ExecContext(countCtx, `
UPDATE cache_stats
SET hit_count = hit_count + 1,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`)
} else {
_, err = c.db.ExecContext(ctx, `
_, err = c.db.ExecContext(countCtx, `
UPDATE cache_stats
SET miss_count = miss_count + 1,
last_updated_at = CURRENT_TIMESTAMP
@@ -531,7 +576,10 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
`)
}
if err != nil {
switch {
case errors.Is(err, context.DeadlineExceeded):
c.addPendingCount(pendingCount, 1)
case err != nil:
c.log.Warn("failed to count cache hit or miss", "hit", hit, "error", err)
}
@@ -545,14 +593,22 @@ func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
return
}
_, err := c.db.ExecContext(ctx, `
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
_, err := c.db.ExecContext(countCtx, `
UPDATE cache_stats
SET upstream_fetch_count = upstream_fetch_count + 1,
upstream_fetch_bytes = upstream_fetch_bytes + ?,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`, fetchBytes)
if err != nil {
switch {
case errors.Is(err, context.DeadlineExceeded):
c.addPendingCount(&c.pendingUpstreamFetches, 1)
c.addPendingCount(&c.pendingUpstreamFetchBytes, fetchBytes)
case err != nil:
c.log.Warn("failed to count upstream fetch",
"fetch_bytes", fetchBytes, "error", err)
}
@@ -560,13 +616,20 @@ func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
// IncrementTransformCount counts one image transcoded by the image processor.
func (c *Cache) IncrementTransformCount(ctx context.Context) {
_, err := c.db.ExecContext(ctx, `
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
_, err := c.db.ExecContext(countCtx, `
UPDATE cache_stats
SET transform_count = transform_count + 1,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`)
if err != nil {
switch {
case errors.Is(err, context.DeadlineExceeded):
c.addPendingCount(&c.pendingTransforms, 1)
case err != nil:
c.log.Warn("failed to count transform", "error", err)
}
}
+141
View File
@@ -0,0 +1,141 @@
package imgcache
import (
"context"
"fmt"
"sync/atomic"
)
// StartPendingCountWrites starts the goroutine that writes the pending
// counts to the database whenever there are some, one UPDATE at a time:
// however many requests pass their deadline before their counts are
// written, at most this one write waits for the database. It is a no-op
// when already started. The goroutine outlives the caller, so it runs with
// its own context, which StopPendingCountWrites cancels to tell it to
// finish.
func (c *Cache) StartPendingCountWrites() {
if c.pendingCountsCancel != nil {
return
}
ctx, cancel := context.WithCancel(context.Background())
c.pendingCountsCancel = cancel
go func() {
c.pendingCountWriteLoop(ctx)
}()
}
// StopPendingCountWrites tells the goroutine StartPendingCountWrites
// started, if it did, to finish, and waits for its write under way to end
// and for it to return. It then writes the pending counts left, so that
// they are written before the database closes. It waits for all of this at
// most until ctx ends, then logs the counts not written and returns an
// error.
func (c *Cache) StopPendingCountWrites(ctx context.Context) error {
var err error
if c.pendingCountsCancel != nil {
c.pendingCountsCancel()
select {
case <-c.pendingCountsDone:
case <-ctx.Done():
// The write under way is left to finish: the database's
// close waits for a query under way.
err = fmt.Errorf("pending counts still being written: %w", ctx.Err())
}
}
if err == nil {
err = c.writePendingCounts(ctx)
}
if err != nil {
c.log.Error("counts not written at shutdown",
"hits", c.pendingHits.Load(),
"misses", c.pendingMisses.Load(),
"upstream_fetches", c.pendingUpstreamFetches.Load(),
"upstream_fetch_bytes", c.pendingUpstreamFetchBytes.Load(),
"transforms", c.pendingTransforms.Load(),
"error", err,
)
return err
}
return nil
}
// pendingCountWriteLoop is the body of the goroutine
// StartPendingCountWrites starts. It returns when ctx is cancelled, but
// does not cut off a write under way then. A write that fails leaves the
// counts pending, for the next write.
func (c *Cache) pendingCountWriteLoop(ctx context.Context) {
defer close(c.pendingCountsDone)
for {
select {
case <-ctx.Done():
return
case <-c.pendingCountsAdded:
}
err := c.writePendingCounts(context.WithoutCancel(ctx))
if err != nil {
c.log.Warn("failed to write pending counts", "error", err)
}
}
}
// addPendingCount adds n to pendingCount, one of the pending counts, and
// wakes the goroutine that writes them. pendingCountsAdded has capacity one,
// so a wakeup already waiting is enough.
func (c *Cache) addPendingCount(pendingCount *atomic.Int64, n int64) {
pendingCount.Add(n)
select {
case c.pendingCountsAdded <- struct{}{}:
default:
}
}
// writePendingCounts adds the pending counts to the cache_stats row in one
// UPDATE, then takes what it wrote out of them, leaving any added
// meanwhile. When the UPDATE fails, they all stay pending.
func (c *Cache) writePendingCounts(ctx context.Context) error {
c.pendingCountsWriteMutex.Lock()
defer c.pendingCountsWriteMutex.Unlock()
hits := c.pendingHits.Load()
misses := c.pendingMisses.Load()
upstreamFetches := c.pendingUpstreamFetches.Load()
upstreamFetchBytes := c.pendingUpstreamFetchBytes.Load()
transforms := c.pendingTransforms.Load()
if hits+misses+upstreamFetches+upstreamFetchBytes+transforms == 0 {
return nil
}
_, err := c.db.ExecContext(ctx, `
UPDATE cache_stats
SET hit_count = hit_count + ?,
miss_count = miss_count + ?,
upstream_fetch_count = upstream_fetch_count + ?,
upstream_fetch_bytes = upstream_fetch_bytes + ?,
transform_count = transform_count + ?,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`, hits, misses, upstreamFetches, upstreamFetchBytes, transforms)
if err != nil {
return fmt.Errorf("failed to write pending counts: %w", err)
}
c.pendingHits.Add(-hits)
c.pendingMisses.Add(-misses)
c.pendingUpstreamFetches.Add(-upstreamFetches)
c.pendingUpstreamFetchBytes.Add(-upstreamFetchBytes)
c.pendingTransforms.Add(-transforms)
return nil
}
@@ -0,0 +1,316 @@
package imgcache
import (
"context"
"database/sql"
"errors"
"testing"
"time"
)
// pendingCountsWait is how long a test waits for the database's connection,
// for a count to reach the database, for a write to start waiting for the
// connection, or for StopPendingCountWrites to return.
const pendingCountsWait = 5 * time.Second
// holdDatabase takes the one connection of cache's database, so that every
// other query waits for it, and returns the func that frees it.
func holdDatabase(t *testing.T, cache *Cache) func() {
t.Helper()
ctx, cancel := context.WithTimeout(t.Context(), pendingCountsWait)
defer cancel()
conn, err := cache.db.Conn(ctx)
if err != nil {
t.Fatalf("failed to take the database connection: %v", err)
}
return func() { _ = conn.Close() }
}
// waitForConnectionWaits waits until db.Stats().WaitCount, the number of
// times a caller has waited for a connection, reaches waits.
func waitForConnectionWaits(t *testing.T, db *sql.DB, waits int64) {
t.Helper()
deadline := time.Now().Add(pendingCountsWait)
for db.Stats().WaitCount < waits {
if time.Now().After(deadline) {
t.Fatalf("callers waited for the database connection %d times, want %d",
db.Stats().WaitCount, waits)
}
time.Sleep(10 * time.Millisecond)
}
}
// waitForCounters waits until the cache_stats row holds want.
func waitForCounters(t *testing.T, cache *Cache, want cacheStatsCounters) {
t.Helper()
deadline := time.Now().Add(pendingCountsWait)
for {
got := readCacheStatsCounters(t, cache)
if got == want {
return
}
if time.Now().After(deadline) {
t.Fatalf("counters = %+v, want %+v", got, want)
}
time.Sleep(10 * time.Millisecond)
}
}
// TestService_Get_ReturnsByItsDeadlineWhileTheDatabaseIsBusy holds the
// database's one connection while a request whose fetch is held reaches its
// deadline. The request must still return by its deadline with the
// deadline's error, and its miss must reach the database once the
// connection is free.
func TestService_Get_ReturnsByItsDeadlineWhileTheDatabaseIsBusy(t *testing.T) {
t.Parallel()
svc, fixtures, fetcher := setupHeldFetchService(t)
svc.cache.StartPendingCountWrites()
defer func() { _ = svc.cache.StopPendingCountWrites(t.Context()) }()
// Room for the request to reach its fetch on a busy host
const timeout = 2 * time.Second
ctx, cancel := context.WithTimeout(t.Context(), timeout)
defer cancel()
deadline, _ := ctx.Deadline()
results := startGet(ctx, svc, photoVariant(fixtures, 85, FitCover))
// The request has made its database reads by the time it fetches.
select {
case <-fetcher.started:
case <-time.After(timeout):
t.Fatal("request did not reach its fetch by its deadline")
}
releaseDatabase := holdDatabase(t, svc.cache)
defer releaseDatabase()
select {
case got := <-results:
t.Logf("Get() returned %v after its deadline, error = %v",
time.Since(deadline), got.err)
if !errors.Is(got.err, context.DeadlineExceeded) {
t.Errorf("Get() error = %v, want %v", got.err, context.DeadlineExceeded)
}
case <-time.After(time.Until(deadline) + time.Second):
t.Fatal("request did not return by its deadline while the database was busy")
}
releaseDatabase()
waitForCounters(t, svc.cache, cacheStatsCounters{missCount: 1})
}
// TestIncrementStats_ManyCountsPastTheirDeadlineLeaveOneWriteWaiting holds the
// database's one connection while many misses are counted past their
// deadline. At most one write may then wait for the connection, and every
// miss must reach the database once it is free.
func TestIncrementStats_ManyCountsPastTheirDeadlineLeaveOneWriteWaiting(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
cache.StartPendingCountWrites()
defer func() { _ = cache.StopPendingCountWrites(t.Context()) }()
releaseDatabase := holdDatabase(t, cache)
defer releaseDatabase()
waitsBefore := cache.db.Stats().WaitCount
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
const misses = 20
for range misses {
cache.IncrementStats(ended, false, 0)
}
// Once one write waits for the connection, any others start waiting too.
waitForConnectionWaits(t, cache.db, waitsBefore+1)
time.Sleep(arrivalWait)
waits := cache.db.Stats().WaitCount - waitsBefore
t.Logf("%d writes waited for the database", waits)
if waits > 1 {
t.Errorf("%d writes waited for the database, want at most 1", waits)
}
releaseDatabase()
waitForCounters(t, cache, cacheStatsCounters{missCount: misses})
}
// TestStats_IncludesPendingCounts counts hits and misses whose writes miss
// their deadline, and checks that Stats adds them to the database's counts
// before they are written.
func TestStats_IncludesPendingCounts(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
_, err := cache.db.ExecContext(t.Context(),
`UPDATE cache_stats SET hit_count = 75, miss_count = 25 WHERE id = 1`)
if err != nil {
t.Fatal(err)
}
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
cache.IncrementStats(ended, true, 0)
cache.IncrementStats(ended, true, 0)
cache.IncrementStats(ended, false, 0)
stats, err := cache.Stats(t.Context())
if err != nil {
t.Fatalf("Stats() error = %v", err)
}
if stats.HitCount != 77 || stats.MissCount != 26 {
t.Errorf("HitCount = %d, MissCount = %d, want 77 and 26",
stats.HitCount, stats.MissCount)
}
want := cacheStatsCounters{hitCount: 75, missCount: 25}
if got := readCacheStatsCounters(t, cache); got != want {
t.Errorf("counters in the database = %+v, want %+v", got, want)
}
}
// TestStopPendingCountWrites_WritesPendingCounts holds the database's one
// connection while the goroutine that writes the pending counts waits for it,
// counts more, then stops that goroutine. The stop must leave the write under
// way to finish rather than cut it off, and every count must reach the
// database once the connection is free.
func TestStopPendingCountWrites_WritesPendingCounts(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
cache.StartPendingCountWrites()
defer func() { _ = cache.StopPendingCountWrites(t.Context()) }()
releaseDatabase := holdDatabase(t, cache)
defer releaseDatabase()
waitsBefore := cache.db.Stats().WaitCount
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
cache.IncrementStats(ended, true, 0)
// The goroutine's write waits for the connection.
waitForConnectionWaits(t, cache.db, waitsBefore+1)
// Counted after that write read the pending counts
cache.IncrementStats(ended, false, 0)
cache.IncrementUpstreamFetch(ended, 1024)
cache.IncrementTransformCount(ended)
stopped := make(chan error, 1)
go func() {
stopped <- cache.StopPendingCountWrites(t.Context())
}()
// Had the stop cut off the write under way, its own write would wait for
// the connection too.
time.Sleep(arrivalWait)
waits := cache.db.Stats().WaitCount - waitsBefore
t.Logf("%d writes waited for the database", waits)
if waits > 1 {
t.Errorf("%d writes waited for the database, want 1: "+
"the stop cut off the write under way", waits)
}
releaseDatabase()
select {
case err := <-stopped:
if err != nil {
t.Fatalf("StopPendingCountWrites() error = %v", err)
}
case <-time.After(pendingCountsWait):
t.Fatal("StopPendingCountWrites() did not return once the database was free")
}
want := cacheStatsCounters{1, 1, 1, 1024, 1}
if got := readCacheStatsCounters(t, cache); got != want {
t.Errorf("counters = %+v, want %+v", got, want)
}
}
// TestStopPendingCountWrites_ReturnsWhenItsContextEnds holds the database's
// one connection while the goroutine that writes the pending counts waits for
// it, then stops that goroutine with a context that ends first. The stop must
// return the context's error, and the write under way must still reach the
// database once the connection is free.
func TestStopPendingCountWrites_ReturnsWhenItsContextEnds(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
cache.StartPendingCountWrites()
defer func() { _ = cache.StopPendingCountWrites(t.Context()) }()
releaseDatabase := holdDatabase(t, cache)
defer releaseDatabase()
waitsBefore := cache.db.Stats().WaitCount
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
cache.IncrementStats(ended, true, 0)
// The goroutine's write waits for the connection.
waitForConnectionWaits(t, cache.db, waitsBefore+1)
stopCtx, cancelStop := context.WithCancel(t.Context())
cancelStop()
stopped := make(chan error, 1)
go func() {
stopped <- cache.StopPendingCountWrites(stopCtx)
}()
select {
case err := <-stopped:
if !errors.Is(err, context.Canceled) {
t.Errorf("StopPendingCountWrites() error = %v, want %v",
err, context.Canceled)
}
case <-time.After(pendingCountsWait):
t.Fatal("StopPendingCountWrites() did not return once its context ended")
}
releaseDatabase()
waitForCounters(t, cache, cacheStatsCounters{hitCount: 1})
}
+22 -14
View File
@@ -100,8 +100,8 @@ func NewService(cfg *ServiceConfig) (*Service, error) {
allowHTTP = cfg.FetcherConfig.AllowHTTP
}
// JPEG XL is to become the default output format, so pixad does not
// start without it.
// JPEG XL is the default output format, so pixad does not start
// without it.
err := imageprocessor.CheckJPEGXLSupport()
if err != nil {
return nil, err
@@ -162,7 +162,7 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
// Fall through to re-process
} else {
// Counted also when the request context has ended meanwhile
s.cache.IncrementStats(context.WithoutCancel(ctx), true, 0)
s.cache.IncrementStats(ctx, true, 0)
return &ImageResponse{
Content: reader,
@@ -179,7 +179,7 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
// failed or the request context has ended meanwhile
response, err := s.processOrWait(ctx, req)
s.cache.IncrementStats(context.WithoutCancel(ctx), false, 0)
s.cache.IncrementStats(ctx, false, 0)
if err != nil {
return nil, err
@@ -300,14 +300,8 @@ func (s *Service) processOrWait(
}
}()
processingCtx := context.WithoutCancel(ctx)
if deadline, ok := ctx.Deadline(); ok {
var cancel context.CancelFunc
processingCtx, cancel = context.WithDeadline(processingCtx, deadline)
defer cancel()
}
processingCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
return s.processFromSourceOrFetch(processingCtx, req, cacheKey)
})
@@ -344,6 +338,20 @@ func (s *Service) processOrWait(
}, nil
}
// withoutCancelKeepingDeadline returns ctx without its cancellation but with
// its deadline, if it has one, and the func that releases the returned
// context.
func withoutCancelKeepingDeadline(
ctx context.Context,
) (context.Context, context.CancelFunc) {
deadline, hasDeadline := ctx.Deadline()
if !hasDeadline {
return context.WithoutCancel(ctx), func() {}
}
return context.WithDeadline(context.WithoutCancel(ctx), deadline)
}
// loadCachedSource opens source content from cache, without reading it, and
// returns it with its size; nil if the cached data is unavailable, empty or
// exceeds maxResponseSize.
@@ -453,7 +461,7 @@ func (s *Service) fetchAndProcess(
fetchBytes := int64(len(sourceData))
// Counted also when the request context has ended meanwhile
s.cache.IncrementUpstreamFetch(context.WithoutCancel(ctx), fetchBytes)
s.cache.IncrementUpstreamFetch(ctx, fetchBytes)
if err != nil {
return nil, fmt.Errorf("failed to read upstream response: %w", err)
@@ -535,7 +543,7 @@ func (s *Service) processAndStore(
processDuration := time.Since(processStart)
// Counted also when the request context has ended meanwhile
s.cache.IncrementTransformCount(context.WithoutCancel(ctx))
s.cache.IncrementTransformCount(ctx)
// Read processed content
processedData, err := io.ReadAll(processResult.Content)
+19 -9
View File
@@ -38,8 +38,10 @@ func ValidateDimension(name string, value int) error {
return nil
}
// sizeFormatRegex matches patterns like "800x600.webp", "0x0.jpeg", "orig.png"
var sizeFormatRegex = regexp.MustCompile(`^(\d+)x(\d+)\.(\w+)$|^(orig)\.(\w+)$`)
// sizeFormatRegex matches patterns like "800x600.webp", "0x0.jpeg", "orig.png",
// and a size with no format, such as "800x600" or "orig"
var sizeFormatRegex = regexp.MustCompile(
`^(\d+)x(\d+)(?:\.(\w+))?$|^(orig)(?:\.(\w+))?$`)
// ParsedURL contains the parsed components of an image proxy URL.
type ParsedURL struct {
@@ -56,12 +58,13 @@ type ParsedURL struct {
}
// ParseImagePath parses the path captured by chi's wildcard:
// <host>/<path>/<size>.<format>
// <host>/<path>/<size>.<format>, or <host>/<path>/<size> for JPEG XL
// This is the primary entry point when using chi routing.
// Examples:
// - cdn.example.com/photos/cat.jpg/800x600.webp
// - cdn.example.com/photos/cat.jpg/0x0.jpeg
// - cdn.example.com/photos/cat.jpg/orig.png
// - cdn.example.com/photos/cat.jpg/800x600
func ParseImagePath(path string) (*ParsedURL, error) {
// Strip leading slash if present (chi may include it)
path = strings.TrimPrefix(path, "/")
@@ -72,7 +75,8 @@ func ParseImagePath(path string) (*ParsedURL, error) {
return parseImageComponents(path)
}
// ParseImageURL parses a full URL path like /v1/image/<host>/<path>/<size>.<format>
// ParseImageURL parses a full URL path like /v1/image/<host>/<path>/<size>.<format>,
// or /v1/image/<host>/<path>/<size> for JPEG XL
// Use ParseImagePath instead when working with chi's wildcard capture.
func ParseImageURL(urlPath string) (*ParsedURL, error) {
// Remove the /v1/image/ prefix
@@ -89,7 +93,8 @@ func ParseImageURL(urlPath string) (*ParsedURL, error) {
return parseImageComponents(remainder)
}
// parseImageComponents parses <host>/<path>/<size>.<format> structure.
// parseImageComponents parses <host>/<path>/<size>.<format>, or
// <host>/<path>/<size> for JPEG XL.
func parseImageComponents(remainder string) (*ParsedURL, error) {
// Check for path traversal before any other processing
err := checkPathTraversal(remainder)
@@ -97,7 +102,7 @@ func parseImageComponents(remainder string) (*ParsedURL, error) {
return nil, err
}
// Find the last path segment which contains size.format
// Find the last path segment, which holds "size" or "size.format"
lastSlash := strings.LastIndex(remainder, "/")
if lastSlash == -1 {
return nil, ErrMissingSize
@@ -212,7 +217,7 @@ func checkPathTraversal(path string) error {
return nil
}
// parseSizeFormat parses strings like "800x600.webp" or "orig.png"
// parseSizeFormat parses strings like "800x600.webp", "orig.png" or "800x600"
func parseSizeFormat(s string) (Size, ImageFormat, error) {
matches := sizeFormatRegex.FindStringSubmatch(s)
if matches == nil {
@@ -225,11 +230,11 @@ func parseSizeFormat(s string) (Size, ImageFormat, error) {
)
if matches[4] == "orig" {
// "orig.format" pattern
// "orig" or "orig.format" pattern
size = Size{Width: 0, Height: 0}
formatStr = matches[5]
} else {
// "WxH.format" pattern
// "WxH" or "WxH.format" pattern
width, err := strconv.Atoi(matches[1])
if err != nil {
return Size{}, "", ErrInvalidSize
@@ -254,6 +259,11 @@ func parseSizeFormat(s string) (Size, ImageFormat, error) {
return Size{}, "", err
}
// A URL that names no format is served, and signed, as JPEG XL
if formatStr == "" {
return size, FormatJXL, nil
}
format, err := parseFormat(formatStr)
if err != nil {
return Size{}, "", err
+2 -1
View File
@@ -89,7 +89,8 @@ func (s *Server) SetupRoutes() {
r.Use(s.refuseDuringMaintenance)
// Main image proxy route
// /v1/image/<host>/<path>/<width>x<height>.<format>
// /v1/image/<host>/<path>/<width>x<height>.<format>, or with no
// format /v1/image/<host>/<path>/<width>x<height>
r.Get("/image/*", s.h.HandleImage())
r.Head("/image/*", s.h.HandleImage())
+1 -1
View File
@@ -100,7 +100,7 @@
<option value="png" {{if eq .FormFormat "png"}}selected{{end}}>PNG</option>
<option value="webp" {{if eq .FormFormat "webp"}}selected{{end}}>WebP</option>
<option value="avif" {{if eq .FormFormat "avif"}}selected{{end}}>AVIF</option>
<option value="jxl" {{if eq .FormFormat "jxl"}}selected{{end}}>JPEG XL</option>
<option value="jxl" {{if or (eq .FormFormat "jxl") (eq .FormFormat "")}}selected{{end}}>JPEG XL</option>
<option value="gif" {{if eq .FormFormat "gif"}}selected{{end}}>GIF</option>
</select>
</div>