src.dackar.RCA.log_pattern_recognition.rca_pattern_search ========================================================= .. py:module:: src.dackar.RCA.log_pattern_recognition.rca_pattern_search .. autoapi-nested-parse:: rca_pattern_search — RCA Temporal Pattern Matching Two-stage pipeline: Stage 1 (offline): KDE-based episode detection from historical event logs Stage 2 (online): Coarse-to-fine similarity retrieval Typical usage:: from dackar.RCA.log_pattern_recognition.rca_pattern_search import ( SearchConfig, IncidentIndex, PatternSearcher, IncidentExtractor, ) cfg = SearchConfig() # --- Stage 1: build historical index (run offline) --- index = IncidentIndex(cfg) index.build_from_history(events_df, rho_query=rho, query_duration=dur) index.save("/path/to/index") # --- Stage 2: query at runtime --- index = IncidentIndex.load("/path/to/index", cfg) extractor = IncidentExtractor(cfg) query_fp = extractor.extract(alarm_log, soe_log, telemetry, "INC_001", t0, t1) searcher = PatternSearcher(index, cfg) results = searcher.search(query_fp) Submodules ---------- .. toctree:: :maxdepth: 1 /autoapi/src/dackar/RCA/log_pattern_recognition/rca_pattern_search/config/index /autoapi/src/dackar/RCA/log_pattern_recognition/rca_pattern_search/density/index /autoapi/src/dackar/RCA/log_pattern_recognition/rca_pattern_search/extractor/index /autoapi/src/dackar/RCA/log_pattern_recognition/rca_pattern_search/indexer/index /autoapi/src/dackar/RCA/log_pattern_recognition/rca_pattern_search/metrics/index /autoapi/src/dackar/RCA/log_pattern_recognition/rca_pattern_search/models/index /autoapi/src/dackar/RCA/log_pattern_recognition/rca_pattern_search/searcher/index Classes ------- .. autoapisummary:: src.dackar.RCA.log_pattern_recognition.rca_pattern_search.PatternSearchConfig src.dackar.RCA.log_pattern_recognition.rca_pattern_search.SearchConfig src.dackar.RCA.log_pattern_recognition.rca_pattern_search.EpisodeDetector src.dackar.RCA.log_pattern_recognition.rca_pattern_search.IncidentExtractor src.dackar.RCA.log_pattern_recognition.rca_pattern_search.IncidentIndex src.dackar.RCA.log_pattern_recognition.rca_pattern_search.HistoricalSignalEpisode src.dackar.RCA.log_pattern_recognition.rca_pattern_search.IncidentFingerprint src.dackar.RCA.log_pattern_recognition.rca_pattern_search.UnifiedEvent src.dackar.RCA.log_pattern_recognition.rca_pattern_search.PatternSearcher Functions --------- .. autoapisummary:: src.dackar.RCA.log_pattern_recognition.rca_pattern_search.combined_score src.dackar.RCA.log_pattern_recognition.rca_pattern_search.emd_similarity src.dackar.RCA.log_pattern_recognition.rca_pattern_search.jaccard src.dackar.RCA.log_pattern_recognition.rca_pattern_search.nlcs Package Contents ---------------- .. py:class:: PatternSearchConfig Operational configuration for the pattern search subsystem. Kept separate from SearchConfig (which tunes the similarity algorithm) and from CrossPatternConfig (Phase 2). enable_signal_episode_search: master switch; when False the subsystem is completely bypassed and historical_signal_episodes.json is not written. index_staleness_window_days: episode index older than this is flagged "stale"; links built against stale results are capped at confidence 0.70 (§4.11). search_config: algorithm tuning (thresholds, weights, window expansion, etc.). .. py:attribute:: enable_signal_episode_search :type: bool :value: False .. py:attribute:: index_staleness_window_days :type: int :value: 30 .. py:attribute:: search_config :type: SearchConfig .. py:class:: SearchConfig Central configuration for the RCA pattern search pipeline. All parameters are tunable and should be validated empirically against historical data before production use. .. py:attribute:: beta :type: float :value: 0.2 .. py:attribute:: delta :type: float :value: 0.5 .. py:attribute:: kde_bandwidth :type: float | str :value: 'auto' .. py:attribute:: freq_threshold :type: int :value: 5 .. py:attribute:: min_jaccard :type: float :value: 0.3 .. py:attribute:: top_k :type: int :value: 5 .. py:attribute:: alpha :type: float :value: 0.3333333333333333 .. py:attribute:: beta_w :type: float :value: 0.3333333333333333 .. py:attribute:: gamma :type: float :value: 0.3333333333333333 .. py:attribute:: weight_profile :type: str :value: 'equal' .. py:attribute:: emd_normalization_mode :type: str :value: 'tv' .. py:method:: __post_init__() .. py:method:: resolve_weights(profile = None) Returns (alpha, beta_w, gamma) for the given profile name. If profile is None, uses self.weight_profile. Raises ValueError for unrecognised profile names. .. py:class:: EpisodeDetector(config) Detects episode boundaries in a continuous historical event log using kernel density estimation (KDE) over event timestamps. The detection threshold is relative to the query incident density, making the method self-calibrating: a busier query requires historically busier windows to qualify as matching episodes. Pipeline (per detect() call): 1. Convert event timestamps to float seconds since earliest event. 2. Evaluate Gaussian KDE on a fine time grid. 3. Threshold: mask = KDE(t) >= delta * rho_query. 4. Extract contiguous masked regions as raw episode boundaries. 5. Apply beta buffer expansion to each boundary. 6. Merge overlapping expanded boundaries. 7. Discard episodes shorter than query_duration / 10. .. py:attribute:: config .. py:method:: compute_reference_density(query_events, window_start, window_end) Computes the reference event density for the query incident. rho_query = N_query / D_query N_query: events with timestamp_start in [window_start, window_end]. D_query: (window_end - window_start).total_seconds() All sources contribute equally. Returns 0.0 if duration <= 0. .. py:method:: detect(historical_events, rho_query, query_duration) Detects episode boundaries in the historical event log. :param historical_events: Full flat event list, all sources merged. No episode_id expected at this stage. :param rho_query: Reference density from compute_reference_density(). Units: events per second. :param query_duration: D_query in seconds. Used for KDE bandwidth and minimum episode duration filter. :returns: List of (episode_start, episode_end) tuples, sorted ascending. These are expanded boundaries (beta already applied) ready for fingerprinting. Empty list if no qualifying episodes found. .. py:method:: bandwidth_scan(historical_events, rho_query, query_duration, bandwidths = None) Multi-scale diagnostic: counts detected episodes at different bandwidths. Helps operators validate episode segmentation by showing how many episodes are detected when the smoothing scale varies. Useful when query timescale may not match historical episode timescales (e.g., fast transient query for slow degradation history, or vice versa). :param historical_events: Full flat event list, all sources merged. :param rho_query: Reference density from compute_reference_density(). Units: events per second. :param query_duration: D_query in seconds. :param bandwidths: Explicit bandwidth list in seconds. If None, defaults to [D/32, D/16, D/8, D/4, D/2, D, 2D, 4D] for broad coverage. :returns: dict mapping bandwidth (float, seconds) to episode count (int), sorted by bandwidth ascending. .. py:method:: _run_detection(t_seconds, t_epoch, rho_query, query_duration, bw) Core detection pipeline: KDE → threshold → extract → expand → merge → filter. Called by detect() and bandwidth_scan() to avoid code duplication. :param t_seconds: Event timestamps as float seconds since t_epoch. :param t_epoch: Reference time (datetime). :param rho_query: Reference density (events/second). :param query_duration: Duration of query window in seconds. :param bw: Bandwidth in seconds (already resolved, > 0). :returns: Expanded, merged, filtered episode boundaries. .. py:method:: assign_episode_ids(historical_events, episode_boundaries) Assigns episode_id to each historical event based on detected boundaries. An event is assigned to the episode whose expanded boundary contains its timestamp_start. Events outside all boundaries retain episode_id = None (background noise). If an event falls within multiple boundaries (should not occur after merging, handled defensively), it is assigned to the first match and a WARNING is logged. Episode ids: "EP_{asset_id}_{index:05d}" where asset_id is the dominant asset among events in that episode. :param historical_events: Full flat event list. :param episode_boundaries: Output of detect(), sorted (start, end) tuples. :returns: New list of UnifiedEvents with episode_id populated. Original list is not mutated. .. py:class:: IncidentExtractor(config) Converts raw source records to UnifiedEvents and derives IncidentFingerprints. Two public entry points: to_unified_events() — normalise all three sources into a flat list extract() — full pipeline: expand window → filter → fingerprint _derive_fingerprint() is a staticmethod so that indexer.py can call it directly on pre-filtered episode event lists without instantiating an extractor. .. py:attribute:: config .. py:method:: to_unified_events(alarm_log, soe_log, telemetry_summaries, *, anomaly_severity_threshold = _DEFAULT_SEVERITY_THRESHOLD) Converts raw schema records to a flat UnifiedEvent list. Applies filtering defaults: - Alarms with state == "suppressed" are excluded. - Anomalies with promoted_to_kg_event == False are excluded. If the field is absent, severity_score >= anomaly_severity_threshold is used as the inclusion gate. - All SOE records are included. :param alarm_log: Dict with key "alarms": list of alarm dicts. :param soe_log: Dict with key "records": list of SOE record dicts. :param telemetry_summaries: List of telemetry summary dicts, each with key "anomalies": list of anomaly dicts. :param anomaly_severity_threshold: Fallback gate when promoted_to_kg_event absent. :returns: Flat list of UnifiedEvent, unsorted. .. py:method:: extract(alarm_log, soe_log, telemetry_summaries, incident_id, window_start, window_end, metadata = None) Full extraction pipeline for a single query incident. Steps: 1. Apply beta buffer to compute expanded window. 2. Call to_unified_events() for all sources. 3. Filter events to those with timestamp_start in expanded window. 4. Compute density over the expanded window (consistent with historical episode density). 5. Derive event_set, event_seq, freq_vec via _derive_fingerprint(). :param alarm_log: Raw alarm log dict. :param soe_log: Raw SOE log dict. :param telemetry_summaries: List of telemetry summary dicts. :param incident_id: Identifier for this query incident. :param window_start: Incident window start (before buffer expansion). :param window_end: Incident window end (before buffer expansion). :param metadata: Optional dict; "known_rca" and "asset_id" are read if present. :returns: IncidentFingerprint with expanded window stored as window_start/end. .. py:method:: _derive_fingerprint(events, freq_threshold) :staticmethod: Derives the three similarity representations from a list of events. High-frequency event types (count > freq_threshold) are excluded from event_set and event_seq but retained in freq_vec. :param events: Events belonging to a single episode or incident. :param freq_threshold: Count above which a type is considered high-frequency. :returns: (event_set, event_seq, freq_vec) .. py:method:: _parse_alarm_log(alarm_log) .. py:method:: _parse_soe_log(soe_log) .. py:method:: _parse_telemetry_summary(summary, severity_threshold) .. py:class:: IncidentIndex(config) Stores and manages pre-computed IncidentFingerprints for the historical episode database. Internal storage: episodes_df — pd.DataFrame, one row per episode. _inverted_index — dict[str, set[str]]: event_type → set of episode_ids. Used for O(|query_event_set|) Jaccard pre-filtering. .. py:attribute:: config .. py:attribute:: episodes_df :type: pandas.DataFrame .. py:attribute:: _inverted_index :type: dict[str, set[str]] .. py:attribute:: emd_normalization_factor :type: Optional[float] :value: None .. py:attribute:: build_timestamp :type: Optional[datetime.datetime] :value: None .. py:attribute:: asset_scope :type: list[str] :value: [] .. py:method:: build_from_history(events_df, rho_query, query_duration) Populates the index from a raw historical events_df. Steps: 1. Convert events_df rows to UnifiedEvent list. 2. Run EpisodeDetector.detect() to find episode boundaries. 3. Group events by episode boundary. 4. Derive IncidentFingerprint for each episode via IncidentExtractor._derive_fingerprint(). 5. Store fingerprints via add_batch(). :param events_df: Raw historical event log. Expected columns: raw_id, asset_id, source, event_type, timestamp_start, timestamp_end. No episode_id column required. :param rho_query: Reference density from the query incident (events per second). Passed to EpisodeDetector. :param query_duration: D_query in seconds. Used for KDE bandwidth and minimum episode duration filter. .. rubric:: Notes events_df is not modified in place. episode_id is derived from the episode window_start (EP_{asset}_{window_start:%Y%m%dT%H%M%S}), so ids are stable across rebuilds. The add path upserts by episode_id: re-building over the same data replaces episodes in place (idempotent), while genuinely new episodes append — call reset() first for a clean rebuild. Stage 1 is intended to be built once against a *representative* query: rho_query/query_duration come from that query and calibrate episode boundaries (delta * rho_query; bandwidth query_duration/4), then the index is persisted via save() (save/load is not a per-query cache). Results are sensitive to these values — see rca_pattern_matching.md for guidance on choosing representative rho_query/query_duration. .. py:method:: reset() Clears all stored episodes and the inverted index. .. py:method:: add(fingerprint) Adds a single fingerprint to the index (upsert by episode_id). If an episode with the same episode_id already exists it is replaced, and the inverted index is rebuilt once to drop the old episode's stale postings. The common insert-new path stays incremental. Less efficient than add_batch() for many insertions because the inverted index is updated per fingerprint rather than rebuilt once. .. py:method:: add_batch(fingerprints) Adds multiple fingerprints in a single operation (upsert by episode_id). Existing episodes whose episode_id appears in the incoming batch are replaced, and duplicate ids within the batch collapse to the last occurrence, so the append path never produces duplicate episode_ids. Rebuilds the inverted index once after all insertions, which is more efficient than repeated add() calls. .. py:method:: get_candidates(query_event_set) Returns episode_ids of historical episodes that share at least one event type with the query. Uses the inverted index for O(|query_event_set|) lookup. Episodes with no event-type overlap cannot have Jaccard > 0 and are excluded. :param query_event_set: event_set from the query IncidentFingerprint. :returns: De-duplicated list of episode_ids sharing >= 1 event type with the query (candidates are accumulated in a set, so no id repeats). .. py:method:: compute_emd_normalization_factor(max_pairs = 1000) Computes the empirical maximum raw L1 distance across historical episode pairs. Used to normalise EMD scores when emd_normalization_mode="empirical_max". Should be called once after build_from_history() and before search(). Algorithm: - Extract all freq_vec dicts from episodes_df. - If N*(N-1)/2 <= max_pairs: evaluate all pairs. - Else: draw max_pairs random pairs without replacement (seeded). - For each pair (a, b): raw_l1 = Σ_t |a.get(t,0) - b.get(t,0)| - Store the maximum observed L1 distance. :param max_pairs: Maximum number of pairs to evaluate. If the index has fewer pairs than this, all pairs are used. :returns: The empirical maximum raw L1 distance (float >= 0). Returns 1.0 if index is empty or contains only one episode (to avoid division by zero and provide a sensible fallback). Side effect: Sets self.emd_normalization_factor to the computed value. .. py:method:: save(path) Persists the index to disk. episodes_df is saved as parquet with complex columns (frozenset, list, dict) JSON-serialised to strings. The inverted index is saved as JSON (sets → sorted lists). EMD metadata (normalization factor) is saved as JSON. All files are written atomically via a .tmp rename. :param path: Directory path. Created if it does not exist. .. py:method:: load(path, config) :classmethod: Loads a persisted index from disk. Reconstructs episodes_df (deserialising complex columns from JSON strings), the inverted index, and EMD metadata if available. :param path: Directory path written by save(). :param config: SearchConfig to attach to the loaded index. :returns: Populated IncidentIndex. :raises FileNotFoundError if core files (parquet, inverted index) are missing.: .. py:method:: _rebuild_inverted_index() Rebuilds the inverted index from scratch from episodes_df. .. py:method:: __len__() .. py:method:: __repr__() .. py:function:: combined_score(j, n, e, alpha, beta_w, gamma) Weighted combination of the three metric scores. Score = alpha · J + beta_w · NLCS + gamma · EMD Inputs are assumed to be valid (alpha + beta_w + gamma ≈ 1). No re-normalisation is applied here; the caller (PatternSearcher) is responsible for passing a coherent weight triple via SearchConfig. .. py:function:: emd_similarity(a, b, normalization_factor = None) Frequency-based similarity between two event count vectors. Captures the repetition signal that Jaccard and NLCS deliberately discard. Implementation uses the Total Variation (TV) distance between the two probability distributions derived by normalising the count vectors. For categorical distributions with unit ground distance between any two distinct types, TV distance equals the (unit-ground) Earth Mover's Distance: TV(P, Q) = 0.5 · Σ_t |P(t) − Q(t)| where P(t) = a[t] / Σa This is always in [0, 1], so: emd_similarity = 1 − TV(P, Q) Alternative (raw-count) normalisation: If normalization_factor is provided the raw L1 distance between the unnormalised count vectors is used instead: raw_emd = Σ_t |a.get(t, 0) − b.get(t, 0)| emd_similarity = max(0.0, 1 − raw_emd / normalization_factor) Suitable when an empirically derived or vocabulary-size-based upper bound is available (see spec open point on normalisation). Edge cases: Both empty → 1.0 (identical — neither has any events) One empty → 0.0 (maximally dissimilar) .. py:function:: jaccard(a, b) Set-based similarity between two event sets. J(A, B) = |A ∩ B| / |A ∪ B| Returns 0.0 when both sets are empty (undefined ratio treated as no similarity rather than perfect similarity, which is the safer default for retrieval purposes). .. py:function:: nlcs(a, b) Normalised Longest Common Subsequence similarity. NLCS(A, B) = |LCS(A, B)| / max(|A|, |B|) Operates on deduplicated ordered sequences so high-frequency events do not dominate the ordering signal. Returns 0.0 when both sequences are empty. .. py:class:: HistoricalSignalEpisode Public output type for PatternSearcher.search(). Represents a single historical signal episode retrieved for a query incident. Carries all three metric scores individually (§5 of the integration plan) and an index_status field that governs cross-pattern linkage eligibility (§4.11). Sentinel (no_episodes_indexed) instances have episode_id == "" and similarity_to_current == 0.0. Callers must check index_status before attempting linkage. .. py:attribute:: episode_id :type: str .. py:attribute:: asset_id :type: str .. py:attribute:: window_start :type: Optional[datetime.datetime] .. py:attribute:: window_end :type: Optional[datetime.datetime] .. py:attribute:: source_types :type: list[str] .. py:attribute:: event_set :type: frozenset[str] .. py:attribute:: event_seq :type: list[str] .. py:attribute:: freq_vec :type: dict[str, int] .. py:attribute:: similarity_to_current :type: float .. py:attribute:: jaccard_score :type: float .. py:attribute:: nlcs_score :type: float .. py:attribute:: emd_score :type: float .. py:attribute:: weight_profile :type: str .. py:attribute:: matched_events :type: set[str] .. py:attribute:: query_only_events :type: set[str] .. py:attribute:: episode_only_events :type: set[str] .. py:attribute:: episode_density :type: float .. py:attribute:: known_rca :type: Optional[str] .. py:attribute:: linked_doc_ids :type: list[str] .. py:attribute:: index_status :type: str .. py:class:: IncidentFingerprint Pre-computed similarity representations for a single incident or detected historical episode. This is the unit of comparison in the retrieval pipeline. Derived from a list of UnifiedEvents by IncidentExtractor.extract() or EpisodeDetector after episode boundary assignment. The three representations serve distinct metrics: event_set → Jaccard (what types occurred, ignoring order/repetition) event_seq → NLCS (what types occurred and in what order) freq_vec → EMD (how many times each type occurred) High-frequency event types (count > freq_threshold) are excluded from event_set and event_seq but retained in freq_vec. .. py:attribute:: episode_id :type: str .. py:attribute:: asset_id :type: str .. py:attribute:: window_start :type: datetime.datetime .. py:attribute:: window_end :type: datetime.datetime .. py:attribute:: density :type: float .. py:attribute:: event_set :type: frozenset[str] .. py:attribute:: event_seq :type: list[str] .. py:attribute:: freq_vec :type: dict[str, int] .. py:attribute:: known_rca :type: str | None :value: None .. py:attribute:: source_types :type: list[str] :value: [] .. py:class:: UnifiedEvent Canonical representation of a single event from any source. All three input sources (alarm, SOE, anomaly) are normalised into this structure before any further processing. Lifecycle: - Created by IncidentExtractor.to_unified_events() - episode_id is None until EpisodeDetector assigns membership - timestamp_end is nullable and not used in current similarity metrics but carried for traceability and future use .. py:attribute:: raw_id :type: str .. py:attribute:: asset_id :type: str .. py:attribute:: source :type: str .. py:attribute:: event_type :type: str .. py:attribute:: timestamp_start :type: datetime.datetime .. py:attribute:: timestamp_end :type: datetime.datetime | None .. py:attribute:: episode_id :type: str | None :value: None .. py:class:: PatternSearcher(index, config) Retrieves the top-k most similar historical episodes for a query incident. Pipeline per search() call: 1. Inverted-index lookup: episode_ids sharing ≥ 1 event type with query. 2. Jaccard pre-filter: discard candidates below config.min_jaccard. 3. NLCS computation on survivors. 4. EMD computation on survivors. 5. Combined score: weighted sum using the resolved weight profile. 6. Rank descending by combined_score; return the top-k HistoricalSignalEpisodes. The coarse-to-fine design avoids computing NLCS and EMD on clearly dissimilar episodes (those failing the Jaccard gate). .. py:attribute:: index .. py:attribute:: config .. py:method:: search(query, weight_profile = None, staleness_window_days = None) Retrieves top-k most similar historical episodes for a query fingerprint. Returns list[HistoricalSignalEpisode] with index_status populated on every result (§4.11): - "indexed" — normal result from a current, populated index - "no_episodes_indexed" — index is empty; returns a single sentinel episode - "stale" — index is older than staleness_window_days; results returned but flagged; link_confidence capped downstream :param query: IncidentFingerprint for the query incident. :param weight_profile: Weight profile override; None → config.weight_profile. :param staleness_window_days: If set and index.build_timestamp is known, episodes from an index older than this are marked "stale". None disables the staleness check. By design the window lives on PatternSearchConfig.index_staleness_window_days; the orchestrator reads it there and passes it in here (PatternSearcher itself only carries a SearchConfig). .. rubric:: Notes - All three metric scores (jaccard, nlcs, emd) are individually visible on every returned HistoricalSignalEpisode (§5). - matched_events, query_only_events, episode_only_events are derived from event_set comparison. .. py:method:: _resolve_weights(weight_profile) Returns (alpha, beta_w, gamma) for the given weight profile name. Delegates to SearchConfig.resolve_weights() which owns the profile registry and "custom" fallback logic. :raises ValueError for unrecognised profile names.: