use anyhow::{Context, Result};
use clap::Subcommand;
use trusty_common::memory_core::dream::{kg_compact_pass, DreamConfig};
use trusty_common::memory_core::palace::Palace;
use trusty_common::memory_core::retrieval::PalaceHandle;
use trusty_common::memory_core::store::kg_redb::KgRedbStats;
use trusty_common::memory_core::store::OpenIntent;
use trusty_common::memory_core::PalaceRegistry;
#[derive(Debug, Subcommand)]
pub enum PalaceAction {
Stats {
name: String,
#[arg(long, value_name = "DAYS", default_value_t = 90)]
history_days: i64,
#[arg(long)]
json: bool,
},
Compact {
name: String,
#[arg(long)]
dry_run: bool,
#[arg(long, value_name = "DAYS", default_value_t = 90)]
history_days: i64,
},
}
pub async fn dispatch(action: PalaceAction) -> Result<()> {
match action {
PalaceAction::Stats {
name,
history_days,
json,
} => handle_palace_stats(name, history_days, json).await,
PalaceAction::Compact {
name,
dry_run,
history_days,
} => handle_palace_compact(name, dry_run, history_days).await,
}
}
pub async fn handle_palace_stats(name: String, history_days: i64, json: bool) -> Result<()> {
let palace = resolve(&name)?;
print!("{}", stats_report(&name, &palace, history_days, json)?);
Ok(())
}
pub(crate) fn stats_report(
name: &str,
palace: &Palace,
history_days: i64,
json: bool,
) -> Result<String> {
let path = palace.data_dir.join("kg.redb");
let stats = KgRedbStats::measure(&path, history_days)
.with_context(|| format!("measure {}", path.display()))?;
if json {
Ok(format!("{}\n", render_json(name, &stats)?))
} else {
Ok(render_text(name, &stats))
}
}
pub async fn handle_palace_compact(name: String, dry_run: bool, history_days: i64) -> Result<()> {
let palace = resolve(&name)?;
print!(
"{}",
compact_report(&name, &palace, dry_run, history_days).await?
);
Ok(())
}
pub(crate) async fn compact_report(
name: &str,
palace: &Palace,
dry_run: bool,
history_days: i64,
) -> Result<String> {
let intent = if dry_run {
OpenIntent::ReadOnlyClient
} else {
OpenIntent::Writer
};
let handle = std::sync::Arc::new(
PalaceHandle::open_with_intent(palace, intent)
.with_context(|| format!("open palace {}", palace.id))?,
);
let cfg = DreamConfig {
prune_history_after_days: history_days,
compact_min_bytes: 0,
..DreamConfig::default()
};
let report = kg_compact_pass(&handle, &cfg, dry_run).await?;
let mut out = format!("palace={name} {}\n", report.summary());
if let Some(backup) = &report.backup {
out.push_str(&format!(" backup: {}\n", backup.display()));
}
out.push_str(&render_text(name, &report.stats));
if dry_run {
out.push_str("nothing was written — re-run without --dry-run to compact\n");
}
Ok(out)
}
fn resolve(name: &str) -> Result<Palace> {
let data_dir = trusty_common::resolve_data_dir("trusty-memory")
.context("resolve trusty-memory data dir")?;
let root = crate::resolve_palace_registry_dir(data_dir);
PalaceRegistry::list_palaces(&root)
.unwrap_or_default()
.into_iter()
.find(|p| p.id.0 == name)
.with_context(|| format!("no palace named '{name}' under {}", root.display()))
}
fn render_text(name: &str, s: &KgRedbStats) -> String {
let mut out = String::new();
out.push_str(&format!("palace={name} kg.redb={}\n", s.path.display()));
if s.from_snapshot {
out.push_str(
" note: a writer holds the live file; these numbers come from a snapshot taken \
just now\n",
);
}
out.push_str(&format!(
" file_bytes {}\n reclaimable_estimate {} ({}%)\n",
s.file_bytes,
s.reclaimable_bytes,
percent(s.reclaimable_bytes, s.file_bytes)
));
out.push_str(&format!(
" triples active={} history={} stale(>{}d)={} stale_bytes={}\n",
s.triples_active,
s.triples_history,
s.history_cutoff_days,
s.triples_history_stale,
s.triples_history_stale_bytes
));
out.push_str(&format!(
" closed-in-place={} superseded_drawers={}\n",
s.triples_closed_in_place, s.superseded_drawers
));
if let Some(dead) = &s.dead_predicate_index {
out.push_str(&format!(
" dead index triples_by_predicate: {} row(s), {} live bytes — reclaimed by \
the next compaction (#6652)\n",
dead.rows,
dead.live_bytes()
));
}
out.push_str(&format!(
" {:<24} {:>10} {:>12} {:>12} {:>12}\n",
"table", "rows", "stored", "metadata", "fragmented"
));
for t in &s.tables {
out.push_str(&format!(
" {:<24} {:>10} {:>12} {:>12} {:>12}\n",
t.name, t.rows, t.stored_bytes, t.metadata_bytes, t.fragmented_bytes
));
}
out
}
fn render_json(name: &str, s: &KgRedbStats) -> Result<String> {
let tables: Vec<serde_json::Value> = s
.tables
.iter()
.map(|t| {
serde_json::json!({
"name": t.name,
"rows": t.rows,
"stored_bytes": t.stored_bytes,
"metadata_bytes": t.metadata_bytes,
"fragmented_bytes": t.fragmented_bytes,
"pages": t.pages,
})
})
.collect();
let v = serde_json::json!({
"palace": name,
"path": s.path,
"from_snapshot": s.from_snapshot,
"file_bytes": s.file_bytes,
"reclaimable_bytes": s.reclaimable_bytes,
"triples_active": s.triples_active,
"triples_closed_in_place": s.triples_closed_in_place,
"triples_history": s.triples_history,
"triples_history_stale": s.triples_history_stale,
"triples_history_stale_bytes": s.triples_history_stale_bytes,
"history_cutoff_days": s.history_cutoff_days,
"superseded_drawers": s.superseded_drawers,
"tables": tables,
});
serde_json::to_string_pretty(&v).context("serialize palace stats")
}
fn percent(part: u64, whole: u64) -> u64 {
part.saturating_mul(100).checked_div(whole).unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
use trusty_common::memory_core::palace::PalaceId;
fn fixture(name: &str, n: usize) -> (tempfile::TempDir, Palace) {
let dir = tempfile::tempdir().expect("tempdir");
let data_dir = dir.path().join(name);
std::fs::create_dir_all(&data_dir).expect("mkdir");
let kg = trusty_common::memory_core::store::kg_redb::KgStoreRedb::open(
&data_dir.join("kg.redb"),
)
.expect("open kg");
for i in 0..n {
kg.assert(&trusty_common::memory_core::store::Triple {
subject: format!("s{i}"),
predicate: "knows".into(),
object: format!("o{i}"),
valid_from: chrono::Utc::now(),
valid_to: None,
confidence: 1.0,
provenance: None,
})
.expect("assert");
}
drop(kg);
let palace = Palace {
id: PalaceId::new(name),
name: name.into(),
description: None,
created_at: chrono::Utc::now(),
data_dir,
};
(dir, palace)
}
#[test]
fn palace_stats_reports_a_hand_built_palace() {
let (_d, palace) = fixture("stats-fixture", 5);
let text = stats_report("stats-fixture", &palace, 90, false).expect("report");
assert!(text.contains("palace=stats-fixture"), "{text}");
assert!(text.contains("triples active=5"), "{text}");
assert!(text.contains("history=0"), "{text}");
assert!(text.contains("file_bytes"), "{text}");
assert!(text.contains("triples_by_object"), "{text}");
let json = stats_report("stats-fixture", &palace, 90, true).expect("json");
let parsed: serde_json::Value = serde_json::from_str(&json).expect("parse");
assert_eq!(parsed["triples_active"], 5);
assert_eq!(parsed["palace"], "stats-fixture");
}
#[tokio::test]
async fn palace_compact_dry_run_writes_nothing() {
let (_d, palace) = fixture("dry-run-fixture", 4);
let kg_path = palace.data_dir.join("kg.redb");
let rows = |p: &std::path::Path| {
KgRedbStats::measure(p, 90)
.expect("measure")
.tables
.iter()
.map(|t| (t.name.clone(), t.rows))
.collect::<Vec<_>>()
};
let before = rows(&kg_path);
let out = compact_report("dry-run-fixture", &palace, true, 90)
.await
.expect("dry run");
assert!(out.contains("dry-run:"), "{out}");
assert!(out.contains("nothing was written"), "{out}");
assert_eq!(before, rows(&kg_path), "the dry run changed kg.redb");
assert!(!palace.data_dir.join("kg.redb.pre-compact.bak").exists());
assert!(!palace.data_dir.join("kg.redb.compacting").exists());
}
}