use std::marker::PhantomData;
use std::sync::Arc;
use std::time::Duration;
use crate::capability::TierCapability;
use crate::continuity::RecoveryReport;
use crate::control::cachelito::{Cachelito, ControlSnapshot};
use crate::entry::{CommitToken, EntryState, Generation, IntentKind};
use crate::error::CacheError;
use crate::identity::CacheContext;
use crate::key::Key;
use crate::key::KeyRef;
use crate::policy::{CacheOperation, CachePolicy, CacheRequest, CacheState, FailMode};
use crate::pool::MemoryPool;
use crate::telemetry::{Operation, OperationRecord, SharedSink};
use crate::tier::tier_trait::CacheTier;
use crate::tier::{TierId, TierRegistry};
const MAX_KEY_SIZE: usize = 256;
const MAX_POPULATE_ATTEMPTS: u32 = 3;
const RETRY_BASE_BACKOFF: Duration = Duration::from_millis(50);
struct PopulationGuard<'a> {
cachelito: &'a Cachelito,
key: Vec<u8>,
armed: bool,
}
impl<'a> PopulationGuard<'a> {
fn new(cachelito: &'a Cachelito, key: &[u8]) -> Self {
Self {
cachelito,
key: key.to_vec(),
armed: true,
}
}
fn disarm(&mut self) {
self.armed = false;
}
}
impl Drop for PopulationGuard<'_> {
fn drop(&mut self) {
if !self.armed {
return;
}
let _ = self
.cachelito
.fail_with_error(&self.key, CacheError::Cancelled);
}
}
struct IntentGuard<'a> {
cachelito: &'a Cachelito,
token: CommitToken,
armed: bool,
}
impl<'a> IntentGuard<'a> {
fn new(cachelito: &'a Cachelito, token: CommitToken) -> Self {
Self {
cachelito,
token,
armed: true,
}
}
fn token(&self) -> &CommitToken {
&self.token
}
fn abort(&mut self, error: CacheError) {
if !self.armed {
return;
}
let _ = self.cachelito.abort(&self.token, error);
self.armed = false;
}
fn abort_proven_clean(&mut self, error: CacheError) {
if !self.armed {
return;
}
let _ = self.cachelito.abort_proven_clean(&self.token, error);
self.armed = false;
}
fn disarm(&mut self) {
self.armed = false;
}
}
impl Drop for IntentGuard<'_> {
fn drop(&mut self) {
if !self.armed {
return;
}
let _ = self.cachelito.abort(&self.token, CacheError::Cancelled);
}
}
fn error_kind(err: CacheError) -> CacheError {
err
}
fn write_provably_stored_nothing(error: CacheError) -> bool {
matches!(
error,
CacheError::SerializationFailed | CacheError::CapacityExhausted
)
}
fn previous_rung(rung: TierId) -> TierId {
match rung.as_usize() {
0 => TierId::L6,
n => TierId::from_usize(n - 1).unwrap_or(TierId::L0),
}
}
impl<K, V, P> CacheManager<K, V, P> {
const MAX_COMMIT_ATTEMPTS: usize = 64;
const MAX_EVICTION_PASSES: usize = 16;
const MAX_POPULATION_ATTEMPTS: usize = 8;
}
const fn decision_tier_placeholder() -> TierId {
TierId::L0
}
const fn attempt_unreachable() -> bool {
true
}
pub struct CacheManager<K, V, P> {
policy: P,
cachelito: Cachelito,
tier_registry: TierRegistry,
tiers: Vec<Arc<dyn CacheTier<V>>>,
_pool: MemoryPool<V>,
wait_timeout: Duration,
authority_tier: TierId,
telemetry: SharedSink,
_key: PhantomData<K>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Freshness {
Fresh,
Stale,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Lookup<V> {
pub value: Option<V>,
pub freshness: Freshness,
}
impl<K, V, P> CacheManager<K, V, P>
where
K: Key + Send + Sync,
V: Clone + Send + Sync + 'static,
P: CachePolicy<K, V>,
{
pub fn new(
policy: P,
cachelito: Cachelito,
tier_registry: TierRegistry,
tiers: Vec<Arc<dyn CacheTier<V>>>,
pool: MemoryPool<V>,
) -> Self {
Self::with_timeout(
policy,
cachelito,
tier_registry,
tiers,
pool,
Duration::from_secs(5),
)
}
pub fn with_timeout(
policy: P,
cachelito: Cachelito,
tier_registry: TierRegistry,
tiers: Vec<Arc<dyn CacheTier<V>>>,
pool: MemoryPool<V>,
wait_timeout: Duration,
) -> Self {
CacheManager {
policy,
cachelito,
tier_registry,
tiers,
_pool: pool,
wait_timeout,
authority_tier: crate::policy::AUTHORITY_TIER,
telemetry: std::sync::Arc::new(crate::telemetry::NoTelemetry),
_key: PhantomData,
}
}
#[must_use]
pub fn with_authority(mut self, authority: TierId) -> Self {
assert!(
authority.is_authority_rung(),
"authority must be the authority rung; {authority} is a cache rung"
);
self.authority_tier = authority;
self
}
#[must_use]
pub fn with_telemetry(mut self, sink: SharedSink) -> Self {
self.telemetry = sink;
self
}
#[must_use]
pub fn with_shared_telemetry(self, sink: SharedSink) -> Self {
self.with_telemetry(sink)
}
#[must_use]
pub const fn authority_tier(&self) -> TierId {
self.authority_tier
}
fn emit(&self, record: OperationRecord) {
self.telemetry.record(&record);
}
async fn observed_read(
&self,
op: Operation,
key_ref: &KeyRef<'_>,
tier_id: TierId,
fut: impl std::future::Future<Output = Result<Option<V>, CacheError>>,
) -> Result<Option<V>, CacheError> {
let started = std::time::Instant::now();
let result = self.bounded(fut).await;
let outcome = match &result {
Ok(Some(_)) => crate::telemetry::Outcome::Served,
Ok(None) => crate::telemetry::Outcome::Miss,
Err(e) => crate::telemetry::Outcome::from_error(*e),
};
let mut record = OperationRecord::begin(op, key_ref)
.with_source(tier_id)
.with_backend(self.backend_of(tier_id))
.with_generation(Generation::new(0))
.finish(started, outcome);
if let Err(e) = result {
record = record.with_failure(e);
}
self.emit(record);
result
}
async fn observed_write(
&self,
op: Operation,
key_ref: &KeyRef<'_>,
tier_id: TierId,
fut: impl std::future::Future<Output = Result<(), CacheError>>,
) -> Result<(), CacheError> {
let started = std::time::Instant::now();
let result = self.bounded(fut).await;
let outcome = match &result {
Ok(()) => crate::telemetry::Outcome::Populated,
Err(e) => crate::telemetry::Outcome::from_error(*e),
};
let mut record = OperationRecord::begin(op, key_ref)
.with_destination(tier_id)
.with_backend(self.backend_of(tier_id))
.finish(started, outcome);
if let Err(e) = result {
record = record.with_failure(e);
}
self.emit(record);
result
}
fn observed_control_only(
&self,
op: Operation,
key_ref: &KeyRef<'_>,
tier_id: TierId,
started: std::time::Instant,
outcome: crate::telemetry::Outcome,
) {
self.emit(
OperationRecord::begin(op, key_ref)
.with_source(tier_id)
.finish(started, outcome),
);
}
fn record_fallback(&self, key_ref: &KeyRef<'_>, tier: TierId, cause: CacheError) {
let started = std::time::Instant::now();
self.emit(
OperationRecord::begin(Operation::GetOrFetch, key_ref)
.with_source(tier)
.with_backend(self.backend_of(tier))
.with_fallback(true)
.with_failure(cause)
.finish(started, crate::telemetry::Outcome::FellBack),
);
}
fn record_telemetry_refresh_failure(
&self,
key_ref: &KeyRef<'_>,
snapshot: &ControlSnapshot,
cause: CacheError,
) {
let started = std::time::Instant::now();
self.emit(
OperationRecord::begin(Operation::Refresh, key_ref)
.with_source(snapshot.tier)
.with_generation(snapshot.generation)
.with_failure(cause)
.finish(started, crate::telemetry::Outcome::FellBack),
);
}
#[must_use]
pub fn continuity(&self) -> crate::continuity::ContinuityReport {
let mut states = std::collections::BTreeMap::new();
for &tier in self.tier_registry.all() {
states.insert(tier, self.continuity_of(tier));
}
crate::continuity::ContinuityReport { states }
}
fn continuity_of(&self, tier: TierId) -> crate::continuity::ContinuityState {
use crate::continuity::ContinuityState;
if !self.has_tier(&tier) {
return ContinuityState::Unavailable;
}
let health = self.tier_registry.tier_health(tier);
match health.consecutive_failures {
0 => ContinuityState::Healthy,
n if n >= crate::continuity::FAILURE_THRESHOLD => ContinuityState::Unavailable,
_ => ContinuityState::Degraded,
}
}
pub fn recover(&self) -> RecoveryReport {
self.recover_older_than(crate::continuity::INTENT_RECOVERY_AGE)
}
pub fn recover_older_than(&self, age: Duration) -> RecoveryReport {
let mut report = RecoveryReport::default();
for (address, intent) in self.cachelito.stale_intents(age.as_nanos() as u64) {
match intent.kind {
IntentKind::Write => {
match self
.cachelito
.abort_intent_by_address(address, CacheError::UncommittedIntent)
{
Ok(true) => report.recovered += 1,
Ok(false) => {}
Err(_) => report.failed += 1,
}
}
IntentKind::Move => {
report.needs_reconciliation += 1;
match self.cachelito.release_population_claim(address) {
Ok(true) => report.released_for_reconciliation += 1,
Ok(false) => {}
Err(_) => report.failed += 1,
}
}
}
}
report
}
pub fn cachelito(&self) -> &Cachelito {
&self.cachelito
}
async fn bounded<T, F>(&self, fut: F) -> Result<T, CacheError>
where
F: std::future::Future<Output = Result<T, CacheError>>,
{
match tokio::time::timeout(self.wait_timeout, fut).await {
Ok(result) => result,
Err(_) => Err(CacheError::Timeout),
}
}
fn encode_key<'a>(
key: &K,
ctx: &CacheContext,
buf: &'a mut [u8; MAX_KEY_SIZE],
) -> Result<KeyRef<'a>, CacheError> {
let mut raw = [0u8; MAX_KEY_SIZE];
let key_len = key.encode(&mut raw)?;
let len = crate::key::frame_tenant_key(&ctx.identity().tenant, &raw[..key_len], buf)?;
Ok(KeyRef(&buf[..len]))
}
fn check_auth(ctx: &CacheContext) -> Result<(), CacheError> {
if !ctx.is_authenticated() {
return Err(CacheError::Unauthenticated);
}
Ok(())
}
fn resolve(
&self,
request: &CacheRequest<K, V>,
ctx: &CacheContext,
) -> crate::policy::PolicyDecision {
let state = CacheState::new();
self.policy.select(request, &state, ctx.identity())
}
fn resolve_with_snapshot(
&self,
request: &CacheRequest<K, V>,
ctx: &CacheContext,
snapshot: &ControlSnapshot,
) -> crate::policy::PolicyDecision {
let health = self.tier_registry.tier_health(snapshot.tier);
let state = CacheState::from_snapshot(snapshot, health);
self.policy.select(request, &state, ctx.identity())
}
fn request_for(op: CacheOperation, key: &K, ctx: &CacheContext) -> CacheRequest<K, V> {
let mut request = CacheRequest::new(op, key.clone());
if let Some(ttl) = ctx.ttl() {
request = request.with_ttl(ttl);
}
request
}
fn authorize(decision: &crate::policy::PolicyDecision) -> Result<(), CacheError> {
if decision.authorized {
Ok(())
} else {
Err(CacheError::Unauthorized)
}
}
fn record_outcome(&self, tier: TierId, result: Result<(), CacheError>) {
match result {
Ok(()) => self.tier_registry.recover(tier),
Err(CacheError::TierUnavailable) => self.tier_registry.fail(tier),
Err(_) => {}
}
}
fn record_read<T>(&self, tier: TierId, result: &Result<T, CacheError>) {
match result {
Ok(_) => self.tier_registry.recover(tier),
Err(CacheError::TierUnavailable) => self.tier_registry.fail(tier),
Err(_) => {}
}
}
pub async fn get(&self, key: &K, ctx: &CacheContext) -> Result<Option<V>, CacheError> {
let started = std::time::Instant::now();
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
if let Err(e) = Self::check_auth(ctx) {
self.observed_control_only(
Operation::Get,
&key_ref,
decision_tier_placeholder(),
started,
crate::telemetry::Outcome::from_error(e),
);
return Err(e);
}
let request = Self::request_for(CacheOperation::Get, key, ctx);
let decision = self.resolve(&request, ctx);
if let Err(e) = Self::authorize(&decision) {
self.observed_control_only(
Operation::Get,
&key_ref,
decision.tier,
started,
crate::telemetry::Outcome::from_error(e),
);
return Err(e);
}
let snapshot = self.cachelito.peek(key_ref.0)?;
match snapshot.state {
EntryState::Ready if !snapshot.expired => {
let Some(tier) = self.bound_tier(&snapshot.tier) else {
let e = CacheError::TierUnavailable;
self.observed_control_only(
Operation::Get,
&key_ref,
snapshot.tier,
started,
crate::telemetry::Outcome::from_error(e),
);
return Err(e);
};
let result = self
.observed_read(Operation::Get, &key_ref, snapshot.tier, tier.get(&key_ref))
.await;
self.record_read(snapshot.tier, &result);
if result.is_ok()
&& let Ok(after) = self.cachelito.peek(key_ref.0)
&& !matches!(after.state, EntryState::Ready)
{
return Err(CacheError::Miss);
}
result
}
EntryState::Ready => {
let _ = self.cachelito.set_state(key_ref.0, EntryState::Stale);
self.observed_control_only(
Operation::Get,
&key_ref,
snapshot.tier,
started,
crate::telemetry::Outcome::Miss,
);
Err(CacheError::Miss)
}
EntryState::InFlight | EntryState::Prepared if !snapshot.population_owner => {
self.wait_for_population(&key_ref, &snapshot).await?;
let after = self.cachelito.peek(key_ref.0)?;
if let Some(err) = after.last_error {
return Err(err);
}
if after.state != EntryState::Ready {
return Err(CacheError::Miss);
}
let Some(tier) = self.bound_tier(&after.tier) else {
return Err(CacheError::TierUnavailable);
};
let result = self
.observed_read(Operation::Get, &key_ref, after.tier, tier.get(&key_ref))
.await;
self.record_read(after.tier, &result);
result
}
_ => {
self.observed_control_only(
Operation::Get,
&key_ref,
snapshot.tier,
started,
crate::telemetry::Outcome::Miss,
);
Err(CacheError::Miss)
}
}
}
pub async fn get_or_fetch<F, Fut>(
&self,
key: &K,
ctx: &CacheContext,
fetch: F,
) -> Result<V, CacheError>
where
F: Fn() -> Fut,
Fut: std::future::Future<Output = Result<V, CacheError>>,
{
Self::check_auth(ctx)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
let request = Self::request_for(CacheOperation::Get, key, ctx);
let decision = self.resolve(&request, ctx);
Self::authorize(&decision)?;
for _attempt in 0..Self::MAX_POPULATION_ATTEMPTS {
let snapshot = self.cachelito.acquire(key_ref.0, decision.tier)?;
let snapshot = if snapshot.state == EntryState::Ready && snapshot.expired {
let _ = self.cachelito.set_state(key_ref.0, EntryState::Stale);
self.cachelito.acquire(key_ref.0, decision.tier)?
} else {
snapshot
};
if snapshot.state == EntryState::Ready && !snapshot.expired {
let Some(tier) = self.bound_tier(&snapshot.tier) else {
return Err(CacheError::TierUnavailable);
};
match self
.observed_read(
Operation::GetOrFetch,
&key_ref,
snapshot.tier,
tier.get(&key_ref),
)
.await
{
Ok(Some(v)) => return Ok(v),
Ok(None) => {}
Err(e) => return Err(e),
}
}
match snapshot.state {
EntryState::InFlight if !snapshot.population_owner => {
return self.wait_and_get(key, ctx, &snapshot).await;
}
EntryState::Prepared if snapshot.population_owner => {
self.wait_for_population(&key_ref, &snapshot).await?;
continue;
}
EntryState::Prepared => {
return self.wait_and_get(key, ctx, &snapshot).await;
}
_ => {
return self
.become_population_owner(
key,
key_ref,
snapshot.generation,
ctx,
&decision,
fetch,
)
.await;
}
}
}
let _ = attempt_unreachable();
Err(CacheError::Timeout)
}
pub async fn set(&self, key: &K, value: V, ctx: &CacheContext) -> Result<(), CacheError> {
Self::check_auth(ctx)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
let snapshot = self.cachelito.peek(key_ref.0)?;
let request = Self::request_for(CacheOperation::Set, key, ctx);
let decision = self.resolve_with_snapshot(&request, ctx, &snapshot);
Self::authorize(&decision)?;
if decision.tier == self.authority_tier {
return Err(CacheError::PolicyDenied);
}
let mut last_error = CacheError::TierUnavailable;
'pass: for _pass in 0..=Self::MAX_EVICTION_PASSES {
let mut rung = decision.tier;
'outer: while rung.is_cache_rung() {
let Some(tier) = self.bound_tier(&rung) else {
rung = previous_rung(rung);
continue;
};
let mut lost_races = 0usize;
loop {
let token = self
.cachelito
.prepare(key_ref.0, None, rung, IntentKind::Write)?;
let mut intent = IntentGuard::new(&self.cachelito, token);
let set_result = self
.observed_write(
Operation::Set,
&key_ref,
rung,
tier.set(&key_ref, value.clone(), ctx.ttl()),
)
.await;
self.record_outcome(rung, set_result);
if let Err(CacheError::WriteIndeterminate) = set_result {
let _ = self.bounded(tier.remove(&key_ref)).await;
}
match set_result {
Ok(()) => {
return match self.cachelito.commit(intent.token(), ctx.ttl()) {
Ok(()) => {
intent.disarm();
Ok(())
}
Err(CacheError::StaleGeneration) => {
lost_races += 1;
if lost_races >= Self::MAX_COMMIT_ATTEMPTS {
let e = CacheError::WriteContended;
drop(intent);
return Err(e);
}
tokio::task::yield_now().await;
continue;
}
Err(e) => {
let _ = tier.remove(&key_ref).await;
intent.abort(e);
Err(e)
}
};
}
Err(
e @ (CacheError::CapacityExhausted
| CacheError::TierUnavailable
| CacheError::Timeout
| CacheError::Corrupted
| CacheError::SerializationFailed
| CacheError::WriteIndeterminate),
) => {
if write_provably_stored_nothing(e) {
intent.abort_proven_clean(e);
} else {
intent.abort(e);
}
last_error = e;
rung = previous_rung(rung);
continue 'outer;
}
Err(e) => {
if write_provably_stored_nothing(e) {
intent.abort_proven_clean(e);
} else {
intent.abort(e);
}
return Err(e);
}
}
}
}
if last_error == CacheError::CapacityExhausted
&& self.try_free_somewhere_in_the_ladder(key_ref.0)
{
continue 'pass;
}
break;
}
Err(last_error)
}
fn try_free_somewhere_in_the_ladder(&self, key: &[u8]) -> bool {
let mut budget = 0usize;
for candidate in crate::policy::cache_ladder() {
let Some(tier) = self.bound_tier(&candidate) else {
continue;
};
if self.try_evict_at_rung(&tier, key, &mut budget) {
return true;
}
}
false
}
fn try_evict_at_rung(
&self,
tier: &Arc<dyn CacheTier<V>>,
key: &[u8],
attempts: &mut usize,
) -> bool {
const MAX_EVICTION_ATTEMPTS: usize = 8;
while *attempts < MAX_EVICTION_ATTEMPTS {
*attempts += 1;
let Some(address) = tier.eviction_candidate() else {
return false;
};
match self.cachelito.reserve_eviction(address) {
Ok(true) => {}
Ok(false) => continue,
Err(_) => return false,
}
let r = matches!(tier.remove_if_address(address), Ok(true));
return r;
}
let _ = key;
false
}
pub async fn invalidate(&self, key: &K, ctx: &CacheContext) -> Result<(), CacheError> {
Self::check_auth(ctx)?;
let request = Self::request_for(CacheOperation::Invalidate, key, ctx);
let decision = self.resolve(&request, ctx);
Self::authorize(&decision)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
self.cachelito.invalidate(key_ref.0)?;
Ok(())
}
pub async fn remove(&self, key: &K, ctx: &CacheContext) -> Result<(), CacheError> {
Self::check_auth(ctx)?;
let request = Self::request_for(CacheOperation::Remove, key, ctx);
let decision = self.resolve(&request, ctx);
Self::authorize(&decision)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
for (idx, tier) in self.tiers.iter().enumerate() {
if idx == TierId::L6.as_usize() {
continue;
}
let _ = tier.remove(&key_ref).await;
}
self.cachelito.release(key_ref.0)?;
Ok(())
}
pub async fn exists(&self, key: &K, ctx: &CacheContext) -> Result<bool, CacheError> {
Self::check_auth(ctx)?;
let request = Self::request_for(CacheOperation::Exists, key, ctx);
let decision = self.resolve(&request, ctx);
Self::authorize(&decision)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
let snapshot = self.cachelito.peek(key_ref.0)?;
if snapshot.state != EntryState::Ready {
return Ok(false);
}
let Some(tier) = self.bound_tier(&snapshot.tier) else {
return Ok(false);
};
self.bounded(tier.contains(&key_ref)).await
}
pub async fn refresh_detailed<F, Fut>(
&self,
key: &K,
ctx: &CacheContext,
fetch: F,
) -> Result<Lookup<V>, CacheError>
where
F: Fn() -> Fut,
Fut: std::future::Future<Output = Result<V, CacheError>>,
{
Self::check_auth(ctx)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
let request = Self::request_for(CacheOperation::Refresh, key, ctx);
let mut decision = self.resolve(&request, ctx);
Self::authorize(&decision)?;
decision.fail_mode = FailMode::Closed;
let before = self.cachelito.peek(key_ref.0)?;
let current = if before.state.is_readable() {
match self.bound_tier(&before.tier) {
Some(tier) => self.bounded(tier.get(&key_ref)).await.unwrap_or_default(),
None => None,
}
} else {
None
};
if before.state == EntryState::Ready && !before.expired {
let _ = self.cachelito.set_state(key_ref.0, EntryState::Stale);
}
let snapshot = self.cachelito.acquire(key_ref.0, decision.tier)?;
if !snapshot.population_owner {
return Ok(Lookup {
value: current,
freshness: Freshness::Stale,
});
}
match self
.become_population_owner(key, key_ref, snapshot.generation, ctx, &decision, fetch)
.await
{
Ok(new_value) => Ok(Lookup {
value: Some(new_value),
freshness: Freshness::Fresh,
}),
Err(e) if current.is_some() => {
self.record_telemetry_refresh_failure(&key_ref, &snapshot, e);
Ok(Lookup {
value: current,
freshness: Freshness::Stale,
})
}
Err(e) => Err(e),
}
}
pub async fn refresh<F, Fut>(
&self,
key: &K,
ctx: &CacheContext,
fetch: F,
) -> Result<Option<V>, CacheError>
where
F: Fn() -> Fut,
Fut: std::future::Future<Output = Result<V, CacheError>>,
{
self.refresh_detailed(key, ctx, fetch)
.await
.map(|l| l.value)
}
pub async fn promote(&self, key: &K, ctx: &CacheContext) -> Result<(), CacheError> {
Self::check_auth(ctx)?;
let request = Self::request_for(CacheOperation::Promote, key, ctx);
let decision = self.resolve(&request, ctx);
Self::authorize(&decision)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
let snapshot = self.cachelito.peek(key_ref.0)?;
if snapshot.state != EntryState::Ready {
return Err(CacheError::Miss);
}
if !snapshot.tier.is_cache_rung() || snapshot.tier.as_usize() == 0 {
return Ok(());
}
let new_tier_id = TierId::from_usize(snapshot.tier.as_usize() - 1)
.ok_or(CacheError::ConfigurationError)?;
self.move_entry(key, &key_ref, &snapshot, new_tier_id, ctx)
.await
}
pub async fn demote(&self, key: &K, ctx: &CacheContext) -> Result<(), CacheError> {
Self::check_auth(ctx)?;
let request = Self::request_for(CacheOperation::Demote, key, ctx);
let decision = self.resolve(&request, ctx);
Self::authorize(&decision)?;
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
let snapshot = self.cachelito.peek(key_ref.0)?;
if snapshot.state != EntryState::Ready {
return Err(CacheError::Miss);
}
if !snapshot.tier.is_cache_rung() || snapshot.tier >= crate::policy::LAST_CACHE_TIER {
return Ok(());
}
let new_tier_id = TierId::from_usize(snapshot.tier.as_usize() + 1)
.ok_or(CacheError::ConfigurationError)?;
self.move_entry(key, &key_ref, &snapshot, new_tier_id, ctx)
.await
}
async fn move_entry(
&self,
key: &K,
key_ref: &KeyRef<'_>,
snapshot: &ControlSnapshot,
new_tier_id: TierId,
ctx: &CacheContext,
) -> Result<(), CacheError> {
let token = self
.cachelito
.prepare(key_ref.0, None, new_tier_id, IntentKind::Move)?;
let mut intent = IntentGuard::new(&self.cachelito, token);
let (Some(from_tier), Some(to_tier)) = (
self.bound_tier(&snapshot.tier),
self.bound_tier(&new_tier_id),
) else {
return Err(CacheError::TierUnavailable);
};
let value = match self.bounded(from_tier.get(key_ref)).await {
Ok(Some(v)) => v,
Ok(None) => {
intent.abort(CacheError::Miss);
return Err(CacheError::Miss);
}
Err(e) => {
intent.abort(e);
return Err(e);
}
};
let mut buf2 = [0u8; MAX_KEY_SIZE];
let key_ref2 = match Self::encode_key(key, ctx, &mut buf2) {
Ok(k) => k,
Err(e) => {
intent.abort(e);
return Err(e);
}
};
let set_result = self
.observed_write(
Operation::Promote,
&key_ref2,
new_tier_id,
to_tier.set(&key_ref2, value, ctx.ttl()),
)
.await;
self.record_outcome(new_tier_id, set_result);
if let Err(e) = set_result {
intent.abort(e);
return Err(e);
}
let remove_result = from_tier.remove(key_ref).await;
self.record_outcome(snapshot.tier, remove_result);
match self.cachelito.commit(intent.token(), ctx.ttl()) {
Ok(()) => {
intent.disarm();
Ok(())
}
Err(e) => {
let _ = to_tier.remove(&key_ref2).await;
intent.abort(e);
Err(e)
}
}
}
#[must_use]
pub fn nearest_bound_rung(&self, wanted: TierId) -> Option<TierId> {
if !wanted.is_cache_rung() {
return None;
}
let mut idx = wanted.as_usize();
loop {
let candidate = TierId::from_index(u8::try_from(idx).ok()?)?;
if self.has_tier(&candidate) {
return Some(candidate);
}
if idx == 0 {
return None;
}
idx -= 1;
}
}
#[must_use]
pub fn bound_tier(&self, tier_id: &TierId) -> Option<Arc<dyn CacheTier<V>>> {
self.tiers.get(tier_id.as_usize()).cloned()
}
#[must_use]
pub fn capabilities(&self) -> std::collections::BTreeMap<TierId, TierCapability> {
TierId::ALL
.iter()
.map(|id| {
let mut cap = if self.has_tier(id) {
self.bound_tier(id)
.map_or_else(TierCapability::unbound, |t| t.capability())
} else {
TierCapability::unbound()
};
if *id == self.authority_tier {
cap.flags |= crate::capability::CapabilityFlags::AUTHORITATIVE;
}
(*id, cap)
})
.collect()
}
#[must_use]
pub fn has_tier(&self, tier_id: &TierId) -> bool {
tier_id.as_usize() < self.tiers.len()
}
#[must_use]
pub fn tier(&self, tier_id: &TierId) -> Option<Arc<dyn CacheTier<V>>> {
self.bound_tier(tier_id)
}
fn backend_of(&self, tier_id: TierId) -> crate::tier::tier_trait::BackendKind {
self.bound_tier(&tier_id)
.map_or(crate::tier::tier_trait::BackendKind::Unavailable, |t| {
t.backend()
})
}
async fn wait_for_population(
&self,
key_ref: &KeyRef<'_>,
snapshot: &ControlSnapshot,
) -> Result<(), CacheError> {
let notified = snapshot.notify.notified();
tokio::pin!(notified);
notified.as_mut().enable();
match self.cachelito.peek(key_ref.0) {
Ok(now) if now.state != EntryState::InFlight && now.state != EntryState::Prepared => {
return Ok(());
}
Ok(_) => {}
Err(_) => {}
}
tokio::select! {
_ = tokio::time::sleep(self.wait_timeout) => Err(CacheError::Timeout),
_ = &mut notified => Ok(()),
}
}
async fn wait_and_get(
&self,
key: &K,
ctx: &CacheContext,
snapshot: &ControlSnapshot,
) -> Result<V, CacheError> {
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
self.wait_for_population(&key_ref, snapshot).await?;
let after = self.cachelito.peek(key_ref.0)?;
if let Some(err) = after.last_error {
return Err(err);
}
if !after.state.is_readable() {
return Err(CacheError::Miss);
}
let Some(tier) = self.bound_tier(&after.tier) else {
return Err(CacheError::TierUnavailable);
};
match self
.observed_read(
Operation::GetOrFetch,
&key_ref,
after.tier,
tier.get(&key_ref),
)
.await?
{
Some(v) => Ok(v),
None => Err(CacheError::Miss),
}
}
async fn become_population_owner<F, Fut>(
&self,
key: &K,
key_ref: KeyRef<'_>,
expected_generation: Generation,
ctx: &CacheContext,
decision: &crate::policy::PolicyDecision,
fetch: F,
) -> Result<V, CacheError>
where
F: Fn() -> Fut,
Fut: std::future::Future<Output = Result<V, CacheError>>,
{
let mut backoff = RETRY_BASE_BACKOFF;
let mut attempt = 0_u32;
let generation = expected_generation;
let mut guard = PopulationGuard::new(&self.cachelito, key_ref.0);
loop {
attempt += 1;
let outcome = self
.populate_once(key, &key_ref, generation, ctx, &fetch, decision)
.await;
match outcome {
Ok(value) => {
guard.disarm();
return Ok(value);
}
Err(e) if attempt < MAX_POPULATE_ATTEMPTS && Self::is_retryable(e) => {
tokio::time::sleep(backoff).await;
backoff = backoff.saturating_mul(2);
}
Err(e) => {
let started = std::time::Instant::now();
self.emit(
OperationRecord::begin(Operation::GetOrFetch, &key_ref)
.with_destination(decision.tier)
.with_backend(self.backend_of(decision.tier))
.with_failure(e)
.finish(started, crate::telemetry::Outcome::from_error(e)),
);
let _ = self.cachelito.fail_with_error(key_ref.0, error_kind(e));
guard.disarm();
if matches!(e, CacheError::StaleGeneration) {
return Err(CacheError::StaleGeneration);
}
if decision.fail_mode == FailMode::Open {
return self.try_fallback_tier(key, ctx, e).await;
}
return Err(e);
}
}
}
}
async fn populate_once<F, Fut>(
&self,
key: &K,
key_ref: &KeyRef<'_>,
expected_generation: Generation,
ctx: &CacheContext,
fetch: &F,
decision: &crate::policy::PolicyDecision,
) -> Result<V, CacheError>
where
F: Fn() -> Fut,
Fut: std::future::Future<Output = Result<V, CacheError>>,
{
let result = tokio::time::timeout(self.wait_timeout, fetch()).await;
let value = match result {
Ok(Ok(value)) => value,
Ok(Err(e)) => return Err(e),
Err(_) => return Err(CacheError::Timeout),
};
let Some(tier_id) = self.nearest_bound_rung(decision.tier) else {
return Err(CacheError::TierUnavailable);
};
let Some(tier) = self.bound_tier(&tier_id) else {
return Err(CacheError::TierUnavailable);
};
let mut buf2 = [0u8; MAX_KEY_SIZE];
let key_ref2 = Self::encode_key(key, ctx, &mut buf2)?;
let set_result = self
.observed_write(
Operation::GetOrFetch,
&key_ref2,
tier_id,
tier.set(&key_ref2, value.clone(), ctx.ttl()),
)
.await;
self.record_outcome(tier_id, set_result);
set_result?;
self.cachelito
.publish(key_ref.0, expected_generation, tier_id, ctx.ttl())?;
Ok(value)
}
fn is_retryable(err: CacheError) -> bool {
matches!(
err,
CacheError::TierUnavailable | CacheError::Timeout | CacheError::PopulationFailed
)
}
async fn try_fallback_tier(
&self,
key: &K,
ctx: &CacheContext,
cause: CacheError,
) -> Result<V, CacheError> {
let mut buf = [0u8; MAX_KEY_SIZE];
let key_ref = Self::encode_key(key, ctx, &mut buf)?;
for tier_id in crate::policy::cache_ladder() {
if !self.has_tier(&tier_id) {
continue;
}
if self.tier_registry.is_circuit_open(tier_id) {
continue;
}
let Some(tier) = self.bound_tier(&tier_id) else {
continue;
};
match self.bounded(tier.get(&key_ref)).await {
Ok(Some(v)) => {
self.record_fallback(&key_ref, tier_id, cause);
return Ok(v);
}
Ok(None) => continue,
Err(_) => continue,
}
}
Err(cause)
}
}