use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use anyhow::{bail, Context, Result};
use chrono::{DateTime, NaiveDateTime, Utc};
use rusqlite::{Connection, OpenFlags};
use trusty_common::memory_core::palace::{Drawer, Palace};
use trusty_common::memory_core::retrieval::{PalaceHandle, VectorBackfillOptions};
use trusty_common::memory_core::store::{OpenIntent, INCOMPATIBLE_SUFFIX};
use trusty_common::memory_core::{memory_content_hash, ContentHash};
use uuid::Uuid;
use super::store_snapshot::{with_store_copy, SCRATCH_PREFIX};
#[path = "legacy_kg_guard.rs"]
pub mod guard;
use guard::{backup_stores, probe_stores, screen_drawers, Backup, CopyFn, Rejected};
pub(crate) const LEGACY_KG_FILE: &str = "kg.db";
const COPIED_SUFFIXES: [&str; 3] = ["", "-wal", "-journal"];
const INDEX_FILE: &str = "index.usearch.redb";
const SQLITE_MAGIC: &[u8; 16] = b"SQLite format 3\0";
#[derive(Debug, Default)]
pub struct LegacyDrawers {
pub total_rows: usize,
pub drawers: Vec<Drawer>,
pub unreadable: Vec<String>,
pub triple_rows: usize,
}
pub fn read_legacy_kg(data_dir: &Path) -> Result<Option<LegacyDrawers>> {
let path = data_dir.join(LEGACY_KG_FILE);
let present = path
.try_exists()
.with_context(|| format!("cannot stat {}", path.display()))?;
if !present {
return Ok(None);
}
let mut header = [0u8; 16];
let n = std::io::Read::read(
&mut std::fs::File::open(&path).with_context(|| format!("open {}", path.display()))?,
&mut header,
)
.with_context(|| format!("read header of {}", path.display()))?;
if n < header.len() || &header != SQLITE_MAGIC {
bail!("{} is not a SQLite database", path.display());
}
let scratch = tempfile::TempDir::with_prefix_in(SCRATCH_PREFIX, std::env::temp_dir())
.context("create scratch dir for the legacy kg.db copy")?;
let copy = scratch.path().join(LEGACY_KG_FILE);
copy_stable(data_dir, scratch.path(), |from, to| std::fs::copy(from, to))?;
let conn = Connection::open_with_flags(
©,
OpenFlags::SQLITE_OPEN_READ_WRITE | OpenFlags::SQLITE_OPEN_NO_MUTEX,
)
.with_context(|| format!("open a copy of {}", path.display()))?;
let mut out = LegacyDrawers::default();
if has_table(&conn, "triples")? {
out.triple_rows = conn
.query_row("SELECT COUNT(*) FROM triples", [], |r| r.get::<_, i64>(0))
.context("count legacy triples")?
.try_into()
.unwrap_or(0);
}
if !has_table(&conn, "drawers")? {
return Ok(Some(out));
}
let mut stmt = conn
.prepare(
"SELECT id, room_id, content, importance, tags, source_file, created_at \
FROM drawers",
)
.context("prepare legacy drawer scan")?;
let rows = stmt
.query_map([], |r| {
Ok(LegacyRow {
id: r.get(0)?,
room_id: r.get(1)?,
content: r.get(2)?,
importance: r.get(3)?,
tags: r.get(4)?,
source_file: r.get(5)?,
created_at: r.get(6)?,
})
})
.context("scan legacy drawers")?;
for row in rows {
out.total_rows += 1;
let row = row.context("read legacy drawer row")?;
let id = row.id.clone().unwrap_or_default();
match row.into_drawer() {
Ok(d) => out.drawers.push(d),
Err(reason) => out.unreadable.push(format!("{id}: {reason}")),
}
}
Ok(Some(out))
}
type Fingerprint = Vec<Option<(u64, Option<std::time::SystemTime>)>>;
fn fingerprint(data_dir: &Path) -> Result<Fingerprint> {
COPIED_SUFFIXES
.iter()
.map(|suffix| {
let p = data_dir.join(format!("{LEGACY_KG_FILE}{suffix}"));
match std::fs::metadata(&p) {
Ok(m) => Ok(Some((m.len(), m.modified().ok()))),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(e) => Err(e).with_context(|| format!("cannot stat {}", p.display())),
}
})
.collect()
}
pub(crate) fn copy_stable(
data_dir: &Path,
dest: &Path,
mut copy: impl FnMut(&Path, &Path) -> std::io::Result<u64>,
) -> Result<()> {
for _ in 0..2 {
let before = fingerprint(data_dir)?;
for (suffix, stat) in COPIED_SUFFIXES.iter().zip(&before) {
let name = format!("{LEGACY_KG_FILE}{suffix}");
let (src, dst) = (data_dir.join(&name), dest.join(&name));
match std::fs::remove_file(&dst) {
Err(e) if e.kind() != std::io::ErrorKind::NotFound => {
return Err(e).with_context(|| format!("clear {}", dst.display()));
}
_ => {}
}
if stat.is_some() {
copy(&src, &dst).with_context(|| format!("copy {}", src.display()))?;
}
}
if fingerprint(data_dir)? == before {
return Ok(());
}
}
bail!(
"{} changed during both copy attempts; refusing to read a possibly torn copy",
data_dir.join(LEGACY_KG_FILE).display()
)
}
fn has_table(conn: &Connection, name: &str) -> Result<bool> {
let n: i64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = ?1",
[name],
|r| r.get(0),
)
.with_context(|| format!("look up legacy table {name}"))?;
Ok(n > 0)
}
struct LegacyRow {
id: Option<String>,
room_id: Option<String>,
content: Option<String>,
importance: Option<f64>,
tags: Option<String>,
source_file: Option<String>,
created_at: Option<String>,
}
impl LegacyRow {
fn into_drawer(self) -> std::result::Result<Drawer, String> {
let id = parse_uuid(self.id.as_deref(), "id")?;
let room_id = parse_uuid(self.room_id.as_deref(), "room_id")?;
let content = self.content.ok_or("content is NULL")?;
let created_at = parse_timestamp(self.created_at.as_deref())?;
let mut d = Drawer::new(room_id, content);
d.id = id;
d.created_at = created_at;
d.importance = self.importance.unwrap_or(0.5) as f32;
d.source_file = self.source_file.map(PathBuf::from);
d.tags = self
.tags
.and_then(|t| serde_json::from_str(&t).ok())
.unwrap_or_default();
Ok(d)
}
}
fn parse_uuid(raw: Option<&str>, column: &str) -> std::result::Result<Uuid, String> {
let raw = raw.ok_or_else(|| format!("{column} is NULL"))?;
Uuid::parse_str(raw).map_err(|e| format!("invalid {column}: {e}"))
}
fn parse_timestamp(raw: Option<&str>) -> std::result::Result<DateTime<Utc>, String> {
let raw = raw.ok_or("created_at is NULL")?;
if let Ok(dt) = DateTime::parse_from_rfc3339(raw) {
return Ok(dt.with_timezone(&Utc));
}
NaiveDateTime::parse_from_str(raw, "%Y-%m-%d %H:%M:%S%.f")
.map(|n| n.and_utc())
.map_err(|e| format!("invalid created_at {raw:?}: {e}"))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IncompatibleFile {
pub path: PathBuf,
pub bytes: u64,
}
pub fn list_incompatible_files(data_dir: &Path) -> Result<Vec<IncompatibleFile>> {
let mut out = Vec::new();
let entries = match std::fs::read_dir(data_dir) {
Ok(e) => e,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(out),
Err(e) => return Err(e).with_context(|| format!("list {}", data_dir.display())),
};
for entry in entries {
let entry = entry.with_context(|| format!("list {}", data_dir.display()))?;
if entry
.file_name()
.to_string_lossy()
.contains(INCOMPATIBLE_SUFFIX)
{
let bytes = entry.metadata().map(|m| m.len()).unwrap_or(0);
out.push(IncompatibleFile {
path: entry.path(),
bytes,
});
}
}
out.sort_by(|a, b| a.path.cmp(&b.path));
Ok(out)
}
#[derive(Debug, Default)]
pub struct LegacyReport {
pub palace: String,
pub dry_run: bool,
pub legacy_present: bool,
pub legacy_rows: usize,
pub unreadable: Vec<String>,
pub legacy_triples: usize,
pub already_live: usize,
pub missing: usize,
pub content_duplicates: Option<usize>,
pub rejected: Vec<Rejected>,
pub backup: Option<Backup>,
pub allow_short: bool,
pub include_content_duplicates: bool,
pub imported: usize,
pub vectors: Option<(usize, usize)>,
pub embed_error: Option<String>,
pub incompatible: Vec<IncompatibleFile>,
}
impl LegacyReport {
pub fn render(&self) -> String {
let mut out = format!(
"palace={} mode={}\n",
self.palace,
if self.dry_run { "dry-run" } else { "apply" }
);
if self.legacy_present {
out.push_str(&format!(
" legacy kg.db: rows={} unreadable={} already_live={} missing={} imported={} \
legacy_triples={} (not imported)\n",
self.legacy_rows,
self.unreadable.len(),
self.already_live,
self.missing,
self.imported,
self.legacy_triples
));
for u in &self.unreadable {
out.push_str(&format!(" unreadable row {u}\n"));
}
if let Some(dups) = self.content_duplicates {
let fate = match (self.dry_run, self.include_content_duplicates) {
(true, _) => "--apply skips them unless --include-content-duplicates",
(false, false) => "skipped; --include-content-duplicates imports them",
(false, true) => "imported by --include-content-duplicates",
};
out.push_str(&format!(
" content_duplicates={dups} (missing drawers whose content a live drawer \
already holds under another id; {fate})\n"
));
}
out.push_str(&guard::render_guards(
self.dry_run,
self.allow_short,
&self.rejected,
self.backup.as_ref(),
));
} else {
out.push_str(" legacy kg.db: none\n");
}
if let Some((repaired, still)) = self.vectors {
out.push_str(&format!(
" vectors: repaired={repaired} still_missing={still}\n"
));
}
if let Some(e) = &self.embed_error {
out.push_str(&format!(
" vectors: FAILED after the import committed: {e}\n"
));
}
let total: u64 = self.incompatible.iter().map(|f| f.bytes).sum();
out.push_str(&format!(
" .v2-incompatible files: {} ({total} bytes; redb 2.x, not readable here, left \
untouched)\n",
self.incompatible.len()
));
if self.dry_run && self.missing > 0 {
out.push_str("nothing was written — stop the daemon and re-run with --apply\n");
}
out
}
}
pub fn scan_report(palace: &Palace, allow_short: bool) -> Result<LegacyReport> {
let data_dir = &palace.data_dir;
let mut report = base_report(palace, true, allow_short)?;
let Some(legacy) = read_legacy_kg(data_dir)? else {
return Ok(report);
};
let (live, live_hashes) = with_store_copy(data_dir, &std::env::temp_dir(), |s| {
let hashes: HashSet<_> = s
.load_drawers()?
.iter()
.map(|d| memory_content_hash(d.content()))
.collect();
Ok((s.load_drawer_ids()?, hashes))
})?
.unwrap_or_default();
fill_legacy_counts(&mut report, &legacy, &live);
let (_, duplicates, rejected) = split_missing(legacy.drawers, &live, &live_hashes, allow_short);
report.content_duplicates = Some(duplicates.len());
report.rejected = rejected;
Ok(report)
}
fn split_missing(
drawers: Vec<Drawer>,
live: &HashSet<Uuid>,
live_hashes: &HashSet<ContentHash>,
allow_short: bool,
) -> (Vec<Drawer>, Vec<Drawer>, Vec<Rejected>) {
let missing = drawers.into_iter().filter(|d| !live.contains(&d.id));
let (passed, rejected) = screen_drawers(missing, allow_short);
let (distinct, duplicates) = passed
.into_iter()
.partition(|d| !live_hashes.contains(&memory_content_hash(d.content())));
(distinct, duplicates, rejected)
}
fn merge_imported(in_memory: &mut Vec<Drawer>, imported: Vec<Drawer>) {
let mut fresh: HashMap<Uuid, Drawer> = imported.into_iter().map(|d| (d.id, d)).collect();
for slot in in_memory.iter_mut() {
if let Some(d) = fresh.remove(&slot.id) {
*slot = d;
}
}
in_memory.extend(fresh.into_values());
}
pub async fn apply_report(
palace: &Palace,
embed: bool,
include_content_duplicates: bool,
allow_short: bool,
) -> Result<LegacyReport> {
let copy: CopyFn = |from, to| std::fs::copy(from, to);
let flags = (embed, include_content_duplicates, allow_short);
apply_report_with(palace, flags, &palace.data_dir, copy).await
}
pub(crate) async fn apply_report_with(
palace: &Palace,
(embed, include_content_duplicates, allow_short): (bool, bool, bool),
backup_parent: &Path,
copy: CopyFn,
) -> Result<LegacyReport> {
let mut report = base_report(palace, false, allow_short)?;
let Some(legacy) = read_legacy_kg(&palace.data_dir)? else {
return Ok(report);
};
probe_stores(&palace.data_dir)?;
let backup = backup_stores(&palace.data_dir, backup_parent, copy)?;
let kept = format!("backup kept at {}", backup.dir.display());
report.backup = Some(backup);
report.include_content_duplicates = include_content_duplicates;
import_legacy(palace, legacy, report, embed)
.await
.context(kept)
}
async fn import_legacy(
palace: &Palace,
legacy: LegacyDrawers,
mut report: LegacyReport,
embed: bool,
) -> Result<LegacyReport> {
let handle = PalaceHandle::open_with_intent(palace, OpenIntent::Writer)
.with_context(|| format!("open palace {} for writing", palace.id))?;
if handle.drawer_load_degraded {
bail!(
"palace {}: the live drawer table loaded degraded; refusing to import over rows \
it could not read",
palace.id
);
}
let live = handle
.kg
.load_drawer_ids()
.context("load the drawer ids in kg.redb")?;
let live_hashes: HashSet<ContentHash> = handle
.kg
.load_drawers()
.context("load the drawers in kg.redb")?
.iter()
.map(|d| memory_content_hash(d.content()))
.collect();
fill_legacy_counts(&mut report, &legacy, &live);
let (mut to_import, duplicates, rejected) =
split_missing(legacy.drawers, &live, &live_hashes, report.allow_short);
report.content_duplicates = Some(duplicates.len());
report.rejected = rejected;
if report.include_content_duplicates {
to_import.extend(duplicates);
}
if !to_import.is_empty() {
handle
.kg
.upsert_drawers_atomic(to_import.clone())
.await
.context("import legacy drawers into kg.redb")?;
report.imported = to_import.len();
merge_imported(&mut handle.drawers.write(), to_import);
}
if embed {
match handle
.backfill_missing_vectors(VectorBackfillOptions {
dry_run: false,
..VectorBackfillOptions::default()
})
.await
{
Ok(v) => report.vectors = Some((v.repaired, v.still_missing_ids.len())),
Err(e) => report.embed_error = Some(format!("{e:#}")),
}
}
Ok(report)
}
fn base_report(palace: &Palace, dry_run: bool, allow_short: bool) -> Result<LegacyReport> {
Ok(LegacyReport {
palace: palace.id.as_str().to_string(),
dry_run,
allow_short,
incompatible: list_incompatible_files(&palace.data_dir)?,
..LegacyReport::default()
})
}
fn fill_legacy_counts(report: &mut LegacyReport, legacy: &LegacyDrawers, live: &HashSet<Uuid>) {
report.legacy_present = true;
report.legacy_rows = legacy.total_rows;
report.unreadable = legacy.unreadable.clone();
report.legacy_triples = legacy.triple_rows;
report.already_live = legacy
.drawers
.iter()
.filter(|d| live.contains(&d.id))
.count();
report.missing = legacy.drawers.len() - report.already_live;
}
pub fn unaccounted_legacy_data(data_dir: &Path, live: &HashSet<Uuid>) -> Option<String> {
let mut reasons = Vec::new();
match read_legacy_kg(data_dir) {
Ok(None) => {}
Ok(Some(l)) => {
let missing =
l.drawers.iter().filter(|d| !live.contains(&d.id)).count() + l.unreadable.len();
if missing > 0 {
reasons.push(format!(
"legacy kg.db holds {missing} drawer(s) absent from the live store"
));
}
if l.triple_rows > 0 {
reasons.push(format!(
"legacy kg.db holds {} triple(s) the live graph never imported",
l.triple_rows
));
}
}
Err(e) => reasons.push(format!("legacy kg.db could not be checked: {e:#}")),
}
match list_incompatible_files(data_dir) {
Ok(f) if !f.is_empty() => reasons.push(format!("{} .v2-incompatible file(s)", f.len())),
Ok(_) => {}
Err(e) => reasons.push(format!("quarantine files could not be listed: {e:#}")),
}
(!reasons.is_empty()).then(|| reasons.join("; "))
}
#[cfg(test)]
#[path = "legacy_kg_tests.rs"]
pub(crate) mod tests;