use std::collections::HashSet;
use petgraph::graph::NodeIndex;
use super::DirGraph;
use crate::datatypes::Value;
use crate::graph::schema::{InternedKey, NodeData};
use crate::graph::storage::column_store::ColumnStore;
use crate::graph::storage::undo::{BucketId, UndoEntry, UndoJournal};
use crate::graph::storage::{GraphRead, GraphWrite};
fn swap_data_scale(a: &mut DirGraph, b: &mut DirGraph) {
std::mem::swap(&mut a.graph, &mut b.graph);
std::mem::swap(&mut a.type_indices, &mut b.type_indices);
std::mem::swap(&mut a.id_indices, &mut b.id_indices);
std::mem::swap(&mut a.property_indices, &mut b.property_indices);
std::mem::swap(&mut a.composite_indices, &mut b.composite_indices);
std::mem::swap(&mut a.range_indices, &mut b.range_indices);
std::mem::swap(&mut a.secondary_label_index, &mut b.secondary_label_index);
std::mem::swap(&mut a.embeddings, &mut b.embeddings);
std::mem::swap(&mut a.timeseries_store, &mut b.timeseries_store);
std::mem::swap(&mut a.unique_indices, &mut b.unique_indices);
}
fn undo_bucket_append(members: Option<&mut Vec<NodeIndex>>, idx: NodeIndex) {
if let Some(members) = members {
if let Some(pos) = members.iter().rposition(|member| *member == idx) {
members.remove(pos);
}
}
}
fn journal_covers(graph: &DirGraph) -> bool {
graph.graph.supports_undo_journal()
}
pub(crate) enum StatementCheckpoint {
None,
Journal {
shell: Box<DirGraph>,
recorded_ops: Option<usize>,
},
Clone { snapshot: Box<DirGraph> },
}
impl StatementCheckpoint {
pub(crate) fn open(graph: &mut DirGraph) -> Self {
if !journal_covers(graph) {
return Self::Clone {
snapshot: Box::new(graph.fork_transaction()),
};
}
let recorded_ops = graph.graph.recorded_ops_len();
let shell = Box::new(graph.schema_shell());
graph.graph.begin_undo();
Self::Journal {
shell,
recorded_ops,
}
}
pub(crate) fn commit(self, graph: &mut DirGraph) {
if matches!(self, Self::Journal { .. }) {
graph.graph.take_undo();
}
}
pub(crate) fn rollback(self, graph: &mut DirGraph) {
match self {
Self::None => {}
Self::Clone { snapshot } => *graph = *snapshot,
Self::Journal {
shell,
recorded_ops,
} => {
let journal = graph.graph.take_undo();
let fallout = journal
.map(|journal| replay(graph, *journal))
.unwrap_or_default();
if let Some(len) = recorded_ops {
graph.graph.truncate_recorded_ops(len);
}
graph.restore_schema_shell(*shell);
graph.rebuild_unique_indices_for_types(&fallout.stale_unique_indices);
graph.invalidate_edge_type_counts_cache();
}
}
}
}
impl DirGraph {
pub(super) fn schema_shell(&mut self) -> DirGraph {
let mut husk = DirGraph::new();
swap_data_scale(self, &mut husk);
let shell = self.clone();
swap_data_scale(self, &mut husk);
shell
}
pub(super) fn restore_schema_shell(&mut self, mut shell: DirGraph) {
swap_data_scale(&mut shell, self);
*self = shell;
}
}
#[derive(Default)]
struct ReplayFallout {
stale_id_indices: HashSet<String>,
stale_unique_indices: HashSet<String>,
}
impl ReplayFallout {
fn node_identity_changed(&mut self, node_type: String) {
self.stale_unique_indices.insert(node_type.clone());
self.stale_id_indices.insert(node_type);
}
}
fn replay(graph: &mut DirGraph, journal: UndoJournal) -> ReplayFallout {
let mut fallout = ReplayFallout::default();
for entry in journal.into_replay_order() {
apply(graph, entry, &mut fallout);
}
for node_type in &fallout.stale_id_indices {
graph.id_indices.remove(node_type);
}
fallout
}
fn apply(graph: &mut DirGraph, entry: UndoEntry, fallout: &mut ReplayFallout) {
match entry {
UndoEntry::NodeAdded { idx, node_type } => undo_node_added(graph, idx, node_type, fallout),
UndoEntry::NodeWeight { idx, prior } => undo_node_weight(graph, idx, prior, fallout),
UndoEntry::NodeRemoved { idx, prior } => undo_node_removed(graph, idx, prior, fallout),
UndoEntry::EdgeAdded { idx } => {
GraphWrite::remove_edge(&mut graph.graph, idx);
}
UndoEntry::EdgeWeight { idx, prior } => {
if let Some(slot) = GraphWrite::edge_weight_mut(&mut graph.graph, idx) {
*slot = prior;
}
}
UndoEntry::EdgeRemoved {
idx,
src,
tgt,
prior,
} => {
let restored = GraphWrite::add_edge(&mut graph.graph, src, tgt, prior);
debug_assert_slot_reused(restored.index(), idx.index(), "edge");
}
UndoEntry::BucketAppended {
bucket,
idx,
bucket_was_new,
} => undo_bucket_appended(graph, bucket, idx, bucket_was_new),
UndoEntry::BucketRemoved { bucket, idx, pos } => {
undo_bucket_removed(graph, bucket, idx, pos)
}
UndoEntry::TimeseriesRemoved { node, prior } => {
graph.timeseries_store.insert(node, *prior);
}
UndoEntry::EmbeddingRemoved {
store_key,
node,
prior,
} => {
if let Some(store) = graph.embeddings.get_mut(&store_key) {
store.restore_embedding(node, &prior);
}
}
UndoEntry::ColumnarCell {
node_type,
row_id,
key,
prior,
} => {
edit_column_master(graph, node_type, fallout, |store| {
store.set(row_id, key, &prior.unwrap_or(Value::Null), None);
});
}
UndoEntry::ColumnarSchemaGrown {
node_type,
prior_schema,
prior_column_count,
} => {
edit_column_master(graph, node_type, fallout, |store| {
store.restore_schema(prior_schema, prior_column_count);
});
}
UndoEntry::ColumnarRowsAppended {
node_type,
prior_row_count,
prior_schema,
prior_column_count,
store_was_new,
} => undo_columnar_rows_appended(
graph,
node_type,
prior_row_count,
prior_schema,
prior_column_count,
store_was_new,
fallout,
),
UndoEntry::ColumnarTitle {
node_type,
row_id,
prior,
} => {
edit_column_master(graph, node_type, fallout, |store| {
store.set_title(row_id, &prior.unwrap_or(Value::Null));
});
}
UndoEntry::ColumnarTombstone { node_type, row_id } => {
edit_column_master(graph, node_type, fallout, |store| {
store.untombstone(row_id);
});
}
}
}
fn undo_node_added(
graph: &mut DirGraph,
idx: NodeIndex,
node_type: InternedKey,
fallout: &mut ReplayFallout,
) {
fallout.node_identity_changed(graph.interner.resolve(node_type).to_string());
debug_assert!(
graph
.graph
.edges_directed(idx, petgraph::Direction::Outgoing)
.next()
.is_none()
&& graph
.graph
.edges_directed(idx, petgraph::Direction::Incoming)
.next()
.is_none(),
"a rolled-back node must be isolated before removal"
);
GraphWrite::remove_node(&mut graph.graph, idx);
}
fn undo_node_weight(
graph: &mut DirGraph,
idx: NodeIndex,
prior: NodeData,
fallout: &mut ReplayFallout,
) {
let restored_type = graph.interner.resolve(prior.node_type).to_string();
if let Some(current) = GraphRead::node_type_of(&graph.graph, idx) {
let current = graph.interner.resolve(current).to_string();
if current != restored_type {
fallout.stale_unique_indices.insert(current);
}
}
fallout.stale_unique_indices.insert(restored_type);
if let Some(slot) = GraphWrite::node_weight_mut(&mut graph.graph, idx) {
*slot = prior;
}
}
fn undo_node_removed(
graph: &mut DirGraph,
idx: NodeIndex,
prior: NodeData,
fallout: &mut ReplayFallout,
) {
let type_name = graph.interner.resolve(prior.node_type).to_string();
let restored = GraphWrite::add_node(&mut graph.graph, prior);
debug_assert_slot_reused(restored.index(), idx.index(), "node");
fallout.node_identity_changed(type_name);
}
fn undo_bucket_appended(
graph: &mut DirGraph,
bucket: BucketId,
idx: NodeIndex,
bucket_was_new: bool,
) {
match bucket {
BucketId::NodeType(name) => {
graph.type_indices.undo_append(&name, idx);
if bucket_was_new {
graph.type_indices.remove(&name);
}
}
BucketId::SecondaryLabel(label) => {
if let Some(members) = graph.secondary_label_index.get_mut(&label) {
members.retain(|member| *member != idx);
if bucket_was_new || members.is_empty() {
graph.secondary_label_index.remove(&label);
}
}
}
BucketId::PropertyValue { key, value } => {
if let Some(value_map) = graph.property_indices.get_mut(&key) {
undo_bucket_append(value_map.get_mut(&value), idx);
if bucket_was_new {
value_map.remove(&value);
}
}
}
BucketId::RangeValue { key, value } => {
if let Some(btree) = graph.range_indices.get_mut(&key) {
undo_bucket_append(btree.get_mut(&value), idx);
if bucket_was_new {
btree.remove(&value);
}
}
}
BucketId::CompositeTuple { key, value } => {
if let Some(comp_map) = graph.composite_indices.get_mut(&key) {
undo_bucket_append(comp_map.get_mut(&value), idx);
if bucket_was_new {
comp_map.remove(&value);
}
}
}
}
}
fn undo_bucket_removed(graph: &mut DirGraph, bucket: BucketId, idx: NodeIndex, pos: usize) {
match bucket {
BucketId::NodeType(name) => {
let members = graph.type_indices.entry_or_default(name);
let pos = pos.min(members.len());
members.insert(pos, idx);
}
BucketId::SecondaryLabel(label) => {
let members = graph.secondary_label_index.entry(label).or_default();
let pos = pos.min(members.len());
members.insert(pos, idx);
}
BucketId::PropertyValue { key, value } => {
if let Some(value_map) = graph.property_indices.get_mut(&key) {
let members = value_map.entry_or_default(&value);
let pos = pos.min(members.len());
members.insert(pos, idx);
}
}
BucketId::RangeValue { key, value } => {
if let Some(btree) = graph.range_indices.get_mut(&key) {
let members = btree.entry_or_default(&value);
let pos = pos.min(members.len());
members.insert(pos, idx);
}
}
BucketId::CompositeTuple { key, value } => {
if let Some(comp_map) = graph.composite_indices.get_mut(&key) {
let members = comp_map.entry_or_default(&value);
let pos = pos.min(members.len());
members.insert(pos, idx);
}
}
}
}
fn edit_column_master(
graph: &mut DirGraph,
node_type: InternedKey,
fallout: &mut ReplayFallout,
edit: impl FnOnce(&mut ColumnStore),
) {
let Some(type_name) = graph.interner.try_resolve(node_type).map(str::to_string) else {
return;
};
if let Some(store) = GraphWrite::column_store_mut(&mut graph.graph, node_type) {
edit(std::sync::Arc::make_mut(store));
}
columnar_type_touched(fallout, type_name);
}
fn undo_columnar_rows_appended(
graph: &mut DirGraph,
node_type: InternedKey,
prior_row_count: u32,
prior_schema: std::sync::Arc<crate::graph::schema::TypeSchema>,
prior_column_count: usize,
store_was_new: bool,
fallout: &mut ReplayFallout,
) {
let Some(type_name) = graph.interner.try_resolve(node_type).map(str::to_string) else {
return;
};
if store_was_new {
GraphWrite::take_column_store(&mut graph.graph, node_type);
} else if let Some(store) = GraphWrite::column_store_mut(&mut graph.graph, node_type) {
let store = std::sync::Arc::make_mut(store);
store.truncate_rows(prior_row_count);
store.restore_schema(prior_schema, prior_column_count);
}
columnar_type_touched(fallout, type_name);
}
fn columnar_type_touched(fallout: &mut ReplayFallout, type_name: String) {
fallout.stale_unique_indices.insert(type_name);
}
#[inline]
fn debug_assert_slot_reused(restored: usize, expected: usize, kind: &str) {
if restored != expected {
debug_assert_eq!(
restored, expected,
"rollback re-inserted a {kind} on slot {restored}, expected {expected}: \
petgraph's LIFO free-list reuse no longer holds"
);
eprintln!(
"[kglite] statement rollback re-inserted a {kind} on slot {restored} \
instead of {expected}; derived indexes for that {kind} may be stale. \
Please report this with the failing query."
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::datatypes::Value;
use petgraph::stable_graph::StableDiGraph;
use std::collections::HashMap;
#[test]
fn petgraph_reuses_freed_node_slots_lifo() {
let mut g: StableDiGraph<u32, u32> = StableDiGraph::new();
let indices: Vec<_> = (0..5).map(|i| g.add_node(i)).collect();
for &idx in &indices[1..4] {
g.remove_node(idx);
}
for &idx in indices[1..4].iter().rev() {
let restored = g.add_node(99);
assert_eq!(
restored, idx,
"reverse-order re-insertion must reuse the vacated slot"
);
}
}
#[test]
fn petgraph_reuses_freed_edge_slots_lifo() {
let mut g: StableDiGraph<u32, u32> = StableDiGraph::new();
let a = g.add_node(0);
let b = g.add_node(1);
let edges: Vec<_> = (0..4).map(|w| g.add_edge(a, b, w)).collect();
for &e in &edges[..3] {
g.remove_edge(e);
}
for &e in edges[..3].iter().rev() {
let restored = g.add_edge(a, b, 99);
assert_eq!(restored, e, "edge slots must be reused in reverse order");
}
}
fn item(graph: &mut DirGraph, id: i64) -> NodeIndex {
graph.insert_node_routed(
Value::Int64(id),
Value::String(format!("item-{id}")),
"Item",
HashMap::new(),
)
}
#[test]
fn schema_shell_copies_schema_and_parks_data() {
let mut graph = DirGraph::new();
let idx = item(&mut graph, 1);
graph
.type_indices
.entry_or_default("Item".to_string())
.push(idx);
graph.upsert_node_type_metadata("Item", HashMap::new());
let shell = graph.schema_shell();
assert_eq!(shell.graph.node_count(), 0, "the backend must be parked");
assert!(shell.type_indices.is_empty(), "type_indices must be parked");
assert!(
shell.node_type_metadata.contains_key("Item"),
"schema-scale metadata must be cloned"
);
assert_eq!(
graph.graph.node_count(),
1,
"the live graph must be handed its data back"
);
assert_eq!(graph.type_indices.len(), 1);
}
#[test]
fn gate_accepts_user_property_indexes() {
let mut graph = DirGraph::new();
assert!(journal_covers(&graph));
graph
.property_indices
.insert(("Item".to_string(), "name".to_string()), Default::default());
graph
.range_indices
.insert(("Item".to_string(), "qty".to_string()), Default::default());
graph.composite_indices.insert(
("Item".to_string(), vec!["name".to_string()]),
Default::default(),
);
assert!(journal_covers(&graph));
}
#[test]
fn undoing_an_append_leaves_a_pre_existing_occurrence() {
let mut members = vec![NodeIndex::new(3), NodeIndex::new(7), NodeIndex::new(3)];
undo_bucket_append(Some(&mut members), NodeIndex::new(3));
assert_eq!(members, vec![NodeIndex::new(3), NodeIndex::new(7)]);
}
}