use lru::LruCache;
use std::collections::HashMap;
use std::hash::{Hash, Hasher};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, LazyLock, Mutex};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct PartitionLoc {
pub data_offset: u64,
pub data_size: u32,
}
impl PartitionLoc {
#[inline]
pub fn new(data_offset: u64, data_size: u32) -> Self {
Self {
data_offset,
data_size,
}
}
#[inline]
pub fn offset_only(data_offset: u64) -> Self {
Self {
data_offset,
data_size: 0,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub struct GenerationIdentity {
pub device: u64,
pub inode: u64,
pub size: u64,
pub generation: u64,
pub mtime_ns: i128,
}
impl GenerationIdentity {
pub fn resolve(path: &std::path::Path, generation: u64) -> Option<Self> {
let (device, inode, size, mtime_ns) = stat_identity(path)?;
Some(Self {
device,
inode,
size,
generation,
mtime_ns,
})
}
}
#[cfg(unix)]
fn stat_identity(path: &std::path::Path) -> Option<(u64, u64, u64, i128)> {
use std::os::unix::fs::MetadataExt;
let md = std::fs::metadata(path).ok()?;
let mtime_ns = md.mtime() as i128 * 1_000_000_000 + md.mtime_nsec() as i128;
Some((md.dev(), md.ino(), md.len(), mtime_ns))
}
#[cfg(not(unix))]
fn stat_identity(path: &std::path::Path) -> Option<(u64, u64, u64, i128)> {
let md = std::fs::metadata(path).ok()?;
let mtime_ns = md
.modified()
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_nanos() as i128)
.unwrap_or(0);
Some((0, 0, md.len(), mtime_ns))
}
const PER_ENTRY_OVERHEAD: usize = 96;
#[inline]
const fn entry_cost(key_len: usize) -> usize {
key_len.saturating_add(PER_ENTRY_OVERHEAD)
}
pub const DEFAULT_GLOBAL_KEY_CACHE_BYTES: usize = 64 * 1024 * 1024;
pub const DEFAULT_GLOBAL_KEY_CACHE_SHARDS: usize = 128;
#[derive(Clone, Copy, Debug)]
struct Entry {
loc: PartitionLoc,
seq: u64,
}
struct Shard {
map: HashMap<GenerationIdentity, LruCache<Box<[u8]>, Entry>>,
current_bytes: usize,
seq: u64,
}
impl Shard {
fn new() -> Self {
Self {
map: HashMap::new(),
current_bytes: 0,
seq: 0,
}
}
#[inline]
fn next_seq(&mut self) -> u64 {
self.seq = self.seq.wrapping_add(1);
self.seq
}
fn total_len(&self) -> usize {
self.map.values().map(LruCache::len).sum()
}
fn get(&mut self, identity: &GenerationIdentity, key: &[u8]) -> Option<PartitionLoc> {
let seq = self.next_seq();
let inner = self.map.get_mut(identity)?;
let entry = inner.get_mut(key)?;
entry.seq = seq;
Some(entry.loc)
}
fn evict_one(&mut self) -> bool {
let mut victim: Option<GenerationIdentity> = None;
let mut min_seq = u64::MAX;
for (id, inner) in self.map.iter() {
if let Some((_, entry)) = inner.peek_lru() {
if victim.is_none() || entry.seq < min_seq {
min_seq = entry.seq;
victim = Some(*id);
}
}
}
let Some(id) = victim else {
return false;
};
let Some(inner) = self.map.get_mut(&id) else {
return false;
};
match inner.pop_lru() {
Some((k, _)) => {
self.current_bytes = self.current_bytes.saturating_sub(entry_cost(k.len()));
if inner.is_empty() {
self.map.remove(&id);
}
true
}
None => false,
}
}
}
#[derive(Debug, Clone, Copy, Default)]
pub(crate) struct GlobalKeyCacheSnapshot {
pub hits: u64,
pub misses: u64,
pub evictions: u64,
pub invalidations: u64,
pub resident_bytes: usize,
pub capacity_bytes: usize,
}
pub struct GlobalKeyOffsetCache {
shards: Box<[Mutex<Shard>]>,
per_shard_bytes: usize,
mask: usize,
disabled: bool,
hits: AtomicU64,
misses: AtomicU64,
evictions: AtomicU64,
invalidations: AtomicU64,
}
impl std::fmt::Debug for GlobalKeyOffsetCache {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("GlobalKeyOffsetCache")
.field("shards", &self.shards.len())
.field("per_shard_bytes", &self.per_shard_bytes)
.field("disabled", &self.disabled)
.field("len", &self.len())
.field("resident_bytes", &self.resident_bytes())
.field("hits", &self.hits.load(Ordering::Relaxed))
.field("misses", &self.misses.load(Ordering::Relaxed))
.field("evictions", &self.evictions.load(Ordering::Relaxed))
.field("invalidations", &self.invalidations.load(Ordering::Relaxed))
.finish()
}
}
static GLOBAL: LazyLock<Arc<GlobalKeyOffsetCache>> = LazyLock::new(|| {
Arc::new(GlobalKeyOffsetCache::with_budget_bytes(
DEFAULT_GLOBAL_KEY_CACHE_BYTES,
))
});
impl GlobalKeyOffsetCache {
pub fn global() -> Arc<GlobalKeyOffsetCache> {
Arc::clone(&GLOBAL)
}
pub fn with_budget_bytes(total_budget_bytes: usize) -> Self {
Self::with_budget_and_shards(total_budget_bytes, DEFAULT_GLOBAL_KEY_CACHE_SHARDS)
}
pub fn with_budget_and_shards(total_budget_bytes: usize, shard_count: usize) -> Self {
let shard_count = shard_count.max(1).next_power_of_two();
let per_shard_bytes = (total_budget_bytes / shard_count).max(entry_cost(0));
let mut shards = Vec::with_capacity(shard_count);
for _ in 0..shard_count {
shards.push(Mutex::new(Shard::new()));
}
Self {
shards: shards.into_boxed_slice(),
per_shard_bytes,
mask: shard_count - 1,
disabled: false,
hits: AtomicU64::new(0),
misses: AtomicU64::new(0),
evictions: AtomicU64::new(0),
invalidations: AtomicU64::new(0),
}
}
pub fn disabled() -> Self {
Self {
shards: vec![Mutex::new(Shard::new())].into_boxed_slice(),
per_shard_bytes: 0,
mask: 0,
disabled: true,
hits: AtomicU64::new(0),
misses: AtomicU64::new(0),
evictions: AtomicU64::new(0),
invalidations: AtomicU64::new(0),
}
}
#[inline]
fn lock(m: &Mutex<Shard>) -> std::sync::MutexGuard<'_, Shard> {
m.lock().unwrap_or_else(|e| e.into_inner())
}
#[inline]
fn shard_for(&self, identity: &GenerationIdentity, key: &[u8]) -> &Mutex<Shard> {
let mut h = std::collections::hash_map::DefaultHasher::new();
identity.hash(&mut h);
key.hash(&mut h);
let idx = (h.finish() as usize) & self.mask;
&self.shards[idx]
}
pub fn get(&self, identity: GenerationIdentity, key: &[u8]) -> Option<PartitionLoc> {
if self.disabled {
return None;
}
let mut guard = Self::lock(self.shard_for(&identity, key));
let found = guard.get(&identity, key);
drop(guard);
match found {
Some(loc) => {
self.hits.fetch_add(1, Ordering::Relaxed);
Some(loc)
}
None => {
self.misses.fetch_add(1, Ordering::Relaxed);
None
}
}
}
pub fn insert(&self, identity: GenerationIdentity, key: &[u8], loc: PartitionLoc) {
if self.disabled {
return;
}
let cost = entry_cost(key.len());
let mut guard = Self::lock(self.shard_for(&identity, key));
let seq = guard.next_seq();
let inner = guard
.map
.entry(identity)
.or_insert_with(LruCache::unbounded);
let replaced = inner.put(key.into(), Entry { loc, seq }).is_some();
if replaced {
guard.current_bytes = guard.current_bytes.saturating_sub(cost);
}
guard.current_bytes = guard.current_bytes.saturating_add(cost);
let mut evicted_here: u64 = 0;
while guard.current_bytes > self.per_shard_bytes && guard.total_len() > 1 {
if !guard.evict_one() {
break;
}
evicted_here += 1;
}
drop(guard);
if evicted_here > 0 {
self.evictions.fetch_add(evicted_here, Ordering::Relaxed);
}
}
pub fn invalidate(&self, identity: GenerationIdentity) -> u64 {
if self.disabled {
return 0;
}
let mut dropped: u64 = 0;
for shard in self.shards.iter() {
let mut guard = Self::lock(shard);
if let Some(inner) = guard.map.remove(&identity) {
let mut reclaimed = 0usize;
let mut n: u64 = 0;
for (k, _) in inner.iter() {
reclaimed = reclaimed.saturating_add(entry_cost(k.len()));
n += 1;
}
guard.current_bytes = guard.current_bytes.saturating_sub(reclaimed);
dropped += n;
}
}
if dropped > 0 {
self.invalidations.fetch_add(dropped, Ordering::Relaxed);
}
dropped
}
pub fn invalidate_all(&self) -> u64 {
if self.disabled {
return 0;
}
let mut dropped: u64 = 0;
for shard in self.shards.iter() {
let mut guard = Self::lock(shard);
dropped = dropped.saturating_add(guard.total_len() as u64);
guard.map.clear();
guard.current_bytes = 0;
}
if dropped > 0 {
self.invalidations.fetch_add(dropped, Ordering::Relaxed);
}
dropped
}
pub fn len(&self) -> usize {
self.shards.iter().map(|m| Self::lock(m).total_len()).sum()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn resident_bytes(&self) -> usize {
self.shards
.iter()
.map(|m| Self::lock(m).current_bytes)
.sum()
}
pub fn budget_bytes(&self) -> usize {
self.per_shard_bytes.saturating_mul(self.shards.len())
}
pub fn hit_count(&self) -> u64 {
self.hits.load(Ordering::Relaxed)
}
pub fn miss_count(&self) -> u64 {
self.misses.load(Ordering::Relaxed)
}
pub fn eviction_count(&self) -> u64 {
self.evictions.load(Ordering::Relaxed)
}
pub fn invalidation_count(&self) -> u64 {
self.invalidations.load(Ordering::Relaxed)
}
pub(crate) fn snapshot(&self) -> GlobalKeyCacheSnapshot {
GlobalKeyCacheSnapshot {
hits: self.hit_count(),
misses: self.miss_count(),
evictions: self.eviction_count(),
invalidations: self.invalidation_count(),
resident_bytes: self.resident_bytes(),
capacity_bytes: self.budget_bytes(),
}
}
}
impl Default for GlobalKeyOffsetCache {
fn default() -> Self {
Self::with_budget_bytes(DEFAULT_GLOBAL_KEY_CACHE_BYTES)
}
}
#[cfg(test)]
#[path = "global_key_offset_tests.rs"]
mod tests;