Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

9 changes: 7 additions & 2 deletions cargo-bazel-lock.json
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
{
"checksum": "6a40a06e42157027df4fb5b43e623e477cd821b17bd6a2670363a11b508a094b",
"checksum": "594b6ba3c63b9701d278637b815cac9903f9c34e2906b60fd6018eb89b58575b",
"crates": {
"adler2 2.0.1": {
"name": "adler2",
Expand Down Expand Up @@ -14087,6 +14087,10 @@
{
"id": "wait-timeout 0.2.1",
"target": "wait_timeout"
},
{
"id": "zstd 0.13.3",
"target": "zstd"
}
],
"selects": {}
Expand Down Expand Up @@ -34676,7 +34680,8 @@
"codspeed-divan-compat 5.0.1",
"cucumber 0.21.1",
"insta 1.48.0",
"wait-timeout 0.2.1"
"wait-timeout 0.2.1",
"zstd 0.13.3"
],
"unused_patches": []
}
6 changes: 5 additions & 1 deletion crates/graphforge-api/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -476,7 +476,11 @@ gf_rust_integration_test(

gf_rust_integration_test(
name = "multi_ontology_certification",
srcs = ["tests/multi_ontology_certification.rs"],
srcs = [
"tests/multi_ontology_certification.rs",
"src/permanent_parquet_test_support.rs",
],
crate_root = "tests/multi_ontology_certification.rs",
crate = ":graphforge_api",
data = _API_TEST_DATA + ["//tests/fixtures/multi-ontology-v1:certification_files"],
rustc_env = {
Expand Down
3 changes: 3 additions & 0 deletions crates/graphforge-api/src/belief_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1708,6 +1708,9 @@ mod tests {
policy: policy.clone(),
})
.unwrap();
crate::permanent_parquet_test_support::assert_file(
&projection.graph.dir.join("topology/nodes.parquet"),
);
assert_ne!(projection.source_generation_uuid(), Uuid::nil());
assert!(!projection.policy_bytes().is_empty());
assert_ne!(projection.policy_fingerprint(), [0; 32]);
Expand Down
16 changes: 16 additions & 0 deletions crates/graphforge-api/src/checkpoints.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3053,6 +3053,22 @@ mod tests {
actor_uuid: None,
})
.unwrap_or_else(|error| panic!("{shape}: {error:?}"));
crate::permanent_parquet_test_support::assert_participants(
directory.path(),
graphforge_storage::WORKSPACE_CAPABILITY_ID,
);
if matches!(shape, "knowledge" | "epistemic") {
crate::permanent_parquet_test_support::assert_participants(
directory.path(),
"knowledge",
);
}
if shape == "epistemic" {
crate::permanent_parquet_test_support::assert_participants(
directory.path(),
"epistemic",
);
}
drop(graph);

let reopened = GraphForge::new(Some(directory.path().to_str().unwrap())).unwrap();
Expand Down
9 changes: 7 additions & 2 deletions crates/graphforge-api/src/knowledge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2984,8 +2984,12 @@ pub(crate) fn snapshot_to_participant(
}

fn write_parquet(batch: &RecordBatch, schema: &SchemaRef) -> Result<Vec<u8>, GfError> {
let mut writer = ArrowWriter::try_new(Vec::new(), Arc::clone(schema), None)
.map_err(|error| GfError::Storage(error.to_string()))?;
let mut writer = ArrowWriter::try_new(
Vec::new(),
Arc::clone(schema),
Some(graphforge_storage::permanent_parquet::writer_properties().build()),
)
.map_err(|error| GfError::Storage(error.to_string()))?;
writer
.write(batch)
.map_err(|error| GfError::Storage(error.to_string()))?;
Expand Down Expand Up @@ -3729,6 +3733,7 @@ mod tests {

let created = graph.create_assertion(request.clone()).unwrap();
assert_eq!(created.stats.rows_produced, 1);
crate::permanent_parquet_test_support::assert_participants(root.path(), "knowledge");
let cancelled = CancellationToken::new();
cancelled.cancel();
assert_eq!(
Expand Down
2 changes: 2 additions & 0 deletions crates/graphforge-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,8 @@ mod multi_ontology;
mod mutation_transaction;
#[cfg(test)]
mod mutation_transaction_fault_tests;
#[cfg(test)]
mod permanent_parquet_test_support;
pub use multi_ontology::{
ActivationProfileChangeRequest, BridgeAdoptionRequest, BridgeCandidate, BridgeDeleteRequest,
BridgeUpdateRequest, CompositionValidationReceipt, ModuleAdoptionRequest, ModuleCandidate,
Expand Down
88 changes: 88 additions & 0 deletions crates/graphforge-api/src/permanent_parquet_test_support.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
//! Inspect authenticated published payloads in existing public lifecycle tests.

use std::path::Path;

use graphforge_storage::ResolvedProjectGeneration;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use parquet::file::metadata::ParquetMetaData;

fn assert_metadata(metadata: &ParquetMetaData, name: &str) -> usize {
let mut columns = 0;
for group in metadata.row_groups() {
assert!(
group.num_rows() <= 1_048_576,
"{name}: default row-group ceiling"
);
for column in group.columns() {
assert!(
matches!(column.compression(), parquet::basic::Compression::ZSTD(_)),
"{name}: permanent column {} lost Zstd",
column.column_path()
);
columns += 1;
}
}
columns
}

pub(crate) fn assert_file(path: &Path) {
let reader =
ParquetRecordBatchReaderBuilder::try_new(std::fs::File::open(path).unwrap()).unwrap();
assert!(assert_metadata(reader.metadata(), &path.display().to_string()) > 0);
}

pub(crate) fn assert_graph(generation: &ResolvedProjectGeneration) {
let inventory = generation.graph_files_inventory().unwrap().unwrap();
let owned = generation
.declared_graph_files_inventory()
.unwrap()
.is_some();
let mut columns = 0;
for entry in inventory
.files
.iter()
.filter(|entry| entry.relative_path.ends_with(".parquet"))
{
let path = if owned {
generation.graph_tree_root().join(&entry.relative_path)
} else {
graphforge_storage::graph_object_path(
generation.container_root(),
&entry.content_sha256,
)
.unwrap()
};
let reader =
ParquetRecordBatchReaderBuilder::try_new(std::fs::File::open(path).unwrap()).unwrap();
columns += assert_metadata(reader.metadata(), &entry.relative_path);
}
assert!(columns > 0, "published graph fixture must contain data");
}

pub(crate) fn assert_participants(root: &Path, capability: &str) {
let generation = graphforge_storage::resolve_project_generation(root).unwrap();
let mut columns = 0;
for participant in generation
.participant_snapshots()
.unwrap()
.into_iter()
.filter(|participant| {
participant.capability_id == capability && participant.encoding == "parquet"
})
{
let reader =
ParquetRecordBatchReaderBuilder::try_new(bytes::Bytes::from(participant.bytes))
.unwrap();
columns += assert_metadata(reader.metadata(), &participant.record_family_id);
if participant.record_family_id == "restoration_transition" {
assert_eq!(
reader.metadata().file_metadata().created_by(),
Some("graphforge-restoration-transition/1")
);
}
}
assert!(
columns > 0,
"{capability}: public fixture must publish Parquet data"
);
}
10 changes: 8 additions & 2 deletions crates/graphforge-api/src/provenance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -280,8 +280,12 @@ fn participant(
}

fn write_parquet(batch: &RecordBatch, schema: &SchemaRef) -> Result<Vec<u8>, GfError> {
let mut writer = ArrowWriter::try_new(Vec::new(), Arc::clone(schema), None)
.map_err(|error| GfError::Storage(error.to_string()))?;
let mut writer = ArrowWriter::try_new(
Vec::new(),
Arc::clone(schema),
Some(graphforge_storage::permanent_parquet::writer_properties().build()),
)
.map_err(|error| GfError::Storage(error.to_string()))?;
writer
.write(batch)
.map_err(|error| GfError::Storage(error.to_string()))?;
Expand Down Expand Up @@ -551,6 +555,8 @@ mod tests {
.unwrap()
.is_some()
);
crate::permanent_parquet_test_support::assert_participants(root.path(), "provenance");
crate::permanent_parquet_test_support::assert_graph(&generation);
let ledger = read_ledger(&generation).unwrap();
assert_eq!(ledger.events.len(), 1);
assert_eq!(ledger.events[0].event_kind, EventKind::CreateNode);
Expand Down
3 changes: 3 additions & 0 deletions crates/graphforge-api/src/search_index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -782,6 +782,9 @@ mod tests {
.unwrap();
let replaced = current_search_artifact(&graph.dir, &key).unwrap().unwrap();
assert_ne!(replaced.path, first.path);
crate::permanent_parquet_test_support::assert_file(
&replaced.path.join(graphforge_storage::VECTOR_DATA_FILE),
);
let rows = read_vector_snapshot(&replaced.path, 2, VectorStoreLimits::default(), || Ok(()))
.unwrap();
assert_eq!(rows.len(), 1);
Expand Down
7 changes: 7 additions & 0 deletions crates/graphforge-api/tests/multi_ontology_certification.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
//! Public Rust certification for the reproducible six-domain M9 fixture.

#[path = "../src/permanent_parquet_test_support.rs"]
#[allow(dead_code)]
mod permanent_parquet_test_support;

use std::collections::BTreeMap;

use graphforge_api::{
Expand Down Expand Up @@ -343,6 +347,9 @@ fn retained_data_migrates_atomically_and_exact_identity_reopens() {
let receipt = graph
.migrate_ontology_module(&request, &preview, None)
.expect("publish retained-data migration");
permanent_parquet_test_support::assert_graph(
&graphforge_storage::resolve_project_generation(project.path()).unwrap(),
);
let replay = graph
.migrate_ontology_module(&request, &preview, None)
.expect("exact idempotent replay");
Expand Down
Loading
Loading