Files
pixa/internal/handlers/image.go
T
clawbot 0233001110 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 06:51:08 +00:00

309 lines
8.3 KiB
Go

package handlers
import (
"errors"
"fmt"
"io"
"net/http"
"net/url"
"strconv"
"time"
"github.com/go-chi/chi/v5"
"sneak.berlin/go/pixa/internal/encurl"
"sneak.berlin/go/pixa/internal/httpfetcher"
"sneak.berlin/go/pixa/internal/imageprocessor"
"sneak.berlin/go/pixa/internal/imgcache"
)
// HandleImage handles the main image proxy route:
// /v1/image/<host>/<path>/<width>x<height>.<format>
func (s *Handlers) HandleImage() http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
req, ok := s.parseImageRequest(w, r)
if !ok {
return
}
// Validate signature if required
err := s.imgSvc.ValidateRequest(req)
if err != nil {
s.log.Warn("signature validation failed",
"host", req.SourceHost,
"path", req.SourcePath,
"error", err,
)
s.respondError(w, "unauthorized", http.StatusUnauthorized)
return
}
// Get cache key for logging
cacheKey := imgcache.CacheKey(req)
// Get the image (from cache or fetch/process)
startTime := time.Now()
resp, err := s.imgSvc.Get(r.Context(), req)
if err != nil {
s.respondImageError(w, req, err)
return
}
defer func() { _ = resp.Content.Close() }()
s.writeImageResponse(w, r, req, resp, cacheKey, startTime)
}
}
// HandleRobotsTxt serves robots.txt to prevent search engine crawling.
func (s *Handlers) HandleRobotsTxt() http.HandlerFunc {
robotsTxt := []byte("User-agent: *\nDisallow: /\n")
return func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "text/plain")
w.Header().Set("Content-Length", strconv.Itoa(len(robotsTxt)))
w.WriteHeader(http.StatusOK)
_, _ = w.Write(robotsTxt)
}
}
// parseImageRequest parses the wildcard path and query parameters into
// an ImageRequest. On invalid input it writes an error response and
// returns false.
func (s *Handlers) parseImageRequest(
w http.ResponseWriter, r *http.Request,
) (*imgcache.ImageRequest, bool) {
// Get the wildcard path from chi
pathParam := chi.URLParam(r, "*")
// Parse the URL path
parsed, err := imgcache.ParseImagePath(pathParam)
if err != nil {
s.log.Warn("failed to parse image URL",
"path", pathParam,
"error", err,
)
s.respondError(w, "invalid image URL: "+err.Error(), http.StatusBadRequest)
return nil, false
}
// Convert to ImageRequest
req := parsed.ToImageRequest()
// Parse signature params from query string. r.URL.Query() would silently
// drop a pair it cannot decode, such as q=80%, so that q would be served
// at 85; a query string that cannot be decoded is refused instead. A
// parameter given more than once is refused too, as only its first value
// would be read.
query, err := url.ParseQuery(r.URL.RawQuery)
if err != nil {
s.respondError(w, fmt.Sprintf("invalid query string %q: %v",
r.URL.RawQuery, err), http.StatusBadRequest)
return nil, false
}
for name, values := range query {
if len(values) > 1 {
s.respondError(w, fmt.Sprintf("invalid %s: given more than once",
name), http.StatusBadRequest)
return nil, false
}
}
req.Signature = query.Get("sig")
req.Expires, err = parseExpires(query)
if err != nil {
s.respondError(w, err.Error(), http.StatusBadRequest)
return nil, false
}
// Parse optional quality and fit params. Only a q missing from the URL is
// 85. A q in the URL that is not a whole number from 1 to 100, an empty
// one included, is refused, checked as the generator checks its quality
// field; that check alone would take an empty q as missing.
qStr := query.Get("q")
if query.Has("q") && qStr == "" {
s.respondError(w, `invalid q: not a number, got ""`,
http.StatusBadRequest)
return nil, false
}
req.Quality, err = parseFormInt(query, "q",
encurl.DefaultQuality, minQuality, maxQuality)
if err != nil {
s.respondError(w, fmt.Sprintf("%v, got %q", err, qStr),
http.StatusBadRequest)
return nil, false
}
// Only a fit missing from the URL is cover. A fit in the URL that is not a
// fit mode is refused by the fit-mode check below; that check would take an
// empty fit as missing, so an empty one is refused here.
req.FitMode = imgcache.FitMode(query.Get("fit"))
if query.Has("fit") && req.FitMode == "" {
s.respondError(w, `invalid fit: not a fit mode, got ""`, http.StatusBadRequest)
return nil, false
}
// Default fit mode if not set
if req.FitMode == "" {
req.FitMode = imgcache.FitCover
}
// Enforce dimension and fit-mode bounds, shared with the encrypted-URL
// route. Dimensions are already bounded by the path parser above; this
// also rejects an unrecognized fit mode with 400 instead of letting it
// reach the processor as a 500.
err = imgcache.ValidateImageRequest(req)
if err != nil {
s.respondError(w, "invalid image request: "+err.Error(),
http.StatusBadRequest)
return nil, false
}
return req, true
}
// parseExpires reads the exp query parameter, a Unix time in seconds. An exp
// missing from the URL gives the zero time, which the signature check takes
// as no expiration. An exp in the URL that is not a whole number, an empty
// one included, is an error naming exp and the value.
func parseExpires(query url.Values) (time.Time, error) {
if !query.Has("exp") {
return time.Time{}, nil
}
expStr := query.Get("exp")
exp, err := strconv.ParseInt(expStr, 10, 64)
if err != nil {
return time.Time{}, fmt.Errorf("%w exp: not a number, got %q",
errInvalidFormField, expStr)
}
return time.Unix(exp, 0), nil
}
// respondImageError maps image retrieval errors to HTTP responses.
func (s *Handlers) respondImageError(
w http.ResponseWriter, req *imgcache.ImageRequest, err error,
) {
s.log.Error("failed to get image",
"host", req.SourceHost,
"path", req.SourcePath,
"error", err,
)
// Check for specific error types
if errors.Is(err, httpfetcher.ErrSSRFBlocked) {
s.respondError(w, "forbidden", http.StatusForbidden)
return
}
if errors.Is(err, httpfetcher.ErrUpstreamError) {
s.respondError(w, "upstream error", http.StatusBadGateway)
return
}
if errors.Is(err, httpfetcher.ErrTooManyConnections) ||
errors.Is(err, imageprocessor.ErrTooManyImages) {
s.respondError(w, "server busy, try again later",
http.StatusServiceUnavailable)
return
}
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,
// handling conditional and HEAD requests.
func (s *Handlers) writeImageResponse(
w http.ResponseWriter, r *http.Request,
req *imgcache.ImageRequest, resp *imgcache.ImageResponse,
cacheKey imgcache.VariantKey, startTime time.Time,
) {
// Set response headers
w.Header().Set("Content-Type", resp.ContentType)
if resp.ContentLength > 0 {
w.Header().Set("Content-Length", strconv.FormatInt(resp.ContentLength, 10))
}
// Cache control headers
w.Header().Set("Cache-Control", cacheControl(req.Expires))
w.Header().Set("X-Pixa-Cache", string(resp.CacheStatus))
if resp.ETag != "" {
w.Header().Set("ETag", resp.ETag)
// Check for conditional request (If-None-Match)
if ifNoneMatch := r.Header.Get("If-None-Match"); ifNoneMatch != "" {
if ifNoneMatch == resp.ETag {
w.WriteHeader(http.StatusNotModified)
return
}
}
}
// Handle HEAD request - return headers only
if r.Method == http.MethodHead {
w.WriteHeader(http.StatusOK)
return
}
// Stream the response
w.WriteHeader(http.StatusOK)
servedBytes, err := io.Copy(w, resp.Content)
if err != nil {
s.log.Error("failed to write response",
"error", err,
)
}
// Log cache status and timing after serving
duration := time.Since(startTime)
s.log.Info("image served",
"cache_key", cacheKey,
"cache_status", resp.CacheStatus,
"duration_ms", duration.Milliseconds(),
"format", req.Format,
"served_bytes", servedBytes,
"fetched_bytes", resp.FetchedBytes,
)
}