mod copy;
mod cursor;
mod judge;
mod probe;
mod run;
mod scope;
use std::sync::Arc;
use log::info;
use crate::{
error::{Error, Result},
host::CompactStore,
};
const CAS_COPY_RETRIES: usize = 8;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum CopyOutcome {
Copied,
Superseded,
Retain,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CompactionType {
Lookup,
Scan,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct CompactionStats {
pub scanned_records: usize,
pub live_copied: usize,
pub superseded: usize,
pub dead_dropped: usize,
pub retained: usize,
pub bytes_freed: u64,
pub new_begin_address: u64,
}
impl CompactionStats {
#[inline]
pub const fn is_empty(&self) -> bool {
self.scanned_records == 0
}
}
struct CompactRunResult {
scanned_records: usize,
live_copied: usize,
superseded: usize,
dead_dropped: usize,
retained: usize,
actual_until: u64,
scope: scope::MetaDeathScope,
}
pub struct LogCompactor<S: CompactStore> {
store: Arc<S>,
cas_retries: usize,
}
impl<S: CompactStore> Clone for LogCompactor<S> {
fn clone(&self) -> Self {
Self {
store: Arc::clone(&self.store),
cas_retries: self.cas_retries,
}
}
}
impl<S: CompactStore> LogCompactor<S> {
pub fn new(store: Arc<S>) -> Self {
Self {
store,
cas_retries: CAS_COPY_RETRIES,
}
}
#[inline]
pub fn with_cas_retries(store: Arc<S>, cas_retries: usize) -> Self {
Self { store, cas_retries }
}
pub async fn compact_with_filter<F>(
&self,
until_address: u64,
comp_type: CompactionType,
is_deleted: F,
) -> Result<CompactionStats>
where
F: FnMut(&[u8], &[u8]) -> bool,
{
let read_only_addr = self.store.read_only_address();
if until_address > read_only_addr {
return Err(Error::UntilAddressOutOfRange {
until_address,
read_only_address: read_only_addr,
});
}
let old_begin = self.store.begin_address();
if until_address <= old_begin {
return Ok(CompactionStats {
new_begin_address: old_begin,
..CompactionStats::default()
});
}
let session = self.store.new_session()?;
let run_res = match comp_type {
CompactionType::Lookup => {
self
.compact_lookup(
&session,
old_begin,
until_address,
read_only_addr,
is_deleted,
)
.await?
}
CompactionType::Scan => {
self
.compact_scan(
&session,
old_begin,
until_address,
read_only_addr,
is_deleted,
)
.await?
}
};
self.store.shift_begin_address(run_res.actual_until).await?;
run_res
.scope
.collect_dead(&*self.store, run_res.actual_until);
let new_begin_address = self.store.begin_address();
let bytes_freed = run_res.actual_until.saturating_sub(old_begin);
info!(
"紧缩完成: 类型={comp_type:?}, 扫描={}, 迁移={}, 弃迁={}, 判死={}, 保留={}, 释放={bytes_freed}B, 新起始地址={new_begin_address:#x}",
run_res.scanned_records,
run_res.live_copied,
run_res.superseded,
run_res.dead_dropped,
run_res.retained
);
Ok(CompactionStats {
scanned_records: run_res.scanned_records,
live_copied: run_res.live_copied,
superseded: run_res.superseded,
dead_dropped: run_res.dead_dropped,
retained: run_res.retained,
bytes_freed,
new_begin_address,
})
}
pub async fn compact_lazy(&self, max_seek_bytes: u64) -> Result<CompactionStats> {
let begin_addr = self.store.begin_address();
let read_only_addr = self.store.read_only_address();
if begin_addr >= read_only_addr || max_seek_bytes == 0 {
return Ok(CompactionStats {
new_begin_address: begin_addr,
..CompactionStats::default()
});
}
let until_address = read_only_addr.min(begin_addr.saturating_add(max_seek_bytes));
self.compact(until_address, CompactionType::Lookup).await
}
#[inline]
pub async fn compact(
&self,
until_address: u64,
comp_type: CompactionType,
) -> Result<CompactionStats> {
self
.compact_with_filter(until_address, comp_type, |_, _| false)
.await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn compaction_stats_default_and_is_empty() {
let stats = CompactionStats::default();
assert!(stats.is_empty());
assert_eq!(stats.scanned_records, 0);
assert_eq!(stats.live_copied, 0);
assert_eq!(stats.dead_dropped, 0);
assert_eq!(stats.bytes_freed, 0);
}
}