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
93 changes: 93 additions & 0 deletions crates/graphforge-api/src/resumable_construction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -498,6 +498,99 @@ mod tests {
.unwrap()
}

/// A published construction ships its adjacency CSR (#1388): a fresh
/// process finds it current in the hydrated workspace without rebuilding,
/// hop queries answer from it, and a corrupted published shard is refused
/// by the open-time digest sweep like every other graph object.
#[test]
fn published_construction_serves_its_adjacency_index_and_refuses_corruption() {
use graphforge_storage::adjacency::AdjacencyFreshnessState;

let directory = tempfile::TempDir::new().unwrap();
let path = directory.path().to_str().unwrap();
let graph = GraphForge::new(Some(path)).unwrap();
let node_ids: Vec<Uuid> = (0..4).map(|_| Uuid::now_v7()).collect();
let edge_ids: Vec<Uuid> = (0..3).map(|_| Uuid::now_v7()).collect();
let mut session = graph.begin_graph_construction(Default::default()).unwrap();
session.append_nodes("nodes", &nodes(&node_ids)).unwrap();
session
.append_edges(
"edges",
&edges(
&edge_ids,
&[
(node_ids[0], node_ids[1]),
(node_ids[1], node_ids[2]),
(node_ids[2], node_ids[3]),
],
),
)
.unwrap();
session.seal_and_publish().unwrap();
drop(session);
let inspection = graph.inspect_adjacency().unwrap();
assert_eq!(inspection.state, AdjacencyFreshnessState::Current);
assert_eq!(inspection.artifact_source_generation, Some(1));
drop(graph);

let reopened = GraphForge::new(Some(path)).unwrap();
assert!(graphforge_storage::adjacency::manifest_path(&reopened.dir()).is_file());
let inspection = reopened.inspect_adjacency().unwrap();
assert_eq!(inspection.state, AdjacencyFreshnessState::Current);
assert_eq!(inspection.artifact_source_generation, Some(1));
assert_construction_relationships(
&reopened,
&[
[node_ids[0], edge_ids[0], node_ids[1]],
[node_ids[1], edge_ids[1], node_ids[2]],
[node_ids[2], edge_ids[2], node_ids[3]],
],
);
let hops = reopened
.execute("MATCH (a)-[r]->(b)-[s]->(c) RETURN count(*) AS n")
.unwrap();
assert_eq!(
hops.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::Int64Array>()
.unwrap()
.value(0),
2
);
drop(reopened);

let generation = graphforge_storage::resolve_project_generation(directory.path()).unwrap();
let inventory = generation.graph_files_inventory().unwrap().unwrap();
let shard = inventory
.files
.iter()
.find(|entry| {
entry.relative_path.starts_with("indexes/adjacency/")
&& entry.relative_path.ends_with(".csr")
})
.expect("published CSR shard");
let object = graphforge_storage::graph_object_path(
generation.container_root(),
&shard.content_sha256,
)
.unwrap();
let mut permissions = std::fs::metadata(&object).unwrap().permissions();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt as _;
permissions.set_mode(0o600);
}
#[cfg(not(unix))]
permissions.set_readonly(false);
std::fs::set_permissions(&object, permissions).unwrap();
let mut bytes = std::fs::read(&object).unwrap();
bytes[0] ^= 1;
std::fs::write(&object, bytes).unwrap();
let error = GraphForge::new(Some(path)).unwrap_err();
assert!(error.to_string().contains("digest"), "{error}");
}

