Skip to main content

uqa_graph/
versioned_store.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Version-tracked graph store (Section 9.3, Paper 2).
8//!
9//! Wraps a [`GraphStore`] and applies [`GraphDelta`]s with a monotonic
10//! version counter. Each apply records an inverse delta so the store
11//! can rewind to an earlier version. Invalidation callbacks fire on
12//! affected edge labels so dependent path indexes can refresh.
13
14use 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    /// Apply a delta to the base store, accumulating an inverse delta
55    /// for rollback. Returns the new version number.
56    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    /// Rewind to the given version by replaying inverse deltas. Errors
117    /// when the target version is in the future or below zero.
118    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    /// Register a callback fired with the set of affected edge labels
188    /// every time `apply` lands a delta that touches at least one edge.
189    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}