1use std::any::Any;
22use std::collections::HashMap;
23use std::fmt;
24use std::fs::File;
25use std::path::{Path, PathBuf};
26use std::sync::Arc;
27
28use arrow::array::RecordBatch;
29use arrow::compute::concat_batches;
30use arrow::datatypes::{DataType, Field, SchemaRef};
31use async_trait::async_trait;
32use datafusion::catalog::{CatalogProvider, SchemaProvider};
33use datafusion::datasource::{MemTable, TableProvider, TableType};
34use datafusion::error::DataFusionError;
35use datafusion::physical_plan::ExecutionPlan;
36use datafusion::prelude::Expr;
37use datafusion_catalog::Session;
38use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
39
40use graphforge_core::OntologyMode;
41use graphforge_ir::RuntimeCatalog;
42use graphforge_ontology::OntologyHandle;
43
44use crate::schemas::{
45 EXPLORATORY_EDGE_SCHEMA, TOPOLOGY_NODES_SCHEMA, TYPED_EDGE_SCHEMA, property_schema,
46};
47
48fn parquet_err(e: impl std::fmt::Display) -> DataFusionError {
53 DataFusionError::External(e.to_string().into())
54}
55
56fn io_err(e: &std::io::Error) -> DataFusionError {
57 DataFusionError::External(e.to_string().into())
58}
59
60fn total_rows(batches: &[RecordBatch]) -> u64 {
62 u64::try_from(batches.iter().map(RecordBatch::num_rows).sum::<usize>()).unwrap_or(u64::MAX)
63}
64
65pub(crate) fn read_parquet_or_empty(
69 path: &Path,
70 schema: SchemaRef,
71) -> Result<Vec<RecordBatch>, DataFusionError> {
72 if !path.exists() {
73 return Ok(vec![RecordBatch::new_empty(schema)]);
74 }
75 let file = File::open(path).map_err(|e| io_err(&e))?;
76 let builder = ParquetRecordBatchReaderBuilder::try_new(file).map_err(parquet_err)?;
77 let file_schema = builder.schema().clone();
78 let reader = builder.build().map_err(parquet_err)?;
79 let batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().map_err(parquet_err)?;
80 if batches.is_empty() {
81 return Ok(vec![RecordBatch::new_empty(file_schema)]);
82 }
83 let merged = concat_batches(&file_schema, &batches)
84 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?;
85 Ok(vec![merged])
86}
87
88pub(crate) fn normalize_topology_nodes(
91 batches: Vec<RecordBatch>,
92) -> Result<Vec<RecordBatch>, DataFusionError> {
93 use arrow::array::{Array, ListArray, UInt32Array};
94 use arrow::datatypes::UInt32Type;
95
96 batches
97 .into_iter()
98 .map(|batch| {
99 if batch.schema().field_with_name("type_ids").is_ok() {
100 return Ok(batch);
101 }
102 let type_idx = batch.schema().index_of("type_id").map_err(|e| {
103 DataFusionError::Execution(format!("legacy node topology missing type_id: {e}"))
104 })?;
105 let primary_ids = batch
106 .column(type_idx)
107 .as_any()
108 .downcast_ref::<UInt32Array>()
109 .ok_or_else(|| DataFusionError::Execution("type_id is not UInt32".into()))?;
110 let nullable_labels = ListArray::from_iter_primitive::<UInt32Type, _, _>(
111 (0..batch.num_rows()).map(|row| Some([Some(primary_ids.value(row))])),
112 );
113 let labels = ListArray::new(
114 Arc::new(Field::new("item", DataType::UInt32, false)),
115 nullable_labels.offsets().clone(),
116 nullable_labels.values().clone(),
117 None,
118 );
119 let mut columns = batch.columns().to_vec();
120 columns.insert(type_idx + 1, Arc::new(labels));
121 RecordBatch::try_new(TOPOLOGY_NODES_SCHEMA.clone(), columns)
122 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))
123 })
124 .collect()
125}
126
127pub fn read_edges(
157 dir: &Path,
158 rel_name: &str,
159 mode: OntologyMode,
160) -> Result<Vec<RecordBatch>, DataFusionError> {
161 if rel_name == "*" && matches!(mode, OntologyMode::Advisory | OntologyMode::Strict) {
165 return read_edges_union(dir, None, None);
166 }
167 if matches!(mode, OntologyMode::Advisory | OntologyMode::Strict) {
173 let mut comps = Path::new(rel_name).components();
174 let single_normal =
175 matches!(comps.next(), Some(std::path::Component::Normal(_))) && comps.next().is_none();
176 if !single_normal {
177 return Err(DataFusionError::Execution(format!(
178 "invalid relation name {rel_name:?}: must be a plain file stem"
179 )));
180 }
181 }
182 let (stem, schema) = match mode {
183 OntologyMode::Exploratory => ("_exploratory", EXPLORATORY_EDGE_SCHEMA.clone()),
184 OntologyMode::Advisory | OntologyMode::Strict => (rel_name, TYPED_EDGE_SCHEMA.clone()),
185 };
186 let path = dir
187 .join("topology")
188 .join("edges")
189 .join(format!("{stem}.parquet"));
190 let batches = read_parquet_or_empty(&path, schema)?;
191 crate::io_stats::record_edge_full_read(total_rows(&batches));
192 Ok(batches)
193}
194
195#[allow(clippy::implicit_hasher)]
215pub fn read_edges_filtered(
216 dir: &Path,
217 rel_name: &str,
218 mode: OntologyMode,
219 edge_ids: &std::collections::HashSet<u64>,
220) -> Result<Vec<RecordBatch>, DataFusionError> {
221 read_edges_filtered_observed(dir, rel_name, mode, edge_ids, None)
222}
223
224#[allow(clippy::implicit_hasher)]
226#[doc(hidden)]
227pub fn read_edges_filtered_observed(
228 dir: &Path,
229 rel_name: &str,
230 mode: OntologyMode,
231 edge_ids: &std::collections::HashSet<u64>,
232 observer: Option<&std::sync::Arc<dyn crate::io_stats::FilteredReadObserver>>,
233) -> Result<Vec<RecordBatch>, DataFusionError> {
234 if rel_name == "*" && matches!(mode, OntologyMode::Advisory | OntologyMode::Strict) {
236 return read_edges_union(dir, Some(edge_ids), observer);
237 }
238 if matches!(mode, OntologyMode::Advisory | OntologyMode::Strict) {
239 let mut comps = Path::new(rel_name).components();
240 let single_normal =
241 matches!(comps.next(), Some(std::path::Component::Normal(_))) && comps.next().is_none();
242 if !single_normal {
243 return Err(DataFusionError::Execution(format!(
244 "invalid relation name {rel_name:?}: must be a plain file stem"
245 )));
246 }
247 }
248 let (stem, schema) = match mode {
249 OntologyMode::Exploratory => ("_exploratory", EXPLORATORY_EDGE_SCHEMA.clone()),
250 OntologyMode::Advisory | OntologyMode::Strict => (rel_name, TYPED_EDGE_SCHEMA.clone()),
251 };
252 let path = dir
253 .join("topology")
254 .join("edges")
255 .join(format!("{stem}.parquet"));
256 read_parquet_filtered_u64(
257 &path,
258 schema,
259 "edge_id",
260 edge_ids,
261 FilteredReadKind::Edge,
262 observer,
263 )
264}
265
266fn read_edges_union(
275 dir: &Path,
276 edge_ids: Option<&std::collections::HashSet<u64>>,
277 observer: Option<&std::sync::Arc<dyn crate::io_stats::FilteredReadObserver>>,
278) -> Result<Vec<RecordBatch>, DataFusionError> {
279 let mut files = crate::mutator::parquet_files_in(dir, "topology/edges")
280 .map_err(|e| DataFusionError::Execution(e.to_string()))?;
281 files.sort();
282 let mut out = Vec::new();
283 for path in files {
284 let stem = path
285 .file_stem()
286 .and_then(|s| s.to_str())
287 .unwrap_or_default()
288 .to_owned();
289 let schema = discover_parquet_schema(&path).unwrap_or_else(|| TYPED_EDGE_SCHEMA.clone());
290 let batches = if let Some(ids) = edge_ids {
291 read_parquet_filtered_u64(
292 &path,
293 schema,
294 "edge_id",
295 ids,
296 FilteredReadKind::Edge,
297 observer,
298 )?
299 } else {
300 let b = read_parquet_or_empty(&path, schema)?;
301 crate::io_stats::record_edge_full_read(total_rows(&b));
302 b
303 };
304 for batch in &batches {
305 if batch.num_rows() > 0 {
306 out.push(tag_rel_type_name(batch, &stem)?);
307 }
308 }
309 }
310 if out.is_empty() {
311 out.push(RecordBatch::new_empty(EXPLORATORY_EDGE_SCHEMA.clone()));
312 }
313 Ok(out)
314}
315
316fn tag_rel_type_name(batch: &RecordBatch, stem: &str) -> Result<RecordBatch, DataFusionError> {
320 if batch.schema().field_with_name("rel_type_name").is_ok() {
321 return Ok(batch.clone());
322 }
323 let names = arrow::array::StringArray::from(vec![stem; batch.num_rows()]);
324 let mut cols: Vec<arrow::array::ArrayRef> = batch.columns().to_vec();
325 cols.push(Arc::new(names));
326 RecordBatch::try_new(EXPLORATORY_EDGE_SCHEMA.clone(), cols)
327 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))
328}
329
330#[derive(Clone, Copy, PartialEq, Eq)]
335enum FilteredReadKind {
336 Edge,
337 Node,
338}
339
340struct FilteredReadObservation {
343 observer: Option<std::sync::Arc<dyn crate::io_stats::FilteredReadObserver>>,
344 table: crate::io_stats::FilteredReadTable,
345 completed: bool,
346}
347
348impl FilteredReadObservation {
349 fn new(
350 observer: Option<&std::sync::Arc<dyn crate::io_stats::FilteredReadObserver>>,
351 kind: FilteredReadKind,
352 ) -> Self {
353 let table = kind.into();
354 if let Some(observer) = &observer {
355 observer.read_started(table);
356 }
357 Self {
358 observer: observer.cloned(),
359 table,
360 completed: false,
361 }
362 }
363
364 fn scanned(&self, rows: u64) {
365 if let Some(observer) = &self.observer {
366 observer.rows_scanned(self.table, rows);
367 }
368 }
369
370 fn pruning(&self, pruning: crate::io_stats::FilteredReadPruning) {
371 if let Some(observer) = &self.observer {
372 observer.pruning(self.table, pruning);
373 }
374 }
375
376 fn complete(&mut self, rows: u64, full: bool) {
377 if let Some(observer) = &self.observer {
378 observer.read_completed(self.table, rows, full);
379 }
380 self.completed = true;
381 }
382}
383
384impl Drop for FilteredReadObservation {
385 fn drop(&mut self) {
386 if !self.completed
387 && let Some(observer) = &self.observer
388 {
389 observer.read_failed(self.table);
390 }
391 }
392}
393
394impl From<FilteredReadKind> for crate::io_stats::FilteredReadTable {
395 fn from(value: FilteredReadKind) -> Self {
396 match value {
397 FilteredReadKind::Edge => Self::Edge,
398 FilteredReadKind::Node => Self::Node,
399 }
400 }
401}
402
403struct DenseNodeSelection {
407 row_groups: Vec<usize>,
408 selection: parquet::arrow::arrow_reader::RowSelection,
409 pages_considered: u64,
410 pages_selected: u64,
411 exact_rows_selected: u64,
412}
413
414struct DenseNodeLayout {
415 group_rows: Vec<usize>,
416 group_pages: Vec<Vec<usize>>,
417 total_rows: usize,
418 pages_considered: u64,
419}
420
421fn dense_node_layout(
423 metadata: &parquet::file::metadata::ParquetMetaData,
424 key_leaf: usize,
425) -> Option<DenseNodeLayout> {
426 use parquet::basic::BoundaryOrder;
427 use parquet::file::page_index::column_index::ColumnIndexMetaData;
428 use parquet::file::statistics::Statistics;
429
430 let total_rows = usize::try_from(metadata.file_metadata().num_rows()).ok()?;
431 if total_rows == 0 || u64::try_from(total_rows).ok()? > i64::MAX as u64 {
432 return None;
433 }
434 let row_groups = metadata.row_groups();
435 let column_indexes = metadata.column_index()?;
436 let offset_indexes = metadata.offset_index()?;
437 if column_indexes.len() != row_groups.len() || offset_indexes.len() != row_groups.len() {
438 return None;
439 }
440
441 let mut group_rows = Vec::with_capacity(row_groups.len());
442 let mut group_pages = Vec::with_capacity(row_groups.len());
443 let mut file_row_offset = 0usize;
444 let mut pages_considered = 0u64;
445
446 for (group_idx, row_group) in row_groups.iter().enumerate() {
447 let rows = usize::try_from(row_group.num_rows()).ok()?;
448 if rows == 0 {
449 return None;
450 }
451 let expected_min = i64::try_from(file_row_offset.checked_add(1)?).ok()?;
452 let expected_max = i64::try_from(file_row_offset.checked_add(rows)?).ok()?;
453 let Statistics::Int64(group_stats) = row_group.column(key_leaf).statistics()? else {
454 return None;
455 };
456 if group_stats.null_count_opt() != Some(0)
457 || group_stats.min_opt() != Some(&expected_min)
458 || group_stats.max_opt() != Some(&expected_max)
459 {
460 return None;
461 }
462
463 let page_index = column_indexes.get(group_idx)?.get(key_leaf)?;
464 if page_index.get_boundary_order() != Some(BoundaryOrder::ASCENDING) {
465 return None;
466 }
467 let ColumnIndexMetaData::INT64(page_stats) = page_index else {
468 return None;
469 };
470 let locations = offset_indexes
471 .get(group_idx)?
472 .get(key_leaf)?
473 .page_locations();
474 if locations.is_empty()
475 || usize::try_from(page_stats.num_pages()).ok()? != locations.len()
476 || (0..locations.len()).any(|page| page_stats.null_count(page) != Some(0))
477 {
478 return None;
479 }
480
481 let mut first_rows = Vec::with_capacity(locations.len());
482 for (page_idx, location) in locations.iter().enumerate() {
483 let first = usize::try_from(location.first_row_index).ok()?;
484 if (page_idx == 0 && first != 0)
485 || first >= rows
486 || first_rows.last().is_some_and(|previous| *previous >= first)
487 {
488 return None;
489 }
490 first_rows.push(first);
491 }
492 for (page_idx, &first) in first_rows.iter().enumerate() {
493 let end = first_rows.get(page_idx + 1).copied().unwrap_or(rows);
494 let page_rows = end.checked_sub(first)?;
495 let page_min =
496 i64::try_from(file_row_offset.checked_add(first)?.checked_add(1)?).ok()?;
497 let page_max =
498 i64::try_from(file_row_offset.checked_add(first)?.checked_add(page_rows)?).ok()?;
499 if page_stats.min_value(page_idx) != Some(&page_min)
500 || page_stats.max_value(page_idx) != Some(&page_max)
501 {
502 return None;
503 }
504 }
505
506 pages_considered = pages_considered.checked_add(u64::try_from(locations.len()).ok()?)?;
507 group_rows.push(rows);
508 group_pages.push(first_rows);
509 file_row_offset = file_row_offset.checked_add(rows)?;
510 }
511 if file_row_offset != total_rows {
512 return None;
513 }
514
515 Some(DenseNodeLayout {
516 group_rows,
517 group_pages,
518 total_rows,
519 pages_considered,
520 })
521}
522
523fn dense_node_selection(
527 metadata: &parquet::file::metadata::ParquetMetaData,
528 key_leaf: usize,
529 sorted_ids: &[u64],
530) -> Option<DenseNodeSelection> {
531 let DenseNodeLayout {
532 group_rows,
533 group_pages,
534 total_rows,
535 pages_considered,
536 } = dense_node_layout(metadata, key_leaf)?;
537
538 let max_id = u64::try_from(total_rows).ok()?;
539 let ordinals: Vec<usize> = sorted_ids
540 .iter()
541 .copied()
542 .filter(|&id| id != 0 && id <= max_id)
543 .map(|id| usize::try_from(id - 1).ok())
544 .collect::<Option<_>>()?;
545 let mut selected_groups = Vec::new();
546 let mut ranges = Vec::with_capacity(ordinals.len());
547 let mut selected_pages = 0u64;
548 let mut ordinal_cursor = 0usize;
549 let mut file_start = 0usize;
550 let mut retained_start = 0usize;
551
552 for (group_idx, &rows) in group_rows.iter().enumerate() {
553 let file_end = file_start.checked_add(rows)?;
554 let first = ordinal_cursor;
555 while ordinal_cursor < ordinals.len() && ordinals[ordinal_cursor] < file_end {
556 ordinal_cursor += 1;
557 }
558 if first != ordinal_cursor {
559 selected_groups.push(group_idx);
560 let mut last_page = None;
561 for &ordinal in &ordinals[first..ordinal_cursor] {
562 let local = ordinal.checked_sub(file_start)?;
563 let selected = retained_start.checked_add(local)?;
564 ranges.push(selected..selected.checked_add(1)?);
565 let page = group_pages[group_idx].partition_point(|&start| start <= local) - 1;
566 if last_page != Some(page) {
567 selected_pages = selected_pages.checked_add(1)?;
568 last_page = Some(page);
569 }
570 }
571 retained_start = retained_start.checked_add(rows)?;
572 }
573 file_start = file_end;
574 }
575
576 Some(DenseNodeSelection {
577 row_groups: selected_groups,
578 selection: parquet::arrow::arrow_reader::RowSelection::from_consecutive_ranges(
579 ranges.into_iter(),
580 retained_start,
581 ),
582 pages_considered,
583 pages_selected: selected_pages,
584 exact_rows_selected: u64::try_from(ordinals.len()).ok()?,
585 })
586}
587
588fn filtered_keys_match(
589 batches: &[RecordBatch],
590 key_column: &str,
591 expected: &std::collections::HashSet<u64>,
592) -> bool {
593 use arrow::array::Array as _;
594
595 let mut actual = std::collections::HashSet::with_capacity(expected.len());
596 let mut rows = 0usize;
597 for batch in batches {
598 let Some(column) = batch.column_by_name(key_column) else {
599 return false;
600 };
601 let Some(ids) = column.as_any().downcast_ref::<arrow::array::UInt64Array>() else {
602 return false;
603 };
604 rows = match rows.checked_add(ids.len()) {
605 Some(rows) => rows,
606 None => return false,
607 };
608 for row in 0..ids.len() {
609 if ids.is_null(row) || !actual.insert(ids.value(row)) {
610 return false;
611 }
612 }
613 }
614 rows == expected.len() && actual == *expected
615}
616
617#[allow(clippy::too_many_lines)]
622fn read_parquet_filtered_u64(
623 path: &Path,
624 fallback_schema: SchemaRef,
625 key_column: &str,
626 ids: &std::collections::HashSet<u64>,
627 kind: FilteredReadKind,
628 observer: Option<&std::sync::Arc<dyn crate::io_stats::FilteredReadObserver>>,
629) -> Result<Vec<RecordBatch>, DataFusionError> {
630 if ids.is_empty() || !path.exists() {
633 return Ok(vec![RecordBatch::new_empty(fallback_schema)]);
634 }
635 read_parquet_filtered_u64_attempt(path, fallback_schema, key_column, ids, kind, observer, true)
636}
637
638#[allow(clippy::too_many_lines, clippy::too_many_arguments)]
639fn read_parquet_filtered_u64_attempt(
640 path: &Path,
641 fallback_schema: SchemaRef,
642 key_column: &str,
643 ids: &std::collections::HashSet<u64>,
644 kind: FilteredReadKind,
645 observer: Option<&std::sync::Arc<dyn crate::io_stats::FilteredReadObserver>>,
646 allow_dense_node_selection: bool,
647) -> Result<Vec<RecordBatch>, DataFusionError> {
648 use parquet::arrow::ProjectionMask;
649 use parquet::arrow::arrow_reader::{
650 ArrowPredicateFn, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowFilter,
651 };
652 use parquet::file::metadata::PageIndexPolicy;
653 use parquet::file::statistics::Statistics;
654
655 let mut observation = FilteredReadObservation::new(observer, kind);
656 let file = File::open(path).map_err(|e| io_err(&e))?;
657 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Optional);
660 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(file, options)
661 .map_err(parquet_err)?;
662
663 let total = builder.metadata().file_metadata().num_rows();
667 let builder_row_groups = u64::try_from(builder.metadata().num_row_groups()).unwrap_or(u64::MAX);
668 if total >= 0 && ids.len() as u64 * 2 > u64::try_from(total).unwrap_or(u64::MAX) {
669 drop(builder);
670 let batches = read_parquet_or_empty(path, fallback_schema.clone())?;
671 let scanned = total_rows(&batches);
675 record_full(kind, scanned);
676 observation.scanned(scanned);
677 let file_schema = batches
682 .first()
683 .map_or_else(|| fallback_schema.clone(), RecordBatch::schema);
684 let key_idx = file_schema
685 .index_of(key_column)
686 .map_err(|e| DataFusionError::Execution(format!("filtered read: {e}")))?;
687 let mut filtered = Vec::with_capacity(batches.len());
688 for batch in &batches {
689 let col = batch
690 .column(key_idx)
691 .as_any()
692 .downcast_ref::<arrow::array::UInt64Array>()
693 .ok_or_else(|| {
694 DataFusionError::Execution("filtered read: key column not UInt64".into())
695 })?;
696 let mask: arrow::array::BooleanArray = {
697 use arrow::array::Array as _;
698 (0..col.len())
699 .map(|i| Some(!col.is_null(i) && ids.contains(&col.value(i))))
700 .collect()
701 };
702 filtered.push(
703 arrow::compute::filter_record_batch(batch, &mask)
704 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?,
705 );
706 }
707 if filtered.is_empty() {
708 filtered.push(RecordBatch::new_empty(fallback_schema));
709 }
710 record_pruning(
711 kind,
712 &observation,
713 crate::io_stats::FilteredReadPruning {
714 strategy: crate::io_stats::FilteredReadStrategy::FullFallback,
715 row_groups_considered: builder_row_groups,
716 row_groups_selected: builder_row_groups,
717 pages_considered: 0,
718 pages_selected: 0,
719 exact_rows_selected: 0,
720 metadata_fallbacks: 0,
721 validation_fallbacks: 0,
722 },
723 );
724 observation.complete(total_rows(&filtered), true);
725 return Ok(filtered);
726 }
727
728 let key_leaf = builder
730 .parquet_schema()
731 .columns()
732 .iter()
733 .position(|c| c.name() == key_column)
734 .ok_or_else(|| {
735 DataFusionError::Execution(format!("filtered read: no column {key_column}"))
736 })?;
737
738 let mut sorted: Vec<u64> = ids.iter().copied().collect();
739 sorted.sort_unstable();
740 let dense_requested =
741 allow_dense_node_selection && kind == FilteredReadKind::Node && key_column == "node_id";
742 let dense = dense_requested
743 .then(|| dense_node_selection(builder.metadata(), key_leaf, &sorted))
744 .flatten();
745 let metadata_fallbacks = u64::from(dense_requested && dense.is_none());
746
747 let (keep, selection, mut pruning) = if let Some(dense) = dense {
750 let selected_groups = u64::try_from(dense.row_groups.len()).unwrap_or(u64::MAX);
751 (
752 dense.row_groups,
753 Some(dense.selection),
754 crate::io_stats::FilteredReadPruning {
755 strategy: crate::io_stats::FilteredReadStrategy::DenseRowSelection,
756 row_groups_considered: builder_row_groups,
757 row_groups_selected: selected_groups,
758 pages_considered: dense.pages_considered,
759 pages_selected: dense.pages_selected,
760 exact_rows_selected: dense.exact_rows_selected,
761 metadata_fallbacks: 0,
762 validation_fallbacks: 0,
763 },
764 )
765 } else {
766 let keep: Vec<usize> = builder
768 .metadata()
769 .row_groups()
770 .iter()
771 .enumerate()
772 .filter(|(_, rg)| match rg.column(key_leaf).statistics() {
773 Some(Statistics::Int64(s)) => match (s.min_opt(), s.max_opt()) {
774 (Some(&min), Some(&max)) => {
775 let lo = u64::try_from(min).unwrap_or(0);
776 let hi = u64::try_from(max).unwrap_or(u64::MAX);
777 sorted.partition_point(|&x| x < lo) < sorted.partition_point(|&x| x <= hi)
778 }
779 _ => true,
780 },
781 _ => true,
782 })
783 .map(|(i, _)| i)
784 .collect();
785 let selected_groups = u64::try_from(keep.len()).unwrap_or(u64::MAX);
786 (
787 keep,
788 None,
789 crate::io_stats::FilteredReadPruning {
790 strategy: crate::io_stats::FilteredReadStrategy::RowGroupPredicate,
791 row_groups_considered: builder_row_groups,
792 row_groups_selected: selected_groups,
793 pages_considered: 0,
794 pages_selected: 0,
795 exact_rows_selected: 0,
796 metadata_fallbacks,
797 validation_fallbacks: 0,
798 },
799 )
800 };
801 let used_dense_selection = selection.is_some();
802
803 let mask = ProjectionMask::leaves(builder.parquet_schema(), [key_leaf]);
805 let owned: std::sync::Arc<std::collections::HashSet<u64>> = std::sync::Arc::new(ids.clone());
806 let scan_observer = observer.cloned();
807 let predicate = ArrowPredicateFn::new(mask, move |batch: RecordBatch| {
808 use arrow::array::Array as _;
809 let col = batch
810 .column(0)
811 .as_any()
812 .downcast_ref::<arrow::array::UInt64Array>()
813 .ok_or_else(|| {
814 arrow::error::ArrowError::CastError("filtered read: key column not UInt64".into())
815 })?;
816 let rows = total_rows(std::slice::from_ref(&batch));
820 record_scanned(kind, rows);
821 if let Some(observer) = &scan_observer {
822 observer.rows_scanned(kind.into(), rows);
823 }
824 Ok((0..col.len())
825 .map(|i| Some(!col.is_null(i) && owned.contains(&col.value(i))))
826 .collect())
827 });
828 let builder = builder.with_row_groups(keep);
829 let builder = if let Some(selection) = selection {
830 builder.with_row_selection(selection)
831 } else {
832 builder
833 };
834 let reader = builder
835 .with_row_filter(RowFilter::new(vec![Box::new(predicate)]))
836 .build()
837 .map_err(parquet_err)?;
838 let batches: Vec<RecordBatch> = reader
839 .collect::<Result<_, _>>()
840 .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?;
841 let returned = total_rows(&batches);
845 record_filtered(kind, returned);
846 if used_dense_selection {
847 let max_id = u64::try_from(total).unwrap_or(0);
848 let expected: std::collections::HashSet<u64> = ids
849 .iter()
850 .copied()
851 .filter(|&id| id != 0 && id <= max_id)
852 .collect();
853 if !filtered_keys_match(&batches, key_column, &expected) {
854 pruning.validation_fallbacks = 1;
855 record_pruning(kind, &observation, pruning);
856 observation.complete(returned, false);
857 return read_parquet_filtered_u64_attempt(
858 path,
859 fallback_schema,
860 key_column,
861 ids,
862 kind,
863 observer,
864 false,
865 );
866 }
867 }
868 record_pruning(kind, &observation, pruning);
869 observation.complete(returned, false);
870 if batches.is_empty() {
871 return Ok(vec![RecordBatch::new_empty(fallback_schema)]);
872 }
873 Ok(batches)
874}
875
876fn record_full(kind: FilteredReadKind, rows: u64) {
878 match kind {
879 FilteredReadKind::Edge => crate::io_stats::record_edge_full_read(rows),
880 FilteredReadKind::Node => crate::io_stats::record_node_full_read(rows),
881 }
882}
883
884fn record_filtered(kind: FilteredReadKind, rows: u64) {
886 match kind {
887 FilteredReadKind::Edge => crate::io_stats::record_edge_filtered_read(rows),
888 FilteredReadKind::Node => crate::io_stats::record_node_filtered_read(rows),
889 }
890}
891
892fn record_scanned(kind: FilteredReadKind, rows: u64) {
895 match kind {
896 FilteredReadKind::Edge => crate::io_stats::record_edge_scanned(rows),
897 FilteredReadKind::Node => crate::io_stats::record_node_scanned(rows),
898 }
899}
900
901fn record_pruning(
904 kind: FilteredReadKind,
905 observation: &FilteredReadObservation,
906 pruning: crate::io_stats::FilteredReadPruning,
907) {
908 if kind == FilteredReadKind::Node {
909 crate::io_stats::record_node_pruning(pruning);
910 }
911 observation.pruning(pruning);
912}
913
914pub fn read_nodes(dir: &Path) -> Result<Vec<RecordBatch>, DataFusionError> {
922 let path = dir.join("topology").join("nodes.parquet");
923 let batches =
924 normalize_topology_nodes(read_parquet_or_empty(&path, TOPOLOGY_NODES_SCHEMA.clone())?)?;
925 crate::io_stats::record_node_full_read(total_rows(&batches));
926 Ok(batches)
927}
928
929#[allow(clippy::implicit_hasher)]
943pub fn read_nodes_filtered(
944 dir: &Path,
945 node_ids: &std::collections::HashSet<u64>,
946) -> Result<Vec<RecordBatch>, DataFusionError> {
947 read_nodes_filtered_observed(dir, node_ids, None)
948}
949
950#[allow(clippy::implicit_hasher)]
952#[doc(hidden)]
953pub fn read_nodes_filtered_observed(
954 dir: &Path,
955 node_ids: &std::collections::HashSet<u64>,
956 observer: Option<&std::sync::Arc<dyn crate::io_stats::FilteredReadObserver>>,
957) -> Result<Vec<RecordBatch>, DataFusionError> {
958 let path = dir.join("topology").join("nodes.parquet");
959 normalize_topology_nodes(read_parquet_filtered_u64(
960 &path,
961 TOPOLOGY_NODES_SCHEMA.clone(),
962 "node_id",
963 node_ids,
964 FilteredReadKind::Node,
965 observer,
966 )?)
967}
968
969pub(crate) fn max_edge_id(dir: &Path) -> Result<u64, DataFusionError> {
979 use arrow::array::{Array, UInt64Array};
980
981 let edges_dir = dir.join("topology").join("edges");
982 let entries = match std::fs::read_dir(&edges_dir) {
983 Ok(rd) => rd,
984 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(0),
985 Err(e) => return Err(io_err(&e)),
986 };
987
988 let mut max = 0u64;
989 for entry in entries {
990 let path = entry.map_err(|e| io_err(&e))?.path();
991 if path.extension().and_then(|s| s.to_str()) != Some("parquet") {
992 continue;
993 }
994 let Some(schema) = discover_parquet_schema(&path) else {
997 continue;
998 };
999 for batch in read_parquet_or_empty(&path, schema)? {
1000 if let Some(col) = batch.column_by_name("edge_id")
1001 && let Some(ids) = col.as_any().downcast_ref::<UInt64Array>()
1002 {
1003 for i in 0..ids.len() {
1004 if !ids.is_null(i) {
1005 max = max.max(ids.value(i));
1006 }
1007 }
1008 }
1009 }
1010 }
1011 Ok(max)
1012}
1013
1014pub fn read_properties(dir: &Path, stem: &str) -> Result<Vec<RecordBatch>, DataFusionError> {
1024 let path = dir.join("properties").join(format!("{stem}.parquet"));
1025 match discover_parquet_schema(&path) {
1026 Some(schema) => read_parquet_or_empty(&path, schema),
1027 None => Ok(Vec::new()),
1028 }
1029}
1030
1031pub fn read_edge_properties(dir: &Path, stem: &str) -> Result<Vec<RecordBatch>, DataFusionError> {
1038 let path = dir.join("edge_properties").join(format!("{stem}.parquet"));
1039 match discover_parquet_schema(&path) {
1040 Some(schema) => read_parquet_or_empty(&path, schema),
1041 None => Ok(Vec::new()),
1042 }
1043}
1044
1045#[must_use]
1050pub fn list_edge_property_stems(dir: &Path) -> Vec<String> {
1051 list_parquet_stems(&dir.join("edge_properties"))
1052}
1053
1054#[must_use]
1058pub fn list_property_stems(dir: &Path) -> Vec<String> {
1059 list_parquet_stems(&dir.join("properties"))
1060}
1061
1062fn list_parquet_stems(dir: &Path) -> Vec<String> {
1064 let Ok(entries) = std::fs::read_dir(dir) else {
1065 return Vec::new();
1066 };
1067 let mut stems: Vec<String> = entries
1068 .filter_map(|entry| {
1069 let path = entry.ok()?.path();
1070 if path.extension().and_then(|e| e.to_str()) != Some("parquet") {
1071 return None;
1072 }
1073 Some(path.file_stem()?.to_str()?.to_owned())
1074 })
1075 .collect();
1076 stems.sort();
1077 stems
1078}
1079
1080#[derive(Debug, Clone)]
1086pub struct TopologyNodeTable {
1087 path: PathBuf,
1088}
1089
1090impl TopologyNodeTable {
1091 #[must_use]
1093 pub fn new(path: PathBuf) -> Self {
1094 Self { path }
1095 }
1096}
1097
1098#[async_trait]
1099impl TableProvider for TopologyNodeTable {
1100 fn as_any(&self) -> &dyn Any {
1101 self
1102 }
1103
1104 fn schema(&self) -> SchemaRef {
1105 TOPOLOGY_NODES_SCHEMA.clone()
1106 }
1107
1108 fn table_type(&self) -> TableType {
1109 TableType::Base
1110 }
1111
1112 async fn scan(
1113 &self,
1114 state: &dyn Session,
1115 projection: Option<&Vec<usize>>,
1116 filters: &[Expr],
1117 limit: Option<usize>,
1118 ) -> Result<Arc<dyn ExecutionPlan>, DataFusionError> {
1119 let batches = normalize_topology_nodes(read_parquet_or_empty(
1120 &self.path,
1121 TOPOLOGY_NODES_SCHEMA.clone(),
1122 )?)?;
1123 let mem = MemTable::try_new(TOPOLOGY_NODES_SCHEMA.clone(), vec![batches])?;
1124 mem.scan(state, projection, filters, limit).await
1125 }
1126}
1127
1128#[derive(Debug, Clone)]
1134pub struct TypedEdgeTable {
1135 path: PathBuf,
1136 schema: SchemaRef,
1137}
1138
1139impl TypedEdgeTable {
1140 #[must_use]
1145 pub fn open(dir: &Path, rel_type_name: &str) -> Self {
1146 let path = dir
1147 .join("topology")
1148 .join("edges")
1149 .join(format!("{rel_type_name}.parquet"));
1150 let schema = if rel_type_name == "_exploratory" {
1151 EXPLORATORY_EDGE_SCHEMA.clone()
1152 } else {
1153 TYPED_EDGE_SCHEMA.clone()
1154 };
1155 Self { path, schema }
1156 }
1157}
1158
1159#[async_trait]
1160impl TableProvider for TypedEdgeTable {
1161 fn as_any(&self) -> &dyn Any {
1162 self
1163 }
1164
1165 fn schema(&self) -> SchemaRef {
1166 self.schema.clone()
1167 }
1168
1169 fn table_type(&self) -> TableType {
1170 TableType::Base
1171 }
1172
1173 async fn scan(
1174 &self,
1175 state: &dyn Session,
1176 projection: Option<&Vec<usize>>,
1177 filters: &[Expr],
1178 limit: Option<usize>,
1179 ) -> Result<Arc<dyn ExecutionPlan>, DataFusionError> {
1180 let batches = read_parquet_or_empty(&self.path, self.schema.clone())?;
1181 let mem = MemTable::try_new(self.schema.clone(), vec![batches])?;
1182 mem.scan(state, projection, filters, limit).await
1183 }
1184}
1185
1186#[derive(Debug, Clone)]
1197pub struct UnionEdgeTable {
1198 dir: PathBuf,
1199}
1200
1201impl UnionEdgeTable {
1202 #[must_use]
1204 pub fn open(dir: &Path) -> Self {
1205 Self {
1206 dir: dir.to_path_buf(),
1207 }
1208 }
1209}
1210
1211#[async_trait]
1212impl TableProvider for UnionEdgeTable {
1213 fn as_any(&self) -> &dyn Any {
1214 self
1215 }
1216
1217 fn schema(&self) -> SchemaRef {
1218 EXPLORATORY_EDGE_SCHEMA.clone()
1219 }
1220
1221 fn table_type(&self) -> TableType {
1222 TableType::Base
1223 }
1224
1225 async fn scan(
1226 &self,
1227 state: &dyn Session,
1228 projection: Option<&Vec<usize>>,
1229 filters: &[Expr],
1230 limit: Option<usize>,
1231 ) -> Result<Arc<dyn ExecutionPlan>, DataFusionError> {
1232 let batches = read_edges_union(&self.dir, None, None)?;
1233 let mem = MemTable::try_new(EXPLORATORY_EDGE_SCHEMA.clone(), vec![batches])?;
1234 mem.scan(state, projection, filters, limit).await
1235 }
1236}
1237
1238#[derive(Debug, Clone)]
1244pub struct PropertyTable {
1245 path: PathBuf,
1246 schema: SchemaRef,
1247}
1248
1249impl PropertyTable {
1250 #[must_use]
1255 pub fn open(dir: &Path, entity_type: &str, schema: SchemaRef) -> Self {
1256 let path = dir
1257 .join("properties")
1258 .join(format!("{entity_type}.parquet"));
1259 Self { path, schema }
1260 }
1261
1262 #[must_use]
1272 pub fn open_discovered(dir: &Path, stem: &str) -> Self {
1273 let path = dir.join("properties").join(format!("{stem}.parquet"));
1274 let schema = discover_parquet_schema(&path)
1275 .unwrap_or_else(|| crate::schemas::PROPERTY_BASE_SCHEMA.clone());
1276 Self { path, schema }
1277 }
1278
1279 #[must_use]
1281 pub fn schema_ref(&self) -> SchemaRef {
1282 self.schema.clone()
1283 }
1284}
1285
1286#[derive(Debug, Clone)]
1296pub struct EdgePropertyTable {
1297 path: PathBuf,
1298 schema: SchemaRef,
1299}
1300
1301impl EdgePropertyTable {
1302 #[must_use]
1311 pub fn open_discovered(dir: &Path, rel_type: &str) -> Self {
1312 let path = dir
1313 .join("edge_properties")
1314 .join(format!("{rel_type}.parquet"));
1315 let schema = discover_parquet_schema(&path)
1316 .unwrap_or_else(|| crate::schemas::EDGE_PROPERTY_BASE_SCHEMA.clone());
1317 Self { path, schema }
1318 }
1319
1320 #[must_use]
1322 pub fn schema_ref(&self) -> SchemaRef {
1323 self.schema.clone()
1324 }
1325}
1326
1327#[async_trait]
1328impl TableProvider for EdgePropertyTable {
1329 fn as_any(&self) -> &dyn Any {
1330 self
1331 }
1332
1333 fn schema(&self) -> SchemaRef {
1334 self.schema.clone()
1335 }
1336
1337 fn table_type(&self) -> TableType {
1338 TableType::Base
1339 }
1340
1341 async fn scan(
1342 &self,
1343 state: &dyn Session,
1344 projection: Option<&Vec<usize>>,
1345 filters: &[Expr],
1346 limit: Option<usize>,
1347 ) -> Result<Arc<dyn ExecutionPlan>, DataFusionError> {
1348 let batches = read_parquet_or_empty(&self.path, self.schema.clone())?;
1349 let mem = MemTable::try_new(self.schema.clone(), vec![batches])?;
1350 mem.scan(state, projection, filters, limit).await
1351 }
1352}
1353
1354pub(crate) fn discover_parquet_schema(path: &Path) -> Option<SchemaRef> {
1357 let file = File::open(path).ok()?;
1358 let builder = ParquetRecordBatchReaderBuilder::try_new(file).ok()?;
1359 Some(builder.schema().clone())
1360}
1361
1362#[async_trait]
1363impl TableProvider for PropertyTable {
1364 fn as_any(&self) -> &dyn Any {
1365 self
1366 }
1367
1368 fn schema(&self) -> SchemaRef {
1369 self.schema.clone()
1370 }
1371
1372 fn table_type(&self) -> TableType {
1373 TableType::Base
1374 }
1375
1376 async fn scan(
1377 &self,
1378 state: &dyn Session,
1379 projection: Option<&Vec<usize>>,
1380 filters: &[Expr],
1381 limit: Option<usize>,
1382 ) -> Result<Arc<dyn ExecutionPlan>, DataFusionError> {
1383 let batches = read_parquet_or_empty(&self.path, self.schema.clone())?;
1384 let mem = MemTable::try_new(self.schema.clone(), vec![batches])?;
1385 mem.scan(state, projection, filters, limit).await
1386 }
1387}
1388
1389struct GraphSchema {
1394 tables: HashMap<String, Arc<dyn TableProvider>>,
1395}
1396
1397impl fmt::Debug for GraphSchema {
1398 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1399 f.debug_struct("GraphSchema")
1400 .field("table_names", &self.table_names())
1401 .finish()
1402 }
1403}
1404
1405impl GraphSchema {
1406 fn new() -> Self {
1407 Self {
1408 tables: HashMap::new(),
1409 }
1410 }
1411
1412 fn register(&mut self, name: impl Into<String>, table: Arc<dyn TableProvider>) {
1413 self.tables.insert(name.into(), table);
1414 }
1415}
1416
1417#[async_trait]
1418impl SchemaProvider for GraphSchema {
1419 fn as_any(&self) -> &dyn Any {
1420 self
1421 }
1422
1423 fn table_names(&self) -> Vec<String> {
1424 let mut names: Vec<String> = self.tables.keys().cloned().collect();
1425 names.sort();
1426 names
1427 }
1428
1429 async fn table(&self, name: &str) -> Result<Option<Arc<dyn TableProvider>>, DataFusionError> {
1430 Ok(self.tables.get(name).cloned())
1431 }
1432
1433 fn table_exist(&self, name: &str) -> bool {
1434 self.tables.contains_key(name)
1435 }
1436}
1437
1438pub struct GraphCatalog {
1447 schema: Arc<GraphSchema>,
1448 prop_names: HashMap<u32, String>,
1453 rel_names: HashMap<u32, String>,
1457 label_names: HashMap<u32, String>,
1462}
1463
1464impl fmt::Debug for GraphCatalog {
1465 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1466 f.debug_struct("GraphCatalog")
1467 .field("schema_names", &self.schema_names())
1468 .finish()
1469 }
1470}
1471
1472impl GraphCatalog {
1473 pub fn open(
1483 dir: &Path,
1484 ontology: Option<&OntologyHandle>,
1485 runtime_catalog: &RuntimeCatalog,
1486 ) -> Result<Self, DataFusionError> {
1487 let mut schema = GraphSchema::new();
1488
1489 let nodes_path = dir.join("topology").join("nodes.parquet");
1491 schema.register(
1492 "topology_nodes",
1493 Arc::new(TopologyNodeTable::new(nodes_path)),
1494 );
1495
1496 if let Some(handle) = ontology {
1498 for rel_name in handle.relation_type_names() {
1499 schema.register(
1500 format!("edges_{rel_name}"),
1501 Arc::new(TypedEdgeTable::open(dir, rel_name)),
1502 );
1503 }
1504 } else {
1505 schema.register(
1514 "edges__exploratory",
1515 Arc::new(TypedEdgeTable::open(dir, "_exploratory")),
1516 );
1517 for rel_name in runtime_catalog.relation_types() {
1518 let typed_path = dir
1519 .join("topology")
1520 .join("edges")
1521 .join(format!("{rel_name}.parquet"));
1522 if typed_path.exists() {
1523 schema.register(
1524 format!("edges_{rel_name}"),
1525 Arc::new(TypedEdgeTable::open(dir, rel_name)),
1526 );
1527 }
1528 }
1529 }
1530
1531 if !schema.table_exist("edges__exploratory") {
1533 schema.register(
1534 "edges__exploratory",
1535 Arc::new(TypedEdgeTable::open(dir, "_exploratory")),
1536 );
1537 }
1538
1539 register_property_tables(dir, ontology, &mut schema);
1541
1542 let prop_names = build_prop_names(ontology, runtime_catalog);
1544 let rel_names = build_rel_names(runtime_catalog);
1545 let label_names = build_label_names(runtime_catalog);
1546
1547 Ok(Self {
1548 schema: Arc::new(schema),
1549 prop_names,
1550 rel_names,
1551 label_names,
1552 })
1553 }
1554
1555 #[must_use]
1558 pub fn prop_names(&self) -> &HashMap<u32, String> {
1559 &self.prop_names
1560 }
1561
1562 #[must_use]
1566 pub fn rel_names(&self) -> &HashMap<u32, String> {
1567 &self.rel_names
1568 }
1569
1570 #[must_use]
1575 pub fn label_names(&self) -> &HashMap<u32, String> {
1576 &self.label_names
1577 }
1578}
1579
1580fn build_prop_names(
1587 _ontology: Option<&OntologyHandle>,
1588 runtime_catalog: &RuntimeCatalog,
1589) -> HashMap<u32, String> {
1590 runtime_catalog
1591 .property_names()
1592 .map(|(id, name)| (id.0, name.to_owned()))
1593 .collect()
1594}
1595
1596fn build_rel_names(runtime_catalog: &RuntimeCatalog) -> HashMap<u32, String> {
1603 runtime_catalog
1604 .relation_type_names_with_ids()
1605 .map(|(id, name)| (id.0, name.to_owned()))
1606 .collect()
1607}
1608
1609fn build_label_names(runtime_catalog: &RuntimeCatalog) -> HashMap<u32, String> {
1617 runtime_catalog
1618 .entity_type_names_with_ids()
1619 .map(|(id, name)| (id.0, name.to_owned()))
1620 .collect()
1621}
1622
1623impl CatalogProvider for GraphCatalog {
1624 fn as_any(&self) -> &dyn Any {
1625 self
1626 }
1627
1628 fn schema_names(&self) -> Vec<String> {
1629 vec!["graph".to_owned()]
1630 }
1631
1632 fn schema(&self, name: &str) -> Option<Arc<dyn SchemaProvider>> {
1633 if name == "graph" {
1634 Some(self.schema.clone())
1635 } else {
1636 None
1637 }
1638 }
1639}
1640
1641fn register_property_tables(
1646 dir: &Path,
1647 ontology: Option<&OntologyHandle>,
1648 schema: &mut GraphSchema,
1649) {
1650 if let Some(handle) = ontology {
1651 for (entity_name, prop_defs) in handle.entity_property_defs() {
1652 let prop_schema = Arc::new(property_schema(entity_name, &prop_defs));
1653 schema.register(
1654 format!("properties_{entity_name}"),
1655 Arc::new(PropertyTable::open(dir, entity_name, prop_schema)),
1656 );
1657 }
1658 } else {
1659 schema.register(
1663 "properties__untyped",
1664 Arc::new(PropertyTable::open_discovered(dir, "_untyped")),
1665 );
1666 }
1667}
1668
1669#[cfg(test)]
1674mod tests {
1675 use super::*;
1676 use arrow::array::{
1677 FixedSizeBinaryArray, StringArray, TimestampMicrosecondArray, UInt32Array, UInt64Array,
1678 };
1679 use arrow::buffer::OffsetBuffer;
1680 use arrow::datatypes::{DataType, Field, Schema};
1681 use datafusion::prelude::SessionContext;
1682 use parquet::arrow::ArrowWriter;
1683 use parquet::file::properties::WriterProperties;
1684 use tempfile::TempDir;
1685
1686 #[test]
1687 fn parquet_and_io_error_helpers_preserve_external_messages() {
1688 let parquet = parquet_err("parquet boom");
1689 assert!(parquet.to_string().contains("parquet boom"));
1690 let io = io_err(&std::io::Error::other("io boom"));
1691 assert!(io.to_string().contains("io boom"));
1692 }
1693
1694 #[derive(Default)]
1695 struct Wave12Observer {
1696 started: std::sync::atomic::AtomicUsize,
1697 scanned: std::sync::atomic::AtomicUsize,
1698 completed: std::sync::atomic::AtomicUsize,
1699 failed: std::sync::atomic::AtomicUsize,
1700 pruning: std::sync::atomic::AtomicUsize,
1701 }
1702
1703 impl crate::io_stats::FilteredReadObserver for Wave12Observer {
1704 fn read_started(&self, _: crate::io_stats::FilteredReadTable) {
1705 self.started
1706 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1707 }
1708
1709 fn rows_scanned(&self, _: crate::io_stats::FilteredReadTable, rows: u64) {
1710 self.scanned
1711 .fetch_add(rows as usize, std::sync::atomic::Ordering::Relaxed);
1712 }
1713
1714 fn read_completed(&self, _: crate::io_stats::FilteredReadTable, rows: u64, _: bool) {
1715 self.completed
1716 .fetch_add(rows as usize, std::sync::atomic::Ordering::Relaxed);
1717 }
1718
1719 fn read_failed(&self, _: crate::io_stats::FilteredReadTable) {
1720 self.failed
1721 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1722 }
1723
1724 fn pruning(
1725 &self,
1726 _: crate::io_stats::FilteredReadTable,
1727 _: crate::io_stats::FilteredReadPruning,
1728 ) {
1729 self.pruning
1730 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1731 }
1732 }
1733
1734 fn write_nodes_parquet(path: &Path) {
1735 let uuid_bytes: Vec<u8> = vec![1u8; 16];
1736 let uuid_arr =
1737 FixedSizeBinaryArray::try_from_iter(std::iter::once(uuid_bytes.clone())).unwrap();
1738 let ts =
1739 TimestampMicrosecondArray::from(vec![0i64]).with_timezone_opt(Some(Arc::from("UTC")));
1740 let labels = arrow::array::ListArray::new(
1741 Arc::new(Field::new("item", DataType::UInt32, false)),
1742 OffsetBuffer::new(vec![0, 1].into()),
1743 Arc::new(UInt32Array::from(vec![0u32])),
1744 None,
1745 );
1746
1747 let batch = RecordBatch::try_new(
1748 TOPOLOGY_NODES_SCHEMA.clone(),
1749 vec![
1750 Arc::new(uuid_arr),
1751 Arc::new(UInt64Array::from(vec![1u64])),
1752 Arc::new(UInt32Array::from(vec![0u32])),
1753 Arc::new(labels),
1754 Arc::new(ts.clone()),
1755 Arc::new(ts),
1756 ],
1757 )
1758 .unwrap();
1759
1760 let file = File::create(path).unwrap();
1761 let mut writer = ArrowWriter::try_new(
1762 file,
1763 TOPOLOGY_NODES_SCHEMA.clone(),
1764 Some(WriterProperties::builder().build()),
1765 )
1766 .unwrap();
1767 writer.write(&batch).unwrap();
1768 writer.close().unwrap();
1769 }
1770
1771 #[test]
1772 fn legacy_scalar_node_labels_normalize_to_singleton_sets() {
1773 let dir = TempDir::new().unwrap();
1774 std::fs::create_dir_all(dir.path().join("topology")).unwrap();
1775 let old_schema = Arc::new(Schema::new(vec![
1776 crate::schemas::uuid_field("node_uuid"),
1777 crate::schemas::id_field("node_id"),
1778 Field::new("type_id", DataType::UInt32, false),
1779 crate::schemas::ts_field("created_at"),
1780 crate::schemas::ts_field("updated_at"),
1781 ]));
1782 let uuid = FixedSizeBinaryArray::try_from_iter([vec![1u8; 16]].into_iter()).unwrap();
1783 let ts =
1784 TimestampMicrosecondArray::from(vec![0i64]).with_timezone_opt(Some(Arc::from("UTC")));
1785 let legacy = RecordBatch::try_new(
1786 old_schema,
1787 vec![
1788 Arc::new(uuid),
1789 Arc::new(UInt64Array::from(vec![1])),
1790 Arc::new(UInt32Array::from(vec![7])),
1791 Arc::new(ts.clone()),
1792 Arc::new(ts),
1793 ],
1794 )
1795 .unwrap();
1796
1797 let file = File::create(dir.path().join("topology/nodes.parquet")).unwrap();
1798 let mut writer = ArrowWriter::try_new(file, legacy.schema(), None).unwrap();
1799 writer.write(&legacy).unwrap();
1800 writer.close().unwrap();
1801
1802 let normalized = read_nodes(dir.path()).unwrap();
1803 assert_eq!(normalized[0].schema(), TOPOLOGY_NODES_SCHEMA.clone());
1804 let labels = normalized[0]
1805 .column_by_name("type_ids")
1806 .unwrap()
1807 .as_any()
1808 .downcast_ref::<arrow::array::ListArray>()
1809 .unwrap();
1810 let values = labels.value(0);
1811 let values = values.as_any().downcast_ref::<UInt32Array>().unwrap();
1812 assert_eq!(values.values(), &[7]);
1813 }
1814
1815 fn write_edge_parquet(path: &Path) {
1816 std::fs::create_dir_all(path.parent().unwrap()).unwrap();
1817 let fsb = |v: Vec<u8>| FixedSizeBinaryArray::try_from_iter(std::iter::once(v)).unwrap();
1818 let ts =
1819 TimestampMicrosecondArray::from(vec![0i64]).with_timezone_opt(Some(Arc::from("UTC")));
1820
1821 let batch = RecordBatch::try_new(
1822 TYPED_EDGE_SCHEMA.clone(),
1823 vec![
1824 Arc::new(fsb(vec![2u8; 16])),
1825 Arc::new(fsb(vec![1u8; 16])),
1826 Arc::new(fsb(vec![3u8; 16])),
1827 Arc::new(UInt64Array::from(vec![1u64])),
1828 Arc::new(UInt64Array::from(vec![1u64])),
1829 Arc::new(UInt64Array::from(vec![2u64])),
1830 Arc::new(ts),
1831 ],
1832 )
1833 .unwrap();
1834
1835 let file = File::create(path).unwrap();
1836 let mut writer = ArrowWriter::try_new(
1837 file,
1838 TYPED_EDGE_SCHEMA.clone(),
1839 Some(WriterProperties::builder().build()),
1840 )
1841 .unwrap();
1842 writer.write(&batch).unwrap();
1843 writer.close().unwrap();
1844 }
1845
1846 #[tokio::test]
1847 async fn topology_node_table_scan_returns_rows() {
1848 let dir = TempDir::new().unwrap();
1849 let nodes_dir = dir.path().join("topology");
1850 std::fs::create_dir_all(&nodes_dir).unwrap();
1851 let path = nodes_dir.join("nodes.parquet");
1852 write_nodes_parquet(&path);
1853
1854 let table = TopologyNodeTable::new(path);
1855 let ctx = SessionContext::new();
1856 ctx.register_table("nodes", Arc::new(table)).unwrap();
1857 let df = ctx.sql("SELECT node_id FROM nodes").await.unwrap();
1858 let batches = df.collect().await.unwrap();
1859 let total: usize = batches.iter().map(|b| b.num_rows()).sum();
1860 assert_eq!(total, 1);
1861 }
1862
1863 #[tokio::test]
1864 async fn topology_node_table_missing_file_returns_empty() {
1865 let table = TopologyNodeTable::new(PathBuf::from("/nonexistent/nodes.parquet"));
1866 let ctx = SessionContext::new();
1867 ctx.register_table("nodes", Arc::new(table)).unwrap();
1868 let df = ctx.sql("SELECT node_id FROM nodes").await.unwrap();
1869 let batches = df.collect().await.unwrap();
1870 let total: usize = batches.iter().map(|b| b.num_rows()).sum();
1871 assert_eq!(total, 0);
1872 }
1873
1874 #[tokio::test]
1875 async fn typed_edge_table_scan_returns_rows() {
1876 let dir = TempDir::new().unwrap();
1877 let edge_path = dir
1878 .path()
1879 .join("topology")
1880 .join("edges")
1881 .join("KNOWS.parquet");
1882 write_edge_parquet(&edge_path);
1883
1884 let table = TypedEdgeTable::open(dir.path(), "KNOWS");
1885 let ctx = SessionContext::new();
1886 ctx.register_table("edges", Arc::new(table)).unwrap();
1887 let df = ctx.sql("SELECT src_id, dst_id FROM edges").await.unwrap();
1888 let batches = df.collect().await.unwrap();
1889 let total: usize = batches.iter().map(|b| b.num_rows()).sum();
1890 assert_eq!(total, 1);
1891 }
1892
1893 #[tokio::test]
1894 async fn typed_edge_table_exploratory_has_rel_type_name_column() {
1895 let table = TypedEdgeTable::open(Path::new("/nonexistent"), "_exploratory");
1896 let schema = table.schema();
1897 assert!(
1898 schema.field_with_name("rel_type_name").is_ok(),
1899 "exploratory schema must have rel_type_name"
1900 );
1901 }
1902
1903 #[tokio::test]
1904 async fn union_edge_table_scan_unions_all_relations() {
1905 let dir = TempDir::new().unwrap();
1906 let edges = dir.path().join("topology").join("edges");
1907 write_typed_edge(&edges.join("KNOWS.parquet"), 1, 1, 2);
1908 write_typed_edge(&edges.join("OWNS.parquet"), 2, 2, 3);
1909
1910 let table = UnionEdgeTable::open(dir.path());
1911 assert_eq!(table.schema(), EXPLORATORY_EDGE_SCHEMA.clone());
1912 let ctx = SessionContext::new();
1913 ctx.register_table("edges", Arc::new(table)).unwrap();
1914 let df = ctx
1915 .sql("SELECT edge_id, rel_type_name FROM edges ORDER BY edge_id")
1916 .await
1917 .unwrap();
1918 let batches = df.collect().await.unwrap();
1919 assert_eq!(row_count(&batches), 2, "both relations' edges unioned");
1920 }
1921
1922 #[tokio::test]
1923 async fn property_table_missing_file_returns_empty_with_correct_schema() {
1924 let schema = Arc::new(property_schema("Person", &[]));
1925 let table = PropertyTable::open(Path::new("/nonexistent"), "Person", schema.clone());
1926 let ctx = SessionContext::new();
1927 ctx.register_table("props", Arc::new(table)).unwrap();
1928 let df = ctx.sql("SELECT node_uuid FROM props").await.unwrap();
1929 let batches = df.collect().await.unwrap();
1930 let total: usize = batches.iter().map(|b| b.num_rows()).sum();
1931 assert_eq!(total, 0);
1932 }
1933
1934 #[tokio::test]
1935 async fn every_table_provider_exposes_base_contract_and_empty_scan() {
1936 let dir = TempDir::new().unwrap();
1937 let providers: Vec<Arc<dyn TableProvider>> = vec![
1938 Arc::new(TopologyNodeTable::new(
1939 dir.path().join("topology/nodes.parquet"),
1940 )),
1941 Arc::new(TypedEdgeTable::open(dir.path(), "KNOWS")),
1942 Arc::new(UnionEdgeTable::open(dir.path())),
1943 Arc::new(PropertyTable::open_discovered(dir.path(), "Person")),
1944 Arc::new(EdgePropertyTable::open_discovered(dir.path(), "KNOWS")),
1945 ];
1946 let ctx = SessionContext::new();
1947 for (index, provider) in providers.into_iter().enumerate() {
1948 assert_eq!(provider.table_type(), TableType::Base);
1949 assert!(
1950 provider.as_any().is::<TopologyNodeTable>()
1951 || provider.as_any().is::<TypedEdgeTable>()
1952 || provider.as_any().is::<UnionEdgeTable>()
1953 || provider.as_any().is::<PropertyTable>()
1954 || provider.as_any().is::<EdgePropertyTable>()
1955 );
1956 let name = format!("provider_{index}");
1957 let expected = provider.schema();
1958 ctx.register_table(&name, provider).unwrap();
1959 let frame = ctx.sql(&format!("SELECT * FROM {name}")).await.unwrap();
1960 assert_eq!(frame.schema().inner(), &expected);
1961 let batches = frame.collect().await.unwrap();
1962 assert_eq!(row_count(&batches), 0);
1963 }
1964 }
1965
1966 #[test]
1967 fn graph_catalog_open_exploratory_registers_tables() {
1968 let dir = TempDir::new().unwrap();
1969 let catalog = RuntimeCatalog::new();
1970 let gc = GraphCatalog::open(dir.path(), None, &catalog).unwrap();
1971 let schema = gc.schema("graph").unwrap();
1972 let names = schema.table_names();
1973 assert!(
1974 names.contains(&"topology_nodes".to_owned()),
1975 "got {names:?}"
1976 );
1977 assert!(
1978 names.contains(&"edges__exploratory".to_owned()),
1979 "got {names:?}"
1980 );
1981 }
1982
1983 #[test]
1984 fn graph_catalog_schema_names() {
1985 let dir = TempDir::new().unwrap();
1986 let catalog = RuntimeCatalog::new();
1987 let gc = GraphCatalog::open(dir.path(), None, &catalog).unwrap();
1988 assert_eq!(gc.schema_names(), vec!["graph"]);
1989 }
1990
1991 fn row_count(batches: &[RecordBatch]) -> usize {
1996 batches.iter().map(RecordBatch::num_rows).sum()
1997 }
1998
1999 #[test]
2000 fn read_edges_strict_returns_typed_rows() {
2001 let dir = TempDir::new().unwrap();
2002 write_edge_parquet(
2003 &dir.path()
2004 .join("topology")
2005 .join("edges")
2006 .join("KNOWS.parquet"),
2007 );
2008
2009 let batches = read_edges(dir.path(), "KNOWS", OntologyMode::Strict).unwrap();
2010 assert_eq!(row_count(&batches), 1);
2011 assert_eq!(batches[0].schema(), TYPED_EDGE_SCHEMA.clone());
2013 assert!(
2014 batches[0]
2015 .schema()
2016 .field_with_name("rel_type_name")
2017 .is_err()
2018 );
2019 }
2020
2021 #[test]
2022 fn read_edges_rejects_path_traversal_rel_name() {
2023 let dir = TempDir::new().unwrap();
2024 for bad in ["../secret", "a/b", "..", "/etc/passwd"] {
2025 let err = read_edges(dir.path(), bad, OntologyMode::Strict).unwrap_err();
2026 assert!(
2027 err.to_string().contains("invalid relation name"),
2028 "expected rejection for {bad:?}, got: {err}"
2029 );
2030 }
2031 assert!(read_edges(dir.path(), "../secret", OntologyMode::Exploratory).is_ok());
2034 }
2035
2036 #[test]
2037 fn read_edges_missing_file_returns_empty_typed_batch() {
2038 let dir = TempDir::new().unwrap();
2039 let batches = read_edges(dir.path(), "KNOWS", OntologyMode::Strict).unwrap();
2041 assert_eq!(row_count(&batches), 0);
2042 assert_eq!(batches[0].schema(), TYPED_EDGE_SCHEMA.clone());
2043 }
2044
2045 #[test]
2046 fn read_edges_exploratory_uses_exploratory_file_and_schema() {
2047 let dir = TempDir::new().unwrap();
2048 write_edge_parquet(
2051 &dir.path()
2052 .join("topology")
2053 .join("edges")
2054 .join("KNOWS.parquet"),
2055 );
2056
2057 let batches = read_edges(dir.path(), "KNOWS", OntologyMode::Exploratory).unwrap();
2058 assert_eq!(row_count(&batches), 0);
2061 assert_eq!(batches[0].schema(), EXPLORATORY_EDGE_SCHEMA.clone());
2062 assert!(batches[0].schema().field_with_name("rel_type_name").is_ok());
2063 }
2064
2065 fn write_typed_edge(path: &Path, edge_id: u64, src_id: u64, dst_id: u64) {
2071 std::fs::create_dir_all(path.parent().unwrap()).unwrap();
2072 let fsb = |v: Vec<u8>| FixedSizeBinaryArray::try_from_iter(std::iter::once(v)).unwrap();
2073 let uuid = |id: u64| {
2076 let mut b = [0u8; 16];
2077 b[..8].copy_from_slice(&id.to_le_bytes());
2078 b.to_vec()
2079 };
2080 let ts =
2081 TimestampMicrosecondArray::from(vec![0i64]).with_timezone_opt(Some(Arc::from("UTC")));
2082 let batch = RecordBatch::try_new(
2083 TYPED_EDGE_SCHEMA.clone(),
2084 vec![
2085 Arc::new(fsb(uuid(edge_id))),
2086 Arc::new(fsb(uuid(src_id))),
2087 Arc::new(fsb(uuid(dst_id))),
2088 Arc::new(UInt64Array::from(vec![edge_id])),
2089 Arc::new(UInt64Array::from(vec![src_id])),
2090 Arc::new(UInt64Array::from(vec![dst_id])),
2091 Arc::new(ts),
2092 ],
2093 )
2094 .unwrap();
2095 let file = File::create(path).unwrap();
2096 let mut writer = ArrowWriter::try_new(file, TYPED_EDGE_SCHEMA.clone(), None).unwrap();
2097 writer.write(&batch).unwrap();
2098 writer.close().unwrap();
2099 }
2100
2101 fn edge_rel_pairs(batches: &[RecordBatch]) -> Vec<(u64, String)> {
2103 use arrow::array::{StringArray, UInt64Array};
2104 let mut out = Vec::new();
2105 for b in batches {
2106 let eids = b.column(3).as_any().downcast_ref::<UInt64Array>().unwrap();
2107 let rels = b
2108 .column_by_name("rel_type_name")
2109 .unwrap()
2110 .as_any()
2111 .downcast_ref::<StringArray>()
2112 .unwrap();
2113 for i in 0..b.num_rows() {
2114 out.push((eids.value(i), rels.value(i).to_owned()));
2115 }
2116 }
2117 out.sort();
2118 out
2119 }
2120
2121 #[test]
2122 fn read_edges_strict_wildcard_unions_all_relations() {
2123 let dir = TempDir::new().unwrap();
2124 let edges = dir.path().join("topology").join("edges");
2125 write_typed_edge(&edges.join("KNOWS.parquet"), 1, 1, 2);
2126 write_typed_edge(&edges.join("OWNS.parquet"), 2, 2, 3);
2127
2128 let batches = read_edges(dir.path(), "*", OntologyMode::Strict).unwrap();
2129 assert_eq!(batches[0].schema(), EXPLORATORY_EDGE_SCHEMA.clone());
2131 assert_eq!(
2132 edge_rel_pairs(&batches),
2133 vec![(1, "KNOWS".to_owned()), (2, "OWNS".to_owned())]
2134 );
2135 }
2136
2137 #[test]
2138 fn read_edges_filtered_strict_wildcard_unions_traversed_ids() {
2139 let dir = TempDir::new().unwrap();
2140 let edges = dir.path().join("topology").join("edges");
2141 write_typed_edge(&edges.join("KNOWS.parquet"), 1, 1, 2);
2142 write_typed_edge(&edges.join("OWNS.parquet"), 2, 2, 3);
2143
2144 let want: std::collections::HashSet<u64> = [2].into_iter().collect();
2145 let one = read_edges_filtered(dir.path(), "*", OntologyMode::Strict, &want).unwrap();
2146 assert_eq!(edge_rel_pairs(&one), vec![(2, "OWNS".to_owned())]);
2147
2148 let both: std::collections::HashSet<u64> = [1, 2].into_iter().collect();
2149 let two = read_edges_filtered(dir.path(), "*", OntologyMode::Strict, &both).unwrap();
2150 assert_eq!(
2151 edge_rel_pairs(&two),
2152 vec![(1, "KNOWS".to_owned()), (2, "OWNS".to_owned())]
2153 );
2154 }
2155
2156 #[test]
2157 fn read_edges_strict_wildcard_empty_dir_is_one_empty_exploratory_batch() {
2158 let dir = TempDir::new().unwrap();
2159 let batches = read_edges(dir.path(), "*", OntologyMode::Strict).unwrap();
2160 assert_eq!(row_count(&batches), 0);
2161 assert_eq!(batches[0].schema(), EXPLORATORY_EDGE_SCHEMA.clone());
2162 }
2163
2164 #[test]
2165 fn read_nodes_returns_rows_and_empty_when_absent() {
2166 let dir = TempDir::new().unwrap();
2167 std::fs::create_dir_all(dir.path().join("topology")).unwrap();
2168
2169 let empty = read_nodes(dir.path()).unwrap();
2171 assert_eq!(row_count(&empty), 0);
2172 assert_eq!(empty[0].schema(), TOPOLOGY_NODES_SCHEMA.clone());
2173
2174 write_nodes_parquet(&dir.path().join("topology").join("nodes.parquet"));
2176 let batches = read_nodes(dir.path()).unwrap();
2177 assert_eq!(row_count(&batches), 1);
2178 assert_eq!(batches[0].schema(), TOPOLOGY_NODES_SCHEMA.clone());
2179 }
2180
2181 #[test]
2182 fn catalog_and_schema_debug_identity_are_stable_and_content_free() {
2183 let schema = GraphSchema::new();
2184 assert_eq!(format!("{schema:?}"), "GraphSchema { table_names: [] }");
2185 assert!(schema.as_any().downcast_ref::<GraphSchema>().is_some());
2186 assert!(!schema.table_exist("missing"));
2187
2188 let dir = TempDir::new().unwrap();
2189 let catalog = GraphCatalog::open(dir.path(), None, &RuntimeCatalog::new()).unwrap();
2190 assert_eq!(
2191 format!("{catalog:?}"),
2192 "GraphCatalog { schema_names: [\"graph\"] }"
2193 );
2194 assert!(catalog.as_any().downcast_ref::<GraphCatalog>().is_some());
2195 assert!(catalog.schema("graph").is_some());
2196 assert!(catalog.schema("private").is_none());
2197 }
2198
2199 #[test]
2200 fn wave12_legacy_normalization_rejects_missing_and_mistyped_primary_labels() {
2201 let missing = RecordBatch::new_empty(Arc::new(Schema::new(vec![Field::new(
2202 "node_id",
2203 DataType::UInt64,
2204 false,
2205 )])));
2206 assert!(normalize_topology_nodes(vec![missing]).is_err());
2207
2208 let wrong_type = RecordBatch::try_new(
2209 Arc::new(Schema::new(vec![Field::new(
2210 "type_id",
2211 DataType::Utf8,
2212 false,
2213 )])),
2214 vec![Arc::new(StringArray::from(vec!["label"]))],
2215 )
2216 .unwrap();
2217 assert!(normalize_topology_nodes(vec![wrong_type]).is_err());
2218 }
2219
2220 #[test]
2221 fn wave12_filtered_key_validation_rejects_missing_wrong_null_and_duplicate_ids() {
2222 let missing = RecordBatch::new_empty(Arc::new(Schema::new(vec![Field::new(
2223 "other",
2224 DataType::UInt64,
2225 false,
2226 )])));
2227 assert!(!filtered_keys_match(
2228 &[missing],
2229 "node_id",
2230 &Default::default()
2231 ));
2232
2233 let wrong = RecordBatch::try_new(
2234 Arc::new(Schema::new(vec![Field::new(
2235 "node_id",
2236 DataType::Utf8,
2237 false,
2238 )])),
2239 vec![Arc::new(StringArray::from(vec!["1"]))],
2240 )
2241 .unwrap();
2242 assert!(!filtered_keys_match(
2243 &[wrong],
2244 "node_id",
2245 &[1].into_iter().collect()
2246 ));
2247
2248 let nullable = RecordBatch::try_new(
2249 Arc::new(Schema::new(vec![Field::new(
2250 "node_id",
2251 DataType::UInt64,
2252 true,
2253 )])),
2254 vec![Arc::new(UInt64Array::from(vec![Some(1), None]))],
2255 )
2256 .unwrap();
2257 assert!(!filtered_keys_match(
2258 &[nullable],
2259 "node_id",
2260 &[1].into_iter().collect()
2261 ));
2262
2263 let duplicate = RecordBatch::try_new(
2264 Arc::new(Schema::new(vec![Field::new(
2265 "node_id",
2266 DataType::UInt64,
2267 false,
2268 )])),
2269 vec![Arc::new(UInt64Array::from(vec![1, 1]))],
2270 )
2271 .unwrap();
2272 assert!(!filtered_keys_match(
2273 &[duplicate],
2274 "node_id",
2275 &[1].into_iter().collect()
2276 ));
2277 }
2278
2279 #[test]
2280 fn wave12_filtered_observation_reports_completion_or_failure_once() {
2281 let observer = Arc::new(Wave12Observer::default());
2282 {
2283 let erased: Arc<dyn crate::io_stats::FilteredReadObserver> = observer.clone();
2284 let mut observation =
2285 FilteredReadObservation::new(Some(&erased), FilteredReadKind::Node);
2286 observation.scanned(3);
2287 observation.pruning(crate::io_stats::FilteredReadPruning {
2288 strategy: crate::io_stats::FilteredReadStrategy::RowGroupPredicate,
2289 row_groups_considered: 1,
2290 row_groups_selected: 1,
2291 pages_considered: 1,
2292 pages_selected: 1,
2293 exact_rows_selected: 1,
2294 metadata_fallbacks: 0,
2295 validation_fallbacks: 0,
2296 });
2297 observation.complete(2, false);
2298 }
2299 {
2300 let erased: Arc<dyn crate::io_stats::FilteredReadObserver> = observer.clone();
2301 let _failed = FilteredReadObservation::new(Some(&erased), FilteredReadKind::Edge);
2302 }
2303 assert_eq!(
2304 observer.started.load(std::sync::atomic::Ordering::Relaxed),
2305 2
2306 );
2307 assert_eq!(
2308 observer.scanned.load(std::sync::atomic::Ordering::Relaxed),
2309 3
2310 );
2311 assert_eq!(
2312 observer
2313 .completed
2314 .load(std::sync::atomic::Ordering::Relaxed),
2315 2
2316 );
2317 assert_eq!(
2318 observer.failed.load(std::sync::atomic::Ordering::Relaxed),
2319 1
2320 );
2321 assert_eq!(
2322 observer.pruning.load(std::sync::atomic::Ordering::Relaxed),
2323 1
2324 );
2325 }
2326
2327 #[test]
2328 fn wave12_typed_relation_names_are_confined_to_one_plain_stem() {
2329 let dir = TempDir::new().unwrap();
2330 for invalid in ["../escape", "nested/name", "."] {
2331 assert!(read_edges(dir.path(), invalid, OntologyMode::Strict).is_err());
2332 assert!(
2333 read_edges_filtered(
2334 dir.path(),
2335 invalid,
2336 OntologyMode::Advisory,
2337 &[1].into_iter().collect(),
2338 )
2339 .is_err()
2340 );
2341 }
2342 }
2343
2344 #[test]
2345 fn wave12_max_edge_id_skips_non_parquet_and_corrupt_parquet_entries() {
2346 let dir = TempDir::new().unwrap();
2347 let edges = dir.path().join("topology/edges");
2348 std::fs::create_dir_all(&edges).unwrap();
2349 std::fs::write(edges.join("note.txt"), b"not parquet").unwrap();
2350 std::fs::write(edges.join("broken.parquet"), b"not parquet").unwrap();
2351 assert_eq!(max_edge_id(dir.path()).unwrap(), 0);
2352 }
2353}