pingora_cache/lib.rs
1// Copyright 2026 Cloudflare, Inc.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! The HTTP caching layer for proxies.
16
17#![allow(clippy::new_without_default)]
18
19use http::{method::Method, request::Parts as ReqHeader, response::Parts as RespHeader};
20use key::{CacheHashKey, CompactCacheKey, HashBinary};
21use lock::WritePermit;
22use log::warn;
23use pingora_error::Result;
24use pingora_http::ResponseHeader;
25use pingora_timeout::timeout;
26use std::time::{Duration, Instant, SystemTime};
27use storage::MissFinishType;
28use strum::IntoStaticStr;
29use trace::{CacheTraceCTX, Span, Tag};
30
31pub mod cache_control;
32pub mod eviction;
33pub mod filters;
34pub mod hashtable;
35pub mod key;
36pub mod lock;
37pub mod max_file_size;
38mod memory;
39pub mod meta;
40pub mod predictor;
41pub mod put;
42pub mod storage;
43pub mod trace;
44mod variance;
45
46use crate::max_file_size::MaxFileSizeTracker;
47pub use key::CacheKey;
48use lock::{CacheKeyLockImpl, LockStatus, Locked};
49pub use memory::MemCache;
50pub use meta::{set_compression_dict_content, set_compression_dict_path};
51pub use meta::{CacheMeta, CacheMetaDefaults};
52pub use storage::{HitHandler, MissHandler, PurgeType, Storage};
53pub use variance::VarianceBuilder;
54
55pub mod prelude {}
56
57/// The state machine for http caching
58///
59/// This object is used to handle the state and transitions for HTTP caching through the life of a
60/// request.
61pub struct HttpCache {
62 phase: CachePhase,
63 // Box the rest so that a disabled HttpCache struct is small
64 inner: Option<Box<HttpCacheInner>>,
65 digest: HttpCacheDigest,
66}
67
68/// This reflects the phase of HttpCache during the lifetime of a request
69#[derive(Clone, Copy, Debug, PartialEq, Eq)]
70pub enum CachePhase {
71 /// Cache disabled, with reason (NeverEnabled if never explicitly used)
72 Disabled(NoCacheReason),
73 /// Cache enabled but nothing is set yet
74 Uninit,
75 /// Cache was enabled, the request decided not to use it
76 // HttpCache.inner_enabled is kept
77 Bypass,
78 /// Awaiting the cache key to be generated
79 CacheKey,
80 /// Cache hit
81 Hit,
82 /// No cached asset is found
83 Miss,
84 /// A staled (expired) asset is found
85 Stale,
86 /// A staled (expired) asset was found, but another request is revalidating it
87 StaleUpdating,
88 /// A staled (expired) asset was found, so a fresh one was fetched
89 Expired,
90 /// A staled (expired) asset was found, and it was revalidated to be fresh
91 Revalidated,
92 /// Revalidated, but deemed uncacheable, so we do not freshen it
93 RevalidatedNoCache(NoCacheReason),
94}
95
96impl CachePhase {
97 /// Convert [CachePhase] as `str`, for logging and debugging.
98 pub fn as_str(&self) -> &'static str {
99 match self {
100 CachePhase::Disabled(_) => "disabled",
101 CachePhase::Uninit => "uninitialized",
102 CachePhase::Bypass => "bypass",
103 CachePhase::CacheKey => "key",
104 CachePhase::Hit => "hit",
105 CachePhase::Miss => "miss",
106 CachePhase::Stale => "stale",
107 CachePhase::StaleUpdating => "stale-updating",
108 CachePhase::Expired => "expired",
109 CachePhase::Revalidated => "revalidated",
110 CachePhase::RevalidatedNoCache(_) => "revalidated-nocache",
111 }
112 }
113}
114
115/// The possible reasons for not caching
116#[derive(Copy, Clone, Debug, PartialEq, Eq)]
117pub enum NoCacheReason {
118 /// Caching is not enabled to begin with
119 NeverEnabled,
120 /// Origin directives indicated this was not cacheable
121 OriginNotCache,
122 /// Response size was larger than the cache's configured maximum asset size
123 ResponseTooLarge,
124 /// Disabling caching due to unknown body size and previously exceeding maximum asset size;
125 /// the asset is otherwise cacheable, but cache needs to confirm the final size of the asset
126 /// before it can mark it as cacheable again.
127 PredictedResponseTooLarge,
128 /// Due to internal caching storage error
129 StorageError,
130 /// Due to other types of internal issues
131 InternalError,
132 /// will be cacheable but skip cache admission now
133 ///
134 /// This happens when the cache predictor predicted that this request is not cacheable, but
135 /// the response turns out to be OK to cache. However, it might be too large to re-enable caching
136 /// for this request
137 Deferred,
138 /// Due to the proxy upstream filter declining the current request from going upstream
139 DeclinedToUpstream,
140 /// Due to the upstream being unreachable or otherwise erroring during proxying
141 UpstreamError,
142 /// The writer of the cache lock sees that the request is not cacheable (Could be OriginNotCache)
143 CacheLockGiveUp,
144 /// This request waited too long for the writer of the cache lock to finish, so this request will
145 /// fetch from the origin without caching
146 CacheLockTimeout,
147 /// Other custom defined reasons
148 Custom(&'static str),
149}
150
151impl NoCacheReason {
152 /// Convert [NoCacheReason] as `str`, for logging and debugging.
153 pub fn as_str(&self) -> &'static str {
154 use NoCacheReason::*;
155 match self {
156 NeverEnabled => "NeverEnabled",
157 OriginNotCache => "OriginNotCache",
158 ResponseTooLarge => "ResponseTooLarge",
159 PredictedResponseTooLarge => "PredictedResponseTooLarge",
160 StorageError => "StorageError",
161 InternalError => "InternalError",
162 Deferred => "Deferred",
163 DeclinedToUpstream => "DeclinedToUpstream",
164 UpstreamError => "UpstreamError",
165 CacheLockGiveUp => "CacheLockGiveUp",
166 CacheLockTimeout => "CacheLockTimeout",
167 Custom(s) => s,
168 }
169 }
170}
171
172/// Information collected about the caching operation that will not be cleared
173#[derive(Debug, Default)]
174pub struct HttpCacheDigest {
175 pub lock_duration: Option<Duration>,
176 // time spent in cache lookup and reading the header
177 pub lookup_duration: Option<Duration>,
178}
179
180/// Convenience function to add a duration to an optional duration
181fn add_duration_to_opt(target_opt: &mut Option<Duration>, to_add: Duration) {
182 *target_opt = Some(target_opt.map_or(to_add, |existing| existing + to_add));
183}
184
185impl HttpCacheDigest {
186 fn add_lookup_duration(&mut self, extra_lookup_duration: Duration) {
187 add_duration_to_opt(&mut self.lookup_duration, extra_lookup_duration)
188 }
189
190 fn add_lock_duration(&mut self, extra_lock_duration: Duration) {
191 add_duration_to_opt(&mut self.lock_duration, extra_lock_duration)
192 }
193}
194
195/// Response cacheable decision
196///
197///
198#[derive(Debug)]
199pub enum RespCacheable {
200 Cacheable(CacheMeta),
201 Uncacheable(NoCacheReason),
202}
203
204impl RespCacheable {
205 /// Whether it is cacheable
206 #[inline]
207 pub fn is_cacheable(&self) -> bool {
208 matches!(*self, Self::Cacheable(_))
209 }
210
211 /// Unwrap [RespCacheable] to get the [CacheMeta] stored
212 /// # Panic
213 /// Panic when this object is not cacheable. Check [Self::is_cacheable()] first.
214 pub fn unwrap_meta(self) -> CacheMeta {
215 match self {
216 Self::Cacheable(meta) => meta,
217 Self::Uncacheable(_) => panic!("expected Cacheable value"),
218 }
219 }
220}
221
222/// Indicators of which level of cache freshness logic to force apply to an asset.
223///
224/// For example, should an existing fresh asset be revalidated or re-retrieved altogether.
225#[derive(Debug, Clone, Copy, PartialEq, Eq)]
226pub enum ForcedFreshness {
227 /// Indicates the asset should be considered stale and revalidated
228 ForceExpired,
229
230 /// Indicates the asset should be considered absent and treated like a miss
231 /// instead of a hit
232 ForceMiss,
233
234 /// Indicates the asset should be considered fresh despite possibly being stale
235 ForceFresh,
236}
237
238/// Freshness state of cache hit asset
239///
240///
241#[derive(Debug, Copy, Clone, IntoStaticStr, PartialEq, Eq)]
242#[strum(serialize_all = "snake_case")]
243pub enum HitStatus {
244 /// The asset's freshness directives indicate it has expired
245 Expired,
246
247 /// The asset was marked as expired, and should be treated as stale
248 ForceExpired,
249
250 /// The asset was marked as absent, and should be treated as a miss
251 ForceMiss,
252
253 /// An error occurred while processing the asset, so it should be treated as
254 /// a miss
255 FailedHitFilter,
256
257 /// The asset is not expired
258 Fresh,
259
260 /// Asset exists but is expired, forced to be a hit
261 ForceFresh,
262}
263
264impl HitStatus {
265 /// For displaying cache hit status
266 pub fn as_str(&self) -> &'static str {
267 self.into()
268 }
269
270 /// Whether cached asset can be served as fresh
271 pub fn is_fresh(&self) -> bool {
272 *self == HitStatus::Fresh || *self == HitStatus::ForceFresh
273 }
274
275 /// Check whether the hit status should be treated as a miss. A forced miss
276 /// is obviously treated as a miss. A hit-filter failure is treated as a
277 /// miss because we can't use the asset as an actual hit. If we treat it as
278 /// expired, we still might not be able to use it even if revalidation
279 /// succeeds.
280 pub fn is_treated_as_miss(self) -> bool {
281 matches!(self, HitStatus::ForceMiss | HitStatus::FailedHitFilter)
282 }
283}
284
285pub struct LockCtx {
286 pub lock: Option<Locked>,
287 pub cache_lock: &'static CacheKeyLockImpl,
288 pub wait_timeout: Option<Duration>,
289}
290
291// Fields like storage handlers that are needed only when cache is enabled (or bypassing).
292struct HttpCacheInnerEnabled {
293 pub meta: Option<CacheMeta>,
294 // when set, even if an asset exists, it would only be considered valid after this timestamp
295 pub valid_after: Option<SystemTime>,
296 pub miss_handler: Option<MissHandler>,
297 pub body_reader: Option<HitHandler>,
298 pub storage: &'static (dyn storage::Storage + Sync), // static for now
299 pub eviction: Option<&'static (dyn eviction::EvictionManager + Sync)>,
300 pub lock_ctx: Option<LockCtx>,
301 pub traces: trace::CacheTraceCTX,
302}
303
304struct HttpCacheInner {
305 // Prefer adding fields to InnerEnabled if possible, these fields are released
306 // when cache is disabled.
307 // If fields are needed after cache disablement, add directly to Inner.
308 pub enabled_ctx: Option<Box<HttpCacheInnerEnabled>>,
309 pub key: Option<CacheKey>,
310 // when set, an asset will be rejected from the cache if it exceeds configured size in bytes
311 pub max_file_size_tracker: Option<MaxFileSizeTracker>,
312 pub predictor: Option<&'static (dyn predictor::CacheablePredictor + Sync)>,
313}
314
315#[derive(Debug, Default)]
316#[non_exhaustive]
317pub struct CacheOptionOverrides {
318 pub wait_timeout: Option<Duration>,
319}
320
321impl HttpCache {
322 /// Create a new [HttpCache].
323 ///
324 /// Caching is not enabled by default.
325 pub fn new() -> Self {
326 HttpCache {
327 phase: CachePhase::Disabled(NoCacheReason::NeverEnabled),
328 inner: None,
329 digest: HttpCacheDigest::default(),
330 }
331 }
332
333 /// Whether the cache is enabled
334 pub fn enabled(&self) -> bool {
335 !matches!(self.phase, CachePhase::Disabled(_) | CachePhase::Bypass)
336 }
337
338 /// Whether the cache is being bypassed
339 pub fn bypassing(&self) -> bool {
340 matches!(self.phase, CachePhase::Bypass)
341 }
342
343 /// Return the [CachePhase]
344 pub fn phase(&self) -> CachePhase {
345 self.phase
346 }
347
348 /// Whether anything was fetched from the upstream
349 ///
350 /// This essentially checks all possible [CachePhase] who need to contact the upstream server
351 pub fn upstream_used(&self) -> bool {
352 use CachePhase::*;
353 match self.phase {
354 Disabled(_) | Bypass | Miss | Expired | Revalidated | RevalidatedNoCache(_) => true,
355 Hit | Stale | StaleUpdating => false,
356 Uninit | CacheKey => false, // invalid states for this call, treat them as false to keep it simple
357 }
358 }
359
360 /// Check whether the backend storage is the type `T`.
361 pub fn storage_type_is<T: 'static>(&self) -> bool {
362 self.inner
363 .as_ref()
364 .and_then(|inner| {
365 inner
366 .enabled_ctx
367 .as_ref()
368 .and_then(|ie| ie.storage.as_any().downcast_ref::<T>())
369 })
370 .is_some()
371 }
372
373 /// Release the cache lock if the current request is a cache writer.
374 ///
375 /// Generally callers should prefer using `disable` when a cache lock should be released
376 /// due to an error to clear all cache context. This function is for releasing the cache lock
377 /// while still keeping the cache around for reading, e.g. when serving stale.
378 pub fn release_write_lock(&mut self, reason: NoCacheReason) {
379 use NoCacheReason::*;
380 if let Some(inner) = self.inner.as_mut() {
381 if let Some(lock_ctx) = inner
382 .enabled_ctx
383 .as_mut()
384 .and_then(|ie| ie.lock_ctx.as_mut())
385 {
386 let lock = lock_ctx.lock.take();
387 if let Some(Locked::Write(permit)) = lock {
388 let lock_status = match reason {
389 // let the next request try to fetch it
390 InternalError | StorageError | Deferred | UpstreamError => {
391 LockStatus::TransientError
392 }
393 // depends on why the proxy upstream filter declined the request,
394 // for now still allow next request try to acquire to avoid thundering herd
395 DeclinedToUpstream => LockStatus::TransientError,
396 // no need for the lock anymore
397 OriginNotCache | ResponseTooLarge | PredictedResponseTooLarge => {
398 LockStatus::GiveUp
399 }
400 Custom(reason) => lock_ctx.cache_lock.custom_lock_status(reason),
401 // should never happen, NeverEnabled shouldn't hold a lock
402 NeverEnabled => panic!("NeverEnabled holds a write lock"),
403 CacheLockGiveUp | CacheLockTimeout => {
404 panic!("CacheLock* are for cache lock readers only")
405 }
406 };
407 lock_ctx
408 .cache_lock
409 .release(inner.key.as_ref().unwrap(), permit, lock_status);
410 }
411 }
412 }
413 }
414
415 /// Disable caching
416 pub fn disable(&mut self, reason: NoCacheReason) {
417 // XXX: compile type enforce?
418 assert!(
419 reason != NoCacheReason::NeverEnabled,
420 "NeverEnabled not allowed as a disable reason"
421 );
422 match self.phase {
423 CachePhase::Disabled(old_reason) => {
424 // replace reason
425 if old_reason == NoCacheReason::NeverEnabled {
426 // safeguard, don't allow replacing NeverEnabled as a reason
427 // TODO: can be promoted to assertion once confirmed nothing is attempting this
428 warn!("Tried to replace cache NeverEnabled with reason: {reason:?}");
429 return;
430 }
431 self.phase = CachePhase::Disabled(reason);
432 }
433 _ => {
434 self.phase = CachePhase::Disabled(reason);
435 self.release_write_lock(reason);
436 // enabled_ctx will be cleared out
437 #[cfg_attr(not(feature = "trace"), allow(unused_mut))]
438 let mut inner_enabled = self
439 .inner_mut()
440 .enabled_ctx
441 .take()
442 .expect("could remove enabled_ctx on disable");
443 // log initial disable reason
444 inner_enabled
445 .traces
446 .cache_span
447 .set_tag(|| trace::Tag::new("disable_reason", reason.as_str()));
448 }
449 }
450 }
451
452 /* The following methods panic when they are used in the wrong phase.
453 * This is better than returning errors as such panics are only caused by coding error, which
454 * should be fixed right away. Tokio runtime only crashes the current task instead of the whole
455 * program when these panics happen. */
456
457 /// Set the cache to bypass
458 ///
459 /// # Panic
460 /// This call is only allowed in [CachePhase::CacheKey] phase (before any cache lookup is performed).
461 /// Use it in any other phase will lead to panic.
462 pub fn bypass(&mut self) {
463 match self.phase {
464 CachePhase::CacheKey => {
465 // before cache lookup / found / miss
466 self.phase = CachePhase::Bypass;
467 self.inner_enabled_mut()
468 .traces
469 .cache_span
470 .set_tag(|| trace::Tag::new("bypassed", true));
471 }
472 _ => panic!("wrong phase to bypass HttpCache {:?}", self.phase),
473 }
474 }
475
476 /// Enable the cache
477 ///
478 /// - `storage`: the cache storage backend that implements [storage::Storage]
479 /// - `eviction`: optionally the eviction manager, without it, nothing will be evicted from the storage
480 /// - `predictor`: optionally a cache predictor. The cache predictor predicts whether something is likely
481 /// to be cacheable or not. This is useful because the proxy can apply different types of optimization to
482 /// cacheable and uncacheable requests.
483 /// - `cache_lock`: optionally a cache lock which handles concurrent lookups to the same asset. Without it
484 /// such lookups will all be allowed to fetch the asset independently.
485 pub fn enable(
486 &mut self,
487 storage: &'static (dyn storage::Storage + Sync),
488 eviction: Option<&'static (dyn eviction::EvictionManager + Sync)>,
489 predictor: Option<&'static (dyn predictor::CacheablePredictor + Sync)>,
490 cache_lock: Option<&'static CacheKeyLockImpl>,
491 option_overrides: Option<CacheOptionOverrides>,
492 ) {
493 match self.phase {
494 CachePhase::Disabled(_) => {
495 self.phase = CachePhase::Uninit;
496
497 let lock_ctx = cache_lock.map(|cache_lock| LockCtx {
498 cache_lock,
499 lock: None,
500 wait_timeout: option_overrides
501 .as_ref()
502 .and_then(|overrides| overrides.wait_timeout),
503 });
504
505 self.inner = Some(Box::new(HttpCacheInner {
506 enabled_ctx: Some(Box::new(HttpCacheInnerEnabled {
507 meta: None,
508 valid_after: None,
509 miss_handler: None,
510 body_reader: None,
511 storage,
512 eviction,
513 lock_ctx,
514 traces: CacheTraceCTX::new(),
515 })),
516 key: None,
517 max_file_size_tracker: None,
518 predictor,
519 }));
520 }
521 _ => panic!("Cannot enable already enabled HttpCache {:?}", self.phase),
522 }
523 }
524
525 /// Set the cache lock implementation.
526 /// # Panic
527 /// Must be called before a cache lock is attempted to be acquired,
528 /// i.e. in the `cache_key_callback` or `cache_hit_filter` phases.
529 pub fn set_cache_lock(
530 &mut self,
531 cache_lock: Option<&'static CacheKeyLockImpl>,
532 option_overrides: Option<CacheOptionOverrides>,
533 ) {
534 match self.phase {
535 CachePhase::Disabled(_)
536 | CachePhase::CacheKey
537 | CachePhase::Stale
538 | CachePhase::Hit => {
539 let inner_enabled = self.inner_enabled_mut();
540 if inner_enabled
541 .lock_ctx
542 .as_ref()
543 .is_some_and(|ctx| ctx.lock.is_some())
544 {
545 panic!("lock already set when resetting cache lock")
546 } else {
547 let lock_ctx = cache_lock.map(|cache_lock| LockCtx {
548 cache_lock,
549 lock: None,
550 wait_timeout: option_overrides.and_then(|overrides| overrides.wait_timeout),
551 });
552 inner_enabled.lock_ctx = lock_ctx;
553 }
554 }
555 _ => panic!("wrong phase: {:?}", self.phase),
556 }
557 }
558
559 // Enable distributed tracing
560 pub fn enable_tracing(&mut self, parent_span: trace::Span) {
561 if let Some(inner_enabled) = self.inner.as_mut().and_then(|i| i.enabled_ctx.as_mut()) {
562 inner_enabled.traces.enable(parent_span);
563 }
564 }
565
566 // Get the cache parent tracing span
567 pub fn get_cache_span(&self) -> Option<trace::SpanHandle> {
568 self.inner
569 .as_ref()
570 .and_then(|i| i.enabled_ctx.as_ref().map(|ie| ie.traces.get_cache_span()))
571 }
572
573 // Get the cache `miss` tracing span
574 pub fn get_miss_span(&self) -> Option<trace::SpanHandle> {
575 self.inner
576 .as_ref()
577 .and_then(|i| i.enabled_ctx.as_ref().map(|ie| ie.traces.get_miss_span()))
578 }
579
580 // Get the cache `hit` tracing span
581 pub fn get_hit_span(&self) -> Option<trace::SpanHandle> {
582 self.inner
583 .as_ref()
584 .and_then(|i| i.enabled_ctx.as_ref().map(|ie| ie.traces.get_hit_span()))
585 }
586
587 // shortcut to access inner fields, panic if phase is disabled
588 #[inline]
589 fn inner_enabled_mut(&mut self) -> &mut HttpCacheInnerEnabled {
590 self.inner.as_mut().unwrap().enabled_ctx.as_mut().unwrap()
591 }
592
593 #[inline]
594 fn inner_enabled(&self) -> &HttpCacheInnerEnabled {
595 self.inner.as_ref().unwrap().enabled_ctx.as_ref().unwrap()
596 }
597
598 // shortcut to access inner fields, panic if cache was never enabled
599 #[inline]
600 fn inner_mut(&mut self) -> &mut HttpCacheInner {
601 self.inner.as_mut().unwrap()
602 }
603
604 #[inline]
605 fn inner(&self) -> &HttpCacheInner {
606 self.inner.as_ref().unwrap()
607 }
608
609 /// Set the cache key
610 /// # Panic
611 /// Cache key is only allowed to be set in its own phase. Set it in other phases will cause panic.
612 pub fn set_cache_key(&mut self, key: CacheKey) {
613 match self.phase {
614 CachePhase::Uninit | CachePhase::CacheKey => {
615 self.phase = CachePhase::CacheKey;
616 self.inner_mut().key = Some(key);
617 }
618 _ => panic!("wrong phase {:?}", self.phase),
619 }
620 }
621
622 /// Return the cache key used for asset lookup
623 /// # Panic
624 /// Can only be called after the cache key is set and the cache is not disabled. Panic otherwise.
625 pub fn cache_key(&self) -> &CacheKey {
626 match self.phase {
627 CachePhase::Disabled(NoCacheReason::NeverEnabled) | CachePhase::Uninit => {
628 panic!("wrong phase {:?}", self.phase)
629 }
630 _ => self
631 .inner()
632 .key
633 .as_ref()
634 .expect("cache key should be set (set_cache_key not called?)"),
635 }
636 }
637
638 /// Return the max size allowed to be cached.
639 pub fn max_file_size_bytes(&self) -> Option<usize> {
640 assert!(
641 !matches!(
642 self.phase,
643 CachePhase::Disabled(NoCacheReason::NeverEnabled)
644 ),
645 "tried to access max file size bytes when cache never enabled"
646 );
647 self.inner()
648 .max_file_size_tracker
649 .as_ref()
650 .map(|t| t.max_file_size_bytes())
651 }
652
653 /// Set the maximum response _body_ size in bytes that will be admitted to the cache.
654 ///
655 /// Response header size should not contribute to the max file size.
656 ///
657 /// To track body bytes, call `track_bytes_for_max_file_size`.
658 pub fn set_max_file_size_bytes(&mut self, max_file_size_bytes: usize) {
659 match self.phase {
660 CachePhase::Disabled(_) => panic!("wrong phase {:?}", self.phase),
661 _ => {
662 self.inner_mut().max_file_size_tracker =
663 Some(MaxFileSizeTracker::new(max_file_size_bytes));
664 }
665 }
666 }
667
668 /// Record body bytes for the max file size tracker.
669 ///
670 /// The `bytes_len` input contributes to a cumulative body byte tracker.
671 ///
672 /// Once the cumulative body bytes exceeds the maximum allowable cache file size (as configured
673 /// by `set_max_file_size_bytes`), then the return value will be false.
674 ///
675 /// Else the return value is true as long as the max file size is not exceeded.
676 /// If max file size was not configured, the return value is always true.
677 pub fn track_body_bytes_for_max_file_size(&mut self, bytes_len: usize) -> bool {
678 // This is intended to be callable when cache has already been disabled,
679 // so that we can re-mark an asset as cacheable if the body size is under limits.
680 assert!(
681 !matches!(
682 self.phase,
683 CachePhase::Disabled(NoCacheReason::NeverEnabled)
684 ),
685 "tried to access max file size bytes when cache never enabled"
686 );
687 self.inner_mut()
688 .max_file_size_tracker
689 .as_mut()
690 .is_none_or(|t| t.add_body_bytes(bytes_len))
691 }
692
693 /// Check if the max file size has been exceeded according to max file size tracker.
694 ///
695 /// Return true if max file size was exceeded.
696 pub fn exceeded_max_file_size(&self) -> bool {
697 assert!(
698 !matches!(
699 self.phase,
700 CachePhase::Disabled(NoCacheReason::NeverEnabled)
701 ),
702 "tried to access max file size bytes when cache never enabled"
703 );
704 self.inner()
705 .max_file_size_tracker
706 .as_ref()
707 .is_some_and(|t| !t.allow_caching())
708 }
709
710 /// Set that cache is found in cache storage.
711 ///
712 /// This function is called after [Self::cache_lookup()] which returns the [CacheMeta] and
713 /// [HitHandler].
714 ///
715 /// The `hit_status` enum allows the caller to force expire assets.
716 pub fn cache_found(&mut self, meta: CacheMeta, hit_handler: HitHandler, hit_status: HitStatus) {
717 // Stale allowed because of cache lock and then retry
718 if !matches!(self.phase, CachePhase::CacheKey | CachePhase::Stale) {
719 panic!("wrong phase {:?}", self.phase)
720 }
721
722 self.phase = match hit_status {
723 HitStatus::Fresh | HitStatus::ForceFresh => CachePhase::Hit,
724 HitStatus::Expired | HitStatus::ForceExpired => CachePhase::Stale,
725 HitStatus::FailedHitFilter | HitStatus::ForceMiss => self.phase,
726 };
727
728 let phase = self.phase;
729 let inner = self.inner_mut();
730
731 let key = inner.key.as_ref().expect("key must be set on hit");
732 let inner_enabled = inner
733 .enabled_ctx
734 .as_mut()
735 .expect("cache_found must be called while cache enabled");
736
737 // The cache lock might not be set for stale hit or hits treated as
738 // misses, so we need to initialize it here
739 let stale = phase == CachePhase::Stale;
740 if stale || hit_status.is_treated_as_miss() {
741 if let Some(lock_ctx) = inner_enabled.lock_ctx.as_mut() {
742 lock_ctx.lock = Some(lock_ctx.cache_lock.lock(key, stale));
743 }
744 }
745
746 if hit_status.is_treated_as_miss() {
747 // Clear the body and meta for hits that are treated as misses
748 inner_enabled.body_reader = None;
749 inner_enabled.meta = None;
750 } else {
751 // Set the metadata appropriately for legit hits
752 inner_enabled.traces.start_hit_span(phase, hit_status);
753 inner_enabled.traces.log_meta_in_hit_span(&meta);
754 if let Some(eviction) = inner_enabled.eviction {
755 // TODO: make access() accept CacheKey
756 let cache_key = key.to_compact();
757 if hit_handler.should_count_access() {
758 let size = hit_handler.get_eviction_weight();
759 eviction.access(&cache_key, size, meta.0.internal.fresh_until);
760 }
761 }
762 inner_enabled.meta = Some(meta);
763 inner_enabled.body_reader = Some(hit_handler);
764 }
765 }
766
767 /// Mark `self` to be cache miss.
768 ///
769 /// This function is called after [Self::cache_lookup()] finds nothing or the caller decides
770 /// not to use the assets found.
771 /// # Panic
772 /// Panic in other phases.
773 pub fn cache_miss(&mut self) {
774 match self.phase {
775 // from CacheKey: set state to miss during cache lookup
776 // from Bypass: response became cacheable, set state to miss to cache
777 // from Stale: waited for cache lock, then retried and found asset was gone
778 CachePhase::CacheKey | CachePhase::Bypass | CachePhase::Stale => {
779 self.phase = CachePhase::Miss;
780 // It's possible that we've set the meta on lookup and have come back around
781 // here after not being able to acquire the cache lock, and our item has since
782 // purged or expired. We should be sure that the meta is not set in this case
783 // as there shouldn't be a meta set for cache misses.
784 self.inner_enabled_mut().meta = None;
785 self.inner_enabled_mut().traces.start_miss_span();
786 }
787 _ => panic!("wrong phase {:?}", self.phase),
788 }
789 }
790
791 /// Return the [HitHandler]
792 /// # Panic
793 /// Call this after [Self::cache_found()], panic in other phases.
794 pub fn hit_handler(&mut self) -> &mut HitHandler {
795 match self.phase {
796 CachePhase::Hit
797 | CachePhase::Stale
798 | CachePhase::StaleUpdating
799 | CachePhase::Revalidated
800 | CachePhase::RevalidatedNoCache(_) => {
801 self.inner_enabled_mut().body_reader.as_mut().unwrap()
802 }
803 _ => panic!("wrong phase {:?}", self.phase),
804 }
805 }
806
807 /// Return the body reader during a cache admission (miss/expired) which decouples the downstream
808 /// read and upstream cache write
809 pub fn miss_body_reader(&mut self) -> Option<&mut HitHandler> {
810 match self.phase {
811 CachePhase::Miss | CachePhase::Expired => {
812 let inner_enabled = self.inner_enabled_mut();
813 if inner_enabled.storage.support_streaming_partial_write() {
814 inner_enabled.body_reader.as_mut()
815 } else {
816 // body_reader could be set even when the storage doesn't support streaming
817 // Expired cache would have the reader set.
818 None
819 }
820 }
821 _ => None,
822 }
823 }
824
825 /// Return whether the underlying storage backend supports streaming partial write.
826 ///
827 /// Returns None if cache is not enabled.
828 pub fn support_streaming_partial_write(&self) -> Option<bool> {
829 self.inner.as_ref().and_then(|inner| {
830 inner
831 .enabled_ctx
832 .as_ref()
833 .map(|c| c.storage.support_streaming_partial_write())
834 })
835 }
836
837 /// Call this when cache hit is fully read.
838 ///
839 /// This call will release resource if any and log the timing in tracing if set.
840 /// # Panic
841 /// Panic in phases where there is no cache hit.
842 pub async fn finish_hit_handler(&mut self) -> Result<()> {
843 match self.phase {
844 CachePhase::Hit
845 | CachePhase::Miss
846 | CachePhase::Expired
847 | CachePhase::Stale
848 | CachePhase::StaleUpdating
849 | CachePhase::Revalidated
850 | CachePhase::RevalidatedNoCache(_) => {
851 let inner = self.inner_mut();
852 let inner_enabled = inner.enabled_ctx.as_mut().expect("cache enabled");
853 if inner_enabled.body_reader.is_none() {
854 // already finished, we allow calling this function more than once
855 return Ok(());
856 }
857 let body_reader = inner_enabled.body_reader.take().unwrap();
858 let key = inner.key.as_ref().unwrap();
859 let result = body_reader
860 .finish(
861 inner_enabled.storage,
862 key,
863 &inner_enabled.traces.hit_span.handle(),
864 )
865 .await;
866 inner_enabled.traces.finish_hit_span();
867 result
868 }
869 _ => panic!("wrong phase {:?}", self.phase),
870 }
871 }
872
873 /// Set the [MissHandler] according to cache_key and meta, can only call once
874 pub async fn set_miss_handler(&mut self) -> Result<()> {
875 match self.phase {
876 // set_miss_handler() needs to be called after set_cache_meta() (which change Stale to Expire).
877 // This is an artificial rule to enforce the state transitions
878 CachePhase::Miss | CachePhase::Expired => {
879 let inner = self.inner_mut();
880 let inner_enabled = inner
881 .enabled_ctx
882 .as_mut()
883 .expect("cache enabled on miss and expired");
884 if inner_enabled.miss_handler.is_some() {
885 panic!("write handler is already set")
886 }
887 let meta = inner_enabled.meta.as_ref().unwrap();
888 let key = inner.key.as_ref().unwrap();
889 let miss_handler = inner_enabled
890 .storage
891 .get_miss_handler(key, meta, &inner_enabled.traces.get_miss_span())
892 .await?;
893
894 inner_enabled.miss_handler = Some(miss_handler);
895
896 if inner_enabled.storage.support_streaming_partial_write() {
897 // If a reader can access partial write, the cache lock can be released here
898 // to let readers start reading the body.
899 if let Some(lock_ctx) = inner_enabled.lock_ctx.as_mut() {
900 let lock = lock_ctx.lock.take();
901 if let Some(Locked::Write(permit)) = lock {
902 lock_ctx.cache_lock.release(key, permit, LockStatus::Done);
903 }
904 }
905 // Downstream read and upstream write can be decoupled
906 let body_reader = inner_enabled
907 .storage
908 .lookup_streaming_write(
909 key,
910 inner_enabled
911 .miss_handler
912 .as_ref()
913 .expect("miss handler already set")
914 .streaming_write_tag(),
915 &inner_enabled.traces.get_miss_span(),
916 )
917 .await?;
918
919 if let Some((_meta, body_reader)) = body_reader {
920 inner_enabled.body_reader = Some(body_reader);
921 } else {
922 // body_reader should exist now because streaming_partial_write is to support it
923 panic!("unable to get body_reader for {:?}", meta);
924 }
925 }
926 Ok(())
927 }
928 _ => panic!("wrong phase {:?}", self.phase),
929 }
930 }
931
932 /// Return the [MissHandler] to write the response body to cache.
933 ///
934 /// `None`: the handler has not been set or already finished
935 pub fn miss_handler(&mut self) -> Option<&mut MissHandler> {
936 match self.phase {
937 CachePhase::Miss | CachePhase::Expired => {
938 self.inner_enabled_mut().miss_handler.as_mut()
939 }
940 _ => panic!("wrong phase {:?}", self.phase),
941 }
942 }
943
944 /// Finish cache admission
945 ///
946 /// If [self] is dropped without calling this, the cache admission is considered incomplete and
947 /// should be cleaned up.
948 ///
949 /// This call will also trigger eviction if set.
950 pub async fn finish_miss_handler(&mut self) -> Result<()> {
951 match self.phase {
952 CachePhase::Miss | CachePhase::Expired => {
953 let inner = self.inner_mut();
954 let inner_enabled = inner
955 .enabled_ctx
956 .as_mut()
957 .expect("cache enabled on miss and expired");
958 if inner_enabled.miss_handler.is_none() {
959 // already finished, we allow calling this function more than once
960 return Ok(());
961 }
962 let miss_handler = inner_enabled.miss_handler.take().unwrap();
963 let size = miss_handler.finish().await?;
964 let key = inner
965 .key
966 .as_ref()
967 .expect("key set by miss or expired phase");
968 if let Some(lock_ctx) = inner_enabled.lock_ctx.as_mut() {
969 let lock = lock_ctx.lock.take();
970 if let Some(Locked::Write(permit)) = lock {
971 // no need to call r.unlock() because release() will call it
972 // r is a guard to make sure the lock is unlocked when this request is dropped
973 lock_ctx.cache_lock.release(key, permit, LockStatus::Done);
974 }
975 }
976 if let Some(eviction) = inner_enabled.eviction {
977 let cache_key = key.to_compact();
978 let meta = inner_enabled.meta.as_ref().unwrap();
979 let evicted = match size {
980 MissFinishType::Created(size) => {
981 eviction.admit(cache_key, size, meta.0.internal.fresh_until)
982 }
983 MissFinishType::Appended(size, max_size) => {
984 eviction.increment_weight(&cache_key, size, max_size)
985 }
986 };
987 // actual eviction can be done async
988 let span = inner_enabled.traces.child("eviction");
989 let handle = span.handle();
990 let storage = inner_enabled.storage;
991 tokio::task::spawn(async move {
992 for item in evicted {
993 if let Err(e) = storage.purge(&item, PurgeType::Eviction, &handle).await
994 {
995 warn!("Failed to purge {item} during eviction for finish miss handler: {e}");
996 }
997 }
998 });
999 }
1000 inner_enabled.traces.finish_miss_span();
1001 Ok(())
1002 }
1003 _ => panic!("wrong phase {:?}", self.phase),
1004 }
1005 }
1006
1007 /// Set the [CacheMeta] of the cache
1008 pub fn set_cache_meta(&mut self, meta: CacheMeta) {
1009 match self.phase {
1010 // TODO: store the staled meta somewhere else for future use?
1011 CachePhase::Stale | CachePhase::Miss => {
1012 let inner_enabled = self.inner_enabled_mut();
1013 // TODO: have a separate expired span?
1014 inner_enabled.traces.log_meta_in_miss_span(&meta);
1015 inner_enabled.meta = Some(meta);
1016 }
1017 _ => panic!("wrong phase {:?}", self.phase),
1018 }
1019 if self.phase == CachePhase::Stale {
1020 self.phase = CachePhase::Expired;
1021 }
1022 }
1023
1024 /// Set the [CacheMeta] of the cache after revalidation.
1025 ///
1026 /// Certain info such as the original cache admission time will be preserved. Others will
1027 /// be replaced by the input `meta`.
1028 pub async fn revalidate_cache_meta(&mut self, mut meta: CacheMeta) -> Result<bool> {
1029 let result = match self.phase {
1030 CachePhase::Stale => {
1031 let inner = self.inner_mut();
1032 let inner_enabled = inner
1033 .enabled_ctx
1034 .as_mut()
1035 .expect("stale phase has cache enabled");
1036 // TODO: we should keep old meta in place, just use new one to update it
1037 // that requires cacheable_filter to take a mut header and just return InternalMeta
1038
1039 // update new meta with old meta's created time
1040 let old_meta = inner_enabled.meta.take().unwrap();
1041 let created = old_meta.0.internal.created;
1042 meta.0.internal.created = created;
1043 // meta.internal.updated was already set to new meta's `created`,
1044 // no need to set `updated` here
1045 // Merge old extensions with new ones. New exts take precedence if they conflict.
1046 let mut extensions = old_meta.0.extensions;
1047 extensions.extend(meta.0.extensions);
1048 meta.0.extensions = extensions;
1049
1050 inner_enabled.meta.replace(meta);
1051
1052 #[cfg_attr(not(feature = "trace"), allow(unused_mut))]
1053 let mut span = inner_enabled.traces.child("update_meta");
1054 let result = inner_enabled
1055 .storage
1056 .update_meta(
1057 inner.key.as_ref().unwrap(),
1058 inner_enabled.meta.as_ref().unwrap(),
1059 &span.handle(),
1060 )
1061 .await;
1062 span.set_tag(|| trace::Tag::new("updated", result.is_ok()));
1063
1064 // regardless of result, release the cache lock
1065 if let Some(lock_ctx) = inner_enabled.lock_ctx.as_mut() {
1066 let lock = lock_ctx.lock.take();
1067 if let Some(Locked::Write(permit)) = lock {
1068 lock_ctx.cache_lock.release(
1069 inner.key.as_ref().expect("key set by stale phase"),
1070 permit,
1071 LockStatus::Done,
1072 );
1073 }
1074 }
1075
1076 result
1077 }
1078 _ => panic!("wrong phase {:?}", self.phase),
1079 };
1080 self.phase = CachePhase::Revalidated;
1081 result
1082 }
1083
1084 /// After a successful revalidation, update certain headers for the cached asset
1085 /// such as `Etag` with the fresh response header `resp`.
1086 pub fn revalidate_merge_header(&mut self, resp: &RespHeader) -> ResponseHeader {
1087 match self.phase {
1088 CachePhase::Stale => {
1089 /*
1090 * https://datatracker.ietf.org/doc/html/rfc9110#section-15.4.5
1091 * 304 response MUST generate ... would have been sent in a 200 ...
1092 * - Content-Location, Date, ETag, and Vary
1093 * - Cache-Control and Expires...
1094 */
1095 let mut old_header = self.inner_enabled().meta.as_ref().unwrap().0.header.clone();
1096 let mut clone_header = |header_name: &'static str| {
1097 for (i, value) in resp.headers.get_all(header_name).iter().enumerate() {
1098 if i == 0 {
1099 old_header
1100 .insert_header(header_name, value)
1101 .expect("can add valid header");
1102 } else {
1103 old_header
1104 .append_header(header_name, value)
1105 .expect("can add valid header");
1106 }
1107 }
1108 };
1109 clone_header("cache-control");
1110 clone_header("expires");
1111 clone_header("cache-tag");
1112 clone_header("cdn-cache-control");
1113 clone_header("etag");
1114 // https://datatracker.ietf.org/doc/html/rfc9111#section-4.3.4
1115 // "...cache MUST update its header fields with the header fields provided in the 304..."
1116 // But if the Vary header changes, the cached response may no longer match the
1117 // incoming request.
1118 //
1119 // For simplicity, ignore changing Vary in revalidation for now.
1120 // TODO: if we support vary during revalidation, there are a few edge cases to
1121 // consider (what if Vary header appears/disappears/changes)?
1122 //
1123 // clone_header("vary");
1124 old_header
1125 }
1126 _ => panic!("wrong phase {:?}", self.phase),
1127 }
1128 }
1129
1130 /// Mark this asset uncacheable after revalidation
1131 pub fn revalidate_uncacheable(&mut self, header: ResponseHeader, reason: NoCacheReason) {
1132 match self.phase {
1133 CachePhase::Stale => {
1134 // replace cache meta header
1135 self.inner_enabled_mut().meta.as_mut().unwrap().0.header = header;
1136 // upstream request done, release write lock
1137 self.release_write_lock(reason);
1138 }
1139 _ => panic!("wrong phase {:?}", self.phase),
1140 }
1141 self.phase = CachePhase::RevalidatedNoCache(reason);
1142 // TODO: remove this asset from cache once finished?
1143 }
1144
1145 /// Mark this asset as stale, but being updated separately from this request.
1146 pub fn set_stale_updating(&mut self) {
1147 match self.phase {
1148 CachePhase::Stale => self.phase = CachePhase::StaleUpdating,
1149 _ => panic!("wrong phase {:?}", self.phase),
1150 }
1151 }
1152
1153 /// Update the variance of the [CacheMeta].
1154 ///
1155 /// Note that this process may change the lookup `key`, and eventually (when the asset is
1156 /// written to storage) invalidate other cached variants under the same primary key as the
1157 /// current asset.
1158 pub fn update_variance(&mut self, variance: Option<HashBinary>) {
1159 // If this is a cache miss, we will simply update the variance in the meta.
1160 //
1161 // If this is an expired response, we will have to consider a few cases:
1162 //
1163 // **Case 1**: Variance was absent, but caller sets it now.
1164 // We will just insert it into the meta. The current asset becomes the primary variant.
1165 // Because the current location of the asset is already the primary variant, nothing else
1166 // needs to be done.
1167 //
1168 // **Case 2**: Variance was present, but it changed or was removed.
1169 // We want the current asset to take over the primary slot, in order to invalidate all
1170 // other variants derived under the old Vary.
1171 //
1172 // **Case 3**: Variance did not change.
1173 // Nothing needs to happen.
1174 let inner = match self.phase {
1175 CachePhase::Miss | CachePhase::Expired => self.inner_mut(),
1176 _ => panic!("wrong phase {:?}", self.phase),
1177 };
1178 let inner_enabled = inner
1179 .enabled_ctx
1180 .as_mut()
1181 .expect("cache enabled on miss and expired");
1182
1183 // Update the variance in the meta
1184 if let Some(variance_hash) = variance.as_ref() {
1185 inner_enabled
1186 .meta
1187 .as_mut()
1188 .unwrap()
1189 .set_variance_key(*variance_hash);
1190 } else {
1191 inner_enabled.meta.as_mut().unwrap().remove_variance();
1192 }
1193
1194 // Change the lookup `key` if necessary, in order to admit asset into the primary slot
1195 // instead of the secondary slot.
1196 let key = inner.key.as_ref().unwrap();
1197 if let Some(old_variance) = key.get_variance_key().as_ref() {
1198 // This is a secondary variant slot.
1199 if Some(*old_variance) != variance.as_ref() {
1200 // This new variance does not match the variance in the cache key we used to look
1201 // up this asset.
1202 // Drop the cache lock to avoid leaving a dangling lock
1203 // (because we locked with the old cache key for the secondary slot)
1204 // TODO: maybe we should try to signal waiting readers to compete for the primary key
1205 // lock instead? we will not be modifying this secondary slot so it's not actually
1206 // ready for readers
1207 if let Some(lock_ctx) = inner_enabled.lock_ctx.as_mut() {
1208 if let Some(Locked::Write(permit)) = lock_ctx.lock.take() {
1209 lock_ctx.cache_lock.release(key, permit, LockStatus::Done);
1210 }
1211 }
1212 // Remove the `variance` from the `key`, so that we admit this asset into the
1213 // primary slot. (`key` is used to tell storage where to write the data.)
1214 inner.key.as_mut().unwrap().remove_variance_key();
1215 }
1216 }
1217 }
1218
1219 /// Return the [CacheMeta] of this asset
1220 ///
1221 /// # Panic
1222 /// Panic in phases which has no cache meta.
1223 pub fn cache_meta(&self) -> &CacheMeta {
1224 match self.phase {
1225 // TODO: allow in Bypass phase?
1226 CachePhase::Stale
1227 | CachePhase::StaleUpdating
1228 | CachePhase::Expired
1229 | CachePhase::Hit
1230 | CachePhase::Revalidated
1231 | CachePhase::RevalidatedNoCache(_) => self.inner_enabled().meta.as_ref().unwrap(),
1232 CachePhase::Miss => {
1233 // this is the async body read case, safe because body_reader is only set
1234 // after meta is retrieved
1235 if self.inner_enabled().body_reader.is_some() {
1236 self.inner_enabled().meta.as_ref().unwrap()
1237 } else {
1238 panic!("wrong phase {:?}", self.phase);
1239 }
1240 }
1241
1242 _ => panic!("wrong phase {:?}", self.phase),
1243 }
1244 }
1245
1246 /// Return the [CacheMeta] of this asset if any
1247 ///
1248 /// Different from [Self::cache_meta()], this function is allowed to be called in
1249 /// [CachePhase::Miss] phase where the cache meta maybe set.
1250 /// # Panic
1251 /// Panic in phases that shouldn't have cache meta.
1252 pub fn maybe_cache_meta(&self) -> Option<&CacheMeta> {
1253 match self.phase {
1254 CachePhase::Miss
1255 | CachePhase::Stale
1256 | CachePhase::StaleUpdating
1257 | CachePhase::Expired
1258 | CachePhase::Hit
1259 | CachePhase::Revalidated
1260 | CachePhase::RevalidatedNoCache(_) => self.inner_enabled().meta.as_ref(),
1261 _ => panic!("wrong phase {:?}", self.phase),
1262 }
1263 }
1264
1265 /// Return the [`CacheKey`] of this asset if any.
1266 ///
1267 /// This is allowed to be called in any phase. If the cache key callback was not called,
1268 /// this will return None.
1269 pub fn maybe_cache_key(&self) -> Option<&CacheKey> {
1270 (!matches!(
1271 self.phase(),
1272 CachePhase::Disabled(NoCacheReason::NeverEnabled) | CachePhase::Uninit
1273 ))
1274 .then(|| self.cache_key())
1275 }
1276
1277 /// Perform the cache lookup from the given cache storage with the given cache key
1278 ///
1279 /// A cache hit will return [CacheMeta] which contains the header and meta info about
1280 /// the cache as well as a [HitHandler] to read the cache hit body.
1281 /// # Panic
1282 /// Panic in other phases.
1283 pub async fn cache_lookup(&mut self) -> Result<Option<(CacheMeta, HitHandler)>> {
1284 match self.phase {
1285 // Stale is allowed here because stale-> cache_lock -> lookup again
1286 CachePhase::CacheKey | CachePhase::Stale => {
1287 let inner = self
1288 .inner
1289 .as_mut()
1290 .expect("Cache phase is checked and should have inner");
1291 let inner_enabled = inner
1292 .enabled_ctx
1293 .as_mut()
1294 .expect("Cache enabled on cache_lookup");
1295 #[cfg_attr(not(feature = "trace"), allow(unused_mut))]
1296 let mut span = inner_enabled.traces.child("lookup");
1297 let key = inner.key.as_ref().unwrap(); // safe, this phase should have cache key
1298 let now = Instant::now();
1299 let result = inner_enabled.storage.lookup(key, &span.handle()).await?;
1300 // one request may have multiple lookups
1301 self.digest.add_lookup_duration(now.elapsed());
1302 let result = result.and_then(|(meta, header)| {
1303 if let Some(ts) = inner_enabled.valid_after {
1304 if meta.created() < ts {
1305 span.set_tag(|| trace::Tag::new("not valid", true));
1306 return None;
1307 }
1308 }
1309 Some((meta, header))
1310 });
1311 if result.is_none() {
1312 if let Some(lock_ctx) = inner_enabled.lock_ctx.as_mut() {
1313 lock_ctx.lock = Some(lock_ctx.cache_lock.lock(key, false));
1314 }
1315 }
1316 span.set_tag(|| trace::Tag::new("found", result.is_some()));
1317 Ok(result)
1318 }
1319 _ => panic!("wrong phase {:?}", self.phase),
1320 }
1321 }
1322
1323 /// Update variance and see if the meta matches the current variance
1324 ///
1325 /// `cache_lookup() -> compute vary hash -> cache_vary_lookup()`
1326 /// This function allows callers to compute vary based on the initial cache hit.
1327 /// `meta` should be the ones returned from the initial cache_lookup()
1328 /// - return true if the meta is the variance.
1329 /// - return false if the current meta doesn't match the variance, need to cache_lookup() again
1330 pub fn cache_vary_lookup(&mut self, variance: HashBinary, meta: &CacheMeta) -> bool {
1331 match self.phase {
1332 // Stale is allowed here because stale-> cache_lock -> lookup again
1333 CachePhase::CacheKey | CachePhase::Stale => {
1334 let inner = self.inner_mut();
1335 // make sure that all variances found are fresher than this asset
1336 // this is because when purging all the variance, only the primary slot is deleted
1337 // the created TS of the primary is the tombstone of all the variances
1338 inner
1339 .enabled_ctx
1340 .as_mut()
1341 .expect("cache enabled")
1342 .valid_after = Some(meta.created());
1343
1344 // update vary
1345 let key = inner.key.as_mut().unwrap();
1346 // if no variance was previously set, then this is the first cache hit
1347 let is_initial_cache_hit = key.get_variance_key().is_none();
1348 key.set_variance_key(variance);
1349 let variance_binary = key.variance_bin();
1350 let matches_variance = meta.variance() == variance_binary;
1351
1352 // We should remove the variance in the lookup `key` if this is the primary variant
1353 // slot. We know this is the primary variant slot if this is the initial cache hit,
1354 // AND the variance in the `key` already matches the `meta`'s.
1355 //
1356 // For the primary variant slot, the storage backend needs to use the primary key
1357 // for both cache lookup and updating the meta. Otherwise it will look for the
1358 // asset in the wrong location during revalidation.
1359 //
1360 // We can recreate the "full" cache key by using the meta's variance, if needed.
1361 if matches_variance && is_initial_cache_hit {
1362 inner.key.as_mut().unwrap().remove_variance_key();
1363 }
1364
1365 matches_variance
1366 }
1367 _ => panic!("wrong phase {:?}", self.phase),
1368 }
1369 }
1370
1371 /// Whether this request is behind a cache lock in order to wait for another request to read the
1372 /// asset.
1373 pub fn is_cache_locked(&self) -> bool {
1374 matches!(
1375 self.inner_enabled()
1376 .lock_ctx
1377 .as_ref()
1378 .and_then(|l| l.lock.as_ref()),
1379 Some(Locked::Read(_))
1380 )
1381 }
1382
1383 /// Whether this request is the leader request to fetch the assets for itself and other requests
1384 /// behind the cache lock.
1385 pub fn is_cache_lock_writer(&self) -> bool {
1386 matches!(
1387 self.inner_enabled()
1388 .lock_ctx
1389 .as_ref()
1390 .and_then(|l| l.lock.as_ref()),
1391 Some(Locked::Write(_))
1392 )
1393 }
1394
1395 /// Take the write lock from this request to transfer it to another one.
1396 /// # Panic
1397 /// Call is_cache_lock_writer() to check first, will panic otherwise.
1398 pub fn take_write_lock(&mut self) -> (WritePermit, &'static CacheKeyLockImpl) {
1399 let lock_ctx = self
1400 .inner_enabled_mut()
1401 .lock_ctx
1402 .as_mut()
1403 .expect("take_write_lock() called without cache lock");
1404 let lock = lock_ctx
1405 .lock
1406 .take()
1407 .expect("take_write_lock() called without lock");
1408 match lock {
1409 Locked::Write(w) => (w, lock_ctx.cache_lock),
1410 Locked::Read(_) => panic!("take_write_lock() called on read lock"),
1411 }
1412 }
1413
1414 /// Set the write lock, which is usually transferred from [Self::take_write_lock()]
1415 ///
1416 /// # Panic
1417 /// Panics if cache lock was not originally configured for this request.
1418 // TODO: it may make sense to allow configuring the CacheKeyLock here too that the write permit
1419 // is associated with
1420 // (The WritePermit comes from the CacheKeyLock and should be used when releasing from the CacheKeyLock,
1421 // shouldn't be possible to give a WritePermit to a request using a different CacheKeyLock)
1422 pub fn set_write_lock(&mut self, write_lock: WritePermit) {
1423 if let Some(lock_ctx) = self.inner_enabled_mut().lock_ctx.as_mut() {
1424 lock_ctx.lock.replace(Locked::Write(write_lock));
1425 }
1426 }
1427
1428 /// Whether this request's cache hit is staled
1429 fn has_staled_asset(&self) -> bool {
1430 matches!(self.phase, CachePhase::Stale | CachePhase::StaleUpdating)
1431 }
1432
1433 /// Whether this asset is staled and stale if error is allowed
1434 pub fn can_serve_stale_error(&self) -> bool {
1435 self.has_staled_asset() && self.cache_meta().serve_stale_if_error(SystemTime::now())
1436 }
1437
1438 /// Whether this asset is staled and stale while revalidate is allowed.
1439 pub fn can_serve_stale_updating(&self) -> bool {
1440 self.has_staled_asset()
1441 && self
1442 .cache_meta()
1443 .serve_stale_while_revalidate(SystemTime::now())
1444 }
1445
1446 /// Wait for the cache read lock to be unlocked
1447 /// # Panic
1448 /// Check [Self::is_cache_locked()], panic if this request doesn't have a read lock.
1449 pub async fn cache_lock_wait(&mut self) -> LockStatus {
1450 let inner_enabled = self.inner_enabled_mut();
1451 #[cfg_attr(not(feature = "trace"), allow(unused_mut))]
1452 let mut span = inner_enabled.traces.child("cache_lock");
1453 // should always call is_cache_locked() before this function, which should guarantee that
1454 // the inner cache has a read lock and lock ctx
1455 let (read_lock, status) = if let Some(lock_ctx) = inner_enabled.lock_ctx.as_mut() {
1456 let lock = lock_ctx.lock.take(); // remove the lock from self
1457 if let Some(Locked::Read(r)) = lock {
1458 let now = Instant::now();
1459 // it's possible for a request to be locked more than once,
1460 // so wait the remainder of our configured timeout
1461 let status = if let Some(wait_timeout) = lock_ctx.wait_timeout {
1462 let wait_timeout =
1463 wait_timeout.saturating_sub(self.lock_duration().unwrap_or(Duration::ZERO));
1464 match timeout(wait_timeout, r.wait()).await {
1465 Ok(()) => r.lock_status(),
1466 Err(_) => LockStatus::WaitTimeout,
1467 }
1468 } else {
1469 r.wait().await;
1470 r.lock_status()
1471 };
1472 self.digest.add_lock_duration(now.elapsed());
1473 (r, status)
1474 } else {
1475 panic!("cache_lock_wait on wrong type of lock")
1476 }
1477 } else {
1478 panic!("cache_lock_wait without cache lock")
1479 };
1480 if let Some(lock_ctx) = self.inner_enabled().lock_ctx.as_ref() {
1481 lock_ctx
1482 .cache_lock
1483 .trace_lock_wait(&mut span, &read_lock, status);
1484 }
1485 status
1486 }
1487
1488 /// How long did this request wait behind the read lock
1489 pub fn lock_duration(&self) -> Option<Duration> {
1490 self.digest.lock_duration
1491 }
1492
1493 /// How long did this request spent on cache lookup and reading the header
1494 pub fn lookup_duration(&self) -> Option<Duration> {
1495 self.digest.lookup_duration
1496 }
1497
1498 /// Delete the asset from the cache storage
1499 /// # Panic
1500 /// Need to be called after the cache key is set. Panic otherwise.
1501 pub async fn purge(&self) -> Result<bool> {
1502 match self.phase {
1503 CachePhase::CacheKey => {
1504 let inner = self.inner();
1505 let inner_enabled = self.inner_enabled();
1506 let span = inner_enabled.traces.child("purge");
1507 let key = inner.key.as_ref().unwrap().to_compact();
1508 Self::purge_impl(inner_enabled.storage, inner_enabled.eviction, &key, span).await
1509 }
1510 _ => panic!("wrong phase {:?}", self.phase),
1511 }
1512 }
1513
1514 /// Delete the asset from the cache storage via a spawned task.
1515 /// Returns corresponding `JoinHandle` of that task.
1516 /// # Panic
1517 /// Need to be called after the cache key is set. Panic otherwise.
1518 pub fn spawn_async_purge(
1519 &self,
1520 context: &'static str,
1521 ) -> tokio::task::JoinHandle<Result<bool>> {
1522 if matches!(self.phase, CachePhase::Disabled(_) | CachePhase::Uninit) {
1523 panic!("wrong phase {:?}", self.phase);
1524 }
1525
1526 let inner_enabled = self.inner_enabled();
1527 let span = inner_enabled.traces.child("purge");
1528 let key = self.inner().key.as_ref().unwrap().to_compact();
1529 let storage = inner_enabled.storage;
1530 let eviction = inner_enabled.eviction;
1531 tokio::task::spawn(async move {
1532 Self::purge_impl(storage, eviction, &key, span)
1533 .await
1534 .map_err(|e| {
1535 warn!("Failed to purge {key} (context: {context}): {e}");
1536 e
1537 })
1538 })
1539 }
1540
1541 #[cfg_attr(not(feature = "trace"), allow(unused_mut))]
1542 async fn purge_impl(
1543 storage: &'static (dyn storage::Storage + Sync),
1544 eviction: Option<&'static (dyn eviction::EvictionManager + Sync)>,
1545 key: &CompactCacheKey,
1546 mut span: Span,
1547 ) -> Result<bool> {
1548 let result = storage
1549 .purge(key, PurgeType::Invalidation, &span.handle())
1550 .await;
1551 let purged = matches!(result, Ok(true));
1552 // need to inform eviction manager if asset was removed
1553 if let Some(eviction) = eviction.as_ref() {
1554 if purged {
1555 eviction.remove(key);
1556 }
1557 }
1558 span.set_tag(|| trace::Tag::new("purged", purged));
1559 result
1560 }
1561
1562 /// Check the cacheable prediction
1563 ///
1564 /// Return true if the predictor is not set
1565 pub fn cacheable_prediction(&self) -> bool {
1566 if let Some(predictor) = self.inner().predictor {
1567 predictor.cacheable_prediction(self.cache_key())
1568 } else {
1569 true
1570 }
1571 }
1572
1573 /// Tell the predictor that this response, which is previously predicted to be uncacheable,
1574 /// is cacheable now.
1575 pub fn response_became_cacheable(&self) {
1576 if let Some(predictor) = self.inner().predictor {
1577 predictor.mark_cacheable(self.cache_key());
1578 }
1579 }
1580
1581 /// Tell the predictor that this response is uncacheable so that it will know next time
1582 /// this request arrives.
1583 pub fn response_became_uncacheable(&self, reason: NoCacheReason) {
1584 if let Some(predictor) = self.inner().predictor {
1585 predictor.mark_uncacheable(self.cache_key(), reason);
1586 }
1587 }
1588
1589 /// Tag all spans as being part of a subrequest.
1590 pub fn tag_as_subrequest(&mut self) {
1591 self.inner_enabled_mut()
1592 .traces
1593 .cache_span
1594 .set_tag(|| Tag::new("is_subrequest", true))
1595 }
1596}