use chrono::{DateTime, NaiveDateTime, Utc};
use sha2::{Digest, Sha256};
use trusty_common::memory_core::store::kg::Triple;
use super::bridge::{KuzuEntityRow, KuzuMemoryRow};
use super::screen::{screen, SecretRule};
pub const SOURCE_TAG_PREFIX: &str = "source:kuzu-memory/";
pub const HASH_TAG_PREFIX: &str = "kuzu-hash:";
pub const STORE_TAG_PREFIX: &str = "kuzu-store:";
pub const PENDING_TAG_PREFIX: &str = "kuzu-pending:";
pub const ORIGIN_TAG: &str = "kuzu-memory";
const COLUMN_TAG_PREFIXES: &[&str] = &[
"memory_type:",
"knowledge_type:",
"source_type:",
"project:",
"agent:",
"user:",
"session:",
"meta:",
];
const MAX_META_TAGS: usize = 8;
const MAX_META_VALUE_CHARS: usize = 64;
const MAX_TOKEN_CHARS: usize = 64;
const DEFAULT_IMPORTANCE: f32 = 0.5;
const DEFAULT_EDGE_CONFIDENCE: f32 = 0.8;
#[derive(Debug, Clone, PartialEq)]
pub struct MappedMemory {
pub source_key: String,
pub memory_id: String,
pub store: String,
pub content: String,
pub hash: String,
pub created_at: Option<DateTime<Utc>>,
pub importance: f32,
pub tags: Vec<String>,
pub refused_tags: Vec<SecretRule>,
}
impl MappedMemory {
pub fn staging_tags(&self) -> Vec<String> {
let mut tags: Vec<String> = self
.tags
.iter()
.filter(|t| !t.starts_with(SOURCE_TAG_PREFIX) && !t.starts_with(HASH_TAG_PREFIX))
.cloned()
.collect();
tags.push(format!("{PENDING_TAG_PREFIX}{}", self.memory_id));
tags
}
}
pub fn is_generated_tag(tag: &str) -> bool {
tag == ORIGIN_TAG
|| [
SOURCE_TAG_PREFIX,
HASH_TAG_PREFIX,
STORE_TAG_PREFIX,
PENDING_TAG_PREFIX,
]
.iter()
.chain(COLUMN_TAG_PREFIXES)
.any(|p| tag.starts_with(p))
}
pub fn merge_tags(current: &[String], fresh: &[String]) -> Vec<String> {
let mut out: Vec<String> = current
.iter()
.filter(|t| !is_generated_tag(t) && !fresh.contains(t))
.cloned()
.collect();
out.extend(fresh.iter().cloned());
out
}
pub fn source_key(memory_id: &str) -> String {
format!("kuzu-memory/{memory_id}")
}
fn is_safe_token(s: &str) -> bool {
!s.is_empty()
&& s.len() <= MAX_TOKEN_CHARS
&& s.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
}
pub fn kuzu_content_hash(content: &str) -> String {
let digest = Sha256::digest(content.trim().to_lowercase().as_bytes());
digest.iter().map(|b| format!("{b:02x}")).collect()
}
pub fn map_memory(row: &KuzuMemoryRow, store: &str) -> Option<MappedMemory> {
let memory_id = row.id.clone()?;
let content = row.content.clone().filter(|c| !c.trim().is_empty())?;
let mut refused_tags = Vec::new();
let mut passes = |v: &str| match screen(v) {
Some(rule) => {
refused_tags.push(rule);
false
}
None => true,
};
let hash = match row.content_hash.as_deref().filter(|h| !h.is_empty()) {
Some(h) if passes(h) => h.to_string(),
_ => kuzu_content_hash(&content),
};
let key = source_key(&memory_id);
let mut tags = vec![
format!("source:{key}"),
format!("{HASH_TAG_PREFIX}{hash}"),
format!("{STORE_TAG_PREFIX}{store}"),
ORIGIN_TAG.to_string(),
];
let columns = [
("memory_type", &row.memory_type),
("knowledge_type", &row.knowledge_type),
("source_type", &row.source_type),
("project", &row.project_tag),
("agent", &row.agent_id),
("user", &row.user_id),
("session", &row.session_id),
];
for (name, value) in columns {
let Some(v) = value.as_deref().map(str::trim).filter(|v| !v.is_empty()) else {
continue;
};
if (name == "agent" && v == "default") || !passes(v) {
continue;
}
tags.push(format!("{name}:{v}"));
}
tags.extend(metadata_tags(row.metadata.as_deref(), &mut passes));
Some(MappedMemory {
source_key: key,
memory_id,
store: store.to_string(),
content,
hash,
created_at: row.created_at.as_deref().and_then(parse_timestamp),
importance: row
.importance
.map(|i| (i as f32).clamp(0.0, 1.0))
.unwrap_or(DEFAULT_IMPORTANCE),
tags,
refused_tags,
})
}
fn metadata_tags(raw: Option<&str>, passes: &mut impl FnMut(&str) -> bool) -> Vec<String> {
let Some(serde_json::Value::Object(map)) = raw.and_then(|r| serde_json::from_str(r).ok())
else {
return Vec::new();
};
let mut out = Vec::new();
for (k, v) in map {
let value = match v {
serde_json::Value::String(s) => s,
serde_json::Value::Number(n) => n.to_string(),
serde_json::Value::Bool(b) => b.to_string(),
_ => continue,
};
if !is_safe_token(&k) || value.is_empty() || value.chars().count() > MAX_META_VALUE_CHARS {
continue;
}
if !passes(&k) || !passes(&value) {
continue;
}
out.push(format!("meta:{k}:{value}"));
if out.len() == MAX_META_TAGS {
break;
}
}
out
}
fn parse_timestamp(s: &str) -> Option<DateTime<Utc>> {
if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
return Some(dt.with_timezone(&Utc));
}
["%Y-%m-%dT%H:%M:%S%.f", "%Y-%m-%d %H:%M:%S%.f"]
.iter()
.find_map(|f| NaiveDateTime::parse_from_str(s, f).ok())
.map(|n| n.and_utc())
}
pub fn entity_subject(entity_id: &str) -> String {
format!("entity:{entity_id}")
}
pub fn drawer_subject(id: uuid::Uuid) -> String {
format!("drawer:{id}")
}
pub fn triple(subject: String, predicate: &str, object: String, confidence: Option<f64>) -> Triple {
Triple {
subject,
predicate: predicate.to_string(),
object,
valid_from: Utc::now(),
valid_to: None,
confidence: confidence
.map(|c| (c as f32).clamp(0.0, 1.0))
.unwrap_or(DEFAULT_EDGE_CONFIDENCE),
provenance: Some(ORIGIN_TAG.to_string()),
}
}
pub fn entity_triples(entity: &KuzuEntityRow) -> (Vec<Triple>, Vec<SecretRule>) {
let (mut out, mut refused) = (Vec::new(), Vec::new());
let Some(id) = entity.id.as_deref().filter(|i| !i.is_empty()) else {
return (out, refused);
};
fn present(v: &Option<String>) -> Option<&str> {
v.as_deref().filter(|s| !s.trim().is_empty())
}
let id_rule = screen(id);
for (predicate, value) in [
("has_name", present(&entity.name)),
("entity_type", present(&entity.entity_type)),
] {
let Some(value) = value else { continue };
match id_rule.or_else(|| screen(value)) {
Some(rule) => refused.push(rule),
None => out.push(triple(
entity_subject(id),
predicate,
value.to_string(),
None,
)),
}
}
(out, refused)
}
pub fn relates_to_predicate(relationship_type: Option<&str>) -> String {
match relationship_type
.map(str::trim)
.filter(|t| is_safe_token(t))
{
Some(t) => format!("relates_to:{t}"),
None => "relates_to".to_string(),
}
}