Skip to main content

graphforge_storage/
graph_projection.rs

1//! Deterministic graph-only workspace projections.
2//!
3//! This module deliberately understands only graph-owned workspace files. It
4//! never enumerates or copies project-generation participants, provenance,
5//! knowledge, epistemic, valid-time, search, or derived-index directories.
6
7use std::collections::{BTreeMap, BTreeSet};
8use std::fs::{self, File};
9use std::path::{Path, PathBuf};
10use std::sync::Arc;
11
12use arrow::array::{
13    Array, ArrayRef, BinaryArray, BooleanArray, FixedSizeBinaryArray, FixedSizeListArray,
14    Float32Array, Float64Array, Int32Array, Int64Array, LargeBinaryArray, LargeListArray,
15    LargeStringArray, ListArray, StringArray, StructArray, UInt32Array, UInt64Array,
16};
17use arrow::compute::{concat_batches, take};
18use arrow::datatypes::{DataType, Field, Schema, TimeUnit};
19use arrow::record_batch::RecordBatch;
20use graphforge_core::GfError;
21use graphforge_core::canonical::{
22    CANONICAL_CONTRACT_VERSION, CanonicalDomain, CanonicalWriter, fingerprint,
23};
24use parquet::arrow::ArrowWriter;
25
26type GraphUuid = [u8; 16];
27type EdgeEndpoints = BTreeMap<GraphUuid, (GraphUuid, GraphUuid)>;
28
29/// Explicit graph identities requested for one projection.
30#[derive(Clone, Debug, Default, Eq, PartialEq)]
31pub struct GraphProjectionSelection {
32    /// Explicit node UUIDs. Selected edges add their endpoints automatically.
33    pub node_uuids: BTreeSet<[u8; 16]>,
34    /// Explicit edge UUIDs. No other edge is induced from selected nodes.
35    pub edge_uuids: BTreeSet<[u8; 16]>,
36}
37
38/// Exact identities materialized into the graph-only target workspace.
39#[derive(Clone, Debug, Eq, PartialEq)]
40pub struct GraphProjectionSummary {
41    /// Canonically ordered node UUIDs, including edge endpoint closure.
42    pub node_uuids: Vec<[u8; 16]>,
43    /// Canonically ordered explicitly selected edge UUIDs.
44    pub edge_uuids: Vec<[u8; 16]>,
45    /// Endpoint UUIDs added beyond the caller's explicit node set.
46    pub endpoint_node_uuids: Vec<[u8; 16]>,
47    /// Domain-separated digest of canonical logical graph tables.
48    pub graph_content_fingerprint: [u8; 32],
49}
50
51/// Materialize one deterministic graph-only workspace projection.
52///
53/// The target must be absent or empty. Node/edge UUIDs and surrogate IDs are
54/// preserved. Every selected edge adds both endpoint nodes; node selection
55/// never induces an edge. Topology and property rows are rewritten in UUID
56/// order. Only core graph files and required ontology/runtime-catalog metadata
57/// are copied; derived indexes and every non-graph domain are excluded.
58///
59/// # Errors
60/// Returns validation for missing identities, unsafe/non-empty targets, or
61/// malformed graph files, and storage errors for I/O/Arrow/Parquet failures.
62pub fn materialize_graph_projection(
63    source: &Path,
64    target: &Path,
65    selection: &GraphProjectionSelection,
66) -> Result<GraphProjectionSummary, GfError> {
67    validate_distinct_paths(source, target)?;
68    validate_graph_empty_target(target)?;
69
70    let nodes_path = source.join("topology/nodes.parquet");
71    let node_ids = uuid_rows(&nodes_path, "node_uuid")?;
72    require_present(&selection.node_uuids, &node_ids, "node")?;
73
74    let edge_files = sorted_parquet_files(&source.join("topology/edges"))?;
75    let edges = edge_endpoints(&edge_files)?;
76    let edge_ids = edges.keys().copied().collect::<BTreeSet<_>>();
77    require_present(&selection.edge_uuids, &edge_ids, "edge")?;
78
79    let mut selected_nodes = selection.node_uuids.clone();
80    for edge_uuid in &selection.edge_uuids {
81        let (src, dst) = edges
82            .get(edge_uuid)
83            .expect("selected edge presence was validated");
84        if !node_ids.contains(src) || !node_ids.contains(dst) {
85            return Err(validation(
86                "selected edge references a missing endpoint node",
87            ));
88        }
89        selected_nodes.insert(*src);
90        selected_nodes.insert(*dst);
91    }
92    let endpoint_node_uuids = selected_nodes
93        .difference(&selection.node_uuids)
94        .copied()
95        .collect::<Vec<_>>();
96
97    clear_graph_empty_target(target)?;
98    fs::create_dir_all(target).map_err(storage)?;
99    project_parquet_file(
100        &nodes_path,
101        &target.join("topology/nodes.parquet"),
102        "node_uuid",
103        &selected_nodes,
104    )?;
105    project_parquet_directory(
106        &source.join("topology/edges"),
107        &target.join("topology/edges"),
108        "edge_uuid",
109        &selection.edge_uuids,
110    )?;
111    project_parquet_directory(
112        &source.join("properties"),
113        &target.join("properties"),
114        "node_uuid",
115        &selected_nodes,
116    )?;
117    project_parquet_directory(
118        &source.join("edge_properties"),
119        &target.join("edge_properties"),
120        "edge_uuid",
121        &selection.edge_uuids,
122    )?;
123    copy_runtime_catalog(source, target)?;
124    for file in [
125        graphforge_core::manifest::MANIFEST_FILE,
126        graphforge_core::manifest::ONTOLOGY_FILE,
127    ] {
128        copy_regular_file_if_present(&source.join(file), &target.join(file))?;
129    }
130    let graph_content_fingerprint = projected_graph_fingerprint(target)?;
131
132    Ok(GraphProjectionSummary {
133        node_uuids: selected_nodes.into_iter().collect(),
134        edge_uuids: selection.edge_uuids.iter().copied().collect(),
135        endpoint_node_uuids,
136        graph_content_fingerprint,
137    })
138}
139
140fn edge_endpoints(files: &[PathBuf]) -> Result<EdgeEndpoints, GfError> {
141    let mut edges = BTreeMap::new();
142    for path in files {
143        for batch in read_parquet(path)? {
144            let edge_ids = uuid_column(&batch, "edge_uuid")?;
145            let sources = uuid_column(&batch, "src_uuid")?;
146            let targets = uuid_column(&batch, "dst_uuid")?;
147            for row in 0..batch.num_rows() {
148                let edge_uuid = uuid_at(edge_ids, row)?;
149                let endpoints = (uuid_at(sources, row)?, uuid_at(targets, row)?);
150                if edges.insert(edge_uuid, endpoints).is_some() {
151                    return Err(validation("graph contains a duplicate edge UUID"));
152                }
153            }
154        }
155    }
156    Ok(edges)
157}
158
159fn uuid_rows(path: &Path, column: &str) -> Result<BTreeSet<GraphUuid>, GfError> {
160    if !path.exists() {
161        return Ok(BTreeSet::new());
162    }
163    let mut rows = BTreeSet::new();
164    for batch in read_parquet(path)? {
165        let values = uuid_column(&batch, column)?;
166        for row in 0..batch.num_rows() {
167            if !rows.insert(uuid_at(values, row)?) {
168                return Err(validation(format!("graph contains a duplicate {column}")));
169            }
170        }
171    }
172    Ok(rows)
173}
174
175fn project_parquet_directory(
176    source: &Path,
177    target: &Path,
178    key: &str,
179    selected: &BTreeSet<[u8; 16]>,
180) -> Result<(), GfError> {
181    for path in sorted_parquet_files(source)? {
182        let name = path
183            .file_name()
184            .ok_or_else(|| validation("graph parquet path has no file name"))?;
185        project_parquet_file(&path, &target.join(name), key, selected)?;
186    }
187    Ok(())
188}
189
190fn project_parquet_file(
191    source: &Path,
192    target: &Path,
193    key: &str,
194    selected: &BTreeSet<[u8; 16]>,
195) -> Result<(), GfError> {
196    if !source.exists() {
197        return Ok(());
198    }
199    let batches = read_parquet(source)?;
200    let schema = batches
201        .first()
202        .map(RecordBatch::schema)
203        .or_else(|| crate::catalog::discover_parquet_schema(source))
204        .ok_or_else(|| validation("graph parquet schema is unavailable"))?;
205    let combined = if batches.is_empty() {
206        RecordBatch::new_empty(Arc::clone(&schema))
207    } else {
208        concat_batches(&schema, &batches).map_err(storage)?
209    };
210    let keys = uuid_column(&combined, key)?;
211    let mut rows = Vec::new();
212    for row in 0..combined.num_rows() {
213        let uuid = uuid_at(keys, row)?;
214        if selected.contains(&uuid) {
215            rows.push((uuid, row));
216        }
217    }
218    rows.sort_unstable();
219    let indices = rows
220        .into_iter()
221        .map(|(_, row)| {
222            u32::try_from(row).map_err(|_| validation("graph projection row index exceeds UInt32"))
223        })
224        .collect::<Result<Vec<_>, _>>()?;
225    let indices = UInt32Array::from(indices);
226    let columns = combined
227        .columns()
228        .iter()
229        .map(|column| take(column.as_ref(), &indices, None).map_err(storage))
230        .collect::<Result<Vec<_>, _>>()?;
231    let projected = RecordBatch::try_new(Arc::clone(&schema), columns).map_err(storage)?;
232    write_parquet(target, &projected)
233}
234
235fn copy_runtime_catalog(source: &Path, target: &Path) -> Result<(), GfError> {
236    let source = source.join("topology/runtime_catalog.parquet");
237    if !source.exists() {
238        return Ok(());
239    }
240    let batches = read_parquet(&source)?;
241    let schema = batches
242        .first()
243        .map(RecordBatch::schema)
244        .ok_or_else(|| validation("runtime catalog has no schema"))?;
245    let batch = concat_batches(&schema, &batches).map_err(storage)?;
246    let canonical = graphforge_ir::RuntimeCatalog::from_record_batch(&batch)?.to_record_batch();
247    let selected = selected_catalog_rows(target, &canonical)?;
248    let indices = UInt32Array::from(
249        selected
250            .into_iter()
251            .map(|row| {
252                u32::try_from(row)
253                    .map_err(|_| validation("runtime catalog row index exceeds UInt32"))
254            })
255            .collect::<Result<Vec<_>, _>>()?,
256    );
257    let columns = canonical
258        .columns()
259        .iter()
260        .map(|column| take(column.as_ref(), &indices, None).map_err(storage))
261        .collect::<Result<Vec<_>, _>>()?;
262    let canonical = RecordBatch::try_new(canonical.schema(), columns).map_err(storage)?;
263    write_parquet(&target.join("topology/runtime_catalog.parquet"), &canonical)
264}
265
266#[allow(
267    clippy::too_many_lines,
268    reason = "catalog dependency closure remains one auditable selection pass"
269)]
270fn selected_catalog_rows(target: &Path, catalog: &RecordBatch) -> Result<Vec<usize>, GfError> {
271    let mut type_ids = BTreeSet::new();
272    let nodes = target.join("topology/nodes.parquet");
273    if nodes.exists() {
274        for batch in read_parquet(&nodes)? {
275            if let Some(column) = batch.column_by_name("type_id") {
276                let values = column
277                    .as_any()
278                    .downcast_ref::<UInt32Array>()
279                    .ok_or_else(|| validation("node type_id is not UInt32"))?;
280                for row in 0..values.len() {
281                    if !values.is_null(row) {
282                        type_ids.insert(values.value(row));
283                    }
284                }
285            }
286            if let Some(column) = batch.column_by_name("type_ids") {
287                let lists = column
288                    .as_any()
289                    .downcast_ref::<ListArray>()
290                    .ok_or_else(|| validation("node type_ids is not List"))?;
291                for row in 0..lists.len() {
292                    if lists.is_null(row) {
293                        continue;
294                    }
295                    let values = lists.value(row);
296                    let values = values
297                        .as_any()
298                        .downcast_ref::<UInt32Array>()
299                        .ok_or_else(|| validation("node type_ids values are not UInt32"))?;
300                    type_ids.extend(values.values().iter().copied());
301                }
302            }
303        }
304    }
305
306    let mut relation_names = BTreeSet::new();
307    for path in sorted_parquet_files(&target.join("topology/edges"))? {
308        let stem = parquet_stem(&path)?;
309        if stem != "_exploratory" {
310            relation_names.insert(stem);
311        }
312        for batch in read_parquet(&path)? {
313            if let Some(column) = batch.column_by_name("rel_type_name") {
314                let values = column
315                    .as_any()
316                    .downcast_ref::<StringArray>()
317                    .ok_or_else(|| validation("edge rel_type_name is not Utf8"))?;
318                for row in 0..values.len() {
319                    if !values.is_null(row) {
320                        relation_names.insert(values.value(row).to_owned());
321                    }
322                }
323            }
324        }
325    }
326
327    let mut property_names = BTreeSet::new();
328    for directory in ["properties", "edge_properties"] {
329        for path in sorted_parquet_files(&target.join(directory))? {
330            let batches = read_parquet(&path)?;
331            if batches.iter().all(|batch| batch.num_rows() == 0) {
332                continue;
333            }
334            let schema = batches
335                .first()
336                .map(RecordBatch::schema)
337                .ok_or_else(|| validation("projected property table has no schema"))?;
338            for field in schema.fields() {
339                if !matches!(
340                    field.name().as_str(),
341                    "node_uuid" | "node_id" | "edge_uuid" | "edge_id"
342                ) {
343                    property_names.insert(field.name().clone());
344                }
345            }
346        }
347    }
348
349    let kinds = string_column(catalog, "entry_kind")?;
350    let names = string_column(catalog, "name")?;
351    let ids = catalog
352        .column_by_name("runtime_id")
353        .and_then(|column| column.as_any().downcast_ref::<UInt32Array>())
354        .ok_or_else(|| validation("runtime catalog runtime_id is not UInt32"))?;
355    let owners = string_column(catalog, "owner_label")?;
356
357    let mut active_owners = BTreeSet::new();
358    for row in 0..catalog.num_rows() {
359        if kinds.value(row) == "entity_type" && type_ids.contains(&ids.value(row)) {
360            active_owners.insert(names.value(row).to_owned());
361        }
362    }
363    active_owners.extend(relation_names.iter().cloned());
364
365    let mut selected = BTreeSet::new();
366    for row in 0..catalog.num_rows() {
367        let keep = match kinds.value(row) {
368            "entity_type" => type_ids.contains(&ids.value(row)),
369            "relation_type" => relation_names.contains(names.value(row)),
370            "property" => {
371                property_names.contains(names.value(row))
372                    && (owners.is_null(row) || active_owners.contains(owners.value(row)))
373            }
374            _ => false,
375        };
376        if keep {
377            selected.insert(row);
378            if kinds.value(row) == "property" && !owners.is_null(row) {
379                let owner = owners.value(row);
380                for owner_row in 0..catalog.num_rows() {
381                    if matches!(kinds.value(owner_row), "entity_type" | "relation_type")
382                        && names.value(owner_row) == owner
383                    {
384                        selected.insert(owner_row);
385                    }
386                }
387            }
388        }
389    }
390    Ok(selected.into_iter().collect())
391}
392
393fn parquet_stem(path: &Path) -> Result<String, GfError> {
394    path.file_stem()
395        .and_then(|value| value.to_str())
396        .map(str::to_owned)
397        .ok_or_else(|| validation("graph parquet path has no UTF-8 stem"))
398}
399
400fn string_column<'a>(batch: &'a RecordBatch, name: &str) -> Result<&'a StringArray, GfError> {
401    batch
402        .column_by_name(name)
403        .and_then(|column| column.as_any().downcast_ref::<StringArray>())
404        .ok_or_else(|| validation(format!("runtime catalog {name} is not Utf8")))
405}
406
407fn read_parquet(path: &Path) -> Result<Vec<RecordBatch>, GfError> {
408    let schema = crate::catalog::discover_parquet_schema(path).ok_or_else(|| {
409        validation(format!(
410            "cannot discover graph schema for {}",
411            path.display()
412        ))
413    })?;
414    crate::catalog::read_parquet_or_empty(path, schema)
415        .map_err(|error| GfError::Storage(error.to_string()))
416}
417
418fn write_parquet(path: &Path, batch: &RecordBatch) -> Result<(), GfError> {
419    let parent = path
420        .parent()
421        .ok_or_else(|| validation("graph parquet target has no parent"))?;
422    fs::create_dir_all(parent).map_err(storage)?;
423    let file = File::create(path).map_err(storage)?;
424    let mut writer = ArrowWriter::try_new(file, batch.schema(), None).map_err(storage)?;
425    writer.write(batch).map_err(storage)?;
426    writer.close().map_err(storage)?;
427    Ok(())
428}
429
430fn projected_graph_fingerprint(root: &Path) -> Result<[u8; 32], GfError> {
431    let mut paths = Vec::new();
432    let nodes = root.join("topology/nodes.parquet");
433    if nodes.exists() {
434        paths.push(nodes);
435    }
436    for directory in ["topology/edges", "properties", "edge_properties"] {
437        paths.extend(sorted_parquet_files(&root.join(directory))?);
438    }
439    let runtime_catalog = root.join("topology/runtime_catalog.parquet");
440    if runtime_catalog.exists() {
441        paths.push(runtime_catalog);
442    }
443    paths.sort();
444
445    let mut writer = CanonicalWriter::new();
446    writer.raw(b"GFGP1").map_err(canonical_error)?;
447    writer
448        .u32(exact_u32(paths.len(), "graph table count")?)
449        .map_err(canonical_error)?;
450    for path in paths {
451        let relative = path
452            .strip_prefix(root)
453            .map_err(|_| validation("graph projection path escaped target"))?
454            .to_str()
455            .ok_or_else(|| validation("graph projection path is not UTF-8"))?;
456        writer.text(relative).map_err(canonical_error)?;
457        let batches = read_parquet(&path)?;
458        let schema = batches
459            .first()
460            .map(RecordBatch::schema)
461            .ok_or_else(|| validation("graph projection table has no schema"))?;
462        let batch = concat_batches(&schema, &batches).map_err(storage)?;
463        let logical = logical_fingerprint_batch(relative, &batch)?;
464        encode_table(&mut writer, &logical)?;
465    }
466    fingerprint(
467        CanonicalDomain::GraphProjection,
468        CANONICAL_CONTRACT_VERSION,
469        &writer.finish(),
470    )
471    .map_err(canonical_error)
472}
473
474fn logical_fingerprint_batch(relative: &str, batch: &RecordBatch) -> Result<RecordBatch, GfError> {
475    let source_schema = batch.schema();
476    let names: Vec<&str> = if relative == "topology/nodes.parquet" {
477        vec!["node_uuid", "type_id", "type_ids"]
478    } else if relative.starts_with("topology/edges/") {
479        let mut names = vec!["edge_uuid", "src_uuid", "dst_uuid"];
480        if batch.column_by_name("rel_type_name").is_some() {
481            names.push("rel_type_name");
482        }
483        names
484    } else if relative == "topology/runtime_catalog.parquet" {
485        vec!["entry_kind", "name", "runtime_id", "owner_label"]
486    } else {
487        source_schema
488            .fields()
489            .iter()
490            .map(|field| field.name().as_str())
491            .collect()
492    };
493    let mut fields = Vec::with_capacity(names.len());
494    let mut columns = Vec::with_capacity(names.len());
495    for name in names {
496        let index = source_schema
497            .index_of(name)
498            .map_err(|_| validation(format!("graph fingerprint field {name} is absent")))?;
499        fields.push(Arc::clone(&source_schema.fields()[index]));
500        columns.push(Arc::clone(batch.column(index)));
501    }
502    let schema = Arc::new(Schema::new_with_metadata(
503        fields,
504        source_schema.metadata().clone(),
505    ));
506    RecordBatch::try_new(schema, columns).map_err(storage)
507}
508
509fn encode_table(writer: &mut CanonicalWriter, batch: &RecordBatch) -> Result<(), GfError> {
510    encode_schema(writer, batch.schema().as_ref())?;
511    writer
512        .u64(exact_u64(batch.num_rows(), "graph row count")?)
513        .map_err(canonical_error)?;
514    let schema = batch.schema();
515    let columns = schema
516        .fields()
517        .iter()
518        .zip(batch.columns())
519        .map(|(field, column)| {
520            let logical = dictionary_value_type(field.data_type());
521            if logical == field.data_type() {
522                Ok((logical, Arc::clone(column)))
523            } else {
524                arrow::compute::cast(column, logical)
525                    .map(|decoded| (logical, decoded))
526                    .map_err(storage)
527            }
528        })
529        .collect::<Result<Vec<_>, _>>()?;
530    for row in 0..batch.num_rows() {
531        for (field, (data_type, column)) in schema.fields().iter().zip(&columns) {
532            encode_value(writer, data_type, column, row, field.is_nullable())?;
533        }
534    }
535    Ok(())
536}
537
538fn encode_schema(writer: &mut CanonicalWriter, schema: &Schema) -> Result<(), GfError> {
539    writer.raw(b"GFS1").map_err(canonical_error)?;
540    writer
541        .u32(exact_u32(schema.fields().len(), "graph field count")?)
542        .map_err(canonical_error)?;
543    for field in &schema.fields {
544        encode_field(writer, field)?;
545    }
546    let ordered = schema.metadata().iter().collect::<BTreeMap<_, _>>();
547    writer
548        .u32(exact_u32(ordered.len(), "graph metadata count")?)
549        .map_err(canonical_error)?;
550    for (key, value) in ordered {
551        writer.text(key).map_err(canonical_error)?;
552        writer.text(value).map_err(canonical_error)?;
553    }
554    Ok(())
555}
556
557fn encode_field(writer: &mut CanonicalWriter, field: &Field) -> Result<(), GfError> {
558    writer.text(field.name()).map_err(canonical_error)?;
559    writer
560        .u8(u8::from(field.is_nullable()))
561        .map_err(canonical_error)?;
562    encode_type(writer, field.data_type())
563}
564
565fn encode_type(writer: &mut CanonicalWriter, data_type: &DataType) -> Result<(), GfError> {
566    match data_type {
567        DataType::Boolean => writer.u8(0x02),
568        DataType::Int32 => writer.u8(0x12),
569        DataType::Int64 => writer.u8(0x13),
570        DataType::UInt32 => writer.u8(0x16),
571        DataType::UInt64 => writer.u8(0x17),
572        DataType::Float32 => writer.u8(0x21),
573        DataType::Float64 => writer.u8(0x22),
574        DataType::Utf8 | DataType::LargeUtf8 => writer.u8(0x30),
575        DataType::Binary | DataType::LargeBinary => writer.u8(0x31),
576        DataType::FixedSizeBinary(width) => {
577            writer.u8(0x32).map_err(canonical_error)?;
578            writer.u32(
579                u32::try_from(*width)
580                    .map_err(|_| validation("negative fixed-size binary width"))?,
581            )
582        }
583        DataType::Timestamp(unit, timezone) => {
584            validate_timezone(timezone.as_deref())?;
585            writer.u8(0x52).map_err(canonical_error)?;
586            writer.u8(time_unit_tag(*unit))
587        }
588        DataType::Time64(unit) => {
589            writer.u8(0x53).map_err(canonical_error)?;
590            writer.u8(time_unit_tag(*unit))
591        }
592        DataType::List(field) | DataType::LargeList(field) => {
593            writer.u8(0x60).map_err(canonical_error)?;
594            encode_field(writer, field)?;
595            return Ok(());
596        }
597        DataType::FixedSizeList(field, length) => {
598            writer.u8(0x61).map_err(canonical_error)?;
599            writer
600                .u32(u32::try_from(*length).map_err(|_| validation("negative fixed-list length"))?)
601                .map_err(canonical_error)?;
602            encode_field(writer, field)?;
603            return Ok(());
604        }
605        DataType::Struct(fields) => {
606            writer.u8(0x62).map_err(canonical_error)?;
607            writer
608                .u32(exact_u32(fields.len(), "struct field count")?)
609                .map_err(canonical_error)?;
610            for field in fields {
611                encode_field(writer, field)?;
612            }
613            return Ok(());
614        }
615        DataType::Dictionary(_, value) => return encode_type(writer, value),
616        other => return Err(validation(format!("unsupported graph Arrow type {other}"))),
617    }
618    .map_err(canonical_error)
619}
620
621fn encode_value(
622    writer: &mut CanonicalWriter,
623    data_type: &DataType,
624    array: &ArrayRef,
625    row: usize,
626    nullable: bool,
627) -> Result<(), GfError> {
628    if array.is_null(row) {
629        if !nullable {
630            return Err(validation("non-nullable graph field contains null"));
631        }
632        writer.u8(0).map_err(canonical_error)?;
633        return Ok(());
634    }
635    writer.u8(1).map_err(canonical_error)?;
636    encode_present_value(writer, data_type, array, row)
637}
638
639#[allow(clippy::too_many_lines)]
640fn encode_present_value(
641    writer: &mut CanonicalWriter,
642    data_type: &DataType,
643    array: &ArrayRef,
644    row: usize,
645) -> Result<(), GfError> {
646    macro_rules! write {
647        ($value:expr) => {
648            $value.map_err(canonical_error)?
649        };
650    }
651    match data_type {
652        DataType::Boolean => {
653            write!(writer.u8(u8::from(downcast::<BooleanArray>(array)?.value(row))));
654        }
655        DataType::Int32 => {
656            write!(writer.raw(&downcast::<Int32Array>(array)?.value(row).to_be_bytes()));
657        }
658        DataType::Int64 => write!(writer.i64(downcast::<Int64Array>(array)?.value(row))),
659        DataType::UInt32 => write!(writer.u32(downcast::<UInt32Array>(array)?.value(row))),
660        DataType::UInt64 => write!(writer.u64(downcast::<UInt64Array>(array)?.value(row))),
661        DataType::Float32 => {
662            write!(writer.u32(normalize_f32(downcast::<Float32Array>(array)?.value(row))));
663        }
664        DataType::Float64 => {
665            write!(writer.u64(normalize_f64(downcast::<Float64Array>(array)?.value(row))));
666        }
667        DataType::Utf8 => write!(writer.text(downcast::<StringArray>(array)?.value(row))),
668        DataType::LargeUtf8 => {
669            write!(writer.text(downcast::<LargeStringArray>(array)?.value(row)));
670        }
671        DataType::Binary => write!(writer.binary(downcast::<BinaryArray>(array)?.value(row))),
672        DataType::LargeBinary => {
673            write!(writer.binary(downcast::<LargeBinaryArray>(array)?.value(row)));
674        }
675        DataType::FixedSizeBinary(_) => {
676            write!(writer.raw(downcast::<FixedSizeBinaryArray>(array)?.value(row)));
677        }
678        DataType::Timestamp(unit, timezone) => {
679            validate_timezone(timezone.as_deref())?;
680            write!(writer.i64(timestamp_value(array, *unit, row)?));
681        }
682        DataType::Time64(unit) => write!(writer.i64(time64_value(array, *unit, row)?)),
683        DataType::List(field) => {
684            encode_list(writer, field, &downcast::<ListArray>(array)?.value(row))?;
685        }
686        DataType::LargeList(field) => {
687            encode_list(
688                writer,
689                field,
690                &downcast::<LargeListArray>(array)?.value(row),
691            )?;
692        }
693        DataType::FixedSizeList(field, _) => {
694            encode_list(
695                writer,
696                field,
697                &downcast::<FixedSizeListArray>(array)?.value(row),
698            )?;
699        }
700        DataType::Struct(fields) => {
701            let values = downcast::<StructArray>(array)?;
702            for (field, child) in fields.iter().zip(values.columns()) {
703                encode_value(writer, field.data_type(), child, row, field.is_nullable())?;
704            }
705        }
706        DataType::Dictionary(_, value) => {
707            let decoded = arrow::compute::cast(array, value).map_err(storage)?;
708            encode_present_value(writer, value, &decoded, row)?;
709        }
710        other => return Err(validation(format!("unsupported graph Arrow value {other}"))),
711    }
712    Ok(())
713}
714
715fn encode_list(
716    writer: &mut CanonicalWriter,
717    field: &Field,
718    values: &ArrayRef,
719) -> Result<(), GfError> {
720    writer
721        .u64(exact_u64(values.len(), "graph list length")?)
722        .map_err(canonical_error)?;
723    for index in 0..values.len() {
724        encode_value(
725            writer,
726            field.data_type(),
727            values,
728            index,
729            field.is_nullable(),
730        )?;
731    }
732    Ok(())
733}
734
735fn dictionary_value_type(data_type: &DataType) -> &DataType {
736    match data_type {
737        DataType::Dictionary(_, value) => value,
738        other => other,
739    }
740}
741
742fn downcast<T: 'static>(array: &ArrayRef) -> Result<&T, GfError> {
743    array
744        .as_any()
745        .downcast_ref::<T>()
746        .ok_or_else(|| validation("graph Arrow array/type mismatch"))
747}
748
749fn timestamp_value(array: &ArrayRef, unit: TimeUnit, row: usize) -> Result<i64, GfError> {
750    Ok(match unit {
751        TimeUnit::Second => downcast::<arrow::array::TimestampSecondArray>(array)?.value(row),
752        TimeUnit::Millisecond => {
753            downcast::<arrow::array::TimestampMillisecondArray>(array)?.value(row)
754        }
755        TimeUnit::Microsecond => {
756            downcast::<arrow::array::TimestampMicrosecondArray>(array)?.value(row)
757        }
758        TimeUnit::Nanosecond => {
759            downcast::<arrow::array::TimestampNanosecondArray>(array)?.value(row)
760        }
761    })
762}
763
764fn time64_value(array: &ArrayRef, unit: TimeUnit, row: usize) -> Result<i64, GfError> {
765    match unit {
766        TimeUnit::Microsecond => {
767            Ok(downcast::<arrow::array::Time64MicrosecondArray>(array)?.value(row))
768        }
769        TimeUnit::Nanosecond => {
770            Ok(downcast::<arrow::array::Time64NanosecondArray>(array)?.value(row))
771        }
772        _ => Err(validation(
773            "Time64 must use microsecond or nanosecond units",
774        )),
775    }
776}
777
778fn validate_timezone(timezone: Option<&str>) -> Result<(), GfError> {
779    if timezone.is_none_or(|value| matches!(value, "UTC" | "Etc/UTC" | "Z" | "+00:00")) {
780        Ok(())
781    } else {
782        Err(validation("graph timestamp timezone is not canonical UTC"))
783    }
784}
785
786const fn time_unit_tag(unit: TimeUnit) -> u8 {
787    match unit {
788        TimeUnit::Second => 0,
789        TimeUnit::Millisecond => 1,
790        TimeUnit::Microsecond => 2,
791        TimeUnit::Nanosecond => 3,
792    }
793}
794
795fn normalize_f32(value: f32) -> u32 {
796    if value.is_nan() {
797        0x7fc0_0000
798    } else if value == 0.0 {
799        0
800    } else {
801        value.to_bits()
802    }
803}
804
805fn normalize_f64(value: f64) -> u64 {
806    if value.is_nan() {
807        0x7ff8_0000_0000_0000
808    } else if value == 0.0 {
809        0
810    } else {
811        value.to_bits()
812    }
813}
814
815fn exact_u32(value: usize, field: &str) -> Result<u32, GfError> {
816    u32::try_from(value).map_err(|_| validation(format!("{field} exceeds UInt32")))
817}
818
819fn exact_u64(value: usize, field: &str) -> Result<u64, GfError> {
820    u64::try_from(value).map_err(|_| validation(format!("{field} exceeds UInt64")))
821}
822
823fn canonical_error(error: impl std::fmt::Display) -> GfError {
824    validation(error.to_string())
825}
826
827fn sorted_parquet_files(directory: &Path) -> Result<Vec<PathBuf>, GfError> {
828    let entries = match fs::read_dir(directory) {
829        Ok(entries) => entries,
830        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
831        Err(error) => return Err(storage(error)),
832    };
833    let mut paths = Vec::new();
834    for entry in entries {
835        let entry = entry.map_err(storage)?;
836        let file_type = entry.file_type().map_err(storage)?;
837        if file_type.is_symlink() {
838            return Err(validation("graph directory contains a symbolic link"));
839        }
840        let path = entry.path();
841        if file_type.is_file()
842            && path.extension().and_then(|value| value.to_str()) == Some("parquet")
843        {
844            paths.push(path);
845        }
846    }
847    paths.sort();
848    Ok(paths)
849}
850
851fn uuid_column<'a>(
852    batch: &'a RecordBatch,
853    name: &str,
854) -> Result<&'a FixedSizeBinaryArray, GfError> {
855    batch
856        .column_by_name(name)
857        .and_then(|column| column.as_any().downcast_ref::<FixedSizeBinaryArray>())
858        .filter(|column| column.value_length() == 16)
859        .ok_or_else(|| validation(format!("graph column {name} must be FixedSizeBinary(16)")))
860}
861
862fn uuid_at(column: &FixedSizeBinaryArray, row: usize) -> Result<[u8; 16], GfError> {
863    if column.is_null(row) {
864        return Err(validation("graph UUID column contains null"));
865    }
866    column
867        .value(row)
868        .try_into()
869        .map_err(|_| validation("graph UUID has invalid width"))
870}
871
872fn require_present(
873    requested: &BTreeSet<[u8; 16]>,
874    available: &BTreeSet<[u8; 16]>,
875    kind: &str,
876) -> Result<(), GfError> {
877    if requested.is_subset(available) {
878        Ok(())
879    } else {
880        Err(validation(format!(
881            "graph projection references a missing {kind} UUID"
882        )))
883    }
884}
885
886fn validate_distinct_paths(source: &Path, target: &Path) -> Result<(), GfError> {
887    let source = source.canonicalize().map_err(storage)?;
888    let target = target
889        .canonicalize()
890        .or_else(|_| {
891            target
892                .parent()
893                .ok_or_else(|| std::io::Error::other("target has no parent"))?
894                .canonicalize()
895                .map(|parent| parent.join(target.file_name().unwrap_or_default()))
896        })
897        .map_err(storage)?;
898    if source == target || target.starts_with(&source) || source.starts_with(&target) {
899        return Err(validation(
900            "graph projection source and target must be disjoint",
901        ));
902    }
903    Ok(())
904}
905
906fn validate_graph_empty_target(target: &Path) -> Result<(), GfError> {
907    match fs::symlink_metadata(target) {
908        Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => {
909            Err(validation("graph projection target must be a directory"))
910        }
911        Ok(_) => {
912            for entry in fs::read_dir(target).map_err(storage)? {
913                let entry = entry.map_err(storage)?;
914                let name = entry.file_name();
915                let name = name
916                    .to_str()
917                    .ok_or_else(|| validation("graph projection target name is not UTF-8"))?;
918                match name {
919                    "topology" => validate_empty_topology(&entry.path())?,
920                    "properties" | "edge_properties" => {
921                        validate_empty_parquet_directory(&entry.path())?;
922                    }
923                    value
924                        if value == graphforge_core::manifest::MANIFEST_FILE
925                            || value == graphforge_core::manifest::ONTOLOGY_FILE =>
926                    {
927                        if !entry.file_type().map_err(storage)?.is_file() {
928                            return Err(validation("graph target metadata is not a regular file"));
929                        }
930                    }
931                    _ => {
932                        return Err(validation(
933                            "graph projection target contains non-graph or non-empty state",
934                        ));
935                    }
936                }
937            }
938            Ok(())
939        }
940        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
941        Err(error) => Err(storage(error)),
942    }
943}
944
945fn validate_empty_topology(directory: &Path) -> Result<(), GfError> {
946    for entry in fs::read_dir(directory).map_err(storage)? {
947        let entry = entry.map_err(storage)?;
948        let name = entry.file_name();
949        let name = name
950            .to_str()
951            .ok_or_else(|| validation("topology target name is not UTF-8"))?;
952        match name {
953            "edges" => validate_empty_parquet_directory(&entry.path())?,
954            "nodes.parquet" => require_empty_parquet(&entry.path())?,
955            "runtime_catalog.parquet" | "generation.json" => {
956                if !entry.file_type().map_err(storage)?.is_file() {
957                    return Err(validation("graph target metadata is not a regular file"));
958                }
959            }
960            _ => return Err(validation("graph projection target topology is not empty")),
961        }
962    }
963    Ok(())
964}
965
966fn validate_empty_parquet_directory(directory: &Path) -> Result<(), GfError> {
967    for entry in fs::read_dir(directory).map_err(storage)? {
968        let entry = entry.map_err(storage)?;
969        let path = entry.path();
970        if !entry.file_type().map_err(storage)?.is_file()
971            || path.extension().and_then(|value| value.to_str()) != Some("parquet")
972        {
973            return Err(validation(
974                "graph projection target graph directory is not empty",
975            ));
976        }
977        require_empty_parquet(&path)?;
978    }
979    Ok(())
980}
981
982fn require_empty_parquet(path: &Path) -> Result<(), GfError> {
983    let rows = read_parquet(path)?
984        .iter()
985        .map(RecordBatch::num_rows)
986        .sum::<usize>();
987    if rows == 0 {
988        Ok(())
989    } else {
990        Err(validation(
991            "graph projection target already contains graph rows",
992        ))
993    }
994}
995
996fn clear_graph_empty_target(target: &Path) -> Result<(), GfError> {
997    if !target.exists() {
998        return Ok(());
999    }
1000    for name in ["topology", "properties", "edge_properties"] {
1001        let path = target.join(name);
1002        if path.exists() {
1003            fs::remove_dir_all(path).map_err(storage)?;
1004        }
1005    }
1006    for name in [
1007        graphforge_core::manifest::MANIFEST_FILE,
1008        graphforge_core::manifest::ONTOLOGY_FILE,
1009    ] {
1010        let path = target.join(name);
1011        if path.exists() {
1012            fs::remove_file(path).map_err(storage)?;
1013        }
1014    }
1015    Ok(())
1016}
1017
1018fn copy_regular_file_if_present(source: &Path, target: &Path) -> Result<(), GfError> {
1019    let metadata = match fs::symlink_metadata(source) {
1020        Ok(metadata) => metadata,
1021        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
1022        Err(error) => return Err(storage(error)),
1023    };
1024    if !metadata.file_type().is_file() {
1025        return Err(validation("graph metadata must be a regular file"));
1026    }
1027    let parent = target
1028        .parent()
1029        .ok_or_else(|| validation("graph metadata target has no parent"))?;
1030    fs::create_dir_all(parent).map_err(storage)?;
1031    fs::copy(source, target).map_err(storage)?;
1032    Ok(())
1033}
1034
1035fn validation(message: impl Into<String>) -> GfError {
1036    GfError::Validation(message.into())
1037}
1038
1039fn storage(error: impl std::fmt::Display) -> GfError {
1040    GfError::Storage(error.to_string())
1041}
1042
1043#[cfg(test)]
1044mod tests {
1045    use std::collections::HashMap;
1046
1047    use arrow::array::{FixedSizeBinaryBuilder, Int64Array, ListArray, UInt32Array, UInt64Array};
1048    use graphforge_core::uuid::Uuid;
1049    use graphforge_core::{OntologyMode, TypeId};
1050    use graphforge_ir::{IrLiteral, RuntimeCatalog};
1051    use parquet::file::properties::WriterProperties;
1052    use tempfile::TempDir;
1053
1054    use super::*;
1055    use crate::{GraphWriter, read_edge_properties, read_nodes, read_properties};
1056
1057    const TS: i64 = 1_700_000_000_000_000;
1058
1059    fn uuid(marker: u8) -> Uuid {
1060        let mut bytes = [0_u8; 16];
1061        bytes[15] = marker;
1062        Uuid::from_bytes(bytes)
1063    }
1064
1065    fn fixture() -> (TempDir, [Uuid; 3], [Uuid; 2]) {
1066        let source = TempDir::new().unwrap();
1067        let nodes = [uuid(3), uuid(1), uuid(2)];
1068        let edges = [uuid(11), uuid(12)];
1069        let mut writer =
1070            GraphWriter::open_at(source.path(), OntologyMode::Exploratory, TS).unwrap();
1071        writer
1072            .create_node_with_labels(nodes[0], &[TypeId(7), TypeId(9)])
1073            .unwrap();
1074        writer.create_node(nodes[1], TypeId(8)).unwrap();
1075        writer.create_node(nodes[2], TypeId(10)).unwrap();
1076        for (index, node) in nodes.iter().enumerate() {
1077            writer
1078                .set_properties(
1079                    node,
1080                    None,
1081                    HashMap::from([("value".into(), IrLiteral::Int(index as i64))]),
1082                )
1083                .unwrap();
1084        }
1085        writer
1086            .create_edge(edges[0], "KNOWS", &nodes[0], &nodes[1])
1087            .unwrap();
1088        writer
1089            .create_edge(edges[1], "KNOWS", &nodes[1], &nodes[2])
1090            .unwrap();
1091        writer
1092            .set_edge_properties(
1093                &edges[0],
1094                Some("KNOWS"),
1095                HashMap::from([("weight".into(), IrLiteral::Float(0.75))]),
1096            )
1097            .unwrap();
1098        writer.flush().unwrap();
1099
1100        let mut catalog = RuntimeCatalog::new();
1101        catalog.intern_label("Person");
1102        catalog.intern_relation_type("KNOWS");
1103        catalog.intern_property("value", Some("Person"));
1104        write_parquet(
1105            &source.path().join("topology/runtime_catalog.parquet"),
1106            &catalog.to_record_batch(),
1107        )
1108        .unwrap();
1109        fs::write(
1110            source.path().join(graphforge_core::manifest::MANIFEST_FILE),
1111            b"ontology: ontology.yaml\n",
1112        )
1113        .unwrap();
1114        fs::write(
1115            source.path().join(graphforge_core::manifest::ONTOLOGY_FILE),
1116            b"version: 1\n",
1117        )
1118        .unwrap();
1119        for excluded in [
1120            "knowledge",
1121            "epistemic",
1122            "provenance",
1123            "valid_time",
1124            "indexes",
1125        ] {
1126            fs::create_dir_all(source.path().join(excluded)).unwrap();
1127            fs::write(
1128                source.path().join(excluded).join("must-not-copy"),
1129                b"secret",
1130            )
1131            .unwrap();
1132        }
1133        (source, nodes, edges)
1134    }
1135
1136    #[test]
1137    fn projection_preserves_graph_rows_closes_endpoints_and_never_induces_edges() {
1138        let (source, nodes, edges) = fixture();
1139        let target = TempDir::new().unwrap();
1140        let summary = materialize_graph_projection(
1141            source.path(),
1142            target.path(),
1143            &GraphProjectionSelection {
1144                node_uuids: BTreeSet::from([*nodes[2].as_bytes()]),
1145                edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1146            },
1147        )
1148        .unwrap();
1149
1150        assert_eq!(
1151            summary.node_uuids,
1152            vec![
1153                *nodes[1].as_bytes(),
1154                *nodes[2].as_bytes(),
1155                *nodes[0].as_bytes(),
1156            ]
1157        );
1158        assert_eq!(
1159            summary.endpoint_node_uuids,
1160            vec![*nodes[1].as_bytes(), *nodes[0].as_bytes()]
1161        );
1162        assert_eq!(summary.edge_uuids, vec![*edges[0].as_bytes()]);
1163
1164        let source_nodes = read_nodes(source.path()).unwrap();
1165        let projected_nodes = read_nodes(target.path()).unwrap();
1166        let source_ids = id_map(&source_nodes[0], "node_uuid", "node_id");
1167        let projected_ids = id_map(&projected_nodes[0], "node_uuid", "node_id");
1168        assert_eq!(projected_ids.len(), 3);
1169        for uuid in &summary.node_uuids {
1170            assert_eq!(projected_ids.get(uuid), source_ids.get(uuid));
1171        }
1172        let labels = projected_nodes[0]
1173            .column_by_name("type_ids")
1174            .unwrap()
1175            .as_any()
1176            .downcast_ref::<ListArray>()
1177            .unwrap();
1178        let selected_row = summary
1179            .node_uuids
1180            .iter()
1181            .position(|uuid| uuid == nodes[0].as_bytes())
1182            .unwrap();
1183        let values = labels.value(selected_row);
1184        assert_eq!(
1185            values
1186                .as_any()
1187                .downcast_ref::<UInt32Array>()
1188                .unwrap()
1189                .values(),
1190            &[7, 9]
1191        );
1192
1193        let projected_edges =
1194            read_parquet(&target.path().join("topology/edges/_exploratory.parquet")).unwrap();
1195        assert_eq!(projected_edges[0].num_rows(), 1);
1196        assert_eq!(
1197            uuid_at(uuid_column(&projected_edges[0], "edge_uuid").unwrap(), 0).unwrap(),
1198            *edges[0].as_bytes()
1199        );
1200        let source_edge_ids = id_map(
1201            &read_parquet(&source.path().join("topology/edges/_exploratory.parquet")).unwrap()[0],
1202            "edge_uuid",
1203            "edge_id",
1204        );
1205        let projected_edge_ids = id_map(&projected_edges[0], "edge_uuid", "edge_id");
1206        assert_eq!(
1207            projected_edge_ids.get(edges[0].as_bytes()),
1208            source_edge_ids.get(edges[0].as_bytes())
1209        );
1210
1211        assert_eq!(
1212            read_properties(target.path(), "_untyped")
1213                .unwrap()
1214                .iter()
1215                .map(RecordBatch::num_rows)
1216                .sum::<usize>(),
1217            3
1218        );
1219        assert_eq!(
1220            read_edge_properties(target.path(), "KNOWS")
1221                .unwrap()
1222                .iter()
1223                .map(RecordBatch::num_rows)
1224                .sum::<usize>(),
1225            1
1226        );
1227        assert!(
1228            target
1229                .path()
1230                .join("topology/runtime_catalog.parquet")
1231                .exists()
1232        );
1233        assert!(
1234            target
1235                .path()
1236                .join(graphforge_core::manifest::MANIFEST_FILE)
1237                .exists()
1238        );
1239        assert!(
1240            target
1241                .path()
1242                .join(graphforge_core::manifest::ONTOLOGY_FILE)
1243                .exists()
1244        );
1245        for excluded in [
1246            "knowledge",
1247            "epistemic",
1248            "provenance",
1249            "valid_time",
1250            "indexes",
1251        ] {
1252            assert!(!target.path().join(excluded).exists());
1253        }
1254    }
1255
1256    #[test]
1257    fn projection_is_canonically_ordered_and_reproducible() {
1258        let (source, nodes, edges) = fixture();
1259        let first = TempDir::new().unwrap();
1260        let second = TempDir::new().unwrap();
1261        let selection = GraphProjectionSelection {
1262            node_uuids: BTreeSet::from([*nodes[2].as_bytes()]),
1263            edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1264        };
1265        let left = materialize_graph_projection(source.path(), first.path(), &selection).unwrap();
1266        let right = materialize_graph_projection(source.path(), second.path(), &selection).unwrap();
1267        assert_eq!(left, right);
1268        assert_ne!(left.graph_content_fingerprint, [0; 32]);
1269        for relative in [
1270            "topology/nodes.parquet",
1271            "topology/edges/_exploratory.parquet",
1272            "properties/_untyped.parquet",
1273            "edge_properties/KNOWS.parquet",
1274            "topology/runtime_catalog.parquet",
1275        ] {
1276            assert_eq!(
1277                fs::read(first.path().join(relative)).unwrap(),
1278                fs::read(second.path().join(relative)).unwrap(),
1279                "non-deterministic output for {relative}"
1280            );
1281        }
1282    }
1283
1284    #[test]
1285    fn projection_fingerprint_ignores_parquet_chunking_and_dictionary_layout() {
1286        let (source, nodes, edges) = fixture();
1287        let baseline_target = TempDir::new().unwrap();
1288        let rewritten_target = TempDir::new().unwrap();
1289        let selection = GraphProjectionSelection {
1290            node_uuids: BTreeSet::from([*nodes[2].as_bytes()]),
1291            edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1292        };
1293        let baseline =
1294            materialize_graph_projection(source.path(), baseline_target.path(), &selection)
1295                .unwrap();
1296
1297        for relative in [
1298            "topology/nodes.parquet",
1299            "topology/edges/_exploratory.parquet",
1300            "properties/_untyped.parquet",
1301            "edge_properties/KNOWS.parquet",
1302            "topology/runtime_catalog.parquet",
1303        ] {
1304            let path = source.path().join(relative);
1305            let batches = read_parquet(&path).unwrap();
1306            let schema = batches[0].schema();
1307            let replacement = path.with_extension("rewritten");
1308            let file = fs::File::create(&replacement).unwrap();
1309            let properties = WriterProperties::builder()
1310                .set_dictionary_enabled(false)
1311                .set_max_row_group_row_count(Some(1))
1312                .build();
1313            let mut writer = ArrowWriter::try_new(file, schema, Some(properties)).unwrap();
1314            for batch in batches {
1315                for row in 0..batch.num_rows() {
1316                    writer.write(&batch.slice(row, 1)).unwrap();
1317                }
1318            }
1319            writer.close().unwrap();
1320            fs::rename(replacement, path).unwrap();
1321        }
1322
1323        let rewritten =
1324            materialize_graph_projection(source.path(), rewritten_target.path(), &selection)
1325                .unwrap();
1326        assert_eq!(
1327            baseline.graph_content_fingerprint,
1328            rewritten.graph_content_fingerprint
1329        );
1330    }
1331
1332    #[test]
1333    fn unrelated_runtime_catalog_entries_do_not_change_projection_identity() {
1334        let (source, nodes, edges) = fixture();
1335        let first = TempDir::new().unwrap();
1336        let second = TempDir::new().unwrap();
1337        let selection = GraphProjectionSelection {
1338            node_uuids: BTreeSet::from([*nodes[0].as_bytes()]),
1339            edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1340        };
1341        let baseline =
1342            materialize_graph_projection(source.path(), first.path(), &selection).unwrap();
1343
1344        let catalog_path = source.path().join("topology/runtime_catalog.parquet");
1345        let batch = read_parquet(&catalog_path).unwrap().remove(0);
1346        let mut catalog = RuntimeCatalog::from_record_batch(&batch).unwrap();
1347        catalog.intern_label("Unrelated");
1348        catalog.intern_relation_type("IGNORES");
1349        catalog.intern_property("noise", Some("Unrelated"));
1350        write_parquet(&catalog_path, &catalog.to_record_batch()).unwrap();
1351
1352        let with_noise =
1353            materialize_graph_projection(source.path(), second.path(), &selection).unwrap();
1354        assert_eq!(
1355            baseline.graph_content_fingerprint,
1356            with_noise.graph_content_fingerprint
1357        );
1358        assert_eq!(
1359            fs::read(first.path().join("topology/runtime_catalog.parquet")).unwrap(),
1360            fs::read(second.path().join("topology/runtime_catalog.parquet")).unwrap()
1361        );
1362        let projected = read_parquet(&second.path().join("topology/runtime_catalog.parquet"))
1363            .unwrap()
1364            .remove(0);
1365        let names = string_column(&projected, "name").unwrap();
1366        assert!(!(0..names.len()).any(|row| names.value(row) == "Unrelated"));
1367        assert!(!(0..names.len()).any(|row| names.value(row) == "IGNORES"));
1368        assert!(!(0..names.len()).any(|row| names.value(row) == "noise"));
1369    }
1370
1371    #[test]
1372    fn existing_graph_empty_hydrated_workspace_is_a_valid_target() {
1373        let (source, nodes, _) = fixture();
1374        let target = TempDir::new().unwrap();
1375        write_parquet(
1376            &target.path().join("topology/nodes.parquet"),
1377            &RecordBatch::new_empty(Arc::clone(&crate::TOPOLOGY_NODES_SCHEMA)),
1378        )
1379        .unwrap();
1380        write_parquet(
1381            &target.path().join("topology/runtime_catalog.parquet"),
1382            &RuntimeCatalog::new().to_record_batch(),
1383        )
1384        .unwrap();
1385        fs::write(
1386            target.path().join("topology/generation.json"),
1387            b"{\"topology_generation\":0,\"search_generation\":0}\n",
1388        )
1389        .unwrap();
1390
1391        let summary = materialize_graph_projection(
1392            source.path(),
1393            target.path(),
1394            &GraphProjectionSelection {
1395                node_uuids: BTreeSet::from([*nodes[0].as_bytes()]),
1396                edge_uuids: BTreeSet::new(),
1397            },
1398        )
1399        .unwrap();
1400        assert_eq!(summary.node_uuids, vec![*nodes[0].as_bytes()]);
1401        assert_eq!(read_nodes(target.path()).unwrap()[0].num_rows(), 1);
1402        assert!(!target.path().join("topology/generation.json").exists());
1403    }
1404
1405    #[test]
1406    fn missing_identity_and_nonempty_target_fail_before_writing() {
1407        let (source, _, _) = fixture();
1408        let target = TempDir::new().unwrap();
1409        let missing = uuid(99);
1410        let error = materialize_graph_projection(
1411            source.path(),
1412            target.path(),
1413            &GraphProjectionSelection {
1414                node_uuids: BTreeSet::from([*missing.as_bytes()]),
1415                edge_uuids: BTreeSet::new(),
1416            },
1417        )
1418        .unwrap_err();
1419        assert!(matches!(error, GfError::Validation(_)));
1420        assert!(fs::read_dir(target.path()).unwrap().next().is_none());
1421
1422        fs::write(target.path().join("owned"), b"keep").unwrap();
1423        let error = materialize_graph_projection(
1424            source.path(),
1425            target.path(),
1426            &GraphProjectionSelection::default(),
1427        )
1428        .unwrap_err();
1429        assert!(matches!(error, GfError::Validation(_)));
1430        assert_eq!(fs::read(target.path().join("owned")).unwrap(), b"keep");
1431    }
1432
1433    #[test]
1434    fn corrupt_property_uuid_is_rejected_instead_of_silently_dropped() {
1435        let source = TempDir::new().unwrap();
1436        let target = TempDir::new().unwrap();
1437        let mut uuids = FixedSizeBinaryBuilder::new(16);
1438        uuids.append_null();
1439        let batch = RecordBatch::try_from_iter([
1440            ("node_uuid", Arc::new(uuids.finish()) as ArrayRef),
1441            ("value", Arc::new(Int64Array::from(vec![1])) as ArrayRef),
1442        ])
1443        .unwrap();
1444        let source_path = source.path().join("properties/Person.parquet");
1445        write_parquet(&source_path, &batch).unwrap();
1446
1447        let error = project_parquet_file(
1448            &source_path,
1449            &target.path().join("properties/Person.parquet"),
1450            "node_uuid",
1451            &BTreeSet::new(),
1452        )
1453        .unwrap_err();
1454        assert!(matches!(error, GfError::Validation(_)));
1455        assert!(error.to_string().contains("UUID column contains null"));
1456        assert!(!target.path().join("properties/Person.parquet").exists());
1457    }
1458
1459    fn id_map(batch: &RecordBatch, uuid_name: &str, id_name: &str) -> BTreeMap<[u8; 16], u64> {
1460        let uuids = uuid_column(batch, uuid_name).unwrap();
1461        let ids = batch
1462            .column_by_name(id_name)
1463            .unwrap()
1464            .as_any()
1465            .downcast_ref::<UInt64Array>()
1466            .unwrap();
1467        (0..batch.num_rows())
1468            .map(|row| (uuid_at(uuids, row).unwrap(), ids.value(row)))
1469            .collect()
1470    }
1471}