From 26dc33708872c4c53fdfc34a472292f0a5b66e60 Mon Sep 17 00:00:00 2001 From: Bartosz Burda Date: Mon, 14 Sep 2026 17:18:21 +0200 Subject: [PATCH 01/14] feat(msgs): identify a fault record by code and reporting source A fault record is now the pair (fault_code, reporting source). The source is the source_id a ReportFault call carried, and it owns the record. ClearFault, GetFault, GetSnapshots and GetRosbag gain a trailing request field source_id naming that owner. Empty means unscoped: the call applies only when exactly one record carries the fault_code, and fails as ambiguous when several do. On GetRosbag the field scopes the fault_code lookup only, the recording_id path is unchanged. Fault.msg and ReportFault.srv do not change shape. Their header comments, FaultEvent.msg, ListFaultsForEntity.srv and the package README are rewritten to describe the per-record model instead of aggregation by code. --- src/ros2_medkit_msgs/README.md | 26 ++++++++++++++----- src/ros2_medkit_msgs/msg/Fault.msg | 15 +++++++---- src/ros2_medkit_msgs/msg/FaultEvent.msg | 3 ++- src/ros2_medkit_msgs/srv/ClearFault.srv | 11 ++++++-- src/ros2_medkit_msgs/srv/GetFault.srv | 8 +++++- src/ros2_medkit_msgs/srv/GetRosbag.srv | 9 ++++++- src/ros2_medkit_msgs/srv/GetSnapshots.srv | 8 +++++- .../srv/ListFaultsForEntity.srv | 3 ++- src/ros2_medkit_msgs/srv/ReportFault.srv | 10 ++++--- 9 files changed, 71 insertions(+), 22 deletions(-) diff --git a/src/ros2_medkit_msgs/README.md b/src/ros2_medkit_msgs/README.md index fe8cf92a7..9ad70535e 100644 --- a/src/ros2_medkit_msgs/README.md +++ b/src/ros2_medkit_msgs/README.md @@ -6,7 +6,7 @@ ROS 2 message and service definitions for the ros2_medkit fault management syste This package provides the interface definitions used by the fault management components: -- **FaultManager** (`ros2_medkit_fault_manager`) - Central fault aggregation and lifecycle management +- **FaultManager** (`ros2_medkit_fault_manager`) - Central fault record store and lifecycle management - **FaultReporter** (`ros2_medkit_fault_reporter`) - Client library for fault reporting - **Gateway** (`ros2_medkit_gateway`) - REST API endpoints for fault access @@ -14,7 +14,15 @@ This package provides the interface definitions used by the fault management com ### Fault.msg -Core fault data model representing an aggregated fault condition with AUTOSAR DEM-style debounce filtering. +Core fault data model representing one fault record with AUTOSAR DEM-style debounce filtering. + +A record is identified by the pair (`fault_code`, owning reporting source). The owner is the +`source_id` a `ReportFault` call carried. Two sources reporting one `fault_code` are two records, +each with its own status, debounce counter, `occurrence_count`, severity and timestamps, and each +cleared on its own. The services that act on a single record (`ClearFault`, `GetFault`, +`GetSnapshots`, `GetRosbag`) carry a `source_id` request field naming the owner. An empty +`source_id` is unscoped: the call applies only when exactly one record carries the `fault_code`, +and fails as ambiguous when several do. | Field | Type | Description | |-------|------|-------------| @@ -26,7 +34,7 @@ Core fault data model representing an aggregated fault condition with AUTOSAR DE | `last_passed` | builtin_interfaces/Time | When fault last reported PASSED (zero = never) | | `occurrence_count` | uint32 | Times this fault has occurred, counted on edges (first FAILED, then each FAILED that arrives while CLEARED). Repeats within one occurrence do not increment it | | `status` | string | Current status (see STATUS_* constants) | -| `reporting_sources` | string[] | List of source identifiers that reported this fault | +| `reporting_sources` | string[] | The reporting source that owns this record, as a one-element list | **Severity Levels:** | Constant | Value | Description | @@ -121,12 +129,14 @@ Query faults with optional filtering. ### ClearFault.srv -Clear/acknowledge a fault. Cleared faults are retained and queryable with `statuses=["CLEARED"]`. +Clear/acknowledge one fault record. Cleared records are retained and queryable with +`statuses=["CLEARED"]`. **Request:** | Field | Type | Description | |-------|------|-------------| -| `fault_code` | string | The fault to clear | +| `fault_code` | string | The fault code to clear | +| `source_id` | string | Reporting source that owns the record. Empty means unscoped: the call applies only when exactly one record carries the `fault_code`, otherwise it fails with a message beginning `ambiguous:` and clears nothing | | `skip_correlation_auto_clear` | bool | When `true`, only the requested fault_code is cleared; symptom faults that the correlation engine would normally auto-clear via `auto_clear_with_root` rules are left untouched. Default `false` (cascade clear). The gateway sets this to `true` on per-entity `DELETE /{entity-path}/faults/{fault_code}` so that an operator with access to one entity cannot cascade-clear correlated symptoms reported by apps in other entities. | **Response:** @@ -136,7 +146,11 @@ Clear/acknowledge a fault. Cleared faults are retained and queryable with `statu | `message` | string | Status or error message | | `auto_cleared_codes` | string[] | Symptom fault codes auto-cleared with the root cause (empty when `skip_correlation_auto_clear=true`) | -> **Note:** `skip_correlation_auto_clear` was added in `ros2_medkit_msgs` post-0.4.0. Adding a request field changes the service type hash, so out-of-tree callers that invoke `/fault_manager/clear_fault` directly (via `ros2 service call` or a generated client) must rebuild against the new `ros2_medkit_msgs` release to keep talking to `fault_manager`. +> **Note:** `skip_correlation_auto_clear` was added in `ros2_medkit_msgs` post-0.4.0, and `source_id` +> after it. `source_id` was added the same way to `GetFault`, `GetSnapshots` and `GetRosbag`. Adding a +> request field changes the service type hash, so out-of-tree callers that invoke those services +> directly (via `ros2 service call` or a generated client) must rebuild against the new +> `ros2_medkit_msgs` release to keep talking to `fault_manager`. ## Usage diff --git a/src/ros2_medkit_msgs/msg/Fault.msg b/src/ros2_medkit_msgs/msg/Fault.msg index aa24ec629..99227411b 100644 --- a/src/ros2_medkit_msgs/msg/Fault.msg +++ b/src/ros2_medkit_msgs/msg/Fault.msg @@ -14,16 +14,19 @@ # # Fault.msg - Core fault data model for ros2_medkit fault management system. # -# A Fault represents an aggregated fault condition identified by a global fault_code. -# Multiple sources can report the same fault_code, and they are aggregated into a -# single Fault with tracked occurrence_count and reporting_sources. +# A Fault is one fault record, identified by the pair (fault_code, reporting source). +# The reporting source is the source_id a ReportFault call carried, and it is the +# owner of the record. Two sources reporting one fault_code are two records, each +# with its own status, debounce counter, occurrence_count, severity and timestamps, +# and each cleared on its own. reporting_sources holds the one owner. # # Debounce model (AUTOSAR DEM-style): # - FAILED events decrement internal counter (towards confirmation) # - PASSED events increment internal counter (towards healing) # - Status reflects the current debounce state -# Global fault identifier (e.g., "MOTOR_OVERHEAT", "SENSOR_FAILURE_001") +# Global fault identifier (e.g., "MOTOR_OVERHEAT", "SENSOR_FAILURE_001"). +# Half of the record identity: the other half is the owning reporting source. string fault_code # Fault severity level (use SEVERITY_* constants) @@ -53,7 +56,9 @@ uint32 occurrence_count # Current fault status (use STATUS_* constants) string status -# List of source identifiers that have reported this fault +# The reporting source that owns this record, as a one-element list. Together with +# fault_code it identifies the record. Services that act on a single record take +# this value as their source_id request field. string[] reporting_sources # Severity level constants diff --git a/src/ros2_medkit_msgs/msg/FaultEvent.msg b/src/ros2_medkit_msgs/msg/FaultEvent.msg index 40cebe6bb..7fa119289 100644 --- a/src/ros2_medkit_msgs/msg/FaultEvent.msg +++ b/src/ros2_medkit_msgs/msg/FaultEvent.msg @@ -35,7 +35,8 @@ string EVENT_CONFIRMED = "fault_confirmed" # the two apart. string EVENT_CLEARED = "fault_cleared" # Emitted when fault data changes without status transition (e.g., last_occurred -# advanced, severity escalated, new source added to reporting_sources). +# advanced, severity escalated). A report from a source that owns no record yet +# is not one of them: it creates a record of its own. # occurrence_count is not one of them: it only changes on an occurrence edge # (the first FAILED, or a FAILED that arrives while the fault is CLEARED). string EVENT_UPDATED = "fault_updated" diff --git a/src/ros2_medkit_msgs/srv/ClearFault.srv b/src/ros2_medkit_msgs/srv/ClearFault.srv index cfa0b778f..340c58ace 100644 --- a/src/ros2_medkit_msgs/srv/ClearFault.srv +++ b/src/ros2_medkit_msgs/srv/ClearFault.srv @@ -14,8 +14,10 @@ # # ClearFault.srv - Clear/acknowledge a fault in the FaultManager. # -# Marks the specified fault as CLEARED. The fault record is retained for historical -# purposes and can be retrieved by calling ListFaults with statuses=["CLEARED"]. +# Marks one fault record as CLEARED. A record is identified by (fault_code, source_id), +# so clearing one owner's record never touches another owner's record of the same code. +# The record is retained for historical purposes and can be retrieved by calling +# ListFaults with statuses=["CLEARED"]. # Cleared faults are retained indefinitely unless storage cleanup is configured. # A FaultEvent with EVENT_CLEARED is published when a fault is successfully cleared. @@ -29,6 +31,11 @@ string fault_code # The fault_code to clear # Set this from scoped per-entity DELETE routes so a viewer of one entity # cannot cascade-clear symptom faults reported by apps in other entities. bool skip_correlation_auto_clear + +# Reporting source that owns the fault record (the source_id used in ReportFault). +# Together with fault_code it identifies one record. Empty: the call applies only +# when exactly one record carries fault_code, otherwise it fails as ambiguous. +string source_id --- # Response fields bool success # True if the fault was found and cleared diff --git a/src/ros2_medkit_msgs/srv/GetFault.srv b/src/ros2_medkit_msgs/srv/GetFault.srv index 337c2b5f6..427af3f7f 100644 --- a/src/ros2_medkit_msgs/srv/GetFault.srv +++ b/src/ros2_medkit_msgs/srv/GetFault.srv @@ -15,13 +15,19 @@ # GetFault.srv - Get single fault by code with environment data. # # Unlike ListFaults.srv (plural) which returns a list of faults for filtering, -# this service returns a single fault with full environment data including +# this service returns a single fault record with full environment data including # snapshots and extended data records for SOVD-compliant fault responses. +# A record is identified by (fault_code, source_id). # Request fields # Fault code to retrieve (e.g., "MOTOR_OVERHEAT") string fault_code + +# Reporting source that owns the fault record (the source_id used in ReportFault). +# Together with fault_code it identifies one record. Empty: the call applies only +# when exactly one record carries fault_code, otherwise it fails as ambiguous. +string source_id --- # Response fields diff --git a/src/ros2_medkit_msgs/srv/GetRosbag.srv b/src/ros2_medkit_msgs/srv/GetRosbag.srv index d8edf0e8e..77890285e 100644 --- a/src/ros2_medkit_msgs/srv/GetRosbag.srv +++ b/src/ros2_medkit_msgs/srv/GetRosbag.srv @@ -14,7 +14,8 @@ # # GetRosbag.srv - Retrieve rosbag file information for a fault. # -# Returns the path to the rosbag file that was recorded when a fault was confirmed. +# Returns the path to the rosbag file that was recorded when a fault record was +# confirmed. A recording is linked per record, identified by (fault_code, source_id). # The rosbag contains time-window recording of configured topics before and after # the fault confirmation event. @@ -29,6 +30,12 @@ string fault_code # fault_code alone keeps its historical meaning, "the newest recording of this # fault". string recording_id + +# Reporting source that owns the fault record (the source_id used in ReportFault). +# Together with fault_code it identifies one record. Empty: the call applies only +# when exactly one record carries fault_code, otherwise it fails as ambiguous. +# It scopes the fault_code lookup only. The recording_id path ignores it. +string source_id --- # Response fields diff --git a/src/ros2_medkit_msgs/srv/GetSnapshots.srv b/src/ros2_medkit_msgs/srv/GetSnapshots.srv index de1bf5f71..fcd8dfe91 100644 --- a/src/ros2_medkit_msgs/srv/GetSnapshots.srv +++ b/src/ros2_medkit_msgs/srv/GetSnapshots.srv @@ -14,7 +14,8 @@ # # GetSnapshots.srv - Retrieve topic snapshots captured when a fault was confirmed. # -# Returns topic data that was captured at the moment a fault transitioned to CONFIRMED. +# Returns topic data that was captured at the moment a fault record transitioned to +# CONFIRMED. Snapshots belong to one record, identified by (fault_code, source_id). # Snapshots provide debugging context by preserving system state at fault time. # Request fields @@ -25,6 +26,11 @@ string fault_code # Optional topic filter. If empty, returns snapshots for all captured topics. # If specified, returns only the snapshot for that specific topic. string topic + +# Reporting source that owns the fault record (the source_id used in ReportFault). +# Together with fault_code it identifies one record. Empty: the call applies only +# when exactly one record carries fault_code, otherwise it fails as ambiguous. +string source_id --- # Response fields diff --git a/src/ros2_medkit_msgs/srv/ListFaultsForEntity.srv b/src/ros2_medkit_msgs/srv/ListFaultsForEntity.srv index c58bdcc14..5ae09449a 100644 --- a/src/ros2_medkit_msgs/srv/ListFaultsForEntity.srv +++ b/src/ros2_medkit_msgs/srv/ListFaultsForEntity.srv @@ -14,7 +14,8 @@ # # ListFaultsForEntity.srv - Request list of faults for a specific entity. # -# Filters faults by checking if entity FQN appears in fault.reporting_sources[]. +# Filters fault records by their owning reporting source: the entity FQN must match +# the single entry of fault.reporting_sources[], exactly or as an FQN suffix. # Used by the gateway to serve entity-specific fault lists. # Request fields diff --git a/src/ros2_medkit_msgs/srv/ReportFault.srv b/src/ros2_medkit_msgs/srv/ReportFault.srv index c6403ca9f..e26fda5d9 100644 --- a/src/ros2_medkit_msgs/srv/ReportFault.srv +++ b/src/ros2_medkit_msgs/srv/ReportFault.srv @@ -15,14 +15,14 @@ # ReportFault.srv - Report a fault event to the FaultManager. # # Called by FaultReporter clients or directly by nodes to report fault conditions. -# The FaultManager aggregates reports by fault_code across all sources and applies -# debounce filtering based on FAILED/PASSED event counts. +# The FaultManager keeps one record per (fault_code, source_id) pair and applies +# debounce filtering per record, based on FAILED/PASSED event counts. # Request fields # Global fault identifier. Convention: UPPER_SNAKE_CASE (e.g., "MOTOR_OVERHEAT"). # Recommended max length: 64 characters. Allowed: A-Z, 0-9, underscore. -# Same fault_code from different sources = single aggregated fault. +# Same fault_code from different sources = one record each, filtered independently. string fault_code # Event type indicating whether the fault condition is detected or cleared. @@ -47,7 +47,9 @@ string description # Identifier of the reporting node/entity. Recommended: fully qualified ROS 2 node name # including namespace (e.g., "/powertrain/engine/temp_sensor"). -# This value is added to the fault's reporting_sources array for multi-source tracking. +# This value owns the record: (fault_code, source_id) is the record identity, and the +# services that read or clear one record take it back as their source_id field. +# Required: a report with an empty source_id is rejected. string source_id --- # Response fields From 9cbdf86557623dc1f9518fea0bf3593efe6b7aaf Mon Sep 17 00:00:00 2001 From: Bartosz Burda Date: Mon, 14 Sep 2026 20:30:08 +0200 Subject: [PATCH 02/14] feat(fault_manager): key every fault record by code and owner A fault record is now identified by the pair (fault_code, owner), where the owner is the source_id of the ReportFault call that created it. Two sources reporting one code are two records: status, debounce counter, occurrence count, timestamps, severity, freeze frame, snapshots, near misses, rosbag links and capture cooldown are all per record, and a clear or a heal driven by one owner never touches another owner's record. Storage: FaultId addresses a record, FaultState carries its owner and emits it as the single entry of reporting_sources, and the virtual API takes a FaultId wherever it used to take a bare code. get_faults_by_code is new and backs the unscoped service resolution. check_time_based_confirmation and reclassify_healed_as_cleared return identities, because one code can hold a moved record for one owner and an untouched one for another. SQLite: faults gains owner with a composite unique index in place of the single-column primary key, freeze_frames is rebuilt the same way, snapshots and rosbag_files gain owner, and the rosbag unique index widens. A legacy database is migrated by an explicit rebuild under one transaction, probing each table on its own and backfilling the owner from the first entry of the old reporting_sources array, which is the only per-source fact the old schema recorded. Re-opening a migrated database changes nothing. Node: ClearFault, GetFault, GetSnapshots and GetRosbag resolve their target record from source_id, or unscoped when exactly one record carries the code, failing with an ambiguous: message that lists the owners otherwise. Every audit row now names the record's owner, so the clear_service, auto_heal, auto_confirm_timer and startup_reclassify literals are gone. Capture, snapshot writing and rosbag linking are per record, and correlation forms every relation between records of one owner. The unit tests that pinned one record per code are rewritten to the record contract rather than deleted, and the new cases cover two owners of one code, the unscoped resolution, the legacy-database migration and the burst that links one recording to two records. --- .../capture_thread_pool.hpp | 13 +- .../correlation/correlation_engine.hpp | 72 +- .../fault_manager_node.hpp | 18 +- .../fault_storage.hpp | 248 +++--- .../rosbag_capture.hpp | 53 +- .../snapshot_capture.hpp | 29 +- .../sqlite_fault_storage.hpp | 63 +- .../src/capture_thread_pool.cpp | 16 +- .../src/correlation/correlation_engine.cpp | 113 +-- .../src/fault_manager_node.cpp | 226 +++-- .../src/fault_storage.cpp | 183 ++-- .../src/rosbag_capture.cpp | 110 +-- .../src/snapshot_capture.cpp | 61 +- .../src/sqlite_fault_storage.cpp | 784 ++++++++++-------- .../test/test_capture_thread_pool.cpp | 87 +- .../test/test_correlation_engine.cpp | 242 ++++-- .../test/test_entity_thresholds.cpp | 43 +- .../test/test_fault_manager.cpp | 669 +++++++++++---- .../test/test_rosbag_capture.cpp | 423 ++++++---- .../test/test_rosbag_storage_parity.cpp | 74 +- .../test/test_snapshot_capture.cpp | 104 ++- .../test/test_sqlite_storage.cpp | 690 ++++++++++----- 22 files changed, 2703 insertions(+), 1618 deletions(-) diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/capture_thread_pool.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/capture_thread_pool.hpp index 7cff1d63d..0c1bb1361 100644 --- a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/capture_thread_pool.hpp +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/capture_thread_pool.hpp @@ -27,6 +27,7 @@ #include #include "rclcpp/logger.hpp" +#include "ros2_medkit_fault_manager/fault_storage.hpp" namespace ros2_medkit_fault_manager { @@ -46,7 +47,7 @@ enum class EnqueueResult { struct EnqueueOutcome { EnqueueResult result; - std::optional evicted_code; ///< Set only for kEvictedOldest. + std::optional evicted_id; ///< Set only for kEvictedOldest. }; /// Bounded worker pool that runs fault-capture jobs off the service thread. @@ -63,7 +64,7 @@ class CaptureThreadPool { /// @param capture_fn Invoked per job on a worker thread. Must be thread-safe /// for pool_size concurrent calls. Exceptions are caught and logged. CaptureThreadPool(std::size_t pool_size, std::size_t queue_depth, QueueFullPolicy full_policy, rclcpp::Logger logger, - std::function capture_fn); + std::function capture_fn); ~CaptureThreadPool(); CaptureThreadPool(const CaptureThreadPool &) = delete; @@ -71,8 +72,8 @@ class CaptureThreadPool { CaptureThreadPool(CaptureThreadPool &&) = delete; CaptureThreadPool & operator=(CaptureThreadPool &&) = delete; - /// Enqueue a capture job. Non-blocking. Thread-safe. - EnqueueOutcome enqueue(const std::string & fault_code); + /// Enqueue a capture job for one fault record. Non-blocking. Thread-safe. + EnqueueOutcome enqueue(const FaultId & id); /// Stop accepting work, let in-flight jobs finish, discard pending, join all /// workers. Idempotent and noexcept. Called by the destructor. @@ -92,11 +93,11 @@ class CaptureThreadPool { const std::size_t queue_depth_; const QueueFullPolicy full_policy_; rclcpp::Logger logger_; - std::function capture_fn_; + std::function capture_fn_; mutable std::mutex queue_mutex_; std::condition_variable cv_; - std::deque queue_; + std::deque queue_; bool stop_{false}; std::atomic dropped_captures_{0}; diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/correlation/correlation_engine.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/correlation/correlation_engine.hpp index f1092138a..7a3007522 100644 --- a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/correlation/correlation_engine.hpp +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/correlation/correlation_engine.hpp @@ -20,10 +20,12 @@ #include #include #include +#include #include #include "ros2_medkit_fault_manager/correlation/pattern_matcher.hpp" #include "ros2_medkit_fault_manager/correlation/types.hpp" +#include "ros2_medkit_fault_manager/fault_storage.hpp" namespace ros2_medkit_fault_manager { namespace correlation { @@ -56,11 +58,15 @@ struct ProcessFaultResult { /// Result of clearing a fault struct ProcessClearResult { - /// List of symptom fault codes that should be auto-cleared - std::vector auto_cleared_codes; + /// The symptom RECORDS that should be auto-cleared. Records, not codes: a root + /// cause of one owner must never auto-clear another owner's record of the same + /// symptom code. + std::vector auto_cleared_symptoms; }; -/// Information about a muted fault (for ListFaults response) +/// Information about a muted record (for ListFaults response). +/// Shape unchanged: MutedFaultInfo carries the CODE, so a muted (code, owner) appears +/// on the wire as its code. struct MutedFaultData { std::string fault_code; std::string root_cause_code; @@ -95,18 +101,24 @@ class CorrelationEngine { /// @param config Correlation configuration (must be enabled and valid) explicit CorrelationEngine(const CorrelationConfig & config); - /// Process an incoming fault - /// @param fault_code The fault code + /// Process an incoming fault record + /// + /// Rules match on the fault CODE - a rule says which codes are root causes and which + /// are symptoms, and says nothing about who reports them. Every RELATION the engine + /// then forms (root to symptom, muting, cluster membership, auto-clear cascade) is + /// between records of the SAME owner, so a root cause of owner X never mutes or + /// auto-clears a symptom of owner Y. + /// @param id The fault record /// @param severity The fault severity (for representative selection) /// @param timestamp When the fault occurred /// @return Processing result indicating whether to mute, correlations, etc. - ProcessFaultResult process_fault(const std::string & fault_code, const std::string & severity, + ProcessFaultResult process_fault(const FaultId & id, const std::string & severity, std::chrono::steady_clock::time_point timestamp = std::chrono::steady_clock::now()); - /// Process a fault being cleared - /// @param fault_code The fault code being cleared - /// @return Result with list of symptoms to auto-clear - ProcessClearResult process_clear(const std::string & fault_code); + /// Process a fault record being cleared + /// @param id The record being cleared + /// @return Result with the symptom records to auto-clear, all of the same owner + ProcessClearResult process_clear(const FaultId & id); /// Get all currently muted faults /// @return List of muted fault data @@ -115,10 +127,10 @@ class CorrelationEngine { /// Get count of muted faults uint32_t get_muted_count() const; - /// Whether a fault code is currently muted as a symptom. - /// @param fault_code Code to test - /// @return True while the code is suppressed by a root cause - bool is_muted(const std::string & fault_code) const; + /// Whether a fault record is currently muted as a symptom. + /// @param id Record to test + /// @return True while the record is suppressed by a root cause of the same owner + bool is_muted(const FaultId & id) const; /// Get all active clusters /// @return List of cluster data @@ -132,18 +144,18 @@ class CorrelationEngine { void cleanup_expired(); private: - /// Check if fault matches a root cause pattern in any hierarchical rule + /// Check if the record's CODE matches a root cause pattern in any hierarchical rule /// @return Rule ID if matched, empty optional otherwise std::optional try_as_root_cause(const std::string & fault_code); - /// Check if fault is a symptom of any pending root cause + /// Check if the record is a symptom of a pending root cause OF THE SAME OWNER /// @return ProcessFaultResult with correlation info if matched - std::optional try_as_symptom(const std::string & fault_code, - std::chrono::steady_clock::time_point timestamp); + std::optional try_as_symptom(const FaultId & id, std::chrono::steady_clock::time_point timestamp); - /// Check if fault matches an auto-cluster rule + /// Check if the record's code matches an auto-cluster rule. A cluster holds records + /// of one owner, so the pending-cluster key carries the owner alongside the rule id. /// @return ProcessFaultResult with cluster info if matched - std::optional try_auto_cluster(const std::string & fault_code, const std::string & severity, + std::optional try_auto_cluster(const FaultId & id, const std::string & severity, std::chrono::steady_clock::time_point timestamp); /// Generate unique cluster ID @@ -154,35 +166,37 @@ class CorrelationEngine { /// Active root causes waiting for symptoms struct PendingRootCause { - std::string fault_code; + FaultId fault_id; std::string rule_id; std::chrono::steady_clock::time_point timestamp; uint32_t window_ms; }; std::vector pending_root_causes_; - /// Mapping from root cause to its symptoms - std::map> root_to_symptoms_; + /// Mapping from a root cause RECORD to its symptom RECORDS, all of the same owner + std::map> root_to_symptoms_; - /// Muted faults (fault_code -> data) - std::map muted_faults_; + /// Muted records (record -> data) + std::map muted_faults_; /// Active clusters (cluster_id -> data) std::map active_clusters_; - /// Mapping from fault code to cluster ID (for faults in clusters) - std::map fault_to_cluster_; + /// Mapping from a record to its cluster ID (for records in clusters) + std::map fault_to_cluster_; /// Pending cluster with steady_clock timestamp for window tracking struct PendingCluster { ClusterData data; + std::string owner; ///< Every member record of this cluster has this owner std::chrono::steady_clock::time_point steady_first_at; std::map fault_severities; ///< fault_code -> severity }; - /// Pending clusters being formed (rule_id -> cluster data) + /// Pending clusters being formed ((rule_id, owner) -> cluster data). Keyed by owner + /// too, so one rule forms one cluster per owner and membership never crosses owners. /// Once min_count is reached, moved to active_clusters_ - std::map pending_clusters_; + std::map, PendingCluster> pending_clusters_; /// Counter for cluster ID generation uint64_t cluster_counter_{0}; diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_manager_node.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_manager_node.hpp index 9157102d9..3813d35a8 100644 --- a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_manager_node.hpp +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_manager_node.hpp @@ -17,9 +17,11 @@ #include #include #include +#include #include +#include #include -#include +#include #include "rclcpp/rclcpp.hpp" #include "ros2_medkit_fault_manager/capture_thread_pool.hpp" @@ -119,6 +121,16 @@ class FaultManagerNode : public rclcpp::Node { /// @return true if entity_id matches any source static bool matches_entity(const std::vector & reporting_sources, const std::string & entity_id); + /// Resolve the record a service request addresses. + /// + /// With a non-empty source_id the answer is that one record. With an empty one the + /// call is unscoped: it applies only when exactly one record carries the code, and + /// with several it fails, because picking one would clear or serve an arbitrary + /// owner's record. @p error is filled on failure and is empty on success, with a + /// message beginning "ambiguous:" listing the owners in the several-records case. + std::optional resolve_target(const std::string & fault_code, const std::string & source_id, + std::string & error) const; + private: /// Create storage backend based on configuration std::unique_ptr create_storage(); @@ -179,7 +191,7 @@ class FaultManagerNode : public rclcpp::Node { /// Shared by the report path and the time-based confirmation timer so a /// confirmation produces the same evidence whichever one produced it. /// @param fault_code Code of the fault that reached CONFIRMED - void capture_on_confirm(const std::string & fault_code); + void capture_on_confirm(const FaultId & id); /// Validate severity value static bool is_valid_severity(uint8_t severity); @@ -262,7 +274,7 @@ class FaultManagerNode : public rclcpp::Node { /// Per-fault cooldown tracking for snapshot recapture std::mutex last_capture_mutex_; - std::unordered_map last_capture_times_; + std::map last_capture_times_; }; } // namespace ros2_medkit_fault_manager diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_storage.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_storage.hpp index 3bc60dc35..b56bc9a48 100644 --- a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_storage.hpp +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/fault_storage.hpp @@ -20,6 +20,7 @@ #include #include #include +#include #include #include "rclcpp/rclcpp.hpp" @@ -87,16 +88,42 @@ bool is_near_miss(bool is_failed_event, const std::string & resulting_status); /// defaults (-1 / 3). Returns true if the config was already valid. bool sanitize_debounce_config(DebounceConfig & config); -/// Internal fault state stored in memory +/// Identity of one fault record: a fault code plus the reporting source that owns it. +/// +/// The owner is the source_id the ReportFault call carried. Two sources reporting one +/// fault_code are two records, and every piece of per-fault state (status, debounce +/// counter, occurrence count, timestamps, severity, freeze frame, snapshots, near +/// misses, rosbag links) belongs to one record. A clear or a heal driven by one owner +/// never touches another owner's record of the same code. +/// +/// An aggregate on purpose: FaultId{code, owner} is the only way to build one, so no +/// call site can silently pass a bare code where a record identity is required. +struct FaultId { + std::string fault_code; + std::string owner; + + bool operator==(const FaultId & other) const { + return fault_code == other.fault_code && owner == other.owner; + } + + /// Ordered by code first, so the records of one code are contiguous in a map and + /// get_faults_by_code is a range scan rather than a full sweep. + bool operator<(const FaultId & other) const { + return std::tie(fault_code, owner) < std::tie(other.fault_code, other.owner); + } +}; + +/// Internal fault state of one record, stored in memory struct FaultState { std::string fault_code; + /// Reporting source that owns this record. Half of the record identity. + std::string owner; uint8_t severity{0}; std::string description; rclcpp::Time first_occurred; rclcpp::Time last_occurred; uint32_t occurrence_count{0}; ///< Count of genuine occurrences (new fault + each re-raise after CLEARED) std::string status; - std::set reporting_sources; // Debounce state (internal, not exposed in Fault.msg) int32_t debounce_counter{0}; ///< FAILED decrements (-1), PASSED increments (+1) @@ -116,6 +143,7 @@ using EventType = ros2_medkit_msgs::srv::ReportFault::Request; /// one moment and only mean anything together. struct SnapshotData { std::string fault_code; + std::string owner; ///< Reporting source that owns the record this reading belongs to std::string topic; std::string message_type; std::string data; ///< JSON-encoded message data @@ -127,13 +155,15 @@ struct SnapshotData { /// Compact freeze-frame captured when a fault confirms: a single JSON object mapping /// each captured topic to its latest value at confirmation time. Unlike per-topic -/// snapshots, a freeze-frame is keyed by fault_code (one row per code) and is RETAINED -/// across clear_fault, so the confirmed-state record persists after acknowledgement. -/// A row exists only for fault codes with a configured capture set; a fault code with -/// no capture configured gets no row at all (lookup returns nullopt, never an empty {}). +/// snapshots, a freeze-frame is keyed by the record identity (one row per +/// (fault_code, owner) pair) and is RETAINED across clear_fault, so the confirmed-state +/// record persists after acknowledgement. A row exists only for fault codes with a +/// configured capture set; a fault code with no capture configured gets no row at all +/// (lookup returns nullopt, never an empty {}). struct FreezeFrameData { std::string fault_code; - std::string data; ///< Compact JSON object: {"": , ...} + std::string owner; ///< Reporting source that owns the record this frame belongs to + std::string data; ///< Compact JSON object: {"": , ...} int64_t captured_at_ns{0}; }; @@ -147,7 +177,7 @@ struct FreezeFrameData { /// off the wire with its own copy of this rule; keep the two in step. std::string rosbag_recording_id(const std::string & file_path); -/// One entry of the near-miss series for a fault code. +/// One entry of the near-miss series of one fault record. /// /// A near miss is a FAILED report that moved the debounce counter WITHOUT the fault /// ending up CONFIRMED - the fault nearly happened. PASSED reports move the counter @@ -157,21 +187,21 @@ std::string rosbag_recording_id(const std::string & file_path); /// /// The series is append-only: one entry per qualifying report, never updated in place, /// and RETAINED across clear_fault, because acknowledging a fault cycle must not erase -/// the record of how often that code approached confirmation. It is bounded per fault -/// code (see set_max_near_misses_per_fault) and evicts the OLDEST entries first, so a +/// the record of how often that record approached confirmation. It is bounded per +/// record (see set_max_near_misses_per_fault) and evicts the OLDEST entries first, so a /// long-running appliance keeps the recent series rather than freezing it at boot. struct NearMissRecord { std::string fault_code; int64_t occurred_at_ns{0}; ///< Timestamp of the report that moved the counter int32_t debounce_counter{0}; ///< Counter value AFTER this report - /// Confirmation threshold the report was evaluated against. With per-entity overrides this is - /// the threshold of the REPORTING SOURCE, while the counter is shared by every source of the - /// code, so it is not by itself the distance to confirmation for the fault as a whole. + /// Confirmation threshold the report was evaluated against. The counter it belongs to is + /// this record's own, so with per-entity overrides both are the reporting source's and the + /// value is the distance to confirmation for the record. int32_t confirmation_threshold{0}; uint8_t severity{0}; ///< Severity carried by the report - std::string source_id; ///< Reporting source + std::string source_id; ///< Reporting source, which is the record's owner /// Fault status after the report was applied. Never CONFIRMED - that is what makes the report a /// near miss. It separates a counter climbing from a resting state (PREFAILED) from one walking @@ -187,6 +217,7 @@ struct NearMissRecord { /// only when its last row goes. struct RosbagFileInfo { std::string fault_code; + std::string owner; ///< Reporting source that owns the record holding this link std::string recording_id; ///< Basename of file_path; shared by every row of a burst std::string file_path; std::string format; ///< "sqlite3" or "mcap" @@ -211,11 +242,12 @@ class FaultStorage { /// @param event_type EVENT_FAILED (0) or EVENT_PASSED (1) /// @param severity Fault severity level (only used for FAILED events) /// @param description Human-readable description (only used for FAILED events) - /// @param source_id Reporting source identifier + /// @param source_id Reporting source identifier. Together with fault_code it names the + /// record this event applies to, and it creates that record when none exists. /// @param timestamp Current time for tracking /// @param config Debounce configuration to apply for this event (resolved per-entity by the node) - /// @return true if this is a new occurrence (new fault or reactivated CLEARED fault), - /// false if existing active fault was updated + /// @return true if this is a new occurrence (new record or reactivated CLEARED record), + /// false if an existing active record was updated virtual bool report_fault_event(const std::string & fault_code, uint8_t event_type, uint8_t severity, const std::string & description, const std::string & source_id, const rclcpp::Time & timestamp, const DebounceConfig & config) = 0; @@ -228,34 +260,44 @@ class FaultStorage { virtual std::vector list_faults(bool filter_by_severity, uint8_t severity, const std::vector & statuses) const = 0; - /// Get a single fault by fault_code - /// @param fault_code The fault code to look up + /// Get a single fault record + /// @param id The record to look up /// @return The fault if found, nullopt otherwise - virtual std::optional get_fault(const std::string & fault_code) const = 0; + virtual std::optional get_fault(const FaultId & id) const = 0; + + /// Every record carrying @p fault_code, one per owning reporting source. + /// + /// The resolution step behind an unscoped service call: with exactly one record the + /// call applies to it, with several it is ambiguous and the owners are the answer. + /// Returns them ordered by owner, so an ambiguity message is stable. + /// @param fault_code The fault code to look up + /// @return Every record carrying the code, empty when none does + virtual std::vector get_faults_by_code(const std::string & fault_code) const = 0; - /// Clear a fault by fault_code (manual acknowledgment). Drops the fault's per-topic snapshots; + /// Clear one fault record (manual acknowledgment). Drops that record's per-topic snapshots; /// the freeze-frame and the near-miss series are RETAINED, because they outlive a single fault - /// cycle and cannot be reconstructed afterwards. - /// @param fault_code The fault code to clear - /// @return true if fault was found and cleared, false if not found - virtual bool clear_fault(const std::string & fault_code) = 0; + /// cycle and cannot be reconstructed afterwards. Another owner's record of the same code is + /// untouched, snapshots included. + /// @param id The record to clear + /// @return true if the record was found and cleared, false if not found + virtual bool clear_fault(const FaultId & id) = 0; - /// Get total number of stored faults + /// Get total number of stored fault records virtual size_t size() const = 0; - /// Check if a fault exists - virtual bool contains(const std::string & fault_code) const = 0; + /// Check if a fault record exists + virtual bool contains(const FaultId & id) const = 0; - /// Check and confirm PREFAILED faults that have been pending too long (time-based confirmation) + /// Check and confirm PREFAILED records that have been pending too long (time-based confirmation) /// @param current_time Current timestamp for age calculation - /// @return Fault codes that were confirmed by this call (so the caller can audit each). - virtual std::vector check_time_based_confirmation(const rclcpp::Time & current_time) = 0; + /// @return The records that were confirmed by this call (so the caller can audit each). + virtual std::vector check_time_based_confirmation(const rclcpp::Time & current_time) = 0; - /// Set maximum snapshots per fault code (0 = unlimited) + /// Set maximum snapshots per fault record (0 = unlimited) virtual void set_max_snapshots_per_fault(size_t /*max_count*/) { } - /// Cap on RECORDINGS retained per fault code (0 = unlimited). + /// Cap on RECORDINGS retained per fault record (0 = unlimited). /// /// Enforced inside store_rosbag_file(s), atomically with the insert: past the cap /// the fault's OLDEST recordings lose their row, and a bag whose last referencing @@ -303,16 +345,15 @@ class FaultStorage { /// @param snapshot The snapshot data to store virtual void store_snapshot(const SnapshotData & snapshot) = 0; - /// Get snapshots for a fault, NEWEST capture set first. + /// Get snapshots of one fault record, NEWEST capture set first. /// /// Ordered by capture_id descending, then captured_at_ns descending, in every /// backend. Readers fold rows into a per-topic map, so insertion order would let /// an older capture's values win on one backend and not the other. - /// @param fault_code The fault code to get snapshots for + /// @param id The record to get snapshots for /// @param topic_filter Optional topic filter (empty = all topics) - /// @return Vector of snapshots for the fault - virtual std::vector get_snapshots(const std::string & fault_code, - const std::string & topic_filter = "") const = 0; + /// @return Vector of snapshots of the record + virtual std::vector get_snapshots(const FaultId & id, const std::string & topic_filter = "") const = 0; /// Highest capture_id any stored snapshot holds, across every fault (0 when none). /// @@ -324,22 +365,23 @@ class FaultStorage { return 0; } - /// Store the compact freeze-frame captured for a fault (JSON dict of topic values). - /// Keyed by fault_code: a later capture for the same code replaces the frame. The frame + /// Store the compact freeze-frame captured for a fault record (JSON dict of topic values). + /// Keyed by the record identity carried on @p frame: a later capture for the same record + /// replaces the frame, and another owner's record of the same code keeps its own. The frame /// is retained across clear_fault so the confirmed-state record survives acknowledgement. - /// Storage is bounded by the number of distinct fault codes (one row per code, replaced - /// in place); rows are never evicted. Faults themselves are never deleted (clear_fault - /// only flips status), so there is currently no delete hook to tie eviction to. - /// @param frame The freeze-frame to store + /// Storage is bounded by the number of distinct records (one row per record, replaced in + /// place); rows are never evicted. Records themselves are never deleted (clear_fault only + /// flips status), so there is currently no delete hook to tie eviction to. + /// @param frame The freeze-frame to store, carrying fault_code and owner virtual void store_freeze_frame(const FreezeFrameData & frame) = 0; - /// Get the freeze-frame captured for a fault, if any. - /// @param fault_code The fault code to look up - /// @return The freeze-frame if one was captured, nullopt otherwise (including fault - /// codes with no capture configured, which never get a row) - virtual std::optional get_freeze_frame(const std::string & fault_code) const = 0; + /// Get the freeze-frame captured for a fault record, if any. + /// @param id The record to look up + /// @return The freeze-frame if one was captured, nullopt otherwise (including records + /// with no capture configured, which never get a row) + virtual std::optional get_freeze_frame(const FaultId & id) const = 0; - /// Set the maximum number of near-miss entries retained per fault code. + /// Set the maximum number of near-miss entries retained per fault record. /// /// Entries beyond the bound are evicted oldest-first, including entries already stored when the /// bound is applied. 0 means unlimited, and so does any bound larger than the storage backend @@ -352,17 +394,18 @@ class FaultStorage { return 0; } - /// Get the near-miss series for a fault code, oldest entry first. - /// The series survives clear_fault; an unknown or never-near-missed code returns empty. - /// @param fault_code The fault code to look up + /// Get the near-miss series of one fault record, oldest entry first. + /// The series survives clear_fault; an unknown or never-near-missed record returns empty. + /// @param id The record to look up /// @return The retained near-miss entries in chronological order - virtual std::vector get_near_misses(const std::string & fault_code) const = 0; + virtual std::vector get_near_misses(const FaultId & id) const = 0; - /// Store rosbag file metadata for a fault - /// @param info The rosbag file info to store (replaces any existing entry for fault_code) + /// Store rosbag file metadata for a fault record + /// @param info The rosbag file info to store, carrying fault_code and owner (replaces any + /// existing link between that record and the same file_path) virtual void store_rosbag_file(const RosbagFileInfo & info) = 0; - /// Store one row per fault of a burst that shares a recording. + /// Store one row per record of a burst that shares a recording. /// /// Implementations MUST be all-or-nothing. The caller treats a throw as "no row /// was written" and discards the recording, so a batch that stored some rows and @@ -381,20 +424,20 @@ class FaultStorage { } } - /// The MOST RECENT recording of a fault, or nullopt. + /// The MOST RECENT recording of a fault record, or nullopt. /// - /// A fault can hold several recordings, so "the" recording is a choice: newest + /// A record can hold several recordings, so "the" recording is a choice: newest /// wins, because a black box is evidence about the machine you are about to /// inspect. Implementations must order deterministically - an unordered pick /// serves an arbitrary recording, which no test catches reliably. - /// @param fault_code The fault code to get rosbag for + /// @param id The record to get rosbag for /// @return Rosbag file info if exists, nullopt otherwise - virtual std::optional get_rosbag_file(const std::string & fault_code) const = 0; + virtual std::optional get_rosbag_file(const FaultId & id) const = 0; - /// Every recording of a fault, newest first. - virtual std::vector get_rosbag_files(const std::string & fault_code) const = 0; + /// Every recording of a fault record, newest first. + virtual std::vector get_rosbag_files(const FaultId & id) const = 0; - /// Every row of one recording - one per fault the recording covers. Backs the + /// Every row of one recording - one per record the recording covers. Backs the /// bulk-data download and the entity authorization scope check, both of which /// start from a recording id and need the faults behind it. virtual std::vector get_rosbag_files_by_recording(const std::string & recording_id) const = 0; @@ -405,23 +448,23 @@ class FaultStorage { /// @return number of rows removed virtual size_t delete_rosbag_recording(const std::string & recording_id) = 0; - /// Delete rosbag file record and the actual file for a fault. Faults from one - /// burst can share a recording; the file is unlinked only with the last record + /// Delete rosbag rows and the actual file for one fault record. Records from one + /// burst can share a recording; the file is unlinked only with the last row /// that references it. - /// @param fault_code The fault code to delete rosbag for - /// @return true if record was deleted, false if not found - virtual bool delete_rosbag_file(const std::string & fault_code) = 0; + /// @param id The record to delete rosbags for + /// @return true if at least one row was deleted, false if none was found + virtual bool delete_rosbag_file(const FaultId & id) = 0; - /// Delete the records of several faults (typically the whole burst behind one + /// Delete the rows of several records (typically the whole burst behind one /// recording). Backends with real transactions (SQLite) remove the rows /// atomically and unlink the file only after the commit, so a crash mid-delete /// never leaves a row pointing at a removed bag. Default: plain loop. - /// @param fault_codes The fault codes to delete rosbag records for - /// @return Number of records actually deleted - virtual size_t delete_rosbag_files(const std::vector & fault_codes) { + /// @param ids The records to delete rosbag rows for + /// @return Number of records for which rows were actually deleted + virtual size_t delete_rosbag_files(const std::vector & ids) { size_t deleted = 0; - for (const auto & code : fault_codes) { - if (delete_rosbag_file(code)) { + for (const auto & id : ids) { + if (delete_rosbag_file(id)) { ++deleted; } } @@ -437,20 +480,22 @@ class FaultStorage { /// @return Vector of rosbag file info virtual std::vector get_all_rosbag_files() const = 0; - /// Get rosbags for all faults associated with an entity + /// Get rosbags of every record owned by an entity /// @param entity_fqn The entity's fully qualified name to filter by - /// @return Vector of rosbag file info for faults reported by this entity + /// @return Vector of rosbag file info for records this entity owns virtual std::vector list_rosbags_for_entity(const std::string & entity_fqn) const = 0; - /// Get all stored faults regardless of status (for filtering) - /// @return Vector of all faults in storage + /// Get all stored fault records regardless of status (for filtering) + /// @return Vector of all records in storage virtual std::vector get_all_faults() const = 0; - /// One-time startup cleanup: reclassify HEALED faults as CLEARED. Called when healing is disabled, - /// so a HEALED row left by a previous (healing-enabled) run does not behave inconsistently under - /// the latch. Default is a no-op (in-memory storage starts empty). - /// @return fault codes of the reclassified faults, so the caller can audit each transition - virtual std::vector reclassify_healed_as_cleared() { + /// One-time startup cleanup: reclassify HEALED records as CLEARED. Called when healing is + /// disabled, so a HEALED row left by a previous (healing-enabled) run does not behave + /// inconsistently under the latch. Default is a no-op (in-memory storage starts empty). + /// @return the reclassified records, so the caller can audit each transition. Identities, + /// not codes: one code can hold a HEALED record for one owner and an untouched + /// CONFIRMED one for another, and auditing by code would claim both moved. + virtual std::vector reclassify_healed_as_cleared() { return {}; } @@ -477,15 +522,17 @@ class InMemoryFaultStorage : public FaultStorage { std::vector list_faults(bool filter_by_severity, uint8_t severity, const std::vector & statuses) const override; - std::optional get_fault(const std::string & fault_code) const override; + std::optional get_fault(const FaultId & id) const override; + + std::vector get_faults_by_code(const std::string & fault_code) const override; - bool clear_fault(const std::string & fault_code) override; + bool clear_fault(const FaultId & id) override; size_t size() const override; - bool contains(const std::string & fault_code) const override; + bool contains(const FaultId & id) const override; - std::vector check_time_based_confirmation(const rclcpp::Time & current_time) override; + std::vector check_time_based_confirmation(const rclcpp::Time & current_time) override; void set_max_snapshots_per_fault(size_t max_count) override; void set_retain_snapshots_on_clear(bool retain) override; @@ -493,32 +540,31 @@ class InMemoryFaultStorage : public FaultStorage { void store_snapshot(const SnapshotData & snapshot) override; void store_snapshots(const std::vector & snapshots) override; - std::vector get_snapshots(const std::string & fault_code, - const std::string & topic_filter = "") const override; + std::vector get_snapshots(const FaultId & id, const std::string & topic_filter = "") const override; int64_t get_max_capture_id() const override; void store_freeze_frame(const FreezeFrameData & frame) override; - std::optional get_freeze_frame(const std::string & fault_code) const override; + std::optional get_freeze_frame(const FaultId & id) const override; void set_max_rosbags_per_fault(size_t max_count) override; size_t set_max_near_misses_per_fault(size_t max_count) override; - std::vector get_near_misses(const std::string & fault_code) const override; + std::vector get_near_misses(const FaultId & id) const override; void store_rosbag_file(const RosbagFileInfo & info) override; /// All-or-nothing, as the base class requires: the batch is built beside the live /// map and swapped in, so a throw leaves the store exactly as it was. void store_rosbag_files(const std::vector & infos) override; - std::optional get_rosbag_file(const std::string & fault_code) const override; - std::vector get_rosbag_files(const std::string & fault_code) const override; + std::optional get_rosbag_file(const FaultId & id) const override; + std::vector get_rosbag_files(const FaultId & id) const override; std::vector get_rosbag_files_by_recording(const std::string & recording_id) const override; - bool delete_rosbag_file(const std::string & fault_code) override; + bool delete_rosbag_file(const FaultId & id) override; size_t delete_rosbag_recording(const std::string & recording_id) override; size_t get_total_rosbag_storage_bytes() const override; std::vector get_all_rosbag_files() const override; std::vector list_rosbags_for_entity(const std::string & entity_fqn) const override; std::vector get_all_faults() const override; - std::vector reclassify_healed_as_cleared() override; + std::vector reclassify_healed_as_cleared() override; private: /// Update fault status based on debounce counter and given config @@ -537,12 +583,12 @@ class InMemoryFaultStorage : public FaultStorage { bool path_referenced(const std::string & file_path) const; mutable std::mutex mutex_; - std::map faults_; + std::map faults_; std::vector snapshots_; - std::map freeze_frames_; ///< fault_code -> freeze-frame (retained across clear) + std::map freeze_frames_; ///< record -> freeze-frame (retained across clear) /// One entry per LINK, mirroring the flat SQLite table rather than a map keyed by - /// fault code - a fault holds several recordings now, and a recording several - /// faults. + /// the record - a record holds several recordings now, and a recording several + /// records. /// /// `seq` is the in-memory twin of SQLite's autoincrement id and is load-bearing, /// not decoration: a burst stamps ONE created_at_ns across every row it writes, so @@ -559,7 +605,7 @@ class InMemoryFaultStorage : public FaultStorage { /// A backend constructed directly (tests, embedders) therefore behaves exactly as /// it always did until someone opts into a history. 0 = unlimited. size_t max_rosbags_per_fault_{1}; - std::map> near_misses_; ///< fault_code -> series (retained across clear) + std::map> near_misses_; ///< record -> series (retained across clear) DebounceConfig config_; size_t max_snapshots_per_fault_{0}; ///< 0 = unlimited bool retain_snapshots_on_clear_{false}; diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/rosbag_capture.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/rosbag_capture.hpp index 10e003606..b1c268332 100644 --- a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/rosbag_capture.hpp +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/rosbag_capture.hpp @@ -100,22 +100,24 @@ class RosbagCapture { /// Check if ring buffer is currently running bool is_running() const; - /// Called when a fault enters PREFAILED state (for lazy_start mode) - /// @param fault_code The fault code that entered PREFAILED - void on_fault_prefailed(const std::string & fault_code); - - /// Called when a fault is confirmed - flushes buffer to bag file. A fault that - /// confirms while the previous fault's post-roll is still running is attached - /// to that recording (same burst, same window) rather than losing its bag. - /// @param fault_code The fault code that was confirmed - void on_fault_confirmed(const std::string & fault_code); - - /// Called when a fault is cleared - deletes its bag record if auto_cleanup. - /// A shared bag survives until its last referencing fault clears; a fault + /// Called when a fault record enters PREFAILED state (for lazy_start mode) + /// @param id The record that entered PREFAILED + void on_fault_prefailed(const FaultId & id); + + /// Called when a fault record is confirmed - flushes buffer to bag file. A record + /// that confirms while the previous one's post-roll is still running is attached + /// to that recording (same burst, same window) rather than losing its bag. Two + /// owners confirming one code in a burst attach to the one in-flight recording, + /// and each gets its own rosbag_files row. + /// @param id The record that was confirmed + void on_fault_confirmed(const FaultId & id); + + /// Called when a fault record is cleared - deletes its bag record if auto_cleanup. + /// A shared bag survives until its last referencing record clears; a record /// cleared during its burst's post-roll is dropped from the in-flight - /// recording state and never gets a record. - /// @param fault_code The fault code that was cleared - void on_fault_cleared(const std::string & fault_code); + /// recording state and never gets a row. + /// @param id The record that was cleared + void on_fault_cleared(const FaultId & id); /// Get current configuration const RosbagConfig & config() const { @@ -216,10 +218,10 @@ class RosbagCapture { /// (falls back to SensorDataQoS when no publisher is known or qos_match is off) rclcpp::QoS resolve_topic_qos(const std::string & topic) const; - /// Compute the entity topic set for a fault (the faulting source node's + /// Compute the entity topic set for a fault record (its owning source node's /// pub/sub topics + /tf, intersected with the subscribed set). Empty set = /// scope unresolved. Never throws; failures degrade to an empty set. - std::set compute_entity_topics(const std::string & fault_code); + std::set compute_entity_topics(const FaultId & id); /// In "entity" mode, compute the set of topics to write for a confirmed fault /// (the faulting source node's pub/sub topics + /tf). Empty set = write all. @@ -228,7 +230,7 @@ class RosbagCapture { /// In "entity" mode, union an attached fault's entity topics into the active /// capture filter so its data reaches the shared bag from the attach onwards /// (empty resolution widens to all topics). Caller holds post_fault_timer_mutex_. - void widen_capture_filter_for(const std::string & fault_code, const std::set & topics); + void widen_capture_filter_for(const FaultId & id, const std::set & topics); /// Whether a topic should be written to the bag given the active entity filter bool should_capture_topic(const std::string & topic) const; @@ -330,9 +332,9 @@ class RosbagCapture { /// Cheap, and checked before the entity scope is resolved: a level-triggered /// reporter re-confirming the same fault would otherwise pay for a fault-store read /// and a full graph enumeration on every repeat, none of which it can use. - bool is_current_recording_primary(const std::string & fault_code) const; + bool is_current_recording_primary(const FaultId & id) const; - bool attach_to_active_recording(const std::string & fault_code, const std::set & entity_topics); + bool attach_to_active_recording(const FaultId & id, const std::set & entity_topics); /// Try to subscribe to a single topic /// @param topic The topic to subscribe to @@ -383,16 +385,17 @@ class RosbagCapture { /// Running state std::atomic running_{false}; - /// Upper bound on how many extra faults one recording is registered for. + /// Upper bound on how many extra records one recording is registered for. static constexpr size_t kMaxAttachedFaults = 32; /// Post-fault recording state - std::string current_fault_code_; + FaultId current_fault_id_; std::string current_bag_path_; - /// Faults confirmed while the post-roll was already running. They share the + /// Records confirmed while the post-roll was already running. They share the /// recording window (one root cause, one burst), so each gets a metadata row - /// pointing at the same bag when it finalises. - std::set attached_fault_codes_; + /// pointing at the same bag when it finalises. Two owners of one code are two + /// entries here and two rows on finalise. + std::set attached_fault_ids_; /// Protects post_fault_timer_, the recording_post_fault_ transitions and the /// state above against concurrent access from on_fault_confirmed() (capture-pool /// thread) and post_fault_timer_callback() / stop() (executor thread). The diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/snapshot_capture.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/snapshot_capture.hpp index 9a33dd2b2..25c7c4889 100644 --- a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/snapshot_capture.hpp +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/snapshot_capture.hpp @@ -164,17 +164,22 @@ class SnapshotCapture { SnapshotCapture(SnapshotCapture &&) = delete; SnapshotCapture & operator=(SnapshotCapture &&) = delete; - /// Capture snapshots for a fault that was just confirmed + /// Capture snapshots for a fault record that was just confirmed + /// + /// Topic resolution stays keyed by the fault CODE (fault_specific and patterns are + /// config about what a code means), while everything written is keyed by the RECORD: + /// the rows, the freeze frame and the mid-capture acknowledgement check. Two owners + /// confirming one code therefore capture the same topics into two separate frames. /// /// If the fault code resolves to no capture set (not in fault_specific, no pattern /// match, no default_topics), the entity-default fallback (when enabled) captures - /// the reporting source node's own published topics instead. A code explicitly + /// the owning source node's own published topics instead. A code explicitly /// listed in fault_specific or matched by a pattern never falls through - a /// present-but-empty topic list is a per-fault opt-out. Only when nothing /// resolves does capture return early: no freeze_frames row is written /// (no empty {} row) and FaultStorage::get_freeze_frame() returns nullopt for it. - /// @param fault_code The fault code that was confirmed - void capture(const std::string & fault_code); + /// @param id The fault record that was confirmed + void capture(const FaultId & id); /// Get current configuration const SnapshotConfig & config() const { @@ -195,12 +200,12 @@ class SnapshotCapture { /// entity-default capture. std::vector resolve_topics(const std::string & fault_code, bool & explicit_match) const; - /// Entity-default fallback: topics published by the fault's reporting source - /// node(s), excluding per-node noise (/rosout, /parameter_events). Non-FQN - /// sources (bare plugin entity ids - the gateway covers those) are skipped, - /// and empty is returned when no source resolves to a live node. Never + /// Entity-default fallback: topics published by the record's owning source + /// node, excluding per-node noise (/rosout, /parameter_events). A non-FQN + /// owner (a bare plugin entity id - the gateway covers those) is skipped, + /// and empty is returned when the owner resolves to no live node. Never /// throws; any failure degrades to empty. - std::vector resolve_entity_topics(const std::string & fault_code) const; + std::vector resolve_entity_topics(const FaultId & id) const; /// Capture a single topic on-demand (creates temporary subscription) /// On success also records the captured value into @p freeze_frame under the topic key. @@ -208,14 +213,14 @@ class SnapshotCapture { /// Appends to @p rows rather than storing: a capture is persisted as one set, /// so the per-fault cap can drop a whole old capture instead of truncating this /// one topic by topic. - bool capture_topic_on_demand(const std::string & fault_code, const std::string & topic, nlohmann::json & freeze_frame, + bool capture_topic_on_demand(const FaultId & id, const std::string & topic, nlohmann::json & freeze_frame, std::vector & rows); /// Capture a topic from background cache /// On success also records the cached value into @p freeze_frame under the topic key. /// @return true if data was available in cache - bool capture_topic_from_cache(const std::string & fault_code, const std::string & topic, - nlohmann::json & freeze_frame, std::vector & rows); + bool capture_topic_from_cache(const FaultId & id, const std::string & topic, nlohmann::json & freeze_frame, + std::vector & rows); /// Initialize background subscriptions for all configured topics void init_background_subscriptions(); diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/sqlite_fault_storage.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/sqlite_fault_storage.hpp index 3565febf4..05bc5741a 100644 --- a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/sqlite_fault_storage.hpp +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/sqlite_fault_storage.hpp @@ -53,15 +53,17 @@ class SqliteFaultStorage : public FaultStorage { std::vector list_faults(bool filter_by_severity, uint8_t severity, const std::vector & statuses) const override; - std::optional get_fault(const std::string & fault_code) const override; + std::optional get_fault(const FaultId & id) const override; - bool clear_fault(const std::string & fault_code) override; + std::vector get_faults_by_code(const std::string & fault_code) const override; + + bool clear_fault(const FaultId & id) override; size_t size() const override; - bool contains(const std::string & fault_code) const override; + bool contains(const FaultId & id) const override; - std::vector check_time_based_confirmation(const rclcpp::Time & current_time) override; + std::vector check_time_based_confirmation(const rclcpp::Time & current_time) override; void set_max_snapshots_per_fault(size_t max_count) override; void set_retain_snapshots_on_clear(bool retain) override; @@ -71,29 +73,28 @@ class SqliteFaultStorage : public FaultStorage { void store_snapshot(const SnapshotData & snapshot) override; void store_snapshots(const std::vector & snapshots) override; - std::vector get_snapshots(const std::string & fault_code, - const std::string & topic_filter = "") const override; + std::vector get_snapshots(const FaultId & id, const std::string & topic_filter = "") const override; int64_t get_max_capture_id() const override; void store_freeze_frame(const FreezeFrameData & frame) override; - std::optional get_freeze_frame(const std::string & fault_code) const override; + std::optional get_freeze_frame(const FaultId & id) const override; size_t set_max_near_misses_per_fault(size_t max_count) override; - std::vector get_near_misses(const std::string & fault_code) const override; + std::vector get_near_misses(const FaultId & id) const override; void store_rosbag_file(const RosbagFileInfo & info) override; void store_rosbag_files(const std::vector & infos) override; - std::optional get_rosbag_file(const std::string & fault_code) const override; - std::vector get_rosbag_files(const std::string & fault_code) const override; + std::optional get_rosbag_file(const FaultId & id) const override; + std::vector get_rosbag_files(const FaultId & id) const override; std::vector get_rosbag_files_by_recording(const std::string & recording_id) const override; - bool delete_rosbag_file(const std::string & fault_code) override; + bool delete_rosbag_file(const FaultId & id) override; size_t delete_rosbag_recording(const std::string & recording_id) override; - size_t delete_rosbag_files(const std::vector & fault_codes) override; + size_t delete_rosbag_files(const std::vector & ids) override; size_t get_total_rosbag_storage_bytes() const override; std::vector get_all_rosbag_files() const override; std::vector list_rosbags_for_entity(const std::string & entity_fqn) const override; std::vector get_all_faults() const override; - std::vector reclassify_healed_as_cleared() override; + std::vector reclassify_healed_as_cleared() override; /// Get the database path const std::string & db_path() const { @@ -120,16 +121,31 @@ class SqliteFaultStorage : public FaultStorage { /// C++ through rosbag_recording_id() so the basename rule has one implementation. void migrate_rosbag_files_add_recording_id(); - /// Whether a fault other than @p fault_code still references @p file_path. - /// One recording can back several faults of the same burst, so the bag must - /// only be unlinked once the last of them is gone. Caller holds mutex_. - - /// Whether any fault at all still references @p file_path. Caller holds mutex_. + /// Move a database whose fault identity is the bare fault_code onto (fault_code, owner). + /// + /// Detected by PRAGMA table_info(faults) lacking `owner`, which is what makes it safe to + /// re-run on every open: once the column is there the migration is done. faults and + /// freeze_frames are REBUILT (their identity is a PRIMARY KEY, which ALTER TABLE cannot + /// change and CREATE TABLE IF NOT EXISTS would silently skip), while snapshots and + /// rosbag_files only gain a column. Every backfill takes the owner from the legacy row's + /// first reporting source, empty when it had none: that is the only per-source fact the + /// old schema holds, so a row with several sources folds onto its first one rather than + /// inventing per-source counters, statuses and timestamps that were never recorded. + void migrate_faults_add_owner(); + + /// Whether @p table already carries the `owner` column. The four tables are created + /// by four independent CREATE TABLE IF NOT EXISTS statements, so a database can hold + /// a legacy one next to one this release just created: each is probed on its own. + bool table_has_owner(const char * table) const; + + /// Whether any row at all still references @p file_path. One recording can back + /// several records of the same burst, so the bag must only be unlinked once the + /// last of them is gone. Caller holds mutex_. bool path_referenced(const std::string & file_path) const; /// store_rosbag_file body without taking mutex_. Caller holds mutex_ and /// unlinks the returned replaced-bag path once the row change is durable. - /// @return file_paths whose last row for this fault the per-fault cap evicted. + /// @return file_paths whose last row for this record the per-record cap evicted. /// The caller unlinks each only after the commit, and only if path_referenced() /// still says nobody holds it. std::vector store_rosbag_file_locked(const RosbagFileInfo & info); @@ -147,7 +163,7 @@ class SqliteFaultStorage : public FaultStorage { /// @param debounce_counter Counter value after the report /// @param config Debounce config the report was evaluated against /// @param severity Severity carried by the report - /// @param source_id Reporting source + /// @param source_id Reporting source, which is the record's owner and scopes the series /// @param resulting_status Fault status after the report was applied void record_near_miss_locked(const std::string & fault_code, int64_t occurred_at_ns, int32_t debounce_counter, const DebounceConfig & config, uint8_t severity, const std::string & source_id, @@ -156,10 +172,9 @@ class SqliteFaultStorage : public FaultStorage { /// Run a plain SQL statement or throw with the SQLite error. Caller holds mutex_. void exec_or_throw(const char * sql); - /// Parse JSON array string to vector of strings - static std::vector parse_json_array(const std::string & json_str); - - /// Serialize vector of strings to JSON array string + /// Serialize vector of strings to JSON array string. The faults table keeps + /// reporting_sources as the serialized form of owner, so a reader of the raw table + /// still sees the field it always saw. static std::string serialize_json_array(const std::vector & vec); std::string db_path_; diff --git a/src/ros2_medkit_fault_manager/src/capture_thread_pool.cpp b/src/ros2_medkit_fault_manager/src/capture_thread_pool.cpp index 3efe8b8fa..8bc8c5618 100644 --- a/src/ros2_medkit_fault_manager/src/capture_thread_pool.cpp +++ b/src/ros2_medkit_fault_manager/src/capture_thread_pool.cpp @@ -22,7 +22,7 @@ namespace ros2_medkit_fault_manager { CaptureThreadPool::CaptureThreadPool(std::size_t pool_size, std::size_t queue_depth, QueueFullPolicy full_policy, - rclcpp::Logger logger, std::function capture_fn) + rclcpp::Logger logger, std::function capture_fn) : queue_depth_(queue_depth == 0 ? 1 : queue_depth) , full_policy_(full_policy) , logger_(std::move(logger)) @@ -44,13 +44,13 @@ CaptureThreadPool::~CaptureThreadPool() { shutdown(); } -EnqueueOutcome CaptureThreadPool::enqueue(const std::string & fault_code) { +EnqueueOutcome CaptureThreadPool::enqueue(const FaultId & id) { std::lock_guard lock(queue_mutex_); if (stop_) { return {EnqueueResult::kRejectedShuttingDown, std::nullopt}; } if (queue_.size() < queue_depth_) { - queue_.push_back(fault_code); + queue_.push_back(id); cv_.notify_one(); return {EnqueueResult::kAccepted, std::nullopt}; } @@ -60,9 +60,9 @@ EnqueueOutcome CaptureThreadPool::enqueue(const std::string & fault_code) { return {EnqueueResult::kDroppedNewest, std::nullopt}; } // kDropOldest: evict the oldest pending job. - std::string evicted = std::move(queue_.front()); + FaultId evicted = std::move(queue_.front()); queue_.pop_front(); - queue_.push_back(fault_code); + queue_.push_back(id); dropped_captures_.fetch_add(1, std::memory_order_relaxed); cv_.notify_one(); return {EnqueueResult::kEvictedOldest, std::move(evicted)}; @@ -103,7 +103,7 @@ std::size_t CaptureThreadPool::pending_size() const { void CaptureThreadPool::worker_loop() { for (;;) { - std::string job; + FaultId job; { std::unique_lock lock(queue_mutex_); cv_.wait(lock, [this] { @@ -120,9 +120,9 @@ void CaptureThreadPool::worker_loop() { capture_fn_(job); } } catch (const std::exception & e) { - RCLCPP_ERROR(logger_, "Capture job for '%s' threw: %s", job.c_str(), e.what()); + RCLCPP_ERROR(logger_, "Capture job for '%s' threw: %s", job.fault_code.c_str(), e.what()); } catch (...) { - RCLCPP_ERROR(logger_, "Capture job for '%s' threw unknown exception", job.c_str()); + RCLCPP_ERROR(logger_, "Capture job for '%s' threw unknown exception", job.fault_code.c_str()); } } } diff --git a/src/ros2_medkit_fault_manager/src/correlation/correlation_engine.cpp b/src/ros2_medkit_fault_manager/src/correlation/correlation_engine.cpp index c7cd12728..5f3a905fd 100644 --- a/src/ros2_medkit_fault_manager/src/correlation/correlation_engine.cpp +++ b/src/ros2_medkit_fault_manager/src/correlation/correlation_engine.cpp @@ -24,10 +24,13 @@ CorrelationEngine::CorrelationEngine(const CorrelationConfig & config) : config_(config), matcher_(std::make_unique(config.patterns)) { } -ProcessFaultResult CorrelationEngine::process_fault(const std::string & fault_code, const std::string & severity, +ProcessFaultResult CorrelationEngine::process_fault(const FaultId & id, const std::string & severity, std::chrono::steady_clock::time_point timestamp) { std::lock_guard lock(mutex_); + // Rules match on the CODE; every relation formed below is between records of the + // same owner. + const std::string & fault_code = id.fault_code; ProcessFaultResult result; // First, clean up expired entries @@ -42,8 +45,8 @@ ProcessFaultResult CorrelationEngine::process_fault(const std::string & fault_co }), pending_root_causes_.end()); - // Check if this fault is a symptom of an existing root cause - auto symptom_result = try_as_symptom(fault_code, timestamp); + // Check if this record is a symptom of an existing root cause of the same owner + auto symptom_result = try_as_symptom(id, timestamp); if (symptom_result) { return *symptom_result; } @@ -59,14 +62,14 @@ ProcessFaultResult CorrelationEngine::process_fault(const std::string & fault_co if (rule.id == *root_cause_rule) { // Add to pending root causes PendingRootCause prc; - prc.fault_code = fault_code; + prc.fault_id = id; prc.rule_id = rule.id; prc.timestamp = timestamp; prc.window_ms = rule.window_ms; pending_root_causes_.push_back(prc); // Initialize symptom list - root_to_symptoms_[fault_code] = {}; + root_to_symptoms_[id] = {}; break; } } @@ -74,8 +77,8 @@ ProcessFaultResult CorrelationEngine::process_fault(const std::string & fault_co return result; } - // Check if this fault matches an auto-cluster rule - auto cluster_result = try_auto_cluster(fault_code, severity, timestamp); + // Check if this record matches an auto-cluster rule + auto cluster_result = try_auto_cluster(id, severity, timestamp); if (cluster_result) { return *cluster_result; } @@ -84,20 +87,22 @@ ProcessFaultResult CorrelationEngine::process_fault(const std::string & fault_co return result; } -ProcessClearResult CorrelationEngine::process_clear(const std::string & fault_code) { +ProcessClearResult CorrelationEngine::process_clear(const FaultId & id) { std::lock_guard lock(mutex_); + const std::string & fault_code = id.fault_code; ProcessClearResult result; - // Check if this is a root cause with symptoms - auto it = root_to_symptoms_.find(fault_code); + // Check if this record is a root cause with symptoms. The lookup is by record, so + // clearing owner X's root cause never reaches owner Y's symptoms. + auto it = root_to_symptoms_.find(id); if (it != root_to_symptoms_.end()) { // Find the rule to check auto_clear_with_root for (const auto & prc : pending_root_causes_) { - if (prc.fault_code == fault_code) { + if (prc.fault_id == id) { for (const auto & rule : config_.rules) { if (rule.id == prc.rule_id && rule.auto_clear_with_root) { - result.auto_cleared_codes = it->second; + result.auto_cleared_symptoms = it->second; break; } } @@ -106,21 +111,21 @@ ProcessClearResult CorrelationEngine::process_clear(const std::string & fault_co } // Also check finalized root causes (not in pending anymore) - if (result.auto_cleared_codes.empty()) { + if (result.auto_cleared_symptoms.empty()) { for (const auto & rule : config_.rules) { if (rule.mode == CorrelationMode::HIERARCHICAL && rule.auto_clear_with_root) { // Check if fault matches this rule's root cause if (matcher_->matches_any(fault_code, rule.root_cause_codes)) { - result.auto_cleared_codes = it->second; + result.auto_cleared_symptoms = it->second; break; } } } } - // Clean up muted faults - for (const auto & symptom_code : it->second) { - muted_faults_.erase(symptom_code); + // Clean up muted records + for (const auto & symptom_id : it->second) { + muted_faults_.erase(symptom_id); } // Remove from root_to_symptoms @@ -129,13 +134,13 @@ ProcessClearResult CorrelationEngine::process_clear(const std::string & fault_co // Remove from pending root causes pending_root_causes_.erase(std::remove_if(pending_root_causes_.begin(), pending_root_causes_.end(), - [&fault_code](const PendingRootCause & prc) { - return prc.fault_code == fault_code; + [&id](const PendingRootCause & prc) { + return prc.fault_id == id; }), pending_root_causes_.end()); - // Check if this fault is part of a cluster - auto cluster_it = fault_to_cluster_.find(fault_code); + // Check if this record is part of a cluster + auto cluster_it = fault_to_cluster_.find(id); if (cluster_it != fault_to_cluster_.end()) { const std::string cluster_id = cluster_it->second; @@ -159,7 +164,7 @@ ProcessClearResult CorrelationEngine::process_clear(const std::string & fault_co // Reassign representative if the cleared fault was the representative if (pending_cluster.representative_code == fault_code) { for (const auto & rule : config_.rules) { - if (rule.id == pending_it->first) { + if (rule.id == pending_it->first.first) { switch (rule.representative) { case Representative::FIRST: { auto & sevs = pending_it->second.fault_severities; @@ -214,7 +219,7 @@ ProcessClearResult CorrelationEngine::process_clear(const std::string & fault_co active_clusters_.erase(active_it); } else if (active_cluster.representative_code == fault_code) { // Sync representative from pending cluster (already updated above) - for (const auto & [rule_id, pending] : pending_clusters_) { + for (const auto & [pending_key, pending] : pending_clusters_) { if (pending.data.cluster_id == cluster_id) { active_cluster.representative_code = pending.data.representative_code; active_cluster.representative_severity = pending.data.representative_severity; @@ -227,8 +232,8 @@ ProcessClearResult CorrelationEngine::process_clear(const std::string & fault_co fault_to_cluster_.erase(cluster_it); } - // Remove from muted faults if it was a symptom - muted_faults_.erase(fault_code); + // Remove from muted records if it was a symptom + muted_faults_.erase(id); return result; } @@ -239,7 +244,7 @@ std::vector CorrelationEngine::get_muted_faults() const { std::vector result; result.reserve(muted_faults_.size()); - for (const auto & [code, data] : muted_faults_) { + for (const auto & [muted_id, data] : muted_faults_) { result.push_back(data); } @@ -251,9 +256,9 @@ uint32_t CorrelationEngine::get_muted_count() const { return static_cast(muted_faults_.size()); } -bool CorrelationEngine::is_muted(const std::string & fault_code) const { +bool CorrelationEngine::is_muted(const FaultId & id) const { std::lock_guard lock(mutex_); - return muted_faults_.find(fault_code) != muted_faults_.end(); + return muted_faults_.find(id) != muted_faults_.end(); } std::vector CorrelationEngine::get_clusters() const { @@ -290,26 +295,26 @@ void CorrelationEngine::cleanup_expired() { pending_root_causes_.end()); // Clean up expired pending clusters - std::vector expired_pending; - for (const auto & [rule_id, pending] : pending_clusters_) { + std::vector> expired_pending; + for (const auto & [key, pending] : pending_clusters_) { // Find rule to get window_ms for (const auto & rule : config_.rules) { - if (rule.id == rule_id) { + if (rule.id == key.first) { auto elapsed = std::chrono::duration_cast(now - pending.steady_first_at).count(); if (elapsed > static_cast(rule.window_ms)) { - expired_pending.push_back(rule_id); + expired_pending.push_back(key); } break; } } } - for (const auto & rule_id : expired_pending) { - auto it = pending_clusters_.find(rule_id); + for (const auto & key : expired_pending) { + auto it = pending_clusters_.find(key); if (it != pending_clusters_.end()) { if (active_clusters_.find(it->second.data.cluster_id) == active_clusters_.end()) { for (const auto & fault_code : it->second.data.fault_codes) { - fault_to_cluster_.erase(fault_code); + fault_to_cluster_.erase(FaultId{fault_code, it->second.owner}); } } pending_clusters_.erase(it); @@ -331,9 +336,15 @@ std::optional CorrelationEngine::try_as_root_cause(const std::strin return std::nullopt; } -std::optional CorrelationEngine::try_as_symptom(const std::string & fault_code, +std::optional CorrelationEngine::try_as_symptom(const FaultId & id, std::chrono::steady_clock::time_point timestamp) { + const std::string & fault_code = id.fault_code; for (const auto & prc : pending_root_causes_) { + // A root cause only explains its OWN reporter's faults. Without this the first + // owner to report a root-cause code would mute every other owner's symptom. + if (prc.fault_id.owner != id.owner) { + continue; + } // Find the rule for (const auto & rule : config_.rules) { if (rule.id != prc.rule_id || rule.mode != CorrelationMode::HIERARCHICAL) { @@ -367,23 +378,23 @@ std::optional CorrelationEngine::try_as_symptom(const std::s // This fault is a symptom! ProcessFaultResult result; result.should_mute = rule.mute_symptoms; - result.root_cause_code = prc.fault_code; + result.root_cause_code = prc.fault_id.fault_code; result.rule_id = rule.id; result.delay_ms = static_cast(elapsed); // Track the symptom (avoid duplicates) - auto & symptoms = root_to_symptoms_[prc.fault_code]; - if (std::find(symptoms.begin(), symptoms.end(), fault_code) == symptoms.end()) { - symptoms.push_back(fault_code); + auto & symptoms = root_to_symptoms_[prc.fault_id]; + if (std::find(symptoms.begin(), symptoms.end(), id) == symptoms.end()) { + symptoms.push_back(id); } if (rule.mute_symptoms) { MutedFaultData muted; muted.fault_code = fault_code; - muted.root_cause_code = prc.fault_code; + muted.root_cause_code = prc.fault_id.fault_code; muted.rule_id = rule.id; muted.delay_ms = result.delay_ms; - muted_faults_[fault_code] = muted; + muted_faults_[id] = muted; } return result; @@ -393,9 +404,9 @@ std::optional CorrelationEngine::try_as_symptom(const std::s return std::nullopt; } -std::optional CorrelationEngine::try_auto_cluster(const std::string & fault_code, - const std::string & severity, +std::optional CorrelationEngine::try_auto_cluster(const FaultId & id, const std::string & severity, std::chrono::steady_clock::time_point timestamp) { + const std::string & fault_code = id.fault_code; for (const auto & rule : config_.rules) { if (rule.mode != CorrelationMode::AUTO_CLUSTER) { continue; @@ -416,8 +427,11 @@ std::optional CorrelationEngine::try_auto_cluster(const std: auto now_system = std::chrono::system_clock::now(); - // Check if we have a pending cluster for this rule - auto pending_it = pending_clusters_.find(rule.id); + // Check if we have a pending cluster for this rule AND this owner. One rule forms + // one cluster per owner, so a burst on owner A never counts owner B's faults + // towards min_count and never mutes them as non-representative members. + const auto pending_key = std::make_pair(rule.id, id.owner); + auto pending_it = pending_clusters_.find(pending_key); if (pending_it != pending_clusters_.end()) { // Check if within time window using steady_clock timestamp auto elapsed = @@ -433,6 +447,7 @@ std::optional CorrelationEngine::try_auto_cluster(const std: if (pending_it == pending_clusters_.end()) { // Start new pending cluster PendingCluster pending; + pending.owner = id.owner; pending.steady_first_at = timestamp; pending.data.cluster_id = generate_cluster_id(rule.id); pending.data.rule_id = rule.id; @@ -445,8 +460,8 @@ std::optional CorrelationEngine::try_auto_cluster(const std: pending.data.first_at = now_system; pending.data.last_at = now_system; - pending_clusters_[rule.id] = pending; - fault_to_cluster_[fault_code] = pending.data.cluster_id; + pending_clusters_[pending_key] = pending; + fault_to_cluster_[id] = pending.data.cluster_id; // Not enough faults yet for a cluster ProcessFaultResult result; @@ -474,7 +489,7 @@ std::optional CorrelationEngine::try_auto_cluster(const std: cluster.fault_codes.push_back(fault_code); pending.fault_severities[fault_code] = severity; cluster.last_at = now_system; - fault_to_cluster_[fault_code] = cluster.cluster_id; + fault_to_cluster_[id] = cluster.cluster_id; // Update representative based on rule's representative selection bool update_representative = false; diff --git a/src/ros2_medkit_fault_manager/src/fault_manager_node.cpp b/src/ros2_medkit_fault_manager/src/fault_manager_node.cpp index 146f0e3b4..d3a57524f 100644 --- a/src/ros2_medkit_fault_manager/src/fault_manager_node.cpp +++ b/src/ros2_medkit_fault_manager/src/fault_manager_node.cpp @@ -275,10 +275,10 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" if (!reclassified.empty()) { if (audit_log_) { const int64_t reclassified_at_ns = get_wall_clock_time().nanoseconds(); - for (const auto & fault_code : reclassified) { - auto fault = storage_->get_fault(fault_code); + for (const auto & reclassified_id : reclassified) { + auto fault = storage_->get_fault(reclassified_id); if (fault) { - audit_transition(kTransitionCleared, *fault, "startup_reclassify", reclassified_at_ns); + audit_transition(kTransitionCleared, *fault, reclassified_id.owner, reclassified_at_ns); } } } @@ -416,13 +416,13 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" auto rosbag_mutex = std::make_shared(); capture_pool_ = std::make_unique( static_cast(capture_pool_size_), static_cast(capture_queue_depth_), - capture_queue_full_policy_, get_logger(), [snap, bag, rosbag_mutex](const std::string & fault_code) { + capture_queue_full_policy_, get_logger(), [snap, bag, rosbag_mutex](const FaultId & id) { if (snap) { - snap->capture(fault_code); + snap->capture(id); } if (bag) { std::lock_guard bag_lock(*rosbag_mutex); - bag->on_fault_confirmed(fault_code); + bag->on_fault_confirmed(id); } }); } @@ -454,10 +454,10 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" // Audit every timer-driven PREFAILED->CONFIRMED transition. Without this the // confirmations are invisible to the audit log's verify(). const int64_t confirmed_at_ns = get_wall_clock_time().nanoseconds(); - for (const auto & fault_code : confirmed) { - auto fault = storage_->get_fault(fault_code); + for (const auto & confirmed_id : confirmed) { + auto fault = storage_->get_fault(confirmed_id); if (fault) { - audit_transition(kTransitionConfirmed, *fault, "auto_confirm_timer", confirmed_at_ns); + audit_transition(kTransitionConfirmed, *fault, confirmed_id.owner, confirmed_at_ns); // A timer-driven confirmation is a confirmation: it has to reach the // event stream and the black box exactly like one raised by a report, // or subscribers see no alarm and no recording is ever made for it. @@ -467,10 +467,10 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" // the trigger subscribers see an alarm that the fault list hides. // Capture is deliberately not gated, matching the report path, where // just_confirmed is set regardless of muting. - if (!correlation_engine_ || !correlation_engine_->is_muted(fault_code)) { + if (!correlation_engine_ || !correlation_engine_->is_muted(confirmed_id)) { publish_fault_event(ros2_medkit_msgs::msg::FaultEvent::EVENT_CONFIRMED, *fault); } - capture_on_confirm(fault_code); + capture_on_confirm(confirmed_id); } } RCLCPP_INFO(get_logger(), "Auto-confirmed %zu PREFAILED fault(s) due to time threshold", confirmed.size()); @@ -687,7 +687,7 @@ void FaultManagerNode::audit_transition(const char * transition, const ros2_medk } } -void FaultManagerNode::capture_on_confirm(const std::string & fault_code) { +void FaultManagerNode::capture_on_confirm(const FaultId & id) { // Capture snapshots/rosbag when a fault confirms, via the bounded pool. // Both callers - the report handler and the auto-confirm timer - run on the // node's single-threaded executor, so confirmations are already serialized; @@ -707,8 +707,8 @@ void FaultManagerNode::capture_on_confirm(const std::string & fault_code) { if (cooldown_enabled) { const auto cooldown = std::chrono::duration(snapshot_recapture_cooldown_sec_); const auto sweep_now = std::chrono::steady_clock::now(); - // Bound the map (issue #441): a storm of distinct fault codes would otherwise - // leave one permanent entry per code. Entries older than the cooldown can never + // Bound the map (issue #441): a storm of distinct records would otherwise + // leave one permanent entry per record. Entries older than the cooldown can never // gate a capture again, so drop them while we hold the lock. for (auto it = last_capture_times_.begin(); it != last_capture_times_.end();) { if (sweep_now - it->second >= cooldown) { @@ -717,16 +717,16 @@ void FaultManagerNode::capture_on_confirm(const std::string & fault_code) { ++it; } } - auto it = last_capture_times_.find(fault_code); + auto it = last_capture_times_.find(id); if (it != last_capture_times_.end()) { on_cooldown = (sweep_now - it->second) < cooldown; } } if (on_cooldown) { - RCLCPP_DEBUG(get_logger(), "Skipping capture for '%s' - cooldown active", fault_code.c_str()); + RCLCPP_DEBUG(get_logger(), "Skipping capture for '%s' - cooldown active", id.fault_code.c_str()); } else { - const EnqueueOutcome outcome = capture_pool_->enqueue(fault_code); + const EnqueueOutcome outcome = capture_pool_->enqueue(id); const auto now = std::chrono::steady_clock::now(); // RCLCPP_WARN_THROTTLE needs a non-const Clock lvalue (Humble/Lyrical // compat); mirror rosbag_capture.cpp's local-copy pattern. Cast the @@ -736,20 +736,20 @@ void FaultManagerNode::capture_on_confirm(const std::string & fault_code) { switch (outcome.result) { case EnqueueResult::kAccepted: if (cooldown_enabled) { - last_capture_times_[fault_code] = now; + last_capture_times_[id] = now; } break; case EnqueueResult::kEvictedOldest: if (cooldown_enabled) { - last_capture_times_[fault_code] = now; - if (outcome.evicted_code) { - last_capture_times_.erase(*outcome.evicted_code); // keep evicted fault retriable + last_capture_times_[id] = now; + if (outcome.evicted_id) { + last_capture_times_.erase(*outcome.evicted_id); // keep evicted record retriable } } RCLCPP_WARN_THROTTLE(get_logger(), throttle_clock, 2000, "Capture queue full (drop_oldest): evicted pending '%s' for '%s' " "(pool=%d, queue=%d, total_dropped=%llu)", - outcome.evicted_code ? outcome.evicted_code->c_str() : "?", fault_code.c_str(), + outcome.evicted_id ? outcome.evicted_id->fault_code.c_str() : "?", id.fault_code.c_str(), capture_pool_size_, capture_queue_depth_, static_cast(capture_pool_->dropped_captures())); break; @@ -759,11 +759,11 @@ void FaultManagerNode::capture_on_confirm(const std::string & fault_code) { RCLCPP_WARN_THROTTLE(get_logger(), throttle_clock, 2000, "Capture queue full (reject_newest): dropped capture for '%s' " "(pool=%d, queue=%d, total_dropped=%llu)", - fault_code.c_str(), capture_pool_size_, capture_queue_depth_, + id.fault_code.c_str(), capture_pool_size_, capture_queue_depth_, static_cast(capture_pool_->dropped_captures())); break; case EnqueueResult::kRejectedShuttingDown: - RCLCPP_DEBUG(get_logger(), "Capture pool shutting down; skipped capture for '%s'", fault_code.c_str()); + RCLCPP_DEBUG(get_logger(), "Capture pool shutting down; skipped capture for '%s'", id.fault_code.c_str()); break; } } @@ -803,12 +803,17 @@ void FaultManagerNode::handle_report_fault( return; } - // Get status before update (if fault exists) - auto fault_before = storage_->get_fault(request->fault_code); + // The record this report addresses. The source_id owns it, so a report from a + // second source of the same code creates a second record instead of merging. + const FaultId id{request->fault_code, request->source_id}; + + // Get status before update (if the record exists) + auto fault_before = storage_->get_fault(id); std::string status_before = fault_before ? fault_before->status : ""; - // Resolve per-entity debounce config (longest-prefix match on source_id) - // TODO(#276): warn when different entities resolve different configs for the same fault_code + // Resolve per-entity debounce config (longest-prefix match on source_id). The config + // and the record now share an owner, so the resolved band always governs the counter + // it is applied to. auto resolved_config = resolve_config(request->source_id); // Report the fault event (use wall clock time, not sim time, for proper timestamps) @@ -818,15 +823,15 @@ void FaultManagerNode::handle_report_fault( response->accepted = true; - // Get updated fault state to publish event - auto fault_after = storage_->get_fault(request->fault_code); + // Get updated record state to publish event + auto fault_after = storage_->get_fault(id); if (fault_after) { // Process through correlation engine (if enabled) // Only process FAILED events with correlation bool should_mute = false; if (correlation_engine_ && request->event_type == ros2_medkit_msgs::srv::ReportFault::Request::EVENT_FAILED) { auto correlation_result = - correlation_engine_->process_fault(request->fault_code, correlation::severity_to_string(request->severity)); + correlation_engine_->process_fault(id, correlation::severity_to_string(request->severity)); should_mute = correlation_result.should_mute; @@ -896,11 +901,11 @@ void FaultManagerNode::handle_report_fault( audit_transition(kTransitionConfirmed, *fault_after, request->source_id, event_time.nanoseconds()); } if (just_healed) { - audit_transition(kTransitionHealed, *fault_after, "auto_heal", event_time.nanoseconds()); + audit_transition(kTransitionHealed, *fault_after, id.owner, event_time.nanoseconds()); } if (just_confirmed) { - capture_on_confirm(request->fault_code); + capture_on_confirm(id); } // Handle PREFAILED state for lazy_start rosbag capture @@ -908,7 +913,7 @@ void FaultManagerNode::handle_report_fault( (!is_new && status_before != ros2_medkit_msgs::msg::Fault::STATUS_PREFAILED && fault_after->status == ros2_medkit_msgs::msg::Fault::STATUS_PREFAILED); if (just_prefailed && rosbag_capture_) { - rosbag_capture_->on_fault_prefailed(request->fault_code); + rosbag_capture_->on_fault_prefailed(id); } } @@ -1006,43 +1011,58 @@ void FaultManagerNode::handle_clear_fault( return; } - // Process through correlation engine first (to get auto-clear list). + // Which record. Resolved before anything is changed, so an ambiguous unscoped call + // clears nothing at all rather than one arbitrary owner's record. + std::string resolve_error; + const auto target = resolve_target(request->fault_code, request->source_id, resolve_error); + if (!target) { + response->success = false; + response->message = resolve_error; + RCLCPP_WARN(get_logger(), "ClearFault rejected: %s", resolve_error.c_str()); + return; + } + + // Process through correlation engine first (to get auto-clear list). Symptoms come + // back as records of the same owner, so a cascade never crosses an owner boundary. // `skip_correlation_auto_clear` lets the caller opt out of cascade-clearing - // correlated symptom fault codes. Per-entity DELETE routes set it to true + // correlated symptom faults. Per-entity DELETE routes set it to true // so they cannot reach across entity boundaries via the correlation graph. - std::vector auto_cleared_codes; + std::vector auto_cleared; if (correlation_engine_ && !request->skip_correlation_auto_clear) { - auto clear_result = correlation_engine_->process_clear(request->fault_code); - auto_cleared_codes = clear_result.auto_cleared_codes; + auto clear_result = correlation_engine_->process_clear(*target); + auto_cleared = clear_result.auto_cleared_symptoms; } - bool cleared = storage_->clear_fault(request->fault_code); + bool cleared = storage_->clear_fault(*target); response->success = cleared; if (cleared) { - // Evict cooldown tracking for cleared fault and auto-cleared symptoms + // Evict cooldown tracking for the cleared record and the auto-cleared symptoms { std::lock_guard lock(last_capture_mutex_); - last_capture_times_.erase(request->fault_code); - for (const auto & symptom_code : auto_cleared_codes) { - last_capture_times_.erase(symptom_code); + last_capture_times_.erase(*target); + for (const auto & symptom_id : auto_cleared) { + last_capture_times_.erase(symptom_id); } } // Auto-clear correlated symptoms - for (const auto & symptom_code : auto_cleared_codes) { - storage_->clear_fault(symptom_code); + std::vector auto_cleared_codes; + auto_cleared_codes.reserve(auto_cleared.size()); + for (const auto & symptom_id : auto_cleared) { + storage_->clear_fault(symptom_id); + auto_cleared_codes.push_back(symptom_id.fault_code); if (audit_log_) { - auto symptom = storage_->get_fault(symptom_code); + auto symptom = storage_->get_fault(symptom_id); if (symptom) { - audit_transition(kTransitionCleared, *symptom, "clear_service", get_wall_clock_time().nanoseconds()); + audit_transition(kTransitionCleared, *symptom, symptom_id.owner, get_wall_clock_time().nanoseconds()); } } - RCLCPP_DEBUG(get_logger(), "Auto-cleared symptom: %s (root cause: %s)", symptom_code.c_str(), + RCLCPP_DEBUG(get_logger(), "Auto-cleared symptom: %s (root cause: %s)", symptom_id.fault_code.c_str(), request->fault_code.c_str()); - // Also cleanup rosbag for auto-cleared faults + // Also cleanup rosbag for auto-cleared records if (rosbag_capture_) { - rosbag_capture_->on_fault_cleared(symptom_code); + rosbag_capture_->on_fault_cleared(symptom_id); } } @@ -1053,19 +1073,19 @@ void FaultManagerNode::handle_clear_fault( response->message = "Fault cleared: " + request->fault_code + " (auto-cleared " + std::to_string(auto_cleared_codes.size()) + " symptoms)"; } - RCLCPP_INFO(get_logger(), "Fault cleared: %s (auto-cleared %zu symptoms)", request->fault_code.c_str(), - auto_cleared_codes.size()); + RCLCPP_INFO(get_logger(), "Fault cleared: %s (source=%s, auto-cleared %zu symptoms)", request->fault_code.c_str(), + target->owner.c_str(), auto_cleared_codes.size()); - // Cleanup rosbag for the main fault (auto_cleanup handled inside RosbagCapture) + // Cleanup rosbag for the cleared record (auto_cleanup handled inside RosbagCapture) if (rosbag_capture_) { - rosbag_capture_->on_fault_cleared(request->fault_code); + rosbag_capture_->on_fault_cleared(*target); } - // Publish EVENT_CLEARED - get the cleared fault to include in event - auto fault = storage_->get_fault(request->fault_code); + // Publish EVENT_CLEARED - get the cleared record to include in event + auto fault = storage_->get_fault(*target); if (fault) { publish_fault_event(ros2_medkit_msgs::msg::FaultEvent::EVENT_CLEARED, *fault, auto_cleared_codes); - audit_transition(kTransitionCleared, *fault, "clear_service", get_wall_clock_time().nanoseconds()); + audit_transition(kTransitionCleared, *fault, target->owner, get_wall_clock_time().nanoseconds()); } } else { response->message = "Fault not found: " + request->fault_code; @@ -1097,8 +1117,16 @@ void FaultManagerNode::handle_get_fault(const std::shared_ptrget_fault(request->fault_code); + // Which record + std::string resolve_error; + const auto target = resolve_target(request->fault_code, request->source_id, resolve_error); + if (!target) { + response->success = false; + response->error_message = resolve_error; + return; + } + + auto fault = storage_->get_fault(*target); if (!fault) { response->success = false; response->error_message = "Fault not found: " + request->fault_code; @@ -1115,7 +1143,7 @@ void FaultManagerNode::handle_get_fault(const std::shared_ptrenvironment_data.extended_data_records = extended_records; // Get freeze frame snapshots from storage - auto stored_snapshots = storage_->get_snapshots(request->fault_code); + auto stored_snapshots = storage_->get_snapshots(*target); for (const auto & stored_snapshot : stored_snapshots) { ros2_medkit_msgs::msg::Snapshot snapshot; snapshot.type = ros2_medkit_msgs::msg::Snapshot::TYPE_FREEZE_FRAME; @@ -1136,7 +1164,7 @@ void FaultManagerNode::handle_get_fault(const std::shared_ptr