image: add an optional public image processing gateway (#11593)

* image: add an optional public image processing gateway

* image: fix representation metadata and processing bounds

* image: restrict passthrough to non-executable media types

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* image: tighten source media-type validation

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* image: write passthrough body on the inner response writer

CodeQL still flagged the passthrough write: the content type was set on
the wrapper while the body reached w.ResponseWriter, so the validated
header could not be associated with the write. Set headers and copy the
body on the same inner writer.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* image: serve processed output on the inner response writer

* image: reject XML source types and unsafe conditional metadata

---------

Co-authored-by: zhaoyuchen <yc.zhao@yinzon.com>
Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
authored and GitHub committed 2026-10-05 19:49:02 +08:00
1 parent 7a19961074
commit 8c67b75190
8 files changed
+1795

No files matched your search

+1
View File
@@ -32,6 +32,7 @@ var Commands = []*Command{
cmdFix,
cmdFuse,
cmdIam,
cmdImage,
cmdMaster,
cmdMasterFollower,
cmdMini,
+85
View File
@@ -0,0 +1,85 @@
package command
import (
"net"
"net/http"
"os"
"strconv"
"time"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/images/gateway"
)
var cmdImage = &Command{
UsageLine: "image -source=http://localhost:8333/public-bucket -imgproxy=http://localhost:8080",
Short: "Start an on-demand gateway for public images",
Long: `Run an optional image gateway in front of public S3 objects. x-oss-process
supports aspect-preserving downscaling, absolute quality,Q_85, and JPEG/PNG/WebP.
Source access is checked before every cache read. Processed images remain in a
bounded memory cache and are never written to S3 or the Filer.
Encoding runs in a separate imgproxy service; configure its pixel, file size,
and process resource limits. This public endpoint does not accept S3 signatures
or read private objects.`,
}
var imageOptions struct {
source, imgproxy, bind *string
port, concurrency, maxDimension *int
cacheMB, sourceMB, resultMB *int64
timeout *time.Duration
}
// init registers the optional endpoint without changing S3, Filer, or Volume services.
func init() {
cmdImage.Run = runImage
imageOptions.source = cmdImage.Flag.String("source", "", "Fixed anonymous S3 HTTP(S) source URL, optionally with a bucket path")
imageOptions.imgproxy = cmdImage.Flag.String("imgproxy", "", "Separate imgproxy HTTP(S) service URL")
imageOptions.bind = cmdImage.Flag.String("ip.bind", "127.0.0.1", "Listen address")
imageOptions.port = cmdImage.Flag.Int("port", 8334, "HTTP port")
imageOptions.concurrency = cmdImage.Flag.Int("concurrency", 8, "Maximum concurrent requests and encoding jobs; excess requests return 429")
imageOptions.maxDimension = cmdImage.Flag.Int("maxDimension", 4096, "Maximum output dimension")
imageOptions.cacheMB = cmdImage.Flag.Int64("cacheCapacityMB", 64, "Processed image memory cache in MiB; 0 disables caching")
imageOptions.sourceMB = cmdImage.Flag.Int64("maxSourceMB", 25, "Maximum source image size in MiB")
imageOptions.resultMB = cmdImage.Flag.Int64("maxResultMB", 10, "Maximum processed image size in MiB")
imageOptions.timeout = cmdImage.Flag.Duration("timeout", 15*time.Second, "Separate timeout budgets for source metadata, shared encoding, and client writes")
}
// runImage starts the gateway, reading signing material from the environment rather than process arguments.
func runImage(cmd *Command, args []string) bool {
if *imageOptions.cacheMB < 0 || *imageOptions.cacheMB > 1<<20 ||
*imageOptions.sourceMB < 1 || *imageOptions.sourceMB > 1024 ||
*imageOptions.resultMB < 1 || *imageOptions.resultMB > 1024 {
glog.Errorf("Invalid image gateway capacity options")
return false
}
handler, err := gateway.New(gateway.Config{
Source: *imageOptions.source, Imgproxy: *imageOptions.imgproxy,
Key: os.Getenv("IMGPROXY_KEY"), Salt: os.Getenv("IMGPROXY_SALT"),
Concurrency: *imageOptions.concurrency, MaxDimension: *imageOptions.maxDimension,
CacheBytes: *imageOptions.cacheMB << 20, MaxSourceBytes: *imageOptions.sourceMB << 20,
MaxResultBytes: *imageOptions.resultMB << 20, Timeout: *imageOptions.timeout,
})
if err != nil {
glog.Errorf("Invalid image gateway configuration: %v", err)
return false
}
if *imageOptions.port < 1 || *imageOptions.port > 65535 {
glog.Errorf("Invalid image gateway port")
return false
}
// Bound the initial source check, shared encoding, and response write phases.
// ReadTimeout also limits draining of ignored request bodies after the handler returns.
server := &http.Server{
Addr: net.JoinHostPort(*imageOptions.bind, strconv.Itoa(*imageOptions.port)), Handler: handler,
ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 10 * time.Second,
WriteTimeout: 3 * *imageOptions.timeout,
IdleTimeout: 60 * time.Second, MaxHeaderBytes: 16 << 10,
}
glog.V(0).Infof("Image gateway listening on %s", server.Addr)
if err = server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
glog.Errorf("Image gateway failed: %v", err)
return false
}
return true
}
+159
View File
@@ -0,0 +1,159 @@
# On-demand image gateway
`weed image` is an optional public image endpoint. Original images remain in
SeaweedFS. On a cache miss, a separate imgproxy service resizes and encodes the
image; results enter a bounded in-process cache. Subsequent requests recheck
anonymous source access before reusing results. Cache eviction and process
restarts do not affect stored objects. No S3 objects, Filer entries, or application
file records are created for processed images.
```mermaid
flowchart LR
Browser --> CDN
CDN --> Gateway[weed image]
Gateway -->|Check source access and revision on every request| S3[SeaweedFS S3]
Gateway -->|Cache miss| imgproxy
imgproxy -->|Anonymous source read| S3
Gateway --> Cache[Bounded memory cache]
```
## Request interface
```text
/image.png?x-oss-process=image/resize,w_640/quality,Q_85/format,webp
/image.png?x-oss-process=image/resize,w_240,h_240,m_lfit,limit_1/format,webp
/image.png?x-oss-process=image/format,jpg/quality,Q_85
```
- `resize` accepts `w`, `h`, `m_lfit`, and `limit_1`. It preserves aspect ratio
and never enlarges the source. One dimension determines the other; two
dimensions specify a bounding box. Any unspecified axis uses `maxDimension`
as its bound, including format-only requests; neither output axis can exceed it.
- `quality,Q_1` through `quality,Q_100` specify absolute quality, defaulting to 85.
Aliyun's relative quality `q` has no equivalent here and returns 400.
- `format` accepts `jpg`, `jpeg`, `png`, and `webp`, defaulting to WebP.
- The default maximum output dimension is 4096. Repeated parameters, unknown
operations, cropping, watermarks, animation processing, automatic format
negotiation, and other OSS operations are unsupported. PNG output is lossless;
quality mainly affects JPEG/WebP. Animated sources use imgproxy's default
first-frame behavior.
- A single `versionId` may select a specific source version. Other query
parameters are rejected.
- GET, HEAD, Range, and conditional requests apply to the returned representation.
Processed images have independent SHA-256 ETags and byte lengths. Derived
responses use ETag validators rather than the source modification date, which
cannot distinguish overwrites within the same second. Without a processing
parameter the gateway reads the original image, including single byte ranges.
- Original GET and HEAD responses reject unsupported or malformed media types,
including HTML and all `image/*+xml` subtypes. Non-XML image types and octet-stream
objects retain their original `Content-Type`, including parameters. This prevents
uploaded executable documents from being served under the gateway's origin.
A 304 may omit `Content-Type`; an explicit unsupported type is still rejected
so conditional responses cannot reclassify previously cached bytes.
This implements a limited subset of OSS image processing parameters, rather than
full Aliyun OSS compatibility. It runs on a separate port and does not change
existing `weed s3`, Filer, or Volume services.
## Running the gateway
Start a separate imgproxy service, then run:
```sh
weed image \
-source=http://s3:8333/public-bucket \
-imgproxy=http://imgproxy:8080 \
-ip.bind=0.0.0.0 -port=8334 \
-cacheCapacityMB=64 -concurrency=8 \
-maxSourceMB=25 -maxResultMB=10 -maxDimension=4096 -timeout=15s
```
`source` fixes the source HTTP(S) URL, optionally including a bucket path or a
bucket domain pointing to S3. With `http://s3:8333/public-bucket`, a client request
for `/a/b.png` reads `http://s3:8333/public-bucket/a/b.png`. Both backend URLs must
omit credentials, queries, and fragments. The imgproxy URL must be a root URL
without a path prefix; prefixed URLs are rejected at startup. Dot segments and backslashes, including
repeatedly escaped forms, are rejected to prevent backend path normalization from
escaping a fixed bucket prefix. The gateway and imgproxy must reach the same source.
Configure imgproxy limits and isolate the encoder with container or process
resource limits. Suggested imgproxy settings:
```text
IMGPROXY_WORKERS=2
IMGPROXY_REQUESTS_QUEUE_SIZE=8
IMGPROXY_MAX_SRC_FILE_SIZE=26214400
IMGPROXY_MAX_SRC_RESOLUTION=25
IMGPROXY_MAX_RESULT_DIMENSION=4096
IMGPROXY_MAX_ANIMATION_FRAMES=1
IMGPROXY_MAX_REDIRECTS=0
IMGPROXY_ALLOWED_PROCESSING_OPTIONS=rs,q,f
IMGPROXY_ALLOWED_SOURCES=http://s3:8333/public-bucket/
```
These variables apply to imgproxy 3.x. In imgproxy 4.x, some source security limits
move to separate source configuration; configure equivalent restrictions according
to that version's documentation. File size limits do not replace pixel and memory
limits. The gateway does not decode images and cannot enforce encoder-side limits.
Set hexadecimal `IMGPROXY_KEY` and `IMGPROXY_SALT` in both services to sign backend
requests. The gateway supports one key/salt pair with full SHA-256 signatures;
imgproxy must use the default 32-byte signature length. Unsigned requests are only
appropriate on an isolated trusted network. Keep signing material out of command
lines, public configuration, and logs.
## Access and caching
The endpoint only accepts anonymous public reads. It rejects client Authorization,
S3 signature parameters, and `x-amz-*` request headers. Browser cookies are ignored
and never forwarded to backends. The gateway does not create administrative
credentials or signed source URLs for private objects. Use the existing S3 endpoint
for private images.
Every derived request performs an anonymous source HEAD, including cache hits,
HEAD, and 304 requests. The encoder uses anonymous GET, so the source must apply
the same public-read policy to HEAD and GET; SeaweedFS S3 authorizes both as object
reads. A source 403 or 404 rejects derived requests even if cached bytes remain.
Source URL, ETag, modification time, length, version ID, and canonical processing
options form the cache key. Sources without ETags are not cached. A revision
change detected during encoding returns 409, allowing a client retry. Use immutable
object names for unversioned sources to avoid races between metadata checks and
frequent source overwrites.
Concurrent misses for the same result share one encoding job. A cancelled waiter
does not cancel work needed by other waiters. If all waiters leave, the job may run
until its independent processing timeout; separate work tokens still bound such
background jobs. Source metadata checks, shared encoding (including its final
source recheck), and client writes each receive a full timeout budget. A slow
initial HEAD does not shorten the caller's wait for shared encoding. The HTTP
server also bounds request body reads to 10 seconds. Requests, result sizes,
duration, and concurrency are bounded. Original
single-range responses are limited by the complete object's size, not just the
selected range. Excess concurrency returns 429; errors are never cached. The cache
is limited by bytes and 1024 entries, defaults to 64 MiB, and can be disabled with
`-cacheCapacityMB=0`. Each instance has its own cache, empty after restart.
Responses default to `Cache-Control: no-cache`, allowing downstream storage with
mandatory revalidation so source access is checked on every request. A CDN must
preserve `x-oss-process` and `versionId` and include the complete query in its cache
key. If you explicitly set a CDN TTL for permanently public immutable objects,
CDN hits bypass the gateway. Revocation or deletion then requires a CDN purge,
otherwise access changes only take effect after TTL expiry. Errors use `no-store`.
Configure HTTPS, CORS, and external rate limits in the existing reverse proxy or
CDN. Route only public image GET/HEAD requests to this gateway. Uploads, listings,
signed reads, and other S3 operations continue using the existing S3 endpoint.
## Validation
```sh
go test -race ./weed/images/gateway
go vet ./weed/images/gateway
CGO_ENABLED=0 go build ./weed
```
Tests cover revocation and deletion, overwrite invalidation, explicit versions,
shared work and cancellation, resource limits, backend failures, signatures and
path escaping, and HEAD/ETag/Range semantics. Before deployment, also test real
SeaweedFS and imgproxy for dimensions, media types, transparency, and regeneration
after clearing the cache.
+65
View File
@@ -0,0 +1,65 @@
package gateway
import (
"container/list"
"sync"
)
type result struct {
data []byte
contentType, etag string
}
type cacheEntry struct {
key string
value *result
size int64
}
// cache is a byte-bounded LRU with an entry limit to bound metadata for small objects.
type cache struct {
mu sync.Mutex
entries map[string]*list.Element
order *list.List
capacity, used int64
}
// newCache creates a disposable in-process cache without creating storage objects.
func newCache(capacity int64) *cache {
return &cache{entries: make(map[string]*list.Element), order: list.New(), capacity: capacity}
}
// get returns immutable results and updates recency.
func (c *cache) get(key string) *result {
c.mu.Lock()
defer c.mu.Unlock()
if e := c.entries[key]; e != nil {
c.order.MoveToFront(e)
return e.Value.(*cacheEntry).value
}
return nil
}
// put accounts for data and index overhead, skipping results larger than capacity.
func (c *cache) put(key string, value *result) {
c.mu.Lock()
defer c.mu.Unlock()
size := int64(len(key) + len(value.data) + len(value.contentType) + len(value.etag) + 128)
if c.capacity == 0 || size > c.capacity {
return
}
if old := c.entries[key]; old != nil {
c.used -= old.Value.(*cacheEntry).size
c.order.Remove(old)
delete(c.entries, key)
}
for c.used+size > c.capacity || len(c.entries) >= 1024 {
old := c.order.Back()
entry := old.Value.(*cacheEntry)
delete(c.entries, entry.key)
c.used -= entry.size
c.order.Remove(old)
}
c.entries[key] = c.order.PushFront(&cacheEntry{key: key, value: value, size: size})
c.used += size
}
+457
View File
@@ -0,0 +1,457 @@
// Package gateway provides an optional on-demand gateway for public S3 images.
package gateway
import (
"bytes"
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/base64"
"encoding/hex"
"errors"
"fmt"
"io"
"mime"
"net/http"
"net/url"
"strconv"
"strings"
"time"
"golang.org/x/sync/singleflight"
)
type Config struct {
Source, Imgproxy, Key, Salt string
Concurrency, MaxDimension int
CacheBytes, MaxSourceBytes, MaxResultBytes int64
Timeout time.Duration
}
type Gateway struct {
source, processor *url.URL
key, salt []byte
config Config
client *http.Client
requests chan struct{}
work chan struct{}
cache *cache
flights singleflight.Group
}
type failure struct {
status int
message string
}
type responseWriter struct {
http.ResponseWriter
timeout time.Duration
started bool
}
// beginWrite starts a separate client write budget after backend processing.
func (w *responseWriter) beginWrite() {
if !w.started {
w.started = true
_ = http.NewResponseController(w.ResponseWriter).SetWriteDeadline(time.Now().Add(w.timeout))
}
}
// WriteHeader starts the write budget and prevents caching protocol errors, including 412 and 416.
// Every response path commits its status before writing bytes through the embedded writer.
func (w *responseWriter) WriteHeader(status int) {
w.beginWrite()
if status >= 400 {
w.Header().Set("Cache-Control", "no-store")
}
w.ResponseWriter.WriteHeader(status)
}
// Error returns a fixed message without source URLs or signing material.
func (f *failure) Error() string { return f.message }
// New validates fixed backends and resource limits; redirects are disabled.
func New(c Config) (*Gateway, error) {
if c.Concurrency < 1 || c.Concurrency > 1024 || c.MaxDimension < 1 || c.MaxDimension > 16384 ||
c.CacheBytes < 0 || c.CacheBytes > 1<<40 || c.MaxSourceBytes < 1 || c.MaxSourceBytes > 1<<30 ||
c.MaxResultBytes < 1 || c.MaxResultBytes > 1<<30 || c.Timeout <= 0 {
return nil, fmt.Errorf("invalid resource limits")
}
parse := func(value string) (*url.URL, error) {
u, err := url.Parse(value)
if err != nil || (u.Scheme != "http" && u.Scheme != "https") || u.Host == "" ||
u.User != nil || u.RawQuery != "" || u.ForceQuery || u.Fragment != "" {
return nil, fmt.Errorf("backends must be HTTP(S) URLs without credentials, queries, or fragments")
}
return u, nil
}
source, err := parse(c.Source)
if err != nil {
return nil, err
}
processor, err := parse(c.Imgproxy)
if err != nil {
return nil, err
}
// imgproxy signs its root-relative processing path; reject ambiguous proxy prefixes.
if processor.Path != "" && processor.Path != "/" || processor.RawPath != "" {
return nil, fmt.Errorf("imgproxy URL must not contain a path prefix")
}
key, keyErr := hex.DecodeString(c.Key)
salt, saltErr := hex.DecodeString(c.Salt)
if keyErr != nil || saltErr != nil || (len(key) == 0) != (len(salt) == 0) {
return nil, fmt.Errorf("imgproxy key and salt must both be provided as hexadecimal values")
}
transport := http.DefaultTransport.(*http.Transport).Clone()
transport.DisableCompression = true
return &Gateway{
source: source, processor: processor, key: key, salt: salt, config: c,
client: &http.Client{Transport: transport, CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }},
requests: make(chan struct{}, c.Concurrency), cache: newCache(c.CacheBytes),
work: make(chan struct{}, c.Concurrency),
}, nil
}
// sourceURL joins the fixed source and object path, ignoring client hosts and credentials.
func (g *Gateway) sourceURL(r *http.Request, query url.Values) *url.URL {
u := *g.source
u.Path = strings.TrimSuffix(u.Path, "/") + r.URL.Path
u.RawPath = strings.TrimSuffix(g.source.EscapedPath(), "/") + r.URL.EscapedPath()
if id := query.Get("versionId"); id != "" {
u.RawQuery = url.Values{"versionId": {id}}.Encode()
}
return &u
}
// ServeHTTP checks anonymous source access before cached or conditional responses.
func (g *Gateway) ServeHTTP(w http.ResponseWriter, r *http.Request) {
wrapped := &responseWriter{ResponseWriter: w, timeout: g.config.Timeout}
w = wrapped
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("X-Content-Type-Options", "nosniff")
if r.Method != http.MethodGet && r.Method != http.MethodHead {
w.Header().Set("Allow", "GET, HEAD")
g.writeError(w, r, &failure{405, "only GET and HEAD are allowed"})
return
}
if !safeObjectPath(r.URL.Path) {
g.writeError(w, r, &failure{400, "object path contains unsafe normalization segments"})
return
}
query, err := url.ParseQuery(r.URL.RawQuery)
if err != nil {
g.writeError(w, r, &failure{400, "invalid query parameters"})
return
}
for name, values := range query {
if (name != "x-oss-process" && name != "versionId") || len(values) != 1 || len(values[0]) > 1024 {
g.writeError(w, r, &failure{400, "unsupported or repeated query parameters"})
return
}
}
if r.Header.Get("Authorization") != "" {
g.writeError(w, r, &failure{403, "this endpoint only serves anonymously readable public images"})
return
}
for name := range r.Header {
if strings.HasPrefix(strings.ToLower(name), "x-amz-") {
g.writeError(w, r, &failure{403, "this endpoint does not accept S3 credentials or encryption headers"})
return
}
}
var o options
_, processing := query["x-oss-process"]
if processing {
o, err = parseOptions(query.Get("x-oss-process"), g.config.MaxDimension)
if err != nil {
g.writeError(w, r, &failure{400, err.Error()})
return
}
}
select {
case g.requests <- struct{}{}:
defer func() { <-g.requests }()
default:
g.writeError(w, r, &failure{429, "image request concurrency limit reached"})
return
}
ctx, cancel := context.WithTimeout(r.Context(), g.config.Timeout)
defer cancel()
source := g.sourceURL(r, query)
if !processing {
g.serveOriginal(wrapped, r.WithContext(ctx), source)
return
}
metadata, err := g.head(ctx, source)
if err != nil {
g.writeError(w, r, err)
return
}
// Only sources with validators can reuse results; access is still checked on every request.
cacheKey := source.String() + "\n" + revision(metadata) + "\n" + o.path()
cacheable := metadata.Get("ETag") != ""
var image *result
if cacheable {
image = g.cache.get(cacheKey)
}
if image == nil {
// Source metadata and shared work each have their own bounded phase.
// A slow HEAD must not consume the time needed to await a valid encoding job.
waitCtx, waitCancel := context.WithTimeout(r.Context(), g.config.Timeout)
defer waitCancel()
cancel()
flight := g.flights.DoChan(cacheKey, func() (interface{}, error) {
if cacheable {
if hit := g.cache.get(cacheKey); hit != nil {
return hit, nil
}
}
// Shared work is independent of any waiter and remains bounded by its own timeout.
// Separate tokens prevent cancelled requests from bypassing the work concurrency limit.
select {
case g.work <- struct{}{}:
defer func() { <-g.work }()
default:
return nil, &failure{429, "image encoding concurrency limit reached"}
}
workCtx, workCancel := context.WithTimeout(context.Background(), g.config.Timeout)
defer workCancel()
processed, processErr := g.process(workCtx, source, o)
if processErr != nil {
return nil, processErr
}
// Recheck access and revision before caching results from an in-flight encoding.
current, checkErr := g.head(workCtx, source)
if checkErr != nil {
return nil, checkErr
}
if revision(current) != revision(metadata) {
return nil, &failure{409, "source changed during processing; retry the request"}
}
if cacheable {
g.cache.put(cacheKey, processed)
}
return processed, nil
})
select {
case outcome := <-flight:
if outcome.Err != nil {
g.writeError(w, r, outcome.Err)
return
}
image = outcome.Val.(*result)
case <-waitCtx.Done():
g.writeError(w, r, &failure{504, "image processing timed out or request was cancelled"})
return
}
}
// Keep the status hook for protocol errors; body writes use the embedded writer directly.
// Attach output metadata to that same writer before ServeContent sends the response.
rw := wrapped
rw.Header().Set("Content-Type", image.contentType)
rw.Header().Set("ETag", image.etag)
wrapped.beginWrite()
// A second-resolution source date cannot distinguish overwrites; use the output ETag only.
http.ServeContent(rw, r, "", time.Time{}, bytes.NewReader(image.data))
}
// safeObjectPath rejects paths that a backend proxy could normalize outside the fixed bucket prefix.
// Check repeatedly escaped dot segments and backslashes while preserving ordinary object escaping.
func safeObjectPath(value string) bool {
for i := 0; i < 4; i++ {
if strings.Contains(value, "\\") {
return false
}
for _, part := range strings.Split(value, "/") {
if part == "." || part == ".." {
return false
}
}
next, err := url.PathUnescape(value)
if err != nil || next == value {
return true
}
value = next
}
return false
}
// allowedSourceType excludes XML documents, including custom image subtypes that can contain scripts.
func allowedSourceType(mediaType string) bool {
if strings.HasSuffix(mediaType, "+xml") {
return false
}
return strings.HasPrefix(mediaType, "image/") || mediaType == "application/octet-stream" || mediaType == "binary/octet-stream"
}
// revision includes S3 version and content validators in the cache key.
func revision(h http.Header) string {
return h.Get("ETag") + "\n" + h.Get("Last-Modified") + "\n" + h.Get("Content-Length") + "\n" + h.Get("x-amz-version-id")
}
// read makes credential-free requests with the supplied timeout and cancellation context.
func (g *Gateway) read(ctx context.Context, method string, u *url.URL, headers http.Header) (*http.Response, error) {
req, err := http.NewRequestWithContext(ctx, method, u.String(), nil)
if err != nil {
return nil, &failure{502, "could not create backend request"}
}
if headers != nil {
req.Header = headers.Clone()
}
resp, err := g.client.Do(req)
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
return nil, &failure{504, "image backend request timed out"}
}
return nil, &failure{502, "image backend request failed"}
}
return resp, nil
}
// head checks anonymous access and source size without forwarding client conditions.
func (g *Gateway) head(ctx context.Context, source *url.URL) (http.Header, error) {
resp, err := g.read(ctx, http.MethodHead, source, nil)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode == 403 || resp.StatusCode == 404 {
return nil, &failure{resp.StatusCode, "source image is not accessible"}
}
if resp.StatusCode != 200 {
return nil, &failure{502, "source image check failed"}
}
if resp.ContentLength < 0 || resp.ContentLength > g.config.MaxSourceBytes {
return nil, &failure{413, "source image exceeds size limit or lacks a content length"}
}
return resp.Header, nil
}
// processorURL builds a fixed-operation imgproxy URL, optionally signed with its key and salt.
func (g *Gateway) processorURL(source *url.URL, o options) *url.URL {
path := o.path() + "/" + base64.RawURLEncoding.EncodeToString([]byte(source.String())) + "." + o.format
signature := "insecure"
if len(g.key) != 0 {
mac := hmac.New(sha256.New, g.key)
mac.Write(g.salt)
mac.Write([]byte(path))
signature = base64.RawURLEncoding.EncodeToString(mac.Sum(nil))
}
u := *g.processor
u.Path = strings.TrimSuffix(u.Path, "/") + "/" + signature + path
u.RawPath = ""
return &u
}
// process reads bounded results, rejecting incorrect media types and unsuccessful responses.
func (g *Gateway) process(ctx context.Context, source *url.URL, o options) (*result, error) {
resp, err := g.read(ctx, http.MethodGet, g.processorURL(source, o), nil)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
return nil, &failure{502, "image encoder returned an error"}
}
contentType, _, err := mime.ParseMediaType(resp.Header.Get("Content-Type"))
expected := map[string]string{"jpg": "image/jpeg", "png": "image/png", "webp": "image/webp"}[o.format]
if err != nil || contentType != expected {
return nil, &failure{502, "image encoder returned an unexpected media type"}
}
data, err := io.ReadAll(io.LimitReader(resp.Body, g.config.MaxResultBytes+1))
if err != nil {
return nil, &failure{502, "could not read processed image"}
}
if len(data) == 0 || int64(len(data)) > g.config.MaxResultBytes {
return nil, &failure{413, "processed image is empty or exceeds size limit"}
}
digest := sha256.Sum256(data)
return &result{data: data, contentType: contentType, etag: fmt.Sprintf("\"%x\"", digest)}, nil
}
// serveOriginal preserves public source conditions and ranges without forwarding credentials.
func (g *Gateway) serveOriginal(w *responseWriter, r *http.Request, source *url.URL) {
headers := make(http.Header)
for _, key := range []string{"Range", "If-Range", "If-Match", "If-None-Match", "If-Modified-Since", "If-Unmodified-Since"} {
if value := r.Header.Get(key); value != "" {
headers.Set(key, value)
}
}
resp, err := g.read(r.Context(), r.Method, source, headers)
if err != nil {
g.writeError(w, r, err)
return
}
defer resp.Body.Close()
if resp.StatusCode != 200 && resp.StatusCode != 206 && resp.StatusCode != 304 {
status := resp.StatusCode
if status < 400 || status >= 500 {
status = 502
}
if status == 416 {
w.Header().Set("Content-Range", resp.Header.Get("Content-Range"))
}
g.writeError(w, r, &failure{status, "source image read failed"})
return
}
sourceSize := resp.ContentLength
if resp.StatusCode == 206 {
// A 206 length describes only the selected range; enforce limits against the complete object.
rangeValue := resp.Header.Get("Content-Range")
_, total, ok := strings.Cut(rangeValue, "/")
sourceSize, err = strconv.ParseInt(total, 10, 64)
if !ok || !strings.HasPrefix(rangeValue, "bytes ") || err != nil || sourceSize < 1 {
g.writeError(w, r, &failure{502, "source range response lacks a valid complete length"})
return
}
}
if sourceSize > g.config.MaxSourceBytes || resp.ContentLength > g.config.MaxSourceBytes {
g.writeError(w, r, &failure{413, "source image exceeds size limit"})
return
}
// Serving executable documents under this origin allows cross-site scripting.
contentType := resp.Header.Get("Content-Type")
mediaType, _, mediaErr := mime.ParseMediaType(contentType)
// A 304 may omit this header, but must not replace cached metadata with an executable type.
if (resp.StatusCode != 304 || contentType != "") && (mediaErr != nil || !allowedSourceType(mediaType)) {
g.writeError(w, r, &failure{502, "source image has an unsupported media type"})
return
}
var data []byte
if r.Method != http.MethodHead && resp.StatusCode != 304 {
data, err = io.ReadAll(io.LimitReader(resp.Body, g.config.MaxSourceBytes+1))
if err != nil || int64(len(data)) > g.config.MaxSourceBytes {
g.writeError(w, r, &failure{502, "source read failed or exceeded size limit"})
return
}
}
// Write through the inner ResponseWriter so the validated content type is
// visibly attached to the writer that receives the body.
rw := w.ResponseWriter
if contentType != "" {
rw.Header().Set("Content-Type", contentType)
}
for _, key := range []string{"Content-Encoding", "Content-Length", "ETag", "Last-Modified", "Accept-Ranges", "Content-Range", "x-amz-version-id"} {
if value := resp.Header.Get(key); value != "" {
rw.Header().Set(key, value)
}
}
w.beginWrite()
rw.WriteHeader(resp.StatusCode)
if len(data) != 0 {
_, _ = io.Copy(rw, bytes.NewReader(data))
}
}
// writeError prevents downstream caching of access and transient errors.
func (g *Gateway) writeError(w http.ResponseWriter, r *http.Request, err error) {
f, ok := err.(*failure)
if !ok {
f = &failure{502, "image processing failed"}
}
w.Header().Set("Cache-Control", "no-store")
w.Header().Del("ETag")
w.Header().Del("Content-Length")
http.Error(w, f.message, f.status)
}
+803
View File
@@ -0,0 +1,803 @@
package gateway
import (
"bytes"
"compress/gzip"
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/base64"
"fmt"
"io"
"net/http"
"net/http/httptest"
"net/url"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
)
const transform = "image/resize,w_640/quality,Q_85/format,webp"
// fixture models S3 access, revisions, and an external encoder for HTTP-level tests.
type fixture struct {
origin, processor *httptest.Server
gateway *Gateway
mu sync.Mutex
status int
etag, version string
contentType string
sourceBytes int64
processStatus int
processType string
data []byte
original []byte
sources []string
heads atomic.Int32
encodes atomic.Int32
started, release chan struct{}
headStarted, headRelease chan struct{}
headPaused atomic.Bool
}
// newFixture creates HTTP backends that fail tests if credentials are forwarded.
func newFixture(t *testing.T) *fixture {
t.Helper()
f := &fixture{status: 200, etag: "\"source-v1\"", sourceBytes: 100, contentType: "image/png", processStatus: 200,
processType: "image/webp", data: []byte("RIFF transformed webp bytes"), original: []byte("original-image")}
f.origin = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
// Pause only the initial source check; the post-encoding recheck remains fast.
if r.Method == http.MethodHead && f.headRelease != nil && f.headPaused.CompareAndSwap(false, true) {
close(f.headStarted)
select {
case <-f.headRelease:
case <-r.Context().Done():
return
}
}
f.mu.Lock()
defer f.mu.Unlock()
if r.Header.Get("Authorization") != "" || r.Header.Get("Cookie") != "" || r.URL.Query().Get("x-oss-process") != "" {
t.Error("source request contained credentials or processing parameters")
}
if r.Method == http.MethodHead {
f.heads.Add(1)
w.Header().Set("Content-Type", f.contentType)
w.Header().Set("ETag", f.etag)
w.Header().Set("Content-Length", fmt.Sprint(f.sourceBytes))
w.Header().Set("Last-Modified", "Sun, 04 Oct 2026 10:00:00 GMT")
if f.version != "" {
w.Header().Set("x-amz-version-id", f.version)
}
w.WriteHeader(f.status)
return
}
if f.status != 200 {
w.WriteHeader(f.status)
return
}
w.Header().Set("Content-Type", f.contentType)
w.Header().Set("ETag", f.etag)
http.ServeContent(w, r, "", time.Time{}, bytes.NewReader(f.original))
}))
f.processor = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
f.encodes.Add(1)
encoded := strings.TrimSuffix(r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:], ".webp")
source, err := base64.RawURLEncoding.DecodeString(encoded)
if err != nil {
t.Errorf("encoder received an invalid source URL: %v", err)
}
f.mu.Lock()
f.sources = append(f.sources, string(source))
status, contentType, data := f.processStatus, f.processType, append([]byte(nil), f.data...)
f.mu.Unlock()
if f.started != nil {
select {
case f.started <- struct{}{}:
default:
}
}
if f.release != nil {
select {
case <-f.release:
case <-r.Context().Done():
return
}
}
w.Header().Set("Content-Type", contentType)
w.WriteHeader(status)
_, _ = w.Write(data)
}))
var err error
f.gateway, err = New(Config{Source: f.origin.URL + "/bucket", Imgproxy: f.processor.URL,
Concurrency: 16, MaxDimension: 4096, CacheBytes: 4096, MaxSourceBytes: 1024, MaxResultBytes: 1024, Timeout: 2 * time.Second})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { f.origin.Close(); f.processor.Close(); f.gateway.client.CloseIdleConnections() })
return f
}
// request sends conditional or range requests to the gateway handler.
func (f *fixture) request(method, path string, headers http.Header) *httptest.ResponseRecorder {
r := httptest.NewRequest(method, "http://untrusted-host"+path, nil)
if headers != nil {
r.Header = headers
}
w := httptest.NewRecorder()
f.gateway.ServeHTTP(w, r)
return w
}
// imagePath escapes the processing query while preserving its slash-separated semantics.
func imagePath() string { return "/image.png?x-oss-process=" + url.QueryEscape(transform) }
// TestCacheRechecksAccess checks cache hits, revocation of public access, and source deletion.
func TestCacheRechecksAccess(t *testing.T) {
f := newFixture(t)
first := f.request("GET", imagePath(), nil)
if first.Code != 200 || first.Header().Get("ETag") == "\"source-v1\"" || first.Header().Get("Content-Type") != "image/webp" {
t.Fatalf("unexpected first processed response: %d %v", first.Code, first.Header())
}
second := f.request("GET", imagePath(), nil)
if !bytes.Equal(first.Body.Bytes(), second.Body.Bytes()) || f.encodes.Load() != 1 || f.heads.Load() != 3 {
t.Fatal("cache hits must recheck access without re-encoding")
}
f.mu.Lock()
f.status = 403
f.mu.Unlock()
denied := f.request("GET", imagePath(), http.Header{"If-None-Match": {first.Header().Get("ETag")}})
if denied.Code != 403 || denied.Header().Get("Cache-Control") != "no-store" {
t.Fatal("conditional request bypassed revoked access")
}
f.mu.Lock()
f.status = 404
f.mu.Unlock()
if f.request("HEAD", imagePath(), nil).Code != 404 {
t.Fatal("deleted source still served a cached result")
}
}
// TestRepresentationHeaders checks HEAD, 304, and Range semantics for processed bytes.
func TestRepresentationHeaders(t *testing.T) {
f := newFixture(t)
full := f.request("GET", imagePath(), nil)
etag := full.Header().Get("ETag")
head := f.request("HEAD", imagePath(), nil)
if head.Code != 200 || head.Body.Len() != 0 || head.Header().Get("Content-Length") != fmt.Sprint(len(f.data)) {
t.Fatal("HEAD used the source length or returned a body")
}
conditional := f.request("GET", imagePath(), http.Header{"If-None-Match": {etag}})
if conditional.Code != 304 || conditional.Body.Len() != 0 {
t.Fatal("output ETag did not produce 304")
}
if f.request("GET", imagePath(), http.Header{"If-None-Match": {"\"source-v1\""}}).Code != 200 {
t.Fatal("source ETag was incorrectly used")
}
ranged := f.request("GET", imagePath(), http.Header{"Range": {"bytes=0-3"}})
if ranged.Code != 206 || ranged.Body.String() != "RIFF" || ranged.Header().Get("Content-Range") != fmt.Sprintf("bytes 0-3/%d", len(f.data)) {
t.Fatal("Range did not apply to processed bytes")
}
invalid := f.request("GET", imagePath(), http.Header{"Range": {"bytes=999-"}})
if invalid.Code != 416 || invalid.Header().Get("Cache-Control") != "no-store" {
t.Fatal("unsatisfiable output range was not rejected")
}
failed := f.request("GET", imagePath(), http.Header{"If-Match": {"\"wrong\""}})
if failed.Code != 412 || failed.Header().Get("Cache-Control") != "no-store" {
t.Fatal("failed precondition response was cacheable")
}
}
// TestChangedSource checks invalidation after overwrites and explicit version reads.
func TestChangedSource(t *testing.T) {
f := newFixture(t)
f.request("GET", imagePath(), nil)
f.mu.Lock()
f.etag = "\"source-v2\""
f.data = []byte("new result")
f.mu.Unlock()
if f.request("GET", imagePath(), nil).Body.String() != "new result" || f.encodes.Load() != 2 {
t.Fatal("source overwrite returned stale cache")
}
if f.request("GET", imagePath()+"&versionId=v1", nil).Code != 200 {
t.Fatal("explicit version read failed")
}
f.mu.Lock()
last := f.sources[len(f.sources)-1]
f.mu.Unlock()
if last != f.origin.URL+"/bucket/image.png?versionId=v1" {
t.Fatalf("version was not forwarded to the source: %s", last)
}
}
// TestSourceChangesDuringEncoding checks that changing sources do not populate stale cache entries.
func TestSourceChangesDuringEncoding(t *testing.T) {
f := newFixture(t)
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
done := make(chan *httptest.ResponseRecorder, 1)
go func() { done <- f.request("GET", imagePath(), nil) }()
<-f.started
f.mu.Lock()
f.etag = "\"replaced\""
f.mu.Unlock()
close(f.release)
if (<-done).Code != 409 {
t.Fatal("source replacement still cached a processed result")
}
if f.request("GET", imagePath(), nil).Code != 200 || f.encodes.Load() != 2 {
t.Fatal("retry did not re-encode")
}
}
// TestRejectUnsafeRequests checks that credentials, arbitrary URLs, duplicate queries, and writes do not reach backends.
func TestRejectUnsafeRequests(t *testing.T) {
f := newFixture(t)
tests := []struct {
method, path string
header http.Header
status int
}{
{"PUT", imagePath(), nil, 405},
{"GET", imagePath(), http.Header{"Authorization": {"AWS4-HMAC-SHA256 secret"}}, 403},
{"GET", imagePath(), http.Header{"X-Amz-Security-Token": {"secret"}}, 403},
{"GET", imagePath() + "&X-Amz-Signature=secret", nil, 400},
{"GET", imagePath() + "&source=http://private-target", nil, 400},
{"GET", imagePath() + "&x-oss-process=image", nil, 400},
{"GET", "/image.png?x-oss-process=%", nil, 400},
{"GET", "/image.png?x-oss-process=", nil, 400},
}
for _, test := range tests {
if response := f.request(test.method, test.path, test.header); response.Code != test.status {
t.Errorf("request %s returned %d; expected %d", test.path, response.Code, test.status)
}
}
if f.encodes.Load() != 0 || f.heads.Load() != 0 {
t.Fatal("invalid request reached a backend")
}
}
// TestLimitsAndBackendFailures checks resource limits, media types, missing sources, and encoder errors.
func TestLimitsAndBackendFailures(t *testing.T) {
tests := []struct {
name string
change func(*fixture)
status int
}{
{"oversized source", func(f *fixture) { f.sourceBytes = 1025 }, 413},
{"oversized result", func(f *fixture) { f.data = make([]byte, 1025) }, 413},
{"empty result", func(f *fixture) { f.data = nil }, 413},
{"incorrect media type", func(f *fixture) { f.processType = "text/html" }, 502},
{"encoder failure", func(f *fixture) { f.processStatus = 500 }, 502},
{"source denied", func(f *fixture) { f.status = 403 }, 403},
{"source missing", func(f *fixture) { f.status = 404 }, 404},
{"source failure", func(f *fixture) { f.status = 500 }, 502},
{"redirect", func(f *fixture) { f.status = 302 }, 502},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
f := newFixture(t)
test.change(f)
for i := 0; i < 2; i++ {
w := f.request("GET", imagePath(), nil)
if w.Code != test.status || w.Header().Get("Cache-Control") != "no-store" {
t.Fatalf("unexpected error response: %d %v", w.Code, w.Header())
}
}
if f.gateway.cache.order.Len() != 0 {
t.Fatal("failed result entered the cache")
}
})
}
}
// TestConcurrentMisses checks that concurrent misses for the same image encode only once.
func TestConcurrentMisses(t *testing.T) {
f := newFixture(t)
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
var wg sync.WaitGroup
wg.Add(8)
for i := 0; i < 8; i++ {
go func() {
defer wg.Done()
if w := f.request("GET", imagePath(), nil); w.Code != 200 {
t.Errorf("concurrent request failed: %d", w.Code)
}
}()
}
<-f.started
close(f.release)
wg.Wait()
if f.encodes.Load() != 1 {
t.Fatalf("encoder calls for concurrent misses: %d", f.encodes.Load())
}
}
// TestConcurrencyAndCancellation checks immediate overload rejection and prompt cancellation of waiting requests.
func TestConcurrencyAndCancellation(t *testing.T) {
f := newFixture(t)
f.gateway.requests = make(chan struct{}, 1)
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
ctx, cancel := context.WithCancel(context.Background())
r := httptest.NewRequest("GET", imagePath(), nil).WithContext(ctx)
done := make(chan struct{})
go func() { f.gateway.ServeHTTP(httptest.NewRecorder(), r); close(done) }()
<-f.started
if f.request("GET", imagePath(), nil).Code != 429 {
t.Fatal("overload did not return 429")
}
cancel()
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("cancelled request continued waiting for encoding")
}
close(f.release)
}
// TestOriginalRead checks source ETag, HEAD, and Range without processing.
func TestOriginalRead(t *testing.T) {
f := newFixture(t)
get := f.request("GET", "/image.png", nil)
if get.Code != 200 || get.Body.String() != "original-image" || get.Header().Get("Content-Type") != "image/png" {
t.Fatal("unexpected source read")
}
ranged := f.request("GET", "/image.png", http.Header{"Range": {"bytes=0-2"}})
if ranged.Code != 206 || ranged.Body.String() != "ori" {
t.Fatal("unexpected source range read")
}
if f.request("GET", "/image.png", http.Header{"If-None-Match": {"\"source-v1\""}}).Code != 304 {
t.Fatal("unexpected source conditional read")
}
if f.request("HEAD", "/image.png", nil).Body.Len() != 0 || f.encodes.Load() != 0 {
t.Fatal("source request unexpectedly encoded an image")
}
}
// TestBrowserCookiesAreNotForwarded checks anonymous reads when browsers send site cookies.
func TestBrowserCookiesAreNotForwarded(t *testing.T) {
f := newFixture(t)
if f.request("GET", imagePath(), http.Header{"Cookie": {"session=secret"}}).Code != 200 {
t.Fatal("site cookie interfered with a public image read")
}
f.mu.Lock()
f.status = 403
f.mu.Unlock()
if f.request("GET", imagePath(), http.Header{"Cookie": {"session=secret"}}).Code != 403 {
t.Fatal("cookie elevated source read permissions")
}
}
// TestNoETagDisablesCache checks that sources without validators are processed on each request.
func TestNoETagDisablesCache(t *testing.T) {
f := newFixture(t)
f.etag = ""
f.request("GET", imagePath(), nil)
f.request("GET", imagePath(), nil)
if f.encodes.Load() != 2 {
t.Fatal("result was reused without a source ETag")
}
}
// TestFixedSourceAndSigning checks path escaping, fixed hosts, and imgproxy HMAC signatures.
func TestFixedSourceAndSigning(t *testing.T) {
f := newFixture(t)
f.gateway.key, f.gateway.salt = []byte("secret"), []byte("salt")
r := httptest.NewRequest("GET", "http://attacker/a%20b/%23.png?versionId=v%2B1", nil)
source := f.gateway.sourceURL(r, r.URL.Query())
if source.String() != f.origin.URL+"/bucket/a%20b/%23.png?versionId=v%2B1" {
t.Fatal("fixed source or object escaping changed")
}
u := f.gateway.processorURL(source, options{width: 640, quality: 85, format: "webp"})
parts := strings.SplitN(strings.TrimPrefix(u.Path, "/"), "/", 2)
mac := hmac.New(sha256.New, []byte("secret"))
mac.Write([]byte("salt/" + parts[1]))
if parts[0] != base64.RawURLEncoding.EncodeToString(mac.Sum(nil)) {
t.Fatal("imgproxy signature did not match the protocol")
}
}
// TestRealHTTPHead checks output length and body suppression for real HTTP HEAD responses.
func TestRealHTTPHead(t *testing.T) {
f := newFixture(t)
s := httptest.NewServer(f.gateway)
defer s.Close()
resp, err := http.Head(s.URL + imagePath())
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
body, _ := io.ReadAll(resp.Body)
if resp.StatusCode != 200 || resp.ContentLength != int64(len(f.data)) || len(body) != 0 {
t.Fatal("real HEAD response violated HTTP semantics")
}
}
// waitUntil bounds synchronization waits so a regression fails rather than hanging.
func waitUntil(t *testing.T, ready func() bool) {
t.Helper()
deadline := time.Now().Add(time.Second)
for !ready() {
if time.Now().After(deadline) {
t.Fatal("timed out waiting for test synchronization")
}
time.Sleep(time.Millisecond)
}
}
// TestCancelledLeaderDoesNotFailWaiter checks that shared work outlives its first caller.
func TestCancelledLeaderDoesNotFailWaiter(t *testing.T) {
f := newFixture(t)
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
leader := make(chan *httptest.ResponseRecorder, 1)
go func() {
w := httptest.NewRecorder()
f.gateway.ServeHTTP(w, httptest.NewRequest("GET", imagePath(), nil).WithContext(ctx))
leader <- w
}()
<-f.started
waiter := make(chan *httptest.ResponseRecorder, 1)
go func() { waiter <- f.request("GET", imagePath(), nil) }()
waitUntil(t, func() bool { return f.heads.Load() >= 2 })
cancel()
if (<-leader).Code != 504 {
t.Fatal("cancelled leader did not leave promptly")
}
close(f.release)
response := <-waiter
if response.Code != 200 || response.Body.String() != string(f.data) || f.encodes.Load() != 1 {
t.Fatalf("leader cancellation failed shared work: status=%d encodes=%d", response.Code, f.encodes.Load())
}
}
// TestCancelledWorkRemainsBounded checks that cancellation cannot free an active encoding slot.
func TestCancelledWorkRemainsBounded(t *testing.T) {
f := newFixture(t)
f.gateway.requests, f.gateway.work = make(chan struct{}, 1), make(chan struct{}, 1)
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
done := make(chan struct{})
go func() {
f.gateway.ServeHTTP(httptest.NewRecorder(), httptest.NewRequest("GET", imagePath(), nil).WithContext(ctx))
close(done)
}()
<-f.started
cancel()
<-done
other := strings.Replace(imagePath(), "image.png", "other.png", 1)
if response := f.request("GET", other, nil); response.Code != 429 || f.encodes.Load() != 1 {
t.Fatal("cancelled request bypassed active encoding concurrency limit")
}
close(f.release)
waitUntil(t, func() bool { return len(f.gateway.work) == 0 })
if f.request("GET", other, nil).Code != 200 {
t.Fatal("completed work did not release its concurrency slot")
}
}
// TestAbandonedWorkTimesOut checks that a job with no remaining callers releases its slot.
func TestAbandonedWorkTimesOut(t *testing.T) {
f := newFixture(t)
f.gateway.config.Timeout = 50 * time.Millisecond
f.gateway.work = make(chan struct{}, 1)
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
f.gateway.ServeHTTP(httptest.NewRecorder(), httptest.NewRequest("GET", imagePath(), nil).WithContext(ctx))
close(done)
}()
<-f.started
cancel()
<-done
waitUntil(t, func() bool { return len(f.gateway.work) == 0 })
close(f.release)
if f.request("GET", imagePath(), nil).Code != 200 || f.encodes.Load() != 2 {
t.Fatal("abandoned encoding did not time out and permit a new job")
}
}
// TestSameSecondOverwriteDoesNotReturn304 checks that source dates cannot validate output bytes.
func TestSameSecondOverwriteDoesNotReturn304(t *testing.T) {
f := newFixture(t)
first := f.request("GET", imagePath(), nil)
f.mu.Lock()
f.etag, f.data = "\"source-v2\"", []byte("replacement bytes")
f.mu.Unlock()
headers := http.Header{"If-Modified-Since": {"Sun, 04 Oct 2026 10:00:00 GMT"}}
response := f.request("GET", imagePath(), headers)
if response.Code != 200 || response.Body.String() != "replacement bytes" || response.Header().Get("Last-Modified") != "" {
t.Fatal("same-second overwrite incorrectly used the source modification date")
}
if response.Header().Get("ETag") == first.Header().Get("ETag") {
t.Fatal("changed output reused the old validator")
}
if f.request("GET", imagePath(), http.Header{"If-None-Match": {response.Header().Get("ETag")}}).Code != 304 {
t.Fatal("output ETag stopped validating unchanged output")
}
}
// TestOriginalRangeLimitsAndErrors checks complete-object limits and unsatisfied range metadata.
func TestOriginalRangeLimitsAndErrors(t *testing.T) {
f := newFixture(t)
f.original = make([]byte, 2048)
if response := f.request("GET", "/image.png", http.Header{"Range": {"bytes=0-0"}}); response.Code != 413 {
t.Fatalf("small range bypassed the complete source size limit: %d", response.Code)
}
f.original = []byte("original-image")
response := f.request("GET", "/image.png", http.Header{"Range": {"bytes=999-"}})
if response.Code != 416 || response.Header().Get("Content-Range") != "bytes */14" || response.Header().Get("Cache-Control") != "no-store" {
t.Fatalf("unsatisfied source range lost its complete length: %d %v", response.Code, response.Header())
}
}
// TestOriginalRangeRequiresCompleteLength rejects unknown or invalid source range totals.
func TestOriginalRangeRequiresCompleteLength(t *testing.T) {
for _, contentRange := range []string{"", "bytes 0-0/*", "bytes 0-0/no-number", "items 0-0/100", "bytes 0-0/0"} {
t.Run(contentRange, func(t *testing.T) {
f := newFixture(t)
s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Range", contentRange)
w.WriteHeader(206)
_, _ = w.Write([]byte("x"))
}))
defer s.Close()
f.gateway.source, _ = url.Parse(s.URL)
if f.request("GET", "/image.png", http.Header{"Range": {"bytes=0-0"}}).Code != 502 {
t.Fatal("source range without a valid complete length was accepted")
}
})
}
}
// TestRejectNormalizedTraversal protects a configured bucket prefix, including escaped paths.
func TestRejectNormalizedTraversal(t *testing.T) {
f := newFixture(t)
for _, path := range []string{"/../other/image.png", "/./image.png", "/%2e%2e/other.png", "/%252e%252e/other.png", "/%25252e%25252e/other.png", "/x%2f..%2fother.png", "/x%252f..%252fother.png", "/x%5c..%5cother.png", "/x%255c..%255cother.png"} {
for _, query := range []string{"", "?x-oss-process=" + url.QueryEscape(transform)} {
if response := f.request("GET", path+query, nil); response.Code != 400 {
t.Errorf("unsafe path accepted: %s status=%d", path, response.Code)
}
}
}
if f.heads.Load() != 0 || f.encodes.Load() != 0 {
t.Fatal("unsafe path reached a backend")
}
if f.request("GET", "/a%20b/%23%25.png?x-oss-process="+url.QueryEscape(transform), nil).Code != 200 {
t.Fatal("ordinary escaped object name was rejected")
}
}
// deadlineRecorder exposes the HTTP server deadline hook without relying on timing-sensitive sockets.
type deadlineRecorder struct {
*httptest.ResponseRecorder
deadline time.Time
calls int
}
// SetWriteDeadline records when the gateway allocates its client write budget.
func (w *deadlineRecorder) SetWriteDeadline(deadline time.Time) error {
w.deadline, w.calls = deadline, w.calls+1
return nil
}
// TestWriteBudgetStartsAfterProcessing checks that slow backends do not consume the write budget.
func TestWriteBudgetStartsAfterProcessing(t *testing.T) {
f := newFixture(t)
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
w := &deadlineRecorder{ResponseRecorder: httptest.NewRecorder()}
done := make(chan struct{})
go func() { f.gateway.ServeHTTP(w, httptest.NewRequest("GET", imagePath(), nil)); close(done) }()
<-f.started
if w.calls != 0 {
t.Fatal("write budget started while backend processing was blocked")
}
released := time.Now()
close(f.release)
<-done
if w.calls != 1 || w.deadline.Before(released.Add(f.gateway.config.Timeout)) || w.Code != 200 {
t.Fatal("response did not receive a separate full write timeout")
}
}
// TestSlowMetadataDoesNotConsumeEncodingWait checks independent source and work budgets.
func TestSlowMetadataDoesNotConsumeEncodingWait(t *testing.T) {
f := newFixture(t)
f.gateway.config.Timeout = 300 * time.Millisecond
f.headStarted, f.headRelease = make(chan struct{}), make(chan struct{})
f.started, f.release = make(chan struct{}, 1), make(chan struct{})
done := make(chan *httptest.ResponseRecorder, 1)
go func() { done <- f.request("GET", imagePath(), nil) }()
<-f.headStarted
time.Sleep(200 * time.Millisecond)
close(f.headRelease)
<-f.started
time.Sleep(200 * time.Millisecond)
close(f.release)
response := <-done
if response.Code != 200 || response.Body.String() != string(f.data) {
t.Fatalf("initial metadata consumed the encoding wait budget: %d", response.Code)
}
}
// TestOriginalContentEncoding checks original encoded bytes, HEAD, ranges, and real client decoding.
func TestOriginalContentEncoding(t *testing.T) {
f := newFixture(t)
var compressed bytes.Buffer
writer := gzip.NewWriter(&compressed)
_, _ = writer.Write(f.original)
if err := writer.Close(); err != nil {
t.Fatal(err)
}
source := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "image/png")
w.Header().Set("Content-Encoding", "gzip")
w.Header().Set("Content-Length", fmt.Sprint(compressed.Len()))
http.ServeContent(w, r, "", time.Time{}, bytes.NewReader(compressed.Bytes()))
}))
defer source.Close()
f.gateway.source, _ = url.Parse(source.URL)
get := f.request("GET", "/encoded.png", nil)
if get.Code != 200 || get.Header().Get("Content-Encoding") != "gzip" || !bytes.Equal(get.Body.Bytes(), compressed.Bytes()) {
t.Fatal("original response lost its content encoding or changed stored bytes")
}
head := f.request("HEAD", "/encoded.png", nil)
if head.Header().Get("Content-Encoding") != "gzip" || head.Header().Get("Content-Length") != fmt.Sprint(compressed.Len()) || head.Body.Len() != 0 {
t.Fatal("encoded original HEAD lost representation metadata")
}
ranged := f.request("GET", "/encoded.png", http.Header{"Range": {"bytes=0-2"}})
if ranged.Code != 206 || ranged.Header().Get("Content-Encoding") != "gzip" || !bytes.Equal(ranged.Body.Bytes(), compressed.Bytes()[:3]) {
t.Fatal("encoded original range did not preserve stored representation")
}
server := httptest.NewServer(f.gateway)
defer server.Close()
response, err := http.Get(server.URL + "/encoded.png")
if err != nil {
t.Fatal(err)
}
defer response.Body.Close()
decoded, err := io.ReadAll(response.Body)
if err != nil || !bytes.Equal(decoded, f.original) {
t.Fatal("real client could not decode the original representation")
}
}
// TestRepeatedSlashesPreserveBucketPrefix checks actual forwarded URLs rather than assuming traversal.
func TestRepeatedSlashesPreserveBucketPrefix(t *testing.T) {
f := newFixture(t)
for _, path := range []string{"//other-bucket/image.png", "/nested//image.png", "/%2fother-bucket/image.png"} {
if response := f.request("GET", path+"?x-oss-process="+url.QueryEscape(transform), nil); response.Code != 200 {
t.Fatalf("valid repeated-slash key was not served: %s %d", path, response.Code)
}
f.mu.Lock()
encodedSource := f.sources[len(f.sources)-1]
f.mu.Unlock()
source, err := url.Parse(encodedSource)
if err != nil || source.Host != strings.TrimPrefix(f.origin.URL, "http://") || !strings.HasPrefix(source.Path, "/bucket/") {
t.Fatalf("repeated slashes escaped the fixed source: %s", encodedSource)
}
}
}
// TestOriginalMediaType checks executable documents cannot be served under this origin.
func TestOriginalMediaType(t *testing.T) {
f := newFixture(t)
for _, mediaType := range []string{"text/html", "image/svg+xml", "image/svg+xml; bad", "image/x-example+xml", "image/x-example+xml; charset=utf-8", "IMAGE/X-EXAMPLE+XML", "application/xhtml+xml", "text/xml", ""} {
f.contentType = mediaType
for _, request := range []struct {
method string
headers http.Header
}{{"GET", nil}, {"HEAD", nil}, {"GET", http.Header{"Range": {"bytes=0-0"}}}} {
if response := f.request(request.method, "/page", request.headers); response.Code != 502 {
t.Fatalf("executable source type %q was served for %s: %d", mediaType, request.method, response.Code)
}
}
}
for _, mediaType := range []string{"image/avif; codecs=\"av01.0.08M.08\"", "image/webp", "image/vnd.microsoft.icon", "application/octet-stream", "binary/octet-stream"} {
f.contentType = mediaType
response := f.request("GET", "/image.png", nil)
if response.Code != 200 || response.Body.String() != "original-image" {
t.Fatalf("safe source type %q was rejected: %d", mediaType, response.Code)
}
if got := response.Header().Get("Content-Type"); got != mediaType {
t.Fatalf("source content type %q forwarded as %q", mediaType, got)
}
}
}
// TestWriteBudgetsForOriginalAndErrors checks that status commits cover all non-encoding body paths.
func TestWriteBudgetsForOriginalAndErrors(t *testing.T) {
f := newFixture(t)
for _, request := range []struct {
name, method string
headers http.Header
status int
}{
{"original", "GET", nil, 200},
{"head", "HEAD", nil, 200},
{"conditional", "GET", http.Header{"If-None-Match": {f.etag}}, 304},
{"range", "GET", http.Header{"Range": {"bytes=0-0"}}, 206},
{"invalid-range", "GET", http.Header{"Range": {"bytes=999-"}}, 416},
{"unsupported-method", "POST", nil, 405},
} {
t.Run(request.name, func(t *testing.T) {
w := &deadlineRecorder{ResponseRecorder: httptest.NewRecorder()}
r := httptest.NewRequest(request.method, "http://untrusted-host/image.png", nil)
if request.headers != nil {
r.Header = request.headers.Clone()
}
started := time.Now()
f.gateway.ServeHTTP(w, r)
if w.Code != request.status || w.calls != 1 || w.deadline.Before(started.Add(f.gateway.config.Timeout)) {
t.Fatalf("response lost its full write budget: status=%d calls=%d deadline=%v", w.Code, w.calls, w.deadline)
}
})
}
}
// TestOriginalConditionalMediaType prevents 304 metadata from reclassifying cached bytes as executable.
func TestOriginalConditionalMediaType(t *testing.T) {
f := newFixture(t)
var contentType atomic.Value
contentType.Store("")
origin := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
// net/http strips Content-Type from 304 server responses; preserve explicit backend metadata on the wire.
connection, buffer, err := w.(http.Hijacker).Hijack()
if err != nil {
t.Errorf("could not write conditional backend response: %v", err)
return
}
defer connection.Close()
_, _ = fmt.Fprint(buffer, "HTTP/1.1 304 Not Modified\r\nETag: \"source-v1\"\r\nConnection: close\r\n")
if value := contentType.Load().(string); value != "" {
_, _ = fmt.Fprintf(buffer, "Content-Type: %s\r\n", value)
}
_, _ = fmt.Fprint(buffer, "\r\n")
if err := buffer.Flush(); err != nil {
t.Errorf("could not flush conditional backend response: %v", err)
}
}))
defer origin.Close()
f.gateway.source, _ = url.Parse(origin.URL)
for _, check := range []struct {
mediaType string
status int
}{
{"", 304},
{"image/png", 304},
{"image/avif; codecs=\"av01.0.08M.08\"", 304},
{"text/html", 502},
{"image/x-example+xml", 502},
{"image/svg+xml; bad", 502},
} {
contentType.Store(check.mediaType)
response := f.request("GET", "/image.png", http.Header{"If-None-Match": {f.etag}})
if response.Code != check.status {
t.Fatalf("conditional type %q returned %d instead of %d", check.mediaType, response.Code, check.status)
}
if check.status == 304 && (response.Body.Len() != 0 || response.Header().Get("Content-Type") != check.mediaType) {
t.Fatal("safe conditional response changed its body or media type")
}
}
}
// TestProcessedErrorsDisableCaching checks standard-library precondition and range failures.
func TestProcessedErrorsDisableCaching(t *testing.T) {
f := newFixture(t)
for _, check := range []struct {
headers http.Header
status int
}{
{http.Header{"If-Match": {"\"other-output\""}}, 412},
{http.Header{"Range": {"bytes=999-"}}, 416},
} {
response := f.request("GET", imagePath(), check.headers)
if response.Code != check.status || response.Header().Get("Cache-Control") != "no-store" {
t.Fatalf("processed protocol error lost no-store: status=%d cache=%q", response.Code, response.Header().Get("Cache-Control"))
}
}
}
+104
View File
@@ -0,0 +1,104 @@
package gateway
import (
"fmt"
"strconv"
"strings"
)
type options struct {
width, height, quality int
format string
}
// parseOptions accepts the supported OSS subset and rejects duplicates and unknown operations.
// quality,Q is absolute quality; relative quality q has no equivalent in imgproxy.
func parseOptions(value string, maxDimension int) (options, error) {
o := options{quality: 85, format: "webp"}
parts := strings.Split(value, "/")
if len(value) > 256 || len(parts) < 2 || parts[0] != "image" {
return o, fmt.Errorf("expected image/processing-operation")
}
seen := make(map[string]bool)
for _, part := range parts[1:] {
fields := strings.Split(part, ",")
if seen[fields[0]] {
return o, fmt.Errorf("repeated processing operation")
}
seen[fields[0]] = true
switch fields[0] {
case "resize":
if len(fields) < 2 {
return o, fmt.Errorf("missing resize dimensions")
}
keys := make(map[string]bool)
for _, f := range fields[1:] {
kv := strings.SplitN(f, "_", 2)
if len(kv) != 2 || keys[kv[0]] {
return o, fmt.Errorf("invalid or repeated resize parameter")
}
keys[kv[0]] = true
switch kv[0] {
case "w", "h":
n, err := strconv.Atoi(kv[1])
if err != nil || n < 1 || n > maxDimension {
return o, fmt.Errorf("dimension exceeds allowed range")
}
if kv[0] == "w" {
o.width = n
} else {
o.height = n
}
case "m":
if kv[1] != "lfit" {
return o, fmt.Errorf("only aspect-preserving lfit is supported")
}
case "limit":
if kv[1] != "1" {
return o, fmt.Errorf("image enlargement is not allowed")
}
default:
return o, fmt.Errorf("unsupported resize parameter")
}
}
if o.width == 0 && o.height == 0 {
return o, fmt.Errorf("missing resize dimensions")
}
case "quality":
if len(fields) != 2 || !strings.HasPrefix(fields[1], "Q_") {
return o, fmt.Errorf("only absolute quality Q is supported")
}
q, err := strconv.Atoi(strings.TrimPrefix(fields[1], "Q_"))
if err != nil || q < 1 || q > 100 {
return o, fmt.Errorf("quality must be between 1 and 100")
}
o.quality = q
case "format":
if len(fields) != 2 {
return o, fmt.Errorf("invalid format parameter")
}
o.format = fields[1]
if o.format == "jpeg" {
o.format = "jpg"
}
if o.format != "jpg" && o.format != "png" && o.format != "webp" {
return o, fmt.Errorf("only JPEG, PNG, and WebP are supported")
}
default:
return o, fmt.Errorf("unsupported image processing operation")
}
}
// Bound both output axes, including format-only and single-dimension requests.
if o.width == 0 {
o.width = maxDimension
}
if o.height == 0 {
o.height = maxDimension
}
return o, nil
}
// path canonicalizes imgproxy operations so equivalent URLs share cached results.
func (o options) path() string {
return fmt.Sprintf("/rs:fit:%d:%d:0:0/q:%d/f:%s", o.width, o.height, o.quality, o.format)
}
+121
View File
@@ -0,0 +1,121 @@
package gateway
import (
"strings"
"testing"
"time"
)
// TestOptionsCoverage checks the supported OSS subset and rejected operations.
func TestOptionsCoverage(t *testing.T) {
valid := []string{transform, "image/resize,w_240,h_640,m_lfit,limit_1/format,webp", "image/format,jpeg/quality,Q_100", "image/format,png", "image/resize,h_640"}
for _, value := range valid {
if _, err := parseOptions(value, 4096); err != nil {
t.Errorf("valid operation rejected %s: %v", value, err)
}
}
invalid := []string{"", "image", "image/", "image/resize", "image/resize,w_0", "image/resize,w_-1", "image/resize,w_4097", "image/resize,w_1,w_2", "image/resize,m_lfit", "image/resize,w_640,m_fill", "image/resize,w_1,limit_0", "image/quality,q_85", "image/quality,Q_0", "image/quality,Q_101", "image/quality,Q_85,Q_90", "image/format,svg", "image/format,webp/format,png", "image/watermark,x", "image/resize,w_1/resize,w_2", strings.Repeat("x", 257)}
for _, value := range invalid {
if _, err := parseOptions(value, 4096); err == nil {
t.Errorf("invalid operation accepted: %s", value)
}
}
a, _ := parseOptions(transform, 4096)
b, _ := parseOptions("image/format,webp/quality,Q_85/resize,w_640", 4096)
if a.path() != b.path() {
t.Fatal("equivalent operations were not canonicalized")
}
}
// TestByteLimitedCache checks byte eviction, replacement, disabled caching, and oversized entries.
func TestByteLimitedCache(t *testing.T) {
c := newCache(400)
value := &result{data: make([]byte, 100)}
c.put("a", value)
c.put("b", value)
if c.get("a") != nil || c.get("b") == nil {
t.Fatal("cache did not evict by byte usage")
}
c.put("b", &result{data: []byte("updated")})
if c.get("b").data[0] != 'u' {
t.Fatal("cache replacement failed")
}
c.put("large", &result{data: make([]byte, 401)})
if c.get("large") != nil || c.used > c.capacity {
t.Fatal("oversized result exceeded cache capacity")
}
disabled := newCache(0)
disabled.put("a", value)
if disabled.get("a") != nil {
t.Fatal("cache disabling failed")
}
large := newCache(1 << 20)
for i := 0; i < 2048; i++ {
large.put(string(rune(i)), value)
}
if len(large.entries) != 1024 {
t.Fatal("small object entry count was not bounded")
}
}
// TestConfigurationRejectsUnsafeBackends checks explicitly configured backends without embedded credentials.
func TestConfigurationRejectsUnsafeBackends(t *testing.T) {
base := Config{Source: "http://localhost:8333/bucket", Imgproxy: "http://localhost:8080", Concurrency: 1, MaxDimension: 4096, CacheBytes: 0, MaxSourceBytes: 1024, MaxResultBytes: 1024, Timeout: time.Second}
for _, source := range []string{"", "file:///etc/passwd", "http://secret:password@localhost", "http://localhost?signature=secret", "http://localhost#fragment"} {
c := base
c.Source = source
if _, err := New(c); err == nil {
t.Errorf("invalid backend accepted: %s", source)
}
}
for _, mutate := range []func(*Config){
func(c *Config) { c.Key = "zz"; c.Salt = "aa" }, func(c *Config) { c.Key = "aa" },
func(c *Config) { c.Concurrency = 0 }, func(c *Config) { c.CacheBytes = -1 },
func(c *Config) { c.MaxSourceBytes = 0 }, func(c *Config) { c.Timeout = 0 },
} {
c := base
mutate(&c)
if _, err := New(c); err == nil {
t.Fatal("invalid resource or signing configuration accepted")
}
}
}
// TestOutputBounds covers width-only, height-only, and format-only output limits.
func TestOutputBounds(t *testing.T) {
for _, test := range []struct {
value string
width, height int
}{
{"image/resize,w_640", 640, 4096},
{"image/resize,h_640", 4096, 640},
{"image/format,png", 4096, 4096},
{"image/resize,w_240,h_640", 240, 640},
} {
o, err := parseOptions(test.value, 4096)
if err != nil || o.width != test.width || o.height != test.height {
t.Fatalf("missing output axis bound for %s: %+v %v", test.value, o, err)
}
}
}
// TestProcessorRootURL rejects ambiguous path prefixes while allowing root URLs.
func TestProcessorRootURL(t *testing.T) {
base := Config{Source: "http://localhost:8333/bucket", Concurrency: 1, MaxDimension: 4096, MaxSourceBytes: 1024, MaxResultBytes: 1024, Timeout: time.Second}
for _, address := range []string{"http://localhost:8080/imgproxy", "http://localhost:8080/prefix/", "http://localhost:8080/%2f"} {
c := base
c.Imgproxy = address
if _, err := New(c); err == nil {
t.Fatalf("prefixed imgproxy URL accepted: %s", address)
}
}
for _, address := range []string{"http://localhost:8080", "http://localhost:8080/"} {
c := base
c.Imgproxy = address
g, err := New(c)
if err != nil {
t.Fatalf("root imgproxy URL rejected: %s %v", address, err)
}
g.client.CloseIdleConnections()
}
}