5 Commits
Author SHA1 Message Date
clawbot add0e2ba6a Keep newFromSmartConfig within the length limit (closes #64)
check / check (push) Successful in 3m41s
Rebasing onto the four settings from
#142 put newFromSmartConfig at 82
lines, over the linter's 80-line limit. validateKnownKeys now returns
early for a nil config (no config file) itself, as lookupValue already
does, so the caller drops its own nil check. Behavior is unchanged.

Model: opus-5-5
2026-09-29 07:32:56 +00:00
clawbot d800aa62a7 Read a cached source only once a processing slot is taken (closes #64)
A request whose source was in the disk cache read the whole file into
memory, then waited for a processing slot, so a burst of new sizes for
one large cached image held one copy per waiting request, with no
ceiling. The service now opens the cached file and hands it to the image
processor, which reads it only after taking its slot. The file's size,
now returned by GetSourceContent, still sends an empty or oversized
cached source to upstream instead. A cached file that fails while being
read now fails the request instead of being fetched again.

Model: opus-5-5
2026-09-29 07:32:56 +00:00
clawbot 9e1f5d4cea Test that a cached source is read only with a processing slot (closes #64)
Failing test: with the only processing slot held, a request for a new
width of a cached image waits for the slot while the cached file is
rewritten; the answer must come from the rewritten file, so the request
read none of the source before it had a slot. Also a test that a fetch
whose request context ends while it waits for a connection shared by all
hosts gives its host's slot back; that one passes already.

Model: opus-5-5
2026-09-29 07:32:45 +00:00
clawbot fffd97d3d8 Bound concurrent image processing and upstream fetches (closes #64)
max_concurrent_processing (default: the number of CPUs Go uses) bounds
the images processed at once, and upstream_connections (default 64) the
fetches from all upstream hosts together, beside the per-host limit. A
request that finds either full waits up to 10 seconds, then gets 503
"server busy, try again later". The processor holds its slot from before
it reads the input until it returns, and takes a free slot even after the
request context has ended; a fetch holds its connection until the
response body is closed, after its image is processed. libvips now starts
with one worker thread per image and no operation cache. Both settings
have PIXA_ variables and are in README.md and config.example.yml.

Model: opus-5-5
2026-09-29 07:32:45 +00:00
clawbot 6cb280919f Test the processing and upstream connection limits (closes #64)
Failing tests for two limits that do not exist yet. Config:
max_concurrent_processing and upstream_connections, their defaults, and
valid and invalid values from the file and the environment. Image
processor: never more images at once than its limit, waiting and then
failing with ErrTooManyImages when no slot frees, and freeing its slot on
every error. Fetcher: connections to all hosts counted together, apart
from the per-host limit, and freed on errors. Both image routes answer
503 when either wait gives up. TestEnvironmentSetsEveryKey sets the two
new variables, as it compares the whole config. The tests do not compile
until the limits exist.

Model: opus-5-5
2026-09-29 07:32:32 +00:00
14 changed files with 59 additions and 745 deletions
+14 -21
View File
@@ -40,8 +40,9 @@ What the [upaas](https://git.eeqj.de/sneak/upaas) app for pixa needs:
- **Port:** pixa listens on container port `8080`. - **Port:** pixa listens on container port `8080`.
- **Volume:** container path `/var/lib/pixa`, where pixa keeps its - **Volume:** container path `/var/lib/pixa`, where pixa keeps its
database and cache. Creating the host directory when it is missing is database and cache. upaas bind-mounts the host path it is given and
upaas's job, tracked in https://git.eeqj.de/sneak/upaas/issues/235. does not create it, so the host directory must exist before the first
deploy.
- **Environment variables:** - **Environment variables:**
- `PIXA_SIGNING_KEY` (required): secret for signed and encrypted URLs - `PIXA_SIGNING_KEY` (required): secret for signed and encrypted URLs
and login, 32+ characters, for example from and login, 32+ characters, for example from
@@ -57,6 +58,10 @@ What the [upaas](https://git.eeqj.de/sneak/upaas) app for pixa needs:
`healthy`. The probe uses the port from `PORT` (default `8080`), so a `healthy`. The probe uses the port from `PORT` (default `8080`), so a
port changed only in a mounted config file is not seen by it: change port changed only in a mounted config file is not seen by it: change
the port with `PORT`. the port with `PORT`.
- **First run:** create the host directory, owned by root or by uid
`65532` and gid `65532`. The server runs as the container's `pixad`
user, which has that uid and gid, and the container gives the
directory to `pixad` when it starts.
## Rationale ## Rationale
@@ -81,10 +86,7 @@ prevent abuse, and allowlisted source hosts for open access.
Multiple source paths may reference the same content blob; the Multiple source paths may reference the same content blob; the
database tracks references rather than using filesystem refcounting. database tracks references rather than using filesystem refcounting.
Toward a target of 1-5k r/s, pixa keeps in memory the content types of In-process caching of request-to-output mappings targets 1-5k r/s.
the 10,000 transformed images most recently cached or served, so a
cache hit on one of them reads only the image file from disk and not
the metadata file stored beside it.
### Routes ### Routes
@@ -243,7 +245,7 @@ variables set by the file's `env:` section are checked the same way.
| `PIXA_METRICS_PASSWORD` | `metrics.password` | Password for `/metrics`; set together with the username | | `PIXA_METRICS_PASSWORD` | `metrics.password` | Password for `/metrics`; set together with the username |
| `PIXA_SENTRY_DSN` | `sentry_dsn` | Sentry DSN for error reporting; empty disables it | | `PIXA_SENTRY_DSN` | `sentry_dsn` | Sentry DSN for error reporting; empty disables it |
| `PIXA_DEBUG` | `debug` | Debug logging and plain-HTTP local development; default `false` | | `PIXA_DEBUG` | `debug` | Debug logging and plain-HTTP local development; default `false` |
| `PIXA_MAINTENANCE_MODE` | `maintenance_mode` | Answer image requests with 503; the health check stays 200; default `false` | | `PIXA_MAINTENANCE_MODE` | `maintenance_mode` | Maintenance flag reported by the health check; default `false` |
Key settings in more detail: Key settings in more detail:
@@ -286,9 +288,8 @@ Key settings in more detail:
bytes; default `52428800` (50 MiB). It also limits the image data pixa bytes; default `52428800` (50 MiB). It also limits the image data pixa
decodes decodes
- `downstream_timeout` — time allowed for answering one client request, as a - `downstream_timeout` — time allowed for answering one client request, as a
duration; default `60s`. The upstream fetch counts toward it, and so do the duration; default `60s`. The upstream fetch counts toward it, so keep it
waits for an upstream connection and for a processing slot (up to 10 seconds longer than `upstream_fetch_timeout`
each), so keep it longer than `upstream_fetch_timeout` plus 20 seconds
- `signing_key` — HMAC secret for URL signatures - `signing_key` — HMAC secret for URL signatures
- `cache_max_bytes` — disk cache size limit in bytes; `0` disables the - `cache_max_bytes` — disk cache size limit in bytes; `0` disables the
disk cache entirely; omitted defaults to 75% of the free space on disk cache entirely; omitted defaults to 75% of the free space on
@@ -297,20 +298,12 @@ Key settings in more detail:
hosts together, on top of `upstream_connections_per_host`; default `64`. A hosts together, on top of `upstream_connections_per_host`; default `64`. A
fetch holds its connection until its image has been processed. A fetch that fetch holds its connection until its image has been processed. A fetch that
finds all of them in use waits up to 10 seconds for one to free up; if none finds all of them in use waits up to 10 seconds for one to free up; if none
does, and `downstream_timeout` has not ended first, the request is answered does, the request is answered 503 with the error
503 with the error `server busy, try again later` `server busy, try again later`
- `max_concurrent_processing` — the most images decoded and encoded at once; - `max_concurrent_processing` — the most images decoded and encoded at once;
default the number of CPUs pixa can use (`GOMAXPROCS`), which follows a default the number of CPUs pixa can use (`GOMAXPROCS`), which follows a
container's CPU limit. A request that finds all of them in use waits up to 10 container's CPU limit. A request that finds all of them in use waits up to 10
seconds for one to free up; if none does, and `downstream_timeout` has not seconds for one to free up; if none does, it is answered 503 the same way
ended first, it is answered 503 the same way
- `maintenance_mode` — while `true`, the image routes (`/v1/image/` and
`/v1/e/`) answer every request with 503, a `Retry-After` header and a JSON
error body. The health check (`/.well-known/healthcheck.json`) still answers
200 and reports `"maintenance_mode": true`. It stays 200 because the image's
Docker `HEALTHCHECK` requests it: a 503 there would make the container
unhealthy, and upaas marks a deploy failed when its container is unhealthy.
The login and URL generator pages and `/metrics` keep working
See `config.example.yml` for all options with defaults. See `config.example.yml` for all options with defaults.
-26
View File
@@ -29,32 +29,6 @@ P2: security: referer blacklist
# Completed Steps # Completed Steps
- 2026-09-29 the container makes `/var/lib/pixa` usable by itself (closes
#159): `deploy/docker-entrypoint.sh` creates the directory if it is missing,
gives the directory and everything in it to `pixad` when the directory or one
of its top-level entries belongs to another user or group, sets its mode to
`750`, then runs the server as `pixad`; data left by an earlier run under
another uid is taken over this way; "Running under upaas" in `README.md` no
longer tells the operator to create or chown the host directory.
- 2026-09-29 variant content types kept in memory (closes #70):
`Cache.metaCache` holds the content types of up to 10,000 variants in an LRU
(`github.com/hashicorp/golang-lru/v2`), filled by `StoreVariant` and by
`GetVariant` after it reads a `.meta` file, where a type `StoreVariant` added
meanwhile is kept over the one read, and never with the
`application/octet-stream` served for a variant without one; for a variant it
holds, `GetVariant` skips the `.meta` read, still opening the variant file and
taking the size from it; eviction removes the entry before deleting the files,
and `GetVariant` removes it when the file will not open; the cap is a
constant, not a setting; the unused `variantMeta` type is gone; `README.md`
describes it.
- 2026-09-29 maintenance mode refuses image requests (closes #71): while
`maintenance_mode` is on, `/v1/image/` and `/v1/e/` answer 503 with a
`Retry-After` header and the JSON error body, from one middleware in
`internal/server/routes.go`; the health check stays 200 and reports
`maintenance_mode`, as the image's Docker `HEALTHCHECK` requests it and upaas
marks a deploy failed when its container is unhealthy; the login and URL
generator pages and `/metrics` keep working; documented in `README.md` and
`config.example.yml`.
- 2026-09-29 bound concurrent image processing and upstream fetches (closes - 2026-09-29 bound concurrent image processing and upstream fetches (closes
#64): `max_concurrent_processing` (default the number of CPUs pixa can use) #64): `max_concurrent_processing` (default the number of CPUs pixa can use)
limits the images decoded and encoded at once, and `upstream_connections` limits the images decoded and encoded at once, and `upstream_connections`
+4 -13
View File
@@ -16,12 +16,6 @@
# Server settings # Server settings
port: 8080 port: 8080
debug: false debug: false
# While true, the image routes (/v1/image/ and /v1/e/) answer every request
# with 503 and a Retry-After header. The health check keeps answering 200 and
# reports maintenance_mode as true. It stays 200 because the image's Docker
# HEALTHCHECK requests it: a 503 there would make the container unhealthy, and
# upaas marks a deploy failed when its container is unhealthy.
maintenance_mode: false maintenance_mode: false
# Data directory for SQLite database and cache files # Data directory for SQLite database and cache files
@@ -80,14 +74,13 @@ upstream_connections_per_host: 20
# Maximum concurrent connections to all upstream hosts together, on top of # Maximum concurrent connections to all upstream hosts together, on top of
# the per-host limit (default: 64). A fetch holds its connection until its # the per-host limit (default: 64). A fetch holds its connection until its
# image has been processed. A fetch that finds none free waits up to 10 # image has been processed. A fetch that finds none free waits up to 10
# seconds for one, and if none frees up the request is answered 503, unless # seconds for one, and if none frees up the request is answered 503.
# downstream_timeout has ended first.
upstream_connections: 64 upstream_connections: 64
# Maximum number of images decoded and encoded at once (default: the # Maximum number of images decoded and encoded at once (default: the
# number of CPUs pixa can use, which follows a container's CPU limit). A # number of CPUs pixa can use, which follows a container's CPU limit). A
# request that finds none free waits up to 10 seconds for one, and if none # request that finds none free waits up to 10 seconds for one, and if none
# frees up it is answered 503, unless downstream_timeout has ended first. # frees up it is answered 503.
# max_concurrent_processing: 4 # max_concurrent_processing: 4
# Time allowed for one fetch from an upstream host (default: 30s) # Time allowed for one fetch from an upstream host (default: 30s)
@@ -97,10 +90,8 @@ upstream_fetch_timeout: 30s
# (1 GiB) (default: 52428800, 50 MiB) # (1 GiB) (default: 52428800, 50 MiB)
upstream_max_response_size: 52428800 upstream_max_response_size: 52428800
# Time allowed for answering one client request (default: 60s). The # Time allowed for answering one client request, the upstream fetch
# upstream fetch counts toward it, and so do the waits for an upstream # included, so keep it longer than upstream_fetch_timeout (default: 60s)
# connection and for a processing slot (up to 10 seconds each), so keep it
# longer than upstream_fetch_timeout plus 20 seconds.
downstream_timeout: 60s downstream_timeout: 60s
# The origin a browser lets read pixa's responses, sent as the CORS # The origin a browser lets read pixa's responses, sent as the CORS
+5 -13
View File
@@ -1,22 +1,14 @@
#!/bin/sh #!/bin/sh
# deploy/docker-entrypoint.sh: the Docker image's ENTRYPOINT. It runs as # deploy/docker-entrypoint.sh: the Docker image's ENTRYPOINT. It runs as
# root only to make /var/lib/pixa usable by pixad: a host directory # root only to give /var/lib/pixa to pixad: a host directory
# bind-mounted there keeps its host owner, often root, and data from an # bind-mounted there keeps its host owner, often root, and pixad could
# earlier run may belong to another uid. The server itself always runs # not write to it. The server itself always runs as pixad.
# as pixad.
set -eu set -eu
main() { main() {
mkdir -p /var/lib/pixa if [ "$(stat -c %U /var/lib/pixa)" != pixad ]; then
# Only the directory and its top-level entries are checked, so a chown pixad:pixad /var/lib/pixa
# normal start does not walk the cache. -depth gives each directory
# to pixad after its contents, so a start stopped part way leaves
# something at the top for the next start to find; -h changes a
# symlink itself, never the file it points to.
if [ -n "$(find /var/lib/pixa -maxdepth 1 \( ! -user pixad -o ! -group pixad \))" ]; then
find /var/lib/pixa -depth -exec chown -h pixad:pixad {} +
fi fi
chmod 750 /var/lib/pixa
exec su-exec pixad /usr/local/bin/pixad "$@" exec su-exec pixad /usr/local/bin/pixad "$@"
} }
-1
View File
@@ -14,7 +14,6 @@ require (
github.com/go-chi/httprate v0.16.0 github.com/go-chi/httprate v0.16.0
github.com/gorilla/csrf v1.7.3 github.com/gorilla/csrf v1.7.3
github.com/gorilla/securecookie v1.1.2 github.com/gorilla/securecookie v1.1.2
github.com/hashicorp/golang-lru/v2 v2.0.7
github.com/prometheus/client_golang v1.23.2 github.com/prometheus/client_golang v1.23.2
github.com/slok/go-http-metrics v0.13.0 github.com/slok/go-http-metrics v0.13.0
github.com/spf13/cobra v1.10.2 github.com/spf13/cobra v1.10.2
-2
View File
@@ -228,8 +228,6 @@ github.com/hashicorp/go-version v1.2.1/go.mod h1:fltr4n8CU8Ke44wwGCBoEymUuxUHl09
github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8=
github.com/hashicorp/golang-lru v0.5.4 h1:YDjusn29QI/Das2iO9M0BHnIbxPeyuCHsjMW+lJfyTc= github.com/hashicorp/golang-lru v0.5.4 h1:YDjusn29QI/Das2iO9M0BHnIbxPeyuCHsjMW+lJfyTc=
github.com/hashicorp/golang-lru v0.5.4/go.mod h1:iADmTwqILo4mZ8BN3D2Q6+9jd8WM5uGBxy+E8yxSoD4= github.com/hashicorp/golang-lru v0.5.4/go.mod h1:iADmTwqILo4mZ8BN3D2Q6+9jd8WM5uGBxy+E8yxSoD4=
github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k=
github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM=
github.com/hashicorp/hcl v1.0.1-vault-7 h1:ag5OxFVy3QYTFTJODRzTKVZ6xvdfLLCA1cy/Y6xGI0I= github.com/hashicorp/hcl v1.0.1-vault-7 h1:ag5OxFVy3QYTFTJODRzTKVZ6xvdfLLCA1cy/Y6xGI0I=
github.com/hashicorp/hcl v1.0.1-vault-7/go.mod h1:XYhtn6ijBSAj6n4YqAaf7RBPS4I06AItNorpy+MoQNM= github.com/hashicorp/hcl v1.0.1-vault-7/go.mod h1:XYhtn6ijBSAj6n4YqAaf7RBPS4I06AItNorpy+MoQNM=
github.com/hashicorp/logutils v1.0.0/go.mod h1:QIAnNjmIWmVIIkWDTG1z5v++HQmx9WQRO+LraFDTW64= github.com/hashicorp/logutils v1.0.0/go.mod h1:QIAnNjmIWmVIIkWDTG1z5v++HQmx9WQRO+LraFDTW64=
+12 -61
View File
@@ -14,7 +14,6 @@ import (
"sync" "sync"
"time" "time"
lru "github.com/hashicorp/golang-lru/v2"
"sneak.berlin/go/pixa/internal/httpfetcher" "sneak.berlin/go/pixa/internal/httpfetcher"
) )
@@ -27,10 +26,6 @@ var (
// HTTP status code for successful fetch. // HTTP status code for successful fetch.
const httpStatusOK = 200 const httpStatusOK = 200
// metaCacheSize is how many variants' content types metaCache holds. A
// variant not among them is served as before, reading its .meta file.
const metaCacheSize = 10000
// CacheConfig holds cache configuration. // CacheConfig holds cache configuration.
type CacheConfig struct { type CacheConfig struct {
StateDir string StateDir string
@@ -54,6 +49,12 @@ type CacheConfig struct {
Logger *slog.Logger Logger *slog.Logger
} }
// variantMeta stores content type for fast cache hits without reading .meta file.
type variantMeta struct {
ContentType string
Size int64
}
// Cache implements the caching layer for the image proxy. // Cache implements the caching layer for the image proxy.
type Cache struct { type Cache struct {
db *sql.DB db *sql.DB
@@ -75,10 +76,9 @@ type Cache struct {
evictionStarted bool evictionStarted bool
evictionStopOnce sync.Once evictionStopOnce sync.Once
// metaCache holds the content types of the variants most recently // In-memory cache of variant metadata (content type, size) to avoid
// stored or served, so a hit does not read the variant's .meta file. // reading .meta files
// It never stands in for the variant file, which is always opened. metaCache map[VariantKey]variantMeta
metaCache *lru.Cache[VariantKey, string]
// contentLocks serializes StoreSource and evictSourceBlob per // contentLocks serializes StoreSource and evictSourceBlob per
// content hash, closing the race window between an eviction's row // content hash, closing the race window between an eviction's row
@@ -101,11 +101,6 @@ func NewCache(db *sql.DB, config CacheConfig) (*Cache, error) {
log = slog.Default() log = slog.Default()
} }
metaCache, err := lru.New[VariantKey, string](metaCacheSize)
if err != nil {
return nil, fmt.Errorf("failed to create variant content type cache: %w", err)
}
c := &Cache{ c := &Cache{
db: db, db: db,
config: config, config: config,
@@ -114,7 +109,7 @@ func NewCache(db *sql.DB, config CacheConfig) (*Cache, error) {
evictionPressure: make(chan struct{}, 1), evictionPressure: make(chan struct{}, 1),
evictionStop: make(chan struct{}), evictionStop: make(chan struct{}),
evictionDone: make(chan struct{}), evictionDone: make(chan struct{}),
metaCache: metaCache, metaCache: make(map[VariantKey]variantMeta),
contentLocks: newContentLock(), contentLocks: newContentLock(),
} }
@@ -182,30 +177,13 @@ func (c *Cache) Lookup(ctx context.Context, req *ImageRequest) (*LookupResult, e
}, nil }, nil
} }
// GetVariant returns a reader, size, and content type for a cached // GetVariant returns a reader, size, and content type for a cached variant.
// variant. The content type comes from metaCache, or else from the
// variant's .meta file and is then kept in metaCache. A variant with
// no .meta file is served as application/octet-stream, which is not
// kept.
func (c *Cache) GetVariant(cacheKey VariantKey) (io.ReadCloser, int64, string, error) { func (c *Cache) GetVariant(cacheKey VariantKey) (io.ReadCloser, int64, string, error) {
if c.disabled { if c.disabled {
return nil, 0, "", ErrNotFound return nil, 0, "", ErrNotFound
} }
contentType, known := c.metaCache.Get(cacheKey) return c.variants.LoadWithMeta(cacheKey)
if !known {
return c.loadVariantWithMeta(cacheKey)
}
reader, size, err := c.variants.LoadWithSize(cacheKey)
if err != nil {
// The file is gone, e.g. deleted outside pixa
c.metaCache.Remove(cacheKey)
return nil, 0, "", err
}
return reader, size, contentType, nil
} }
// StoreSource stores fetched source content and metadata. On a // StoreSource stores fetched source content and metadata. On a
@@ -308,8 +286,6 @@ func (c *Cache) StoreVariant(
return err return err
} }
c.metaCache.Add(cacheKey, contentType)
_, err = c.db.ExecContext(ctx, ` _, err = c.db.ExecContext(ctx, `
INSERT INTO variant_content (cache_key, size_bytes, content_type) INSERT INTO variant_content (cache_key, size_bytes, content_type)
VALUES (?, ?, ?) VALUES (?, ?, ?)
@@ -523,31 +499,6 @@ func (c *Cache) IncrementTransformCount(ctx context.Context) {
} }
} }
// loadVariantWithMeta is GetVariant for a variant metaCache does not
// hold: it reads the content type from the variant's .meta file and
// keeps it in metaCache, unless a StoreVariant has put one there
// meanwhile, as the store's is newer. A read that finds no .meta file,
// as one can between a store's writing of the variant file and of its
// .meta file, serves application/octet-stream and keeps nothing, so
// metaCache only ever holds a type read from a .meta file or passed to
// StoreVariant.
func (c *Cache) loadVariantWithMeta(
cacheKey VariantKey,
) (io.ReadCloser, int64, string, error) {
reader, size, contentType, err := c.variants.LoadWithMeta(cacheKey)
if err != nil {
return nil, 0, "", err
}
if contentType == "" {
return reader, size, fallbackContentType, nil
}
c.metaCache.ContainsOrAdd(cacheKey, contentType)
return reader, size, contentType, nil
}
// writeMetadataSidecar writes the JSON metadata sidecar of a stored source. // writeMetadataSidecar writes the JSON metadata sidecar of a stored source.
// A failure is logged and is otherwise non-fatal; the metadata is in the // A failure is logged and is otherwise non-fatal; the metadata is in the
// database. // database.
+3 -7
View File
@@ -39,8 +39,8 @@ const tempFilePrefix = ".tmp-"
// to each variant file. // to each variant file.
const variantMetaSuffix = ".meta" const variantMetaSuffix = ".meta"
// fallbackContentType is the content type given to a variant file that // fallbackContentType is recorded when a reconciled variant file has
// has no readable .meta sidecar, when it is served or reconciled. // no readable .meta sidecar.
const fallbackContentType = "application/octet-stream" const fallbackContentType = "application/octet-stream"
// UsageBytes returns the total number of bytes of cache content // UsageBytes returns the total number of bytes of cache content
@@ -271,9 +271,7 @@ func (c *Cache) sourceCandidates(ctx context.Context) ([]evictionCandidate, erro
// evictVariant removes one variant: accounting row first, then the // evictVariant removes one variant: accounting row first, then the
// content and .meta files, so the database never references a deleted // content and .meta files, so the database never references a deleted
// file. The metaCache entry goes before the files; a GetVariant that // file.
// read them just before may put it back, and the next GetVariant then
// fails to open the file and removes it again.
func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error { func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error {
_, err := c.db.ExecContext(ctx, _, err := c.db.ExecContext(ctx,
`DELETE FROM variant_content WHERE cache_key = ?`, string(cacheKey)) `DELETE FROM variant_content WHERE cache_key = ?`, string(cacheKey))
@@ -281,8 +279,6 @@ func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error {
return fmt.Errorf("failed to delete variant accounting row: %w", err) return fmt.Errorf("failed to delete variant accounting row: %w", err)
} }
c.metaCache.Remove(cacheKey)
err = c.variants.DeleteWithMeta(cacheKey) err = c.variants.DeleteWithMeta(cacheKey)
if err != nil { if err != nil {
return err return err
@@ -1,349 +0,0 @@
package imgcache
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"sync"
"testing"
"time"
)
// webpRequest returns a request for a 100x100 WebP variant of path.
func webpRequest(path string) *ImageRequest {
return &ImageRequest{
SourceHost: testHostCDN,
SourcePath: path,
Size: Size{Width: 100, Height: 100},
Format: FormatWebP,
Quality: 85,
FitMode: FitCover,
}
}
// assertVariantServed checks that GetVariant serves key with the given
// content and the image/webp content type storeEvictionTestVariant stores.
func assertVariantServed(t *testing.T, cache *Cache, key VariantKey, content []byte) {
t.Helper()
reader, size, contentType, err := cache.GetVariant(key)
if err != nil {
t.Fatalf("GetVariant(%s) error = %v", key, err)
}
defer func() { _ = reader.Close() }()
got, err := io.ReadAll(reader)
if err != nil {
t.Fatalf("reading variant %s: %v", key, err)
}
if !bytes.Equal(got, content) {
t.Errorf("GetVariant(%s) content = %q, want %q", key, got, content)
}
if size != int64(len(content)) {
t.Errorf("GetVariant(%s) size = %d, want %d", key, size, len(content))
}
if contentType != testContentTypeWebP {
t.Errorf("GetVariant(%s) content type = %q, want %q",
key, contentType, testContentTypeWebP)
}
}
// assertVariantNotFound checks that GetVariant refuses key with
// ErrNotFound.
func assertVariantNotFound(t *testing.T, cache *Cache, key VariantKey) {
t.Helper()
reader, _, _, err := cache.GetVariant(key)
if err == nil {
_ = reader.Close()
}
if !errors.Is(err, ErrNotFound) {
t.Errorf("GetVariant(%s) error = %v, want ErrNotFound", key, err)
}
}
// assertLookupMisses checks that Lookup reports request as a miss.
func assertLookupMisses(t *testing.T, cache *Cache, request *ImageRequest) {
t.Helper()
lookup, err := cache.Lookup(t.Context(), request)
if err != nil {
t.Fatalf("Lookup(%s) error = %v", request.SourcePath, err)
}
if lookup.Hit {
t.Errorf("Lookup(%s) is a hit, want a miss", request.SourcePath)
}
}
// renameFile renames the file at from to to.
func renameFile(t *testing.T, from, to string) {
t.Helper()
err := os.Rename(from, to)
if err != nil {
t.Fatalf("renaming %s: %v", from, err)
}
}
// TestSecondHitDoesNotReadMetaFile checks that once a variant has been
// stored or read, a hit takes its content type from memory: with the
// .meta file deleted, GetVariant must still return the stored content
// type rather than the application/octet-stream it uses without one.
func TestSecondHitDoesNotReadMetaFile(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<20)
content := []byte("webp variant bytes")
storeEvictionTestVariant(t, cache, testVariantKeyOne, content)
// A second Cache on the same state directory starts with nothing in
// memory, as pixad does after a restart, so its first read uses the
// .meta file.
restarted, err := NewCache(cache.db, cache.config)
if err != nil {
t.Fatalf("NewCache() error = %v", err)
}
assertVariantServed(t, restarted, testVariantKeyOne, content)
err = os.Remove(cache.variants.keyToPath(testVariantKeyOne) + ".meta")
if err != nil {
t.Fatalf("removing .meta file: %v", err)
}
assertVariantServed(t, cache, testVariantKeyOne, content)
assertVariantServed(t, restarted, testVariantKeyOne, content)
}
// TestReadDuringStoreKeepsStoredContentType checks that a GetVariant
// which began before StoreVariant finished cannot replace the content
// type the store kept in memory. Such a read can find the variant file
// but not yet its .meta file, and so gets application/octet-stream. The
// test deletes the .meta file after the store, then runs the part of
// GetVariant that comes after its check of memory.
func TestReadDuringStoreKeepsStoredContentType(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<20)
content := []byte("webp variant bytes")
storeEvictionTestVariant(t, cache, testVariantKeyOne, content)
err := os.Remove(cache.variants.keyToPath(testVariantKeyOne) + ".meta")
if err != nil {
t.Fatalf("removing .meta file: %v", err)
}
reader, _, contentType, err := cache.loadVariantWithMeta(testVariantKeyOne)
if err != nil {
t.Fatalf("loadVariantWithMeta(%s) error = %v", testVariantKeyOne, err)
}
_ = reader.Close()
t.Logf("the read without a .meta file got content type %q", contentType)
assertVariantServed(t, cache, testVariantKeyOne, content)
}
// TestReadOfOlderMetaFileKeepsStoredContentType checks that a read
// which got its content type from a .meta file that StoreVariant had
// not yet rewritten cannot replace the type the store kept in memory.
// The test writes such a .meta file, with a different content type,
// after the store, then runs the part of GetVariant that comes after
// its check of memory.
func TestReadOfOlderMetaFileKeepsStoredContentType(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<20)
content := []byte("webp variant bytes")
storeEvictionTestVariant(t, cache, testVariantKeyOne, content)
olderMeta, err := json.Marshal(VariantMeta{
ContentType: testContentTypeJPEG,
Size: int64(len(content)),
})
if err != nil {
t.Fatalf("encoding .meta file: %v", err)
}
metaPath := cache.variants.keyToPath(testVariantKeyOne) + ".meta"
err = os.WriteFile(metaPath, olderMeta, StorageFilePerm)
if err != nil {
t.Fatalf("writing .meta file: %v", err)
}
reader, _, contentType, err := cache.loadVariantWithMeta(testVariantKeyOne)
if err != nil {
t.Fatalf("loadVariantWithMeta(%s) error = %v", testVariantKeyOne, err)
}
_ = reader.Close()
if contentType != testContentTypeJPEG {
t.Fatalf("loadVariantWithMeta(%s) content type = %q, want %q from the .meta file",
testVariantKeyOne, contentType, testContentTypeJPEG)
}
kept, _ := cache.metaCache.Get(testVariantKeyOne)
if kept != testContentTypeWebP {
t.Errorf("content type in memory = %q, want the stored %q",
kept, testContentTypeWebP)
}
}
// TestFailedReadDuringStoreKeepsStoredContentType checks that a read
// which found no .meta file cannot leave application/octet-stream in
// memory, even when another read has removed the content type
// StoreVariant kept there. In this order: the store; a read that finds
// the variant in memory but cannot open its file, and so removes it
// from memory; a read that opened the variant file before the store
// wrote its .meta file. Later hits must get the stored content type.
func TestFailedReadDuringStoreKeepsStoredContentType(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<20)
content := []byte("webp variant bytes")
variantPath := cache.variants.keyToPath(testVariantKeyOne)
metaPath := variantPath + ".meta"
storeEvictionTestVariant(t, cache, testVariantKeyOne, content)
renameFile(t, variantPath, variantPath+".hidden")
assertVariantNotFound(t, cache, testVariantKeyOne)
renameFile(t, variantPath+".hidden", variantPath)
renameFile(t, metaPath, metaPath+".hidden")
reader, _, contentType, err := cache.GetVariant(testVariantKeyOne)
if err != nil {
t.Fatalf("GetVariant(%s) error = %v", testVariantKeyOne, err)
}
_ = reader.Close()
t.Logf("the read without a .meta file got content type %q", contentType)
renameFile(t, metaPath+".hidden", metaPath)
assertVariantServed(t, cache, testVariantKeyOne, content)
}
// TestEvictedVariantIsNotServed checks that a variant the evictor
// removed is a miss and cannot be read, although it had been stored
// and served before.
func TestEvictedVariantIsNotServed(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1500)
oldRequest := webpRequest("/old.jpg")
newRequest := webpRequest("/new.jpg")
oldKey := CacheKey(oldRequest)
newKey := CacheKey(newRequest)
oldContent := bytes.Repeat([]byte{0x01}, 1000)
newContent := bytes.Repeat([]byte{0x02}, 1000)
storeEvictionTestVariant(t, cache, oldKey, oldContent)
storeEvictionTestVariant(t, cache, newKey, newContent)
assertVariantServed(t, cache, oldKey, oldContent)
setVariantLastAccessed(t, cache, oldKey, time.Now().Add(-time.Hour))
err := cache.EvictToLimit(t.Context())
if err != nil {
t.Fatalf("EvictToLimit() error = %v", err)
}
assertLookupMisses(t, cache, oldRequest)
assertVariantNotFound(t, cache, oldKey)
assertVariantServed(t, cache, newKey, newContent)
}
// TestVariantDeletedFromDiskIsNotServed checks that a variant whose
// file was deleted by something other than the evictor cannot be read,
// and is a miss afterwards.
func TestVariantDeletedFromDiskIsNotServed(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<20)
request := webpRequest("/deleted.jpg")
key := CacheKey(request)
content := []byte("webp variant bytes")
storeEvictionTestVariant(t, cache, key, content)
assertVariantServed(t, cache, key, content)
err := os.Remove(cache.variants.keyToPath(key))
if err != nil {
t.Fatalf("removing variant file: %v", err)
}
assertVariantNotFound(t, cache, key)
assertLookupMisses(t, cache, request)
}
// TestConcurrentVariantStoreReadAndEvict stores, reads and evicts
// variants from several goroutines at once, for the race detector.
func TestConcurrentVariantStoreReadAndEvict(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<20)
ctx := t.Context()
var wg sync.WaitGroup
for goroutine := range 8 {
wg.Go(func() {
key := VariantKey(fmt.Sprintf("aabbccdd01%02d", goroutine))
content := []byte(key)
for range 20 {
err := cache.StoreVariant(
ctx, key, bytes.NewReader(content), testContentTypeWebP)
if err != nil {
t.Errorf("StoreVariant(%s) error = %v", key, err)
return
}
reader, _, contentType, err := cache.GetVariant(key)
if err != nil {
t.Errorf("GetVariant(%s) error = %v", key, err)
return
}
_ = reader.Close()
if contentType != testContentTypeWebP {
t.Errorf("GetVariant(%s) content type = %q, want %q",
key, contentType, testContentTypeWebP)
}
err = cache.evictVariant(ctx, key)
if err != nil {
t.Errorf("evictVariant(%s) error = %v", key, err)
return
}
assertVariantNotFound(t, cache, key)
}
})
}
wg.Wait()
}
+12 -24
View File
@@ -506,44 +506,32 @@ func (s *VariantStorage) Load(key VariantKey) (io.ReadCloser, error) {
return f, nil return f, nil
} }
// LoadWithSize returns a reader and file size for the content at the // LoadWithMeta returns a reader, size, and content type for the content at
// given key. // the given key.
func (s *VariantStorage) LoadWithSize(key VariantKey) (io.ReadCloser, int64, error) { func (s *VariantStorage) LoadWithMeta(
key VariantKey,
) (io.ReadCloser, int64, string, error) {
path := s.keyToPath(key) path := s.keyToPath(key)
metaPath := path + ".meta"
f, err := os.Open(path) //nolint:gosec // path derived from cache key f, err := os.Open(path) //nolint:gosec // path derived from cache key
if err != nil { if err != nil {
if os.IsNotExist(err) { if os.IsNotExist(err) {
return nil, 0, ErrNotFound return nil, 0, "", ErrNotFound
} }
return nil, 0, fmt.Errorf("failed to open content: %w", err) return nil, 0, "", fmt.Errorf("failed to open content: %w", err)
} }
stat, err := f.Stat() stat, err := f.Stat()
if err != nil { if err != nil {
_ = f.Close() _ = f.Close()
return nil, 0, fmt.Errorf("failed to stat content: %w", err) return nil, 0, "", fmt.Errorf("failed to stat content: %w", err)
} }
return f, stat.Size(), nil // Load metadata for content type
} contentType := "application/octet-stream" // fallback
// LoadWithMeta returns a reader, size, and content type for the content at
// the given key. The content type is read from the .meta file, and is
// empty when that file is missing or unreadable.
func (s *VariantStorage) LoadWithMeta(
key VariantKey,
) (io.ReadCloser, int64, string, error) {
f, size, err := s.LoadWithSize(key)
if err != nil {
return nil, 0, "", err
}
var contentType string
metaPath := s.keyToPath(key) + ".meta"
metaData, err := os.ReadFile(metaPath) //nolint:gosec // path derived from cache key metaData, err := os.ReadFile(metaPath) //nolint:gosec // path derived from cache key
if err == nil { if err == nil {
@@ -553,7 +541,7 @@ func (s *VariantStorage) LoadWithMeta(
} }
} }
return f, size, contentType, nil return f, stat.Size(), contentType, nil
} }
// Exists checks if content exists at the given key. // Exists checks if content exists at the given key.
@@ -24,7 +24,6 @@ const (
testHostExample = "example.com" testHostExample = "example.com"
testPathCat = "/photos/cat.jpg" testPathCat = "/photos/cat.jpg"
testContentTypeJPEG = "image/jpeg" testContentTypeJPEG = "image/jpeg"
testContentTypeWebP = "image/webp"
testHeaderContentType = "Content-Type" testHeaderContentType = "Content-Type"
) )
@@ -18,7 +18,6 @@ import (
"sneak.berlin/go/pixa/internal/database" "sneak.berlin/go/pixa/internal/database"
"sneak.berlin/go/pixa/internal/globals" "sneak.berlin/go/pixa/internal/globals"
"sneak.berlin/go/pixa/internal/handlers" "sneak.berlin/go/pixa/internal/handlers"
"sneak.berlin/go/pixa/internal/healthcheck"
"sneak.berlin/go/pixa/internal/logger" "sneak.berlin/go/pixa/internal/logger"
"sneak.berlin/go/pixa/internal/middleware" "sneak.berlin/go/pixa/internal/middleware"
) )
@@ -72,15 +71,8 @@ func newTestServer(t *testing.T) *Server {
t.Fatalf("database.New() error = %v", err) t.Fatalf("database.New() error = %v", err)
} }
hc, err := healthcheck.New(lc, healthcheck.Params{
Globals: &globals.Globals{}, Config: cfg, Logger: log, Database: db,
})
if err != nil {
t.Fatalf("healthcheck.New() error = %v", err)
}
h, err := handlers.New(lc, handlers.Params{ h, err := handlers.New(lc, handlers.Params{
Logger: log, Healthcheck: hc, Database: db, Config: cfg, Logger: log, Database: db, Config: cfg,
}) })
if err != nil { if err != nil {
t.Fatalf("handlers.New() error = %v", err) t.Fatalf("handlers.New() error = %v", err)
@@ -1,169 +0,0 @@
package server
import (
"encoding/json"
"net/http"
"net/http/httptest"
"strconv"
"testing"
"sneak.berlin/go/pixa/internal/healthcheck"
)
// unsignedImagePath is an image URL that carries no signature.
const unsignedImagePath = "/v1/image/cdn.example.com/cat.jpg/100x100.jpeg"
// TestMaintenanceModeRefusesImageRequests verifies that while maintenance
// mode is on, both image routes answer 503 Service Unavailable with a
// Retry-After header and the JSON error body the image handlers send.
func TestMaintenanceModeRefusesImageRequests(t *testing.T) {
t.Parallel()
s := newTestServer(t)
s.config.MaintenanceMode = true
requests := []struct {
method string
path string
}{
{http.MethodGet, unsignedImagePath},
{http.MethodHead, unsignedImagePath},
{http.MethodGet, "/v1/e/token/cat.jpg"},
}
for _, tc := range requests {
t.Run(tc.method+" "+tc.path, func(t *testing.T) {
t.Parallel()
rec := httptest.NewRecorder()
s.ServeHTTP(rec, httptest.NewRequestWithContext(
t.Context(), tc.method, tc.path, nil))
t.Logf("status %d, body %s", rec.Code, rec.Body.String())
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want %d",
rec.Code, http.StatusServiceUnavailable)
}
retryAfter := rec.Header().Get("Retry-After")
seconds, err := strconv.Atoi(retryAfter)
if err != nil || seconds <= 0 {
t.Errorf("Retry-After = %q, want a positive number of seconds",
retryAfter)
}
// A HEAD response carries no body.
if tc.method == http.MethodHead {
return
}
var body struct {
Error string `json:"error"`
Status int `json:"status"`
Timestamp string `json:"timestamp"`
}
err = json.NewDecoder(rec.Body).Decode(&body)
if err != nil {
t.Fatalf("body is not JSON: %v", err)
}
if body.Error == "" || body.Status != http.StatusServiceUnavailable ||
body.Timestamp == "" {
t.Errorf("body = %+v, want an error, status %d and a timestamp",
body, http.StatusServiceUnavailable)
}
})
}
}
// TestImageRequestsServedWithoutMaintenanceMode verifies that while
// maintenance mode is off, image requests reach the image handlers instead
// of the 503. The handlers refuse an unsigned image URL with 401 and a token
// they cannot decrypt with 400, so either status shows a request got through.
func TestImageRequestsServedWithoutMaintenanceMode(t *testing.T) {
t.Parallel()
s := newTestServer(t)
s.config.MaintenanceMode = false
requests := []struct {
method string
path string
want int
}{
{http.MethodGet, unsignedImagePath, http.StatusUnauthorized},
{http.MethodHead, unsignedImagePath, http.StatusUnauthorized},
{http.MethodGet, "/v1/e/token/cat.jpg", http.StatusBadRequest},
}
for _, tc := range requests {
t.Run(tc.method+" "+tc.path, func(t *testing.T) {
t.Parallel()
rec := httptest.NewRecorder()
s.ServeHTTP(rec, httptest.NewRequestWithContext(
t.Context(), tc.method, tc.path, nil))
t.Logf("status %d, body %s", rec.Code, rec.Body.String())
if rec.Code != tc.want {
t.Errorf("status = %d, want %d from the image handler",
rec.Code, tc.want)
}
})
}
}
// TestMaintenanceModeKeepsOtherRoutes verifies that while maintenance mode
// is on, the health check still answers 200 and reports it, and the login
// page and /metrics still answer 200. The image's Docker HEALTHCHECK
// requests the health check: a 503 there would make the container
// unhealthy, and upaas marks a deploy failed when its container is
// unhealthy.
func TestMaintenanceModeKeepsOtherRoutes(t *testing.T) {
t.Parallel()
s := newTestServer(t)
s.config.MaintenanceMode = true
// /metrics is routed only when its username is set.
s.config.MetricsUsername = "metrics"
s.config.MetricsPassword = "metrics-password"
s.SetupRoutes()
rec := httptest.NewRecorder()
s.ServeHTTP(rec, httptest.NewRequestWithContext(t.Context(),
http.MethodGet, "/.well-known/healthcheck.json", nil))
t.Logf("health check status %d, body %s", rec.Code, rec.Body.String())
if rec.Code != http.StatusOK {
t.Fatalf("health check status = %d, want %d", rec.Code, http.StatusOK)
}
var health healthcheck.Response
err := json.NewDecoder(rec.Body).Decode(&health)
if err != nil || !health.Maintenance {
t.Errorf("health check maintenance_mode = %v (error %v), want true",
health.Maintenance, err)
}
rec = httptest.NewRecorder()
s.ServeHTTP(rec, clientRequest(t, http.MethodGet, nil, firstClient, ""))
if rec.Code != http.StatusOK {
t.Errorf("login page status = %d, want %d", rec.Code, http.StatusOK)
}
req := httptest.NewRequestWithContext(t.Context(),
http.MethodGet, "/metrics", nil)
req.SetBasicAuth(s.config.MetricsUsername, s.config.MetricsPassword)
rec = httptest.NewRecorder()
s.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Errorf("/metrics status = %d, want %d", rec.Code, http.StatusOK)
}
}
+3 -44
View File
@@ -1,9 +1,7 @@
package server package server
import ( import (
"encoding/json"
"net/http" "net/http"
"strconv"
"time" "time"
sentryhttp "github.com/getsentry/sentry-go/http" sentryhttp "github.com/getsentry/sentry-go/http"
@@ -19,10 +17,6 @@ import (
// make per minute; the next is refused with 429 Too Many Requests. // make per minute; the next is refused with 429 Too Many Requests.
const LoginAttemptsPerMinute = 5 const LoginAttemptsPerMinute = 5
// MaintenanceRetryAfterSeconds is the Retry-After, in seconds, sent with
// the 503 that the image routes answer while maintenance mode is on.
const MaintenanceRetryAfterSeconds = 300
// SetupRoutes configures all HTTP routes. // SetupRoutes configures all HTTP routes.
func (s *Server) SetupRoutes() { func (s *Server) SetupRoutes() {
s.router = chi.NewRouter() s.router = chi.NewRouter()
@@ -74,23 +68,15 @@ func (s *Server) SetupRoutes() {
s.router.Get("/logout", s.h.HandleLogout()) s.router.Get("/logout", s.h.HandleLogout())
// Image routes, refused while maintenance mode is on. Only these: the
// image's Docker HEALTHCHECK requests the health check, a 503 there
// would make the container unhealthy, and upaas marks a deploy failed
// when its container is unhealthy.
s.router.Group(func(r chi.Router) {
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>
r.Get("/v1/image/*", s.h.HandleImage()) s.router.Get("/v1/image/*", s.h.HandleImage())
r.Head("/v1/image/*", s.h.HandleImage()) s.router.Head("/v1/image/*", s.h.HandleImage())
// Encrypted image URL route // Encrypted image URL route
// The trailing filename (e.g., /img.jpg) is ignored but helps // The trailing filename (e.g., /img.jpg) is ignored but helps
// browsers with content type // browsers with content type
r.Get("/v1/e/{token}/*", s.h.HandleImageEnc()) s.router.Get("/v1/e/{token}/*", s.h.HandleImageEnc())
})
// Metrics endpoint with auth // Metrics endpoint with auth
if s.config.MetricsUsername != "" { if s.config.MetricsUsername != "" {
@@ -100,30 +86,3 @@ func (s *Server) SetupRoutes() {
}) })
} }
} }
// refuseDuringMaintenance answers a request with 503 Service Unavailable,
// a Retry-After header and a JSON error body while maintenance mode is on,
// and passes it on otherwise. The body has the fields of the JSON errors
// the image handlers send.
func (s *Server) refuseDuringMaintenance(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !s.MaintenanceMode() {
next.ServeHTTP(w, r)
return
}
w.Header().Set("Retry-After", strconv.Itoa(MaintenanceRetryAfterSeconds))
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusServiceUnavailable)
err := json.NewEncoder(w).Encode(map[string]any{
"error": "down for maintenance, try again later",
"status": http.StatusServiceUnavailable,
"timestamp": time.Now().UTC().Format(time.RFC3339),
})
if err != nil {
s.log.Error("json encode error", "error", err)
}
})
}