use anyhow::{Context, Result};
use std::collections::HashSet;
use trusty_common::memory_core::store::kg::{KnowledgeGraph, Triple};
use trusty_common::memory_core::store::OpenIntent;
use super::kg_rebuild::{scan_active_triples, STRUCTURAL_PREFIXES};
use crate::kg_extract::{canonical_entity, is_stop_token, AUTO_PROVENANCE};
use crate::AppState;
#[derive(Debug, Clone)]
pub struct TwinRepoint {
pub old: Triple,
pub new: Triple,
}
fn render(r: &TwinRepoint) -> String {
format!(
"({} {} {}) => ({} {} {})",
r.old.subject, r.old.predicate, r.old.object, r.new.subject, r.new.predicate, r.new.object
)
}
fn key(t: &Triple) -> (String, String, String) {
(t.subject.clone(), t.predicate.clone(), t.object.clone())
}
fn canonical_term(term: &str) -> Option<&str> {
if STRUCTURAL_PREFIXES.iter().any(|p| term.starts_with(p)) {
return None;
}
canonical_entity(term)
}
pub fn twin_repoints(active: &[Triple]) -> Vec<TwinRepoint> {
let mut out: Vec<TwinRepoint> = Vec::new();
for t in active {
if t.provenance.as_deref() != Some(AUTO_PROVENANCE) {
continue;
}
let subject = canonical_term(&t.subject);
let object = canonical_term(&t.object);
if subject.is_none() && object.is_none() {
continue;
}
if (subject.is_none() && is_stop_token(&t.subject))
|| (object.is_none() && is_stop_token(&t.object))
{
continue;
}
let mut new = t.clone();
if let Some(s) = subject {
new.subject = s.to_string();
}
if let Some(o) = object {
new.object = o.to_string();
}
out.push(TwinRepoint {
old: t.clone(),
new,
});
}
out.sort_by_key(|a| key(&a.old));
out
}
#[derive(Debug, Clone)]
pub struct PalaceMergeSummary {
pub palace_id: String,
pub selected: Vec<String>,
pub merged: Vec<String>,
pub failed: Vec<(String, String)>,
pub error: Option<String>,
}
impl PalaceMergeSummary {
pub fn failure_count(&self) -> usize {
if !self.failed.is_empty() {
self.failed.len()
} else {
usize::from(self.error.is_some())
}
}
}
pub async fn report_merge(
state: &AppState,
palace_filter: Option<&str>,
apply: bool,
) -> Result<usize> {
let summaries = merge_palaces(state, palace_filter, apply).await?;
let mut total = 0usize;
let failures: usize = summaries
.iter()
.map(PalaceMergeSummary::failure_count)
.sum();
for s in &summaries {
if let Some(e) = &s.error {
eprintln!("[merge-error] palace={} error={}", s.palace_id, e);
if s.failed.is_empty() {
continue;
}
}
if apply {
for moved in &s.merged {
println!("[merge] palace={} repointed {}", s.palace_id, moved);
}
for (moved, err) in &s.failed {
eprintln!(
"[merge-FAILED] palace={} repoint={} error={}",
s.palace_id, moved, err
);
}
total += s.merged.len();
} else {
for moved in &s.selected {
println!("[merge] palace={} would repoint {}", s.palace_id, moved);
}
total += s.selected.len();
}
}
if apply {
println!("kg-rebuild merge: {total} triples repointed, {failures} failed");
} else {
println!("kg-rebuild merge: {total} triples would be repointed (dry run)");
}
Ok(failures)
}
pub async fn merge_palaces(
state: &AppState,
palace_filter: Option<&str>,
apply: bool,
) -> Result<Vec<PalaceMergeSummary>> {
let mut out: Vec<PalaceMergeSummary> = Vec::new();
let palaces = trusty_common::memory_core::PalaceRegistry::list_palaces(&state.data_root)
.with_context(|| format!("list palaces under {}", state.data_root.display()))?;
for palace in palaces {
let id = palace.id.0.clone();
if palace_filter.is_some_and(|filter| filter != id) {
continue;
}
let summary = merge_one(state, &id, apply)
.await
.unwrap_or_else(|e| PalaceMergeSummary {
palace_id: id.clone(),
selected: Vec::new(),
merged: Vec::new(),
failed: Vec::new(),
error: Some(format!("{e:#}")),
});
out.push(summary);
}
Ok(out)
}
async fn merge_one(state: &AppState, palace_id: &str, apply: bool) -> Result<PalaceMergeSummary> {
let kg_path = state.data_root.join(palace_id).join("kg.db");
let intent = if apply {
OpenIntent::Writer
} else {
OpenIntent::ReadOnlyClient
};
let kg = KnowledgeGraph::open_with_intent(&kg_path, intent)
.with_context(|| format!("open kg for palace {palace_id}"))?;
let active = scan_active_triples(&kg, palace_id).await?;
let selected = twin_repoints(&active);
let mut summary = PalaceMergeSummary {
palace_id: palace_id.to_string(),
selected: selected.iter().map(render).collect(),
merged: Vec::new(),
failed: Vec::new(),
error: None,
};
if !apply {
return Ok(summary);
}
let mut live: HashSet<(String, String, String)> = active.iter().map(key).collect();
for r in &selected {
match repoint(&kg, r, &live).await {
Ok(()) => {
live.insert(key(&r.new));
summary.merged.push(render(r));
}
Err(e) => summary.failed.push((render(r), format!("{e:#}"))),
}
}
if !summary.failed.is_empty() {
summary.error = Some(format!(
"{} of {} re-point(s) failed",
summary.failed.len(),
selected.len()
));
}
Ok(summary)
}
async fn repoint(
kg: &KnowledgeGraph,
r: &TwinRepoint,
live: &HashSet<(String, String, String)>,
) -> Result<()> {
if !live.contains(&key(&r.new)) {
kg.assert(r.new.clone())
.await
.with_context(|| format!("assert cleaned twin {}", render(r)))?;
}
let closed = kg
.retract_triple(&r.old.subject, &r.old.predicate, &r.old.object)
.await
.with_context(|| format!("retract punctuated twin {}", render(r)))?;
anyhow::ensure!(
closed > 0,
"punctuated row vanished before its retract, so the fact now stands at both nodes: {}",
render(r)
);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::commands::kg_rebuild::stale_subject_candidates;
use serde_json::json;
use trusty_common::memory_core::palace::PalaceId;
fn triple(subject: &str, predicate: &str, object: &str, provenance: Option<&str>) -> Triple {
Triple {
subject: subject.to_string(),
predicate: predicate.to_string(),
object: object.to_string(),
valid_from: chrono::Utc::now(),
valid_to: None,
confidence: 0.6,
provenance: provenance.map(|p| p.to_string()),
}
}
#[test]
fn twin_repoints_moves_both_positions() {
let active = vec![
triple("`redb`", "uses", "mmap", Some(AUTO_PROVENANCE)),
triple("trusty-memory", "uses", "`redb`", Some(AUTO_PROVENANCE)),
triple("(sled)", "is-a", "*store*", Some(AUTO_PROVENANCE)),
triple("redb", "is-a", "database", Some(AUTO_PROVENANCE)),
];
let got = twin_repoints(&active);
let moved: Vec<(String, String)> = got
.iter()
.map(|r| {
(
format!("{} {}", r.old.subject, r.old.object),
format!("{} {}", r.new.subject, r.new.object),
)
})
.collect();
assert_eq!(
moved,
vec![
("(sled) *store*".to_string(), "sled store".to_string()),
("`redb` mmap".to_string(), "redb mmap".to_string()),
(
"trusty-memory `redb`".to_string(),
"trusty-memory redb".to_string()
),
],
"subject-only, object-only and both-position twins must all move, \
and an already-clean triple must not"
);
assert!(
got.iter().all(|r| r.old.predicate == r.new.predicate
&& r.old.confidence == r.new.confidence
&& r.old.provenance == r.new.provenance),
"a re-point changes identity, never the fact's metadata"
);
}
#[test]
fn twin_repoints_skips_namespaces() {
let drawer = format!("drawer:{}", uuid::Uuid::new_v4());
let active = vec![
triple("tag:v1.", "tags", &drawer, Some(AUTO_PROVENANCE)),
triple(
"topic:(beta)",
"mentioned-in",
&drawer,
Some(AUTO_PROVENANCE),
),
triple("room:*general*", "contains", &drawer, Some(AUTO_PROVENANCE)),
];
assert!(
twin_repoints(&active).is_empty(),
"namespaced terms must never be re-pointed"
);
}
#[test]
fn twin_repoints_spares_a_manual_triple() {
let active = vec![
triple("`redb`", "is-a", "database", Some("kg_assert")),
triple("`sled`", "is-a", "database", None),
];
assert!(
twin_repoints(&active).is_empty(),
"only auto-extracted triples may be re-pointed"
);
}
#[test]
fn twin_repoints_leaves_the_purge_its_stopwords() {
let active = vec![
triple("(\"the", "is-a", "thing", Some(AUTO_PROVENANCE)),
triple("`redb`", "is-a", "database", Some(AUTO_PROVENANCE)),
triple("`redb`", "is-a", "(\"the", Some(AUTO_PROVENANCE)),
triple("(\"the", "uses", "`redb`", Some(AUTO_PROVENANCE)),
];
let merged: Vec<(String, String)> = twin_repoints(&active)
.iter()
.map(|r| (r.new.subject.clone(), r.new.object.clone()))
.collect();
assert_eq!(
merged,
vec![("redb".to_string(), "database".to_string())],
"a stopword in either position keeps the whole triple out of this pass"
);
assert_eq!(
stale_subject_candidates(&active),
vec!["(\"the".to_string()],
"the two passes must partition the punctuated subjects, not overlap"
);
}
#[test]
fn merge_summary_counts_each_failure_once() {
let mut s = PalaceMergeSummary {
palace_id: "a".to_string(),
selected: vec!["x".to_string(), "y".to_string()],
merged: vec!["x".to_string()],
failed: vec![("y".to_string(), "store is read-only".to_string())],
error: Some("1 of 2 re-point(s) failed".to_string()),
};
assert_eq!(s.failure_count(), 1, "the error merely summarises the one");
s.failed.clear();
assert_eq!(s.failure_count(), 1, "a palace-level error is one failure");
s.error = None;
assert_eq!(s.failure_count(), 0);
}
async fn palace_fixture(data_root: std::path::PathBuf) -> Result<AppState> {
trusty_common::memory_core::retrieval::seed_shared_embedder_with_mock();
unsafe {
std::env::set_var("TRUSTY_SKIP_PALACE_ENFORCEMENT", "1");
}
let state = AppState::new(data_root);
state.set_ready();
let _ = crate::tools::dispatch_tool(&state, "palace_create", json!({"name": "a"})).await?;
Ok(state)
}
async fn objects_of(kg: &KnowledgeGraph, subject: &str) -> Result<Vec<String>> {
let mut got: Vec<String> = kg
.query_active(subject)
.await?
.into_iter()
.map(|t| t.object)
.collect();
got.sort();
Ok(got)
}
#[tokio::test]
async fn merge_repoints_both_positions_and_keeps_the_cleaned_nodes_triples() -> Result<()> {
let tmp = tempfile::tempdir()?;
let state = palace_fixture(tmp.path().to_path_buf()).await?;
let handle = state
.registry
.open_palace(&state.data_root, &PalaceId::new("a"))?;
for t in [
triple("redb", "is-a", "database", Some(AUTO_PROVENANCE)),
triple("`redb`", "uses", "mmap", Some(AUTO_PROVENANCE)),
triple("trusty-memory", "uses", "`redb`", Some(AUTO_PROVENANCE)),
triple("trusty-memory", "uses", "tokio", Some(AUTO_PROVENANCE)),
triple("`sled`", "is-a", "store", Some("kg_assert")),
] {
handle.kg.assert(t).await?;
}
let applied = merge_palaces(&state, Some("a"), true).await?;
assert_eq!(applied.len(), 1);
assert!(
applied[0].failed.is_empty(),
"no re-point should have failed"
);
assert!(applied[0].error.is_none());
assert_eq!(applied[0].merged.len(), 2, "two twins must have moved");
assert_eq!(
objects_of(&handle.kg, "redb").await?,
vec!["database".to_string(), "mmap".to_string()],
"the merged node must keep its OWN pre-existing triple and gain the re-pointed one"
);
assert!(
handle.kg.query_active("`redb`").await?.is_empty(),
"the punctuated node must be gone from the subject position"
);
assert_eq!(
objects_of(&handle.kg, "trusty-memory").await?,
vec!["redb".to_string(), "tokio".to_string()],
"the object position must move onto the cleaned node without \
taking its sibling objects at the same predicate down"
);
assert_eq!(
objects_of(&handle.kg, "`sled`").await?,
vec!["store".to_string()],
"a manually asserted punctuated triple must be untouched"
);
let again = merge_palaces(&state, Some("a"), true).await?;
assert!(
again[0].selected.is_empty(),
"the merge must be idempotent, got {:?}",
again[0].selected
);
Ok(())
}
#[tokio::test]
async fn merge_dry_run_writes_nothing() -> Result<()> {
let tmp = tempfile::tempdir()?;
let state = palace_fixture(tmp.path().to_path_buf()).await?;
let handle = state
.registry
.open_palace(&state.data_root, &PalaceId::new("a"))?;
handle
.kg
.assert(triple("`redb`", "uses", "mmap", Some(AUTO_PROVENANCE)))
.await?;
let dry = merge_palaces(&state, Some("a"), false).await?;
assert_eq!(dry[0].selected.len(), 1, "the twin must be reported");
assert!(dry[0].merged.is_empty(), "a dry run repoints nothing");
assert!(
!handle.kg.query_active("`redb`").await?.is_empty(),
"a dry run must leave the punctuated node standing"
);
assert!(
handle.kg.query_active("redb").await?.is_empty(),
"a dry run must not create the cleaned node either"
);
assert_eq!(report_merge(&state, Some("a"), false).await?, 0);
Ok(())
}
#[tokio::test]
async fn merge_palaces_propagates_an_unreadable_data_root() -> Result<()> {
let tmp = tempfile::tempdir()?;
let data_root = tmp.path().join("not-a-directory");
std::fs::write(&data_root, b"")?;
let state = AppState::new(data_root);
let err = merge_palaces(&state, None, false)
.await
.expect_err("an unlistable data root must fail the run, not report an empty one");
let rendered = format!("{err:#}");
assert!(
rendered.contains("list palaces"),
"the error must name what could not be read, got {rendered}"
);
Ok(())
}
#[tokio::test]
async fn merge_reports_a_failed_repoint_and_leaves_the_fact_readable() -> Result<()> {
let seed_root = tempfile::tempdir()?;
let seeded = palace_fixture(seed_root.path().to_path_buf()).await?;
let handle = seeded
.registry
.open_palace(&seeded.data_root, &PalaceId::new("a"))?;
handle
.kg
.assert(triple("`redb`", "uses", "mmap", Some(AUTO_PROVENANCE)))
.await?;
let tmp = tempfile::tempdir()?;
let palace_dir = tmp.path().join("a");
std::fs::create_dir_all(&palace_dir)?;
for entry in std::fs::read_dir(seeded.data_root.join("a"))? {
let entry = entry?;
if entry.file_type()?.is_file() {
std::fs::copy(entry.path(), palace_dir.join(entry.file_name()))?;
}
}
let kg_path = palace_dir.join("kg.db");
let state = AppState::new(tmp.path().to_path_buf());
state.set_ready();
let live = redb::Database::create(palace_dir.join("kg.redb"))
.context("hold the palace's redb lock")?;
let read_only = KnowledgeGraph::open_with_intent(&kg_path, OpenIntent::ReadOnlyClient)?;
let applied = merge_palaces(&state, Some("a"), true).await?;
assert_eq!(applied.len(), 1);
assert_eq!(applied[0].selected.len(), 1, "the twin must still be found");
assert!(
applied[0].merged.is_empty(),
"a re-point that could not write must never be reported as merged"
);
assert_eq!(applied[0].failed.len(), 1, "the failure must be recorded");
assert!(
applied[0].failed[0].1.contains("read-only"),
"the failure must carry its error text, got {:?}",
applied[0].failed[0].1
);
assert!(
applied[0].error.is_some(),
"a palace with a failed re-point must carry an error"
);
assert!(
report_merge(&state, Some("a"), true).await? > 0,
"the failure must reach the count kg_rebuild bails on"
);
drop(read_only);
drop(live);
let after = KnowledgeGraph::open_with_intent(&kg_path, OpenIntent::Writer)?;
assert_eq!(
objects_of(&after, "`redb`").await?,
vec!["mmap".to_string()],
"a re-point that failed must leave the fact readable at the punctuated node"
);
Ok(())
}
}