use std::collections::{BTreeSet, HashMap};
use std::time::{Duration, Instant};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum VacuumError {
AlreadyRunning,
ConflictingLock(String),
InvalidId(String),
Internal(String),
}
impl std::fmt::Display for VacuumError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::AlreadyRunning => write!(f, "Vacuum already running"),
Self::ConflictingLock(msg) => write!(f, "Conflicting lock: {msg}"),
Self::InvalidId(msg) => write!(f, "Invalid page/slot id: {msg}"),
Self::Internal(msg) => write!(f, "Internal vacuum error: {msg}"),
}
}
}
impl std::error::Error for VacuumError {}
pub type VacuumResult<T> = Result<T, VacuumError>;
#[derive(Debug, Clone, PartialEq)]
pub enum VacuumTrigger {
Manual,
FragmentationThreshold(f64),
TimeBased(Duration),
Combined {
threshold: f64,
interval: Duration,
},
}
impl Default for VacuumTrigger {
fn default() -> Self {
Self::FragmentationThreshold(0.25)
}
}
#[derive(Debug, Clone)]
pub struct VacuumConfig {
pub trigger: VacuumTrigger,
pub online: bool,
pub max_pages_per_pass: usize,
pub min_freed_pages: usize,
pub page_step_delay: Duration,
}
impl Default for VacuumConfig {
fn default() -> Self {
Self {
trigger: VacuumTrigger::default(),
online: true,
max_pages_per_pass: 0,
min_freed_pages: 1,
page_step_delay: Duration::from_millis(1),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PageSpaceInfo {
pub page_id: u64,
pub total_bytes: u64,
pub used_bytes: u64,
}
impl PageSpaceInfo {
pub fn new(page_id: u64, total_bytes: u64, used_bytes: u64) -> Self {
Self {
page_id,
total_bytes,
used_bytes,
}
}
pub fn dead_bytes(&self) -> u64 {
self.total_bytes.saturating_sub(self.used_bytes)
}
pub fn fragmentation(&self) -> f64 {
if self.total_bytes == 0 {
return 0.0;
}
self.dead_bytes() as f64 / self.total_bytes as f64
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum VacuumLockState {
Idle,
Running,
Bypassed,
}
pub struct VacuumLockGuard<'a> {
engine: &'a mut VacuumEngine,
}
impl<'a> Drop for VacuumLockGuard<'a> {
fn drop(&mut self) {
self.engine.lock_state = VacuumLockState::Idle;
}
}
#[derive(Debug, Clone)]
pub struct VacuumStats {
pub bytes_reclaimed: u64,
pub pages_compacted: usize,
pub slots_removed: usize,
pub duration: Duration,
pub pre_vacuum_dead_bytes: u64,
pub post_vacuum_dead_bytes: u64,
pub pre_vacuum_fragmentation: f64,
pub post_vacuum_fragmentation: f64,
}
impl VacuumStats {
pub fn made_progress(&self) -> bool {
self.bytes_reclaimed > 0 || self.pages_compacted > 0
}
}
#[derive(Debug, Default, Clone)]
pub struct VacuumMetrics {
pub total_runs: u64,
pub total_bytes_reclaimed: u64,
pub total_pages_compacted: u64,
pub total_duration: Duration,
pub run_history: Vec<VacuumStats>,
}
impl VacuumMetrics {
const MAX_HISTORY: usize = 32;
pub fn record(&mut self, stats: VacuumStats) {
self.total_runs += 1;
self.total_bytes_reclaimed += stats.bytes_reclaimed;
self.total_pages_compacted += stats.pages_compacted as u64;
self.total_duration += stats.duration;
self.run_history.push(stats);
if self.run_history.len() > Self::MAX_HISTORY {
self.run_history.remove(0);
}
}
}
#[derive(Debug)]
struct VacuumSchedule {
last_vacuum: Option<Instant>,
}
impl VacuumSchedule {
fn new() -> Self {
Self { last_vacuum: None }
}
fn should_run(&self, trigger: &VacuumTrigger, fragmentation: f64) -> bool {
match trigger {
VacuumTrigger::Manual => false,
VacuumTrigger::FragmentationThreshold(threshold) => fragmentation >= *threshold,
VacuumTrigger::TimeBased(interval) => self
.last_vacuum
.map(|t| t.elapsed() >= *interval)
.unwrap_or(true),
VacuumTrigger::Combined {
threshold,
interval,
} => {
let frag_trigger = fragmentation >= *threshold;
let time_trigger = self
.last_vacuum
.map(|t| t.elapsed() >= *interval)
.unwrap_or(true);
frag_trigger || time_trigger
}
}
}
fn record_run(&mut self) {
self.last_vacuum = Some(Instant::now());
}
}
#[derive(Debug, Default)]
struct FreeList {
freed_pages: BTreeSet<u64>,
freed_slots: BTreeSet<u64>,
}
impl FreeList {
fn add_page(&mut self, page_id: u64) {
self.freed_pages.insert(page_id);
}
fn add_slot(&mut self, slot_id: u64) {
self.freed_slots.insert(slot_id);
}
fn freed_page_count(&self) -> usize {
self.freed_pages.len()
}
fn freed_slot_count(&self) -> usize {
self.freed_slots.len()
}
fn drain_pages(&mut self, limit: usize) -> Vec<u64> {
if limit == 0 {
self.freed_pages
.iter()
.cloned()
.collect::<Vec<_>>()
.tap(|_| {
self.freed_pages.clear();
})
} else {
let pages: Vec<u64> = self.freed_pages.iter().cloned().take(limit).collect();
for p in &pages {
self.freed_pages.remove(p);
}
pages
}
}
fn drain_slots(&mut self) -> Vec<u64> {
let slots: Vec<u64> = self.freed_slots.iter().cloned().collect();
self.freed_slots.clear();
slots
}
}
trait Tap: Sized {
fn tap<F: FnOnce(&Self)>(self, f: F) -> Self {
f(&self);
self
}
}
impl<T> Tap for Vec<T> {}
#[derive(Debug, Clone)]
pub struct SpaceReport {
pub total_bytes: u64,
pub used_bytes: u64,
pub dead_bytes: u64,
pub fragmentation: f64,
pub freed_page_count: usize,
pub freed_slot_count: usize,
}
impl SpaceReport {
fn compute(pages: &HashMap<u64, PageSpaceInfo>, free_list: &FreeList) -> Self {
let total_bytes: u64 = pages.values().map(|p| p.total_bytes).sum();
let used_bytes: u64 = pages.values().map(|p| p.used_bytes).sum();
let dead_bytes = total_bytes.saturating_sub(used_bytes);
let fragmentation = if total_bytes == 0 {
0.0
} else {
dead_bytes as f64 / total_bytes as f64
};
Self {
total_bytes,
used_bytes,
dead_bytes,
fragmentation,
freed_page_count: free_list.freed_page_count(),
freed_slot_count: free_list.freed_slot_count(),
}
}
}
pub struct VacuumEngine {
config: VacuumConfig,
free_list: FreeList,
page_info: HashMap<u64, PageSpaceInfo>,
schedule: VacuumSchedule,
lock_state: VacuumLockState,
metrics: VacuumMetrics,
}
impl VacuumEngine {
pub fn new(config: VacuumConfig) -> Self {
Self {
config,
free_list: FreeList::default(),
page_info: HashMap::new(),
schedule: VacuumSchedule::new(),
lock_state: VacuumLockState::Idle,
metrics: VacuumMetrics::default(),
}
}
pub fn register_page(&mut self, info: PageSpaceInfo) {
self.page_info.insert(info.page_id, info);
}
pub fn mark_page_freed(&mut self, page_id: u64) {
self.free_list.add_page(page_id);
}
pub fn mark_slot_freed(&mut self, slot_id: u64) {
self.free_list.add_slot(slot_id);
}
pub fn space_report(&self) -> SpaceReport {
SpaceReport::compute(&self.page_info, &self.free_list)
}
pub fn fragmentation_ratio(&self) -> f64 {
self.space_report().fragmentation
}
pub fn lock_state(&self) -> &VacuumLockState {
&self.lock_state
}
pub fn bypass(&mut self) {
self.lock_state = VacuumLockState::Bypassed;
}
pub fn release_bypass(&mut self) {
if self.lock_state == VacuumLockState::Bypassed {
self.lock_state = VacuumLockState::Idle;
}
}
pub fn should_run_now(&self) -> bool {
if self.lock_state != VacuumLockState::Idle {
return false;
}
if self.free_list.freed_page_count() < self.config.min_freed_pages {
return false;
}
let frag = self.fragmentation_ratio();
self.schedule.should_run(&self.config.trigger, frag)
}
pub fn vacuum(&mut self) -> VacuumResult<VacuumStats> {
match self.lock_state {
VacuumLockState::Running => return Err(VacuumError::AlreadyRunning),
VacuumLockState::Bypassed => {
return Err(VacuumError::ConflictingLock(
"Vacuum is currently bypassed".to_string(),
))
}
VacuumLockState::Idle => {}
}
self.lock_state = VacuumLockState::Running;
let pre_report = self.space_report();
let start = Instant::now();
let limit = self.config.max_pages_per_pass;
let pages_to_compact = self.free_list.drain_pages(limit);
let pages_compacted = pages_to_compact.len();
let mut bytes_reclaimed: u64 = 0;
for page_id in pages_to_compact {
let reclaimed = self
.page_info
.get(&page_id)
.map(|p| p.total_bytes)
.unwrap_or(0);
bytes_reclaimed += reclaimed;
self.page_info.remove(&page_id);
if self.config.online {
let _ = self.config.page_step_delay;
}
}
let slots_removed = self.free_list.drain_slots().len();
let duration = start.elapsed();
let post_report = self.space_report();
let stats = VacuumStats {
bytes_reclaimed,
pages_compacted,
slots_removed,
duration,
pre_vacuum_dead_bytes: pre_report.dead_bytes,
post_vacuum_dead_bytes: post_report.dead_bytes,
pre_vacuum_fragmentation: pre_report.fragmentation,
post_vacuum_fragmentation: post_report.fragmentation,
};
self.metrics.record(stats.clone());
self.schedule.record_run();
self.lock_state = VacuumLockState::Idle;
Ok(stats)
}
pub fn maybe_vacuum(&mut self) -> VacuumResult<Option<VacuumStats>> {
if self.should_run_now() {
self.vacuum().map(Some)
} else {
Ok(None)
}
}
pub fn metrics(&self) -> &VacuumMetrics {
&self.metrics
}
pub fn pending_freed_pages(&self) -> usize {
self.free_list.freed_page_count()
}
pub fn pending_freed_slots(&self) -> usize {
self.free_list.freed_slot_count()
}
pub fn into_metrics(self) -> VacuumMetrics {
self.metrics
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
fn make_engine() -> VacuumEngine {
VacuumEngine::new(VacuumConfig::default())
}
#[test]
fn test_register_page_and_fragmentation() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(1, 4096, 2048));
let ratio = eng.fragmentation_ratio();
assert!((ratio - 0.5).abs() < 1e-9);
}
#[test]
fn test_zero_fragmentation_on_full_page() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(1, 4096, 4096));
assert_eq!(eng.fragmentation_ratio(), 0.0);
}
#[test]
fn test_full_fragmentation_on_empty_page() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(1, 4096, 0));
assert_eq!(eng.fragmentation_ratio(), 1.0);
}
#[test]
fn test_fragmentation_empty_engine() {
let eng = make_engine();
assert_eq!(eng.fragmentation_ratio(), 0.0);
}
#[test]
fn test_mark_page_freed_increments_pending() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(10, 4096, 1024));
eng.mark_page_freed(10);
assert_eq!(eng.pending_freed_pages(), 1);
}
#[test]
fn test_mark_slot_freed_increments_pending() {
let mut eng = make_engine();
eng.mark_slot_freed(100);
assert_eq!(eng.pending_freed_slots(), 1);
}
#[test]
fn test_mark_multiple_pages_freed() {
let mut eng = make_engine();
for i in 0..5u64 {
eng.register_page(PageSpaceInfo::new(i, 4096, 0));
eng.mark_page_freed(i);
}
assert_eq!(eng.pending_freed_pages(), 5);
}
#[test]
fn test_space_report_totals() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(0, 4096, 1000));
eng.register_page(PageSpaceInfo::new(1, 4096, 3000));
let report = eng.space_report();
assert_eq!(report.total_bytes, 8192);
assert_eq!(report.used_bytes, 4000);
assert_eq!(report.dead_bytes, 4192);
}
#[test]
fn test_space_report_freed_counts() {
let mut eng = make_engine();
eng.mark_page_freed(1);
eng.mark_slot_freed(200);
let report = eng.space_report();
assert_eq!(report.freed_page_count, 1);
assert_eq!(report.freed_slot_count, 1);
}
#[test]
fn test_vacuum_reclaims_freed_pages() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(5, 4096, 0));
eng.mark_page_freed(5);
assert_eq!(eng.pending_freed_pages(), 1);
let stats = eng.vacuum().expect("vacuum failed");
assert_eq!(stats.pages_compacted, 1);
assert_eq!(stats.bytes_reclaimed, 4096);
assert_eq!(eng.pending_freed_pages(), 0);
}
#[test]
fn test_vacuum_removes_slots() {
let mut eng = make_engine();
eng.mark_slot_freed(10);
eng.mark_slot_freed(20);
let stats = eng.vacuum().expect("vacuum failed");
assert_eq!(stats.slots_removed, 2);
assert_eq!(eng.pending_freed_slots(), 0);
}
#[test]
fn test_vacuum_pre_post_fragmentation() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(0, 4096, 0));
eng.mark_page_freed(0);
let stats = eng.vacuum().expect("vacuum failed");
assert!(stats.pre_vacuum_fragmentation >= 1.0 || stats.pre_vacuum_fragmentation >= 0.0);
assert_eq!(stats.post_vacuum_fragmentation, 0.0);
}
#[test]
fn test_vacuum_updates_metrics() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(1, 4096, 0));
eng.mark_page_freed(1);
eng.vacuum().expect("vacuum failed");
assert_eq!(eng.metrics().total_runs, 1);
assert_eq!(eng.metrics().total_pages_compacted, 1);
assert!(eng.metrics().total_bytes_reclaimed > 0);
}
#[test]
fn test_vacuum_history_stored() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(1, 4096, 0));
eng.mark_page_freed(1);
eng.vacuum().expect("vacuum failed");
assert_eq!(eng.metrics().run_history.len(), 1);
}
#[test]
fn test_max_pages_per_pass_limits_compaction() {
let config = VacuumConfig {
max_pages_per_pass: 2,
min_freed_pages: 1,
trigger: VacuumTrigger::Manual,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
for i in 0..5u64 {
eng.register_page(PageSpaceInfo::new(i, 4096, 0));
eng.mark_page_freed(i);
}
let stats = eng.vacuum().expect("vacuum failed");
assert_eq!(stats.pages_compacted, 2);
assert_eq!(eng.pending_freed_pages(), 3);
}
#[test]
fn test_already_running_error() {
let mut eng = make_engine();
eng.lock_state = VacuumLockState::Running;
let err = eng.vacuum().unwrap_err();
assert_eq!(err, VacuumError::AlreadyRunning);
}
#[test]
fn test_bypassed_error() {
let mut eng = make_engine();
eng.bypass();
let err = eng.vacuum().unwrap_err();
matches!(err, VacuumError::ConflictingLock(_));
}
#[test]
fn test_bypass_and_release() {
let mut eng = make_engine();
eng.bypass();
assert_eq!(*eng.lock_state(), VacuumLockState::Bypassed);
eng.release_bypass();
assert_eq!(*eng.lock_state(), VacuumLockState::Idle);
}
#[test]
fn test_should_run_manual_trigger_never_auto() {
let config = VacuumConfig {
trigger: VacuumTrigger::Manual,
min_freed_pages: 1,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.register_page(PageSpaceInfo::new(0, 4096, 0));
eng.mark_page_freed(0);
assert!(!eng.should_run_now());
}
#[test]
fn test_should_run_fragmentation_threshold_fires() {
let config = VacuumConfig {
trigger: VacuumTrigger::FragmentationThreshold(0.20),
min_freed_pages: 1,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.register_page(PageSpaceInfo::new(0, 100, 50)); eng.mark_page_freed(0);
assert!(eng.should_run_now());
}
#[test]
fn test_should_run_below_threshold_does_not_fire() {
let config = VacuumConfig {
trigger: VacuumTrigger::FragmentationThreshold(0.80),
min_freed_pages: 1,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.register_page(PageSpaceInfo::new(0, 100, 70)); eng.mark_page_freed(0);
assert!(!eng.should_run_now());
}
#[test]
fn test_should_run_not_enough_freed_pages() {
let config = VacuumConfig {
trigger: VacuumTrigger::FragmentationThreshold(0.10),
min_freed_pages: 5,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.register_page(PageSpaceInfo::new(0, 100, 0)); eng.mark_page_freed(0); assert!(!eng.should_run_now()); }
#[test]
fn test_should_run_bypassed_returns_false() {
let config = VacuumConfig {
trigger: VacuumTrigger::FragmentationThreshold(0.10),
min_freed_pages: 1,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.register_page(PageSpaceInfo::new(0, 100, 0));
eng.mark_page_freed(0);
eng.bypass();
assert!(!eng.should_run_now());
}
#[test]
fn test_time_based_trigger_fires_immediately_on_first_run() {
let config = VacuumConfig {
trigger: VacuumTrigger::TimeBased(Duration::from_secs(3600)),
min_freed_pages: 1,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.mark_page_freed(0);
assert!(eng.should_run_now());
}
#[test]
fn test_combined_trigger_fragmentation_path() {
let config = VacuumConfig {
trigger: VacuumTrigger::Combined {
threshold: 0.20,
interval: Duration::from_secs(3600),
},
min_freed_pages: 1,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.register_page(PageSpaceInfo::new(0, 100, 50)); eng.mark_page_freed(0);
assert!(eng.should_run_now());
}
#[test]
fn test_maybe_vacuum_does_not_run_when_not_needed() {
let config = VacuumConfig {
trigger: VacuumTrigger::Manual,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
let result = eng.maybe_vacuum().expect("unexpected error");
assert!(result.is_none());
}
#[test]
fn test_maybe_vacuum_runs_when_needed() {
let config = VacuumConfig {
trigger: VacuumTrigger::TimeBased(Duration::from_secs(0)),
min_freed_pages: 1,
..Default::default()
};
let mut eng = VacuumEngine::new(config);
eng.register_page(PageSpaceInfo::new(0, 4096, 0));
eng.mark_page_freed(0);
let result = eng.maybe_vacuum().expect("unexpected error");
assert!(result.is_some());
}
#[test]
fn test_page_space_info_dead_bytes() {
let info = PageSpaceInfo::new(1, 4096, 1024);
assert_eq!(info.dead_bytes(), 3072);
}
#[test]
fn test_page_space_info_fragmentation() {
let info = PageSpaceInfo::new(1, 1000, 250);
assert!((info.fragmentation() - 0.75).abs() < 1e-9);
}
#[test]
fn test_page_space_info_zero_total() {
let info = PageSpaceInfo::new(1, 0, 0);
assert_eq!(info.fragmentation(), 0.0);
}
#[test]
fn test_vacuum_error_display() {
assert_eq!(
VacuumError::AlreadyRunning.to_string(),
"Vacuum already running"
);
assert!(VacuumError::ConflictingLock("x".to_string())
.to_string()
.contains("x"));
assert!(VacuumError::InvalidId("bad".to_string())
.to_string()
.contains("bad"));
assert!(VacuumError::Internal("oops".to_string())
.to_string()
.contains("oops"));
}
#[test]
fn test_into_metrics_returns_cumulative() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(0, 4096, 0));
eng.mark_page_freed(0);
eng.vacuum().expect("vacuum failed");
let metrics = eng.into_metrics();
assert_eq!(metrics.total_runs, 1);
}
#[test]
fn test_lock_released_after_vacuum() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(0, 4096, 0));
eng.mark_page_freed(0);
eng.vacuum().expect("vacuum failed");
assert_eq!(*eng.lock_state(), VacuumLockState::Idle);
}
#[test]
fn test_second_vacuum_runs_after_first() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(0, 4096, 0));
eng.mark_page_freed(0);
eng.vacuum().expect("first vacuum failed");
eng.register_page(PageSpaceInfo::new(1, 4096, 0));
eng.mark_page_freed(1);
let stats = eng.vacuum().expect("second vacuum failed");
assert_eq!(stats.pages_compacted, 1);
assert_eq!(eng.metrics().total_runs, 2);
}
#[test]
fn test_vacuum_with_no_freed_pages() {
let mut eng = make_engine();
eng.register_page(PageSpaceInfo::new(0, 4096, 2048));
let stats = eng.vacuum().expect("vacuum failed");
assert_eq!(stats.pages_compacted, 0);
assert_eq!(stats.bytes_reclaimed, 0);
assert_eq!(stats.slots_removed, 0);
}
}