src.dackar.knowledge_graph.kg_schema_builder_workflow

Attributes

LOGGER

Classes

GraphBatch

In-memory accumulator for nodes and edges before bulk Neo4j ingestion.

Functions

load_toml_schema(path)

Load a TOML schema file and return its contents as a dict.

load_and_merge_schemas(schema_paths)

Load one or more TOML schema files and merge them into a single schema dict.

_schema_props(spec)

Extract the property list from a node or relation schema spec.

_node_primary_key(spec)

Determine the primary key property name for a node spec.

_label_candidates(name)

Generate a list of label name variants to try when resolving a schema label.

resolve_node_label(schema, *candidates)

Resolve the first candidate name that exists as a node label in the schema.

resolve_node_label_strict(schema, *candidates)

Resolve the first candidate name that exists as a node label in the schema.

relation_endpoint_map(schema)

Build a lookup from relation type name to its (from_entity, to_entity) pair.

_is_primitive(value)

Return True if value is a Neo4j-native scalar type (str, int, float, or bool).

sanitize_value(value)

Convert an arbitrary Python value to a Neo4j-compatible property value.

sanitize_props(props)

Sanitize all values in a property dict, dropping None entries.

generate_ddl_from_schema(schema)

Generate Neo4j DDL statements (constraints and indexes) from a merged schema.

apply_schema_constraints(client, schema_paths[, database])

Load schemas and apply all DDL constraints and indexes to the database.

_prefix(value, prefix)

Build a namespaced node id by prepending prefix to value.

_truthy(value)

Return True if value is non-empty and not a sentinel null-like string.

_norm_text_key(value)

Normalise a text string into a stable, lowercase key.

_safe_rel(g, src, dst, preferred, fallback[, attrs])

Add an edge using preferred relation type, falling back to fallback if not in schema.

_safe_ptr_entity_rel(g, src, dst[, attrs])

ProcessedTextRecord -> entity relation chooser.

_safe_ptr_failure_mode_rel(g, src, dst[, attrs])

Add a ProcessedTextRecord → failure_mode edge using the best available relation type.

build_graph_from_workflow_artifacts([schema_paths, ...])

Translate all RCA workflow artifacts into a graph of nodes and edges.

ingest_graph_toml(client, nodes, edges[, database])

Bulk-upsert a node list and edge list into Neo4j.

Module Contents

src.dackar.knowledge_graph.kg_schema_builder_workflow.LOGGER[source]
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 .toml schema 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 .toml files.

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 raises KeyError instead of silently falling back when no candidate matches. Use this wherever a miss should be a hard error (e.g. building the labels registry 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_entity and to_entity in 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:

True if 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). - list of primitives → kept as-is; mixed/complex lists → JSON string. - None and 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 None entries.

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 was None removed.

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 UNIQUE constraint on id is created. Additionally, a CREATE INDEX statement is emitted for each property that carries "indexed": true in 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(), and Py2Neo.query() into a single convenience call.

Parameters:
  • client (dackar.knowledge_graph.py2neo.Py2Neo) – Active Py2Neo connection.

  • 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 id string; 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]])

schema[source]
nodes: Dict[str, Dict[str, Any]][source]
edges: Dict[Tuple[str, str, str], Dict[str, Any]][source]
relation_map[source]
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 id property).

  • label (str) – Neo4j label for the node.

  • attrs (Optional[Dict[str, Any]]) – Optional property dict; id is 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. When False, an undeclared type raises ValueError.

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 False and 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 to Py2Neo.upsert_nodes_batch() and Py2Neo.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_lubrication rather than FM: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 an OPCTX/PM context 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", or None if 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:

True if 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 None if value is None or 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 GraphBatch to 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 GraphBatch to 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 event artifact dict.

  • kg_context (Optional[Dict[str, Any]]) – Parsed kg_context artifact dict (components, failure modes, past events).

  • telemetry_summary (Optional[Dict[str, Any]]) – Parsed telemetry_summary artifact dict.

  • evidence_bundle (Optional[Dict[str, Any]]) – Parsed evidence_bundle artifact dict.

  • causality_candidates (Optional[Dict[str, Any]]) – Parsed causality_candidates artifact dict.

  • rca_card (Optional[Dict[str, Any]]) – Parsed rca_card artifact dict.

  • operational_context (Optional[Dict[str, Any]]) – Parsed operational_context artifact dict.

  • pm_compliance (Optional[Dict[str, Any]]) – Parsed pm_compliance artifact dict.

  • documents (Optional[Sequence[Dict[str, Any]]]) – List of document descriptor dicts.

  • processed_text_records (Optional[Sequence[Dict[str, Any]]]) – List of processed_text_record dicts.

Returns:

A two-tuple (nodes, edges) — lists of dicts ready for bulk ingestion via ingest_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 by Py2Neo.upsert_edges_batch(), skipping each call when the corresponding list is empty.

Parameters:
Return type:

None