1use std::collections::{HashMap, HashSet};
51use std::fmt;
52use std::fs;
53use std::path::{Path, PathBuf};
54use std::sync::Arc;
55use std::time::{SystemTime, UNIX_EPOCH};
56
57use arrow::array::{
58 ArrayRef, BooleanBuilder, FixedSizeBinaryArray, Float64Array, Float64Builder, Int64Builder,
59 RecordBatch, StringArray, StringBuilder, TimestampMicrosecondArray,
60 TimestampMicrosecondBuilder, UInt32Array, UInt64Array,
61};
62use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
63use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
64
65use graphforge_core::uuid::{Uuid, to_bytes};
66use graphforge_core::{GfError, OntologyMode, TypeId};
67use graphforge_ir::IrLiteral;
68
69pub type PendingNodeMatch = ([u8; 16], u64, u32, Vec<u32>, HashMap<String, IrLiteral>);
71
72use crate::schemas::{
73 EXPLORATORY_EDGE_SCHEMA, TOPOLOGY_NODES_SCHEMA, TYPED_EDGE_SCHEMA, uuid_field,
74};
75
76const EXPLORATORY_STEM: &str = "_exploratory";
78const UNTYPED_STEM: &str = "_untyped";
80const NODE_PROPERTY_UUID_FIELD: &str = "node_uuid";
83const EDGE_PROPERTY_UUID_FIELD: &str = "edge_uuid";
85
86fn io_err(e: &std::io::Error) -> GfError {
91 GfError::Storage(e.to_string())
92}
93
94fn pq_err(e: impl fmt::Display) -> GfError {
95 GfError::Storage(e.to_string())
96}
97
98fn max_u64_column(batches: &[RecordBatch], col: &str) -> u64 {
102 use arrow::array::Array;
103 let mut max = 0u64;
104 for batch in batches {
105 if let Some(c) = batch.column_by_name(col)
106 && let Some(ids) = c.as_any().downcast_ref::<UInt64Array>()
107 {
108 for i in 0..ids.len() {
109 if !ids.is_null(i) {
110 max = max.max(ids.value(i));
111 }
112 }
113 }
114 }
115 max
116}
117
118struct NodeRow {
123 node_uuid: [u8; 16],
124 node_id: u64,
125 type_id: u32,
126 type_ids: Vec<u32>,
127}
128
129struct EdgeRow {
130 edge_uuid: [u8; 16],
131 src_uuid: [u8; 16],
132 dst_uuid: [u8; 16],
133 edge_id: u64,
134 src_id: u64,
135 dst_id: u64,
136 rel_type_name: Option<String>,
139}
140
141struct PropRow {
142 node_uuid: [u8; 16],
143 props: HashMap<String, IrLiteral>,
144}
145
146struct EdgePropRow {
147 edge_uuid: [u8; 16],
148 props: HashMap<String, IrLiteral>,
149}
150
151trait PropRowLike {
159 fn uuid_bytes(&self) -> &[u8; 16];
160 fn props(&self) -> &HashMap<String, IrLiteral>;
161 fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral>;
162 fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self;
163}
164
165impl PropRowLike for PropRow {
166 fn uuid_bytes(&self) -> &[u8; 16] {
167 &self.node_uuid
168 }
169 fn props(&self) -> &HashMap<String, IrLiteral> {
170 &self.props
171 }
172 fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral> {
173 &mut self.props
174 }
175 fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self {
176 Self {
177 node_uuid: uuid,
178 props,
179 }
180 }
181}
182
183impl PropRowLike for EdgePropRow {
184 fn uuid_bytes(&self) -> &[u8; 16] {
185 &self.edge_uuid
186 }
187 fn props(&self) -> &HashMap<String, IrLiteral> {
188 &self.props
189 }
190 fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral> {
191 &mut self.props
192 }
193 fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self {
194 Self {
195 edge_uuid: uuid,
196 props,
197 }
198 }
199}
200
201pub struct GraphWriter {
209 dir: PathBuf,
210 mode: OntologyMode,
211 now_micros: i64,
214 next_node_id: u64,
215 next_edge_id: u64,
216 uuid_to_node_id: HashMap<[u8; 16], u64>,
219 nodes: Vec<NodeRow>,
220 edges: HashMap<String, Vec<EdgeRow>>,
222 properties: HashMap<String, Vec<PropRow>>,
224 edge_properties: HashMap<String, Vec<EdgePropRow>>,
227 pending_delta: Vec<crate::adjacency_delta::DeltaEdge>,
230}
231
232impl GraphWriter {
233 pub fn open(dir: &Path, mode: OntologyMode) -> Result<Self, GfError> {
241 let now_micros = SystemTime::now()
242 .duration_since(UNIX_EPOCH)
243 .map_or(0, |d| i64::try_from(d.as_micros()).unwrap_or(i64::MAX));
244 Self::open_at(dir, mode, now_micros)
245 }
246
247 pub fn open_at(dir: &Path, mode: OntologyMode, now_micros: i64) -> Result<Self, GfError> {
253 fs::create_dir_all(dir).map_err(|e| io_err(&e))?;
254 let max_node_id =
258 max_u64_column(&crate::catalog::read_nodes(dir).map_err(pq_err)?, "node_id");
259 let max_edge_id = crate::catalog::max_edge_id(dir).map_err(pq_err)?;
260 Ok(Self {
261 dir: dir.to_path_buf(),
262 mode,
263 now_micros,
264 next_node_id: max_node_id + 1,
265 next_edge_id: max_edge_id + 1,
266 uuid_to_node_id: HashMap::new(),
267 nodes: Vec::new(),
268 edges: HashMap::new(),
269 properties: HashMap::new(),
270 edge_properties: HashMap::new(),
271 pending_delta: Vec::new(),
272 })
273 }
274
275 pub fn create_node(&mut self, node_uuid: Uuid, type_id: TypeId) -> Result<u64, GfError> {
280 self.create_node_with_labels(node_uuid, &[type_id])
281 }
282
283 pub fn create_node_with_labels(
289 &mut self,
290 node_uuid: Uuid,
291 type_ids: &[TypeId],
292 ) -> Result<u64, GfError> {
293 let bytes = to_bytes(&node_uuid);
294 let node_id = self.next_node_id;
295 self.next_node_id += 1;
296 self.uuid_to_node_id.insert(bytes, node_id);
298 self.nodes.push(NodeRow {
299 node_uuid: bytes,
300 node_id,
301 type_id: type_ids.first().map_or(u32::MAX, |id| id.0),
302 type_ids: type_ids.iter().map(|id| id.0).collect(),
303 });
304 Ok(node_id)
305 }
306
307 pub fn register_existing_node(&mut self, node_uuid: Uuid, node_id: u64) {
317 self.uuid_to_node_id.insert(to_bytes(&node_uuid), node_id);
318 }
319
320 #[must_use]
324 pub fn node_id_for_uuid(&self, node_uuid: &Uuid) -> Option<u64> {
325 self.uuid_to_node_id.get(&to_bytes(node_uuid)).copied()
326 }
327
328 pub fn create_edge(
337 &mut self,
338 edge_uuid: Uuid,
339 rel_type: &str,
340 src_uuid: &Uuid,
341 dst_uuid: &Uuid,
342 ) -> Result<u64, GfError> {
343 let src_bytes = to_bytes(src_uuid);
344 let dst_bytes = to_bytes(dst_uuid);
345 let src_id = *self.uuid_to_node_id.get(&src_bytes).ok_or_else(|| {
346 GfError::Storage(format!(
347 "create_edge: source {} has no node_id; call create_node first",
348 graphforge_core::uuid::to_string(src_uuid)
349 ))
350 })?;
351 let dst_id = *self.uuid_to_node_id.get(&dst_bytes).ok_or_else(|| {
352 GfError::Storage(format!(
353 "create_edge: destination {} has no node_id; call create_node first",
354 graphforge_core::uuid::to_string(dst_uuid)
355 ))
356 })?;
357
358 let edge_id = self.next_edge_id;
359 self.next_edge_id += 1;
360
361 let (stem, rel_type_name) = match self.mode {
362 OntologyMode::Exploratory => (EXPLORATORY_STEM.to_owned(), Some(rel_type.to_owned())),
363 OntologyMode::Advisory | OntologyMode::Strict => (rel_type.to_owned(), None),
364 };
365
366 self.edges.entry(stem).or_default().push(EdgeRow {
367 edge_uuid: to_bytes(&edge_uuid),
368 src_uuid: src_bytes,
369 dst_uuid: dst_bytes,
370 edge_id,
371 src_id,
372 dst_id,
373 rel_type_name,
374 });
375 Ok(edge_id)
376 }
377
378 pub fn set_properties(
387 &mut self,
388 node_uuid: &Uuid,
389 entity_type: Option<&str>,
390 props: HashMap<String, IrLiteral>,
391 ) -> Result<(), GfError> {
392 let stem = match (self.mode, entity_type) {
393 (OntologyMode::Advisory | OntologyMode::Strict, Some(t)) => t.to_owned(),
394 _ => UNTYPED_STEM.to_owned(),
395 };
396 self.properties.entry(stem).or_default().push(PropRow {
397 node_uuid: to_bytes(node_uuid),
398 props,
399 });
400 Ok(())
401 }
402
403 pub fn set_edge_properties(
416 &mut self,
417 edge_uuid: &Uuid,
418 rel_type: Option<&str>,
419 props: HashMap<String, IrLiteral>,
420 ) -> Result<(), GfError> {
421 let stem = rel_type.unwrap_or(UNTYPED_STEM).to_owned();
422 self.edge_properties
423 .entry(stem)
424 .or_default()
425 .push(EdgePropRow {
426 edge_uuid: to_bytes(edge_uuid),
427 props,
428 });
429 Ok(())
430 }
431
432 #[must_use]
446 pub fn contains_pending_node(&self, node_uuid: &[u8; 16]) -> bool {
447 self.nodes.iter().any(|r| &r.node_uuid == node_uuid)
448 }
449
450 #[must_use]
452 pub fn pending_node_labels(&self, targets: &HashSet<[u8; 16]>) -> HashSet<u32> {
453 self.nodes
454 .iter()
455 .filter(|row| targets.contains(&row.node_uuid))
456 .flat_map(|row| row.type_ids.iter().copied())
457 .collect()
458 }
459
460 pub fn pending_nodes_batch(&self) -> Result<RecordBatch, GfError> {
467 let n = self.nodes.len();
468 if n == 0 {
469 return Ok(RecordBatch::new_empty(TOPOLOGY_NODES_SCHEMA.clone()));
470 }
471 let uuids =
472 FixedSizeBinaryArray::try_from_iter(self.nodes.iter().map(|r| r.node_uuid.to_vec()))
473 .map_err(pq_err)?;
474 let node_ids = UInt64Array::from(self.nodes.iter().map(|r| r.node_id).collect::<Vec<_>>());
475 let type_ids = UInt32Array::from(self.nodes.iter().map(|r| r.type_id).collect::<Vec<_>>());
476 let nullable_label_sets =
477 arrow::array::ListArray::from_iter_primitive::<arrow::datatypes::UInt32Type, _, _>(
478 self.nodes
479 .iter()
480 .map(|row| Some(row.type_ids.iter().copied().map(Some))),
481 );
482 let label_sets = arrow::array::ListArray::new(
483 Arc::new(Field::new("item", DataType::UInt32, false)),
484 nullable_label_sets.offsets().clone(),
485 nullable_label_sets.values().clone(),
486 None,
487 );
488 let ts = self.timestamp_array(n);
489 RecordBatch::try_new(
490 TOPOLOGY_NODES_SCHEMA.clone(),
491 vec![
492 Arc::new(uuids),
493 Arc::new(node_ids),
494 Arc::new(type_ids),
495 Arc::new(label_sets),
496 Arc::new(ts.clone()),
497 Arc::new(ts),
498 ],
499 )
500 .map_err(pq_err)
501 }
502
503 #[must_use]
505 #[allow(clippy::type_complexity)]
506 pub fn find_pending_node(
507 &self,
508 labels: &[u32],
509 properties: &[(String, IrLiteral)],
510 ) -> Option<PendingNodeMatch> {
511 self.find_pending_nodes(labels, properties)
512 .into_iter()
513 .next()
514 }
515
516 #[must_use]
518 pub fn find_pending_nodes(
519 &self,
520 labels: &[u32],
521 properties: &[(String, IrLiteral)],
522 ) -> Vec<PendingNodeMatch> {
523 self.nodes
524 .iter()
525 .filter_map(|node| {
526 if !labels.iter().all(|wanted| node.type_ids.contains(wanted)) {
527 return None;
528 }
529 let props = self
530 .properties
531 .values()
532 .flatten()
533 .filter(|row| row.node_uuid == node.node_uuid)
534 .flat_map(|row| {
535 row.props
536 .iter()
537 .map(|(key, value)| (key.clone(), value.clone()))
538 })
539 .collect::<HashMap<_, _>>();
540 properties
541 .iter()
542 .all(|(name, value)| props.get(name) == Some(value))
543 .then(|| {
544 (
545 node.node_uuid,
546 node.node_id,
547 node.type_id,
548 node.type_ids.clone(),
549 props,
550 )
551 })
552 })
553 .collect()
554 }
555
556 #[must_use]
558 pub fn contains_pending_edge(&self, edge_uuid: &[u8; 16]) -> bool {
559 self.edges
560 .values()
561 .any(|rows| rows.iter().any(|r| &r.edge_uuid == edge_uuid))
562 }
563
564 #[must_use]
566 #[allow(clippy::type_complexity)]
567 pub fn find_pending_edge(
568 &self,
569 rel_type: &str,
570 src: &[u8; 16],
571 dst: &[u8; 16],
572 undirected: bool,
573 properties: &[(String, IrLiteral)],
574 ) -> Option<([u8; 16], [u8; 16], [u8; 16], HashMap<String, IrLiteral>)> {
575 self.edges.iter().find_map(|(stem, edges)| {
576 edges.iter().find_map(|edge| {
577 let edge_type = edge.rel_type_name.as_deref().unwrap_or(stem);
578 let direct = edge.src_uuid == *src && edge.dst_uuid == *dst;
579 let reverse = edge.src_uuid == *dst && edge.dst_uuid == *src;
580 if edge_type != rel_type || !(direct || undirected && reverse) {
581 return None;
582 }
583 let props = self
584 .edge_properties
585 .values()
586 .flatten()
587 .filter(|row| row.edge_uuid == edge.edge_uuid)
588 .flat_map(|row| {
589 row.props
590 .iter()
591 .map(|(key, value)| (key.clone(), value.clone()))
592 })
593 .collect::<HashMap<_, _>>();
594 properties
595 .iter()
596 .all(|(name, value)| props.get(name) == Some(value))
597 .then_some((edge.edge_uuid, edge.src_uuid, edge.dst_uuid, props))
598 })
599 })
600 }
601
602 #[must_use]
610 pub fn pending_incident_edge_uuids<S: std::hash::BuildHasher>(
611 &self,
612 nodes: &HashSet<[u8; 16], S>,
613 ) -> Vec<[u8; 16]> {
614 self.edges
615 .values()
616 .flatten()
617 .filter(|r| nodes.contains(&r.src_uuid) || nodes.contains(&r.dst_uuid))
618 .map(|r| r.edge_uuid)
619 .collect()
620 }
621
622 pub fn cancel_nodes<S: std::hash::BuildHasher>(
627 &mut self,
628 targets: &HashSet<[u8; 16], S>,
629 ) -> u64 {
630 let before = self.nodes.len();
631 self.nodes.retain(|r| !targets.contains(&r.node_uuid));
632 let dropped = (before - self.nodes.len()) as u64;
633 self.properties.retain(|_, rows| {
636 rows.retain(|r| !targets.contains(&r.node_uuid));
637 !rows.is_empty()
638 });
639 self.uuid_to_node_id
640 .retain(|uuid, _| !targets.contains(uuid));
641 dropped
642 }
643
644 pub fn cancel_edges<S: std::hash::BuildHasher>(
647 &mut self,
648 targets: &HashSet<[u8; 16], S>,
649 ) -> u64 {
650 let mut dropped = 0u64;
651 self.edges.retain(|_, rows| {
652 let before = rows.len();
653 rows.retain(|r| !targets.contains(&r.edge_uuid));
654 dropped += (before - rows.len()) as u64;
655 !rows.is_empty()
656 });
657 self.edge_properties.retain(|_, rows| {
658 rows.retain(|r| !targets.contains(&r.edge_uuid));
659 !rows.is_empty()
660 });
661 dropped
662 }
663
664 pub fn merge_pending_node_props(
669 &mut self,
670 node_uuid: &[u8; 16],
671 entity_type: Option<&str>,
672 props: HashMap<String, IrLiteral>,
673 ) {
674 let stem = match (self.mode, entity_type) {
675 (OntologyMode::Advisory | OntologyMode::Strict, Some(t)) => t.to_owned(),
676 _ => UNTYPED_STEM.to_owned(),
677 };
678 let rows = self.properties.entry(stem).or_default();
679 if let Some(row) = rows.iter_mut().find(|r| &r.node_uuid == node_uuid) {
680 row.props.extend(props);
681 } else {
682 rows.push(PropRow {
683 node_uuid: *node_uuid,
684 props,
685 });
686 }
687 }
688
689 pub fn add_pending_node_labels(&mut self, node_uuid: &[u8; 16], labels: &[u32]) -> u64 {
691 let Some(row) = self
692 .nodes
693 .iter_mut()
694 .find(|row| &row.node_uuid == node_uuid)
695 else {
696 return 0;
697 };
698 let before = row.type_ids.len();
699 row.type_ids.extend(labels.iter().copied());
700 row.type_ids.sort_unstable();
701 row.type_ids.dedup();
702 (row.type_ids.len() - before) as u64
703 }
704
705 pub fn remove_pending_node_labels(&mut self, node_uuid: &[u8; 16], labels: &[u32]) -> u64 {
709 let Some(row) = self
710 .nodes
711 .iter_mut()
712 .find(|row| &row.node_uuid == node_uuid)
713 else {
714 return 0;
715 };
716 let before = row.type_ids.len();
717 row.type_ids.retain(|label| !labels.contains(label));
718 (before - row.type_ids.len()) as u64
719 }
720
721 pub fn merge_pending_edge_props(
725 &mut self,
726 edge_uuid: &[u8; 16],
727 rel_type: Option<&str>,
728 props: HashMap<String, IrLiteral>,
729 ) {
730 let stem = rel_type.unwrap_or(UNTYPED_STEM).to_owned();
731 let rows = self.edge_properties.entry(stem).or_default();
732 if let Some(row) = rows.iter_mut().find(|r| &r.edge_uuid == edge_uuid) {
733 row.props.extend(props);
734 } else {
735 rows.push(EdgePropRow {
736 edge_uuid: *edge_uuid,
737 props,
738 });
739 }
740 }
741
742 pub fn remove_pending_node_props(&mut self, node_uuid: &[u8; 16], keys: &HashSet<String>) {
747 for rows in self.properties.values_mut() {
748 for row in rows.iter_mut().filter(|r| &r.node_uuid == node_uuid) {
749 row.props.retain(|k, _| !keys.contains(k));
750 }
751 }
752 }
753
754 pub fn remove_pending_edge_props(&mut self, edge_uuid: &[u8; 16], keys: &HashSet<String>) {
757 for rows in self.edge_properties.values_mut() {
758 for row in rows.iter_mut().filter(|r| &r.edge_uuid == edge_uuid) {
759 row.props.retain(|k, _| !keys.contains(k));
760 }
761 }
762 }
763
764 pub fn flush(&mut self) -> Result<(), GfError> {
782 let mut staged = RewriteBatch::new();
783 self.flush_into(&mut staged)?;
784 let pending = self.take_pending_delta();
785 if let Some(generation) = crate::generation::commit_topology_aware(staged, &self.dir)? {
786 self.write_segment_best_effort(generation, &pending);
791 }
792 Ok(())
793 }
794
795 #[must_use]
799 pub fn take_pending_delta(&mut self) -> Vec<crate::adjacency_delta::DeltaEdge> {
800 let mut edges = std::mem::take(&mut self.pending_delta);
801 edges.sort_unstable_by_key(|e| e.edge_id);
805 edges
806 }
807
808 pub fn write_segment_best_effort(
813 &self,
814 generation: u64,
815 edges: &[crate::adjacency_delta::DeltaEdge],
816 ) {
817 if crate::adjacency::adjacency_dir(&self.dir).exists() {
818 let _ = crate::adjacency_delta::write_delta_segment(&self.dir, generation, edges);
819 }
820 }
821
822 pub fn flush_into(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
837 self.flush_nodes(staged)?;
838 self.flush_edges(staged)?;
839 self.flush_properties(staged)?;
840 self.flush_edge_properties(staged)?;
841 Ok(())
842 }
843
844 fn flush_nodes(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
845 if self.nodes.is_empty() {
846 return Ok(());
847 }
848 let topology = self.dir.join("topology");
849 fs::create_dir_all(&topology).map_err(|e| io_err(&e))?;
850
851 let batch = self.pending_nodes_batch()?;
852
853 let path = topology.join("nodes.parquet");
857 let read_path = staged
858 .staged_temp(&path)
859 .map_or_else(|| path.clone(), Path::to_path_buf);
860 let existing = crate::catalog::normalize_topology_nodes(
861 crate::catalog::read_parquet_or_empty(&read_path, TOPOLOGY_NODES_SCHEMA.clone())
862 .map_err(pq_err)?,
863 )
864 .map_err(pq_err)?;
865 let merged = concat_with_existing(&TOPOLOGY_NODES_SCHEMA, existing, batch)?;
866 staged.restage(&path, TOPOLOGY_NODES_SCHEMA.clone(), &merged)?;
867 self.nodes.clear();
868 Ok(())
869 }
870
871 fn flush_edges(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
872 if self.edges.is_empty() {
873 return Ok(());
874 }
875 let edges_dir = self.dir.join("topology").join("edges");
876 fs::create_dir_all(&edges_dir).map_err(|e| io_err(&e))?;
877
878 let buffered: Vec<(String, Vec<EdgeRow>)> = self.edges.drain().collect();
880 for (stem, rows) in buffered {
881 let exploratory = stem == EXPLORATORY_STEM;
882 for r in &rows {
885 self.pending_delta.push(crate::adjacency_delta::DeltaEdge {
886 rel_type_name: if exploratory {
887 r.rel_type_name.clone().unwrap_or_default()
888 } else {
889 stem.clone()
890 },
891 edge_id: r.edge_id,
892 src_id: r.src_id,
893 dst_id: r.dst_id,
894 });
895 }
896 let schema = if exploratory {
897 EXPLORATORY_EDGE_SCHEMA.clone()
898 } else {
899 TYPED_EDGE_SCHEMA.clone()
900 };
901 let batch = self.edge_batch(&rows, &schema, exploratory)?;
902 let path = edges_dir.join(format!("{stem}.parquet"));
905 let read_path = staged
906 .staged_temp(&path)
907 .map_or_else(|| path.clone(), Path::to_path_buf);
908 let existing = crate::catalog::read_parquet_or_empty(&read_path, schema.clone())
909 .map_err(pq_err)?;
910 let merged = concat_with_existing(&schema, existing, batch)?;
911 staged.restage(&path, schema, &merged)?;
912 }
913 Ok(())
914 }
915
916 fn edge_batch(
917 &self,
918 rows: &[EdgeRow],
919 schema: &SchemaRef,
920 exploratory: bool,
921 ) -> Result<RecordBatch, GfError> {
922 let edge_uuids =
923 FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.edge_uuid.to_vec()))
924 .map_err(pq_err)?;
925 let src_uuids =
926 FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.src_uuid.to_vec()))
927 .map_err(pq_err)?;
928 let dst_uuids =
929 FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.dst_uuid.to_vec()))
930 .map_err(pq_err)?;
931 let edge_ids = UInt64Array::from(rows.iter().map(|r| r.edge_id).collect::<Vec<_>>());
932 let src_ids = UInt64Array::from(rows.iter().map(|r| r.src_id).collect::<Vec<_>>());
933 let dst_ids = UInt64Array::from(rows.iter().map(|r| r.dst_id).collect::<Vec<_>>());
934 let ts = self.timestamp_array(rows.len());
935 let mut cols: Vec<ArrayRef> = vec![
936 Arc::new(edge_uuids),
937 Arc::new(src_uuids),
938 Arc::new(dst_uuids),
939 Arc::new(edge_ids),
940 Arc::new(src_ids),
941 Arc::new(dst_ids),
942 Arc::new(ts),
943 ];
944 if exploratory {
945 let names = StringArray::from(
946 rows.iter()
947 .map(|r| r.rel_type_name.clone().unwrap_or_default())
948 .collect::<Vec<_>>(),
949 );
950 cols.push(Arc::new(names));
951 }
952 RecordBatch::try_new(schema.clone(), cols).map_err(pq_err)
953 }
954
955 fn flush_properties(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
956 if self.properties.is_empty() {
957 return Ok(());
958 }
959 let buffered: Vec<(String, Vec<PropRow>)> = self.properties.drain().collect();
960 for (stem, new_rows) in buffered {
961 let existing = read_props_through(staged, &node_props_path(&self.dir, &stem))?;
966 let mut rows = decode_property_rows(&existing)?;
967 rows.extend(new_rows);
968 let (schema, cols) = build_property_columns(&stem, &rows)?;
969 stage_property_file(staged, &self.dir, "properties", &stem, schema, cols)?;
970 }
971 Ok(())
972 }
973
974 fn flush_edge_properties(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
978 if self.edge_properties.is_empty() {
979 return Ok(());
980 }
981 let buffered: Vec<(String, Vec<EdgePropRow>)> = self.edge_properties.drain().collect();
982 for (stem, new_rows) in buffered {
983 let existing = read_props_through(staged, &edge_props_path(&self.dir, &stem))?;
984 let mut rows = decode_edge_property_rows(&existing)?;
985 rows.extend(new_rows);
986 let (schema, cols) = build_property_columns_keyed(
987 EDGE_PROPERTY_UUID_FIELD,
988 "graphforge.rel_type",
989 &stem,
990 &rows,
991 )?;
992 stage_property_file(staged, &self.dir, "edge_properties", &stem, schema, cols)?;
993 }
994 Ok(())
995 }
996
997 fn timestamp_array(&self, n: usize) -> TimestampMicrosecondArray {
998 TimestampMicrosecondArray::from(vec![self.now_micros; n])
999 .with_timezone_opt(Some(Arc::from("UTC")))
1000 }
1001}
1002
1003#[derive(Clone, PartialEq, Eq)]
1010enum ColType {
1011 Int,
1012 Float,
1013 Bool,
1014 Str,
1015 HetScalar,
1016 Duration,
1017 DateTime,
1018 Date,
1019 LocalDateTime,
1020 Time,
1021 ZonedTime,
1022 ZonedDateTime,
1023 List(Box<ColType>),
1025}
1026
1027impl ColType {
1028 fn of(lit: &IrLiteral) -> Option<Self> {
1029 match lit {
1030 IrLiteral::Null
1031 | IrLiteral::Map(_)
1033 | IrLiteral::Uuid(_) => None,
1035 IrLiteral::Int(_) => Some(Self::Int),
1036 IrLiteral::Float(_) => Some(Self::Float),
1037 IrLiteral::Bool(_) => Some(Self::Bool),
1038 IrLiteral::Str(_) => Some(Self::Str),
1039 IrLiteral::Duration { .. } => Some(Self::Duration),
1040 IrLiteral::DateTime(_) => Some(Self::DateTime),
1041 IrLiteral::Date(_) => Some(Self::Date),
1042 IrLiteral::LocalDateTime { .. } => Some(Self::LocalDateTime),
1043 IrLiteral::Time(_) => Some(Self::Time),
1044 IrLiteral::ZonedTime { .. } => Some(Self::ZonedTime),
1045 IrLiteral::ZonedDateTime { .. } => Some(Self::ZonedDateTime),
1046 IrLiteral::List(items) => items
1050 .iter()
1051 .find_map(Self::of)
1052 .map(|inner| Self::List(Box::new(inner))),
1053 }
1054 }
1055
1056 fn data_type(&self) -> DataType {
1057 match self {
1058 Self::Int => DataType::Int64,
1059 Self::Float => DataType::Float64,
1060 Self::Bool => DataType::Boolean,
1061 Self::Str => DataType::Utf8,
1063 Self::HetScalar => DataType::Struct(heterogeneous_scalar_fields()),
1064 Self::Duration => DataType::Struct(crate::schemas::duration_struct_fields()),
1065 Self::DateTime => DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1066 Self::Date => DataType::Struct(crate::schemas::date_struct_fields()),
1067 Self::LocalDateTime => DataType::Struct(crate::schemas::localdatetime_struct_fields()),
1068 Self::Time => DataType::Time64(TimeUnit::Nanosecond),
1069 Self::ZonedTime => DataType::Struct(crate::schemas::time_struct_fields()),
1070 Self::ZonedDateTime => DataType::Struct(crate::schemas::datetime_struct_fields()),
1071 Self::List(inner) => {
1072 DataType::List(Arc::new(Field::new("item", inner.data_type(), true)))
1073 }
1074 }
1075 }
1076
1077 fn is_scalar(&self) -> bool {
1078 matches!(
1079 self,
1080 Self::Int | Self::Float | Self::Bool | Self::Str | Self::HetScalar
1081 )
1082 }
1083}
1084
1085fn heterogeneous_scalar_fields() -> arrow::datatypes::Fields {
1086 arrow::datatypes::Fields::from(vec![
1087 Field::new("__het_tag", DataType::Int8, false),
1088 Field::new("__het_int", DataType::Int64, true),
1089 Field::new("__het_float", DataType::Float64, true),
1090 Field::new("__het_str", DataType::Utf8, true),
1091 Field::new("__het_bool", DataType::Boolean, true),
1092 ])
1093}
1094
1095fn build_property_columns(
1098 entity_type: &str,
1099 rows: &[PropRow],
1100) -> Result<(Schema, Vec<ArrayRef>), GfError> {
1101 build_property_columns_keyed(
1102 NODE_PROPERTY_UUID_FIELD,
1103 "graphforge.entity_type",
1104 entity_type,
1105 rows,
1106 )
1107}
1108
1109fn build_property_columns_keyed<R: PropRowLike>(
1120 uuid_field_name: &str,
1121 meta_key: &str,
1122 meta_value: &str,
1123 rows: &[R],
1124) -> Result<(Schema, Vec<ArrayRef>), GfError> {
1125 let mut order: Vec<String> = Vec::new();
1130 let mut seen: HashSet<String> = HashSet::new();
1131 let mut col_types: HashMap<String, ColType> = HashMap::new();
1132 for row in rows {
1133 for (name, lit) in row.props() {
1134 reject_map_property_value(name, lit)?;
1135 if seen.insert(name.clone()) {
1136 order.push(name.clone());
1137 }
1138 if let Some(t) = ColType::of(lit) {
1143 col_types
1144 .entry(name.clone())
1145 .and_modify(|existing| {
1146 if *existing != t {
1147 *existing = if existing.is_scalar() && t.is_scalar() {
1148 ColType::HetScalar
1149 } else {
1150 ColType::Str
1151 };
1152 }
1153 })
1154 .or_insert(t);
1155 }
1156 }
1157 }
1158
1159 let mut fields: Vec<Field> = vec![uuid_field(uuid_field_name)];
1160 for name in &order {
1161 let ct = col_types.get(name).cloned().unwrap_or(ColType::Str);
1162 fields.push(Field::new(name, ct.data_type(), true));
1163 }
1164 let meta: HashMap<String, String> = [(meta_key.to_owned(), meta_value.to_owned())]
1165 .into_iter()
1166 .collect();
1167 let schema = Schema::new(fields).with_metadata(meta);
1168
1169 let mut cols: Vec<ArrayRef> = Vec::with_capacity(order.len() + 1);
1170 let uuids = FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.uuid_bytes().to_vec()))
1171 .map_err(pq_err)?;
1172 cols.push(Arc::new(uuids));
1173
1174 for name in &order {
1175 let ct = col_types.get(name).cloned().unwrap_or(ColType::Str);
1176 cols.push(build_property_array(name, ct, rows));
1177 }
1178 Ok((schema, cols))
1179}
1180
1181fn reject_map_property_value(name: &str, lit: &IrLiteral) -> Result<(), GfError> {
1182 if contains_uuid_literal(lit) {
1183 return Err(GfError::Validation(format!(
1184 "property `{name}` cannot store typed UUID query parameters"
1185 )));
1186 }
1187 if contains_map_literal(lit) {
1188 return Err(GfError::Storage(format!(
1189 "property `{name}` cannot store map values"
1190 )));
1191 }
1192 Ok(())
1193}
1194
1195fn contains_map_literal(lit: &IrLiteral) -> bool {
1196 match lit {
1197 IrLiteral::Map(_) => true,
1198 IrLiteral::List(items) => items.iter().any(contains_map_literal),
1199 _ => false,
1200 }
1201}
1202
1203fn contains_uuid_literal(lit: &IrLiteral) -> bool {
1204 match lit {
1205 IrLiteral::Uuid(_) => true,
1206 IrLiteral::List(items) => items.iter().any(contains_uuid_literal),
1207 IrLiteral::Map(entries) => entries
1208 .iter()
1209 .any(|(_, value)| contains_uuid_literal(value)),
1210 _ => false,
1211 }
1212}
1213
1214#[allow(
1218 clippy::too_many_lines,
1219 reason = "one builder arm per ColType; the per-type append loops read clearest inline"
1220)]
1221fn build_property_array<R: PropRowLike>(name: &str, ct: ColType, rows: &[R]) -> ArrayRef {
1222 match ct {
1223 ColType::Int => {
1224 let mut b = Int64Builder::new();
1225 for row in rows {
1226 match row.props().get(name) {
1227 Some(IrLiteral::Int(v)) => b.append_value(*v),
1228 _ => b.append_null(),
1229 }
1230 }
1231 Arc::new(b.finish())
1232 }
1233 ColType::Float => {
1234 let mut b = Float64Builder::new();
1235 for row in rows {
1236 match row.props().get(name) {
1237 Some(IrLiteral::Float(v)) => b.append_value(*v),
1238 _ => b.append_null(),
1239 }
1240 }
1241 Arc::new(b.finish())
1242 }
1243 ColType::Bool => {
1244 let mut b = BooleanBuilder::new();
1245 for row in rows {
1246 match row.props().get(name) {
1247 Some(IrLiteral::Bool(v)) => b.append_value(*v),
1248 _ => b.append_null(),
1249 }
1250 }
1251 Arc::new(b.finish())
1252 }
1253 ColType::Duration => {
1254 use arrow::array::{Int64Builder, StructArray};
1258 use arrow::buffer::NullBuffer;
1259 let (mut mb, mut db, mut sb, mut nb) = (
1260 Int64Builder::new(),
1261 Int64Builder::new(),
1262 Int64Builder::new(),
1263 Int64Builder::new(),
1264 );
1265 let mut valid = Vec::with_capacity(rows.len());
1266 for row in rows {
1267 if let Some(IrLiteral::Duration {
1268 months,
1269 days,
1270 seconds,
1271 nanos,
1272 }) = row.props().get(name)
1273 {
1274 mb.append_value(*months);
1275 db.append_value(*days);
1276 sb.append_value(*seconds);
1277 nb.append_value(*nanos);
1278 valid.push(true);
1279 } else {
1280 mb.append_null();
1281 db.append_null();
1282 sb.append_null();
1283 nb.append_null();
1284 valid.push(false);
1285 }
1286 }
1287 Arc::new(StructArray::new(
1288 crate::schemas::duration_struct_fields(),
1289 vec![
1290 Arc::new(mb.finish()),
1291 Arc::new(db.finish()),
1292 Arc::new(sb.finish()),
1293 Arc::new(nb.finish()),
1294 ],
1295 Some(NullBuffer::from(valid)),
1296 ))
1297 }
1298 ColType::DateTime => {
1299 let mut b = TimestampMicrosecondBuilder::new();
1300 for row in rows {
1301 match row.props().get(name) {
1302 Some(IrLiteral::DateTime(v)) => b.append_value(*v),
1303 _ => b.append_null(),
1304 }
1305 }
1306 Arc::new(b.finish().with_timezone_opt(Some(Arc::from("UTC"))))
1307 }
1308 ColType::Date => {
1309 use arrow::array::{Int64Builder, StructArray};
1312 use arrow::buffer::NullBuffer;
1313 let mut b = Int64Builder::new();
1314 let mut valid = Vec::with_capacity(rows.len());
1315 for row in rows {
1316 if let Some(IrLiteral::Date(v)) = row.props().get(name) {
1317 b.append_value(*v);
1318 valid.push(true);
1319 } else {
1320 b.append_null();
1321 valid.push(false);
1322 }
1323 }
1324 Arc::new(StructArray::new(
1325 crate::schemas::date_struct_fields(),
1326 vec![Arc::new(b.finish())],
1327 Some(NullBuffer::from(valid)),
1328 ))
1329 }
1330 ColType::LocalDateTime => {
1331 use arrow::array::{Int64Builder, StructArray, Time64NanosecondBuilder};
1334 use arrow::buffer::NullBuffer;
1335 let (mut date_b, mut time_b) = (Int64Builder::new(), Time64NanosecondBuilder::new());
1336 let mut valid = Vec::with_capacity(rows.len());
1337 for row in rows {
1338 if let Some(IrLiteral::LocalDateTime { days, nanos }) = row.props().get(name) {
1339 date_b.append_value(*days);
1340 time_b.append_value(*nanos);
1341 valid.push(true);
1342 } else {
1343 date_b.append_null();
1344 time_b.append_null();
1345 valid.push(false);
1346 }
1347 }
1348 Arc::new(StructArray::new(
1349 crate::schemas::localdatetime_struct_fields(),
1350 vec![Arc::new(date_b.finish()), Arc::new(time_b.finish())],
1351 Some(NullBuffer::from(valid)),
1352 ))
1353 }
1354 ColType::Time => {
1355 use arrow::array::Time64NanosecondBuilder;
1356 let mut b = Time64NanosecondBuilder::new();
1357 for row in rows {
1358 match row.props().get(name) {
1359 Some(IrLiteral::Time(v)) => b.append_value(*v),
1360 _ => b.append_null(),
1361 }
1362 }
1363 Arc::new(b.finish())
1364 }
1365 ColType::ZonedTime => {
1366 use arrow::array::{Int32Builder, StructArray, Time64NanosecondBuilder};
1368 use arrow::buffer::NullBuffer;
1369 let (mut time_b, mut off_b) = (Time64NanosecondBuilder::new(), Int32Builder::new());
1370 let mut valid = Vec::with_capacity(rows.len());
1371 for row in rows {
1372 if let Some(IrLiteral::ZonedTime { nanos, offset }) = row.props().get(name) {
1373 time_b.append_value(*nanos);
1374 off_b.append_value(*offset);
1375 valid.push(true);
1376 } else {
1377 time_b.append_null();
1378 off_b.append_null();
1379 valid.push(false);
1380 }
1381 }
1382 Arc::new(StructArray::new(
1383 crate::schemas::time_struct_fields(),
1384 vec![Arc::new(time_b.finish()), Arc::new(off_b.finish())],
1385 Some(NullBuffer::from(valid)),
1386 ))
1387 }
1388 ColType::ZonedDateTime => {
1389 use arrow::array::{
1392 Int32Builder, Int64Builder, StringBuilder, StructArray, Time64NanosecondBuilder,
1393 };
1394 use arrow::buffer::NullBuffer;
1395 let (mut date_b, mut time_b, mut off_b, mut zone_b) = (
1396 Int64Builder::new(),
1397 Time64NanosecondBuilder::new(),
1398 Int32Builder::new(),
1399 StringBuilder::new(),
1400 );
1401 let mut valid = Vec::with_capacity(rows.len());
1402 for row in rows {
1403 if let Some(IrLiteral::ZonedDateTime {
1404 days,
1405 nanos,
1406 offset,
1407 zone,
1408 }) = row.props().get(name)
1409 {
1410 date_b.append_value(*days);
1411 time_b.append_value(*nanos);
1412 off_b.append_value(*offset);
1413 zone_b.append_option(zone.as_deref());
1415 valid.push(true);
1416 } else {
1417 date_b.append_null();
1418 time_b.append_null();
1419 off_b.append_null();
1420 zone_b.append_null();
1421 valid.push(false);
1422 }
1423 }
1424 Arc::new(StructArray::new(
1425 crate::schemas::datetime_struct_fields(),
1426 vec![
1427 Arc::new(date_b.finish()),
1428 Arc::new(time_b.finish()),
1429 Arc::new(off_b.finish()),
1430 Arc::new(zone_b.finish()),
1431 ],
1432 Some(NullBuffer::from(valid)),
1433 ))
1434 }
1435 ColType::Str => {
1436 let mut b = StringBuilder::new();
1437 for row in rows {
1438 match row.props().get(name) {
1439 Some(IrLiteral::Null) | None => b.append_null(),
1440 Some(other) => b.append_value(literal_to_string(other)),
1441 }
1442 }
1443 Arc::new(b.finish())
1444 }
1445 ColType::HetScalar => build_heterogeneous_scalar_array(name, rows),
1446 ColType::List(inner) => {
1447 use arrow::array::ListArray;
1452 use arrow::buffer::{NullBuffer, OffsetBuffer};
1453 let mut elem_rows: Vec<PropRow> = Vec::new();
1454 let mut offsets: Vec<i32> = vec![0];
1455 let mut valid = Vec::with_capacity(rows.len());
1456 for row in rows {
1457 if let Some(IrLiteral::List(items)) = row.props().get(name) {
1458 for it in items {
1459 let mut props = HashMap::with_capacity(1);
1460 props.insert("item".to_string(), it.clone());
1461 elem_rows.push(PropRow {
1462 node_uuid: [0u8; 16],
1463 props,
1464 });
1465 }
1466 valid.push(true);
1467 } else {
1468 valid.push(false);
1469 }
1470 offsets.push(i32::try_from(elem_rows.len()).unwrap_or(i32::MAX));
1471 }
1472 let child = build_property_array("item", (*inner).clone(), &elem_rows);
1473 let field = Arc::new(Field::new("item", inner.data_type(), true));
1474 Arc::new(ListArray::new(
1475 field,
1476 OffsetBuffer::new(offsets.into()),
1477 child,
1478 Some(NullBuffer::from(valid)),
1479 ))
1480 }
1481 }
1482}
1483
1484fn build_heterogeneous_scalar_array<R: PropRowLike>(name: &str, rows: &[R]) -> ArrayRef {
1485 use arrow::array::{BooleanBuilder, Float64Builder, Int8Builder, Int64Builder, StructArray};
1486 use arrow::buffer::NullBuffer;
1487
1488 let mut tags = Int8Builder::new();
1489 let mut ints = Int64Builder::new();
1490 let mut floats = Float64Builder::new();
1491 let mut strings = StringBuilder::new();
1492 let mut bools = BooleanBuilder::new();
1493 let mut valid = Vec::with_capacity(rows.len());
1494 for row in rows {
1495 let value = row.props().get(name);
1496 let tag = match value {
1497 Some(IrLiteral::Int(value)) => {
1498 ints.append_value(*value);
1499 floats.append_null();
1500 strings.append_null();
1501 bools.append_null();
1502 Some(0)
1503 }
1504 Some(IrLiteral::Float(value)) => {
1505 ints.append_null();
1506 floats.append_value(*value);
1507 strings.append_null();
1508 bools.append_null();
1509 Some(1)
1510 }
1511 Some(IrLiteral::Str(value)) => {
1512 ints.append_null();
1513 floats.append_null();
1514 strings.append_value(value);
1515 bools.append_null();
1516 Some(2)
1517 }
1518 Some(IrLiteral::Bool(value)) => {
1519 ints.append_null();
1520 floats.append_null();
1521 strings.append_null();
1522 bools.append_value(*value);
1523 Some(3)
1524 }
1525 _ => {
1526 ints.append_null();
1527 floats.append_null();
1528 strings.append_null();
1529 bools.append_null();
1530 None
1531 }
1532 };
1533 tags.append_value(tag.unwrap_or_default());
1534 valid.push(tag.is_some());
1535 }
1536 Arc::new(StructArray::new(
1537 heterogeneous_scalar_fields(),
1538 vec![
1539 Arc::new(tags.finish()),
1540 Arc::new(ints.finish()),
1541 Arc::new(floats.finish()),
1542 Arc::new(strings.finish()),
1543 Arc::new(bools.finish()),
1544 ],
1545 Some(NullBuffer::from(valid)),
1546 ))
1547}
1548
1549fn literal_to_string(lit: &IrLiteral) -> String {
1551 match lit {
1552 IrLiteral::Null => String::new(),
1553 IrLiteral::Bool(b) => b.to_string(),
1554 IrLiteral::Int(i) => i.to_string(),
1555 IrLiteral::Float(f) => f.to_string(),
1556 IrLiteral::Str(s) => s.clone(),
1557 IrLiteral::Uuid(bytes) => {
1558 let mut encoded = String::with_capacity(32);
1559 for byte in bytes {
1560 std::fmt::Write::write_fmt(&mut encoded, format_args!("{byte:02x}"))
1561 .expect("writing to a String cannot fail");
1562 }
1563 encoded
1564 }
1565 IrLiteral::Duration {
1568 months,
1569 days,
1570 seconds,
1571 nanos,
1572 } => format!("{months}mo{days}d{seconds}s{nanos}ns"),
1573 IrLiteral::DateTime(t) => t.to_string(),
1574 IrLiteral::Date(d) => d.to_string(),
1575 IrLiteral::LocalDateTime { days, nanos } => format!("{days}d{nanos}ns"),
1578 IrLiteral::Time(nanos) => format!("{nanos}ns"),
1579 IrLiteral::ZonedTime { nanos, offset } => format!("{nanos}ns{offset:+}s"),
1580 IrLiteral::ZonedDateTime {
1581 days,
1582 nanos,
1583 offset,
1584 zone,
1585 } => format!(
1586 "{days}d{nanos}ns{offset:+}s{}",
1587 zone.as_deref().unwrap_or("")
1588 ),
1589 IrLiteral::List(items) => {
1591 let parts: Vec<String> = items.iter().map(literal_to_string).collect();
1592 format!("[{}]", parts.join(","))
1593 }
1594 IrLiteral::Map(entries) => {
1595 let parts: Vec<String> = entries
1596 .iter()
1597 .map(|(key, value)| format!("{key}:{}", literal_to_string(value)))
1598 .collect();
1599 format!("{{{}}}", parts.join(","))
1600 }
1601 }
1602}
1603
1604fn decode_property_rows(batches: &[RecordBatch]) -> Result<Vec<PropRow>, GfError> {
1616 let mut out = Vec::new();
1617 for batch in batches {
1618 decode_property_batch(batch, NODE_PROPERTY_UUID_FIELD, |node_uuid, props| {
1619 out.push(PropRow { node_uuid, props });
1620 })?;
1621 }
1622 Ok(out)
1623}
1624
1625fn decode_edge_property_rows(batches: &[RecordBatch]) -> Result<Vec<EdgePropRow>, GfError> {
1629 let mut out = Vec::new();
1630 for batch in batches {
1631 decode_property_batch(batch, EDGE_PROPERTY_UUID_FIELD, |edge_uuid, props| {
1632 out.push(EdgePropRow { edge_uuid, props });
1633 })?;
1634 }
1635 Ok(out)
1636}
1637
1638pub fn read_entity_property_keys(
1643 dir: &Path,
1644 stem: &str,
1645 uuid: &[u8; 16],
1646 is_edge: bool,
1647) -> Result<HashSet<String>, GfError> {
1648 let path = if is_edge {
1649 edge_props_path(dir, stem)
1650 } else {
1651 node_props_path(dir, stem)
1652 };
1653 let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
1654 return Ok(HashSet::new());
1655 };
1656 let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1657 let rows = if is_edge {
1658 decode_edge_property_rows(&batches)?
1659 .into_iter()
1660 .map(|row| (row.edge_uuid, row.props))
1661 .collect::<Vec<_>>()
1662 } else {
1663 decode_property_rows(&batches)?
1664 .into_iter()
1665 .map(|row| (row.node_uuid, row.props))
1666 .collect::<Vec<_>>()
1667 };
1668 Ok(rows
1669 .into_iter()
1670 .find_map(|(row_uuid, props)| (row_uuid == *uuid).then(|| props.into_keys().collect()))
1671 .unwrap_or_default())
1672}
1673
1674pub fn read_entity_properties(
1678 dir: &Path,
1679 stem: &str,
1680 uuid: &[u8; 16],
1681 is_edge: bool,
1682) -> Result<HashMap<String, IrLiteral>, GfError> {
1683 let path = if is_edge {
1684 edge_props_path(dir, stem)
1685 } else {
1686 node_props_path(dir, stem)
1687 };
1688 let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
1689 return Ok(HashMap::new());
1690 };
1691 let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1692 let rows = if is_edge {
1693 decode_edge_property_rows(&batches)?
1694 .into_iter()
1695 .map(|row| (row.edge_uuid, row.props))
1696 .collect::<Vec<_>>()
1697 } else {
1698 decode_property_rows(&batches)?
1699 .into_iter()
1700 .map(|row| (row.node_uuid, row.props))
1701 .collect::<Vec<_>>()
1702 };
1703 Ok(rows
1704 .into_iter()
1705 .find_map(|(row_uuid, props)| (row_uuid == *uuid).then_some(props))
1706 .unwrap_or_default())
1707}
1708
1709pub fn read_node_property_rows(
1711 dir: &Path,
1712 stem: &str,
1713) -> Result<HashMap<[u8; 16], HashMap<String, IrLiteral>>, GfError> {
1714 let path = node_props_path(dir, stem);
1715 if !path.try_exists().map_err(|error| io_err(&error))? {
1716 return Ok(HashMap::new());
1717 }
1718 let file = fs::File::open(&path).map_err(|error| io_err(&error))?;
1719 let schema = ParquetRecordBatchReaderBuilder::try_new(file)
1720 .map_err(pq_err)?
1721 .schema()
1722 .clone();
1723 let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1724 Ok(decode_property_rows(&batches)?
1725 .into_iter()
1726 .map(|row| (row.node_uuid, row.props))
1727 .collect())
1728}
1729
1730pub fn count_entity_properties<S: std::hash::BuildHasher>(
1733 dir: &Path,
1734 targets: &HashSet<[u8; 16], S>,
1735 is_edge: bool,
1736) -> Result<u64, GfError> {
1737 if targets.is_empty() {
1738 return Ok(0);
1739 }
1740 let property_dir = dir.join(if is_edge {
1741 "edge_properties"
1742 } else {
1743 "properties"
1744 });
1745 let Ok(entries) = std::fs::read_dir(property_dir) else {
1746 return Ok(0);
1747 };
1748 let mut count = 0u64;
1749 for entry in entries {
1750 let path = entry
1751 .map_err(|error| GfError::Storage(error.to_string()))?
1752 .path();
1753 if path
1754 .extension()
1755 .is_none_or(|extension| extension != "parquet")
1756 {
1757 continue;
1758 }
1759 let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
1760 continue;
1761 };
1762 let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1763 if is_edge {
1764 for row in decode_edge_property_rows(&batches)? {
1765 if targets.contains(&row.edge_uuid) {
1766 count += row.props.len() as u64;
1767 }
1768 }
1769 } else {
1770 for row in decode_property_rows(&batches)? {
1771 if targets.contains(&row.node_uuid) {
1772 count += row.props.len() as u64;
1773 }
1774 }
1775 }
1776 }
1777 Ok(count)
1778}
1779
1780#[allow(
1787 clippy::too_many_lines,
1788 reason = "one decode arm per Arrow type, plus per-shape struct dispatch; clearest inline"
1789)]
1790fn decode_property_batch(
1791 batch: &RecordBatch,
1792 uuid_field_name: &str,
1793 mut emit: impl FnMut([u8; 16], HashMap<String, IrLiteral>),
1794) -> Result<(), GfError> {
1795 use arrow::array::Array;
1796
1797 let schema = batch.schema();
1798 let uuid_col = batch
1799 .column_by_name(uuid_field_name)
1800 .and_then(|c| c.as_any().downcast_ref::<FixedSizeBinaryArray>())
1801 .ok_or_else(|| {
1802 GfError::Storage(format!("property file missing {uuid_field_name} column"))
1803 })?;
1804 for r in 0..batch.num_rows() {
1805 let mut uuid = [0u8; 16];
1806 uuid.copy_from_slice(uuid_col.value(r));
1807 let mut props: HashMap<String, IrLiteral> = HashMap::new();
1808 for (c, field) in schema.fields().iter().enumerate() {
1809 if field.name() == uuid_field_name {
1810 continue;
1811 }
1812 let col = batch.column(c);
1813 if col.is_null(r) {
1814 continue; }
1816 let lit = decode_value(col, field, r)?;
1817 props.insert(field.name().clone(), lit);
1818 }
1819 emit(uuid, props);
1820 }
1821 Ok(())
1822}
1823
1824#[allow(
1829 clippy::too_many_lines,
1830 reason = "one decode arm per Arrow type, plus per-shape struct dispatch; clearest inline"
1831)]
1832fn decode_value(
1833 col: &arrow::array::ArrayRef,
1834 field: &arrow::datatypes::Field,
1835 r: usize,
1836) -> Result<IrLiteral, GfError> {
1837 use arrow::array::{
1838 Array, BooleanArray, Int32Array, Int64Array, ListArray, StructArray, Time64NanosecondArray,
1839 };
1840 Ok(match field.data_type() {
1841 DataType::Int64 => IrLiteral::Int(downcast::<Int64Array>(col, field)?.value(r)),
1842 DataType::Float64 => IrLiteral::Float(downcast::<Float64Array>(col, field)?.value(r)),
1843 DataType::Boolean => IrLiteral::Bool(downcast::<BooleanArray>(col, field)?.value(r)),
1844 DataType::Utf8 => IrLiteral::Str(downcast::<StringArray>(col, field)?.value(r).to_owned()),
1845 DataType::Struct(fields) => {
1850 let s = downcast::<StructArray>(col, field)?;
1851 let names: Vec<&str> = fields.iter().map(|f| f.name().as_str()).collect();
1852 match names.as_slice() {
1853 [
1854 "__het_tag",
1855 "__het_int",
1856 "__het_float",
1857 "__het_str",
1858 "__het_bool",
1859 ] => {
1860 let tag = s
1861 .column(0)
1862 .as_any()
1863 .downcast_ref::<arrow::array::Int8Array>()
1864 .ok_or_else(|| GfError::Storage("heterogeneous tag not Int8".into()))?
1865 .value(r);
1866 match tag {
1867 0 => IrLiteral::Int(
1868 s.column(1)
1869 .as_any()
1870 .downcast_ref::<Int64Array>()
1871 .ok_or_else(|| {
1872 GfError::Storage("heterogeneous int not Int64".into())
1873 })?
1874 .value(r),
1875 ),
1876 1 => IrLiteral::Float(
1877 s.column(2)
1878 .as_any()
1879 .downcast_ref::<Float64Array>()
1880 .ok_or_else(|| {
1881 GfError::Storage("heterogeneous float not Float64".into())
1882 })?
1883 .value(r),
1884 ),
1885 2 => IrLiteral::Str(
1886 s.column(3)
1887 .as_any()
1888 .downcast_ref::<StringArray>()
1889 .ok_or_else(|| {
1890 GfError::Storage("heterogeneous string not Utf8".into())
1891 })?
1892 .value(r)
1893 .to_owned(),
1894 ),
1895 3 => IrLiteral::Bool(
1896 s.column(4)
1897 .as_any()
1898 .downcast_ref::<BooleanArray>()
1899 .ok_or_else(|| {
1900 GfError::Storage("heterogeneous bool not Boolean".into())
1901 })?
1902 .value(r),
1903 ),
1904 _ => {
1905 return Err(GfError::Storage(format!(
1906 "unsupported heterogeneous property tag {tag}"
1907 )));
1908 }
1909 }
1910 }
1911 ["months", "days", "seconds", "nanos"] => {
1912 let i64_at = |idx: usize| -> Result<i64, GfError> {
1913 Ok(s.column(idx)
1914 .as_any()
1915 .downcast_ref::<Int64Array>()
1916 .ok_or_else(|| {
1917 GfError::Storage("duration struct child not Int64".into())
1918 })?
1919 .value(r))
1920 };
1921 IrLiteral::Duration {
1922 months: i64_at(0)?,
1923 days: i64_at(1)?,
1924 seconds: i64_at(2)?,
1925 nanos: i64_at(3)?,
1926 }
1927 }
1928 ["epoch_day"] => IrLiteral::Date(
1929 s.column(0)
1930 .as_any()
1931 .downcast_ref::<Int64Array>()
1932 .ok_or_else(|| GfError::Storage("date epoch_day not Int64".into()))?
1933 .value(r),
1934 ),
1935 ["date", "time"] => {
1936 let days = s
1937 .column(0)
1938 .as_any()
1939 .downcast_ref::<Int64Array>()
1940 .ok_or_else(|| GfError::Storage("localdatetime date not Int64".into()))?
1941 .value(r);
1942 let nanos = s
1943 .column(1)
1944 .as_any()
1945 .downcast_ref::<Time64NanosecondArray>()
1946 .ok_or_else(|| {
1947 GfError::Storage("localdatetime time not Time64(ns)".into())
1948 })?
1949 .value(r);
1950 IrLiteral::LocalDateTime { days, nanos }
1951 }
1952 ["time", "offset"] => {
1953 let nanos = s
1954 .column(0)
1955 .as_any()
1956 .downcast_ref::<Time64NanosecondArray>()
1957 .ok_or_else(|| GfError::Storage("time not Time64(ns)".into()))?
1958 .value(r);
1959 let offset = s
1960 .column(1)
1961 .as_any()
1962 .downcast_ref::<Int32Array>()
1963 .ok_or_else(|| GfError::Storage("time offset not Int32".into()))?
1964 .value(r);
1965 IrLiteral::ZonedTime { nanos, offset }
1966 }
1967 ["date", "time", "offset", "zone"] => {
1968 let days = s
1969 .column(0)
1970 .as_any()
1971 .downcast_ref::<Int64Array>()
1972 .ok_or_else(|| GfError::Storage("datetime date not Int64".into()))?
1973 .value(r);
1974 let nanos = s
1975 .column(1)
1976 .as_any()
1977 .downcast_ref::<Time64NanosecondArray>()
1978 .ok_or_else(|| GfError::Storage("datetime time not Time64(ns)".into()))?
1979 .value(r);
1980 let offset = s
1981 .column(2)
1982 .as_any()
1983 .downcast_ref::<Int32Array>()
1984 .ok_or_else(|| GfError::Storage("datetime offset not Int32".into()))?
1985 .value(r);
1986 let zone_col = s
1987 .column(3)
1988 .as_any()
1989 .downcast_ref::<StringArray>()
1990 .ok_or_else(|| GfError::Storage("datetime zone not Utf8".into()))?;
1991 let zone = (!zone_col.is_null(r)).then(|| zone_col.value(r).to_owned());
1993 IrLiteral::ZonedDateTime {
1994 days,
1995 nanos,
1996 offset,
1997 zone,
1998 }
1999 }
2000 _ => {
2001 return Err(GfError::Storage(format!(
2002 "property column {} has unsupported struct shape {names:?}",
2003 field.name()
2004 )));
2005 }
2006 }
2007 }
2008 DataType::Time64(TimeUnit::Nanosecond) => {
2009 IrLiteral::Time(downcast::<Time64NanosecondArray>(col, field)?.value(r))
2010 }
2011 DataType::Timestamp(TimeUnit::Microsecond, _) => {
2012 IrLiteral::DateTime(downcast::<TimestampMicrosecondArray>(col, field)?.value(r))
2013 }
2014 DataType::List(inner) => {
2018 let larr = downcast::<ListArray>(col, field)?;
2019 let elems = larr.value(r);
2020 let mut items = Vec::with_capacity(elems.len());
2021 for j in 0..elems.len() {
2022 if elems.is_null(j) {
2023 items.push(IrLiteral::Null);
2024 } else {
2025 items.push(decode_value(&elems, inner, j)?);
2026 }
2027 }
2028 IrLiteral::List(items)
2029 }
2030 other => {
2031 return Err(GfError::Storage(format!(
2032 "property column {} has unsupported type {other:?}",
2033 field.name()
2034 )));
2035 }
2036 })
2037}
2038
2039fn downcast<'a, A: 'static>(
2042 col: &'a arrow::array::ArrayRef,
2043 field: &arrow::datatypes::Field,
2044) -> Result<&'a A, GfError> {
2045 col.as_any().downcast_ref::<A>().ok_or_else(|| {
2046 GfError::Storage(format!(
2047 "property column {} could not be read as its declared type",
2048 field.name()
2049 ))
2050 })
2051}
2052
2053fn apply_property_updates<R: PropRowLike>(
2075 mut rows: Vec<R>,
2076 updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2077) -> (Vec<R>, u64) {
2078 let mut index: HashMap<[u8; 16], usize> = HashMap::with_capacity(rows.len());
2081 for (i, row) in rows.iter().enumerate() {
2082 index.entry(*row.uuid_bytes()).or_insert(i);
2083 }
2084 let mut touched = 0u64;
2085 for (uuid, new_props) in updates {
2086 if new_props.is_empty() {
2087 continue;
2088 }
2089 touched += 1;
2090 if let Some(&i) = index.get(uuid) {
2091 rows[i].props_mut().extend(new_props.clone());
2092 } else {
2093 index.insert(*uuid, rows.len());
2094 rows.push(R::from_parts(*uuid, new_props.clone()));
2095 }
2096 }
2097 (rows, touched)
2098}
2099
2100fn apply_property_removals<R: PropRowLike>(
2106 mut rows: Vec<R>,
2107 removals: &HashMap<[u8; 16], HashSet<String>>,
2108) -> (Vec<R>, u64) {
2109 let mut index: HashMap<[u8; 16], usize> = HashMap::with_capacity(rows.len());
2110 for (i, row) in rows.iter().enumerate() {
2111 index.entry(*row.uuid_bytes()).or_insert(i);
2112 }
2113 let mut touched = 0u64;
2114 for (uuid, keys) in removals {
2115 if keys.is_empty() {
2116 continue;
2117 }
2118 touched += 1;
2119 if let Some(&i) = index.get(uuid) {
2120 let props = rows[i].props_mut();
2121 for k in keys {
2122 props.remove(k);
2123 }
2124 }
2125 }
2126 (rows, touched)
2127}
2128
2129#[allow(clippy::implicit_hasher)]
2143pub fn stage_set_node_properties(
2144 staged: &mut RewriteBatch,
2145 dir: &Path,
2146 stem: &str,
2147 updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2148) -> Result<u64, GfError> {
2149 let existing = read_props_through(staged, &node_props_path(dir, stem))?;
2150 let rows = decode_property_rows(&existing)?;
2151 let (rows, touched) = apply_property_updates(rows, updates);
2152 stage_node_property_file(staged, dir, stem, &rows)?;
2153 Ok(touched)
2154}
2155
2156#[allow(clippy::implicit_hasher)] pub fn stage_remove_node_properties(
2163 staged: &mut RewriteBatch,
2164 dir: &Path,
2165 stem: &str,
2166 removals: &HashMap<[u8; 16], HashSet<String>>,
2167) -> Result<u64, GfError> {
2168 let existing = read_props_through(staged, &node_props_path(dir, stem))?;
2169 let rows = decode_property_rows(&existing)?;
2170 let (rows, touched) = apply_property_removals(rows, removals);
2171 stage_node_property_file(staged, dir, stem, &rows)?;
2172 Ok(touched)
2173}
2174
2175#[allow(clippy::implicit_hasher)] pub fn stage_set_edge_properties(
2184 staged: &mut RewriteBatch,
2185 dir: &Path,
2186 rel_stem: &str,
2187 updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2188) -> Result<u64, GfError> {
2189 let existing = read_props_through(staged, &edge_props_path(dir, rel_stem))?;
2190 let rows = decode_edge_property_rows(&existing)?;
2191 let (rows, touched) = apply_property_updates(rows, updates);
2192 stage_edge_property_file(staged, dir, rel_stem, &rows)?;
2193 Ok(touched)
2194}
2195
2196#[allow(clippy::implicit_hasher)] pub fn stage_remove_edge_properties(
2204 staged: &mut RewriteBatch,
2205 dir: &Path,
2206 rel_stem: &str,
2207 removals: &HashMap<[u8; 16], HashSet<String>>,
2208) -> Result<u64, GfError> {
2209 let existing = read_props_through(staged, &edge_props_path(dir, rel_stem))?;
2210 let rows = decode_edge_property_rows(&existing)?;
2211 let (rows, touched) = apply_property_removals(rows, removals);
2212 stage_edge_property_file(staged, dir, rel_stem, &rows)?;
2213 Ok(touched)
2214}
2215
2216#[allow(clippy::implicit_hasher)] pub fn set_node_properties(
2223 dir: &Path,
2224 stem: &str,
2225 updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2226) -> Result<u64, GfError> {
2227 let mut staged = RewriteBatch::new();
2228 let touched = stage_set_node_properties(&mut staged, dir, stem, updates)?;
2229 crate::generation::commit_topology_aware(staged, dir)?;
2230 Ok(touched)
2231}
2232
2233#[allow(clippy::implicit_hasher)] pub fn remove_node_properties(
2240 dir: &Path,
2241 stem: &str,
2242 removals: &HashMap<[u8; 16], HashSet<String>>,
2243) -> Result<u64, GfError> {
2244 let mut staged = RewriteBatch::new();
2245 let touched = stage_remove_node_properties(&mut staged, dir, stem, removals)?;
2246 crate::generation::commit_topology_aware(staged, dir)?;
2247 Ok(touched)
2248}
2249
2250#[allow(clippy::implicit_hasher)] pub fn set_edge_properties_rewrite(
2257 dir: &Path,
2258 rel_stem: &str,
2259 updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2260) -> Result<u64, GfError> {
2261 let mut staged = RewriteBatch::new();
2262 let touched = stage_set_edge_properties(&mut staged, dir, rel_stem, updates)?;
2263 staged.commit()?;
2264 Ok(touched)
2265}
2266
2267#[allow(clippy::implicit_hasher)] pub fn remove_edge_properties(
2274 dir: &Path,
2275 rel_stem: &str,
2276 removals: &HashMap<[u8; 16], HashSet<String>>,
2277) -> Result<u64, GfError> {
2278 let mut staged = RewriteBatch::new();
2279 let touched = stage_remove_edge_properties(&mut staged, dir, rel_stem, removals)?;
2280 staged.commit()?;
2281 Ok(touched)
2282}
2283
2284fn stage_node_property_file(
2290 staged: &mut RewriteBatch,
2291 dir: &Path,
2292 stem: &str,
2293 rows: &[PropRow],
2294) -> Result<(), GfError> {
2295 if rows.is_empty() {
2296 return Ok(());
2297 }
2298 let (schema, cols) = build_property_columns(stem, rows)?;
2299 stage_property_file(staged, dir, "properties", stem, schema, cols)
2300}
2301
2302fn stage_edge_property_file(
2305 staged: &mut RewriteBatch,
2306 dir: &Path,
2307 stem: &str,
2308 rows: &[EdgePropRow],
2309) -> Result<(), GfError> {
2310 if rows.is_empty() {
2311 return Ok(());
2312 }
2313 let (schema, cols) =
2314 build_property_columns_keyed(EDGE_PROPERTY_UUID_FIELD, "graphforge.rel_type", stem, rows)?;
2315 stage_property_file(staged, dir, "edge_properties", stem, schema, cols)
2316}
2317
2318fn stage_property_file(
2323 staged: &mut RewriteBatch,
2324 dir: &Path,
2325 subdir: &str,
2326 stem: &str,
2327 schema: Schema,
2328 cols: Vec<ArrayRef>,
2329) -> Result<(), GfError> {
2330 let batch = RecordBatch::try_new(Arc::new(schema), cols).map_err(pq_err)?;
2331 staged.restage(
2332 &dir.join(subdir).join(format!("{stem}.parquet")),
2333 batch.schema(),
2334 &batch,
2335 )
2336}
2337
2338fn node_props_path(dir: &Path, stem: &str) -> PathBuf {
2340 dir.join("properties").join(format!("{stem}.parquet"))
2341}
2342
2343fn edge_props_path(dir: &Path, stem: &str) -> PathBuf {
2345 dir.join("edge_properties").join(format!("{stem}.parquet"))
2346}
2347
2348fn read_props_through(staged: &RewriteBatch, path: &Path) -> Result<Vec<RecordBatch>, GfError> {
2352 let read_path = staged
2353 .staged_temp(path)
2354 .map_or_else(|| path.to_path_buf(), Path::to_path_buf);
2355 match crate::catalog::discover_parquet_schema(&read_path) {
2356 Some(schema) => crate::catalog::read_parquet_or_empty(&read_path, schema).map_err(pq_err),
2357 None => Ok(Vec::new()),
2358 }
2359}
2360
2361use crate::staging::RewriteBatch;
2366
2367fn concat_with_existing(
2372 schema: &SchemaRef,
2373 existing: Vec<RecordBatch>,
2374 new: RecordBatch,
2375) -> Result<RecordBatch, GfError> {
2376 let mut all = existing;
2377 all.push(new);
2378 arrow::compute::concat_batches(schema, &all).map_err(pq_err)
2379}
2380
2381#[cfg(test)]
2386mod tests {
2387 use std::fs::File;
2388
2389 use super::*;
2390 use graphforge_core::uuid::new_v7;
2391 use tempfile::TempDir;
2392
2393 const TS: i64 = 1_700_000_000_000_000;
2394
2395 #[test]
2396 fn create_node_persists_complete_label_set_and_primary_label() {
2397 let dir = TempDir::new().unwrap();
2398 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2399 w.create_node_with_labels(new_v7(), &[TypeId(4), TypeId(9)])
2400 .unwrap();
2401 w.flush().unwrap();
2402
2403 let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
2404 let batch = &nodes[0];
2405 let primary = batch
2406 .column_by_name("type_id")
2407 .unwrap()
2408 .as_any()
2409 .downcast_ref::<UInt32Array>()
2410 .unwrap();
2411 assert_eq!(primary.value(0), 4);
2412 let sets = batch
2413 .column_by_name("type_ids")
2414 .unwrap()
2415 .as_any()
2416 .downcast_ref::<arrow::array::ListArray>()
2417 .unwrap();
2418 let labels = sets.value(0);
2419 let labels = labels.as_any().downcast_ref::<UInt32Array>().unwrap();
2420 assert_eq!(labels.values(), &[4, 9]);
2421 }
2422
2423 #[test]
2424 fn surrogate_ids_are_monotonic_from_one() {
2425 let dir = TempDir::new().unwrap();
2426 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2427 assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 1);
2428 assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 2);
2429 assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 3);
2430 }
2431
2432 #[test]
2433 fn create_edge_with_unknown_endpoint_errors() {
2434 let dir = TempDir::new().unwrap();
2435 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2436 let a = new_v7();
2437 w.create_node(a, TypeId(0)).unwrap();
2438 let b = new_v7();
2440 let e = w.create_edge(new_v7(), "KNOWS", &a, &b);
2441 assert!(matches!(e, Err(GfError::Storage(_))), "got {e:?}");
2442
2443 let unknown_source = new_v7();
2444 let source_error = w.create_edge(new_v7(), "KNOWS", &unknown_source, &a);
2445 assert!(matches!(&source_error, Err(GfError::Storage(_))));
2446 assert!(source_error.unwrap_err().to_string().contains("source"));
2447 }
2448
2449 #[test]
2450 fn register_existing_node_resolves_edge_endpoint() {
2451 let dir = TempDir::new().unwrap();
2455 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2456 let a = new_v7();
2458 w.register_existing_node(a, 42);
2459 let b = new_v7();
2461 assert_eq!(w.create_node(b, TypeId(0)).unwrap(), 1);
2462 let edge_id = w
2464 .create_edge(new_v7(), "KNOWS", &a, &b)
2465 .expect("edge with a registered endpoint resolves");
2466 assert_eq!(edge_id, 1);
2467 }
2468
2469 #[test]
2470 fn register_existing_node_does_not_write_a_node_row() {
2471 let dir = TempDir::new().unwrap();
2474 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2475 w.register_existing_node(new_v7(), 7);
2476 let created = new_v7();
2477 w.create_node(created, TypeId(0)).unwrap();
2478 w.flush().unwrap();
2479 let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
2481 let rows: usize = nodes.iter().map(arrow::array::RecordBatch::num_rows).sum();
2482 assert_eq!(
2483 rows, 1,
2484 "register_existing_node must not persist a node row"
2485 );
2486 }
2487
2488 #[test]
2489 fn empty_flush_creates_no_directories() {
2490 let dir = TempDir::new().unwrap();
2491 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2492 w.flush().unwrap();
2493 assert!(!dir.path().join("topology").exists());
2494 assert!(!dir.path().join("properties").exists());
2495 }
2496
2497 #[test]
2498 fn null_first_property_column_infers_later_concrete_type() {
2499 use arrow::array::Array;
2500 use arrow::datatypes::DataType;
2501 use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
2502
2503 let dir = TempDir::new().unwrap();
2504 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2505 let a = new_v7();
2506 let b = new_v7();
2507 let c = new_v7();
2508 w.create_node(a, TypeId(0)).unwrap();
2509 w.create_node(b, TypeId(0)).unwrap();
2510 w.create_node(c, TypeId(0)).unwrap();
2511 w.set_properties(
2513 &a,
2514 None,
2515 HashMap::from([("score".to_owned(), IrLiteral::Null)]),
2516 )
2517 .unwrap();
2518 w.set_properties(
2519 &b,
2520 None,
2521 HashMap::from([("score".to_owned(), IrLiteral::Int(10))]),
2522 )
2523 .unwrap();
2524 w.set_properties(
2525 &c,
2526 None,
2527 HashMap::from([("score".to_owned(), IrLiteral::Int(20))]),
2528 )
2529 .unwrap();
2530 w.flush().unwrap();
2531
2532 let path = dir.path().join("properties").join("_untyped.parquet");
2533 let file = File::open(&path).unwrap();
2534 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2535 let schema = builder.schema().clone();
2536 let score_fields: Vec<_> = schema
2538 .fields()
2539 .iter()
2540 .filter(|f| f.name() == "score")
2541 .collect();
2542 assert_eq!(score_fields.len(), 1, "expected a single score column");
2543 assert_eq!(
2544 score_fields[0].data_type(),
2545 &DataType::Int64,
2546 "null-first then Int should infer Int64, not Utf8"
2547 );
2548
2549 let mut reader = builder.build().unwrap();
2551 let batch = reader.next().unwrap().unwrap();
2552 assert_eq!(batch.num_rows(), 3);
2553 let scores = batch
2554 .column(schema.index_of("score").unwrap())
2555 .as_any()
2556 .downcast_ref::<arrow::array::Int64Array>()
2557 .unwrap();
2558 assert!(scores.is_null(0));
2559 assert_eq!(scores.value(1), 10);
2560 assert_eq!(scores.value(2), 20);
2561 }
2562
2563 #[test]
2564 fn mixed_type_property_column_uses_tagged_scalars() {
2565 use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
2566
2567 let dir = TempDir::new().unwrap();
2568 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2569 let a = new_v7();
2570 let b = new_v7();
2571 w.create_node(a, TypeId(0)).unwrap();
2572 w.create_node(b, TypeId(0)).unwrap();
2573 w.set_properties(
2574 &a,
2575 None,
2576 HashMap::from([("x".to_owned(), IrLiteral::Int(1))]),
2577 )
2578 .unwrap();
2579 w.set_properties(
2580 &b,
2581 None,
2582 HashMap::from([("x".to_owned(), IrLiteral::Str("two".to_owned()))]),
2583 )
2584 .unwrap();
2585 w.flush().unwrap();
2586
2587 let path = dir.path().join("properties").join("_untyped.parquet");
2588 let file = File::open(&path).unwrap();
2589 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2590 let schema = builder.schema().clone();
2591 let x = schema.field_with_name("x").unwrap();
2592 assert_eq!(
2593 x.data_type(),
2594 &DataType::Struct(heterogeneous_scalar_fields())
2595 );
2596 }
2597
2598 #[test]
2599 fn every_heterogeneous_scalar_tag_round_trips_exactly() {
2600 let dir = TempDir::new().unwrap();
2601 let cases = [
2602 IrLiteral::Int(-1),
2603 IrLiteral::Float(2.25),
2604 IrLiteral::Str("three".into()),
2605 IrLiteral::Bool(true),
2606 ];
2607 let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2608 let mut expected = HashMap::new();
2609 for value in cases {
2610 let node = new_v7();
2611 writer.create_node(node, TypeId(0)).unwrap();
2612 writer
2613 .set_properties(
2614 &node,
2615 None,
2616 HashMap::from([("mixed".into(), value.clone())]),
2617 )
2618 .unwrap();
2619 expected.insert(to_bytes(&node), value);
2620 }
2621 writer.flush().unwrap();
2622
2623 let reopened = read_node_props(dir.path(), "_untyped");
2624 assert_eq!(reopened.len(), expected.len());
2625 for (node, value) in expected {
2626 assert_eq!(reopened[&node].get("mixed"), Some(&value));
2627 }
2628 }
2629
2630 fn read_node_props(dir: &Path, stem: &str) -> HashMap<[u8; 16], HashMap<String, IrLiteral>> {
2637 read_node_property_rows(dir, stem).unwrap()
2638 }
2639
2640 fn read_edge_props(dir: &Path, stem: &str) -> HashMap<[u8; 16], HashMap<String, IrLiteral>> {
2641 let batches = crate::catalog::read_edge_properties(dir, stem).unwrap();
2642 let mut out = HashMap::new();
2643 for row in decode_edge_property_rows(&batches).unwrap() {
2644 out.insert(row.edge_uuid, row.props);
2645 }
2646 out
2647 }
2648
2649 #[test]
2650 fn set_node_properties_sets_new_and_overwrites_existing() {
2651 let dir = TempDir::new().unwrap();
2652 assert!(
2653 read_node_property_rows(dir.path(), "_untyped")
2654 .unwrap()
2655 .is_empty()
2656 );
2657 fs::create_dir_all(dir.path().join("properties")).unwrap();
2658 fs::write(dir.path().join("properties/_untyped.parquet"), b"invalid").unwrap();
2659 assert!(read_node_property_rows(dir.path(), "_untyped").is_err());
2660 fs::remove_file(dir.path().join("properties/_untyped.parquet")).unwrap();
2661 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2662 let a = new_v7();
2663 w.create_node(a, TypeId(0)).unwrap();
2664 w.set_properties(
2665 &a,
2666 None,
2667 HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2668 )
2669 .unwrap();
2670 w.flush().unwrap();
2671
2672 let ab = to_bytes(&a);
2673 let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2674 let updates = HashMap::from([(
2676 ab,
2677 HashMap::from([
2678 ("age".to_owned(), IrLiteral::Int(31)),
2679 ("name".to_owned(), IrLiteral::Str("Al".to_owned())),
2680 ]),
2681 )]);
2682 let touched = set_node_properties(dir.path(), "_untyped", &updates).unwrap();
2683 assert_eq!(touched, 1);
2684 assert_eq!(
2685 crate::generation::read_search_generation(dir.path()).unwrap(),
2686 search_generation + 1
2687 );
2688
2689 let props = read_node_props(dir.path(), "_untyped");
2690 assert_eq!(props[&ab]["age"], IrLiteral::Int(31));
2691 assert_eq!(props[&ab]["name"], IrLiteral::Str("Al".to_owned()));
2692 }
2693
2694 #[test]
2695 fn set_node_properties_inserts_row_for_propertyless_node() {
2696 let dir = TempDir::new().unwrap();
2698 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2699 let a = new_v7();
2700 w.create_node(a, TypeId(0)).unwrap();
2701 w.flush().unwrap(); let ab = to_bytes(&a);
2704 let updates =
2705 HashMap::from([(ab, HashMap::from([("age".to_owned(), IrLiteral::Int(42))]))]);
2706 let touched = set_node_properties(dir.path(), "_untyped", &updates).unwrap();
2707 assert_eq!(touched, 1);
2708
2709 let props = read_node_props(dir.path(), "_untyped");
2710 assert_eq!(props[&ab]["age"], IrLiteral::Int(42));
2711 }
2712
2713 #[test]
2714 fn set_node_properties_routes_by_stem_in_strict_mode() {
2715 let dir = TempDir::new().unwrap();
2717 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2718 let a = new_v7();
2719 w.create_node(a, TypeId(1)).unwrap();
2720 w.set_properties(
2721 &a,
2722 Some("Person"),
2723 HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2724 )
2725 .unwrap();
2726 w.flush().unwrap();
2727
2728 let ab = to_bytes(&a);
2729 let updates =
2730 HashMap::from([(ab, HashMap::from([("age".to_owned(), IrLiteral::Int(99))]))]);
2731 set_node_properties(dir.path(), "Person", &updates).unwrap();
2732
2733 assert!(
2734 dir.path()
2735 .join("properties")
2736 .join("Person.parquet")
2737 .exists()
2738 );
2739 let props = read_node_props(dir.path(), "Person");
2740 assert_eq!(props[&ab]["age"], IrLiteral::Int(99));
2741 }
2742
2743 #[test]
2744 fn remove_node_properties_drops_key_and_column_when_last() {
2745 let dir = TempDir::new().unwrap();
2746 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2747 let a = new_v7();
2748 w.create_node(a, TypeId(0)).unwrap();
2749 w.set_properties(
2750 &a,
2751 None,
2752 HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2753 )
2754 .unwrap();
2755 w.flush().unwrap();
2756
2757 let ab = to_bytes(&a);
2758 let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2759 let removals = HashMap::from([(ab, HashSet::from(["age".to_owned()]))]);
2760 let touched = remove_node_properties(dir.path(), "_untyped", &removals).unwrap();
2761 assert_eq!(touched, 1);
2762 assert_eq!(
2763 crate::generation::read_search_generation(dir.path()).unwrap(),
2764 search_generation + 1
2765 );
2766
2767 let props = read_node_props(dir.path(), "_untyped");
2770 assert!(props.get(&ab).map_or(true, HashMap::is_empty));
2771 }
2772
2773 #[test]
2774 fn remove_node_properties_missing_key_and_uuid_are_noops() {
2775 let dir = TempDir::new().unwrap();
2776 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2777 let a = new_v7();
2778 w.create_node(a, TypeId(0)).unwrap();
2779 w.set_properties(
2780 &a,
2781 None,
2782 HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2783 )
2784 .unwrap();
2785 w.flush().unwrap();
2786
2787 let ab = to_bytes(&a);
2788 let removals = HashMap::from([
2790 (ab, HashSet::from(["nope".to_owned()])),
2791 (to_bytes(&new_v7()), HashSet::from(["age".to_owned()])),
2792 ]);
2793 remove_node_properties(dir.path(), "_untyped", &removals).unwrap();
2794
2795 let props = read_node_props(dir.path(), "_untyped");
2797 assert_eq!(props[&ab]["age"], IrLiteral::Int(30));
2798 }
2799
2800 #[test]
2801 fn set_and_remove_edge_properties_round_trip() {
2802 let dir = TempDir::new().unwrap();
2803 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2804 let a = new_v7();
2805 let b = new_v7();
2806 w.create_node(a, TypeId(0)).unwrap();
2807 w.create_node(b, TypeId(0)).unwrap();
2808 let e = new_v7();
2809 w.create_edge(e, "KNOWS", &a, &b).unwrap();
2810 w.set_edge_properties(
2811 &e,
2812 Some("KNOWS"),
2813 HashMap::from([("since".to_owned(), IrLiteral::Int(2019))]),
2814 )
2815 .unwrap();
2816 w.flush().unwrap();
2817
2818 let eb = to_bytes(&e);
2819 let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2820 let updates = HashMap::from([(
2822 eb,
2823 HashMap::from([("since".to_owned(), IrLiteral::Int(2020))]),
2824 )]);
2825 assert_eq!(
2826 set_edge_properties_rewrite(dir.path(), "KNOWS", &updates).unwrap(),
2827 1
2828 );
2829 assert_eq!(
2830 read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
2831 IrLiteral::Int(2020)
2832 );
2833
2834 let removals = HashMap::from([(eb, HashSet::from(["since".to_owned()]))]);
2836 assert_eq!(
2837 remove_edge_properties(dir.path(), "KNOWS", &removals).unwrap(),
2838 1
2839 );
2840 let props = read_edge_props(dir.path(), "KNOWS");
2841 assert!(props.get(&eb).map_or(true, HashMap::is_empty));
2842 assert_eq!(
2843 crate::generation::read_search_generation(dir.path()).unwrap(),
2844 search_generation
2845 );
2846 }
2847
2848 #[test]
2849 fn set_node_properties_empty_map_writes_nothing() {
2850 let dir = TempDir::new().unwrap();
2851 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2852 w.create_node(new_v7(), TypeId(0)).unwrap();
2853 w.flush().unwrap();
2854
2855 let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2856 let touched = set_node_properties(dir.path(), "_untyped", &HashMap::new()).unwrap();
2857 assert_eq!(touched, 0);
2858 assert_eq!(
2859 crate::generation::read_search_generation(dir.path()).unwrap(),
2860 search_generation
2861 );
2862 assert!(
2864 !dir.path()
2865 .join("properties")
2866 .join("_untyped.parquet")
2867 .exists()
2868 );
2869 }
2870
2871 #[test]
2872 fn staged_set_is_invisible_until_commit_across_stems() {
2873 let dir = TempDir::new().unwrap();
2876 let (a, e) = (new_v7(), new_v7());
2877 let b = new_v7();
2878 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2879 w.create_node(a, TypeId(0)).unwrap();
2880 w.create_node(b, TypeId(0)).unwrap();
2881 w.create_edge(e, "KNOWS", &a, &b).unwrap();
2882 w.set_properties(
2883 &a,
2884 None,
2885 HashMap::from([("name".to_owned(), IrLiteral::Str("old".into()))]),
2886 )
2887 .unwrap();
2888 w.set_edge_properties(
2889 &e,
2890 Some("KNOWS"),
2891 HashMap::from([("since".to_owned(), IrLiteral::Int(2000))]),
2892 )
2893 .unwrap();
2894 w.flush().unwrap();
2895
2896 let (ab, eb) = (to_bytes(&a), to_bytes(&e));
2897 let node_updates = HashMap::from([(
2898 ab,
2899 HashMap::from([("name".to_owned(), IrLiteral::Str("new".into()))]),
2900 )]);
2901 let edge_updates = HashMap::from([(
2902 eb,
2903 HashMap::from([("since".to_owned(), IrLiteral::Int(2024))]),
2904 )]);
2905
2906 let mut staged = RewriteBatch::new();
2907 let touched = stage_set_node_properties(&mut staged, dir.path(), "_untyped", &node_updates)
2908 .unwrap()
2909 + stage_set_edge_properties(&mut staged, dir.path(), "KNOWS", &edge_updates).unwrap();
2910 assert_eq!(touched, 2, "one node + one edge written");
2911
2912 assert_eq!(
2914 read_node_props(dir.path(), "_untyped")[&ab]["name"],
2915 IrLiteral::Str("old".into())
2916 );
2917 assert_eq!(
2918 read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
2919 IrLiteral::Int(2000)
2920 );
2921
2922 staged.commit().unwrap();
2923 assert_eq!(
2924 read_node_props(dir.path(), "_untyped")[&ab]["name"],
2925 IrLiteral::Str("new".into())
2926 );
2927 assert_eq!(
2928 read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
2929 IrLiteral::Int(2024)
2930 );
2931 }
2932
2933 #[test]
2938 fn pending_nodes_batch_is_canonical_and_non_consuming() {
2939 let dir = TempDir::new().unwrap();
2940 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2941 let empty = w.pending_nodes_batch().unwrap();
2942 assert_eq!(empty.schema(), TOPOLOGY_NODES_SCHEMA.clone());
2943 assert_eq!(empty.num_rows(), 0);
2944
2945 let node = new_v7();
2946 w.create_node_with_labels(node, &[TypeId(3), TypeId(7)])
2947 .unwrap();
2948 for _ in 0..2 {
2949 let batch = w.pending_nodes_batch().unwrap();
2950 assert_eq!(batch.schema(), TOPOLOGY_NODES_SCHEMA.clone());
2951 assert_eq!(batch.num_rows(), 1);
2952 assert_eq!(
2953 batch
2954 .column_by_name("node_id")
2955 .unwrap()
2956 .as_any()
2957 .downcast_ref::<UInt64Array>()
2958 .unwrap()
2959 .value(0),
2960 1
2961 );
2962 }
2963 assert!(w.contains_pending_node(&to_bytes(&node)));
2964 }
2965
2966 #[test]
2967 fn cancel_nodes_drops_rows_props_and_mapping() {
2968 let dir = TempDir::new().unwrap();
2969 let (a, b) = (new_v7(), new_v7());
2970 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2971 w.create_node(a, TypeId(0)).unwrap();
2972 w.create_node(b, TypeId(0)).unwrap();
2973 w.set_properties(
2974 &a,
2975 None,
2976 HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
2977 )
2978 .unwrap();
2979
2980 assert!(w.contains_pending_node(&to_bytes(&a)));
2981 let dropped = w.cancel_nodes(&HashSet::from([to_bytes(&a)]));
2982 assert_eq!(dropped, 1);
2983 assert!(!w.contains_pending_node(&to_bytes(&a)));
2984
2985 let err = w.create_edge(new_v7(), "KNOWS", &b, &a);
2987 assert!(err.is_err(), "edge to a cancelled node must fail");
2988
2989 w.flush().unwrap();
2990 let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
2992 let total: usize = nodes.iter().map(RecordBatch::num_rows).sum();
2993 assert_eq!(total, 1, "only b persisted");
2994 assert!(
2995 !read_node_props(dir.path(), "_untyped").contains_key(&to_bytes(&a)),
2996 "cancelled node's props never hit disk"
2997 );
2998 }
2999
3000 #[test]
3001 fn cancel_edges_drops_rows_and_edge_props() {
3002 let dir = TempDir::new().unwrap();
3003 let (a, b) = (new_v7(), new_v7());
3004 let (e1, e2) = (new_v7(), new_v7());
3005 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3006 w.create_node(a, TypeId(0)).unwrap();
3007 w.create_node(b, TypeId(0)).unwrap();
3008 w.create_edge(e1, "KNOWS", &a, &b).unwrap();
3009 w.create_edge(e2, "KNOWS", &b, &a).unwrap();
3010 w.set_edge_properties(
3011 &e1,
3012 Some("KNOWS"),
3013 HashMap::from([("since".to_owned(), IrLiteral::Int(2020))]),
3014 )
3015 .unwrap();
3016
3017 assert!(w.contains_pending_edge(&to_bytes(&e1)));
3018 assert_eq!(w.cancel_edges(&HashSet::from([to_bytes(&e1)])), 1);
3019 assert!(!w.contains_pending_edge(&to_bytes(&e1)));
3020
3021 w.flush().unwrap();
3022 assert!(
3023 !read_edge_props(dir.path(), "KNOWS").contains_key(&to_bytes(&e1)),
3024 "cancelled edge's props never hit disk"
3025 );
3026 let edges = crate::catalog::read_parquet_or_empty(
3027 &dir.path().join("topology/edges/_exploratory.parquet"),
3028 EXPLORATORY_EDGE_SCHEMA.clone(),
3029 )
3030 .unwrap();
3031 let total: usize = edges.iter().map(RecordBatch::num_rows).sum();
3032 assert_eq!(total, 1, "only e2 persisted");
3033 }
3034
3035 #[test]
3036 fn pending_incident_edge_uuids_sees_buffered_edges() {
3037 let dir = TempDir::new().unwrap();
3038 let (a, b, c) = (new_v7(), new_v7(), new_v7());
3039 let e_ab = new_v7();
3040 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3041 w.create_node(a, TypeId(0)).unwrap();
3042 w.create_node(b, TypeId(0)).unwrap();
3043 w.create_node(c, TypeId(0)).unwrap();
3044 w.create_edge(e_ab, "KNOWS", &a, &b).unwrap();
3045
3046 let hits = w.pending_incident_edge_uuids(&HashSet::from([to_bytes(&b)]));
3048 assert_eq!(hits, vec![to_bytes(&e_ab)]);
3049 assert!(
3050 w.pending_incident_edge_uuids(&HashSet::from([to_bytes(&c)]))
3051 .is_empty()
3052 );
3053 }
3054
3055 #[test]
3056 fn pending_query_and_label_edits_are_exact_before_flush_and_reopen() {
3057 let dir = TempDir::new().unwrap();
3058 let (alice, bob, edge) = (new_v7(), new_v7(), new_v7());
3059 let (alice_bytes, bob_bytes, edge_bytes) =
3060 (to_bytes(&alice), to_bytes(&bob), to_bytes(&edge));
3061 let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
3062 writer
3063 .create_node_with_labels(alice, &[TypeId(3), TypeId(7)])
3064 .unwrap();
3065 writer.create_node(bob, TypeId(3)).unwrap();
3066 writer
3067 .set_properties(
3068 &alice,
3069 Some("Person"),
3070 HashMap::from([
3071 ("name".into(), IrLiteral::Str("Alice".into())),
3072 ("age".into(), IrLiteral::Int(42)),
3073 ]),
3074 )
3075 .unwrap();
3076
3077 assert_eq!(
3078 writer.pending_node_labels(&HashSet::from([alice_bytes, bob_bytes])),
3079 HashSet::from([3, 7])
3080 );
3081 let matched = writer
3082 .find_pending_node(&[3, 7], &[("name".into(), IrLiteral::Str("Alice".into()))])
3083 .unwrap();
3084 assert_eq!(matched.0, alice_bytes);
3085 assert_eq!(matched.2, 3);
3086 assert_eq!(matched.3, vec![3, 7]);
3087 assert_eq!(matched.4["age"], IrLiteral::Int(42));
3088 assert!(writer.find_pending_node(&[9], &[]).is_none());
3089 assert!(
3090 writer
3091 .find_pending_node(&[3], &[("name".into(), IrLiteral::Str("Bob".into()))])
3092 .is_none()
3093 );
3094
3095 assert_eq!(writer.add_pending_node_labels(&alice_bytes, &[7, 9]), 1);
3096 assert_eq!(writer.add_pending_node_labels(&[0xff; 16], &[1]), 0);
3097 assert_eq!(writer.remove_pending_node_labels(&alice_bytes, &[7, 99]), 1);
3098 assert_eq!(writer.remove_pending_node_labels(&[0xff; 16], &[1]), 0);
3099 assert_eq!(
3100 writer.pending_node_labels(&HashSet::from([alice_bytes])),
3101 HashSet::from([3, 9])
3102 );
3103
3104 writer.create_edge(edge, "KNOWS", &alice, &bob).unwrap();
3105 writer
3106 .set_edge_properties(
3107 &edge,
3108 Some("KNOWS"),
3109 HashMap::from([("since".into(), IrLiteral::Int(2024))]),
3110 )
3111 .unwrap();
3112 let direct = writer
3113 .find_pending_edge(
3114 "KNOWS",
3115 &alice_bytes,
3116 &bob_bytes,
3117 false,
3118 &[("since".into(), IrLiteral::Int(2024))],
3119 )
3120 .unwrap();
3121 assert_eq!(direct.0, edge_bytes);
3122 assert_eq!(direct.1, alice_bytes);
3123 assert_eq!(direct.2, bob_bytes);
3124 assert_eq!(direct.3["since"], IrLiteral::Int(2024));
3125 assert!(
3126 writer
3127 .find_pending_edge("KNOWS", &bob_bytes, &alice_bytes, false, &[])
3128 .is_none()
3129 );
3130 assert!(
3131 writer
3132 .find_pending_edge("KNOWS", &bob_bytes, &alice_bytes, true, &[])
3133 .is_some()
3134 );
3135 assert!(
3136 writer
3137 .find_pending_edge("IGNORES", &alice_bytes, &bob_bytes, false, &[])
3138 .is_none()
3139 );
3140
3141 writer.flush().unwrap();
3142 let mut reopened = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS + 1).unwrap();
3143 assert_eq!(reopened.create_node(new_v7(), TypeId(3)).unwrap(), 3);
3144 assert_eq!(
3145 read_node_props(dir.path(), "Person")[&alice_bytes]["name"],
3146 IrLiteral::Str("Alice".into())
3147 );
3148 assert_eq!(
3149 read_edge_props(dir.path(), "KNOWS")[&edge_bytes]["since"],
3150 IrLiteral::Int(2024)
3151 );
3152 }
3153
3154 #[test]
3155 fn merge_and_remove_pending_props_edit_buffered_rows() {
3156 let dir = TempDir::new().unwrap();
3157 let a = new_v7();
3158 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3159 w.create_node(a, TypeId(0)).unwrap();
3160 w.set_properties(
3161 &a,
3162 None,
3163 HashMap::from([
3164 ("name".to_owned(), IrLiteral::Str("old".into())),
3165 ("age".to_owned(), IrLiteral::Int(30)),
3166 ]),
3167 )
3168 .unwrap();
3169
3170 w.merge_pending_node_props(
3172 &to_bytes(&a),
3173 None,
3174 HashMap::from([
3175 ("name".to_owned(), IrLiteral::Str("new".into())),
3176 ("city".to_owned(), IrLiteral::Str("Oslo".into())),
3177 ]),
3178 );
3179 w.remove_pending_node_props(
3181 &to_bytes(&a),
3182 &HashSet::from(["age".to_owned(), "absent".to_owned()]),
3183 );
3184
3185 w.flush().unwrap();
3186 let props = &read_node_props(dir.path(), "_untyped")[&to_bytes(&a)];
3187 assert_eq!(props["name"], IrLiteral::Str("new".into()));
3188 assert_eq!(props["city"], IrLiteral::Str("Oslo".into()));
3189 assert!(!props.contains_key("age"), "removed before flush");
3190 }
3191
3192 #[test]
3193 fn flush_into_composes_with_staged_delete_in_one_batch() {
3194 let dir = TempDir::new().unwrap();
3198 let (a, b) = (new_v7(), new_v7());
3199 let mut seed = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3200 seed.create_node(a, TypeId(0)).unwrap();
3201 seed.create_node(b, TypeId(0)).unwrap();
3202 seed.set_properties(
3203 &a,
3204 None,
3205 HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
3206 )
3207 .unwrap();
3208 seed.flush().unwrap();
3209
3210 let mut staged = RewriteBatch::new();
3211 let removed = crate::mutator::stage_delete_nodes(
3212 &mut staged,
3213 dir.path(),
3214 &HashSet::from([to_bytes(&a)]),
3215 )
3216 .unwrap();
3217 assert_eq!(removed, 1);
3218
3219 let d = new_v7();
3220 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3221 w.create_node(d, TypeId(0)).unwrap();
3222 w.flush_into(&mut staged).unwrap();
3223
3224 let staged_nodes = staged
3227 .staged_paths()
3228 .filter(|p| p.ends_with("topology/nodes.parquet"))
3229 .count();
3230 assert_eq!(staged_nodes, 1, "net content, no double-stage");
3231
3232 let pre: usize = crate::catalog::read_nodes(dir.path())
3234 .unwrap()
3235 .iter()
3236 .map(RecordBatch::num_rows)
3237 .sum();
3238 assert_eq!(pre, 2);
3239
3240 staged.commit().unwrap();
3241 let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
3242 let total: usize = nodes.iter().map(RecordBatch::num_rows).sum();
3243 assert_eq!(total, 2, "b survives, a deleted, d created");
3244 assert!(
3245 !read_node_props(dir.path(), "_untyped").contains_key(&to_bytes(&a)),
3246 "deleted node's props gone"
3247 );
3248 }
3249
3250 #[test]
3251 fn every_persisted_property_family_round_trips_through_parquet_reopen() {
3252 let dir = TempDir::new().unwrap();
3253 let node = new_v7();
3254 let propertyless = new_v7();
3255 let values = HashMap::from([
3256 ("int".into(), IrLiteral::Int(-7)),
3257 ("float".into(), IrLiteral::Float(2.5)),
3258 ("bool".into(), IrLiteral::Bool(true)),
3259 ("str".into(), IrLiteral::Str("value".into())),
3260 (
3261 "duration".into(),
3262 IrLiteral::Duration {
3263 months: 1,
3264 days: -2,
3265 seconds: 3,
3266 nanos: 4,
3267 },
3268 ),
3269 ("datetime".into(), IrLiteral::DateTime(TS)),
3270 ("date".into(), IrLiteral::Date(19_000)),
3271 (
3272 "local_datetime".into(),
3273 IrLiteral::LocalDateTime {
3274 days: 19_001,
3275 nanos: 123,
3276 },
3277 ),
3278 ("time".into(), IrLiteral::Time(456)),
3279 (
3280 "zoned_time".into(),
3281 IrLiteral::ZonedTime {
3282 nanos: 789,
3283 offset: -21_600,
3284 },
3285 ),
3286 (
3287 "zoned_datetime".into(),
3288 IrLiteral::ZonedDateTime {
3289 days: 19_002,
3290 nanos: 987,
3291 offset: 3_600,
3292 zone: Some("Europe/Paris".into()),
3293 },
3294 ),
3295 (
3296 "offset_datetime".into(),
3297 IrLiteral::ZonedDateTime {
3298 days: 19_003,
3299 nanos: 654,
3300 offset: 0,
3301 zone: None,
3302 },
3303 ),
3304 (
3305 "ints".into(),
3306 IrLiteral::List(vec![IrLiteral::Int(1), IrLiteral::Null, IrLiteral::Int(3)]),
3307 ),
3308 (
3309 "dates".into(),
3310 IrLiteral::List(vec![IrLiteral::Date(19_004), IrLiteral::Date(19_005)]),
3311 ),
3312 ("empty".into(), IrLiteral::List(Vec::new())),
3313 ("null".into(), IrLiteral::Null),
3314 ]);
3315 let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3316 writer.create_node(node, TypeId(0)).unwrap();
3317 writer.create_node(propertyless, TypeId(0)).unwrap();
3318 writer.set_properties(&node, None, values.clone()).unwrap();
3319 writer
3320 .set_properties(
3321 &propertyless,
3322 None,
3323 values
3324 .keys()
3325 .map(|name| (name.clone(), IrLiteral::Null))
3326 .collect(),
3327 )
3328 .unwrap();
3329 writer.flush().unwrap();
3330
3331 let reopened = read_node_props(dir.path(), "_untyped");
3332 let actual = reopened.get(&to_bytes(&node)).unwrap();
3333 for (name, expected) in &values {
3334 if matches!(expected, IrLiteral::Null) {
3335 assert!(!actual.contains_key(name));
3336 } else if name == "empty" {
3337 assert_eq!(actual.get(name), Some(&IrLiteral::Str("[]".into())));
3338 } else {
3339 assert_eq!(actual.get(name), Some(expected), "property {name}");
3340 }
3341 }
3342 assert!(
3343 reopened
3344 .get(&to_bytes(&propertyless))
3345 .is_none_or(HashMap::is_empty)
3346 );
3347 }
3348
3349 #[test]
3350 fn property_literal_rendering_and_nested_invalid_values_are_deterministic() {
3351 let uuid = [0xabu8; 16];
3352 let cases = [
3353 (IrLiteral::Null, "".into()),
3354 (IrLiteral::Bool(true), "true".into()),
3355 (IrLiteral::Int(-2), "-2".into()),
3356 (IrLiteral::Float(1.25), "1.25".into()),
3357 (IrLiteral::Str("s".into()), "s".into()),
3358 (IrLiteral::Uuid(uuid), "ab".repeat(16)),
3359 (
3360 IrLiteral::Duration {
3361 months: 1,
3362 days: 2,
3363 seconds: 3,
3364 nanos: 4,
3365 },
3366 "1mo2d3s4ns".into(),
3367 ),
3368 (IrLiteral::DateTime(5), "5".into()),
3369 (IrLiteral::Date(6), "6".into()),
3370 (
3371 IrLiteral::LocalDateTime { days: 7, nanos: 8 },
3372 "7d8ns".into(),
3373 ),
3374 (IrLiteral::Time(9), "9ns".into()),
3375 (
3376 IrLiteral::ZonedTime {
3377 nanos: 10,
3378 offset: -1,
3379 },
3380 "10ns-1s".into(),
3381 ),
3382 (
3383 IrLiteral::ZonedDateTime {
3384 days: 11,
3385 nanos: 12,
3386 offset: 13,
3387 zone: Some("UTC".into()),
3388 },
3389 "11d12ns+13sUTC".into(),
3390 ),
3391 (
3392 IrLiteral::List(vec![IrLiteral::Int(1), IrLiteral::Str("x".into())]),
3393 "[1,x]".into(),
3394 ),
3395 (
3396 IrLiteral::Map(vec![("a".into(), IrLiteral::Bool(false))]),
3397 "{a:false}".into(),
3398 ),
3399 ];
3400 for (literal, expected) in cases {
3401 assert_eq!(literal_to_string(&literal), expected);
3402 }
3403
3404 for invalid in [
3405 IrLiteral::Uuid(uuid),
3406 IrLiteral::List(vec![IrLiteral::Uuid(uuid)]),
3407 IrLiteral::Map(vec![("nested".into(), IrLiteral::Uuid(uuid))]),
3408 ] {
3409 assert_eq!(
3410 reject_map_property_value("p", &invalid).unwrap_err().code(),
3411 "GF_VALIDATION"
3412 );
3413 }
3414 for invalid in [
3415 IrLiteral::Map(vec![]),
3416 IrLiteral::List(vec![IrLiteral::Map(vec![])]),
3417 ] {
3418 assert_eq!(
3419 reject_map_property_value("p", &invalid).unwrap_err().code(),
3420 "GF_IO"
3421 );
3422 }
3423 }
3424
3425 #[test]
3426 fn persisted_property_decoder_rejects_unsupported_shape_type_and_dynamic_array() {
3427 use arrow::array::{Int32Array, Int64Array, StructArray, UInt8Array};
3428 use arrow::datatypes::{DataType, Field, Fields};
3429
3430 let unsupported: arrow::array::ArrayRef = Arc::new(UInt8Array::from(vec![1]));
3431 let unsupported_field = Field::new("unsupported", DataType::UInt8, false);
3432 assert!(decode_value(&unsupported, &unsupported_field, 0).is_err());
3433
3434 let fields: Fields = vec![Field::new("other", DataType::Int32, false)].into();
3435 let structure: arrow::array::ArrayRef = Arc::new(StructArray::new(
3436 fields.clone(),
3437 vec![Arc::new(Int32Array::from(vec![1]))],
3438 None,
3439 ));
3440 let structure_field = Field::new("structure", DataType::Struct(fields), false);
3441 assert!(decode_value(&structure, &structure_field, 0).is_err());
3442
3443 let wrong_dynamic: arrow::array::ArrayRef = Arc::new(Int64Array::from(vec![1]));
3444 let declared = Field::new("declared", DataType::UInt64, false);
3445 assert!(decode_value(&wrong_dynamic, &declared, 0).is_err());
3446 }
3447}