use std::{
cmp,
sync::Arc,
thread,
time::{Duration, Instant}
};
use parking_lot::{Condvar, Mutex, MutexGuard};
use crate::ConnPool;
pub enum Vacuum {
Nothing,
Once(usize),
Repeat(usize)
}
pub struct Params {
pub autorun: bool,
pub page_count: Box<dyn FnMut(usize) -> Vacuum + Send>,
pub cooldown: Duration
}
impl Default for Params {
fn default() -> Self {
let page_count = |npages| {
let n = cmp::min(npages / 2, 256);
if n > 16 {
Vacuum::Repeat(n)
} else {
Vacuum::Nothing
}
};
Self {
autorun: false,
page_count: Box::new(page_count),
cooldown: Duration::from_secs(30)
}
}
}
#[derive(Default)]
struct Mutable {
shutdown: bool,
do_vacuum: bool
}
impl Mutable {
fn new() -> Mutex<Self> {
Mutex::new(Self::default())
}
}
#[derive(Default)]
struct Shared {
rw: Mutex<Mutable>,
signal: Condvar
}
impl Shared {
fn new() -> Arc<Self> {
let slf = Self {
rw: Mutable::new(),
signal: Condvar::new()
};
Arc::new(slf)
}
fn lock(&self) -> MutexGuard<'_, Mutable> {
self.rw.lock()
}
#[inline]
pub(crate) fn with_lock<F, R>(&self, f: F) -> R
where
F: FnOnce(&mut Mutable) -> R
{
let mut rw = self.lock();
f(&mut rw)
}
}
pub struct VacCtx {
sh: Arc<Shared>,
jh: Option<thread::JoinHandle<()>>
}
impl VacCtx {
pub fn vacuum(&self) {
self.sh.with_lock(|rw| {
rw.do_vacuum = true;
self.sh.signal.notify_one();
});
}
pub const fn take_jh(&mut self) -> Option<thread::JoinHandle<()>> {
self.jh.take()
}
}
impl Drop for VacCtx {
fn drop(&mut self) {
self.sh.with_lock(|rw| {
rw.shutdown = true;
self.sh.signal.notify_one();
});
if let Some(jh) = self.jh.take() {
let _ = jh.join();
}
}
}
#[must_use]
pub fn run(cp: ConnPool, params: Params) -> VacCtx {
let sh = Shared::new();
let sh2 = Arc::clone(&sh);
let jh = thread::spawn(move || {
autovacuum(&sh2, &cp, params);
});
VacCtx { sh, jh: Some(jh) }
}
#[allow(clippy::significant_drop_tightening)]
fn autovacuum(
sh: &Arc<Shared>,
cp: &ConnPool,
Params {
autorun,
mut page_count,
cooldown
}: Params
) {
let mut last_run: Option<Instant> = None;
let mut rw = sh.lock();
rw.do_vacuum = autorun;
loop {
if rw.shutdown {
break;
}
if let Some(lastrun) = last_run {
let now = Instant::now();
let elapsed = now - lastrun;
if elapsed < cooldown {
let dur = cooldown.checked_sub(elapsed).unwrap();
if sh.signal.wait_for(&mut rw, dur).timed_out() {
last_run = None;
}
continue;
}
last_run = None;
}
if rw.do_vacuum {
rw.do_vacuum = false;
let Ok(nfree) = cp.freelist_count() else {
continue;
};
match page_count(nfree) {
Vacuum::Nothing => {
last_run = Some(Instant::now());
}
Vacuum::Once(n) => {
let wrconn = cp.writer();
if wrconn.incremental_vacuum(Some(n)).is_err() {
}
last_run = Some(Instant::now());
}
Vacuum::Repeat(n) => {
let wrconn = cp.writer();
if wrconn.incremental_vacuum(Some(n)).is_err() {
}
rw.do_vacuum = true;
last_run = Some(Instant::now());
}
}
} else {
sh.signal.wait(&mut rw);
}
}
}