src.dackar.RCA.equipment_similarity.kg_equipment_poller

kg_equipment_poller — KGEquipmentPoller.

One-time population script: queries all element_usage → element_definition nodes from Neo4j (with linked failure modes), builds spec text via EquipmentSpecBuilder, and upserts into EquipmentSpecStore.

Run at KG setup time and re-run whenever the KG is updated.

Usage:

from py2neo import Graph
from storage.chroma_store import ChromaRecordStore
from equipment_similarity.equipment_spec_store import EquipmentSpecStore
from equipment_similarity.kg_equipment_poller import KGEquipmentPoller

neo4j = Graph("bolt://localhost:7687", auth=("neo4j", "password"))
chroma = ChromaRecordStore(persist_directory="./chroma_store")
spec_store = EquipmentSpecStore(chroma)
poller = KGEquipmentPoller(client=neo4j)
n = poller.poll_and_upsert(spec_store, batch_size=100)
print(f"Upserted {n} component specs")

Attributes

logger

JsonDict

_SPEC_QUERY

Classes

KGEquipmentPoller

Polls Neo4j for all element_usage / element_definition pairs and

Module Contents

src.dackar.RCA.equipment_similarity.kg_equipment_poller.logger[source]
src.dackar.RCA.equipment_similarity.kg_equipment_poller.JsonDict[source]
src.dackar.RCA.equipment_similarity.kg_equipment_poller._SPEC_QUERY = Multiline-String[source]
Show Value
"""
MATCH (c:element_usage)-[:instance_of]->(def:element_definition)
OPTIONAL MATCH (c)-[:subject_to|applies_to]-(fm:failure_mode)
RETURN
    c.id                    AS component_id,
    c.name                  AS component_name,
    def.domain_category     AS domain_category,
    def.structural_kind     AS structural_kind,
    def.nominal_size        AS nominal_size,
    def.design_pressure     AS design_pressure,
    def.design_temperature  AS design_temperature,
    def.material_spec       AS material_spec,
    def.manufacturer        AS manufacturer,
    def.model_number        AS model_number,
    collect(DISTINCT fm.name)             AS failure_mode_names,
    collect(DISTINCT fm.failure_mechanism) AS failure_mechanisms
ORDER BY component_id
"""
class src.dackar.RCA.equipment_similarity.kg_equipment_poller.KGEquipmentPoller(client, database=None)[source]

Polls Neo4j for all element_usage / element_definition pairs and upserts their spec embeddings into EquipmentSpecStore.

Parameters:
  • client (Any) – A py2neo.Graph (or any object with a run(query, **params) method returning an iterable of record-like objects).

  • database (Optional[str]) – Optional Neo4j database name. If None, the default database for the connection is used.

client[source]
database = None[source]
_builder[source]
poll_and_upsert(spec_store, batch_size=100)[source]

Query all components from the KG and upsert their specs.

Parameters:
Returns:

Total number of component specs upserted.

Return type:

int

Notes

This is upsert-only. Each component’s spec is written by a stable record_id (equip_spec::{component_id}), so re-running refreshes existing components and adds new ones, but specs for components that were removed from the KG are not deleted from the Chroma collection — they linger as stale vectors. A full refresh that drops removed components requires rebuilding the equipment_specs collection from scratch (delete the collection, then re-poll).

_fetch_rows()[source]

Execute the spec query and return a list of row dicts.

Supports py2neo-style Graph objects (graph.run(query)) and any object with a query(cypher, params, db=database) method (matches the KGContextBuilder client interface).

Return type:

List[JsonDict]

static _has_substantive_spec_data(definition_props, failure_mode_names, failure_mechanisms)[source]
Parameters:
  • definition_props (JsonDict)

  • failure_mode_names (List[str])

  • failure_mechanisms (List[str])

Return type:

bool