use std::path::{Path, PathBuf};
use std::sync::Arc;
use arrow::array::{RecordBatch, StringArray, UInt64Array};
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use graphforge_core::GfError;
use crate::adjacency::{
ALL_RELATIONS_STEM, BuildEntry, CsrIndex, Direction, adjacency_dir, csr_from_entries,
usable_stem,
};
use crate::schemas::ADJACENCY_DELTA_SCHEMA;
use crate::staging::RewriteBatch;
pub const MAX_DELTA_CHAIN: u64 = 64;
fn storage_err(e: impl std::fmt::Display) -> GfError {
GfError::Storage(e.to_string())
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct DeltaEdge {
pub rel_type_name: String,
pub edge_id: u64,
pub src_id: u64,
pub dst_id: u64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct DeltaSegment {
pub generation: u64,
pub edges: Vec<DeltaEdge>,
}
#[must_use]
pub fn delta_dir(project_dir: &Path) -> PathBuf {
adjacency_dir(project_dir).join("deltas")
}
#[must_use]
pub fn delta_path(project_dir: &Path, generation: u64) -> PathBuf {
delta_dir(project_dir).join(format!("{generation}.parquet"))
}
pub fn write_delta_segment(
project_dir: &Path,
generation: u64,
edges: &[DeltaEdge],
) -> Result<(), GfError> {
std::fs::create_dir_all(delta_dir(project_dir)).map_err(storage_err)?;
let rel_types: StringArray = edges
.iter()
.map(|e| Some(e.rel_type_name.as_str()))
.collect();
let edge_ids: UInt64Array = edges.iter().map(|e| e.edge_id).collect();
let src_ids: UInt64Array = edges.iter().map(|e| e.src_id).collect();
let dst_ids: UInt64Array = edges.iter().map(|e| e.dst_id).collect();
let batch = RecordBatch::try_new(
Arc::clone(&ADJACENCY_DELTA_SCHEMA),
vec![
Arc::new(rel_types),
Arc::new(edge_ids),
Arc::new(src_ids),
Arc::new(dst_ids),
],
)
.map_err(storage_err)?;
let mut staged = RewriteBatch::new();
staged.stage(
&delta_path(project_dir, generation),
Arc::clone(&ADJACENCY_DELTA_SCHEMA),
&batch,
)?;
staged.commit()
}
pub fn read_delta_segment(project_dir: &Path, generation: u64) -> Result<DeltaSegment, GfError> {
let path = delta_path(project_dir, generation);
let file = std::fs::File::open(&path).map_err(storage_err)?;
let reader = ParquetRecordBatchReaderBuilder::try_new(file)
.map_err(storage_err)?
.build()
.map_err(storage_err)?;
let mut edges = Vec::new();
for batch in reader {
let batch = batch.map_err(storage_err)?;
if batch.schema().fields() != ADJACENCY_DELTA_SCHEMA.fields() {
return Err(GfError::Storage(format!(
"adjacency delta {} has unexpected schema",
path.display()
)));
}
let rel_types = batch
.column(0)
.as_any()
.downcast_ref::<StringArray>()
.ok_or_else(|| GfError::Storage("delta: rel_type_name not Utf8".to_owned()))?;
let cols: Vec<&UInt64Array> = (1..=3)
.map(|i| {
batch
.column(i)
.as_any()
.downcast_ref::<UInt64Array>()
.ok_or_else(|| GfError::Storage("delta: id column not UInt64".to_owned()))
})
.collect::<Result<_, _>>()?;
for i in 0..batch.num_rows() {
edges.push(DeltaEdge {
rel_type_name: rel_types.value(i).to_owned(),
edge_id: cols[0].value(i),
src_id: cols[1].value(i),
dst_id: cols[2].value(i),
});
}
}
Ok(DeltaSegment { generation, edges })
}
#[must_use]
pub fn read_delta_chain(project_dir: &Path, base: u64, current: u64) -> Option<Vec<DeltaSegment>> {
if current <= base {
return Some(Vec::new());
}
if current - base > MAX_DELTA_CHAIN {
return None;
}
let mut chain = Vec::with_capacity(usize::try_from(current - base).unwrap_or(0));
for g in (base + 1)..=current {
if !delta_path(project_dir, g).exists() {
return None; }
match read_delta_segment(project_dir, g) {
Ok(seg) => chain.push(seg),
Err(_) => return None, }
}
Some(chain)
}
pub fn discard_segment(project_dir: &Path, generation: u64) {
let _ = std::fs::remove_file(delta_path(project_dir, generation));
}
pub fn prune_delta_segments(project_dir: &Path, up_to: u64) {
let dir = delta_dir(project_dir);
let Ok(entries) = std::fs::read_dir(&dir) else {
return;
};
for entry in entries.flatten() {
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) != Some("parquet") {
continue;
}
let parsed = path
.file_stem()
.and_then(|s| s.to_str())
.and_then(|s| s.parse::<u64>().ok());
if let Some(g) = parsed
&& g <= up_to
{
let _ = std::fs::remove_file(&path);
}
}
}
#[must_use]
pub fn apply_delta_segments(
base: &CsrIndex,
stem: &str,
direction: Direction,
chain: &[DeltaSegment],
) -> CsrIndex {
let mut entries: Vec<BuildEntry> = base_entries(base, direction);
let take_all = stem == ALL_RELATIONS_STEM;
for seg in chain {
for e in &seg.edges {
if take_all || (e.rel_type_name == stem && usable_stem(&e.rel_type_name)) {
entries.push((e.src_id, e.edge_id, e.dst_id));
}
}
}
csr_from_entries(&entries, direction)
}
fn base_entries(base: &CsrIndex, direction: Direction) -> Vec<BuildEntry> {
let mut entries = Vec::with_capacity(base.edge_ids.len());
for key in 0..base.node_count() {
let lo = base.offsets[usize::try_from(key).unwrap_or(0)];
let hi = base.offsets[usize::try_from(key + 1).unwrap_or(0)];
for j in lo..hi {
let j = usize::try_from(j).unwrap_or(0);
let (edge, neighbor) = (base.edge_ids[j], base.neighbor_ids[j]);
entries.push(match direction {
Direction::Out => (key, edge, neighbor), Direction::In => (neighbor, edge, key), });
}
}
entries
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn seg(generation: u64, edges: &[(&str, u64, u64, u64)]) -> DeltaSegment {
DeltaSegment {
generation,
edges: edges
.iter()
.map(|&(rel, edge_id, src_id, dst_id)| DeltaEdge {
rel_type_name: rel.to_owned(),
edge_id,
src_id,
dst_id,
})
.collect(),
}
}
#[test]
fn apply_equals_full_rebuild() {
let base_edges: Vec<BuildEntry> = vec![
(0, 1, 1),
(0, 2, 1), (1, 3, 2),
(2, 4, 0),
(2, 5, 2), ];
let chain = vec![
seg(8, &[("KNOWS", 6, 1, 3), ("KNOWS", 7, 3, 0)]),
seg(9, &[("OWNS", 8, 0, 2), ("../evil", 9, 3, 1)]),
];
let delta_entries: Vec<BuildEntry> = vec![(1, 6, 3), (3, 7, 0), (0, 8, 2), (3, 9, 1)];
for direction in [Direction::Out, Direction::In] {
let base_all = csr_from_entries(&base_edges, direction);
let mut all = base_edges.clone();
all.extend_from_slice(&delta_entries);
let expected_all = csr_from_entries(&all, direction);
assert_eq!(
apply_delta_segments(&base_all, ALL_RELATIONS_STEM, direction, &chain),
expected_all,
"_all {direction:?}"
);
let base_knows: Vec<BuildEntry> = vec![(0, 1, 1), (0, 2, 1), (1, 3, 2)];
let base_knows_csr = csr_from_entries(&base_knows, direction);
let mut knows = base_knows.clone();
knows.push((1, 6, 3));
knows.push((3, 7, 0));
let expected_knows = csr_from_entries(&knows, direction);
assert_eq!(
apply_delta_segments(&base_knows_csr, "KNOWS", direction, &chain),
expected_knows,
"KNOWS {direction:?}"
);
}
}
#[test]
fn empty_chain_returns_the_base_unchanged() {
let base = csr_from_entries(&[(0, 1, 1), (1, 2, 0)], Direction::Out);
assert_eq!(
apply_delta_segments(&base, ALL_RELATIONS_STEM, Direction::Out, &[]),
base
);
}
#[test]
fn segment_round_trips_including_empty() {
let dir = TempDir::new().unwrap();
let edges = vec![
DeltaEdge {
rel_type_name: "KNOWS".into(),
edge_id: 6,
src_id: 1,
dst_id: 3,
},
DeltaEdge {
rel_type_name: "OWNS".into(),
edge_id: 7,
src_id: 0,
dst_id: 2,
},
];
write_delta_segment(dir.path(), 6, &edges).unwrap();
assert_eq!(read_delta_segment(dir.path(), 6).unwrap().edges, edges);
write_delta_segment(dir.path(), 7, &[]).unwrap();
assert!(read_delta_segment(dir.path(), 7).unwrap().edges.is_empty());
}
#[test]
fn chain_is_some_only_when_contiguous_and_bounded() {
let dir = TempDir::new().unwrap();
write_delta_segment(dir.path(), 6, &[]).unwrap();
write_delta_segment(dir.path(), 7, &[]).unwrap();
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()); }
#[test]
fn prune_removes_only_consumed_segments() {
let dir = TempDir::new().unwrap();
for g in 6..=9 {
write_delta_segment(dir.path(), g, &[]).unwrap();
}
prune_delta_segments(dir.path(), 7);
assert!(!delta_path(dir.path(), 6).exists());
assert!(!delta_path(dir.path(), 7).exists());
assert!(delta_path(dir.path(), 8).exists()); assert!(delta_path(dir.path(), 9).exists());
prune_delta_segments(dir.path(), 7);
assert!(delta_path(dir.path(), 8).exists());
}
}