src.dackar.RCA.orchestrators.causality_engine_v32¶
causality_engine_v32 — Rule-based causality engine, TSKR-aware production variant.
Role in the pipeline¶
This engine extends v31 with two additional scoring dimensions:
TSKR temporal patterns — Allen interval algebra classifies each anomaly window relative to the failure event; latency alignment and recurrence profiles are folded into the temporal score.
NER entity normalisation —
EntityNormalizerreconciles free-text component mentions in FMEA/KG records against the structured plant vocabulary, improving candidate matching precision.
Relationship to v31¶
causality_engine_v31 is the baseline engine and is intentionally retained
alongside this module. Running both engines on the same inputs provides an
independent validation baseline: v31 results represent the purely
structural/evidence view, while v32 adds temporal reasoning. Comparing the
two candidate rankings helps verify that TSKR enrichment improves rather than
regresses root-cause identification, and surfaces edge cases such as
delayed-onset failure modes where temporal penalties may be inappropriate.
Intended usage: pass RuleBasedCausalityEngineV32 as the causality_engine
argument of RCAReasoningOrchestrator for production runs, and
RuleBasedCausalityEngineV31 for baseline validation passes.
Attributes¶
Classes¶
TSKR-aware deterministic causality engine with explicit screening metadata. |
Module Contents¶
- src.dackar.RCA.orchestrators.causality_engine_v32._CRITICAL_BARRIER_KEYWORDS = ('reactor protection', 'reactor trip', 'trip logic', 'reactor shutdown', 'containment...[source]¶
- src.dackar.RCA.orchestrators.causality_engine_v32._HIGH_BARRIER_KEYWORDS = ('core cooling', 'emergency core cooling', 'residual heat removal', 'decay heat removal',...[source]¶
- src.dackar.RCA.orchestrators.causality_engine_v32._CRITICAL_RISK_KEYWORDS = ('reactor protection', 'reactor trip', 'trip logic', 'reactor shutdown', 'containment...[source]¶
- src.dackar.RCA.orchestrators.causality_engine_v32._HIGH_RISK_KEYWORDS = ('core cooling', 'emergency core cooling', 'residual heat removal', 'decay heat removal',...[source]¶
- src.dackar.RCA.orchestrators.causality_engine_v32._DEFAULT_SCORING_PROFILES: Dict[str, Dict[str, float]][source]¶
- class src.dackar.RCA.orchestrators.causality_engine_v32.RuleBasedCausalityEngineV32(config=None)[source]¶
TSKR-aware deterministic causality engine with explicit screening metadata.
- Parameters:
config (Optional[CausalityEngineConfigV32])
- _CAUSAL_CATEGORIES: List[str] = ['A', 'B', 'C', 'D', 'E', 'F', 'G', 'H', 'I', 'J', 'K', 'L'][source]¶
- _RULEOUT_REASON_CODES: List[str] = ['physically_impossible', 'timeline_inconsistent', 'barrier_held', 'no_supporting_data',...[source]¶
- generate(event, telemetry_summary, kg_context, tskr_patterns, operational_context, pm_compliance, run_context)[source]¶
Generate ranked causal candidate hypotheses for the event.
Combines failure-mode candidates (from the KG neighbourhood, scored on temporal/logical/documentary streams) with historical event analogs, assigns cause categories, and ranks them by composite score.
- Parameters:
event (JsonDict) – Target abnormal event.
telemetry_summary (JsonDict) – Telemetry anomaly summary for the event window.
kg_context (JsonDict) – KG neighbourhood (components, failure modes, past events).
tskr_patterns (Optional[JsonDict]) – TSKR chain-position patterns keyed by target, or None.
operational_context (Optional[JsonDict]) – Optional supporting artifacts, or None.
pm_compliance (Optional[JsonDict]) – Optional supporting artifacts, or None.
run_context (JsonDict) – Orchestrator run context.
- Returns:
Candidate hypotheses conforming to
schemas/causality_candidates.json(each with scores, a cause category, and temporal evidence).- Return type:
- _build_failure_mode_candidates(event, event_time, telemetry_summary, kg_context, tskr_index, pm_compliance, past_event_index, common_cause_index, operational_context=None, sf_index=None)[source]¶
- _build_past_event_candidates(event, event_time, telemetry_summary, kg_context, tskr_index, pm_compliance, past_event_index, common_cause_index, operational_context=None, sf_index=None)[source]¶
- _evidence_posture(support_score, contradiction_score, contextual_score, retrieved_hit_count=0)[source]¶
Classify the evidence posture for a candidate.
Distinguishes “evidence against” (contradicted) from “no data retrieved” (no_data) — these have different implications for corrective action scope. A “weak” posture means documents were retrieved but none were strongly for or against the hypothesis. “no_data” means the retrieval layer returned nothing for this candidate, so the hypothesis is neither supported nor contradicted by the document corpus.
- Parameters:
support_score (float)
contradiction_score (float)
contextual_score (float)
retrieved_hit_count (int)
- Return type:
str
- _temporal_posture(temporal_score, temporal_precedence, latency_consistency, temporal_contradiction)[source]¶
- Parameters:
temporal_score (float)
temporal_precedence (float)
latency_consistency (float)
temporal_contradiction (bool)
- Return type:
str
- refine_with_evidence(causality_candidates, evidence_bundle, kg_context=None, signal_evidence=None, entity_normalizer_cfg=None, coverage_summary=None, allen_relation_map=None, protection_logic_context=None)[source]¶
Re-score candidates with retrieved evidence and auxiliary signals.
Folds each candidate’s supporting/contradicting evidence, signal-episode chain scores, Allen temporal relations, and protection-logic barrier state into an updated composite score and evidence posture, returning a new candidates payload (the input is not mutated).
- Parameters:
causality_candidates (JsonDict) – Candidate hypotheses from
generate().evidence_bundle (JsonDict) – Retrieved evidence whose per-candidate summary drives re-scoring.
kg_context (Optional[JsonDict]) – Optional KG neighbourhood supplying failure modes for entity normalization, or None.
signal_evidence (Optional[JsonDict]) – Optional per-candidate signal-episode chain scores, or None.
entity_normalizer_cfg (Optional[Dict[str, Any]]) – Optional entity-normalizer configuration overrides, or None.
coverage_summary (Optional[JsonDict]) – Optional evidence-coverage summary shaping the coverage factor.
allen_relation_map (Optional[JsonDict]) – Optional Allen temporal-relation map between components, or None.
protection_logic_context (Optional[JsonDict]) – Optional protection-logic (barrier) context, or None.
- Returns:
A refined candidates payload conforming to
schemas/causality_candidates.json.- Return type:
- classmethod _canonical_candidate_key(*, component_id, mechanism_id, category, chain_position, event_scope_id)[source]¶
- Parameters:
component_id (Optional[str])
mechanism_id (Optional[str])
category (Optional[str])
chain_position (Optional[str])
event_scope_id (Optional[str])
- Return type:
str
- static _canonical_tuple(*, component_id, mechanism_id, category, chain_position)[source]¶
- Parameters:
component_id (Optional[str])
mechanism_id (Optional[str])
category (Optional[str])
chain_position (Optional[str])
- Return type:
- static _chain_position_from_signal_dag(position_type)[source]¶
Map a signal-DAG
position_typeonto the candidatechain_positionvocabulary.The telemetry-propagation DAG classifies a candidate’s anomaly as a root / common-cause root (upstream initiator), an intermediate node, or a convergence confluence (a downstream node where multiple chains meet — a symptom). This maps that view onto the coarser initiating/contributing/consequence vocabulary used for analyst-facing chain-position reasoning.
- Parameters:
position_type (Optional[str])
- Return type:
Optional[str]
- static _chain_position_from_relation(relation)[source]¶
- Parameters:
relation (Optional[str])
- Return type:
str
- _chain_position_for_candidate(*, relation, temporal_precedence, temporal_contradiction)[source]¶
- Parameters:
relation (Optional[str])
temporal_precedence (float)
temporal_contradiction (bool)
- Return type:
Tuple[str, str]
- classmethod _infer_category_from_text(text, default='A')[source]¶
- Parameters:
text (str)
default (str)
- Return type:
Tuple[str, List[str]]
- classmethod _infer_primary_category_for_past_event(*, pe)[source]¶
- Parameters:
pe (JsonDict)
- Return type:
Tuple[str, List[str]]
- classmethod _assess_category_applicability(*, kg_context, operational_context, external_oe_unavailable)[source]¶
- classmethod _build_metamodel_scaffolds(*, retained_candidates, filtered_out_candidates, event_analogs, kg_context, operational_context, external_oe_unavailable)[source]¶
- static _has_external_oe_signal(summary_lookup)[source]¶
- Parameters:
summary_lookup (Dict[str, JsonDict])
- Return type:
bool
- _apply_uncertainty_propagation(candidate)[source]¶
- Parameters:
candidate (JsonDict)
- Return type:
None
- static _coverage_quality_profile(coverage_summary)[source]¶
- Parameters:
coverage_summary (Optional[JsonDict])
- Return type:
Tuple[float, List[str]]
- static _build_allen_component_index(allen_relation_map)[source]¶
Index Allen relation map nodes by component_id for fast per-candidate lookup.
- Returns:
causal_scores {component_id → best allen_base_score among causal nodes} causal_relation {component_id → allen_relation_to_event of the best node} follow_ids set of component_ids that have at least one ‘follows’ node
- Parameters:
allen_relation_map (Optional[JsonDict])
- Return type:
Tuple[Dict[str, float], Dict[str, str], set[str]]
- static _apply_allen_temporal_blend(candidate, causal_scores, causal_relation, follow_ids, weights)[source]¶
Blend Allen base score into candidate temporal score in-place.
- Blend formula (α = 0.25):
new_temporal = 0.75 × old_temporal + 0.25 × allen_score (when match found)
Allen can both raise and lower the temporal score depending on whether the Allen base score is above or below the TSKR-derived baseline. This allows candidates with weak Allen relations (e.g. OVERLAPS with a low allen_base_score) to score lower than candidates with strong relations (e.g. PRECEDES with a high allen_base_score), as intended. When the component has a ‘follows’ node, temporal_contradiction is set True. composite_raw and composite_score are updated by the temporal weight delta.
- Parameters:
candidate (JsonDict)
causal_scores (Dict[str, float])
causal_relation (Dict[str, str])
follow_ids (set[str])
weights (Dict[str, float])
- Return type:
None
- static _apply_score_confidence_interval(candidate)[source]¶
Compute a per-candidate score confidence interval from data-degradation signals.
Five scoring dimensions are assessed; each contributes 1/5 to the interval width when its primary data source is absent or proxy-derived:
structural — physical_plausibility gate ran in degraded mode temporal — temporal_score_quality is “proxy” (no Allen causal match) telemetry — telemetry sub-score is zero (no telemetry signal available) evidence — candidate is observationally_ungrounded (no affects-class evidence) governance — barrier_logic gate ran in degraded mode
width = n_degraded / 5 (0.0 → narrow, 1.0 → very wide) lower = max(0.0, composite_score − width/2) upper = min(1.0, composite_score + width/2)
Writes candidate[“score_confidence_interval”].
- Parameters:
candidate (JsonDict)
- Return type:
None
- static _build_plc_barrier_index(protection_logic_context)[source]¶
Parse protection_logic_context into lookup structures.
- Returns:
sf_state_index {sf_id → barrier_state} from barrier_states[] logic_signal_ids set of signal/component IDs from logic_set
input_signals and output_signals (all logic_sets)
- Parameters:
protection_logic_context (Optional[JsonDict])
- Return type:
Tuple[Dict[str, str], set[str]]
- classmethod _operating_point_score(*, operational_context, primary_causal_category, fm_superclass, fm_name)[source]¶
Return (score 0–1, rationale_note) for the operating-point dimension.
Returns (0.0, “not_assessed”) when operational_context is None or mode is absent — never penalises candidates for missing data.
Only Category E candidates receive the power-level modifier. Train OOS bonus applies to standby-mechanism keywords for all categories.
- Parameters:
operational_context (Optional[JsonDict])
primary_causal_category (str)
fm_superclass (Optional[str])
fm_name (Optional[str])
- Return type:
Tuple[float, str]
- static _build_sensitivity_table(*, candidates, coverage_summary, top_n=5)[source]¶
Step 5 — sensitivity table: estimate composite-score delta per candidate if each currently missing/not_assessed data source were available at full quality.
The estimate re-computes the coverage_factor with the target family set to ‘complete’, then scales the composite_raw by the ratio of the new factor to the current one (capped at 1.0). This is an upper-bound estimate, not a precise prediction.
- _apply_coverage_quality_adjustment(candidate, *, coverage_factor, coverage_flags)[source]¶
- Parameters:
candidate (JsonDict)
coverage_factor (float)
coverage_flags (List[str])
- Return type:
None
- _apply_category_minimum_evidence_gate(candidate)[source]¶
- Parameters:
candidate (JsonDict)
- Return type:
None
- _apply_physical_plausibility_gate(candidate, plc_logic_signal_ids=None, plc_sf_state=None)[source]¶
- Parameters:
candidate (JsonDict)
plc_logic_signal_ids (Optional[set[str]])
plc_sf_state (Optional[Dict[str, str]])
- Return type:
None
- _apply_timeline_consistency_gate(candidate)[source]¶
- Parameters:
candidate (JsonDict)
- Return type:
None
- _apply_barrier_logic_gate(candidate, plc_sf_state=None)[source]¶
- Parameters:
candidate (JsonDict)
plc_sf_state (Optional[Dict[str, str]])
- Return type:
None
- static _governance_weight_for_fm(superclass)[source]¶
- Parameters:
superclass (Optional[str])
- Return type:
float
- _scoring_profile_for_fm(category)[source]¶
Return the full weight profile for a causal category (Step 2 / Phase 4c).
Looks up self.config.scoring_profiles by category letter; falls back to the ‘A’ (equipment_origin) profile when the category is unrecognised. Returns a copy so callers cannot mutate the config.
- Parameters:
category (str)
- Return type:
Dict[str, float]
- _alarm_signal_for_candidate(component_id, operational_context, components)[source]¶
Derive an alarm-based structural corroboration signal for a candidate.
Iterates
operational_context.recent_alarmsand checks whether each alarm’ssystem_affectedmatches the candidate’s component or any component in the same KG subgraph neighborhood.Match tiers (highest wins, not additive to avoid gaming): - Direct:
system_affectedequals the candidatecomponent_idexactly, or is a prefix/substring of it (plant tag convention).
Neighborhood:
system_affectedmatches any other component in the subgraph (e.g. an upstream component that feeds the failing one).
Alarm weight = priority weight × acknowledgement factor: - Unacknowledged alarms (
acknowledged_atis null): full weight - Acknowledged alarms: 0.5× (condition was noted but may be ongoing)Returns a float in [0.0, 1.0] — the maximum weighted alarm signal across all alarms. Zero when no alarms match or when
operational_contexthas norecent_alarms.
- _symptom_match_score(event, fm, telemetry_summary)[source]¶
Score [0, 1] for how well the event’s observed symptoms match what this failure mode is expected to produce.
0.5 = neutral (no symptom data available in either direction) >0.5 = symptoms consistent with this failure mode <0.5 = symptoms inconsistent with this failure mode
- Two sub-signals combined by available weight:
Anomaly pattern match (weight 0.6): dominant observed pattern vs. fm.expected_anomaly_pattern. Observed pattern is taken from the most frequently occurring anomaly pattern in telemetry (more objective), falling back to event.symptom_signature.anomaly_pattern.
Symptom type overlap (weight 0.4): F1-score between event’s symptom_types and fm.expected_symptom_types.
When a sub-signal has no data, its weight is excluded and the remaining signal is used alone. When neither sub-signal has data, returns 0.5.
- classmethod _pattern_similarity_score(expected_pattern, observed_pattern)[source]¶
- Parameters:
expected_pattern (str)
observed_pattern (str)
- Return type:
float
- static _dominant_telemetry_pattern(telemetry_summary)[source]¶
Return the most frequently occurring anomaly pattern across all signals, or None if no anomalies are present.
- static _recency_factor(time_distance_days)[source]¶
Map time_distance_days to a [0.55, 1.0] recency multiplier.
None (unknown age) receives a conservative 0.75 — neither penalised nor given full credit. Values are intentionally coarse so that small differences in document age do not create artificial score cliffs.
- Return type:
float
- _governance_details(pm_compliance, fm_name=None, fm_superclass=None, component_name=None, component_id=None, fm_id=None)[source]¶
Candidate-specific governance score from PM compliance data, with full trace.
Matching priority (per failed check): 1. Structural —
check.component_id == component_id: the check isscoped to this candidate’s component in the CMMS; no keyword heuristics needed.
FM-level —
fm_id in check.applicable_fm_ids: the PM task explicitly targets this failure mode (e.g., a surveillance test for a specific trip function); this narrows a component-level check to a single FM.Keyword fallback —
check_typekeywords matched againstfm_name + component_nametext; used only when neither structural field is available (legacy or synthetic data withoutcomponent_id).
Returns a dict containing: -
score: float in [0.5, 0.95] -pm_data_available: bool -total_checks: int — total PM checks in the compliance record -failed_check_count: int — asset-level failed checks -relevant_failed_checks: list of dicts, one per matched failed check,each containing
check_type,check_id,wo_id(if present),overdue_by_days,matched_keywords, andmatch_method(“component_id”, “applicable_fm_ids”, or “keyword”)count_boost: float — score contribution from number of relevant failuresoverdue_boost: float — score contribution from overdue dayscandidate_text: str — the lowercased text used for keyword matching
Score semantics: - 0.5 (neutral): no PM data, all checks passed, or no checks relevant to
this candidate. Never below 0.5 — PM alone cannot exonerate a candidate.
> 0.5: at least one failed check is relevant; scaled by count + overdue.
Maximum 0.95 — PM alone is never conclusive.
- Return type:
Dict[str, Any]
- _governance_score(pm_compliance, fm_name=None, fm_superclass=None, component_name=None, component_id=None, fm_id=None)[source]¶
Thin wrapper — returns only the score float from
_governance_details().- Return type:
float
- static _pm_check_matched_keywords(check, candidate_text)[source]¶
Return the set of keywords from this check type that appear in candidate_text.
Splits on whitespace AND hyphens/underscores so that hyphenated compounds like “in-leakage” are tokenised as [“in”, “leakage”] and the keyword “leakage” correctly matches.
- Parameters:
check (JsonDict)
candidate_text (str)
- Return type:
set
- static _pm_check_relevant(check, candidate_text)[source]¶
Return True if a PM check type matches keywords in the candidate’s text.
- static _governance_rationale(gov)[source]¶
Render the governance details dict as a traceable rationale string.
Examples
No PM data: “No PM compliance data available; score=0.5 (neutral).”
All passed: “All 4 PM checks passed on asset; score=0.5 (neutral).”
- No match: “3 asset-level PM failures; none relevant to candidate
(candidate_text=’bearing wear pump-1a’); score=0.5 (neutral).”
- Match: “2 relevant failed PM checks: lubrication (keywords: bearing,
wear; WO=WO-123; overdue=45d), inspection (keywords: corrosion; overdue=0d); count_boost=0.15, overdue_boost=0.05; score=0.75.”
- Parameters:
gov (Dict[str, Any])
- Return type:
str
- _build_safety_function_index(kg_context)[source]¶
Build a {component_id: [sf_dict, …]} lookup from kg_context.safety_functions.
- _affected_safety_functions_for_candidate(component_id, sf_index, impact_type='direct')[source]¶
Return the list of safety function dicts linked to component_id via sf_index.
Deduplicates by sf_id. Returns an empty list when the component has no associated safety functions or sf_index is empty (e.g. the KG has no safety_function nodes, or the feature was disabled in KGContextBuilderConfig).
- _barrier_signal_from_safety_functions(affected_safety_functions)[source]¶
- Parameters:
affected_safety_functions (List[JsonDict])
- Return type:
float
- static _apply_risk_significance_to_governance(*, governance_score, risk_significance_scalar)[source]¶
- Parameters:
governance_score (float)
risk_significance_scalar (float)
- Return type:
tuple[float, float]
- _recurrence_score_from_features(same_failure_mode_event_count, same_component_event_count, same_asset_event_count, unresolved_fm_count=0, unresolved_component_count=0, weighted_unresolved_fm_boost=None)[source]¶
- Parameters:
unresolved_fm_count (int)
unresolved_component_count (int)
weighted_unresolved_fm_boost (Optional[float])
- _recurrence_features_for_candidate(candidate, event, past_event_index, hypothesis_component_id=None, hypothesis_failure_mode_id=None)[source]¶
- _SUPPORT_DEPENDENCY_EDGE_FAMILIES = ('support', 'connects_port', 'connector', 'power', 'supplies', 'supply', 'cool',...[source]¶
- _common_cause_score_from_features(shared_dependency_signal, shared_upstream_signal, symptom_convergence_signal, governance_commonality_signal, train_oos_signal=0.0)[source]¶
- Parameters:
train_oos_signal (float)
- _common_cause_features_for_candidate(candidate, kg_context, telemetry_summary, pm_compliance, common_cause_index, candidate_component_id=None, operational_context=None)[source]¶
- _lookup_tskr_pattern(tskr_index, target_id)[source]¶
Return the highest-confidence pattern for target_id, or None.
- _pattern_latency_alignment(pattern)[source]¶
- Parameters:
pattern (Optional[JsonDict])
- Return type:
float
- _pattern_temporal_contradiction(pattern)[source]¶
- Parameters:
pattern (Optional[JsonDict])
- Return type:
bool
- _refresh_candidate_confidence_and_thresholds(candidate)[source]¶
- Parameters:
candidate (JsonDict)
- Return type:
None
- _update_score_rationale_for_refinement(candidate, *, support_score, contradiction_score, contextual_score, prior_evidence_score, authority_tier, authority_weight)[source]¶
- Parameters:
candidate (JsonDict)
support_score (float)
contradiction_score (float)
contextual_score (float)
prior_evidence_score (float)
authority_tier (Optional[str])
authority_weight (float)
- Return type:
None