uqa_graph/
versioned_store.rs1use std::collections::BTreeSet;
15
16use crate::delta::{DeltaOp, GraphDelta};
17use crate::store::{GraphStore, GraphStoreError, GraphStoreResult};
18
19type InvalidationCallback = Box<dyn Fn(&BTreeSet<String>) + Send + Sync>;
20
21pub struct VersionedGraphStore<'a, G: GraphStore> {
22 base: &'a mut G,
23 graph: String,
24 version: u64,
25 deltas: Vec<GraphDelta>,
26 inverse_deltas: Vec<GraphDelta>,
27 on_invalidate: Vec<InvalidationCallback>,
28}
29
30impl<'a, G: GraphStore + Clone> VersionedGraphStore<'a, G> {
31 pub fn new(base: &'a mut G, graph: impl Into<String>) -> Self {
32 Self {
33 base,
34 graph: graph.into(),
35 version: 0,
36 deltas: Vec::new(),
37 inverse_deltas: Vec::new(),
38 on_invalidate: Vec::new(),
39 }
40 }
41
42 pub fn version(&self) -> u64 {
43 self.version
44 }
45
46 pub fn base(&self) -> &G {
47 &*self.base
48 }
49
50 pub fn base_mut(&mut self) -> &mut G {
51 self.base
52 }
53
54 pub fn apply(&mut self, delta: GraphDelta) -> GraphStoreResult<u64> {
57 if !self.base.has_graph(&self.graph) {
58 return Err(GraphStoreError::UnknownGraph(self.graph.clone()));
59 }
60 let next_version = self.version.checked_add(1).ok_or_else(|| {
61 GraphStoreError::IdExhausted("graph version counter overflow".to_string())
62 })?;
63 let mut candidate = self.base.clone();
64 let mut inverse = GraphDelta::new();
65 let mut affected_labels = delta.affected_edge_labels();
66 for op in delta.ops() {
67 match op {
68 DeltaOp::AddVertex(vertex) => {
69 candidate.add_vertex(vertex.clone(), &self.graph)?;
70 inverse.remove_vertex(vertex.vertex_id);
71 }
72 DeltaOp::RemoveVertex(vertex_id) => {
73 let existing = candidate.get_vertex(*vertex_id).cloned();
74 let incident_edges: Vec<_> = candidate
75 .edges_in_graph(&self.graph)?
76 .into_iter()
77 .filter(|edge| edge.source_id == *vertex_id || edge.target_id == *vertex_id)
78 .collect();
79 affected_labels.extend(incident_edges.iter().map(|edge| edge.label.clone()));
80 candidate.remove_vertex(*vertex_id, &self.graph)?;
81 for edge in incident_edges {
82 inverse.add_edge(edge);
83 }
84 if let Some(v) = existing {
85 inverse.add_vertex(v);
86 }
87 }
88 DeltaOp::AddEdge(edge) => {
89 candidate.add_edge(edge.clone(), &self.graph)?;
90 inverse.remove_edge(edge.edge_id);
91 }
92 DeltaOp::RemoveEdge(edge_id) => {
93 let existing = candidate.get_edge(*edge_id).cloned();
94 if let Some(edge) = &existing {
95 affected_labels.insert(edge.label.clone());
96 }
97 candidate.remove_edge(*edge_id, &self.graph)?;
98 if let Some(e) = existing {
99 inverse.add_edge(e);
100 }
101 }
102 }
103 }
104 *self.base = candidate;
105 self.version = next_version;
106 self.deltas.push(delta);
107 self.inverse_deltas.push(inverse);
108 if !affected_labels.is_empty() {
109 for callback in &self.on_invalidate {
110 callback(&affected_labels);
111 }
112 }
113 Ok(self.version)
114 }
115
116 pub fn rollback(&mut self, to_version: u64) -> GraphStoreResult<()> {
119 if to_version > self.version {
120 return Err(GraphStoreError::InvalidMutation(format!(
121 "cannot rollback to version {to_version} (current: {})",
122 self.version
123 )));
124 }
125 let mut candidate = self.base.clone();
126 let mut remaining_version = self.version;
127 let mut inverse_count = 0usize;
128 while remaining_version > to_version {
129 let offset = inverse_count.checked_add(1).ok_or_else(|| {
130 GraphStoreError::CorruptGraph("version history index overflow".to_string())
131 })?;
132 let inverse = self
133 .inverse_deltas
134 .get(
135 self.inverse_deltas
136 .len()
137 .checked_sub(offset)
138 .ok_or_else(|| {
139 GraphStoreError::CorruptGraph(
140 "version history is shorter than the current graph version"
141 .to_string(),
142 )
143 })?,
144 )
145 .ok_or_else(|| {
146 GraphStoreError::CorruptGraph(
147 "version history is shorter than the current graph version".to_string(),
148 )
149 })?;
150 for op in inverse.ops().iter().rev() {
151 match op {
152 DeltaOp::AddVertex(vertex) => {
153 candidate.add_vertex(vertex.clone(), &self.graph)?;
154 }
155 DeltaOp::RemoveVertex(vertex_id) => {
156 candidate.remove_vertex(*vertex_id, &self.graph)?;
157 }
158 DeltaOp::AddEdge(edge) => {
159 candidate.add_edge(edge.clone(), &self.graph)?;
160 }
161 DeltaOp::RemoveEdge(edge_id) => {
162 candidate.remove_edge(*edge_id, &self.graph)?;
163 }
164 }
165 }
166 remaining_version = remaining_version.checked_sub(1).ok_or_else(|| {
167 GraphStoreError::CorruptGraph("graph version underflow".to_string())
168 })?;
169 inverse_count = inverse_count.checked_add(1).ok_or_else(|| {
170 GraphStoreError::CorruptGraph("version history index overflow".to_string())
171 })?;
172 }
173 *self.base = candidate;
174 self.version = remaining_version;
175 let new_len = self
176 .inverse_deltas
177 .len()
178 .checked_sub(inverse_count)
179 .ok_or_else(|| {
180 GraphStoreError::CorruptGraph("version history truncation underflow".to_string())
181 })?;
182 self.inverse_deltas.truncate(new_len);
183 self.deltas.truncate(new_len);
184 Ok(())
185 }
186
187 pub fn on_invalidate<F>(&mut self, callback: F)
190 where
191 F: Fn(&BTreeSet<String>) + Send + Sync + 'static,
192 {
193 self.on_invalidate.push(Box::new(callback));
194 }
195}