Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/operation/operation.md
Original file line number Diff line number Diff line change
Expand Up @@ -2056,6 +2056,12 @@ flowchart TD
- `cache.reval_dropped`: Counter, revalidation jobs dropped because the queue was full or body read failed
- `cache.reval_error`: Counter, background revalidation fetch failures
- `cache.reval_duration`: Histogram, end-to-end duration of each background revalidation job
- `cache.reval_backend_dispatch`: Counter, revalidation fetches dispatched directly to the route's
resolved backend instead of looping back through skipper's own listener (force mode with a static
backend only). These calls bypass skipper's proxy pipeline, so they are *not* reflected in the
usual per-route/backend metrics (`MeasureBackend*`) or the access log - this counter is the only
metrics-level signal for that traffic. Self-loopback revalidations don't increment it, but do show
up in the standard backend metrics/access log via their inner hop through the listener.

**L2 (if Valkey or Redis is configured):**

Expand Down
25 changes: 16 additions & 9 deletions filters/cache/cache_control.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,14 +8,15 @@ import (
)

type cacheDirectives struct {
noStore bool
noCache bool
private bool
mustRevalidate bool
proxyRevalidate bool
public bool
maxAge int64 // -1 = not present; 0 means max-age=0
sMaxAge int64 // -1 = not present
noStore bool
noCache bool
private bool
mustRevalidate bool
proxyRevalidate bool
public bool
maxAge int64 // -1 = not present; 0 means max-age=0
sMaxAge int64 // -1 = not present
staleWhileRevalidate int64 // -1 = not present (RFC 5861)
}

