use std::cell::RefCell;
use std::collections::HashMap;
use std::fs;
use std::path::{Path, PathBuf};
use std::rc::Rc;
use crate::error::{Error, Result};
use crate::store::NodeId;
use super::cache::DerivedCache;
use super::dag::{CacheEffect, CacheNote, OutputCache};
use super::node::NodeKind;
use super::write_atomic;
pub const DEFAULT_PROMOTE_BUDGET_BYTES: u64 = 256 * 1024 * 1024;
pub const DEFAULT_PROMOTE_MAX_NODE_BYTES: u64 = 8 * 1024 * 1024;
pub const DEFAULT_PROMOTE_HORIZON: u64 = 64;
pub const DEFAULT_PROMOTE_RENT_PER_BYTE: u64 = 1;
pub const DEFAULT_PROMOTE_MIN_HITS: u32 = 2;
#[derive(Debug, Clone, Copy)]
pub struct PromotePolicy {
pub enabled: bool,
pub budget_bytes: u64,
pub max_node_bytes: u64,
pub horizon: u64,
pub rent_per_byte: u64,
pub min_hits: u32,
}
impl Default for PromotePolicy {
fn default() -> Self {
PromotePolicy {
enabled: false,
budget_bytes: DEFAULT_PROMOTE_BUDGET_BYTES,
max_node_bytes: DEFAULT_PROMOTE_MAX_NODE_BYTES,
horizon: DEFAULT_PROMOTE_HORIZON,
rent_per_byte: DEFAULT_PROMOTE_RENT_PER_BYTE,
min_hits: DEFAULT_PROMOTE_MIN_HITS,
}
}
}
pub struct PromotedStore {
root: PathBuf,
}
impl std::fmt::Debug for PromotedStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PromotedStore")
.field("root", &self.root)
.finish_non_exhaustive()
}
}
impl PromotedStore {
pub fn open(root: impl AsRef<Path>) -> Result<Self> {
let root = root.as_ref().to_path_buf();
fs::create_dir_all(&root)?;
Ok(PromotedStore { root })
}
pub fn root(&self) -> &Path {
&self.root
}
fn bytes_path(&self, id: &NodeId) -> PathBuf {
self.root.join(id.to_hex())
}
fn digest_path(&self, id: &NodeId) -> PathBuf {
self.root.join(format!("{}.b3", id.to_hex()))
}
pub fn get(&self, id: &NodeId) -> Result<Option<Vec<u8>>> {
let bytes = match fs::read(self.bytes_path(id)) {
Ok(b) => b,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(Error::io(format!("reading promoted entry {id}: {e}"))),
};
let sidecar = match fs::read(self.digest_path(id)) {
Ok(b) => b,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(Error::io(format!("reading promoted sidecar {id}: {e}"))),
};
if sidecar.as_slice() != blake3::hash(&bytes).as_bytes() {
return Err(Error::integrity_mismatch(format!(
"promoted entry {id} does not match its sidecar digest"
)));
}
Ok(Some(bytes))
}
pub fn put(&mut self, id: &NodeId, bytes: &[u8]) -> Result<u64> {
write_atomic(&self.bytes_path(id), bytes)?;
write_atomic(&self.digest_path(id), blake3::hash(bytes).as_bytes())?;
Ok(bytes.len() as u64)
}
pub fn remove(&mut self, id: &NodeId) -> Result<u64> {
let mut reclaimed = 0u64;
for path in [self.bytes_path(id), self.digest_path(id)] {
match fs::metadata(&path) {
Ok(meta) => reclaimed = reclaimed.saturating_add(meta.len()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(Error::io(format!("stat promoted entry {id}: {e}"))),
}
match fs::remove_file(&path) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(Error::io(format!("removing promoted entry {id}: {e}"))),
}
}
Ok(reclaimed)
}
pub fn contains(&self, id: &NodeId) -> Result<bool> {
Ok(self.bytes_path(id).exists())
}
pub fn total_bytes(&self) -> Result<u64> {
let mut total = 0u64;
for entry in fs::read_dir(&self.root)? {
let entry = entry?;
if entry.file_name().to_string_lossy().starts_with('.') {
continue;
}
let meta = entry.metadata()?;
if meta.is_file() {
total = total.saturating_add(meta.len());
}
}
Ok(total)
}
pub fn output_bytes(&self) -> Result<u64> {
let mut total = 0u64;
for entry in fs::read_dir(&self.root)? {
let entry = entry?;
let name = entry.file_name();
let name = name.to_string_lossy();
if name.starts_with('.') || name.ends_with(".b3") {
continue;
}
let meta = entry.metadata()?;
if meta.is_file() {
total = total.saturating_add(meta.len());
}
}
Ok(total)
}
pub fn clear(&mut self) -> Result<u64> {
let mut reclaimed = 0u64;
for entry in fs::read_dir(&self.root)? {
let entry = entry?;
let path = entry.path();
if entry.file_type()?.is_file() {
if !entry.file_name().to_string_lossy().starts_with('.') {
reclaimed = reclaimed.saturating_add(entry.metadata()?.len());
}
fs::remove_file(&path)?;
}
}
Ok(reclaimed)
}
}
fn kind_weight(kind: NodeKind) -> u64 {
match kind {
NodeKind::PdfStreamDecoded | NodeKind::PackageMemberDecoded => 4, NodeKind::PackageOpcModel
| NodeKind::DocxModel
| NodeKind::DocxStory
| NodeKind::EpubModel
| NodeKind::EpubContent
| NodeKind::OdtModel
| NodeKind::OdtContent => 8, NodeKind::ContentOperators | NodeKind::TextRuns | NodeKind::PagePreview => 2,
_ => 1,
}
}
fn cost_of(kind: NodeKind, bytes: u64, measured_micros: u64) -> u64 {
measured_micros.max(bytes.saturating_mul(kind_weight(kind)))
}
#[derive(Debug, Clone, Copy)]
struct PromoEntry {
bytes: u64,
cost_units: u64,
hits: u32,
born: u64,
last_seen: u64,
promoted: bool,
}
#[derive(Default)]
struct GovernorInner {
policy: PromotePolicy,
now: u64,
entries: HashMap<NodeId, PromoEntry>,
promoted_bytes: u64,
initialized: bool,
}
#[derive(Clone, Default)]
pub struct Governor {
inner: Rc<RefCell<GovernorInner>>,
}
impl std::fmt::Debug for Governor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let inner = self.inner.borrow();
f.debug_struct("Governor")
.field("policy", &inner.policy)
.field("now", &inner.now)
.field("entries", &inner.entries.len())
.field("promoted_bytes", &inner.promoted_bytes)
.finish()
}
}
impl Governor {
pub fn with_policy(policy: PromotePolicy) -> Self {
let g = Governor::default();
g.inner.borrow_mut().policy = policy;
g
}
pub fn policy(&self) -> PromotePolicy {
self.inner.borrow().policy
}
pub fn set_policy(&self, policy: PromotePolicy) {
self.inner.borrow_mut().policy = policy;
}
pub fn enabled(&self) -> bool {
self.inner.borrow().policy.enabled
}
fn begin(&self) -> bool {
let mut inner = self.inner.borrow_mut();
if inner.initialized {
false
} else {
inner.initialized = true;
true
}
}
fn set_promoted_bytes(&self, n: u64) {
self.inner.borrow_mut().promoted_bytes = n;
}
fn promoted_bytes(&self) -> u64 {
self.inner.borrow().promoted_bytes
}
fn entry(&self, id: &NodeId) -> Option<PromoEntry> {
self.inner.borrow().entries.get(id).copied()
}
fn mark_promoted(&self, id: &NodeId) {
let mut inner = self.inner.borrow_mut();
let Some(e) = inner.entries.get(id) else {
return;
};
if e.promoted {
return;
}
let bytes = e.bytes;
if let Some(e) = inner.entries.get_mut(id) {
e.promoted = true;
}
inner.promoted_bytes = inner.promoted_bytes.saturating_add(bytes);
}
fn mark_evicted(&self, id: &NodeId) {
let mut inner = self.inner.borrow_mut();
let Some(e) = inner.entries.get(id) else {
return;
};
if !e.promoted {
return;
}
let bytes = e.bytes;
if let Some(e) = inner.entries.get_mut(id) {
e.promoted = false;
}
inner.promoted_bytes = inner.promoted_bytes.saturating_sub(bytes);
}
fn record_stored(&self, kind: NodeKind, id: &NodeId, bytes: u64, wall_micros: u64) {
let mut inner = self.inner.borrow_mut();
inner.now = inner.now.saturating_add(1);
let now = inner.now;
let cost = cost_of(kind, bytes, wall_micros);
let entry = inner.entries.entry(*id).or_insert(PromoEntry {
bytes,
cost_units: cost,
hits: 0,
born: now,
last_seen: now,
promoted: false,
});
entry.bytes = bytes;
entry.cost_units = cost;
entry.last_seen = now;
}
fn credit_hit(&self, id: &NodeId) -> Option<PromoEntry> {
let mut inner = self.inner.borrow_mut();
inner.now = inner.now.saturating_add(1);
let now = inner.now;
let entry = inner.entries.get_mut(id)?;
entry.hits = entry.hits.saturating_add(1);
entry.last_seen = now;
Some(*entry)
}
fn should_promote(&self, e: &PromoEntry) -> bool {
let inner = self.inner.borrow();
should_promote(&inner.policy, e, inner.now)
}
fn eviction_plan(&self, need_bytes: u64) -> Vec<NodeId> {
let inner = self.inner.borrow();
let budget = inner.policy.budget_bytes;
let mut used = inner.promoted_bytes;
if used.saturating_add(need_bytes) <= budget {
return Vec::new();
}
let mut victims: Vec<(&NodeId, u64)> = inner
.entries
.iter()
.filter(|(_, e)| e.promoted)
.map(|(id, e)| (id, e.last_seen))
.collect();
victims.sort_by_key(|(_, seen)| *seen);
let mut plan = Vec::new();
for (id, _) in victims {
if used.saturating_add(need_bytes) <= budget {
break;
}
plan.push(*id);
used = used.saturating_sub(inner.entries.get(id).map_or(0, |e| e.bytes));
}
plan
}
}
fn should_promote(p: &PromotePolicy, e: &PromoEntry, now: u64) -> bool {
if !p.enabled {
return false;
}
if e.bytes > p.max_node_bytes {
return false; }
if e.hits < p.min_hits {
return false; }
let age = now.saturating_sub(e.born).max(1);
let future_hits = (e.hits as u64).saturating_mul(p.horizon) / age;
let saved = future_hits.saturating_mul(e.cost_units);
let rent = e
.bytes
.saturating_mul(p.rent_per_byte)
.saturating_mul(p.horizon);
saved > e.cost_units.saturating_add(rent)
}
pub struct GovernedCache {
ephemeral: DerivedCache,
promoted: PromotedStore,
gov: Governor,
}
impl std::fmt::Debug for GovernedCache {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("GovernedCache")
.field("promoted", &self.promoted)
.field("gov", &self.gov)
.finish_non_exhaustive()
}
}
impl GovernedCache {
pub fn open(store_root: impl AsRef<Path>, gov: Governor) -> Result<Self> {
let root = store_root.as_ref();
let ephemeral = DerivedCache::open(root.join("cache"))?;
let mut promoted = PromotedStore::open(root.join("promoted"))?;
if gov.begin() {
let on_disk = promoted.output_bytes()?;
if on_disk > gov.policy().budget_bytes {
promoted.clear()?;
gov.set_promoted_bytes(0);
} else {
gov.set_promoted_bytes(on_disk);
}
}
Ok(GovernedCache {
ephemeral,
promoted,
gov,
})
}
pub fn promoted(&self) -> &PromotedStore {
&self.promoted
}
fn maybe_promote(&mut self, id: &NodeId) {
let Some(entry) = self.gov.entry(id) else {
return;
};
if entry.promoted || entry.bytes > self.gov.policy().max_node_bytes {
return;
}
if !self.gov.should_promote(&entry) {
return;
}
let Ok(Some(bytes)) = self.ephemeral.get(id) else {
return;
};
let need = bytes.len() as u64;
if need > self.gov.policy().budget_bytes {
return;
}
for victim in self.gov.eviction_plan(need) {
if self.promoted.remove(&victim).is_ok() {
self.gov.mark_evicted(&victim);
}
}
if self.gov.promoted_bytes().saturating_add(need) > self.gov.policy().budget_bytes {
return;
}
if self.promoted.put(id, &bytes).is_ok() {
self.gov.mark_promoted(id);
}
}
}
impl OutputCache for GovernedCache {
fn get(&self, id: &NodeId) -> Result<Option<Vec<u8>>> {
match self.promoted.get(id) {
Ok(Some(bytes)) => return Ok(Some(bytes)),
Ok(None) => {}
Err(e) => return Err(e),
}
self.ephemeral.get(id)
}
fn put(&mut self, id: &NodeId, bytes: &[u8]) -> Result<()> {
self.ephemeral.put(id, bytes)?;
Ok(())
}
fn note(&mut self, note: CacheNote<'_>) {
if !self.gov.enabled() {
return;
}
match note.effect {
CacheEffect::Stored => {
self.gov
.record_stored(note.kind, note.id, note.bytes, note.wall_micros)
}
CacheEffect::Hit => {
self.gov.credit_hit(note.id);
self.maybe_promote(note.id);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::field::dag::{CacheEffect, CacheNote, OutputCache};
fn temp_root(label: &str) -> PathBuf {
let mut p = std::env::temp_dir();
p.push(format!(
"vole-promote-{label}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
p
}
fn stored(kind: NodeKind, id: &NodeId, bytes: u64) -> CacheNote<'_> {
CacheNote {
effect: CacheEffect::Stored,
kind,
id,
bytes,
wall_micros: 0,
}
}
fn hit(kind: NodeKind, id: &NodeId, bytes: u64) -> CacheNote<'_> {
CacheNote {
effect: CacheEffect::Hit,
kind,
id,
bytes,
wall_micros: 0,
}
}
fn promoted_policy(min_hits: u32, budget_bytes: u64) -> PromotePolicy {
PromotePolicy {
enabled: true,
budget_bytes,
min_hits,
rent_per_byte: 0,
..PromotePolicy::default()
}
}
#[test]
fn promotion_is_disabled_by_default_writes_nothing() {
let root = temp_root("off");
let gov = Governor::default();
assert!(!gov.enabled());
let mut cache = GovernedCache::open(&root, gov).unwrap();
let id = NodeId::from_bytes([3u8; 32]);
cache.put(&id, b"bytes").unwrap();
cache.note(stored(NodeKind::Literal, &id, 5));
cache.note(hit(NodeKind::Literal, &id, 5));
assert_eq!(cache.promoted.total_bytes().unwrap(), 0);
assert_eq!(cache.get(&id).unwrap().as_deref(), Some(&b"bytes"[..]));
fs::remove_dir_all(&root).ok();
}
#[test]
fn a_promoted_hit_is_served_after_the_ephemeral_cache_is_cleared() {
let root = temp_root("serve");
let gov = Governor::with_policy(promoted_policy(1, 1 << 20));
let mut cache = GovernedCache::open(&root, gov).unwrap();
let id = NodeId::from_bytes([7u8; 32]);
let cold = b"promoted intermediate output".to_vec();
cache.put(&id, &cold).unwrap();
cache.note(stored(
NodeKind::PackageMemberDecoded,
&id,
cold.len() as u64,
));
cache.note(hit(NodeKind::PackageMemberDecoded, &id, cold.len() as u64));
assert!(cache.promoted.contains(&id).unwrap());
cache.ephemeral.clear().unwrap();
assert_eq!(cache.ephemeral.get(&id).unwrap(), None);
assert_eq!(cache.get(&id).unwrap(), Some(cold));
fs::remove_dir_all(&root).ok();
}
#[test]
fn a_promoted_entry_for_a_changed_id_is_a_miss() {
let root = temp_root("changed");
let gov = Governor::with_policy(promoted_policy(1, 1 << 20));
let mut cache = GovernedCache::open(&root, gov).unwrap();
let id = NodeId::from_bytes([1u8; 32]);
cache.put(&id, b"aaa").unwrap();
cache.note(stored(NodeKind::DocxStory, &id, 3));
cache.note(hit(NodeKind::DocxStory, &id, 3));
assert!(cache.promoted.contains(&id).unwrap());
let other = NodeId::from_bytes([2u8; 32]);
assert_eq!(cache.get(&other).unwrap(), None);
fs::remove_dir_all(&root).ok();
}
#[test]
fn eviction_respects_the_byte_budget() {
let root = temp_root("evict");
let gov = Governor::with_policy(promoted_policy(1, 5));
let mut cache = GovernedCache::open(&root, gov).unwrap();
let a = NodeId::from_bytes([1u8; 32]);
let b = NodeId::from_bytes([2u8; 32]);
cache.put(&a, b"aaa").unwrap();
cache.note(stored(NodeKind::Literal, &a, 3));
cache.note(hit(NodeKind::Literal, &a, 3));
assert!(cache.promoted.contains(&a).unwrap());
cache.put(&b, b"bbb").unwrap();
cache.note(stored(NodeKind::Literal, &b, 3));
cache.note(hit(NodeKind::Literal, &b, 3));
assert!(cache.promoted.contains(&b).unwrap());
assert!(!cache.promoted.contains(&a).unwrap());
assert!(cache.promoted.output_bytes().unwrap() <= 5);
fs::remove_dir_all(&root).ok();
}
#[test]
fn oversized_nodes_are_never_promoted() {
let root = temp_root("oversize");
let gov = Governor::with_policy(PromotePolicy {
max_node_bytes: 4,
..promoted_policy(1, 1 << 20)
});
let mut cache = GovernedCache::open(&root, gov).unwrap();
let id = NodeId::from_bytes([5u8; 32]);
cache.put(&id, b"giant final answer").unwrap();
cache.note(stored(NodeKind::DocumentExact, &id, 18));
cache.note(hit(NodeKind::DocumentExact, &id, 18));
assert_eq!(cache.promoted.total_bytes().unwrap(), 0);
fs::remove_dir_all(&root).ok();
}
#[test]
fn poisoned_promoted_bytes_fail_closed_never_wrong() {
let root = temp_root("poison");
let mut store = PromotedStore::open(&root).unwrap();
let id = NodeId::from_bytes([9u8; 32]);
store.put(&id, b"correct bytes").unwrap();
fs::write(store.bytes_path(&id), b"wrong!!").unwrap();
assert_eq!(
store.get(&id).unwrap_err().class(),
crate::ErrorClass::IntegrityMismatch
);
fs::remove_dir_all(&root).ok();
}
}