Skip to main content

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}