type requestCacheDirectives struct {
Expand Down Expand Up @@ -68,7 +69,7 @@ func parseRequestCacheControl(h http.Header) requestCacheDirectives {
// Uses Header.Values to handle multiple header lines; matches names
// case-insensitively per RFC 9111 §5.2.
func parseCacheControl(h http.Header) cacheDirectives {
d := cacheDirectives{maxAge: -1, sMaxAge: -1}
d := cacheDirectives{maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

please add line breaks to make it more readable

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

done

for _, line := range h.Values("Cache-Control") {
for token := range strings.SplitSeq(line, ",") {
parts := strings.SplitN(strings.TrimSpace(token), "=", 2)
Expand Down Expand Up @@ -98,6 +99,12 @@ func parseCacheControl(h http.Header) cacheDirectives {
d.sMaxAge = v
}
}
case "stale-while-revalidate":
if len(parts) == 2 {
if v, err := strconv.ParseInt(strings.TrimSpace(parts[1]), 10, 64); err == nil {
d.staleWhileRevalidate = v
}
}
}
}
}
Expand Down
30 changes: 17 additions & 13 deletions filters/cache/cache_control_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,19 +59,23 @@ func TestParseCacheControl(t *testing.T) {
header http.Header
want cacheDirectives
}{
{"no-store", http.Header{"Cache-Control": {"no-store"}}, cacheDirectives{noStore: true, maxAge: -1, sMaxAge: -1}},
{"no-cache", http.Header{"Cache-Control": {"no-cache"}}, cacheDirectives{noCache: true, maxAge: -1, sMaxAge: -1}},
{"private", http.Header{"Cache-Control": {"private"}}, cacheDirectives{private: true, maxAge: -1, sMaxAge: -1}},
{"must-revalidate", http.Header{"Cache-Control": {"must-revalidate"}}, cacheDirectives{mustRevalidate: true, maxAge: -1, sMaxAge: -1}},
{"comma-separated", http.Header{"Cache-Control": {"no-store, must-revalidate"}}, cacheDirectives{noStore: true, mustRevalidate: true, maxAge: -1, sMaxAge: -1}},
{"multiple lines", http.Header{"Cache-Control": {"no-cache", "must-revalidate"}}, cacheDirectives{noCache: true, mustRevalidate: true, maxAge: -1, sMaxAge: -1}},
{"case-insensitive", http.Header{"Cache-Control": {"NO-STORE"}}, cacheDirectives{noStore: true, maxAge: -1, sMaxAge: -1}},
{"value suffix stripped", http.Header{"Cache-Control": {`no-cache="x-private"`}}, cacheDirectives{noCache: true, maxAge: -1, sMaxAge: -1}},
{"empty", http.Header{}, cacheDirectives{maxAge: -1, sMaxAge: -1}},
{"max-age=3600", http.Header{"Cache-Control": {"max-age=3600"}}, cacheDirectives{maxAge: 3600, sMaxAge: -1}},
{"s-maxage=60", http.Header{"Cache-Control": {"s-maxage=60"}}, cacheDirectives{maxAge: -1, sMaxAge: 60}},
{"max-age=3.4", http.Header{"Cache-Control": {"max-age=3.4"}}, cacheDirectives{maxAge: -1, sMaxAge: -1}}, // malformed: ParseInt fails, sentinel unchanged
{"max-age=0,max-age=5", http.Header{"Cache-Control": {"max-age=0, max-age=5"}}, cacheDirectives{maxAge: 5, sMaxAge: -1}}, // last-write-wins; no duplicate guard
{"no-store", http.Header{"Cache-Control": {"no-store"}}, cacheDirectives{noStore: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"no-cache", http.Header{"Cache-Control": {"no-cache"}}, cacheDirectives{noCache: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"private", http.Header{"Cache-Control": {"private"}}, cacheDirectives{private: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"must-revalidate", http.Header{"Cache-Control": {"must-revalidate"}}, cacheDirectives{mustRevalidate: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"comma-separated", http.Header{"Cache-Control": {"no-store, must-revalidate"}}, cacheDirectives{noStore: true, mustRevalidate: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"multiple lines", http.Header{"Cache-Control": {"no-cache", "must-revalidate"}}, cacheDirectives{noCache: true, mustRevalidate: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"case-insensitive", http.Header{"Cache-Control": {"NO-STORE"}}, cacheDirectives{noStore: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"value suffix stripped", http.Header{"Cache-Control": {`no-cache="x-private"`}}, cacheDirectives{noCache: true, maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"empty", http.Header{}, cacheDirectives{maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"max-age=3600", http.Header{"Cache-Control": {"max-age=3600"}}, cacheDirectives{maxAge: 3600, sMaxAge: -1, staleWhileRevalidate: -1}},
{"s-maxage=60", http.Header{"Cache-Control": {"s-maxage=60"}}, cacheDirectives{maxAge: -1, sMaxAge: 60, staleWhileRevalidate: -1}},
{"max-age=3.4", http.Header{"Cache-Control": {"max-age=3.4"}}, cacheDirectives{maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}}, // malformed: ParseInt fails, sentinel unchanged
{"max-age=0,max-age=5", http.Header{"Cache-Control": {"max-age=0, max-age=5"}}, cacheDirectives{maxAge: 5, sMaxAge: -1, staleWhileRevalidate: -1}}, // last-write-wins; no duplicate guard
{"stale-while-revalidate=3600", http.Header{"Cache-Control": {"stale-while-revalidate=3600"}}, cacheDirectives{maxAge: -1, sMaxAge: -1, staleWhileRevalidate: 3600}},
{"stale-while-revalidate=0", http.Header{"Cache-Control": {"stale-while-revalidate=0"}}, cacheDirectives{maxAge: -1, sMaxAge: -1, staleWhileRevalidate: 0}},
{"stale-while-revalidate malformed", http.Header{"Cache-Control": {"stale-while-revalidate=bad"}}, cacheDirectives{maxAge: -1, sMaxAge: -1, staleWhileRevalidate: -1}},
{"max-age and stale-while-revalidate", http.Header{"Cache-Control": {"max-age=300, stale-while-revalidate=60"}}, cacheDirectives{maxAge: 300, sMaxAge: -1, staleWhileRevalidate: 60}},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
Expand Down
121 changes: 104 additions & 17 deletions filters/cache/filter.go
Original file line number Diff line number Diff line change
Expand Up @@ -246,7 +246,7 @@ func (s *cacheSpec) revalidationWorker() {
if job.filter != nil {
s.metrics.MeasureSince("cache.reval_wait_duration", job.enqueuedAt)
start := time.Now()
job.filter.doRevalidate(job.key, job.req)
job.filter.doRevalidate(job.key, job.req, job.backendURL)
s.metrics.MeasureSince("cache.reval_duration", start)
}
case <-s.ctx.Done():
Expand Down Expand Up @@ -277,6 +277,7 @@ func (s *cacheSpec) updateMetrics() {
type revalJob struct {
key string
req *http.Request // cloned via Request.Clone
backendURL string // non-empty: dial here directly instead of looping through listenAddr
filter *cacheFilter // instance whose doRevalidate to call
enqueuedAt time.Time // wall-clock time the job entered the queue; used to measure wait time
}
Expand Down Expand Up @@ -411,11 +412,13 @@ func (f *cacheFilter) Request(ctx filters.FilterContext) {
Request: ctx.Request(), // link response to originating request per net/http convention
}
ctx.Serve(notModified)
f.enqueueRevalidation(key, ctx.Request())
revalReq, backendURL := revalidationDispatch(ctx, f.rfcMode)
f.enqueueRevalidation(key, revalReq, backendURL)
return
}
ctx.Serve(headBodyOmitted(method, rsp))
f.enqueueRevalidation(key, ctx.Request())
revalReq, backendURL := revalidationDispatch(ctx, f.rfcMode)
f.enqueueRevalidation(key, revalReq, backendURL)
return
}

Expand Down Expand Up @@ -516,10 +519,7 @@ func (f *cacheFilter) coalesce(ctx filters.FilterContext, key string) {
}, nil
}
ttl, shouldStore := f.resolveTTL(resp.StatusCode, resp.Header, directives)
swr := f.swrWindow
if resp.StatusCode != http.StatusOK {
swr = 0
}
swr := f.resolveSWR(resp.StatusCode, directives)
cia := correctedInitialAge(requestTime, responseTime, resp.Header)
coalescedHeader := resp.Header.Clone()
stripHopByHop(coalescedHeader)
Expand Down Expand Up @@ -679,6 +679,7 @@ func (f *cacheFilter) Response(ctx filters.FilterContext) {
if !shouldStore {
return
}
swr := f.resolveSWR(rsp.StatusCode, directives)

if ctx.Request().Header.Get("Authorization") != "" && !directives.public && !directives.mustRevalidate {
return
Expand Down Expand Up @@ -706,7 +707,7 @@ func (f *cacheFilter) Response(ctx filters.FilterContext) {
sentinel := &Entry{
CreatedAt: time.Now(),
TTL: ttl,
StaleWhileRevalidate: f.swrWindow,
StaleWhileRevalidate: swr,
VaryHeaders: varyNames,
}
if err := f.storage.Set(ctx.Request().Context(), "vary:"+baseKey, sentinel); err != nil {
Expand All @@ -715,10 +716,6 @@ func (f *cacheFilter) Response(ctx filters.FilterContext) {
}
}

swr := f.swrWindow
if rsp.StatusCode != http.StatusOK {
swr = 0
}
responseTime := time.Now()
var requestTime time.Time
if rt, ok := ctx.StateBag()[stateBagRequestTime].(time.Time); ok {
Expand Down Expand Up @@ -749,15 +746,67 @@ func (f *cacheFilter) Response(ctx filters.FilterContext) {
}
}

// revalidationDispatch decides how background revalidation should reach the
// origin: directly at the route's resolved backend when safe (force mode,
// static network backend — ctx.BackendUrl() non-empty), or via self-loopback
// through skipper's own listener otherwise (RFC mode, or a load-balanced /
// dynamic backend with no single resolvable URL — BackendUrl() is empty for
// both per its own doc comment).
//
// The two branches deliberately use different request objects. Direct
// dispatch uses ctx.Request() — already transformed by earlier filters
// (modPath, setRequestHeader, ...) into exactly what should be sent to the
// backend. Using the pre-filter request here would be wrong: it would send
// the client-facing path/host straight to the backend, bypassing whatever
// those filters did. The self-loopback fallback still needs revalidationRequest
// (ctx.OriginalRequest(), falling back to ctx.Request()) because it has to
// re-match a route, which requires the pre-mutation path.
//
// Note: RFC mode always uses self-loopback, never direct dispatch, even when
// ctx.BackendUrl() is non-empty. A Response()-filter positioned after cache()
// could rewrite Cache-Control; only the self-loopback path re-runs the full
// filter chain and sees that rewritten value before doRevalidate reads it for
// TTL and stale-while-revalidate purposes (resolveTTL, resolveSWR). Direct
// dispatch would read the pre-rewrite upstream headers instead.
func revalidationDispatch(ctx filters.FilterContext, rfcMode bool) (req *http.Request, backendURL string) {
if !rfcMode {
if b := ctx.BackendUrl(); b != "" {
return ctx.Request(), b
}
}
return revalidationRequest(ctx), ""
}

// revalidationRequest returns the request to replay for background revalidation.
// It must be ctx.OriginalRequest(), not ctx.Request(): by the time this filter's
// Request() runs, earlier filters (e.g. modPath stripping a path prefix before
// forwarding) may have already mutated ctx.Request() in place. doRevalidate loops
// the request back through skipper's own listener so the full filter chain reruns
// on it; replaying the already-mutated request can no longer match the route that
// produced that mutation, so revalidation permanently fails with a routing error
// and a stale entry (e.g. a cached error) never gets refreshed.
// ctx.OriginalRequest() can be nil per its interface contract, so fall back to
// ctx.Request() rather than passing nil into enqueueRevalidation.
func revalidationRequest(ctx filters.FilterContext) *http.Request {
if orig := ctx.OriginalRequest(); orig != nil {
return orig
}
return ctx.Request()
}

// enqueueRevalidation sends a revalidation job to the background worker.
// The request is cloned in the calling goroutine before orig is released.
// backendURL, when non-empty, tells doRevalidate to dispatch directly to that
// backend instead of looping back through skipper's own listener (see
// revalidationDispatch).
// If the queue is full the job is dropped and reval_dropped is incremented.
// The closure captures f.doRevalidate so the spec-level worker respects this route's config.
func (f *cacheFilter) enqueueRevalidation(key string, orig *http.Request) {
func (f *cacheFilter) enqueueRevalidation(key string, orig *http.Request, backendURL string) {
cloned := orig.Clone(context.Background())
job := revalJob{
key: key,
req: cloned,
backendURL: backendURL,
filter: f,
enqueuedAt: time.Now(),
}
Expand All @@ -773,11 +822,29 @@ func (f *cacheFilter) enqueueRevalidation(key string, orig *http.Request) {
// doRevalidate revalidates key against the upstream. It sends a conditional
// request (If-None-Match / If-Modified-Since) when the stored entry carries
// validators; a 304 response reuses the stored payload and merges new headers.
func (f *cacheFilter) doRevalidate(key string, req *http.Request) {
// backendURL, when non-empty, is dialed directly (incrementing
// cache.reval_backend_dispatch, since that call bypasses skipper's own proxy
// pipeline and so isn't captured by the usual backend/access-log metrics);
// otherwise the request loops back through skipper's own listener
// (f.listenAddr) so the full filter chain reruns on it, including the
// standard backend metrics and access log for that inner hop.
func (f *cacheFilter) doRevalidate(key string, req *http.Request, backendURL string) {
f.revalSF.Do(key, func() (any, error) { //nolint:errcheck
req.Header.Set(revalidateHeader, "1")
req.URL.Scheme = "http"
req.URL.Host = f.listenAddr
if backendURL != "" {
u, err := url.Parse(backendURL)
if err != nil {
f.metrics.IncCounter("cache.reval_error")
log.WithFields(log.Fields{
"backendURL": backendURL,
}).WithError(err).Warn("cache: invalid backend URL for direct revalidation dispatch")
return nil, nil
}
req.URL.Scheme, req.URL.Host = u.Scheme, u.Host
f.metrics.IncCounter("cache.reval_backend_dispatch")
} else {
req.URL.Scheme, req.URL.Host = "http", f.listenAddr
}
req.RequestURI = ""

if stored, err := f.storage.Get(context.Background(), key); err == nil && stored != nil {
Expand Down Expand Up @@ -852,7 +919,7 @@ func (f *cacheFilter) doRevalidate(key string, req *http.Request) {
Payload: body,
CreatedAt: responseTime,
TTL: ttl,
StaleWhileRevalidate: f.swrWindow,
StaleWhileRevalidate: f.resolveSWR(statusCode, directives),
StaleIfError: f.staleIfError,
ETag: responseHeader.Get("ETag"),
LastModified: responseHeader.Get("Last-Modified"),
Expand Down Expand Up @@ -922,6 +989,26 @@ func (f *cacheFilter) resolveTTL(statusCode int, header http.Header, directives
return ttl, true
}

// resolveSWR decides the stale-while-revalidate window for a stored entry.
//
// Force mode (rfcMode==false): operator swrWindow is authoritative, matching
// resolveTTL's force-mode behavior.
//
// RFC mode (rfcMode==true): honors the response's stale-while-revalidate
// directive (RFC 5861) when present; otherwise there is no SWR window.
func (f *cacheFilter) resolveSWR(statusCode int, directives cacheDirectives) time.Duration {
if statusCode != http.StatusOK {
return 0
}
if !f.rfcMode {
return f.swrWindow
}
if directives.staleWhileRevalidate >= 0 {
return time.Duration(directives.staleWhileRevalidate) * time.Second
}
return 0
}

// cacheKey builds a deterministic cache key from the route ID and request.
// routeID is included so entries from different routes never collide when all
// routes share the same LRUStorage instance.
Expand Down
Loading
Loading