use super::probe::{LockProbe, Probe, WriterProbe};
use crate::ui_state::Clock;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, MutexGuard, PoisonError};
use std::time::{Duration, Instant};
const TTL: Duration = Duration::from_secs(2);
pub(super) struct TtlCache<P, C: Clock> {
inner: P,
clock: C,
cache: Mutex<HashMap<PathBuf, (Instant, Probe)>>,
}
impl<P, C: Clock> TtlCache<P, C> {
pub(super) fn new(inner: P, clock: C) -> Self {
let cache = Mutex::new(HashMap::new());
Self {
inner,
clock,
cache,
}
}
fn entries(&self) -> MutexGuard<'_, HashMap<PathBuf, (Instant, Probe)>> {
self.cache.lock().unwrap_or_else(PoisonError::into_inner)
}
fn cached(&self, target: &Path, compute: impl FnOnce() -> Probe) -> Probe {
let now = self.clock.now();
if let Some(probe) = self.fresh(target, now) {
return probe;
}
let probe = compute();
self.entries().insert(target.to_path_buf(), (now, probe));
probe
}
fn fresh(&self, target: &Path, now: Instant) -> Option<Probe> {
let &(at, probe) = self.entries().get(target)?;
(now.saturating_duration_since(at) < TTL).then_some(probe)
}
pub(super) fn invalidate(&self, target: &Path) {
self.entries().remove(target);
}
}
impl<P: LockProbe, C: Clock> LockProbe for TtlCache<P, C> {
fn lock_state(&self, inbox_dir: &Path) -> Probe {
self.cached(inbox_dir, || self.inner.lock_state(inbox_dir))
}
}
impl<P: WriterProbe, C: Clock> WriterProbe for TtlCache<P, C> {
fn writer_state(&self, path: &Path) -> Probe {
self.cached(path, || self.inner.writer_state(path))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_support::FakeClock;
use std::cell::Cell;
struct CountingProbe {
answer: Probe,
lock_calls: Cell<usize>,
writer_calls: Cell<usize>,
}
impl CountingProbe {
fn new(answer: Probe) -> Self {
Self {
answer,
lock_calls: Cell::new(0),
writer_calls: Cell::new(0),
}
}
}
impl LockProbe for CountingProbe {
fn lock_state(&self, _dir: &Path) -> Probe {
self.lock_calls.set(self.lock_calls.get() + 1);
self.answer
}
}
impl WriterProbe for CountingProbe {
fn writer_state(&self, _path: &Path) -> Probe {
self.writer_calls.set(self.writer_calls.get() + 1);
self.answer
}
}
fn dir() -> &'static Path {
Path::new("/ws/inbox/agent-1")
}
#[test]
fn first_observation_is_a_miss_then_within_ttl_is_a_hit() {
let clock = FakeClock::new();
let cache = TtlCache::new(CountingProbe::new(Probe::Held), clock.handle());
assert_eq!(cache.lock_state(dir()), Probe::Held);
clock.advance(Duration::from_millis(1999));
assert_eq!(cache.lock_state(dir()), Probe::Held);
assert_eq!(cache.inner.lock_calls.get(), 1, "one observation, one hit");
}
#[test]
fn expired_entry_is_recomputed() {
let clock = FakeClock::new();
let cache = TtlCache::new(CountingProbe::new(Probe::Free), clock.handle());
assert_eq!(cache.lock_state(dir()), Probe::Free);
clock.advance(TTL);
assert_eq!(cache.lock_state(dir()), Probe::Free);
assert_eq!(cache.inner.lock_calls.get(), 2, "stale entry re-observed");
}
#[test]
fn invalidate_forces_the_next_read_to_recompute() {
let clock = FakeClock::new();
let cache = TtlCache::new(CountingProbe::new(Probe::Held), clock.handle());
assert_eq!(cache.lock_state(dir()), Probe::Held);
cache.invalidate(dir());
assert_eq!(cache.lock_state(dir()), Probe::Held);
assert_eq!(cache.inner.lock_calls.get(), 2, "eviction re-observed");
cache.invalidate(Path::new("/never/cached"));
}
#[test]
fn writer_question_is_cached_independently_of_the_lock() {
let clock = FakeClock::new();
let cache = TtlCache::new(CountingProbe::new(Probe::Held), clock.handle());
let file = Path::new("/ws/steps/agent-1/003/response.json");
assert_eq!(cache.writer_state(file), Probe::Held);
assert_eq!(cache.writer_state(file), Probe::Held);
assert_eq!(
cache.inner.writer_calls.get(),
1,
"writer cached by its key"
);
assert_eq!(cache.inner.lock_calls.get(), 0, "lock question untouched");
}
}