src.dackar.RCA.equipment_similarity.kg_equipment_poller ======================================================= .. py:module:: src.dackar.RCA.equipment_similarity.kg_equipment_poller .. autoapi-nested-parse:: 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 ---------- .. autoapisummary:: src.dackar.RCA.equipment_similarity.kg_equipment_poller.logger src.dackar.RCA.equipment_similarity.kg_equipment_poller.JsonDict src.dackar.RCA.equipment_similarity.kg_equipment_poller._SPEC_QUERY Classes ------- .. autoapisummary:: src.dackar.RCA.equipment_similarity.kg_equipment_poller.KGEquipmentPoller Module Contents --------------- .. py:data:: logger .. py:data:: JsonDict .. py:data:: _SPEC_QUERY :value: Multiline-String .. raw:: html
Show Value .. code-block:: python """ 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 """ .. raw:: html
.. py:class:: KGEquipmentPoller(client, database = None) Polls Neo4j for all element_usage / element_definition pairs and upserts their spec embeddings into EquipmentSpecStore. :param client: A ``py2neo.Graph`` (or any object with a ``run(query, **params)`` method returning an iterable of record-like objects). :param database: Optional Neo4j database name. If ``None``, the default database for the connection is used. .. py:attribute:: client .. py:attribute:: database :value: None .. py:attribute:: _builder .. py:method:: poll_and_upsert(spec_store, batch_size = 100) Query all components from the KG and upsert their specs. :param spec_store: ``EquipmentSpecStore`` to upsert into. :param batch_size: Number of records per Chroma upsert call. Larger batches are faster but use more memory. :returns: Total number of component specs upserted. :rtype: int .. rubric:: 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). .. py:method:: _fetch_rows() 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). .. py:method:: _has_substantive_spec_data(definition_props, failure_mode_names, failure_mechanisms) :staticmethod: