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
2 changes: 0 additions & 2 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -335,8 +335,6 @@ jobs:
runtime_entity_labels::tests::catalog_persistence_replaces_read_only_cas_aliases
route_component::tests::
route_component::owned::tests::legacy_cas_materialization_translates_routes_before_creating_paths
uuid_membership::tests::owned_manifest_
uuid_membership::tests::readonly_shared_runs_preserve_uuid_snapshot_authentication
graph_files::tests::graph_file_copy_
adjacency::builder::tests::admitted_reserved_relations_build_portable_distinct_indexes
mutator::tests::mapped_reserved_delete_preserves_other_routes_and_reopens
Expand Down
6 changes: 0 additions & 6 deletions crates/graphforge-api/src/branches/import_graph.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,12 +55,6 @@ fn incorporate_with_policy(
return Err(conflict());
}
}
if !graphforge_storage::uuid_membership_index_present(&destination.dir()) {
graphforge_storage::rebuild_uuid_membership_indexes(
&destination.dir(),
graphforge_storage::UuidIndexBuildLimits::default(),
)?;
}
let mut labels = BTreeMap::<String, Vec<IrLiteral>>::new();
for (kind, query) in [
(
Expand Down
181 changes: 85 additions & 96 deletions crates/graphforge-api/src/bulk_construction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -282,123 +282,119 @@ impl ValidatedBulkEdges {
}
}

fn open_membership_index(
impl GraphForge {
/// The identity probe for the committed topology generation, cached on the
/// facade until the generation moves. It reads the published node and edge
/// Parquet, so a lookup is bounded by the candidates, not the graph.
pub(crate) fn cached_identity_probe(
&self,
) -> Result<
std::sync::MutexGuard<'_, Option<graphforge_storage::TopologyIdentityProbe>>,
graphforge_core::GfError,
> {
let current_generation = graphforge_storage::read_topology_generation(&self.dir())?;
let mut cached = self.identity_probe.lock().map_err(|_| {
graphforge_core::GfError::Storage("identity probe lock poisoned".into())
})?;
if cached
.as_ref()
.is_some_and(|probe| probe.topology_generation() != current_generation)
{
*cached = None;
}
if cached.is_none() {
let dir = self.dir();
let files = dir.topology_files()?;
*cached = Some(graphforge_storage::TopologyIdentityProbe::open(
&dir,
&files,
current_generation,
)?);
}
Ok(cached)
}
}

fn open_identity_probe(
graph: &GraphForge,
input_kind: BulkInputKind,
) -> Result<
std::sync::MutexGuard<'_, Option<graphforge_storage::UuidMembershipIndex>>,
std::sync::MutexGuard<'_, Option<graphforge_storage::TopologyIdentityProbe>>,
BulkValidationError,
> {
let current_generation =
graphforge_storage::read_topology_generation(&graph.dir()).map_err(|error| {
contract_error(
input_kind,
BulkValidationReason::ProjectState,
&error.to_string(),
)
})?;
let mut cached = graph.uuid_membership_index.lock().map_err(|error| {
graph.cached_identity_probe().map_err(|error| {
contract_error(
input_kind,
BulkValidationReason::ProjectState,
&error.to_string(),
)
})?;
if cached
.as_ref()
.is_some_and(|index| index.topology_generation() != current_generation)
{
*cached = None;
}
if !graphforge_storage::uuid_membership_index_present(&graph.dir()) {
let has_nodes =
graphforge_storage::node_topology_present(&graph.dir()).map_err(|error| {
contract_error(
input_kind,
BulkValidationReason::ProjectState,
&error.to_string(),
)
})?;
let has_edges = std::fs::read_dir(graph.dir().join("topology/edges"))
.ok()
.is_some_and(|mut entries| entries.any(|entry| entry.is_ok()));
if has_nodes || has_edges {
return Err(contract_error(
input_kind,
BulkValidationReason::ProjectState,
"UUID membership index is missing; run the bounded storage rebuild before ingest",
));
}
return Ok(cached);
}
if cached.is_none() {
*cached = Some(
graphforge_storage::UuidMembershipIndex::open(&graph.dir()).map_err(|error| {
contract_error(
input_kind,
BulkValidationReason::ProjectState,
&error.to_string(),
)
})?,
);
}
Ok(cached)
})
}

fn existing_edge_context(
graph: &GraphForge,
endpoint_candidates: &[Uuid],
edge_candidates: Option<&[Uuid]>,
) -> Result<(HashSet<Uuid>, HashSet<Uuid>), BulkValidationError> {
let mut index = open_membership_index(graph, BulkInputKind::Edge)?;
let known_nodes = indexed_existing(
index.as_mut(),
endpoint_candidates,
graphforge_storage::UuidIndexKind::Node,
BulkInputKind::Edge,
)?;
let mut probe = open_identity_probe(graph, BulkInputKind::Edge)?;
let known_nodes = live_nodes(probe.as_mut(), endpoint_candidates, BulkInputKind::Edge)?;
let Some(edge_candidates) = edge_candidates else {
return Ok((known_nodes, HashSet::new()));
};
let mut existing = indexed_existing(
index.as_mut(),
edge_candidates,
graphforge_storage::UuidIndexKind::Edge,
BulkInputKind::Edge,
)?;
existing.extend(indexed_existing(
index.as_mut(),
edge_candidates,
graphforge_storage::UuidIndexKind::Node,
BulkInputKind::Edge,
)?);
let existing = taken_identities(probe.as_mut(), edge_candidates, BulkInputKind::Edge)?;
Ok((known_nodes, existing))
}

/// Probe `candidates` (sorted, deduplicated) and return the subset the index
/// already holds. Membership only: callers never iterate the result.
fn indexed_existing(
index: Option<&mut graphforge_storage::UuidMembershipIndex>,
/// Probe `candidates` (sorted, deduplicated) and return the live nodes among
/// them. Membership only: callers never iterate the result.
fn live_nodes(
probe: Option<&mut graphforge_storage::TopologyIdentityProbe>,
candidates: &[Uuid],
input_kind: BulkInputKind,
) -> Result<HashSet<Uuid>, BulkValidationError> {
let Some(probe) = probe else {
return Ok(HashSet::new());
};
let (found, _) = probe
.probe(graphforge_storage::UuidIndexKind::Node, candidates)
.map_err(|error| {
contract_error(
input_kind,
BulkValidationReason::ProjectState,
&error.to_string(),
)
})?;
Ok(selected(candidates, &found))
}

/// The subset of `candidates` that is already spent: a live node, a live edge
/// or a deleted entity. Node and edge UUIDs share one namespace and a deleted
/// UUID is never reused, exactly as the commit enforces.
fn taken_identities(
probe: Option<&mut graphforge_storage::TopologyIdentityProbe>,
candidates: &[Uuid],
index_kind: graphforge_storage::UuidIndexKind,
input_kind: BulkInputKind,
) -> Result<HashSet<Uuid>, BulkValidationError> {
let Some(index) = index else {
let Some(probe) = probe else {
return Ok(HashSet::new());
};
let (found, _) = index.probe(index_kind, candidates).map_err(|error| {
let (found, _) = probe.taken(candidates).map_err(|error| {
contract_error(
input_kind,
BulkValidationReason::ProjectState,
&error.to_string(),
)
})?;
Ok(candidates
Ok(selected(candidates, &found))
}

fn selected(candidates: &[Uuid], found: &[bool]) -> HashSet<Uuid> {
candidates
.iter()
.copied()
.zip(found)
.filter_map(|(uuid, present)| present.then_some(uuid))
.collect())
.collect()
}

/// Non-null UUIDs of `field`, sorted and deduplicated: the same set, in the
Expand Down Expand Up @@ -444,27 +440,21 @@ fn candidate_endpoint_uuids(batches: &[RecordBatch]) -> Result<Vec<Uuid>, BulkVa

#[cfg(test)]
fn indexed_uuid_count(graph: &GraphForge, kind: graphforge_storage::UuidIndexKind) -> u64 {
graphforge_storage::UuidMembershipIndex::open(&graph.dir())
.expect("published graph has an authenticated UUID membership index")
.count(kind)
let dir = graph.dir();
let files = dir.topology_files().expect("published topology files");
graphforge_storage::TopologyIdentityProbe::open(
&dir,
&files,
graphforge_storage::read_topology_generation(&dir).expect("topology generation"),
)
.expect("published graph has readable topology")
.count(kind)
}

pub(crate) fn register_existing_endpoints(
writer: &mut graphforge_storage::GraphWriter,
dir: &std::path::Path,
endpoints: &BTreeSet<Uuid>,
) -> Result<(), super::GfError> {
if !graphforge_storage::uuid_membership_index_present(dir) {
graphforge_storage::rebuild_uuid_membership_indexes(
dir,
graphforge_storage::UuidIndexBuildLimits::default(),
)?;
}
if !graphforge_storage::uuid_membership_index_is_fresh(dir)? {
return Err(super::GfError::Storage(
"bulk endpoint UUID index is stale".into(),
));
}
let requested = endpoints.iter().copied().collect::<Vec<_>>();
writer
.register_existing_endpoints(&requested)
Expand Down Expand Up @@ -580,8 +570,7 @@ mod tests {
.unwrap();
let missing = uuid(70_001);
let failure =
register_existing_endpoints(&mut writer, &graph.dir(), &BTreeSet::from([missing]))
.unwrap_err();
register_existing_endpoints(&mut writer, &BTreeSet::from([missing])).unwrap_err();
assert_eq!(
failure.to_string(),
"validation error: bulk edge endpoint disappeared before publication"
Expand Down
17 changes: 3 additions & 14 deletions crates/graphforge-api/src/bulk_construction/normalization.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ use super::{
RecordBatch, Schema, SchemaRef, Sha256, StringArray, StructArray, Time64NanosecondArray,
TimestampMicrosecondArray, UInt8Array, UInt16Array, UInt32Array, Uuid, ValidatedBulkEdges,
ValidatedBulkNodes, batch_error, candidate_endpoint_uuids, candidate_uuids, contract_error,
existing_edge_context, field_error, indexed_existing, open_membership_index, row_error,
existing_edge_context, field_error, open_identity_probe, row_error, taken_identities,
};

const NODE_REQUIRED: [(&str, DataType, bool); 2] = [
Expand Down Expand Up @@ -1095,19 +1095,8 @@ impl GraphForge {
let mut existing = HashSet::new();
if reject_existing {
let candidates = candidate_uuids(batches, BulkInputKind::Node, "node_uuid")?;
let mut index = open_membership_index(self, BulkInputKind::Node)?;
existing = indexed_existing(
index.as_mut(),
&candidates,
graphforge_storage::UuidIndexKind::Node,
BulkInputKind::Node,
)?;
existing.extend(indexed_existing(
index.as_mut(),
&candidates,
graphforge_storage::UuidIndexKind::Edge,
BulkInputKind::Node,
)?);
let mut probe = open_identity_probe(self, BulkInputKind::Node)?;
existing = taken_identities(probe.as_mut(), &candidates, BulkInputKind::Node)?;
}
let mut observed = HashSet::new();
let mut rows = Vec::new();
Expand Down
39 changes: 10 additions & 29 deletions crates/graphforge-api/src/bulk_construction/publication.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,8 @@ use super::{
Arc, BTreeSet, BulkEdgePublicationError, BulkEdgeRow, BulkInputKind, BulkNodePublicationError,
BulkNodeRow, BulkValidationReason, DataType, Digest, Field, FixedSizeBinaryArray, GraphForge,
HashMap, OperationId, RecordBatch, Schema, SchemaRef, Sha256, StringArray, UInt64Array, Uuid,
ValidatedBulkNodes, contract_metadata, indexed_existing, open_membership_index,
register_existing_endpoints, row_error,
ValidatedBulkNodes, contract_metadata, open_identity_probe, register_existing_endpoints,
row_error, taken_identities,
};

fn bulk_node_generation_uuid(operation_uuid: OperationId, rows: &[BulkNodeRow]) -> Uuid {
Expand Down Expand Up @@ -221,19 +221,8 @@ impl GraphForge {
let mut candidates = normalized.identities().collect::<Vec<_>>();
candidates.sort_unstable();
candidates.dedup();
let mut index = open_membership_index(self, BulkInputKind::Node)?;
let mut existing = indexed_existing(
index.as_mut(),
&candidates,
graphforge_storage::UuidIndexKind::Node,
BulkInputKind::Node,
)?;
existing.extend(indexed_existing(
index.as_mut(),
&candidates,
graphforge_storage::UuidIndexKind::Edge,
BulkInputKind::Node,
)?);
let mut probe = open_identity_probe(self, BulkInputKind::Node)?;
let existing = taken_identities(probe.as_mut(), &candidates, BulkInputKind::Node)?;
if normalized
.rows
.iter()
Expand Down Expand Up @@ -408,19 +397,8 @@ impl GraphForge {
.collect::<Vec<_>>();
candidates.sort_unstable();
candidates.dedup();
let mut index = open_membership_index(self, BulkInputKind::Edge)?;
let mut existing = indexed_existing(
index.as_mut(),
&candidates,
graphforge_storage::UuidIndexKind::Edge,
BulkInputKind::Edge,
)?;
existing.extend(indexed_existing(
index.as_mut(),
&candidates,
graphforge_storage::UuidIndexKind::Node,
BulkInputKind::Edge,
)?);
let mut probe = open_identity_probe(self, BulkInputKind::Edge)?;
let existing = taken_identities(probe.as_mut(), &candidates, BulkInputKind::Edge)?;
if let Some(row) = normalized
.rows
.iter()
Expand Down Expand Up @@ -457,7 +435,7 @@ impl GraphForge {
.iter()
.flat_map(|row| [row.source_uuid, row.target_uuid])
.collect::<BTreeSet<_>>();
register_existing_endpoints(&mut writer, &self.dir(), &endpoints)?;
register_existing_endpoints(&mut writer, &endpoints)?;
for row in &normalized.rows {
next_catalog.intern_relation_type(&row.rel_type)?;
writer.create_edge(
Expand Down Expand Up @@ -547,3 +525,6 @@ impl GraphForge {

#[cfg(test)]
mod tests;

#[cfg(test)]
mod deleted_identity_tests;
Loading