1use 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#[derive(Clone, Debug, Default, Eq, PartialEq)]
31pub struct GraphProjectionSelection {
32 pub node_uuids: BTreeSet<[u8; 16]>,
34 pub edge_uuids: BTreeSet<[u8; 16]>,
36}
37
38#[derive(Clone, Debug, Eq, PartialEq)]
40pub struct GraphProjectionSummary {
41 pub node_uuids: Vec<[u8; 16]>,
43 pub edge_uuids: Vec<[u8; 16]>,
45 pub endpoint_node_uuids: Vec<[u8; 16]>,
47 pub graph_content_fingerprint: [u8; 32],
49}
50
51pub 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}