src.dackar.knowledge_graph.kg_schema_builder_workflow¶
Attributes¶
Classes¶
In-memory accumulator for nodes and edges before bulk Neo4j ingestion. |
Functions¶
|
Load a TOML schema file and return its contents as a dict. |
|
Load one or more TOML schema files and merge them into a single schema dict. |
|
Extract the property list from a node or relation schema spec. |
|
Determine the primary key property name for a node spec. |
|
Generate a list of label name variants to try when resolving a schema label. |
|
Resolve the first candidate name that exists as a node label in the schema. |
|
Resolve the first candidate name that exists as a node label in the schema. |
|
Build a lookup from relation type name to its (from_entity, to_entity) pair. |
|
Return True if value is a Neo4j-native scalar type (str, int, float, or bool). |
|
Convert an arbitrary Python value to a Neo4j-compatible property value. |
|
Sanitize all values in a property dict, dropping |
|
Generate Neo4j DDL statements (constraints and indexes) from a merged schema. |
|
Load schemas and apply all DDL constraints and indexes to the database. |
|
Build a namespaced node id by prepending prefix to value. |
|
Return True if value is non-empty and not a sentinel null-like string. |
|
Normalise a text string into a stable, lowercase key. |
|
Add an edge using preferred relation type, falling back to fallback if not in schema. |
|
ProcessedTextRecord -> entity relation chooser. |
|
Add a ProcessedTextRecord → failure_mode edge using the best available relation type. |
|
Translate all RCA workflow artifacts into a graph of nodes and edges. |
|
Bulk-upsert a node list and edge list into Neo4j. |
Module Contents¶
- src.dackar.knowledge_graph.kg_schema_builder_workflow.load_toml_schema(path)[source]¶
Load a TOML schema file and return its contents as a dict.
- Parameters:
path (Union[str, pathlib.Path]) – Filesystem path to the
.tomlschema file.- Returns:
Parsed TOML document as a nested dictionary.
- Return type:
Dict[str, Any]
- src.dackar.knowledge_graph.kg_schema_builder_workflow.load_and_merge_schemas(schema_paths)[source]¶
Load one or more TOML schema files and merge them into a single schema dict.
The merged dict has top-level keys
"node"and"relation". Duplicate keys across files raise an error to prevent silent overwrites.- Parameters:
schema_paths (Union[str, pathlib.Path, Iterable[Union[str, pathlib.Path]]]) – A single path or an iterable of paths to
.tomlfiles.- Returns:
A merged schema dict with
{"node": {...}, "relation": {...}}.- Raises:
ValueError – If the same node or relation key appears in more than one file.
- Return type:
Dict[str, Any]
- src.dackar.knowledge_graph.kg_schema_builder_workflow._schema_props(spec)[source]¶
Extract the property list from a node or relation schema spec.
Checks both
"node_properties"and"properties"keys for compatibility with different TOML schema conventions.- Parameters:
spec (Dict[str, Any]) – A single node or relation entry from the merged schema.
- Returns:
List of property descriptor dicts, or an empty list if none are defined.
- Return type:
List[Dict[str, Any]]
- src.dackar.knowledge_graph.kg_schema_builder_workflow._node_primary_key(spec)[source]¶
Determine the primary key property name for a node spec.
Selects the first non-optional property; falls back to
"id"if present in the property list, then to the first listed property name, and finally to the hard-coded default"id".- Parameters:
spec (Dict[str, Any]) – A single node entry from the merged schema.
- Returns:
The property name to use as the primary key.
- Return type:
str
- src.dackar.knowledge_graph.kg_schema_builder_workflow._label_candidates(name)[source]¶
Generate a list of label name variants to try when resolving a schema label.
Produces the original name, its lowercase form, a PascalCase conversion, and a snake_case conversion (from PascalCase input) to account for naming conventions used across different TOML schemas.
- Parameters:
name (str) – Base label name.
- Returns:
List of candidate strings (original, lowercase, PascalCase, snake_case), deduplicated while preserving order.
- Return type:
List[str]
- src.dackar.knowledge_graph.kg_schema_builder_workflow.resolve_node_label(schema, *candidates)[source]¶
Resolve the first candidate name that exists as a node label in the schema.
For each candidate, tries the original name and common variants (lowercase, PascalCase, snake_case). If none match, returns the first candidate unchanged as a safe fallback.
- Parameters:
schema (Dict[str, Any]) – Merged schema dict (must contain a
"node"key).*candidates (str) – One or more preferred label names, in priority order.
- Returns:
The matched label string from the schema, or candidates[0] if no match is found.
- Return type:
str
- src.dackar.knowledge_graph.kg_schema_builder_workflow.resolve_node_label_strict(schema, *candidates)[source]¶
Resolve the first candidate name that exists as a node label in the schema.
Identical to
resolve_node_label()but raisesKeyErrorinstead of silently falling back when no candidate matches. Use this wherever a miss should be a hard error (e.g. building thelabelsregistry at graph construction time).- Parameters:
schema (Dict[str, Any]) – Merged schema dict (must contain a
"node"key).*candidates (str) – One or more preferred label names, in priority order.
- Returns:
The matched label string from the schema.
- Raises:
KeyError – If no variant of any candidate is found in the schema’s node registry. The error message lists the available node labels so the caller can diagnose schema/TOML mismatches immediately.
- Return type:
str
- src.dackar.knowledge_graph.kg_schema_builder_workflow.relation_endpoint_map(schema)[source]¶
Build a lookup from relation type name to its (from_entity, to_entity) pair.
Only relations that declare both
from_entityandto_entityin the schema are included.- Parameters:
schema (Dict[str, Any]) – Merged schema dict (must contain a
"relation"key).- Returns:
Dict mapping each relation name to a
(source_entity, target_entity)tuple of entity type strings.- Return type:
Dict[str, Tuple[str, str]]
- src.dackar.knowledge_graph.kg_schema_builder_workflow._is_primitive(value)[source]¶
Return True if value is a Neo4j-native scalar type (str, int, float, or bool).
- Parameters:
value (Any) – Any Python object.
- Returns:
Trueif value can be stored directly as a Neo4j property scalar.- Return type:
bool
- src.dackar.knowledge_graph.kg_schema_builder_workflow.sanitize_value(value)[source]¶
Convert an arbitrary Python value to a Neo4j-compatible property value.
Conversion rules: -
datetime/date→ ISO-8601 string. -dict→ JSON string (sorted keys, UTF-8). -listof primitives → kept as-is; mixed/complex lists → JSON string. -Noneand primitive scalars → returned unchanged. - Anything else →str(value).- Parameters:
value (Any) – Python value to sanitize.
- Returns:
A value safe for storage as a Neo4j node or relationship property.
- Return type:
Any
- src.dackar.knowledge_graph.kg_schema_builder_workflow.sanitize_props(props)[source]¶
Sanitize all values in a property dict, dropping
Noneentries.- Parameters:
props (Dict[str, Any]) – Raw property dict potentially containing non-Neo4j-native types.
- Returns:
New dict with all values converted by
sanitize_value()and keys whose value wasNoneremoved.- Return type:
Dict[str, Any]
- src.dackar.knowledge_graph.kg_schema_builder_workflow.generate_ddl_from_schema(schema)[source]¶
Generate Neo4j DDL statements (constraints and indexes) from a merged schema.
For every node label a
UNIQUEconstraint onidis created. Additionally, aCREATE INDEXstatement is emitted for each property that carries"indexed": truein its spec.- Parameters:
schema (Dict[str, Any]) – Merged schema dict as returned by
load_and_merge_schemas().- Returns:
List of Cypher DDL strings ready to be executed against Neo4j.
- Raises:
ValueError – If a node label or indexed-property name is not a safe Neo4j identifier (see
dackar.knowledge_graph.py2neo._safe_token()).- Return type:
List[str]
- src.dackar.knowledge_graph.kg_schema_builder_workflow.apply_schema_constraints(client, schema_paths, database=None)[source]¶
Load schemas and apply all DDL constraints and indexes to the database.
Combines
load_and_merge_schemas(),generate_ddl_from_schema(), andPy2Neo.query()into a single convenience call.- Parameters:
client (dackar.knowledge_graph.py2neo.Py2Neo) – Active
Py2Neoconnection.schema_paths (Union[str, pathlib.Path, Iterable[Union[str, pathlib.Path]]]) – One or more paths to TOML schema files.
database (Optional[str]) – Target database name; uses the driver default when
None.
- Return type:
None
- class src.dackar.knowledge_graph.kg_schema_builder_workflow.GraphBatch(schema=None)[source]¶
In-memory accumulator for nodes and edges before bulk Neo4j ingestion.
Nodes are keyed by their
idstring; duplicate additions are merged. Edges are keyed by(source_id, target_id, rel_type)and support optional schema-based endpoint validation.- Parameters:
schema (Optional[Dict[str, Any]])
- add_node(node_id, label, attrs=None)[source]¶
Add or merge a node into the batch.
If a node with node_id already exists its attributes are updated with the new values (shallow merge after sanitization).
- Parameters:
node_id (str) – Unique identifier for the node (used as the
idproperty).label (str) – Neo4j label for the node.
attrs (Optional[Dict[str, Any]]) – Optional property dict;
idis always set from node_id.
- Returns:
The node_id string, for convenience when chaining calls.
- Return type:
str
- add_edge(src, dst, rel_type, attrs=None, allow_untyped=True)[source]¶
Add or update a directed edge in the batch.
Endpoint labels are looked up from already-added nodes. When rel_type is declared in the schema its expected endpoint types are validated.
- Parameters:
src (str) – Node id of the source node (must already be in the batch).
dst (str) – Node id of the target node (must already be in the batch).
rel_type (str) – Relationship type name.
attrs (Optional[Dict[str, Any]]) – Optional property dict for the relationship.
allow_untyped (bool) – When
True, relationship types not declared in the schema are accepted. WhenFalse, an undeclared type raisesValueError.
- Raises:
KeyError – If src or dst has not been added to the batch yet.
ValueError – If the schema declares endpoint types for rel_type and the actual node labels do not match, or if allow_untyped is
Falseand rel_type is not in the schema.
- Return type:
None
- as_lists()[source]¶
Return the accumulated nodes and edges as plain lists.
- Returns:
A two-tuple
(nodes, edges)where each element is a list of dicts suitable for passing toPy2Neo.upsert_nodes_batch()andPy2Neo.upsert_edges_batch()respectively.- Return type:
Tuple[List[Dict[str, Any]], List[Dict[str, Any]]]
- src.dackar.knowledge_graph.kg_schema_builder_workflow._prefix(value, prefix)[source]¶
Build a namespaced node id by prepending prefix to value.
Strips whitespace, collapses internal spaces to underscores, and skips empty values. Normalising spaces ensures the resulting ID is safe to use unquoted in Cypher —
FM:loss_of_lubricationrather thanFM:loss of lubrication. Casing is preserved so that structured IDs (e.g.CMP-001) are not altered.Strings already namespaced with this prefix (
"<prefix>:...") are returned unchanged to avoid double-prefixing. A colon appearing elsewhere in the value (e.g. a time"12:30"or anOPCTX/PMcontext value) no longer suppresses namespacing, so distinct raw values can no longer collide onto the same un-prefixed id.- Parameters:
value (Optional[str]) – Raw identifier string (e.g.
"pump-101"or"loss of lubrication").prefix (str) – Namespace prefix (e.g.
"ASSET").
- Returns:
A string like
"ASSET:pump-101", orNoneif value is falsy or blank after stripping.- Return type:
Optional[str]
- src.dackar.knowledge_graph.kg_schema_builder_workflow._truthy(value)[source]¶
Return True if value is non-empty and not a sentinel null-like string.
Treats
"unknown","none", and"null"as falsy in addition to standard Python falsy values.- Parameters:
value (Any) – Any value to test.
- Returns:
Trueif the value is considered meaningfully present.- Return type:
bool
- src.dackar.knowledge_graph.kg_schema_builder_workflow._norm_text_key(value)[source]¶
Normalise a text string into a stable, lowercase key.
Collapses internal whitespace, strips leading/trailing whitespace, and lowercases the result. Used to derive deterministic node ids from free-text labels.
- Parameters:
value (Optional[str]) – Input string to normalise, or
None.- Returns:
Normalised string, or
Noneif value isNoneor blank.- Return type:
Optional[str]
- src.dackar.knowledge_graph.kg_schema_builder_workflow._safe_rel(g, src, dst, preferred, fallback, attrs=None)[source]¶
Add an edge using preferred relation type, falling back to fallback if not in schema.
- Parameters:
g (GraphBatch) – The
GraphBatchto add the edge to.src (str) – Source node id.
dst (str) – Target node id.
preferred (str) – Preferred relationship type name.
fallback (str) – Fallback relationship type name used when preferred is not declared in the schema’s relation map.
attrs (Optional[Dict[str, Any]]) – Optional property dict for the relationship.
- Return type:
None
- src.dackar.knowledge_graph.kg_schema_builder_workflow._safe_ptr_entity_rel(g, src, dst, attrs=None)[source]¶
ProcessedTextRecord -> entity relation chooser. Avoid using document-scoped relations like ‘mentions’ when the schema types them as condition_report/work_order -> element_usage.
- Parameters:
g (GraphBatch)
src (str)
dst (str)
attrs (Optional[Dict[str, Any]])
- Return type:
None
- src.dackar.knowledge_graph.kg_schema_builder_workflow._safe_ptr_failure_mode_rel(g, src, dst, attrs=None)[source]¶
Add a ProcessedTextRecord → failure_mode edge using the best available relation type.
Tries
"supports_hypothesis","references_failure_mode", and"caused_by"in order, picking the first whose endpoint types match the actual node labels. Falls back to an untyped"references_failure_mode"edge if none match.- Parameters:
g (GraphBatch) – The
GraphBatchto add the edge to.src (str) – Source node id (a ProcessedTextRecord node).
dst (str) – Target node id (a failure_mode node).
attrs (Optional[Dict[str, Any]]) – Optional property dict for the relationship.
- Return type:
None
- src.dackar.knowledge_graph.kg_schema_builder_workflow.build_graph_from_workflow_artifacts(schema_paths=None, *, event=None, kg_context=None, telemetry_summary=None, evidence_bundle=None, causality_candidates=None, rca_card=None, operational_context=None, pm_compliance=None, documents=None, processed_text_records=None)[source]¶
Translate all RCA workflow artifacts into a graph of nodes and edges.
Iterates over every supplied artifact, creates typed nodes for each entity (assets, components, failure modes, events, telemetry signals, anomalies, causal candidates, evidence snippets, etc.), and links them with labelled directed edges. All artifacts are optional; only those provided contribute nodes and edges.
- Parameters:
schema_paths (Optional[Union[str, pathlib.Path, Iterable[Union[str, pathlib.Path]]]]) – Optional path(s) to TOML schema file(s) used for label resolution and endpoint validation. An empty schema is used when
None.event (Optional[Dict[str, Any]]) – Parsed
eventartifact dict.kg_context (Optional[Dict[str, Any]]) – Parsed
kg_contextartifact dict (components, failure modes, past events).telemetry_summary (Optional[Dict[str, Any]]) – Parsed
telemetry_summaryartifact dict.evidence_bundle (Optional[Dict[str, Any]]) – Parsed
evidence_bundleartifact dict.causality_candidates (Optional[Dict[str, Any]]) – Parsed
causality_candidatesartifact dict.rca_card (Optional[Dict[str, Any]]) – Parsed
rca_cardartifact dict.operational_context (Optional[Dict[str, Any]]) – Parsed
operational_contextartifact dict.pm_compliance (Optional[Dict[str, Any]]) – Parsed
pm_complianceartifact dict.documents (Optional[Sequence[Dict[str, Any]]]) – List of document descriptor dicts.
processed_text_records (Optional[Sequence[Dict[str, Any]]]) – List of
processed_text_recorddicts.
- Returns:
A two-tuple
(nodes, edges)— lists of dicts ready for bulk ingestion viaingest_graph_toml().- Return type:
Tuple[List[Dict[str, Any]], List[Dict[str, Any]]]
- src.dackar.knowledge_graph.kg_schema_builder_workflow.ingest_graph_toml(client, nodes, edges, database=None)[source]¶
Bulk-upsert a node list and edge list into Neo4j.
A thin wrapper that calls
Py2Neo.upsert_nodes_batch()followed byPy2Neo.upsert_edges_batch(), skipping each call when the corresponding list is empty.- Parameters:
client (dackar.knowledge_graph.py2neo.Py2Neo) – Active
Py2Neoconnection.nodes (List[Dict[str, Any]]) – List of node dicts as returned by
build_graph_from_workflow_artifacts().edges (List[Dict[str, Any]]) – List of edge dicts as returned by
build_graph_from_workflow_artifacts().database (Optional[str]) – Target database name; uses the driver default when
None.
- Return type:
None