use std::collections::{HashMap, HashSet};
use petgraph::graph::NodeIndex;
use crate::datatypes::{DataFrame, Value};
use crate::graph::mutation::maintain::{add_connections, add_nodes, detach_delete_nodes};
use crate::graph::schema::{DirGraph, InternedKey};
use crate::graph::storage::{GraphRead, GraphWrite};
use crate::graph::wal::{MutationOp, WalFrame};
type NodeKey = (String, Value);
type EdgeKey = (String, String, Value, String, Value);
type UpsertRow = (Value, Value, HashMap<String, Value>);
enum NodeNet {
Upsert {
title: Value,
props: Vec<(String, Value)>,
},
Remove,
}
enum EdgeNet {
Upsert { props: Vec<(String, Value)> },
Remove,
}
type LabelNet = HashMap<NodeKey, Vec<String>>;
pub fn apply_frames(
graph: &mut DirGraph,
frames: &[WalFrame],
after_lsn: u64,
) -> Result<u64, String> {
graph
.prepare_disk_mutation()
.map_err(|e| format!("disk mutation lease failed: {e}"))?;
let mut nodes: HashMap<NodeKey, NodeNet> = HashMap::new();
let mut edges: HashMap<EdgeKey, EdgeNet> = HashMap::new();
let mut labels: LabelNet = HashMap::new();
let mut max_lsn = after_lsn;
let mut any = false;
for frame in frames {
if frame.lsn <= after_lsn {
continue;
}
any = true;
max_lsn = max_lsn.max(frame.lsn);
for op in &frame.ops {
match op {
MutationOp::UpsertNode {
node_type,
id,
title,
properties,
} => {
nodes.insert(
(node_type.clone(), id.clone()),
NodeNet::Upsert {
title: title.clone(),
props: properties.clone(),
},
);
}
MutationOp::RemoveNode { node_type, id } => {
nodes.insert((node_type.clone(), id.clone()), NodeNet::Remove);
}
MutationOp::UpsertEdge {
conn_type,
src_type,
src_id,
tgt_type,
tgt_id,
properties,
} => {
edges.insert(
(
conn_type.clone(),
src_type.clone(),
src_id.clone(),
tgt_type.clone(),
tgt_id.clone(),
),
EdgeNet::Upsert {
props: properties.clone(),
},
);
}
MutationOp::RemoveEdge {
conn_type,
src_type,
src_id,
tgt_type,
tgt_id,
} => {
edges.insert(
(
conn_type.clone(),
src_type.clone(),
src_id.clone(),
tgt_type.clone(),
tgt_id.clone(),
),
EdgeNet::Remove,
);
}
MutationOp::SetNodeLabels {
node_type,
id,
labels: set,
} => {
labels.insert((node_type.clone(), id.clone()), set.clone());
}
}
}
}
if any {
apply_net(graph, nodes, edges, labels)?;
}
Ok(max_lsn)
}
fn apply_net(
graph: &mut DirGraph,
nodes: HashMap<NodeKey, NodeNet>,
edges: HashMap<EdgeKey, EdgeNet>,
labels: LabelNet,
) -> Result<(), String> {
let removed_nodes: HashSet<NodeKey> = nodes
.iter()
.filter(|(_, v)| matches!(v, NodeNet::Remove))
.map(|(k, _)| k.clone())
.collect();
apply_node_upserts(graph, &nodes)?;
apply_label_sets(graph, &labels, &removed_nodes);
apply_edge_upserts(graph, &edges, &removed_nodes)?;
apply_edge_removes(graph, &edges);
apply_node_removes(graph, &nodes);
Ok(())
}
fn apply_node_upserts(
graph: &mut DirGraph,
nodes: &HashMap<NodeKey, NodeNet>,
) -> Result<(), String> {
let mut node_groups: HashMap<&str, NodeRows> = HashMap::new();
for ((node_type, id), net) in nodes {
if let NodeNet::Upsert { title, props } = net {
let g = node_groups.entry(node_type.as_str()).or_default();
for (k, _) in props {
g.note_column(k);
}
g.rows
.push((id.clone(), title.clone(), props.iter().cloned().collect()));
}
}
for (node_type, group) in node_groups {
let (framed, exact) = split_faithful_columns(&group.columns, &group.rows);
if fixed_columns_are_faithful(&group.rows) {
upsert_node_rows(graph, node_type, &framed, &exact, &group.rows)?;
} else {
for part in partition_by_fixed_shapes(&group.rows) {
upsert_node_rows(graph, node_type, &framed, &exact, &part)?;
}
}
}
Ok(())
}
fn upsert_node_rows(
graph: &mut DirGraph,
node_type: &str,
framed: &[String],
exact: &[String],
rows: &[UpsertRow],
) -> Result<(), String> {
declare_exact_node_columns(graph, node_type, exact, rows);
let df = build_dataframe(&["id", "title"], framed, rows)?;
add_nodes(
graph,
df,
node_type.to_string(),
"id".to_string(),
Some("title".to_string()),
Some("replace".to_string()),
)?;
apply_exact_node_props(graph, node_type, exact, rows);
Ok(())
}
fn declare_exact_node_columns(
graph: &mut DirGraph,
node_type: &str,
exact: &[String],
rows: &[UpsertRow],
) {
if exact.is_empty() {
return;
}
let mut declared: HashMap<String, String> = HashMap::new();
for col in exact {
declared.insert(
col.clone(),
declared_type_name(rows.iter().filter_map(|(_, _, p)| p.get(col))),
);
}
graph.upsert_node_type_metadata(node_type, declared);
let keys: Vec<_> = exact
.iter()
.map(|col| graph.interner.get_or_intern(col))
.collect();
graph.ensure_type_schema_keys(node_type, &keys);
}
fn apply_exact_node_props(
graph: &mut DirGraph,
node_type: &str,
exact: &[String],
rows: &[UpsertRow],
) {
if exact.is_empty() {
return;
}
for (id, _, props) in rows {
let Some(idx) = graph.lookup_by_id(node_type, id) else {
continue;
};
for col in exact {
match props.get(col) {
None | Some(Value::Null) => continue,
Some(value) => {
let key = graph.interner.get_or_intern(col);
graph.ensure_type_schema_keys(node_type, &[key]);
GraphWrite::set_node_property(&mut graph.graph, idx, key, value.clone());
}
}
}
}
}
fn apply_label_sets(graph: &mut DirGraph, labels: &LabelNet, removed_nodes: &HashSet<NodeKey>) {
for (key @ (node_type, id), target) in labels {
if removed_nodes.contains(key) {
continue;
}
let Some(idx) = graph.lookup_by_id(node_type, id) else {
continue;
};
for stale in graph.secondary_label_names(idx) {
if !target.contains(&stale) {
let key = graph.interner.get_or_intern(&stale);
let _ = graph.remove_node_label(idx, key);
}
}
for label in target {
let key = graph.interner.get_or_intern(label);
graph.add_node_label(idx, key);
}
}
}
fn apply_edge_upserts(
graph: &mut DirGraph,
edges: &HashMap<EdgeKey, EdgeNet>,
removed_nodes: &HashSet<NodeKey>,
) -> Result<(), String> {
let mut edge_groups: HashMap<(&str, &str, &str), EdgeRows> = HashMap::new();
for ((conn, src_type, src_id, tgt_type, tgt_id), net) in edges {
if let EdgeNet::Upsert { props } = net {
if removed_nodes.contains(&(src_type.clone(), src_id.clone()))
|| removed_nodes.contains(&(tgt_type.clone(), tgt_id.clone()))
{
continue;
}
let g = edge_groups
.entry((conn.as_str(), src_type.as_str(), tgt_type.as_str()))
.or_default();
for (k, _) in props {
g.note_column(k);
}
g.rows.push((
src_id.clone(),
tgt_id.clone(),
props.iter().cloned().collect(),
));
}
}
for ((conn, src_type, tgt_type), group) in edge_groups {
let (framed, exact) = split_faithful_columns(&group.columns, &group.rows);
let key = EdgeGroup {
conn,
src_type,
tgt_type,
};
if fixed_columns_are_faithful(&group.rows) {
upsert_edge_rows(graph, key, &framed, &exact, &group.rows)?;
} else {
for part in partition_by_fixed_shapes(&group.rows) {
upsert_edge_rows(graph, key, &framed, &exact, &part)?;
}
}
}
Ok(())
}
#[derive(Clone, Copy)]
struct EdgeGroup<'a> {
conn: &'a str,
src_type: &'a str,
tgt_type: &'a str,
}
fn upsert_edge_rows(
graph: &mut DirGraph,
group: EdgeGroup<'_>,
framed: &[String],
exact: &[String],
rows: &[UpsertRow],
) -> Result<(), String> {
let df = build_dataframe(&["src_id", "tgt_id"], framed, rows)?;
add_connections(
graph,
df,
group.conn.to_string(),
group.src_type.to_string(),
"src_id".to_string(),
group.tgt_type.to_string(),
"tgt_id".to_string(),
None,
None,
Some("replace".to_string()),
)?;
apply_exact_edge_props(graph, group, exact, rows);
Ok(())
}
fn apply_exact_edge_props(
graph: &mut DirGraph,
group: EdgeGroup<'_>,
exact: &[String],
rows: &[UpsertRow],
) {
let EdgeGroup {
conn,
src_type,
tgt_type,
} = group;
if exact.is_empty() {
return;
}
let mut declared: HashMap<String, String> = HashMap::new();
for col in exact {
declared.insert(
col.clone(),
declared_type_name(rows.iter().filter_map(|(_, _, p)| p.get(col))),
);
}
graph.upsert_connection_type_metadata(conn, src_type, tgt_type, declared);
let conn_key = InternedKey::from_str(conn);
for (src_id, tgt_id, props) in rows {
let (Some(src), Some(tgt)) = (
graph.lookup_by_id(src_type, src_id),
graph.lookup_by_id(tgt_type, tgt_id),
) else {
continue;
};
let Some(eidx) = graph
.graph
.edges_connecting(src, tgt)
.find(|er| er.weight().connection_type == conn_key)
.map(|er| er.id())
else {
continue;
};
for col in exact {
match props.get(col) {
None | Some(Value::Null) => continue,
Some(value) => {
let key = graph.interner.get_or_intern(col);
if let Some(edge) = GraphWrite::edge_weight_mut(&mut graph.graph, eidx) {
match edge.properties.iter_mut().find(|(ek, _)| *ek == key) {
Some((_, existing)) => *existing = value.clone(),
None => edge.properties.push((key, value.clone())),
}
}
}
}
}
}
}
fn split_faithful_columns(columns: &[String], rows: &[UpsertRow]) -> (Vec<String>, Vec<String>) {
let mut framed = Vec::with_capacity(columns.len());
let mut exact = Vec::new();
for col in columns {
if column_is_faithful(rows.iter().filter_map(|(_, _, p)| p.get(col))) {
framed.push(col.clone());
} else {
exact.push(col.clone());
}
}
(framed, exact)
}
fn column_is_faithful<'a>(values: impl Iterator<Item = &'a Value>) -> bool {
let mut shape: Option<FrameShape> = None;
for value in values {
if matches!(value, Value::Null) {
continue;
}
match (frame_shape(value), shape) {
(None, _) => return false,
(Some(s), None) => shape = Some(s),
(Some(s), Some(seen)) if s == seen => {}
(Some(_), Some(_)) => return false,
}
}
true
}
fn fixed_columns_are_faithful(rows: &[UpsertRow]) -> bool {
column_is_faithful(rows.iter().map(|(a, _, _)| a))
&& column_is_faithful(rows.iter().map(|(_, b, _)| b))
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum FixedKind {
Null,
Shaped(FrameShape),
Shapeless,
}
fn fixed_kind(value: &Value) -> FixedKind {
match value {
Value::Null => FixedKind::Null,
other => match frame_shape(other) {
Some(shape) => FixedKind::Shaped(shape),
None => FixedKind::Shapeless,
},
}
}
fn partition_by_fixed_shapes(rows: &[UpsertRow]) -> Vec<Vec<UpsertRow>> {
type PartKey = (FixedKind, FixedKind);
let fits = |part: FixedKind, row: FixedKind| {
part == FixedKind::Null || row == FixedKind::Null || part == row
};
let mut parts: Vec<(PartKey, Vec<UpsertRow>)> = Vec::new();
for row in rows {
let key = (fixed_kind(&row.0), fixed_kind(&row.1));
let slot = parts
.iter_mut()
.find(|(k, _)| fits(k.0, key.0) && fits(k.1, key.1));
match slot {
Some((k, part)) => {
if k.0 == FixedKind::Null {
k.0 = key.0;
}
if k.1 == FixedKind::Null {
k.1 = key.1;
}
part.push(row.clone());
}
None => parts.push((key, vec![row.clone()])),
}
}
parts.into_iter().map(|(_, part)| part).collect()
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum FrameShape {
UniqueId,
Int64,
Float64,
String,
Boolean,
DateTime,
Timestamp,
List,
Map,
}
fn frame_shape(value: &Value) -> Option<FrameShape> {
Some(match value {
Value::UniqueId(_) => FrameShape::UniqueId,
Value::Int64(_) => FrameShape::Int64,
Value::Float64(_) => FrameShape::Float64,
Value::String(_) => FrameShape::String,
Value::Boolean(_) => FrameShape::Boolean,
Value::DateTime(_) => FrameShape::DateTime,
Value::Timestamp(_) => FrameShape::Timestamp,
Value::List(_) => FrameShape::List,
Value::Map(_) => FrameShape::Map,
_ => return None,
})
}
fn declared_type_name<'a>(values: impl Iterator<Item = &'a Value>) -> String {
let mut seen: Option<&'static str> = None;
for value in values {
if matches!(value, Value::Null) {
continue;
}
match seen {
None => seen = Some(value.type_name()),
Some(name) if name == value.type_name() => {}
Some(_) => return "mixed".to_string(),
}
}
seen.unwrap_or("mixed").to_string()
}
fn apply_edge_removes(graph: &mut DirGraph, edges: &HashMap<EdgeKey, EdgeNet>) {
let mut removed_edges = 0usize;
for ((conn, src_type, src_id, tgt_type, tgt_id), net) in edges {
if !matches!(net, EdgeNet::Remove) {
continue;
}
let (Some(src), Some(tgt)) = (
graph.lookup_by_id(src_type, src_id),
graph.lookup_by_id(tgt_type, tgt_id),
) else {
continue;
};
let conn_key = InternedKey::from_str(conn);
let eidx = graph
.graph
.edges_connecting(src, tgt)
.find(|er| er.weight().connection_type == conn_key)
.map(|er| er.id());
if let Some(eidx) = eidx {
GraphWrite::remove_edge(&mut graph.graph, eidx);
removed_edges += 1;
}
}
if removed_edges > 0 {
graph.invalidate_edge_type_counts_cache();
graph.connection_types.clear();
}
}
fn apply_node_removes(graph: &mut DirGraph, nodes: &HashMap<NodeKey, NodeNet>) {
let mut to_delete: HashSet<NodeIndex> = HashSet::new();
for ((node_type, id), net) in nodes {
if matches!(net, NodeNet::Remove) {
if let Some(idx) = graph.lookup_by_id(node_type, id) {
to_delete.insert(idx);
}
}
}
if !to_delete.is_empty() {
detach_delete_nodes(graph, &to_delete);
}
}
#[derive(Default)]
struct NodeRows {
columns: Vec<String>,
seen: std::collections::HashSet<String>,
rows: Vec<UpsertRow>,
}
#[derive(Default)]
struct EdgeRows {
columns: Vec<String>,
seen: std::collections::HashSet<String>,
rows: Vec<UpsertRow>,
}
impl NodeRows {
fn note_column(&mut self, name: &str) {
if self.seen.insert(name.to_string()) {
self.columns.push(name.to_string());
}
}
}
impl EdgeRows {
fn note_column(&mut self, name: &str) {
if self.seen.insert(name.to_string()) {
self.columns.push(name.to_string());
}
}
}
fn build_dataframe(
fixed: &[&str],
prop_columns: &[String],
rows: &[UpsertRow],
) -> Result<DataFrame, String> {
let mut columns: Vec<String> = fixed.iter().map(|s| s.to_string()).collect();
columns.extend(prop_columns.iter().cloned());
let out_rows: Vec<Vec<Value>> = rows
.iter()
.map(|(a, b, props)| {
let mut row = Vec::with_capacity(columns.len());
row.push(a.clone());
row.push(b.clone());
for col in prop_columns {
row.push(props.get(col).cloned().unwrap_or(Value::Null));
}
row
})
.collect();
DataFrame::from_cypher_rows(columns, out_rows)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::graph::storage::GraphRead;
fn frame(lsn: u64, ops: Vec<MutationOp>) -> WalFrame {
WalFrame { lsn, ops }
}
fn upsert_node(id: i64, title: &str, props: Vec<(&str, Value)>) -> MutationOp {
MutationOp::UpsertNode {
node_type: "Person".into(),
id: Value::Int64(id),
title: Value::String(title.into()),
properties: props.into_iter().map(|(k, v)| (k.to_string(), v)).collect(),
}
}
fn knows(src: i64, tgt: i64) -> MutationOp {
MutationOp::UpsertEdge {
conn_type: "KNOWS".into(),
src_type: "Person".into(),
src_id: Value::Int64(src),
tgt_type: "Person".into(),
tgt_id: Value::Int64(tgt),
properties: vec![],
}
}
fn prop(g: &mut DirGraph, id: i64, key: &str) -> Option<Value> {
let idx = g.lookup_by_id("Person", &Value::Int64(id))?;
g.graph
.node_view(idx)
.and_then(|n| n.get_field_ref(key).map(|c| c.into_owned()))
}
#[test]
fn replays_upserts_and_edge() {
let mut g = DirGraph::new();
let frames = vec![frame(
1,
vec![
upsert_node(1, "Alice", vec![("age", Value::Int64(30))]),
upsert_node(2, "Bob", vec![]),
knows(1, 2),
],
)];
let max = apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(max, 1);
assert_eq!(g.graph.node_count(), 2);
assert_eq!(g.graph.edge_count(), 1);
assert_eq!(prop(&mut g, 1, "age"), Some(Value::Int64(30)));
}
#[test]
fn later_upsert_replaces_properties() {
let mut g = DirGraph::new();
let frames = vec![
frame(
1,
vec![upsert_node(1, "Alice", vec![("age", Value::Int64(30))])],
),
frame(
2,
vec![upsert_node(1, "Alice", vec![("age", Value::Int64(41))])],
),
];
apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(
g.graph.node_count(),
1,
"same (type,id) is upserted, not duplicated"
);
assert_eq!(prop(&mut g, 1, "age"), Some(Value::Int64(41)));
}
#[test]
fn remove_node_deletes_it_and_its_edges() {
let mut g = DirGraph::new();
let frames = vec![
frame(
1,
vec![
upsert_node(1, "Alice", vec![]),
upsert_node(2, "Bob", vec![]),
knows(1, 2),
],
),
frame(
2,
vec![MutationOp::RemoveNode {
node_type: "Person".into(),
id: Value::Int64(2),
}],
),
];
apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(g.graph.node_count(), 1);
assert_eq!(
g.graph.edge_count(),
0,
"incident edge removed with the node"
);
assert!(g.lookup_by_id("Person", &Value::Int64(2)).is_none());
}
#[test]
fn replayed_node_removal_prunes_the_embedding_store() {
let mut g = DirGraph::new();
apply_frames(
&mut g,
&[frame(
1,
vec![
upsert_node(1, "Alice", vec![]),
upsert_node(2, "Bob", vec![]),
],
)],
0,
)
.unwrap();
let report = crate::graph::embeddings::set_embeddings(
&mut g,
"Person",
"name",
None,
[
(Value::Int64(1), vec![1.0f32, 0.0]),
(Value::Int64(2), vec![0.0, 1.0]),
],
)
.expect("seed embeddings");
assert_eq!(report.embeddings_stored, 2);
let doomed = g
.lookup_by_id("Person", &Value::Int64(2))
.expect("Bob is present");
apply_frames(
&mut g,
&[frame(
2,
vec![MutationOp::RemoveNode {
node_type: "Person".into(),
id: Value::Int64(2),
}],
)],
1,
)
.unwrap();
let store = g
.embeddings
.get(&("Person".to_string(), "name_emb".to_string()))
.expect("store");
assert_eq!(store.len(), 1);
assert_eq!(store.get_embedding(doomed.index()), None);
assert_eq!(store.validate_shape(), Ok(()));
}
#[test]
fn remove_edge_keeps_endpoints() {
let mut g = DirGraph::new();
let frames = vec![
frame(
1,
vec![
upsert_node(1, "Alice", vec![]),
upsert_node(2, "Bob", vec![]),
knows(1, 2),
],
),
frame(
2,
vec![MutationOp::RemoveEdge {
conn_type: "KNOWS".into(),
src_type: "Person".into(),
src_id: Value::Int64(1),
tgt_type: "Person".into(),
tgt_id: Value::Int64(2),
}],
),
];
apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(g.graph.node_count(), 2, "endpoints survive an edge remove");
assert_eq!(g.graph.edge_count(), 0);
}
#[test]
fn frames_at_or_below_checkpoint_are_skipped() {
let mut g = DirGraph::new();
let frames = vec![
frame(1, vec![upsert_node(1, "Old", vec![])]),
frame(2, vec![upsert_node(2, "New", vec![])]),
];
let max = apply_frames(&mut g, &frames, 1).unwrap();
assert_eq!(max, 2);
assert!(g.lookup_by_id("Person", &Value::Int64(1)).is_none());
assert!(g.lookup_by_id("Person", &Value::Int64(2)).is_some());
}
fn labels_of(g: &mut DirGraph, id: i64) -> Vec<String> {
let idx = g
.lookup_by_id("Person", &Value::Int64(id))
.expect("node must exist");
g.node_labels(idx)
.into_iter()
.map(|k| g.interner.resolve(k).to_string())
.collect()
}
fn set_labels(id: i64, labels: &[&str]) -> MutationOp {
MutationOp::SetNodeLabels {
node_type: "Person".into(),
id: Value::Int64(id),
labels: labels.iter().map(|s| s.to_string()).collect(),
}
}
#[test]
fn replay_restores_secondary_labels_in_exact_order() {
let mut g = DirGraph::new();
let frames = vec![frame(
1,
vec![
upsert_node(1, "Alice", vec![("age", Value::Int64(30))]),
set_labels(1, &["Manager", "Employee"]),
],
)];
apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(
labels_of(&mut g, 1),
vec!["Person", "Employee", "Manager"],
"primary first, then secondaries sorted by name"
);
assert_eq!(prop(&mut g, 1, "age"), Some(Value::Int64(30)));
assert!(g.has_secondary_labels, "fast-skip flag must be set");
assert_eq!(g.nodes_with_label("Employee").len(), 1);
}
#[test]
fn replay_removes_labels_the_log_dropped() {
let mut g = DirGraph::new();
apply_frames(
&mut g,
&[frame(
1,
vec![upsert_node(1, "Alice", vec![]), set_labels(1, &["A", "B"])],
)],
0,
)
.unwrap();
assert_eq!(labels_of(&mut g, 1), vec!["Person", "A", "B"]);
apply_frames(&mut g, &[frame(2, vec![set_labels(1, &["B"])])], 1).unwrap();
assert_eq!(labels_of(&mut g, 1), vec!["Person", "B"]);
assert!(
g.nodes_with_label("A").is_empty(),
"the dropped label must leave no index residue"
);
}
#[test]
fn replay_to_an_empty_label_set_clears_the_flag() {
let mut g = DirGraph::new();
apply_frames(
&mut g,
&[
frame(
1,
vec![upsert_node(1, "Alice", vec![]), set_labels(1, &["A"])],
),
frame(2, vec![set_labels(1, &[])]),
],
0,
)
.unwrap();
assert_eq!(labels_of(&mut g, 1), vec!["Person"]);
assert!(!g.has_secondary_labels);
}
#[test]
fn property_upsert_does_not_clobber_labels() {
for reversed in [false, true] {
let mut ops = vec![
upsert_node(1, "Alice", vec![]),
set_labels(1, &["Employee"]),
upsert_node(1, "Alice", vec![("age", Value::Int64(41))]),
];
if reversed {
ops.swap(1, 2);
}
let mut g = DirGraph::new();
apply_frames(&mut g, &[frame(1, ops)], 0).unwrap();
assert_eq!(
labels_of(&mut g, 1),
vec!["Person", "Employee"],
"{reversed}"
);
assert_eq!(prop(&mut g, 1, "age"), Some(Value::Int64(41)), "{reversed}");
}
}
#[test]
fn label_set_for_a_removed_node_is_skipped() {
let mut g = DirGraph::new();
let frames = vec![frame(
1,
vec![
upsert_node(1, "Alice", vec![]),
set_labels(1, &["Employee"]),
MutationOp::RemoveNode {
node_type: "Person".into(),
id: Value::Int64(1),
},
],
)];
apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(g.graph.node_count(), 0);
assert!(g.nodes_with_label("Employee").is_empty());
}
#[test]
fn replaying_labels_twice_is_idempotent() {
let frames = vec![frame(
1,
vec![
upsert_node(1, "Alice", vec![]),
set_labels(1, &["Employee", "Manager"]),
],
)];
let mut g = DirGraph::new();
apply_frames(&mut g, &frames, 0).unwrap();
apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(labels_of(&mut g, 1), vec!["Person", "Employee", "Manager"]);
assert_eq!(
g.nodes_with_label("Employee").len(),
1,
"no duplicate bucket entry"
);
}
#[test]
fn replays_onto_a_mapped_graph() {
use crate::graph::storage::mode::{new_dir_graph_in_mode, StorageMode};
let mut g = new_dir_graph_in_mode(StorageMode::Mapped, None).unwrap();
assert!(g.graph.is_mapped(), "fixture must really be mapped");
let frames = vec![
frame(
1,
vec![
upsert_node(1, "Alice", vec![("age", Value::Int64(30))]),
upsert_node(2, "Bob", vec![]),
knows(1, 2),
set_labels(1, &["Employee"]),
],
),
frame(
2,
vec![MutationOp::RemoveNode {
node_type: "Person".into(),
id: Value::Int64(2),
}],
),
];
apply_frames(&mut g, &frames, 0).unwrap();
assert!(g.graph.is_mapped(), "replay must not switch the backend");
assert_eq!(g.graph.node_count(), 1);
assert_eq!(g.graph.edge_count(), 0, "edge went with the removed node");
assert_eq!(labels_of(&mut g, 1), vec!["Person", "Employee"]);
assert_eq!(prop(&mut g, 1, "age"), Some(Value::Int64(30)));
}
#[test]
fn mixed_typed_property_keeps_every_value_type() {
use chrono::NaiveDate;
let date = NaiveDate::from_ymd_opt(2020, 1, 2).unwrap();
let cases: Vec<(i64, Value)> = vec![
(1, Value::Int64(1)),
(2, Value::String("two".into())),
(3, Value::Float64(3.5)),
(4, Value::Boolean(true)),
(5, Value::DateTime(date)),
];
let mut g = DirGraph::new();
let frames: Vec<WalFrame> = cases
.iter()
.enumerate()
.map(|(i, (id, v))| {
frame(
i as u64 + 1,
vec![upsert_node(*id, "n", vec![("mixedish", v.clone())])],
)
})
.collect();
apply_frames(&mut g, &frames, 0).unwrap();
for (id, expected) in &cases {
assert_eq!(
prop(&mut g, *id, "mixedish").as_ref(),
Some(expected),
"node {id}"
);
}
}
#[test]
fn single_typed_properties_keep_their_types_through_the_frame() {
use chrono::NaiveDate;
let props = vec![
("i", Value::Int64(7)),
("f", Value::Float64(0.5)),
("s", Value::String("x".into())),
("b", Value::Boolean(true)),
(
"d",
Value::DateTime(NaiveDate::from_ymd_opt(2020, 1, 2).unwrap()),
),
("l", Value::List(vec![Value::Int64(1), Value::Int64(2)])),
];
let mut g = DirGraph::new();
apply_frames(
&mut g,
&[frame(1, vec![upsert_node(1, "a", props.clone())])],
0,
)
.unwrap();
for (key, expected) in props {
assert_eq!(prop(&mut g, 1, key).as_ref(), Some(&expected), "{key}");
}
let meta = g.get_node_type_metadata("Person").cloned().unwrap();
assert!(
!meta.values().any(|t| t == "mixed"),
"single-typed columns must stay framed: {meta:?}"
);
}
#[test]
fn int_and_float_under_one_property_do_not_promote() {
let mut g = DirGraph::new();
apply_frames(
&mut g,
&[frame(
1,
vec![
upsert_node(1, "a", vec![("n", Value::Int64(2))]),
upsert_node(2, "b", vec![("n", Value::Float64(2.5))]),
],
)],
0,
)
.unwrap();
assert_eq!(prop(&mut g, 1, "n"), Some(Value::Int64(2)));
assert_eq!(prop(&mut g, 2, "n"), Some(Value::Float64(2.5)));
}
#[test]
fn point_property_survives_replay_as_a_point() {
let mut g = DirGraph::new();
apply_frames(
&mut g,
&[frame(
1,
vec![upsert_node(
1,
"a",
vec![(
"loc",
Value::Point {
lat: 59.9,
lon: 10.7,
},
)],
)],
)],
0,
)
.unwrap();
assert_eq!(
prop(&mut g, 1, "loc"),
Some(Value::Point {
lat: 59.9,
lon: 10.7
})
);
}
#[test]
fn mixed_property_folds_with_later_ops_on_the_same_node() {
let mut g = DirGraph::new();
apply_frames(
&mut g,
&[
frame(
1,
vec![
upsert_node(1, "a", vec![("m", Value::Int64(1))]),
upsert_node(2, "b", vec![("m", Value::String("two".into()))]),
],
),
frame(
2,
vec![upsert_node(
1,
"a",
vec![("m", Value::Boolean(false)), ("age", Value::Int64(41))],
)],
),
],
0,
)
.unwrap();
assert_eq!(prop(&mut g, 1, "m"), Some(Value::Boolean(false)));
assert_eq!(prop(&mut g, 1, "age"), Some(Value::Int64(41)));
assert_eq!(prop(&mut g, 2, "m"), Some(Value::String("two".into())));
assert_eq!(
g.get_node_type_metadata("Person").unwrap().get("m"),
Some(&"mixed".to_string()),
"a heterogeneous property is declared 'mixed', not left undeclared"
);
}
#[test]
fn nodes_whose_ids_differ_in_type_keep_their_ids() {
let mut g = DirGraph::new();
let string_id = MutationOp::UpsertNode {
node_type: "Person".into(),
id: Value::String("x".into()),
title: Value::String("b".into()),
properties: vec![("tag".to_string(), Value::String("str-id".into()))],
};
apply_frames(
&mut g,
&[frame(
1,
vec![
upsert_node(1, "a", vec![("tag", Value::String("int-id".into()))]),
string_id,
],
)],
0,
)
.unwrap();
assert_eq!(g.graph.node_count(), 2);
let idx = g
.lookup_by_id("Person", &Value::Int64(1))
.expect("the integer id must still be an integer");
assert_eq!(g.graph.get_node_id(idx), Some(Value::Int64(1)));
let idx = g
.lookup_by_id("Person", &Value::String("x".into()))
.expect("the string id must survive alongside it");
assert_eq!(g.graph.get_node_id(idx), Some(Value::String("x".into())));
}
#[test]
fn nodes_whose_titles_differ_in_type_keep_their_titles() {
let mut g = DirGraph::new();
let numeric_title = MutationOp::UpsertNode {
node_type: "Person".into(),
id: Value::Int64(2),
title: Value::Int64(5),
properties: vec![],
};
apply_frames(
&mut g,
&[frame(1, vec![upsert_node(1, "a", vec![]), numeric_title])],
0,
)
.unwrap();
let title = |g: &mut DirGraph, id: i64| {
let idx = g.lookup_by_id("Person", &Value::Int64(id)).unwrap();
g.graph.get_node_title(idx)
};
assert_eq!(title(&mut g, 1), Some(Value::String("a".into())));
assert_eq!(title(&mut g, 2), Some(Value::Int64(5)));
}
#[test]
fn edges_reach_endpoints_whose_ids_differ_in_type() {
let mut g = DirGraph::new();
let string_node = MutationOp::UpsertNode {
node_type: "Person".into(),
id: Value::String("x".into()),
title: Value::String("b".into()),
properties: vec![],
};
let edge = MutationOp::UpsertEdge {
conn_type: "KNOWS".into(),
src_type: "Person".into(),
src_id: Value::Int64(1),
tgt_type: "Person".into(),
tgt_id: Value::String("x".into()),
properties: vec![],
};
apply_frames(
&mut g,
&[frame(
1,
vec![upsert_node(1, "a", vec![]), string_node, knows(1, 2), edge],
)],
0,
)
.unwrap();
assert_eq!(g.graph.node_count(), 3, "no stub under a stringified id");
assert_eq!(g.graph.edge_count(), 2, "both edges land");
let src = g.lookup_by_id("Person", &Value::Int64(1)).unwrap();
for tgt_id in [Value::Int64(2), Value::String("x".into())] {
let tgt = g
.lookup_by_id("Person", &tgt_id)
.unwrap_or_else(|| panic!("endpoint {tgt_id:?} must exist"));
assert!(
g.graph.find_edge(src, tgt).is_some(),
"the edge to {tgt_id:?} must connect that node"
);
}
}
#[test]
fn mixed_typed_property_keeps_its_types_on_a_mapped_graph() {
use crate::graph::storage::mode::{new_dir_graph_in_mode, StorageMode};
let mut g = new_dir_graph_in_mode(StorageMode::Mapped, None).unwrap();
assert!(g.graph.is_mapped(), "fixture must really be mapped");
apply_frames(
&mut g,
&[frame(
1,
vec![
upsert_node(1, "a", vec![("m", Value::Int64(1))]),
upsert_node(2, "b", vec![("m", Value::String("two".into()))]),
],
)],
0,
)
.unwrap();
assert!(g.graph.is_mapped(), "replay must not switch the backend");
assert_eq!(prop(&mut g, 1, "m"), Some(Value::Int64(1)));
assert_eq!(prop(&mut g, 2, "m"), Some(Value::String("two".into())));
}
#[test]
fn mixed_typed_edge_property_keeps_every_value_type() {
let mut g = DirGraph::new();
let knows_with = |src: i64, tgt: i64, v: Value| MutationOp::UpsertEdge {
conn_type: "KNOWS".into(),
src_type: "Person".into(),
src_id: Value::Int64(src),
tgt_type: "Person".into(),
tgt_id: Value::Int64(tgt),
properties: vec![("w".to_string(), v)],
};
apply_frames(
&mut g,
&[frame(
1,
vec![
upsert_node(1, "a", vec![]),
upsert_node(2, "b", vec![]),
upsert_node(3, "c", vec![]),
knows_with(1, 2, Value::Int64(7)),
knows_with(1, 3, Value::String("heavy".into())),
],
)],
0,
)
.unwrap();
let w = |g: &mut DirGraph, src: i64, tgt: i64| -> Option<Value> {
let s = g.lookup_by_id("Person", &Value::Int64(src))?;
let t = g.lookup_by_id("Person", &Value::Int64(tgt))?;
let e = g.graph.find_edge(s, t)?;
g.graph
.edge_weight(e)?
.properties
.iter()
.find(|(k, _)| *k == InternedKey::from_str("w"))
.map(|(_, v)| v.clone())
};
assert_eq!(w(&mut g, 1, 2), Some(Value::Int64(7)));
assert_eq!(w(&mut g, 1, 3), Some(Value::String("heavy".into())));
}
#[test]
fn replaying_twice_is_idempotent() {
let frames = vec![frame(
1,
vec![
upsert_node(1, "Alice", vec![("age", Value::Int64(30))]),
upsert_node(2, "Bob", vec![]),
knows(1, 2),
],
)];
let mut g = DirGraph::new();
apply_frames(&mut g, &frames, 0).unwrap();
apply_frames(&mut g, &frames, 0).unwrap();
assert_eq!(g.graph.node_count(), 2, "idempotent — no duplicate nodes");
assert_eq!(g.graph.edge_count(), 1, "idempotent — no duplicate edge");
}
}