Skip to main content

graphforge_storage/
catalog.rs

1//! DataFusion [`TableProvider`] and [`CatalogProvider`] implementations.
2//!
3//! Each GraphForge graph directory (`project/`) maps to a [`GraphCatalog`] which
4//! presents its Parquet files as DataFusion tables under the address
5//! `graph.graph.<table_name>`:
6//!
7//! | Table name | File | Schema |
8//! |---|---|---|
9//! | `topology_nodes` | `topology/nodes.parquet` | `TOPOLOGY_NODES_SCHEMA` |
10//! | `edges_TYPENAME` | `topology/edges/TYPENAME.parquet` | `TYPED_EDGE_SCHEMA` |
11//! | `edges__exploratory` | `topology/edges/_exploratory.parquet` | `EXPLORATORY_EDGE_SCHEMA` |
12//! | `properties_ENTITY` | `properties/ENTITY.parquet` | `property_schema(entity, defs)` |
13//!
14//! # Scan implementation
15//!
16//! Scans are implemented via DataFusion's [`MemTable`]: the Parquet file is read
17//! into memory at query time and wrapped in a `MemTable` which handles projection
18//! and filter application.  This is correct and simple for M12; lower-level
19//! pushdown can be added in a later milestone.
20
21use 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
48// ---------------------------------------------------------------------------
49// Parquet I/O helpers
50// ---------------------------------------------------------------------------
51
52fn 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
60/// Total rows across `batches`, for the [`io_stats`](crate::io_stats) counters.
61fn total_rows(batches: &[RecordBatch]) -> u64 {
62    u64::try_from(batches.iter().map(RecordBatch::num_rows).sum::<usize>()).unwrap_or(u64::MAX)
63}
64
65/// Read all row groups from a Parquet file into a single [`RecordBatch`].
66///
67/// Returns an empty batch (correct schema, zero rows) if the file does not exist.
68pub(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
88/// Normalize legacy scalar-label topology batches to the current multi-label
89/// schema. Existing `type_id` values become singleton `type_ids` lists.
90pub(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
127// ---------------------------------------------------------------------------
128// Direct readers (catalog-free)
129// ---------------------------------------------------------------------------
130//
131// Physical execution nodes (e.g. `VarLenExpandExec`, #580) need to read the
132// edge / node tables directly from the project directory: the DataFusion
133// `TaskContext` they execute with exposes neither the `GraphCatalog` nor the
134// project path, so the path is baked into the node at lowering time and the
135// node reads the Parquet itself.  These helpers expose the same on-disk layout
136// and schemas the `GraphWriter` produces, reusing [`read_parquet_or_empty`]
137// (which returns a correctly-typed empty batch when the file is absent).
138
139/// Read all edge rows for relation `rel_name` from the project at `dir`.
140///
141/// The on-disk layout mirrors [`GraphWriter`](crate::GraphWriter):
142/// - **Strict / Advisory** — typed edges in `topology/edges/<rel_name>.parquet`
143///   ([`TYPED_EDGE_SCHEMA`]).
144/// - **Exploratory** — all edges in `topology/edges/_exploratory.parquet`
145///   ([`EXPLORATORY_EDGE_SCHEMA`], carrying a `rel_type_name` column).  The
146///   returned batch is **not** filtered by `rel_name`; callers that need a
147///   single relation must filter on `rel_type_name` themselves.
148///
149/// Returns a single (possibly empty) [`RecordBatch`] with the mode-appropriate
150/// schema; a missing file yields an empty batch rather than an error.
151///
152/// # Errors
153/// Returns [`DataFusionError::Execution`] if `rel_name` is not a plain file
154/// stem (contains path separators or `..`), and propagates Parquet / Arrow
155/// errors encountered while reading.
156pub fn read_edges(
157    dir: &Path,
158    rel_name: &str,
159    mode: OntologyMode,
160) -> Result<Vec<RecordBatch>, DataFusionError> {
161    // Untyped wildcard in a typed project (#823): `"*"` means "all relation
162    // types", served as a union over every edge file rather than a literal
163    // (nonexistent) `*.parquet`. Exploratory already reads the shared file.
164    if rel_name == "*" && matches!(mode, OntologyMode::Advisory | OntologyMode::Strict) {
165        return read_edges_union(dir, None, None);
166    }
167    // `rel_name` becomes a path component in Strict/Advisory mode; require a
168    // single plain file stem so it can't traverse outside `topology/edges/`
169    // (rejects path separators, `..`, absolute prefixes, and empty names).
170    // Exploratory mode uses a fixed stem, so the caller-supplied name never
171    // reaches the filesystem there.
172    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/// Like [`read_edges`] but returns only rows whose `edge_id` is in
196/// `edge_ids` — the traversal's lazy edge-record read (#830): on an adjacency
197/// Hit, only the traversed edges' records are needed, not the whole file.
198///
199/// Two pruning layers before decode:
200/// 1. **Row groups** whose `edge_id` min/max statistics cannot contain any
201///    requested id are skipped entirely (edge files are globally
202///    edge_id-ascending, so groups partition the id range).
203/// 2. A Parquet **row filter** on `edge_id` within surviving groups (with the
204///    page index enabled when present, this also skips whole pages).
205///
206/// Short-circuits: an empty `edge_ids` never opens the file (one empty batch);
207/// a requested set covering more than half the file falls back to the plain
208/// full read (the filter would cost more than it saves). Contract parity with
209/// [`read_edges`]: always at least one (possibly empty) batch with the
210/// mode-appropriate schema; a missing file yields an empty batch.
211///
212/// # Errors
213/// Same as [`read_edges`], plus Parquet filter construction failures.
214#[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/// [`read_edges_filtered`] with optional aggregate-only operator attribution.
225#[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    // Untyped wildcard union (#823) — the lazy #709 read over all relations.
235    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
266/// Read the union of every relation's edges (#823): the "all relation types"
267/// read for an untyped traversal in a typed project. Enumerates every
268/// `topology/edges/*.parquet` (stem order, for deterministic adjacency/BFS),
269/// reads each (filtered to `edge_ids` when given — the lazy #709 read), and
270/// normalizes every batch to [`EXPLORATORY_EDGE_SCHEMA`] by tagging a typed
271/// file's rows with `rel_type_name = <file stem>` (a file already carrying the
272/// column — a stray `_exploratory.parquet` — passes through). Always returns at
273/// least one (possibly empty) `EXPLORATORY_EDGE_SCHEMA` batch.
274fn 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
316/// Normalize a typed-edge batch to [`EXPLORATORY_EDGE_SCHEMA`] by appending a
317/// constant `rel_type_name = stem` column. A batch already carrying the column
318/// (an `_exploratory` file) is returned unchanged.
319fn 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/// Which [`io_stats`](crate::io_stats) counters a filtered read attributes to —
331/// `read_parquet_filtered_u64` is keyed generically but the counters are split
332/// by table so the benchmark can prove edge *and* node reads are
333/// neighborhood-proportional.
334#[derive(Clone, Copy, PartialEq, Eq)]
335enum FilteredReadKind {
336    Edge,
337    Node,
338}
339
340/// Ensures every observed read has exactly one terminal completion/failure
341/// event, including errors from reader construction and batch decoding.
342struct 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
403/// Exact row selection for the canonical dense `node_id == row_ordinal + 1`
404/// layout. The selection is relative to the concatenation of `row_groups`, as
405/// required by Parquet after row-group filtering.
406struct 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
421/// Prove the canonical dense node layout from row-group and page metadata.
422fn 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
523/// Map requested ids to exact row ordinals after proving the canonical dense
524/// layout. Any incomplete or surprising metadata fails closed to the
525/// conservative predicate path.
526fn 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/// Read a Parquet file keeping only rows whose `key_column` (UInt64) value is
618/// in `ids`, with row-group pruning on the column's min/max statistics. `kind`
619/// selects the [`io_stats`](crate::io_stats) counters: the >50% fallback is a
620/// full read, the pushdown path a filtered read, of the named table.
621#[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    // Empty request or missing file: never open / one empty batch (contract
631    // parity with `read_parquet_or_empty`).
632    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    // Optional, NOT required: with_page_index(true) errors on files lacking a
658    // page index; Optional enables page-level skipping when one is present.
659    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    // Fallback: a large requested fraction makes the Parquet-level filter
664    // overhead a net loss — read plainly, then trim in memory so the public
665    // "only the requested ids" contract still holds.
666    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        // The fallback scanned the whole file before trimming, so record it as
672        // a full read (its row count is the full file, not the trimmed result):
673        // a fallback must not masquerade as a cheap filtered read.
674        let scanned = total_rows(&batches);
675        record_full(kind, scanned);
676        observation.scanned(scanned);
677        // Resolve the key column against the batches' ACTUAL schema (the
678        // on-disk file's), not `fallback_schema`: a column-shifted file would
679        // otherwise turn the index lookup into an out-of-bounds panic at
680        // `RecordBatch::column` instead of a graceful error.
681        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    // The key column's leaf index (flat schemas: leaf index == field index).
729    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    // Exact ordinal selection is node-only. Edges and noncanonical node files
748    // retain the conservative row-group min/max behavior.
749    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        // Missing or non-Int64 statistics keep the group (never prune blind).
767        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    // 2) Row filter on the key column within surviving groups.
804    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        // The predicate only sees rows in pages the page index did not skip, so
817        // this counts the decode-cost footprint (#838): flat for a clustered id
818        // set, ~whole-file for a scattered one.
819        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    // Predicate-pushdown path: record the rows actually materialized after
842    // row-group + row-filter pruning (the neighborhood-proportional cost #767
843    // measures), whether or not any survived.
844    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
876/// Attribute a full read of `rows` to the table named by `kind`.
877fn 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
884/// Attribute a filtered (predicate-pushdown) read of `rows` to `kind`'s table.
885fn 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
892/// Attribute `rows` evaluated by the pushdown predicate (the decode footprint
893/// after page-index skipping) to `kind`'s table.
894fn 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
901/// Record aggregate pruning work globally and, when installed, against the
902/// calling physical hop.
903fn 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
914/// Read all node rows from `topology/nodes.parquet` in the project at `dir`.
915///
916/// Returns a single (possibly empty) [`RecordBatch`] with
917/// [`TOPOLOGY_NODES_SCHEMA`]; a missing file yields an empty batch.
918///
919/// # Errors
920/// Propagates Parquet / Arrow errors encountered while reading.
921pub 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/// Like [`read_nodes`] but returns only rows whose `node_id` is in `node_ids` —
930/// the traversal's lazy node-record read (#838): on an adjacency Hit only the
931/// reached destination nodes' records are needed to project the destination
932/// columns, not the whole node table. Canonical dense files use exact physical
933/// row selection; legacy, gapped, or noncanonical files retain conservative
934/// row-group pruning plus a membership predicate.
935///
936/// Contract parity with [`read_nodes`]: always at least one (possibly empty)
937/// batch with [`TOPOLOGY_NODES_SCHEMA`]; an empty `node_ids` or a missing file
938/// never opens the file.
939///
940/// # Errors
941/// Same as [`read_nodes`], plus Parquet filter construction failures.
942#[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/// [`read_nodes_filtered`] with optional aggregate-only operator attribution.
951#[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
969/// Return the largest `edge_id` surrogate across every edge file under
970/// `topology/edges/` (both typed `<rel>.parquet` and `_exploratory.parquet`),
971/// or `0` if there are no edge files yet.
972///
973/// Used by [`GraphWriter`](crate::GraphWriter) to continue surrogate assignment
974/// from the on-disk maximum when appending across separate write sessions.
975///
976/// # Errors
977/// Propagates Parquet / Arrow errors encountered while reading an edge file.
978pub(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        // The `edge_id` column exists in both the typed and exploratory edge
995        // schemas, so the on-disk file schema (discovered) reads it either way.
996        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
1014/// Read `properties/<stem>.parquet` for the project at `dir`, discovering its
1015/// (dynamic) schema from the file. Returns an **empty `Vec`** when the file is
1016/// absent — so a caller decoding rows sees zero pre-existing property rows.
1017///
1018/// Production reads go through the staged-batch read-through in `writer`
1019/// (#792); the node-hydration path (`nodes(p)`, #1024) also reads through it.
1020///
1021/// # Errors
1022/// Propagates Parquet / Arrow errors encountered while reading.
1023pub 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
1031/// Edge analogue of [`read_properties`]: read `edge_properties/<stem>.parquet`
1032/// (keyed by `edge_uuid`), discovering its dynamic schema from the file. Returns
1033/// an **empty `Vec`** when the file is absent.
1034///
1035/// # Errors
1036/// Propagates Parquet / Arrow errors encountered while reading.
1037pub 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/// Stems (relation names) of every `edge_properties/<stem>.parquet` under
1046/// `dir`, **sorted** so schema unions built from them are deterministic
1047/// (#1023). Empty when the directory is absent — a project with no persisted
1048/// edge properties.
1049#[must_use]
1050pub fn list_edge_property_stems(dir: &Path) -> Vec<String> {
1051    list_parquet_stems(&dir.join("edge_properties"))
1052}
1053
1054/// Stems (entity type names, or `_untyped`) of every
1055/// `properties/<stem>.parquet` under `dir`, **sorted** for deterministic
1056/// schema unions (#1024). Empty when the directory is absent.
1057#[must_use]
1058pub fn list_property_stems(dir: &Path) -> Vec<String> {
1059    list_parquet_stems(&dir.join("properties"))
1060}
1061
1062/// Sorted `<stem>` names of the `<stem>.parquet` files directly under `dir`.
1063fn 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// ---------------------------------------------------------------------------
1081// TopologyNodeTable
1082// ---------------------------------------------------------------------------
1083
1084/// [`TableProvider`] for `topology/nodes.parquet`.
1085#[derive(Debug, Clone)]
1086pub struct TopologyNodeTable {
1087    path: PathBuf,
1088}
1089
1090impl TopologyNodeTable {
1091    /// Create a table backed by the given Parquet file path.
1092    #[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// ---------------------------------------------------------------------------
1129// TypedEdgeTable
1130// ---------------------------------------------------------------------------
1131
1132/// [`TableProvider`] for `topology/edges/TYPENAME.parquet`.
1133#[derive(Debug, Clone)]
1134pub struct TypedEdgeTable {
1135    path: PathBuf,
1136    schema: SchemaRef,
1137}
1138
1139impl TypedEdgeTable {
1140    /// Open the edge table for `rel_type_name` inside `dir`.
1141    ///
1142    /// - `"_exploratory"` → schema includes `rel_type_name` column
1143    /// - any other name → [`TYPED_EDGE_SCHEMA`]
1144    #[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// ---------------------------------------------------------------------------
1187// UnionEdgeTable
1188// ---------------------------------------------------------------------------
1189
1190/// [`TableProvider`] over the union of every relation's edge file (#823) — the
1191/// scan source for an **untyped** single-hop pattern (`(a)-[]->(b)`) in a typed
1192/// project, where the `_exploratory` table does not exist. Materializes
1193/// [`read_edges_union`] into a [`MemTable`]; the schema is always
1194/// [`EXPLORATORY_EDGE_SCHEMA`] (each row tagged with its source relation), so it
1195/// is a drop-in for the exploratory edge scan the untyped lowering already uses.
1196#[derive(Debug, Clone)]
1197pub struct UnionEdgeTable {
1198    dir: PathBuf,
1199}
1200
1201impl UnionEdgeTable {
1202    /// Open a union edge table over `dir`'s `topology/edges/`.
1203    #[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// ---------------------------------------------------------------------------
1239// PropertyTable
1240// ---------------------------------------------------------------------------
1241
1242/// [`TableProvider`] for `properties/ENTITY_TYPE.parquet`.
1243#[derive(Debug, Clone)]
1244pub struct PropertyTable {
1245    path: PathBuf,
1246    schema: SchemaRef,
1247}
1248
1249impl PropertyTable {
1250    /// Open a property table.
1251    ///
1252    /// If the file does not yet exist scans return an empty batch with the
1253    /// correct schema.
1254    #[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    /// Open a property table for `stem` (an entity type name, or `"_untyped"`),
1263    /// discovering the column schema from the Parquet file on disk.
1264    ///
1265    /// The exploratory `_untyped.parquet` file's schema is inferred at write
1266    /// time from the observed property literals, so the read path cannot know it
1267    /// statically — it must be read back from the file. When the file does not
1268    /// exist yet, falls back to [`PROPERTY_BASE_SCHEMA`] (just `node_uuid`), so a
1269    /// join against an as-yet-unwritten property table yields zero property rows
1270    /// rather than an error.
1271    #[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    /// The property column schema (including the `node_uuid` join key).
1280    #[must_use]
1281    pub fn schema_ref(&self) -> SchemaRef {
1282        self.schema.clone()
1283    }
1284}
1285
1286// ---------------------------------------------------------------------------
1287// EdgePropertyTable
1288// ---------------------------------------------------------------------------
1289
1290/// [`TableProvider`] for `edge_properties/REL_TYPE.parquet` (#784).
1291///
1292/// The edge analogue of [`PropertyTable`], keyed by `edge_uuid` and read from
1293/// the dedicated `edge_properties/` directory so a relation type cannot collide
1294/// with a same-named node label under `properties/`.
1295#[derive(Debug, Clone)]
1296pub struct EdgePropertyTable {
1297    path: PathBuf,
1298    schema: SchemaRef,
1299}
1300
1301impl EdgePropertyTable {
1302    /// Open an edge-property table for `rel_type`, discovering the column schema
1303    /// from the Parquet file on disk.
1304    ///
1305    /// The per-relation schema is inferred at write time from the observed
1306    /// property literals, so the read path reads it back from the file. When the
1307    /// file does not exist yet, falls back to [`EDGE_PROPERTY_BASE_SCHEMA`] (just
1308    /// `edge_uuid`), so a join against an as-yet-unwritten edge-property table
1309    /// yields zero property rows rather than an error.
1310    #[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    /// The property column schema (including the `edge_uuid` join key).
1321    #[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
1354/// Read just the Arrow schema of a Parquet file, or `None` if it is absent or
1355/// unreadable.
1356pub(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
1389// ---------------------------------------------------------------------------
1390// GraphSchema — inner schema provider
1391// ---------------------------------------------------------------------------
1392
1393struct 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
1438// ---------------------------------------------------------------------------
1439// GraphCatalog
1440// ---------------------------------------------------------------------------
1441
1442/// DataFusion [`CatalogProvider`] for a GraphForge project directory.
1443///
1444/// Exposes topology, edge, and property tables under the `"graph"` schema.
1445/// Construct via [`GraphCatalog::open`].
1446pub struct GraphCatalog {
1447    schema: Arc<GraphSchema>,
1448    /// Reverse map `PropId.0` → property name, merged from the ontology and the
1449    /// runtime catalog at [`open`](Self::open) time. The relational lowering
1450    /// layer borrows it to resolve numeric `PropertyAccess` IDs to real column
1451    /// names without re-plumbing the ontology/runtime catalog separately.
1452    prop_names: HashMap<u32, String>,
1453    /// Reverse map `TypeId.0` → relation-type name, merged from the ontology and
1454    /// the runtime catalog. Lets the lowering layer resolve a `TypedEdgeScan`'s
1455    /// relation name in exploratory mode (where the ontology map is empty).
1456    rel_names: HashMap<u32, String>,
1457    /// Reverse map `TypeId.0` → entity-type (node label) name, from the runtime
1458    /// catalog. Lets the lowering layer render a real label for an unlabelled
1459    /// `MATCH (n) RETURN n` in exploratory mode (where the ontology map is
1460    /// empty) by resolving the node's stored `type_id` (#889).
1461    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    /// Open a GraphForge project directory as a DataFusion catalog.
1474    ///
1475    /// - `dir`: project root (contains `topology/`, `properties/`, etc.)
1476    /// - `ontology`: compiled ontology, or `None` in exploratory mode
1477    /// - `runtime_catalog`: runtime type catalog for exploratory mode
1478    ///
1479    /// # Errors
1480    ///
1481    /// Propagates I/O errors encountered while registering tables.
1482    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        // ---- topology nodes ----
1490        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        // ---- typed edge tables ----
1497        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            // No ontology. Exploratory-written edges all land in the single
1506            // `_exploratory.parquet` (tagged with `rel_type_name`); register that
1507            // catch-all. For runtime relation types, register a per-relation
1508            // `edges_<rel>` table ONLY when its typed file actually exists on disk
1509            // (e.g. data written in strict/advisory mode then reloaded with just a
1510            // runtime catalog). Registering `edges_<rel>` for a relation whose
1511            // data is really in `_exploratory.parquet` would make the read path
1512            // scan a non-existent typed file and return 0 rows.
1513            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        // Always register the exploratory fallback (advisory mode uses it too).
1532        if !schema.table_exist("edges__exploratory") {
1533            schema.register(
1534                "edges__exploratory",
1535                Arc::new(TypedEdgeTable::open(dir, "_exploratory")),
1536            );
1537        }
1538
1539        // ---- property tables ----
1540        register_property_tables(dir, ontology, &mut schema);
1541
1542        // ---- name maps (for read-path property + relation resolution) ----
1543        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    /// Reverse map `PropId.0` → property name (ontology + runtime catalog),
1556    /// used by the relational lowering layer to resolve `PropertyAccess`.
1557    #[must_use]
1558    pub fn prop_names(&self) -> &HashMap<u32, String> {
1559        &self.prop_names
1560    }
1561
1562    /// Reverse map `RuntimeTypeId.0` → relation-type name from the runtime
1563    /// catalog. The relational lowerer tags these keys before merging them with
1564    /// ontology TypeIds so the two zero-based ID spaces cannot collide.
1565    #[must_use]
1566    pub fn rel_names(&self) -> &HashMap<u32, String> {
1567        &self.rel_names
1568    }
1569
1570    /// Reverse map `TypeId.0` → entity-type (node label) name, from the runtime
1571    /// catalog. Used by the relational lowering layer to render a node value's
1572    /// label for an unlabelled match — including in exploratory mode, where the
1573    /// ontology map is empty (#889).
1574    #[must_use]
1575    pub fn label_names(&self) -> &HashMap<u32, String> {
1576        &self.label_names
1577    }
1578}
1579
1580/// Build the `PropId.0 → name` map for resolving `PropertyAccess`.
1581///
1582/// The binder interns every observed property into the [`RuntimeCatalog`] and
1583/// emits its runtime `PropId` (it does **not** emit ontology property IDs — they
1584/// live in a separate ID space). So the runtime catalog is the single
1585/// authoritative source for `PropId → name`, in all ontology modes.
1586fn 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
1596/// Build the `TypeId.0 → relation-name` map from the runtime catalog.
1597///
1598/// Only the runtime-catalog side is needed here: the binder resolves relation
1599/// types ontology-first (so an ontology-sourced `TypeId` is already covered by
1600/// the lowerer's ontology map) and falls back to the `RuntimeCatalog` only in
1601/// exploratory mode (or advisory misses) — the case this map fills.
1602fn 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
1609/// Build the `TypeId.0 → label-name` map from the runtime catalog.
1610///
1611/// As with [`build_rel_names`], only the runtime-catalog side is needed: the
1612/// binder resolves labels ontology-first, so ontology-sourced label `TypeId`s
1613/// are already covered by the lowerer's ontology map; this fills the exploratory
1614/// case (no ontology), letting an unlabelled `MATCH (n) RETURN n` recover the
1615/// node's label name from its stored `type_id` (#889).
1616fn 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
1641// ---------------------------------------------------------------------------
1642// Property table registration helper
1643// ---------------------------------------------------------------------------
1644
1645fn 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        // Exploratory: properties are written to a single `_untyped.parquet`
1660        // whose column schema is inferred at write time, so register it with the
1661        // schema discovered from disk (just `node_uuid` until it is written).
1662        schema.register(
1663            "properties__untyped",
1664            Arc::new(PropertyTable::open_discovered(dir, "_untyped")),
1665        );
1666    }
1667}
1668
1669// ---------------------------------------------------------------------------
1670// Tests
1671// ---------------------------------------------------------------------------
1672
1673#[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    // -----------------------------------------------------------------------
1992    // Direct readers (read_edges / read_nodes) — #580
1993    // -----------------------------------------------------------------------
1994
1995    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        // Strict mode reads the typed schema (no rel_type_name column).
2012        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        // Exploratory mode uses a fixed stem, so a traversal-looking rel_name is
2032        // harmless (never reaches the path) — it must NOT error.
2033        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        // No edge file written.
2040        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        // Exploratory edges live in `_exploratory.parquet`; a typed `KNOWS.parquet`
2049        // must be ignored in this mode.
2050        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        // The exploratory file does not exist → empty batch with the exploratory
2059        // schema (which carries rel_type_name), NOT the typed KNOWS rows.
2060        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    // -----------------------------------------------------------------------
2066    // Untyped wildcard union read (#823)
2067    // -----------------------------------------------------------------------
2068
2069    /// Write a one-row typed edge file with the given surrogates.
2070    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        // A 16-byte uuid from an id, non-panicking for any u64 (the tests only
2074        // assert on edge_id/rel_type_name, but keep the helper id-range-safe).
2075        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    /// `(edge_id, rel_type_name)` pairs from EXPLORATORY-schema batches.
2102    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        // Union schema carries rel_type_name, each row tagged with its file stem.
2130        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        // Missing file first → empty batch with the topology schema.
2170        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        // Then write one node and read it back.
2175        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}