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        let batches = read_parquet(&path)?;
310        if stem != "_exploratory" && batches.iter().any(|batch| batch.num_rows() != 0) {
311            relation_names.insert(stem);
312        }
313        for batch in batches {
314            if let Some(column) = batch.column_by_name("rel_type_name") {
315                let values = column
316                    .as_any()
317                    .downcast_ref::<StringArray>()
318                    .ok_or_else(|| validation("edge rel_type_name is not Utf8"))?;
319                for row in 0..values.len() {
320                    if !values.is_null(row) {
321                        relation_names.insert(values.value(row).to_owned());
322                    }
323                }
324            }
325        }
326    }
327
328    let mut property_names = BTreeSet::new();
329    for directory in ["properties", "edge_properties"] {
330        for path in sorted_parquet_files(&target.join(directory))? {
331            let batches = read_parquet(&path)?;
332            if batches.iter().all(|batch| batch.num_rows() == 0) {
333                continue;
334            }
335            let schema = batches
336                .first()
337                .map(RecordBatch::schema)
338                .ok_or_else(|| validation("projected property table has no schema"))?;
339            for field in schema.fields() {
340                if !matches!(
341                    field.name().as_str(),
342                    "node_uuid" | "node_id" | "edge_uuid" | "edge_id"
343                ) {
344                    property_names.insert(field.name().clone());
345                }
346            }
347        }
348    }
349
350    let kinds = string_column(catalog, "entry_kind")?;
351    let names = string_column(catalog, "name")?;
352    let ids = catalog
353        .column_by_name("runtime_id")
354        .and_then(|column| column.as_any().downcast_ref::<UInt32Array>())
355        .ok_or_else(|| validation("runtime catalog runtime_id is not UInt32"))?;
356    let owners = string_column(catalog, "owner_label")?;
357
358    let mut active_owners = BTreeSet::new();
359    for row in 0..catalog.num_rows() {
360        if kinds.value(row) == "entity_type" && type_ids.contains(&ids.value(row)) {
361            active_owners.insert(names.value(row).to_owned());
362        }
363    }
364    active_owners.extend(relation_names.iter().cloned());
365
366    let mut selected = BTreeSet::new();
367    for row in 0..catalog.num_rows() {
368        let keep = match kinds.value(row) {
369            "entity_type" => type_ids.contains(&ids.value(row)),
370            "relation_type" => relation_names.contains(names.value(row)),
371            "property" => {
372                property_names.contains(names.value(row))
373                    && (owners.is_null(row) || active_owners.contains(owners.value(row)))
374            }
375            _ => false,
376        };
377        if keep {
378            selected.insert(row);
379            if kinds.value(row) == "property" && !owners.is_null(row) {
380                let owner = owners.value(row);
381                for owner_row in 0..catalog.num_rows() {
382                    if matches!(kinds.value(owner_row), "entity_type" | "relation_type")
383                        && names.value(owner_row) == owner
384                    {
385                        selected.insert(owner_row);
386                    }
387                }
388            }
389        }
390    }
391    Ok(selected.into_iter().collect())
392}
393
394fn parquet_stem(path: &Path) -> Result<String, GfError> {
395    path.file_stem()
396        .and_then(|value| value.to_str())
397        .map(str::to_owned)
398        .ok_or_else(|| validation("graph parquet path has no UTF-8 stem"))
399}
400
401fn string_column<'a>(batch: &'a RecordBatch, name: &str) -> Result<&'a StringArray, GfError> {
402    batch
403        .column_by_name(name)
404        .and_then(|column| column.as_any().downcast_ref::<StringArray>())
405        .ok_or_else(|| validation(format!("runtime catalog {name} is not Utf8")))
406}
407
408fn read_parquet(path: &Path) -> Result<Vec<RecordBatch>, GfError> {
409    let schema = crate::catalog::discover_parquet_schema(path).ok_or_else(|| {
410        validation(format!(
411            "cannot discover graph schema for {}",
412            path.display()
413        ))
414    })?;
415    crate::catalog::read_parquet_or_empty(path, schema)
416        .map_err(|error| GfError::Storage(error.to_string()))
417}
418
419fn write_parquet(path: &Path, batch: &RecordBatch) -> Result<(), GfError> {
420    let parent = path
421        .parent()
422        .ok_or_else(|| validation("graph parquet target has no parent"))?;
423    fs::create_dir_all(parent).map_err(storage)?;
424    let file = File::create(path).map_err(storage)?;
425    let mut writer = ArrowWriter::try_new(file, batch.schema(), None).map_err(storage)?;
426    writer.write(batch).map_err(storage)?;
427    writer.close().map_err(storage)?;
428    Ok(())
429}
430
431fn projected_graph_fingerprint(root: &Path) -> Result<[u8; 32], GfError> {
432    let mut paths = Vec::new();
433    let nodes = root.join("topology/nodes.parquet");
434    if nodes.exists() {
435        paths.push(nodes);
436    }
437    for directory in ["topology/edges", "properties", "edge_properties"] {
438        paths.extend(sorted_parquet_files(&root.join(directory))?);
439    }
440    let runtime_catalog = root.join("topology/runtime_catalog.parquet");
441    if runtime_catalog.exists() {
442        paths.push(runtime_catalog);
443    }
444    paths.sort();
445
446    let mut writer = CanonicalWriter::new();
447    writer.raw(b"GFGP1").map_err(canonical_error)?;
448    writer
449        .u32(exact_u32(paths.len(), "graph table count")?)
450        .map_err(canonical_error)?;
451    for path in paths {
452        let relative = path
453            .strip_prefix(root)
454            .map_err(|_| validation("graph projection path escaped target"))?
455            .to_str()
456            .ok_or_else(|| validation("graph projection path is not UTF-8"))?;
457        writer.text(relative).map_err(canonical_error)?;
458        let batches = read_parquet(&path)?;
459        let schema = batches
460            .first()
461            .map(RecordBatch::schema)
462            .ok_or_else(|| validation("graph projection table has no schema"))?;
463        let batch = concat_batches(&schema, &batches).map_err(storage)?;
464        let logical = logical_fingerprint_batch(relative, &batch)?;
465        encode_table(&mut writer, &logical)?;
466    }
467    fingerprint(
468        CanonicalDomain::GraphProjection,
469        CANONICAL_CONTRACT_VERSION,
470        &writer.finish(),
471    )
472    .map_err(canonical_error)
473}
474
475fn logical_fingerprint_batch(relative: &str, batch: &RecordBatch) -> Result<RecordBatch, GfError> {
476    let source_schema = batch.schema();
477    let names: Vec<&str> = if relative == "topology/nodes.parquet" {
478        vec!["node_uuid", "type_id", "type_ids"]
479    } else if relative.starts_with("topology/edges/") {
480        let mut names = vec!["edge_uuid", "src_uuid", "dst_uuid"];
481        if batch.column_by_name("rel_type_name").is_some() {
482            names.push("rel_type_name");
483        }
484        names
485    } else if relative == "topology/runtime_catalog.parquet" {
486        vec!["entry_kind", "name", "runtime_id", "owner_label"]
487    } else {
488        source_schema
489            .fields()
490            .iter()
491            .map(|field| field.name().as_str())
492            .collect()
493    };
494    let mut fields = Vec::with_capacity(names.len());
495    let mut columns = Vec::with_capacity(names.len());
496    for name in names {
497        let index = source_schema
498            .index_of(name)
499            .map_err(|_| validation(format!("graph fingerprint field {name} is absent")))?;
500        fields.push(Arc::clone(&source_schema.fields()[index]));
501        columns.push(Arc::clone(batch.column(index)));
502    }
503    let schema = Arc::new(Schema::new_with_metadata(
504        fields,
505        source_schema.metadata().clone(),
506    ));
507    RecordBatch::try_new(schema, columns).map_err(storage)
508}
509
510fn encode_table(writer: &mut CanonicalWriter, batch: &RecordBatch) -> Result<(), GfError> {
511    encode_schema(writer, batch.schema().as_ref())?;
512    writer
513        .u64(exact_u64(batch.num_rows(), "graph row count")?)
514        .map_err(canonical_error)?;
515    let schema = batch.schema();
516    let columns = schema
517        .fields()
518        .iter()
519        .zip(batch.columns())
520        .map(|(field, column)| {
521            let logical = dictionary_value_type(field.data_type());
522            if logical == field.data_type() {
523                Ok((logical, Arc::clone(column)))
524            } else {
525                arrow::compute::cast(column, logical)
526                    .map(|decoded| (logical, decoded))
527                    .map_err(storage)
528            }
529        })
530        .collect::<Result<Vec<_>, _>>()?;
531    for row in 0..batch.num_rows() {
532        for (field, (data_type, column)) in schema.fields().iter().zip(&columns) {
533            encode_value(writer, data_type, column, row, field.is_nullable())?;
534        }
535    }
536    Ok(())
537}
538
539fn encode_schema(writer: &mut CanonicalWriter, schema: &Schema) -> Result<(), GfError> {
540    writer.raw(b"GFS1").map_err(canonical_error)?;
541    writer
542        .u32(exact_u32(schema.fields().len(), "graph field count")?)
543        .map_err(canonical_error)?;
544    for field in &schema.fields {
545        encode_field(writer, field)?;
546    }
547    let ordered = schema.metadata().iter().collect::<BTreeMap<_, _>>();
548    writer
549        .u32(exact_u32(ordered.len(), "graph metadata count")?)
550        .map_err(canonical_error)?;
551    for (key, value) in ordered {
552        writer.text(key).map_err(canonical_error)?;
553        writer.text(value).map_err(canonical_error)?;
554    }
555    Ok(())
556}
557
558fn encode_field(writer: &mut CanonicalWriter, field: &Field) -> Result<(), GfError> {
559    writer.text(field.name()).map_err(canonical_error)?;
560    writer
561        .u8(u8::from(field.is_nullable()))
562        .map_err(canonical_error)?;
563    encode_type(writer, field.data_type())
564}
565
566fn encode_type(writer: &mut CanonicalWriter, data_type: &DataType) -> Result<(), GfError> {
567    match data_type {
568        DataType::Boolean => writer.u8(0x02),
569        DataType::Int32 => writer.u8(0x12),
570        DataType::Int64 => writer.u8(0x13),
571        DataType::UInt32 => writer.u8(0x16),
572        DataType::UInt64 => writer.u8(0x17),
573        DataType::Float32 => writer.u8(0x21),
574        DataType::Float64 => writer.u8(0x22),
575        DataType::Utf8 | DataType::LargeUtf8 => writer.u8(0x30),
576        DataType::Binary | DataType::LargeBinary => writer.u8(0x31),
577        DataType::FixedSizeBinary(width) => {
578            writer.u8(0x32).map_err(canonical_error)?;
579            writer.u32(
580                u32::try_from(*width)
581                    .map_err(|_| validation("negative fixed-size binary width"))?,
582            )
583        }
584        DataType::Timestamp(unit, timezone) => {
585            validate_timezone(timezone.as_deref())?;
586            writer.u8(0x52).map_err(canonical_error)?;
587            writer.u8(time_unit_tag(*unit))
588        }
589        DataType::Time64(unit) => {
590            writer.u8(0x53).map_err(canonical_error)?;
591            writer.u8(time_unit_tag(*unit))
592        }
593        DataType::List(field) | DataType::LargeList(field) => {
594            writer.u8(0x60).map_err(canonical_error)?;
595            encode_field(writer, field)?;
596            return Ok(());
597        }
598        DataType::FixedSizeList(field, length) => {
599            writer.u8(0x61).map_err(canonical_error)?;
600            writer
601                .u32(u32::try_from(*length).map_err(|_| validation("negative fixed-list length"))?)
602                .map_err(canonical_error)?;
603            encode_field(writer, field)?;
604            return Ok(());
605        }
606        DataType::Struct(fields) => {
607            writer.u8(0x62).map_err(canonical_error)?;
608            writer
609                .u32(exact_u32(fields.len(), "struct field count")?)
610                .map_err(canonical_error)?;
611            for field in fields {
612                encode_field(writer, field)?;
613            }
614            return Ok(());
615        }
616        DataType::Dictionary(_, value) => return encode_type(writer, value),
617        other => return Err(validation(format!("unsupported graph Arrow type {other}"))),
618    }
619    .map_err(canonical_error)
620}
621
622fn encode_value(
623    writer: &mut CanonicalWriter,
624    data_type: &DataType,
625    array: &ArrayRef,
626    row: usize,
627    nullable: bool,
628) -> Result<(), GfError> {
629    if array.is_null(row) {
630        if !nullable {
631            return Err(validation("non-nullable graph field contains null"));
632        }
633        writer.u8(0).map_err(canonical_error)?;
634        return Ok(());
635    }
636    writer.u8(1).map_err(canonical_error)?;
637    encode_present_value(writer, data_type, array, row)
638}
639
640#[allow(clippy::too_many_lines)]
641fn encode_present_value(
642    writer: &mut CanonicalWriter,
643    data_type: &DataType,
644    array: &ArrayRef,
645    row: usize,
646) -> Result<(), GfError> {
647    macro_rules! write {
648        ($value:expr) => {
649            $value.map_err(canonical_error)?
650        };
651    }
652    match data_type {
653        DataType::Boolean => {
654            write!(writer.u8(u8::from(downcast::<BooleanArray>(array)?.value(row))));
655        }
656        DataType::Int32 => {
657            write!(writer.raw(&downcast::<Int32Array>(array)?.value(row).to_be_bytes()));
658        }
659        DataType::Int64 => write!(writer.i64(downcast::<Int64Array>(array)?.value(row))),
660        DataType::UInt32 => write!(writer.u32(downcast::<UInt32Array>(array)?.value(row))),
661        DataType::UInt64 => write!(writer.u64(downcast::<UInt64Array>(array)?.value(row))),
662        DataType::Float32 => {
663            write!(writer.u32(normalize_f32(downcast::<Float32Array>(array)?.value(row))));
664        }
665        DataType::Float64 => {
666            write!(writer.u64(normalize_f64(downcast::<Float64Array>(array)?.value(row))));
667        }
668        DataType::Utf8 => write!(writer.text(downcast::<StringArray>(array)?.value(row))),
669        DataType::LargeUtf8 => {
670            write!(writer.text(downcast::<LargeStringArray>(array)?.value(row)));
671        }
672        DataType::Binary => write!(writer.binary(downcast::<BinaryArray>(array)?.value(row))),
673        DataType::LargeBinary => {
674            write!(writer.binary(downcast::<LargeBinaryArray>(array)?.value(row)));
675        }
676        DataType::FixedSizeBinary(_) => {
677            write!(writer.raw(downcast::<FixedSizeBinaryArray>(array)?.value(row)));
678        }
679        DataType::Timestamp(unit, timezone) => {
680            validate_timezone(timezone.as_deref())?;
681            write!(writer.i64(timestamp_value(array, *unit, row)?));
682        }
683        DataType::Time64(unit) => write!(writer.i64(time64_value(array, *unit, row)?)),
684        DataType::List(field) => {
685            encode_list(writer, field, &downcast::<ListArray>(array)?.value(row))?;
686        }
687        DataType::LargeList(field) => {
688            encode_list(
689                writer,
690                field,
691                &downcast::<LargeListArray>(array)?.value(row),
692            )?;
693        }
694        DataType::FixedSizeList(field, _) => {
695            encode_list(
696                writer,
697                field,
698                &downcast::<FixedSizeListArray>(array)?.value(row),
699            )?;
700        }
701        DataType::Struct(fields) => {
702            let values = downcast::<StructArray>(array)?;
703            for (field, child) in fields.iter().zip(values.columns()) {
704                encode_value(writer, field.data_type(), child, row, field.is_nullable())?;
705            }
706        }
707        DataType::Dictionary(_, value) => {
708            let decoded = arrow::compute::cast(array, value).map_err(storage)?;
709            encode_present_value(writer, value, &decoded, row)?;
710        }
711        other => return Err(validation(format!("unsupported graph Arrow value {other}"))),
712    }
713    Ok(())
714}
715
716fn encode_list(
717    writer: &mut CanonicalWriter,
718    field: &Field,
719    values: &ArrayRef,
720) -> Result<(), GfError> {
721    writer
722        .u64(exact_u64(values.len(), "graph list length")?)
723        .map_err(canonical_error)?;
724    for index in 0..values.len() {
725        encode_value(
726            writer,
727            field.data_type(),
728            values,
729            index,
730            field.is_nullable(),
731        )?;
732    }
733    Ok(())
734}
735
736fn dictionary_value_type(data_type: &DataType) -> &DataType {
737    match data_type {
738        DataType::Dictionary(_, value) => value,
739        other => other,
740    }
741}
742
743fn downcast<T: 'static>(array: &ArrayRef) -> Result<&T, GfError> {
744    array
745        .as_any()
746        .downcast_ref::<T>()
747        .ok_or_else(|| validation("graph Arrow array/type mismatch"))
748}
749
750fn timestamp_value(array: &ArrayRef, unit: TimeUnit, row: usize) -> Result<i64, GfError> {
751    Ok(match unit {
752        TimeUnit::Second => downcast::<arrow::array::TimestampSecondArray>(array)?.value(row),
753        TimeUnit::Millisecond => {
754            downcast::<arrow::array::TimestampMillisecondArray>(array)?.value(row)
755        }
756        TimeUnit::Microsecond => {
757            downcast::<arrow::array::TimestampMicrosecondArray>(array)?.value(row)
758        }
759        TimeUnit::Nanosecond => {
760            downcast::<arrow::array::TimestampNanosecondArray>(array)?.value(row)
761        }
762    })
763}
764
765fn time64_value(array: &ArrayRef, unit: TimeUnit, row: usize) -> Result<i64, GfError> {
766    match unit {
767        TimeUnit::Microsecond => {
768            Ok(downcast::<arrow::array::Time64MicrosecondArray>(array)?.value(row))
769        }
770        TimeUnit::Nanosecond => {
771            Ok(downcast::<arrow::array::Time64NanosecondArray>(array)?.value(row))
772        }
773        _ => Err(validation(
774            "Time64 must use microsecond or nanosecond units",
775        )),
776    }
777}
778
779fn validate_timezone(timezone: Option<&str>) -> Result<(), GfError> {
780    if timezone.is_none_or(|value| matches!(value, "UTC" | "Etc/UTC" | "Z" | "+00:00")) {
781        Ok(())
782    } else {
783        Err(validation("graph timestamp timezone is not canonical UTC"))
784    }
785}
786
787const fn time_unit_tag(unit: TimeUnit) -> u8 {
788    match unit {
789        TimeUnit::Second => 0,
790        TimeUnit::Millisecond => 1,
791        TimeUnit::Microsecond => 2,
792        TimeUnit::Nanosecond => 3,
793    }
794}
795
796fn normalize_f32(value: f32) -> u32 {
797    if value.is_nan() {
798        0x7fc0_0000
799    } else if value == 0.0 {
800        0
801    } else {
802        value.to_bits()
803    }
804}
805
806fn normalize_f64(value: f64) -> u64 {
807    if value.is_nan() {
808        0x7ff8_0000_0000_0000
809    } else if value == 0.0 {
810        0
811    } else {
812        value.to_bits()
813    }
814}
815
816fn exact_u32(value: usize, field: &str) -> Result<u32, GfError> {
817    u32::try_from(value).map_err(|_| validation(format!("{field} exceeds UInt32")))
818}
819
820fn exact_u64(value: usize, field: &str) -> Result<u64, GfError> {
821    u64::try_from(value).map_err(|_| validation(format!("{field} exceeds UInt64")))
822}
823
824fn canonical_error(error: impl std::fmt::Display) -> GfError {
825    validation(error.to_string())
826}
827
828fn sorted_parquet_files(directory: &Path) -> Result<Vec<PathBuf>, GfError> {
829    let entries = match fs::read_dir(directory) {
830        Ok(entries) => entries,
831        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
832        Err(error) => return Err(storage(error)),
833    };
834    let mut paths = Vec::new();
835    for entry in entries {
836        let entry = entry.map_err(storage)?;
837        let file_type = entry.file_type().map_err(storage)?;
838        if file_type.is_symlink() {
839            return Err(validation("graph directory contains a symbolic link"));
840        }
841        let path = entry.path();
842        if file_type.is_file()
843            && path.extension().and_then(|value| value.to_str()) == Some("parquet")
844        {
845            paths.push(path);
846        }
847    }
848    paths.sort();
849    Ok(paths)
850}
851
852fn uuid_column<'a>(
853    batch: &'a RecordBatch,
854    name: &str,
855) -> Result<&'a FixedSizeBinaryArray, GfError> {
856    batch
857        .column_by_name(name)
858        .and_then(|column| column.as_any().downcast_ref::<FixedSizeBinaryArray>())
859        .filter(|column| column.value_length() == 16)
860        .ok_or_else(|| validation(format!("graph column {name} must be FixedSizeBinary(16)")))
861}
862
863fn uuid_at(column: &FixedSizeBinaryArray, row: usize) -> Result<[u8; 16], GfError> {
864    if column.is_null(row) {
865        return Err(validation("graph UUID column contains null"));
866    }
867    column
868        .value(row)
869        .try_into()
870        .map_err(|_| validation("graph UUID has invalid width"))
871}
872
873fn require_present(
874    requested: &BTreeSet<[u8; 16]>,
875    available: &BTreeSet<[u8; 16]>,
876    kind: &str,
877) -> Result<(), GfError> {
878    if requested.is_subset(available) {
879        Ok(())
880    } else {
881        Err(validation(format!(
882            "graph projection references a missing {kind} UUID"
883        )))
884    }
885}
886
887fn validate_distinct_paths(source: &Path, target: &Path) -> Result<(), GfError> {
888    let source = source.canonicalize().map_err(storage)?;
889    let target = target
890        .canonicalize()
891        .or_else(|_| {
892            target
893                .parent()
894                .ok_or_else(|| std::io::Error::other("target has no parent"))?
895                .canonicalize()
896                .map(|parent| parent.join(target.file_name().unwrap_or_default()))
897        })
898        .map_err(storage)?;
899    if source == target || target.starts_with(&source) || source.starts_with(&target) {
900        return Err(validation(
901            "graph projection source and target must be disjoint",
902        ));
903    }
904    Ok(())
905}
906
907fn validate_graph_empty_target(target: &Path) -> Result<(), GfError> {
908    match fs::symlink_metadata(target) {
909        Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => {
910            Err(validation("graph projection target must be a directory"))
911        }
912        Ok(_) => {
913            for entry in fs::read_dir(target).map_err(storage)? {
914                let entry = entry.map_err(storage)?;
915                let name = entry.file_name();
916                let name = name
917                    .to_str()
918                    .ok_or_else(|| validation("graph projection target name is not UTF-8"))?;
919                match name {
920                    "topology" => validate_empty_topology(&entry.path())?,
921                    "properties" | "edge_properties" => {
922                        validate_empty_parquet_directory(&entry.path())?;
923                    }
924                    value
925                        if value == graphforge_core::manifest::MANIFEST_FILE
926                            || value == graphforge_core::manifest::ONTOLOGY_FILE =>
927                    {
928                        if !entry.file_type().map_err(storage)?.is_file() {
929                            return Err(validation("graph target metadata is not a regular file"));
930                        }
931                    }
932                    _ => {
933                        return Err(validation(
934                            "graph projection target contains non-graph or non-empty state",
935                        ));
936                    }
937                }
938            }
939            Ok(())
940        }
941        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
942        Err(error) => Err(storage(error)),
943    }
944}
945
946fn validate_empty_topology(directory: &Path) -> Result<(), GfError> {
947    for entry in fs::read_dir(directory).map_err(storage)? {
948        let entry = entry.map_err(storage)?;
949        let name = entry.file_name();
950        let name = name
951            .to_str()
952            .ok_or_else(|| validation("topology target name is not UTF-8"))?;
953        match name {
954            "edges" => validate_empty_parquet_directory(&entry.path())?,
955            "nodes.parquet" => require_empty_parquet(&entry.path())?,
956            "runtime_catalog.parquet" | "generation.json" => {
957                if !entry.file_type().map_err(storage)?.is_file() {
958                    return Err(validation("graph target metadata is not a regular file"));
959                }
960            }
961            _ => return Err(validation("graph projection target topology is not empty")),
962        }
963    }
964    Ok(())
965}
966
967fn validate_empty_parquet_directory(directory: &Path) -> Result<(), GfError> {
968    for entry in fs::read_dir(directory).map_err(storage)? {
969        let entry = entry.map_err(storage)?;
970        let path = entry.path();
971        if !entry.file_type().map_err(storage)?.is_file()
972            || path.extension().and_then(|value| value.to_str()) != Some("parquet")
973        {
974            return Err(validation(
975                "graph projection target graph directory is not empty",
976            ));
977        }
978        require_empty_parquet(&path)?;
979    }
980    Ok(())
981}
982
983fn require_empty_parquet(path: &Path) -> Result<(), GfError> {
984    let rows = read_parquet(path)?
985        .iter()
986        .map(RecordBatch::num_rows)
987        .sum::<usize>();
988    if rows == 0 {
989        Ok(())
990    } else {
991        Err(validation(
992            "graph projection target already contains graph rows",
993        ))
994    }
995}
996
997fn clear_graph_empty_target(target: &Path) -> Result<(), GfError> {
998    if !target.exists() {
999        return Ok(());
1000    }
1001    for name in ["topology", "properties", "edge_properties"] {
1002        let path = target.join(name);
1003        if path.exists() {
1004            fs::remove_dir_all(path).map_err(storage)?;
1005        }
1006    }
1007    for name in [
1008        graphforge_core::manifest::MANIFEST_FILE,
1009        graphforge_core::manifest::ONTOLOGY_FILE,
1010    ] {
1011        let path = target.join(name);
1012        if path.exists() {
1013            fs::remove_file(path).map_err(storage)?;
1014        }
1015    }
1016    Ok(())
1017}
1018
1019fn copy_regular_file_if_present(source: &Path, target: &Path) -> Result<(), GfError> {
1020    let metadata = match fs::symlink_metadata(source) {
1021        Ok(metadata) => metadata,
1022        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
1023        Err(error) => return Err(storage(error)),
1024    };
1025    if !metadata.file_type().is_file() {
1026        return Err(validation("graph metadata must be a regular file"));
1027    }
1028    let parent = target
1029        .parent()
1030        .ok_or_else(|| validation("graph metadata target has no parent"))?;
1031    fs::create_dir_all(parent).map_err(storage)?;
1032    fs::copy(source, target).map_err(storage)?;
1033    Ok(())
1034}
1035
1036fn validation(message: impl Into<String>) -> GfError {
1037    GfError::Validation(message.into())
1038}
1039
1040fn storage(error: impl std::fmt::Display) -> GfError {
1041    GfError::Storage(error.to_string())
1042}
1043
1044#[cfg(test)]
1045mod tests {
1046    use std::collections::HashMap;
1047    use std::sync::Arc;
1048
1049    use arrow::array::{
1050        ArrayRef, BinaryArray, BooleanArray, FixedSizeBinaryBuilder, Float32Array, Float64Array,
1051        Int32Array, Int64Array, LargeBinaryArray, LargeStringArray, ListArray, StringArray,
1052        Time64MicrosecondArray, Time64NanosecondArray, TimestampMicrosecondArray,
1053        TimestampMillisecondArray, TimestampNanosecondArray, TimestampSecondArray, UInt32Array,
1054        UInt64Array,
1055    };
1056    use graphforge_core::uuid::Uuid;
1057    use graphforge_core::{OntologyMode, TypeId};
1058    use graphforge_ir::{IrLiteral, RuntimeCatalog};
1059    use parquet::file::properties::WriterProperties;
1060    use tempfile::TempDir;
1061
1062    use super::*;
1063    use crate::{GraphWriter, read_edge_properties, read_nodes, read_properties};
1064
1065    const TS: i64 = 1_700_000_000_000_000;
1066
1067    fn uuid(marker: u8) -> Uuid {
1068        let mut bytes = [0_u8; 16];
1069        bytes[15] = marker;
1070        Uuid::from_bytes(bytes)
1071    }
1072
1073    fn fixture() -> (TempDir, [Uuid; 3], [Uuid; 2]) {
1074        let source = TempDir::new().unwrap();
1075        let nodes = [uuid(3), uuid(1), uuid(2)];
1076        let edges = [uuid(11), uuid(12)];
1077        let mut writer =
1078            GraphWriter::open_at(source.path(), OntologyMode::Exploratory, TS).unwrap();
1079        writer
1080            .create_node_with_labels(nodes[0], &[TypeId(7), TypeId(9)])
1081            .unwrap();
1082        writer.create_node(nodes[1], TypeId(8)).unwrap();
1083        writer.create_node(nodes[2], TypeId(10)).unwrap();
1084        for (index, node) in nodes.iter().enumerate() {
1085            writer
1086                .set_properties(
1087                    node,
1088                    None,
1089                    HashMap::from([("value".into(), IrLiteral::Int(index as i64))]),
1090                )
1091                .unwrap();
1092        }
1093        writer
1094            .create_edge(edges[0], "KNOWS", &nodes[0], &nodes[1])
1095            .unwrap();
1096        writer
1097            .create_edge(edges[1], "KNOWS", &nodes[1], &nodes[2])
1098            .unwrap();
1099        writer
1100            .set_edge_properties(
1101                &edges[0],
1102                Some("KNOWS"),
1103                HashMap::from([("weight".into(), IrLiteral::Float(0.75))]),
1104            )
1105            .unwrap();
1106        writer.flush().unwrap();
1107
1108        let mut catalog = RuntimeCatalog::new();
1109        catalog.intern_label("Person");
1110        catalog.intern_relation_type("KNOWS");
1111        catalog.intern_property("value", Some("Person"));
1112        write_parquet(
1113            &source.path().join("topology/runtime_catalog.parquet"),
1114            &catalog.to_record_batch(),
1115        )
1116        .unwrap();
1117        fs::write(
1118            source.path().join(graphforge_core::manifest::MANIFEST_FILE),
1119            b"ontology: ontology.yaml\n",
1120        )
1121        .unwrap();
1122        fs::write(
1123            source.path().join(graphforge_core::manifest::ONTOLOGY_FILE),
1124            b"version: 1\n",
1125        )
1126        .unwrap();
1127        for excluded in [
1128            "knowledge",
1129            "epistemic",
1130            "provenance",
1131            "valid_time",
1132            "indexes",
1133        ] {
1134            fs::create_dir_all(source.path().join(excluded)).unwrap();
1135            fs::write(
1136                source.path().join(excluded).join("must-not-copy"),
1137                b"secret",
1138            )
1139            .unwrap();
1140        }
1141        (source, nodes, edges)
1142    }
1143
1144    #[test]
1145    fn projection_preserves_graph_rows_closes_endpoints_and_never_induces_edges() {
1146        let (source, nodes, edges) = fixture();
1147        let target = TempDir::new().unwrap();
1148        let summary = materialize_graph_projection(
1149            source.path(),
1150            target.path(),
1151            &GraphProjectionSelection {
1152                node_uuids: BTreeSet::from([*nodes[2].as_bytes()]),
1153                edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1154            },
1155        )
1156        .unwrap();
1157
1158        assert_eq!(
1159            summary.node_uuids,
1160            vec![
1161                *nodes[1].as_bytes(),
1162                *nodes[2].as_bytes(),
1163                *nodes[0].as_bytes(),
1164            ]
1165        );
1166        assert_eq!(
1167            summary.endpoint_node_uuids,
1168            vec![*nodes[1].as_bytes(), *nodes[0].as_bytes()]
1169        );
1170        assert_eq!(summary.edge_uuids, vec![*edges[0].as_bytes()]);
1171
1172        let source_nodes = read_nodes(source.path()).unwrap();
1173        let projected_nodes = read_nodes(target.path()).unwrap();
1174        let source_ids = id_map(&source_nodes[0], "node_uuid", "node_id");
1175        let projected_ids = id_map(&projected_nodes[0], "node_uuid", "node_id");
1176        assert_eq!(projected_ids.len(), 3);
1177        for uuid in &summary.node_uuids {
1178            assert_eq!(projected_ids.get(uuid), source_ids.get(uuid));
1179        }
1180        let labels = projected_nodes[0]
1181            .column_by_name("type_ids")
1182            .unwrap()
1183            .as_any()
1184            .downcast_ref::<ListArray>()
1185            .unwrap();
1186        let selected_row = summary
1187            .node_uuids
1188            .iter()
1189            .position(|uuid| uuid == nodes[0].as_bytes())
1190            .unwrap();
1191        let values = labels.value(selected_row);
1192        assert_eq!(
1193            values
1194                .as_any()
1195                .downcast_ref::<UInt32Array>()
1196                .unwrap()
1197                .values(),
1198            &[7, 9]
1199        );
1200
1201        let projected_edges =
1202            read_parquet(&target.path().join("topology/edges/_exploratory.parquet")).unwrap();
1203        assert_eq!(projected_edges[0].num_rows(), 1);
1204        assert_eq!(
1205            uuid_at(uuid_column(&projected_edges[0], "edge_uuid").unwrap(), 0).unwrap(),
1206            *edges[0].as_bytes()
1207        );
1208        let source_edge_ids = id_map(
1209            &read_parquet(&source.path().join("topology/edges/_exploratory.parquet")).unwrap()[0],
1210            "edge_uuid",
1211            "edge_id",
1212        );
1213        let projected_edge_ids = id_map(&projected_edges[0], "edge_uuid", "edge_id");
1214        assert_eq!(
1215            projected_edge_ids.get(edges[0].as_bytes()),
1216            source_edge_ids.get(edges[0].as_bytes())
1217        );
1218
1219        assert_eq!(
1220            read_properties(target.path(), "_untyped")
1221                .unwrap()
1222                .iter()
1223                .map(RecordBatch::num_rows)
1224                .sum::<usize>(),
1225            3
1226        );
1227        assert_eq!(
1228            read_edge_properties(target.path(), "KNOWS")
1229                .unwrap()
1230                .iter()
1231                .map(RecordBatch::num_rows)
1232                .sum::<usize>(),
1233            1
1234        );
1235        assert!(
1236            target
1237                .path()
1238                .join("topology/runtime_catalog.parquet")
1239                .exists()
1240        );
1241        assert!(
1242            target
1243                .path()
1244                .join(graphforge_core::manifest::MANIFEST_FILE)
1245                .exists()
1246        );
1247        assert!(
1248            target
1249                .path()
1250                .join(graphforge_core::manifest::ONTOLOGY_FILE)
1251                .exists()
1252        );
1253        for excluded in [
1254            "knowledge",
1255            "epistemic",
1256            "provenance",
1257            "valid_time",
1258            "indexes",
1259        ] {
1260            assert!(!target.path().join(excluded).exists());
1261        }
1262    }
1263
1264    #[test]
1265    fn projection_is_canonically_ordered_and_reproducible() {
1266        let (source, nodes, edges) = fixture();
1267        let first = TempDir::new().unwrap();
1268        let second = TempDir::new().unwrap();
1269        let selection = GraphProjectionSelection {
1270            node_uuids: BTreeSet::from([*nodes[2].as_bytes()]),
1271            edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1272        };
1273        let left = materialize_graph_projection(source.path(), first.path(), &selection).unwrap();
1274        let right = materialize_graph_projection(source.path(), second.path(), &selection).unwrap();
1275        assert_eq!(left, right);
1276        assert_ne!(left.graph_content_fingerprint, [0; 32]);
1277        for relative in [
1278            "topology/nodes.parquet",
1279            "topology/edges/_exploratory.parquet",
1280            "properties/_untyped.parquet",
1281            "edge_properties/KNOWS.parquet",
1282            "topology/runtime_catalog.parquet",
1283        ] {
1284            assert_eq!(
1285                fs::read(first.path().join(relative)).unwrap(),
1286                fs::read(second.path().join(relative)).unwrap(),
1287                "non-deterministic output for {relative}"
1288            );
1289        }
1290    }
1291
1292    #[test]
1293    fn projection_fingerprint_ignores_parquet_chunking_and_dictionary_layout() {
1294        let (source, nodes, edges) = fixture();
1295        let baseline_target = TempDir::new().unwrap();
1296        let rewritten_target = TempDir::new().unwrap();
1297        let selection = GraphProjectionSelection {
1298            node_uuids: BTreeSet::from([*nodes[2].as_bytes()]),
1299            edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1300        };
1301        let baseline =
1302            materialize_graph_projection(source.path(), baseline_target.path(), &selection)
1303                .unwrap();
1304
1305        for relative in [
1306            "topology/nodes.parquet",
1307            "topology/edges/_exploratory.parquet",
1308            "properties/_untyped.parquet",
1309            "edge_properties/KNOWS.parquet",
1310            "topology/runtime_catalog.parquet",
1311        ] {
1312            let path = source.path().join(relative);
1313            let batches = read_parquet(&path).unwrap();
1314            let schema = batches[0].schema();
1315            let replacement = path.with_extension("rewritten");
1316            let file = fs::File::create(&replacement).unwrap();
1317            let properties = WriterProperties::builder()
1318                .set_dictionary_enabled(false)
1319                .set_max_row_group_row_count(Some(1))
1320                .build();
1321            let mut writer = ArrowWriter::try_new(file, schema, Some(properties)).unwrap();
1322            for batch in batches {
1323                for row in 0..batch.num_rows() {
1324                    writer.write(&batch.slice(row, 1)).unwrap();
1325                }
1326            }
1327            writer.close().unwrap();
1328            fs::rename(replacement, path).unwrap();
1329        }
1330
1331        let rewritten =
1332            materialize_graph_projection(source.path(), rewritten_target.path(), &selection)
1333                .unwrap();
1334        assert_eq!(
1335            baseline.graph_content_fingerprint,
1336            rewritten.graph_content_fingerprint
1337        );
1338    }
1339
1340    #[test]
1341    fn unrelated_runtime_catalog_entries_do_not_change_projection_identity() {
1342        let (source, nodes, edges) = fixture();
1343        let first = TempDir::new().unwrap();
1344        let second = TempDir::new().unwrap();
1345        let selection = GraphProjectionSelection {
1346            node_uuids: BTreeSet::from([*nodes[0].as_bytes()]),
1347            edge_uuids: BTreeSet::from([*edges[0].as_bytes()]),
1348        };
1349        let baseline =
1350            materialize_graph_projection(source.path(), first.path(), &selection).unwrap();
1351
1352        let catalog_path = source.path().join("topology/runtime_catalog.parquet");
1353        let batch = read_parquet(&catalog_path).unwrap().remove(0);
1354        let mut catalog = RuntimeCatalog::from_record_batch(&batch).unwrap();
1355        catalog.intern_label("Unrelated");
1356        catalog.intern_relation_type("IGNORES");
1357        catalog.intern_property("noise", Some("Unrelated"));
1358        write_parquet(&catalog_path, &catalog.to_record_batch()).unwrap();
1359
1360        let with_noise =
1361            materialize_graph_projection(source.path(), second.path(), &selection).unwrap();
1362        assert_eq!(
1363            baseline.graph_content_fingerprint,
1364            with_noise.graph_content_fingerprint
1365        );
1366        assert_eq!(
1367            fs::read(first.path().join("topology/runtime_catalog.parquet")).unwrap(),
1368            fs::read(second.path().join("topology/runtime_catalog.parquet")).unwrap()
1369        );
1370        let projected = read_parquet(&second.path().join("topology/runtime_catalog.parquet"))
1371            .unwrap()
1372            .remove(0);
1373        let names = string_column(&projected, "name").unwrap();
1374        assert!(!(0..names.len()).any(|row| names.value(row) == "Unrelated"));
1375        assert!(!(0..names.len()).any(|row| names.value(row) == "IGNORES"));
1376        assert!(!(0..names.len()).any(|row| names.value(row) == "noise"));
1377    }
1378
1379    #[test]
1380    fn typed_projection_keeps_exact_owned_catalog_and_reopens_graph_rows() {
1381        let source = TempDir::new().unwrap();
1382        let (alice, bob, excluded) = (uuid(31), uuid(32), uuid(33));
1383        let (knows, ignores) = (uuid(41), uuid(42));
1384        let mut catalog = RuntimeCatalog::new();
1385        let person = catalog.intern_label("Person");
1386        let company = catalog.intern_label("Company");
1387        catalog.intern_relation_type("KNOWS");
1388        catalog.intern_relation_type("IGNORES");
1389        catalog.intern_property("name", Some("Person"));
1390        catalog.intern_property("global", None);
1391        catalog.intern_property("since", Some("KNOWS"));
1392        catalog.intern_property("noise", Some("Company"));
1393
1394        let mut writer = GraphWriter::open_at(source.path(), OntologyMode::Strict, TS).unwrap();
1395        writer.create_node(alice, TypeId(person.0)).unwrap();
1396        writer.create_node(bob, TypeId(person.0)).unwrap();
1397        writer.create_node(excluded, TypeId(company.0)).unwrap();
1398        writer.create_edge(knows, "KNOWS", &alice, &bob).unwrap();
1399        writer
1400            .create_edge(ignores, "IGNORES", &alice, &excluded)
1401            .unwrap();
1402        writer
1403            .set_properties(
1404                &alice,
1405                Some("Person"),
1406                HashMap::from([
1407                    ("name".into(), IrLiteral::Str("Alice".into())),
1408                    ("global".into(), IrLiteral::Bool(true)),
1409                ]),
1410            )
1411            .unwrap();
1412        writer
1413            .set_properties(
1414                &excluded,
1415                Some("Company"),
1416                HashMap::from([("noise".into(), IrLiteral::Str("exclude".into()))]),
1417            )
1418            .unwrap();
1419        writer
1420            .set_edge_properties(
1421                &knows,
1422                Some("KNOWS"),
1423                HashMap::from([("since".into(), IrLiteral::Int(2020))]),
1424            )
1425            .unwrap();
1426        writer.flush().unwrap();
1427        write_parquet(
1428            &source.path().join("topology/runtime_catalog.parquet"),
1429            &catalog.to_record_batch(),
1430        )
1431        .unwrap();
1432
1433        let target = TempDir::new().unwrap();
1434        let summary = materialize_graph_projection(
1435            source.path(),
1436            target.path(),
1437            &GraphProjectionSelection {
1438                node_uuids: BTreeSet::from([*alice.as_bytes()]),
1439                edge_uuids: BTreeSet::from([*knows.as_bytes()]),
1440            },
1441        )
1442        .unwrap();
1443        assert_eq!(summary.node_uuids, vec![*alice.as_bytes(), *bob.as_bytes()]);
1444        assert_eq!(summary.edge_uuids, vec![*knows.as_bytes()]);
1445        assert_eq!(summary.endpoint_node_uuids, vec![*bob.as_bytes()]);
1446
1447        let projected = read_parquet(&target.path().join("topology/runtime_catalog.parquet"))
1448            .unwrap()
1449            .remove(0);
1450        let kinds = string_column(&projected, "entry_kind").unwrap();
1451        let names = string_column(&projected, "name").unwrap();
1452        let owners = string_column(&projected, "owner_label").unwrap();
1453        let inventory = (0..projected.num_rows())
1454            .map(|row| {
1455                (
1456                    kinds.value(row).to_owned(),
1457                    names.value(row).to_owned(),
1458                    (!owners.is_null(row)).then(|| owners.value(row).to_owned()),
1459                )
1460            })
1461            .collect::<BTreeSet<_>>();
1462        assert_eq!(
1463            inventory,
1464            BTreeSet::from([
1465                ("entity_type".into(), "Person".into(), None),
1466                ("relation_type".into(), "KNOWS".into(), None),
1467                ("property".into(), "global".into(), None),
1468                ("property".into(), "name".into(), Some("Person".into())),
1469                ("property".into(), "since".into(), Some("KNOWS".into())),
1470            ])
1471        );
1472
1473        let reopened_nodes = read_nodes(target.path()).unwrap();
1474        assert_eq!(
1475            reopened_nodes
1476                .iter()
1477                .map(RecordBatch::num_rows)
1478                .sum::<usize>(),
1479            2
1480        );
1481        assert_eq!(
1482            read_properties(target.path(), "Person").unwrap()[0].num_rows(),
1483            1
1484        );
1485        assert_eq!(
1486            read_edge_properties(target.path(), "KNOWS").unwrap()[0].num_rows(),
1487            1
1488        );
1489        let projected_again = projected_graph_fingerprint(target.path()).unwrap();
1490        assert_eq!(projected_again, summary.graph_content_fingerprint);
1491    }
1492
1493    #[test]
1494    fn existing_graph_empty_hydrated_workspace_is_a_valid_target() {
1495        let (source, nodes, _) = fixture();
1496        let target = TempDir::new().unwrap();
1497        write_parquet(
1498            &target.path().join("topology/nodes.parquet"),
1499            &RecordBatch::new_empty(Arc::clone(&crate::TOPOLOGY_NODES_SCHEMA)),
1500        )
1501        .unwrap();
1502        write_parquet(
1503            &target.path().join("topology/runtime_catalog.parquet"),
1504            &RuntimeCatalog::new().to_record_batch(),
1505        )
1506        .unwrap();
1507        fs::write(
1508            target.path().join("topology/generation.json"),
1509            b"{\"topology_generation\":0,\"search_generation\":0}\n",
1510        )
1511        .unwrap();
1512
1513        let summary = materialize_graph_projection(
1514            source.path(),
1515            target.path(),
1516            &GraphProjectionSelection {
1517                node_uuids: BTreeSet::from([*nodes[0].as_bytes()]),
1518                edge_uuids: BTreeSet::new(),
1519            },
1520        )
1521        .unwrap();
1522        assert_eq!(summary.node_uuids, vec![*nodes[0].as_bytes()]);
1523        assert_eq!(read_nodes(target.path()).unwrap()[0].num_rows(), 1);
1524        assert!(!target.path().join("topology/generation.json").exists());
1525    }
1526
1527    #[test]
1528    fn empty_projection_target_validation_rejects_nonregular_metadata_and_graph_entries() {
1529        let target = TempDir::new().unwrap();
1530        fs::create_dir(target.path().join(graphforge_core::manifest::MANIFEST_FILE)).unwrap();
1531        assert_eq!(
1532            validate_graph_empty_target(target.path())
1533                .unwrap_err()
1534                .code(),
1535            "GF_VALIDATION"
1536        );
1537
1538        let target = TempDir::new().unwrap();
1539        let properties = target.path().join("properties");
1540        fs::create_dir(&properties).unwrap();
1541        fs::write(properties.join("not-parquet.txt"), b"preserve").unwrap();
1542        assert_eq!(
1543            validate_graph_empty_target(target.path())
1544                .unwrap_err()
1545                .code(),
1546            "GF_VALIDATION"
1547        );
1548        assert_eq!(
1549            fs::read(properties.join("not-parquet.txt")).unwrap(),
1550            b"preserve"
1551        );
1552
1553        let target = TempDir::new().unwrap();
1554        let edges = target.path().join("topology/edges");
1555        fs::create_dir_all(&edges).unwrap();
1556        fs::create_dir(edges.join("nested.parquet")).unwrap();
1557        assert_eq!(
1558            validate_graph_empty_target(target.path())
1559                .unwrap_err()
1560                .code(),
1561            "GF_VALIDATION"
1562        );
1563    }
1564
1565    #[test]
1566    fn missing_identity_and_nonempty_target_fail_before_writing() {
1567        let (source, _, _) = fixture();
1568        let target = TempDir::new().unwrap();
1569        let missing = uuid(99);
1570        let error = materialize_graph_projection(
1571            source.path(),
1572            target.path(),
1573            &GraphProjectionSelection {
1574                node_uuids: BTreeSet::from([*missing.as_bytes()]),
1575                edge_uuids: BTreeSet::new(),
1576            },
1577        )
1578        .unwrap_err();
1579        assert!(matches!(error, GfError::Validation(_)));
1580        assert!(fs::read_dir(target.path()).unwrap().next().is_none());
1581
1582        fs::write(target.path().join("owned"), b"keep").unwrap();
1583        let error = materialize_graph_projection(
1584            source.path(),
1585            target.path(),
1586            &GraphProjectionSelection::default(),
1587        )
1588        .unwrap_err();
1589        assert!(matches!(error, GfError::Validation(_)));
1590        assert_eq!(fs::read(target.path().join("owned")).unwrap(), b"keep");
1591    }
1592
1593    #[test]
1594    fn projection_rejects_overlapping_and_nonempty_targets_without_mutation() {
1595        let (source, _, _) = fixture();
1596        let empty = GraphProjectionSelection::default();
1597        let same_error =
1598            materialize_graph_projection(source.path(), source.path(), &empty).unwrap_err();
1599        assert_eq!(same_error.code(), "GF_VALIDATION");
1600        assert!(same_error.to_string().contains("must be disjoint"));
1601
1602        let child = source.path().join("projection-child");
1603        fs::create_dir(&child).unwrap();
1604        fs::write(child.join("sentinel"), b"child").unwrap();
1605        let child_error = materialize_graph_projection(source.path(), &child, &empty).unwrap_err();
1606        assert_eq!(child_error.code(), "GF_VALIDATION");
1607        assert_eq!(fs::read(child.join("sentinel")).unwrap(), b"child");
1608
1609        let ancestor = TempDir::new().unwrap();
1610        let nested_source = ancestor.path().join("source");
1611        fs::create_dir(&nested_source).unwrap();
1612        let ancestor_error =
1613            materialize_graph_projection(&nested_source, ancestor.path(), &empty).unwrap_err();
1614        assert_eq!(ancestor_error.code(), "GF_VALIDATION");
1615        assert!(nested_source.exists());
1616
1617        let regular_root = TempDir::new().unwrap();
1618        let regular_target = regular_root.path().join("target");
1619        fs::write(&regular_target, b"regular").unwrap();
1620        assert_eq!(
1621            materialize_graph_projection(source.path(), &regular_target, &empty)
1622                .unwrap_err()
1623                .code(),
1624            "GF_VALIDATION"
1625        );
1626        assert_eq!(fs::read(&regular_target).unwrap(), b"regular");
1627
1628        let unexpected = TempDir::new().unwrap();
1629        fs::write(unexpected.path().join("knowledge.parquet"), b"owned").unwrap();
1630        assert_eq!(
1631            materialize_graph_projection(source.path(), unexpected.path(), &empty)
1632                .unwrap_err()
1633                .code(),
1634            "GF_VALIDATION"
1635        );
1636        assert_eq!(
1637            fs::read(unexpected.path().join("knowledge.parquet")).unwrap(),
1638            b"owned"
1639        );
1640
1641        let bad_topology = TempDir::new().unwrap();
1642        fs::create_dir(bad_topology.path().join("topology")).unwrap();
1643        fs::write(bad_topology.path().join("topology/unknown"), b"keep").unwrap();
1644        assert_eq!(
1645            materialize_graph_projection(source.path(), bad_topology.path(), &empty)
1646                .unwrap_err()
1647                .code(),
1648            "GF_VALIDATION"
1649        );
1650        assert_eq!(
1651            fs::read(bad_topology.path().join("topology/unknown")).unwrap(),
1652            b"keep"
1653        );
1654
1655        let bad_properties = TempDir::new().unwrap();
1656        fs::create_dir(bad_properties.path().join("properties")).unwrap();
1657        fs::write(
1658            bad_properties.path().join("properties/not-parquet"),
1659            b"keep",
1660        )
1661        .unwrap();
1662        assert_eq!(
1663            materialize_graph_projection(source.path(), bad_properties.path(), &empty)
1664                .unwrap_err()
1665                .code(),
1666            "GF_VALIDATION"
1667        );
1668        assert_eq!(
1669            fs::read(bad_properties.path().join("properties/not-parquet")).unwrap(),
1670            b"keep"
1671        );
1672
1673        let nonempty = TempDir::new().unwrap();
1674        let mut node_uuid = FixedSizeBinaryBuilder::new(16);
1675        node_uuid.append_value(uuid(90).as_bytes()).unwrap();
1676        let batch = RecordBatch::try_from_iter([
1677            ("node_uuid", Arc::new(node_uuid.finish()) as ArrayRef),
1678            ("value", Arc::new(Int64Array::from(vec![1])) as ArrayRef),
1679        ])
1680        .unwrap();
1681        let nonempty_path = nonempty.path().join("properties/Person.parquet");
1682        write_parquet(&nonempty_path, &batch).unwrap();
1683        let before = fs::read(&nonempty_path).unwrap();
1684        assert_eq!(
1685            materialize_graph_projection(source.path(), nonempty.path(), &empty)
1686                .unwrap_err()
1687                .code(),
1688            "GF_VALIDATION"
1689        );
1690        assert_eq!(fs::read(&nonempty_path).unwrap(), before);
1691    }
1692
1693    #[test]
1694    fn corrupt_property_uuid_is_rejected_instead_of_silently_dropped() {
1695        let source = TempDir::new().unwrap();
1696        let target = TempDir::new().unwrap();
1697        let mut uuids = FixedSizeBinaryBuilder::new(16);
1698        uuids.append_null();
1699        let batch = RecordBatch::try_from_iter([
1700            ("node_uuid", Arc::new(uuids.finish()) as ArrayRef),
1701            ("value", Arc::new(Int64Array::from(vec![1])) as ArrayRef),
1702        ])
1703        .unwrap();
1704        let source_path = source.path().join("properties/Person.parquet");
1705        write_parquet(&source_path, &batch).unwrap();
1706
1707        let error = project_parquet_file(
1708            &source_path,
1709            &target.path().join("properties/Person.parquet"),
1710            "node_uuid",
1711            &BTreeSet::new(),
1712        )
1713        .unwrap_err();
1714        assert!(matches!(error, GfError::Validation(_)));
1715        assert!(error.to_string().contains("UUID column contains null"));
1716        assert!(!target.path().join("properties/Person.parquet").exists());
1717    }
1718
1719    #[test]
1720    fn canonical_graph_encoding_normalizes_floats_timezones_and_nested_types() {
1721        assert_eq!(normalize_f32(f32::NAN), 0x7fc0_0000);
1722        assert_eq!(normalize_f32(-0.0), 0);
1723        assert_eq!(normalize_f64(f64::NAN), 0x7ff8_0000_0000_0000);
1724        assert_eq!(normalize_f64(-0.0), 0);
1725        assert_eq!(time_unit_tag(TimeUnit::Second), 0);
1726        assert_eq!(time_unit_tag(TimeUnit::Millisecond), 1);
1727        assert_eq!(time_unit_tag(TimeUnit::Microsecond), 2);
1728        assert_eq!(time_unit_tag(TimeUnit::Nanosecond), 3);
1729        for timezone in [
1730            None,
1731            Some("UTC"),
1732            Some("Etc/UTC"),
1733            Some("Z"),
1734            Some("+00:00"),
1735        ] {
1736            assert!(validate_timezone(timezone).is_ok());
1737        }
1738        assert_eq!(
1739            validate_timezone(Some("America/Denver"))
1740                .unwrap_err()
1741                .code(),
1742            "GF_VALIDATION"
1743        );
1744
1745        let supported = [
1746            DataType::Boolean,
1747            DataType::Int32,
1748            DataType::Int64,
1749            DataType::UInt32,
1750            DataType::UInt64,
1751            DataType::Float32,
1752            DataType::Float64,
1753            DataType::Utf8,
1754            DataType::LargeUtf8,
1755            DataType::Binary,
1756            DataType::LargeBinary,
1757            DataType::FixedSizeBinary(16),
1758            DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1759            DataType::Time64(TimeUnit::Nanosecond),
1760            DataType::List(Arc::new(Field::new("item", DataType::Utf8, true))),
1761            DataType::FixedSizeList(Arc::new(Field::new("item", DataType::UInt64, false)), 2),
1762            DataType::Struct(vec![Field::new("name", DataType::Utf8, false)].into()),
1763            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
1764        ];
1765        for data_type in supported {
1766            let mut writer = CanonicalWriter::new();
1767            encode_type(&mut writer, &data_type).unwrap();
1768            assert!(!writer.finish().is_empty());
1769        }
1770        let mut writer = CanonicalWriter::new();
1771        assert_eq!(
1772            encode_type(&mut writer, &DataType::Date32)
1773                .unwrap_err()
1774                .code(),
1775            "GF_VALIDATION"
1776        );
1777
1778        let micros: ArrayRef = Arc::new(TimestampMicrosecondArray::from(vec![TS]));
1779        assert_eq!(
1780            timestamp_value(&micros, TimeUnit::Microsecond, 0).unwrap(),
1781            TS
1782        );
1783        let time: ArrayRef = Arc::new(Time64MicrosecondArray::from(vec![123_i64]));
1784        assert_eq!(time64_value(&time, TimeUnit::Microsecond, 0).unwrap(), 123);
1785        assert_eq!(
1786            time64_value(&time, TimeUnit::Second, 0).unwrap_err().code(),
1787            "GF_VALIDATION"
1788        );
1789    }
1790
1791    #[test]
1792    fn canonical_value_encoding_traverses_every_supported_nested_arrow_shape() {
1793        use arrow::array::{
1794            FixedSizeListArray, Int32Array, LargeListArray, StringDictionaryBuilder, StructArray,
1795        };
1796        use arrow::datatypes::Int32Type;
1797
1798        let large: ArrayRef = Arc::new(LargeListArray::from_iter_primitive::<Int32Type, _, _>([
1799            Some(vec![Some(1), None, Some(2)]),
1800        ]));
1801        let fixed: ArrayRef = Arc::new(FixedSizeListArray::from_iter_primitive::<Int32Type, _, _>(
1802            [Some(vec![Some(3), Some(4)])],
1803            2,
1804        ));
1805        let struct_fields: arrow::datatypes::Fields =
1806            vec![Field::new("value", DataType::Int32, false)].into();
1807        let structure: ArrayRef = Arc::new(StructArray::new(
1808            struct_fields.clone(),
1809            vec![Arc::new(Int32Array::from(vec![5]))],
1810            None,
1811        ));
1812        let mut dictionary_builder = StringDictionaryBuilder::<Int32Type>::new();
1813        dictionary_builder.append("six").unwrap();
1814        let dictionary: ArrayRef = Arc::new(dictionary_builder.finish());
1815
1816        for (data_type, array) in [
1817            (large.data_type().clone(), large),
1818            (fixed.data_type().clone(), fixed),
1819            (DataType::Struct(struct_fields), structure),
1820            (dictionary.data_type().clone(), dictionary),
1821        ] {
1822            let mut writer = CanonicalWriter::new();
1823            encode_present_value(&mut writer, &data_type, &array, 0).unwrap();
1824            assert!(!writer.finish().is_empty());
1825        }
1826    }
1827
1828    #[test]
1829    fn canonical_value_encoding_rejects_type_mismatch_and_nonnullable_null() {
1830        let floats: ArrayRef = Arc::new(Float32Array::from(vec![Some(-0.0), Some(f32::NAN)]));
1831        let doubles: ArrayRef = Arc::new(Float64Array::from(vec![Some(-0.0), Some(f64::NAN)]));
1832        let strings: ArrayRef = Arc::new(StringArray::from(vec![Some("value"), None]));
1833        let mut writer = CanonicalWriter::new();
1834        encode_present_value(&mut writer, &DataType::Float32, &floats, 0).unwrap();
1835        encode_present_value(&mut writer, &DataType::Float32, &floats, 1).unwrap();
1836        encode_present_value(&mut writer, &DataType::Float64, &doubles, 0).unwrap();
1837        encode_present_value(&mut writer, &DataType::Float64, &doubles, 1).unwrap();
1838        encode_value(&mut writer, &DataType::Utf8, &strings, 0, false).unwrap();
1839        encode_value(&mut writer, &DataType::Utf8, &strings, 1, true).unwrap();
1840        assert!(!writer.finish().is_empty());
1841
1842        let mut writer = CanonicalWriter::new();
1843        assert_eq!(
1844            encode_value(&mut writer, &DataType::Utf8, &strings, 1, false)
1845                .unwrap_err()
1846                .code(),
1847            "GF_VALIDATION"
1848        );
1849        let mut writer = CanonicalWriter::new();
1850        assert_eq!(
1851            encode_present_value(&mut writer, &DataType::UInt64, &strings, 0)
1852                .unwrap_err()
1853                .code(),
1854            "GF_VALIDATION"
1855        );
1856    }
1857
1858    #[test]
1859    fn canonical_value_encoding_covers_every_scalar_and_time_representation() {
1860        let values: Vec<(DataType, ArrayRef)> = vec![
1861            (DataType::Boolean, Arc::new(BooleanArray::from(vec![true]))),
1862            (DataType::Int32, Arc::new(Int32Array::from(vec![-7]))),
1863            (DataType::Int64, Arc::new(Int64Array::from(vec![-9]))),
1864            (DataType::UInt32, Arc::new(UInt32Array::from(vec![7]))),
1865            (DataType::UInt64, Arc::new(UInt64Array::from(vec![9]))),
1866            (DataType::Utf8, Arc::new(StringArray::from(vec!["small"]))),
1867            (
1868                DataType::LargeUtf8,
1869                Arc::new(LargeStringArray::from(vec!["large"])),
1870            ),
1871            (
1872                DataType::Binary,
1873                Arc::new(BinaryArray::from_vec(vec![b"small".as_slice()])),
1874            ),
1875            (
1876                DataType::LargeBinary,
1877                Arc::new(LargeBinaryArray::from_vec(vec![b"large".as_slice()])),
1878            ),
1879            (
1880                DataType::Timestamp(TimeUnit::Second, None),
1881                Arc::new(TimestampSecondArray::from(vec![1_i64])),
1882            ),
1883            (
1884                DataType::Timestamp(TimeUnit::Millisecond, Some("Z".into())),
1885                Arc::new(TimestampMillisecondArray::from(vec![2_i64]).with_timezone("Z")),
1886            ),
1887            (
1888                DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1889                Arc::new(TimestampMicrosecondArray::from(vec![3_i64]).with_timezone("UTC")),
1890            ),
1891            (
1892                DataType::Timestamp(TimeUnit::Nanosecond, Some("Etc/UTC".into())),
1893                Arc::new(TimestampNanosecondArray::from(vec![4_i64]).with_timezone("Etc/UTC")),
1894            ),
1895            (
1896                DataType::Time64(TimeUnit::Microsecond),
1897                Arc::new(Time64MicrosecondArray::from(vec![5_i64])),
1898            ),
1899            (
1900                DataType::Time64(TimeUnit::Nanosecond),
1901                Arc::new(Time64NanosecondArray::from(vec![6_i64])),
1902            ),
1903        ];
1904        let mut encodings = Vec::new();
1905        for (data_type, array) in values {
1906            let mut writer = CanonicalWriter::new();
1907            encode_present_value(&mut writer, &data_type, &array, 0).unwrap();
1908            let encoded = writer.finish();
1909            assert!(!encoded.is_empty(), "{data_type} must emit canonical bytes");
1910            encodings.push(encoded);
1911        }
1912        assert_eq!(encodings.len(), 15);
1913
1914        let seconds: ArrayRef = Arc::new(TimestampSecondArray::from(vec![11_i64]));
1915        let millis: ArrayRef = Arc::new(TimestampMillisecondArray::from(vec![12_i64]));
1916        let nanos: ArrayRef = Arc::new(TimestampNanosecondArray::from(vec![13_i64]));
1917        assert_eq!(timestamp_value(&seconds, TimeUnit::Second, 0).unwrap(), 11);
1918        assert_eq!(
1919            timestamp_value(&millis, TimeUnit::Millisecond, 0).unwrap(),
1920            12
1921        );
1922        assert_eq!(
1923            timestamp_value(&nanos, TimeUnit::Nanosecond, 0).unwrap(),
1924            13
1925        );
1926        let time_nanos: ArrayRef = Arc::new(Time64NanosecondArray::from(vec![14_i64]));
1927        assert_eq!(
1928            time64_value(&time_nanos, TimeUnit::Nanosecond, 0).unwrap(),
1929            14
1930        );
1931    }
1932
1933    #[test]
1934    fn projection_identity_and_path_validation_matrix_fails_before_mutation() {
1935        let one = [1_u8; 16];
1936        let two = [2_u8; 16];
1937        assert!(require_present(&BTreeSet::new(), &BTreeSet::new(), "node").is_ok());
1938        assert!(
1939            require_present(&BTreeSet::from([one]), &BTreeSet::from([one, two]), "node").is_ok()
1940        );
1941        assert!(
1942            require_present(&BTreeSet::from([two]), &BTreeSet::from([one]), "edge")
1943                .unwrap_err()
1944                .to_string()
1945                .contains("missing edge UUID")
1946        );
1947
1948        let root = TempDir::new().unwrap();
1949        let source = root.path().join("source");
1950        let sibling = root.path().join("sibling");
1951        std::fs::create_dir(&source).unwrap();
1952        std::fs::create_dir(&sibling).unwrap();
1953        assert!(validate_distinct_paths(&source, &sibling).is_ok());
1954        assert!(validate_distinct_paths(&source, &source).is_err());
1955        assert!(validate_distinct_paths(&source, &source.join("child")).is_err());
1956        assert!(validate_distinct_paths(&source.join("child"), &source).is_err());
1957
1958        assert!(
1959            uuid_rows(&source.join("missing.parquet"), "node_uuid")
1960                .unwrap()
1961                .is_empty()
1962        );
1963        let wrong = RecordBatch::try_from_iter([(
1964            "node_uuid",
1965            Arc::new(UInt64Array::from(vec![1_u64])) as ArrayRef,
1966        )])
1967        .unwrap();
1968        assert!(uuid_column(&wrong, "node_uuid").is_err());
1969        assert!(uuid_column(&wrong, "missing").is_err());
1970
1971        let mut nullable = FixedSizeBinaryBuilder::new(16);
1972        nullable.append_null();
1973        let nullable = nullable.finish();
1974        assert!(uuid_at(&nullable, 0).is_err());
1975    }
1976
1977    #[test]
1978    fn wave10_projection_private_bounds_and_missing_inventory_are_exact() {
1979        assert!(
1980            sorted_parquet_files(Path::new("definitely-absent"))
1981                .unwrap()
1982                .is_empty()
1983        );
1984        assert!(exact_u32(usize::MAX, "field count").is_err());
1985        assert!(validate_timezone(Some("America/Denver")).is_err());
1986
1987        let values: ArrayRef = Arc::new(arrow::array::Int64Array::from(vec![1]));
1988        assert!(time64_value(&values, TimeUnit::Second, 0).is_err());
1989        assert!(downcast::<arrow::array::StringArray>(&values).is_err());
1990    }
1991
1992    #[cfg(unix)]
1993    #[test]
1994    fn wave10_projection_inventory_rejects_symbolic_links() {
1995        use std::os::unix::fs::symlink;
1996
1997        let root = tempfile::tempdir().unwrap();
1998        let target = root.path().join("target.parquet");
1999        fs::write(&target, b"caller").unwrap();
2000        symlink(&target, root.path().join("linked.parquet")).unwrap();
2001        assert!(sorted_parquet_files(root.path()).is_err());
2002        assert_eq!(fs::read(target).unwrap(), b"caller");
2003    }
2004
2005    #[test]
2006    fn wave13_projection_rejects_duplicate_graph_identities() {
2007        let root = TempDir::new().unwrap();
2008        let duplicate = uuid(41);
2009        let mut node_uuids = FixedSizeBinaryBuilder::new(16);
2010        node_uuids.append_value(duplicate.as_bytes()).unwrap();
2011        node_uuids.append_value(duplicate.as_bytes()).unwrap();
2012        let nodes =
2013            RecordBatch::try_from_iter([("node_uuid", Arc::new(node_uuids.finish()) as ArrayRef)])
2014                .unwrap();
2015        let nodes_path = root.path().join("nodes.parquet");
2016        write_parquet(&nodes_path, &nodes).unwrap();
2017        assert!(uuid_rows(&nodes_path, "node_uuid").is_err());
2018
2019        let mut edge_uuids = FixedSizeBinaryBuilder::new(16);
2020        let mut sources = FixedSizeBinaryBuilder::new(16);
2021        let mut targets = FixedSizeBinaryBuilder::new(16);
2022        for _ in 0..2 {
2023            edge_uuids.append_value(duplicate.as_bytes()).unwrap();
2024            sources.append_value(uuid(42).as_bytes()).unwrap();
2025            targets.append_value(uuid(43).as_bytes()).unwrap();
2026        }
2027        let edges = RecordBatch::try_from_iter([
2028            ("edge_uuid", Arc::new(edge_uuids.finish()) as ArrayRef),
2029            ("src_uuid", Arc::new(sources.finish()) as ArrayRef),
2030            ("dst_uuid", Arc::new(targets.finish()) as ArrayRef),
2031        ])
2032        .unwrap();
2033        let edges_path = root.path().join("edges.parquet");
2034        write_parquet(&edges_path, &edges).unwrap();
2035        assert!(edge_endpoints(&[edges_path]).is_err());
2036    }
2037
2038    #[test]
2039    fn wave13_projection_path_shape_and_cleanup_guards_are_structured() {
2040        let root = TempDir::new().unwrap();
2041        let missing = root.path().join("missing.parquet");
2042        assert!(
2043            project_parquet_file(
2044                &missing,
2045                &root.path().join("unused.parquet"),
2046                "node_uuid",
2047                &BTreeSet::new(),
2048            )
2049            .is_ok()
2050        );
2051        assert!(copy_regular_file_if_present(&missing, &root.path().join("copy")).is_ok());
2052        assert!(clear_graph_empty_target(&root.path().join("absent-target")).is_ok());
2053
2054        let metadata_directory = root.path().join("metadata-directory");
2055        fs::create_dir(&metadata_directory).unwrap();
2056        assert!(
2057            copy_regular_file_if_present(&metadata_directory, &root.path().join("metadata-copy"))
2058                .is_err()
2059        );
2060
2061        let source_file = root.path().join("manifest-source");
2062        let target_file = root.path().join("nested/manifest-copy");
2063        fs::write(&source_file, b"manifest").unwrap();
2064        copy_regular_file_if_present(&source_file, &target_file).unwrap();
2065        assert_eq!(fs::read(&target_file).unwrap(), b"manifest");
2066
2067        let target = root.path().join("clear-target");
2068        for directory in ["topology", "properties", "edge_properties"] {
2069            fs::create_dir_all(target.join(directory)).unwrap();
2070        }
2071        for name in [
2072            graphforge_core::manifest::MANIFEST_FILE,
2073            graphforge_core::manifest::ONTOLOGY_FILE,
2074        ] {
2075            fs::write(target.join(name), b"metadata").unwrap();
2076        }
2077        clear_graph_empty_target(&target).unwrap();
2078        assert!(fs::read_dir(&target).unwrap().next().is_none());
2079    }
2080
2081    #[test]
2082    fn wave13_projection_target_metadata_must_be_regular_files() {
2083        let target = TempDir::new().unwrap();
2084        fs::create_dir(target.path().join(graphforge_core::manifest::MANIFEST_FILE)).unwrap();
2085        assert!(validate_graph_empty_target(target.path()).is_err());
2086
2087        let topology = TempDir::new().unwrap();
2088        fs::create_dir(topology.path().join("generation.json")).unwrap();
2089        assert!(validate_empty_topology(topology.path()).is_err());
2090
2091        let graph_directory = TempDir::new().unwrap();
2092        fs::create_dir(graph_directory.path().join("nested.parquet")).unwrap();
2093        assert!(validate_empty_parquet_directory(graph_directory.path()).is_err());
2094
2095        assert_ne!(normalize_f32(1.25), 0);
2096        assert_ne!(normalize_f64(1.25), 0);
2097        assert_eq!(
2098            dictionary_value_type(&DataType::Dictionary(
2099                Box::new(DataType::Int32),
2100                Box::new(DataType::Utf8),
2101            )),
2102            &DataType::Utf8
2103        );
2104    }
2105
2106    fn id_map(batch: &RecordBatch, uuid_name: &str, id_name: &str) -> BTreeMap<[u8; 16], u64> {
2107        let uuids = uuid_column(batch, uuid_name).unwrap();
2108        let ids = batch
2109            .column_by_name(id_name)
2110            .unwrap()
2111            .as_any()
2112            .downcast_ref::<UInt64Array>()
2113            .unwrap();
2114        (0..batch.num_rows())
2115            .map(|row| (uuid_at(uuids, row).unwrap(), ids.value(row)))
2116            .collect()
2117    }
2118}