use std::borrow::Cow;
use std::fs::File;
use std::io::Result;
use std::path::Path;
use std::path::PathBuf;
use std::sync::atomic::AtomicU8;
use std::sync::atomic::Ordering::Relaxed;
use std::sync::Arc;
use crate::cache_dir::CacheDir;
use crate::multiplicative_hash::MultiplicativeHash;
use crate::trigger::PeriodicTrigger;
use crate::Key;
use crate::KISMET_TEMPORARY_SUBDIRECTORY as TEMP_SUBDIR;
const MAINTENANCE_SCALE: usize = 2;
const PRIMARY_MIXER: MultiplicativeHash =
MultiplicativeHash::new_keyed(b"kismet: primary shard mixer");
const SECONDARY_MIXER: MultiplicativeHash =
MultiplicativeHash::new_keyed(b"kismet: secondary shard mixer");
#[derive(Clone, Debug)]
pub struct Cache {
load_estimates: Arc<[AtomicU8]>,
base_dir: PathBuf,
trigger: PeriodicTrigger,
num_shards: usize,
shard_capacity: usize,
}
#[inline]
fn format_id(shard: usize) -> String {
format!(".kismet_{:04x}", shard)
}
struct Shard {
id: usize,
shard_dir: PathBuf,
trigger: PeriodicTrigger,
capacity: usize,
}
impl Shard {
fn replace_shard(self, id: usize) -> Shard {
let mut shard_dir = self.shard_dir;
shard_dir.pop();
shard_dir.push(&format_id(id));
Shard {
id,
shard_dir,
trigger: self.trigger,
capacity: self.capacity,
}
}
fn file_exists(&mut self, name: &str) -> bool {
self.shard_dir.push(name);
let result = std::fs::metadata(&self.shard_dir);
self.shard_dir.pop();
result.is_ok()
}
}
impl CacheDir for Shard {
#[inline]
fn temp_dir(&self) -> Cow<Path> {
let mut dir = self.shard_dir.clone();
dir.push(TEMP_SUBDIR);
Cow::from(dir)
}
#[inline]
fn base_dir(&self) -> Cow<Path> {
Cow::from(&self.shard_dir)
}
#[inline]
fn trigger(&self) -> &PeriodicTrigger {
&self.trigger
}
#[inline]
fn capacity(&self) -> usize {
self.capacity
}
}
impl Cache {
pub fn new(base_dir: PathBuf, mut num_shards: usize, mut total_capacity: usize) -> Cache {
if num_shards < 2 {
num_shards = 2;
}
if total_capacity < num_shards {
total_capacity = num_shards;
}
let mut load_estimates = Vec::with_capacity(num_shards);
load_estimates.resize_with(num_shards, || AtomicU8::new(0));
let shard_capacity =
(total_capacity / num_shards) + ((total_capacity % num_shards) != 0) as usize;
let trigger =
PeriodicTrigger::new(shard_capacity.min(total_capacity / MAINTENANCE_SCALE) as u64);
Cache {
load_estimates: load_estimates.into_boxed_slice().into(),
base_dir,
trigger,
num_shards,
shard_capacity,
}
}
fn random_shard_id(&self) -> usize {
use rand::Rng;
rand::thread_rng().gen_range(0..self.num_shards)
}
fn other_shard_id(&self, base: usize, mut other: usize) -> usize {
if base != other {
return other;
}
other += 1;
if other < self.num_shards {
other
} else {
0
}
}
fn shard_ids(&self, key: Key) -> (usize, usize) {
let h1 = PRIMARY_MIXER.map(key.hash, self.num_shards);
let h2 = SECONDARY_MIXER.map(key.secondary_hash, self.num_shards);
(h1, self.other_shard_id(h1, h2))
}
fn sort_by_load(&self, (h1, h2): (usize, usize)) -> (usize, usize) {
let load1 = self.load_estimates[h1].load(Relaxed) as usize;
let load2 = self.load_estimates[h2].load(Relaxed) as usize;
let capacity = self.shard_capacity;
if load1.clamp(0, capacity) <= load2.clamp(0, capacity) {
(h1, h2)
} else {
(h2, h1)
}
}
fn shard(&self, shard_id: usize) -> Shard {
let mut dir = self.base_dir.clone();
dir.push(&format_id(shard_id));
Shard {
id: shard_id,
shard_dir: dir,
trigger: self.trigger,
capacity: self.shard_capacity,
}
}
pub fn get(&self, key: Key) -> Result<Option<File>> {
let (h1, h2) = self.shard_ids(key);
let shard = self.shard(h1);
if let Some(file) = shard.get(key.name)? {
Ok(Some(file))
} else {
shard.replace_shard(h2).get(key.name)
}
}
pub fn temp_dir(&self, key: Option<Key>) -> Result<Cow<Path>> {
let shard_id = match key {
Some(key) => self.sort_by_load(self.shard_ids(key)).0,
None => self.random_shard_id(),
};
let shard = self.shard(shard_id);
if self.trigger.event() {
shard.cleanup_temp_directory()?;
}
Ok(Cow::from(shard.ensure_temp_dir()?.into_owned()))
}
fn update_estimate(&self, shard_id: usize, update: Option<u64>) {
let target = &self.load_estimates[shard_id];
match update {
Some(remaining) => {
let update = remaining.clamp(0, u8::MAX as u64 - 1) as u8;
target.store(update + 1, Relaxed);
}
None => {
let _ = target.fetch_update(Relaxed, Relaxed, |i| {
if i < u8::MAX {
Some(i + 1)
} else {
None
}
});
}
};
}
fn force_maintain_shard(&self, shard: Shard) -> Result<()> {
let update = shard.maintain()?.clamp(0, u8::MAX as u64) as u8;
self.load_estimates[shard.id].store(update, Relaxed);
Ok(())
}
fn maintain_random_other_shard(&self, base: Shard) -> Result<()> {
let shard_id = self.other_shard_id(base.id, self.random_shard_id());
self.force_maintain_shard(base.replace_shard(shard_id))
}
pub fn set(&self, key: Key, value: &Path) -> Result<()> {
let (h1, h2) = self.sort_by_load(self.shard_ids(key));
let mut shard = self.shard(h2);
if !shard.file_exists(key.name) {
shard = shard.replace_shard(h1);
}
let update = shard.set(key.name, value)?;
self.update_estimate(h1, update);
if update.is_some() {
self.maintain_random_other_shard(shard)?;
} else if self.load_estimates[h1].load(Relaxed) as usize / 2 > self.shard_capacity {
self.force_maintain_shard(shard)?;
}
Ok(())
}
pub fn put(&self, key: Key, value: &Path) -> Result<()> {
let (h1, h2) = self.sort_by_load(self.shard_ids(key));
let mut shard = self.shard(h2);
if !shard.file_exists(key.name) {
shard = shard.replace_shard(h1);
}
let update = shard.put(key.name, value)?;
self.update_estimate(h1, update);
if update.is_some() {
self.maintain_random_other_shard(shard)?;
} else if self.load_estimates[h1].load(Relaxed) as usize / 2 > self.shard_capacity {
self.force_maintain_shard(shard)?;
}
Ok(())
}
pub fn touch(&self, key: Key) -> Result<bool> {
let (h1, h2) = self.shard_ids(key);
let shard = self.shard(h1);
if shard.touch(key.name)? {
return Ok(true);
}
shard.replace_shard(h2).touch(key.name)
}
}
#[test]
fn smoke_test() {
use tempfile::NamedTempFile;
use test_dir::{DirBuilder, TestDir};
const PAYLOAD_MULTIPLIER: usize = 113;
let temp = TestDir::temp();
let cache = Cache::new(temp.path("."), 3, 9);
for i in 0..200 {
let name = format!("{}", i);
let temp_dir = cache.temp_dir(None).expect("temp_dir must succeed");
let tmp = NamedTempFile::new_in(temp_dir).expect("new temp file must succeed");
std::fs::write(tmp.path(), format!("{}", PAYLOAD_MULTIPLIER * i))
.expect("write must succeed");
if (i % 2) != 0 {
cache
.put(Key::new(&name, i as u64, i as u64 + 42), tmp.path())
.expect("put must succeed");
} else {
cache
.set(Key::new(&name, i as u64, i as u64 + 42), tmp.path())
.expect("set must succeed");
}
}
let present: usize = (0..200)
.map(|i| {
let name = format!("{}", i);
match cache
.get(Key::new(&name, i as u64, i as u64 + 42))
.expect("get must succeed")
{
Some(mut file) => {
use std::io::Read;
let mut buf = Vec::new();
file.read_to_end(&mut buf).expect("read must succeed");
assert_eq!(buf, format!("{}", PAYLOAD_MULTIPLIER * i).into_bytes());
1
}
None => 0,
}
})
.sum();
assert!(present >= 9);
assert!(present <= 18);
}
#[test]
fn test_set() {
use std::io::{Read, Write};
use tempfile::NamedTempFile;
use test_dir::{DirBuilder, TestDir};
let temp = TestDir::temp();
let cache = Cache::new(temp.path("."), 0, 0);
{
let tmp = NamedTempFile::new_in(cache.temp_dir(None).expect("temp_dir must succeed"))
.expect("new temp file must succeed");
tmp.as_file().write_all(b"v1").expect("write must succeed");
cache
.set(Key::new("entry", 1, 2), tmp.path())
.expect("initial set must succeed");
}
{
let mut cached = cache
.get(Key::new("entry", 1, 2))
.expect("must succeed")
.expect("must be found");
let mut dst = Vec::new();
cached.read_to_end(&mut dst).expect("read must succeed");
assert_eq!(&dst, b"v1");
}
{
let tmp = NamedTempFile::new_in(cache.temp_dir(None).expect("temp_dir must succeed"))
.expect("new temp file must succeed");
tmp.as_file().write_all(b"v2").expect("write must succeed");
cache
.set(Key::new("entry", 1, 2), tmp.path())
.expect("overwrite must succeed");
}
{
let mut cached = cache
.get(Key::new("entry", 1, 2))
.expect("must succeed")
.expect("must be found");
let mut dst = Vec::new();
cached.read_to_end(&mut dst).expect("read must succeed");
assert_eq!(&dst, b"v2");
}
}
#[test]
fn test_put() {
use std::io::{Read, Write};
use tempfile::NamedTempFile;
use test_dir::{DirBuilder, TestDir};
let temp = TestDir::temp();
let cache = Cache::new(temp.path("."), 0, 0);
{
let tmp = NamedTempFile::new_in(cache.temp_dir(None).expect("temp_dir must succeed"))
.expect("new temp file must succeed");
tmp.as_file().write_all(b"v1").expect("write must succeed");
cache
.set(Key::new("entry", 1, 2), tmp.path())
.expect("initial set must succeed");
}
{
let tmp = NamedTempFile::new_in(cache.temp_dir(None).expect("temp_dir must succeed"))
.expect("new temp file must succeed");
tmp.as_file().write_all(b"v2").expect("write must succeed");
cache
.put(Key::new("entry", 1, 2), tmp.path())
.expect("put must succeed");
}
{
let mut cached = cache
.get(Key::new("entry", 1, 2))
.expect("must succeed")
.expect("must be found");
let mut dst = Vec::new();
cached.read_to_end(&mut dst).expect("read must succeed");
assert_eq!(&dst, b"v1");
}
}
#[test]
fn test_touch() {
use std::io::Read;
use tempfile::NamedTempFile;
use test_dir::{DirBuilder, TestDir};
const PAYLOAD_MULTIPLIER: usize = 113;
let temp = TestDir::temp();
let cache = Cache::new(temp.path("."), 2, 600);
for i in 0..2000 {
assert_eq!(
cache
.touch(Key::new("0", 0, 42))
.expect("touch must succeed"),
i > 0
);
let name = format!("{}", i);
let temp_dir = cache.temp_dir(None).expect("temp_dir must succeed");
let tmp = NamedTempFile::new_in(temp_dir).expect("new temp file must succeed");
std::fs::write(tmp.path(), format!("{}", PAYLOAD_MULTIPLIER * i))
.expect("write must succeed");
cache
.put(Key::new(&name, i as u64, i as u64 + 42), tmp.path())
.expect("put must succeed");
if i == 0 {
std::thread::sleep(std::time::Duration::from_secs(2));
}
}
let mut file = cache
.get(Key::new("0", 0, 42))
.expect("get must succeed")
.expect("file must be found");
let mut buf = Vec::new();
file.read_to_end(&mut buf).expect("read must succeed");
assert_eq!(buf, b"0");
}