use wbase::time::now_ms;
use whasher::{HashMap, new_hash_map};
use super::{
CompactRunResult, CopyOutcome, LogCompactor, cursor::ScanCursor, scope::MetaDeathScope,
};
use crate::{error::Result, host::CompactStore};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct CandidateRecord {
addr: u64,
is_dead: bool,
}
impl<S: CompactStore> LogCompactor<S> {
pub(super) async fn compact_lookup<F>(
&self,
session: &S::Session,
begin_addr: u64,
until_address: u64,
read_only_addr: u64,
mut is_deleted: F,
) -> Result<CompactRunResult>
where
F: FnMut(&[u8], &[u8]) -> bool,
{
let mut cursor = ScanCursor::new(self.store.hlog(), begin_addr, until_address);
let mut scanned_records: usize = 0;
let mut live_copied: usize = 0;
let mut superseded: usize = 0;
let mut dead_dropped: usize = 0;
let mut retained: usize = 0;
let mut scope = MetaDeathScope::new();
let mut actual_until = begin_addr;
let mut retain_floor: Option<u64> = None;
let now = now_ms();
while let Some((curr_addr, rec)) = cursor.pull().await? {
actual_until = curr_addr + rec.physical_size() as u64;
if actual_until > read_only_addr {
actual_until = curr_addr;
break;
}
scanned_records += 1;
let is_tombstone = rec.is_tombstone();
let key = rec.key;
let val = if is_tombstone { &[] } else { rec.value };
scope.observe(&*self.store, is_tombstone, key, val, curr_addr);
if self
.judge_dead(session, is_tombstone, key, val, now, &mut is_deleted)
.await
{
dead_dropped += 1;
self.store.index().delete(key, curr_addr);
} else {
let latest = self
.find_latest_address(session, key, Some((curr_addr, false)))
.await?;
if let Some(latest) = latest
&& latest.main_addr == curr_addr
&& !latest.is_tombstone
{
match self
.conditional_copy_to_tail(session, key, val, curr_addr, latest.index_addr)
.await?
{
CopyOutcome::Copied => live_copied += 1,
CopyOutcome::Superseded => superseded += 1,
CopyOutcome::Retain => {
retained += 1;
retain_floor = Some(match retain_floor {
Some(floor) => floor.min(curr_addr),
None => curr_addr,
});
}
}
} else {
superseded += 1;
}
}
}
let cursor_addr = cursor.cursor();
if cursor_addr <= read_only_addr {
actual_until = cursor_addr;
}
if let Some(floor) = retain_floor {
actual_until = actual_until.min(floor);
}
Ok(CompactRunResult {
scanned_records,
live_copied,
superseded,
dead_dropped,
retained,
actual_until,
scope,
})
}
pub(super) async fn compact_scan<F>(
&self,
session: &S::Session,
begin_addr: u64,
until_address: u64,
read_only_addr: u64,
mut is_deleted: F,
) -> Result<CompactRunResult>
where
F: FnMut(&[u8], &[u8]) -> bool,
{
let mut cursor = ScanCursor::new(self.store.hlog(), begin_addr, until_address);
let mut scanned_records: usize = 0;
let mut live_copied: usize = 0;
let mut superseded: usize = 0;
let mut dead_dropped: usize = 0;
let mut retained: usize = 0;
let mut candidates: HashMap<Box<[u8]>, CandidateRecord> = new_hash_map();
let mut scope = MetaDeathScope::new();
let mut actual_until = begin_addr;
let mut retain_floor: Option<u64> = None;
let now = now_ms();
while let Some((curr_addr, rec)) = cursor.pull().await? {
actual_until = curr_addr + rec.physical_size() as u64;
if actual_until > read_only_addr {
actual_until = curr_addr;
break;
}
scanned_records += 1;
let is_tombstone = rec.is_tombstone();
let key = rec.key;
let raw_val = if is_tombstone { &[] } else { rec.value };
scope.observe(&*self.store, is_tombstone, key, raw_val, curr_addr);
let is_dead = self
.judge_dead(session, is_tombstone, key, raw_val, now, &mut is_deleted)
.await;
match candidates.get_mut(key) {
Some(cand) => {
superseded += 1;
cand.addr = curr_addr;
cand.is_dead = is_dead;
}
None => {
candidates.insert(
Box::from(key),
CandidateRecord {
addr: curr_addr,
is_dead,
},
);
}
}
}
let cursor_addr = cursor.cursor();
if cursor_addr <= read_only_addr {
actual_until = cursor_addr;
}
cursor = ScanCursor::new(self.store.hlog(), actual_until, read_only_addr);
while !candidates.is_empty()
&& let Some((curr_addr, rec)) = cursor.pull().await?
{
if curr_addr + rec.physical_size() as u64 > read_only_addr {
break;
}
if candidates.remove(rec.key).is_some() {
superseded += 1;
}
}
let now = now_ms();
for (key, cand) in candidates {
if cand.is_dead || self.is_stale_subkey(&key) {
self.store.index().delete(&key, cand.addr);
dead_dropped += 1;
continue;
}
let latest = self
.find_latest_address(session, &key, Some((cand.addr, false)))
.await?;
let Some(latest) = latest.filter(|l| l.main_addr == cand.addr && !l.is_tombstone) else {
superseded += 1;
continue;
};
let record = self.store.hlog().read_record(cand.addr).await?;
let val = record.value()?;
if self
.judge_dead(session, false, &key, val, now, &mut is_deleted)
.await
{
self.store.index().delete(&key, cand.addr);
dead_dropped += 1;
continue;
}
match self
.conditional_copy_to_tail(session, &key, val, cand.addr, latest.index_addr)
.await?
{
CopyOutcome::Copied => live_copied += 1,
CopyOutcome::Superseded => superseded += 1,
CopyOutcome::Retain => {
retained += 1;
retain_floor = Some(match retain_floor {
Some(floor) => floor.min(cand.addr),
None => cand.addr,
});
}
}
}
if let Some(floor) = retain_floor {
actual_until = actual_until.min(floor);
}
Ok(CompactRunResult {
scanned_records,
live_copied,
superseded,
dead_dropped,
retained,
actual_until,
scope,
})
}
}