8 Commits
Author SHA1 Message Date
clawbot 31ed20ec9e Count a late hit, and report no totals on a disabled cache (closes #56)
check / check (push) Successful in 3m12s
A hit is now counted with the request context detached from its
cancellation, as the miss already is, so a hit served after the client
left or the request timed out still moves the hit count. On a disabled
disk cache, Stats reports no items and no size, whatever rows an earlier
run left in the database; the size already followed that rule.

Model: opus-5-5
2026-09-29 02:46:31 +00:00
clawbot 3daa416f4b Test that a late hit and a disabled cache are counted right (closes #56)
A hit served after the request context has ended must still move the hit
count, and a disabled disk cache must report no items and no size even
when its database holds rows from an earlier run.

Model: opus-5-5
2026-09-29 02:46:31 +00:00
clawbot 614bcd9eb6 Count interrupted misses and the upstream bytes they read (closes #56)
The miss and transform counters are written with context.WithoutCancel,
so a client disconnect or the request timeout during or after the work
no longer loses them. A failed read of the upstream body now returns
the bytes read before the error, so an over-size or cut-off body still
moves the upstream fetch counters. processFromSourceOrFetch passes the
cached source's length directly instead of through a local named
fetchBytes.

Model: opus-5-5
2026-09-29 02:46:30 +00:00
clawbot 64108c05bb Test that cache stats count interrupted misses (closes #56)
Checks every cache_stats counter after a miss whose request context
ends during or after the upstream fetch, and after one whose upstream
body is over the size limit.

Model: opus-5-5
2026-09-29 02:46:30 +00:00
clawbot ae8b45e93f Make the cache stats count what is cached, fetched and transcoded (closes #56)
Stats read request_cache and output_content, which nothing writes, so
TotalItems and TotalSizeBytes were always 0. They now count
source_content plus variant_content, the size through UsageBytes; a
failed query is still logged at warn. Get counts a miss after the work,
passing the bytes fetched from upstream (0 for a cached source; still
counted when the fetched source then fails), so upstream_fetch_count and
upstream_fetch_bytes move. transform_count is incremented after each
successful image processor call. request_cache and output_content stay
in the schema; dropping them is a separate decision.

Model: opus-5-5
2026-09-29 02:46:30 +00:00
clawbot 19e018b037 Test that the cache stats counters and totals move (closes #56)
Failing tests, committed ahead of the fix. Stats totals are checked
after storing a source image and two processed variants. A walk through
Service.Get (a miss that fetches, a hit, a miss that reuses the cached
source, a source failing the magic byte check, a source not found)
checks every cache_stats counter after each step. The warn-log test for
the Stats queries now drops source_content and variant_content, the
tables Stats will read.

Model: opus-5-5
2026-09-29 02:46:30 +00:00
clawbot ed3f8770e6 Give pixad a fixed uid and gid 65532 (closes #151)
check / check (push) Successful in 20s
adduser took the first free uid, 1000, and the entrypoint gives a
bind-mounted /var/lib/pixa to pixad, so on the host a person's login
account ended up owning pixa's database and cache. The image now creates
the pixad group with gid 65532 and the pixad user with uid 65532, which
host login and system accounts do not use. The first-run step of
"Running under upaas" in README.md names the uid and gid.

Model: opus-5-5
2026-09-29 04:44:49 +02:00
clawbot 2afe61e301 Keep max-age within an expiring image URL's lifetime (closes #63)
check / check (push) Successful in 12s
Both image routes sent Cache-Control: public, max-age=31536000,
immutable unconditionally, so a browser or proxy could keep serving an
image for a year after its signed or encrypted URL had expired. max-age
is now the whole seconds left until the URL expires, never negative and
at most one year; a URL with no expiry keeps one year. The 304 answer
uses the same value. An encrypted URL's expiry now reaches
ImageRequest.Expires through ToImageRequest. immutable stays: freshness
now ends no later than the URL's expiry. README.md documents the header.

Model: opus-5-5
2026-09-29 03:51:58 +02:00
12 changed files with 656 additions and 53 deletions
+6 -2
View File
@@ -68,8 +68,12 @@ RUN apk add --no-cache \
COPY --from=builder /pixad /usr/local/bin/pixad COPY --from=builder /pixad /usr/local/bin/pixad
COPY deploy/docker-entrypoint.sh /usr/local/bin/docker-entrypoint.sh COPY deploy/docker-entrypoint.sh /usr/local/bin/docker-entrypoint.sh
# Create non-root user, config directory, and data directory # Create non-root user, config directory, and data directory. pixad
RUN adduser -D -H -s /sbin/nologin pixad && \ # gets uid and gid 65532, which host login and system accounts do not
# use: a bind-mounted /var/lib/pixa is given to pixad, and on the host
# it must not belong to a person's account.
RUN addgroup -g 65532 pixad && \
adduser -D -H -s /sbin/nologin -u 65532 -G pixad pixad && \
mkdir -p /var/lib/pixa /etc/pixa && \ mkdir -p /var/lib/pixa /etc/pixa && \
chown pixad:pixad /var/lib/pixa chown pixad:pixad /var/lib/pixa
+11 -2
View File
@@ -58,8 +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. It may be owned by root: the - **First run:** create the host directory, owned by root or by uid
container gives it to its `pixad` user when it starts. `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
@@ -100,6 +102,13 @@ than once, is refused with 400.
- `<format>`: one of `orig`, `png`, `jpeg`, `webp` - `<format>`: one of `orig`, `png`, `jpeg`, `webp`
- `<size>`: `orig` or `<width>x<height>` (e.g. `800x600`) - `<size>`: `orig` or `<width>x<height>` (e.g. `800x600`)
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),
`max-age` is the whole seconds left until then, at most one year, so no browser
or proxy cache keeps the image after pixa would refuse the URL. A URL with no
expiry gets one year. `immutable` only stops a client revalidating while its
copy is fresh.
The login form (`POST /`) is limited to 5 attempts per minute per client The login form (`POST /`) is limited to 5 attempts per minute per client
address, counting an IPv6 client by its /64; an attempt over the limit is address, counting an IPv6 client by its /64; an attempt over the limit is
refused with 429 and a `Retry-After` header. Behind a reverse proxy the client refused with 429 and a `Retry-After` header. Behind a reverse proxy the client
+24
View File
@@ -30,6 +30,30 @@ exhaustion
# Completed Steps # Completed Steps
- 2026-09-29 fixed uid and gid for `pixad` (closes #151): the image creates the
`pixad` group with gid 65532 and the `pixad` user with uid 65532, instead of
the first free uid 1000, so a bind-mounted `/var/lib/pixa` given to `pixad`
is not owned on the host by a person's login account; the first-run step of
"Running under upaas" in `README.md` names the uid and gid.
- 2026-09-29 `max-age` never outlives an expiring URL (closes #63): both image
routes build `Cache-Control` from the request's `Expires`, which an encrypted
URL's expiry now fills too; `max-age` is one year, or the whole seconds left
until the `exp` of a `/v1/image/` URL or the expiry of an encrypted URL when
that is sooner, never negative; an allowlisted host's URL that has an `exp`
follows it too; `immutable` stays, as freshness now ends at the expiry;
documented in `README.md`.
- 2026-09-28 cache stats report real numbers (closes #56): `Cache.Stats`
counts the cached source images and processed variants (`source_content`
plus `variant_content`) and takes their size from `Cache.UsageBytes`,
instead of reading `request_cache` and `output_content`, which nothing
writes; those two tables are left in the schema; a disabled disk cache
reports no items and no size. A hit is counted even when the request
context has ended. A miss is counted after it is served or fails, also
when the request context has ended by then, with the bytes it read from
upstream, so `upstream_fetch_count` and `upstream_fetch_bytes` move,
including for an upstream body that fails partway or a fetched source
that then fails the magic byte check; `transform_count` counts each image
the image processor transcodes.
- 2026-09-28 strip metadata from processed images (closes #82): every output is - 2026-09-28 strip metadata from processed images (closes #82): every output is
exported with govips' `StripMetadata`, so it carries no EXIF, XMP, IPTC or ICC exported with govips' `StripMetadata`, so it carries no EXIF, XMP, IPTC or ICC
profile; the image is first turned upright with `AutoRotate` (before sizes are profile; the image is first turned upright with `AutoRotate` (before sizes are
+8 -1
View File
@@ -103,7 +103,8 @@ func (g *Generator) Parse(token string) (*Payload, error) {
} }
// ToImageRequest converts the payload to an ImageRequest. // ToImageRequest converts the payload to an ImageRequest.
// Applies default values for omitted optional fields. // Applies default values for omitted optional fields. An ExpiresAt of 0, a URL
// that never expires, gives the zero Expires.
func (p *Payload) ToImageRequest() *imgcache.ImageRequest { func (p *Payload) ToImageRequest() *imgcache.ImageRequest {
format := p.Format format := p.Format
if format == "" { if format == "" {
@@ -120,6 +121,11 @@ func (p *Payload) ToImageRequest() *imgcache.ImageRequest {
fitMode = DefaultFitMode fitMode = DefaultFitMode
} }
var expires time.Time
if p.ExpiresAt != 0 {
expires = time.Unix(p.ExpiresAt, 0)
}
return &imgcache.ImageRequest{ return &imgcache.ImageRequest{
SourceHost: p.SourceHost, SourceHost: p.SourceHost,
SourcePath: p.SourcePath, SourcePath: p.SourcePath,
@@ -131,6 +137,7 @@ func (p *Payload) ToImageRequest() *imgcache.ImageRequest {
Format: format, Format: format,
Quality: quality, Quality: quality,
FitMode: fitMode, FitMode: fitMode,
Expires: expires,
} }
} }
+19 -1
View File
@@ -220,6 +220,24 @@ func (s *Handlers) respondImageError(
s.respondError(w, "internal error", http.StatusInternalServerError) s.respondError(w, "internal error", http.StatusInternalServerError)
} }
// cacheControl returns the Cache-Control header for an image served through a
// URL that expires at expires, or never when expires is the zero time. A cache
// may keep the image for a year, but not past the URL's expiry, after which
// pixa refuses the URL. The seconds left are rounded down and never negative.
// immutable only stops revalidation while the image is fresh, so it also ends
// at the expiry.
func cacheControl(expires time.Time) string {
const oneYear = 365 * 24 * time.Hour
maxAge := oneYear
if !expires.IsZero() {
maxAge = min(maxAge, max(time.Until(expires), 0))
}
return fmt.Sprintf("public, max-age=%d, immutable", int64(maxAge/time.Second))
}
// writeImageResponse writes headers and streams the image content, // writeImageResponse writes headers and streams the image content,
// handling conditional and HEAD requests. // handling conditional and HEAD requests.
func (s *Handlers) writeImageResponse( func (s *Handlers) writeImageResponse(
@@ -235,7 +253,7 @@ func (s *Handlers) writeImageResponse(
} }
// Cache control headers // Cache control headers
w.Header().Set("Cache-Control", "public, max-age=31536000, immutable") w.Header().Set("Cache-Control", cacheControl(req.Expires))
w.Header().Set("X-Pixa-Cache", string(resp.CacheStatus)) w.Header().Set("X-Pixa-Cache", string(resp.CacheStatus))
if resp.ETag != "" { if resp.ETag != "" {
@@ -0,0 +1,203 @@
package handlers
import (
"image/color"
"log/slog"
"net/http"
"net/http/httptest"
"strconv"
"strings"
"testing"
"testing/fstest"
"time"
"github.com/go-chi/chi/v5"
"sneak.berlin/go/pixa/internal/encurl"
"sneak.berlin/go/pixa/internal/imgcache"
)
// photoPath is the path of the JPEG that newSignedHostServer serves.
const photoPath = "/images/photo.jpg"
// newSignedHostServer returns a router for both image routes, and the Handlers
// behind it, whose fetcher serves a JPEG at photoPath on signedHost. signedHost
// is not on the allowlist, so a /v1/image/ URL for it is served only with a
// valid signature.
func newSignedHostServer(t *testing.T) (*Handlers, http.Handler) {
t.Helper()
cache, err := imgcache.NewCache(setupTestDB(t), imgcache.CacheConfig{
StateDir: t.TempDir(),
CacheTTL: time.Hour,
NegativeTTL: 5 * time.Minute,
})
if err != nil {
t.Fatalf("imgcache.NewCache() error = %v", err)
}
jpegData := generateTestJPEG(t, 100, 100, color.RGBA{255, 0, 0, 255})
svc, err := imgcache.NewService(&imgcache.ServiceConfig{
Cache: cache,
Fetcher: newMockFetcher(fstest.MapFS{
signedHost + photoPath: &fstest.MapFile{Data: jpegData},
}),
SigningKey: testSigningKey,
})
if err != nil {
t.Fatalf("imgcache.NewService() error = %v", err)
}
encGen, err := encurl.NewGenerator(testSigningKey)
if err != nil {
t.Fatalf("encurl.NewGenerator() error = %v", err)
}
h := &Handlers{
log: slog.New(slog.DiscardHandler),
imgSvc: svc,
encGen: encGen,
}
r := chi.NewRouter()
r.Get("/v1/image/*", h.HandleImage())
r.Get("/v1/e/{token}/*", h.HandleImageEnc())
return h, r
}
// getMaxAge sends a GET for target to srv, requires a 200, and returns the
// max-age of the response's Cache-Control header, which must read
// "public, max-age=<seconds>, immutable".
func getMaxAge(t *testing.T, srv http.Handler, target string) int {
t.Helper()
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, target, nil)
rec := httptest.NewRecorder()
srv.ServeHTTP(rec, req)
header := rec.Header().Get("Cache-Control")
t.Logf("GET %s: %d, Cache-Control: %s", target, rec.Code, header)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
}
value, hasPrefix := strings.CutPrefix(header, "public, max-age=")
value, hasSuffix := strings.CutSuffix(value, ", immutable")
maxAge, err := strconv.Atoi(value)
if !hasPrefix || !hasSuffix || err != nil {
t.Fatalf("Cache-Control = %q, want public, max-age=<seconds>, immutable",
header)
}
return maxAge
}
// TestHandleImage_SignedURL_MaxAgeEndsAtExp verifies that an image served
// through a signed URL expiring in 60 seconds may be cached for at most those
// 60 seconds. The lower bound of 50 shows the max-age is the time left, not 0.
func TestHandleImage_SignedURL_MaxAgeEndsAtExp(t *testing.T) {
t.Parallel()
h, srv := newSignedHostServer(t)
signedURL, err := h.imgSvc.GenerateSignedURL("", &imgcache.ImageRequest{
SourceHost: signedHost,
SourcePath: photoPath,
Size: imgcache.Size{Width: 50, Height: 50},
Format: imgcache.FormatJPEG,
}, time.Minute)
if err != nil {
t.Fatalf("GenerateSignedURL() error = %v", err)
}
maxAge := getMaxAge(t, srv, signedURL)
if maxAge < 50 || maxAge > 60 {
t.Errorf("max-age = %d, want 50 to 60", maxAge)
}
}
// TestHandleImage_AllowlistedHost_MaxAge verifies the max-age of an image from
// an allowlisted host, which is served without checking sig or exp. A URL with
// no exp may be cached for a year. A URL whose exp has passed is the one request
// that reaches the header after its expiry, and must get 0, never less.
func TestHandleImage_AllowlistedHost_MaxAge(t *testing.T) {
t.Parallel()
pastExp := strconv.FormatInt(time.Now().Add(-time.Hour).Unix(), 10)
tests := []struct {
name string
query string
wantMaxAge int
}{
{"no exp", "", 31536000},
{"exp already past", "?exp=" + pastExp, 0},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
fix := setupTestHandler(t)
r := chi.NewRouter()
r.Get("/v1/image/*", fix.handler.HandleImage())
maxAge := getMaxAge(t, r,
"/v1/image/"+fix.goodHost+"/images/photo.jpg/50x50.jpeg"+tt.query)
if maxAge != tt.wantMaxAge {
t.Errorf("max-age = %d, want %d", maxAge, tt.wantMaxAge)
}
})
}
}
// TestHandleImageEnc_MaxAge verifies that an image served through an encrypted
// URL with a 60 second TTL may be cached for at most those 60 seconds, that one
// with a two-year TTL may be cached for a year, and that one made without a
// TTL, which never expires, may be cached for a year.
func TestHandleImageEnc_MaxAge(t *testing.T) {
t.Parallel()
tests := []struct {
name string
expiresAt int64
wantAtLeast int
wantAtMost int
}{
{"60 second TTL", time.Now().Add(time.Minute).Unix(), 50, 60},
{"two-year TTL", time.Now().Add(2 * 365 * 24 * time.Hour).Unix(), 31536000, 31536000},
{"no TTL", 0, 31536000, 31536000},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
h, srv := newSignedHostServer(t)
token, err := h.encGen.Generate(&encurl.Payload{
SourceHost: signedHost,
SourcePath: photoPath,
Width: 50,
Height: 50,
Format: imgcache.FormatJPEG,
ExpiresAt: tt.expiresAt,
})
if err != nil {
t.Fatalf("Generate() error = %v", err)
}
maxAge := getMaxAge(t, srv, "/v1/e/"+token+"/img.jpg")
if maxAge < tt.wantAtLeast || maxAge > tt.wantAtMost {
t.Errorf("max-age = %d, want %d to %d",
maxAge, tt.wantAtLeast, tt.wantAtMost)
}
})
}
}
@@ -14,9 +14,9 @@ import (
) )
// signedHost is not on the allowlist setupTestHandler builds, so a request // signedHost is not on the allowlist setupTestHandler builds, so a request
// for it needs a valid signature. No image is served for it: a request that // for it needs a valid signature. setupTestHandler serves no image for it: a
// passes the signature check gets 502 from the failed fetch, and one that // request that passes the signature check gets 502 from the failed fetch, and
// fails the check gets 401. // one that fails the check gets 401.
const signedHost = "signed.example.com" const signedHost = "signed.example.com"
// getImage sends a GET for target to the image route of fix and returns the // getImage sends a GET for target to the image route of fix and returns the
+2 -2
View File
@@ -89,8 +89,8 @@ func (s *Handlers) HandleImageEnc() http.HandlerFunc {
w.Header().Set("Content-Length", strconv.FormatInt(resp.ContentLength, 10)) w.Header().Set("Content-Length", strconv.FormatInt(resp.ContentLength, 10))
} }
// Cache headers - encrypted URLs can be cached since they're immutable // Cache headers: max-age ends at the URL's expiry
w.Header().Set("Cache-Control", "public, max-age=31536000, immutable") w.Header().Set("Cache-Control", cacheControl(req.Expires))
w.Header().Set("X-Pixa-Cache", string(resp.CacheStatus)) w.Header().Set("X-Pixa-Cache", string(resp.CacheStatus))
// Stream the response // Stream the response
+27 -12
View File
@@ -419,19 +419,21 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
return nil, fmt.Errorf("failed to get cache stats: %w", err) return nil, fmt.Errorf("failed to get cache stats: %w", err)
} }
// Get actual item count and total size from content tables // Count and size the cached source images and processed variants. A
err = c.db.QueryRowContext(ctx, // disabled cache holds none, whatever rows an earlier run left.
`SELECT COUNT(*) FROM request_cache`, if !c.disabled {
).Scan(&stats.TotalItems) err = c.db.QueryRowContext(ctx, `
if err != nil { SELECT (SELECT COUNT(*) FROM source_content)
c.log.Warn("failed to count cache items for stats", "error", err) + (SELECT COUNT(*) FROM variant_content)
} `).Scan(&stats.TotalItems)
if err != nil {
c.log.Warn("failed to count cache items for stats", "error", err)
}
err = c.db.QueryRowContext(ctx, stats.TotalSizeBytes, err = c.UsageBytes(ctx)
`SELECT COALESCE(SUM(size_bytes), 0) FROM output_content`, if err != nil {
).Scan(&stats.TotalSizeBytes) c.log.Warn("failed to sum cache size for stats", "error", err)
if err != nil { }
c.log.Warn("failed to sum cache size for stats", "error", err)
} }
// Compute hit rate as a ratio // Compute hit rate as a ratio
@@ -481,6 +483,19 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
} }
} }
// IncrementTransformCount counts one image transcoded by the image processor.
func (c *Cache) IncrementTransformCount(ctx context.Context) {
_, err := c.db.ExecContext(ctx, `
UPDATE cache_stats
SET transform_count = transform_count + 1,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`)
if err != nil {
c.log.Warn("failed to count transform", "error", err)
}
}
// 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.
+6 -3
View File
@@ -95,7 +95,8 @@ type ImageRequest struct {
FitMode FitMode FitMode FitMode
// Signature is the HMAC signature for non-allowlisted hosts // Signature is the HMAC signature for non-allowlisted hosts
Signature string Signature string
// Expires is the signature expiration timestamp // Expires is when the URL expires: the exp of a signed URL, or the expiry
// of an encrypted URL; the zero time if it has none
Expires time.Time Expires time.Time
// AllowHTTP indicates whether HTTP (non-TLS) is allowed for this request // AllowHTTP indicates whether HTTP (non-TLS) is allowed for this request
AllowHTTP bool AllowHTTP bool
@@ -162,9 +163,11 @@ type ImageCache interface {
// CacheStats contains cache statistics // CacheStats contains cache statistics
type CacheStats struct { type CacheStats struct {
// TotalItems is the number of cached items // TotalItems is the number of cached source images plus processed
// variants
TotalItems int64 TotalItems int64
// TotalSizeBytes is the total size of cached content // TotalSizeBytes is the total size of cached source images and
// processed variants
TotalSizeBytes int64 TotalSizeBytes int64
// HitCount is the number of cache hits // HitCount is the number of cache hits
HitCount int64 HitCount int64
+32 -26
View File
@@ -143,7 +143,8 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
s.log.Error("failed to get cached variant", "key", result.CacheKey, "error", err) s.log.Error("failed to get cached variant", "key", result.CacheKey, "error", err)
// Fall through to re-process // Fall through to re-process
} else { } else {
s.cache.IncrementStats(ctx, true, 0) // Counted also when the request context has ended meanwhile
s.cache.IncrementStats(context.WithoutCancel(ctx), true, 0)
return &ImageResponse{ return &ImageResponse{
Content: reader, Content: reader,
@@ -155,12 +156,15 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
} }
} }
// Cache miss - check if we have source content cached // Cache miss - process the cached source or fetch it, then count the
// miss with the bytes it fetched from upstream, also when it failed or
// the request context has ended meanwhile
cacheKey := CacheKey(req) cacheKey := CacheKey(req)
s.cache.IncrementStats(ctx, false, 0) response, fetchedBytes, err := s.processFromSourceOrFetch(ctx, req, cacheKey)
s.cache.IncrementStats(context.WithoutCancel(ctx), false, fetchedBytes)
response, err := s.processFromSourceOrFetch(ctx, req, cacheKey)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@@ -268,22 +272,20 @@ func (s *Service) loadCachedSource(contentHash ContentHash) []byte {
} }
// processFromSourceOrFetch processes an image, using cached source content // processFromSourceOrFetch processes an image, using cached source content
// if available. // if available. It also returns the number of bytes fetched from upstream,
// as fetchAndProcess does, or 0 when the cached source was used.
func (s *Service) processFromSourceOrFetch( func (s *Service) processFromSourceOrFetch(
ctx context.Context, ctx context.Context,
req *ImageRequest, req *ImageRequest,
cacheKey VariantKey, cacheKey VariantKey,
) (*ImageResponse, error) { ) (*ImageResponse, int64, error) {
// Check if we have cached source content // Check if we have cached source content
contentHash, _, err := s.cache.LookupSource(ctx, req) contentHash, _, err := s.cache.LookupSource(ctx, req)
if err != nil { if err != nil {
s.log.Warn("source lookup failed", "error", err) s.log.Warn("source lookup failed", "error", err)
} }
var ( var sourceData []byte
sourceData []byte
fetchBytes int64
)
if contentHash != "" { if contentHash != "" {
s.log.Debug("using cached source", "hash", contentHash) s.log.Debug("using cached source", "hash", contentHash)
@@ -292,26 +294,25 @@ func (s *Service) processFromSourceOrFetch(
// Fetch from upstream if we don't have source data or it's empty // Fetch from upstream if we don't have source data or it's empty
if len(sourceData) == 0 { if len(sourceData) == 0 {
resp, err := s.fetchAndProcess(ctx, req, cacheKey) return s.fetchAndProcess(ctx, req, cacheKey)
if err != nil {
return nil, err
}
return resp, nil
} }
// Process using cached source // Process using cached source; nothing was fetched from upstream
fetchBytes = int64(len(sourceData)) resp, err := s.processAndStore(
ctx, req, cacheKey, sourceData, int64(len(sourceData)),
)
return s.processAndStore(ctx, req, cacheKey, sourceData, fetchBytes) return resp, 0, err
} }
// fetchAndProcess fetches from upstream, processes, and caches the result. // fetchAndProcess fetches from upstream, processes, and caches the result.
// It also returns the number of bytes read from upstream, including when
// reading the response or a later step fails.
func (s *Service) fetchAndProcess( func (s *Service) fetchAndProcess(
ctx context.Context, ctx context.Context,
req *ImageRequest, req *ImageRequest,
cacheKey VariantKey, cacheKey VariantKey,
) (*ImageResponse, error) { ) (*ImageResponse, int64, error) {
// Fetch from upstream // Fetch from upstream
sourceURL := req.SourceURL() sourceURL := req.SourceURL()
@@ -330,20 +331,20 @@ func (s *Service) fetchAndProcess(
} }
} }
return nil, fmt.Errorf("upstream fetch failed: %w", err) return nil, 0, fmt.Errorf("upstream fetch failed: %w", err)
} }
defer func() { _ = fetchResult.Content.Close() }() defer func() { _ = fetchResult.Content.Close() }()
// Read and validate the source content // Read and validate the source content
sourceData, err := io.ReadAll(fetchResult.Content) sourceData, err := io.ReadAll(fetchResult.Content)
fetchBytes := int64(len(sourceData))
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to read upstream response: %w", err) return nil, fetchBytes, fmt.Errorf("failed to read upstream response: %w", err)
} }
// Calculate download bitrate // Calculate download bitrate
fetchBytes := int64(len(sourceData))
var downloadRate string var downloadRate string
if fetchResult.FetchDurationMs > 0 { if fetchResult.FetchDurationMs > 0 {
@@ -368,7 +369,7 @@ func (s *Service) fetchAndProcess(
// Validate magic bytes match content type // Validate magic bytes match content type
err = magic.ValidateMagicBytes(sourceData, fetchResult.ContentType) err = magic.ValidateMagicBytes(sourceData, fetchResult.ContentType)
if err != nil { if err != nil {
return nil, fmt.Errorf("content validation failed: %w", err) return nil, fetchBytes, fmt.Errorf("content validation failed: %w", err)
} }
// Store source content // Store source content
@@ -378,7 +379,9 @@ func (s *Service) fetchAndProcess(
// Continue even if caching fails // Continue even if caching fails
} }
return s.processAndStore(ctx, req, cacheKey, sourceData, fetchBytes) resp, err := s.processAndStore(ctx, req, cacheKey, sourceData, fetchBytes)
return resp, fetchBytes, err
} }
// processAndStore processes an image and stores the result. // processAndStore processes an image and stores the result.
@@ -406,6 +409,9 @@ func (s *Service) processAndStore(
processDuration := time.Since(processStart) processDuration := time.Since(processStart)
// Counted also when the request context has ended meanwhile
s.cache.IncrementTransformCount(context.WithoutCancel(ctx))
// Read processed content // Read processed content
processedData, err := io.ReadAll(processResult.Content) processedData, err := io.ReadAll(processResult.Content)
_ = processResult.Content.Close() _ = processResult.Content.Close()
+315 -1
View File
@@ -4,6 +4,9 @@ import (
"bytes" "bytes"
"context" "context"
"database/sql" "database/sql"
"image/color"
"io"
"io/fs"
"log/slog" "log/slog"
"math" "math"
"strings" "strings"
@@ -11,6 +14,7 @@ import (
"time" "time"
"sneak.berlin/go/pixa/internal/database" "sneak.berlin/go/pixa/internal/database"
"sneak.berlin/go/pixa/internal/httpfetcher"
) )
func setupStatsTestDB(t *testing.T) *sql.DB { func setupStatsTestDB(t *testing.T) *sql.DB {
@@ -125,7 +129,7 @@ func TestStats_LogsFailedCountQueries(t *testing.T) {
} }
_, err = db.ExecContext(t.Context(), _, err = db.ExecContext(t.Context(),
`DROP TABLE request_cache; DROP TABLE output_content`) `DROP TABLE source_content; DROP TABLE variant_content`)
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -183,3 +187,313 @@ func TestIncrementStats_LogsFailedUpdates(t *testing.T) {
} }
} }
} }
// TestStats_TotalsCountSourcesAndVariants verifies that TotalItems and
// TotalSizeBytes cover the stored source images and processed variants.
func TestStats_TotalsCountSourcesAndVariants(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
storeEvictionTestSource(t, cache, testHostCDN, testPathCat,
bytes.Repeat([]byte{0xAA}, 1000))
storeEvictionTestVariant(t, cache, testVariantKeyOne,
bytes.Repeat([]byte{0xAB}, 500))
storeEvictionTestVariant(t, cache, testVariantKeyTwo,
bytes.Repeat([]byte{0xAC}, 250))
stats, err := cache.Stats(t.Context())
if err != nil {
t.Fatalf("Stats() error = %v", err)
}
if stats.TotalItems != 3 {
t.Errorf("TotalItems = %d, want 3 (1 source, 2 variants)", stats.TotalItems)
}
if stats.TotalSizeBytes != 1750 {
t.Errorf("TotalSizeBytes = %d, want 1750 (1000+500+250)",
stats.TotalSizeBytes)
}
}
// TestStats_DisabledCacheReportsNoItems verifies that a disabled disk cache
// reports no items and no size, even when its database still holds the
// rows of an earlier run with the disk cache enabled.
func TestStats_DisabledCacheReportsNoItems(t *testing.T) {
t.Parallel()
enabled, _ := newEvictionTestCache(t, 1<<30)
storeEvictionTestSource(t, enabled, testHostCDN, testPathCat,
bytes.Repeat([]byte{0xAA}, 1000))
storeEvictionTestVariant(t, enabled, testVariantKeyOne,
bytes.Repeat([]byte{0xAB}, 500))
disabled, err := NewCache(enabled.db, CacheConfig{
StateDir: t.TempDir(),
CacheTTL: time.Hour,
NegativeTTL: 5 * time.Minute,
DisableDiskCache: true,
})
if err != nil {
t.Fatal(err)
}
stats, err := disabled.Stats(t.Context())
if err != nil {
t.Fatalf("Stats() error = %v", err)
}
if stats.TotalItems != 0 || stats.TotalSizeBytes != 0 {
t.Errorf("TotalItems = %d, TotalSizeBytes = %d, want 0 and 0",
stats.TotalItems, stats.TotalSizeBytes)
}
}
// cacheStatsCounters holds the counters of the cache_stats row, in column
// order.
type cacheStatsCounters struct {
hitCount int64
missCount int64
upstreamFetchCount int64
upstreamFetchBytes int64
transformCount int64
}
// readCacheStatsCounters reads the counters of the cache_stats row.
func readCacheStatsCounters(t *testing.T, cache *Cache) cacheStatsCounters {
t.Helper()
var got cacheStatsCounters
err := cache.db.QueryRowContext(t.Context(), `
SELECT hit_count, miss_count, upstream_fetch_count,
upstream_fetch_bytes, transform_count
FROM cache_stats WHERE id = 1
`).Scan(&got.hitCount, &got.missCount, &got.upstreamFetchCount,
&got.upstreamFetchBytes, &got.transformCount)
if err != nil {
t.Fatalf("failed to read cache_stats: %v", err)
}
return got
}
// TestService_Get_CountsStats walks Get through a miss that fetches the
// source, a hit, a miss that reuses the cached source, and two misses whose
// source cannot be used, checking every cache_stats counter after each.
func TestService_Get_CountsStats(t *testing.T) {
t.Parallel()
svc, fixtures := SetupTestService(t)
// NewTestFS builds the same files the test service's fetcher serves.
testFS, _ := NewTestFS(t)
photo, err := fs.ReadFile(testFS, fixtures.GoodHostJPEG)
if err != nil {
t.Fatal(err)
}
fake, err := fs.ReadFile(testFS, fixtures.InvalidFile)
if err != nil {
t.Fatal(err)
}
photoBytes, fakeBytes := int64(len(photo)), int64(len(fake))
// want is hits, misses, upstream fetches, upstream bytes, transforms.
steps := []struct {
name string
path string
size int
wantErr bool
want cacheStatsCounters
}{
{"miss that fetches the source", testPathPhoto, 50, false,
cacheStatsCounters{0, 1, 1, photoBytes, 1}},
{"hit", testPathPhoto, 50, false,
cacheStatsCounters{1, 1, 1, photoBytes, 1}},
{"miss that reuses the cached source", testPathPhoto, 25, false,
cacheStatsCounters{1, 2, 1, photoBytes, 2}},
{"miss whose source fails the magic byte check", "/images/fake.jpg", 50, true,
cacheStatsCounters{1, 3, 2, photoBytes + fakeBytes, 2}},
{"miss whose source is not found", "/images/nonexistent.jpg", 50, true,
cacheStatsCounters{1, 4, 2, photoBytes + fakeBytes, 2}},
}
for _, step := range steps {
resp, err := svc.Get(t.Context(), &ImageRequest{
SourceHost: fixtures.GoodHost,
SourcePath: step.path,
Size: Size{Width: step.size, Height: step.size},
Format: FormatJPEG,
Quality: 85,
FitMode: FitCover,
})
if (err != nil) != step.wantErr {
t.Fatalf("%s: Get() error = %v, want error %t", step.name, err, step.wantErr)
}
if err == nil {
_ = resp.Content.Close()
}
got := readCacheStatsCounters(t, svc.cache)
if got != step.want {
t.Fatalf("after the %s: counters = %+v, want %+v", step.name, got, step.want)
}
}
}
// fakeUpstream answers every fetch with itself as a JPEG body. The body
// serves data, then calls cancel, when set, and returns err; io.EOF ends
// the body normally.
type fakeUpstream struct {
data *bytes.Reader
cancel context.CancelFunc
err error
}
func (u *fakeUpstream) Fetch(
context.Context, string,
) (*httpfetcher.FetchResult, error) {
return &httpfetcher.FetchResult{
Content: io.NopCloser(u),
ContentLength: -1,
ContentType: testContentTypeJPEG,
}, nil
}
func (u *fakeUpstream) Read(p []byte) (int, error) {
if u.data.Len() > 0 {
return u.data.Read(p)
}
if u.cancel != nil {
u.cancel()
}
return 0, u.err
}
// TestService_Get_CountsInterruptedMisses checks every cache_stats counter
// after a miss whose request context ends during or after the upstream
// fetch, and after a miss whose upstream body is over the size limit.
func TestService_Get_CountsInterruptedMisses(t *testing.T) {
t.Parallel()
photo := generateTestJPEG(t, 100, 100, color.RGBA{255, 0, 0, 255})
half := len(photo) / 2
// want is hits, misses, upstream fetches, upstream bytes, transforms.
tests := []struct {
name string
served int // bytes of the photo the upstream body serves
cancel bool // whether the body then ends the request context
readErr error // what the body then returns
wantErr bool
want cacheStatsCounters
}{
{"request context ends during the fetch", half, true, context.Canceled, true,
cacheStatsCounters{0, 1, 1, int64(half), 0}},
{"request context ends after the fetch", len(photo), true, io.EOF, false,
cacheStatsCounters{0, 1, 1, int64(len(photo)), 1}},
{"upstream body over the size limit", half, false,
httpfetcher.ErrResponseTooLarge, true,
cacheStatsCounters{0, 1, 1, int64(half), 0}},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
svc, fixtures := SetupTestService(t)
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
upstream := &fakeUpstream{
data: bytes.NewReader(photo[:tc.served]),
err: tc.readErr,
}
if tc.cancel {
upstream.cancel = cancel
}
svc.fetcher = upstream
resp, err := svc.Get(ctx, &ImageRequest{
SourceHost: fixtures.GoodHost,
SourcePath: testPathPhoto,
Size: Size{Width: 50, Height: 50},
Format: FormatJPEG,
Quality: 85,
FitMode: FitCover,
})
t.Logf("Get() error = %v", err)
if (err != nil) != tc.wantErr {
t.Fatalf("Get() error = %v, want error %t", err, tc.wantErr)
}
if err == nil {
_ = resp.Content.Close()
}
got := readCacheStatsCounters(t, svc.cache)
if got != tc.want {
t.Errorf("counters = %+v, want %+v", got, tc.want)
}
})
}
}
// TestService_Get_CountsHitAfterRequestEnds checks every cache_stats counter
// after a hit served with a request context that has already ended: only
// the hit count moves.
func TestService_Get_CountsHitAfterRequestEnds(t *testing.T) {
t.Parallel()
svc, fixtures := SetupTestService(t)
req := &ImageRequest{
SourceHost: fixtures.GoodHost,
SourcePath: testPathPhoto,
Size: Size{Width: 50, Height: 50},
Format: FormatJPEG,
Quality: 85,
FitMode: FitCover,
}
// A first request caches the variant.
resp, err := svc.Get(t.Context(), req)
if err != nil {
t.Fatalf("first Get() error = %v", err)
}
_ = resp.Content.Close()
want := readCacheStatsCounters(t, svc.cache)
want.hitCount++
ctx, cancel := context.WithCancel(t.Context())
cancel()
resp, err = svc.Get(ctx, req)
if err != nil {
t.Fatalf("Get() with an ended request context: error = %v", err)
}
_ = resp.Content.Close()
if resp.CacheStatus != CacheHit {
t.Fatalf("CacheStatus = %v, want %v", resp.CacheStatus, CacheHit)
}
got := readCacheStatsCounters(t, svc.cache)
if got != want {
t.Errorf("counters = %+v, want %+v", got, want)
}
}