1use 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
34pub 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#[derive(Clone, Debug, PartialEq, Eq)]
44pub struct DeltaEdge {
45 pub rel_type_name: String,
47 pub edge_id: u64,
49 pub src_id: u64,
51 pub dst_id: u64,
53}
54
55#[derive(Clone, Debug, PartialEq, Eq)]
57pub struct DeltaSegment {
58 pub generation: u64,
60 pub edges: Vec<DeltaEdge>,
62}
63
64#[must_use]
66pub fn delta_dir(project_dir: &Path) -> PathBuf {
67 adjacency_dir(project_dir).join("deltas")
68}
69
70#[must_use]
72pub fn delta_path(project_dir: &Path, generation: u64) -> PathBuf {
73 delta_dir(project_dir).join(format!("{generation}.parquet"))
74}
75
76pub 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
115pub 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#[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; }
183 match read_delta_segment(project_dir, g) {
184 Ok(seg) => chain.push(seg),
185 Err(_) => return None, }
187 }
188 Some(chain)
189}
190
191pub fn discard_segment(project_dir: &Path, generation: u64) {
199 let _ = std::fs::remove_file(delta_path(project_dir, generation));
200}
201
202pub 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#[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
262fn 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), Direction::In => (neighbor, edge, key), });
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 #[test]
306 fn apply_equals_full_rebuild() {
307 let base_edges: Vec<BuildEntry> = vec![
309 (0, 1, 1),
310 (0, 2, 1), (1, 3, 2),
312 (2, 4, 0),
313 (2, 5, 2), ];
315 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 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 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 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 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); assert_eq!(read_delta_chain(dir.path(), 5, 7).unwrap().len(), 2); assert!(read_delta_chain(dir.path(), 5, 8).is_none()); assert!(read_delta_chain(dir.path(), 4, 7).is_none()); assert!(read_delta_chain(dir.path(), 0, MAX_DELTA_CHAIN + 1).is_none()); }
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()); assert!(delta_path(dir.path(), 9).exists());
410 prune_delta_segments(dir.path(), 7);
412 assert!(delta_path(dir.path(), 8).exists());
413 }
414}