use anyhow::{Context, Result};
use std::path::PathBuf;
use std::sync::Arc;
use super::config::{COMPACT_MIN_RECLAIM_PERCENT, DreamConfig};
use crate::memory_core::retrieval::PalaceHandle;
use crate::memory_core::store::kg_redb::copy_swap::{
self, CompactFaultHook, CompactPlan, PreparedCompaction,
};
use crate::memory_core::store::kg_redb::stats::{KgRedbStats, history_cutoff_ms};
use crate::memory_core::timeouts;
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct KgCompactReport {
pub stats: KgRedbStats,
pub dry_run: bool,
pub skipped: Option<String>,
pub bytes_before: u64,
pub bytes_after: u64,
pub rows_copied: u64,
pub history_rows_pruned: u64,
pub backup: Option<PathBuf>,
}
impl KgCompactReport {
pub fn ran(&self) -> bool {
self.skipped.is_none() && !self.dry_run
}
pub fn bytes_reclaimed(&self) -> u64 {
self.bytes_before.saturating_sub(self.bytes_after)
}
pub fn summary(&self) -> String {
match &self.skipped {
Some(reason) => format!("skipped: {reason}"),
None if self.dry_run => format!(
"dry-run: would prune {} stale history row(s) and reclaim ~{} bytes of {}",
self.stats.triples_history_stale, self.stats.reclaimable_bytes, self.bytes_before
),
None => format!(
"compacted {} -> {} bytes ({} reclaimed), pruned {} history row(s)",
self.bytes_before,
self.bytes_after,
self.bytes_reclaimed(),
self.history_rows_pruned
),
}
}
}
pub async fn kg_compact_pass(
handle: &Arc<PalaceHandle>,
config: &DreamConfig,
dry_run: bool,
) -> Result<KgCompactReport> {
kg_compact_pass_with_hook(handle, config, dry_run, None).await
}
pub async fn kg_compact_pass_with_hook(
handle: &Arc<PalaceHandle>,
config: &DreamConfig,
dry_run: bool,
hook: Option<CompactFaultHook>,
) -> Result<KgCompactReport> {
let Some(path) = kg_redb_path(handle) else {
return Ok(skipped(
empty_stats(),
dry_run,
"palace has no on-disk data directory",
));
};
let days = config.effective_prune_history_days();
let measure_path = path.clone();
let stats = tokio::task::spawn_blocking(move || KgRedbStats::measure(&measure_path, days))
.await
.context("join kg.redb measurement")??;
if let Some(reason) = gate(&stats, config) {
return Ok(skipped(stats, dry_run, &reason));
}
if dry_run {
let bytes_before = stats.file_bytes;
return Ok(KgCompactReport {
stats,
dry_run: true,
skipped: None,
bytes_before,
bytes_after: bytes_before,
rows_copied: 0,
history_rows_pruned: 0,
backup: None,
});
}
let plan = CompactPlan {
history_cutoff_ms: Some(history_cutoff_ms(
chrono::Utc::now().timestamp_millis(),
days,
)),
keep_backup: config.compact_keep_backup,
};
let store = handle.kg.redb_store().clone();
let prepare_hook = hook.clone();
let prepared: PreparedCompaction = tokio::task::spawn_blocking(move || {
copy_swap::prepare(&store, plan, prepare_hook.as_ref())
})
.await
.context("join kg.redb rewrite")??;
let outcome = {
let _write_guard = timeouts::lock_with_timeout(
&handle.write_mutex,
timeouts::write_lock_timeout(),
handle.id.as_str(),
)
.await?;
let store = handle.kg.redb_store().clone();
let commit_hook = hook.clone();
tokio::task::spawn_blocking(move || prepared.commit(&store, commit_hook.as_ref()))
.await
.context("join kg.redb swap")??
};
Ok(KgCompactReport {
stats,
dry_run: false,
skipped: None,
bytes_before: outcome.bytes_before,
bytes_after: outcome.bytes_after,
rows_copied: outcome.rows_copied,
history_rows_pruned: outcome.history_rows_pruned,
backup: outcome.backup,
})
}
fn gate(stats: &KgRedbStats, config: &DreamConfig) -> Option<String> {
if !config.compact {
return Some("dream.compact is disabled".to_string());
}
if stats.file_bytes < config.compact_min_bytes {
return Some(format!(
"kg.redb is {} bytes, under the {}-byte dream.compact_min_bytes floor",
stats.file_bytes, config.compact_min_bytes
));
}
let threshold = stats.file_bytes.saturating_mul(COMPACT_MIN_RECLAIM_PERCENT) / 100;
if stats.reclaimable_bytes < threshold {
return Some(format!(
"only ~{} of {} bytes are reclaimable, under the {COMPACT_MIN_RECLAIM_PERCENT}% \
floor",
stats.reclaimable_bytes, stats.file_bytes
));
}
None
}
fn kg_redb_path(handle: &Arc<PalaceHandle>) -> Option<PathBuf> {
handle.data_dir.as_ref().map(|d| d.join("kg.redb"))
}
fn skipped(stats: KgRedbStats, dry_run: bool, reason: &str) -> KgCompactReport {
let bytes = stats.file_bytes;
KgCompactReport {
stats,
dry_run,
skipped: Some(reason.to_string()),
bytes_before: bytes,
bytes_after: bytes,
rows_copied: 0,
history_rows_pruned: 0,
backup: None,
}
}
fn empty_stats() -> KgRedbStats {
KgRedbStats {
path: PathBuf::new(),
file_bytes: 0,
from_snapshot: false,
tables: Vec::new(),
triples_active: 0,
triples_closed_in_place: 0,
triples_history: 0,
triples_history_stale: 0,
triples_history_stale_bytes: 0,
history_cutoff_days: 0,
superseded_drawers: 0,
dead_predicate_index: None,
reclaimable_bytes: 0,
}
}