use std::path::Path;
use kimetsu_core::KimetsuResult;
use kimetsu_core::ids::RunId;
use kimetsu_core::memory::{MemoryKind, MemoryScope};
use rusqlite::params;
use crate::project::{add_memory, invalidate_memory, load_project, load_project_readonly};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MemoryExport {
pub text: String,
pub scope: String,
pub kind: String,
pub confidence: f32,
pub created_at: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Pack {
pub kimetsu_pack: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exported_at: Option<String>,
#[serde(default)]
pub memory_count: usize,
pub memories: Vec<MemoryExport>,
}
#[derive(Debug, Clone, Default)]
pub struct PackRef {
pub name: Option<String>,
pub version: Option<String>,
}
pub fn parse_pack_or_array(json: &str) -> KimetsuResult<(PackRef, Vec<MemoryExport>)> {
if let Ok(pack) = serde_json::from_str::<Pack>(json) {
return Ok((
PackRef {
name: pack.name,
version: pack.version,
},
pack.memories,
));
}
let entries: Vec<MemoryExport> = serde_json::from_str(json)
.map_err(|e| format!("pack: not a Pack envelope or a memory array: {e}"))?;
Ok((PackRef::default(), entries))
}
pub fn redact_context_suffix(text: &str) -> &str {
let trimmed = text.trim_end();
if let Some(pos) = find_trailing_context_paren(trimmed) {
let candidate = trimmed[..pos].trim_end();
if !candidate.is_empty() {
return candidate;
}
}
text
}
pub fn redact_tags_prefix(text: &str) -> &str {
let trimmed = text.trim_start();
if let Some(rest) = trimmed.strip_prefix("[tags: ") {
if let Some(close) = rest.find(']') {
let after = rest[close + 1..].trim_start();
if !after.is_empty() {
return after;
}
}
}
text
}
pub fn apply_export_redaction(
entry: MemoryExport,
redact: bool,
redact_tags: bool,
) -> MemoryExport {
if !redact && !redact_tags {
return entry;
}
let mut text: &str = &entry.text;
let after_tags: String;
let after_ctx: String;
if redact_tags {
let stripped = redact_tags_prefix(text);
after_tags = stripped.to_string();
text = &after_tags;
}
if redact {
let stripped = redact_context_suffix(text);
after_ctx = stripped.to_string();
text = &after_ctx;
}
MemoryExport {
text: text.to_string(),
..entry
}
}
fn find_trailing_context_paren(s: &str) -> Option<usize> {
if !s.ends_with(')') {
return None;
}
let bytes = s.as_bytes();
let close = s.len() - 1;
let prefix = b" (context: ";
for start in (0..close).rev() {
if start + prefix.len() > close {
continue;
}
if &bytes[start..start + prefix.len()] == prefix {
return Some(start);
}
}
None
}
#[derive(Debug, Clone, Default)]
pub struct ImportSummary {
pub imported: usize,
pub deduped: usize,
pub superseded: usize,
pub quarantined: usize,
}
thread_local! {
static IMPORT_PROVENANCE: std::cell::RefCell<Option<serde_json::Value>> =
const { std::cell::RefCell::new(None) };
}
pub(crate) struct ImportProvenanceScope;
impl ImportProvenanceScope {
pub(crate) fn new(v: serde_json::Value) -> Self {
IMPORT_PROVENANCE.with(|c| *c.borrow_mut() = Some(v));
ImportProvenanceScope
}
}
impl Drop for ImportProvenanceScope {
fn drop(&mut self) {
IMPORT_PROVENANCE.with(|c| *c.borrow_mut() = None);
}
}
pub(crate) fn build_provenance(run_id: RunId, text: &str) -> serde_json::Value {
IMPORT_PROVENANCE.with(|c| {
if let Some(src) = c.borrow().as_ref() {
let mut v = src.clone();
if let Some(obj) = v.as_object_mut() {
obj.insert("run_id".into(), serde_json::json!(run_id.to_string()));
obj.insert("text".into(), serde_json::json!(text));
}
v
} else {
serde_json::json!({
"source": "manual_cli",
"run_id": run_id.to_string(),
"text": text,
})
}
})
}
#[derive(Debug, Clone, Default, serde::Serialize)]
pub struct ScrubReport {
pub total: usize,
pub kinds: std::collections::BTreeMap<String, usize>,
}
impl ScrubReport {
pub fn is_clean(&self) -> bool {
self.total == 0
}
pub fn summary(&self) -> String {
if self.total == 0 {
return "no credentials or PII found".to_string();
}
let parts: Vec<String> = self.kinds.iter().map(|(k, n)| format!("{k}×{n}")).collect();
format!("scrubbed {}: {}", self.total, parts.join(", "))
}
}
pub fn export_memories(
start: &Path,
scope: Option<MemoryScope>,
kind: Option<MemoryKind>,
redact: bool,
redact_tags: bool,
) -> KimetsuResult<(Vec<MemoryExport>, ScrubReport)> {
let (sql, params_vec): (&str, Vec<String>) = match (scope.as_ref(), kind.as_ref()) {
(Some(s), Some(k)) => (
"SELECT scope, kind, text, confidence, created_at
FROM memories
WHERE invalidated_at IS NULL
AND superseded_by IS NULL
AND lower(scope) = lower(?1)
AND lower(kind) = lower(?2)
ORDER BY created_at DESC",
vec![s.to_string(), k.to_string()],
),
(Some(s), None) => (
"SELECT scope, kind, text, confidence, created_at
FROM memories
WHERE invalidated_at IS NULL
AND superseded_by IS NULL
AND lower(scope) = lower(?1)
ORDER BY created_at DESC",
vec![s.to_string()],
),
(None, Some(k)) => (
"SELECT scope, kind, text, confidence, created_at
FROM memories
WHERE invalidated_at IS NULL
AND superseded_by IS NULL
AND lower(kind) = lower(?1)
ORDER BY created_at DESC",
vec![k.to_string()],
),
(None, None) => (
"SELECT scope, kind, text, confidence, created_at
FROM memories
WHERE invalidated_at IS NULL
AND superseded_by IS NULL
ORDER BY created_at DESC",
vec![],
),
};
let (_paths, _config, conn) = load_project(start)?;
let mut stmt = conn.prepare(sql)?;
let refs: Vec<&dyn rusqlite::ToSql> = params_vec
.iter()
.map(|s| s as &dyn rusqlite::ToSql)
.collect();
let rows = stmt.query_map(refs.as_slice(), |row| {
Ok(MemoryExport {
scope: row.get(0)?,
kind: row.get(1)?,
text: row.get(2)?,
confidence: row.get::<_, f64>(3)? as f32,
created_at: row.get(4)?,
})
})?;
let mut out = Vec::new();
let mut report = ScrubReport::default();
for row in rows {
let mut entry = apply_export_redaction(row?, redact, redact_tags);
let scrubbed = crate::redact::scrub_for_export(&entry.text);
for m in &scrubbed.matches {
*report.kinds.entry(m.kind.to_string()).or_insert(0) += 1;
report.total += 1;
}
entry.text = scrubbed.text;
out.push(entry);
}
Ok((out, report))
}
pub fn import_memories(
start: &Path,
entries: &[MemoryExport],
scope_override: Option<MemoryScope>,
) -> KimetsuResult<ImportSummary> {
let mut summary = ImportSummary::default();
let pre_existing_ids: std::collections::HashSet<String> = {
match load_project_readonly(start) {
Ok((_paths, _config, conn)) => {
let mut stmt = conn
.prepare("SELECT memory_id FROM memories WHERE invalidated_at IS NULL")
.unwrap_or_else(|_| conn.prepare("SELECT memory_id FROM memories").unwrap());
stmt.query_map([], |row| row.get::<_, String>(0))
.map(|rows| rows.filter_map(|r| r.ok()).collect())
.unwrap_or_default()
}
Err(_) => std::collections::HashSet::new(),
}
};
let mut this_batch_ids: std::collections::HashSet<String> = std::collections::HashSet::new();
for entry in entries {
let scope = if let Some(ref ov) = scope_override {
*ov
} else {
match entry.scope.parse::<MemoryScope>() {
Ok(s) => s,
Err(_) => {
eprintln!(
"kimetsu-brain import: skipping entry with unknown scope `{}`",
entry.scope
);
summary.deduped += 1;
continue;
}
}
};
let kind = match entry.kind.parse::<MemoryKind>() {
Ok(k) => k,
Err(_) => {
eprintln!(
"kimetsu-brain import: skipping entry with unknown kind `{}`",
entry.kind
);
summary.deduped += 1;
continue;
}
};
match add_memory(start, scope, kind, &entry.text) {
Ok(id) => {
if pre_existing_ids.contains(&id) || !this_batch_ids.insert(id) {
summary.deduped += 1;
} else {
summary.imported += 1;
}
}
Err(e) => {
eprintln!("kimetsu-brain import: failed to add memory: {e}");
summary.deduped += 1;
}
}
}
Ok(summary)
}
pub fn quarantine_memories(
start: &Path,
entries: &[MemoryExport],
scope_override: Option<MemoryScope>,
pack: Option<&PackRef>,
) -> KimetsuResult<ImportSummary> {
let mut summary = ImportSummary::default();
let existing: std::collections::HashSet<String> = match load_project_readonly(start) {
Ok((_paths, _config, conn)) => conn
.prepare("SELECT normalized_text FROM memories WHERE invalidated_at IS NULL")
.and_then(|mut stmt| {
stmt.query_map([], |row| row.get::<_, String>(0))
.map(|rows| rows.filter_map(|r| r.ok()).collect())
})
.unwrap_or_default(),
Err(_) => std::collections::HashSet::new(),
};
let origin = match pack {
Some(p) => format!(
"pack {}@{}",
p.name.as_deref().unwrap_or("unknown"),
p.version.as_deref().unwrap_or("?")
),
None => "an import".to_string(),
};
let rationale = format!(
"Quarantined on import from {origin}. Imported memories are held for \
review rather than entering retrieval, because a poisoned memory \
persists across every future session. Accept only what you would have \
written yourself."
);
let mut seen_in_batch = std::collections::HashSet::new();
for entry in entries {
let scope = match scope_override {
Some(ov) => ov,
None => match entry.scope.parse::<MemoryScope>() {
Ok(s) => s,
Err(_) => {
eprintln!(
"kimetsu-brain import: skipping entry with unknown scope `{}`",
entry.scope
);
summary.deduped += 1;
continue;
}
},
};
let kind = match entry.kind.parse::<MemoryKind>() {
Ok(k) => k,
Err(_) => {
eprintln!(
"kimetsu-brain import: skipping entry with unknown kind `{}`",
entry.kind
);
summary.deduped += 1;
continue;
}
};
let normalized = kimetsu_core::memory::normalize_memory_text(&entry.text);
if existing.contains(&normalized) || !seen_in_batch.insert(normalized) {
summary.deduped += 1;
continue;
}
match crate::project::propose_memory(
start,
scope,
kind,
&entry.text,
entry.confidence,
&rationale,
) {
Ok(_) => summary.quarantined += 1,
Err(e) => {
eprintln!("kimetsu-brain import: failed to quarantine memory: {e}");
summary.deduped += 1;
}
}
}
Ok(summary)
}
pub fn import_pack(
start: &Path,
entries: &[MemoryExport],
scope_override: Option<MemoryScope>,
replace: bool,
pack: Option<&PackRef>,
quarantine: bool,
) -> KimetsuResult<ImportSummary> {
let mut superseded = 0usize;
if replace {
let scopes = pack_target_scopes(entries, scope_override);
let reason = match pack {
Some(p) => format!(
"replaced_by_pack:{}@{}",
p.name.as_deref().unwrap_or("unknown"),
p.version.as_deref().unwrap_or("?")
),
None => "replaced_by_import".to_string(),
};
for id in active_memory_ids_in_scopes(start, &scopes)? {
invalidate_memory(start, &id, Some(&reason))?;
superseded += 1;
}
}
let scrubbed: Vec<MemoryExport> = entries
.iter()
.map(|e| {
let mut e = e.clone();
e.text = crate::redact::scrub_for_export(&e.text).text;
e
})
.collect();
let _prov = pack.map(|p| {
ImportProvenanceScope::new(serde_json::json!({
"source": "pack",
"pack_name": p.name,
"pack_version": p.version,
}))
});
let mut summary = if quarantine {
quarantine_memories(start, &scrubbed, scope_override, pack)?
} else {
import_memories(start, &scrubbed, scope_override)?
};
summary.superseded = superseded;
Ok(summary)
}
fn pack_target_scopes(
entries: &[MemoryExport],
scope_override: Option<MemoryScope>,
) -> Vec<MemoryScope> {
if let Some(ov) = scope_override {
return vec![ov];
}
let mut seen = std::collections::HashSet::new();
let mut out = Vec::new();
for e in entries {
if let Ok(s) = e.scope.parse::<MemoryScope>() {
if seen.insert(s.to_string()) {
out.push(s);
}
}
}
out
}
fn active_memory_ids_in_scopes(start: &Path, scopes: &[MemoryScope]) -> KimetsuResult<Vec<String>> {
if scopes.is_empty() {
return Ok(Vec::new());
}
let (_p, _c, conn) = load_project_readonly(start)?;
let mut ids = Vec::new();
for sc in scopes {
let mut stmt = conn.prepare(
"SELECT memory_id FROM memories
WHERE scope = ?1 AND invalidated_at IS NULL AND superseded_by IS NULL",
)?;
let rows = stmt.query_map(params![sc.to_string()], |r| r.get::<_, String>(0))?;
for r in rows {
ids.push(r?);
}
}
Ok(ids)
}