#[test]
fn cancelled_public_construction_preserves_current() {
let graph = GraphForge::new(None).unwrap();
Expand Down
43 changes: 31 additions & 12 deletions crates/graphforge-api/tests/permanent_storage_budgets.rs
Original file line number Diff line number Diff line change
Expand Up @@ -955,17 +955,22 @@ fn assess(f: Fixture) {
construct(&source, f, &nodes, &edges);
let construction_ns = started.elapsed().as_nanos();
let graph = GraphForge::new(source.to_str()).unwrap();
let unindexed = graph.storage_attribution().unwrap();
assert_eq!(
unindexed.categories[&graphforge_storage::ArtifactCategory::Adjacency].allocated_bytes,
0
// ADR 0037: construction publishes the adjacency CSR with the generation, so
// a freshly constructed project is already indexed. This previously asserted
// the category was empty, which was true only while adjacency did not exist
// until some reader built it.
let constructed = graph.storage_attribution().unwrap();
assert!(
constructed.categories[&graphforge_storage::ArtifactCategory::Adjacency].allocated_bytes
> 0,
"construction must publish adjacency with the generation (ADR 0037)"
);
validate_storage_budgets(
Fixture {
adjacency: false,
adjacency: true,
..f
},
&unindexed,
&constructed,
);
// Candidate codecs inspect construction's published payload generation.
// Adjacency publication may switch to a generation-owned graph tree.
Expand All @@ -983,7 +988,12 @@ fn assess(f: Fixture) {
let identity_experiment = identity_padding_experiment(&source);
let manifest_experiment = bucket_manifest_experiment(&source);
if f.heterogeneous {
assert_eq!(manifest_experiment["authenticated_lookups_checked"], 352);
// 352 before ADR 0037. Construction now publishes the adjacency CSR with
// the generation, so `indexes/adjacency/**` are graph files, the manifest
// Merkle tree covers them, and the authenticated walk visits their leaves
// and the interior nodes above them. The count is a pinned observation of
// that tree's shape, not an invariant, so it moves when the file set does.
assert_eq!(manifest_experiment["authenticated_lookups_checked"], 373);
assert!(
manifest_experiment["source_manifest_allocated_bytes"]
.as_u64()
Expand All @@ -998,18 +1008,27 @@ fn assess(f: Fixture) {
let adjacency_experiment = f.adjacency.then(|| adjacency_codec_experiment(&source));
let storage = graph.storage_attribution().unwrap();
storage.validate_for_qualification().unwrap();
validate_storage_budgets(f, &storage);
assert_eq!(
validate_storage_budgets(
Fixture {
adjacency: true,
..f
},
&storage,
);
// Before ADR 0037 this asserted presence == `f.adjacency`, because an
// unbuilt project carried no adjacency. Construction now always publishes
// it, so the remaining property is that an explicit rebuild does not
// remove it and does not push the category past its budget.
assert!(
storage.categories[&graphforge_storage::ArtifactCategory::Adjacency].allocated_bytes > 0,
f.adjacency,
"unbuilt and built adjacency are distinct measured capabilities"
"a constructed project carries adjacency whether or not it was rebuilt"
);
drop(graph);
let whole_project = whole_project_file_census(&source);
let fingerprint = round_trip(root.path(), &source, f, &nodes, &edges);
println!(
"PERMANENT_STORAGE_ASSESSMENT {}",
json!({"fixture":f.name,"nodes":f.nodes,"edges":f.edges,"routes":f.routes,"random_ids":matches!(f.identifiers, Identifiers::Random),"properties":f.properties,"heterogeneous_schemas":f.heterogeneous,"adjacency_built":f.adjacency,"unindexed_allocated_bytes":unindexed.allocated_bytes,"unindexed_categories":unindexed.categories,
json!({"fixture":f.name,"nodes":f.nodes,"edges":f.edges,"routes":f.routes,"random_ids":matches!(f.identifiers, Identifiers::Random),"properties":f.properties,"heterogeneous_schemas":f.heterogeneous,"adjacency_built":f.adjacency,"constructed_allocated_bytes":constructed.allocated_bytes,"constructed_categories":constructed.categories,
"construction_elapsed_ns":construction_ns,"semantic_fingerprint":fingerprint,"whole_project":whole_project,
"permanent": {"logical_bytes":storage.logical_bytes,"physical_logical_bytes":storage.physical_logical_bytes,"allocated_bytes":storage.allocated_bytes,"physical_objects":storage.physical_objects,"categories":storage.categories},
"adjacency_codec_experiment":adjacency_experiment,"parquet_experiment":experiment,"identity_padding_experiment":identity_experiment,"manifest_bucket_experiment":manifest_experiment})
Expand Down
34 changes: 27 additions & 7 deletions crates/graphforge-storage/src/adjacency.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@

mod builder;
mod codec;
pub(crate) use builder::build_adjacency_index_for_edge_files;
pub use builder::{
ADJACENCY_SPILL_DIR_NAME, AdjacencyBuildMetrics, AdjacencyBuildOptions,
DEFAULT_ADJACENCY_CHUNK_ROWS, DEFAULT_ADJACENCY_MERGE_FAN_IN, build_adjacency_index,
Expand Down Expand Up @@ -1170,21 +1171,40 @@ fn capture_adjacency_inventory(
crate::AuthenticatedPropertyInventory::from_inventory_at_root(root, inventory, None)
}

/// The `(relation, path)` edge tables an adjacency build streams, resolved
/// through the admitted inventory when one is supplied, else the legacy raw
/// `topology/edges/` layout.
pub(crate) fn resolve_adjacency_edge_files(
project_dir: &Path,
inventory: Option<&crate::AuthenticatedPropertyInventory>,
) -> Result<Vec<(String, PathBuf)>, GfError> {
match inventory {
Some(inventory) => Ok(inventory.edge_files(None)),
None => crate::mutator::edge_parquet_files(project_dir, None),
}
}

fn for_each_adjacency_edge_file(
project_dir: &Path,
inventory: Option<&crate::AuthenticatedPropertyInventory>,
batch_size: usize,
on_batch: &mut dyn FnMut(&str, bool, &RecordBatch) -> Result<(), GfError>,
) -> Result<(), GfError> {
let paths = match inventory {
Some(inventory) => inventory.edge_files(None),
None => crate::mutator::edge_parquet_files(project_dir, None)?,
};
let paths = resolve_adjacency_edge_files(project_dir, inventory)?;
for_each_adjacency_edge_path(&paths, batch_size, on_batch)
}

/// Stream every `(relation, path)` edge table in `paths` as projected batches.
pub(crate) fn for_each_adjacency_edge_path(
paths: &[(String, PathBuf)],
batch_size: usize,
on_batch: &mut dyn FnMut(&str, bool, &RecordBatch) -> Result<(), GfError>,
) -> Result<(), GfError> {
for (stem, path) in paths {
// An unreadable edge file must FAIL the build, not be skipped: a
// manifest written without it would stamp the current generation and
// make an index missing a relation's edges look fresh.
let _schema = match crate::catalog::discover_parquet_schema_detailed(&path) {
let _schema = match crate::catalog::discover_parquet_schema_detailed(path) {
Ok(schema) => schema,
Err(detail) => {
return Err(GfError::Storage(format!(
Expand All @@ -1199,8 +1219,8 @@ fn for_each_adjacency_edge_file(
} else {
&["edge_id", "src_id", "dst_id"]
};
stream_projected_parquet_batches(&path, columns, batch_size, &mut |batch| {
on_batch(&stem, exploratory, &batch)
stream_projected_parquet_batches(path, columns, batch_size, &mut |batch| {
on_batch(stem, exploratory, &batch)
})?;
}
Ok(())
Expand Down
41 changes: 33 additions & 8 deletions crates/graphforge-storage/src/adjacency/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
use super::{
ALL_RELATIONS_STEM, AdjacencyManifestRow, BuildEntry, DEFAULT_ADJACENCY_BATCH_SIZE,
DEFAULT_CSR_SHARD_EDGES, DEFAULT_CSR_SHARD_NODES, Direction, ShardedCsrWriter, adjacency_dir,
csr_path, for_each_adjacency_edge_file, named_column, storage_err, string_column,
csr_path, for_each_adjacency_edge_path, named_column, storage_err, string_column,
uint64_column, usable_stem, write_manifest,
};
use graphforge_core::GfError;
Expand Down Expand Up @@ -242,6 +242,34 @@ pub fn build_adjacency_index_from_inventory(
checkpoint()?;
// Generation BEFORE the scan — see the race note in the doc comment.
let generation = crate::generation::read_topology_generation(source_project_dir)?;
let edge_files = super::resolve_adjacency_edge_files(source_project_dir, inventory)?;
build_adjacency_index_for_edge_files(
artifact_project_dir,
&edge_files,
generation,
built_at_micros,
options,
checkpoint,
)
}

/// Build the index from an explicit `(relation, path)` edge-table list already
/// admitted by the caller, stamping `topology_generation` into the manifest.
///
/// This is the publish-side entry (#1388): construction encoding and portable
/// import hand over the exact edge tables of the generation they are about to
/// publish, so the CSR ships inside that generation instead of being rebuilt
/// by every query process. The relation names are the semantic routes the
/// caller resolved; no inventory or counter file is consulted here.
pub(crate) fn build_adjacency_index_for_edge_files(
artifact_project_dir: &Path,
edge_files: &[(String, PathBuf)],
topology_generation: u64,
built_at_micros: i64,
options: &AdjacencyBuildOptions,
mut checkpoint: impl FnMut() -> Result<(), GfError>,
) -> Result<(Vec<AdjacencyManifestRow>, AdjacencyBuildMetrics), GfError> {
let generation = topology_generation;
let options = options.effective();

let adjacency = adjacency_dir(artifact_project_dir);
Expand All @@ -261,8 +289,7 @@ pub fn build_adjacency_index_from_inventory(

let build_result = (|| {
let mut groups = stream_build_groups(
source_project_dir,
inventory,
edge_files,
&options,
&mut spill,
&mut metrics,
Expand Down Expand Up @@ -782,8 +809,7 @@ fn merge_keyed_runs(
}

fn stream_build_groups(
project_dir: &Path,
inventory: Option<&crate::AuthenticatedPropertyInventory>,
edge_files: &[(String, PathBuf)],
options: &AdjacencyBuildOptions,
spill: &mut SpillSession,
metrics: &mut AdjacencyBuildMetrics,
Expand All @@ -797,9 +823,8 @@ fn stream_build_groups(
EntryGroup::with_label(ALL_RELATIONS_STEM),
);

for_each_adjacency_edge_file(
project_dir,
inventory,
for_each_adjacency_edge_path(
edge_files,
options.batch_size,
&mut |stem, exploratory, batch| {
checkpoint()?;
Expand Down
Loading