use yo_common::{Code, Error, Result, Rng};
use yo_index::RawMap;
use crate::access::{Lfu, Policy};
use crate::cold::{self, Blocks};
use crate::demote::Doorkeeper;
use crate::evict;
use crate::value::{self, Encoding, Kind};
pub const WINDOW: usize = 8192;
pub const WALK: usize = 1024;
pub const BARREN: usize = 16;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Faulted {
Missing,
Warm,
Served,
Promoted,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct Stats {
pub demoted: u64,
pub promoted: u64,
pub faults: u64,
pub served: u64,
pub bytes_out: u64,
pub bytes_in: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct Relief {
pub moved: usize,
pub freed: usize,
}
impl Relief {
#[must_use]
pub const fn made_room(self) -> bool {
self.moved > 0 || self.freed > 0
}
}
#[must_use]
pub fn worth_demoting(rec: &[u8]) -> bool {
let m = value::Meta::from_byte(rec[0]);
if m.is_cold() || m.kind() != Kind::String || m.encoding() == Encoding::Int {
return false;
}
rec.len() > value::cold_record_len(m.has_expiry())
}
pub struct Tier<B: Blocks> {
blocks: B,
door: Doorkeeper,
scratch: cold::Scratch,
pool: evict::Pool,
keybuf: Vec<u8>,
rng: Rng,
stats: Stats,
}
impl<B: Blocks> Tier<B> {
pub fn new(blocks: B) -> Tier<B> {
Tier::with_window(blocks, WINDOW)
}
pub fn with_window(blocks: B, window: usize) -> Tier<B> {
Tier {
blocks,
door: Doorkeeper::new(window),
scratch: cold::Scratch::new(),
pool: evict::Pool::new(),
keybuf: Vec::new(),
rng: Rng::new(0x5eed_1234_9abc_def0),
stats: Stats::default(),
}
}
#[must_use]
pub const fn stats(&self) -> Stats {
self.stats
}
#[must_use]
pub fn store_bytes(&self) -> u64 {
self.blocks.bytes()
}
pub const fn blocks(&self) -> &B {
&self.blocks
}
pub const fn blocks_mut(&mut self) -> &mut B {
&mut self.blocks
}
#[must_use]
pub fn memory_bytes(&self) -> usize {
self.door.memory_bytes()
+ self.scratch.memory_bytes()
+ self.pool.memory_bytes()
+ self.keybuf.capacity()
}
pub fn demote(&mut self, map: &mut RawMap, key: &[u8]) -> Result<bool> {
let Some(addr) = map.find(key) else {
return Ok(false);
};
let rec = map.value_at(addr);
if !worth_demoting(rec) {
return Ok(false);
}
let m = value::Meta::from_byte(rec[0]);
let (kind, enc) = (m.kind(), m.encoding());
let expire_at = value::expire_at(rec);
let was = value::access(rec).unwrap_or_default();
let value::Str::Bytes(bytes) = value::read(rec) else {
return Ok(false);
};
let len = bytes.len() as u32;
let chain = cold::write(&mut self.blocks, bytes, &mut self.scratch)?;
let wrote = map.set_with(
key,
value::cold_record_len(expire_at.is_some()),
|_| {},
|out| {
value::write_cold_record(out, kind, enc, chain.at, len, expire_at);
value::set_access(out, was);
value::has_expiry(out)
},
);
debug_assert!(wrote.is_some(), "the key was found a moment ago");
self.stats.demoted += 1;
self.stats.bytes_out += u64::from(len);
Ok(true)
}
pub fn stash(&mut self, bytes: &[u8]) -> Result<cold::Chain> {
let chain = cold::write(&mut self.blocks, bytes, &mut self.scratch)?;
self.stats.demoted += 1;
self.stats.bytes_out += chain.len;
Ok(chain)
}
pub fn fetch(&mut self, chain: cold::Chain, out: &mut Vec<u8>) -> Result<()> {
out.clear();
out.reserve(chain.len as usize);
self.blocks.release();
{
let reader = cold::Reader::open(&self.blocks, chain)?;
for piece in reader.range(0, reader.len()) {
out.extend_from_slice(piece?);
}
}
self.stats.faults += 1;
self.stats.bytes_in += chain.len;
self.stats.promoted += 1;
Ok(())
}
pub fn fault(&mut self, map: &mut RawMap, key: &[u8], out: &mut Vec<u8>) -> Result<Faulted> {
self.read(map, key, out, true)
}
pub fn thaw(&mut self, map: &mut RawMap, key: &[u8], out: &mut Vec<u8>) -> Result<Faulted> {
self.read(map, key, out, false)
}
fn read(
&mut self,
map: &mut RawMap,
key: &[u8],
out: &mut Vec<u8>,
ask: bool,
) -> Result<Faulted> {
let Some(addr) = map.find(key) else {
return Ok(Faulted::Missing);
};
let rec = map.value_at(addr);
let Some(c) = value::cold(rec) else {
return Ok(Faulted::Warm);
};
let m = value::Meta::from_byte(rec[0]);
if m.kind().is_body() {
return Err(Error::new(
Code::Invalid,
"a demoted body cannot be read back as a string",
)
.with_detail(m.kind().name().to_string()));
}
let enc = m.encoding();
let expire_at = value::expire_at(rec);
let was = value::access(rec).unwrap_or_default();
out.clear();
out.reserve(c.len as usize);
let chain = cold::Chain {
at: c.at,
len: u64::from(c.len),
};
self.blocks.release();
{
let reader = cold::Reader::open(&self.blocks, chain)?;
for piece in reader.range(0, reader.len()) {
out.extend_from_slice(piece?);
}
}
self.stats.faults += 1;
self.stats.bytes_in += u64::from(c.len);
if ask && !self.door.admit(RawMap::hash_of(key)) {
self.stats.served += 1;
return Ok(Faulted::Served);
}
let wrote = map.set_with(
key,
value::record_len(enc, out.len(), expire_at.is_some()),
|_| {},
|dst| {
value::write_record(dst, enc, out, expire_at);
value::set_access(dst, was);
value::has_expiry(dst)
},
);
debug_assert!(wrote.is_some(), "the key was found a moment ago");
self.stats.promoted += 1;
Ok(Faulted::Promoted)
}
pub fn relieve(
&mut self,
map: &mut RawMap,
budget: usize,
policy: Policy,
now_ms: u64,
lfu: Lfu,
) -> Result<Relief> {
let start = map.memory_bytes();
let mut moved = 0;
let mut barren = 0;
while map.memory_bytes() > budget {
let round = self.round(map, policy, now_ms, lfu)?;
while map.memory_bytes() > budget && map.compact_hard().is_some() {}
if round == 0 {
barren += 1;
if barren == BARREN {
break;
}
continue;
}
barren = 0;
moved += round;
}
Ok(Relief {
moved,
freed: start.saturating_sub(map.memory_bytes()),
})
}
fn round(&mut self, map: &mut RawMap, policy: Policy, now_ms: u64, lfu: Lfu) -> Result<usize> {
self.pool.clear();
let r = self.rng.next_u64();
let pool = &mut self.pool;
let mut seen = 0usize;
let mut found = 0usize;
map.sample(r, |k, v, _| {
seen += 1;
if worth_demoting(v) {
pool.offer(k, evict::score(v, policy, now_ms, lfu));
found += 1;
}
found < evict::CANDIDATES && seen < WALK
});
let mut moved = 0;
let mut kb = core::mem::take(&mut self.keybuf);
while let Some(k) = self.pool.take() {
kb.clear();
kb.extend_from_slice(k);
if self.demote(map, &kb)? {
moved += 1;
}
}
self.keybuf = kb;
Ok(moved)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::access::Access;
use yo_common::{Addr, Code, Error, Space};
struct Mem {
blobs: Vec<Vec<u8>>,
reads: std::cell::Cell<usize>,
}
impl Mem {
fn new() -> Mem {
Mem {
blobs: Vec::new(),
reads: std::cell::Cell::new(0),
}
}
}
impl Blocks for Mem {
fn put(&mut self, bytes: &[u8]) -> Result<Addr> {
self.blobs.push(bytes.to_vec());
Ok(Addr::new(Space::Log, (self.blobs.len() - 1) as u64))
}
fn get(&self, at: Addr) -> Result<&[u8]> {
self.reads.set(self.reads.get() + 1);
self.blobs
.get(at.offset() as usize)
.map(Vec::as_slice)
.ok_or_else(|| Error::new(Code::NotFound, "no such block"))
}
fn bytes(&self) -> u64 {
self.blobs.iter().map(|b| b.len() as u64).sum()
}
}
fn tier() -> Tier<Mem> {
Tier::new(Mem::new())
}
fn map_with(key: &[u8], val: &[u8], expire_at: Option<u64>) -> RawMap {
let mut m = RawMap::new();
put(&mut m, key, val, expire_at);
m
}
fn put(m: &mut RawMap, key: &[u8], val: &[u8], expire_at: Option<u64>) {
let enc = Encoding::of(val);
let len = value::record_len(enc, val.len(), expire_at.is_some());
m.set_with(
key,
len,
|_| {},
|out| {
value::write_record(out, enc, val, expire_at);
value::has_expiry(out)
},
);
}
fn fault_twice(t: &mut Tier<Mem>, m: &mut RawMap, key: &[u8]) -> (Faulted, Faulted, Vec<u8>) {
let mut out = Vec::new();
let first = t.fault(m, key, &mut out).expect("a first read");
let second = t.fault(m, key, &mut out).expect("a second read");
(first, second, out)
}
#[test]
fn a_value_goes_out_to_the_file_and_the_record_shrinks_to_a_pointer() {
let val = vec![b'x'; 4000];
let mut m = map_with(b"k", &val, None);
let before = m.value_at(m.find(b"k").expect("there")).len();
let mut t = tier();
assert!(t.demote(&mut m, b"k").expect("demoted"));
let rec = m.value_at(m.find(b"k").expect("still there"));
assert!(rec.len() < before / 100, "the record did not shrink");
assert_eq!(value::cold(rec).expect("cold").len, 4000);
assert_eq!(t.stats().demoted, 1);
assert_eq!(t.stats().bytes_out, 4000);
}
#[test]
fn the_questions_that_do_not_want_the_bytes_are_still_answered_in_memory() {
let val = vec![b'y'; 900];
let deadline = Some(1_900_000_000_000);
let mut m = map_with(b"k", &val, deadline);
let mut t = tier();
t.demote(&mut m, b"k").expect("demoted");
let rec = m.value_at(m.find(b"k").expect("there"));
assert_eq!(value::str_len(rec), Some(900));
assert_eq!(value::kind(rec), Kind::String);
assert_eq!(value::Meta::from_byte(rec[0]).encoding(), Encoding::Raw);
assert_eq!(value::expire_at(rec), deadline);
assert_eq!(t.blocks().reads.get(), 0, "answering those read the device");
}
#[test]
fn a_value_too_short_to_be_worth_moving_is_left_where_it_is() {
let mut m = map_with(b"k", b"hello-world!", None);
let mut t = tier();
assert!(!t.demote(&mut m, b"k").expect("asked"));
assert!(value::cold(m.value_at(m.find(b"k").expect("there"))).is_none());
}
#[test]
fn an_int_encoded_value_is_never_moved() {
let mut m = map_with(b"k", b"1234567890123", None);
let mut t = tier();
assert!(!t.demote(&mut m, b"k").expect("asked"));
}
#[test]
fn a_key_that_is_not_there_is_a_no_and_not_an_error() {
let mut m = RawMap::new();
let mut t = tier();
assert!(!t.demote(&mut m, b"nothing").expect("asked"));
let mut out = Vec::new();
assert_eq!(
t.fault(&mut m, b"nothing", &mut out).expect("asked"),
Faulted::Missing
);
}
#[test]
fn demoting_twice_is_a_no_the_second_time() {
let val = vec![b'z'; 500];
let mut m = map_with(b"k", &val, None);
let mut t = tier();
assert!(t.demote(&mut m, b"k").expect("demoted"));
assert!(!t.demote(&mut m, b"k").expect("asked again"));
assert_eq!(t.stats().demoted, 1);
}
#[test]
fn a_resident_key_is_warm_and_the_buffer_is_left_alone() {
let mut m = map_with(b"k", b"a value long enough to matter", None);
let mut t = tier();
let mut out = vec![1, 2, 3];
assert_eq!(
t.fault(&mut m, b"k", &mut out).expect("read"),
Faulted::Warm
);
assert_eq!(out, vec![1, 2, 3], "a warm read touched the buffer");
assert_eq!(t.stats().faults, 0);
}
#[test]
fn the_first_read_serves_from_the_file_and_the_second_brings_it_back() {
let val = vec![b'q'; 3000];
let mut m = map_with(b"k", &val, None);
let mut t = tier();
t.demote(&mut m, b"k").expect("demoted");
let (first, second, out) = fault_twice(&mut t, &mut m, b"k");
assert_eq!(first, Faulted::Served, "one read earned a slot in memory");
assert_eq!(second, Faulted::Promoted);
assert_eq!(out, val);
assert_eq!(t.stats().faults, 2);
assert_eq!(t.stats().served, 1);
assert_eq!(t.stats().promoted, 1);
let mut again = Vec::new();
assert_eq!(
t.fault(&mut m, b"k", &mut again).expect("read"),
Faulted::Warm
);
assert_eq!(
value::read(m.value_at(m.find(b"k").expect("there"))).len(),
3000
);
}
#[test]
fn a_scan_over_cold_data_promotes_nothing() {
let mut m = RawMap::new();
let val = vec![b'c'; 700];
for i in 0..64u32 {
put(&mut m, &i.to_le_bytes(), &val, None);
}
let mut t = tier();
for i in 0..64u32 {
t.demote(&mut m, &i.to_le_bytes()).expect("demoted");
}
let mut out = Vec::new();
for i in 0..64u32 {
t.fault(&mut m, &i.to_le_bytes(), &mut out).expect("read");
}
assert_eq!(
t.stats().promoted,
0,
"a single pass over cold keys pulled some back in"
);
assert_eq!(t.stats().served, 64);
}
#[test]
fn the_deadline_and_the_access_field_survive_a_round_trip() {
let val = vec![b'r'; 1200];
let deadline = Some(1_888_777_666_555);
let mut m = map_with(b"k", &val, deadline);
let a = Access::lru(1_000_000);
{
let addr = m.find(b"k").expect("there");
value::set_access(m.value_at_mut(addr), a);
}
let mut t = tier();
t.demote(&mut m, b"k").expect("demoted");
assert_eq!(
value::access(m.value_at(m.find(b"k").expect("there"))),
Some(a),
"demotion looked like a use"
);
let (_, _, out) = fault_twice(&mut t, &mut m, b"k");
assert_eq!(out, val);
let rec = m.value_at(m.find(b"k").expect("there"));
assert_eq!(value::expire_at(rec), deadline);
assert_eq!(value::access(rec), Some(a));
}
#[test]
fn a_value_bigger_than_one_chunk_makes_the_trip_as_well() {
let val: Vec<u8> = (0..cold::CHUNK * 2 + 77).map(|i| (i % 251) as u8).collect();
let mut m = map_with(b"big", &val, None);
let mut t = tier();
assert!(t.demote(&mut m, b"big").expect("demoted"));
let (_, _, out) = fault_twice(&mut t, &mut m, b"big");
assert_eq!(out, val);
}
#[test]
fn relieve_moves_values_out_until_the_map_fits() {
let mut m = RawMap::new();
let val = vec![b'p'; 2000];
for i in 0..4_000u32 {
put(&mut m, &i.to_le_bytes(), &val, None);
}
let full = m.memory_bytes();
let budget = full / 2;
let mut t = tier();
let moved = t
.relieve(
&mut m,
budget,
Policy::AllKeysLru,
2_000_000,
Lfu::default(),
)
.expect("relieved");
assert!(moved.moved > 0, "nothing was moved");
assert!(
m.memory_bytes() <= budget,
"still {} bytes against a budget of {budget}",
m.memory_bytes()
);
assert_eq!(m.len(), 4_000);
}
#[test]
fn one_unlucky_round_does_not_end_the_sweep() {
let mut m = RawMap::new();
let val = vec![b'u'; 2000];
for i in 0..4_000u32 {
put(&mut m, &i.to_le_bytes(), &val, None);
}
let mut t = tier();
t.relieve(&mut m, 1, Policy::AllKeysLru, 2_000_000, Lfu::default())
.expect("relieved");
let cold = (0..4_000u32)
.filter(|i| {
let addr = m.find(&i.to_le_bytes()).expect("still there");
value::cold(m.value_at(addr)).is_some()
})
.count();
assert!(
cold > 3_900,
"only {cold} of 4000 were moved, so the sweep gave up early"
);
}
#[test]
fn the_memory_the_map_holds_actually_goes_down() {
let mut m = RawMap::new();
let val = vec![b'v'; 2000];
for i in 0..4_000u32 {
put(&mut m, &i.to_le_bytes(), &val, None);
}
let before = m.memory_bytes();
let mut t = tier();
t.relieve(&mut m, 1, Policy::AllKeysLru, 2_000_000, Lfu::default())
.expect("relieved");
let mut bare = RawMap::new();
let stub = vec![b'v'; 4];
for i in 0..4_000u32 {
put(&mut bare, &i.to_le_bytes(), &stub, None);
}
let floor = bare.memory_bytes();
assert!(
m.memory_bytes() <= floor,
"{before} bytes went to {}, and the floor is {floor}",
m.memory_bytes()
);
}
#[test]
fn relieve_gives_up_rather_than_spinning_when_nothing_is_worth_moving() {
let mut m = RawMap::new();
for i in 0..200u32 {
put(&mut m, &i.to_le_bytes(), b"tiny", None);
}
let mut t = tier();
let moved = t
.relieve(&mut m, 1, Policy::AllKeysLru, 2_000_000, Lfu::default())
.expect("asked");
assert_eq!(moved, Relief::default());
}
#[test]
fn a_sweep_that_moves_nothing_and_frees_a_segment_still_says_it_made_room() {
let mut m = RawMap::new();
let val = vec![b'v'; 4096];
for i in 0..2_000u32 {
put(&mut m, &i.to_le_bytes(), &val, None);
}
let mut t = tier();
for i in 0..2_000u32 {
assert!(
t.demote(&mut m, &i.to_le_bytes()).expect("demoted"),
"key {i} did not go out"
);
}
let before = m.memory_bytes();
let r = t
.relieve(
&mut m,
before - 1,
Policy::AllKeysLru,
2_000_000,
Lfu::default(),
)
.expect("swept");
assert_eq!(r.moved, 0, "there was nothing left in memory to move");
assert!(
r.freed > 0,
"compaction gave nothing back, so this checked nothing"
);
assert!(
r.made_room(),
"a sweep that freed {} said it did not",
r.freed
);
assert_eq!(m.len(), 2_000, "a sweep that lost keys");
}
#[test]
fn what_relieve_moved_still_reads_back_byte_for_byte() {
let mut m = RawMap::new();
let mut want = Vec::new();
for i in 0..4_000u32 {
let val: Vec<u8> = (0..900).map(|j| (i as usize + j) as u8).collect();
put(&mut m, &i.to_le_bytes(), &val, None);
want.push(val);
}
let budget = m.memory_bytes() / 2;
let mut t = tier();
let moved = t
.relieve(
&mut m,
budget,
Policy::AllKeysLru,
2_000_000,
Lfu::default(),
)
.expect("relieved");
assert!(
moved.moved > 0,
"nothing was moved, so this checked nothing"
);
let mut out = Vec::new();
for (i, val) in want.iter().enumerate() {
let key = (i as u32).to_le_bytes();
match t.fault(&mut m, &key, &mut out).expect("read") {
Faulted::Warm => {
let rec = m.value_at(m.find(&key).expect("there"));
assert_eq!(value::read(rec), value::Str::Bytes(val));
}
Faulted::Served | Faulted::Promoted => assert_eq!(&out, val),
Faulted::Missing => panic!("key {i} went missing"),
}
}
}
}