Compare commits
2
Commits
e455901de3
...
1789fd4584
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1789fd4584 | ||
|
|
1616e91a6a |
@@ -92,8 +92,69 @@ the metadata file stored beside it.
|
||||
|
||||
### Routes
|
||||
|
||||
pixa answers these routes; any other path answers 404. A path in this list asked
|
||||
with a method the list does not give answers 405, except `/static/<file>`, which
|
||||
answers any method as it answers `GET`. A browser's CORS preflight request
|
||||
(`OPTIONS` with `Origin` and `Access-Control-Request-Method` headers) to any
|
||||
path under `/v1/` answers 200, in maintenance mode too.
|
||||
|
||||
- `GET /` — the login page, or the URL generator page with a login session
|
||||
(see Encrypted URLs). Needs: nothing. Answers: 200.
|
||||
- `POST /` — log in with the signing key typed into the login page. Needs: the
|
||||
login page's form (below). Answers: 303 to `/` with a login session cookie
|
||||
that lasts 30 days for the right key; 200 with the login page and an error for
|
||||
a wrong key; 429 over the login limit (below).
|
||||
- `POST /generate` — make an encrypted URL from the generator page's form.
|
||||
Needs: a login session and the generator page's form (below); without a login
|
||||
session it answers 303 to `/`. Answers: 200 with the page showing the URL; 400
|
||||
with the page naming a field that is not valid; 500 when the URL cannot be
|
||||
made.
|
||||
- `GET /logout` — end the login session. Needs: nothing. Answers: 303 to `/`.
|
||||
- `GET` or `HEAD` `/v1/image/<host>/<path>/<size>.<format>` — an image, fetched,
|
||||
resized and converted (below). Needs: a signature, unless the host is
|
||||
allowlisted (see Source Hosts). Answers: 200; 304 when `If-None-Match` matches
|
||||
the image's `ETag`; 400 for a URL or parameter that is not valid; 401 for a
|
||||
missing or wrong signature, a missing `exp` or an `exp` in the past; 403 when
|
||||
the upstream host, or a host it redirects to, is `localhost`, ends in
|
||||
`.localhost` or `.local`, or has an address in a blocked network (see
|
||||
`blocked_networks`); 502 when the upstream answered with an error status, and
|
||||
for 5 minutes after that for the same source URL; 503 when pixa is busy or in
|
||||
maintenance mode; 500 for any other failure.
|
||||
- `GET /v1/e/<token>/<name>` — an image through an encrypted URL (see Encrypted
|
||||
URLs). Needs: nothing but the URL. Answers: 200; 400 for a token that does not
|
||||
decrypt, or that asks for a size or fit that is not valid; 410 once it has
|
||||
expired; 504 when the upstream has not sent its response headers within
|
||||
`upstream_fetch_timeout`, but 500 when that time runs out while the image
|
||||
itself is still arriving; 403, 502, 503 and 500 as for `/v1/image/`.
|
||||
- `GET /robots.txt` — asks every crawler to stay away (`Disallow: /`). Needs:
|
||||
nothing. Answers: 200.
|
||||
- `GET /.well-known/healthcheck.json` — JSON with `status` (`ok`), `now`,
|
||||
`uptime_seconds`, `uptime_human`, `version`, `appname` and
|
||||
`maintenance_mode`. Needs: nothing. Answers: 200, always.
|
||||
- `GET /static/<file>` — the script the login and generator pages load. Needs:
|
||||
nothing. Answers: 200, or 404 for a file that does not exist.
|
||||
- `GET /metrics` — Prometheus metrics (see Architecture). Needs: HTTP basic
|
||||
authentication with `metrics.username` and `metrics.password`. Answers: 200;
|
||||
401 without them; 404 when they are not set, as the route then does not exist.
|
||||
|
||||
Both `POST` routes accept only a form that pixa's own page served: the page puts
|
||||
a token in the form and sets a cookie to match, and a request without both is
|
||||
refused with 403, so another site cannot submit the form from a visitor's
|
||||
browser. The login and generator pages are meant to be opened over HTTPS: while
|
||||
`debug` is off, a form sent from a page opened over plain HTTP is refused with
|
||||
403, and while it is on, so is one sent from a page opened over HTTPS. Plain
|
||||
HTTP is for development on the browser's own machine: the login session cookie
|
||||
is always marked `Secure`, and over plain HTTP a browser keeps such a cookie
|
||||
only for its own machine (`localhost`), if at all. A form is also refused with
|
||||
403 when the page's host is not the `Host` header pixa receives, so a reverse
|
||||
proxy in front of pixa must pass that header on unchanged. A form body over
|
||||
1 MiB is refused with 413. The image routes answer the errors listed for them
|
||||
with JSON holding `error`, `status` and `timestamp`.
|
||||
|
||||
An image URL has this form:
|
||||
|
||||
```
|
||||
/v1/image/<host>/<path>/<size>.<format>?sig=<signature>&exp=<expiration>
|
||||
/v1/image/<host>/<path>/<size>.<format>?sig=<signature>&exp=<expiration>&q=<quality>&fit=<fit>
|
||||
```
|
||||
|
||||
Images are only fetched from origins using TLS with valid certificates, unless
|
||||
@@ -106,6 +167,11 @@ than once, is refused with 400.
|
||||
- `<format>`: one of `orig` (or `original`), `jpeg` (or `jpg`), `png`, `webp`,
|
||||
`avif`, `gif`
|
||||
- `<size>`: `orig` or `<width>x<height>` (e.g. `800x600`)
|
||||
- `sig` and `exp`: the signature and its expiry, needed unless the host is
|
||||
allowlisted (see Signature Specification)
|
||||
- `q` and `fit`: the output quality and how the image is fitted to `<size>`,
|
||||
both optional (values under Signature Specification). Both are part of what
|
||||
is cached, so each value of either is a separate cached image.
|
||||
|
||||
An image is served with `Cache-Control: public, max-age=<seconds>, immutable`.
|
||||
When the URL has an expiry (an `exp`, or the TTL of an encrypted URL),
|
||||
@@ -139,6 +205,39 @@ its own `X-Forwarded-For`, whether it connects directly or through the proxy,
|
||||
because its own address is trusted too. Setting `trusted_proxies` to only the
|
||||
address pixa sees for requests that come through the proxy closes this.
|
||||
|
||||
### Encrypted URLs
|
||||
|
||||
An encrypted URL is an image URL made on pixa's own web page by someone who
|
||||
knows the signing key. It works for any upstream host, allowlisted or not,
|
||||
without a signature, and whoever gets it can neither read the source URL from it
|
||||
nor change what it asks for.
|
||||
|
||||
1. Open `/` in a browser over HTTPS (or over plain HTTP while `debug` is on, see
|
||||
Routes) and log in with the signing key (`signing_key`). The login session
|
||||
lasts 30 days, or until `/logout`.
|
||||
2. On the generator page, give the source image's URL, the width and height, the
|
||||
format, quality and fit, and how long the URL lasts, then submit the form
|
||||
(`POST /generate`). Width and height both empty or `0` keep the original
|
||||
size; if only one of them is empty or `0`, that side is scaled to keep the
|
||||
image's proportions.
|
||||
3. The page shows the URL, `https://<host>/v1/e/<token>/img.<format>`, and when
|
||||
it expires. `<host>` is the host the page was opened on, and the URL starts
|
||||
with `http` instead while `debug` is on. The name after the token is ignored
|
||||
and only gives the URL a file extension, `jpg` for `orig`.
|
||||
|
||||
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 URL's scheme is not kept: the image is fetched like any other (see
|
||||
Routes), and the blocked networks still apply.
|
||||
|
||||
How long the URL lasts is chosen on the page, from 1 minute to 1 year, or
|
||||
never. The expiry is fixed in the token when the URL is made and cannot be
|
||||
changed or revoked afterwards. Until then the image is served with a `max-age`
|
||||
that ends at the expiry (see Routes); after it the URL answers 410
|
||||
`URL has expired`. A URL made to last forever stops working only when
|
||||
`signing_key` changes: changing it makes every encrypted URL already handed out
|
||||
answer 400, and ends every login session.
|
||||
|
||||
### Image Metadata
|
||||
|
||||
pixa decodes and re-encodes every image it serves, and removes all metadata from
|
||||
@@ -237,6 +336,17 @@ startup naming it, as an unknown config key does. The one other accepted
|
||||
name is `PIXA_CONFIG_PATH`, the config file's path (like `--config`). The
|
||||
variables set by the file's `env:` section are checked the same way.
|
||||
|
||||
pixa reads at most one config file: the one given with `--config` (or `-c`),
|
||||
otherwise the one `PIXA_CONFIG_PATH` names, otherwise the first of these that
|
||||
pixa finds: `/etc/pixa/config.yml`, `/etc/pixa/config.yaml`,
|
||||
`~/.config/pixa/config.yml`, `~/.config/pixa/config.yaml`, then `config.yml`
|
||||
and `config.yaml` in the working directory. A named file that does not exist,
|
||||
cannot be read or does not parse aborts startup. Of the files pixa looks for on
|
||||
its own, one it finds but cannot read or parse aborts startup; one it cannot
|
||||
find, for any reason, is passed over without a message, even when the file is
|
||||
there in a directory pixa may not enter. With no file, pixa uses the environment
|
||||
and the defaults.
|
||||
|
||||
| Variable | Config key | Meaning |
|
||||
| ------------------------------------ | ------------------------------- | ---------------------------------------------------------------------------- |
|
||||
| `PIXA_SIGNING_KEY` | `signing_key` | Required: secret for signed and encrypted URLs and login, 32+ characters |
|
||||
@@ -322,12 +432,12 @@ Key settings in more detail:
|
||||
seconds for one to free up; if none does, and `downstream_timeout` has not
|
||||
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
|
||||
`/v1/e/`) answer every request for an image 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.
|
||||
|
||||
|
||||
@@ -29,6 +29,23 @@ P2: security: referer blacklist
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 2026-10-04 routes, encrypted URLs and config file documented (closes #75):
|
||||
"Routes" in `README.md` lists every route with its method, purpose, what it
|
||||
needs and the status codes it answers with, and says `q` and `fit` are part
|
||||
of what is cached; "Encrypted URLs" covers logging in, making one on the
|
||||
generator page, how long it lasts and the 410 once it has expired;
|
||||
"Configuration" gives the order in which pixa looks for its config file;
|
||||
`config.example.yml` lists `db_url` and `env` and gives every key's default;
|
||||
`scripts/manual-test.sh` is left to #97.
|
||||
- 2026-10-04 shutdown stops cache eviction in progress (closes #102):
|
||||
`StartEviction` runs the eviction goroutine with its own context, which
|
||||
`StopEviction` cancels, so a pass in progress stops at its next database
|
||||
call, file, row or eviction candidate instead of running to completion, and
|
||||
no pass starts after it, so a stop logs at most one warning;
|
||||
`StopEviction` takes a context and, when that context ends before the
|
||||
goroutine exits, stops waiting and returns its error; the handlers' stop hook
|
||||
passes fx's stop context, so an eviction still running when fx's stop
|
||||
deadline ends fails the stop and makes the exit code 1.
|
||||
- 2026-10-04 dead code in `internal/imgcache` is gone (closes #73): `Purge`,
|
||||
which only returned an error and which nothing called, is no longer part of
|
||||
the `ImageCache` interface or `Service`; the `SignatureValidator`,
|
||||
@@ -432,7 +449,5 @@ P2: security: referer blacklist
|
||||
- integration tests for the image proxy flow
|
||||
- load tests to verify the 1k to 5k req/s target
|
||||
- P2: documentation
|
||||
- configuration options
|
||||
- API endpoints
|
||||
- deployment guide
|
||||
- example nginx or caddy reverse proxy config
|
||||
|
||||
+27
-9
@@ -12,27 +12,38 @@
|
||||
# Durations are Go duration strings such as 30s or 2m and must be
|
||||
# positive; a bare number has no unit and aborts startup. Sizes are a
|
||||
# whole number of bytes.
|
||||
#
|
||||
# A key left out takes the default its comment gives.
|
||||
|
||||
# Server settings
|
||||
# Port to listen on (default: 8080)
|
||||
port: 8080
|
||||
|
||||
# Debug logging and plain-HTTP local development (default: 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.
|
||||
# for an image 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. (default: false)
|
||||
maintenance_mode: false
|
||||
|
||||
# Data directory for SQLite database and cache files
|
||||
# (default: /var/lib/pixa)
|
||||
state_dir: ./data
|
||||
|
||||
# SQLite database URL (default:
|
||||
# file:<state_dir>/state.sqlite3?_journal_mode=WAL). An empty value aborts
|
||||
# startup; leave the key out to use the default.
|
||||
# db_url: "file:./data/state.sqlite3?_journal_mode=WAL"
|
||||
|
||||
# Image proxy settings
|
||||
# HMAC signing key for URL signatures (required, at least 32 characters)
|
||||
# Generate with: openssl rand -base64 32
|
||||
signing_key: "CHANGE_ME_generate_with_openssl_rand_base64_32"
|
||||
|
||||
# Hosts that don't require signatures
|
||||
# Hosts that don't require signatures (default: none)
|
||||
# Use "." prefix for wildcard subdomain matching (e.g., ".example.com" matches "cdn.example.com")
|
||||
allowlist_hosts:
|
||||
- s3.sneak.cloud
|
||||
@@ -45,7 +56,7 @@ allowlist_hosts:
|
||||
# SSRF protection. These are added to the always-enforced built-in ranges
|
||||
# (loopback, RFC 1918 private, link-local, CGNAT, benchmark, NAT64, and
|
||||
# similar), never replacing them. Each entry must be a valid CIDR in IPv4
|
||||
# or IPv6 form; an invalid entry aborts startup.
|
||||
# or IPv6 form; an invalid entry aborts startup. (default: none)
|
||||
# blocked_networks:
|
||||
# - 100.64.0.0/10
|
||||
# - 2001:db8::/32
|
||||
@@ -72,6 +83,7 @@ allowlist_hosts:
|
||||
# - 2001:db8::/32
|
||||
|
||||
# Allow HTTP upstream (only for testing, always use HTTPS in production)
|
||||
# (default: false)
|
||||
allow_http: false
|
||||
|
||||
# Maximum concurrent connections per upstream host (default: 20)
|
||||
@@ -121,10 +133,16 @@ access_control_allow_origin: "*"
|
||||
# with a minimum of 500 MiB.
|
||||
# cache_max_bytes: 10737418240
|
||||
|
||||
# Sentry error reporting (optional)
|
||||
# Sentry DSN for error reporting (default: empty, which turns it off)
|
||||
sentry_dsn: ""
|
||||
|
||||
# Metrics endpoint authentication (optional)
|
||||
# Username and password for /metrics, set together (default: unset). Metrics
|
||||
# are measured and /metrics is served only when both are set.
|
||||
# metrics:
|
||||
# username: "admin"
|
||||
# password: "secret"
|
||||
|
||||
# Environment variables set while this file loads, as described at the top
|
||||
# (default: none)
|
||||
# env:
|
||||
# PIXA_DEBUG: "true"
|
||||
|
||||
@@ -58,23 +58,16 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) {
|
||||
}
|
||||
|
||||
lc.Append(fx.Hook{
|
||||
// The eviction goroutine must outlive OnStart, so it cannot
|
||||
// inherit this hook's context. It makes its own instead, which
|
||||
// leaves it uncancellable: an in-flight pass runs to completion
|
||||
// during OnStop regardless of the shutdown deadline. Making the
|
||||
// loop cancellable changes shutdown semantics and is tracked
|
||||
// separately in issue #102, rather than being folded into the
|
||||
// lint-conformance change that surfaced it.
|
||||
//nolint:contextcheck // see issue #102
|
||||
//nolint:contextcheck // the eviction loop outlives OnStart; OnStop cancels it
|
||||
OnStart: func(_ context.Context) error {
|
||||
return s.initImageService()
|
||||
},
|
||||
OnStop: func(_ context.Context) error {
|
||||
if s.imgCache != nil {
|
||||
s.imgCache.StopEviction()
|
||||
OnStop: func(ctx context.Context) error {
|
||||
if s.imgCache == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return nil
|
||||
return s.imgCache.StopEviction(ctx)
|
||||
},
|
||||
})
|
||||
|
||||
|
||||
@@ -11,7 +11,6 @@ import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
lru "github.com/hashicorp/golang-lru/v2"
|
||||
@@ -69,11 +68,11 @@ type Cache struct {
|
||||
|
||||
// Eviction machinery. The channels are created in NewCache so
|
||||
// stores can signal write pressure without racing StartEviction.
|
||||
// evictionCancel, set by StartEviction, cancels the eviction
|
||||
// goroutine's context.
|
||||
evictionPressure chan struct{}
|
||||
evictionStop chan struct{}
|
||||
evictionDone chan struct{}
|
||||
evictionStarted bool
|
||||
evictionStopOnce sync.Once
|
||||
evictionCancel context.CancelFunc
|
||||
|
||||
// metaCache holds the content types of the variants most recently
|
||||
// stored or served, so a hit does not read the variant's .meta file.
|
||||
@@ -112,7 +111,6 @@ func NewCache(db *sql.DB, config CacheConfig) (*Cache, error) {
|
||||
log: log,
|
||||
disabled: config.DisableDiskCache,
|
||||
evictionPressure: make(chan struct{}, 1),
|
||||
evictionStop: make(chan struct{}),
|
||||
evictionDone: make(chan struct{}),
|
||||
metaCache: metaCache,
|
||||
contentLocks: newContentLock(),
|
||||
|
||||
@@ -117,7 +117,10 @@ func (c *Cache) EvictToLimit(ctx context.Context) error {
|
||||
|
||||
// evictBatch fetches one batch of LRU candidates across variants and
|
||||
// source blobs and evicts them oldest-first until excessBytes are
|
||||
// freed or the batch is exhausted. It returns the bytes freed.
|
||||
// freed or the batch is exhausted. It returns the bytes freed. A
|
||||
// candidate that fails once ctx is cancelled (every one started after
|
||||
// that fails at its first database call) ends the batch with ctx's
|
||||
// error, without a warning.
|
||||
func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error) {
|
||||
candidates, err := c.evictionCandidates(ctx)
|
||||
if err != nil {
|
||||
@@ -133,6 +136,10 @@ func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error
|
||||
|
||||
err := c.evictCandidate(ctx, candidate)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return freed, ctx.Err()
|
||||
}
|
||||
|
||||
c.log.Warn("failed to evict cache entry",
|
||||
"cache_key", candidate.cacheKey,
|
||||
"content_hash", candidate.contentHash,
|
||||
@@ -424,37 +431,44 @@ func (c *Cache) notifyWritePressure() {
|
||||
// startup and again on every periodic tick thereafter, and evicts to
|
||||
// the configured limit on the given periodic interval and on
|
||||
// write-pressure notifications. It is a no-op on a disabled cache or
|
||||
// when already started.
|
||||
// when already started. The goroutine outlives the caller, so it runs
|
||||
// with its own context, which StopEviction cancels.
|
||||
func (c *Cache) StartEviction(interval time.Duration) {
|
||||
if c.disabled || c.evictionStarted {
|
||||
if c.disabled || c.evictionCancel != nil {
|
||||
return
|
||||
}
|
||||
|
||||
c.evictionStarted = true
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
c.evictionCancel = cancel
|
||||
|
||||
go c.evictionLoop(interval)
|
||||
go c.evictionLoop(ctx, interval)
|
||||
}
|
||||
|
||||
// StopEviction stops the background eviction goroutine and waits for
|
||||
// it to exit. It is safe to call when eviction was never started, and
|
||||
// safe to call more than once.
|
||||
func (c *Cache) StopEviction() {
|
||||
if !c.evictionStarted {
|
||||
return
|
||||
// StopEviction cancels the background eviction goroutine, which
|
||||
// interrupts a pass in progress, and waits for it to exit or for ctx to
|
||||
// end, whichever comes first. In the second case it returns an error
|
||||
// wrapping ctx's error. It is safe to call when eviction was never
|
||||
// started, and safe to call more than once.
|
||||
func (c *Cache) StopEviction(ctx context.Context) error {
|
||||
if c.evictionCancel == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
c.evictionStopOnce.Do(func() {
|
||||
close(c.evictionStop)
|
||||
<-c.evictionDone
|
||||
})
|
||||
c.evictionCancel()
|
||||
|
||||
select {
|
||||
case <-c.evictionDone:
|
||||
return nil
|
||||
case <-ctx.Done():
|
||||
return fmt.Errorf("cache eviction still running: %w", ctx.Err())
|
||||
}
|
||||
}
|
||||
|
||||
// evictionLoop is the body of the background eviction goroutine.
|
||||
func (c *Cache) evictionLoop(interval time.Duration) {
|
||||
// evictionLoop is the body of the background eviction goroutine. It
|
||||
// returns when ctx is cancelled, and starts no pass after that.
|
||||
func (c *Cache) evictionLoop(ctx context.Context, interval time.Duration) {
|
||||
defer close(c.evictionDone)
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
c.runReconciliationPass(ctx)
|
||||
c.runEvictionPass(ctx)
|
||||
|
||||
@@ -463,7 +477,7 @@ func (c *Cache) evictionLoop(interval time.Duration) {
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-c.evictionStop:
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
// Reconciliation walks the cache directories, so it only
|
||||
@@ -484,8 +498,13 @@ func (c *Cache) evictionLoop(interval time.Duration) {
|
||||
}
|
||||
|
||||
// runEvictionPass runs one eviction pass, logging failures instead of
|
||||
// propagating them (the loop must keep running).
|
||||
// propagating them (the loop must keep running). It does nothing once
|
||||
// ctx is cancelled.
|
||||
func (c *Cache) runEvictionPass(ctx context.Context) {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
|
||||
err := c.EvictToLimit(ctx)
|
||||
if err != nil {
|
||||
c.log.Warn("cache eviction pass failed", "error", err)
|
||||
@@ -493,8 +512,13 @@ func (c *Cache) runEvictionPass(ctx context.Context) {
|
||||
}
|
||||
|
||||
// runReconciliationPass runs one reconciliation pass, logging failures
|
||||
// instead of propagating them (the loop must keep running).
|
||||
// instead of propagating them (the loop must keep running). It does
|
||||
// nothing once ctx is cancelled.
|
||||
func (c *Cache) runReconciliationPass(ctx context.Context) {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
|
||||
err := c.reconcileAccounting(ctx)
|
||||
if err != nil {
|
||||
c.log.Warn("cache accounting reconciliation failed", "error", err)
|
||||
@@ -511,7 +535,8 @@ func (c *Cache) runReconciliationPass(ctx context.Context) {
|
||||
// know (and rows whose files are gone), and sweeps stale temp files
|
||||
// left behind by crashed writes. Running it periodically, not just
|
||||
// once, bounds how long such drift can accumulate unaccounted for on a
|
||||
// long-running process to one eviction interval.
|
||||
// long-running process to one eviction interval. Once ctx is cancelled,
|
||||
// it stops at the next file or row and returns ctx's error.
|
||||
func (c *Cache) reconcileAccounting(ctx context.Context) error {
|
||||
if c.disabled {
|
||||
return nil
|
||||
@@ -546,6 +571,10 @@ func (c *Cache) reconcileVariantFiles(ctx context.Context) error {
|
||||
return filepath.WalkDir(
|
||||
c.variants.baseDir,
|
||||
func(path string, entry fs.DirEntry, err error) error {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
if err != nil || entry.IsDir() {
|
||||
return err
|
||||
}
|
||||
@@ -634,6 +663,10 @@ func (c *Cache) reconcileVariantRows(ctx context.Context) error {
|
||||
}
|
||||
|
||||
for _, key := range keys {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
if c.variants.Exists(key) {
|
||||
continue
|
||||
}
|
||||
@@ -699,6 +732,10 @@ func (c *Cache) reconcileSourceFiles(ctx context.Context) error {
|
||||
return filepath.WalkDir(
|
||||
c.srcContent.baseDir,
|
||||
func(path string, entry fs.DirEntry, err error) error {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
if err != nil || entry.IsDir() {
|
||||
return err
|
||||
}
|
||||
@@ -760,6 +797,10 @@ func (c *Cache) reconcileSourceRows(ctx context.Context) error {
|
||||
}
|
||||
|
||||
for _, hash := range hashes {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
if c.srcContent.Exists(hash) {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -4,9 +4,12 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"database/sql"
|
||||
"errors"
|
||||
"io/fs"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -657,7 +660,7 @@ func TestEvictionRunsUnderWritePressure(t *testing.T) {
|
||||
// An interval far longer than the test ensures only write
|
||||
// pressure can trigger eviction here.
|
||||
cache.StartEviction(time.Hour)
|
||||
defer cache.StopEviction()
|
||||
defer func() { _ = cache.StopEviction(t.Context()) }()
|
||||
|
||||
keys := []VariantKey{
|
||||
testVariantKeyOne, testVariantKeyTwo, testVariantKeyThree,
|
||||
@@ -690,7 +693,7 @@ func TestEvictionRunsOnPeriodicSchedule(t *testing.T) {
|
||||
// write-pressure notification fires and only the periodic ticker
|
||||
// can trigger eviction.
|
||||
cache.StartEviction(100 * time.Millisecond)
|
||||
defer cache.StopEviction()
|
||||
defer func() { _ = cache.StopEviction(t.Context()) }()
|
||||
|
||||
keys := []VariantKey{
|
||||
testVariantKeyOne, testVariantKeyTwo, testVariantKeyThree,
|
||||
@@ -751,7 +754,7 @@ func TestStartEvictionReconcilesAccountingWithDisk(t *testing.T) {
|
||||
}
|
||||
|
||||
cache.StartEviction(time.Hour)
|
||||
defer cache.StopEviction()
|
||||
defer func() { _ = cache.StopEviction(t.Context()) }()
|
||||
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
|
||||
@@ -808,7 +811,7 @@ func TestPeriodicReconciliationAdoptsFileThatAppearsAfterStartup(t *testing.T) {
|
||||
const interval = 100 * time.Millisecond
|
||||
|
||||
cache.StartEviction(interval)
|
||||
defer cache.StopEviction()
|
||||
defer func() { _ = cache.StopEviction(t.Context()) }()
|
||||
|
||||
// Let startup reconciliation run and settle on an empty cache
|
||||
// before introducing the untracked file, so the adoption we assert
|
||||
@@ -862,6 +865,242 @@ func TestPeriodicReconciliationAdoptsFileThatAppearsAfterStartup(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestStopEvictionInterruptsPassInProgress holds the test database's
|
||||
// only connection, so the startup reconciliation pass waits for it, and
|
||||
// checks that StopEviction stops that pass instead of waiting for the
|
||||
// connection to come free, and that the stop logs one warning: the
|
||||
// interrupted reconciliation's, with no eviction pass started after it.
|
||||
func TestStopEvictionInterruptsPassInProgress(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cache, _ := newEvictionTestCache(t, 1<<30)
|
||||
|
||||
var logBuf bytes.Buffer
|
||||
|
||||
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
|
||||
|
||||
conn, err := cache.db.Conn(t.Context())
|
||||
if err != nil {
|
||||
t.Fatalf("failed to take the database connection: %v", err)
|
||||
}
|
||||
|
||||
defer func() { _ = conn.Close() }()
|
||||
|
||||
cache.StartEviction(time.Hour)
|
||||
|
||||
// The pass is in progress once it waits for the connection.
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
|
||||
for cache.db.Stats().WaitCount == 0 {
|
||||
if time.Now().After(deadline) {
|
||||
t.Fatal("the reconciliation pass never waited for the database")
|
||||
}
|
||||
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
err = cache.StopEviction(ctx)
|
||||
t.Logf("StopEviction() error = %v", err)
|
||||
|
||||
if err != nil {
|
||||
t.Fatalf("StopEviction() error = %v, want nil: the pass waiting for "+
|
||||
"the database did not stop", err)
|
||||
}
|
||||
|
||||
t.Logf("log output: %s", logBuf.String())
|
||||
|
||||
warnings := strings.Count(logBuf.String(), `"level":"WARN"`)
|
||||
if warnings != 1 {
|
||||
t.Errorf("the stop logged %d warnings, want 1", warnings)
|
||||
}
|
||||
}
|
||||
|
||||
// TestEvictToLimitStopsAtNextCandidateOnceCancelled cancels the context
|
||||
// while the oldest of three source blobs is being evicted, and checks that
|
||||
// EvictToLimit then returns context.Canceled without evicting the other
|
||||
// two or logging a warning for either of them.
|
||||
func TestEvictToLimitStopsAtNextCandidateOnceCancelled(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cache, _ := newEvictionTestCache(t, 1)
|
||||
|
||||
var logBuf bytes.Buffer
|
||||
|
||||
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
|
||||
|
||||
hashes := []ContentHash{
|
||||
storeEvictionTestSource(t, cache, "cancel.example.com", "/a.jpg",
|
||||
bytes.Repeat([]byte{0x61}, 1000)),
|
||||
storeEvictionTestSource(t, cache, "cancel.example.com", "/b.jpg",
|
||||
bytes.Repeat([]byte{0x62}, 1000)),
|
||||
storeEvictionTestSource(t, cache, "cancel.example.com", "/c.jpg",
|
||||
bytes.Repeat([]byte{0x63}, 1000)),
|
||||
}
|
||||
|
||||
base := time.Now().Add(-time.Hour)
|
||||
|
||||
for i, hash := range hashes {
|
||||
setSourceLastAccessed(t, cache, hash, base.Add(time.Duration(i)*time.Minute))
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
defer cancel()
|
||||
|
||||
cache.evictSourceBlobTestHook = func(ContentHash) { cancel() }
|
||||
|
||||
err := cache.EvictToLimit(ctx)
|
||||
t.Logf("EvictToLimit() error = %v", err)
|
||||
t.Logf("log output: %s", logBuf.String())
|
||||
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Errorf("EvictToLimit() error = %v, want context.Canceled", err)
|
||||
}
|
||||
|
||||
if cache.srcContent.Exists(hashes[0]) {
|
||||
t.Errorf("source blob %s, evicted when the context was cancelled, "+
|
||||
"is still on disk", hashes[0])
|
||||
}
|
||||
|
||||
for _, hash := range hashes[1:] {
|
||||
if !cache.srcContent.Exists(hash) {
|
||||
t.Errorf("source blob %s was evicted after the context was cancelled", hash)
|
||||
}
|
||||
}
|
||||
|
||||
if strings.Contains(logBuf.String(), `"level":"WARN"`) {
|
||||
t.Errorf("EvictToLimit logged a warning after the context was cancelled")
|
||||
}
|
||||
|
||||
assertNoDanglingReferences(t, cache)
|
||||
}
|
||||
|
||||
// TestStopEvictionReturnsWhenItsContextEnds pauses an eviction pass where
|
||||
// cancellation cannot reach it, after a source blob's rows are deleted and
|
||||
// before its file is removed, and checks that StopEviction returns its
|
||||
// context's error when that context ends instead of waiting for the pass.
|
||||
// Once the pass goes on, the goroutine exits and no row points at a
|
||||
// missing file.
|
||||
func TestStopEvictionReturnsWhenItsContextEnds(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cache, _ := newEvictionTestCache(t, 1)
|
||||
|
||||
paused := make(chan struct{})
|
||||
resume := make(chan struct{})
|
||||
|
||||
cache.evictSourceBlobTestHook = func(ContentHash) {
|
||||
close(paused)
|
||||
<-resume
|
||||
}
|
||||
|
||||
hash := storeEvictionTestSource(t, cache, "stop.example.com", "/a.jpg",
|
||||
bytes.Repeat([]byte{0x61}, 1000))
|
||||
|
||||
cache.StartEviction(time.Hour)
|
||||
|
||||
select {
|
||||
case <-paused:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("the eviction pass never reached the source blob")
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 50*time.Millisecond)
|
||||
defer cancel()
|
||||
|
||||
err := cache.StopEviction(ctx)
|
||||
t.Logf("StopEviction() error = %v", err)
|
||||
|
||||
if !errors.Is(err, context.DeadlineExceeded) {
|
||||
t.Errorf("StopEviction() error = %v, want context.DeadlineExceeded", err)
|
||||
}
|
||||
|
||||
close(resume)
|
||||
|
||||
err = cache.StopEviction(t.Context())
|
||||
if err != nil {
|
||||
t.Fatalf("second StopEviction() error = %v, want nil", err)
|
||||
}
|
||||
|
||||
assertNoDanglingReferences(t, cache)
|
||||
|
||||
if cache.srcContent.Exists(hash) {
|
||||
t.Errorf("source blob %s is still on disk after its rows were deleted", hash)
|
||||
}
|
||||
}
|
||||
|
||||
// TestReconciliationWalksStopOnceCancelled checks that both directory
|
||||
// walks of a reconciliation pass return the context's error once it is
|
||||
// cancelled, leaving in place a stale temp file they would otherwise
|
||||
// remove.
|
||||
func TestReconciliationWalksStopOnceCancelled(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cache, _ := newEvictionTestCache(t, 1<<30)
|
||||
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
|
||||
staleTime := time.Now().Add(-2 * staleTempFileAge)
|
||||
|
||||
walks := map[string]func(context.Context) error{
|
||||
cache.variants.baseDir: cache.reconcileVariantFiles,
|
||||
cache.srcContent.baseDir: cache.reconcileSourceFiles,
|
||||
}
|
||||
|
||||
for dir, walk := range walks {
|
||||
tempFile := filepath.Join(dir, tempFilePrefix+"stale")
|
||||
|
||||
err := os.WriteFile(tempFile, []byte("partial"), 0o600)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to write temp file: %v", err)
|
||||
}
|
||||
|
||||
err = os.Chtimes(tempFile, staleTime, staleTime)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to backdate temp file: %v", err)
|
||||
}
|
||||
|
||||
err = walk(ctx)
|
||||
t.Logf("walk of %s: error = %v", dir, err)
|
||||
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Errorf("walk of %s: error = %v, want context.Canceled", dir, err)
|
||||
}
|
||||
|
||||
_, err = os.Stat(tempFile)
|
||||
if err != nil {
|
||||
t.Errorf("walk of %s went on after cancellation: %v", dir, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestReconciliationPassLogsNoWarningOnceCancelled checks that a
|
||||
// reconciliation pass run with an already cancelled context logs no
|
||||
// warning, so a periodic tick the loop takes after a stop adds no
|
||||
// warning to the one from the pass the stop interrupted.
|
||||
func TestReconciliationPassLogsNoWarningOnceCancelled(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cache, _ := newEvictionTestCache(t, 1<<30)
|
||||
|
||||
var logBuf bytes.Buffer
|
||||
|
||||
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
|
||||
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
|
||||
cache.runReconciliationPass(ctx)
|
||||
t.Logf("log output: %s", logBuf.String())
|
||||
|
||||
if strings.Contains(logBuf.String(), `"level":"WARN"`) {
|
||||
t.Errorf("runReconciliationPass logged a warning with a cancelled context")
|
||||
}
|
||||
}
|
||||
|
||||
// TestEvictSourceBlobExcludesConcurrentStoreOfIdenticalContent exercises
|
||||
// the exact TOCTOU window between evictSourceBlob's row-deletion
|
||||
// transaction commit and its content file unlink: a concurrent
|
||||
|
||||
Reference in New Issue
Block a user