use std::collections::{BTreeMap, BTreeSet, HashMap};
use serde_json::Value;
use topodb::{Db, EdgeId, NodeId, NodeRecord, Op, PropValue, Props, Scope, ScopeSet, TopoError};
use crate::{
merge_required_prop, normalize_edge_type, scopes_to_scope_set, ALIAS_EDGE_TYPE, ALIAS_LABEL,
ALIAS_NAME_PROP, ENTITY_LABEL, ENTITY_NAME_PROP, MEMORY_CONTENT_HASH_PROP, MEMORY_CONTENT_PROP,
MEMORY_LABEL, MEMORY_SUPERSEDED_AT_PROP,
};
pub const DEFAULT_REMEMBER_EDGE_TYPE: &str = "about";
#[derive(Debug)]
pub enum ComposeError {
Invalid(String),
Engine(TopoError),
}
impl From<TopoError> for ComposeError {
fn from(e: TopoError) -> Self {
ComposeError::Engine(e)
}
}
pub fn entity_dedup_key(name: &str) -> String {
name.split_whitespace()
.collect::<Vec<_>>()
.join(" ")
.to_lowercase()
}
pub fn normalize_content(content: &str) -> String {
content.split_whitespace().collect::<Vec<_>>().join(" ")
}
pub fn content_hash(content: &str) -> String {
let normalized = normalize_content(content);
let mut h: u64 = 0xcbf2_9ce4_8422_2325;
for b in normalized.as_bytes() {
h ^= *b as u64;
h = h.wrapping_mul(0x0000_0100_0000_01b3);
}
format!("{h:016x}")
}
pub fn memory_props(content: &str, extra: Option<&Value>) -> Result<Props, String> {
if let Some(Value::Object(map)) = extra {
for reserved in [MEMORY_CONTENT_HASH_PROP, MEMORY_SUPERSEDED_AT_PROP] {
if map.contains_key(reserved) {
return Err(format!(
"props must not include {reserved:?}: it is maintained by the engine write path"
));
}
}
}
let mut props = merge_required_prop(
MEMORY_CONTENT_PROP,
PropValue::Str(content.to_string()),
extra,
)?;
props.insert(
MEMORY_CONTENT_HASH_PROP.into(),
PropValue::Str(content_hash(content)),
);
Ok(props)
}
pub fn resolve_entities_by_name(
db: &Db,
scopes: &ScopeSet,
name: &str,
) -> Result<Vec<NodeRecord>, TopoError> {
let value = PropValue::Str(name.to_string());
let mut out = db.nodes_by_prop_normalized(scopes, ENTITY_LABEL, ENTITY_NAME_PROP, &value)?;
let aliases = match db.nodes_by_prop_normalized(scopes, ALIAS_LABEL, ALIAS_NAME_PROP, &value) {
Ok(hits) => hits,
Err(TopoError::Rejected(_)) => Vec::new(),
Err(e) => return Err(e),
};
for alias in aliases {
for edge in db.edges_from(scopes, alias.id, None, Some(ALIAS_EDGE_TYPE), true)? {
if let Some(canonical) = db.node(scopes, edge.to) {
if canonical.label == ENTITY_LABEL {
out.push(canonical);
}
}
}
}
out.sort_by_key(|n| n.id);
out.dedup_by_key(|n| n.id);
Ok(out)
}
pub fn find_existing_entity(
db: &Db,
lookup: &ScopeSet,
name: &str,
) -> Result<Option<NodeRecord>, TopoError> {
match resolve_entities_by_name(db, lookup, name) {
Ok(hits) => Ok(hits.into_iter().min_by_key(|n| n.id)),
Err(TopoError::Rejected(_)) => Ok(None),
Err(e) => Err(e),
}
}
pub fn existing_memory(
db: &Db,
write_scope: Scope,
content: &str,
) -> Result<Option<NodeId>, TopoError> {
let hash = content_hash(content);
let want = normalize_content(content);
let scope_set = scopes_to_scope_set(&[write_scope]);
let candidates = db.nodes_by_prop(
&scope_set,
MEMORY_LABEL,
MEMORY_CONTENT_HASH_PROP,
&PropValue::Str(hash),
)?;
Ok(candidates
.into_iter()
.filter(|n| {
!n.props.contains_key(MEMORY_SUPERSEDED_AT_PROP)
&& matches!(n.props.get(MEMORY_CONTENT_PROP), Some(PropValue::Str(c)) if normalize_content(c) == want)
})
.min_by_key(|n| n.id)
.map(|n| n.id))
}
fn plan_supersede(
db: &Db,
write_scope: Scope,
ids: &[String],
now_ms: i64,
) -> Result<(Vec<Op>, Vec<String>), ComposeError> {
let mut ops = Vec::new();
let mut marked = Vec::new();
if ids.is_empty() {
return Ok((ops, marked));
}
let scope_set = scopes_to_scope_set(&[write_scope]);
let mut seen = BTreeSet::new();
for raw in ids {
let id: NodeId = raw
.parse()
.map_err(|e| ComposeError::Invalid(format!("invalid node id {raw:?}: {e}")))?;
if !seen.insert(id) {
continue;
}
let node = db.node(&scope_set, id).ok_or_else(|| {
ComposeError::Invalid(format!(
"supersedes id {raw} is not a node in the write scope"
))
})?;
if node.label != MEMORY_LABEL {
return Err(ComposeError::Invalid(format!(
"supersedes id {raw} is a {}, not a Memory",
node.label
)));
}
if node.props.contains_key(MEMORY_SUPERSEDED_AT_PROP) {
continue;
}
let mut props: BTreeMap<String, Option<PropValue>> = BTreeMap::new();
props.insert(
MEMORY_SUPERSEDED_AT_PROP.into(),
Some(PropValue::Int(now_ms)),
);
ops.push(Op::SetNodeProps { id, props });
for e in db.edges_from(&scope_set, id, None, None, true)? {
ops.push(Op::CloseEdge {
id: e.id,
valid_to: None,
});
}
marked.push(id.to_string());
}
Ok((ops, marked))
}
pub struct RememberRequest {
pub content: String,
pub entities: Vec<String>,
pub edge_type: Option<String>,
pub supersedes: Vec<String>,
pub props: Option<Value>,
}
impl RememberRequest {
pub fn validate(&self) -> Result<String, String> {
let ty = normalize_edge_type(
self.edge_type
.as_deref()
.unwrap_or(DEFAULT_REMEMBER_EDGE_TYPE),
)?;
if self.entities.is_empty() {
return Err(
"entities must contain at least one name — use create_memory for a deliberately unlinked note".into(),
);
}
if self.entities.iter().any(|n| n.trim().is_empty()) {
return Err("entity names must be non-empty".into());
}
Ok(ty)
}
}
pub struct PlannedEntity {
pub name: String,
pub id: NodeId,
pub created: bool,
}
pub struct RememberPlan {
pub ops: Vec<Op>,
pub memory_id: NodeId,
pub deduplicated: bool,
pub new_memory: Option<String>,
pub new_entities: Vec<(NodeId, String)>,
pub entities: Vec<PlannedEntity>,
pub edge_ids: Vec<String>,
pub superseded: Vec<String>,
}
pub fn plan_remember(
db: &Db,
write_scope: Scope,
lookup: &ScopeSet,
now_ms: i64,
req: &RememberRequest,
) -> Result<RememberPlan, ComposeError> {
let ty = req.validate().map_err(ComposeError::Invalid)?;
let memory_props_result =
memory_props(&req.content, req.props.as_ref()).map_err(ComposeError::Invalid)?;
let mut supersedes_ids = BTreeSet::new();
for raw in &req.supersedes {
let id: NodeId = raw
.parse()
.map_err(|e| ComposeError::Invalid(format!("invalid node id {raw:?}: {e}")))?;
supersedes_ids.insert(id);
}
let existing = existing_memory(db, write_scope, &req.content)?;
let deduplicated = existing.is_some() && !supersedes_ids.contains(&existing.unwrap());
let memory_id = if deduplicated {
existing.unwrap()
} else {
NodeId::new()
};
let (supersede_ops, superseded) = plan_supersede(db, write_scope, &req.supersedes, now_ms)?;
struct Resolved {
name: String,
id: NodeId,
created: bool,
op: Option<Op>,
}
let mut seen = BTreeSet::new();
let mut resolved: Vec<Resolved> = Vec::new();
for name in req
.entities
.iter()
.filter(|n| seen.insert(entity_dedup_key(n)))
{
match find_existing_entity(db, lookup, name)? {
Some(node) => resolved.push(Resolved {
name: name.clone(),
id: node.id,
created: false,
op: None,
}),
None => {
let id = NodeId::new();
let props =
merge_required_prop(ENTITY_NAME_PROP, PropValue::Str(name.clone()), None)
.map_err(ComposeError::Invalid)?;
resolved.push(Resolved {
name: name.clone(),
id,
created: true,
op: Some(Op::CreateNode {
id,
scope: write_scope,
label: ENTITY_LABEL.into(),
props,
}),
});
}
}
}
let mut ops: Vec<Op> = Vec::new();
let mut new_memory = None;
if !deduplicated {
ops.push(Op::CreateNode {
id: memory_id,
scope: write_scope,
label: MEMORY_LABEL.into(),
props: memory_props_result,
});
new_memory = Some(req.content.clone());
}
let mut seen_ids = BTreeSet::new();
resolved.retain(|r| seen_ids.insert(r.id));
let already_linked: HashMap<NodeId, EdgeId> = if deduplicated {
let scope_set = scopes_to_scope_set(&[write_scope]);
db.edges_from(&scope_set, memory_id, None, Some(ty.as_str()), true)?
.into_iter()
.map(|e| (e.to, e.id))
.collect()
} else {
HashMap::new()
};
let mut entities = Vec::with_capacity(resolved.len());
let mut edge_ids = Vec::with_capacity(resolved.len());
let mut new_entities = Vec::new();
for r in resolved {
if let Some(op) = r.op {
new_entities.push((r.id, r.name.clone()));
ops.push(op);
}
let edge_id = match already_linked.get(&r.id) {
Some(existing_edge) => existing_edge.to_string(),
None => {
let id = EdgeId::new();
ops.push(Op::CreateEdge {
id,
scope: write_scope,
ty: ty.clone().into(),
from: memory_id,
to: r.id,
props: Props::new(),
valid_from: None,
});
id.to_string()
}
};
edge_ids.push(edge_id);
entities.push(PlannedEntity {
name: r.name,
id: r.id,
created: r.created,
});
}
ops.extend(supersede_ops);
Ok(RememberPlan {
ops,
memory_id,
deduplicated,
new_memory,
new_entities,
entities,
edge_ids,
superseded,
})
}