use velesdb_core::agent::AgentMemory;
use velesdb_core::collection::graph::GraphEdge;
use velesdb_core::Database;
use super::enumeration::{enumerate_by_cursor, AGENT_COLLECTIONS};
#[derive(Debug, Clone, Copy)]
enum Subsystem {
Semantic,
Episodic,
Procedural,
}
fn subsystem_of(collection: &str) -> Result<Subsystem, crate::MemoryError> {
if collection == AGENT_COLLECTIONS[0] {
Ok(Subsystem::Semantic)
} else if collection == AGENT_COLLECTIONS[1] {
Ok(Subsystem::Episodic)
} else if collection == AGENT_COLLECTIONS[2] {
Ok(Subsystem::Procedural)
} else {
Err(velesdb_core::Error::Query(format!(
"`{collection}` is not an agent memory collection; edges are exported \
per subsystem, and the subsystems are {AGENT_COLLECTIONS:?}"
))
.into())
}
}
#[derive(Debug, Clone, Copy)]
enum Direction {
Outgoing,
Incoming,
}
fn edges_at(
memory: &AgentMemory,
subsystem: Subsystem,
direction: Direction,
id: u64,
) -> Result<Vec<GraphEdge>, crate::MemoryError> {
let result = match (subsystem, direction) {
(Subsystem::Semantic, Direction::Outgoing) => memory.semantic().relations(id),
(Subsystem::Semantic, Direction::Incoming) => memory.semantic().incoming_relations(id),
(Subsystem::Episodic, Direction::Outgoing) => memory.episodic().relations(id),
(Subsystem::Episodic, Direction::Incoming) => memory.episodic().incoming_relations(id),
(Subsystem::Procedural, Direction::Outgoing) => memory.procedural().relations(id),
(Subsystem::Procedural, Direction::Incoming) => memory.procedural().incoming_relations(id),
};
result.map_err(crate::MemoryError::from)
}
fn require_derived_id(edge: GraphEdge) -> Result<GraphEdge, crate::MemoryError> {
let derived = velesdb_core::hash_edge_id(edge.source(), edge.target(), edge.label());
if edge.id() == derived {
return Ok(edge);
}
Err(velesdb_core::Error::Query(format!(
"edge {} ({} -{}-> {}) carries an id its triple does not derive (expected \
{derived}); reinserting it would rederive the expected id and give one \
logical edge two identities, so the export stops here rather than \
renumbering it",
edge.id(),
edge.source(),
edge.label(),
edge.target(),
))
.into())
}
fn collect(
memory: &AgentMemory,
subsystem: Subsystem,
ids: &[u64],
direction: Direction,
) -> Result<Vec<GraphEdge>, crate::MemoryError> {
let mut out: Vec<GraphEdge> = Vec::new();
let mut seen: std::collections::HashSet<u64> = std::collections::HashSet::new();
for &id in ids {
for edge in edges_at(memory, subsystem, direction, id)? {
let edge = require_derived_id(edge)?;
if !seen.insert(edge.id()) {
return Err(velesdb_core::Error::Query(format!(
"edge {} was reported twice by a single {direction:?} walk of \
one collection; edge ids are unique per collection, so the \
edge index is inconsistent and no export from it can be \
trusted",
edge.id(),
))
.into());
}
out.push(edge);
}
}
Ok(out)
}
fn live_ids(db: &Database, collection: &str, batch: usize) -> Result<Vec<u64>, crate::MemoryError> {
Ok(enumerate_by_cursor(db, collection, batch)?
.into_iter()
.map(|fact| fact.id)
.collect())
}
pub fn export_edges(
memory: &AgentMemory,
db: &Database,
collection: &str,
batch: usize,
) -> Result<Vec<GraphEdge>, crate::MemoryError> {
let subsystem = subsystem_of(collection)?;
collect(
memory,
subsystem,
&live_ids(db, collection, batch)?,
Direction::Outgoing,
)
}
pub fn export_edges_verified(
memory: &AgentMemory,
db: &Database,
collection: &str,
batch: usize,
) -> Result<Vec<GraphEdge>, crate::MemoryError> {
let subsystem = subsystem_of(collection)?;
let ids = live_ids(db, collection, batch)?;
let exported = collect(memory, subsystem, &ids, Direction::Outgoing)?;
let crossed = collect(memory, subsystem, &ids, Direction::Incoming)?;
let (left, right) = (comparable(&exported), comparable(&crossed));
if left != right {
return Err(velesdb_core::Error::Query(format!(
"the outgoing and incoming edge indexes disagree over the same {} live \
facts: {} tuples out, {} tuples in, {} present in only one of them. \
An export cannot be called lossless while the two indexes describe \
different graphs",
ids.len(),
left.len(),
right.len(),
left.symmetric_difference(&right).count(),
))
.into());
}
Ok(exported)
}
pub(super) fn same_edge_tuples(left: &[GraphEdge], right: &[GraphEdge]) -> Result<(), String> {
let (left, right) = (comparable(left), comparable(right));
if left == right {
return Ok(());
}
Err(format!(
"{} tuples on one side, {} on the other, {} present in only one of them",
left.len(),
right.len(),
left.symmetric_difference(&right).count(),
))
}
pub(super) type CanonicalEdge = (u64, u64, u64, String, String);
pub(super) fn canonical_edge(edge: &GraphEdge) -> CanonicalEdge {
let properties: std::collections::BTreeMap<_, _> = edge.properties().iter().collect();
(
edge.id(),
edge.source(),
edge.target(),
edge.label().to_owned(),
serde_json::to_string(&properties)
.unwrap_or_else(|err| format!("<properties failed to render: {err}>")),
)
}
fn comparable(edges: &[GraphEdge]) -> std::collections::BTreeSet<CanonicalEdge> {
edges.iter().map(canonical_edge).collect()
}
pub fn cross_check_edges(
memory: &AgentMemory,
db: &Database,
collection: &str,
batch: usize,
) -> Result<Vec<GraphEdge>, crate::MemoryError> {
let subsystem = subsystem_of(collection)?;
collect(
memory,
subsystem,
&live_ids(db, collection, batch)?,
Direction::Incoming,
)
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct EdgeReinsertion {
pub inserted: u64,
}
pub fn reinsert_edges(
memory: &AgentMemory,
collection: &str,
edges: &[GraphEdge],
) -> Result<EdgeReinsertion, crate::MemoryError> {
let subsystem = subsystem_of(collection)?;
let mut inserted = 0u64;
for edge in edges {
let properties: serde_json::Map<String, serde_json::Value> = edge
.properties()
.iter()
.map(|(key, value)| (key.clone(), value.clone()))
.collect();
let properties = if properties.is_empty() {
None
} else {
Some(properties)
};
let returned = relate_on(
memory,
subsystem,
(edge.source(), edge.target()),
edge.label(),
properties.as_ref(),
)?;
if returned != edge.id() {
return Err(velesdb_core::Error::Query(format!(
"edge {} ({} -{}-> {}) was reinserted under id {returned}; the \
destination derives edge ids differently from the source, so no \
edge in this export can be trusted to keep its identity",
edge.id(),
edge.source(),
edge.label(),
edge.target(),
))
.into());
}
inserted += 1;
}
Ok(EdgeReinsertion { inserted })
}
fn relate_on(
memory: &AgentMemory,
subsystem: Subsystem,
endpoints: (u64, u64),
label: &str,
properties: Option<&serde_json::Map<String, serde_json::Value>>,
) -> Result<u64, crate::MemoryError> {
let (from, to) = endpoints;
let result = match subsystem {
Subsystem::Semantic => memory.semantic().relate(from, to, label, properties),
Subsystem::Episodic => memory.episodic().relate(from, to, label, properties),
Subsystem::Procedural => memory.procedural().relate(from, to, label, properties),
};
result.map_err(crate::MemoryError::from)
}