Skip to main content

graphforge_storage/
adjacency_delta.rs

1//! Incremental adjacency: per-commit **delta segments** (#765).
2//!
3//! Each pure-append topology commit writes one tiny Parquet file
4//! `indexes/adjacency/deltas/<generation>.parquet` holding exactly the edges it
5//! created (possibly zero rows, for a node-only flush). The adjacency provider
6//! serves a fresh view as *base CSR ⊎ a contiguous chain of segments* covering
7//! `(G_base, G_cur]`, so newly-created edges are visible without a full rebuild.
8//! Anything that breaks the chain — a DELETE, a crash between commit and segment
9//! write, an unreadable segment, or a chain longer than [`MAX_DELTA_CHAIN`] —
10//! reads as stale and falls back to the existing full-rebuild path.
11//!
12//! The merge ([`apply_delta_segments`]) reconstructs the base entries from the
13//! loaded CSR, concatenates the (filtered) delta entries, and re-runs the
14//! builder's [`csr_from_entries`](crate::adjacency::csr_from_entries). Because
15//! `edge_id` is a unique surrogate, the result is byte-identical to a full
16//! rebuild from `topology/` for the same edges — that equivalence is the
17//! correctness contract (acceptance criterion 1), pinned by the tests below.
18
19use std::path::{Path, PathBuf};
20use std::sync::Arc;
21
22use arrow::array::{RecordBatch, StringArray, UInt64Array};
23use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
24
25use graphforge_core::GfError;
26
27use crate::adjacency::{
28    ALL_RELATIONS_STEM, BuildEntry, CsrIndex, Direction, adjacency_dir, csr_from_entries,
29    usable_stem,
30};
31use crate::schemas::ADJACENCY_DELTA_SCHEMA;
32use crate::staging::RewriteBatch;
33
34/// Longest delta chain served before the index reads stale and is rebuilt.
35/// Bounds both per-view merge cost and `deltas/` accumulation between rebuilds.
36pub const MAX_DELTA_CHAIN: u64 = 64;
37
38fn storage_err(e: impl std::fmt::Display) -> GfError {
39    GfError::Storage(e.to_string())
40}
41
42/// One created edge in a delta segment, in creation (ascending `edge_id`) order.
43#[derive(Clone, Debug, PartialEq, Eq)]
44pub struct DeltaEdge {
45    /// Typed file stem, or the exploratory row's `rel_type_name`.
46    pub rel_type_name: String,
47    /// Edge surrogate (globally ascending across an intact chain).
48    pub edge_id: u64,
49    /// Source node surrogate.
50    pub src_id: u64,
51    /// Destination node surrogate.
52    pub dst_id: u64,
53}
54
55/// The edges created by the commit that bumped the counter to `generation`.
56#[derive(Clone, Debug, PartialEq, Eq)]
57pub struct DeltaSegment {
58    /// Post-bump topology generation this segment is tagged with.
59    pub generation: u64,
60    /// Created edges, in creation order. Empty for a node-only commit.
61    pub edges: Vec<DeltaEdge>,
62}
63
64/// `indexes/adjacency/deltas/` within `project_dir`.
65#[must_use]
66pub fn delta_dir(project_dir: &Path) -> PathBuf {
67    adjacency_dir(project_dir).join("deltas")
68}
69
70/// Path of the delta segment for `generation`.
71#[must_use]
72pub fn delta_path(project_dir: &Path, generation: u64) -> PathBuf {
73    delta_dir(project_dir).join(format!("{generation}.parquet"))
74}
75
76/// Write the segment for `generation` (creating `deltas/` if needed). A
77/// zero-edge segment is still written so a node-only commit keeps the chain
78/// contiguous. Atomic via temp-file + rename ([`RewriteBatch`]).
79///
80/// # Errors
81/// [`GfError::Storage`] on any directory-create, encode, or rename failure.
82pub fn write_delta_segment(
83    project_dir: &Path,
84    generation: u64,
85    edges: &[DeltaEdge],
86) -> Result<(), GfError> {
87    std::fs::create_dir_all(delta_dir(project_dir)).map_err(storage_err)?;
88    let rel_types: StringArray = edges
89        .iter()
90        .map(|e| Some(e.rel_type_name.as_str()))
91        .collect();
92    let edge_ids: UInt64Array = edges.iter().map(|e| e.edge_id).collect();
93    let src_ids: UInt64Array = edges.iter().map(|e| e.src_id).collect();
94    let dst_ids: UInt64Array = edges.iter().map(|e| e.dst_id).collect();
95    let batch = RecordBatch::try_new(
96        Arc::clone(&ADJACENCY_DELTA_SCHEMA),
97        vec![
98            Arc::new(rel_types),
99            Arc::new(edge_ids),
100            Arc::new(src_ids),
101            Arc::new(dst_ids),
102        ],
103    )
104    .map_err(storage_err)?;
105
106    let mut staged = RewriteBatch::new();
107    staged.stage(
108        &delta_path(project_dir, generation),
109        Arc::clone(&ADJACENCY_DELTA_SCHEMA),
110        &batch,
111    )?;
112    staged.commit()
113}
114
115/// Read one delta segment from disk. A missing file is an error here — callers
116/// that tolerate absence ([`read_delta_chain`]) check existence first.
117///
118/// # Errors
119/// [`GfError::Storage`] if the file is missing, unreadable, or has the wrong
120/// schema.
121pub fn read_delta_segment(project_dir: &Path, generation: u64) -> Result<DeltaSegment, GfError> {
122    let path = delta_path(project_dir, generation);
123    let file = std::fs::File::open(&path).map_err(storage_err)?;
124    let reader = ParquetRecordBatchReaderBuilder::try_new(file)
125        .map_err(storage_err)?
126        .build()
127        .map_err(storage_err)?;
128    let mut edges = Vec::new();
129    for batch in reader {
130        let batch = batch.map_err(storage_err)?;
131        if batch.schema().fields() != ADJACENCY_DELTA_SCHEMA.fields() {
132            return Err(GfError::Storage(format!(
133                "adjacency delta {} has unexpected schema",
134                path.display()
135            )));
136        }
137        let rel_types = batch
138            .column(0)
139            .as_any()
140            .downcast_ref::<StringArray>()
141            .ok_or_else(|| GfError::Storage("delta: rel_type_name not Utf8".to_owned()))?;
142        let cols: Vec<&UInt64Array> = (1..=3)
143            .map(|i| {
144                batch
145                    .column(i)
146                    .as_any()
147                    .downcast_ref::<UInt64Array>()
148                    .ok_or_else(|| GfError::Storage("delta: id column not UInt64".to_owned()))
149            })
150            .collect::<Result<_, _>>()?;
151        for i in 0..batch.num_rows() {
152            edges.push(DeltaEdge {
153                rel_type_name: rel_types.value(i).to_owned(),
154                edge_id: cols[0].value(i),
155                src_id: cols[1].value(i),
156                dst_id: cols[2].value(i),
157            });
158        }
159    }
160    Ok(DeltaSegment { generation, edges })
161}
162
163/// Read the contiguous chain of segments covering `(base, current]`.
164///
165/// Returns `Some(chain)` only when the chain is *intact and bounded*: every
166/// generation `base+1 ..= current` has a readable segment and the span is at
167/// most [`MAX_DELTA_CHAIN`]. Returns `Some(vec![])` when `current <= base`
168/// (nothing to apply). Returns `None` on any gap, unreadable segment, or
169/// over-cap span — all of which the provider treats as stale (full rebuild).
170#[must_use]
171pub fn read_delta_chain(project_dir: &Path, base: u64, current: u64) -> Option<Vec<DeltaSegment>> {
172    if current <= base {
173        return Some(Vec::new());
174    }
175    if current - base > MAX_DELTA_CHAIN {
176        return None;
177    }
178    let mut chain = Vec::with_capacity(usize::try_from(current - base).unwrap_or(0));
179    for g in (base + 1)..=current {
180        if !delta_path(project_dir, g).exists() {
181            return None; // gap ⇒ chain broken
182        }
183        match read_delta_segment(project_dir, g) {
184            Ok(seg) => chain.push(seg),
185            Err(_) => return None, // unreadable ⇒ chain broken
186        }
187    }
188    Some(chain)
189}
190
191/// Best-effort removal of the segment at `generation` — for a commit that
192/// bumps the counter but writes **no** segment (any DELETE / non-pure-append
193/// statement). This makes "no segment here" an affirmative invariant: the chain
194/// can never contain a file the bumping commit did not author, even across a
195/// counter reset that left a stale file at that generation. A failure is
196/// harmless — an unexpected file at `generation` only breaks the chain there,
197/// forcing a (correct) rebuild.
198pub fn discard_segment(project_dir: &Path, generation: u64) {
199    let _ = std::fs::remove_file(delta_path(project_dir, generation));
200}
201
202/// Remove every delta segment with generation `<= up_to` (consumed by a rebuild
203/// or compaction at `up_to`). Best-effort: a failed unlink leaves dead weight a
204/// later prune removes, never incorrect data. Segments `> up_to` (written by a
205/// concurrent append during the build) survive, so the new base + those is
206/// immediately fresh.
207pub fn prune_delta_segments(project_dir: &Path, up_to: u64) {
208    let dir = delta_dir(project_dir);
209    let Ok(entries) = std::fs::read_dir(&dir) else {
210        return;
211    };
212    for entry in entries.flatten() {
213        let path = entry.path();
214        if path.extension().and_then(|s| s.to_str()) != Some("parquet") {
215            continue;
216        }
217        let parsed = path
218            .file_stem()
219            .and_then(|s| s.to_str())
220            .and_then(|s| s.parse::<u64>().ok());
221        if let Some(g) = parsed
222            && g <= up_to
223        {
224            let _ = std::fs::remove_file(&path);
225        }
226    }
227}
228
229/// Merge a contiguous delta `chain` onto `base` (a loaded CSR for `direction`)
230/// and return the overlaid CSR — equal to a full rebuild over the base's edges
231/// plus the chain's, for `stem`.
232///
233/// `stem == _all` takes every delta edge; a per-relation `stem` takes edges
234/// whose `rel_type_name == stem` (and `usable_stem`, matching the builder, so a
235/// hostile relation name can never be materialized at a per-relation path).
236///
237/// Implementation: reconstruct the base `(src, edge, dst)` entries from the CSR,
238/// concatenate the filtered delta entries, and re-run
239/// [`csr_from_entries`](crate::adjacency::csr_from_entries). Since `edge_id` is
240/// unique, the sorted result is identical to building from the union directly —
241/// independent of input order, so it is robust even if a segment's edge_ids do
242/// not strictly exceed the base's.
243#[must_use]
244pub fn apply_delta_segments(
245    base: &CsrIndex,
246    stem: &str,
247    direction: Direction,
248    chain: &[DeltaSegment],
249) -> CsrIndex {
250    let mut entries: Vec<BuildEntry> = base_entries(base, direction);
251    let take_all = stem == ALL_RELATIONS_STEM;
252    for seg in chain {
253        for e in &seg.edges {
254            if take_all || (e.rel_type_name == stem && usable_stem(&e.rel_type_name)) {
255                entries.push((e.src_id, e.edge_id, e.dst_id));
256            }
257        }
258    }
259    csr_from_entries(&entries, direction)
260}
261
262/// Reconstruct the `(src, edge, dst)` build entries a CSR for `direction` was
263/// built from: an `out` CSR is keyed by `src` with `dst` neighbors, an `in` CSR
264/// by `dst` with `src` neighbors.
265fn base_entries(base: &CsrIndex, direction: Direction) -> Vec<BuildEntry> {
266    let mut entries = Vec::with_capacity(base.edge_ids.len());
267    for key in 0..base.node_count() {
268        let lo = base.offsets[usize::try_from(key).unwrap_or(0)];
269        let hi = base.offsets[usize::try_from(key + 1).unwrap_or(0)];
270        for j in lo..hi {
271            let j = usize::try_from(j).unwrap_or(0);
272            let (edge, neighbor) = (base.edge_ids[j], base.neighbor_ids[j]);
273            entries.push(match direction {
274                Direction::Out => (key, edge, neighbor), // key=src, neighbor=dst
275                Direction::In => (neighbor, edge, key),  // key=dst, neighbor=src
276            });
277        }
278    }
279    entries
280}
281
282#[cfg(test)]
283mod tests {
284    use super::*;
285    use tempfile::TempDir;
286
287    fn seg(generation: u64, edges: &[(&str, u64, u64, u64)]) -> DeltaSegment {
288        DeltaSegment {
289            generation,
290            edges: edges
291                .iter()
292                .map(|&(rel, edge_id, src_id, dst_id)| DeltaEdge {
293                    rel_type_name: rel.to_owned(),
294                    edge_id,
295                    src_id,
296                    dst_id,
297                })
298                .collect(),
299        }
300    }
301
302    /// The merge equals a full rebuild for every (stem, direction) — the #765
303    /// acceptance-1 contract, over a fixture with parallel edges, a self-loop,
304    /// a new node beyond the base, and an unusable relation name.
305    #[test]
306    fn apply_equals_full_rebuild() {
307        // Base edges (src, edge, dst), edge_ids 1..=5.
308        let base_edges: Vec<BuildEntry> = vec![
309            (0, 1, 1),
310            (0, 2, 1), // parallel edge 0->1
311            (1, 3, 2),
312            (2, 4, 0),
313            (2, 5, 2), // self-loop 2->2
314        ];
315        // Delta edges, edge_ids 6..=9 (strictly after the base), incl. a new
316        // node 3 and an unusable relation name.
317        let chain = vec![
318            seg(8, &[("KNOWS", 6, 1, 3), ("KNOWS", 7, 3, 0)]),
319            seg(9, &[("OWNS", 8, 0, 2), ("../evil", 9, 3, 1)]),
320        ];
321        let delta_entries: Vec<BuildEntry> = vec![(1, 6, 3), (3, 7, 0), (0, 8, 2), (3, 9, 1)];
322
323        for direction in [Direction::Out, Direction::In] {
324            // _all overlay: every delta edge participates.
325            let base_all = csr_from_entries(&base_edges, direction);
326            let mut all = base_edges.clone();
327            all.extend_from_slice(&delta_entries);
328            let expected_all = csr_from_entries(&all, direction);
329            assert_eq!(
330                apply_delta_segments(&base_all, ALL_RELATIONS_STEM, direction, &chain),
331                expected_all,
332                "_all {direction:?}"
333            );
334
335            // Per-relation KNOWS overlay: only KNOWS delta edges (6, 7); the
336            // unusable "../evil" name never reaches a per-rel stem.
337            let base_knows: Vec<BuildEntry> = vec![(0, 1, 1), (0, 2, 1), (1, 3, 2)];
338            let base_knows_csr = csr_from_entries(&base_knows, direction);
339            let mut knows = base_knows.clone();
340            knows.push((1, 6, 3));
341            knows.push((3, 7, 0));
342            let expected_knows = csr_from_entries(&knows, direction);
343            assert_eq!(
344                apply_delta_segments(&base_knows_csr, "KNOWS", direction, &chain),
345                expected_knows,
346                "KNOWS {direction:?}"
347            );
348        }
349    }
350
351    #[test]
352    fn empty_chain_returns_the_base_unchanged() {
353        let base = csr_from_entries(&[(0, 1, 1), (1, 2, 0)], Direction::Out);
354        assert_eq!(
355            apply_delta_segments(&base, ALL_RELATIONS_STEM, Direction::Out, &[]),
356            base
357        );
358    }
359
360    #[test]
361    fn segment_round_trips_including_empty() {
362        let dir = TempDir::new().unwrap();
363        let edges = vec![
364            DeltaEdge {
365                rel_type_name: "KNOWS".into(),
366                edge_id: 6,
367                src_id: 1,
368                dst_id: 3,
369            },
370            DeltaEdge {
371                rel_type_name: "OWNS".into(),
372                edge_id: 7,
373                src_id: 0,
374                dst_id: 2,
375            },
376        ];
377        write_delta_segment(dir.path(), 6, &edges).unwrap();
378        assert_eq!(read_delta_segment(dir.path(), 6).unwrap().edges, edges);
379
380        // Node-only commit: an empty segment keeps the chain contiguous.
381        write_delta_segment(dir.path(), 7, &[]).unwrap();
382        assert!(read_delta_segment(dir.path(), 7).unwrap().edges.is_empty());
383    }
384
385    #[test]
386    fn chain_is_some_only_when_contiguous_and_bounded() {
387        let dir = TempDir::new().unwrap();
388        // Generations 6 and 7 present.
389        write_delta_segment(dir.path(), 6, &[]).unwrap();
390        write_delta_segment(dir.path(), 7, &[]).unwrap();
391
392        assert_eq!(read_delta_chain(dir.path(), 5, 5).unwrap().len(), 0); // none needed
393        assert_eq!(read_delta_chain(dir.path(), 5, 7).unwrap().len(), 2); // contiguous
394        assert!(read_delta_chain(dir.path(), 5, 8).is_none()); // gap at 8
395        assert!(read_delta_chain(dir.path(), 4, 7).is_none()); // gap at 5
396        assert!(read_delta_chain(dir.path(), 0, MAX_DELTA_CHAIN + 1).is_none()); // over cap
397    }
398
399    #[test]
400    fn prune_removes_only_consumed_segments() {
401        let dir = TempDir::new().unwrap();
402        for g in 6..=9 {
403            write_delta_segment(dir.path(), g, &[]).unwrap();
404        }
405        prune_delta_segments(dir.path(), 7);
406        assert!(!delta_path(dir.path(), 6).exists());
407        assert!(!delta_path(dir.path(), 7).exists());
408        assert!(delta_path(dir.path(), 8).exists()); // written after the build stamp
409        assert!(delta_path(dir.path(), 9).exists());
410        // Idempotent.
411        prune_delta_segments(dir.path(), 7);
412        assert!(delta_path(dir.path(), 8).exists());
413    }
414}