Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
410c167c70 | ||
|
|
bdde021b45 | ||
|
|
db91ab29e6 |
+6
-3
@@ -32,10 +32,13 @@ RUN script/bootstrap --cgo
|
|||||||
COPY . .
|
COPY . .
|
||||||
|
|
||||||
# Without -v first; on a failure, again with -v for the details, and
|
# Without -v first; on a failure, again with -v for the details, and
|
||||||
# the step fails even if the second run passes.
|
# the step fails even if the second run passes. -parallel 4: by default
|
||||||
RUN go test -count=1 -timeout 90s -race -cover ./... || \
|
# 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 ---"; \
|
{ 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
|
# 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
|
# make BuildKit build them first, so this stage runs only when lint and
|
||||||
|
|||||||
@@ -83,10 +83,12 @@ which answers 200 whenever pixa is running, in maintenance mode too (see
|
|||||||
`maintenance_mode`).
|
`maintenance_mode`).
|
||||||
|
|
||||||
On SIGTERM or SIGINT pixa stops accepting connections, gives the requests in
|
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,
|
progress and the images being processed 5 seconds to finish, then writes to the
|
||||||
or with 1 when images were still being processed after those 5 seconds or
|
database the counts of cache hits, misses, fetches and conversions that requests
|
||||||
another part of pixa failed to stop. A request not finished by then is cut off.
|
could not write by their deadline, and exits: with 0, or with 1 when images were
|
||||||
`docker stop` waits 10 seconds before it kills the container.
|
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
|
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
|
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
|
with the page naming a field that is not valid; 500 when the URL cannot be
|
||||||
made.
|
made.
|
||||||
- `GET /logout` — end the login session. Needs: nothing. Answers: 303 to `/`.
|
- `GET /logout` — end the login session. Needs: nothing. Answers: 303 to `/`.
|
||||||
- `GET` or `HEAD` `/v1/image/<host>/<path>/<size>.<format>` — an image, fetched,
|
- `GET` or `HEAD` `/v1/image/<host>/<path>/<size>.<format>`, or
|
||||||
resized and converted (below). Needs: a signature, unless the host is
|
`/v1/image/<host>/<path>/<size>` with no format — an image, fetched, resized
|
||||||
allowlisted (see Source Hosts). Answers: 200; 304 when `If-None-Match` matches
|
and converted (below). Needs: a signature, unless the host is allowlisted (see
|
||||||
the image's `ETag`; 400 for a URL or parameter that is not valid, or for the
|
Source Hosts). Answers: 200; 304 when `If-None-Match` matches the image's
|
||||||
format `auto` an `Accept` header that is not valid; 406 for the format `auto`
|
`ETag`; 400 for a URL or parameter that is not valid, or for the format `auto`
|
||||||
when `Accept` allows none of the formats it chooses from; 401 for a missing or
|
an `Accept` header that is not valid; 406 for the format `auto` when `Accept`
|
||||||
wrong signature, a missing `exp` or an `exp` in the past; 403 when the
|
allows none of the formats it chooses from; 401 for a missing or wrong
|
||||||
request's `Referer` names a host in `referer_blocklist`, checked before the
|
signature, a missing `exp` or an `exp` in the past; 403 when the request's
|
||||||
signature, the cache and the upstream fetch; 403 when the upstream host, or a
|
`Referer` names a host in `referer_blocklist`, checked before the signature,
|
||||||
host it redirects to, is `localhost`, ends in `.localhost` or `.local`, or has
|
the cache and the upstream fetch; 403 when the upstream host, or a host it
|
||||||
an address in a blocked network (see `blocked_networks`); 502 when the
|
redirects to, is `localhost`, ends in `.localhost` or `.local`, or has an
|
||||||
upstream answered with an error status, and for 5 minutes after that for the
|
address in a blocked network (see `blocked_networks`); 502 when the upstream
|
||||||
same source URL; 503 when pixa is busy or in maintenance mode; 500 for any
|
answered with an error status, and for 5 minutes after that for the same
|
||||||
other failure.
|
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
|
- `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
|
(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
|
`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
|
MiB is refused with 413. The image routes answer the errors listed for them with
|
||||||
JSON holding `error`, `status` and `timestamp`.
|
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>.<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
|
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.
|
once, is refused with 400.
|
||||||
|
|
||||||
- `<format>`: one of `orig` (or `original`), `jpeg` (or `jpg`), `png`, `webp`,
|
- `<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`)
|
- `<size>`: `orig` or `<width>x<height>` (e.g. `800x600`)
|
||||||
- `sig` and `exp`: the signature and its expiry, needed unless the host is
|
- `sig` and `exp`: the signature and its expiry, needed unless the host is
|
||||||
allowlisted (see Signature Specification)
|
allowlisted (see Signature Specification)
|
||||||
@@ -321,13 +327,15 @@ nor change what it asks for.
|
|||||||
lasts 30 days, or until `/logout`.
|
lasts 30 days, or until `/logout`.
|
||||||
2. On the generator page, give the source image's URL, the width and height, the
|
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
|
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
|
(`POST /generate`). The format is JPEG XL unless another is chosen; a form
|
||||||
size; if only one of them is empty or `0`, that side is scaled to keep the
|
sent with an empty format, or none, also makes a JPEG XL URL. Width and
|
||||||
image's proportions.
|
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
|
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
|
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
|
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,
|
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
|
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
|
- `width` — requested width in pixels, `0` for original
|
||||||
- `height` — requested height in pixels, `0` for original
|
- `height` — requested height in pixels, `0` for original
|
||||||
- `format` — output format, one of those listed under Routes, with `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
|
signed as `orig`, `jpg` as `jpeg`, and no format as `jxl`, so a URL with no
|
||||||
format chosen for the request
|
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
|
- `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
|
signature expires; a request whose `exp` is not a whole number, an empty
|
||||||
`exp=` included, is refused with 400
|
`exp=` included, is refused with 400
|
||||||
|
|||||||
@@ -30,6 +30,43 @@ P2: security: per-IP rate limiting on the image routes
|
|||||||
|
|
||||||
# Completed Steps
|
# 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
|
- 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`
|
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
|
or the container's, is detected as `image/jxl`, which the upstream fetch
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ import (
|
|||||||
// Default values for optional fields.
|
// Default values for optional fields.
|
||||||
const (
|
const (
|
||||||
DefaultQuality = 85
|
DefaultQuality = 85
|
||||||
DefaultFormat = imgcache.FormatOriginal
|
DefaultFormat = imgcache.FormatJXL
|
||||||
DefaultFitMode = imgcache.FitCover
|
DefaultFitMode = imgcache.FitCover
|
||||||
|
|
||||||
// HKDF salt for URL encryption key derivation
|
// HKDF salt for URL encryption key derivation
|
||||||
@@ -37,7 +37,7 @@ type Payload struct {
|
|||||||
SourceQuery string `cbor:"q,omitempty"` // optional
|
SourceQuery string `cbor:"q,omitempty"` // optional
|
||||||
Width int `cbor:"w,omitempty"` // 0 = original
|
Width int `cbor:"w,omitempty"` // 0 = original
|
||||||
Height int `cbor:"ht,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
|
Quality int `cbor:"ql,omitempty"` // default: 85
|
||||||
FitMode imgcache.FitMode `cbor:"fm,omitempty"` // default: cover
|
FitMode imgcache.FitMode `cbor:"fm,omitempty"` // default: cover
|
||||||
ExpiresAt int64 `cbor:"e,omitempty"` // 0 = never expires
|
ExpiresAt int64 `cbor:"e,omitempty"` // 0 = never expires
|
||||||
|
|||||||
@@ -369,10 +369,15 @@ func (s *Handlers) buildGeneratedURL(r *http.Request, token, format string) stri
|
|||||||
scheme = "http"
|
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
|
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
|
return scheme + "://" + r.Host + "/v1/e/" + url.PathEscape(token) + "/img." + ext
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ package handlers
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"time"
|
"time"
|
||||||
@@ -71,16 +72,21 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
lc.Append(fx.Hook{
|
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 {
|
OnStart: func(_ context.Context) error {
|
||||||
return s.initImageService()
|
return s.initImageService()
|
||||||
},
|
},
|
||||||
|
// The pending counts are written here, before the database's stop
|
||||||
|
// hook closes it.
|
||||||
OnStop: func(ctx context.Context) error {
|
OnStop: func(ctx context.Context) error {
|
||||||
if s.imgCache == nil {
|
if s.imgCache == nil {
|
||||||
return 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.
|
// write-pressure passes. No-op when the disk cache is disabled.
|
||||||
cache.StartEviction(imgcache.DefaultEvictionInterval)
|
cache.StartEviction(imgcache.DefaultEvictionInterval)
|
||||||
|
|
||||||
|
// Writes the counts requests could not write by their deadline
|
||||||
|
cache.StartPendingCountWrites()
|
||||||
|
|
||||||
// Create the fetcher config
|
// Create the fetcher config
|
||||||
fetcherCfg := httpfetcher.DefaultConfig()
|
fetcherCfg := httpfetcher.DefaultConfig()
|
||||||
fetcherCfg.AllowHTTP = s.config.AllowHTTP
|
fetcherCfg.AllowHTTP = s.config.AllowHTTP
|
||||||
|
|||||||
@@ -18,7 +18,8 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// HandleImage handles the main image proxy route:
|
// 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 {
|
func (s *Handlers) HandleImage() http.HandlerFunc {
|
||||||
return func(w http.ResponseWriter, r *http.Request) {
|
return func(w http.ResponseWriter, r *http.Request) {
|
||||||
if s.refuseBlockedReferer(w, r) {
|
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))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
outputFormat := req.Format
|
||||||
if outputFormat == FormatOriginal || outputFormat == "" {
|
if outputFormat == FormatOriginal {
|
||||||
outputFormat = p.formatFromString(inputFormat)
|
outputFormat = p.formatFromString(inputFormat)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -504,36 +504,21 @@ func (p *ImageProcessor) encode(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
var params vips.ExportParams
|
|
||||||
|
|
||||||
switch format {
|
switch format {
|
||||||
case FormatJPEG:
|
case FormatJPEG:
|
||||||
params = vips.ExportParams{
|
return exportJPEG(img, quality)
|
||||||
Format: vips.ImageTypeJPEG,
|
|
||||||
Quality: quality,
|
|
||||||
}
|
|
||||||
|
|
||||||
case FormatPNG:
|
case FormatPNG:
|
||||||
params = vips.ExportParams{
|
return exportPNG(img)
|
||||||
Format: vips.ImageTypePNG,
|
|
||||||
}
|
|
||||||
|
|
||||||
case FormatGIF:
|
case FormatGIF:
|
||||||
params = vips.ExportParams{
|
return exportGIF(img)
|
||||||
Format: vips.ImageTypeGIF,
|
|
||||||
}
|
|
||||||
|
|
||||||
case FormatWebP:
|
case FormatWebP:
|
||||||
params = vips.ExportParams{
|
return exportWebP(img, quality)
|
||||||
Format: vips.ImageTypeWEBP,
|
|
||||||
Quality: quality,
|
|
||||||
}
|
|
||||||
|
|
||||||
case FormatAVIF:
|
case FormatAVIF:
|
||||||
params = vips.ExportParams{
|
return exportAVIF(img, quality)
|
||||||
Format: vips.ImageTypeAVIF,
|
|
||||||
Quality: quality,
|
|
||||||
}
|
|
||||||
|
|
||||||
case FormatJXL:
|
case FormatJXL:
|
||||||
return exportJXL(img, quality)
|
return exportJXL(img, quality)
|
||||||
@@ -544,17 +529,90 @@ func (p *ImageProcessor) encode(
|
|||||||
default:
|
default:
|
||||||
return nil, fmt.Errorf("%w: %s", ErrUnsupportedOutputFormat, format)
|
return nil, fmt.Errorf("%w: %s", ErrUnsupportedOutputFormat, format)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Drop EXIF, XMP, IPTC and the ICC profile. govips ignores this for
|
// govips sends libvips Go's zero value for some settings an export leaves
|
||||||
// GIF, which carries none of them.
|
// out, such as no compression at all for PNG, so each export below sets
|
||||||
params.StripMetadata = true
|
// every setting whose zero value is not what pixa wants. Stripping metadata
|
||||||
|
// drops EXIF, XMP, IPTC and the ICC profile.
|
||||||
|
|
||||||
output, _, err := img.Export(¶ms)
|
// exportJPEG encodes img as JPEG at quality, without metadata. The settings
|
||||||
if err != nil {
|
// it leaves out are at libvips' defaults.
|
||||||
return nil, err
|
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
|
// jxlResolution is the resolution every JPEG XL image is saved with, in
|
||||||
|
|||||||
@@ -11,6 +11,8 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
lru "github.com/hashicorp/golang-lru/v2"
|
lru "github.com/hashicorp/golang-lru/v2"
|
||||||
@@ -79,6 +81,26 @@ type Cache struct {
|
|||||||
evictionDone chan struct{}
|
evictionDone chan struct{}
|
||||||
evictionCancel context.CancelFunc
|
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
|
// metaCache holds the content types of the variants most recently
|
||||||
// stored or served, so a hit does not read the variant's .meta file.
|
// 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.
|
// It never stands in for the variant file, which is always opened.
|
||||||
@@ -139,6 +161,9 @@ func newCache(
|
|||||||
metaCache: metaCache,
|
metaCache: metaCache,
|
||||||
contentLocks: newContentLock(),
|
contentLocks: newContentLock(),
|
||||||
|
|
||||||
|
pendingCountsAdded: make(chan struct{}, 1),
|
||||||
|
pendingCountsDone: make(chan struct{}),
|
||||||
|
|
||||||
reconciliationPageSize: defaultReconciliationPageSize,
|
reconciliationPageSize: defaultReconciliationPageSize,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -474,12 +499,21 @@ func (c *Cache) CleanExpired(ctx context.Context) error {
|
|||||||
func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
|
func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
|
||||||
var stats CacheStats
|
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
|
// Fetch hit/miss counts from the stats table
|
||||||
err := c.db.QueryRowContext(ctx, `
|
err := c.db.QueryRowContext(ctx, `
|
||||||
SELECT hit_count, miss_count
|
SELECT hit_count, miss_count
|
||||||
FROM cache_stats WHERE id = 1
|
FROM cache_stats WHERE id = 1
|
||||||
`).Scan(&stats.HitCount, &stats.MissCount)
|
`).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) {
|
if err != nil && !errors.Is(err, sql.ErrNoRows) {
|
||||||
return nil, fmt.Errorf("failed to get cache stats: %w", err)
|
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
|
// 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) {
|
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
|
||||||
|
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
|
pendingCount := &c.pendingMisses
|
||||||
|
|
||||||
if hit {
|
if hit {
|
||||||
_, err = c.db.ExecContext(ctx, `
|
pendingCount = &c.pendingHits
|
||||||
|
|
||||||
|
_, err = c.db.ExecContext(countCtx, `
|
||||||
UPDATE cache_stats
|
UPDATE cache_stats
|
||||||
SET hit_count = hit_count + 1,
|
SET hit_count = hit_count + 1,
|
||||||
last_updated_at = CURRENT_TIMESTAMP
|
last_updated_at = CURRENT_TIMESTAMP
|
||||||
WHERE id = 1
|
WHERE id = 1
|
||||||
`)
|
`)
|
||||||
} else {
|
} else {
|
||||||
_, err = c.db.ExecContext(ctx, `
|
_, err = c.db.ExecContext(countCtx, `
|
||||||
UPDATE cache_stats
|
UPDATE cache_stats
|
||||||
SET miss_count = miss_count + 1,
|
SET miss_count = miss_count + 1,
|
||||||
last_updated_at = CURRENT_TIMESTAMP
|
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)
|
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
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
_, err := c.db.ExecContext(ctx, `
|
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
_, err := c.db.ExecContext(countCtx, `
|
||||||
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 + ?,
|
||||||
last_updated_at = CURRENT_TIMESTAMP
|
last_updated_at = CURRENT_TIMESTAMP
|
||||||
WHERE id = 1
|
WHERE id = 1
|
||||||
`, fetchBytes)
|
`, 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",
|
c.log.Warn("failed to count upstream fetch",
|
||||||
"fetch_bytes", fetchBytes, "error", err)
|
"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.
|
// IncrementTransformCount counts one image transcoded by the image processor.
|
||||||
func (c *Cache) IncrementTransformCount(ctx context.Context) {
|
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
|
UPDATE cache_stats
|
||||||
SET transform_count = transform_count + 1,
|
SET transform_count = transform_count + 1,
|
||||||
last_updated_at = CURRENT_TIMESTAMP
|
last_updated_at = CURRENT_TIMESTAMP
|
||||||
WHERE id = 1
|
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)
|
c.log.Warn("failed to count transform", "error", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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})
|
||||||
|
}
|
||||||
@@ -100,8 +100,8 @@ func NewService(cfg *ServiceConfig) (*Service, error) {
|
|||||||
allowHTTP = cfg.FetcherConfig.AllowHTTP
|
allowHTTP = cfg.FetcherConfig.AllowHTTP
|
||||||
}
|
}
|
||||||
|
|
||||||
// JPEG XL is to become the default output format, so pixad does not
|
// JPEG XL is the default output format, so pixad does not start
|
||||||
// start without it.
|
// without it.
|
||||||
err := imageprocessor.CheckJPEGXLSupport()
|
err := imageprocessor.CheckJPEGXLSupport()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -162,7 +162,7 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
|
|||||||
// Fall through to re-process
|
// Fall through to re-process
|
||||||
} else {
|
} else {
|
||||||
// Counted also when the request context has ended meanwhile
|
// 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{
|
return &ImageResponse{
|
||||||
Content: reader,
|
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
|
// failed or the request context has ended meanwhile
|
||||||
response, err := s.processOrWait(ctx, req)
|
response, err := s.processOrWait(ctx, req)
|
||||||
|
|
||||||
s.cache.IncrementStats(context.WithoutCancel(ctx), false, 0)
|
s.cache.IncrementStats(ctx, false, 0)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -300,14 +300,8 @@ func (s *Service) processOrWait(
|
|||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
processingCtx := context.WithoutCancel(ctx)
|
processingCtx, cancel := withoutCancelKeepingDeadline(ctx)
|
||||||
|
defer cancel()
|
||||||
if deadline, ok := ctx.Deadline(); ok {
|
|
||||||
var cancel context.CancelFunc
|
|
||||||
|
|
||||||
processingCtx, cancel = context.WithDeadline(processingCtx, deadline)
|
|
||||||
defer cancel()
|
|
||||||
}
|
|
||||||
|
|
||||||
return s.processFromSourceOrFetch(processingCtx, req, cacheKey)
|
return s.processFromSourceOrFetch(processingCtx, req, cacheKey)
|
||||||
})
|
})
|
||||||
@@ -344,6 +338,20 @@ func (s *Service) processOrWait(
|
|||||||
}, nil
|
}, 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
|
// 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.
|
||||||
@@ -453,7 +461,7 @@ func (s *Service) fetchAndProcess(
|
|||||||
fetchBytes := int64(len(sourceData))
|
fetchBytes := int64(len(sourceData))
|
||||||
|
|
||||||
// Counted also when the request context has ended meanwhile
|
// Counted also when the request context has ended meanwhile
|
||||||
s.cache.IncrementUpstreamFetch(context.WithoutCancel(ctx), fetchBytes)
|
s.cache.IncrementUpstreamFetch(ctx, fetchBytes)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("failed to read upstream response: %w", err)
|
return nil, fmt.Errorf("failed to read upstream response: %w", err)
|
||||||
@@ -535,7 +543,7 @@ func (s *Service) processAndStore(
|
|||||||
processDuration := time.Since(processStart)
|
processDuration := time.Since(processStart)
|
||||||
|
|
||||||
// Counted also when the request context has ended meanwhile
|
// Counted also when the request context has ended meanwhile
|
||||||
s.cache.IncrementTransformCount(context.WithoutCancel(ctx))
|
s.cache.IncrementTransformCount(ctx)
|
||||||
|
|
||||||
// Read processed content
|
// Read processed content
|
||||||
processedData, err := io.ReadAll(processResult.Content)
|
processedData, err := io.ReadAll(processResult.Content)
|
||||||
|
|||||||
@@ -38,8 +38,10 @@ func ValidateDimension(name string, value int) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// sizeFormatRegex matches patterns like "800x600.webp", "0x0.jpeg", "orig.png"
|
// sizeFormatRegex matches patterns like "800x600.webp", "0x0.jpeg", "orig.png",
|
||||||
var sizeFormatRegex = regexp.MustCompile(`^(\d+)x(\d+)\.(\w+)$|^(orig)\.(\w+)$`)
|
// 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.
|
// ParsedURL contains the parsed components of an image proxy URL.
|
||||||
type ParsedURL struct {
|
type ParsedURL struct {
|
||||||
@@ -56,12 +58,13 @@ type ParsedURL struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ParseImagePath parses the path captured by chi's wildcard:
|
// 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.
|
// This is the primary entry point when using chi routing.
|
||||||
// Examples:
|
// Examples:
|
||||||
// - cdn.example.com/photos/cat.jpg/800x600.webp
|
// - cdn.example.com/photos/cat.jpg/800x600.webp
|
||||||
// - cdn.example.com/photos/cat.jpg/0x0.jpeg
|
// - cdn.example.com/photos/cat.jpg/0x0.jpeg
|
||||||
// - cdn.example.com/photos/cat.jpg/orig.png
|
// - cdn.example.com/photos/cat.jpg/orig.png
|
||||||
|
// - cdn.example.com/photos/cat.jpg/800x600
|
||||||
func ParseImagePath(path string) (*ParsedURL, error) {
|
func ParseImagePath(path string) (*ParsedURL, error) {
|
||||||
// Strip leading slash if present (chi may include it)
|
// Strip leading slash if present (chi may include it)
|
||||||
path = strings.TrimPrefix(path, "/")
|
path = strings.TrimPrefix(path, "/")
|
||||||
@@ -72,7 +75,8 @@ func ParseImagePath(path string) (*ParsedURL, error) {
|
|||||||
return parseImageComponents(path)
|
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.
|
// Use ParseImagePath instead when working with chi's wildcard capture.
|
||||||
func ParseImageURL(urlPath string) (*ParsedURL, error) {
|
func ParseImageURL(urlPath string) (*ParsedURL, error) {
|
||||||
// Remove the /v1/image/ prefix
|
// Remove the /v1/image/ prefix
|
||||||
@@ -89,7 +93,8 @@ func ParseImageURL(urlPath string) (*ParsedURL, error) {
|
|||||||
return parseImageComponents(remainder)
|
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) {
|
func parseImageComponents(remainder string) (*ParsedURL, error) {
|
||||||
// Check for path traversal before any other processing
|
// Check for path traversal before any other processing
|
||||||
err := checkPathTraversal(remainder)
|
err := checkPathTraversal(remainder)
|
||||||
@@ -97,7 +102,7 @@ func parseImageComponents(remainder string) (*ParsedURL, error) {
|
|||||||
return nil, err
|
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, "/")
|
lastSlash := strings.LastIndex(remainder, "/")
|
||||||
if lastSlash == -1 {
|
if lastSlash == -1 {
|
||||||
return nil, ErrMissingSize
|
return nil, ErrMissingSize
|
||||||
@@ -212,7 +217,7 @@ func checkPathTraversal(path string) error {
|
|||||||
return nil
|
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) {
|
func parseSizeFormat(s string) (Size, ImageFormat, error) {
|
||||||
matches := sizeFormatRegex.FindStringSubmatch(s)
|
matches := sizeFormatRegex.FindStringSubmatch(s)
|
||||||
if matches == nil {
|
if matches == nil {
|
||||||
@@ -225,11 +230,11 @@ func parseSizeFormat(s string) (Size, ImageFormat, error) {
|
|||||||
)
|
)
|
||||||
|
|
||||||
if matches[4] == "orig" {
|
if matches[4] == "orig" {
|
||||||
// "orig.format" pattern
|
// "orig" or "orig.format" pattern
|
||||||
size = Size{Width: 0, Height: 0}
|
size = Size{Width: 0, Height: 0}
|
||||||
formatStr = matches[5]
|
formatStr = matches[5]
|
||||||
} else {
|
} else {
|
||||||
// "WxH.format" pattern
|
// "WxH" or "WxH.format" pattern
|
||||||
width, err := strconv.Atoi(matches[1])
|
width, err := strconv.Atoi(matches[1])
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return Size{}, "", ErrInvalidSize
|
return Size{}, "", ErrInvalidSize
|
||||||
@@ -254,6 +259,11 @@ func parseSizeFormat(s string) (Size, ImageFormat, error) {
|
|||||||
return Size{}, "", err
|
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)
|
format, err := parseFormat(formatStr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return Size{}, "", err
|
return Size{}, "", err
|
||||||
|
|||||||
@@ -89,7 +89,8 @@ func (s *Server) SetupRoutes() {
|
|||||||
r.Use(s.refuseDuringMaintenance)
|
r.Use(s.refuseDuringMaintenance)
|
||||||
|
|
||||||
// Main image proxy route
|
// 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.Get("/image/*", s.h.HandleImage())
|
||||||
r.Head("/image/*", s.h.HandleImage())
|
r.Head("/image/*", s.h.HandleImage())
|
||||||
|
|
||||||
|
|||||||
@@ -100,7 +100,7 @@
|
|||||||
<option value="png" {{if eq .FormFormat "png"}}selected{{end}}>PNG</option>
|
<option value="png" {{if eq .FormFormat "png"}}selected{{end}}>PNG</option>
|
||||||
<option value="webp" {{if eq .FormFormat "webp"}}selected{{end}}>WebP</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="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>
|
<option value="gif" {{if eq .FormFormat "gif"}}selected{{end}}>GIF</option>
|
||||||
</select>
|
</select>
|
||||||
</div>
|
</div>
|
||||||
|
|||||||
Reference in New Issue
Block a user