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 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(®ular_target, b"regular").unwrap();
1620 assert_eq!(
1621 materialize_graph_projection(source.path(), ®ular_target, &empty)
1622 .unwrap_err()
1623 .code(),
1624 "GF_VALIDATION"
1625 );
1626 assert_eq!(fs::read(®ular_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(µs, 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}