delphi-pump-analytics / analysis /knowledge_retriever.py
pandeydigant31's picture
Initial deploy: DELPHI Pump Analytics (renamed from MURPHY)
199bfa3 verified
Raw
History Blame Contribute Delete
17.6 kB
"""
Domain Knowledge Retriever for CSH2 LLM Analysis
Pattern-based retrieval: maps sensor data patterns to relevant failure modes,
reasoning chains, and operational context from the YAML knowledge base.
At this scale (8 failure modes, 6 chains), rule-based retrieval is more precise
than embedding-based RAG — no vector DB needed.
"""
import logging
from pathlib import Path
from typing import Any, Dict, List, Optional, Set
try:
import yaml
except ImportError:
yaml = None
logger = logging.getLogger(__name__)
KNOWLEDGE_DIR = Path(__file__).parent / "knowledge"
# Trigger thresholds (physically meaningful boundaries)
LOW_PRESSURE_THRESHOLD_BAR = 350 # H35 target — below this is underperforming
HIGH_DISCHARGE_TEMP_K = 100 # Unusual for LH2 (inlet ~25K)
LOW_FLOW_THRESHOLD_KG_MIN = 0.5 # < half of demonstrated 1.13 kg/min
MOTOR_SPEED_ACTIVE_THRESHOLD = 83 # Scaled RPM (raw ~50 × 5/3 display multiplier) — pump is running
class DomainKnowledgeRetriever:
"""Retrieves structured domain knowledge from YAML files.
Loads failure_modes.yaml, reasoning_chains.yaml, operational_context.yaml
at startup. Uses pattern matching on cycle_stats to determine which
knowledge chunks to inject into LLM prompts.
"""
def __init__(self, knowledge_dir: Optional[Path] = None):
self._dir = Path(knowledge_dir) if knowledge_dir else KNOWLEDGE_DIR
if yaml is None:
raise ImportError("PyYAML required. Install with: pip install pyyaml")
self._failure_modes = self._load_yaml("failure_modes.yaml")
self._reasoning_chains = self._load_yaml("reasoning_chains.yaml")
self._operational_context = self._load_yaml("operational_context.yaml")
def _load_yaml(self, filename: str) -> Dict:
path = self._dir / filename
if not path.exists():
logger.warning(f"Knowledge file not found: {path}")
return {}
with open(path, "r") as f:
return yaml.safe_load(f) or {}
# ------------------------------------------------------------------
# Public API
# ------------------------------------------------------------------
def get_context_for_cycle(
self, cycle_stats: Dict, plateaus: Optional[Dict] = None
) -> Dict[str, Any]:
"""Determine which knowledge chunks are relevant based on cycle data.
Args:
cycle_stats: Dict from DataAnalyzer.analyze_cycle()
plateaus: Optional dict mapping tag names to plateau periods
Returns:
Dict with categorized knowledge chunks ready for formatting
"""
context: Dict[str, Any] = {
"failure_modes": [],
"reasoning_chains": [],
"sensor_notes": [],
"regime_caveats": [],
"system_personality": [],
"benchmarks": {},
}
# Always-include sections
context["system_personality"] = self._get_system_personality()
context["sensor_notes"] = self._get_sensor_notes_for_stats(cycle_stats)
context["benchmarks"] = self._get_phase1_benchmarks()
context["regime_caveats"] = self._get_regime_caveats()
# Pattern-based trigger detection
triggers = self._detect_triggers(cycle_stats, plateaus)
if "low_pressure_build" in triggers:
context["failure_modes"].extend(
self._get_fm_by_ids(["FM-001", "FM-002", "FM-007", "FM-008"])
)
context["reasoning_chains"].extend(
self._get_rc_by_ids(["RC-001"])
)
if "high_discharge_temp" in triggers:
context["failure_modes"].extend(
self._get_fm_by_ids(["FM-005", "FM-006"])
)
context["reasoning_chains"].extend(
self._get_rc_by_ids(["RC-003"])
)
if "flow_anomaly" in triggers:
context["failure_modes"].extend(
self._get_fm_by_ids(["FM-001", "FM-002", "FM-007"])
)
context["reasoning_chains"].extend(
self._get_rc_by_ids(["RC-004"])
)
if "pressure_oscillation" in triggers:
context["failure_modes"].extend(
self._get_fm_by_ids(["FM-004"])
)
context["reasoning_chains"].extend(
self._get_rc_by_ids(["RC-002"])
)
if "motor_anomaly" in triggers:
context["failure_modes"].extend(
self._get_fm_by_ids(["FM-003"])
)
context["reasoning_chains"].extend(
self._get_rc_by_ids(["RC-005"])
)
# Deduplicate (triggers may overlap)
context["failure_modes"] = _deduplicate(context["failure_modes"], "id")
context["reasoning_chains"] = _deduplicate(context["reasoning_chains"], "id")
return context
def get_context_for_comparison(
self, stats_a: Dict, stats_b: Dict
) -> Dict[str, Any]:
"""Get knowledge context for a two-cycle comparison.
Runs trigger detection on both cycles and takes the union.
"""
ctx_a = self.get_context_for_cycle(stats_a)
ctx_b = self.get_context_for_cycle(stats_b)
merged: Dict[str, Any] = {
"failure_modes": _deduplicate(
ctx_a["failure_modes"] + ctx_b["failure_modes"], "id"
),
"reasoning_chains": _deduplicate(
ctx_a["reasoning_chains"] + ctx_b["reasoning_chains"], "id"
),
"sensor_notes": _deduplicate(
ctx_a["sensor_notes"] + ctx_b["sensor_notes"], "sensor"
),
"regime_caveats": ctx_a["regime_caveats"], # Same for both
"system_personality": ctx_a["system_personality"],
"benchmarks": ctx_a["benchmarks"],
}
return merged
def format_for_prompt(self, context: Dict[str, Any]) -> str:
"""Format retrieved knowledge into text blocks for LLM prompt injection."""
sections = []
# System personality
if context.get("system_personality"):
lines = ["SYSTEM PERSONALITY (known recurring issues & strengths):"]
for item in context["system_personality"]:
prefix = f" [Priority {item['priority']}]" if "priority" in item else " -"
lines.append(f"{prefix} {item['description']}")
if item.get("details"):
lines.append(f" {item['details']}")
sections.append("\n".join(lines))
# Sensor trustworthiness
if context.get("sensor_notes"):
lines = ["SENSOR TRUSTWORTHINESS NOTES:"]
for note in context["sensor_notes"]:
lines.append(f" {note['sensor']}:")
lines.append(f" Issue: {note['notes']}")
lines.append(f" Workaround: {note['workaround']}")
sections.append("\n".join(lines))
# Failure modes (plain-language, no internal IDs)
if context.get("failure_modes"):
lines = ["RELEVANT FAILURE PATTERNS (from CSH2 knowledge base):"]
for fm in context["failure_modes"]:
priority = fm.get('priority', 'N/A')
lines.append(f" **{fm['name']}** ({priority} priority)")
lines.append(f" Primary symptom: {fm['primary_symptom']}")
mechanism = fm.get("mechanism", "")
if len(mechanism) > 250:
mechanism = mechanism[:250] + "..."
lines.append(f" Mechanism: {mechanism}")
if fm.get("candidate_causes"):
causes = ", ".join(
f"{c['cause']} ({c['prior']:.0%})"
for c in fm["candidate_causes"][:3]
)
lines.append(f" Top causes: {causes}")
if fm.get("recommended_action"):
lines.append(f" Action: {fm['recommended_action']}")
sections.append("\n".join(lines))
# Reasoning chains (plain-language, no internal IDs)
if context.get("reasoning_chains"):
lines = ["DIAGNOSTIC REASONING STEPS:"]
for rc in context["reasoning_chains"]:
lines.append(f" **{rc['name']}**")
lines.append(f" Trigger: {rc['trigger']}")
if rc.get("steps"):
steps_text = "; ".join(
f"Step {s['step']}: {s['description']}"
for s in rc["steps"][:4]
)
lines.append(f" Steps: {steps_text}")
sections.append("\n".join(lines))
# Regime caveats
if context.get("regime_caveats"):
lines = ["OPERATIONAL CAVEATS:"]
for caveat in context["regime_caveats"]:
lines.append(f" {caveat['name']}: {caveat['description']}")
sections.append("\n".join(lines))
# Benchmarks
benchmarks = context.get("benchmarks", {})
if benchmarks:
lines = [
"PHASE 1B BENCHMARKS (reference values — Sept 24-26 ccH2 fills):",
f" Specific energy: {benchmarks.get('specific_energy', '0.26 kWh/kg')}",
f" Flow rate: {benchmarks.get('flow_rate', '1.13-1.25 kg/min at 65% motor')}",
f" Peak pressure: {benchmarks.get('peak_pressure', '370 bar')}",
f" Delivered temp: {benchmarks.get('delivered_temp', '55K')}",
f" Uptime: {benchmarks.get('uptime', '100% (6 fills, 4 back-to-back)')}",
]
sections.append("\n".join(lines))
return "\n\n".join(sections)
# ------------------------------------------------------------------
# Trigger detection
# ------------------------------------------------------------------
def _detect_triggers(
self, stats: Dict, plateaus: Optional[Dict] = None
) -> Set[str]:
"""Detect which diagnostic triggers are active based on sensor patterns."""
triggers: Set[str] = set()
# --- Low pressure build ---
dp = stats.get("discharge_pressure", {})
peak_pressure = dp.get("peak", 0)
if 0 < peak_pressure < LOW_PRESSURE_THRESHOLD_BAR:
triggers.add("low_pressure_build")
# Also check if PT130 plateaus below target
if plateaus:
pt130_plateaus = plateaus.get("PT130", [])
if pt130_plateaus:
max_plateau = max(
(p.get("value", 0) for p in pt130_plateaus), default=0
)
if 0 < max_plateau < LOW_PRESSURE_THRESHOLD_BAR:
triggers.add("low_pressure_build")
# --- High discharge temperature ---
tt130 = stats.get("temperature_TT130", {})
tt130_peak = tt130.get("peak", 0)
tt130_avg = tt130.get("avg", 0)
if tt130_peak > HIGH_DISCHARGE_TEMP_K:
triggers.add("high_discharge_temp")
if tt130_avg > 0 and tt130_peak > tt130_avg * 1.3 and tt130_peak > 60:
triggers.add("high_discharge_temp")
# --- Flow anomaly ---
flow = stats.get("flow", {})
avg_flow = flow.get("avg_flow", -1)
motor = stats.get("motor", {})
avg_speed = motor.get("avg_speed", 0)
if avg_flow >= 0 and avg_flow < LOW_FLOW_THRESHOLD_KG_MIN and avg_speed > MOTOR_SPEED_ACTIVE_THRESHOLD:
triggers.add("flow_anomaly")
# --- Pressure oscillation ---
# Distinguish oscillation from a normal smooth ramp.
# Key insight: a ramp from ambient to 370 bar has peak >> initial.
# True oscillation happens when the system is already at high
# pressure and fluctuates. We trigger ONLY when:
# 1. Peak pressure > 200 bar (system was pressurized)
# 2. Initial pressure > 30% of peak (NOT a ramp from ambient)
# 3. Std is non-trivial (> 20 bar — real fluctuation)
dp_initial = dp.get("initial", 0)
dp_std = dp.get("std", 0)
if peak_pressure > 200 and dp_initial > peak_pressure * 0.3 and dp_std > 20:
triggers.add("pressure_oscillation")
return triggers
# ------------------------------------------------------------------
# Knowledge accessors
# ------------------------------------------------------------------
def _get_fm_by_ids(self, fm_ids: List[str]) -> List[Dict]:
"""Get failure modes by ID, formatted for context injection."""
modes = self._failure_modes.get("failure_modes", [])
results = []
for fm in modes:
if fm.get("id") in fm_ids:
results.append({
"id": fm.get("id"),
"name": fm.get("name"),
"priority": fm.get("system_priority"),
"primary_symptom": fm.get("symptoms", {}).get("primary", ""),
"mechanism": fm.get("physical_mechanism", ""),
"candidate_causes": [
{"cause": c["cause"], "prior": c["prior"]}
for c in fm.get("candidate_causes", [])
],
"recommended_action": fm.get("recommended_actions", {}).get("immediate", ""),
})
return results
def _get_rc_by_ids(self, rc_ids: List[str]) -> List[Dict]:
"""Get reasoning chains by ID, formatted for context injection."""
chains = self._reasoning_chains.get("reasoning_chains", [])
results = []
for chain in chains:
if chain.get("id") in rc_ids:
results.append({
"id": chain.get("id"),
"name": chain.get("name"),
"trigger": chain.get("trigger"),
"related_fms": chain.get("related_failure_modes", []),
"steps": [
{
"step": s.get("step"),
"action": s.get("action"),
"description": s.get("description"),
}
for s in chain.get("steps", [])
],
})
return results
def _get_system_personality(self) -> List[Dict]:
"""Get system personality: recurring issues + strengths."""
personality = self._operational_context.get("system_personality", {})
items = []
for issue in personality.get("recurring_issues", []):
items.append({
"priority": issue.get("priority"),
"description": issue.get("description"),
"details": issue.get("details"),
"type": "recurring_issue",
})
for strength in personality.get("strengths", []):
items.append({
"description": strength,
"details": None,
"type": "strength",
})
return items
def _get_sensor_notes_for_stats(self, stats: Dict) -> List[Dict]:
"""Get sensor trustworthiness notes for sensors present in the data."""
sensors = self._operational_context.get("sensor_trustworthiness", {})
present_sensors = set()
# Detect which sensors are in the stats
if "discharge_pressure" in stats:
present_sensors.add("PT130")
if "supply_pressure" in stats:
present_sensors.add("PT110")
if "temperature_TT110" in stats:
present_sensors.add("TT110")
if "temperature_TT130" in stats:
present_sensors.add("TT130")
if "flow" in stats:
present_sensors.add("FT140")
if "motor" in stats:
present_sensors.add("M130_Speed_vs_MC130")
results = []
for sensor_key, info in sensors.items():
# Match sensor key to present sensors
if sensor_key in present_sensors or any(
s in sensor_key for s in present_sensors
):
results.append({
"sensor": sensor_key,
"notes": info.get("notes", "").strip(),
"workaround": info.get("workaround", "").strip(),
})
return results
def _get_regime_caveats(self) -> List[Dict]:
"""Get all regime caveats."""
caveats = self._operational_context.get("regime_caveats", {})
return [
{
"name": name,
"description": info.get("description", ""),
}
for name, info in caveats.items()
]
def _get_phase1_benchmarks(self) -> Dict[str, str]:
"""Get Phase 1B benchmark values."""
return {
"specific_energy": "0.26 kWh/kg",
"flow_rate": "1.13-1.25 kg/min at 65% motor",
"peak_pressure": "370 bar (Phase 1B ccH2 fills)",
"delivered_temp": "55K",
"uptime": "100% (6 fills, 4 back-to-back)",
}
# ------------------------------------------------------------------
# Utility
# ------------------------------------------------------------------
def _deduplicate(items: List[Dict], key: str) -> List[Dict]:
"""Deduplicate list of dicts by a key field."""
seen: set = set()
result = []
for item in items:
val = item.get(key)
if val not in seen:
seen.add(val)
result.append(item)
return result