diff --git a/.github/copilot-instructions.md b/.github/copilot-instructions.md index b8b36b43d..697495548 100644 --- a/.github/copilot-instructions.md +++ b/.github/copilot-instructions.md @@ -13,7 +13,7 @@ Multi-package colcon workspace under `src/`: | Package | Purpose | |---------|---------| | `ros2_medkit_gateway` | HTTP gateway - REST server, discovery, entity management, handlers, plugin framework | -| `ros2_medkit_fault_manager` | Fault aggregation with SQLite, AUTOSAR DEM-style debounce, rosbag/snapshot capture | +| `ros2_medkit_fault_manager` | Fault records per (code, reporting source) in SQLite, AUTOSAR DEM-style debounce, rosbag/snapshot capture | | `ros2_medkit_fault_reporter` | Client library for nodes to report faults | | `ros2_medkit_diagnostic_bridge` | Bridges `/diagnostics` topic to fault manager | | `ros2_medkit_serialization` | Runtime JSON <-> ROS 2 message serialization via dynmsg | diff --git a/docs/api/messages.rst b/docs/api/messages.rst index 1b796d66e..a4b6b4860 100644 --- a/docs/api/messages.rst +++ b/docs/api/messages.rst @@ -15,7 +15,12 @@ Messages Fault.msg ~~~~~~~~~ -Core fault data model representing an aggregated fault condition. +Core fault data model representing one fault record. + +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. .. code-block:: text @@ -49,7 +54,8 @@ Core fault data model representing an aggregated fault condition. # Current fault status (PREFAILED, PREPASSED, CONFIRMED, HEALED, CLEARED) 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. string[] reporting_sources **Severity Constants:** @@ -146,7 +152,10 @@ prefix (for example ``/robot1/fault_manager/events``). MutedFaultInfo.msg ~~~~~~~~~~~~~~~~~~ -Information about correlated (muted) symptom faults. +Information about a correlated (muted) symptom record. One entry per muted RECORD: +a root cause mutes the symptoms of its own reporting source only, so two sources +muted on one ``fault_code`` produce two entries carrying that code, told apart by +``source_id``. .. code-block:: text @@ -154,6 +163,7 @@ Information about correlated (muted) symptom faults. string root_cause_code # Root cause that triggered muting string rule_id # Correlation rule ID that matched uint32 delay_ms # Time delay from root cause [ms] + string source_id # Reporting source that owns the muted record ClusterInfo.msg ~~~~~~~~~~~~~~~ @@ -189,7 +199,8 @@ Report a fault event to the FaultManager. uint8 event_type # EVENT_FAILED (0) or EVENT_PASSED (1) uint8 severity # Fault.SEVERITY_* constant (for FAILED events) string description # Human-readable description - string source_id # Fully qualified node name (e.g., "/powertrain/temp_sensor") + string source_id # Fully qualified node name (e.g., "/powertrain/temp_sensor"). + # Owns the record: (fault_code, source_id) is the record identity. **Response:** @@ -215,7 +226,7 @@ Report a fault event to the FaultManager. ClearFault.srv ~~~~~~~~~~~~~~ -Clear/acknowledge a fault. +Clear/acknowledge one fault record. **Request:** @@ -223,15 +234,23 @@ Clear/acknowledge a fault. string fault_code # Fault code to clear bool skip_correlation_auto_clear # Opt out of correlation cascade clear + string source_id # Reporting source that owns the record **Response:** .. code-block:: text - bool success # True if fault was found and cleared + bool success # True if the record was found and cleared string message # Status message or error description string[] auto_cleared_codes # Symptoms auto-cleared with root cause +``source_id`` names the owner of the record to clear. Leaving it empty is unscoped: the +call applies only when exactly one record carries the ``fault_code``, and fails otherwise +with ``success=false`` and a ``message`` beginning ``ambiguous:`` that lists the owners. +An ambiguous call clears nothing. ``GetFault``, ``GetSnapshots`` and ``GetRosbag`` carry +the same field and resolve the same way. On ``GetRosbag`` it scopes the ``fault_code`` +lookup only, and the ``recording_id`` path ignores it. + When ``skip_correlation_auto_clear`` is ``false`` (default), clearing a root-cause fault also clears every symptom that the correlation engine attributes to it via ``auto_clear_with_root`` rules; the cleared @@ -246,11 +265,12 @@ cluster-wide clearing still works. .. note:: - 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``. + ``skip_correlation_auto_clear`` was added in ``ros2_medkit_msgs`` post-0.4.0, + and ``source_id`` after it (the same field was added 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``. ListFaults.srv ~~~~~~~~~~~~~~ diff --git a/docs/config/fault-manager.rst b/docs/config/fault-manager.rst index 400de56c0..4c0bca2fa 100644 --- a/docs/config/fault-manager.rst +++ b/docs/config/fault-manager.rst @@ -1,7 +1,8 @@ Fault Manager Configuration =========================== -The ``ros2_medkit_fault_manager`` node aggregates and manages faults from multiple sources. +The ``ros2_medkit_fault_manager`` node keeps and manages one fault record per +(``fault_code``, reporting source) pair. This page documents all configuration parameters. .. contents:: Table of Contents @@ -114,25 +115,26 @@ A **near miss** is a FAILED report that moved the debounce counter without the f CONFIRMED - the fault nearly happened. PASSED reports move the counter in the healing direction (the fault receding) and are not near misses. -The fault manager appends one entry per near miss to a per-fault-code series, holding the -timestamp, the counter value after the report, the confirmation threshold, the severity, the -reporting source and the fault status the report left behind. +The fault manager appends one entry per near miss to the series of the fault record it moved (one +series per fault code and reporting source), holding the timestamp, the counter value after the +report, the confirmation threshold, the severity, the reporting source and the fault status the +report left behind. The status matters when reading the series. The HEALED latch holds the status all the way from the healing threshold down to the confirmation threshold, so reports on the way back into a fault that does confirm are also near misses by the definition above. Entries recording ``PREFAILED`` are approaches from a resting state; entries recording ``HEALED`` are a counter walking back down under -the latch. The recorded confirmation threshold belongs to the reporting source, while the counter -is shared by all sources of that fault code, so with per-entity thresholds it is not on its own the -distance to confirmation. The series is **retained when the fault is cleared**, because acknowledging one -fault cycle must not erase how often that code approached confirmation across cycles. +the latch. The recorded confirmation threshold and the counter both belong to the record's own +reporting source, so with per-entity thresholds each entry still reads as the record's distance to +confirmation. The series is **retained when the fault is cleared**, because acknowledging one +fault cycle must not erase how often that record approached confirmation across cycles. .. code-block:: yaml fault_manager: ros__parameters: near_miss: - max_per_fault: 200 # Entries kept per fault code (0 = unlimited) + max_per_fault: 200 # Entries kept per fault record (0 = unlimited) .. list-table:: :header-rows: 1 @@ -143,7 +145,7 @@ fault cycle must not erase how often that code approached confirmation across cy - Description * - ``near_miss.max_per_fault`` - ``200`` - - Near-miss entries retained per fault code. When the bound is reached the **oldest** + - Near-miss entries retained per fault record. When the bound is reached the **oldest** entries are evicted, the same direction as ``snapshots.max_per_fault`` and the rosbag cap: a series frozen at boot says nothing about whether the rate of near misses is changing. Set to 0 for unlimited, accepting growth with the reporting rate. @@ -208,8 +210,8 @@ threshold overrides: - The ``source_id`` is the identifier passed in ``ReportFault`` service requests, typically the fully qualified name of the reporting ROS 2 node (e.g., ``/sensors/lidar/front_node``). - You can inspect actual ``source_id`` values in the ``reporting_sources`` field of existing - faults via ``GET /api/v1/faults``. + It also owns the record it creates, so you can inspect actual ``source_id`` values in the + one-element ``reporting_sources`` field of existing faults via ``GET /api/v1/faults``. - The ``source_id`` from ``ReportFault`` requests is matched against configured prefixes. - The **longest matching prefix** wins. For example, ``/sensors/lidar/front`` matches ``/sensors/lidar`` over ``/sensors``. @@ -219,9 +221,11 @@ threshold overrides: .. note:: - When multiple entities report the same ``fault_code``, each event applies the - thresholds resolved from that event's ``source_id``. This means the debounce - behavior follows the reporting entity, not the fault. + When multiple entities report the same ``fault_code``, each of them owns its own + record and each event applies the thresholds resolved from that event's + ``source_id``. The debounce counter belongs to the same record, so an entity's + configured band governs exactly the counter its own reports move: one entity's + reports can neither confirm nor heal another entity's fault. ``auto_confirm_after_sec`` is global-only and cannot be overridden per-entity. Critical faults skip debounce and confirm on their first occurrence; that is @@ -246,8 +250,8 @@ Basic Snapshot Settings max_message_size: 65536 # Max message size in bytes (64KB) default_topics: [] # Topics to capture for all faults config_file: "" # Path to YAML config file - recapture_cooldown_sec: 60.0 # Min seconds between snapshot captures per fault - max_per_fault: 10 # Max snapshots stored per fault code (0 = unlimited) + recapture_cooldown_sec: 60.0 # Min seconds between snapshot captures per fault record + max_per_fault: 10 # Max snapshots stored per fault record (0 = unlimited) capture_pool_size: 2 # Max concurrent capture threads (>= 1) capture_queue_depth: 16 # Max pending captures before policy applies (>= 1) capture_queue_full_policy: reject_newest # reject_newest | drop_oldest @@ -284,11 +288,11 @@ Basic Snapshot Settings - Path to YAML file with fault-specific snapshot configurations. * - ``snapshots.recapture_cooldown_sec`` - ``60.0`` - - Minimum seconds between snapshot captures for the same fault code. + - Minimum seconds between snapshot captures for the same fault record. Prevents snapshot storms when a fault is reported repeatedly. Set to 0 to disable. * - ``snapshots.max_per_fault`` - ``10`` - - Maximum number of snapshot rows stored per fault code. One confirmation + - Maximum number of snapshot rows stored per fault record. One confirmation writes one row per configured topic, and those rows are evicted together: past the limit the OLDEST capture set is dropped whole. A capture larger than the cap is kept anyway rather than torn, since half a freeze frame is @@ -339,7 +343,7 @@ Capture continuous rosbag recordings around fault events. max_buffer_mb: 256 # Ring-buffer RAM cap max_bag_size_mb: 50 # Max size per bag file max_total_storage_mb: 500 # Max total storage - max_bags_per_fault: 1 # Recordings kept per fault code + max_bags_per_fault: 1 # Recordings kept per fault record auto_cleanup: true # Auto-delete old bags .. list-table:: @@ -416,7 +420,7 @@ Capture continuous rosbag recordings around fault events. whole burst's bag at a time (oldest first). * - ``rosbag.max_bags_per_fault`` - ``1`` - - How many recordings one fault code keeps. Past the cap the oldest is + - How many recordings one fault record keeps. Past the cap the oldest is unlinked, so the default reproduces the historical behaviour exactly: a new recording replaces the previous one. ``0`` means unlimited, bounded only by ``max_total_storage_mb``. ``3`` is a reasonable value for a fault diff --git a/docs/roadmap.rst b/docs/roadmap.rst index d17bee1ee..a62b6ed93 100644 --- a/docs/roadmap.rst +++ b/docs/roadmap.rst @@ -123,7 +123,7 @@ and central aggregation. **Key features:** - [x] **Two-level filtering**: FaultReporter (local) + FaultManager (central) -- [x] **Multi-source aggregation**: Same fault code from multiple sources combined into single entry +- [x] **Per-source fault records**: Same fault code from multiple sources kept as one record each, addressed by (code, source) - [x] **Persistent storage**: Fault state survives restarts - [x] **REST API + SSE**: Real-time fault monitoring via HTTP - [x] **Backwards compatibility**: Integration with ``diagnostic_updater`` diff --git a/docs/tutorials/docker.rst b/docs/tutorials/docker.rst index 84c042e67..b07d8363c 100644 --- a/docs/tutorials/docker.rst +++ b/docs/tutorials/docker.rst @@ -39,7 +39,7 @@ Images are available for all supported ROS 2 distributions: Each image includes the gateway and all open-core packages: - ``ros2_medkit_gateway`` - HTTP REST server -- ``ros2_medkit_fault_manager`` - Fault aggregation and management +- ``ros2_medkit_fault_manager`` - Fault record keeping and lifecycle management - ``ros2_medkit_fault_reporter`` - Client library for fault reporting - ``ros2_medkit_diagnostic_bridge`` - Bridges ``/diagnostics`` to fault manager - ``ros2_medkit_serialization`` - Runtime JSON/ROS 2 serialization diff --git a/docs/tutorials/snapshots.rst b/docs/tutorials/snapshots.rst index 1070aa373..85af31467 100644 --- a/docs/tutorials/snapshots.rst +++ b/docs/tutorials/snapshots.rst @@ -118,11 +118,11 @@ Configure snapshot capture via fault manager parameters: - Use background subscriptions (caches latest message) * - ``snapshots.recapture_cooldown_sec`` - ``60.0`` - - Minimum seconds between snapshot captures for the same fault code. + - Minimum seconds between snapshot captures for the same fault record. Prevents snapshot storms when a fault is reported repeatedly. Set to 0 to disable. * - ``snapshots.max_per_fault`` - ``10`` - - Maximum number of snapshot rows stored per fault code. One confirmation + - Maximum number of snapshot rows stored per fault record. One confirmation writes one row per configured topic, and those rows are evicted together: past the limit the OLDEST capture set is dropped whole. A capture larger than the cap is kept anyway rather than torn, since half a freeze frame @@ -565,7 +565,7 @@ Rosbag Configuration Options burst's bag at a time. * - ``snapshots.rosbag.max_bags_per_fault`` - ``1`` - - Recordings kept per fault code; ``0`` means unlimited. Past the cap the + - Recordings kept per fault record (``0`` means unlimited). Past the cap the fault's oldest recording is dropped, and a bag is deleted only once no fault still references it (a burst shares one recording). ``1`` is the historical behaviour - each re-confirmation replaces the previous bag; diff --git a/src/ros2_medkit_action_status_bridge/include/ros2_medkit_action_status_bridge/action_status_bridge_node.hpp b/src/ros2_medkit_action_status_bridge/include/ros2_medkit_action_status_bridge/action_status_bridge_node.hpp index 64585aa45..b7d902a82 100644 --- a/src/ros2_medkit_action_status_bridge/include/ros2_medkit_action_status_bridge/action_status_bridge_node.hpp +++ b/src/ros2_medkit_action_status_bridge/include/ros2_medkit_action_status_bridge/action_status_bridge_node.hpp @@ -132,8 +132,9 @@ class ActionStatusBridgeNode : public rclcpp::Node { /// Get (creating on first use) the FaultReporter for an action. The reporter's /// source_id is fixed when first created: the resolved server FQN if discovery /// has settled, otherwise the action name as a fallback so the fault still fires - /// on time. It is never re-attributed afterwards (reporting_sources is - /// append-only and the per-entity scope filter is strict-AND). + /// on time. It is never re-attributed afterwards: the fault manager keeps one + /// record per (fault code, source), so a swapped source would open a second + /// record and leave the provisional one raised with nothing left to heal it. ros2_medkit_fault_reporter::FaultReporter * reporter_for(const std::string & action_name); /// Resolve the action server's node FQN from its status-topic publisher, for diff --git a/src/ros2_medkit_action_status_bridge/src/action_status_bridge_node.cpp b/src/ros2_medkit_action_status_bridge/src/action_status_bridge_node.cpp index 0329e9ea7..8a9d449b5 100644 --- a/src/ros2_medkit_action_status_bridge/src/action_status_bridge_node.cpp +++ b/src/ros2_medkit_action_status_bridge/src/action_status_bridge_node.cpp @@ -476,11 +476,12 @@ ros2_medkit_fault_reporter::FaultReporter * ActionStatusBridgeNode::reporter_for // so the fault still fires on time. Reporting the fault on time takes priority // over entity attribution when discovery is slow. // - // The source is NOT re-attributed later: reporting_sources is append-only on the - // manager side and the per-entity /faults scope filter is strict-AND, so a - // provisional action-name source cannot be swapped for the FQN afterwards. - // Correct attribution for the slow-discovery case is a separate concern (it - // needs a way to supersede a provisional source). + // The source is NOT re-attributed later: the fault manager keeps one record per + // (fault code, source), so swapping the provisional action-name source for the + // FQN would open a second record under the FQN and leave the provisional one + // raised, with no reporter left to send its PASSED. Correct attribution for the + // slow-discovery case is a separate concern (it needs a way to supersede a + // provisional source). const std::string fqn = server_fqn_for_action(action_name); const std::string & source_id = fqn.empty() ? action_name : fqn; auto reporter = std::make_unique(this->shared_from_this(), source_id); diff --git a/src/ros2_medkit_action_status_bridge/test/test_action_status_bridge.cpp b/src/ros2_medkit_action_status_bridge/test/test_action_status_bridge.cpp index 7bf4aa249..fdbbb589b 100644 --- a/src/ros2_medkit_action_status_bridge/test/test_action_status_bridge.cpp +++ b/src/ros2_medkit_action_status_bridge/test/test_action_status_bridge.cpp @@ -283,7 +283,7 @@ TEST_F(ActionStatusBridgeTest, ReporterFor_StickyCreatedOnceNeverSwapped) { ActionStatusBridgeTestAccess access(node.get()); // No publisher exists, so the server FQN is unresolved and the reporter falls // back to the action name. It must be created once and reused, never swapped: - // reporting_sources is append-only, so a provisional source cannot be undone. + // a swapped source would open a second fault record and strand the first. const void * first = access.reporter_identity("/nav"); EXPECT_NE(first, nullptr); EXPECT_EQ(access.reporter_identity("/nav"), first); diff --git a/src/ros2_medkit_diagnostic_bridge/test/test_integration.test.py b/src/ros2_medkit_diagnostic_bridge/test/test_integration.test.py index 64931ef1e..d2b5e521e 100644 --- a/src/ros2_medkit_diagnostic_bridge/test/test_integration.test.py +++ b/src/ros2_medkit_diagnostic_bridge/test/test_integration.test.py @@ -153,7 +153,7 @@ def publish_diagnostic(self, name, level, message='Test message', hardware_id='t # Give time for message to be processed time.sleep(0.3) - def publish_until(self, name, level, expected_code, *, predicate=None, + def publish_until(self, name, level, expected_code, *, predicate=None, owner=None, message='Test message', hardware_id='test_hw', statuses=None, timeout=25.0): """ Republish a diagnostic until a matching fault satisfies *predicate*. @@ -167,6 +167,10 @@ def publish_until(self, name, level, expected_code, *, predicate=None, Re-emission also drives WARN/PASSED past their occurrence threshold within the filter window deterministically. + A record is the pair (fault code, reporting source), so one code can + name several records. With *owner* set, only the record that source + owns is considered. + Returns the matching fault once *predicate* holds; fails at the deadline. """ @@ -179,15 +183,17 @@ def predicate(_fault): self.publish_diagnostic(name, level, message, hardware_id) fault = next( (f for f in self.list_faults(statuses=statuses) - if f.fault_code == expected_code), + if f.fault_code == expected_code + and (owner is None or list(f.reporting_sources) == [owner])), None, ) if fault is not None: last = fault if predicate(fault): return fault + owned = f' owned by {owner}' if owner is not None else '' self.fail( - f'Fault {expected_code} not satisfied within {timeout}s ' + f'Fault {expected_code}{owned} not satisfied within {timeout}s ' f'(last seen: {last})' ) @@ -253,26 +259,30 @@ def test_05_fault_code_generation(self): diag_name, DiagnosticStatus.ERROR, expected_code, ) - def test_06_hardware_id_is_forwarded_as_source(self): - """Test that hardware_id is forwarded into fault reporting_sources.""" + def test_06_each_forwarded_hardware_id_owns_its_own_record(self): + """Test that each forwarded hardware_id reports the code as its own record.""" self.publish_until( 'shared_sensor', DiagnosticStatus.STALE, 'SHARED_SENSOR', message='No data from source A', hardware_id='/my_lidar_driver', - predicate=lambda f: '/my_lidar_driver' in f.reporting_sources, + owner='/my_lidar_driver', ) - - fault = self.publish_until( + self.publish_until( 'shared_sensor', DiagnosticStatus.STALE, 'SHARED_SENSOR', message='No data from source B', hardware_id='/my_camera_driver', - predicate=lambda f: '/my_lidar_driver' in f.reporting_sources - and '/my_camera_driver' in f.reporting_sources, + owner='/my_camera_driver', ) - self.assertEqual(fault.severity, Fault.SEVERITY_CRITICAL) - self.assertIn('/my_lidar_driver', fault.reporting_sources) - self.assertIn('/my_camera_driver', fault.reporting_sources) + # One diagnostic name forwarded under two hardware ids is two records + # of one code, each owned by the hardware id it came from. + records = [f for f in self.list_faults() if f.fault_code == 'SHARED_SENSOR'] + self.assertEqual( + sorted(list(f.reporting_sources) for f in records), + [['/my_camera_driver'], ['/my_lidar_driver']], + ) + for record in records: + self.assertEqual(record.severity, Fault.SEVERITY_CRITICAL, record.reporting_sources) def test_07_empty_hardware_id_falls_back_to_bridge_source(self): """Test that empty hardware_id falls back to the bridge node FQN.""" diff --git a/src/ros2_medkit_fault_manager/README.md b/src/ros2_medkit_fault_manager/README.md index 977f5a8c4..052bd94ca 100644 --- a/src/ros2_medkit_fault_manager/README.md +++ b/src/ros2_medkit_fault_manager/README.md @@ -4,9 +4,64 @@ Central fault manager node for the ros2_medkit fault management system. ## Overview -The FaultManager node provides a central point for fault aggregation and lifecycle management. -It receives fault reports from multiple sources, aggregates them by `fault_code`, and provides -query and clearing interfaces. +The FaultManager node provides a central point for fault record keeping and lifecycle management. +It receives fault reports from multiple sources, keeps one record per `(fault_code, reporting +source)` pair, and provides query and clearing interfaces. + +A record is identified by its `fault_code` and by the `source_id` the `ReportFault` call carried. +That source is the record's **owner**. Two sources reporting one `fault_code` are two records, each +with its own status, debounce counter, occurrence count, severity, timestamps, freeze frame, +snapshots, near-miss series and rosbag links. A clear or an auto-heal driven by one owner never +touches another owner's record. The services that act on a single record (`~/clear_fault`, +`~/get_fault`, `~/get_snapshots`, `~/get_rosbag`) carry a `source_id` request field naming the +owner. Leaving it empty is unscoped and applies only when exactly one record carries the +`fault_code`, otherwise the call fails with a message beginning `ambiguous:` that lists the +owners and changes nothing. + +### Upgrading a database written before the owner column + +A SQLite store written by an earlier release is keyed by the bare `fault_code` and is rebuilt on +the first open. The owner of each migrated record is the first entry of that row's legacy +`reporting_sources` value. A row that listed several sources folds onto its first one, because the +old schema recorded no per-source counter, status or timestamps to split it by, and inventing them +would put numbers in the store that no report ever produced. + +The rebuild is one way. A fault manager from an earlier release opened on a migrated store aborts +on the first report of a fault code the store does not hold yet +(`NOT NULL constraint failed: faults.owner`), and on a fault code two sources share it writes one +source's data into the other source's row. If a rollback may be needed, keep a copy of `faults.db` +from before the upgrade and restore that copy instead of pointing the earlier release at the +migrated file. + +The recovery reads that column as text rather than requiring it to be valid JSON, because earlier +builds escaped only the quote, the backslash and `\b \f \n \r \t`, so a `source_id` carrying any +other control byte was written as text no JSON parser accepts. The mapping is exact: + +| stored `reporting_sources` | migrated owner | +|---|---| +| empty | `legacy` | +| `sensor_a` (a bare word, not an array) | `sensor_a` | +| `[]` | `legacy` | +| `["a<0x01>b"]` | `a<0x01>b` (the control byte is kept) | +| anything that is neither a JSON array nor a bare word (an object wrapper, a value behind a byte-order mark) | `legacy` | + +A migrated record never carries an empty owner. Where no source can be read the record gets the +synthetic owner `legacy` and one warning names it, so it stays addressable through `source_id` like +any other record: `ClearFault` and `GetFault` take `legacy` and act on it alone, and it appears +under that name in an `ambiguous:` refusal. New rows are always written as valid JSON, control +bytes escaped as `\u00XX`. + +The empty owner therefore means exactly one thing: an evidence row (freeze frame, snapshot, rosbag +link) not yet assigned to a record. Such a row takes the owner of its fault code when that code has +exactly one owner. When it has several, or none, the row keeps the empty owner and is named in a +warning rather than handed to an owner the database cannot prove it belongs to. That backfill runs +on any open that finds a row it can actually assign, not only on the open that adds the column, so +a store left half-migrated by an interrupted run is healed the next time it is opened, while a row +nothing can resolve is left alone instead of being retried forever. + +The table names `faults_new` and `freeze_frames_new` are reserved for the rebuild. Anything found +under those names is debris from a rebuild that did not finish and is dropped on open, with a +warning naming the table. ## Quick Start @@ -46,15 +101,16 @@ ros2 service call /fault_manager/clear_fault ros2_medkit_msgs/srv/ClearFault \ ## Features -- **Multi-source aggregation**: Same `fault_code` from different sources creates a single fault +- **Per-source records**: Same `fault_code` from different sources creates one record per source, + each filtered, cleared and healed on its own - **Occurrence tracking**: Counts outages, not reports - the count starts at one and rises only - when a cleared fault is raised again - and tracks all reporting sources + when a cleared record is raised again by its own owner - **Severity escalation**: Fault severity is updated if a higher severity is reported - **Persistent storage**: SQLite backend ensures faults survive node restarts - **Debounce filtering** (optional): AUTOSAR DEM-style counter-based fault confirmation with per-entity threshold overrides - **Snapshot capture**: Captures topic data when faults are confirmed for debugging (the value snapshots are deleted when the fault is cleared, unless `snapshots.retain_on_clear` is set) -- **Near-miss series**: Appends one entry per FAILED report that moved the debounce counter without confirming, bounded per fault code and retained when the fault is cleared -- **Freeze-frame retention**: One compact JSON freeze-frame per fault code, retained across `clear_fault` (see below) +- **Near-miss series**: Appends one entry per FAILED report that moved the debounce counter without confirming, bounded per record and retained when the record is cleared +- **Freeze-frame retention**: One compact JSON freeze-frame per record, retained across `clear_fault` (see below) - **Fault correlation** (optional): Root cause analysis with symptom muting and auto-clear - **Tamper-evident audit log** (optional): Append-only, hash-chained record of fault state transitions for verifiable history @@ -69,13 +125,13 @@ ros2 service call /fault_manager/clear_fault ros2_medkit_msgs/srv/ClearFault \ | `healing_threshold` | int | `3` | Counter value at which faults are healed | | `auto_confirm_after_sec` | double | `0.0` | Auto-confirm PREFAILED faults after timeout (0 = disabled) | | `entity_thresholds.config_file` | string | `""` | Path to YAML file with per-entity debounce threshold overrides | -| `near_miss.max_per_fault` | int | `200` | Near-miss entries retained per fault code, oldest evicted first (0 = unlimited) | +| `near_miss.max_per_fault` | int | `200` | Near-miss entries retained per fault record, oldest evicted first (0 = unlimited) | ### Snapshot Parameters Snapshots capture topic data when faults are confirmed for post-mortem debugging. -Each confirm also writes a **freeze-frame**: a single compact JSON object mapping every captured topic to its value at confirmation time, keyed by fault code. It differs from per-topic snapshots in two ways: snapshots are deleted when the fault is cleared, while the freeze-frame is retained across `clear_fault` (once the snapshots are gone, `~/get_fault` serves the retained frame so the confirmed-state record stays available after acknowledgement); and a re-confirm that captures nothing (e.g. source publishers down) never overwrites an existing non-empty frame. A fault code with no configured capture set gets no freeze-frame row; a configured capture that samples nothing on its first run records an empty `{}` frame. Freeze-frame storage is bounded by the number of distinct fault codes (one row per code, replaced in place) and rows are never evicted. +Each confirm also writes a **freeze-frame**: a single compact JSON object mapping every captured topic to its value at confirmation time, keyed by the record. Which topics are captured is decided by the fault CODE (`fault_specific` and `patterns` are configuration about what a code means), while what is written belongs to the record, so two owners confirming one code capture the same topics into two separate frames. It differs from per-topic snapshots in two ways: snapshots are deleted when the record is cleared, while the freeze-frame is retained across `clear_fault` (once the snapshots are gone, `~/get_fault` serves the retained frame so the confirmed-state record stays available after acknowledgement), and a re-confirm that captures nothing (e.g. source publishers down) never overwrites an existing non-empty frame. A fault code with no configured capture set gets no freeze-frame row. A configured capture that samples nothing on its first run records an empty `{}` frame. Freeze-frame storage is bounded by the number of distinct records (one row per record, replaced in place) and rows are never evicted. Under a fault storm, captures are bounded by a worker pool (`capture_pool_size`) draining a bounded queue (`capture_queue_depth`); excess captures are dropped per `capture_queue_full_policy` and logged (throttled). The pool is shared and is created when snapshots **or** rosbag is enabled, so these parameters bound both. `capture_pool_size` parallelizes freeze-frame snapshot capture only - rosbag stays single-writer regardless of pool size, and correlated faults confirming inside one post-roll window share a single recording. @@ -89,7 +145,7 @@ That single-writer property also shapes what each fault of a burst gets. Nothing | `snapshots.max_message_size` | int | `65536` | Maximum message size in bytes (larger messages skipped) | | `snapshots.default_topics` | string[] | `[]` | Topics to capture for all faults | | `snapshots.config_file` | string | `""` | Path to YAML config for `fault_specific` and `patterns` | -| `snapshots.recapture_cooldown_sec` | double | `60.0` | Min seconds between captures for the same fault code. | +| `snapshots.recapture_cooldown_sec` | double | `60.0` | Min seconds between captures for the same fault record. | | `snapshots.max_per_fault` | int | `10` | Max snapshots retained per fault. | | `snapshots.capture_pool_size` | int | `2` | Max concurrent capture threads under a fault storm (>= 1). Parallelizes snapshot capture only; rosbag stays single-writer. | | `snapshots.capture_queue_depth` | int | `16` | Max pending captures before the full-queue policy applies (>= 1). | @@ -158,8 +214,9 @@ counter walking back down under the latch. Without the field the two cannot be t rows written before the field existed. With per-entity thresholds the recorded `confirmation_threshold` is the one belonging to the -**reporting source**, while the debounce counter is shared by every source of that fault code. It -is therefore not by itself the distance to confirmation for the fault as a whole. +**reporting source**. The counter it describes is that source's own record, so the pair is the +distance to confirmation for the record. The counter is no longer shared between sources of one +fault code: each source's reports move only its own record. Entries are kept and evicted in **arrival order**, not by their timestamps. Reporters carry their own clocks, so a report can arrive carrying a timestamp behind one already stored; ordering the @@ -176,11 +233,11 @@ retained snapshots, because it records the most recent confirmation while the sn to earlier ones. `~/get_snapshots` returns one entry per topic and serves the newest capture of that topic, whichever storage backend is in use. -The bound is **per fault code, not per database**. Fault codes are unbounded in cardinality, so a -reporter emitting a stream of distinct codes still grows the table; the bound caps what any single -code costs, not the total. +The bound is **per record, not per database**. Fault codes are unbounded in cardinality and so is +the set of reporting sources, so a reporter emitting a stream of distinct codes still grows the +table. The bound caps what any single record costs, not the total. -Retention is **bounded per fault code** by `near_miss.max_per_fault` (default 200), evicting the +Retention is **bounded per record** by `near_miss.max_per_fault` (default 200), evicting the **oldest** entries first. That is the same direction as `snapshots.max_per_fault` and the rosbag cap, and for the same reason: a series frozen at boot says nothing about whether the rate is changing, and the evidence a technician wants is the evidence from the fault happening now. Set it @@ -193,7 +250,7 @@ storage API or the database file. ## Advanced: Tamper-Evident Audit Log -An optional append-only, hash-chained audit log records every fault state transition (`occurred`, `confirmed`, `healed`, `cleared`) so the fault history is independently verifiable. Auto-recovery (a fault reaching the healing threshold via PASSED events) is recorded as a distinct `healed` row with source `auto_heal`, so the fault's END is in the timeline and is not confused with a manual `cleared`. The manager has no acknowledge action separate from clearing, so `~/clear_fault` is recorded as `cleared` (clear == ack); there is no `ack` kind. The log also records its own lifecycle with `logging_activated` / `logging_deactivated` markers at start and stop. It is **off by default** because it adds a write and storage cost per transition. +An optional append-only, hash-chained audit log records every fault state transition (`occurred`, `confirmed`, `healed`, `cleared`) so the fault history is independently verifiable. Auto-recovery (a record reaching the healing threshold via PASSED events) is recorded as a distinct `healed` row, so the record's END is in the timeline and is not confused with a manual `cleared`. Every row's `source` is the record's owner, the automatic transitions included: the `transition` column already says what moved the row, and with several owners per fault code the source has to answer whose record moved. The manager has no acknowledge action separate from clearing, so `~/clear_fault` is recorded as `cleared` (clear == ack), and there is no `ack` kind. The log also records its own lifecycle with `logging_activated` / `logging_deactivated` markers at start and stop. It is **off by default** because it adds a write and storage cost per transition. Each transition appends one immutable row holding `record_hash = sha256(prev_hash + canonical(event))` (OpenSSL EVP SHA-256), the `prev_hash` it links to, and a monotonic `seq`. The hash is computed once at insert and never recomputed. A persisted chain head lets the chain resume across restarts. The log is stored in its own SQLite database (separate from the fault store) and is treated as append-only: the manager only ever inserts rows, and `BEFORE UPDATE` / `BEFORE DELETE` triggers reject out-of-band edits (the guarded rotation prune excepted). @@ -441,7 +498,9 @@ ros2 service call /fault_manager/list_faults ros2_medkit_msgs/srv/ListFaults \ Response includes: - `muted_count`: Number of muted symptom faults - `cluster_count`: Number of active fault clusters -- `muted_faults[]`: Details of muted faults (when `include_muted=true`) +- `muted_faults[]`: Details of muted records (when `include_muted=true`), one entry per muted + record. Muting is per record, so a symptom muted under one owner's root cause never hides + another owner's record of the same code, and each entry names its owner in `source_id`. - `clusters[]`: Details of active clusters (when `include_clusters=true`) ### REST API (via Gateway) @@ -462,7 +521,8 @@ Response fields: "fault_code": "MOTOR_COMM_FL", "root_cause_code": "ESTOP_001", "rule_id": "estop_cascade", - "delay_ms": 50 + "delay_ms": 50, + "source_id": "/powertrain/motor_controller" } ], "clusters": [ diff --git a/src/ros2_medkit_fault_manager/config/fault_manager.yaml b/src/ros2_medkit_fault_manager/config/fault_manager.yaml index 9817ec293..51ff6bab1 100644 --- a/src/ros2_medkit_fault_manager/config/fault_manager.yaml +++ b/src/ros2_medkit_fault_manager/config/fault_manager.yaml @@ -12,8 +12,8 @@ fault_manager: # Near-miss series: one appended entry per FAILED report that moved the # debounce counter without confirming. Retained across clear_fault, bounded - # per fault code, oldest evicted first. 0 = unlimited (grows with the - # reporting rate). + # per record - one (fault_code, reporting source) pair - oldest evicted + # first. 0 = unlimited (grows with the reporting rate). # near_miss.max_per_fault: 200 # Healing OFF by default: a recovery signal (e.g. action SUCCEEDED) does not @@ -31,7 +31,7 @@ fault_manager: # to entity-scoped capture and is crash-safe (falls back / self-disables if no # storage backend is usable). Faults confirming during another fault's # post-roll (duration_after_sec) share that recording, one metadata entry per - # fault; at most 32 attach on top of the first, the rest get no entry (WARN). + # record. At most 32 attach on top of the first, the rest get no entry (WARN). snapshots.rosbag.enabled: false # Capture concurrency bound under fault storms (issue #441). These govern the diff --git a/src/ros2_medkit_fault_manager/config/snapshots.yaml b/src/ros2_medkit_fault_manager/config/snapshots.yaml index a4efadfed..3c7f07c51 100644 --- a/src/ros2_medkit_fault_manager/config/snapshots.yaml +++ b/src/ros2_medkit_fault_manager/config/snapshots.yaml @@ -8,9 +8,13 @@ # 2. patterns - regex pattern match on fault code # 3. default_topics - fallback for all other faults # 4. entity-default - zero-config fallback when nothing above matches: -# the reporting source node's own published topics are captured +# the owning reporting source node's own published topics are captured # (parameter snapshots.entity_default, on by default) # +# Resolution is by fault CODE: this file says what a code means, not who reported +# it. What is written belongs to the RECORD, so two sources reporting one code +# capture the same topics into two separate freeze frames. +# # A fault code listed in fault_specific (or matched by a pattern) with an # EMPTY topic list is an explicit per-fault opt-out: nothing is captured and # the entity-default fallback is skipped for that code. @@ -105,7 +109,7 @@ rosbag: # Options: # "entity" - DEFAULT. Subscribe broadly for pre-roll, but on fault confirmation # write only the faulting source node's topics (+ /tf, /tf_static), - # resolved from the fault's reporting source. Useful black box with + # resolved from the record's owning reporting source. Useful black box with # zero per-topic config. Falls back to the full buffer if the source # cannot be resolved to a live node. # "config" - Use same topics as JSON snapshots (from this config file) @@ -154,7 +158,9 @@ rosbag: # Storage path for bag files (default: "" = system temp directory) # Empty string uses /tmp/rosbag_snapshots/ # Bag files are named: fault_{fault_code}_{timestamp}/ and that directory name is - # the recording's public id - the last segment of its bulk-data URL. + # the recording's public id - the last segment of its bulk-data URL. The name carries + # the code only: a burst of records sharing one window shares one recording, and each + # of them gets its own metadata row pointing at it. storage_path: "" # Maximum size per bag file in MB (default: 50) @@ -166,10 +172,10 @@ rosbag: # max_bags_per_fault below only decides how the budget is shared out. max_total_storage_mb: 500 - # Recordings kept per fault code (default: 1, 0 = unlimited) - # A fault that keeps re-confirming leaves a trail of black boxes instead of only - # the latest one. Past the cap the fault's OLDEST recording is dropped, and the bag - # is deleted once no fault still references it (a burst shares one recording). + # Recordings kept per record (default: 1, 0 = unlimited) + # A record that keeps re-confirming leaves a trail of black boxes instead of only + # the latest one. Past the cap the record's OLDEST recording is dropped, and the bag + # is deleted once no record still references it (a burst shares one recording). # # 1 is the historical behaviour: each re-confirmation replaces the previous bag. # 3 is a good starting point for an intermittent fault you are chasing. diff --git a/src/ros2_medkit_fault_manager/design/index.rst b/src/ros2_medkit_fault_manager/design/index.rst index e37d62acb..932be7acf 100644 --- a/src/ros2_medkit_fault_manager/design/index.rst +++ b/src/ros2_medkit_fault_manager/design/index.rst @@ -135,15 +135,15 @@ Main Components - Future implementations can be added in Issue #8: Fault Persistence Options 3. **InMemoryFaultStorage** - Thread-safe in-memory implementation of FaultStorage - - Uses ``std::map`` keyed by ``fault_code`` for O(log n) lookups + - Uses ``std::map`` keyed by ``FaultId`` for O(log n) lookups, ordered by code first - Protected by ``std::mutex`` for concurrent service request handling - - Aggregates reports from multiple sources into single fault entries - - Implements severity escalation (higher severity overwrites lower) - - Tracks occurrence counts and all reporting sources + - Keeps one record per reporting source, filtered and cleared independently + - Implements severity escalation per record (higher severity overwrites lower) + - Tracks occurrence counts per record -4. **FaultState** - Internal representation of a fault entry +4. **FaultState** - Internal representation of one fault record - Maps directly to ``ros2_medkit_msgs::msg::Fault`` via ``to_msg()`` - - Uses ``std::set`` for reporting_sources to ensure uniqueness + - Carries its owner and emits it as the single entry of reporting_sources - Tracks first and last occurrence timestamps - Manages fault status lifecycle with debounce (PREFAILED -> CONFIRMED -> CLEARED) @@ -164,8 +164,9 @@ Reports a new fault or updates an existing one. - **Input validation**: fault_code and source_id cannot be empty, event_type must be valid - **Event types**: FAILED (fault detected) or PASSED (fault condition cleared) - **Debounce**: FAILED events decrement counter, PASSED events increment counter -- **Aggregation**: Same fault_code from different sources creates a single fault entry -- **Severity escalation**: Fault severity is updated if a higher severity is reported +- **Record identity**: (fault_code, source_id). The same fault_code from different sources + creates one record per source, each debounced, cleared and healed on its own +- **Severity escalation**: A record's severity is updated if its own owner reports a higher one - **Returns**: ``accepted=true`` if event was processed ~/list_faults @@ -196,21 +197,29 @@ All ``FaultStorage`` public methods acquire a mutex lock to ensure thread safety when handling concurrent service requests. This is essential since ROS 2 service callbacks may execute on different threads. -Fault Aggregation -~~~~~~~~~~~~~~~~~ +Fault Records +~~~~~~~~~~~~~ + +A fault record is identified by the pair ``(fault_code, source_id)``, where the +``source_id`` is the one the ``ReportFault`` call carried. That source is the record's +owner. Repeated reports from one source update that source's record. A report from a +source that owns no record for the code opens a new one. This provides: -Multiple reports of the same ``fault_code`` (from same or different sources) are -aggregated into a single fault entry. This provides: +- **Deduplication**: Prevents fault flooding from one source's repeated reports +- **Source attribution**: Each record names the one source that owns it +- **Occurrence counting**: Tracks how many times a record was raised -- **Deduplication**: Prevents fault flooding from repeated reports -- **Source tracking**: Identifies all sources reporting the same fault -- **Occurrence counting**: Tracks how many times a fault was reported +Keeping the records apart is what makes the per-source state mean anything: the debounce +counter, the status, the severity, the occurrence count and the timestamps all belong to +one reporter, so one reporter recovering cannot walk another reporter's fault back from +confirmation, and one reporter's CRITICAL cannot escalate another's record. Severity Escalation ~~~~~~~~~~~~~~~~~~~ -When a fault is re-reported with a higher severity, the stored severity is updated. -This ensures the fault reflects the worst-case condition. Severity levels are ordered: +When a record is re-reported by its owner with a higher severity, the stored severity is +updated. This ensures the record reflects the worst-case condition its own reporter saw. +Severity levels are ordered: ``INFO(0) < WARN(1) < ERROR(2) < CRITICAL(3)``. Status Lifecycle (Debounce Model) 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..517fe2fe5 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,13 +58,20 @@ 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). +/// Carries both halves of the record identity: two owners muted on one code are two +/// entries of that code, told apart by owner. struct MutedFaultData { std::string fault_code; + /// Reporting source that owns the muted record. Filled from the map key, so it + /// cannot drift from the record this entry describes. + std::string owner; std::string root_cause_code; std::string rule_id; uint32_t delay_ms{0}; @@ -95,18 +104,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 +130,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 +147,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 +169,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..c3da1941e 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(); @@ -178,8 +190,8 @@ class FaultManagerNode : public rclcpp::Node { /// Enqueue snapshot + rosbag capture for a fault that has just confirmed. /// 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); + /// @param id The record (fault code and owning source) that reached CONFIRMED + 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..18436c69e 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; - /// Clear a fault by fault_code (manual acknowledgment). Drops the fault'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; - - /// Get total number of stored faults + /// 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 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. 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 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), and 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, and 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..a78a29bc6 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; @@ -264,7 +266,7 @@ class RosbagCapture { /// Make @p rows durable for the finished bag at @p bag_path, or discard the bag. /// /// Both finalisation paths end here, so a bag that cannot be looked up never - /// survives on disk: retrieval is keyed by fault code and the quota enumerates + /// survives on disk: retrieval goes through the rows and the quota enumerates /// rows, so a directory with no row is unreachable, uncounted, and can never be /// evicted to make room. /// @@ -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..0eed63cde 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 @@ -88,9 +88,9 @@ struct RosbagConfig { /// (reliable/transient-local where offered) instead of forcing best-effort bool qos_match{true}; - /// Recordings retained per fault code (0 = unlimited, bounded only by - /// max_total_storage_mb). Keep-newest: past the cap the fault's OLDEST recording - /// loses its row, and the bag goes once no fault references it. 1 reproduces the + /// Recordings retained per fault record (0 = unlimited, bounded only by + /// max_total_storage_mb). Keep-newest: past the cap the record's OLDEST recording + /// loses its row, and the bag goes once no record references it. 1 reproduces the /// pre-#620 behaviour, where a re-confirm replaced the previous recording. size_t max_bags_per_fault{1}; @@ -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..ba210203e 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 @@ -25,6 +25,16 @@ namespace ros2_medkit_fault_manager { +/// Owner given to a record migrated out of a database that recorded no readable +/// reporting source. +/// +/// A migrated record is never left with an empty owner. It would be unreachable +/// through the source_id every single-record service takes, and it would leave the +/// child backfill looking for work on every open. The empty owner means one thing +/// only, a child row (freeze frame, snapshot, rosbag link) not yet assigned to a +/// record, and keeping that meaning unambiguous is what this constant is for. +inline constexpr const char * kLegacyOwner = "legacy"; + /// SQLite-based fault storage implementation with persistence /// Thread-safe implementation using mutex protection class SqliteFaultStorage : public FaultStorage { @@ -53,15 +63,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; + + 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; @@ -71,29 +83,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 +131,68 @@ 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 per table by PRAGMA table_info lacking `owner`, which is what makes it safe + /// to re-run on every open. 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. The fault owner + /// comes from the legacy row's first reporting source, empty when none can be read: 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. + /// + /// It never refuses to finish. A row whose content cannot be read leaves an empty owner + /// and a log line, because a migration that aborts on a row takes the whole database with + /// it: the rebuild rolls back and the next open reaches the same row again, forever. + void migrate_faults_add_owner(); + + /// Copy every legacy `faults` row into `faults_new`, deriving the owner in C++. + /// + /// In C++ rather than one INSERT ... SELECT because the derivation must not depend on the + /// column being valid JSON, and SQL's json_extract raises rather than returning NULL on + /// text it cannot parse. Caller holds the migration transaction. + void copy_legacy_fault_rows(); + + /// Give every ownerless child row the owner of its fault, where that is unambiguous. + /// + /// Runs on every open, not only when a column was just added: a child table can be left + /// with empty owners by an open interrupted between the faults rebuild and the backfill, + /// or by a database whose child table gained the column while faults was still legacy. + /// Such a row is invisible to every (code, owner) read, so leaving it is losing evidence. + /// A code carrying several owners is NOT resolved - guessing would hand one owner another + /// owner's evidence - and the row keeps its empty owner with a log line naming it. + void backfill_child_owners(); + + /// Whether any child row is waiting for an owner the database can actually prove, so the + /// backfill above has work to do. Deliberately the backfill's own condition and not just + /// "owner is empty": a row whose code has several owners or none can never be assigned, + /// and treating it as pending would re-enter the write transaction on every open forever. + /// Tables without the column yet are skipped: the migration adds it. + bool has_assignable_child_rows() const; + + /// Whether @p table exists, asked of sqlite_master. Used to tell a leftover scratch table + /// apart from a clean start, which DROP TABLE IF EXISTS on its own cannot report. + bool table_exists(const char * table) const; + + /// Drop one of the migration's scratch tables, warning first when it was really there. + /// `faults_new` and `freeze_frames_new` are reserved for this procedure, so debris under + /// those names is always an unfinished rebuild and never somebody's data. + void drop_scratch_table(const char * table); + + /// 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 +210,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 +219,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..baf32bbdb 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, and 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,8 +244,12 @@ std::vector CorrelationEngine::get_muted_faults() const { std::vector result; result.reserve(muted_faults_.size()); - for (const auto & [code, data] : muted_faults_) { - result.push_back(data); + for (const auto & [muted_id, data] : muted_faults_) { + // Owner taken from the key rather than from the stored copy, so the entry always + // names the record it is filed under. + MutedFaultData entry = data; + entry.owner = muted_id.owner; + result.push_back(entry); } return result; @@ -251,9 +260,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 +299,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 +340,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 +382,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 +408,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 +431,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 +451,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 +464,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 +493,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..b0924ac75 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; @@ -864,8 +869,9 @@ void FaultManagerNode::handle_report_fault( } just_confirmed = true; } else if (!is_new && fault_after->status == ros2_medkit_msgs::msg::Fault::STATUS_CONFIRMED) { - // Fault was already CONFIRMED, data updated (last_occurred, severity, sources). - // Not occurrence_count: a re-report inside one occurrence does not touch it. + // Record was already CONFIRMED, data updated (last_occurred, severity). Not + // reporting_sources: it names the one owner and never grows. Not + // occurrence_count: a re-report inside one occurrence does not touch it. if (!should_mute) { publish_fault_event(ros2_medkit_msgs::msg::FaultEvent::EVENT_UPDATED, *fault_after); } @@ -896,11 +902,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 +914,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); } } @@ -937,10 +943,11 @@ void FaultManagerNode::handle_list_faults( response->muted_count = correlation_engine_->get_muted_count(); response->cluster_count = correlation_engine_->get_cluster_count(); - auto muted_faults = correlation_engine_->get_muted_faults(); - - // Include muted faults details if requested + // Include muted faults details if requested. One entry per muted RECORD, carrying + // both halves of its identity, so two owners muted on one code are two entries of + // that code told apart by source_id. if (request->include_muted) { + const auto muted_faults = correlation_engine_->get_muted_faults(); response->muted_faults.reserve(muted_faults.size()); for (const auto & muted : muted_faults) { ros2_medkit_msgs::msg::MutedFaultInfo info; @@ -948,20 +955,20 @@ void FaultManagerNode::handle_list_faults( info.root_cause_code = muted.root_cause_code; info.rule_id = muted.rule_id; info.delay_ms = muted.delay_ms; + info.source_id = muted.owner; response->muted_faults.push_back(info); } } else { - // Build a set of muted codes for O(1) lookup. - std::unordered_set muted_codes; - for (const auto & muted : muted_faults) { - muted_codes.insert(muted.fault_code); - } - - // Remove all faults whose code is muted, in one pass. + // Hide the muted RECORDS, never every record of a muted code. A root cause mutes + // the symptoms of its own owner, so another owner's record of that same code was + // never suppressed and hiding it would drop a live fault from the default view. auto & faults = response->faults; faults.erase(std::remove_if(faults.begin(), faults.end(), - [&muted_codes](const ros2_medkit_msgs::msg::Fault & fault) { - return muted_codes.count(fault.fault_code) > 0; + [this](const ros2_medkit_msgs::msg::Fault & fault) { + return correlation_engine_->is_muted( + FaultId{fault.fault_code, fault.reporting_sources.empty() + ? std::string() + : fault.reporting_sources.front()}); }), faults.end()); } @@ -1006,43 +1013,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 +1075,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 +1119,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 +1145,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 +1166,7 @@ void FaultManagerNode::handle_get_fault(const std::shared_ptr