use std::collections::{HashMap, HashSet};
use std::fmt;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use arrow::array::{
ArrayRef, BooleanBuilder, FixedSizeBinaryArray, Float64Array, Float64Builder, Int64Builder,
RecordBatch, StringArray, StringBuilder, TimestampMicrosecondArray,
TimestampMicrosecondBuilder, UInt32Array, UInt64Array,
};
use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use graphforge_core::uuid::{Uuid, to_bytes};
use graphforge_core::{GfError, OntologyMode, TypeId};
use graphforge_ir::IrLiteral;
pub type PendingNodeMatch = ([u8; 16], u64, u32, Vec<u32>, HashMap<String, IrLiteral>);
use crate::schemas::{
EXPLORATORY_EDGE_SCHEMA, TOPOLOGY_NODES_SCHEMA, TYPED_EDGE_SCHEMA, uuid_field,
};
const EXPLORATORY_STEM: &str = "_exploratory";
const UNTYPED_STEM: &str = "_untyped";
const NODE_PROPERTY_UUID_FIELD: &str = "node_uuid";
const EDGE_PROPERTY_UUID_FIELD: &str = "edge_uuid";
fn io_err(e: &std::io::Error) -> GfError {
GfError::Storage(e.to_string())
}
fn pq_err(e: impl fmt::Display) -> GfError {
GfError::Storage(e.to_string())
}
fn max_u64_column(batches: &[RecordBatch], col: &str) -> u64 {
use arrow::array::Array;
let mut max = 0u64;
for batch in batches {
if let Some(c) = batch.column_by_name(col)
&& let Some(ids) = c.as_any().downcast_ref::<UInt64Array>()
{
for i in 0..ids.len() {
if !ids.is_null(i) {
max = max.max(ids.value(i));
}
}
}
}
max
}
struct NodeRow {
node_uuid: [u8; 16],
node_id: u64,
type_id: u32,
type_ids: Vec<u32>,
}
struct EdgeRow {
edge_uuid: [u8; 16],
src_uuid: [u8; 16],
dst_uuid: [u8; 16],
edge_id: u64,
src_id: u64,
dst_id: u64,
rel_type_name: Option<String>,
}
struct PropRow {
node_uuid: [u8; 16],
props: HashMap<String, IrLiteral>,
}
struct EdgePropRow {
edge_uuid: [u8; 16],
props: HashMap<String, IrLiteral>,
}
trait PropRowLike {
fn uuid_bytes(&self) -> &[u8; 16];
fn props(&self) -> &HashMap<String, IrLiteral>;
fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral>;
fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self;
}
impl PropRowLike for PropRow {
fn uuid_bytes(&self) -> &[u8; 16] {
&self.node_uuid
}
fn props(&self) -> &HashMap<String, IrLiteral> {
&self.props
}
fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral> {
&mut self.props
}
fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self {
Self {
node_uuid: uuid,
props,
}
}
}
impl PropRowLike for EdgePropRow {
fn uuid_bytes(&self) -> &[u8; 16] {
&self.edge_uuid
}
fn props(&self) -> &HashMap<String, IrLiteral> {
&self.props
}
fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral> {
&mut self.props
}
fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self {
Self {
edge_uuid: uuid,
props,
}
}
}
pub struct GraphWriter {
dir: PathBuf,
mode: OntologyMode,
now_micros: i64,
next_node_id: u64,
next_edge_id: u64,
uuid_to_node_id: HashMap<[u8; 16], u64>,
nodes: Vec<NodeRow>,
edges: HashMap<String, Vec<EdgeRow>>,
properties: HashMap<String, Vec<PropRow>>,
edge_properties: HashMap<String, Vec<EdgePropRow>>,
pending_delta: Vec<crate::adjacency_delta::DeltaEdge>,
}
impl GraphWriter {
pub fn open(dir: &Path, mode: OntologyMode) -> Result<Self, GfError> {
let now_micros = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| i64::try_from(d.as_micros()).unwrap_or(i64::MAX));
Self::open_at(dir, mode, now_micros)
}
pub fn open_at(dir: &Path, mode: OntologyMode, now_micros: i64) -> Result<Self, GfError> {
fs::create_dir_all(dir).map_err(|e| io_err(&e))?;
let max_node_id =
max_u64_column(&crate::catalog::read_nodes(dir).map_err(pq_err)?, "node_id");
let max_edge_id = crate::catalog::max_edge_id(dir).map_err(pq_err)?;
Ok(Self {
dir: dir.to_path_buf(),
mode,
now_micros,
next_node_id: max_node_id + 1,
next_edge_id: max_edge_id + 1,
uuid_to_node_id: HashMap::new(),
nodes: Vec::new(),
edges: HashMap::new(),
properties: HashMap::new(),
edge_properties: HashMap::new(),
pending_delta: Vec::new(),
})
}
pub fn create_node(&mut self, node_uuid: Uuid, type_id: TypeId) -> Result<u64, GfError> {
self.create_node_with_labels(node_uuid, &[type_id])
}
pub fn create_node_with_labels(
&mut self,
node_uuid: Uuid,
type_ids: &[TypeId],
) -> Result<u64, GfError> {
let bytes = to_bytes(&node_uuid);
let node_id = self.next_node_id;
self.next_node_id += 1;
self.uuid_to_node_id.insert(bytes, node_id);
self.nodes.push(NodeRow {
node_uuid: bytes,
node_id,
type_id: type_ids.first().map_or(u32::MAX, |id| id.0),
type_ids: type_ids.iter().map(|id| id.0).collect(),
});
Ok(node_id)
}
pub fn register_existing_node(&mut self, node_uuid: Uuid, node_id: u64) {
self.uuid_to_node_id.insert(to_bytes(&node_uuid), node_id);
}
#[must_use]
pub fn node_id_for_uuid(&self, node_uuid: &Uuid) -> Option<u64> {
self.uuid_to_node_id.get(&to_bytes(node_uuid)).copied()
}
pub fn create_edge(
&mut self,
edge_uuid: Uuid,
rel_type: &str,
src_uuid: &Uuid,
dst_uuid: &Uuid,
) -> Result<u64, GfError> {
let src_bytes = to_bytes(src_uuid);
let dst_bytes = to_bytes(dst_uuid);
let src_id = *self.uuid_to_node_id.get(&src_bytes).ok_or_else(|| {
GfError::Storage(format!(
"create_edge: source {} has no node_id; call create_node first",
graphforge_core::uuid::to_string(src_uuid)
))
})?;
let dst_id = *self.uuid_to_node_id.get(&dst_bytes).ok_or_else(|| {
GfError::Storage(format!(
"create_edge: destination {} has no node_id; call create_node first",
graphforge_core::uuid::to_string(dst_uuid)
))
})?;
let edge_id = self.next_edge_id;
self.next_edge_id += 1;
let (stem, rel_type_name) = match self.mode {
OntologyMode::Exploratory => (EXPLORATORY_STEM.to_owned(), Some(rel_type.to_owned())),
OntologyMode::Advisory | OntologyMode::Strict => (rel_type.to_owned(), None),
};
self.edges.entry(stem).or_default().push(EdgeRow {
edge_uuid: to_bytes(&edge_uuid),
src_uuid: src_bytes,
dst_uuid: dst_bytes,
edge_id,
src_id,
dst_id,
rel_type_name,
});
Ok(edge_id)
}
pub fn set_properties(
&mut self,
node_uuid: &Uuid,
entity_type: Option<&str>,
props: HashMap<String, IrLiteral>,
) -> Result<(), GfError> {
let stem = match (self.mode, entity_type) {
(OntologyMode::Advisory | OntologyMode::Strict, Some(t)) => t.to_owned(),
_ => UNTYPED_STEM.to_owned(),
};
self.properties.entry(stem).or_default().push(PropRow {
node_uuid: to_bytes(node_uuid),
props,
});
Ok(())
}
pub fn set_edge_properties(
&mut self,
edge_uuid: &Uuid,
rel_type: Option<&str>,
props: HashMap<String, IrLiteral>,
) -> Result<(), GfError> {
let stem = rel_type.unwrap_or(UNTYPED_STEM).to_owned();
self.edge_properties
.entry(stem)
.or_default()
.push(EdgePropRow {
edge_uuid: to_bytes(edge_uuid),
props,
});
Ok(())
}
#[must_use]
pub fn contains_pending_node(&self, node_uuid: &[u8; 16]) -> bool {
self.nodes.iter().any(|r| &r.node_uuid == node_uuid)
}
#[must_use]
pub fn pending_node_labels(&self, targets: &HashSet<[u8; 16]>) -> HashSet<u32> {
self.nodes
.iter()
.filter(|row| targets.contains(&row.node_uuid))
.flat_map(|row| row.type_ids.iter().copied())
.collect()
}
pub fn pending_nodes_batch(&self) -> Result<RecordBatch, GfError> {
let n = self.nodes.len();
if n == 0 {
return Ok(RecordBatch::new_empty(TOPOLOGY_NODES_SCHEMA.clone()));
}
let uuids =
FixedSizeBinaryArray::try_from_iter(self.nodes.iter().map(|r| r.node_uuid.to_vec()))
.map_err(pq_err)?;
let node_ids = UInt64Array::from(self.nodes.iter().map(|r| r.node_id).collect::<Vec<_>>());
let type_ids = UInt32Array::from(self.nodes.iter().map(|r| r.type_id).collect::<Vec<_>>());
let nullable_label_sets =
arrow::array::ListArray::from_iter_primitive::<arrow::datatypes::UInt32Type, _, _>(
self.nodes
.iter()
.map(|row| Some(row.type_ids.iter().copied().map(Some))),
);
let label_sets = arrow::array::ListArray::new(
Arc::new(Field::new("item", DataType::UInt32, false)),
nullable_label_sets.offsets().clone(),
nullable_label_sets.values().clone(),
None,
);
let ts = self.timestamp_array(n);
RecordBatch::try_new(
TOPOLOGY_NODES_SCHEMA.clone(),
vec![
Arc::new(uuids),
Arc::new(node_ids),
Arc::new(type_ids),
Arc::new(label_sets),
Arc::new(ts.clone()),
Arc::new(ts),
],
)
.map_err(pq_err)
}
#[must_use]
#[allow(clippy::type_complexity)]
pub fn find_pending_node(
&self,
labels: &[u32],
properties: &[(String, IrLiteral)],
) -> Option<PendingNodeMatch> {
self.find_pending_nodes(labels, properties)
.into_iter()
.next()
}
#[must_use]
pub fn find_pending_nodes(
&self,
labels: &[u32],
properties: &[(String, IrLiteral)],
) -> Vec<PendingNodeMatch> {
self.nodes
.iter()
.filter_map(|node| {
if !labels.iter().all(|wanted| node.type_ids.contains(wanted)) {
return None;
}
let props = self
.properties
.values()
.flatten()
.filter(|row| row.node_uuid == node.node_uuid)
.flat_map(|row| {
row.props
.iter()
.map(|(key, value)| (key.clone(), value.clone()))
})
.collect::<HashMap<_, _>>();
properties
.iter()
.all(|(name, value)| props.get(name) == Some(value))
.then(|| {
(
node.node_uuid,
node.node_id,
node.type_id,
node.type_ids.clone(),
props,
)
})
})
.collect()
}
#[must_use]
pub fn contains_pending_edge(&self, edge_uuid: &[u8; 16]) -> bool {
self.edges
.values()
.any(|rows| rows.iter().any(|r| &r.edge_uuid == edge_uuid))
}
#[must_use]
#[allow(clippy::type_complexity)]
pub fn find_pending_edge(
&self,
rel_type: &str,
src: &[u8; 16],
dst: &[u8; 16],
undirected: bool,
properties: &[(String, IrLiteral)],
) -> Option<([u8; 16], [u8; 16], [u8; 16], HashMap<String, IrLiteral>)> {
self.edges.iter().find_map(|(stem, edges)| {
edges.iter().find_map(|edge| {
let edge_type = edge.rel_type_name.as_deref().unwrap_or(stem);
let direct = edge.src_uuid == *src && edge.dst_uuid == *dst;
let reverse = edge.src_uuid == *dst && edge.dst_uuid == *src;
if edge_type != rel_type || !(direct || undirected && reverse) {
return None;
}
let props = self
.edge_properties
.values()
.flatten()
.filter(|row| row.edge_uuid == edge.edge_uuid)
.flat_map(|row| {
row.props
.iter()
.map(|(key, value)| (key.clone(), value.clone()))
})
.collect::<HashMap<_, _>>();
properties
.iter()
.all(|(name, value)| props.get(name) == Some(value))
.then_some((edge.edge_uuid, edge.src_uuid, edge.dst_uuid, props))
})
})
}
#[must_use]
pub fn pending_incident_edge_uuids<S: std::hash::BuildHasher>(
&self,
nodes: &HashSet<[u8; 16], S>,
) -> Vec<[u8; 16]> {
self.edges
.values()
.flatten()
.filter(|r| nodes.contains(&r.src_uuid) || nodes.contains(&r.dst_uuid))
.map(|r| r.edge_uuid)
.collect()
}
pub fn cancel_nodes<S: std::hash::BuildHasher>(
&mut self,
targets: &HashSet<[u8; 16], S>,
) -> u64 {
let before = self.nodes.len();
self.nodes.retain(|r| !targets.contains(&r.node_uuid));
let dropped = (before - self.nodes.len()) as u64;
self.properties.retain(|_, rows| {
rows.retain(|r| !targets.contains(&r.node_uuid));
!rows.is_empty()
});
self.uuid_to_node_id
.retain(|uuid, _| !targets.contains(uuid));
dropped
}
pub fn cancel_edges<S: std::hash::BuildHasher>(
&mut self,
targets: &HashSet<[u8; 16], S>,
) -> u64 {
let mut dropped = 0u64;
self.edges.retain(|_, rows| {
let before = rows.len();
rows.retain(|r| !targets.contains(&r.edge_uuid));
dropped += (before - rows.len()) as u64;
!rows.is_empty()
});
self.edge_properties.retain(|_, rows| {
rows.retain(|r| !targets.contains(&r.edge_uuid));
!rows.is_empty()
});
dropped
}
pub fn merge_pending_node_props(
&mut self,
node_uuid: &[u8; 16],
entity_type: Option<&str>,
props: HashMap<String, IrLiteral>,
) {
let stem = match (self.mode, entity_type) {
(OntologyMode::Advisory | OntologyMode::Strict, Some(t)) => t.to_owned(),
_ => UNTYPED_STEM.to_owned(),
};
let rows = self.properties.entry(stem).or_default();
if let Some(row) = rows.iter_mut().find(|r| &r.node_uuid == node_uuid) {
row.props.extend(props);
} else {
rows.push(PropRow {
node_uuid: *node_uuid,
props,
});
}
}
pub fn add_pending_node_labels(&mut self, node_uuid: &[u8; 16], labels: &[u32]) -> u64 {
let Some(row) = self
.nodes
.iter_mut()
.find(|row| &row.node_uuid == node_uuid)
else {
return 0;
};
let before = row.type_ids.len();
row.type_ids.extend(labels.iter().copied());
row.type_ids.sort_unstable();
row.type_ids.dedup();
(row.type_ids.len() - before) as u64
}
pub fn remove_pending_node_labels(&mut self, node_uuid: &[u8; 16], labels: &[u32]) -> u64 {
let Some(row) = self
.nodes
.iter_mut()
.find(|row| &row.node_uuid == node_uuid)
else {
return 0;
};
let before = row.type_ids.len();
row.type_ids.retain(|label| !labels.contains(label));
(before - row.type_ids.len()) as u64
}
pub fn merge_pending_edge_props(
&mut self,
edge_uuid: &[u8; 16],
rel_type: Option<&str>,
props: HashMap<String, IrLiteral>,
) {
let stem = rel_type.unwrap_or(UNTYPED_STEM).to_owned();
let rows = self.edge_properties.entry(stem).or_default();
if let Some(row) = rows.iter_mut().find(|r| &r.edge_uuid == edge_uuid) {
row.props.extend(props);
} else {
rows.push(EdgePropRow {
edge_uuid: *edge_uuid,
props,
});
}
}
pub fn remove_pending_node_props(&mut self, node_uuid: &[u8; 16], keys: &HashSet<String>) {
for rows in self.properties.values_mut() {
for row in rows.iter_mut().filter(|r| &r.node_uuid == node_uuid) {
row.props.retain(|k, _| !keys.contains(k));
}
}
}
pub fn remove_pending_edge_props(&mut self, edge_uuid: &[u8; 16], keys: &HashSet<String>) {
for rows in self.edge_properties.values_mut() {
for row in rows.iter_mut().filter(|r| &r.edge_uuid == edge_uuid) {
row.props.retain(|k, _| !keys.contains(k));
}
}
}
pub fn flush(&mut self) -> Result<(), GfError> {
let mut staged = RewriteBatch::new();
self.flush_into(&mut staged)?;
let pending = self.take_pending_delta();
if let Some(generation) = crate::generation::commit_topology_aware(staged, &self.dir)? {
self.write_segment_best_effort(generation, &pending);
}
Ok(())
}
#[must_use]
pub fn take_pending_delta(&mut self) -> Vec<crate::adjacency_delta::DeltaEdge> {
let mut edges = std::mem::take(&mut self.pending_delta);
edges.sort_unstable_by_key(|e| e.edge_id);
edges
}
pub fn write_segment_best_effort(
&self,
generation: u64,
edges: &[crate::adjacency_delta::DeltaEdge],
) {
if crate::adjacency::adjacency_dir(&self.dir).exists() {
let _ = crate::adjacency_delta::write_delta_segment(&self.dir, generation, edges);
}
}
pub fn flush_into(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
self.flush_nodes(staged)?;
self.flush_edges(staged)?;
self.flush_properties(staged)?;
self.flush_edge_properties(staged)?;
Ok(())
}
fn flush_nodes(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
if self.nodes.is_empty() {
return Ok(());
}
let topology = self.dir.join("topology");
fs::create_dir_all(&topology).map_err(|e| io_err(&e))?;
let batch = self.pending_nodes_batch()?;
let path = topology.join("nodes.parquet");
let read_path = staged
.staged_temp(&path)
.map_or_else(|| path.clone(), Path::to_path_buf);
let existing = crate::catalog::normalize_topology_nodes(
crate::catalog::read_parquet_or_empty(&read_path, TOPOLOGY_NODES_SCHEMA.clone())
.map_err(pq_err)?,
)
.map_err(pq_err)?;
let merged = concat_with_existing(&TOPOLOGY_NODES_SCHEMA, existing, batch)?;
staged.restage(&path, TOPOLOGY_NODES_SCHEMA.clone(), &merged)?;
self.nodes.clear();
Ok(())
}
fn flush_edges(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
if self.edges.is_empty() {
return Ok(());
}
let edges_dir = self.dir.join("topology").join("edges");
fs::create_dir_all(&edges_dir).map_err(|e| io_err(&e))?;
let buffered: Vec<(String, Vec<EdgeRow>)> = self.edges.drain().collect();
for (stem, rows) in buffered {
let exploratory = stem == EXPLORATORY_STEM;
for r in &rows {
self.pending_delta.push(crate::adjacency_delta::DeltaEdge {
rel_type_name: if exploratory {
r.rel_type_name.clone().unwrap_or_default()
} else {
stem.clone()
},
edge_id: r.edge_id,
src_id: r.src_id,
dst_id: r.dst_id,
});
}
let schema = if exploratory {
EXPLORATORY_EDGE_SCHEMA.clone()
} else {
TYPED_EDGE_SCHEMA.clone()
};
let batch = self.edge_batch(&rows, &schema, exploratory)?;
let path = edges_dir.join(format!("{stem}.parquet"));
let read_path = staged
.staged_temp(&path)
.map_or_else(|| path.clone(), Path::to_path_buf);
let existing = crate::catalog::read_parquet_or_empty(&read_path, schema.clone())
.map_err(pq_err)?;
let merged = concat_with_existing(&schema, existing, batch)?;
staged.restage(&path, schema, &merged)?;
}
Ok(())
}
fn edge_batch(
&self,
rows: &[EdgeRow],
schema: &SchemaRef,
exploratory: bool,
) -> Result<RecordBatch, GfError> {
let edge_uuids =
FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.edge_uuid.to_vec()))
.map_err(pq_err)?;
let src_uuids =
FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.src_uuid.to_vec()))
.map_err(pq_err)?;
let dst_uuids =
FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.dst_uuid.to_vec()))
.map_err(pq_err)?;
let edge_ids = UInt64Array::from(rows.iter().map(|r| r.edge_id).collect::<Vec<_>>());
let src_ids = UInt64Array::from(rows.iter().map(|r| r.src_id).collect::<Vec<_>>());
let dst_ids = UInt64Array::from(rows.iter().map(|r| r.dst_id).collect::<Vec<_>>());
let ts = self.timestamp_array(rows.len());
let mut cols: Vec<ArrayRef> = vec![
Arc::new(edge_uuids),
Arc::new(src_uuids),
Arc::new(dst_uuids),
Arc::new(edge_ids),
Arc::new(src_ids),
Arc::new(dst_ids),
Arc::new(ts),
];
if exploratory {
let names = StringArray::from(
rows.iter()
.map(|r| r.rel_type_name.clone().unwrap_or_default())
.collect::<Vec<_>>(),
);
cols.push(Arc::new(names));
}
RecordBatch::try_new(schema.clone(), cols).map_err(pq_err)
}
fn flush_properties(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
if self.properties.is_empty() {
return Ok(());
}
let buffered: Vec<(String, Vec<PropRow>)> = self.properties.drain().collect();
for (stem, new_rows) in buffered {
let existing = read_props_through(staged, &node_props_path(&self.dir, &stem))?;
let mut rows = decode_property_rows(&existing)?;
rows.extend(new_rows);
let (schema, cols) = build_property_columns(&stem, &rows)?;
stage_property_file(staged, &self.dir, "properties", &stem, schema, cols)?;
}
Ok(())
}
fn flush_edge_properties(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
if self.edge_properties.is_empty() {
return Ok(());
}
let buffered: Vec<(String, Vec<EdgePropRow>)> = self.edge_properties.drain().collect();
for (stem, new_rows) in buffered {
let existing = read_props_through(staged, &edge_props_path(&self.dir, &stem))?;
let mut rows = decode_edge_property_rows(&existing)?;
rows.extend(new_rows);
let (schema, cols) = build_property_columns_keyed(
EDGE_PROPERTY_UUID_FIELD,
"graphforge.rel_type",
&stem,
&rows,
)?;
stage_property_file(staged, &self.dir, "edge_properties", &stem, schema, cols)?;
}
Ok(())
}
fn timestamp_array(&self, n: usize) -> TimestampMicrosecondArray {
TimestampMicrosecondArray::from(vec![self.now_micros; n])
.with_timezone_opt(Some(Arc::from("UTC")))
}
}
#[derive(Clone, PartialEq, Eq)]
enum ColType {
Int,
Float,
Bool,
Str,
HetScalar,
Duration,
DateTime,
Date,
LocalDateTime,
Time,
ZonedTime,
ZonedDateTime,
List(Box<ColType>),
}
impl ColType {
fn of(lit: &IrLiteral) -> Option<Self> {
match lit {
IrLiteral::Null
| IrLiteral::Map(_)
| IrLiteral::Uuid(_) => None,
IrLiteral::Int(_) => Some(Self::Int),
IrLiteral::Float(_) => Some(Self::Float),
IrLiteral::Bool(_) => Some(Self::Bool),
IrLiteral::Str(_) => Some(Self::Str),
IrLiteral::Duration { .. } => Some(Self::Duration),
IrLiteral::DateTime(_) => Some(Self::DateTime),
IrLiteral::Date(_) => Some(Self::Date),
IrLiteral::LocalDateTime { .. } => Some(Self::LocalDateTime),
IrLiteral::Time(_) => Some(Self::Time),
IrLiteral::ZonedTime { .. } => Some(Self::ZonedTime),
IrLiteral::ZonedDateTime { .. } => Some(Self::ZonedDateTime),
IrLiteral::List(items) => items
.iter()
.find_map(Self::of)
.map(|inner| Self::List(Box::new(inner))),
}
}
fn data_type(&self) -> DataType {
match self {
Self::Int => DataType::Int64,
Self::Float => DataType::Float64,
Self::Bool => DataType::Boolean,
Self::Str => DataType::Utf8,
Self::HetScalar => DataType::Struct(heterogeneous_scalar_fields()),
Self::Duration => DataType::Struct(crate::schemas::duration_struct_fields()),
Self::DateTime => DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
Self::Date => DataType::Struct(crate::schemas::date_struct_fields()),
Self::LocalDateTime => DataType::Struct(crate::schemas::localdatetime_struct_fields()),
Self::Time => DataType::Time64(TimeUnit::Nanosecond),
Self::ZonedTime => DataType::Struct(crate::schemas::time_struct_fields()),
Self::ZonedDateTime => DataType::Struct(crate::schemas::datetime_struct_fields()),
Self::List(inner) => {
DataType::List(Arc::new(Field::new("item", inner.data_type(), true)))
}
}
}
fn is_scalar(&self) -> bool {
matches!(
self,
Self::Int | Self::Float | Self::Bool | Self::Str | Self::HetScalar
)
}
}
fn heterogeneous_scalar_fields() -> arrow::datatypes::Fields {
arrow::datatypes::Fields::from(vec![
Field::new("__het_tag", DataType::Int8, false),
Field::new("__het_int", DataType::Int64, true),
Field::new("__het_float", DataType::Float64, true),
Field::new("__het_str", DataType::Utf8, true),
Field::new("__het_bool", DataType::Boolean, true),
])
}
fn build_property_columns(
entity_type: &str,
rows: &[PropRow],
) -> Result<(Schema, Vec<ArrayRef>), GfError> {
build_property_columns_keyed(
NODE_PROPERTY_UUID_FIELD,
"graphforge.entity_type",
entity_type,
rows,
)
}
fn build_property_columns_keyed<R: PropRowLike>(
uuid_field_name: &str,
meta_key: &str,
meta_value: &str,
rows: &[R],
) -> Result<(Schema, Vec<ArrayRef>), GfError> {
let mut order: Vec<String> = Vec::new();
let mut seen: HashSet<String> = HashSet::new();
let mut col_types: HashMap<String, ColType> = HashMap::new();
for row in rows {
for (name, lit) in row.props() {
reject_map_property_value(name, lit)?;
if seen.insert(name.clone()) {
order.push(name.clone());
}
if let Some(t) = ColType::of(lit) {
col_types
.entry(name.clone())
.and_modify(|existing| {
if *existing != t {
*existing = if existing.is_scalar() && t.is_scalar() {
ColType::HetScalar
} else {
ColType::Str
};
}
})
.or_insert(t);
}
}
}
let mut fields: Vec<Field> = vec![uuid_field(uuid_field_name)];
for name in &order {
let ct = col_types.get(name).cloned().unwrap_or(ColType::Str);
fields.push(Field::new(name, ct.data_type(), true));
}
let meta: HashMap<String, String> = [(meta_key.to_owned(), meta_value.to_owned())]
.into_iter()
.collect();
let schema = Schema::new(fields).with_metadata(meta);
let mut cols: Vec<ArrayRef> = Vec::with_capacity(order.len() + 1);
let uuids = FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.uuid_bytes().to_vec()))
.map_err(pq_err)?;
cols.push(Arc::new(uuids));
for name in &order {
let ct = col_types.get(name).cloned().unwrap_or(ColType::Str);
cols.push(build_property_array(name, ct, rows));
}
Ok((schema, cols))
}
fn reject_map_property_value(name: &str, lit: &IrLiteral) -> Result<(), GfError> {
if contains_uuid_literal(lit) {
return Err(GfError::Validation(format!(
"property `{name}` cannot store typed UUID query parameters"
)));
}
if contains_map_literal(lit) {
return Err(GfError::Storage(format!(
"property `{name}` cannot store map values"
)));
}
Ok(())
}
fn contains_map_literal(lit: &IrLiteral) -> bool {
match lit {
IrLiteral::Map(_) => true,
IrLiteral::List(items) => items.iter().any(contains_map_literal),
_ => false,
}
}
fn contains_uuid_literal(lit: &IrLiteral) -> bool {
match lit {
IrLiteral::Uuid(_) => true,
IrLiteral::List(items) => items.iter().any(contains_uuid_literal),
IrLiteral::Map(entries) => entries
.iter()
.any(|(_, value)| contains_uuid_literal(value)),
_ => false,
}
}
#[allow(
clippy::too_many_lines,
reason = "one builder arm per ColType; the per-type append loops read clearest inline"
)]
fn build_property_array<R: PropRowLike>(name: &str, ct: ColType, rows: &[R]) -> ArrayRef {
match ct {
ColType::Int => {
let mut b = Int64Builder::new();
for row in rows {
match row.props().get(name) {
Some(IrLiteral::Int(v)) => b.append_value(*v),
_ => b.append_null(),
}
}
Arc::new(b.finish())
}
ColType::Float => {
let mut b = Float64Builder::new();
for row in rows {
match row.props().get(name) {
Some(IrLiteral::Float(v)) => b.append_value(*v),
_ => b.append_null(),
}
}
Arc::new(b.finish())
}
ColType::Bool => {
let mut b = BooleanBuilder::new();
for row in rows {
match row.props().get(name) {
Some(IrLiteral::Bool(v)) => b.append_value(*v),
_ => b.append_null(),
}
}
Arc::new(b.finish())
}
ColType::Duration => {
use arrow::array::{Int64Builder, StructArray};
use arrow::buffer::NullBuffer;
let (mut mb, mut db, mut sb, mut nb) = (
Int64Builder::new(),
Int64Builder::new(),
Int64Builder::new(),
Int64Builder::new(),
);
let mut valid = Vec::with_capacity(rows.len());
for row in rows {
if let Some(IrLiteral::Duration {
months,
days,
seconds,
nanos,
}) = row.props().get(name)
{
mb.append_value(*months);
db.append_value(*days);
sb.append_value(*seconds);
nb.append_value(*nanos);
valid.push(true);
} else {
mb.append_null();
db.append_null();
sb.append_null();
nb.append_null();
valid.push(false);
}
}
Arc::new(StructArray::new(
crate::schemas::duration_struct_fields(),
vec![
Arc::new(mb.finish()),
Arc::new(db.finish()),
Arc::new(sb.finish()),
Arc::new(nb.finish()),
],
Some(NullBuffer::from(valid)),
))
}
ColType::DateTime => {
let mut b = TimestampMicrosecondBuilder::new();
for row in rows {
match row.props().get(name) {
Some(IrLiteral::DateTime(v)) => b.append_value(*v),
_ => b.append_null(),
}
}
Arc::new(b.finish().with_timezone_opt(Some(Arc::from("UTC"))))
}
ColType::Date => {
use arrow::array::{Int64Builder, StructArray};
use arrow::buffer::NullBuffer;
let mut b = Int64Builder::new();
let mut valid = Vec::with_capacity(rows.len());
for row in rows {
if let Some(IrLiteral::Date(v)) = row.props().get(name) {
b.append_value(*v);
valid.push(true);
} else {
b.append_null();
valid.push(false);
}
}
Arc::new(StructArray::new(
crate::schemas::date_struct_fields(),
vec![Arc::new(b.finish())],
Some(NullBuffer::from(valid)),
))
}
ColType::LocalDateTime => {
use arrow::array::{Int64Builder, StructArray, Time64NanosecondBuilder};
use arrow::buffer::NullBuffer;
let (mut date_b, mut time_b) = (Int64Builder::new(), Time64NanosecondBuilder::new());
let mut valid = Vec::with_capacity(rows.len());
for row in rows {
if let Some(IrLiteral::LocalDateTime { days, nanos }) = row.props().get(name) {
date_b.append_value(*days);
time_b.append_value(*nanos);
valid.push(true);
} else {
date_b.append_null();
time_b.append_null();
valid.push(false);
}
}
Arc::new(StructArray::new(
crate::schemas::localdatetime_struct_fields(),
vec![Arc::new(date_b.finish()), Arc::new(time_b.finish())],
Some(NullBuffer::from(valid)),
))
}
ColType::Time => {
use arrow::array::Time64NanosecondBuilder;
let mut b = Time64NanosecondBuilder::new();
for row in rows {
match row.props().get(name) {
Some(IrLiteral::Time(v)) => b.append_value(*v),
_ => b.append_null(),
}
}
Arc::new(b.finish())
}
ColType::ZonedTime => {
use arrow::array::{Int32Builder, StructArray, Time64NanosecondBuilder};
use arrow::buffer::NullBuffer;
let (mut time_b, mut off_b) = (Time64NanosecondBuilder::new(), Int32Builder::new());
let mut valid = Vec::with_capacity(rows.len());
for row in rows {
if let Some(IrLiteral::ZonedTime { nanos, offset }) = row.props().get(name) {
time_b.append_value(*nanos);
off_b.append_value(*offset);
valid.push(true);
} else {
time_b.append_null();
off_b.append_null();
valid.push(false);
}
}
Arc::new(StructArray::new(
crate::schemas::time_struct_fields(),
vec![Arc::new(time_b.finish()), Arc::new(off_b.finish())],
Some(NullBuffer::from(valid)),
))
}
ColType::ZonedDateTime => {
use arrow::array::{
Int32Builder, Int64Builder, StringBuilder, StructArray, Time64NanosecondBuilder,
};
use arrow::buffer::NullBuffer;
let (mut date_b, mut time_b, mut off_b, mut zone_b) = (
Int64Builder::new(),
Time64NanosecondBuilder::new(),
Int32Builder::new(),
StringBuilder::new(),
);
let mut valid = Vec::with_capacity(rows.len());
for row in rows {
if let Some(IrLiteral::ZonedDateTime {
days,
nanos,
offset,
zone,
}) = row.props().get(name)
{
date_b.append_value(*days);
time_b.append_value(*nanos);
off_b.append_value(*offset);
zone_b.append_option(zone.as_deref());
valid.push(true);
} else {
date_b.append_null();
time_b.append_null();
off_b.append_null();
zone_b.append_null();
valid.push(false);
}
}
Arc::new(StructArray::new(
crate::schemas::datetime_struct_fields(),
vec![
Arc::new(date_b.finish()),
Arc::new(time_b.finish()),
Arc::new(off_b.finish()),
Arc::new(zone_b.finish()),
],
Some(NullBuffer::from(valid)),
))
}
ColType::Str => {
let mut b = StringBuilder::new();
for row in rows {
match row.props().get(name) {
Some(IrLiteral::Null) | None => b.append_null(),
Some(other) => b.append_value(literal_to_string(other)),
}
}
Arc::new(b.finish())
}
ColType::HetScalar => build_heterogeneous_scalar_array(name, rows),
ColType::List(inner) => {
use arrow::array::ListArray;
use arrow::buffer::{NullBuffer, OffsetBuffer};
let mut elem_rows: Vec<PropRow> = Vec::new();
let mut offsets: Vec<i32> = vec![0];
let mut valid = Vec::with_capacity(rows.len());
for row in rows {
if let Some(IrLiteral::List(items)) = row.props().get(name) {
for it in items {
let mut props = HashMap::with_capacity(1);
props.insert("item".to_string(), it.clone());
elem_rows.push(PropRow {
node_uuid: [0u8; 16],
props,
});
}
valid.push(true);
} else {
valid.push(false);
}
offsets.push(i32::try_from(elem_rows.len()).unwrap_or(i32::MAX));
}
let child = build_property_array("item", (*inner).clone(), &elem_rows);
let field = Arc::new(Field::new("item", inner.data_type(), true));
Arc::new(ListArray::new(
field,
OffsetBuffer::new(offsets.into()),
child,
Some(NullBuffer::from(valid)),
))
}
}
}
fn build_heterogeneous_scalar_array<R: PropRowLike>(name: &str, rows: &[R]) -> ArrayRef {
use arrow::array::{BooleanBuilder, Float64Builder, Int8Builder, Int64Builder, StructArray};
use arrow::buffer::NullBuffer;
let mut tags = Int8Builder::new();
let mut ints = Int64Builder::new();
let mut floats = Float64Builder::new();
let mut strings = StringBuilder::new();
let mut bools = BooleanBuilder::new();
let mut valid = Vec::with_capacity(rows.len());
for row in rows {
let value = row.props().get(name);
let tag = match value {
Some(IrLiteral::Int(value)) => {
ints.append_value(*value);
floats.append_null();
strings.append_null();
bools.append_null();
Some(0)
}
Some(IrLiteral::Float(value)) => {
ints.append_null();
floats.append_value(*value);
strings.append_null();
bools.append_null();
Some(1)
}
Some(IrLiteral::Str(value)) => {
ints.append_null();
floats.append_null();
strings.append_value(value);
bools.append_null();
Some(2)
}
Some(IrLiteral::Bool(value)) => {
ints.append_null();
floats.append_null();
strings.append_null();
bools.append_value(*value);
Some(3)
}
_ => {
ints.append_null();
floats.append_null();
strings.append_null();
bools.append_null();
None
}
};
tags.append_value(tag.unwrap_or_default());
valid.push(tag.is_some());
}
Arc::new(StructArray::new(
heterogeneous_scalar_fields(),
vec![
Arc::new(tags.finish()),
Arc::new(ints.finish()),
Arc::new(floats.finish()),
Arc::new(strings.finish()),
Arc::new(bools.finish()),
],
Some(NullBuffer::from(valid)),
))
}
fn literal_to_string(lit: &IrLiteral) -> String {
match lit {
IrLiteral::Null => String::new(),
IrLiteral::Bool(b) => b.to_string(),
IrLiteral::Int(i) => i.to_string(),
IrLiteral::Float(f) => f.to_string(),
IrLiteral::Str(s) => s.clone(),
IrLiteral::Uuid(bytes) => {
let mut encoded = String::with_capacity(32);
for byte in bytes {
std::fmt::Write::write_fmt(&mut encoded, format_args!("{byte:02x}"))
.expect("writing to a String cannot fail");
}
encoded
}
IrLiteral::Duration {
months,
days,
seconds,
nanos,
} => format!("{months}mo{days}d{seconds}s{nanos}ns"),
IrLiteral::DateTime(t) => t.to_string(),
IrLiteral::Date(d) => d.to_string(),
IrLiteral::LocalDateTime { days, nanos } => format!("{days}d{nanos}ns"),
IrLiteral::Time(nanos) => format!("{nanos}ns"),
IrLiteral::ZonedTime { nanos, offset } => format!("{nanos}ns{offset:+}s"),
IrLiteral::ZonedDateTime {
days,
nanos,
offset,
zone,
} => format!(
"{days}d{nanos}ns{offset:+}s{}",
zone.as_deref().unwrap_or("")
),
IrLiteral::List(items) => {
let parts: Vec<String> = items.iter().map(literal_to_string).collect();
format!("[{}]", parts.join(","))
}
IrLiteral::Map(entries) => {
let parts: Vec<String> = entries
.iter()
.map(|(key, value)| format!("{key}:{}", literal_to_string(value)))
.collect();
format!("{{{}}}", parts.join(","))
}
}
}
fn decode_property_rows(batches: &[RecordBatch]) -> Result<Vec<PropRow>, GfError> {
let mut out = Vec::new();
for batch in batches {
decode_property_batch(batch, NODE_PROPERTY_UUID_FIELD, |node_uuid, props| {
out.push(PropRow { node_uuid, props });
})?;
}
Ok(out)
}
fn decode_edge_property_rows(batches: &[RecordBatch]) -> Result<Vec<EdgePropRow>, GfError> {
let mut out = Vec::new();
for batch in batches {
decode_property_batch(batch, EDGE_PROPERTY_UUID_FIELD, |edge_uuid, props| {
out.push(EdgePropRow { edge_uuid, props });
})?;
}
Ok(out)
}
pub fn read_entity_property_keys(
dir: &Path,
stem: &str,
uuid: &[u8; 16],
is_edge: bool,
) -> Result<HashSet<String>, GfError> {
let path = if is_edge {
edge_props_path(dir, stem)
} else {
node_props_path(dir, stem)
};
let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
return Ok(HashSet::new());
};
let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
let rows = if is_edge {
decode_edge_property_rows(&batches)?
.into_iter()
.map(|row| (row.edge_uuid, row.props))
.collect::<Vec<_>>()
} else {
decode_property_rows(&batches)?
.into_iter()
.map(|row| (row.node_uuid, row.props))
.collect::<Vec<_>>()
};
Ok(rows
.into_iter()
.find_map(|(row_uuid, props)| (row_uuid == *uuid).then(|| props.into_keys().collect()))
.unwrap_or_default())
}
pub fn read_entity_properties(
dir: &Path,
stem: &str,
uuid: &[u8; 16],
is_edge: bool,
) -> Result<HashMap<String, IrLiteral>, GfError> {
let path = if is_edge {
edge_props_path(dir, stem)
} else {
node_props_path(dir, stem)
};
let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
return Ok(HashMap::new());
};
let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
let rows = if is_edge {
decode_edge_property_rows(&batches)?
.into_iter()
.map(|row| (row.edge_uuid, row.props))
.collect::<Vec<_>>()
} else {
decode_property_rows(&batches)?
.into_iter()
.map(|row| (row.node_uuid, row.props))
.collect::<Vec<_>>()
};
Ok(rows
.into_iter()
.find_map(|(row_uuid, props)| (row_uuid == *uuid).then_some(props))
.unwrap_or_default())
}
pub fn read_node_property_rows(
dir: &Path,
stem: &str,
) -> Result<HashMap<[u8; 16], HashMap<String, IrLiteral>>, GfError> {
let path = node_props_path(dir, stem);
if !path.try_exists().map_err(|error| io_err(&error))? {
return Ok(HashMap::new());
}
let file = fs::File::open(&path).map_err(|error| io_err(&error))?;
let schema = ParquetRecordBatchReaderBuilder::try_new(file)
.map_err(pq_err)?
.schema()
.clone();
let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
Ok(decode_property_rows(&batches)?
.into_iter()
.map(|row| (row.node_uuid, row.props))
.collect())
}
pub fn count_entity_properties<S: std::hash::BuildHasher>(
dir: &Path,
targets: &HashSet<[u8; 16], S>,
is_edge: bool,
) -> Result<u64, GfError> {
if targets.is_empty() {
return Ok(0);
}
let property_dir = dir.join(if is_edge {
"edge_properties"
} else {
"properties"
});
let Ok(entries) = std::fs::read_dir(property_dir) else {
return Ok(0);
};
let mut count = 0u64;
for entry in entries {
let path = entry
.map_err(|error| GfError::Storage(error.to_string()))?
.path();
if path
.extension()
.is_none_or(|extension| extension != "parquet")
{
continue;
}
let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
continue;
};
let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
if is_edge {
for row in decode_edge_property_rows(&batches)? {
if targets.contains(&row.edge_uuid) {
count += row.props.len() as u64;
}
}
} else {
for row in decode_property_rows(&batches)? {
if targets.contains(&row.node_uuid) {
count += row.props.len() as u64;
}
}
}
}
Ok(count)
}
#[allow(
clippy::too_many_lines,
reason = "one decode arm per Arrow type, plus per-shape struct dispatch; clearest inline"
)]
fn decode_property_batch(
batch: &RecordBatch,
uuid_field_name: &str,
mut emit: impl FnMut([u8; 16], HashMap<String, IrLiteral>),
) -> Result<(), GfError> {
use arrow::array::Array;
let schema = batch.schema();
let uuid_col = batch
.column_by_name(uuid_field_name)
.and_then(|c| c.as_any().downcast_ref::<FixedSizeBinaryArray>())
.ok_or_else(|| {
GfError::Storage(format!("property file missing {uuid_field_name} column"))
})?;
for r in 0..batch.num_rows() {
let mut uuid = [0u8; 16];
uuid.copy_from_slice(uuid_col.value(r));
let mut props: HashMap<String, IrLiteral> = HashMap::new();
for (c, field) in schema.fields().iter().enumerate() {
if field.name() == uuid_field_name {
continue;
}
let col = batch.column(c);
if col.is_null(r) {
continue; }
let lit = decode_value(col, field, r)?;
props.insert(field.name().clone(), lit);
}
emit(uuid, props);
}
Ok(())
}
#[allow(
clippy::too_many_lines,
reason = "one decode arm per Arrow type, plus per-shape struct dispatch; clearest inline"
)]
fn decode_value(
col: &arrow::array::ArrayRef,
field: &arrow::datatypes::Field,
r: usize,
) -> Result<IrLiteral, GfError> {
use arrow::array::{
Array, BooleanArray, Int32Array, Int64Array, ListArray, StructArray, Time64NanosecondArray,
};
Ok(match field.data_type() {
DataType::Int64 => IrLiteral::Int(downcast::<Int64Array>(col, field)?.value(r)),
DataType::Float64 => IrLiteral::Float(downcast::<Float64Array>(col, field)?.value(r)),
DataType::Boolean => IrLiteral::Bool(downcast::<BooleanArray>(col, field)?.value(r)),
DataType::Utf8 => IrLiteral::Str(downcast::<StringArray>(col, field)?.value(r).to_owned()),
DataType::Struct(fields) => {
let s = downcast::<StructArray>(col, field)?;
let names: Vec<&str> = fields.iter().map(|f| f.name().as_str()).collect();
match names.as_slice() {
[
"__het_tag",
"__het_int",
"__het_float",
"__het_str",
"__het_bool",
] => {
let tag = s
.column(0)
.as_any()
.downcast_ref::<arrow::array::Int8Array>()
.ok_or_else(|| GfError::Storage("heterogeneous tag not Int8".into()))?
.value(r);
match tag {
0 => IrLiteral::Int(
s.column(1)
.as_any()
.downcast_ref::<Int64Array>()
.ok_or_else(|| {
GfError::Storage("heterogeneous int not Int64".into())
})?
.value(r),
),
1 => IrLiteral::Float(
s.column(2)
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| {
GfError::Storage("heterogeneous float not Float64".into())
})?
.value(r),
),
2 => IrLiteral::Str(
s.column(3)
.as_any()
.downcast_ref::<StringArray>()
.ok_or_else(|| {
GfError::Storage("heterogeneous string not Utf8".into())
})?
.value(r)
.to_owned(),
),
3 => IrLiteral::Bool(
s.column(4)
.as_any()
.downcast_ref::<BooleanArray>()
.ok_or_else(|| {
GfError::Storage("heterogeneous bool not Boolean".into())
})?
.value(r),
),
_ => {
return Err(GfError::Storage(format!(
"unsupported heterogeneous property tag {tag}"
)));
}
}
}
["months", "days", "seconds", "nanos"] => {
let i64_at = |idx: usize| -> Result<i64, GfError> {
Ok(s.column(idx)
.as_any()
.downcast_ref::<Int64Array>()
.ok_or_else(|| {
GfError::Storage("duration struct child not Int64".into())
})?
.value(r))
};
IrLiteral::Duration {
months: i64_at(0)?,
days: i64_at(1)?,
seconds: i64_at(2)?,
nanos: i64_at(3)?,
}
}
["epoch_day"] => IrLiteral::Date(
s.column(0)
.as_any()
.downcast_ref::<Int64Array>()
.ok_or_else(|| GfError::Storage("date epoch_day not Int64".into()))?
.value(r),
),
["date", "time"] => {
let days = s
.column(0)
.as_any()
.downcast_ref::<Int64Array>()
.ok_or_else(|| GfError::Storage("localdatetime date not Int64".into()))?
.value(r);
let nanos = s
.column(1)
.as_any()
.downcast_ref::<Time64NanosecondArray>()
.ok_or_else(|| {
GfError::Storage("localdatetime time not Time64(ns)".into())
})?
.value(r);
IrLiteral::LocalDateTime { days, nanos }
}
["time", "offset"] => {
let nanos = s
.column(0)
.as_any()
.downcast_ref::<Time64NanosecondArray>()
.ok_or_else(|| GfError::Storage("time not Time64(ns)".into()))?
.value(r);
let offset = s
.column(1)
.as_any()
.downcast_ref::<Int32Array>()
.ok_or_else(|| GfError::Storage("time offset not Int32".into()))?
.value(r);
IrLiteral::ZonedTime { nanos, offset }
}
["date", "time", "offset", "zone"] => {
let days = s
.column(0)
.as_any()
.downcast_ref::<Int64Array>()
.ok_or_else(|| GfError::Storage("datetime date not Int64".into()))?
.value(r);
let nanos = s
.column(1)
.as_any()
.downcast_ref::<Time64NanosecondArray>()
.ok_or_else(|| GfError::Storage("datetime time not Time64(ns)".into()))?
.value(r);
let offset = s
.column(2)
.as_any()
.downcast_ref::<Int32Array>()
.ok_or_else(|| GfError::Storage("datetime offset not Int32".into()))?
.value(r);
let zone_col = s
.column(3)
.as_any()
.downcast_ref::<StringArray>()
.ok_or_else(|| GfError::Storage("datetime zone not Utf8".into()))?;
let zone = (!zone_col.is_null(r)).then(|| zone_col.value(r).to_owned());
IrLiteral::ZonedDateTime {
days,
nanos,
offset,
zone,
}
}
_ => {
return Err(GfError::Storage(format!(
"property column {} has unsupported struct shape {names:?}",
field.name()
)));
}
}
}
DataType::Time64(TimeUnit::Nanosecond) => {
IrLiteral::Time(downcast::<Time64NanosecondArray>(col, field)?.value(r))
}
DataType::Timestamp(TimeUnit::Microsecond, _) => {
IrLiteral::DateTime(downcast::<TimestampMicrosecondArray>(col, field)?.value(r))
}
DataType::List(inner) => {
let larr = downcast::<ListArray>(col, field)?;
let elems = larr.value(r);
let mut items = Vec::with_capacity(elems.len());
for j in 0..elems.len() {
if elems.is_null(j) {
items.push(IrLiteral::Null);
} else {
items.push(decode_value(&elems, inner, j)?);
}
}
IrLiteral::List(items)
}
other => {
return Err(GfError::Storage(format!(
"property column {} has unsupported type {other:?}",
field.name()
)));
}
})
}
fn downcast<'a, A: 'static>(
col: &'a arrow::array::ArrayRef,
field: &arrow::datatypes::Field,
) -> Result<&'a A, GfError> {
col.as_any().downcast_ref::<A>().ok_or_else(|| {
GfError::Storage(format!(
"property column {} could not be read as its declared type",
field.name()
))
})
}
fn apply_property_updates<R: PropRowLike>(
mut rows: Vec<R>,
updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
) -> (Vec<R>, u64) {
let mut index: HashMap<[u8; 16], usize> = HashMap::with_capacity(rows.len());
for (i, row) in rows.iter().enumerate() {
index.entry(*row.uuid_bytes()).or_insert(i);
}
let mut touched = 0u64;
for (uuid, new_props) in updates {
if new_props.is_empty() {
continue;
}
touched += 1;
if let Some(&i) = index.get(uuid) {
rows[i].props_mut().extend(new_props.clone());
} else {
index.insert(*uuid, rows.len());
rows.push(R::from_parts(*uuid, new_props.clone()));
}
}
(rows, touched)
}
fn apply_property_removals<R: PropRowLike>(
mut rows: Vec<R>,
removals: &HashMap<[u8; 16], HashSet<String>>,
) -> (Vec<R>, u64) {
let mut index: HashMap<[u8; 16], usize> = HashMap::with_capacity(rows.len());
for (i, row) in rows.iter().enumerate() {
index.entry(*row.uuid_bytes()).or_insert(i);
}
let mut touched = 0u64;
for (uuid, keys) in removals {
if keys.is_empty() {
continue;
}
touched += 1;
if let Some(&i) = index.get(uuid) {
let props = rows[i].props_mut();
for k in keys {
props.remove(k);
}
}
}
(rows, touched)
}
#[allow(clippy::implicit_hasher)]
pub fn stage_set_node_properties(
staged: &mut RewriteBatch,
dir: &Path,
stem: &str,
updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
) -> Result<u64, GfError> {
let existing = read_props_through(staged, &node_props_path(dir, stem))?;
let rows = decode_property_rows(&existing)?;
let (rows, touched) = apply_property_updates(rows, updates);
stage_node_property_file(staged, dir, stem, &rows)?;
Ok(touched)
}
#[allow(clippy::implicit_hasher)] pub fn stage_remove_node_properties(
staged: &mut RewriteBatch,
dir: &Path,
stem: &str,
removals: &HashMap<[u8; 16], HashSet<String>>,
) -> Result<u64, GfError> {
let existing = read_props_through(staged, &node_props_path(dir, stem))?;
let rows = decode_property_rows(&existing)?;
let (rows, touched) = apply_property_removals(rows, removals);
stage_node_property_file(staged, dir, stem, &rows)?;
Ok(touched)
}
#[allow(clippy::implicit_hasher)] pub fn stage_set_edge_properties(
staged: &mut RewriteBatch,
dir: &Path,
rel_stem: &str,
updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
) -> Result<u64, GfError> {
let existing = read_props_through(staged, &edge_props_path(dir, rel_stem))?;
let rows = decode_edge_property_rows(&existing)?;
let (rows, touched) = apply_property_updates(rows, updates);
stage_edge_property_file(staged, dir, rel_stem, &rows)?;
Ok(touched)
}
#[allow(clippy::implicit_hasher)] pub fn stage_remove_edge_properties(
staged: &mut RewriteBatch,
dir: &Path,
rel_stem: &str,
removals: &HashMap<[u8; 16], HashSet<String>>,
) -> Result<u64, GfError> {
let existing = read_props_through(staged, &edge_props_path(dir, rel_stem))?;
let rows = decode_edge_property_rows(&existing)?;
let (rows, touched) = apply_property_removals(rows, removals);
stage_edge_property_file(staged, dir, rel_stem, &rows)?;
Ok(touched)
}
#[allow(clippy::implicit_hasher)] pub fn set_node_properties(
dir: &Path,
stem: &str,
updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
) -> Result<u64, GfError> {
let mut staged = RewriteBatch::new();
let touched = stage_set_node_properties(&mut staged, dir, stem, updates)?;
crate::generation::commit_topology_aware(staged, dir)?;
Ok(touched)
}
#[allow(clippy::implicit_hasher)] pub fn remove_node_properties(
dir: &Path,
stem: &str,
removals: &HashMap<[u8; 16], HashSet<String>>,
) -> Result<u64, GfError> {
let mut staged = RewriteBatch::new();
let touched = stage_remove_node_properties(&mut staged, dir, stem, removals)?;
crate::generation::commit_topology_aware(staged, dir)?;
Ok(touched)
}
#[allow(clippy::implicit_hasher)] pub fn set_edge_properties_rewrite(
dir: &Path,
rel_stem: &str,
updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
) -> Result<u64, GfError> {
let mut staged = RewriteBatch::new();
let touched = stage_set_edge_properties(&mut staged, dir, rel_stem, updates)?;
staged.commit()?;
Ok(touched)
}
#[allow(clippy::implicit_hasher)] pub fn remove_edge_properties(
dir: &Path,
rel_stem: &str,
removals: &HashMap<[u8; 16], HashSet<String>>,
) -> Result<u64, GfError> {
let mut staged = RewriteBatch::new();
let touched = stage_remove_edge_properties(&mut staged, dir, rel_stem, removals)?;
staged.commit()?;
Ok(touched)
}
fn stage_node_property_file(
staged: &mut RewriteBatch,
dir: &Path,
stem: &str,
rows: &[PropRow],
) -> Result<(), GfError> {
if rows.is_empty() {
return Ok(());
}
let (schema, cols) = build_property_columns(stem, rows)?;
stage_property_file(staged, dir, "properties", stem, schema, cols)
}
fn stage_edge_property_file(
staged: &mut RewriteBatch,
dir: &Path,
stem: &str,
rows: &[EdgePropRow],
) -> Result<(), GfError> {
if rows.is_empty() {
return Ok(());
}
let (schema, cols) =
build_property_columns_keyed(EDGE_PROPERTY_UUID_FIELD, "graphforge.rel_type", stem, rows)?;
stage_property_file(staged, dir, "edge_properties", stem, schema, cols)
}
fn stage_property_file(
staged: &mut RewriteBatch,
dir: &Path,
subdir: &str,
stem: &str,
schema: Schema,
cols: Vec<ArrayRef>,
) -> Result<(), GfError> {
let batch = RecordBatch::try_new(Arc::new(schema), cols).map_err(pq_err)?;
staged.restage(
&dir.join(subdir).join(format!("{stem}.parquet")),
batch.schema(),
&batch,
)
}
fn node_props_path(dir: &Path, stem: &str) -> PathBuf {
dir.join("properties").join(format!("{stem}.parquet"))
}
fn edge_props_path(dir: &Path, stem: &str) -> PathBuf {
dir.join("edge_properties").join(format!("{stem}.parquet"))
}
fn read_props_through(staged: &RewriteBatch, path: &Path) -> Result<Vec<RecordBatch>, GfError> {
let read_path = staged
.staged_temp(path)
.map_or_else(|| path.to_path_buf(), Path::to_path_buf);
match crate::catalog::discover_parquet_schema(&read_path) {
Some(schema) => crate::catalog::read_parquet_or_empty(&read_path, schema).map_err(pq_err),
None => Ok(Vec::new()),
}
}
use crate::staging::RewriteBatch;
fn concat_with_existing(
schema: &SchemaRef,
existing: Vec<RecordBatch>,
new: RecordBatch,
) -> Result<RecordBatch, GfError> {
let mut all = existing;
all.push(new);
arrow::compute::concat_batches(schema, &all).map_err(pq_err)
}
#[cfg(test)]
mod tests {
use std::fs::File;
use super::*;
use graphforge_core::uuid::new_v7;
use tempfile::TempDir;
const TS: i64 = 1_700_000_000_000_000;
#[test]
fn create_node_persists_complete_label_set_and_primary_label() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
w.create_node_with_labels(new_v7(), &[TypeId(4), TypeId(9)])
.unwrap();
w.flush().unwrap();
let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
let batch = &nodes[0];
let primary = batch
.column_by_name("type_id")
.unwrap()
.as_any()
.downcast_ref::<UInt32Array>()
.unwrap();
assert_eq!(primary.value(0), 4);
let sets = batch
.column_by_name("type_ids")
.unwrap()
.as_any()
.downcast_ref::<arrow::array::ListArray>()
.unwrap();
let labels = sets.value(0);
let labels = labels.as_any().downcast_ref::<UInt32Array>().unwrap();
assert_eq!(labels.values(), &[4, 9]);
}
#[test]
fn surrogate_ids_are_monotonic_from_one() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 1);
assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 2);
assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 3);
}
#[test]
fn create_edge_with_unknown_endpoint_errors() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
let a = new_v7();
w.create_node(a, TypeId(0)).unwrap();
let b = new_v7();
let e = w.create_edge(new_v7(), "KNOWS", &a, &b);
assert!(matches!(e, Err(GfError::Storage(_))), "got {e:?}");
let unknown_source = new_v7();
let source_error = w.create_edge(new_v7(), "KNOWS", &unknown_source, &a);
assert!(matches!(&source_error, Err(GfError::Storage(_))));
assert!(source_error.unwrap_err().to_string().contains("source"));
}
#[test]
fn register_existing_node_resolves_edge_endpoint() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
let a = new_v7();
w.register_existing_node(a, 42);
let b = new_v7();
assert_eq!(w.create_node(b, TypeId(0)).unwrap(), 1);
let edge_id = w
.create_edge(new_v7(), "KNOWS", &a, &b)
.expect("edge with a registered endpoint resolves");
assert_eq!(edge_id, 1);
}
#[test]
fn register_existing_node_does_not_write_a_node_row() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
w.register_existing_node(new_v7(), 7);
let created = new_v7();
w.create_node(created, TypeId(0)).unwrap();
w.flush().unwrap();
let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
let rows: usize = nodes.iter().map(arrow::array::RecordBatch::num_rows).sum();
assert_eq!(
rows, 1,
"register_existing_node must not persist a node row"
);
}
#[test]
fn empty_flush_creates_no_directories() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
w.flush().unwrap();
assert!(!dir.path().join("topology").exists());
assert!(!dir.path().join("properties").exists());
}
#[test]
fn null_first_property_column_infers_later_concrete_type() {
use arrow::array::Array;
use arrow::datatypes::DataType;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let a = new_v7();
let b = new_v7();
let c = new_v7();
w.create_node(a, TypeId(0)).unwrap();
w.create_node(b, TypeId(0)).unwrap();
w.create_node(c, TypeId(0)).unwrap();
w.set_properties(
&a,
None,
HashMap::from([("score".to_owned(), IrLiteral::Null)]),
)
.unwrap();
w.set_properties(
&b,
None,
HashMap::from([("score".to_owned(), IrLiteral::Int(10))]),
)
.unwrap();
w.set_properties(
&c,
None,
HashMap::from([("score".to_owned(), IrLiteral::Int(20))]),
)
.unwrap();
w.flush().unwrap();
let path = dir.path().join("properties").join("_untyped.parquet");
let file = File::open(&path).unwrap();
let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
let schema = builder.schema().clone();
let score_fields: Vec<_> = schema
.fields()
.iter()
.filter(|f| f.name() == "score")
.collect();
assert_eq!(score_fields.len(), 1, "expected a single score column");
assert_eq!(
score_fields[0].data_type(),
&DataType::Int64,
"null-first then Int should infer Int64, not Utf8"
);
let mut reader = builder.build().unwrap();
let batch = reader.next().unwrap().unwrap();
assert_eq!(batch.num_rows(), 3);
let scores = batch
.column(schema.index_of("score").unwrap())
.as_any()
.downcast_ref::<arrow::array::Int64Array>()
.unwrap();
assert!(scores.is_null(0));
assert_eq!(scores.value(1), 10);
assert_eq!(scores.value(2), 20);
}
#[test]
fn mixed_type_property_column_uses_tagged_scalars() {
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let a = new_v7();
let b = new_v7();
w.create_node(a, TypeId(0)).unwrap();
w.create_node(b, TypeId(0)).unwrap();
w.set_properties(
&a,
None,
HashMap::from([("x".to_owned(), IrLiteral::Int(1))]),
)
.unwrap();
w.set_properties(
&b,
None,
HashMap::from([("x".to_owned(), IrLiteral::Str("two".to_owned()))]),
)
.unwrap();
w.flush().unwrap();
let path = dir.path().join("properties").join("_untyped.parquet");
let file = File::open(&path).unwrap();
let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
let schema = builder.schema().clone();
let x = schema.field_with_name("x").unwrap();
assert_eq!(
x.data_type(),
&DataType::Struct(heterogeneous_scalar_fields())
);
}
#[test]
fn every_heterogeneous_scalar_tag_round_trips_exactly() {
let dir = TempDir::new().unwrap();
let cases = [
IrLiteral::Int(-1),
IrLiteral::Float(2.25),
IrLiteral::Str("three".into()),
IrLiteral::Bool(true),
];
let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let mut expected = HashMap::new();
for value in cases {
let node = new_v7();
writer.create_node(node, TypeId(0)).unwrap();
writer
.set_properties(
&node,
None,
HashMap::from([("mixed".into(), value.clone())]),
)
.unwrap();
expected.insert(to_bytes(&node), value);
}
writer.flush().unwrap();
let reopened = read_node_props(dir.path(), "_untyped");
assert_eq!(reopened.len(), expected.len());
for (node, value) in expected {
assert_eq!(reopened[&node].get("mixed"), Some(&value));
}
}
fn read_node_props(dir: &Path, stem: &str) -> HashMap<[u8; 16], HashMap<String, IrLiteral>> {
read_node_property_rows(dir, stem).unwrap()
}
fn read_edge_props(dir: &Path, stem: &str) -> HashMap<[u8; 16], HashMap<String, IrLiteral>> {
let batches = crate::catalog::read_edge_properties(dir, stem).unwrap();
let mut out = HashMap::new();
for row in decode_edge_property_rows(&batches).unwrap() {
out.insert(row.edge_uuid, row.props);
}
out
}
#[test]
fn set_node_properties_sets_new_and_overwrites_existing() {
let dir = TempDir::new().unwrap();
assert!(
read_node_property_rows(dir.path(), "_untyped")
.unwrap()
.is_empty()
);
fs::create_dir_all(dir.path().join("properties")).unwrap();
fs::write(dir.path().join("properties/_untyped.parquet"), b"invalid").unwrap();
assert!(read_node_property_rows(dir.path(), "_untyped").is_err());
fs::remove_file(dir.path().join("properties/_untyped.parquet")).unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let a = new_v7();
w.create_node(a, TypeId(0)).unwrap();
w.set_properties(
&a,
None,
HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
)
.unwrap();
w.flush().unwrap();
let ab = to_bytes(&a);
let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
let updates = HashMap::from([(
ab,
HashMap::from([
("age".to_owned(), IrLiteral::Int(31)),
("name".to_owned(), IrLiteral::Str("Al".to_owned())),
]),
)]);
let touched = set_node_properties(dir.path(), "_untyped", &updates).unwrap();
assert_eq!(touched, 1);
assert_eq!(
crate::generation::read_search_generation(dir.path()).unwrap(),
search_generation + 1
);
let props = read_node_props(dir.path(), "_untyped");
assert_eq!(props[&ab]["age"], IrLiteral::Int(31));
assert_eq!(props[&ab]["name"], IrLiteral::Str("Al".to_owned()));
}
#[test]
fn set_node_properties_inserts_row_for_propertyless_node() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let a = new_v7();
w.create_node(a, TypeId(0)).unwrap();
w.flush().unwrap();
let ab = to_bytes(&a);
let updates =
HashMap::from([(ab, HashMap::from([("age".to_owned(), IrLiteral::Int(42))]))]);
let touched = set_node_properties(dir.path(), "_untyped", &updates).unwrap();
assert_eq!(touched, 1);
let props = read_node_props(dir.path(), "_untyped");
assert_eq!(props[&ab]["age"], IrLiteral::Int(42));
}
#[test]
fn set_node_properties_routes_by_stem_in_strict_mode() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
let a = new_v7();
w.create_node(a, TypeId(1)).unwrap();
w.set_properties(
&a,
Some("Person"),
HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
)
.unwrap();
w.flush().unwrap();
let ab = to_bytes(&a);
let updates =
HashMap::from([(ab, HashMap::from([("age".to_owned(), IrLiteral::Int(99))]))]);
set_node_properties(dir.path(), "Person", &updates).unwrap();
assert!(
dir.path()
.join("properties")
.join("Person.parquet")
.exists()
);
let props = read_node_props(dir.path(), "Person");
assert_eq!(props[&ab]["age"], IrLiteral::Int(99));
}
#[test]
fn remove_node_properties_drops_key_and_column_when_last() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let a = new_v7();
w.create_node(a, TypeId(0)).unwrap();
w.set_properties(
&a,
None,
HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
)
.unwrap();
w.flush().unwrap();
let ab = to_bytes(&a);
let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
let removals = HashMap::from([(ab, HashSet::from(["age".to_owned()]))]);
let touched = remove_node_properties(dir.path(), "_untyped", &removals).unwrap();
assert_eq!(touched, 1);
assert_eq!(
crate::generation::read_search_generation(dir.path()).unwrap(),
search_generation + 1
);
let props = read_node_props(dir.path(), "_untyped");
assert!(props.get(&ab).map_or(true, HashMap::is_empty));
}
#[test]
fn remove_node_properties_missing_key_and_uuid_are_noops() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let a = new_v7();
w.create_node(a, TypeId(0)).unwrap();
w.set_properties(
&a,
None,
HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
)
.unwrap();
w.flush().unwrap();
let ab = to_bytes(&a);
let removals = HashMap::from([
(ab, HashSet::from(["nope".to_owned()])),
(to_bytes(&new_v7()), HashSet::from(["age".to_owned()])),
]);
remove_node_properties(dir.path(), "_untyped", &removals).unwrap();
let props = read_node_props(dir.path(), "_untyped");
assert_eq!(props[&ab]["age"], IrLiteral::Int(30));
}
#[test]
fn set_and_remove_edge_properties_round_trip() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let a = new_v7();
let b = new_v7();
w.create_node(a, TypeId(0)).unwrap();
w.create_node(b, TypeId(0)).unwrap();
let e = new_v7();
w.create_edge(e, "KNOWS", &a, &b).unwrap();
w.set_edge_properties(
&e,
Some("KNOWS"),
HashMap::from([("since".to_owned(), IrLiteral::Int(2019))]),
)
.unwrap();
w.flush().unwrap();
let eb = to_bytes(&e);
let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
let updates = HashMap::from([(
eb,
HashMap::from([("since".to_owned(), IrLiteral::Int(2020))]),
)]);
assert_eq!(
set_edge_properties_rewrite(dir.path(), "KNOWS", &updates).unwrap(),
1
);
assert_eq!(
read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
IrLiteral::Int(2020)
);
let removals = HashMap::from([(eb, HashSet::from(["since".to_owned()]))]);
assert_eq!(
remove_edge_properties(dir.path(), "KNOWS", &removals).unwrap(),
1
);
let props = read_edge_props(dir.path(), "KNOWS");
assert!(props.get(&eb).map_or(true, HashMap::is_empty));
assert_eq!(
crate::generation::read_search_generation(dir.path()).unwrap(),
search_generation
);
}
#[test]
fn set_node_properties_empty_map_writes_nothing() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
w.create_node(new_v7(), TypeId(0)).unwrap();
w.flush().unwrap();
let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
let touched = set_node_properties(dir.path(), "_untyped", &HashMap::new()).unwrap();
assert_eq!(touched, 0);
assert_eq!(
crate::generation::read_search_generation(dir.path()).unwrap(),
search_generation
);
assert!(
!dir.path()
.join("properties")
.join("_untyped.parquet")
.exists()
);
}
#[test]
fn staged_set_is_invisible_until_commit_across_stems() {
let dir = TempDir::new().unwrap();
let (a, e) = (new_v7(), new_v7());
let b = new_v7();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
w.create_node(a, TypeId(0)).unwrap();
w.create_node(b, TypeId(0)).unwrap();
w.create_edge(e, "KNOWS", &a, &b).unwrap();
w.set_properties(
&a,
None,
HashMap::from([("name".to_owned(), IrLiteral::Str("old".into()))]),
)
.unwrap();
w.set_edge_properties(
&e,
Some("KNOWS"),
HashMap::from([("since".to_owned(), IrLiteral::Int(2000))]),
)
.unwrap();
w.flush().unwrap();
let (ab, eb) = (to_bytes(&a), to_bytes(&e));
let node_updates = HashMap::from([(
ab,
HashMap::from([("name".to_owned(), IrLiteral::Str("new".into()))]),
)]);
let edge_updates = HashMap::from([(
eb,
HashMap::from([("since".to_owned(), IrLiteral::Int(2024))]),
)]);
let mut staged = RewriteBatch::new();
let touched = stage_set_node_properties(&mut staged, dir.path(), "_untyped", &node_updates)
.unwrap()
+ stage_set_edge_properties(&mut staged, dir.path(), "KNOWS", &edge_updates).unwrap();
assert_eq!(touched, 2, "one node + one edge written");
assert_eq!(
read_node_props(dir.path(), "_untyped")[&ab]["name"],
IrLiteral::Str("old".into())
);
assert_eq!(
read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
IrLiteral::Int(2000)
);
staged.commit().unwrap();
assert_eq!(
read_node_props(dir.path(), "_untyped")[&ab]["name"],
IrLiteral::Str("new".into())
);
assert_eq!(
read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
IrLiteral::Int(2024)
);
}
#[test]
fn pending_nodes_batch_is_canonical_and_non_consuming() {
let dir = TempDir::new().unwrap();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
let empty = w.pending_nodes_batch().unwrap();
assert_eq!(empty.schema(), TOPOLOGY_NODES_SCHEMA.clone());
assert_eq!(empty.num_rows(), 0);
let node = new_v7();
w.create_node_with_labels(node, &[TypeId(3), TypeId(7)])
.unwrap();
for _ in 0..2 {
let batch = w.pending_nodes_batch().unwrap();
assert_eq!(batch.schema(), TOPOLOGY_NODES_SCHEMA.clone());
assert_eq!(batch.num_rows(), 1);
assert_eq!(
batch
.column_by_name("node_id")
.unwrap()
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap()
.value(0),
1
);
}
assert!(w.contains_pending_node(&to_bytes(&node)));
}
#[test]
fn cancel_nodes_drops_rows_props_and_mapping() {
let dir = TempDir::new().unwrap();
let (a, b) = (new_v7(), new_v7());
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
w.create_node(a, TypeId(0)).unwrap();
w.create_node(b, TypeId(0)).unwrap();
w.set_properties(
&a,
None,
HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
)
.unwrap();
assert!(w.contains_pending_node(&to_bytes(&a)));
let dropped = w.cancel_nodes(&HashSet::from([to_bytes(&a)]));
assert_eq!(dropped, 1);
assert!(!w.contains_pending_node(&to_bytes(&a)));
let err = w.create_edge(new_v7(), "KNOWS", &b, &a);
assert!(err.is_err(), "edge to a cancelled node must fail");
w.flush().unwrap();
let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
let total: usize = nodes.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total, 1, "only b persisted");
assert!(
!read_node_props(dir.path(), "_untyped").contains_key(&to_bytes(&a)),
"cancelled node's props never hit disk"
);
}
#[test]
fn cancel_edges_drops_rows_and_edge_props() {
let dir = TempDir::new().unwrap();
let (a, b) = (new_v7(), new_v7());
let (e1, e2) = (new_v7(), new_v7());
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
w.create_node(a, TypeId(0)).unwrap();
w.create_node(b, TypeId(0)).unwrap();
w.create_edge(e1, "KNOWS", &a, &b).unwrap();
w.create_edge(e2, "KNOWS", &b, &a).unwrap();
w.set_edge_properties(
&e1,
Some("KNOWS"),
HashMap::from([("since".to_owned(), IrLiteral::Int(2020))]),
)
.unwrap();
assert!(w.contains_pending_edge(&to_bytes(&e1)));
assert_eq!(w.cancel_edges(&HashSet::from([to_bytes(&e1)])), 1);
assert!(!w.contains_pending_edge(&to_bytes(&e1)));
w.flush().unwrap();
assert!(
!read_edge_props(dir.path(), "KNOWS").contains_key(&to_bytes(&e1)),
"cancelled edge's props never hit disk"
);
let edges = crate::catalog::read_parquet_or_empty(
&dir.path().join("topology/edges/_exploratory.parquet"),
EXPLORATORY_EDGE_SCHEMA.clone(),
)
.unwrap();
let total: usize = edges.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total, 1, "only e2 persisted");
}
#[test]
fn pending_incident_edge_uuids_sees_buffered_edges() {
let dir = TempDir::new().unwrap();
let (a, b, c) = (new_v7(), new_v7(), new_v7());
let e_ab = new_v7();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
w.create_node(a, TypeId(0)).unwrap();
w.create_node(b, TypeId(0)).unwrap();
w.create_node(c, TypeId(0)).unwrap();
w.create_edge(e_ab, "KNOWS", &a, &b).unwrap();
let hits = w.pending_incident_edge_uuids(&HashSet::from([to_bytes(&b)]));
assert_eq!(hits, vec![to_bytes(&e_ab)]);
assert!(
w.pending_incident_edge_uuids(&HashSet::from([to_bytes(&c)]))
.is_empty()
);
}
#[test]
fn pending_query_and_label_edits_are_exact_before_flush_and_reopen() {
let dir = TempDir::new().unwrap();
let (alice, bob, edge) = (new_v7(), new_v7(), new_v7());
let (alice_bytes, bob_bytes, edge_bytes) =
(to_bytes(&alice), to_bytes(&bob), to_bytes(&edge));
let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
writer
.create_node_with_labels(alice, &[TypeId(3), TypeId(7)])
.unwrap();
writer.create_node(bob, TypeId(3)).unwrap();
writer
.set_properties(
&alice,
Some("Person"),
HashMap::from([
("name".into(), IrLiteral::Str("Alice".into())),
("age".into(), IrLiteral::Int(42)),
]),
)
.unwrap();
assert_eq!(
writer.pending_node_labels(&HashSet::from([alice_bytes, bob_bytes])),
HashSet::from([3, 7])
);
let matched = writer
.find_pending_node(&[3, 7], &[("name".into(), IrLiteral::Str("Alice".into()))])
.unwrap();
assert_eq!(matched.0, alice_bytes);
assert_eq!(matched.2, 3);
assert_eq!(matched.3, vec![3, 7]);
assert_eq!(matched.4["age"], IrLiteral::Int(42));
assert!(writer.find_pending_node(&[9], &[]).is_none());
assert!(
writer
.find_pending_node(&[3], &[("name".into(), IrLiteral::Str("Bob".into()))])
.is_none()
);
assert_eq!(writer.add_pending_node_labels(&alice_bytes, &[7, 9]), 1);
assert_eq!(writer.add_pending_node_labels(&[0xff; 16], &[1]), 0);
assert_eq!(writer.remove_pending_node_labels(&alice_bytes, &[7, 99]), 1);
assert_eq!(writer.remove_pending_node_labels(&[0xff; 16], &[1]), 0);
assert_eq!(
writer.pending_node_labels(&HashSet::from([alice_bytes])),
HashSet::from([3, 9])
);
writer.create_edge(edge, "KNOWS", &alice, &bob).unwrap();
writer
.set_edge_properties(
&edge,
Some("KNOWS"),
HashMap::from([("since".into(), IrLiteral::Int(2024))]),
)
.unwrap();
let direct = writer
.find_pending_edge(
"KNOWS",
&alice_bytes,
&bob_bytes,
false,
&[("since".into(), IrLiteral::Int(2024))],
)
.unwrap();
assert_eq!(direct.0, edge_bytes);
assert_eq!(direct.1, alice_bytes);
assert_eq!(direct.2, bob_bytes);
assert_eq!(direct.3["since"], IrLiteral::Int(2024));
assert!(
writer
.find_pending_edge("KNOWS", &bob_bytes, &alice_bytes, false, &[])
.is_none()
);
assert!(
writer
.find_pending_edge("KNOWS", &bob_bytes, &alice_bytes, true, &[])
.is_some()
);
assert!(
writer
.find_pending_edge("IGNORES", &alice_bytes, &bob_bytes, false, &[])
.is_none()
);
writer.flush().unwrap();
let mut reopened = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS + 1).unwrap();
assert_eq!(reopened.create_node(new_v7(), TypeId(3)).unwrap(), 3);
assert_eq!(
read_node_props(dir.path(), "Person")[&alice_bytes]["name"],
IrLiteral::Str("Alice".into())
);
assert_eq!(
read_edge_props(dir.path(), "KNOWS")[&edge_bytes]["since"],
IrLiteral::Int(2024)
);
}
#[test]
fn merge_and_remove_pending_props_edit_buffered_rows() {
let dir = TempDir::new().unwrap();
let a = new_v7();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
w.create_node(a, TypeId(0)).unwrap();
w.set_properties(
&a,
None,
HashMap::from([
("name".to_owned(), IrLiteral::Str("old".into())),
("age".to_owned(), IrLiteral::Int(30)),
]),
)
.unwrap();
w.merge_pending_node_props(
&to_bytes(&a),
None,
HashMap::from([
("name".to_owned(), IrLiteral::Str("new".into())),
("city".to_owned(), IrLiteral::Str("Oslo".into())),
]),
);
w.remove_pending_node_props(
&to_bytes(&a),
&HashSet::from(["age".to_owned(), "absent".to_owned()]),
);
w.flush().unwrap();
let props = &read_node_props(dir.path(), "_untyped")[&to_bytes(&a)];
assert_eq!(props["name"], IrLiteral::Str("new".into()));
assert_eq!(props["city"], IrLiteral::Str("Oslo".into()));
assert!(!props.contains_key("age"), "removed before flush");
}
#[test]
fn flush_into_composes_with_staged_delete_in_one_batch() {
let dir = TempDir::new().unwrap();
let (a, b) = (new_v7(), new_v7());
let mut seed = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
seed.create_node(a, TypeId(0)).unwrap();
seed.create_node(b, TypeId(0)).unwrap();
seed.set_properties(
&a,
None,
HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
)
.unwrap();
seed.flush().unwrap();
let mut staged = RewriteBatch::new();
let removed = crate::mutator::stage_delete_nodes(
&mut staged,
dir.path(),
&HashSet::from([to_bytes(&a)]),
)
.unwrap();
assert_eq!(removed, 1);
let d = new_v7();
let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
w.create_node(d, TypeId(0)).unwrap();
w.flush_into(&mut staged).unwrap();
let staged_nodes = staged
.staged_paths()
.filter(|p| p.ends_with("topology/nodes.parquet"))
.count();
assert_eq!(staged_nodes, 1, "net content, no double-stage");
let pre: usize = crate::catalog::read_nodes(dir.path())
.unwrap()
.iter()
.map(RecordBatch::num_rows)
.sum();
assert_eq!(pre, 2);
staged.commit().unwrap();
let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
let total: usize = nodes.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total, 2, "b survives, a deleted, d created");
assert!(
!read_node_props(dir.path(), "_untyped").contains_key(&to_bytes(&a)),
"deleted node's props gone"
);
}
#[test]
fn every_persisted_property_family_round_trips_through_parquet_reopen() {
let dir = TempDir::new().unwrap();
let node = new_v7();
let propertyless = new_v7();
let values = HashMap::from([
("int".into(), IrLiteral::Int(-7)),
("float".into(), IrLiteral::Float(2.5)),
("bool".into(), IrLiteral::Bool(true)),
("str".into(), IrLiteral::Str("value".into())),
(
"duration".into(),
IrLiteral::Duration {
months: 1,
days: -2,
seconds: 3,
nanos: 4,
},
),
("datetime".into(), IrLiteral::DateTime(TS)),
("date".into(), IrLiteral::Date(19_000)),
(
"local_datetime".into(),
IrLiteral::LocalDateTime {
days: 19_001,
nanos: 123,
},
),
("time".into(), IrLiteral::Time(456)),
(
"zoned_time".into(),
IrLiteral::ZonedTime {
nanos: 789,
offset: -21_600,
},
),
(
"zoned_datetime".into(),
IrLiteral::ZonedDateTime {
days: 19_002,
nanos: 987,
offset: 3_600,
zone: Some("Europe/Paris".into()),
},
),
(
"offset_datetime".into(),
IrLiteral::ZonedDateTime {
days: 19_003,
nanos: 654,
offset: 0,
zone: None,
},
),
(
"ints".into(),
IrLiteral::List(vec![IrLiteral::Int(1), IrLiteral::Null, IrLiteral::Int(3)]),
),
(
"dates".into(),
IrLiteral::List(vec![IrLiteral::Date(19_004), IrLiteral::Date(19_005)]),
),
("empty".into(), IrLiteral::List(Vec::new())),
("null".into(), IrLiteral::Null),
]);
let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
writer.create_node(node, TypeId(0)).unwrap();
writer.create_node(propertyless, TypeId(0)).unwrap();
writer.set_properties(&node, None, values.clone()).unwrap();
writer
.set_properties(
&propertyless,
None,
values
.keys()
.map(|name| (name.clone(), IrLiteral::Null))
.collect(),
)
.unwrap();
writer.flush().unwrap();
let reopened = read_node_props(dir.path(), "_untyped");
let actual = reopened.get(&to_bytes(&node)).unwrap();
for (name, expected) in &values {
if matches!(expected, IrLiteral::Null) {
assert!(!actual.contains_key(name));
} else if name == "empty" {
assert_eq!(actual.get(name), Some(&IrLiteral::Str("[]".into())));
} else {
assert_eq!(actual.get(name), Some(expected), "property {name}");
}
}
assert!(
reopened
.get(&to_bytes(&propertyless))
.is_none_or(HashMap::is_empty)
);
}
#[test]
fn property_literal_rendering_and_nested_invalid_values_are_deterministic() {
let uuid = [0xabu8; 16];
let cases = [
(IrLiteral::Null, "".into()),
(IrLiteral::Bool(true), "true".into()),
(IrLiteral::Int(-2), "-2".into()),
(IrLiteral::Float(1.25), "1.25".into()),
(IrLiteral::Str("s".into()), "s".into()),
(IrLiteral::Uuid(uuid), "ab".repeat(16)),
(
IrLiteral::Duration {
months: 1,
days: 2,
seconds: 3,
nanos: 4,
},
"1mo2d3s4ns".into(),
),
(IrLiteral::DateTime(5), "5".into()),
(IrLiteral::Date(6), "6".into()),
(
IrLiteral::LocalDateTime { days: 7, nanos: 8 },
"7d8ns".into(),
),
(IrLiteral::Time(9), "9ns".into()),
(
IrLiteral::ZonedTime {
nanos: 10,
offset: -1,
},
"10ns-1s".into(),
),
(
IrLiteral::ZonedDateTime {
days: 11,
nanos: 12,
offset: 13,
zone: Some("UTC".into()),
},
"11d12ns+13sUTC".into(),
),
(
IrLiteral::List(vec![IrLiteral::Int(1), IrLiteral::Str("x".into())]),
"[1,x]".into(),
),
(
IrLiteral::Map(vec![("a".into(), IrLiteral::Bool(false))]),
"{a:false}".into(),
),
];
for (literal, expected) in cases {
assert_eq!(literal_to_string(&literal), expected);
}
for invalid in [
IrLiteral::Uuid(uuid),
IrLiteral::List(vec![IrLiteral::Uuid(uuid)]),
IrLiteral::Map(vec![("nested".into(), IrLiteral::Uuid(uuid))]),
] {
assert_eq!(
reject_map_property_value("p", &invalid).unwrap_err().code(),
"GF_VALIDATION"
);
}
for invalid in [
IrLiteral::Map(vec![]),
IrLiteral::List(vec![IrLiteral::Map(vec![])]),
] {
assert_eq!(
reject_map_property_value("p", &invalid).unwrap_err().code(),
"GF_IO"
);
}
}
#[test]
fn persisted_property_decoder_rejects_unsupported_shape_type_and_dynamic_array() {
use arrow::array::{Int32Array, Int64Array, StructArray, UInt8Array};
use arrow::datatypes::{DataType, Field, Fields};
let unsupported: arrow::array::ArrayRef = Arc::new(UInt8Array::from(vec![1]));
let unsupported_field = Field::new("unsupported", DataType::UInt8, false);
assert!(decode_value(&unsupported, &unsupported_field, 0).is_err());
let fields: Fields = vec![Field::new("other", DataType::Int32, false)].into();
let structure: arrow::array::ArrayRef = Arc::new(StructArray::new(
fields.clone(),
vec![Arc::new(Int32Array::from(vec![1]))],
None,
));
let structure_field = Field::new("structure", DataType::Struct(fields), false);
assert!(decode_value(&structure, &structure_field, 0).is_err());
let wrong_dynamic: arrow::array::ArrayRef = Arc::new(Int64Array::from(vec![1]));
let declared = Field::new("declared", DataType::UInt64, false);
assert!(decode_value(&wrong_dynamic, &declared, 0).is_err());
}
}